817 lines
28 KiB
Python
817 lines
28 KiB
Python
"""
|
|
系统异常(UC024):技术故障判定、默默重试、工单两列、会话兜底、运维告警、后台重试。
|
|
|
|
本文件职责:只处理大模型 / TMS / 建群加人的超时、报错、空结果。
|
|
禁止:把无报价、字段不齐、没配人当成系统异常;禁止为记异常去建空工单;
|
|
禁止在回调线程等完整 LLM / 查 TMS。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from dataclasses import dataclass, field
|
|
from datetime import datetime
|
|
from typing import Any, Callable, Optional
|
|
from zoneinfo import ZoneInfo
|
|
|
|
from agent.policy import inquiry_copy as copy
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_SHANGHAI = ZoneInfo("Asia/Shanghai")
|
|
|
|
TYPE_LLM = copy.TYPE_LLM
|
|
TYPE_TMS = copy.TYPE_TMS
|
|
TYPE_GROUP = copy.TYPE_GROUP
|
|
STEP_LLM_NO_TICKET = copy.STEP_LLM_NO_TICKET
|
|
STEP_LLM = copy.STEP_LLM
|
|
STEP_TMS_QUOTE = copy.STEP_TMS_QUOTE
|
|
STEP_TMS_CABIN = copy.STEP_TMS_CABIN
|
|
STEP_GROUP = copy.STEP_GROUP
|
|
|
|
MAX_AUTO_RETRIES = 3
|
|
|
|
EnqueueFn = Callable[..., tuple[bool, str]]
|
|
ReplyFn = Callable[[str], Any]
|
|
OkFn = Callable[[Any], bool]
|
|
|
|
|
|
@dataclass
|
|
class SystemExceptionEvent:
|
|
"""一次已确认的系统异常(自动重试 3 次仍失败之后)。"""
|
|
|
|
exception_type: str
|
|
step: str
|
|
reason: str
|
|
service: str
|
|
error_code: str
|
|
work_order_no: str = ""
|
|
ticket_status: str = ""
|
|
conversation_kind: str = "private"
|
|
conversation_target: str = ""
|
|
missing_names: list[str] = field(default_factory=list)
|
|
extra: dict[str, Any] = field(default_factory=dict)
|
|
|
|
|
|
def is_tms_tech_failure(tms: dict[str, Any] | None) -> bool:
|
|
"""
|
|
TMS 技术故障:接口失败/超时。
|
|
|
|
正常返回但没报价、字段不齐未出站,都不算。
|
|
"""
|
|
row = tms or {}
|
|
if row.get("ok"):
|
|
return False
|
|
cls = str(row.get("classification") or "").strip()
|
|
if cls in {"TMS_NO_QUOTE_1002", "TMS_QUERY_INCOMPLETE"}:
|
|
return False
|
|
return True
|
|
|
|
|
|
def should_mark_ticket(*, work_order_no: str, ticket_status: str) -> bool:
|
|
"""有正式工单且不是已关闭,才写两列。"""
|
|
if not (work_order_no or "").strip():
|
|
return False
|
|
return (ticket_status or "").strip() != "已关闭"
|
|
|
|
|
|
def ticket_exception_type(ticket: Any) -> str:
|
|
"""从主账工单(对象或字典)读当前系统异常类型。"""
|
|
if ticket is None:
|
|
return ""
|
|
if isinstance(ticket, dict):
|
|
return str(
|
|
ticket.get("system_exception")
|
|
or ticket.get("systemException")
|
|
or ""
|
|
).strip()
|
|
return str(getattr(ticket, "system_exception", "") or "").strip()
|
|
|
|
|
|
def ticket_is_paused(ticket: Any) -> bool:
|
|
"""工单两列有值即暂停这一票。"""
|
|
return bool(ticket_exception_type(ticket))
|
|
|
|
|
|
def extract_work_order_nos(text: str) -> list[str]:
|
|
"""从原话抽出正式工单号(WO + 日期流水)。不认虚拟号。"""
|
|
import re
|
|
|
|
found = re.findall(r"\bWO\d{8,}\b", text or "", flags=re.IGNORECASE)
|
|
out: list[str] = []
|
|
seen: set[str] = set()
|
|
for raw in found:
|
|
no = raw.upper()
|
|
if no in seen:
|
|
continue
|
|
seen.add(no)
|
|
out.append(no)
|
|
return out
|
|
|
|
|
|
def call_with_retries(
|
|
fn: Callable[[], Any],
|
|
*,
|
|
is_ok: OkFn,
|
|
attempts: int = MAX_AUTO_RETRIES,
|
|
) -> Any:
|
|
"""
|
|
默默重试。中途成功立刻返回;会话里不说「正在重试」。
|
|
|
|
调用方在 Worker 槽使用,禁止放进企微回调线程。
|
|
"""
|
|
last: Any = None
|
|
times = max(1, int(attempts))
|
|
for idx in range(times):
|
|
try:
|
|
last = fn()
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("系统异常自动重试捕获异常 n=%s/%s err=%s", idx + 1, times, exc)
|
|
last = {"ok": False, "error": str(exc), "classification": "TECH_EXCEPTION"}
|
|
continue
|
|
if is_ok(last):
|
|
return last
|
|
logger.info("系统异常自动重试未成功 n=%s/%s", idx + 1, times)
|
|
return last
|
|
|
|
|
|
def _now_text(now: datetime | None) -> str:
|
|
stamp = now or datetime.now(_SHANGHAI)
|
|
return stamp.strftime("%Y-%m-%d %H:%M:%S")
|
|
|
|
|
|
def _test_prefix() -> bool:
|
|
try:
|
|
from agent.config import get_settings
|
|
|
|
return (get_settings().ytd_env or "").strip().lower() == "test"
|
|
except Exception: # noqa: BLE001
|
|
return False
|
|
|
|
|
|
def notify_app_ready() -> bool:
|
|
"""询价消息服务已配才推告警。"""
|
|
from agent.policy.deal_outcome import notify_app_ready as deal_ready
|
|
|
|
return deal_ready()
|
|
|
|
|
|
def _enqueue_default(
|
|
*,
|
|
touser: str,
|
|
content: str,
|
|
dedupe_key: str,
|
|
payload: dict[str, Any],
|
|
) -> tuple[bool, str]:
|
|
from agent.channel.queue import get_message_store
|
|
|
|
return get_message_store().enqueue_outbound(
|
|
touser=touser,
|
|
content=content,
|
|
dedupe_key=dedupe_key,
|
|
payload=payload,
|
|
)
|
|
|
|
|
|
def _list_notify_users(ledger: Any) -> list[dict[str, Any]]:
|
|
lister = getattr(ledger, "list_exception_notify_users", None)
|
|
if not callable(lister):
|
|
return []
|
|
try:
|
|
rows = list(lister() or [])
|
|
except Exception:
|
|
logger.exception("列出系统异常通知人失败")
|
|
return []
|
|
out: list[dict[str, Any]] = []
|
|
seen: set[str] = set()
|
|
for row in rows:
|
|
if not isinstance(row, dict):
|
|
continue
|
|
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
|
|
if not uid or uid in seen:
|
|
continue
|
|
seen.add(uid)
|
|
out.append(row)
|
|
return out
|
|
|
|
|
|
def confirm_system_exception(
|
|
*,
|
|
event: SystemExceptionEvent,
|
|
ledger: Any = None,
|
|
reply: Optional[ReplyFn] = None,
|
|
enqueue: Optional[EnqueueFn] = None,
|
|
notify_ready: Optional[bool] = None,
|
|
now: Optional[datetime] = None,
|
|
) -> dict[str, Any]:
|
|
"""
|
|
3 次仍失败后:回话、有单则写两列、推有权限的人。
|
|
|
|
无正式工单:不建单、不 mark。已关闭:不 mark。没有可推的人:照样回话。
|
|
"""
|
|
no = (event.work_order_no or "").strip()
|
|
if reply is not None:
|
|
reply(copy.system_exception_fallback(event.exception_type, names=event.missing_names))
|
|
|
|
marked = False
|
|
if should_mark_ticket(work_order_no=no, ticket_status=event.ticket_status) and ledger is not None:
|
|
marker = getattr(ledger, "mark_system_exception", None)
|
|
if callable(marker):
|
|
payload = {
|
|
"step": event.step,
|
|
"conversation_kind": event.conversation_kind,
|
|
"conversation_target": event.conversation_target,
|
|
"missing_names": list(event.missing_names),
|
|
"service": event.service,
|
|
"error_code": event.error_code,
|
|
"extra": dict(event.extra or {}),
|
|
}
|
|
try:
|
|
result = marker(
|
|
work_order_no=no,
|
|
system_exception=event.exception_type,
|
|
system_exception_reason=(event.reason or "").strip() or event.exception_type,
|
|
step=event.step,
|
|
payload=payload,
|
|
) or {}
|
|
marked = bool(result.get("ok", True))
|
|
except Exception:
|
|
logger.exception("写入系统异常失败 wo=%s", no)
|
|
marked = False
|
|
|
|
ready = notify_app_ready() if notify_ready is None else bool(notify_ready)
|
|
notified = 0
|
|
users = _list_notify_users(ledger) if ledger is not None else []
|
|
if ready and users:
|
|
occurred = _now_text(now)
|
|
test_prefix = _test_prefix()
|
|
payload = copy.system_exception_alert_payload(
|
|
work_order_no=no,
|
|
exception_type=event.exception_type,
|
|
service=event.service,
|
|
error_code=event.error_code,
|
|
step=event.step,
|
|
occurred_at=occurred,
|
|
test_prefix=test_prefix,
|
|
)
|
|
text = copy.system_exception_alert_text(
|
|
work_order_no=no,
|
|
exception_type=event.exception_type,
|
|
service=event.service,
|
|
error_code=event.error_code,
|
|
step=event.step,
|
|
occurred_at=occurred,
|
|
test_prefix=test_prefix,
|
|
)
|
|
sender = enqueue or _enqueue_default
|
|
for row in users:
|
|
uid = str(row.get("wecomId") or row.get("wecom_id") or "").strip()
|
|
if not uid:
|
|
continue
|
|
try:
|
|
ok, info = sender(
|
|
touser=uid,
|
|
content=text,
|
|
dedupe_key=f"sys-exc:{no or 'none'}:{event.step}:{uid}:{occurred}",
|
|
payload=dict(payload),
|
|
)
|
|
if ok:
|
|
notified += 1
|
|
else:
|
|
logger.warning("系统异常告警入队失败 uid=%s info=%s", uid, info)
|
|
except Exception:
|
|
logger.exception("系统异常告警入队异常 uid=%s", uid)
|
|
|
|
return {
|
|
"ok": True,
|
|
"marked": marked,
|
|
"notified": notified,
|
|
"work_order_no": no,
|
|
"exception_type": event.exception_type,
|
|
"step": event.step,
|
|
}
|
|
|
|
|
|
def clear_system_exception(*, ledger: Any, work_order_no: str) -> dict[str, Any]:
|
|
"""重试成功后清两列。"""
|
|
no = (work_order_no or "").strip()
|
|
if not no:
|
|
return {"ok": False, "error": "blank"}
|
|
clearer = getattr(ledger, "clear_system_exception", None)
|
|
if not callable(clearer):
|
|
return {"ok": False, "error": "no_clear"}
|
|
return clearer(work_order_no=no) or {"ok": False}
|
|
|
|
|
|
def reply_paused_ticket(*, ticket: Any, reply: ReplyFn) -> None:
|
|
"""销售/群提到暂停工单号:只重复兜底,不重跑。"""
|
|
kind = ticket_exception_type(ticket) or TYPE_TMS
|
|
names: list[str] = []
|
|
if isinstance(ticket, dict):
|
|
payload = ticket.get("system_exception_payload") or ticket.get("systemExceptionPayload") or {}
|
|
if isinstance(payload, dict):
|
|
raw = payload.get("missing_names") or []
|
|
names = [str(x).strip() for x in raw if str(x).strip()]
|
|
else:
|
|
payload = getattr(ticket, "system_exception_payload", {}) or {}
|
|
if isinstance(payload, dict):
|
|
raw = payload.get("missing_names") or []
|
|
names = [str(x).strip() for x in raw if str(x).strip()]
|
|
reply(copy.system_exception_fallback(kind, names=names))
|
|
|
|
|
|
def query_tms_resilient(ledger: Any, **kwargs: Any) -> dict[str, Any]:
|
|
"""TMS 查价:最多 3 次,ok=true(含无价)即成功。"""
|
|
getter = getattr(ledger, "query_tms", None)
|
|
if not callable(getter):
|
|
return {"ok": False, "classification": "TMS_TECH_FAILURE", "has_price": False, "quote": {}}
|
|
out = call_with_retries(
|
|
lambda: getter(**kwargs),
|
|
is_ok=lambda x: bool((x or {}).get("ok")),
|
|
)
|
|
return out if isinstance(out, dict) else {"ok": False, "classification": "TMS_TECH_FAILURE"}
|
|
|
|
|
|
def confirm_tms_quote_if_tech(
|
|
*,
|
|
ledger: Any,
|
|
reply: ReplyFn,
|
|
work_order_no: str,
|
|
tms: dict[str, Any],
|
|
sender_id: str,
|
|
ticket_status: str = "询价中",
|
|
facts: dict[str, Any] | None = None,
|
|
) -> bool:
|
|
"""技术失败则标 TMS异常并回话告警。返回是否已按系统异常处理。"""
|
|
if not is_tms_tech_failure(tms):
|
|
return False
|
|
err = str(tms.get("classification") or tms.get("error") or "TMS_TECH_FAILURE").strip()
|
|
confirm_system_exception(
|
|
event=SystemExceptionEvent(
|
|
exception_type=TYPE_TMS,
|
|
step=STEP_TMS_QUOTE,
|
|
reason="TMS /quote/v2 接口调用失败",
|
|
service="TMS报价",
|
|
error_code=err or "TMS_TECH_FAILURE",
|
|
work_order_no=work_order_no,
|
|
ticket_status=ticket_status,
|
|
conversation_kind="private",
|
|
conversation_target=sender_id,
|
|
extra={"facts": dict(facts or tms.get("facts") or {})},
|
|
),
|
|
ledger=ledger,
|
|
reply=reply,
|
|
)
|
|
return True
|
|
|
|
|
|
def confirm_group_if_tech(
|
|
*,
|
|
ledger: Any,
|
|
reply: ReplyFn,
|
|
work_order_no: str,
|
|
sender_id: str,
|
|
missing_names: list[str],
|
|
reason: str,
|
|
error_code: str,
|
|
ticket_status: str = "",
|
|
extra: dict[str, Any] | None = None,
|
|
chat_id: str = "",
|
|
) -> None:
|
|
"""建群或加人技术失败(已重试 3 次)后确认异常。"""
|
|
confirm_system_exception(
|
|
event=SystemExceptionEvent(
|
|
exception_type=TYPE_GROUP,
|
|
step=STEP_GROUP,
|
|
reason=(reason or "").strip() or "企微建群/加人失败",
|
|
service="企微建群" if not chat_id else "企微加人",
|
|
error_code=(error_code or "GROUP_FAIL").strip(),
|
|
work_order_no=work_order_no,
|
|
ticket_status=ticket_status or "已报价",
|
|
conversation_kind="private",
|
|
conversation_target=sender_id,
|
|
missing_names=list(missing_names),
|
|
extra=dict(extra or {}),
|
|
),
|
|
ledger=ledger,
|
|
reply=reply,
|
|
)
|
|
|
|
|
|
def confirm_llm_if_tech(
|
|
*,
|
|
ledger: Any,
|
|
reply: ReplyFn,
|
|
sender_id: str,
|
|
reason: str,
|
|
work_order_no: str = "",
|
|
ticket_status: str = "",
|
|
conversation_kind: str = "private",
|
|
conversation_target: str = "",
|
|
extra: dict[str, Any] | None = None,
|
|
) -> dict[str, Any]:
|
|
"""大模型超时/空/报错。无工单不建单。"""
|
|
no = (work_order_no or "").strip()
|
|
step = STEP_LLM if no else STEP_LLM_NO_TICKET
|
|
return confirm_system_exception(
|
|
event=SystemExceptionEvent(
|
|
exception_type=TYPE_LLM,
|
|
step=step,
|
|
reason=(reason or "").strip() or "模型返回空结果",
|
|
service="大模型识别",
|
|
error_code="EMPTY" if "空" in (reason or "") else "LLM_FAIL",
|
|
work_order_no=no,
|
|
ticket_status=ticket_status,
|
|
conversation_kind=conversation_kind,
|
|
conversation_target=conversation_target or sender_id,
|
|
extra=dict(extra or {}),
|
|
),
|
|
ledger=ledger,
|
|
reply=reply,
|
|
)
|
|
|
|
|
|
def handle_extract_tech_fail(
|
|
*,
|
|
snap: dict[str, Any],
|
|
ledger: Any,
|
|
reply: ReplyFn,
|
|
sender_id: str,
|
|
sess: Any = None,
|
|
) -> str:
|
|
"""抽字段模型超时/报错:标大模型异常并回话。返回阶段名,未失败返回空串。"""
|
|
if not isinstance(snap, dict) or not snap.get("tech_fail"):
|
|
return ""
|
|
wo = str(getattr(sess, "work_order_no", "") or "") if sess is not None else ""
|
|
st = str(getattr(sess, "status", "") or "") if sess is not None else ""
|
|
confirm_llm_if_tech(
|
|
ledger=ledger,
|
|
reply=reply,
|
|
sender_id=sender_id,
|
|
reason=str(snap.get("tech_error") or "模型返回空结果"),
|
|
work_order_no=wo,
|
|
ticket_status=st,
|
|
)
|
|
return "system_exception"
|
|
|
|
|
|
def _send_to_conversation(*, kind: str, target: str, text: str, extra: dict[str, Any] | None = None) -> str:
|
|
"""恢复/兜底发到卡住的会话:私聊走 outbox,群走 appchat。可带卡片 extra。"""
|
|
body = (text or "").strip()
|
|
dest = (target or "").strip()
|
|
if not body or not dest:
|
|
return ""
|
|
if (kind or "").strip() == "group":
|
|
try:
|
|
from agent.channel.outbox.sender import WeComAppClient
|
|
from agent.config import get_settings
|
|
|
|
s = get_settings()
|
|
client = WeComAppClient(
|
|
corp_id=s.wecom_corp_id,
|
|
secret=s.wecom_secret,
|
|
agent_id=int(s.wecom_agent_id or "0") or 0,
|
|
)
|
|
client.send_group(chat_id=dest, content=body)
|
|
return dest
|
|
except Exception:
|
|
logger.exception("系统异常群回话失败 chat=%s", dest)
|
|
return ""
|
|
try:
|
|
from agent.channel.queue import get_message_store
|
|
|
|
ok, info = get_message_store().enqueue_outbound(
|
|
touser=dest,
|
|
content=body,
|
|
dedupe_key=f"sys-exc-talk:{dest}:{hash(body) & 0xFFFF}",
|
|
payload=dict(extra or {}),
|
|
)
|
|
return str(info or "") if ok else ""
|
|
except Exception:
|
|
logger.exception("系统异常私聊回话失败 uid=%s", dest)
|
|
return ""
|
|
|
|
|
|
def run_admin_retry(work_order_no: str, *, ledger: Any = None) -> dict[str, Any]:
|
|
"""
|
|
后台点重试:只重跑失败步。成功先发恢复话术再继续;再失败再标再告警。
|
|
|
|
在 Worker HTTP 槽调用,禁止放进回调线程。
|
|
"""
|
|
no = (work_order_no or "").strip()
|
|
book = ledger
|
|
if book is None:
|
|
from agent.policy.air_text_flow import build_ledger
|
|
|
|
book = build_ledger()
|
|
getter = getattr(book, "get_ticket", None)
|
|
ticket = getter(work_order_no=no) if callable(getter) else None
|
|
if ticket is None:
|
|
getter = getattr(book, "get_for_agent", None)
|
|
ticket = getter(work_order_no=no, sender_id="") if callable(getter) else None
|
|
if not ticket_is_paused(ticket):
|
|
return {"ok": False, "error": "not_paused"}
|
|
if isinstance(ticket, dict):
|
|
step = str(ticket.get("system_exception_step") or "").strip()
|
|
payload = ticket.get("system_exception_payload") or {}
|
|
status = str(ticket.get("status") or "")
|
|
facts = dict(ticket.get("facts") or {})
|
|
line = str(ticket.get("business_line") or ticket.get("businessLine") or "")
|
|
else:
|
|
step = str(getattr(ticket, "system_exception_step", "") or "").strip()
|
|
payload = dict(getattr(ticket, "system_exception_payload", {}) or {})
|
|
status = str(getattr(ticket, "status", "") or "")
|
|
facts = dict(getattr(ticket, "facts", {}) or {})
|
|
line = str(getattr(ticket, "business_line", "") or "")
|
|
if not isinstance(payload, dict):
|
|
payload = {}
|
|
kind = str(payload.get("conversation_kind") or "private")
|
|
target = str(payload.get("conversation_target") or "")
|
|
extra = dict(payload.get("extra") or {})
|
|
if extra.get("facts") and not facts:
|
|
facts = {str(k): str(v) for k, v in dict(extra.get("facts") or {}).items()}
|
|
line = (line or str(extra.get("business_line") or "")).upper()
|
|
|
|
def reply(text: str, extra_msg: dict[str, Any] | None = None) -> str:
|
|
return _send_to_conversation(kind=kind, target=target, text=text, extra=extra_msg)
|
|
|
|
if step == STEP_TMS_QUOTE or not step:
|
|
tms = query_tms_resilient(book, work_order_no=no, facts=facts, business_line=line)
|
|
if is_tms_tech_failure(tms):
|
|
confirm_tms_quote_if_tech(
|
|
ledger=book,
|
|
reply=reply,
|
|
work_order_no=no,
|
|
tms=tms,
|
|
sender_id=target,
|
|
ticket_status=status,
|
|
facts=facts,
|
|
)
|
|
return {"ok": False, "error": "still_fail", "step": STEP_TMS_QUOTE}
|
|
reply(copy.system_exception_recovery(STEP_TMS_QUOTE))
|
|
clear_system_exception(ledger=book, work_order_no=no)
|
|
sender = target if kind != "group" else str(extra.get("sender_id") or "")
|
|
_resume_tms_quote(
|
|
ledger=book,
|
|
work_order_no=no,
|
|
facts=facts,
|
|
line=line,
|
|
sender_id=sender,
|
|
tms=tms,
|
|
reply=reply,
|
|
)
|
|
return {"ok": True, "step": STEP_TMS_QUOTE}
|
|
|
|
if step == STEP_GROUP:
|
|
extra = dict(payload.get("extra") or extra)
|
|
userids = [str(x).strip() for x in (extra.get("userids") or []) if str(x).strip()]
|
|
chat_id = str(extra.get("chat_id") or payload.get("chat_id") or "").strip()
|
|
try:
|
|
from agent.channel.outbox.sender import WeComAppClient
|
|
from agent.config import get_settings
|
|
|
|
s = get_settings()
|
|
client = WeComAppClient(
|
|
corp_id=s.wecom_corp_id,
|
|
secret=s.wecom_secret,
|
|
agent_id=int(s.wecom_agent_id or "0") or 0,
|
|
)
|
|
if chat_id and userids:
|
|
result = call_with_retries(
|
|
lambda: client.add_group_members(chat_id=chat_id, userids=userids) or {},
|
|
is_ok=lambda x: bool((x or {}).get("ok")),
|
|
)
|
|
elif userids:
|
|
name = no or str(extra.get("display_no") or "询价协同")
|
|
result = call_with_retries(
|
|
lambda: client.create_group(name=name, userids=userids) or {},
|
|
is_ok=lambda x: bool((x or {}).get("ok") and (x or {}).get("chat_id")),
|
|
)
|
|
else:
|
|
result = {"ok": False, "error": "no_users"}
|
|
except Exception as exc: # noqa: BLE001
|
|
result = {"ok": False, "error": str(exc)}
|
|
if not (result or {}).get("ok"):
|
|
confirm_group_if_tech(
|
|
ledger=book,
|
|
reply=reply,
|
|
work_order_no=no,
|
|
sender_id=target,
|
|
missing_names=[str(x) for x in (payload.get("missing_names") or extra.get("product_names") or [])],
|
|
reason=str((result or {}).get("error") or "企微建群/加人失败"),
|
|
error_code=str((result or {}).get("error") or "GROUP_FAIL"),
|
|
ticket_status=status,
|
|
extra=extra,
|
|
chat_id=chat_id,
|
|
)
|
|
return {"ok": False, "error": "still_fail", "step": STEP_GROUP}
|
|
new_chat = str((result or {}).get("chat_id") or chat_id).strip()
|
|
binder = getattr(book, "bind_collab_group", None)
|
|
if new_chat and callable(binder) and not chat_id:
|
|
binder(
|
|
work_order_no=no,
|
|
chat_id=new_chat,
|
|
member_ids=userids,
|
|
product_names=[str(x) for x in (extra.get("product_names") or [])],
|
|
product_ids=[str(x) for x in (extra.get("product_ids") or [])],
|
|
)
|
|
reply(copy.system_exception_recovery(STEP_GROUP))
|
|
clear_system_exception(ledger=book, work_order_no=no)
|
|
return {"ok": True, "step": STEP_GROUP}
|
|
|
|
if step in {STEP_LLM, STEP_LLM_NO_TICKET}:
|
|
reply(copy.system_exception_recovery(STEP_LLM))
|
|
if no:
|
|
clear_system_exception(ledger=book, work_order_no=no)
|
|
reply("请把刚才那句话再发一次。")
|
|
return {"ok": True, "step": step}
|
|
|
|
if step == STEP_TMS_CABIN:
|
|
return _resume_cabin_step(
|
|
ledger=book,
|
|
work_order_no=no,
|
|
extra=extra,
|
|
target=target,
|
|
status=status,
|
|
reply=reply,
|
|
)
|
|
|
|
return {"ok": False, "error": "unknown_step", "step": step}
|
|
|
|
|
|
def _flow_for_line(line: str) -> Any:
|
|
"""按业务线取私聊流程,用来重试后把报价卡/书签补回去。"""
|
|
key = (line or "").upper()
|
|
if key == "LAND":
|
|
from agent.policy.land_text_flow import get_land_text_flow
|
|
|
|
return get_land_text_flow()
|
|
if key == "AIR":
|
|
from agent.policy.air_text_flow import get_air_text_flow
|
|
|
|
return get_air_text_flow()
|
|
from agent.policy.sea_text_flow import get_sea_text_flow
|
|
|
|
return get_sea_text_flow()
|
|
|
|
|
|
def _resume_tms_quote(
|
|
*,
|
|
ledger: Any,
|
|
work_order_no: str,
|
|
facts: dict[str, Any],
|
|
line: str,
|
|
sender_id: str,
|
|
tms: dict[str, Any],
|
|
reply: ReplyFn,
|
|
) -> None:
|
|
"""重试查价成功:写入报价并按原流程出卡,销售能继续拉群。"""
|
|
from agent.policy.air_text_flow import FlowSession
|
|
|
|
no = (work_order_no or "").strip()
|
|
biz = (line or "SEA").upper()
|
|
if biz not in {"SEA", "LAND", "AIR"}:
|
|
biz = "SEA"
|
|
flow = _flow_for_line(biz)
|
|
sender = (sender_id or "").strip()
|
|
if not tms.get("has_price"):
|
|
miss = getattr(flow, "_emit_tms_miss", None)
|
|
if callable(miss):
|
|
miss(reply, no)
|
|
else:
|
|
reply(f"工单{no}暂无 TMS 报价。")
|
|
if sender:
|
|
sess = FlowSession(
|
|
sender_id=sender,
|
|
thread_id=sender,
|
|
phase="tms_miss",
|
|
business_line=biz,
|
|
facts=dict(facts or {}),
|
|
work_order_no=no,
|
|
status="询价中",
|
|
allowed=(copy.BTN_PULL_COLLAB,) if biz != "AIR" else (),
|
|
wait_version=1,
|
|
)
|
|
flow._save(sess)
|
|
return
|
|
quote = dict(tms.get("quote") or {})
|
|
up = ledger.upsert_quote(work_order_no=no, quote=quote, to_status="已报价")
|
|
if not up.get("ok"):
|
|
reply("主账写入报价失败,不能当作已报价。")
|
|
return
|
|
allowed = (copy.BTN_SKIP_COLLAB,) if biz == "AIR" else (copy.BTN_PULL_COLLAB, copy.BTN_SKIP_COLLAB)
|
|
if sender:
|
|
sess = FlowSession(
|
|
sender_id=sender,
|
|
thread_id=sender,
|
|
phase="wait_collab",
|
|
business_line=biz,
|
|
facts=dict(facts or {}),
|
|
work_order_no=no,
|
|
quote=quote,
|
|
status="已报价",
|
|
allowed=allowed,
|
|
wait_version=1,
|
|
)
|
|
copy.remember_tms_quote(sess, quote)
|
|
flow._save(sess)
|
|
if biz == "AIR":
|
|
reply(
|
|
copy.tms_hit_card(work_order_no=no, quote=quote, status="已报价"),
|
|
copy.tms_hit_wecom_payload(work_order_no=no, quote=quote, status="已报价"),
|
|
)
|
|
return
|
|
emit = getattr(flow, "_emit_sea_hit", None)
|
|
if callable(emit):
|
|
emit(
|
|
reply,
|
|
work_order_no=no,
|
|
quote=quote,
|
|
status="已报价",
|
|
business_line=biz,
|
|
)
|
|
return
|
|
reply(f"工单{no}已重新查到报价。")
|
|
|
|
|
|
def _resume_cabin_step(
|
|
*,
|
|
ledger: Any,
|
|
work_order_no: str,
|
|
extra: dict[str, Any],
|
|
target: str,
|
|
status: str,
|
|
reply: ReplyFn,
|
|
) -> dict[str, Any]:
|
|
"""重试空运锁舱或释放:payload 里有 lock_id 就是释放。"""
|
|
no = (work_order_no or "").strip()
|
|
lock_id = str(extra.get("lock_id") or "").strip()
|
|
option_no = str(extra.get("option_no") or "").strip() or "AIR-OPT-01"
|
|
instruction = str(extra.get("instruction") or "")
|
|
sender = str(extra.get("sender_id") or target or "").strip()
|
|
if lock_id:
|
|
releaser = getattr(ledger, "release_cabin", None)
|
|
out = call_with_retries(
|
|
lambda: releaser(
|
|
work_order_no=no,
|
|
lock_id=lock_id,
|
|
sender_id=sender,
|
|
instruction_text=instruction,
|
|
)
|
|
or {}
|
|
if callable(releaser)
|
|
else {"ok": False, "error": "no_release"},
|
|
is_ok=lambda x: bool((x or {}).get("ok")),
|
|
)
|
|
service = "空运舱位"
|
|
reason = "空运释放舱位接口调用失败"
|
|
else:
|
|
locker = getattr(ledger, "lock_cabin", None)
|
|
out = call_with_retries(
|
|
lambda: locker(
|
|
work_order_no=no,
|
|
option_no=option_no,
|
|
sender_id=sender,
|
|
instruction_text=instruction,
|
|
)
|
|
or {}
|
|
if callable(locker)
|
|
else {"ok": False, "error": "no_lock"},
|
|
is_ok=lambda x: bool((x or {}).get("ok")),
|
|
)
|
|
service = "空运锁舱"
|
|
reason = "空运锁舱接口调用失败"
|
|
if not (out or {}).get("ok"):
|
|
confirm_system_exception(
|
|
event=SystemExceptionEvent(
|
|
exception_type=TYPE_TMS,
|
|
step=STEP_TMS_CABIN,
|
|
reason=reason,
|
|
service=service,
|
|
error_code=str((out or {}).get("failReason") or (out or {}).get("error") or "CABIN_FAIL"),
|
|
work_order_no=no,
|
|
ticket_status=status,
|
|
conversation_kind="group",
|
|
conversation_target=target,
|
|
extra=dict(extra),
|
|
),
|
|
ledger=ledger,
|
|
reply=reply,
|
|
)
|
|
return {"ok": False, "error": "still_fail", "step": STEP_TMS_CABIN}
|
|
reply(copy.system_exception_recovery(STEP_TMS_CABIN))
|
|
clear_system_exception(ledger=ledger, work_order_no=no)
|
|
if not lock_id:
|
|
exec_no = str((out or {}).get("lockId") or (out or {}).get("tmsRequestId") or "").strip()
|
|
patcher = getattr(ledger, "patch_collab_facts", None)
|
|
if exec_no and callable(patcher):
|
|
patcher(
|
|
work_order_no=no,
|
|
facts={"__air_lock_id": exec_no, "__air_lock_option": option_no, "__pending_air_lock": ""},
|
|
sender_id=sender,
|
|
)
|
|
reply(copy.air_lock_ok(work_order_no=no, operator="同事", option_no=option_no, exec_no=exec_no, quote_no=""))
|
|
else:
|
|
reply(copy.air_release_ok(work_order_no=no, operator="同事", option_no="", exec_no=str((out or {}).get("releaseId") or lock_id), quote_no=""))
|
|
return {"ok": True, "step": STEP_TMS_CABIN}
|