Files
inquiry_robot/inquiry-agent/agent/policy/air_group_ops.py
T

1202 lines
43 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
空运协同群业务:激活当前工单、补含电含磁、航线报价/附件、锁舱/释放、成交跟进。
本文件职责:群里副作用(主账 + BOT 回话)。海运应用群仍走 sea_group_ops。
禁止:自动建群/拉人/解散;拿当前工单去猜锁舱工单号。
锁舱按 TMS /air/cabin/lock 拼字段出站;没有 TMS 报价编号也要锁,用工单字段补 quoteId/optionNo。
"""
from __future__ import annotations
import json
import logging
from typing import Any, Optional
from agent.llm.extract_text import extract_air_collab_fields, extract_product_quote
from agent.policy import inquiry_copy as copy
from agent.policy.sea_group_ops import (
_archive_file,
_persist_quote_sheet,
_staff_display_name,
emit_group_deal,
filename_is_excel_pdf,
group_client_of,
halt_closed_group,
hydrate_ticket,
is_handoff_ticket,
is_sales,
is_terminal,
merge_quote_fees,
notify_quote_attachment_saved,
send_filled_sheet,
send_group_text,
ticket_no,
ticket_status,
ticket_view,
)
logger = logging.getLogger(__name__)
TERMINAL = {"已成交", "未成交", "已关闭", "转人工"}
OPEN_ACTIVATE = {"询价中", "已报价", "协商中"}
AIRLINE_ROLE = "air"
_PENDING_QUOTE = "__pending_air_quote"
_PENDING_LOCK = "__pending_air_lock"
_LOCK_ID = "__air_lock_id"
_LOCK_OPTION = "__air_lock_option"
_LOCK_QUOTE = "__air_lock_quote"
_LOCK_EXEC = "__air_lock_exec"
def _ticket_flight_date(ticket: object) -> str:
"""锁舱 flightDate 用工单报价日期,和主账出站一致。"""
view = ticket_view(ticket)
facts = dict(view.get("facts") or {})
return str(
facts.get("报价日期")
or view.get("quoteDate")
or view.get("quote_date")
or ""
).strip()
def is_airline(flow: Any, sender_id: str) -> bool:
"""员工管理岗位=航线·空运且启用。"""
uid = (sender_id or "").strip()
if not uid:
return False
book = getattr(flow, "ledger", None) or getattr(flow, "_ledger", None)
getter = getattr(book, "get_staff_by_wecom_id", None) if book else None
if not callable(getter):
return False
try:
row = getter(wecom_id=uid) or {}
except Exception:
logger.exception("查航线岗失败 wecom=%s", uid)
return False
code = str(row.get("roleCode") or row.get("role_code") or "").strip().lower()
status = str(row.get("status") or "active").strip().lower()
return code == AIRLINE_ROLE and status not in {"disabled", "0", "停用"}
def operator_name(flow: Any, sender_id: str) -> str:
book = getattr(flow, "ledger", None) or getattr(flow, "_ledger", None)
getter = getattr(book, "get_staff_by_wecom_id", None) if book else None
if callable(getter):
try:
row = getter(wecom_id=sender_id) or {}
name = str(row.get("name") or "").strip()
if name:
return name
except Exception:
logger.warning("取操作人姓名失败 wecom=%s", sender_id)
return (sender_id or "").strip() or "航线"
def is_air_ticket(ticket: object) -> bool:
"""按群找回的工单必须是空运。陆运旧单不能当空运当前单补含电含磁。"""
view = ticket_view(ticket)
line = str(view.get("business_line") or view.get("businessLine") or "").strip().upper()
return not line or line == "AIR"
def can_fill_fields(flow: Any, ticket: object, sender_id: str) -> bool:
return is_sales(ticket, sender_id) or is_airline(flow, sender_id)
def activate_ticket(
*,
flow: Any,
chat_id: str,
work_order_no: str,
speak: bool = True,
) -> tuple[Optional[object], str]:
"""
把本群当前工单切到指定空运未结案单。
返回 (工单, phase)。失败时工单为 None。
航线已经用附件或文字落过价时,再发工单号只回到成交跟进,不再发空白模板。
"""
book = flow.ledger
no = (work_order_no or "").strip()
client = group_client_of(flow)
getter = getattr(book, "activate_air_group", None)
if callable(getter):
out = getter(chat_id=chat_id, work_order_no=no) or {}
else:
out = _activate_fallback(book, chat_id=chat_id, work_order_no=no)
err = str(out.get("error") or "")
if err == "ticket_not_found":
if speak:
send_group_text(client, chat_id, copy.air_ticket_not_found(no))
return None, "ticket_missing"
if err == "not_air":
if speak:
send_group_text(client, chat_id, copy.air_ticket_not_air(no))
return None, "not_air"
if err == "terminal":
st = str(out.get("status") or "")
if st == "转人工":
return None, "handoff_ignore"
if speak:
send_group_text(
client, chat_id, copy.air_ticket_closed(no, st)
)
return None, "ticket_closed"
if err == "occupied":
if speak:
send_group_text(client, chat_id, copy.air_ticket_occupied(no))
return None, "occupied"
if not out.get("ok"):
if speak:
send_group_text(client, chat_id, copy.air_ticket_not_found(no))
return None, "ticket_missing"
ticket = hydrate_ticket(flow, _ticket_from_activate(book, no, out))
if is_handoff_ticket(ticket):
return None, "handoff_ignore"
# 先落默认协同,再决定要不要发摘要。发文件切单不说话,工单上仍要有含电/含磁。
_seed_air_activate_collab(flow, ticket)
if speak and _air_already_quoted_for_deal(ticket):
view = ticket_view(ticket)
view["collab_chat_id"] = chat_id
emit_group_deal(flow=flow, ticket=view, quote=dict(view.get("quote") or {}))
_watch_archive_room(chat_id, no)
return ticket, "activated_deal"
if speak:
send_air_brief_and_sheet(flow=flow, ticket=ticket, chat_id=chat_id)
_watch_archive_room(chat_id, no)
from agent.policy.event_exception import WAIT_AIR_QUOTE, start_wait
start_wait(flow.ledger, work_order_no=no, wait_kind=WAIT_AIR_QUOTE)
return ticket, "activated"
def _air_already_quoted_for_deal(ticket: object) -> bool:
"""
航线报价已经落账,再激活应回成交跟进。
附件报价会把来源标成「航线附件报价」,但报价编号仍可能留着 TMS。
只看 source 会误当成还没报价。还只有 TMS 标准价时仍发模板。
"""
status = ticket_status(ticket)
if status == "协商中":
return True
if status != "已报价":
return False
quote = dict(ticket_view(ticket).get("quote") or {})
label = str(quote.get("source_label") or quote.get("sourceLabel") or "")
if "航线" in label or "产品" in label:
return True
return copy.quote_is_sales_adjust(quote)
def _collab_slot_empty(val: object) -> bool:
"""协同项没写,或只写了横杠。横杠是摘要占位,不能当成已经填过。"""
return str(val or "").strip() in {"", "-", "—", "-"}
def _seed_air_activate_collab(flow: Any, ticket: object) -> dict[str, str]:
"""
空运群激活时预填协同,并记到主账。
私聊已经写过的含电/含磁先带上。
货物品名正好是「普货」、该项还空着,默认否。
群里已经记下的值不覆盖。只补空项。
"""
from agent.schema.field_validate import (
AIR_COLLAB_FIELDS,
air_general_cargo_collab_defaults,
pick_collab_from_facts,
)
view = ticket_view(ticket)
no = ticket_no(ticket)
facts = dict(view.get("facts") or {})
# 主账有时品名在货物名称列,事实里没带「品名」。群摘要仍能看到普货。
if not str(facts.get("品名") or facts.get("货物品名") or "").strip():
cargo = str(view.get("goodsName") or view.get("goods_name") or "").strip()
if cargo:
facts["品名"] = cargo
seeded = pick_collab_from_facts(facts, AIR_COLLAB_FIELDS)
for key, val in air_general_cargo_collab_defaults(facts).items():
if _collab_slot_empty(seeded.get(key)):
seeded[key] = val
collab = dict(view.get("collab_facts") or {})
to_write: dict[str, str] = {}
for key, val in seeded.items():
if val and _collab_slot_empty(collab.get(key)):
collab[key] = val
to_write[key] = val
if to_write:
patcher = getattr(getattr(flow, "ledger", None), "patch_collab_facts", None)
if callable(patcher):
patcher(
work_order_no=no,
facts=to_write,
sender_id=str(view.get("sales_wecom_id") or view.get("salesWecomId") or ""),
)
return collab
def send_air_brief_and_sheet(
*,
flow: Any,
ticket: object,
chat_id: str = "",
quote_override: Optional[dict[str, Any]] = None,
tms_quote: Optional[dict[str, Any]] = None,
product_quote: Optional[dict[str, Any]] = None,
) -> str:
"""
空运群再发一版摘要 + 匹配到的空运模板。对齐海运 send_brief_and_sheet。
必须先发出工单信息,模板等摘要发出后再发,禁止后台线程抢先发文件。
填表报价必须带上 TMS 总销售价:群票上的 latest quote 常只有增点费/拆板费明细,
若直接拿去出单,预估费用会落成明细加总(如 200)而不是 salePrice(如 7918.25)。
"""
ticket = hydrate_ticket(flow, ticket)
view = ticket_view(ticket)
no = ticket_no(ticket)
room = str(chat_id or view.get("collab_chat_id") or view.get("collabChatId") or "").strip()
if room and not view.get("collab_chat_id"):
view["collab_chat_id"] = room
facts = dict(view.get("facts") or {})
collab = _seed_air_activate_collab(flow, view)
view["collab_facts"] = collab
view_q = dict(view.get("quote") or {})
override = dict(quote_override or {})
tms = dict(tms_quote) if tms_quote is not None else {}
manual = dict(product_quote or {})
if tms_quote is None:
tms = copy.resolve_tms_quote(None, view)
if not tms and view_q and copy.quote_is_tms(view_q):
tms = view_q
if not manual and copy.quote_is_product(override or view_q):
manual = override or view_q
quote = dict(override or view_q or tms or manual)
# 出单专用:把 TMS salePrice 盖回薄报价,避免预估费用=明细 SUM。
sheet_quote = _sheet_quote_with_tms_sale_total(quote, tms=tms, facts=facts)
sales_name, airline_names, mention_ids = resolve_air_mentions(flow, view)
send_group_text(
group_client_of(flow),
room,
copy.air_group_brief(
work_order_no=no,
facts=facts,
quote=quote,
collab_facts=collab,
has_price=bool(tms.get("total") or tms.get("quoteId")),
tms_quote=tms,
product_quote=manual,
sales_name=sales_name,
airline_names=airline_names,
),
mention_userids=mention_ids,
)
_send_air_activate_sheet(flow=flow, ticket=view, chat_id=room, quote=sheet_quote)
return "brief_and_sheet"
def _air_usable_sale_total(quote: Optional[dict[str, Any]]) -> str:
"""
空运出单可用的总销售价:顶层 total / salePrice,或选定/首条线路 salePrice。
「-」与空串不算。
"""
src = dict(quote or {})
for key in ("total", "total_amount", "salePrice", "totalSalePrice"):
val = str(src.get(key) or "").strip()
if val and val != "-":
return val
selected = src.get("selected_air_option") or src.get("selectedAirOption")
if isinstance(selected, dict):
val = str(selected.get("salePrice") or selected.get("sale_price") or "").strip()
if val and val != "-":
return val
raw_opts = src.get("airOptions") or src.get("air_options") or []
if isinstance(raw_opts, list) and raw_opts and isinstance(raw_opts[0], dict):
val = str(raw_opts[0].get("salePrice") or raw_opts[0].get("sale_price") or "").strip()
if val and val != "-":
return val
return ""
def _money_plain_number(text: str) -> Optional[float]:
"""从「CNY 7918.25」类文本抽出数字;解不开返回 None。"""
raw = str(text or "").strip()
if not raw or raw == "-":
return None
digits = "".join(ch if (ch.isdigit() or ch in ".-") else " " for ch in raw).split()
if not digits:
return None
try:
return float(digits[-1].replace(",", ""))
except ValueError:
return None
def _total_matches_fee_row_sum(quote: dict[str, Any], total_text: str) -> bool:
"""
顶层 total 是否只是明细金额相加(如增点费100+拆板费100=200)。
是则说明缺了 TMS 总销售价,不能当真预估费用。
"""
total_num = _money_plain_number(total_text)
if total_num is None:
return False
rows = quote.get("fee_rows") or quote.get("fee_lines") or quote.get("feeItems") or []
if isinstance(rows, str) and rows.strip():
try:
rows = json.loads(rows)
except json.JSONDecodeError:
rows = []
if not isinstance(rows, list) or not rows:
return False
summed = 0.0
hit = 0
for item in rows:
if not isinstance(item, dict):
continue
amt = _money_plain_number(
str(
item.get("amount")
or item.get("value")
or item.get("salePrice")
or item.get("price")
or ""
)
)
if amt is None:
continue
summed += amt
hit += 1
if hit <= 0:
return False
return abs(total_num - summed) < 0.01
def _sheet_quote_with_tms_sale_total(
quote: dict[str, Any],
*,
tms: Optional[dict[str, Any]] = None,
facts: Optional[dict[str, Any]] = None,
) -> dict[str, Any]:
"""
群填表用的报价包:保留产品/航线改过的明细,但预估费用必须能落到 TMS 总销售价。
协同群 latest quote 常只有 departureCharges 投影出的附加费行,total 空或被加成 200;
销售 1V1 会话里带了 salePrice,所以同一工单两边会不一致。这里出单前对齐。
"""
from agent.schema.air_route_options import project_group_air_quote
out = dict(quote or {})
source = dict(tms or {}) or dict(out)
if not source:
return out
projected = project_group_air_quote(source, dict(facts or {}))
tms_total = _air_usable_sale_total(projected)
cur_total = _air_usable_sale_total(out)
# 补齐线路对象 / 计费重,便于主账从 airOptions.salePrice 再兜底一次。
for key in (
"airOptions",
"air_options",
"selected_air_option",
"selectedAirOption",
"chargeableWeightKg",
"chargeable_weight_kg",
"currency",
"valid_until",
):
cur = out.get(key)
empty = cur in (None, "", "-", [], {})
proj_val = projected.get(key)
if empty and proj_val not in (None, "", "-", [], {}):
out[key] = proj_val
if not tms_total:
return out
need_tms_total = (not cur_total) or cur_total == "-" or _total_matches_fee_row_sum(out, cur_total)
if need_tms_total:
out["total"] = str(projected.get("total") or tms_total).strip()
if projected.get("currency") and not str(out.get("currency") or "").strip():
out["currency"] = projected.get("currency")
return out
def resolve_air_mentions(flow: Any, ticket: object) -> tuple[str, list[str], list[str]]:
"""
空运群摘要必须 @ 到本单销售和航线·空运,对上海运 @ 销售/@ 产品。
"""
view = ticket_view(ticket)
sales_id = str(view.get("sales_wecom_id") or view.get("salesWecomId") or "").strip()
sales_name = str(view.get("sales_name") or view.get("salesName") or "").strip()
if not sales_name or sales_name == "销售":
sales_name = _staff_display_name(flow, sales_id) if sales_id else ""
if not sales_name:
sales_name = sales_id
air_ids: list[str] = []
air_names: list[str] = []
book = getattr(flow, "ledger", None)
lister = getattr(book, "list_staff_by_role", None) if book is not None else None
rows: list[dict[str, Any]] = []
if callable(lister):
try:
rows = list(lister(role_code="air") or [])
except Exception:
logger.exception("列航线岗失败")
for row in rows:
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
name = str(row.get("name") or "").strip() or _staff_display_name(flow, uid)
if not uid:
continue
air_ids.append(uid)
if name and name not in air_names:
air_names.append(name)
mention_ids: list[str] = []
if sales_id:
mention_ids.append(sales_id)
for uid in air_ids:
if uid not in mention_ids:
mention_ids.append(uid)
return sales_name, air_names, mention_ids
def _send_air_activate_sheet(
*,
flow: Any,
ticket: object,
chat_id: str,
quote: dict[str, Any],
) -> None:
"""填匹配到的空运模板并往群里发。失败只补一句,不回滚激活。"""
client = group_client_of(flow)
try:
painted = send_filled_sheet(
flow=flow,
ticket=ticket,
business_line="AIR",
chat_id=chat_id,
quote_override=quote,
) or {}
except Exception:
logger.exception("空运群发模板失败 chat=%s", chat_id)
send_group_text(client, chat_id, copy.group_sheet_send_fail())
return
if painted.get("ok"):
return
err = str(painted.get("error") or "")
if err == "NO_TEMPLATE":
send_group_text(client, chat_id, painted.get("message") or copy.QUOTE_TEMPLATE_MISSING)
return
if err in {"no_renderer_or_chat", "no_file_client"}:
send_group_text(client, chat_id, copy.group_sheet_send_fail())
def _watch_archive_room(chat_id: str, work_order_no: str) -> None:
"""激活成功后听这个群的存档,否则发文件(不能 @)系统听不见。"""
try:
from agent.channel.wshoto.rooms import watch_collab_room
watch_collab_room(chat_id, work_order_no=work_order_no)
except Exception:
logger.exception("空运群登记存档监听失败 chat=%s wo=%s", chat_id, work_order_no)
def resolve_cabin_ticket(
*,
flow: Any,
chat_id: str,
work_order_no: str,
allow_terminal: bool,
) -> tuple[Optional[object], str]:
"""
锁舱/释放按句中工单号取单,不猜当前单。
未结案:激活并切到本群。已在别的群:occupied,不锁不切。
已成交/未成交:不能再激活,但释放仍要读到这张单(成交拦住、未成交放舱)。
"""
book = flow.ledger
no = (work_order_no or "").strip()
client = group_client_of(flow)
getter = getattr(book, "get_ticket", None)
raw = getter(work_order_no=no) if callable(getter) else None
if raw is None:
send_group_text(client, chat_id, copy.air_ticket_not_found(no))
return None, "ticket_missing"
ticket = hydrate_ticket(flow, raw)
view = ticket_view(ticket)
line = str(view.get("business_line") or view.get("businessLine") or "").upper()
if line and line != "AIR":
send_group_text(client, chat_id, copy.air_ticket_not_air(no))
return None, "not_air"
other = str(view.get("collab_chat_id") or view.get("collabChatId") or "").strip()
if other and other != chat_id:
send_group_text(client, chat_id, copy.air_ticket_occupied(no))
return None, "occupied"
if is_terminal(ticket):
if allow_terminal:
return ticket, "terminal_loaded"
if ticket_status(ticket) == "转人工":
return None, "handoff_ignore"
send_group_text(
client, chat_id, copy.air_ticket_closed(no, ticket_status(ticket))
)
return None, "ticket_closed"
return activate_ticket(flow=flow, chat_id=chat_id, work_order_no=no, speak=False)
def apply_air_fields(*, flow: Any, ticket: object, text: str, sender_id: str, chat_id: str) -> str:
"""
补含电/含磁。对齐海运:写入后重发群摘要 + 空运模板,不当一句「已记下」。
"""
ticket = hydrate_ticket(flow, ticket)
no = ticket_no(ticket)
client = group_client_of(flow)
if not is_air_ticket(ticket):
send_group_text(client, chat_id, copy.air_need_work_order())
return "need_work_order"
stopped = halt_closed_group(flow=flow, ticket=ticket, chat_id=chat_id)
if stopped:
return stopped
if not can_fill_fields(flow, ticket, sender_id):
send_group_text(client, chat_id, copy.air_fields_denied())
return "fields_denied"
fields = extract_air_collab_fields(text, allow_b=True)
if not fields:
send_group_text(client, chat_id, copy.air_fields_not_recognized())
return "fields_empty"
flow.ledger.patch_collab_facts(work_order_no=no, facts=fields, sender_id=sender_id)
# 必须回读本单。按群再查会拿到同群更早的陆运/海运旧单(测服 WO202609230007)。
getter = getattr(flow.ledger, "get_ticket", None)
reloaded = getter(work_order_no=no) if callable(getter) else None
ticket2 = hydrate_ticket(flow, reloaded or ticket)
view = ticket_view(ticket2)
collab = dict(view.get("collab_facts") or {})
collab.update(fields)
view["collab_facts"] = collab
view["collab_chat_id"] = chat_id
send_air_brief_and_sheet(flow=flow, ticket=view, chat_id=chat_id)
return "fields_saved"
def apply_air_quote(
*,
flow: Any,
ticket: object,
text: str,
sender_id: str,
chat_id: str,
injected_quote: Optional[dict[str, Any]] = None,
) -> str:
"""
航线文字报价。对齐海运产品报价:同句可带含电含磁;有 TMS 价先问怎么用;
落定后重发摘要 + 空运模板 + 成交跟进。
"""
ticket = hydrate_ticket(flow, ticket)
no = ticket_no(ticket)
client = group_client_of(flow)
stopped = halt_closed_group(flow=flow, ticket=ticket, chat_id=chat_id)
if stopped:
return stopped
if not is_airline(flow, sender_id):
send_group_text(client, chat_id, copy.air_airline_only_quote())
return "not_airline"
fields = extract_air_collab_fields(text)
if fields:
flow.ledger.patch_collab_facts(work_order_no=no, facts=fields, sender_id=sender_id)
quote = injected_quote or extract_product_quote(text) or {}
if not quote:
send_group_text(client, chat_id, copy.air_quote_not_recognized())
return "quote_empty"
view = ticket_view(ticket)
current = dict(view.get("quote") or {})
if _has_tms_quote(current) and not _is_product_quote(current):
_park_air_pending(flow, no, quote, sender_id)
send_group_text(client, chat_id, copy.group_ask_tms_or_product())
return "wait_tms_choice"
from agent.schema.quote_adjust import quote_fee_rows
if quote_fee_rows(current):
quote = merge_quote_fees(current, quote)
return _commit_quote(flow, ticket, quote, chat_id, sender_id)
def apply_air_file(
*,
flow: Any,
ticket: object,
sender_id: str,
filename: str,
media: dict[str, Any],
chat_id: str,
) -> str:
"""
空运群附件:对齐海运。只收 Excel/PDF。航线的当正式报价并出成交跟进;别人的只存档。
发文件不能 @,马上回「附件已保存」。不拆文件金额。已有 TMS 报价编号要留下。
"""
ticket = hydrate_ticket(flow, ticket)
no = ticket_no(ticket)
client = group_client_of(flow)
stopped = halt_closed_group(
flow=flow, ticket=ticket, chat_id=chat_id, file=True
)
if stopped:
return stopped
if not filename_is_excel_pdf(filename):
send_group_text(client, chat_id, copy.group_only_excel_pdf())
return "file_type_rejected"
_archive_file(flow, no, filename, media, sender_id)
notify_quote_attachment_saved(flow, chat_id, no, count=1)
if not is_airline(flow, sender_id):
return "file_archived"
view = ticket_view(ticket)
quote = dict(view.get("quote") or {})
quote["source_label"] = "航线附件报价"
quote["file_name"] = filename
if not quote.get("total"):
quote["total"] = "-"
saved = flow.ledger.upsert_product_quote(work_order_no=no, quote=quote, to_status="已报价") or {}
if not saved.get("ok"):
send_group_text(client, chat_id, "主账写入报价失败。")
return "ledger_fail"
view["collab_chat_id"] = chat_id
view["quote"] = quote
from agent.policy.event_exception import clear_wait
clear_wait(flow.ledger, work_order_no=no)
emit_group_deal(flow=flow, ticket=view, quote=quote)
return "airline_file_quote"
def bound_air_ticket(flow: Any, chat_id: str) -> Optional[object]:
"""这群当前绑的是空运单才算空运协同群。海运应用群不要走这里。"""
hit = current_ticket(flow, chat_id)
if hit is None:
return None
line = str(ticket_view(hit).get("business_line") or "").upper()
return hit if line == "AIR" else None
def apply_air_tms_choice(
*,
flow: Any,
ticket: object,
sender_id: str,
chat_id: str,
keep_tms: bool,
) -> str:
"""
航线选:在 TMS 上改,或不用 TMS。对齐海运:用刚记下的价出摘要 + 空运模板 + 成交跟进。
"""
ticket = hydrate_ticket(flow, ticket)
client = group_client_of(flow)
no = ticket_no(ticket)
stopped = halt_closed_group(flow=flow, ticket=ticket, chat_id=chat_id)
if stopped:
return stopped
if not is_airline(flow, sender_id):
send_group_text(client, chat_id, copy.air_airline_only_quote())
return "not_airline"
pending = _load_air_pending(ticket, no)
if not pending:
send_group_text(client, chat_id, copy.group_need_price_before_choice())
return "quote_choice_no_pending"
view = ticket_view(ticket)
base = dict(view.get("quote") or {})
if keep_tms:
if pending.get("fee_rows") or pending.get("fee_lines"):
merged = merge_quote_fees(base, pending)
else:
merged = dict(base)
merged.update(pending)
if base.get("quoteId") and not pending.get("quoteId"):
merged["quoteId"] = base.get("quoteId")
if base.get("airOptions") and not pending.get("airOptions"):
merged["airOptions"] = base.get("airOptions")
merged["source_label"] = "航线在TMS上改"
else:
merged = dict(pending)
merged["source_label"] = "航线报价"
_clear_air_pending(flow, no, sender_id)
return _commit_quote(flow, ticket, merged, chat_id, sender_id, tms_quote=base if keep_tms else {})
def apply_air_confirm_tms(*, flow: Any, ticket: object, sender_id: str, chat_id: str) -> str:
"""航线确认沿用 TMS:模板入后台,回附件已保存,再出成交跟进。对齐海运。"""
ticket = hydrate_ticket(flow, ticket)
if not is_airline(flow, sender_id):
send_group_text(group_client_of(flow), chat_id, copy.air_airline_only_quote())
return "not_airline"
view = ticket_view(ticket)
quote = dict(view.get("quote") or {})
if not _has_tms_quote(quote):
send_group_text(group_client_of(flow), chat_id, copy.group_confirm_need_tms())
return "need_tms"
quote["source_label"] = "航线确认TMS"
no = ticket_no(ticket)
_clear_air_pending(flow, no, sender_id)
# 确认 TMS 出单:同样用总销售价,避免只落附加费明细时预估费用=200。
sheet_quote = _sheet_quote_with_tms_sale_total(
quote, tms=quote, facts=dict(view.get("facts") or {})
)
saved = flow.ledger.upsert_product_quote(work_order_no=no, quote=sheet_quote, to_status="已报价") or {}
if not saved.get("ok"):
send_group_text(group_client_of(flow), chat_id, "主账写入报价失败。")
return "ledger_fail"
view["collab_chat_id"] = chat_id
view["quote"] = sheet_quote
_persist_quote_sheet(
flow=flow,
ticket=view,
sess=None,
quote=sheet_quote,
sender_id=sender_id,
business_line="AIR",
)
notify_quote_attachment_saved(flow, chat_id, no, count=1)
from agent.policy.event_exception import clear_wait
clear_wait(flow.ledger, work_order_no=no)
emit_group_deal(flow=flow, ticket=view, quote=quote)
return "quote_confirmed"
def apply_air_deal(
*,
flow: Any,
ticket: object,
sender_id: str,
outcome: str,
chat_id: str,
text: str = "",
) -> str:
if not is_sales(ticket, sender_id):
send_group_text(group_client_of(flow), chat_id, copy.air_sales_only_deal())
return "not_sales"
stopped = halt_closed_group(flow=flow, ticket=ticket, chat_id=chat_id)
if stopped:
return stopped
from agent.policy.deal_outcome import finish_deal_outcome, harvest_inline_lost_reason
from agent.policy.sea_group_ops import facts_of_group
view = ticket_view(ticket)
inline = harvest_inline_lost_reason(text) if outcome == "未成交" else ""
result = finish_deal_outcome(
ledger=flow.ledger,
work_order_no=ticket_no(ticket),
outcome=outcome,
facts=facts_of_group(view),
quote=dict(view.get("quote") or {}),
business_line="AIR",
sales_name=str(view.get("sales_name") or ""),
lost_reason=inline,
)
if not result.get("ok"):
send_group_text(group_client_of(flow), chat_id, "主账状态更新失败,成交结果未入账。")
return "ledger_fail"
phase = str(result.get("phase") or "done")
ack_reason = inline if phase == "done" and outcome == "未成交" else ""
send_group_text(
group_client_of(flow),
chat_id,
copy.deal_ack(
outcome,
ticket_no(ticket),
str(result.get("next_followup_at") or ""),
lost_reason=ack_reason,
),
)
return phase
def apply_air_lock(
*,
flow: Any,
ticket: object,
sender_id: str,
chat_id: str,
instruction: str,
option_no: str = "",
airline: str = "",
) -> str:
if not is_airline(flow, sender_id):
send_group_text(group_client_of(flow), chat_id, copy.air_airline_only_cabin())
return "not_airline"
# 锁舱是航线主动再试:即使工单因上次失败标了系统异常,也要再打 TMS,
# 不能只回「TMS接口异常,IT运维排查中」。
no = ticket_no(ticket)
name = operator_name(flow, sender_id)
view_quote = dict(ticket_view(ticket).get("quote") or {})
air_quote = dict((view_quote.get("segment_quotes") or {}).get("AIR") or {})
quote = air_quote or view_quote
options = _air_options(quote)
quote_id = _quote_id(quote)
picked = (option_no or "").strip()
if not picked and len(options) > 1:
facts = dict(ticket_view(ticket).get("collab_facts") or {})
if str(facts.get(_PENDING_LOCK) or "") != "1":
flow.ledger.patch_collab_facts(
work_order_no=no, facts={_PENDING_LOCK: "1"}, sender_id=sender_id
)
send_group_text(group_client_of(flow), chat_id, copy.air_ask_lock_option(options))
return "wait_lock_option"
if not picked and len(options) == 1:
picked = str(options[0].get("optionNo") or options[0].get("option_no") or "").strip()
locker = getattr(flow.ledger, "lock_cabin", None)
if not callable(locker):
send_group_text(
group_client_of(flow),
chat_id,
copy.air_lock_fail(work_order_no=no, operator=name, reason="主账未开放锁舱"),
)
return "lock_unavailable"
from agent.policy.system_exception import (
STEP_TMS_LOCK,
TYPE_TMS,
SystemExceptionEvent,
call_with_retries,
confirm_system_exception,
)
out = call_with_retries(
lambda: locker(
work_order_no=no,
option_no=picked,
sender_id=sender_id,
operator_name=name,
instruction_text=instruction,
airline=(airline or "").strip(),
)
or {},
is_ok=lambda x: bool((x or {}).get("ok")),
)
if not out.get("ok"):
if str(out.get("error") or "") == "no_quote":
send_group_text(group_client_of(flow), chat_id, copy.air_lock_need_quote(no))
return "no_quote"
why = copy.air_cabin_fail_reason(
tms_reason=str(out.get("failReason") or out.get("error") or ""),
flight_date=_ticket_flight_date(ticket),
)
send_group_text(
group_client_of(flow),
chat_id,
copy.air_lock_fail(work_order_no=no, operator=name, reason=why),
)
confirm_system_exception(
event=SystemExceptionEvent(
exception_type=TYPE_TMS,
step=STEP_TMS_LOCK,
reason="空运锁舱接口调用失败",
service="空运锁舱",
error_code=why or "LOCK_FAIL",
work_order_no=no,
ticket_status=ticket_status(ticket),
conversation_kind="group",
conversation_target=chat_id,
extra={
"option_no": picked,
"instruction": instruction,
"sender_id": sender_id,
"operator_name": name,
"flight_date": _ticket_flight_date(ticket),
},
),
ledger=flow.ledger,
)
return "system_exception"
exec_no = str(out.get("lockId") or out.get("tmsRequestId") or "").strip()
quote_no = str(out.get("quoteId") or quote_id).strip()
flow.ledger.patch_collab_facts(
work_order_no=no,
facts={
_LOCK_ID: exec_no,
_LOCK_OPTION: picked,
_LOCK_QUOTE: quote_no,
_LOCK_EXEC: exec_no,
_PENDING_LOCK: "",
},
sender_id=sender_id,
)
send_group_text(
group_client_of(flow),
chat_id,
copy.air_lock_ok(
work_order_no=no,
operator=name,
option_no=picked,
exec_no=exec_no,
quote_no=quote_no,
),
)
from agent.policy.system_exception import clear_system_exception, ticket_is_paused
if ticket_is_paused(ticket):
clear_system_exception(ledger=flow.ledger, work_order_no=no)
view = ticket_view(ticket)
view["collab_chat_id"] = chat_id
emit_group_deal(flow=flow, ticket=view, quote=dict(view.get("quote") or {}))
return "lock_ok"
def _release_should_push_deal(ticket: object) -> bool:
"""
放舱成功后要不要再推成交跟进。
未成交、已关闭、已成交、转人工:成交结果已经定过,只回释放结果。
还在等销售确认时仍要推一张。
"""
return ticket_status(ticket) not in TERMINAL
def apply_air_release(
*,
flow: Any,
ticket: object,
sender_id: str,
chat_id: str,
instruction: str,
) -> str:
if not is_airline(flow, sender_id):
send_group_text(group_client_of(flow), chat_id, copy.air_airline_only_cabin())
return "not_airline"
# 释放同理:暂停单也允许再打 TMS,回话用释放失败原因而不是接口异常兜底。
no = ticket_no(ticket)
if ticket_status(ticket) == "已成交":
send_group_text(group_client_of(flow), chat_id, copy.air_release_deal_blocked(no))
return "release_deal_blocked"
facts = dict(ticket_view(ticket).get("collab_facts") or {})
lock_id = str(facts.get(_LOCK_ID) or "").strip()
name = operator_name(flow, sender_id)
if not lock_id:
send_group_text(group_client_of(flow), chat_id, copy.air_release_need_lock(no))
return "release_need_lock"
releaser = getattr(flow.ledger, "release_cabin", None)
if not callable(releaser):
send_group_text(
group_client_of(flow),
chat_id,
copy.air_release_fail(work_order_no=no, operator=name, reason="主账未开放释放"),
)
return "release_unavailable"
from agent.policy.system_exception import (
STEP_TMS_RELEASE,
TYPE_TMS,
SystemExceptionEvent,
call_with_retries,
confirm_system_exception,
)
out = call_with_retries(
lambda: releaser(
work_order_no=no,
lock_id=lock_id,
sender_id=sender_id,
operator_name=name,
instruction_text=instruction,
)
or {},
is_ok=lambda x: bool((x or {}).get("ok")),
)
if not out.get("ok"):
why = copy.air_cabin_fail_reason(
tms_reason=str(out.get("failReason") or out.get("error") or ""),
flight_date=_ticket_flight_date(ticket),
)
send_group_text(
group_client_of(flow),
chat_id,
copy.air_release_fail(work_order_no=no, operator=name, reason=why),
)
confirm_system_exception(
event=SystemExceptionEvent(
exception_type=TYPE_TMS,
step=STEP_TMS_RELEASE,
reason="空运释放舱位接口调用失败",
service="空运舱位",
error_code=why or "RELEASE_FAIL",
work_order_no=no,
ticket_status=ticket_status(ticket),
conversation_kind="group",
conversation_target=chat_id,
extra={
"lock_id": lock_id,
"instruction": instruction,
"sender_id": sender_id,
"operator_name": name,
"flight_date": _ticket_flight_date(ticket),
},
),
ledger=flow.ledger,
)
return "system_exception"
exec_no = str(out.get("releaseId") or out.get("tmsRequestId") or lock_id).strip()
send_group_text(
group_client_of(flow),
chat_id,
copy.air_release_ok(
work_order_no=no,
operator=name,
option_no=str(facts.get(_LOCK_OPTION) or ""),
exec_no=exec_no,
quote_no=str(facts.get(_LOCK_QUOTE) or ""),
),
)
from agent.policy.system_exception import clear_system_exception, ticket_is_paused
if ticket_is_paused(ticket):
clear_system_exception(ledger=flow.ledger, work_order_no=no)
if _release_should_push_deal(ticket):
view = ticket_view(ticket)
view["collab_chat_id"] = chat_id
emit_group_deal(flow=flow, ticket=view, quote=dict(view.get("quote") or {}))
return "release_ok"
def current_ticket(flow: Any, chat_id: str) -> Optional[object]:
book = flow.ledger
hit = book.find_by_collab_chat(chat_id=chat_id)
if hit is None:
return None
return hydrate_ticket(flow, hit)
def _commit_quote(
flow: Any,
ticket: object,
quote: dict[str, Any],
chat_id: str,
sender_id: str,
tms_quote: Optional[dict[str, Any]] = None,
) -> str:
"""报价落定:写入主账,重发摘要 + 空运模板,再出成交跟进。"""
_ = sender_id
no = ticket_no(ticket)
up = flow.ledger.upsert_product_quote(work_order_no=no, quote=quote, to_status="已报价")
if not up.get("ok"):
send_group_text(group_client_of(flow), chat_id, "主账写入报价失败。")
return "ledger_fail"
from agent.policy.event_exception import clear_wait
clear_wait(flow.ledger, work_order_no=no)
ticket2 = current_ticket(flow, chat_id) or ticket
view = ticket_view(ticket2)
view["collab_chat_id"] = chat_id
view["quote"] = quote
send_air_brief_and_sheet(
flow=flow,
ticket=view,
chat_id=chat_id,
quote_override=quote,
tms_quote=tms_quote if tms_quote is not None else {},
product_quote=quote,
)
emit_group_deal(flow=flow, ticket=view, quote=quote)
return "quoted"
def peek_air_pending(ticket: object) -> dict[str, Any]:
"""这张单是否还在等航线选怎么用 TMS。"""
return _load_air_pending(ticket, ticket_no(ticket))
def _park_air_pending(flow: Any, work_order_no: str, quote: dict[str, Any], sender_id: str) -> None:
"""先记下航线刚报的价,不写主账报价版本。Redis 正式;协同字段兜底给单测。"""
try:
from agent.redis_coord.pending_quote import save_pending_product_quote
save_pending_product_quote(work_order_no, quote)
except Exception:
logger.exception("待确认航线报价写入失败 wo=%s", work_order_no)
flow.ledger.patch_collab_facts(
work_order_no=work_order_no,
facts={_PENDING_QUOTE: json.dumps(quote, ensure_ascii=False)},
sender_id=sender_id,
)
def _load_air_pending(ticket: object, work_order_no: str) -> dict[str, Any]:
try:
from agent.redis_coord.pending_quote import load_pending_product_quote
parked = load_pending_product_quote(work_order_no)
if parked:
return parked
except Exception:
logger.exception("待确认航线报价读取失败 wo=%s", work_order_no)
raw = str((ticket_view(ticket).get("collab_facts") or {}).get(_PENDING_QUOTE) or "").strip()
if not raw:
return {}
try:
data = json.loads(raw)
except json.JSONDecodeError:
return {}
return dict(data) if isinstance(data, dict) else {}
def _clear_air_pending(flow: Any, work_order_no: str, sender_id: str) -> None:
try:
from agent.redis_coord.pending_quote import clear_pending_product_quote
clear_pending_product_quote(work_order_no)
except Exception:
logger.exception("待确认航线报价清理失败 wo=%s", work_order_no)
flow.ledger.patch_collab_facts(
work_order_no=work_order_no,
facts={_PENDING_QUOTE: ""},
sender_id=sender_id,
)
def _has_tms_quote(quote: dict[str, Any]) -> bool:
"""当前价来自 TMS 才要问;航线/产品/销售调价不再替用户选。"""
src = dict(quote or {})
if not (src.get("total") or src.get("quoteId") or src.get("quote_id")):
return False
return copy.quote_is_tms(src)
def _is_product_quote(quote: dict[str, Any]) -> bool:
src = str(quote.get("source_label") or "")
return "航线" in src or "产品" in src
def _quote_id(quote: dict[str, Any]) -> str:
return str(quote.get("quoteId") or quote.get("quote_id") or "").strip()
def _air_options(quote: dict[str, Any]) -> list[dict[str, Any]]:
raw = quote.get("airOptions") or quote.get("air_options")
if isinstance(raw, str) and raw.strip().startswith("["):
try:
raw = json.loads(raw)
except json.JSONDecodeError:
raw = []
if isinstance(raw, list):
return [x for x in raw if isinstance(x, dict)]
blob = str(quote.get("air_options_json") or "").strip()
if blob.startswith("["):
try:
parsed = json.loads(blob)
except json.JSONDecodeError:
return []
return [x for x in parsed if isinstance(x, dict)]
return []
def _activate_fallback(book: Any, *, chat_id: str, work_order_no: str) -> dict[str, Any]:
"""内存账本未实现 activate 时的兜底,单测应走正式方法。"""
getter = getattr(book, "get_ticket", None)
ticket = getter(work_order_no=work_order_no) if callable(getter) else None
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
return {"ok": True, "work_order_no": work_order_no}
def _ticket_from_activate(book: Any, no: str, out: dict[str, Any]) -> object:
getter = getattr(book, "get_ticket", None)
if callable(getter):
hit = getter(work_order_no=no)
if hit is not None:
return hit
return out