103 lines
3.4 KiB
Python
103 lines
3.4 KiB
Python
"""
|
||
意图:继续历史工单。
|
||
|
||
本文件职责:把 DeepSeek A 认出的续办意图交给流程,按工单号回到当时停点。
|
||
工单号优先用路由结论;模型漏填时再从原话机读 WO 号(不在这里判断意图)。
|
||
禁止:在本文件抽询价字段;禁止回调线程跑完整 LLM;禁止直连 TMS。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
|
||
from agent.channel.queue import get_message_store
|
||
from agent.channel.wecom.models import InboundMessage
|
||
from agent.policy import inquiry_copy as copy
|
||
from agent.policy.air_text_flow import get_air_text_flow
|
||
from agent.policy.land_text_flow import get_land_text_flow
|
||
from agent.policy.sea_text_flow import get_sea_text_flow
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
def _pick_continue_flow(*, sender_id: str, work_order_no: str, flow=None):
|
||
"""
|
||
续办按工单选空运/海运流程。
|
||
|
||
海运书签或主账业务线是 SEA 必须走海运 Owner。
|
||
一律走空运会把已报价海运重放成空运协同卡(件数/毛重/包装)。
|
||
单测传入 flow 时不改选路,保持原断言。
|
||
"""
|
||
if flow is not None:
|
||
return flow
|
||
sea = get_sea_text_flow()
|
||
no = (work_order_no or "").strip()
|
||
hit = sea.session_by_ticket(sender_id, no) if no else None
|
||
if hit and (hit.business_line or "").upper() == "SEA":
|
||
return sea
|
||
land = get_land_text_flow()
|
||
land_hit = land.session_by_ticket(sender_id, no) if no else None
|
||
if land_hit and (land_hit.business_line or "").upper() == "LAND":
|
||
return land
|
||
getter = getattr(sea.ledger, "get_for_agent", None)
|
||
if callable(getter) and no:
|
||
try:
|
||
view = getter(work_order_no=no, sender_id=sender_id) or {}
|
||
except Exception:
|
||
logger.warning("continue_thread 读主账业务线失败 wo=%s", no)
|
||
view = {}
|
||
line = str(view.get("business_line") or "").upper()
|
||
if line == "SEA":
|
||
return sea
|
||
if line == "LAND":
|
||
return land
|
||
return get_air_text_flow()
|
||
|
||
|
||
def handle_continue_thread(
|
||
message: InboundMessage,
|
||
*,
|
||
thread_id: str = "",
|
||
flow=None,
|
||
store=None,
|
||
) -> str:
|
||
"""
|
||
续办一张历史工单:回到该单当时节点的卡片/话术。
|
||
|
||
thread_id:RouteDecision 里的工单号。副作用:outbox 回复;读主账。
|
||
"""
|
||
if not message.has_sender():
|
||
logger.warning("continue_thread 拒绝:无 sender_id")
|
||
return "reject_no_sender"
|
||
msg_store = store or get_message_store()
|
||
wo = (thread_id or "").strip() or copy.extract_work_order_no(message.content or "")
|
||
engine = _pick_continue_flow(
|
||
sender_id=message.sender_id, work_order_no=wo, flow=flow
|
||
)
|
||
replies: list[str] = []
|
||
|
||
def reply(text: str, extra: dict | None = None) -> str:
|
||
replies.append(text)
|
||
ok, info = msg_store.enqueue_outbound(
|
||
touser=message.sender_id,
|
||
content=text,
|
||
dedupe_key=f"cont:{message.message_id}:{len(replies)}",
|
||
payload=dict(extra or {}),
|
||
)
|
||
logger.info("continue_thread 出站 ok=%s info=%s", ok, info)
|
||
return info if ok else ""
|
||
|
||
phase = engine.resume_historical(
|
||
sender_id=message.sender_id,
|
||
work_order_no=wo,
|
||
reply=reply,
|
||
text=message.content or "",
|
||
)
|
||
logger.info(
|
||
"continue_thread phase=%s sender=%s wo=%s",
|
||
phase,
|
||
message.sender_id,
|
||
wo or "-",
|
||
)
|
||
return phase
|