Files
inquiry_robot/inquiry-agent/agent/policy/multi_group_ops.py
T
jillion886andCursor 43e633f359 修好多段含空运的群跟单、锁舱话术、相册报价通道,以及附件上传成功后电脑端关不掉 H5。
多段激活后无需再带工单号;含空运群回复走机器人;复制补问清单时空货好时间不再吃成下一行编号。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-10-08 15:23:06 +08:00

612 lines
22 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
多段联运的群内报价与成交跟进。
本文件职责:按发言人岗位把报价归到对应段;一段报完就出该段报价单;
全部段都出了报价单,才在群里发一张整票成交跟进,各段价格不加总。
含空运时,原有群发出工单号后在这里激活:@人、发各段需求和模板。
调用:群入站 Handler,在 Worker 线程。禁止查 TMS 运价;禁止按段各建工单。
"""
from __future__ import annotations
import json
import logging
from contextlib import contextmanager
from typing import Any, Callable, Iterator, Optional
from agent.llm.extract_text import extract_collab_fields, extract_product_quote
from agent.policy import inquiry_copy as copy
from agent.policy.multi_segments import Segment, group_brief, includes_air, segments_from_payload
from agent.policy.sea_group_ops import (
apply_group_deal,
group_client_of,
is_sales,
override_group_client,
send_filled_sheet,
send_group_text,
ticket_no,
ticket_view,
)
from agent.schema.field_validate import collab_fields_for_ticket
logger = logging.getLogger(__name__)
_ROLE_MODE = {"sea": "SEA", "land": "LAND", "air": "AIR"}
SheetFn = Callable[..., dict[str, Any]]
@contextmanager
def bot_sender_when_air_multi(ticket: object) -> Iterator[None]:
"""
含空运的多段用的是机器人原来的群。应用接口发不进这个群,报价和模板都会没反应。
这一次出站改走机器人桥,发完还原。没有空运段的多段仍用应用自己建的群,这里不动。
只影响当前线程,不改流程上缓存的客户端。
"""
if not includes_air(load_segments(ticket)):
yield
return
from agent.channel.aibot.reply import AibotReplyClient
logger.info("multi.outbound 含空运走机器人桥 wo=%s", ticket_no(ticket))
with override_group_client(AibotReplyClient()):
yield
def ticket_is_multi(ticket: object) -> bool:
"""主账业务线是多段联运。"""
view = ticket_view(ticket)
line = str(view.get("business_line") or view.get("businessLine") or "").upper()
return line == "MULTI"
def load_segments(ticket: object) -> list[Segment]:
"""从工单事实里的 segments_json 读回各段。坏数据当没有段。"""
view = ticket_view(ticket)
facts = dict(view.get("facts") or {})
raw = facts.get("segments_json") or ""
if not raw:
return []
try:
payload = json.loads(raw)
except json.JSONDecodeError:
logger.warning("多段 segments_json 无法解析 wo=%s", ticket_no(ticket))
return []
return segments_from_payload(payload)
def segment_quotes(ticket: object) -> dict[str, dict[str, Any]]:
"""已落到各段的报价。键是 SEA/LAND/AIR。"""
quote = dict(ticket_view(ticket).get("quote") or {})
raw = quote.get("segment_quotes") or {}
if not isinstance(raw, dict):
return {}
out: dict[str, dict[str, Any]] = {}
for key, val in raw.items():
mode = str(key or "").upper()
if mode in _ROLE_MODE.values() and isinstance(val, dict):
out[mode] = dict(val)
return out
def sender_mode(flow: Any, sender_id: str) -> str:
"""
发言人负责哪一段。
海运产品归海运,陆运产品归陆运,航线·空运归空运。对不上返回空。
"""
book = getattr(flow, "ledger", None) or getattr(flow, "_ledger", None)
finder = getattr(book, "get_staff_by_wecom_id", None) if book is not None else None
if not callable(finder):
return ""
try:
row = finder(wecom_id=sender_id) or {}
except Exception:
logger.exception("多段查岗位失败 wecom=%s", sender_id)
return ""
role = str(row.get("roleCode") or row.get("role_code") or "").strip().lower()
return _ROLE_MODE.get(role, "")
def handle_multi_inbound(
*,
flow: Any,
ticket: object,
chat_id: str,
sender_id: str,
text: str,
injected_quote: Optional[dict[str, Any]] = None,
sheet_fn: Optional[SheetFn] = None,
) -> str:
"""
群里一条多段消息。
只含工单号:绑定本群并发摘要和各段模板。
有价格:归到发言人那一段并出该段报价单。全部出齐才成交跟进。
副作用:写主账报价、可能发群消息和模板。不查 TMS。
"""
room = (chat_id or "").strip()
no = ticket_no(ticket)
segments = load_segments(ticket)
if not segments:
send_group_text(group_client_of(flow), room, f"工单{no}没有可报价的运输段。")
return "multi_no_segments"
only_no = copy.extract_work_order_no(text or "") == no and _only_work_order(text or "")
if only_no or not (ticket_view(ticket).get("collab_chat_id") or ""):
_bind_room(flow, ticket, room, sender_id)
ticket = _reload(flow, no) or ticket
cabin = _cabin_intent(text or "")
if cabin and includes_air(segments):
return _cabin(flow, ticket, room, sender_id, text or "", cabin)
if only_no:
_send_activation(flow, ticket, room, segments, sender_id=sender_id, sheet_fn=sheet_fn)
return "multi_activated"
mode = sender_mode(flow, sender_id)
target = next((seg for seg in segments if seg.mode == mode), None)
fields = {}
if target is not None:
keys = collab_fields_for_ticket(target.facts, target.mode)
fields = extract_collab_fields(text or "", allowed_keys=keys)
if fields:
_save_collab(flow, ticket, target, fields)
ticket = _reload(flow, no) or ticket
segments = load_segments(ticket)
target = next((seg for seg in segments if seg.mode == mode), target)
quote = extract_product_quote(text or "", injected=injected_quote)
if not quote:
if fields:
return "multi_collab_saved"
if is_sales(ticket, sender_id):
return _maybe_deal(flow, ticket, sender_id, text)
send_group_text(group_client_of(flow), room, copy.group_quote_not_recognized())
return "quote_not_recognized"
if target is None:
send_group_text(group_client_of(flow), room, "没有对上你负责的运输段,请让对应岗位报价。")
return "quote_wrong_role"
return _save_segment_quote(
flow=flow,
ticket=ticket,
chat_id=room,
segment=target,
segments=segments,
quote=quote,
sheet_fn=sheet_fn,
)
def deal_text(*, work_order_no: str, segments: list[Segment], quotes: dict[str, dict[str, Any]]) -> str:
"""整票成交跟进文字。每段各自的价格、有效期,海运陆运再加时效。不加总。"""
no = (work_order_no or "").strip()
lines = [f"询价工单 {no}", copy.DEAL_SUBTITLE]
for seg in segments:
quote = quotes.get(seg.mode) or {}
lines.append(f"段{seg.index}:{seg.word()}")
rows = copy.deal_horizontal_rows(quote=quote, facts=seg.facts, business_line=seg.mode)
lines.extend(f"{row['keyname']}:{row['value']}" for row in rows)
lines.append(copy.DEAL_HINT)
return "\n".join(lines)
def deal_payload(
*,
work_order_no: str,
segments: list[Segment],
quotes: dict[str, dict[str, Any]],
card_version: int = 1,
) -> dict[str, Any]:
"""成交卡一行一段报价,不出现各段加总。按钮仍是整票一次确认。"""
base = copy.deal_wecom_payload(
work_order_no=work_order_no,
quote={"source_label": "分段报价"},
business_line="MULTI",
card_version=card_version,
)
rows = []
for seg in segments:
quote = quotes.get(seg.mode) or {}
total = copy.money_text(
str(quote.get("total") or "-").strip() or "-",
str(quote.get("currency") or ""),
)
rows.append({"keyname": f"段{seg.index}{seg.word()}", "value": copy._clip(total, 26)})
card = dict(base.get("template_card") or {})
card["horizontal_content_list"] = rows
card["source"] = {"desc": "分段报价", "desc_color": 0}
return {"msgtype": "template_card", "template_card": card}
def _save_segment_quote(
*,
flow: Any,
ticket: object,
chat_id: str,
segment: Segment,
segments: list[Segment],
quote: dict[str, Any],
sheet_fn: Optional[SheetFn],
) -> str:
"""写入这一段报价并发该段报价单。其他段未齐时不发成交跟进。"""
no = ticket_no(ticket)
stored = segment_quotes(ticket)
stored[segment.mode] = dict(quote)
version = int(dict(ticket_view(ticket).get("quote") or {}).get("deal_version") or 0)
all_ready = all(seg.mode in stored and _has_money(stored[seg.mode]) for seg in segments)
if all_ready:
version += 1
payload = {
"segment_quotes": stored,
"source_label": "分段报价",
"deal_version": version,
}
writer = getattr(flow.ledger, "upsert_quote", None)
if callable(writer):
writer(
work_order_no=no,
quote=payload,
to_status="已报价" if all_ready else "",
)
ticket = _reload(flow, no) or ticket
_send_sheet(
flow=flow,
ticket=ticket,
chat_id=chat_id,
segment=segment,
quote=quote,
sheet_fn=sheet_fn,
)
client = group_client_of(flow)
if not all_ready:
send_group_text(client, chat_id, f"段{segment.index}{segment.word()}报价单已发出。其他段报完后一起跟进成交。")
return "segment_quoted"
text = deal_text(work_order_no=no, segments=segments, quotes=stored)
send_group_text(client, chat_id, text)
carder = getattr(client, "send_group_card", None)
if callable(carder):
carder(
chat_id=chat_id,
template_card=deal_payload(
work_order_no=no,
segments=segments,
quotes=stored,
card_version=version,
).get("template_card")
or {},
)
logger.info("multi.deal wo=%s segments=%s", no, [seg.mode for seg in segments])
return "multi_deal_sent"
def _cabin(flow: Any, ticket: object, chat_id: str, sender_id: str, text: str, intent: str) -> str:
"""
多段含空运时的锁舱、释放。
锁舱暂不支持:只回复失败原因,不打 TMS,也不标系统异常。
释放仍走空运规则。调用方已确认本单含空运段。群消息在 Worker 线程发出。
"""
from agent.policy.air_group_ops import (
apply_air_release,
is_airline,
operator_name,
)
if intent == "release":
return apply_air_release(
flow=flow,
ticket=ticket,
sender_id=sender_id,
chat_id=chat_id,
instruction=text,
)
if not is_airline(flow, sender_id):
send_group_text(group_client_of(flow), chat_id, copy.air_airline_only_cabin())
return "not_airline"
no = ticket_no(ticket)
send_group_text(
group_client_of(flow),
chat_id,
copy.multi_air_lock_unsupported(
work_order_no=no,
operator=operator_name(flow, sender_id),
),
)
logger.info("multi.lock unsupported wo=%s sender=%s", no, sender_id)
return "multi_lock_unsupported"
def _cabin_intent(text: str) -> str:
raw = text or ""
if "释放" in raw and "舱" in raw:
return "release"
if "锁舱" in raw:
return "lock"
return ""
def _send_activation(
flow: Any,
ticket: object,
chat_id: str,
segments: list[Segment],
*,
sender_id: str = "",
sheet_fn: Optional[SheetFn],
) -> None:
"""原有群激活:@具体销售和各段产品,发核对同款字段,每段一份模板。不含 TMS。"""
no = ticket_no(ticket)
sales_name, products, mention_ids = _activation_mentions(
flow, ticket, segments, sender_id=sender_id
)
brief = group_brief(
work_order_no=no,
segments=segments,
sales_name=sales_name,
mention_names=products,
)
send_group_text(group_client_of(flow), chat_id, brief, mention_userids=mention_ids)
for seg in segments:
_send_sheet(flow=flow, ticket=ticket, chat_id=chat_id, segment=seg, quote={}, sheet_fn=sheet_fn)
def _activation_mentions(
flow: Any,
ticket: object,
segments: list[Segment],
sender_id: str = "",
) -> tuple[str, list[str], list[str]]:
"""激活时要 @ 的人:具体销售 + 各段对得上的产品岗;空运按航线岗名单。"""
from agent.policy.sea_group_ops import resolve_group_mentions
sales_name, products, mention_ids = resolve_group_mentions(
flow, ticket, None, fallback_sales_id=sender_id
)
book = getattr(flow, "ledger", None) or getattr(flow, "_ledger", None)
for seg in segments:
for row in _staff_rows(book, seg):
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
name = str(row.get("name") or "").strip() or uid
if uid and uid not in mention_ids:
mention_ids.append(uid)
if name and name not in products and name != sales_name:
products.append(name)
return sales_name, products, mention_ids
def _staff_rows(book: Any, seg: Segment) -> list[dict[str, Any]]:
"""一段要对上的产品人员。线路没填时,陆运/空运按岗位名单,不猜港口。"""
if book is None:
return []
try:
if seg.mode == "SEA":
matcher = getattr(book, "match_sea_staff", None)
if not callable(matcher):
return []
rows = matcher(
origin=str(seg.facts.get("起运港") or ""),
destination=str(seg.facts.get("目的港") or ""),
) or []
return [dict(x) for x in rows if isinstance(x, dict)]
if seg.mode == "LAND":
matcher = getattr(book, "match_land_staff", None)
rows = []
if callable(matcher):
rows = matcher(route_category=str(seg.facts.get("线路类别") or "")) or []
hits = [dict(x) for x in rows if isinstance(x, dict)]
if hits:
return hits
lister = getattr(book, "list_staff_by_role", None)
if not callable(lister):
return []
return [dict(x) for x in (lister(role_code="land") or []) if isinstance(x, dict)]
if seg.mode == "AIR":
lister = getattr(book, "list_staff_by_role", None)
if not callable(lister):
return []
return [dict(x) for x in (lister(role_code="air") or []) if isinstance(x, dict)]
except Exception:
logger.exception("多段激活匹配产品失败 mode=%s", seg.mode)
return []
def _send_sheet(
*,
flow: Any,
ticket: object,
chat_id: str,
segment: Segment,
quote: dict[str, Any],
sheet_fn: Optional[SheetFn],
) -> None:
"""按这一段的运输方式匹配模板。失败只记日志,不挡住报价归段。"""
sender = sheet_fn or send_filled_sheet
try:
sender(
flow=flow,
ticket=_sheet_ticket(ticket, segment),
quote_override=dict(quote or {}),
business_line=segment.mode,
chat_id=chat_id,
)
except Exception:
logger.exception("多段模板发送失败 wo=%s mode=%s", ticket_no(ticket), segment.mode)
def _sheet_ticket(ticket: object, segment: Segment) -> object:
"""填表时用这一段的询价字段,避免海运模板吃到空运字段。"""
view = ticket_view(ticket)
facts = _facts_for_segment_sheet(segment)
sales_id = str(view.get("sales_wecom_id") or view.get("salesWecomId") or "").strip()
if isinstance(ticket, dict):
cloned = dict(ticket)
cloned["facts"] = facts
cloned["business_line"] = segment.mode
cloned["businessLine"] = segment.mode
if sales_id:
cloned["sales_wecom_id"] = sales_id
cloned["salesWecomId"] = sales_id
return cloned
return {
"work_order_no": view.get("work_order_no"),
"workOrderNo": view.get("work_order_no"),
"facts": facts,
"business_line": segment.mode,
"businessLine": segment.mode,
"collab_chat_id": view.get("collab_chat_id"),
"sales_wecom_id": sales_id,
"salesWecomId": sales_id,
"quote": {},
}
def _facts_for_segment_sheet(segment: Segment) -> dict[str, str]:
"""
一段填表用的字段。
多段不查 TMS,陆运常缺线路类别。
后台陆运拼车模板关键词是「国内拼车」,询价字段是「国内运输拼车」,
不补别名就套不上表,只会发出空运模板。
"""
facts = {k: str(v).strip() for k, v in dict(segment.facts).items() if str(v or "").strip()}
if segment.mode != "LAND":
return facts
from agent.schema.land_options import SAMPLE_PAIRS, canonicalize_transport_type
land_type = canonicalize_transport_type(str(facts.get("运输类型") or ""))
if land_type:
facts["运输类型"] = land_type
if land_type and not str(facts.get("线路类别") or "").strip():
for pair_type, pair_route in SAMPLE_PAIRS:
if pair_type == land_type:
facts["线路类别"] = pair_route
break
hint = _LAND_SHEET_HINTS.get(land_type, "")
if hint:
facts["模板匹配"] = hint
return facts
_LAND_SHEET_HINTS = {
"国内运输拼车": "国内拼车 国内零担",
"国内运输整车": "国内整车",
"跨境整车": "跨境整车",
"跨境集拼": "跨境集拼 跨境零担",
"中港整车": "中港",
"中港零担/集拼": "中港",
}
def _save_collab(flow: Any, ticket: object, segment: Segment, fields: dict[str, str]) -> None:
"""协同字段写进对应段,不挡报价。"""
segments = load_segments(ticket)
for seg in segments:
if seg.mode != segment.mode:
continue
seg.facts.update({k: str(v) for k, v in fields.items() if str(v).strip()})
_write_segments(flow, ticket_no(ticket), segments)
writer = getattr(flow.ledger, "patch_collab_facts", None)
if callable(writer):
writer(work_order_no=ticket_no(ticket), facts=fields, sender_id="")
def _bind_room(flow: Any, ticket: object, chat_id: str, sender_id: str) -> None:
"""把原有群绑到这张多段工单。已绑同一个群则不重复建。"""
view = ticket_view(ticket)
if str(view.get("collab_chat_id") or "") == chat_id:
return
binder = getattr(flow.ledger, "bind_collab_group", None)
members = [str(x) for x in (view.get("collab_member_ids") or [])]
if sender_id and sender_id not in members:
members.append(sender_id)
sales = str(view.get("sales_wecom_id") or view.get("salesWecomId") or sender_id or "")
if sales and sales not in members:
members.append(sales)
if callable(binder):
binder(
work_order_no=ticket_no(ticket),
chat_id=chat_id,
member_ids=members,
product_names=list(view.get("product_names") or []),
product_ids=list(view.get("product_ids") or []),
)
def _maybe_deal(flow: Any, ticket: object, sender_id: str, text: str) -> str:
"""销售在群里确认成交。还有段没报价单时不确认。"""
from agent.routing.group_intent import (
INTENT_DEAL,
INTENT_LOST,
INTENT_NEGOTIATE,
classify_group_text,
)
intent = classify_group_text(text or "")
outcome = {INTENT_DEAL: "已成交", INTENT_LOST: "未成交", INTENT_NEGOTIATE: "协商中"}.get(intent)
if not outcome:
return "multi_noop"
segments = load_segments(ticket)
quotes = segment_quotes(ticket)
if not all(seg.mode in quotes and _has_money(quotes[seg.mode]) for seg in segments):
send_group_text(
group_client_of(flow),
str(ticket_view(ticket).get("collab_chat_id") or ""),
"还有运输段没有报价单,全部报完后再确认成交。",
)
return "deal_wait_segments"
return apply_group_deal(
flow=flow,
ticket=ticket,
sender_id=sender_id,
outcome=outcome,
text=text,
)
def _write_segments(flow: Any, work_order_no: str, segments: list[Segment]) -> None:
from agent.policy.multi_segments import segments_to_payload
writer = getattr(flow.ledger, "update_facts", None)
if not callable(writer):
return
current = {}
getter = getattr(flow.ledger, "get_ticket", None)
if callable(getter):
ticket = getter(work_order_no=work_order_no)
current = dict(ticket_view(ticket).get("facts") or {})
current["segments_json"] = json.dumps(segments_to_payload(segments), ensure_ascii=False)
writer(work_order_no=work_order_no, facts=current)
def _reload(flow: Any, work_order_no: str) -> object:
getter = getattr(flow.ledger, "get_ticket", None)
if callable(getter):
return getter(work_order_no=work_order_no)
return None
def _has_money(quote: dict[str, Any]) -> bool:
return bool(str(quote.get("total") or "").strip() or quote.get("fee_rows") or quote.get("fee_lines"))
def _only_work_order(text: str) -> bool:
"""去掉工单号和 @询价小助手 后几乎没字,当作激活。与空运群同一套,避免 WO 数字被当成报价。"""
from agent.routing.air_group_intent import _only_work_order as air_only
return air_only(text)
def send_segment_templates(*, flow: Any, ticket: object, chat_id: str, sheet_fn: Optional[SheetFn] = None) -> None:
"""拉群或激活后,按段各发一份匹配模板。没有段则不发。"""
for seg in load_segments(ticket):
_send_sheet(
flow=flow,
ticket=ticket,
chat_id=chat_id,
segment=seg,
quote={},
sheet_fn=sheet_fn,
)
def multi_has_air(ticket: object) -> bool:
"""这张多段单含空运段,锁舱和含电含磁才适用。"""
return includes_air(load_segments(ticket))