""" 空运文字询价确定性流程(模型只负责抽字段,本文件只编排与执行)。 本文件职责:补问到必填齐 → 核对确定后建单 → 查 TMS(一条出卡,多条 H5 选线路)。 对不上机场也不能跳过查价、本地说暂无报价。 调用:handler 在 Worker/inbox 线程;禁止回调线程同步跑。 禁止:直连 TMS;用图状态冒充主账成功;用户文字冒充文件已生成。 """ from __future__ import annotations import logging import threading import time import uuid from dataclasses import asdict, dataclass, field from typing import Any, Callable, Optional from agent.ledger.http_ledger import HttpLedger from agent.ledger.memory_ledger import MemoryLedger from agent.llm.extract_text import extract_inquiry_snapshot from agent.policy import inquiry_copy as copy from agent.schema.field_validate import validate_required_fields from agent.schema.land_options import looks_like_confirm from agent.schema.tms_air_query import assemble_air_query logger = logging.getLogger(__name__) # 可返回 outbox id;查价前要等询价确认真正发出。 ReplyFn = Callable[..., Optional[str]] @dataclass class FlowSession: """一单会话书签(运行账,不是六态真相)。""" sender_id: str thread_id: str wait_version: int = 0 phase: str = "idle" business_line: str = "" facts: dict[str, str] = field(default_factory=dict) work_order_no: str = "" first_or_same: str = "first" history_work_order_no: str = "" quote: dict[str, Any] = field(default_factory=dict) tms_quote: dict[str, Any] = field(default_factory=dict) # 陆运 H5 重选要用的整包线路。选定后 tms_quote 只留已选那条,供群摘要使用。 land_route_quote: dict[str, Any] = field(default_factory=dict) status: str = "" allowed: tuple[str, ...] = () immutable_text: str = "" collab_chat_id: str = "" collab_facts: dict[str, str] = field(default_factory=dict) invited_userids: tuple[str, ...] = () invited_sales_name: str = "" invited_product_names: tuple[str, ...] = () deal_version: int = 0 collab_ended: bool = False pending_product_quote: dict[str, Any] = field(default_factory=dict) clarify_missing_sig: str = "" clarify_retry: int = 0 # clarify=3 次反问入口;intent=销售要求人来接管。空=没有待点的转人工卡。 handoff_offer: str = "" handoff_committed: bool = False virtual_work_order_no: str = "" def _session_payload(sess: FlowSession) -> dict[str, Any]: """书签转 JSON。tuple 改 list,方便 HTTP/Worker 共用。""" data = asdict(sess) data["allowed"] = list(sess.allowed or ()) data["invited_userids"] = list(sess.invited_userids or ()) data["invited_product_names"] = list(sess.invited_product_names or ()) data["facts"] = {str(k): str(v) for k, v in (sess.facts or {}).items()} return data def _session_from_payload(raw: dict[str, Any]) -> FlowSession: """Redis JSON → 书签。缺字段用默认,避免旧键打崩。""" return FlowSession( sender_id=str(raw.get("sender_id") or ""), thread_id=str(raw.get("thread_id") or raw.get("sender_id") or ""), wait_version=int(raw.get("wait_version") or 0), phase=str(raw.get("phase") or "idle"), business_line=str(raw.get("business_line") or ""), facts={str(k): str(v) for k, v in dict(raw.get("facts") or {}).items()}, work_order_no=str(raw.get("work_order_no") or ""), first_or_same=str(raw.get("first_or_same") or "first"), history_work_order_no=str(raw.get("history_work_order_no") or ""), quote=dict(raw.get("quote") or {}), tms_quote=dict(raw.get("tms_quote") or raw.get("tmsQuote") or {}), land_route_quote=dict(raw.get("land_route_quote") or raw.get("landRouteQuote") or {}), status=str(raw.get("status") or ""), allowed=tuple(str(x) for x in (raw.get("allowed") or ())), immutable_text=str(raw.get("immutable_text") or ""), collab_chat_id=str(raw.get("collab_chat_id") or ""), collab_facts={str(k): str(v) for k, v in dict(raw.get("collab_facts") or {}).items()}, invited_userids=tuple(str(x) for x in (raw.get("invited_userids") or ())), invited_sales_name=str(raw.get("invited_sales_name") or ""), invited_product_names=tuple(str(x) for x in (raw.get("invited_product_names") or ())), deal_version=int(raw.get("deal_version") or 0), collab_ended=bool(raw.get("collab_ended")), pending_product_quote=dict(raw.get("pending_product_quote") or {}), clarify_missing_sig=str(raw.get("clarify_missing_sig") or ""), clarify_retry=int(raw.get("clarify_retry") or 0), handoff_offer=str(raw.get("handoff_offer") or ""), handoff_committed=bool(raw.get("handoff_committed")), virtual_work_order_no=str(raw.get("virtual_work_order_no") or ""), ) class AirTextInquiryFlow: """ 空运文字询价 Owner。海运走 SeaTextInquiryFlow,陆运走 LandTextInquiryFlow。 同一 sender 新开询价会覆盖「当前书签」(旧单留在主账)。 价格卡按工单号另存一份当时会话:点哪张卡就生成哪张卡当时的报价单。 """ bookmark_kind = "air" # 明细文本真正发出后再发卡,避免销售先看到半截引用条。 _detail_before_card_gap_sec = 0.8 def __init__(self, ledger: Any = None, group_client: Any = None, *, bookmark_store: Any = None) -> None: self._ledger = ledger or build_ledger() self._lock = threading.Lock() self._by_sender: dict[str, FlowSession] = {} self._by_ticket: dict[str, FlowSession] = {} self._group_client = group_client self._bookmark_store = bookmark_store @property def ledger(self) -> Any: return self._ledger @ledger.setter def ledger(self, value: Any) -> None: self._ledger = value def group_client(self) -> Any: """空运群 BOT 出站;单测注入假客户端。""" return self._group_client def session_of(self, sender_id: str) -> Optional[FlowSession]: sid = (sender_id or "").strip() with self._lock: hit = self._by_sender.get(sid) if hit is not None: return hit loaded = self._load_bookmark(sid) if loaded is None: return None self._save(loaded) return loaded def session_by_work_order(self, work_order_no: str) -> Optional[FlowSession]: """用工单号找回书签,不校验 sender。群成交卡校验用。""" no = (work_order_no or "").strip() if not no: return None with self._lock: return self._by_ticket.get(no) def session_by_ticket(self, sender_id: str, work_order_no: str) -> Optional[FlowSession]: """ 只查进程内该书签,不读主账。 续办选空运/海运流程用:海运已跳过协同时必须找回 wait_adopt, 不能因为主账把已报价海运一律收成 wait_collab 就套空运协同卡。 """ no = (work_order_no or "").strip().upper() if not (sender_id or "").strip() or not no: return None with self._lock: hit = self._by_ticket.get(no) if hit and hit.sender_id == sender_id: return hit for key, sess in self._by_ticket.items(): if (key or "").upper() == no and sess.sender_id == sender_id: return sess return None def session_for_action( self, sender_id: str, card_meta: Optional[dict[str, Any]] = None, *, text: str = "", ) -> Optional[FlowSession]: """ 点卡:用工单号找回出卡当时的会话;对不上才用该销售当前书签。 必须校验 sender_id,禁止用别人的工单号串会话。 """ no = copy.work_order_from_card_meta(card_meta, text=text) if not no: return self.session_of(sender_id) with self._lock: hit = self._by_ticket.get(no) if hit and hit.sender_id == sender_id: return hit # 卡上有工单号就只动这一单;进程刚重启时从主账收回,禁止落到别人/最新单。 return self._load_ticket_session(sender_id, no) def _save(self, sess: FlowSession) -> None: with self._lock: self._by_sender[sess.sender_id] = sess no = (sess.work_order_no or "").strip() if no: self._by_ticket[no] = sess vt = (sess.virtual_work_order_no or "").strip() if vt: self._by_ticket[vt] = sess self._persist_bookmark(sess) def _handoff_state(self, sess: Optional[FlowSession]) -> dict[str, Any]: """新建书签时带上转人工计数,避免补问一轮就把 3 次清掉。已冻结的票不当新书签底。""" from agent.policy.handoff_human import is_handoff_frozen if sess is None or is_handoff_frozen(sess): return {} return { "clarify_missing_sig": sess.clarify_missing_sig, "clarify_retry": sess.clarify_retry, "handoff_offer": sess.handoff_offer, "handoff_committed": sess.handoff_committed, "virtual_work_order_no": sess.virtual_work_order_no, "collab_chat_id": sess.collab_chat_id, "work_order_no": sess.work_order_no, "status": sess.status, "quote": dict(sess.quote or {}), "tms_quote": dict(sess.tms_quote or {}), "land_route_quote": dict(sess.land_route_quote or {}), "allowed": sess.allowed, "invited_userids": sess.invited_userids, "invited_sales_name": sess.invited_sales_name, "invited_product_names": sess.invited_product_names, } def _block_if_handoff_committed(self, sess: Optional[FlowSession], reply: ReplyFn) -> bool: """真正转人工之后这一票不再跟。返回 True 表示已经拦住。""" from agent.policy.handoff_human import is_handoff_frozen if sess is None: return False if is_handoff_frozen(sess): no = (sess.work_order_no or sess.virtual_work_order_no or "").strip() reply(f"工单 {no} 已转人工,本单不再自动跟进。如需新询价请直接说明新的运输需求。") return True return False def _detach_frozen_bookmark( self, sess: Optional[FlowSession], sender_id: str ) -> Optional[FlowSession]: """ 真正转人工后丢掉当前私聊书签,好开新询价。 旧单仍留在工单号索引,群里带号仍能对上。 """ from agent.policy.handoff_human import is_handoff_frozen if sess is None or not is_handoff_frozen(sess): return sess logger.info( "handoff.detach_frozen sender=%s wo=%s vt=%s", sender_id, sess.work_order_no or "-", sess.virtual_work_order_no or "-", ) self.drop_current_bookmark(sender_id) return None def offer_handoff(self, *, sender_id: str, text: str, reply: ReplyFn) -> str: """ 路由判定销售明确要求人来接管。运输方式不明则继续补问运输方式。 """ from agent.llm.extract_text import detect_transport_mode from agent.policy.handoff_human import is_closed_for_handoff sess = self.session_of(sender_id) if sess and is_closed_for_handoff(sess.status): reply("该工单已结束,不能再转人工。") return "handoff_closed" if sess and self._block_if_handoff_committed(sess, reply): return "handoff_done" mode = (sess.business_line if sess else "") or detect_transport_mode(text) or "" mode = (mode or "").upper() if not mode: reply(copy.ASK_TRANSPORT) self._save( FlowSession( sender_id=sender_id, thread_id=sender_id, phase="need_mode", facts=dict(sess.facts) if sess else {}, immutable_text=(sess.immutable_text if sess and sess.immutable_text else text), **self._handoff_state(sess), ) ) return "need_mode" if mode == "SEA": from agent.policy.sea_text_flow import get_sea_text_flow return get_sea_text_flow()._offer_handoff_now( sender_id=sender_id, reason="intent", reply=reply, text=text ) if mode == "LAND": from agent.policy.land_text_flow import get_land_text_flow return get_land_text_flow()._offer_handoff_now( sender_id=sender_id, reason="intent", reply=reply, text=text ) return self._offer_handoff_now( sender_id=sender_id, reason="intent", reply=reply, text=text ) def _offer_handoff_now( self, *, sender_id: str, reason: str, reply: ReplyFn, text: str = "", sess: Optional[FlowSession] = None, ) -> str: """ 发出转人工入口。reason=clarify|intent。 点卡/空运群发号之前不冻结。 """ from agent.policy.handoff_human import ( has_collab_group, ticket_no_for_handoff_offer, transport_label, ) sess = sess or self.session_of(sender_id) if sess is None: sess = FlowSession( sender_id=sender_id, thread_id=sender_id, phase="handoff_offer", business_line=self.bookmark_kind.upper() if self.bookmark_kind != "sea" else "SEA", immutable_text=text, ) if self.bookmark_kind == "air": sess.business_line = "AIR" elif self.bookmark_kind == "land": sess.business_line = "LAND" else: sess.business_line = "SEA" line = (sess.business_line or "").upper() if has_collab_group(sess): reply(copy.HANDOFF_ALREADY_COLLAB) return "handoff_already" why = copy.HANDOFF_CLARIFY_REASON if reason == "clarify" else copy.HANDOFF_INTENT_REASON if line == "AIR": no = ticket_no_for_handoff_offer(sess) reply(copy.handoff_air_hint(work_order_no=no)) sess.handoff_offer = reason sess.phase = "handoff_offer" self._save(sess) return "handoff_offer" token = (sess.work_order_no or "").strip() or sender_id reply( copy.handoff_offer_text(reason=why, transport_mode=transport_label(line)), copy.handoff_offer_wecom_payload(reason=why, task_token=token), ) sess.handoff_offer = reason sess.phase = "handoff_offer" sess.allowed = (copy.BTN_PULL_COLLAB,) self._save(sess) return "handoff_offer" def _clarify_or_handoff( self, *, sess: Optional[FlowSession], keys: list[str], new_sess: FlowSession, reply: ReplyFn, clarify_text: str, ) -> str: """缺项:第一次不计数;满 3 次先发(3/3)再出转人工入口,否则照常补问。""" from agent.policy.handoff_human import inherit_pre_create_ticket, note_clarify_round if sess is not None: new_sess.clarify_missing_sig = sess.clarify_missing_sig new_sess.clarify_retry = sess.clarify_retry new_sess.handoff_offer = sess.handoff_offer new_sess.handoff_committed = sess.handoff_committed new_sess.virtual_work_order_no = sess.virtual_work_order_no inherit_pre_create_ticket(sess, new_sess) decision = note_clarify_round(new_sess, keys) reply(copy.with_clarify_progress(clarify_text, new_sess.clarify_retry)) if decision == "handoff": return self._offer_handoff_now( sender_id=new_sess.sender_id, reason="clarify", reply=reply, sess=new_sess, ) self._save(new_sess) return new_sess.phase or "clarify" def _detach_ticketed_for_new_inquiry(self, sess: Optional[FlowSession]) -> Optional[FlowSession]: """ 已经出过正式工单号:这一句按新一轮收字段,不带上一单。 测服 WO202609200012 出号后,下一句海运把旧贸易条款写进 WO202609200013。 点按钮、未成交原因、调价已在前面处理。虚拟号还在补问,要接着填。 """ if sess is None: return None from agent.policy.handoff_human import is_virtual_no no = (sess.work_order_no or "").strip() phase = (sess.phase or "").strip() if not no or is_virtual_no(no): return sess if phase in {"clarify", "need_mode", "handoff_offer", "wait_lost_reason", "wait_air_route"}: return sess logger.info( "%s.new_over_old phase=%s wo=%s", self.bookmark_kind, phase or "-", no, ) return None def drop_current_bookmark(self, sender_id: str) -> Optional[FlowSession]: """ 发图新开:丢掉当前书签,已出号会话仍留在 _by_ticket。 未出号的补问/核对草稿从此不再被文字当「当前单」。 HTTP 与 Worker 必须同时丢掉 Redis 里那份,否则补「托盘」会问运输方式。 """ sid = (sender_id or "").strip() if not sid: return None with self._lock: old = self._by_sender.pop(sid, None) self._delete_bookmark(sid) return old def _persist_bookmark(self, sess: FlowSession) -> None: from agent.redis_coord.flow_bookmark import save_flow_bookmark save_flow_bookmark( self.bookmark_kind, sess.sender_id, _session_payload(sess), store=self._bookmark_store, ) def _load_bookmark(self, sender_id: str) -> Optional[FlowSession]: from agent.redis_coord.flow_bookmark import load_flow_bookmark raw = load_flow_bookmark( self.bookmark_kind, sender_id, store=self._bookmark_store, ) if not raw: return None return _session_from_payload(raw) def _delete_bookmark(self, sender_id: str) -> None: from agent.redis_coord.flow_bookmark import delete_flow_bookmark delete_flow_bookmark( self.bookmark_kind, sender_id, store=self._bookmark_store, ) def _load_ticket_session(self, sender_id: str, work_order_no: str) -> Optional[FlowSession]: """ 点卡时进程内没有该书签:只读主账收回这一单。 找不到或不是本人:返回 None,调用方丢弃,绝不改成当前最新单。 """ getter = getattr(self._ledger, "get_for_agent", None) if not callable(getter): logger.warning("card 点到 %s 但无主账可读", work_order_no) return None view = getter(work_order_no=work_order_no, sender_id=sender_id) or {} if not view.get("found") or view.get("owned") is False: logger.warning( "card 点到 %s 主账未收回 found=%s owned=%s", work_order_no, view.get("found"), view.get("owned"), ) return None if (str(view.get("business_line") or "")).upper() == "MULTI": # 多段联运由 MultiTextInquiryFlow 接卡。这里收回会变成 tms_miss,点拉群没回声。 logger.info("card 多段工单不走空运收回 wo=%s", work_order_no) return None sess = self._session_from_ledger(sender_id=sender_id, view=view) self._save(sess) logger.info( "card 从主账收回 wo=%s phase=%s", sess.work_order_no, sess.phase, ) return sess def resume_historical( self, *, sender_id: str, work_order_no: str, reply: ReplyFn, ) -> str: """ 按工单号回到当时停点并重出对应卡片。 先用进程内该书签;没有再读主账推断节点。不重跑 TMS、不新建工单。 """ no = (work_order_no or "").strip().upper() if not no: reply("请告诉我要继续的工单号,例如 WO202609140034。") return "need_work_order" with self._lock: sess = self._by_ticket.get(no) if sess and sess.sender_id == sender_id: keep = sess else: keep = None if keep: self._align_resume_phase(keep) return self._replay_wait(keep, reply, source="session") getter = getattr(self._ledger, "get_for_agent", None) if not callable(getter): reply(f"找不到工单 {no},请核对工单号后再试。") return "ticket_missing" view = getter(work_order_no=no, sender_id=sender_id) or {} if not view.get("found"): reply(f"找不到工单 {no},请核对工单号后再试。") return "ticket_missing" if view.get("owned") is False: reply("这张工单不属于当前账号,不能继续。") return "ticket_forbidden" sess = self._session_from_ledger(sender_id=sender_id, view=view) self._align_resume_phase(sess) self._save(sess) return self._replay_wait(sess, reply, source="ledger") def _session_from_ledger(self, *, sender_id: str, view: dict[str, Any]) -> FlowSession: """主账快照收成运行书签。wait_phase 由主账推断,不是六态本身。""" no = str(view.get("work_order_no") or "").strip() facts = dict(view.get("facts") or {}) quote = dict(view.get("quote") or {}) status = str(view.get("status") or "").strip() phase = str(view.get("wait_phase") or "").strip() or self._infer_wait_phase( status=status, quote=quote ) allowed: tuple[str, ...] = () if phase == "wait_collab": allowed = (copy.BTN_SKIP_COLLAB,) elif phase == "wait_adopt": allowed = (copy.BTN_ADOPT_EXCEL, copy.BTN_ADOPT_PDF, copy.BTN_REJECT) elif phase == "wait_deal": allowed = (copy.BTN_DEAL, copy.BTN_LOST, copy.BTN_NEGOTIATE) elif phase == "wait_lost_reason": allowed = () return FlowSession( sender_id=sender_id, thread_id=sender_id, phase=phase, business_line=str(view.get("business_line") or "AIR"), facts=facts, work_order_no=no, quote=quote, status=status, allowed=allowed, wait_version=1, invited_sales_name=str(view.get("sales_name") or view.get("salesName") or ""), collab_chat_id=str(view.get("collab_chat_id") or ""), ) def _infer_wait_phase(self, *, status: str, quote: dict[str, Any]) -> str: if status in {"已成交", "已关闭"}: return "done" if status == "未成交": return "wait_lost_reason" if status == "协商中": return "wait_deal" if status == "已报价" and quote: return "wait_adopt" return "tms_miss" def _align_resume_phase(self, sess: FlowSession) -> None: """ 报价文件已经发给销售,续办必须回到成交跟进。 主账六态仍是「已报价」,不能据此重放到协同卡或生成报价单。 Redis 就绪键由出文件时写入;主账 wait_phase=wait_deal 也认。 """ if sess.phase not in {"wait_collab", "wait_adopt", "wait_file"}: return issued = bool((sess.quote or {}).get("quoteFileIssued")) if issued or self._quote_file_already_out(sess.work_order_no): logger.info("resume 报价文件已出,回到成交跟进 wo=%s was=%s", sess.work_order_no, sess.phase) self._promote_wait_deal(sess) def _replay_wait(self, sess: FlowSession, reply: ReplyFn, *, source: str) -> str: """按书签重出当时节点,不改六态。""" logger.info( "resume replay wo=%s phase=%s source=%s", sess.work_order_no, sess.phase, source, ) if sess.phase == "wait_adopt": mode = copy.scheme_transport_mode(sess.business_line) reply( copy.scheme_card( work_order_no=sess.work_order_no, status=sess.status or "已报价", facts=sess.facts, quote=sess.quote, transport_mode=mode, ), copy.scheme_wecom_payload( work_order_no=sess.work_order_no, facts=sess.facts, quote=sess.quote, status=sess.status or "已报价", transport_mode=mode, ), ) return "wait_adopt" if sess.phase == "wait_collab": self._emit_air_hit( reply, work_order_no=sess.work_order_no, quote=sess.quote, first_or_same=sess.first_or_same, history_work_order_no=sess.history_work_order_no, status=sess.status or "已报价", ) return "wait_collab" if sess.phase == "tms_miss": self._emit_tms_miss(reply, sess.work_order_no) return "tms_miss" if sess.phase == "wait_file": self._refresh_wait_file(sess) if sess.phase == "wait_file": reply("报价单正在生成,请稍候。") return "wait_file" if sess.phase == "wait_deal": # 成交卡已经绑在报价文件那条出站上。这里再发一张独立卡, # 会在文件还在上传时被另一条出站线程抢走,销售又先看到成交跟进。 return "wait_deal" if sess.phase == "wait_lost_reason": reply(copy.deal_ack("未成交", sess.work_order_no)) return "wait_lost_reason" if sess.phase == "done": reply(copy.deal_ack(sess.status or "已关闭", sess.work_order_no)) return "done" reply(copy.inquiry_card(work_order_no=sess.work_order_no, facts=sess.facts)) return sess.phase or "replay" def match_button(self, text: str, sess: Optional[FlowSession]) -> str: """精确匹配当前允许的按钮文案 / 企微 EventKey;不是按钮返回空。""" raw = (text or "").strip() if not sess or not raw: return "" canonical = copy.canonical_button(raw) if canonical and canonical in sess.allowed: return canonical for btn in sess.allowed: if raw == btn or raw == f"【{btn}】": return btn return "" def on_text( self, *, sender_id: str, text: str, reply: ReplyFn, injected_facts: Optional[dict[str, Any]] = None, injected_mode: str = "", b_result: Optional[dict[str, Any]] = None, allow_b: bool = False, chat_fn: Optional[Any] = None, inbound_raw: Optional[dict[str, Any]] = None, ) -> str: """ 处理一句私聊文字。返回阶段名供单测断言。 先认按钮(系统交互),再抽字段(模型/注入/B)。 点价格卡必须按卡片工单号找回当时会话,不能默认最新一单。 """ if not (sender_id or "").strip(): return "reject_no_sender" sess = self.session_for_action(sender_id, inbound_raw, text=text) sess = self._detach_frozen_bookmark(sess, sender_id) btn = self.match_button(text, sess) if btn: return self.on_button( sender_id=sender_id, action=btn, reply=reply, card_meta=inbound_raw, ) # 未成交后下一句任意文字即原因;完整新询价(注入字段)则空着原因开新单。 if ( sess and sess.phase == "wait_lost_reason" and not (injected_facts and injected_mode) ): return self._capture_lost_reason(sess, text, reply) if sess: from agent.policy.quote_adjust_ops import try_private_adjust adjusted = try_private_adjust(self, sess, text, reply) if adjusted: return adjusted if sess and sess.phase == "wait_confirm" and looks_like_confirm(text): # 确定查价不再抽字段,避免「确定」被模型改成空 facts。 logger.info("air.confirm sender=%s", sender_id) return self._confirm_or_hold( sender_id=sender_id, facts=dict(sess.facts), text=sess.immutable_text or text, reply=reply, sess=sess, ) if (not sess or sess.phase != "wait_confirm") and looks_like_confirm(text) and not sess: reply(copy.AIR_CONFIRM_EXPIRED) return "confirm_expired" sess = self._detach_ticketed_for_new_inquiry(sess) snap = extract_inquiry_snapshot( text, injected_facts=injected_facts, injected_mode=injected_mode, b_result=b_result, allow_b=allow_b, chat_fn=chat_fn, ) from agent.policy.system_exception import handle_extract_tech_fail llm_phase = handle_extract_tech_fail( snap=snap, ledger=self._ledger, reply=reply, sender_id=sender_id, sess=sess, ) if llm_phase: return llm_phase mode = (snap.get("business_line") or "").upper() facts = dict(snap.get("facts") or {}) if sess and sess.phase in {"clarify", "need_mode", "wait_confirm"}: # 补问轮:旧字段打底;运输方式沿用上一句(飞=空运) from agent.llm.extract_text import detect_transport_mode from agent.schema.field_validate import harvest_oral_measures merged = dict(sess.facts) merged.update(facts) facts = harvest_oral_measures(text, merged) mode = mode or (sess.business_line or "").upper() if not mode: mode = detect_transport_mode( f"{sess.immutable_text or ''}\n{text}" ) if not mode: reply(copy.ASK_TRANSPORT) self._save( FlowSession( sender_id=sender_id, thread_id=sender_id, phase="need_mode", business_line="", facts=facts, immutable_text=sess.immutable_text if sess and sess.immutable_text else text, ) ) return "need_mode" if mode != "AIR": if mode == "LAND": from agent.policy.land_text_flow import get_land_text_flow # 本句已是陆运:丢掉空运补问草稿。 self.drop_current_bookmark(sender_id) return get_land_text_flow().on_text( sender_id=sender_id, text=text, reply=reply, injected_facts=facts, injected_mode=mode, b_result=b_result, allow_b=allow_b, chat_fn=chat_fn, inbound_raw=inbound_raw, ) if mode == "SEA": # 海运交给海运 Owner,空运不再回「本轮先开通空运」 from agent.policy.sea_text_flow import get_sea_text_flow self.drop_current_bookmark(sender_id) return get_sea_text_flow().on_text( sender_id=sender_id, text=text, reply=reply, injected_facts=facts, injected_mode=mode, b_result=b_result, allow_b=allow_b, chat_fn=chat_fn, inbound_raw=inbound_raw, ) reply(copy.ASK_TRANSPORT) return "need_mode" check = validate_required_fields(facts=facts, business_line="AIR") facts = dict(check["facts"]) tms_ready = assemble_air_query(facts) facts = dict(tms_ready["facts"]) if not check["ok"] or not tms_ready["ok"]: keys = list(check["missing"] or []) for key in tms_ready.get("missing") or []: if key not in keys: keys.append(key) new_sess = FlowSession( sender_id=sender_id, thread_id=sender_id, phase="clarify", business_line="AIR", facts=facts, immutable_text=sess.immutable_text if sess and sess.immutable_text else text, wait_version=(sess.wait_version + 1) if sess else 1, ) return self._clarify_or_handoff( sess=sess, keys=keys, new_sess=new_sess, reply=reply, clarify_text=copy.ask_clarify( facts=facts, missing_keys=keys, transport_mode="空运", ), ) return self._hold_for_confirm( sender_id=sender_id, facts=facts, text=sess.immutable_text if sess and sess.immutable_text else text, reply=reply, sess=sess, ) def _hold_for_confirm( self, *, sender_id: str, facts: dict[str, str], text: str, reply: ReplyFn, sess: Optional[FlowSession], ) -> str: """必填已齐:出核对卡,等确定再建单。""" reply(copy.air_confirm_card(facts)) self._save( FlowSession( sender_id=sender_id, thread_id=sender_id, phase="wait_confirm", business_line="AIR", facts=facts, immutable_text=text, wait_version=(sess.wait_version + 1) if sess else 1, ) ) return "wait_confirm" def _confirm_or_hold( self, *, sender_id: str, facts: dict[str, str], text: str, reply: ReplyFn, sess: FlowSession, ) -> str: """确定意图:字段仍齐则建单查价。""" check = validate_required_fields(facts=facts, business_line="AIR") tms_ready = assemble_air_query(check["facts"]) if not check["ok"] or not tms_ready.get("ok"): keys = list(check["missing"] or []) for key in tms_ready.get("missing") or []: if key not in keys: keys.append(key) new_sess = FlowSession( sender_id=sender_id, thread_id=sender_id, phase="clarify", business_line="AIR", facts=dict(tms_ready.get("facts") or facts), immutable_text=sess.immutable_text or text, wait_version=sess.wait_version + 1, ) return self._clarify_or_handoff( sess=sess, keys=keys, new_sess=new_sess, reply=reply, clarify_text=copy.ask_clarify( facts=new_sess.facts, missing_keys=keys, transport_mode="空运", ), ) return self._create_and_quote( sender_id=sender_id, facts=dict(tms_ready["facts"]), text=text, reply=reply, ) def _create_and_quote( self, *, sender_id: str, facts: dict[str, str], text: str, reply: ReplyFn, ) -> str: # 先落已确认字段,避免建单失败后补问轮把毛重/体积弄丢 from agent.policy.handoff_human import void_virtual_no old = self.session_of(sender_id) if old: void_virtual_no(old) self._save( FlowSession( sender_id=sender_id, thread_id=sender_id, phase="quoting", business_line="AIR", facts=facts, immutable_text=text, ) ) created = self._ledger.create_ticket( sender_id=sender_id, business_line="AIR", facts=facts, # 每次出询价确认都新开一单;不用内容指纹,避免相同需求复用旧工单号 idempotency_key=f"create:{sender_id}:{uuid.uuid4().hex}", ) no = str(created.get("work_order_no") or "").strip() if not no: reply("建单失败,请稍后把刚才那条再发一次。重量和体积已经记下。") return "ledger_fail" first, history_wo = self._resolve_similar(created, facts=facts, current_no=no) out_id = reply(copy.inquiry_card(work_order_no=no, facts=facts, first_or_same=first)) or "" # 询价确认必须先到销售手机,再打 TMS。否则查得快时暂无报价卡会抢先发出。 self._wait_inquiry_sent(str(out_id), work_order_no=no) from agent.policy.system_exception import ( confirm_tms_quote_if_tech, query_tms_resilient, ) tms = query_tms_resilient(self._ledger, work_order_no=no, facts=facts) sess = FlowSession( sender_id=sender_id, thread_id=sender_id, business_line="AIR", facts=facts, work_order_no=no, first_or_same=first, history_work_order_no=history_wo, immutable_text=text, wait_version=1, ) if not tms.get("ok"): # 没拿到 TMS 回包(超时/熔断/主账未出站)不能冒充「暂无报价」 if confirm_tms_quote_if_tech( ledger=self._ledger, reply=reply, work_order_no=no, tms=tms, sender_id=sender_id, ticket_status="询价中", ): self.drop_current_bookmark(sender_id) return "system_exception" if tms.get("classification") == "TMS_QUERY_INCOMPLETE": logger.warning("tms.query 未出站 wo=%s class=%s", no, tms.get("classification")) reply("TMS 查询失败,请稍后重试。工单已创建:" + no) sess.phase = "tms_tech" self._save(sess) return "tms_tech" if not tms.get("has_price"): # ok=true 才是 TMS 已查过、确认没价 self._emit_tms_miss(reply, no) sess.phase = "tms_miss" sess.status = "询价中" sess.allowed = () self._save(sess) return "tms_miss" quote = tms["quote"] from agent.schema.air_route_options import ( needs_route_pick, parse_air_options, project_selected_quote, ) options = parse_air_options(quote) if options and not needs_route_pick(quote): quote = project_selected_quote(quote, options[0]) up = self._ledger.upsert_quote(work_order_no=no, quote=quote, to_status="已报价") if not up.get("ok"): reply("主账写入报价失败,不能当作已报价。") sess.phase = "ledger_fail" self._save(sess) return "ledger_fail" sess.quote = quote copy.remember_tms_quote(sess, quote) sess.tms_quote = dict(quote) sess.status = "已报价" if needs_route_pick(quote): return self._offer_air_routes( sess=sess, work_order_no=no, quote=quote, reply=reply, ) sess.phase = "wait_collab" sess.allowed = (copy.BTN_SKIP_COLLAB,) sess.wait_version += 1 self._save(sess) self._emit_air_hit( reply, work_order_no=no, quote=quote, first_or_same=first, history_work_order_no=history_wo, status="已报价", ) return "wait_collab" def _offer_air_routes( self, *, sess: FlowSession, work_order_no: str, quote: dict[str, Any], reply: ReplyFn, ) -> str: """多条线路:企微只发入口卡,选线在 H5。""" from agent.channel.h5.air_route_store import build_air_route_url, issue_air_route_token from agent.schema.air_route_options import parse_air_options token = issue_air_route_token( sender_id=sess.sender_id, work_order_no=work_order_no, quote=quote, facts=dict(sess.facts), ) url = build_air_route_url(token) count = len(parse_air_options(quote)) sess.phase = "wait_air_route" sess.allowed = () sess.wait_version += 1 self._save(sess) reply( copy.air_route_entry_text(work_order_no=work_order_no, count=count), copy.air_route_entry_card(work_order_no=work_order_no, count=count, url=url), ) return "wait_air_route" def select_air_route( self, *, sender_id: str, index: int, reply: ReplyFn, ) -> str: """ H5 选定一条线路:写回起运港,出价格明细卡。 已采用过则作废旧报价单后再出新卡。 """ sess = self.session_of(sender_id) if sess is None or not (sess.work_order_no or "").strip(): reply("线路选择已过期,请回企微重新询价。") return "route_expired" from agent.schema.air_route_options import ( apply_selected_option, parse_air_options, project_selected_quote, ) raw_quote = dict(sess.tms_quote or sess.quote or {}) options = parse_air_options(raw_quote) if index < 0 or index >= len(options): reply("没有这条线路,请重新打开页面选择。") return "route_invalid" option = options[index] facts = apply_selected_option(dict(sess.facts), option) sess.facts = facts quote = project_selected_quote(raw_quote, option) prev_phase = (sess.phase or "").strip() already_adopted = prev_phase in {"wait_deal", "wait_file", "wait_adopt"} if already_adopted: clearer = getattr(self._ledger, "void_quote_files", None) if callable(clearer): clearer(work_order_no=sess.work_order_no) from agent.redis_coord.quote_file import clear_quote_file_ready clear_quote_file_ready(sess.work_order_no) patcher = getattr(self._ledger, "update_facts", None) if callable(patcher): patcher(work_order_no=sess.work_order_no, facts=facts) self._ledger.upsert_quote(work_order_no=sess.work_order_no, quote=quote, to_status="已报价") sess.quote = quote copy.remember_tms_quote(sess, quote) sess.status = "已报价" sess.wait_version += 1 if already_adopted: # 旧报价单已作废,直接出方案卡让销售重新采用。 sess.phase = "wait_adopt" sess.allowed = (copy.BTN_ADOPT_EXCEL, copy.BTN_ADOPT_PDF, copy.BTN_REJECT) self._save(sess) mode = copy.scheme_transport_mode(sess.business_line) reply( copy.scheme_card( work_order_no=sess.work_order_no, status="已报价", facts=facts, quote=quote, transport_mode=mode, ), copy.scheme_wecom_payload( work_order_no=sess.work_order_no, facts=facts, quote=quote, status="已报价", transport_mode=mode, ), ) return "wait_adopt" sess.phase = "wait_collab" sess.allowed = (copy.BTN_SKIP_COLLAB,) self._save(sess) self._emit_air_hit( reply, work_order_no=sess.work_order_no, quote=quote, first_or_same=sess.first_or_same, history_work_order_no=sess.history_work_order_no, status="已报价", ) return "wait_collab" def _emit_air_hit( self, reply: ReplyFn, *, work_order_no: str, quote: dict[str, Any], first_or_same: str = "first", history_work_order_no: str = "", status: str = "已报价", ) -> None: """ 先发完整价格明细文本,等它真正发出,再发带按钮的卡。 企微引用条超长会截成省略号。海运已走这条;空运多项费用同样不能塞进卡。 只等这一条出站,不锁其它工单。 """ src = dict(quote or {}) out_id = reply( copy.sea_quote_detail_text( work_order_no=work_order_no, quote=src, first_or_same=first_or_same, history_work_order_no=history_work_order_no, status=status, business_line="AIR", ) ) or "" self._wait_outbound_sent( str(out_id), work_order_no=work_order_no, what="价格明细文本" ) gap = float(getattr(self, "_detail_before_card_gap_sec", 0.8) or 0) if out_id and gap > 0: time.sleep(gap) reply( copy.tms_hit_card( work_order_no=work_order_no, quote=src, first_or_same=first_or_same, history_work_order_no=history_work_order_no, status=status, ), copy.tms_hit_wecom_payload( work_order_no=work_order_no, quote=src, first_or_same=first_or_same, history_work_order_no=history_work_order_no, status=status, ), ) def _wait_inquiry_sent(self, out_id: str, *, work_order_no: str) -> None: """等询价确认出站完成后再查价。""" self._wait_outbound_sent(out_id, work_order_no=work_order_no, what="询价确认") def _wait_outbound_sent(self, out_id: str, *, work_order_no: str, what: str) -> None: """ 等指定出站真正发出。只等这一条,不锁其它工单。 单测 reply 不回 id 则跳过。超时仍继续,避免卡死 inbox。 """ raw = (out_id or "").strip() if not raw: return from agent.channel.queue import get_message_store store = get_message_store() waiter = getattr(store, "wait_outbound_done", None) if not callable(waiter): return ok = waiter(raw, timeout_sec=30.0) if not ok: logger.warning( "%s出站未在时限内完成 wo=%s out=%s,仍继续", what, work_order_no, raw, ) def _emit_tms_miss(self, reply: ReplyFn, work_order_no: str) -> None: """ TMS 已查过且确认没价:文字兜底 + 暂无报价卡片。 禁止在没打 query_tms 前回这个。 """ from agent.config import get_settings reply( copy.tms_miss(work_order_no=work_order_no), copy.tms_miss_wecom_payload( work_order_no=work_order_no, action_url=(get_settings().public_base_url or "").strip(), ), ) # 卡片发出后再单独发拉群句,不塞进卡、不当卡片失败兜底。 reply(copy.tms_miss_followup(work_order_no=work_order_no)) def _resolve_similar( self, created: dict[str, Any], *, facts: dict[str, str], current_no: str, ) -> tuple[str, str]: """ 全库相同需求:建单结果已带则用;否则再问主账。检索失败当首次,不影响出卡。 """ history = str(created.get("similar_work_order_no") or "").strip() first = str(created.get("first_or_same") or "first") if first == "same" and history and history != current_no: return "same", history finder = getattr(self._ledger, "find_similar_open", None) if not callable(finder): return "first", "" try: similar = finder( facts=facts, business_line="AIR", exclude_work_order_no=current_no ) or {} except TypeError: similar = finder(facts=facts, business_line="AIR") or {} except Exception: logger.warning("find_similar 失败,按首次询价展示 current=%s", current_no) return "first", "" history = str(similar.get("work_order_no") or "").strip() if similar.get("found") and history and history != current_no: return "same", history return "first", "" def on_button( self, *, sender_id: str, action: str, reply: ReplyFn, card_meta: Optional[dict[str, Any]] = None, ) -> str: """处理已允许的按钮;错版本/未允许/已点过直接丢弃。""" sess = self.session_for_action(sender_id, card_meta) card_wo = copy.work_order_from_card_meta(card_meta) if sess and sess.phase == "wait_file": self._refresh_wait_file(sess) # 成交卡已经发到销售手机,HTTP 进程可能还停在 wait_file,不能丢弃 if ( sess and action in {copy.BTN_DEAL, copy.BTN_LOST, copy.BTN_NEGOTIATE} and sess.phase == "wait_file" ): self._promote_wait_deal(sess) if not sess or action not in sess.allowed: logger.info( "card 丢弃 sender=%s action=%s card_wo=%s sess_wo=%s phase=%s", sender_id, action, card_wo or "-", (sess.work_order_no if sess else "-"), (sess.phase if sess else "-"), ) return "card_discard" if card_wo and sess.work_order_no and card_wo != sess.work_order_no: logger.warning( "card 工单不一致 card_wo=%s sess_wo=%s,拒绝串单", card_wo, sess.work_order_no, ) return "card_discard" if action == copy.BTN_SKIP_COLLAB: # 先摘掉按钮,连点第二次不再进生成 sess.allowed = tuple(x for x in sess.allowed if x != copy.BTN_SKIP_COLLAB) self._save(sess) return self._skip_collab(sess, reply, card_meta=card_meta) if action == copy.BTN_REJECT: self._grey_clicked_card(reply, card_meta) from agent.policy.event_exception import clear_wait clear_wait(self._ledger, work_order_no=sess.work_order_no) sess.phase = "wait_adjust" sess.allowed = () sess.wait_version += 1 self._save(sess) reply(copy.reject_ack()) return "wait_adjust" if action == copy.BTN_ADOPT_EXCEL: return self._adopt_file(sess, reply, kind="Excel", card_meta=card_meta) if action == copy.BTN_ADOPT_PDF: return self._adopt_file(sess, reply, kind="PDF", card_meta=card_meta) if action == copy.BTN_DEAL: return self._close(sess, reply, "已成交", card_meta=card_meta) if action == copy.BTN_LOST: return self._close(sess, reply, "未成交", card_meta=card_meta) if action == copy.BTN_NEGOTIATE: return self._close(sess, reply, "协商中", card_meta=card_meta) return "card_unknown" def _skip_collab( self, sess: FlowSession, reply: ReplyFn, card_meta: Optional[dict[str, Any]] = None, ) -> str: code = str((card_meta or {}).get("response_code") or "").strip() if code: reply( copy.BTN_SKIP_COLLAB_DONE, copy.skip_collab_card_update(response_code=code), ) reply(copy.skip_collab_ack()) sess.phase = "wait_adopt" sess.allowed = (copy.BTN_ADOPT_EXCEL, copy.BTN_ADOPT_PDF, copy.BTN_REJECT) sess.wait_version += 1 self._save(sess) mode = copy.scheme_transport_mode(sess.business_line) reply( copy.scheme_card( work_order_no=sess.work_order_no, status=sess.status or "已报价", facts=sess.facts, quote=sess.quote, transport_mode=mode, ), copy.scheme_wecom_payload( work_order_no=sess.work_order_no, facts=sess.facts, quote=sess.quote, status=sess.status or "已报价", transport_mode=mode, ), ) from agent.policy.event_exception import WAIT_ADOPT, start_wait start_wait(self._ledger, work_order_no=sess.work_order_no, wait_kind=WAIT_ADOPT) return "wait_adopt" def _grey_clicked_card( self, reply: ReplyFn, card_meta: Optional[dict[str, Any]], *, replace_name: str = "", ) -> None: """ 点过模板卡后整卡置灰,三个按钮都不可再点。 企微 update_template_card 一次更新后整卡按钮失效;无 ResponseCode 则不更新。 """ code = str((card_meta or {}).get("response_code") or "").strip() if not code: return reply( copy.BTN_SKIP_COLLAB_DONE, copy.skip_collab_card_update(response_code=code, replace_name=replace_name), ) def _adopt_file( self, sess: FlowSession, reply: ReplyFn, *, kind: str, card_meta: Optional[dict[str, Any]] = None, ) -> str: """ 点生成 Excel/PDF:按询价需求匹配后台报价模板再出单。 关键词未命中走默认模板;双无则只回话术,不进成交卡。 Excel 当场发文件;PDF 主账只填 xlsx,转 PDF 进 LibreOffice 单槽。 """ self._grey_clicked_card(reply, card_meta) if kind == "Excel": reply(copy.adopt_excel_ack()) else: reply(copy.adopt_pdf_ack()) fmt = "xlsx" if kind == "Excel" else "pdf" renderer = getattr(self._ledger, "render_quote", None) if not callable(renderer): reply(copy.quote_template_missing()) return "wait_adopt" result = renderer( work_order_no=sess.work_order_no, business_line=sess.business_line or "AIR", facts=dict(sess.facts or {}), quote=dict(sess.quote or {}), format=fmt, sales_wecom_id=sess.sender_id, ) or {} if not result.get("ok"): reply(str(result.get("message") or copy.quote_template_missing())) return "wait_adopt" filename = str(result.get("fileName") or f"{sess.work_order_no}_quote.{fmt}") from agent.policy.event_exception import clear_wait clear_wait(self._ledger, work_order_no=sess.work_order_no) file_extra = self._quote_file_payload(result, filename=filename) if kind == "PDF" and result.get("needsPdf"): queued = self._enqueue_quote_pdf( result, work_order_no=sess.work_order_no, sender_id=sess.sender_id, filename=filename if filename.lower().endswith(".pdf") else f"{sess.work_order_no}_quote.pdf", quote=dict(sess.quote or {}), facts=dict(sess.facts or {}), ) if not queued: fallback = str(result.get("fileName") or f"{sess.work_order_no}_quote.xlsx") if not fallback.lower().endswith(".xlsx"): fallback = f"{sess.work_order_no}_quote.xlsx" reply("PDF 转换排队失败,已按后台模板发送 Excel 报价单。") return self._emit_file_then_deal( sess, reply, kind="Excel", filename=fallback, file_extra=file_extra, ) # 已入队:等 Worker 出文件后再出成交卡,这里不提前发卡 sess.phase = "wait_file" sess.allowed = () sess.wait_version += 1 self._save(sess) return "wait_file" return self._emit_file_then_deal( sess, reply, kind=kind, filename=filename, file_extra=file_extra ) def _emit_file_then_deal( self, sess: FlowSession, reply: ReplyFn, *, kind: str, filename: str, file_extra: dict[str, Any], ) -> str: """ 报价文件与成交跟进走同一条出站:先传文件,成功后再发卡。 两张消息分两条出站时,另一条消费线程会在上传期间把成交卡发出去。 单测 reply 不回 id 则不等待。 """ extra = dict(file_extra or {}) extra["followup_card"] = copy.deal_wecom_payload( work_order_no=sess.work_order_no, quote=sess.quote, facts=sess.facts, business_line=sess.business_line or "AIR", ).get("template_card") or {} reply( "\n".join( [ copy.file_ready( kind=kind, work_order_no=sess.work_order_no, filename=filename ), copy.deal_card( work_order_no=sess.work_order_no, status=sess.status or "已报价", quote=sess.quote, facts=sess.facts, business_line=sess.business_line or "AIR", ), ] ), extra, ) from agent.redis_coord.quote_file import mark_quote_file_ready mark_quote_file_ready(sess.work_order_no) sess.phase = "wait_deal" sess.allowed = (copy.BTN_DEAL, copy.BTN_LOST, copy.BTN_NEGOTIATE) sess.wait_version += 1 self._save(sess) from agent.policy.event_exception import WAIT_DEAL, start_wait start_wait(self._ledger, work_order_no=sess.work_order_no, wait_kind=WAIT_DEAL) return "wait_deal" def _emit_deal(self, reply: ReplyFn, sess: FlowSession) -> None: """ 报价单发出后出成交跟进卡:文字兜底 + 企微模板卡。 字段以报价当时的 quote 为准,不编历史报价/工单状态/来源行。 """ reply( copy.deal_card( work_order_no=sess.work_order_no, status=sess.status or "已报价", quote=sess.quote, facts=sess.facts, business_line=sess.business_line or "AIR", ), copy.deal_wecom_payload( work_order_no=sess.work_order_no, quote=sess.quote, facts=sess.facts, business_line=sess.business_line or "AIR", ), ) def _quote_file_payload(self, result: dict[str, Any], *, filename: str) -> dict[str, Any]: """出站文件载荷:优先本机路径,否则主账回的 Base64。""" return { "msgtype": "file", "filename": filename, "filepath": str(result.get("filePath") or ""), "file_b64": str(result.get("fileBase64") or ""), } def _enqueue_quote_pdf( self, result: dict[str, Any], *, work_order_no: str, sender_id: str, filename: str, quote: Optional[dict[str, Any]] = None, facts: Optional[dict[str, str]] = None, ) -> bool: """ 把主账已填好的 xlsx 交给 Worker LibreOffice 单槽转 PDF。 失败返回 False,调用方改发 Excel,不在 inbox 线程跑 soffice。 """ raw_b64 = str(result.get("fileBase64") or "").strip() if not raw_b64: return False try: import base64 import tempfile from pathlib import Path from agent.jobs.pdf_job import enqueue_pdf_job data = base64.b64decode(raw_b64) work_dir = Path(tempfile.gettempdir()) / "inquiry_quote" work_dir.mkdir(parents=True, exist_ok=True) xlsx_path = work_dir / f"{work_order_no}_quote.xlsx" xlsx_path.write_bytes(data) queued = enqueue_pdf_job( xlsx_path=str(xlsx_path), inquiry_no=work_order_no, quote_version=1, idempotency_key=f"pdf:{work_order_no}:{uuid.uuid4().hex[:8]}", touser=sender_id, filename=filename, quote=quote, facts=facts, ) return bool(queued.accepted) except Exception: logger.exception("报价 PDF 入队失败 wo=%s", work_order_no) return False def _close( self, sess: FlowSession, reply: ReplyFn, status: str, card_meta: Optional[dict[str, Any]] = None, ) -> str: """点成交/未成交/协商中:先整卡置灰,再入主账。""" self._grey_clicked_card(reply, card_meta, replace_name=status) from agent.policy.deal_outcome import finish_deal_outcome result = finish_deal_outcome( ledger=self._ledger, work_order_no=sess.work_order_no, outcome=status, facts=sess.facts, quote=sess.quote, business_line=sess.business_line, sales_name=sess.invited_sales_name, ) if not result.get("ok"): reply("主账状态更新失败,成交结果未入账。") return "ledger_fail" phase = str(result.get("phase") or "done") sess.status = status sess.phase = phase if phase == "wait_deal": sess.allowed = (copy.BTN_DEAL, copy.BTN_LOST, copy.BTN_NEGOTIATE) else: sess.allowed = () sess.wait_version += 1 self._save(sess) reply(copy.deal_ack(status, sess.work_order_no, str(result.get("next_followup_at") or ""))) return phase def _capture_lost_reason(self, sess: FlowSession, text: str, reply: ReplyFn) -> str: """私聊未成交后收下原因,不再追问。""" from agent.policy.deal_outcome import capture_lost_reason result = capture_lost_reason( ledger=self._ledger, work_order_no=sess.work_order_no, reason=text, ) if not result.get("ok"): reply("未成交原因没有记下,请再发一句。") return "wait_lost_reason" sess.phase = "done" sess.allowed = () self._save(sess) reply(copy.lost_reason_ack(text)) return "done" def mark_wait_deal(self, sender_id: str, work_order_no: str) -> None: """ Worker 出完报价文件后:先写 Redis,再切本进程书签。 HTTP 与 Worker 不是同一份内存。只改本进程会让销售一直看到「请稍候」。 找不到本进程书签也要写 Redis,让收消息那边能切到成交跟进。 """ no = (work_order_no or "").strip() if not no: return from agent.redis_coord.quote_file import mark_quote_file_ready mark_quote_file_ready(no) if not sender_id: return with self._lock: sess = self._by_sender.get(sender_id) if sess is None or sess.work_order_no != no: sess = self._by_ticket.get(no) if sess is None: return self._promote_wait_deal(sess) def _promote_wait_deal(self, sess: FlowSession) -> None: """书签切到成交跟进,允许点成交/未成交/协商中。""" sess.phase = "wait_deal" sess.allowed = (copy.BTN_DEAL, copy.BTN_LOST, copy.BTN_NEGOTIATE) sess.wait_version += 1 self._save(sess) from agent.policy.event_exception import WAIT_DEAL, start_wait start_wait(self._ledger, work_order_no=sess.work_order_no, wait_kind=WAIT_DEAL) def _refresh_wait_file(self, sess: FlowSession) -> None: """ wait_file 时看文件是否已经发出。 Worker 写 Redis;本进程重启后也可看出站成交卡/文件去重键。 """ if sess.phase != "wait_file": return no = (sess.work_order_no or "").strip() if not no: return if self._quote_file_already_out(no): self._promote_wait_deal(sess) def _quote_file_already_out(self, work_order_no: str) -> bool: """文件或成交卡是否已经入过出站。""" from agent.redis_coord.quote_file import is_quote_file_ready if is_quote_file_ready(work_order_no): return True try: from agent.channel.queue import get_message_store store = get_message_store() probe = getattr(store, "has_outbound_dedupe", None) if not callable(probe): return False return bool( probe(f"deal-card:{work_order_no}:") or probe(f"pdf-file:{work_order_no}:") ) except Exception: logger.exception("判断报价文件是否已出站失败 wo=%s", work_order_no) return False _FLOW: Optional[AirTextInquiryFlow] = None _FLOW_GUARD = threading.Lock() def build_ledger() -> Any: """生产打 8180;单测 LEDGER_BACKEND=memory。""" from agent.config import get_settings backend = (get_settings().ledger_backend or "http").strip().lower() if backend == "memory": return MemoryLedger() return HttpLedger() def get_air_text_flow() -> AirTextInquiryFlow: """进程内一份流程(与 inbox 线程共用);测试可 new 独立实例。""" global _FLOW with _FLOW_GUARD: if _FLOW is None: _FLOW = AirTextInquiryFlow() return _FLOW def reset_air_text_flow_for_test() -> AirTextInquiryFlow: """单测重置全局实例(强制内存账本)。""" global _FLOW with _FLOW_GUARD: _FLOW = AirTextInquiryFlow(MemoryLedger()) return _FLOW