770 lines
30 KiB
Python
770 lines
30 KiB
Python
"""
|
|
多段联运私聊:补问、核对、建单、拉群。
|
|
|
|
本文件职责:一段一段收齐必填后才建一张工单;不去 TMS 查价。
|
|
不含空运时,销售点「拉产品进群」后按各段线路把人拉进同一个群,并发群摘要。
|
|
含空运时只提示手动拉人到原有群。报价和成交跟进不在本文件。
|
|
调用:文字询价 Handler 在 Worker 线程;禁止回调线程同步跑。
|
|
禁止:直连 TMS;用两段同一种运输方式建单;把协同字段当成必填。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import logging
|
|
import threading
|
|
from typing import Any, Optional
|
|
|
|
from agent.policy import inquiry_copy as copy
|
|
from agent.policy.air_text_flow import (
|
|
FlowSession,
|
|
ReplyFn,
|
|
_session_from_payload,
|
|
_session_payload,
|
|
build_ledger,
|
|
)
|
|
from agent.policy.multi_segments import (
|
|
Segment,
|
|
classify_multi_text,
|
|
confirm_text,
|
|
group_brief,
|
|
includes_air,
|
|
merge_supplement,
|
|
refresh_segments,
|
|
segments_from_payload,
|
|
segments_to_payload,
|
|
split_segments,
|
|
)
|
|
from agent.schema.land_options import (
|
|
looks_like_confirm,
|
|
sample_pair_allowed,
|
|
)
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# 只有补问/核对还能续同一张多段草稿。已出号的拉群停点不能拦住新询价。
|
|
_FRONT = frozenset({"clarify", "wait_confirm"})
|
|
|
|
|
|
class MultiTextInquiryFlow:
|
|
"""多段联运私聊 Owner。一张工单、一个群,段存在书签 quote.multi_segments。"""
|
|
|
|
bookmark_kind = "multi"
|
|
|
|
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
|
|
|
|
def session_of(self, sender_id: str) -> Optional[FlowSession]:
|
|
"""当前销售的多段书签。进程内没有就读 Redis 书签。"""
|
|
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_for_action(
|
|
self,
|
|
sender_id: str,
|
|
card_meta: Optional[dict[str, Any]] = None,
|
|
*,
|
|
text: str = "",
|
|
) -> Optional[FlowSession]:
|
|
"""
|
|
点卡:用工单号找回这一单。HTTP 进程往往没有 Worker 内存书签。
|
|
|
|
卡上有工单号时,不对上就不拿当前草稿,避免串到别人/最新单。
|
|
进程刚重启再从主账收回 MULTI 工单。
|
|
"""
|
|
sid = (sender_id or "").strip()
|
|
if not sid:
|
|
return None
|
|
no = copy.work_order_from_card_meta(card_meta, text=text)
|
|
if no:
|
|
with self._lock:
|
|
hit = self._by_ticket.get(no)
|
|
if hit and hit.sender_id == sid:
|
|
return hit
|
|
sess = self.session_of(sid)
|
|
if sess and (sess.business_line or "").upper() == "MULTI":
|
|
if not no:
|
|
return sess
|
|
if (sess.work_order_no or "").strip() == no:
|
|
return sess
|
|
if no:
|
|
return self._load_ticket_session(sid, no)
|
|
return None
|
|
|
|
def on_button(
|
|
self,
|
|
*,
|
|
sender_id: str,
|
|
action: str,
|
|
reply: ReplyFn,
|
|
card_meta: Optional[dict[str, Any]] = None,
|
|
) -> str:
|
|
"""
|
|
处理多段拉群卡。错工单/未允许直接丢弃,不改主账。
|
|
|
|
副作用:点「拉产品进群协同」时建一个群并发摘要、各段模板。
|
|
调用:卡片 Handler,HTTP 回调线程只入队后由 Worker/本流程执行;本函数会打企微建群。
|
|
"""
|
|
sess = self.session_for_action(sender_id, card_meta)
|
|
card_wo = copy.work_order_from_card_meta(card_meta)
|
|
if not sess or action not in (sess.allowed or ()):
|
|
logger.info(
|
|
"multi 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(
|
|
"multi card 工单不一致 card_wo=%s sess_wo=%s,拒绝串单",
|
|
card_wo,
|
|
sess.work_order_no,
|
|
)
|
|
return "card_discard"
|
|
if action == copy.BTN_PULL_COLLAB:
|
|
return self._pull_group(sess, reply)
|
|
return "card_discard"
|
|
|
|
def drop_current_bookmark(self, sender_id: str) -> Optional[FlowSession]:
|
|
"""销售改发单段时丢掉未完成的多段草稿。已出号的仍留在工单号索引。"""
|
|
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 match_button(self, text: str, sess: Optional[FlowSession]) -> str:
|
|
"""拉群按钮。没到可点阶段就不当按钮。"""
|
|
if sess is None:
|
|
return ""
|
|
action = copy.canonical_button(text)
|
|
if action == copy.BTN_PULL_COLLAB and action in (sess.allowed or ()):
|
|
return action
|
|
return ""
|
|
|
|
def offer_handoff(self, *, sender_id: str, text: str, reply: ReplyFn) -> str:
|
|
"""多段还在收字段时不转人工,请销售先把各段发齐。"""
|
|
_ = sender_id, text
|
|
reply("多段询价请先把各段需求发齐并回复确定。需要人来接管时,请在工单创建后再说。")
|
|
return "multi_handoff_later"
|
|
|
|
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: Any = None,
|
|
inbound_raw: Optional[dict[str, Any]] = None,
|
|
) -> str:
|
|
"""
|
|
处理一条多段私聊。
|
|
|
|
副作用:写书签、可能建主账工单、点拉群时建企微群。不查 TMS。
|
|
新开多段仍按原文切段;补问中的口语句走 injected_* / DeepSeek B,与单聊同一套抽取。
|
|
已出号等拉群时,再发一段陆运+海运视为新询价,不能回旧工单的手动拉群提示。
|
|
"""
|
|
_ = inbound_raw
|
|
sid = (sender_id or "").strip()
|
|
raw = text or ""
|
|
sess = self.session_of(sid)
|
|
if sess and self.match_button(raw, sess):
|
|
return self._pull_group(sess, reply)
|
|
verdict = classify_multi_text(raw)
|
|
if verdict == "same_mode":
|
|
reply(copy.MULTI_SAME_MODE)
|
|
return "reject_same_mode"
|
|
if verdict == "multi":
|
|
# 补问里再写段1/段2是补字段,不能当新开把已填冲掉。
|
|
if sess and (sess.phase or "") == "clarify":
|
|
return self._continue_clarify(
|
|
sess,
|
|
raw,
|
|
reply,
|
|
injected_facts=injected_facts,
|
|
injected_mode=injected_mode,
|
|
b_result=b_result,
|
|
allow_b=allow_b,
|
|
chat_fn=chat_fn,
|
|
)
|
|
return self._open(
|
|
sid,
|
|
split_segments(raw).segments,
|
|
reply,
|
|
immutable=raw,
|
|
allow_b=allow_b,
|
|
chat_fn=chat_fn,
|
|
)
|
|
if sess and sess.phase == "wait_confirm" and looks_like_confirm(raw):
|
|
return self._create_ticket(sess, reply)
|
|
if sess and sess.phase in {"wait_collab", "collab_group"}:
|
|
if (sess.collab_chat_id or "").strip():
|
|
reply(copy.sea_group_exists(work_order_no=sess.work_order_no))
|
|
return "collab_group"
|
|
reply(copy.multi_pull_text(work_order_no=sess.work_order_no))
|
|
return "wait_collab"
|
|
if sess and sess.phase == "wait_manual_group":
|
|
reply(copy.multi_air_manual_hint(work_order_no=sess.work_order_no))
|
|
return "wait_manual_group"
|
|
if sess and sess.phase in _FRONT and not looks_like_confirm(raw):
|
|
if sess.phase == "clarify":
|
|
return self._continue_clarify(
|
|
sess,
|
|
raw,
|
|
reply,
|
|
injected_facts=injected_facts,
|
|
injected_mode=injected_mode,
|
|
b_result=b_result,
|
|
allow_b=allow_b,
|
|
chat_fn=chat_fn,
|
|
)
|
|
if sess.phase == "wait_confirm":
|
|
# 核对卡写了「内容有误请重新发」:销售会只改正字段,按补问叠上去。
|
|
return self._continue_clarify(
|
|
sess,
|
|
raw,
|
|
reply,
|
|
injected_facts=injected_facts,
|
|
injected_mode=injected_mode,
|
|
b_result=b_result,
|
|
allow_b=allow_b,
|
|
chat_fn=chat_fn,
|
|
)
|
|
# 核对卡不再设时效。走到这里说明不是「确定」也不是补字段。
|
|
reply(copy.MULTI_SAME_MODE if verdict == "same_mode" else "请按段发多段需求,例如段1海运、段2空运。")
|
|
return "multi_need_text"
|
|
|
|
def _open(
|
|
self,
|
|
sender_id: str,
|
|
segments: list[Segment],
|
|
reply: ReplyFn,
|
|
*,
|
|
immutable: str,
|
|
allow_b: bool = False,
|
|
chat_fn: Any = None,
|
|
) -> str:
|
|
"""
|
|
新的多段原文:缺必填就一次说完,齐了出核对。不建单。
|
|
|
|
附件口播经常没有「起运港:」这种标签,要按段走和单聊一样的抽取,
|
|
不能只认标签行,否则会整段待补充。
|
|
"""
|
|
filled = self._enrich_segments(segments, allow_b=allow_b, chat_fn=chat_fn)
|
|
ready, updated, gaps = refresh_segments(filled)
|
|
if not ready:
|
|
self._reply_gaps(updated, gaps, reply)
|
|
self._save(self._draft(sender_id, updated, phase="clarify", immutable=immutable))
|
|
return "clarify"
|
|
reply(confirm_text(updated))
|
|
self._save(self._draft(sender_id, updated, phase="wait_confirm", immutable=immutable))
|
|
return "wait_confirm"
|
|
|
|
def _reply_gaps(self, segments: list[Segment], gaps: str, reply: ReplyFn) -> None:
|
|
"""
|
|
补问正文。陆运选项走复制清单,不再另发填写样例图。
|
|
|
|
副作用:出站文字。只在 Worker 里调。
|
|
"""
|
|
reply(gaps)
|
|
|
|
def _enrich_segments(
|
|
self,
|
|
segments: list[Segment],
|
|
*,
|
|
allow_b: bool,
|
|
chat_fn: Any,
|
|
) -> list[Segment]:
|
|
"""
|
|
各段正文再抽一遍字段。标签行已经在切段时收过,这里只补口语/附件里没写冒号的值。
|
|
|
|
没有正文就不调模型,避免补问合并后空跑。
|
|
"""
|
|
if not segments:
|
|
return segments
|
|
if not allow_b and chat_fn is None:
|
|
return segments
|
|
from agent.llm.extract_text import extract_inquiry_snapshot
|
|
|
|
out: list[Segment] = []
|
|
for seg in segments:
|
|
body = (seg.source or "").strip()
|
|
if not body:
|
|
out.append(seg)
|
|
continue
|
|
snap = extract_inquiry_snapshot(
|
|
body,
|
|
injected_mode=seg.mode,
|
|
allow_b=allow_b,
|
|
chat_fn=chat_fn,
|
|
)
|
|
extra = {
|
|
str(key): str(val).strip()
|
|
for key, val in dict(snap.get("facts") or {}).items()
|
|
if str(val or "").strip()
|
|
}
|
|
subtype = str(snap.get("land_subtype") or "").strip()
|
|
if seg.mode == "LAND" and subtype and not str(extra.get("运输类型") or "").strip():
|
|
extra["运输类型"] = subtype
|
|
facts = dict(seg.facts)
|
|
for key, val in extra.items():
|
|
if not str(facts.get(key) or "").strip():
|
|
facts[key] = val
|
|
out.append(
|
|
Segment(index=seg.index, mode=seg.mode, facts=facts, source=seg.source)
|
|
)
|
|
return out
|
|
|
|
def _continue_clarify(
|
|
self,
|
|
sess: FlowSession,
|
|
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: Any = None,
|
|
) -> str:
|
|
"""
|
|
补问中的下一条:能整段重发就替换,否则按段号或唯一缺口段补上。
|
|
|
|
口语补句走单聊同一套抽取,不在这里猜「广州-深圳」。
|
|
"""
|
|
from agent.llm.extract_text import extract_inquiry_snapshot
|
|
from agent.policy.system_exception import handle_extract_tech_fail
|
|
|
|
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,
|
|
)
|
|
llm_phase = handle_extract_tech_fail(
|
|
snap=snap,
|
|
ledger=self._ledger,
|
|
reply=reply,
|
|
sender_id=sess.sender_id,
|
|
sess=sess,
|
|
)
|
|
if llm_phase:
|
|
return llm_phase
|
|
extra = {
|
|
str(key): str(val).strip()
|
|
for key, val in dict(snap.get("facts") or {}).items()
|
|
if str(val or "").strip()
|
|
}
|
|
merged = merge_supplement(self._segments(sess), text, extra_facts=extra)
|
|
return self._open(
|
|
sess.sender_id,
|
|
merged,
|
|
reply,
|
|
immutable=sess.immutable_text or text,
|
|
allow_b=allow_b,
|
|
chat_fn=chat_fn,
|
|
)
|
|
|
|
def _restart(self, sender_id: str, text: str, reply: ReplyFn) -> str:
|
|
"""核对之后没有回复确定,而是重发内容:整段重来,不在旧核对上改一个字。"""
|
|
verdict = classify_multi_text(text)
|
|
if verdict == "same_mode":
|
|
reply(copy.MULTI_SAME_MODE)
|
|
return "reject_same_mode"
|
|
if verdict != "multi":
|
|
reply("请把多段需求再发一次。")
|
|
return "wait_confirm"
|
|
return self._open(sender_id, split_segments(text).segments, reply, immutable=text)
|
|
|
|
def _create_ticket(self, sess: FlowSession, reply: ReplyFn) -> str:
|
|
"""回复确定后建一张 MULTI 工单。含空运只给手动拉群提示,不含空运出拉群卡。"""
|
|
segments = self._segments(sess)
|
|
ready, updated, gaps = refresh_segments(segments)
|
|
if not ready:
|
|
self._reply_gaps(updated, gaps, reply)
|
|
sess.phase = "clarify"
|
|
self._put_segments(sess, updated)
|
|
self._save(sess)
|
|
return "clarify"
|
|
payload = segments_to_payload(updated)
|
|
created = self._ledger.create_ticket(
|
|
sender_id=sess.sender_id,
|
|
business_line="MULTI",
|
|
facts={"segments_json": json.dumps(payload, ensure_ascii=False)},
|
|
idempotency_key=f"multi:{sess.sender_id}:{sess.wait_version}",
|
|
)
|
|
if not created.get("ok"):
|
|
reply("工单创建失败,请稍后再回复确定。")
|
|
return "wait_confirm"
|
|
no = str(created.get("work_order_no") or "").strip()
|
|
sess.work_order_no = no
|
|
sess.status = "询价中"
|
|
self._put_segments(sess, updated)
|
|
if includes_air(updated):
|
|
sess.phase = "wait_manual_group"
|
|
sess.allowed = ()
|
|
self._save(sess)
|
|
reply(copy.multi_air_manual_hint(work_order_no=no))
|
|
logger.info("multi.created air wo=%s segments=%s", no, len(updated))
|
|
return "wait_manual_group"
|
|
sess.phase = "wait_collab"
|
|
sess.allowed = (copy.BTN_PULL_COLLAB,)
|
|
self._save(sess)
|
|
reply(copy.multi_pull_text(work_order_no=no), copy.multi_pull_wecom_payload(work_order_no=no))
|
|
logger.info("multi.created pull wo=%s segments=%s", no, len(updated))
|
|
return "wait_collab"
|
|
|
|
def _pull_group(self, sess: FlowSession, reply: ReplyFn) -> str:
|
|
"""
|
|
点拉群:海运段按港口拉海运产品,陆运段按线路拉陆运产品,进同一个群。
|
|
|
|
任一段没有人,或陆运类型+线路不在填写样例里,不建群,按钮仍可点。
|
|
"""
|
|
no = (sess.work_order_no or "").strip()
|
|
if (sess.collab_chat_id or "").strip():
|
|
reply(copy.sea_group_exists(work_order_no=no))
|
|
return "collab_group"
|
|
segments = self._segments(sess)
|
|
if includes_air(segments):
|
|
reply(copy.multi_air_manual_hint(work_order_no=no))
|
|
return "wait_manual_group"
|
|
for seg in segments:
|
|
if seg.mode != "LAND":
|
|
continue
|
|
land_type = str(seg.facts.get("运输类型") or "").strip()
|
|
route = str(seg.facts.get("线路类别") or "").strip()
|
|
if not sample_pair_allowed(land_type, route):
|
|
reply(copy.land_group_pair_rejected())
|
|
return sess.phase or "wait_collab"
|
|
userids = [sess.sender_id]
|
|
product_names: list[str] = []
|
|
product_ids: list[str] = []
|
|
for seg in segments:
|
|
staff = self._staff_for(seg)
|
|
if not staff:
|
|
reply(copy.land_group_no_staff() if seg.mode == "LAND" else copy.sea_group_no_staff())
|
|
return sess.phase or "wait_collab"
|
|
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(copy.sea_group_no_staff())
|
|
return sess.phase or "wait_collab"
|
|
for uid in self._archive_seats():
|
|
if uid not in userids:
|
|
userids.append(uid)
|
|
client = self._client()
|
|
from agent.policy.system_exception import call_with_retries
|
|
|
|
created = call_with_retries(
|
|
lambda: client.create_group(name=no, userids=userids) or {},
|
|
is_ok=lambda x: bool((x or {}).get("ok") and (x or {}).get("chat_id")),
|
|
)
|
|
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()
|
|
sales_name = self._sales_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,
|
|
)
|
|
sess.collab_chat_id = chat_id
|
|
sess.phase = "collab_group"
|
|
sess.allowed = ()
|
|
sess.invited_product_names = tuple(product_names)
|
|
sess.invited_sales_name = sales_name
|
|
self._save(sess)
|
|
from agent.policy.sea_group_ops import send_group_text
|
|
|
|
brief = group_brief(
|
|
work_order_no=no,
|
|
segments=segments,
|
|
sales_name=sales_name,
|
|
mention_names=product_names,
|
|
)
|
|
send_group_text(client, chat_id, brief, mention_userids=product_ids)
|
|
ticket = self._ledger.get_ticket(work_order_no=no) if hasattr(self._ledger, "get_ticket") else None
|
|
if ticket is not None:
|
|
from agent.policy.multi_group_ops import send_segment_templates
|
|
|
|
send_segment_templates(flow=self, ticket=ticket, chat_id=chat_id)
|
|
reply(
|
|
copy.sea_group_created(
|
|
work_order_no=no,
|
|
sales_name=sales_name,
|
|
product_names=product_names,
|
|
)
|
|
)
|
|
logger.info("multi.group wo=%s chat=%s products=%s", no, chat_id, product_names)
|
|
return "collab_group"
|
|
|
|
def _staff_for(self, seg: Segment) -> list[dict[str, Any]]:
|
|
"""按这一段自己的线路找产品。空运段不在这里拉人。"""
|
|
if seg.mode == "SEA":
|
|
matcher = getattr(self._ledger, "match_sea_staff", None)
|
|
if not callable(matcher):
|
|
return []
|
|
try:
|
|
rows = matcher(
|
|
origin=str(seg.facts.get("起运港") or ""),
|
|
destination=str(seg.facts.get("目的港") or ""),
|
|
) or []
|
|
except Exception:
|
|
logger.exception("多段匹配海运产品失败")
|
|
return []
|
|
return [dict(x) for x in rows if isinstance(x, dict)]
|
|
if seg.mode == "LAND":
|
|
matcher = getattr(self._ledger, "match_land_staff", None)
|
|
if not callable(matcher):
|
|
return []
|
|
try:
|
|
rows = matcher(route_category=str(seg.facts.get("线路类别") or "")) or []
|
|
except Exception:
|
|
logger.exception("多段匹配陆运产品失败")
|
|
return []
|
|
return [dict(x) for x in rows if isinstance(x, dict)]
|
|
return []
|
|
|
|
def _client(self) -> Any:
|
|
"""建群客户端。单测注入;生产与海运共用企微应用客户端。"""
|
|
if self._group_client is not None:
|
|
return self._group_client
|
|
from agent.policy.sea_text_flow import _default_group_client
|
|
|
|
self._group_client = _default_group_client()
|
|
return self._group_client
|
|
|
|
def _archive_seats(self) -> list[str]:
|
|
"""会话存档账号。查本账本,没有也不挡建群。"""
|
|
from agent.config import get_settings
|
|
from agent.policy.sea_text_flow import ARCHIVE_SEAT_NAMES
|
|
|
|
ids: list[str] = []
|
|
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 not callable(finder):
|
|
return ids
|
|
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)
|
|
return ids
|
|
|
|
def _sales_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 {}
|
|
except Exception:
|
|
row = {}
|
|
name = str(row.get("name") or "").strip()
|
|
if name:
|
|
return name
|
|
return (sender_id or "").strip() or "销售"
|
|
|
|
def _draft(self, sender_id: str, segments: list[Segment], *, phase: str, immutable: str) -> FlowSession:
|
|
sess = FlowSession(
|
|
sender_id=sender_id,
|
|
thread_id=sender_id,
|
|
phase=phase,
|
|
business_line="MULTI",
|
|
immutable_text=immutable,
|
|
wait_version=1,
|
|
)
|
|
self._put_segments(sess, segments)
|
|
return sess
|
|
|
|
def _segments(self, sess: FlowSession) -> list[Segment]:
|
|
return segments_from_payload((sess.quote or {}).get("multi_segments"))
|
|
|
|
def _put_segments(self, sess: FlowSession, segments: list[Segment]) -> None:
|
|
quote = dict(sess.quote or {})
|
|
quote["multi_segments"] = segments_to_payload(segments)
|
|
sess.quote = quote
|
|
|
|
def _load_ticket_session(self, sender_id: str, work_order_no: str) -> Optional[FlowSession]:
|
|
"""
|
|
点卡时本进程没有书签:只读主账收回这一张 MULTI 工单。
|
|
|
|
找不到、不是本人、不是多段:返回 None,绝不改成当前最新单或空运停点。
|
|
"""
|
|
getter = getattr(self._ledger, "get_for_agent", None)
|
|
if not callable(getter):
|
|
logger.warning("multi 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(
|
|
"multi 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":
|
|
return None
|
|
sess = self._session_from_ledger(sender_id=sender_id, view=view)
|
|
self._save(sess)
|
|
logger.info("multi card 从主账收回 wo=%s phase=%s", sess.work_order_no, sess.phase)
|
|
return sess
|
|
|
|
def _session_from_ledger(self, *, sender_id: str, view: dict[str, Any]) -> FlowSession:
|
|
"""
|
|
主账快照收成多段书签。询价中/被空运误标成 tms_miss 都按等拉群处理。
|
|
|
|
段列表优先 quote.multi_segments,没有就解 facts.segments_json。
|
|
"""
|
|
no = str(view.get("work_order_no") or "").strip()
|
|
facts = dict(view.get("facts") or {})
|
|
quote = dict(view.get("quote") or {})
|
|
raw_segs = quote.get("multi_segments")
|
|
if not raw_segs:
|
|
blob = facts.get("segments_json") or facts.get("segmentsJson") or ""
|
|
if isinstance(blob, str) and blob.strip():
|
|
try:
|
|
parsed = json.loads(blob)
|
|
except json.JSONDecodeError:
|
|
parsed = None
|
|
if isinstance(parsed, list):
|
|
quote["multi_segments"] = parsed
|
|
raw_segs = parsed
|
|
segments = segments_from_payload(raw_segs)
|
|
status = str(view.get("status") or "").strip()
|
|
chat = str(view.get("collab_chat_id") or view.get("collabChatId") or "").strip()
|
|
phase = str(view.get("wait_phase") or view.get("waitPhase") or "").strip()
|
|
allowed: tuple[str, ...] = ()
|
|
if chat:
|
|
phase = "collab_group"
|
|
elif includes_air(segments):
|
|
phase = "wait_manual_group"
|
|
elif status in {"已成交", "已关闭"}:
|
|
phase = "done"
|
|
elif status == "未成交":
|
|
phase = "wait_lost_reason"
|
|
elif status == "协商中" or phase == "wait_deal":
|
|
phase = "wait_deal"
|
|
allowed = (copy.BTN_DEAL, copy.BTN_LOST, copy.BTN_NEGOTIATE)
|
|
else:
|
|
# 询价中、tms_miss、wait_collab:销售还没拉群
|
|
phase = "wait_collab"
|
|
allowed = (copy.BTN_PULL_COLLAB,)
|
|
return FlowSession(
|
|
sender_id=sender_id,
|
|
thread_id=sender_id,
|
|
phase=phase,
|
|
business_line="MULTI",
|
|
facts={str(k): str(v) for k, v in facts.items()},
|
|
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=chat,
|
|
)
|
|
|
|
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
|
|
self._persist_bookmark(sess)
|
|
|
|
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)
|
|
|
|
|
|
_FLOW: Optional[MultiTextInquiryFlow] = None
|
|
_FLOW_GUARD = threading.Lock()
|
|
|
|
|
|
def get_multi_text_flow() -> MultiTextInquiryFlow:
|
|
"""进程内一份多段流程。单测请自己 new,不要碰这份单例。"""
|
|
global _FLOW
|
|
with _FLOW_GUARD:
|
|
if _FLOW is None:
|
|
_FLOW = MultiTextInquiryFlow()
|
|
return _FLOW
|
|
|
|
|
|
def reset_multi_text_flow_for_test() -> MultiTextInquiryFlow:
|
|
"""单测重置全局实例(强制内存账本)。"""
|
|
global _FLOW
|
|
from agent.ledger.memory_ledger import MemoryLedger
|
|
|
|
with _FLOW_GUARD:
|
|
_FLOW = MultiTextInquiryFlow(MemoryLedger())
|
|
return _FLOW
|