Files
inquiry_robot/inquiry-agent/agent/ledger/memory_ledger.py
T

994 lines
41 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.
"""
主账内存假账本(单测 / 无 8180 时)。
本文件职责:实现建单(每次新开)、相似检索、TMS fixture、报价写入、六态流转,语义对齐 Java 主账。
禁止:冒充已打真实 TMS;禁止被生产 Worker 当默认 Owner(生产必须 HTTP 主账)。
"""
from __future__ import annotations
import logging
import threading
from dataclasses import dataclass, field
from datetime import datetime
from typing import Any, Optional
from zoneinfo import ZoneInfo
from agent.schema.fingerprint import inquiry_fingerprint
logger = logging.getLogger(__name__)
_SHANGHAI = ZoneInfo("Asia/Shanghai")
def _memory_quote_can_lock(quote: dict[str, Any]) -> bool:
"""单测锁舱:有 TMS 编号,或有人工总价/费用,才算有报价。不编造编号。"""
src = dict(quote or {})
label = str(src.get("source_label") or src.get("sourceLabel") or "")
if "确认" in label and "TMS" in label.upper():
return bool(str(src.get("quoteId") or "").strip()) and not str(src.get("quoteId") or "").upper().startswith("QT-")
if str(src.get("total") or "").strip() not in {"", "-"}:
return True
if src.get("fee_rows") or src.get("fee_lines") or src.get("feeItems") or src.get("departureCharges"):
return True
quote_id = str(src.get("quoteId") or "").strip()
return bool(quote_id) and not quote_id.upper().startswith("QT-WO")
@dataclass
class MemoryTicket:
"""内存工单。"""
work_order_no: str
sales_wecom_id: str
business_line: str
facts: dict[str, str]
fingerprint: str
status: str = "询价中"
quote: dict[str, Any] = field(default_factory=dict)
quote_version: int = 0
collab_chat_id: str = ""
collab_facts: dict[str, str] = field(default_factory=dict)
collab_member_ids: list[str] = field(default_factory=list)
product_ids: list[str] = field(default_factory=list)
product_names: list[str] = field(default_factory=list)
collab_ended: bool = False
attachments: list[dict[str, Any]] = field(default_factory=list)
quote_file_ready: bool = False
lost_reason: str = ""
next_followup_at: str = ""
deal_notify_summary: str = ""
event_exception: str = ""
event_wait_kind: str = ""
event_wait_started_at: str = ""
event_wait_alerted: int = 0
system_exception: str = ""
system_exception_reason: str = ""
system_exception_step: str = ""
system_exception_payload: dict[str, Any] = field(default_factory=dict)
class MemoryLedger:
"""
线程安全内存主账。
TMS fixture:facts.tms_force=1002 或目的港含「无价」→ 无价;否则有价。
"""
def __init__(self) -> None:
self._lock = threading.Lock()
self._tickets: dict[str, MemoryTicket] = {}
self._seq = 0
# 单测默认同业务有一张默认模板,避免「生成 Excel」误走失败话术
self._templates: list[dict[str, Any]] = [
{
"templateId": "mem-default-air",
"name": "空运默认",
"bizType": "空运",
"keywords": [],
"isDefault": True,
},
{
"templateId": "mem-default-land",
"name": "陆运默认",
"bizType": "陆运",
"keywords": [],
"isDefault": True,
},
{
"templateId": "mem-default-sea",
"name": "海运默认",
"bizType": "海运",
"keywords": [],
"isDefault": True,
},
]
self.last_render: dict[str, Any] = {}
self._sea_staff: list[dict[str, Any]] = []
self._land_staff: list[dict[str, Any]] = []
self._staff_by_wecom: dict[str, dict[str, Any]] = {}
self._staff_by_name: dict[str, dict[str, Any]] = {}
self._exception_notify_users: list[dict[str, Any]] = []
self._timeouts: dict[str, float] = {}
def _next_no(self) -> str:
self._seq += 1
day = datetime.now(_SHANGHAI).strftime("%Y%m%d")
return f"WO{day}{self._seq:04d}"
def find_similar_open(
self,
*,
facts: dict[str, str],
business_line: str,
exclude_work_order_no: str = "",
) -> dict[str, Any]:
"""
全库相同需求检索:需求指纹一模一样才算,不限六态。
多条只回最新一条;exclude_work_order_no 排除刚建的本单。
只供卡片展示,不复用当前工单号。
"""
fp = inquiry_fingerprint(facts, business_line=business_line)
skip = (exclude_work_order_no or "").strip()
latest: MemoryTicket | None = None
with self._lock:
for t in self._tickets.values():
if t.fingerprint != fp:
continue
if skip and t.work_order_no == skip:
continue
latest = t
if latest is None:
return {"found": False, "first_or_same": "first", "work_order_no": "", "status": ""}
return {
"found": True,
"first_or_same": "same",
"work_order_no": latest.work_order_no,
"status": latest.status,
}
def create_ticket(
self,
*,
sender_id: str,
business_line: str,
facts: dict[str, str],
idempotency_key: str = "",
) -> dict[str, Any]:
_ = idempotency_key
# 出询价确认必须新开一单,相同需求不复用未关闭工单
fp = inquiry_fingerprint(facts, business_line=business_line)
with self._lock:
no = self._next_no()
ticket = MemoryTicket(
work_order_no=no,
sales_wecom_id=sender_id,
business_line=business_line,
facts=dict(facts),
fingerprint=fp,
)
self._tickets[no] = ticket
similar = self.find_similar_open(
facts=facts, business_line=business_line, exclude_work_order_no=no
)
logger.info(
"memory_ledger.create %s first_or_same=%s similar=%s",
no,
similar.get("first_or_same"),
similar.get("work_order_no"),
)
return {
"ok": True,
"created": True,
"first_or_same": similar.get("first_or_same") or "first",
"found": bool(similar.get("found")),
"similar_work_order_no": similar.get("work_order_no") or "",
"work_order_no": no,
"status": "询价中",
}
def query_tms(
self,
*,
work_order_no: str,
facts: dict[str, str],
business_line: str = "",
) -> dict[str, Any]:
_ = business_line
dest = str(facts.get("目的港") or "")
force = str(facts.get("tms_force") or "")
if force == "1002" or force == "1003" or "无价" in dest or "歧义" in dest:
return {
"ok": True,
"classification": "TMS_NO_QUOTE_1002",
"has_price": False,
"quote": {},
}
if force == "tech":
return {"ok": False, "classification": "TMS_TECH_FAILURE", "has_price": False, "quote": {}}
quote = {
"freight": "USD 1,200",
"other_fee": "USD 80",
"total": "USD 1,280",
"eta": "3-5 天",
"valid_until": "2026-09-20",
"source": "TMS",
"source_label": "TMS 标准报价",
"fee_lines": [
{"name": "空运费", "amount": "USD 1,200"},
{"name": "其他附加费", "amount": "USD 80"},
],
}
if str(facts.get("tms_land_multi") or "") == "1":
quote["isMulti"] = "2"
quote["business_line"] = "LAND"
quote["landSegments"] = [
{"segmentSequence": 1, "routeLegSequence": 1, "origin": "上海", "destination": "谅山", "costPrice": 7500, "currency": "CNY", "transitDays": 48},
{"segmentSequence": 1, "routeLegSequence": 2, "origin": "谅山", "destination": "胡志明", "costPrice": 9600, "currency": "CNY", "transitDays": 72},
{"segmentSequence": 2, "routeLegSequence": 1, "origin": "上海", "destination": "深圳", "costPrice": 7500, "currency": "CNY", "transitDays": 48},
{"segmentSequence": 2, "routeLegSequence": 2, "origin": "深圳", "destination": "胡志明", "costPrice": 9600, "currency": "CNY", "transitDays": 72},
{"segmentSequence": 3, "routeLegSequence": 1, "origin": "上海", "destination": "广西", "costPrice": 7500, "currency": "CNY", "transitDays": 48},
{"segmentSequence": 3, "routeLegSequence": 2, "origin": "广西", "destination": "胡志明", "costPrice": 9600, "currency": "CNY", "transitDays": 72},
]
quote["fee_lines"] = [
{"name": "运输费", "feeName": "运输费", "amount": 7500, "currency": "CNY", "segmentSequence": 1, "routeLegSequence": 1},
{"name": "运输费", "feeName": "运输费", "amount": 9600, "currency": "CNY", "segmentSequence": 1, "routeLegSequence": 2},
{"name": "运输费", "feeName": "运输费", "amount": 7500, "currency": "CNY", "segmentSequence": 2, "routeLegSequence": 1},
{"name": "运输费", "feeName": "运输费", "amount": 9600, "currency": "CNY", "segmentSequence": 2, "routeLegSequence": 2},
{"name": "运输费", "feeName": "运输费", "amount": 7500, "currency": "CNY", "segmentSequence": 3, "routeLegSequence": 1},
{"name": "运输费", "feeName": "运输费", "amount": 9600, "currency": "CNY", "segmentSequence": 3, "routeLegSequence": 2},
]
if str(facts.get("tms_land_single") or "") == "1":
quote["isMulti"] = "1"
quote["business_line"] = "LAND"
quote["landSegments"] = [
{"segmentSequence": 1, "routeLegSequence": 1, "origin": "上海", "destination": "谅山", "costPrice": 7500, "currency": "CNY", "transitDays": 48},
{"segmentSequence": 1, "routeLegSequence": 2, "origin": "谅山", "destination": "胡志明", "costPrice": 9600, "currency": "CNY", "transitDays": 72},
]
quote["fee_lines"] = [
{"name": "运输费", "feeName": "运输费", "amount": 7500, "currency": "CNY", "segmentSequence": 1, "routeLegSequence": 1},
{"name": "运输费", "feeName": "运输费", "amount": 9600, "currency": "CNY", "segmentSequence": 1, "routeLegSequence": 2},
]
if str(facts.get("tms_multi") or "") == "1":
quote["isMulti"] = "2"
quote["airOptions"] = [
{
"originAirportCode": "ZUH",
"transferAirportCode": "COC",
"destinationAirportCode": "CRK",
"salePrice": 183370,
"currency": "CNY",
"departureCharges": [
{"chargeName": "运输费", "currency": "CNY", "unitPrice": 11}
],
},
{
"originAirportCode": "SZX",
"transferAirportCode": "",
"destinationAirportCode": "CRK",
"salePrice": 193370,
"currency": "CNY",
"departureCharges": [
{"chargeName": "运输费", "currency": "CNY", "unitPrice": 12}
],
},
]
return {
"ok": True,
"classification": "HAS_PRICE",
"has_price": True,
"quote": quote,
}
def upsert_quote(
self,
*,
work_order_no: str,
quote: dict[str, Any],
to_status: str = "",
) -> dict[str, Any]:
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
ticket.quote_version += 1
ticket.quote = dict(quote)
if to_status:
ticket.status = to_status
return {
"ok": True,
"quote_version": ticket.quote_version,
"status": ticket.status,
}
def update_facts(self, *, work_order_no: str, facts: dict[str, str]) -> dict[str, Any]:
"""选定线路后写回起运港等询价事实。"""
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
ticket.facts = dict(facts)
return {"ok": True}
def void_quote_files(self, *, work_order_no: str) -> dict[str, Any]:
"""采用后改选:旧报价单不再当本单有效文件。"""
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
ticket.quote_file_ready = False
if ticket.quote:
ticket.quote = dict(ticket.quote)
ticket.quote["quoteFileIssued"] = False
ticket.quote["quoteFileVoided"] = True
return {"ok": True}
def transition(
self,
*,
work_order_no: str,
to_status: str,
remark: str = "",
lost_reason: str = "",
next_followup_at: str = "",
deal_notify_summary: str = "",
) -> dict[str, Any]:
_ = remark
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
ticket.status = to_status
if to_status == "未成交":
ticket.next_followup_at = ""
if lost_reason:
ticket.lost_reason = lost_reason
elif to_status == "已成交":
ticket.next_followup_at = ""
if deal_notify_summary:
ticket.deal_notify_summary = deal_notify_summary
elif to_status == "协商中":
ticket.next_followup_at = next_followup_at
return {"ok": True, "status": ticket.status}
def patch_lost_reason(self, *, work_order_no: str, lost_reason: str) -> dict[str, Any]:
reason = (lost_reason or "").strip()
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
if ticket.status != "未成交":
return {"ok": False, "error": "not_lost", "status": ticket.status}
if not reason:
return {"ok": False, "error": "reason_blank"}
ticket.lost_reason = reason
return {"ok": True, "work_order_no": work_order_no, "lost_reason": reason}
def claim_due_negotiate(self) -> list[dict[str, Any]]:
now = datetime.now(_SHANGHAI).strftime("%Y-%m-%d %H:%M:%S")
hits: list[dict[str, Any]] = []
with self._lock:
for ticket in self._tickets.values():
due = (ticket.next_followup_at or "").strip()
if ticket.status != "协商中" or not due or due > now:
continue
ticket.next_followup_at = ""
hits.append(
{
"work_order_no": ticket.work_order_no,
"sales_wecom_id": ticket.sales_wecom_id,
"sales_name": "",
"business_line": ticket.business_line,
"collab_chat_id": ticket.collab_chat_id,
"status": ticket.status,
"facts": dict(ticket.facts),
"quote": dict(ticket.quote),
}
)
return hits
def get_ticket(self, *, work_order_no: str) -> MemoryTicket | None:
with self._lock:
return self._tickets.get(work_order_no)
def set_quote_templates(self, templates: list[dict[str, Any]]) -> None:
"""单测注入后台模板列表;空列表=未维护关键词也未维护默认。"""
with self._lock:
self._templates = list(templates)
def render_quote(
self,
*,
work_order_no: str,
business_line: str,
facts: dict[str, Any],
quote: Optional[dict[str, Any]] = None,
format: str = "xlsx",
sales_wecom_id: str = "",
) -> dict[str, Any]:
"""
按询价需求匹配内存模板。不填真实 xlsx,只回假字节供流程/单测。
"""
from agent.quote_templates.match import FAIL_MESSAGE, pick_quote_template
_ = sales_wecom_id
with self._lock:
templates = list(self._templates)
picked = pick_quote_template(
templates=templates,
business_line=business_line,
facts=facts,
quote=quote,
)
if picked is None:
out = {"ok": False, "error": "NO_TEMPLATE", "message": FAIL_MESSAGE, "facts": dict(facts or {})}
self.last_render = out
return out
from agent.quote_templates.filename import quote_sheet_filename
suffix = ".pdf" if str(format).lower() == "pdf" else ".xlsx"
out = {
"ok": True,
"matchType": picked.match_type,
"templateId": picked.template_id,
"templateName": picked.name,
"fileName": quote_sheet_filename(work_order_no, picked.name, suffix=suffix)
or f"{work_order_no}_quote{suffix}",
"fileBase64": "UEsDBAoAAAAAA",
"needsPdf": False,
"facts": dict(facts or {}),
}
with self._lock:
ticket = self._tickets.get((work_order_no or "").strip())
if ticket is not None:
ticket.quote_file_ready = True
self.last_render = out
return out
def get_for_agent(self, *, work_order_no: str, sender_id: str) -> dict[str, Any]:
"""
智能体续办读取:工单是否存在、是否本人、字段/报价/推断停点。
"""
no = (work_order_no or "").strip()
ticket = self.get_ticket(work_order_no=no)
if ticket is None:
return {"ok": False, "found": False, "owned": False, "work_order_no": no}
owned = ticket.sales_wecom_id == sender_id
status = ticket.status
quote = dict(ticket.quote or {})
wait_phase = "done"
if status in {"已成交", "已关闭"}:
wait_phase = "done"
elif status == "未成交":
wait_phase = "wait_lost_reason" if not (ticket.lost_reason or "").strip() else "done"
elif status == "协商中":
wait_phase = "wait_deal"
elif status == "已报价" and quote:
# 文件已发给销售 → 成交跟进;否则海运协同卡、空运方案卡
if quote.get("quoteFileVoided"):
wait_phase = "wait_adopt"
elif ticket.quote_file_ready or quote.get("quoteFileIssued"):
wait_phase = "wait_deal"
elif (ticket.business_line or "").upper() == "SEA":
wait_phase = "wait_collab"
else:
wait_phase = "wait_adopt"
elif (ticket.business_line or "").upper() == "MULTI":
# 多段建单后仍是询价中;已建群就跟进报价,否则等销售点拉群。
wait_phase = "collab_group" if (ticket.collab_chat_id or "").strip() else "wait_collab"
else:
wait_phase = "tms_miss"
return {
"ok": True,
"found": True,
"owned": owned,
"work_order_no": ticket.work_order_no,
"status": status,
"business_line": ticket.business_line,
"sales_wecom_id": ticket.sales_wecom_id,
"salesWecomId": ticket.sales_wecom_id,
"sales_name": str(self._staff_by_wecom.get(ticket.sales_wecom_id, {}).get("name") or ticket.sales_wecom_id or ""),
"salesName": str(self._staff_by_wecom.get(ticket.sales_wecom_id, {}).get("name") or ticket.sales_wecom_id or ""),
"facts": dict(ticket.facts),
"quote": quote,
"wait_phase": wait_phase,
"collab_chat_id": ticket.collab_chat_id,
"collab_facts": dict(ticket.collab_facts),
"lost_reason": ticket.lost_reason,
"lostReason": ticket.lost_reason,
"deal_notify_summary": ticket.deal_notify_summary,
"next_followup_at": ticket.next_followup_at,
"system_exception": ticket.system_exception,
"systemException": ticket.system_exception,
"system_exception_reason": ticket.system_exception_reason,
"systemExceptionReason": ticket.system_exception_reason,
"system_exception_step": ticket.system_exception_step,
"system_exception_payload": dict(ticket.system_exception_payload),
}
def upsert_exception_notify_user(self, row: dict[str, Any]) -> None:
"""单测注入「系统异常消息通知」接收人(后台账号已绑企微)。"""
data = dict(row)
uid = str(data.get("wecomId") or data.get("wecom_id") or "").strip()
if not uid:
return
data["wecomId"] = uid
with self._lock:
self._exception_notify_users = [
x for x in self._exception_notify_users
if str(x.get("wecomId") or "") != uid
]
self._exception_notify_users.append(data)
def list_exception_notify_users(self) -> list[dict[str, Any]]:
"""角色勾了系统异常消息通知、且账号已绑企微的人。"""
with self._lock:
return [dict(x) for x in self._exception_notify_users]
def mark_system_exception(
self,
*,
work_order_no: str,
system_exception: str,
system_exception_reason: str,
step: str = "",
payload: dict[str, Any] | None = None,
) -> dict[str, Any]:
"""写入工单两列与重试上下文;已关闭拒绝。不改六态。"""
no = (work_order_no or "").strip()
with self._lock:
ticket = self._tickets.get(no)
if ticket is None:
return {"ok": False, "error": "not_found"}
if ticket.status == "已关闭":
return {"ok": False, "error": "closed"}
ticket.system_exception = (system_exception or "").strip()
ticket.system_exception_reason = (system_exception_reason or "").strip()
ticket.system_exception_step = (step or "").strip()
ticket.system_exception_payload = dict(payload or {})
return {"ok": True, "status": ticket.status}
def clear_system_exception(self, *, work_order_no: str) -> dict[str, Any]:
"""重试成功后清两列。"""
no = (work_order_no or "").strip()
with self._lock:
ticket = self._tickets.get(no)
if ticket is None:
return {"ok": False, "error": "not_found"}
ticket.system_exception = ""
ticket.system_exception_reason = ""
ticket.system_exception_step = ""
ticket.system_exception_payload = {}
return {"ok": True}
def upsert_staff(self, row: dict[str, Any]) -> None:
"""单测注入任意岗位(含会话存档账号「询价机器人」)。"""
data = dict(row)
uid = str(data.get("wecomId") or data.get("wecom_id") or "").strip()
name = str(data.get("name") or "").strip()
with self._lock:
if uid:
self._staff_by_wecom[uid] = data
if name:
self._staff_by_name[name] = data
def set_sea_staff(self, rows: list[dict[str, Any]]) -> None:
"""单测注入海运产品岗。"""
with self._lock:
self._sea_staff = [dict(x) for x in rows]
for row in self._sea_staff:
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
name = str(row.get("name") or "").strip()
if uid:
self._staff_by_wecom[uid] = dict(row)
if name:
self._staff_by_name[name] = dict(row)
def match_sea_staff(self, *, origin: str, destination: str) -> list[dict[str, Any]]:
"""
岗位=产品·海运、启用:写了「海运·全部线路」必拉;
写了具体港口的,对得上本单起运港/目的港才拉。
"""
origin = (origin or "").strip()
dest = (destination or "").strip()
hits: list[dict[str, Any]] = []
with self._lock:
for row in self._sea_staff:
if str(row.get("roleCode") or row.get("role_code") or "") != "sea":
continue
if str(row.get("status") or "active") != "active":
continue
routes = str(row.get("routes") or "")
if "海运·全部线路" in routes or "海运全线" in routes:
hits.append(dict(row))
continue
if origin and origin in routes:
hits.append(dict(row))
continue
if dest and dest in routes:
hits.append(dict(row))
return hits
def set_land_staff(self, rows: list[dict[str, Any]]) -> None:
"""单测注入陆运产品岗。"""
with self._lock:
self._land_staff = [dict(x) for x in rows]
for row in self._land_staff:
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
name = str(row.get("name") or "").strip()
if uid:
self._staff_by_wecom[uid] = dict(row)
if name:
self._staff_by_name[name] = dict(row)
def match_land_staff(self, *, route_category: str) -> list[dict[str, Any]]:
"""
岗位=产品·陆运、启用:负责线路整词含本单线路类别才拉。
不按斜杠切开「中港/中亚/中欧」。
"""
want = (route_category or "").strip()
hits: list[dict[str, Any]] = []
if not want:
return hits
with self._lock:
for row in self._land_staff:
if str(row.get("roleCode") or row.get("role_code") or "") != "land":
continue
if str(row.get("status") or "active") != "active":
continue
tokens = [
p.strip()
for p in str(row.get("routes") or "").replace(",", ",").replace("、", ",").split(",")
if p.strip()
]
if want in tokens:
hits.append(dict(row))
return hits
def get_staff_by_wecom_id(self, *, wecom_id: str) -> dict[str, Any]:
"""按企微 userid 取员工展示名。"""
with self._lock:
return dict(self._staff_by_wecom.get((wecom_id or "").strip()) or {})
def list_staff_by_role(self, *, role_code: str) -> list[dict[str, Any]]:
"""按岗位列出启用员工。空运群 @ 航线用。"""
want = (role_code or "").strip().lower()
hits: list[dict[str, Any]] = []
with self._lock:
for row in self._staff_by_wecom.values():
code = str(row.get("roleCode") or row.get("role_code") or "").strip().lower()
status = str(row.get("status") or "active").strip().lower()
if code == want and status not in {"disabled", "0", "停用"}:
hits.append(dict(row))
return hits
def find_staff_by_name(self, *, name: str) -> dict[str, Any]:
"""按员工姓名取企微 userid。"""
with self._lock:
return dict(self._staff_by_name.get((name or "").strip()) or {})
def bind_collab_group(
self,
*,
work_order_no: str,
chat_id: str,
member_ids: list[str],
product_names: list[str] | None = None,
product_ids: list[str] | None = None,
) -> dict[str, Any]:
"""记下协同群和产品姓名,不改六态。重启后群摘要还要靠这些名字 @ 到人。"""
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
room = (chat_id or "").strip()
# 一群同时只跟一张单:陆运/海运旧绑定也要清掉,否则群里补字段会写到旧单。
self._unbind_other_collab_rooms(room=room, keep_work_order_no=work_order_no)
ticket.collab_chat_id = room
ticket.collab_member_ids = [str(x) for x in member_ids]
ticket.product_ids = [str(x) for x in (product_ids or [])]
ticket.product_names = [str(x).strip() for x in (product_names or []) if str(x).strip()]
return {"ok": True, "chat_id": ticket.collab_chat_id}
def end_collab_group(self, *, work_order_no: str) -> dict[str, Any]:
"""协同结束:不改六态,这单不能再拉群。"""
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
ticket.collab_ended = True
ticket.collab_facts["__ended"] = "1"
return {"ok": True, "collab_ended": True}
def attach_file(
self,
*,
work_order_no: str,
file_name: str,
object_key: str = "",
content_type: str = "",
sender_id: str = "",
) -> dict[str, Any]:
"""把群附件记到工单。不改六态。"""
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
row = {
"file_name": file_name,
"object_key": object_key or file_name,
"content_type": content_type,
"sender_id": sender_id,
}
ticket.attachments.append(row)
return {"ok": True, "file_name": file_name}
def upsert_product_quote(
self,
*,
work_order_no: str,
quote: dict[str, Any],
to_status: str = "已报价",
remark: str = "",
kind: str = "",
) -> dict[str, Any]:
"""手工报价 / 销售调价盖 TMS 价。remark/kind 只给 HTTP 主账记流水。"""
_ = remark, kind
return self.upsert_quote(
work_order_no=work_order_no, quote=quote, to_status=to_status
)
def patch_collab_facts(
self,
*,
work_order_no: str,
facts: dict[str, str],
sender_id: str = "",
) -> dict[str, Any]:
"""合并协同字段。不改六态、不写报价版本。"""
_ = sender_id
with self._lock:
ticket = self._tickets.get(work_order_no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
for key, raw in (facts or {}).items():
name = str(key or "").strip()
if not name:
continue
val = str(raw or "").strip()
if val:
ticket.collab_facts[name] = val
else:
ticket.collab_facts.pop(name, None)
return {"ok": True, "facts": dict(ticket.collab_facts)}
def _unbind_other_collab_rooms(self, *, room: str, keep_work_order_no: str) -> None:
"""
同一群只留 keep 这一张单。调用方必须已持 self._lock。
陆运/海运旧单也要摘掉。测服 WO202609230023 补协同写到了同群陆运 WO202609230007。
"""
room_id = (room or "").strip()
keep = (keep_work_order_no or "").strip()
if not room_id:
return
for other in self._tickets.values():
if other.work_order_no == keep:
continue
if other.collab_chat_id == room_id:
other.collab_chat_id = ""
def find_by_collab_chat(self, *, chat_id: str) -> MemoryTicket | None:
"""用工微群 chat_id 找回工单。"""
raw = (chat_id or "").strip()
if not raw:
return None
with self._lock:
for ticket in self._tickets.values():
if ticket.collab_chat_id == raw:
return ticket
return None
def activate_air_group(self, *, chat_id: str, work_order_no: str) -> dict[str, Any]:
"""
空运群当前工单:一单一群,未结案才能激活。
不改六态。
"""
no = (work_order_no or "").strip()
room = (chat_id or "").strip()
with self._lock:
ticket = self._tickets.get(no)
if ticket is None:
return {"ok": False, "error": "ticket_not_found"}
if (ticket.business_line or "").upper() != "AIR":
return {"ok": False, "error": "not_air", "status": ticket.status}
if ticket.status in {"已成交", "未成交", "已关闭"}:
return {"ok": False, "error": "terminal", "status": ticket.status}
if ticket.collab_chat_id and ticket.collab_chat_id != room:
return {"ok": False, "error": "occupied", "other_chat_id": ticket.collab_chat_id}
self._unbind_other_collab_rooms(room=room, keep_work_order_no=no)
ticket.collab_chat_id = room
return {
"ok": True,
"work_order_no": no,
"chat_id": room,
"status": ticket.status,
"business_line": ticket.business_line,
}
def lock_cabin(
self,
*,
work_order_no: str,
option_no: str,
sender_id: str = "",
operator_name: str = "",
instruction_text: str = "",
airline: str = "",
flight_no: str = "",
) -> dict[str, Any]:
"""单测可注入 _cabin_lock_result;默认成功。"""
_ = sender_id, operator_name, instruction_text, airline, flight_no
forced = getattr(self, "_cabin_lock_result", None)
if isinstance(forced, dict):
return dict(forced)
no = (work_order_no or "").strip()
ticket = self.get_ticket(work_order_no=no)
quote = dict(ticket.quote or {}) if ticket else {}
if not _memory_quote_can_lock(quote):
return {
"ok": False,
"error": "no_quote",
"failReason": "这张单还没有报价,不能锁舱",
}
quote_id = str(quote.get("quoteId") or "").strip()
if quote_id.upper().startswith("QT-WO"):
quote_id = ""
return {
"ok": True,
"status": "LOCKED",
"lockId": f"CA-{no}-LOCK",
"quoteId": quote_id,
"optionNo": option_no,
"tmsRequestId": f"TMS-LOCK-{no}",
}
def release_cabin(
self,
*,
work_order_no: str,
lock_id: str,
sender_id: str = "",
operator_name: str = "",
instruction_text: str = "",
) -> dict[str, Any]:
_ = lock_id, sender_id, operator_name, instruction_text
forced = getattr(self, "_cabin_release_result", None)
if isinstance(forced, dict):
return dict(forced)
no = (work_order_no or "").strip()
return {
"ok": True,
"status": "RELEASED",
"releaseId": f"CA-{no}-RLSE",
"tmsRequestId": f"TMS-RLSE-{no}",
}
def start_event_wait(self, *, work_order_no: str, wait_kind: str) -> dict[str, Any]:
"""开始盯等待。同种类不重置起点。不改六态。"""
no = (work_order_no or "").strip()
kind = (wait_kind or "").strip()
with self._lock:
ticket = self._tickets.get(no)
if ticket is None:
return {"ok": False, "error": "not_found"}
if ticket.status in {"已成交", "已关闭", "转人工"} or not kind:
return {"ok": False, "error": "not_watchable", "status": ticket.status}
if ticket.event_wait_kind == kind and ticket.event_wait_started_at:
return {"ok": True, "work_order_no": no, "eventWaitKind": kind}
ticket.event_wait_kind = kind
ticket.event_wait_started_at = datetime.now(_SHANGHAI).strftime("%Y-%m-%d %H:%M:%S")
ticket.event_wait_alerted = 0
return {"ok": True, "work_order_no": no, "eventWaitKind": kind}
def clear_event_wait(self, *, work_order_no: str) -> dict[str, Any]:
no = (work_order_no or "").strip()
with self._lock:
ticket = self._tickets.get(no)
if ticket is None:
return {"ok": False, "error": "not_found"}
ticket.event_wait_kind = ""
ticket.event_wait_started_at = ""
ticket.event_wait_alerted = 0
return {"ok": True, "work_order_no": no}
def list_event_waits(self) -> list[dict[str, Any]]:
out: list[dict[str, Any]] = []
with self._lock:
for ticket in self._tickets.values():
if not (ticket.event_wait_kind or "").strip():
continue
if int(ticket.event_wait_alerted or 0) != 0:
continue
if ticket.status in {"已成交", "已关闭", "转人工"}:
continue
watchers: list[dict[str, Any]] = []
seen: set[str] = set()
sales_id = (ticket.sales_wecom_id or "").strip()
if sales_id:
row = dict(self._staff_by_wecom.get(sales_id) or {})
watchers.append(
{
"wecomId": sales_id,
"name": str(row.get("name") or ticket.sales_wecom_id),
"roleCode": str(row.get("roleCode") or row.get("role_code") or "sales"),
"notifyEventException": row.get("notifyEventException", 0),
}
)
seen.add(sales_id)
for uid in ticket.collab_member_ids:
key = str(uid).strip()
if not key or key in seen:
continue
row = dict(self._staff_by_wecom.get(key) or {})
if not row:
continue
watchers.append(
{
"wecomId": key,
"name": str(row.get("name") or key),
"roleCode": str(row.get("roleCode") or row.get("role_code") or ""),
"notifyEventException": row.get("notifyEventException", 0),
}
)
seen.add(key)
out.append(
{
"work_order_no": ticket.work_order_no,
"workOrderNo": ticket.work_order_no,
"status": ticket.status,
"business_line": ticket.business_line,
"businessLine": ticket.business_line,
"event_wait_kind": ticket.event_wait_kind,
"eventWaitKind": ticket.event_wait_kind,
"event_wait_started_at": ticket.event_wait_started_at,
"eventWaitStartedAt": ticket.event_wait_started_at,
"sales_wecom_id": ticket.sales_wecom_id,
"salesWecomId": ticket.sales_wecom_id,
"sales_name": str(self._staff_by_wecom.get(sales_id, {}).get("name") or "销售"),
"collab_chat_id": ticket.collab_chat_id,
"collabChatId": ticket.collab_chat_id,
"collab_member_ids": list(ticket.collab_member_ids),
"facts": dict(ticket.facts),
"eventException": ticket.event_exception,
"watchers": watchers,
}
)
return out
def record_event_exception(self, *, work_order_no: str, event_exception: str) -> dict[str, Any]:
from agent.policy.event_exception import append_summary
no = (work_order_no or "").strip()
with self._lock:
ticket = self._tickets.get(no)
if ticket is None:
return {"ok": False, "error": "not_found"}
ticket.event_exception = append_summary(ticket.event_exception, event_exception)
ticket.event_wait_alerted = 1
return {
"ok": True,
"work_order_no": no,
"eventException": ticket.event_exception,
}
def get_timeout_by_code(self, *, timeout_code: str) -> dict[str, Any]:
code = (timeout_code or "").strip()
hours = self._timeouts.get(code)
if hours is None:
return {"found": False, "timeoutCode": code}
return {"found": True, "timeoutCode": code, "timeoutHours": hours}