点名海运时丢掉旧陆运补问草稿,避免补字段误发填写样例图。 Co-authored-by: Cursor <cursoragent@cursor.com>
901 lines
34 KiB
Python
901 lines
34 KiB
Python
"""
|
|
海运文字询价确定性流程。
|
|
|
|
本文件职责:海运七项补问 → 建单出确认卡 → 主账查 TMS → 有价/无价卡 → 跳过协同或拉群。
|
|
空运 Owner 仍在 air_text_flow,禁止把海运再堆回去。
|
|
调用:handler 在 Worker/inbox 线程;禁止回调线程同步跑。
|
|
禁止:直连 TMS;本地拆箱型;用户文字冒充群报价版本。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import threading
|
|
from typing import Any, 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.policy.air_text_flow import (
|
|
AirTextInquiryFlow,
|
|
FlowSession,
|
|
ReplyFn,
|
|
build_ledger,
|
|
)
|
|
from agent.schema.field_validate import validate_required_fields
|
|
from agent.schema.tms_sea_query import assemble_sea_query
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# 会话存档账号:每个协同群必拉。私聊「已邀请」不写这些人。
|
|
# 姓名与后台员工管理一致;userid 以员工表/环境变量为准,源码不写死。
|
|
ARCHIVE_SEAT_NAMES = ("询价机器人",)
|
|
|
|
|
|
class SeaTextInquiryFlow(AirTextInquiryFlow):
|
|
"""
|
|
海运文字询价 Owner。
|
|
|
|
复用空运后半段(跳过协同 / 出方案卡 / Excel PDF / 成交),只改查价前与卡片、拉群。
|
|
"""
|
|
|
|
bookmark_kind = "sea"
|
|
|
|
def __init__(self, ledger: Any = None, group_client: Any = None) -> None:
|
|
super().__init__(ledger=ledger)
|
|
self._group_client = group_client
|
|
|
|
def group_client(self) -> Any:
|
|
"""
|
|
企微应用群客户端。测服正式不注入,现场再取默认 WeCom 客户端。
|
|
|
|
单测可注入 FakeGroupClient。取到后挂在流程上,后续补字段/发文件共用。
|
|
"""
|
|
if self._group_client is None:
|
|
self._group_client = _default_group_client()
|
|
return self._group_client
|
|
|
|
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:
|
|
"""
|
|
处理一句海运私聊文字。返回阶段名供单测断言。
|
|
|
|
先认按钮,再抽字段。运输方式必须是 SEA 才建单。
|
|
"""
|
|
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
|
|
# 上一单已出号:新一轮从本句重收,禁止把旧贸易条款/分类带进下一张工单。
|
|
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,
|
|
)
|
|
mode = (snap.get("business_line") or "").upper()
|
|
facts = dict(snap.get("facts") or {})
|
|
if sess and sess.phase in {"clarify", "need_mode"}:
|
|
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 != "SEA":
|
|
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,
|
|
)
|
|
# 空运交给空运 Owner,避免海运流程半做空运
|
|
from agent.policy.air_text_flow import get_air_text_flow
|
|
|
|
self.drop_current_bookmark(sender_id)
|
|
return get_air_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,
|
|
)
|
|
|
|
from agent.schema.sea_options import drop_invented_door_trade_terms, harvest_sea_class
|
|
|
|
facts = harvest_sea_class(text, facts)
|
|
# 旧陆运/补问草稿可能把「门到门」塞进贸易条款;原话没写条款则丢掉。
|
|
facts = drop_invented_door_trade_terms(text, facts)
|
|
check = validate_required_fields(facts=facts, business_line="SEA")
|
|
facts = dict(check["facts"])
|
|
tms_ready = assemble_sea_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="SEA",
|
|
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._create_and_quote(
|
|
sender_id=sender_id,
|
|
facts=facts,
|
|
text=text,
|
|
reply=reply,
|
|
)
|
|
|
|
def _create_and_quote(
|
|
self,
|
|
*,
|
|
sender_id: str,
|
|
facts: dict[str, str],
|
|
text: str,
|
|
reply: ReplyFn,
|
|
) -> str:
|
|
"""建海运工单、先出确认卡、再查 TMS。"""
|
|
import uuid
|
|
|
|
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="SEA",
|
|
facts=facts,
|
|
immutable_text=text,
|
|
)
|
|
)
|
|
created = self._ledger.create_ticket(
|
|
sender_id=sender_id,
|
|
business_line="SEA",
|
|
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,
|
|
transport_mode="海运",
|
|
)
|
|
) or ""
|
|
self._wait_inquiry_sent(str(out_id), work_order_no=no)
|
|
|
|
tms_facts = dict(assemble_sea_query(facts).get("tms_facts") or facts)
|
|
tms = self._ledger.query_tms(work_order_no=no, facts=tms_facts, business_line="SEA")
|
|
sess = FlowSession(
|
|
sender_id=sender_id,
|
|
thread_id=sender_id,
|
|
business_line="SEA",
|
|
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"):
|
|
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"):
|
|
self._emit_tms_miss(reply, no)
|
|
sess.phase = "tms_miss"
|
|
sess.status = "询价中"
|
|
sess.allowed = (copy.BTN_PULL_COLLAB,)
|
|
self._save(sess)
|
|
return "tms_miss"
|
|
|
|
quote = tms["quote"]
|
|
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.status = "已报价"
|
|
sess.phase = "wait_collab"
|
|
sess.allowed = (copy.BTN_PULL_COLLAB, copy.BTN_SKIP_COLLAB)
|
|
sess.wait_version += 1
|
|
self._save(sess)
|
|
self._emit_sea_hit(
|
|
reply,
|
|
work_order_no=no,
|
|
quote=quote,
|
|
first_or_same=first,
|
|
history_work_order_no=history_wo,
|
|
status="已报价",
|
|
business_line=sess.business_line or "SEA",
|
|
)
|
|
return "wait_collab"
|
|
|
|
def _resolve_similar(
|
|
self,
|
|
created: dict[str, Any],
|
|
*,
|
|
facts: dict[str, str],
|
|
current_no: str,
|
|
business_line: str = "SEA",
|
|
) -> 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)
|
|
line = (business_line or "SEA").strip() or "SEA"
|
|
if not callable(finder):
|
|
return "first", ""
|
|
try:
|
|
similar = finder(
|
|
facts=facts, business_line=line, exclude_work_order_no=current_no
|
|
) or {}
|
|
except TypeError:
|
|
similar = finder(facts=facts, business_line=line) 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 _emit_tms_miss(self, reply: ReplyFn, work_order_no: str) -> None:
|
|
"""海运无价:一张带拉群按钮的卡,不再发手动拉人群句。"""
|
|
from agent.config import get_settings
|
|
|
|
reply(
|
|
copy.sea_tms_miss(work_order_no=work_order_no),
|
|
copy.sea_tms_miss_wecom_payload(
|
|
work_order_no=work_order_no,
|
|
action_url=(get_settings().public_base_url or "").strip(),
|
|
),
|
|
)
|
|
|
|
def _session_from_ledger(self, *, sender_id: str, view: dict[str, Any]) -> FlowSession:
|
|
sess = super()._session_from_ledger(sender_id=sender_id, view=view)
|
|
if sess.business_line in {"SEA", "LAND"} and sess.phase == "wait_collab":
|
|
sess.allowed = (copy.BTN_PULL_COLLAB, copy.BTN_SKIP_COLLAB)
|
|
if sess.business_line in {"SEA", "LAND"} and sess.phase == "tms_miss":
|
|
sess.allowed = (copy.BTN_PULL_COLLAB,)
|
|
sess.collab_chat_id = str(view.get("collab_chat_id") or "")
|
|
sess.collab_facts = dict(view.get("collab_facts") or {})
|
|
return sess
|
|
|
|
def _emit_sea_hit(
|
|
self,
|
|
reply: ReplyFn,
|
|
*,
|
|
work_order_no: str,
|
|
quote: dict[str, Any],
|
|
first_or_same: str = "first",
|
|
history_work_order_no: str = "",
|
|
status: str = "已报价",
|
|
business_line: str = "",
|
|
) -> None:
|
|
"""
|
|
先发完整价格明细文本,再发带按钮的卡。
|
|
|
|
企微只发卡时,引用条超长会被截断;文本不受限。
|
|
"""
|
|
line = (business_line or getattr(self, "business_line", "") or "SEA").upper()
|
|
if line not in {"SEA", "LAND"}:
|
|
line = "SEA"
|
|
src = dict(quote or {})
|
|
src["business_line"] = line
|
|
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=line,
|
|
)
|
|
)
|
|
reply(
|
|
copy.sea_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,
|
|
business_line=line,
|
|
),
|
|
copy.sea_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,
|
|
business_line=line,
|
|
),
|
|
)
|
|
|
|
def _replay_wait(self, sess: FlowSession, reply: ReplyFn, *, source: str) -> str:
|
|
if sess.business_line == "SEA" and sess.phase == "wait_collab":
|
|
self._emit_sea_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 "已报价",
|
|
business_line=sess.business_line or "SEA",
|
|
)
|
|
return "wait_collab"
|
|
if sess.business_line == "SEA" and sess.phase == "tms_miss":
|
|
self._emit_tms_miss(reply, sess.work_order_no)
|
|
return "tms_miss"
|
|
return super()._replay_wait(sess, reply, source=source)
|
|
|
|
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)
|
|
if sess and action == copy.BTN_PULL_COLLAB and action in sess.allowed:
|
|
if sess.handoff_offer and not sess.handoff_committed:
|
|
return self._commit_handoff_pull(sess, reply, card_meta=card_meta)
|
|
return self._pull_collab_group(sess, reply, card_meta=card_meta)
|
|
if action in {copy.BTN_DEAL, copy.BTN_LOST, copy.BTN_NEGOTIATE}:
|
|
from agent.policy.sea_group_ops import explain_deal_click
|
|
|
|
ticket_sess = sess
|
|
if ticket_sess is None:
|
|
wo = copy.work_order_from_card_meta(card_meta)
|
|
ticket_sess = self.session_by_work_order(wo)
|
|
if ticket_sess and (ticket_sess.collab_chat_id or "").strip():
|
|
reason = explain_deal_click(
|
|
flow=self,
|
|
sender_id=sender_id,
|
|
action=action,
|
|
card_meta=card_meta,
|
|
sess=ticket_sess,
|
|
)
|
|
if reason:
|
|
return reason
|
|
return super().on_button(
|
|
sender_id=sender_id, action=action, reply=reply, card_meta=card_meta
|
|
)
|
|
|
|
def _pull_collab_group(
|
|
self,
|
|
sess: FlowSession,
|
|
reply: ReplyFn,
|
|
card_meta: Optional[dict[str, Any]] = None,
|
|
) -> str:
|
|
"""
|
|
点「拉产品进群协同」:按线路拉产品人,新建群,私聊回执,按钮置灰。
|
|
|
|
没人可拉 / 建群失败:说明原因,按钮保持可点。已有群:不新建,置灰。
|
|
"""
|
|
no = (sess.work_order_no or "").strip()
|
|
if sess.collab_ended:
|
|
reply(copy.group_cannot_recreate(work_order_no=no))
|
|
return "collab_ended"
|
|
getter = getattr(self._ledger, "get_ticket", None)
|
|
ticket = getter(work_order_no=no) if callable(getter) else None
|
|
if ticket is not None:
|
|
facts = dict(getattr(ticket, "collab_facts", {}) or {})
|
|
if bool(getattr(ticket, "collab_ended", False)) or str(
|
|
facts.get("__ended") or ""
|
|
).strip() in {"1", "true"}:
|
|
sess.collab_ended = True
|
|
self._save(sess)
|
|
reply(copy.group_cannot_recreate(work_order_no=no))
|
|
return "collab_ended"
|
|
existing = (sess.collab_chat_id or "").strip()
|
|
if not existing and ticket is not None:
|
|
existing = str(getattr(ticket, "collab_chat_id", "") or "")
|
|
if existing:
|
|
sess.collab_chat_id = existing
|
|
sess.allowed = ()
|
|
sess.phase = "collab_group"
|
|
self._save(sess)
|
|
self._grey_clicked_card(reply, card_meta)
|
|
self._ensure_archive_seats_in_group(existing)
|
|
try:
|
|
from agent.channel.wshoto.rooms import watch_collab_room
|
|
|
|
watch_collab_room(existing, work_order_no=no)
|
|
except Exception:
|
|
logger.exception("登记已有存档监听群失败 chat=%s", existing)
|
|
reply(copy.sea_group_exists(work_order_no=no))
|
|
return "collab_group"
|
|
|
|
staff = self._match_sea_staff(sess)
|
|
if not staff:
|
|
reply(self._group_no_staff_copy())
|
|
return sess.phase or "wait_collab"
|
|
|
|
userids = [sess.sender_id]
|
|
product_names: list[str] = []
|
|
product_ids: list[str] = []
|
|
for row in staff:
|
|
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
|
|
name = str(row.get("name") or "").strip() or uid
|
|
if not uid:
|
|
continue
|
|
if uid not in userids:
|
|
userids.append(uid)
|
|
if uid not in product_ids:
|
|
product_ids.append(uid)
|
|
product_names.append(name)
|
|
if not product_ids:
|
|
reply(self._group_no_staff_copy())
|
|
return sess.phase or "wait_collab"
|
|
|
|
# 销售 + 产品之后,必须再拉会话存档账号;没有也不挡建群。
|
|
for uid in self._archive_seat_userids():
|
|
if uid not in userids:
|
|
userids.append(uid)
|
|
|
|
client = self.group_client()
|
|
created = client.create_group(name=no, userids=userids) or {}
|
|
if not created.get("ok") or not created.get("chat_id"):
|
|
reply(copy.sea_group_create_fail(reason=str(created.get("error") or "")))
|
|
return sess.phase or "wait_collab"
|
|
|
|
chat_id = str(created.get("chat_id") or "").strip()
|
|
try:
|
|
from agent.channel.wshoto.rooms import watch_collab_room
|
|
|
|
watch_collab_room(chat_id, work_order_no=no)
|
|
except Exception:
|
|
logger.exception("登记存档监听群失败 chat=%s", chat_id)
|
|
sales_name = self._sales_display_name(sess.sender_id)
|
|
binder = getattr(self._ledger, "bind_collab_group", None)
|
|
if callable(binder):
|
|
binder(
|
|
work_order_no=no,
|
|
chat_id=chat_id,
|
|
member_ids=userids,
|
|
product_names=product_names,
|
|
product_ids=product_ids,
|
|
)
|
|
|
|
from agent.schema.field_validate import collab_fields_for_ticket, pick_collab_from_facts
|
|
|
|
line = (sess.business_line or "SEA").upper()
|
|
keys = collab_fields_for_ticket(sess.facts, line)
|
|
seeded = pick_collab_from_facts(sess.facts, keys)
|
|
if seeded:
|
|
merged = dict(sess.collab_facts or {})
|
|
for k, v in seeded.items():
|
|
if v and not str(merged.get(k) or "").strip():
|
|
merged[k] = v
|
|
sess.collab_facts = merged
|
|
self.apply_collab_facts(
|
|
work_order_no=no,
|
|
facts=seeded,
|
|
sender_id=sess.sender_id,
|
|
allowed_keys=keys,
|
|
)
|
|
|
|
ticket_view = ticket if isinstance(ticket, dict) else {}
|
|
tms = copy.resolve_tms_quote(sess, ticket_view)
|
|
brief = copy.sea_group_brief(
|
|
work_order_no=no,
|
|
facts=sess.facts,
|
|
quote=sess.quote,
|
|
collab_facts=sess.collab_facts,
|
|
sales_name=sales_name,
|
|
product_names=product_names,
|
|
has_price=bool(tms or sess.quote),
|
|
tms_quote=tms,
|
|
product_quote=sess.quote if copy.quote_is_product(sess.quote) else None,
|
|
transport_mode=self._group_brief_mode(sess),
|
|
)
|
|
client.send_group(chat_id=chat_id, content=brief)
|
|
sess.collab_chat_id = chat_id
|
|
from agent.policy.sea_group_ops import send_filled_sheet
|
|
|
|
send_filled_sheet(
|
|
flow=self,
|
|
ticket=None,
|
|
sess=sess,
|
|
business_line=(sess.business_line or "SEA"),
|
|
)
|
|
sess.invited_userids = tuple(userids)
|
|
sess.invited_sales_name = sales_name
|
|
sess.invited_product_names = tuple(product_names)
|
|
sess.phase = "collab_group"
|
|
sess.allowed = ()
|
|
sess.wait_version += 1
|
|
self._save(sess)
|
|
self._grey_clicked_card(reply, card_meta)
|
|
reply(
|
|
copy.sea_group_created(
|
|
work_order_no=no,
|
|
sales_name=sales_name,
|
|
product_names=product_names,
|
|
)
|
|
)
|
|
return "collab_group"
|
|
|
|
def _match_handoff_staff(self, sess: FlowSession) -> list[dict[str, Any]]:
|
|
"""
|
|
转人工拉人:线路/港口已齐走原匹配;未齐则拉该运输方式全部启用产品。
|
|
"""
|
|
line = (sess.business_line or "SEA").upper()
|
|
facts = dict(sess.facts or {})
|
|
if line == "LAND":
|
|
if str(facts.get("线路类别") or "").strip():
|
|
return self._match_sea_staff(sess)
|
|
lister = getattr(self._ledger, "list_staff_by_role", None)
|
|
if not callable(lister):
|
|
return []
|
|
try:
|
|
rows = lister(role_code="land") or []
|
|
except Exception:
|
|
logger.exception("转人工拉全部陆运产品失败")
|
|
return []
|
|
return [dict(x) for x in rows if isinstance(x, dict)]
|
|
origin = str(facts.get("起运港") or facts.get("起运地") or "").strip()
|
|
dest = str(facts.get("目的港") or facts.get("目的地") or "").strip()
|
|
if origin or dest:
|
|
return self._match_sea_staff(sess)
|
|
lister = getattr(self._ledger, "list_staff_by_role", None)
|
|
if not callable(lister):
|
|
return []
|
|
try:
|
|
rows = lister(role_code="sea") or []
|
|
except Exception:
|
|
logger.exception("转人工拉全部海运产品失败")
|
|
return []
|
|
return [dict(x) for x in rows if isinstance(x, dict)]
|
|
|
|
def _commit_handoff_pull(
|
|
self,
|
|
sess: FlowSession,
|
|
reply: ReplyFn,
|
|
card_meta: Optional[dict[str, Any]] = None,
|
|
) -> str:
|
|
"""
|
|
点转人工「拉产品进群」:建群、发转人工摘要、正式单改「转人工」。
|
|
虚拟单不写主账。不发 TMS、不发模板、不跟成交。
|
|
"""
|
|
from agent.policy.handoff_human import (
|
|
HANDOFF_STATUS,
|
|
ensure_virtual_no,
|
|
is_virtual_no,
|
|
)
|
|
|
|
if sess.collab_ended or (sess.collab_chat_id or "").strip():
|
|
reply(copy.HANDOFF_ALREADY_COLLAB)
|
|
return "handoff_already"
|
|
display_no = ensure_virtual_no(sess) if not (sess.work_order_no or "").strip() else sess.work_order_no
|
|
staff = self._match_handoff_staff(sess)
|
|
if not staff:
|
|
reply(self._group_no_staff_copy())
|
|
return sess.phase or "handoff_offer"
|
|
userids = [sess.sender_id]
|
|
product_names: list[str] = []
|
|
product_ids: list[str] = []
|
|
for row in staff:
|
|
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
|
|
name = str(row.get("name") or "").strip() or uid
|
|
if not uid:
|
|
continue
|
|
if uid not in userids:
|
|
userids.append(uid)
|
|
if uid not in product_ids:
|
|
product_ids.append(uid)
|
|
product_names.append(name)
|
|
if not product_ids:
|
|
reply(self._group_no_staff_copy())
|
|
return sess.phase or "handoff_offer"
|
|
for uid in self._archive_seat_userids():
|
|
if uid not in userids:
|
|
userids.append(uid)
|
|
client = self.group_client()
|
|
created = client.create_group(name=display_no, userids=userids) or {}
|
|
if not created.get("ok") or not created.get("chat_id"):
|
|
reply(copy.sea_group_create_fail(reason=str(created.get("error") or "")))
|
|
return sess.phase or "handoff_offer"
|
|
chat_id = str(created.get("chat_id") or "").strip()
|
|
brief = copy.handoff_group_brief(
|
|
work_order_no=display_no,
|
|
facts=sess.facts,
|
|
business_line=sess.business_line,
|
|
product_names=product_names,
|
|
)
|
|
client.send_group(chat_id=chat_id, content=brief)
|
|
real = (sess.work_order_no or "").strip()
|
|
if real and not is_virtual_no(real):
|
|
trans = getattr(self._ledger, "transition", None)
|
|
if callable(trans):
|
|
trans(
|
|
work_order_no=real,
|
|
to_status=HANDOFF_STATUS,
|
|
remark=copy.HANDOFF_CLARIFY_REASON
|
|
if sess.handoff_offer == "clarify"
|
|
else copy.HANDOFF_INTENT_REASON,
|
|
)
|
|
binder = getattr(self._ledger, "bind_collab_group", None)
|
|
if callable(binder):
|
|
binder(
|
|
work_order_no=real,
|
|
chat_id=chat_id,
|
|
member_ids=userids,
|
|
product_names=product_names,
|
|
product_ids=product_ids,
|
|
)
|
|
sess.status = HANDOFF_STATUS
|
|
sess.collab_chat_id = chat_id
|
|
sess.handoff_committed = True
|
|
sess.handoff_offer = ""
|
|
sess.phase = "handoff_done"
|
|
sess.allowed = ()
|
|
sess.invited_product_names = tuple(product_names)
|
|
self._save(sess)
|
|
self._grey_clicked_card(reply, card_meta)
|
|
sales_name = self._sales_display_name(sess.sender_id)
|
|
reply(
|
|
copy.sea_group_created(
|
|
work_order_no=display_no,
|
|
sales_name=sales_name,
|
|
product_names=product_names,
|
|
)
|
|
)
|
|
return "handoff_done"
|
|
|
|
def _match_sea_staff(self, sess: FlowSession) -> list[dict[str, Any]]:
|
|
"""主账按起运/目的匹配启用的海运产品岗。"""
|
|
matcher = getattr(self._ledger, "match_sea_staff", None)
|
|
if not callable(matcher):
|
|
return []
|
|
try:
|
|
rows = matcher(
|
|
origin=str((sess.facts or {}).get("起运港") or ""),
|
|
destination=str((sess.facts or {}).get("目的港") or ""),
|
|
) or []
|
|
except Exception:
|
|
logger.exception("匹配海运产品人员失败 wo=%s", sess.work_order_no)
|
|
return []
|
|
return [dict(x) for x in rows if isinstance(x, dict)]
|
|
|
|
def _group_no_staff_copy(self) -> str:
|
|
"""没人可拉时的私聊。海运用海运岗话术;陆运覆盖成陆运岗。"""
|
|
return copy.sea_group_no_staff()
|
|
|
|
def _group_brief_mode(self, sess: FlowSession) -> str:
|
|
"""群摘要运输方式展示。"""
|
|
_ = sess
|
|
return "海运"
|
|
|
|
def _archive_seat_userids(self) -> list[str]:
|
|
"""
|
|
会话存档账号 userid 列表。
|
|
|
|
先读环境变量 WECOM_ARCHIVE_SEAT_USER_IDS,再按姓名「询价机器人」查员工管理。
|
|
查不到只记日志,不把建群判失败。
|
|
"""
|
|
ids: list[str] = []
|
|
from agent.config import get_settings
|
|
|
|
seat = str(getattr(get_settings(), "wecom_archive_seat_user_ids", "") or "").strip()
|
|
for part in seat.replace(";", ",").split(","):
|
|
uid = part.strip()
|
|
if uid and uid not in ids:
|
|
ids.append(uid)
|
|
finder = getattr(self._ledger, "find_staff_by_name", None)
|
|
if callable(finder):
|
|
for name in ARCHIVE_SEAT_NAMES:
|
|
try:
|
|
row = finder(name=name) or {}
|
|
except Exception:
|
|
logger.exception("查会话存档账号失败 name=%s", name)
|
|
continue
|
|
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
|
|
if uid and uid not in ids:
|
|
ids.append(uid)
|
|
if not ids:
|
|
logger.warning("协同群缺少会话存档账号:员工管理未找到「询价机器人」且未配置 WECOM_ARCHIVE_SEAT_USER_IDS")
|
|
return ids
|
|
|
|
def _ensure_archive_seats_in_group(self, chat_id: str) -> None:
|
|
"""已有群补拉会话存档账号;失败只记日志,不改销售侧已有群提示。"""
|
|
seats = self._archive_seat_userids()
|
|
if not chat_id or not seats:
|
|
return
|
|
client = self.group_client()
|
|
adder = getattr(client, "add_group_members", None)
|
|
if not callable(adder):
|
|
return
|
|
try:
|
|
adder(chat_id=chat_id, userids=seats)
|
|
except Exception:
|
|
logger.exception("已有协同群补拉会话存档失败 chat=%s", chat_id)
|
|
|
|
def _sales_display_name(self, sender_id: str) -> str:
|
|
"""私聊「已邀请」用的销售名;没有展示名就用 userid。"""
|
|
finder = getattr(self._ledger, "get_staff_by_wecom_id", None)
|
|
if callable(finder):
|
|
try:
|
|
row = finder(wecom_id=sender_id) or {}
|
|
name = str(row.get("name") or "").strip()
|
|
if name:
|
|
return name
|
|
except Exception:
|
|
logger.warning("查销售姓名失败 sender=%s", sender_id)
|
|
return sender_id
|
|
|
|
def apply_collab_facts(
|
|
self,
|
|
*,
|
|
work_order_no: str,
|
|
facts: dict[str, str],
|
|
sender_id: str,
|
|
allowed_keys: tuple[str, ...] | None = None,
|
|
) -> dict[str, Any]:
|
|
"""
|
|
群里静默写入协同字段。不改六态、不回复、不出报价版本。
|
|
|
|
只合并非空协同键。
|
|
"""
|
|
from agent.schema.field_validate import SEA_COLLAB_FIELDS
|
|
|
|
no = (work_order_no or "").strip()
|
|
keys = tuple(allowed_keys) if allowed_keys is not None else SEA_COLLAB_FIELDS
|
|
cleaned = {
|
|
k: str(v).strip()
|
|
for k, v in (facts or {}).items()
|
|
if k in keys and str(v or "").strip()
|
|
}
|
|
if not no or not cleaned:
|
|
return {"ok": False, "updated": []}
|
|
with self._lock:
|
|
sess = self._by_ticket.get(no)
|
|
if sess:
|
|
merged = dict(sess.collab_facts)
|
|
merged.update(cleaned)
|
|
sess.collab_facts = merged
|
|
self._save(sess)
|
|
writer = getattr(self._ledger, "patch_collab_facts", None)
|
|
if callable(writer):
|
|
writer(work_order_no=no, facts=cleaned, sender_id=sender_id)
|
|
logger.info("sea.collab_fields wo=%s keys=%s by=%s", no, sorted(cleaned), sender_id)
|
|
return {"ok": True, "updated": sorted(cleaned)}
|
|
|
|
|
|
def _default_group_client() -> Any:
|
|
"""生产走企微应用建群;未配置时返回会失败的空客户端。"""
|
|
from agent.channel.outbox.sender import WeComAppClient
|
|
from agent.config import get_settings
|
|
|
|
settings = get_settings()
|
|
return WeComAppClient(
|
|
corp_id=settings.wecom_corp_id,
|
|
secret=settings.wecom_secret,
|
|
agent_id=int(settings.wecom_agent_id or "1000010"),
|
|
)
|
|
|
|
|
|
_FLOW: Optional[SeaTextInquiryFlow] = None
|
|
_FLOW_GUARD = threading.Lock()
|
|
|
|
|
|
def get_sea_text_flow() -> SeaTextInquiryFlow:
|
|
"""进程内一份海运流程;测试可 new 独立实例。"""
|
|
global _FLOW
|
|
with _FLOW_GUARD:
|
|
if _FLOW is None:
|
|
_FLOW = SeaTextInquiryFlow()
|
|
return _FLOW
|
|
|
|
|
|
def reset_sea_text_flow_for_test() -> SeaTextInquiryFlow:
|
|
"""单测重置全局海运实例(强制内存账本)。"""
|
|
global _FLOW
|
|
with _FLOW_GUARD:
|
|
_FLOW = SeaTextInquiryFlow(MemoryLedger())
|
|
return _FLOW
|