diff --git a/inquiry-agent/agent/handlers/air_group_collab.py b/inquiry-agent/agent/handlers/air_group_collab.py index 8a1544a..5955e8f 100644 --- a/inquiry-agent/agent/handlers/air_group_collab.py +++ b/inquiry-agent/agent/handlers/air_group_collab.py @@ -158,6 +158,14 @@ def handle_air_group( if is_handoff_ticket(ticket): logger.info("air_group 转人工不再跟 chat=%s", chat_id) return "handoff_ignore" + from agent.policy.system_exception import reply_paused_ticket, ticket_is_paused + + if ticket_is_paused(ticket): + reply_paused_ticket( + ticket=ticket, + reply=lambda text: send_group_text(group_client_of(engine), chat_id, text), + ) + return "system_exception_paused" return _handle_bound( engine, message, diff --git a/inquiry-agent/agent/handlers/group_collab.py b/inquiry-agent/agent/handlers/group_collab.py index d05a5da..9b61fcc 100644 --- a/inquiry-agent/agent/handlers/group_collab.py +++ b/inquiry-agent/agent/handlers/group_collab.py @@ -100,6 +100,15 @@ def handle_group_inbound( logger.info("group_inbound 忽略非成员 sender=%s chat=%s", message.sender_id, chat_id) return "group_ignore" + from agent.policy.system_exception import reply_paused_ticket, ticket_is_paused + + if ticket_is_paused(ticket): + reply_paused_ticket( + ticket=ticket, + reply=lambda text: send_group_text(group_client_of(engine), chat_id, text), + ) + return "system_exception_paused" + sess = _session_of_ticket(engine, ticket) if is_handoff_ticket(ticket, sess): logger.info( diff --git a/inquiry-agent/agent/handlers/text_inquiry.py b/inquiry-agent/agent/handlers/text_inquiry.py index ce7db55..3b81bf3 100644 --- a/inquiry-agent/agent/handlers/text_inquiry.py +++ b/inquiry-agent/agent/handlers/text_inquiry.py @@ -256,6 +256,28 @@ def handle_text_inquiry( logger.info("text_inquiry 出站 ok=%s info=%s phase_reply=%s", ok, info, body[:80]) return info if ok else "" + from agent.policy.system_exception import ( + extract_work_order_nos, + reply_paused_ticket, + ticket_is_paused, + ) + + ledger = getattr(engine, "_ledger", None) + for no in extract_work_order_nos(message.content or ""): + ticket = None + if ledger is not None: + getter = getattr(ledger, "get_ticket", None) + if callable(getter): + ticket = getter(work_order_no=no) + if ticket is None: + getter = getattr(ledger, "get_for_agent", None) + if callable(getter): + ticket = getter(work_order_no=no, sender_id=message.sender_id) + if ticket_is_paused(ticket): + reply_paused_ticket(ticket=ticket, reply=reply) + logger.info("text_inquiry 暂停工单只回兜底 wo=%s sender=%s", no, message.sender_id) + return "system_exception_paused" + if prefer_handoff: phase = engine.offer_handoff( sender_id=message.sender_id, diff --git a/inquiry-agent/agent/jobs/wake_consumer.py b/inquiry-agent/agent/jobs/wake_consumer.py index f06f67a..68ed342 100644 --- a/inquiry-agent/agent/jobs/wake_consumer.py +++ b/inquiry-agent/agent/jobs/wake_consumer.py @@ -73,8 +73,19 @@ def claim_and_apply_wake( kind=str(data.get("kind") or ""), detail=dict(data.get("detail") or {}), ) - graph.wake_from_worker(wake) - logger.info("wake 已处理 id=%s thread=%s kind=%s", entry_id, wake.thread_id, wake.kind) + if wake.kind == "system_exception_retry": + from agent.policy.system_exception import run_admin_retry + + result = run_admin_retry(wake.inquiry_no) + logger.info( + "系统异常重试 id=%s wo=%s ok=%s", + entry_id, + wake.inquiry_no, + result.get("ok"), + ) + else: + graph.wake_from_worker(wake) + logger.info("wake 已处理 id=%s thread=%s kind=%s", entry_id, wake.thread_id, wake.kind) except Exception: # noqa: BLE001 logger.exception("wake 处理失败 id=%s", entry_id) finally: diff --git a/inquiry-agent/agent/ledger/http_ledger.py b/inquiry-agent/agent/ledger/http_ledger.py index 59ce518..5a10d5a 100644 --- a/inquiry-agent/agent/ledger/http_ledger.py +++ b/inquiry-agent/agent/ledger/http_ledger.py @@ -433,3 +433,47 @@ class HttpLedger: body, idem=f"manual-quote:{work_order_no}:{sig}", ) + + 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]: + """写入工单系统异常两列。已关闭由主账拒绝。""" + out = self._post( + "/inquiry/agent/ticket/markSystemException", + { + "workOrderNo": work_order_no, + "systemException": system_exception, + "systemExceptionReason": system_exception_reason, + "step": step, + "payload": dict(payload or {}), + }, + ) + if out.get("ok") is False: + return out + out.setdefault("ok", True) + return out + + def clear_system_exception(self, *, work_order_no: str) -> dict[str, Any]: + """重试成功后清两列。""" + out = self._post( + "/inquiry/agent/ticket/clearSystemException", + {"workOrderNo": work_order_no}, + ) + if out.get("ok") is False: + return out + out.setdefault("ok", True) + return out + + def list_exception_notify_users(self) -> list[dict[str, Any]]: + """角色勾了系统异常消息通知、且账号已绑企微。""" + out = self._post("/inquiry/agent/account/listExceptionNotify", {}) + rows = out.get("staff") or out.get("result") or [] + if isinstance(rows, list): + return [dict(x) for x in rows if isinstance(x, dict)] + return [] diff --git a/inquiry-agent/agent/ledger/memory_ledger.py b/inquiry-agent/agent/ledger/memory_ledger.py index 0ad3b88..a5816f1 100644 --- a/inquiry-agent/agent/ledger/memory_ledger.py +++ b/inquiry-agent/agent/ledger/memory_ledger.py @@ -44,6 +44,10 @@ class MemoryTicket: lost_reason: str = "" next_followup_at: str = "" deal_notify_summary: str = "" + system_exception: str = "" + system_exception_reason: str = "" + system_exception_step: str = "" + system_exception_payload: dict[str, Any] = field(default_factory=dict) class MemoryLedger: @@ -72,6 +76,7 @@ class MemoryLedger: 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]] = [] def _next_no(self) -> str: self._seq += 1 @@ -373,8 +378,69 @@ class MemoryLedger: "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) diff --git a/inquiry-agent/agent/llm/extract_text.py b/inquiry-agent/agent/llm/extract_text.py index f6e3ba9..a024a0c 100644 --- a/inquiry-agent/agent/llm/extract_text.py +++ b/inquiry-agent/agent/llm/extract_text.py @@ -145,6 +145,19 @@ def extract_inquiry_snapshot( b_payload = invoke_extract_fields( text, allow_network=allow_b, chat_fn=chat_fn ) + if isinstance(b_payload, dict) and not b_payload.get("ok"): + err = str(b_payload.get("error") or "").strip() + if err not in {"b_network_disabled", "extract_fields_shell_not_implemented"} and ( + allow_b or chat_fn is not None + ): + return { + "business_line": mode, + "land_subtype": "", + "facts": {}, + "source": "llm_tech_fail", + "tech_fail": True, + "tech_error": err or "LLM_FAIL", + } if b_payload and b_payload.get("ok"): merged = dict(b_payload.get("facts") or {}) merged.update(facts) diff --git a/inquiry-agent/agent/policy/air_group_ops.py b/inquiry-agent/agent/policy/air_group_ops.py index ef99167..541c272 100644 --- a/inquiry-agent/agent/policy/air_group_ops.py +++ b/inquiry-agent/agent/policy/air_group_ops.py @@ -586,6 +586,14 @@ def apply_air_lock( if not is_airline(flow, sender_id): send_group_text(group_client_of(flow), chat_id, copy.air_airline_only_cabin()) return "not_airline" + from agent.policy.system_exception import reply_paused_ticket, ticket_is_paused + + if ticket_is_paused(ticket): + reply_paused_ticket( + ticket=ticket, + reply=lambda text: send_group_text(group_client_of(flow), chat_id, text), + ) + return "system_exception_paused" no = ticket_no(ticket) name = operator_name(flow, sender_id) quote = ticket_view(ticket).get("quote") or {} @@ -612,24 +620,43 @@ def apply_air_lock( copy.air_lock_fail(work_order_no=no, operator=name, reason="主账未开放锁舱"), ) return "lock_unavailable" - out = locker( - work_order_no=no, - option_no=picked, - sender_id=sender_id, - operator_name=name, - instruction_text=instruction, - ) or {} - if not out.get("ok"): - send_group_text( - group_client_of(flow), - chat_id, - copy.air_lock_fail( - work_order_no=no, - operator=name, - reason=str(out.get("failReason") or out.get("error") or "锁舱失败"), - ), + from agent.policy.system_exception import ( + STEP_TMS_CABIN, + TYPE_TMS, + SystemExceptionEvent, + call_with_retries, + confirm_system_exception, + ) + + out = call_with_retries( + lambda: locker( + work_order_no=no, + option_no=picked, + sender_id=sender_id, + operator_name=name, + instruction_text=instruction, ) - return "lock_fail" + or {}, + is_ok=lambda x: bool((x or {}).get("ok")), + ) + if not out.get("ok"): + confirm_system_exception( + event=SystemExceptionEvent( + exception_type=TYPE_TMS, + step=STEP_TMS_CABIN, + reason="空运锁舱接口调用失败", + service="空运锁舱", + error_code=str(out.get("failReason") or out.get("error") or "LOCK_FAIL"), + work_order_no=no, + ticket_status=ticket_status(ticket), + conversation_kind="group", + conversation_target=chat_id, + extra={"option_no": picked, "instruction": instruction, "sender_id": sender_id}, + ), + ledger=flow.ledger, + reply=lambda text: send_group_text(group_client_of(flow), chat_id, text), + ) + return "system_exception" exec_no = str(out.get("lockId") or out.get("tmsRequestId") or "").strip() quote_no = str(out.get("quoteId") or quote_id).strip() flow.ledger.patch_collab_facts( @@ -671,6 +698,14 @@ def apply_air_release( if not is_airline(flow, sender_id): send_group_text(group_client_of(flow), chat_id, copy.air_airline_only_cabin()) return "not_airline" + from agent.policy.system_exception import reply_paused_ticket, ticket_is_paused + + if ticket_is_paused(ticket): + reply_paused_ticket( + ticket=ticket, + reply=lambda text: send_group_text(group_client_of(flow), chat_id, text), + ) + return "system_exception_paused" no = ticket_no(ticket) if ticket_status(ticket) == "已成交": send_group_text(group_client_of(flow), chat_id, copy.air_release_deal_blocked(no)) @@ -689,24 +724,43 @@ def apply_air_release( copy.air_release_fail(work_order_no=no, operator=name, reason="主账未开放释放"), ) return "release_unavailable" - out = releaser( - work_order_no=no, - lock_id=lock_id, - sender_id=sender_id, - operator_name=name, - instruction_text=instruction, - ) or {} - if not out.get("ok"): - send_group_text( - group_client_of(flow), - chat_id, - copy.air_release_fail( - work_order_no=no, - operator=name, - reason=str(out.get("failReason") or out.get("error") or "释放失败"), - ), + from agent.policy.system_exception import ( + STEP_TMS_CABIN, + TYPE_TMS, + SystemExceptionEvent, + call_with_retries, + confirm_system_exception, + ) + + out = call_with_retries( + lambda: releaser( + work_order_no=no, + lock_id=lock_id, + sender_id=sender_id, + operator_name=name, + instruction_text=instruction, ) - return "release_fail" + or {}, + is_ok=lambda x: bool((x or {}).get("ok")), + ) + if not out.get("ok"): + confirm_system_exception( + event=SystemExceptionEvent( + exception_type=TYPE_TMS, + step=STEP_TMS_CABIN, + reason="空运释放舱位接口调用失败", + service="空运舱位", + error_code=str(out.get("failReason") or out.get("error") or "RELEASE_FAIL"), + work_order_no=no, + ticket_status=ticket_status(ticket), + conversation_kind="group", + conversation_target=chat_id, + extra={"lock_id": lock_id, "instruction": instruction, "sender_id": sender_id}, + ), + ledger=flow.ledger, + reply=lambda text: send_group_text(group_client_of(flow), chat_id, text), + ) + return "system_exception" exec_no = str(out.get("releaseId") or out.get("tmsRequestId") or lock_id).strip() send_group_text( group_client_of(flow), diff --git a/inquiry-agent/agent/policy/air_text_flow.py b/inquiry-agent/agent/policy/air_text_flow.py index 14ba473..24bcf6d 100644 --- a/inquiry-agent/agent/policy/air_text_flow.py +++ b/inquiry-agent/agent/policy/air_text_flow.py @@ -726,6 +726,17 @@ class AirTextInquiryFlow: allow_b=allow_b, chat_fn=chat_fn, ) + from agent.policy.system_exception import handle_extract_tech_fail + + llm_phase = handle_extract_tech_fail( + snap=snap, + ledger=self._ledger, + reply=reply, + sender_id=sender_id, + sess=sess, + ) + if llm_phase: + return llm_phase mode = (snap.get("business_line") or "").upper() facts = dict(snap.get("facts") or {}) if sess and sess.phase in {"clarify", "need_mode"}: @@ -869,7 +880,12 @@ class AirTextInquiryFlow: # 询价确认必须先到销售手机,再打 TMS。否则查得快时暂无报价卡会抢先发出。 self._wait_inquiry_sent(str(out_id), work_order_no=no) - tms = self._ledger.query_tms(work_order_no=no, facts=facts) + from agent.policy.system_exception import ( + confirm_tms_quote_if_tech, + query_tms_resilient, + ) + + tms = query_tms_resilient(self._ledger, work_order_no=no, facts=facts) sess = FlowSession( sender_id=sender_id, thread_id=sender_id, @@ -883,6 +899,16 @@ class AirTextInquiryFlow: ) if not tms.get("ok"): # 没拿到 TMS 回包(超时/熔断/主账未出站)不能冒充「暂无报价」 + if confirm_tms_quote_if_tech( + ledger=self._ledger, + reply=reply, + work_order_no=no, + tms=tms, + sender_id=sender_id, + ticket_status="询价中", + ): + self.drop_current_bookmark(sender_id) + return "system_exception" if tms.get("classification") == "TMS_QUERY_INCOMPLETE": logger.warning("tms.query 未出站 wo=%s class=%s", no, tms.get("classification")) reply("TMS 查询失败,请稍后重试。工单已创建:" + no) diff --git a/inquiry-agent/agent/policy/inquiry_copy.py b/inquiry-agent/agent/policy/inquiry_copy.py index d357b25..f8ef3a4 100644 --- a/inquiry-agent/agent/policy/inquiry_copy.py +++ b/inquiry-agent/agent/policy/inquiry_copy.py @@ -2473,3 +2473,142 @@ def air_release_deal_blocked(work_order_no: str) -> str: def air_release_need_lock(work_order_no: str) -> str: no = (work_order_no or "").strip() return f"工单{no}还没有有效锁舱,不能释放。" + + +# —— 系统异常(UC024):销售/群兜底与恢复;告警卡只给有权限的人 —— +TYPE_LLM = "大模型异常" +TYPE_TMS = "TMS异常" +TYPE_GROUP = "建群异常" + +STEP_LLM_NO_TICKET = "llm_no_ticket" +STEP_LLM = "llm" +STEP_TMS_QUOTE = "tms_quote" +STEP_TMS_CABIN = "tms_cabin" +STEP_GROUP = "group" + +SYS_EXC_LLM_FALLBACK = "大模型服务异常,IT排查中..." +SYS_EXC_TMS_FALLBACK = "TMS接口异常,IT运维排查中..." +SYS_EXC_RECOVERY_TMS = "TMS已恢复正常使用,感谢您的耐心等待。" +SYS_EXC_RECOVERY_LLM = "大模型已恢复正常使用,感谢您的耐心等待。" +SYS_EXC_RECOVERY_CABIN = "空运舱位接口已恢复正常使用,感谢您的耐心等待。" +SYS_EXC_RECOVERY_GROUP = "建群已恢复正常,感谢您的耐心等待。" + +IMPACT_LLM_NO_TICKET = "尚未建单,不落后台工单" +IMPACT_LLM = "当前识别暂停,已有字段与工单状态保持不变" +IMPACT_TMS_QUOTE = "当前询价流程暂停,识别字段与原业务节点保持不变" +IMPACT_TMS_CABIN = "当前锁舱/释放暂停,工单状态保持不变" +IMPACT_GROUP = "当前建群/加人暂停,已建群不新建第二个群" + + +def system_exception_fallback(exception_type: str, *, names: list[str] | None = None) -> str: + """判定系统异常后,给卡住的会话的兜底话术。""" + kind = (exception_type or "").strip() + if kind == TYPE_LLM: + return SYS_EXC_LLM_FALLBACK + if kind == TYPE_GROUP: + people = [str(x).strip() for x in (names or []) if str(x).strip()] + if not people: + return "无法添加同事,请联系管理员检查企业微信配置。" + joined = "、".join(people) + return f"无法添加 {joined} 同事,请联系管理员检查企业微信配置。" + return SYS_EXC_TMS_FALLBACK + + +def system_exception_recovery(step: str) -> str: + """后台点重试成功后,先发这一句,再补没做完的事。""" + key = (step or "").strip() + if key == STEP_LLM or key == STEP_LLM_NO_TICKET: + return SYS_EXC_RECOVERY_LLM + if key == STEP_TMS_CABIN: + return SYS_EXC_RECOVERY_CABIN + if key == STEP_GROUP: + return SYS_EXC_RECOVERY_GROUP + return SYS_EXC_RECOVERY_TMS + + +def system_exception_impact(step: str) -> str: + """告警卡「影响范围」,按卡住的步骤。""" + key = (step or "").strip() + if key == STEP_LLM_NO_TICKET: + return IMPACT_LLM_NO_TICKET + if key == STEP_LLM: + return IMPACT_LLM + if key == STEP_TMS_CABIN: + return IMPACT_TMS_CABIN + if key == STEP_GROUP: + return IMPACT_GROUP + return IMPACT_TMS_QUOTE + + +def system_exception_alert_payload( + *, + work_order_no: str, + exception_type: str, + service: str, + error_code: str, + step: str, + occurred_at: str, + test_prefix: bool = False, +) -> dict[str, Any]: + """ + 系统异常通知卡:走询价消息服务,不走询价小助手。 + + 五行对齐现网示例:异常类型、接口/服务、错误码、影响范围、发生时间。 + """ + no = (work_order_no or "").strip() or "无" + title = "系统异常通知" + if test_prefix: + title = "[test] " + title + rows = [ + {"keyname": "异常类型", "value": (exception_type or "").strip() or "-"}, + {"keyname": "接口/服务", "value": (service or "").strip() or "-"}, + {"keyname": "错误码", "value": (error_code or "").strip() or "-"}, + {"keyname": "影响范围", "value": system_exception_impact(step)}, + {"keyname": "发生时间", "value": (occurred_at or "").strip() or "-"}, + ] + from agent.config import get_settings + + base = str(getattr(get_settings(), "public_base_url", "") or "").strip().rstrip("/") + action_url = base or "https://work.weixin.qq.com" + return { + "msgtype": "template_card", + "wecom_app": "notify", + "template_card": { + "card_type": "text_notice", + "source": {"desc": "询价机器人系统告警", "desc_color": 0}, + "main_title": { + "title": title, + "desc": f"关联工单 {no}", + }, + "horizontal_content_list": rows, + "card_action": {"type": 1, "url": action_url}, + }, + } + + +def system_exception_alert_text( + *, + work_order_no: str, + exception_type: str, + service: str, + error_code: str, + step: str, + occurred_at: str, + test_prefix: bool = False, +) -> str: + """告警卡纯文本兜底(通道降级时仍能看懂)。""" + no = (work_order_no or "").strip() or "无" + head = "系统异常通知" + if test_prefix: + head = "[test] " + head + lines = [ + "询价机器人系统告警", + head, + f"关联工单 {no}", + f"异常类型:{(exception_type or '').strip() or '-'}", + f"接口/服务:{(service or '').strip() or '-'}", + f"错误码:{(error_code or '').strip() or '-'}", + f"影响范围:{system_exception_impact(step)}", + f"发生时间:{(occurred_at or '').strip() or '-'}", + ] + return "\n".join(lines) diff --git a/inquiry-agent/agent/policy/land_text_flow.py b/inquiry-agent/agent/policy/land_text_flow.py index c1e055a..351bf8d 100644 --- a/inquiry-agent/agent/policy/land_text_flow.py +++ b/inquiry-agent/agent/policy/land_text_flow.py @@ -123,6 +123,17 @@ class LandTextInquiryFlow(SeaTextInquiryFlow): allow_b=allow_b, chat_fn=chat_fn, ) + from agent.policy.system_exception import handle_extract_tech_fail + + llm_phase = handle_extract_tech_fail( + snap=snap, + ledger=self._ledger, + reply=reply, + sender_id=sender_id, + sess=sess, + ) + if llm_phase: + return llm_phase mode = (snap.get("business_line") or "").upper() snap_facts = dict(snap.get("facts") or {}) facts = dict(snap_facts) @@ -308,7 +319,11 @@ class LandTextInquiryFlow(SeaTextInquiryFlow): ) or "" self._wait_inquiry_sent(str(out_id), work_order_no=no) - tms = self._ledger.query_tms(work_order_no=no, facts=facts, business_line="LAND") + from agent.policy.system_exception import query_tms_resilient + + tms = query_tms_resilient( + self._ledger, work_order_no=no, facts=facts, business_line="LAND" + ) sess = FlowSession( sender_id=sender_id, thread_id=sender_id, @@ -335,7 +350,11 @@ class LandTextInquiryFlow(SeaTextInquiryFlow): reply(copy.LAND_CONFIRM_EXPIRED) return "land_confirm_expired" reply(copy.LAND_TMS_RETRYING.format(wo=no)) - tms = self._ledger.query_tms(work_order_no=no, facts=facts, business_line="LAND") + from agent.policy.system_exception import query_tms_resilient + + tms = query_tms_resilient( + self._ledger, work_order_no=no, facts=facts, business_line="LAND" + ) return self._apply_land_tms(sess=sess, tms=tms, reply=reply, retrying=True) def _apply_land_tms( @@ -350,6 +369,18 @@ class LandTextInquiryFlow(SeaTextInquiryFlow): _ = retrying no = (sess.work_order_no or "").strip() if not tms.get("ok"): + from agent.policy.system_exception import confirm_tms_quote_if_tech + + if confirm_tms_quote_if_tech( + ledger=self._ledger, + reply=reply, + work_order_no=no, + tms=tms, + sender_id=sess.sender_id, + ticket_status="询价中", + ): + self.drop_current_bookmark(sess.sender_id) + return "system_exception" reply("TMS 查询失败,请稍后重试。工单已创建:" + no) sess.phase = "tms_tech" self._save(sess) diff --git a/inquiry-agent/agent/policy/sea_text_flow.py b/inquiry-agent/agent/policy/sea_text_flow.py index a4e52bb..631c306 100644 --- a/inquiry-agent/agent/policy/sea_text_flow.py +++ b/inquiry-agent/agent/policy/sea_text_flow.py @@ -113,6 +113,17 @@ class SeaTextInquiryFlow(AirTextInquiryFlow): allow_b=allow_b, chat_fn=chat_fn, ) + from agent.policy.system_exception import handle_extract_tech_fail + + llm_phase = handle_extract_tech_fail( + snap=snap, + ledger=self._ledger, + reply=reply, + sender_id=sender_id, + sess=sess, + ) + if llm_phase: + return llm_phase mode = (snap.get("business_line") or "").upper() facts = dict(snap.get("facts") or {}) if sess and sess.phase in {"clarify", "need_mode"}: @@ -263,7 +274,14 @@ class SeaTextInquiryFlow(AirTextInquiryFlow): self._wait_inquiry_sent(str(out_id), work_order_no=no) tms_facts = dict(assemble_sea_query(facts).get("tms_facts") or facts) - tms = self._ledger.query_tms(work_order_no=no, facts=tms_facts, business_line="SEA") + from agent.policy.system_exception import ( + confirm_tms_quote_if_tech, + query_tms_resilient, + ) + + tms = query_tms_resilient( + self._ledger, work_order_no=no, facts=tms_facts, business_line="SEA" + ) sess = FlowSession( sender_id=sender_id, thread_id=sender_id, @@ -276,6 +294,16 @@ class SeaTextInquiryFlow(AirTextInquiryFlow): wait_version=1, ) if not tms.get("ok"): + if confirm_tms_quote_if_tech( + ledger=self._ledger, + reply=reply, + work_order_no=no, + tms=tms, + sender_id=sender_id, + ticket_status="询价中", + ): + self.drop_current_bookmark(sender_id) + return "system_exception" if tms.get("classification") == "TMS_QUERY_INCOMPLETE": logger.warning("tms.query 未出站 wo=%s class=%s", no, tms.get("classification")) reply("TMS 查询失败,请稍后重试。工单已创建:" + no) @@ -549,9 +577,29 @@ class SeaTextInquiryFlow(AirTextInquiryFlow): userids.append(uid) client = self.group_client() - created = client.create_group(name=no, userids=userids) or {} + from agent.policy.system_exception import call_with_retries, confirm_group_if_tech + + 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 ""))) + getter = getattr(self._ledger, "get_ticket", None) + st = "" + if callable(getter): + tic = getter(work_order_no=no) + st = str(getattr(tic, "status", "") or (tic.get("status") if isinstance(tic, dict) else "") or "") + confirm_group_if_tech( + ledger=self._ledger, + reply=reply, + work_order_no=no, + sender_id=sess.sender_id, + missing_names=product_names, + reason=str(created.get("error") or "企微建群失败"), + error_code=str(created.get("error") or "GROUP_CREATE"), + ticket_status=st, + extra={"userids": userids, "product_names": product_names}, + ) return sess.phase or "wait_collab" chat_id = str(created.get("chat_id") or "").strip() @@ -707,9 +755,32 @@ class SeaTextInquiryFlow(AirTextInquiryFlow): if uid not in userids: userids.append(uid) client = self.group_client() - created = client.create_group(name=display_no, userids=userids) or {} + from agent.policy.system_exception import call_with_retries, confirm_group_if_tech + + created = call_with_retries( + lambda: client.create_group(name=display_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 ""))) + real = (sess.work_order_no or "").strip() + mark_no = real if real.upper().startswith("WO") else "" + st = "" + if mark_no: + getter = getattr(self._ledger, "get_ticket", None) + if callable(getter): + tic = getter(work_order_no=mark_no) + st = str(getattr(tic, "status", "") or (tic.get("status") if isinstance(tic, dict) else "") or "") + confirm_group_if_tech( + ledger=self._ledger, + reply=reply, + work_order_no=mark_no, + sender_id=sess.sender_id, + missing_names=product_names, + reason=str(created.get("error") or "企微建群失败"), + error_code=str(created.get("error") or "GROUP_CREATE"), + ticket_status=st, + extra={"userids": userids, "product_names": product_names, "display_no": display_no}, + ) return sess.phase or "handoff_offer" chat_id = str(created.get("chat_id") or "").strip() brief = copy.handoff_group_brief( diff --git a/inquiry-agent/agent/policy/system_exception.py b/inquiry-agent/agent/policy/system_exception.py new file mode 100644 index 0000000..0071c94 --- /dev/null +++ b/inquiry-agent/agent/policy/system_exception.py @@ -0,0 +1,816 @@ +""" +系统异常(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} diff --git a/inquiry-agent/agent/routing/dispatch.py b/inquiry-agent/agent/routing/dispatch.py index b328718..15e2e4e 100644 --- a/inquiry-agent/agent/routing/dispatch.py +++ b/inquiry-agent/agent/routing/dispatch.py @@ -123,6 +123,36 @@ def _wrap_route_failed(message: InboundMessage) -> None: ) +def _wrap_route_model_failed(message: InboundMessage) -> None: + """路由模型超时/报错/空结果:大模型异常(无工单不建单)。""" + from agent.policy.system_exception import confirm_llm_if_tech + from agent.ledger.http_ledger import HttpLedger + from agent.config import get_settings + from agent.channel.queue import get_message_store + + def reply(text: str) -> None: + get_message_store().enqueue_outbound( + touser=message.sender_id, + content=text, + dedupe_key=f"route-llm:{message.message_id}", + ) + + ledger = None + try: + if (get_settings().ledger_backend or "").lower() != "memory": + ledger = HttpLedger() + except Exception: + ledger = None + confirm_llm_if_tech( + ledger=ledger, + reply=reply, + sender_id=message.sender_id, + reason="模型返回空结果", + conversation_kind="group" if getattr(message, "chat_id", "") else "private", + conversation_target=getattr(message, "chat_id", "") or message.sender_id, + ) + + # 意图 → handler(空壳或联调 echo);扩展时只改表,不改 ChannelRuntime 分支森林 HANDLER_REGISTRY: dict[str, HandlerFn] = { "echo_text": _wrap_echo, @@ -142,7 +172,7 @@ HANDLER_REGISTRY: dict[str, HandlerFn] = { "reject_no_sender": _noop, "ignore_empty": _noop, "stub_needs_model": _wrap_echo, - "route_model_failed": _wrap_route_failed, + "route_model_failed": _wrap_route_model_failed, "route_parse_failed": _wrap_route_failed, } diff --git a/inquiry-agent/tests/test_land_dispatch.py b/inquiry-agent/tests/test_land_dispatch.py index 2d8ec6d..3545639 100644 --- a/inquiry-agent/tests/test_land_dispatch.py +++ b/inquiry-agent/tests/test_land_dispatch.py @@ -183,27 +183,31 @@ class LandDispatchTests(unittest.TestCase): last = list(store._outbox.values())[-1].content self.assertNotIn("你好,我是询价机器人", last) - def test_tms_retry_misrouted_as_other(self) -> None: + def test_tms_tech_fail_pauses_until_admin_retry(self) -> None: + from agent.policy.inquiry_copy import SYS_EXC_TMS_FALLBACK from tests.test_land_text_inquiry import CountingLedger - class Flaky(CountingLedger): + class AlwaysFail(CountingLedger): def __init__(self) -> None: super().__init__() self.query_calls = 0 def query_tms(self, **kwargs): self.query_calls += 1 - if self.query_calls == 1: - return { - "ok": False, - "classification": "TMS_TECH_FAILURE", - "has_price": False, - "quote": {}, - } - return super().query_tms(**kwargs) + return { + "ok": False, + "classification": "TMS_TECH_FAILURE", + "has_price": False, + "quote": {}, + } land = reset_land_text_flow_for_test() - land.ledger = Flaky() + ledger = AlwaysFail() + land.ledger = ledger + get_sea_text_flow().ledger = ledger + from agent.policy.air_text_flow import get_air_text_flow + + get_air_text_flow().ledger = ledger facts = { "起运港": "广州", "目的港": "深圳", @@ -231,18 +235,27 @@ class LandDispatchTests(unittest.TestCase): InboundMessage(sender_id="ld5", message_id="m-l5c", content="确定"), store=self.store, ) - self.assertEqual(get_land_text_flow().session_of("ld5").phase, "tms_tech") - wo = get_land_text_flow().session_of("ld5").work_order_no + self.assertGreaterEqual(ledger.query_calls, 3) + wo = ledger.last_created_no + ticket = ledger.get_ticket(work_order_no=wo) + self.assertEqual(ticket.system_exception, "TMS异常") + bodies = "\n".join(i.content for i in self.store._outbox.values()) + self.assertIn(SYS_EXC_TMS_FALLBACK, bodies) + self.assertIsNone(get_land_text_flow().session_of("ld5")) name = dispatch_inbound( InboundMessage(sender_id="ld5", message_id="m-l5d", content="确定"), RouteDecision(intent="ordinary_text_other"), ) - self.assertEqual(name, "ordinary_text_inquiry") + self.assertEqual(name, "ordinary_text_other") sess = get_land_text_flow().session_of("ld5") - self.assertEqual(sess.phase, "wait_collab") - self.assertEqual(sess.work_order_no, wo) + if sess is not None: + self.assertNotEqual(sess.phase, "wait_collab") + handle_text_inquiry( + InboundMessage(sender_id="ld5", message_id="m-l5e", content=f"看看 {wo}"), + store=self.store, + ) last = list(self.store._outbox.values())[-1].content - self.assertNotIn("你好,我是询价机器人", last) + self.assertIn(SYS_EXC_TMS_FALLBACK, last) def test_sea_still_creates_immediately(self) -> None: msg = InboundMessage(sender_id="s-sea", message_id="m-sea-keep", content="海运询价") diff --git a/inquiry-agent/tests/test_land_text_inquiry.py b/inquiry-agent/tests/test_land_text_inquiry.py index d395926..bedb46b 100644 --- a/inquiry-agent/tests/test_land_text_inquiry.py +++ b/inquiry-agent/tests/test_land_text_inquiry.py @@ -333,13 +333,7 @@ class LandTextConfirmTests(unittest.TestCase): injected_mode="LAND", ) phase = flow.on_text(sender_id="u1", text="确定", reply=reply) - self.assertEqual(phase, "tms_tech") - self.assertEqual(ledger.create_calls, 1) wo = ledger.last_created_no - replies.clear() - flow._by_sender.clear() - flow._by_ticket.clear() - phase = flow.on_text(sender_id="u1", text="确定", reply=reply) self.assertEqual(phase, "wait_collab") self.assertEqual(ledger.create_calls, 1) self.assertEqual(ledger.query_calls, 2) diff --git a/inquiry-agent/tests/test_sea_text_inquiry.py b/inquiry-agent/tests/test_sea_text_inquiry.py index 72e9d42..e124309 100644 --- a/inquiry-agent/tests/test_sea_text_inquiry.py +++ b/inquiry-agent/tests/test_sea_text_inquiry.py @@ -622,9 +622,12 @@ class SeaTmsCardTests(unittest.TestCase): injected_facts=facts, injected_mode="SEA", ) - self.assertEqual(phase, "tms_tech") - self.assertIn("TMS 查询失败", "\n".join(self.replies)) + self.assertEqual(phase, "system_exception") + self.assertIn("TMS接口异常,IT运维排查中", "\n".join(self.replies)) self.assertNotIn(copy.BTN_PULL_COLLAB, "\n".join(self.replies)) + tickets = list(self.flow._ledger._tickets.values()) + self.assertTrue(tickets) + self.assertEqual(tickets[-1].system_exception, "TMS异常") class SeaSkipAndPullTests(unittest.TestCase): diff --git a/inquiry-agent/tests/test_system_exception.py b/inquiry-agent/tests/test_system_exception.py new file mode 100644 index 0000000..d343e2a --- /dev/null +++ b/inquiry-agent/tests/test_system_exception.py @@ -0,0 +1,294 @@ +""" +系统异常:话术、判定边界、无单不落库、暂停、告警卡、重试清列。 +""" + +from __future__ import annotations + +import unittest +from datetime import datetime +from zoneinfo import ZoneInfo + +from agent.ledger.memory_ledger import MemoryLedger +from agent.policy import inquiry_copy as copy +from agent.policy.system_exception import ( + STEP_GROUP, + STEP_LLM, + STEP_LLM_NO_TICKET, + STEP_TMS_CABIN, + STEP_TMS_QUOTE, + TYPE_GROUP, + TYPE_LLM, + TYPE_TMS, + SystemExceptionEvent, + call_with_retries, + confirm_system_exception, + extract_work_order_nos, + is_tms_tech_failure, + should_mark_ticket, + ticket_is_paused, +) + + +class TestSystemExceptionCopy(unittest.TestCase): + def test_fallback_three_types(self) -> None: + self.assertEqual(copy.system_exception_fallback(TYPE_LLM), "大模型服务异常,IT排查中...") + self.assertEqual(copy.system_exception_fallback(TYPE_TMS), "TMS接口异常,IT运维排查中...") + self.assertEqual( + copy.system_exception_fallback(TYPE_GROUP, names=["赵六"]), + "无法添加 赵六 同事,请联系管理员检查企业微信配置。", + ) + self.assertIn("李四、王五", copy.system_exception_fallback(TYPE_GROUP, names=["李四", "王五"])) + + def test_recovery_four_steps(self) -> None: + self.assertEqual(copy.system_exception_recovery(STEP_TMS_QUOTE), copy.SYS_EXC_RECOVERY_TMS) + self.assertEqual(copy.system_exception_recovery(STEP_LLM), copy.SYS_EXC_RECOVERY_LLM) + self.assertEqual(copy.system_exception_recovery(STEP_TMS_CABIN), copy.SYS_EXC_RECOVERY_CABIN) + self.assertEqual(copy.system_exception_recovery(STEP_GROUP), copy.SYS_EXC_RECOVERY_GROUP) + + def test_impact_by_step(self) -> None: + self.assertEqual(copy.system_exception_impact(STEP_LLM_NO_TICKET), copy.IMPACT_LLM_NO_TICKET) + self.assertEqual(copy.system_exception_impact(STEP_LLM), copy.IMPACT_LLM) + self.assertEqual(copy.system_exception_impact(STEP_TMS_QUOTE), copy.IMPACT_TMS_QUOTE) + self.assertEqual(copy.system_exception_impact(STEP_TMS_CABIN), copy.IMPACT_TMS_CABIN) + self.assertEqual(copy.system_exception_impact(STEP_GROUP), copy.IMPACT_GROUP) + + +class TestSystemExceptionRules(unittest.TestCase): + def test_tms_no_price_is_not_tech(self) -> None: + self.assertFalse( + is_tms_tech_failure( + {"ok": True, "classification": "TMS_NO_QUOTE_1002", "has_price": False} + ) + ) + self.assertFalse( + is_tms_tech_failure({"ok": False, "classification": "TMS_QUERY_INCOMPLETE"}) + ) + self.assertTrue( + is_tms_tech_failure({"ok": False, "classification": "TMS_TECH_FAILURE"}) + ) + self.assertTrue(is_tms_tech_failure({"ok": False, "error": "timeout"})) + + def test_closed_or_blank_not_marked(self) -> None: + self.assertFalse(should_mark_ticket(work_order_no="", ticket_status="")) + self.assertFalse(should_mark_ticket(work_order_no="WO1", ticket_status="已关闭")) + self.assertTrue(should_mark_ticket(work_order_no="WO1", ticket_status="询价中")) + self.assertTrue(should_mark_ticket(work_order_no="WO1", ticket_status="转人工")) + self.assertTrue(should_mark_ticket(work_order_no="WO1", ticket_status="已成交")) + + def test_extract_work_order_nos(self) -> None: + self.assertEqual(extract_work_order_nos("请看 WO202609210001"), ["WO202609210001"]) + self.assertEqual(extract_work_order_nos("虚拟 VT202609210001 不算"), []) + + def test_retries_succeed_on_second(self) -> None: + hits = {"n": 0} + + def flaky() -> dict: + hits["n"] += 1 + if hits["n"] < 2: + return {"ok": False} + return {"ok": True, "value": 1} + + out = call_with_retries(flaky, is_ok=lambda x: bool(x.get("ok")), attempts=3) + self.assertTrue(out.get("ok")) + self.assertEqual(hits["n"], 2) + + def test_retries_keep_last_failure(self) -> None: + def always_fail() -> dict: + return {"ok": False, "classification": "TMS_TECH_FAILURE"} + + out = call_with_retries(always_fail, is_ok=lambda x: bool(x.get("ok")), attempts=3) + self.assertFalse(out.get("ok")) + self.assertEqual(out.get("classification"), "TMS_TECH_FAILURE") + + +class TestConfirmSystemException(unittest.TestCase): + def setUp(self) -> None: + self.ledger = MemoryLedger() + created = self.ledger.create_ticket( + sender_id="sales1", + business_line="SEA", + facts={"起运港": "上海", "目的港": "洛杉矶", "品名": "衣服"}, + ) + self.wo = str(created.get("work_order_no") or "") + self.ledger.upsert_exception_notify_user({"name": "运维甲", "wecomId": "ops-a"}) + + def test_no_ticket_replies_and_alerts_but_does_not_create(self) -> None: + replies: list[str] = [] + queued: list[dict] = [] + + def enqueue(*, touser, content, dedupe_key, payload): + queued.append({"touser": touser, "payload": payload, "content": content}) + return True, "ok" + + before = len(self.ledger._tickets) + out = confirm_system_exception( + event=SystemExceptionEvent( + exception_type=TYPE_LLM, + step=STEP_LLM_NO_TICKET, + reason="模型返回空结果", + service="大模型识别", + error_code="EMPTY", + ), + ledger=self.ledger, + reply=replies.append, + enqueue=enqueue, + notify_ready=True, + now=datetime(2026, 9, 21, 10, 0, 0, tzinfo=ZoneInfo("Asia/Shanghai")), + ) + self.assertTrue(out.get("ok")) + self.assertFalse(out.get("marked")) + self.assertEqual(len(self.ledger._tickets), before) + self.assertEqual(replies[0], copy.SYS_EXC_LLM_FALLBACK) + self.assertEqual(queued[0]["touser"], "ops-a") + card = queued[0]["payload"]["template_card"] + self.assertEqual(card["source"]["desc"], "询价机器人系统告警") + self.assertEqual(card["main_title"]["desc"], "关联工单 无") + values = {row["keyname"]: row["value"] for row in card["horizontal_content_list"]} + self.assertEqual(values["影响范围"], copy.IMPACT_LLM_NO_TICKET) + + def test_mark_ticket_and_pause(self) -> None: + replies: list[str] = [] + out = confirm_system_exception( + event=SystemExceptionEvent( + exception_type=TYPE_TMS, + step=STEP_TMS_QUOTE, + reason="TMS /quote/v2 接口调用失败", + service="TMS报价", + error_code="QUOTE_CALCULATION", + work_order_no=self.wo, + ticket_status="询价中", + ), + ledger=self.ledger, + reply=replies.append, + enqueue=lambda **kwargs: (True, "ok"), + notify_ready=False, + ) + self.assertTrue(out.get("marked")) + ticket = self.ledger.get_ticket(work_order_no=self.wo) + self.assertEqual(ticket.status, "询价中") + self.assertEqual(ticket.system_exception, TYPE_TMS) + self.assertEqual(ticket.system_exception_reason, "TMS /quote/v2 接口调用失败") + self.assertTrue(ticket_is_paused(ticket)) + self.assertEqual(replies[0], copy.SYS_EXC_TMS_FALLBACK) + + def test_closed_ticket_not_marked(self) -> None: + self.ledger.transition(work_order_no=self.wo, to_status="已关闭") + out = confirm_system_exception( + event=SystemExceptionEvent( + exception_type=TYPE_TMS, + step=STEP_TMS_QUOTE, + reason="x", + service="TMS报价", + error_code="TIMEOUT", + work_order_no=self.wo, + ticket_status="已关闭", + ), + ledger=self.ledger, + reply=lambda _t: None, + notify_ready=False, + ) + self.assertFalse(out.get("marked")) + ticket = self.ledger.get_ticket(work_order_no=self.wo) + self.assertFalse(ticket_is_paused(ticket)) + + def test_later_ticket_does_not_inherit_pre_ticket_exception(self) -> None: + confirm_system_exception( + event=SystemExceptionEvent( + exception_type=TYPE_LLM, + step=STEP_LLM_NO_TICKET, + reason="空结果", + service="大模型识别", + error_code="EMPTY", + ), + ledger=self.ledger, + reply=lambda _t: None, + notify_ready=False, + ) + ticket = self.ledger.get_ticket(work_order_no=self.wo) + self.assertFalse(ticket_is_paused(ticket)) + + def test_admin_retry_tms_success_clears(self) -> None: + self.ledger.mark_system_exception( + work_order_no=self.wo, + system_exception=TYPE_TMS, + system_exception_reason="TMS /quote/v2 接口调用失败", + step=STEP_TMS_QUOTE, + payload={"conversation_kind": "private", "conversation_target": "sales1"}, + ) + from agent.policy.system_exception import run_admin_retry + + talks: list[str] = [] + + def fake_send(*, kind, target, text, extra=None): + talks.append(text) + + import agent.policy.system_exception as se + + orig = se._send_to_conversation + se._send_to_conversation = fake_send + try: + out = run_admin_retry(self.wo, ledger=self.ledger) + finally: + se._send_to_conversation = orig + self.assertTrue(out.get("ok")) + ticket = self.ledger.get_ticket(work_order_no=self.wo) + self.assertFalse(ticket_is_paused(ticket)) + self.assertEqual(ticket.status, "已报价") + self.assertTrue(ticket.quote) + self.assertTrue(any("TMS已恢复" in x for x in talks)) + joined = "\n".join(talks) + self.assertTrue("USD" in joined or "拉产品" in joined or "已报价" in joined) + + def test_admin_retry_cabin_lock_runs_lock(self) -> None: + self.ledger.mark_system_exception( + work_order_no=self.wo, + system_exception=TYPE_TMS, + system_exception_reason="空运锁舱接口调用失败", + step=STEP_TMS_CABIN, + payload={ + "conversation_kind": "group", + "conversation_target": "chat-air", + "extra": {"option_no": "AIR-OPT-01", "instruction": "锁舱"}, + }, + ) + hits = {"n": 0} + orig_lock = self.ledger.lock_cabin + + def counted(**kwargs): + hits["n"] += 1 + return orig_lock(**kwargs) + + self.ledger.lock_cabin = counted # type: ignore[method-assign] + from agent.policy.system_exception import run_admin_retry + + talks: list[str] = [] + import agent.policy.system_exception as se + + orig = se._send_to_conversation + se._send_to_conversation = lambda **kwargs: talks.append(str(kwargs.get("text") or "")) + try: + out = run_admin_retry(self.wo, ledger=self.ledger) + finally: + se._send_to_conversation = orig + self.assertTrue(out.get("ok")) + self.assertGreaterEqual(hits["n"], 1) + ticket = self.ledger.get_ticket(work_order_no=self.wo) + self.assertFalse(ticket_is_paused(ticket)) + self.assertTrue(any("空运舱位接口已恢复" in x for x in talks)) + + def test_extract_tech_fail_when_model_errors(self) -> None: + from agent.llm.extract_text import extract_inquiry_snapshot + + class Boom: + ok = False + error = "timeout" + tool_calls = [] + content = "" + + snap = extract_inquiry_snapshot("上海到洛杉矶衣服", allow_b=True, chat_fn=lambda **_k: Boom()) + self.assertTrue(snap.get("tech_fail")) + self.assertEqual(snap.get("tech_error"), "timeout") + + +if __name__ == "__main__": + unittest.main() diff --git a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/controller/InquiryAgentController.java b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/controller/InquiryAgentController.java index a5722eb..4380ead 100644 --- a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/controller/InquiryAgentController.java +++ b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/controller/InquiryAgentController.java @@ -204,6 +204,37 @@ public class InquiryAgentController { return Result.OK(rendered); } + @ApiOperation("智能体写入系统异常两列(已关闭拒绝,不改六态)") + @PostMapping("/ticket/markSystemException") + public Result> markSystemException(@RequestBody Map body, + HttpServletRequest http) { + if (!authed(http)) { + return Result.error(401, "unauthorized"); + } + return Result.OK(workOrderService.markSystemException(body)); + } + + @ApiOperation("智能体清除系统异常两列(重试成功后)") + @PostMapping("/ticket/clearSystemException") + public Result> clearSystemException(@RequestBody Map body, + HttpServletRequest http) { + if (!authed(http)) { + return Result.error(401, "unauthorized"); + } + return Result.OK(workOrderService.clearSystemException(firstStr(body, "workOrderNo", "work_order_no"))); + } + + @ApiOperation("系统异常告警接收人:角色勾了通知且账号已绑企微") + @PostMapping("/account/listExceptionNotify") + public Result> listExceptionNotify(HttpServletRequest http) { + if (!authed(http)) { + return Result.error(401, "unauthorized"); + } + Map out = new java.util.LinkedHashMap<>(); + out.put("staff", workOrderService.listExceptionNotifyUsers()); + return Result.OK(out); + } + @ApiOperation("经主账查真 TMS(有价回报价 / 无价回暂无报价 / 技术失败)") @PostMapping("/tms/query") public Result> tmsQuery(@RequestBody Map body, diff --git a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/entity/InquiryTicket.java b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/entity/InquiryTicket.java index 844c99a..3380714 100644 --- a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/entity/InquiryTicket.java +++ b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/entity/InquiryTicket.java @@ -122,6 +122,19 @@ public class InquiryTicket implements Serializable { @ApiModelProperty("系统异常原因") private String systemExceptionReason; + /** + * 后台点重试时重跑哪一步:tms_quote / tms_cabin / llm / group。 + * 不改六态;成功后与两列一并清空。 + */ + @ApiModelProperty("系统异常步骤") + private String systemExceptionStep; + + /** + * 重试上下文 JSON:会话目标、未加上的人、业务线等。 + */ + @ApiModelProperty("系统异常重试上下文") + private String systemExceptionPayload; + @ApiModelProperty("关单原因") private String closeReason; diff --git a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/mapper/InquiryExceptionNotifyMapper.java b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/mapper/InquiryExceptionNotifyMapper.java new file mode 100644 index 0000000..c47df0b --- /dev/null +++ b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/mapper/InquiryExceptionNotifyMapper.java @@ -0,0 +1,29 @@ +package org.jeecg.modules.inquiry.mapper; + +import org.apache.ibatis.annotations.Mapper; +import org.apache.ibatis.annotations.Select; + +import java.util.List; +import java.util.Map; + +/** + * 系统异常告警接收人:后台账号角色勾了「系统异常消息通知」,且工号(企微 userid)非空。 + *

+ * 不按海陆空过滤。inquiry 模块不依赖 system-biz,直接读同库 Jeecg 表。 + */ +@Mapper +public interface InquiryExceptionNotifyMapper { + + /** + * 启用账号 + 角色 description 含系统异常通知稳定码 + work_no 已绑企微。 + */ + @Select("SELECT DISTINCT u.work_no AS wecomId, u.realname AS name " + + "FROM sys_role r " + + "INNER JOIN sys_user_role ur ON ur.role_id = r.id " + + "INNER JOIN sys_user u ON u.id = ur.user_id " + + "WHERE u.del_flag = 0 AND u.status = 1 " + + "AND u.work_no IS NOT NULL AND TRIM(u.work_no) <> '' " + + "AND (r.description LIKE '%系统异常·消息通知%' " + + "OR r.description LIKE '%系统异常消息通知%')") + List> listBoundAccounts(); +} diff --git a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/IInquiryWorkOrderService.java b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/IInquiryWorkOrderService.java index 29a5637..506aa06 100644 --- a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/IInquiryWorkOrderService.java +++ b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/IInquiryWorkOrderService.java @@ -41,10 +41,25 @@ public interface IInquiryWorkOrderService extends IService { void handleAbnormal(TicketAbnormalRequest req, String operator); /** - * 系统异常重试:累加 retry_count,记日志;真正重跑任务由 Worker 后续认领(本期只落主账状态)。 + * 系统异常重试:已关闭拒绝;记下重试意图后唤醒智能体(快回,不在本线程重跑)。 */ void retryAbnormal(TicketAbnormalRequest req, String operator); + /** + * 智能体写入系统异常两列与重试上下文。已关闭拒绝。不改六态。 + */ + Map markSystemException(Map body); + + /** + * 重试成功后清两列。 + */ + Map clearSystemException(String workOrderNo); + + /** + * 后台账号:角色勾了系统异常消息通知,且工号已绑企微。 + */ + java.util.List> listExceptionNotifyUsers(); + /** * 智能体建单:每次出询价确认都新开一单,相同需求不复用未关闭工单。 */ diff --git a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/impl/InquiryWorkOrderServiceImpl.java b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/impl/InquiryWorkOrderServiceImpl.java index d0055ca..60d645e 100644 --- a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/impl/InquiryWorkOrderServiceImpl.java +++ b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/service/impl/InquiryWorkOrderServiceImpl.java @@ -20,8 +20,10 @@ import org.jeecg.modules.inquiry.mapper.InquiryAttachmentMapper; import org.jeecg.modules.inquiry.mapper.InquiryExceptionCaseMapper; import org.jeecg.modules.inquiry.mapper.InquiryQuoteMapper; import org.jeecg.modules.inquiry.mapper.InquiryStaffRoleMapper; +import org.jeecg.modules.inquiry.mapper.InquiryExceptionNotifyMapper; import org.jeecg.modules.inquiry.mapper.InquiryTicketMapper; import org.jeecg.modules.inquiry.mapper.InquiryTicketStatusLogMapper; +import org.jeecg.modules.inquiry.support.AgentSystemExceptionRetryClient; import org.jeecg.modules.inquiry.service.IInquiryWorkOrderService; import org.jeecg.modules.inquiry.support.InquiryQuoteDisplaySupport; import org.jeecg.modules.inquiry.support.InquiryTicketDisplaySupport; @@ -69,6 +71,10 @@ public class InquiryWorkOrderServiceImpl extends ServiceImpl pageView(Page page, Map filters) { @@ -300,6 +306,8 @@ public class InquiryWorkOrderServiceImpl extends ServiceImpl openCases = exceptionCaseMapper.selectList(new LambdaQueryWrapper() .eq(InquiryExceptionCase::getWorkOrderNo, ticket.getWorkOrderNo()) .eq(InquiryExceptionCase::getCaseType, "SYSTEM") @@ -339,6 +353,110 @@ public class InquiryWorkOrderServiceImpl extends ServiceImpl markSystemException(Map body) { + Map out = new LinkedHashMap<>(); + String no = firstBody(body, "workOrderNo", "work_order_no"); + InquiryTicket ticket = findByWorkOrderNo(no); + if (ticket == null) { + out.put("ok", false); + out.put("error", "not_found"); + return out; + } + if (TicketStatuses.CLOSED.equals(ticket.getStatus())) { + out.put("ok", false); + out.put("error", "closed"); + return out; + } + String type = firstBody(body, "systemException", "system_exception"); + String reason = firstBody(body, "systemExceptionReason", "system_exception_reason"); + String step = firstBody(body, "step", "systemExceptionStep", "system_exception_step"); + Object payloadObj = body == null ? null : body.get("payload"); + if (payloadObj == null && body != null) { + payloadObj = body.get("systemExceptionPayload"); + } + String payloadJson = payloadObj == null ? null + : (payloadObj instanceof String ? (String) payloadObj : JSON.toJSONString(payloadObj)); + ticket.setSystemException(type); + ticket.setSystemExceptionReason(reason); + ticket.setSystemExceptionStep(step); + ticket.setSystemExceptionPayload(payloadJson); + ticket.setUpdateTime(new Date()); + updateById(ticket); + InquiryExceptionCase created = new InquiryExceptionCase(); + created.setTicketId(ticket.getId()); + created.setWorkOrderNo(ticket.getWorkOrderNo()); + created.setCaseType("SYSTEM"); + created.setTitle(StringUtils.defaultIfBlank(type, "系统异常")); + created.setReason(reason); + created.setStatus("OPEN"); + created.setRetryCount(0); + created.setCreateTime(new Date()); + created.setUpdateTime(new Date()); + exceptionCaseMapper.insert(created); + appendLog(ticket, ticket.getStatus(), ticket.getStatus(), "system_exception", "agent", reason); + out.put("ok", true); + out.put("status", ticket.getStatus()); + return out; + } + + @Override + @Transactional(rollbackFor = Exception.class) + public Map clearSystemException(String workOrderNo) { + Map out = new LinkedHashMap<>(); + InquiryTicket ticket = findByWorkOrderNo(StringUtils.trimToEmpty(workOrderNo)); + if (ticket == null) { + out.put("ok", false); + out.put("error", "not_found"); + return out; + } + ticket.setSystemException(null); + ticket.setSystemExceptionReason(null); + ticket.setSystemExceptionStep(null); + ticket.setSystemExceptionPayload(null); + ticket.setUpdateTime(new Date()); + updateById(ticket); + List openCases = exceptionCaseMapper.selectList(new LambdaQueryWrapper() + .eq(InquiryExceptionCase::getWorkOrderNo, ticket.getWorkOrderNo()) + .eq(InquiryExceptionCase::getCaseType, "SYSTEM") + .in(InquiryExceptionCase::getStatus, "OPEN", "RETRYING")); + Date now = new Date(); + for (InquiryExceptionCase c : openCases) { + c.setStatus("HANDLED"); + c.setHandledBy("agent"); + c.setHandledAt(now); + c.setUpdateTime(now); + exceptionCaseMapper.updateById(c); + } + appendLog(ticket, ticket.getStatus(), ticket.getStatus(), "system_exception_clear", "agent", "重试成功"); + out.put("ok", true); + return out; + } + + @Override + public List> listExceptionNotifyUsers() { + List> rows = exceptionNotifyMapper.listBoundAccounts(); + return rows == null ? Collections.emptyList() : rows; + } + + private static String firstBody(Map body, String... keys) { + if (body == null) { + return ""; + } + for (String key : keys) { + Object v = body.get(key); + if (v != null && StringUtils.isNotBlank(String.valueOf(v))) { + return String.valueOf(v).trim(); + } + } + return ""; } /** @@ -699,6 +817,19 @@ public class InquiryWorkOrderServiceImpl extends ServiceImpl excPayload = new LinkedHashMap<>(); + if (StringUtils.isNotBlank(ticket.getSystemExceptionPayload())) { + JSONObject parsed = JSON.parseObject(ticket.getSystemExceptionPayload()); + if (parsed != null) { + excPayload.putAll(parsed); + } + } + out.put("system_exception_payload", excPayload); return out; } diff --git a/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/support/AgentSystemExceptionRetryClient.java b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/support/AgentSystemExceptionRetryClient.java new file mode 100644 index 0000000..9d33efa --- /dev/null +++ b/inquiry-api/jeecg-module-inquiry/src/main/java/org/jeecg/modules/inquiry/support/AgentSystemExceptionRetryClient.java @@ -0,0 +1,122 @@ +package org.jeecg.modules.inquiry.support; + +import com.alibaba.fastjson.JSON; +import com.alibaba.fastjson.JSONObject; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; + +import java.io.OutputStream; +import java.net.HttpURLConnection; +import java.net.URL; +import java.nio.charset.StandardCharsets; +import java.util.LinkedHashMap; +import java.util.Map; + +/** + * 后台点「重试」后唤醒智能体:只入队,不等 TMS / 大模型 / 建群跑完。 + *

+ * 打 AGENT_HTTP_BASE_URL /internal/bridge/wake,kind=system_exception_retry。 + * 禁止在本类同步等 Python Graph。 + */ +@Slf4j +@Component +public class AgentSystemExceptionRetryClient { + + @Value("${inquiry.agent.base-url:${AGENT_HTTP_BASE_URL:http://127.0.0.1:8910}}") + private String agentBaseUrl; + + @Value("${inquiry.agent.callback-token:${AGENT_CALLBACK_TOKEN:}}") + private String agentToken; + + @Value("${inquiry.agent.retry-timeout-ms:8000}") + private int timeoutMs; + + /** + * 入队系统异常重试。返回 null 表示智能体已收下;非空为失败原因。 + */ + public String enqueueRetry(String workOrderNo) { + String no = StringUtils.trimToEmpty(workOrderNo); + if (StringUtils.isBlank(no)) { + return "工单号为空"; + } + String base = StringUtils.removeEnd(StringUtils.defaultIfBlank(agentBaseUrl, ""), "/"); + if (StringUtils.isBlank(base)) { + return "智能体地址未配置(AGENT_HTTP_BASE_URL)"; + } + String url = base + "/internal/bridge/wake"; + HttpURLConnection conn = null; + try { + conn = (HttpURLConnection) new URL(url).openConnection(); + conn.setRequestMethod("POST"); + conn.setConnectTimeout(Math.max(2000, timeoutMs)); + conn.setReadTimeout(Math.max(2000, timeoutMs)); + conn.setDoOutput(true); + conn.setRequestProperty("Content-Type", "application/json;charset=UTF-8"); + conn.setRequestProperty("Accept", "application/json"); + if (StringUtils.isNotBlank(agentToken)) { + conn.setRequestProperty("X-Agent-Callback-Token", agentToken); + } + Map body = new LinkedHashMap<>(); + body.put("thread_id", no); + body.put("inquiry_no", no); + body.put("kind", "system_exception_retry"); + body.put("quote_version", 0); + body.put("wait_version", 0); + byte[] bytes = JSON.toJSONBytes(body); + try (OutputStream out = conn.getOutputStream()) { + out.write(bytes); + } + int code = conn.getResponseCode(); + String respText = readQuietly(conn, code >= 200 && code < 300); + if (code == 202 || code == 200) { + JSONObject resp = safeJson(respText); + if (resp != null && "accepted".equals(resp.getString("status"))) { + return null; + } + if (resp != null && StringUtils.isNotBlank(resp.getString("entry_id"))) { + return null; + } + return "智能体未确认入队"; + } + return "无法提交重试: HTTP " + code; + } catch (Exception e) { + log.warn("系统异常重试唤醒失败 wo={} err={}", no, e.getMessage()); + return "无法提交重试: " + e.getMessage(); + } finally { + if (conn != null) { + conn.disconnect(); + } + } + } + + private static String readQuietly(HttpURLConnection conn, boolean ok) { + try { + java.io.InputStream in = ok ? conn.getInputStream() : conn.getErrorStream(); + if (in == null) { + return ""; + } + java.io.ByteArrayOutputStream buf = new java.io.ByteArrayOutputStream(); + byte[] chunk = new byte[512]; + int n; + while ((n = in.read(chunk)) > 0) { + buf.write(chunk, 0, n); + } + return new String(buf.toByteArray(), StandardCharsets.UTF_8); + } catch (Exception e) { + return ""; + } + } + + private static JSONObject safeJson(String text) { + if (StringUtils.isBlank(text)) { + return null; + } + try { + return JSON.parseObject(text); + } catch (Exception e) { + return null; + } + } +} diff --git a/inquiry-api/jeecg-module-inquiry/src/main/resources/db/migration_inq_ticket_system_exception_ctx_v1.sql b/inquiry-api/jeecg-module-inquiry/src/main/resources/db/migration_inq_ticket_system_exception_ctx_v1.sql new file mode 100644 index 0000000..1e3f3a2 --- /dev/null +++ b/inquiry-api/jeecg-module-inquiry/src/main/resources/db/migration_inq_ticket_system_exception_ctx_v1.sql @@ -0,0 +1,33 @@ +-- 系统异常重试上下文:步骤 + JSON。仅改 inquiry_robot。 +-- 两列展示字段 system_exception / system_exception_reason 已存在。 +-- 可重复执行:已有列则跳过。 + +SET NAMES utf8mb4; + +SET @db := DATABASE(); + +SET @exists_step := ( + SELECT COUNT(*) FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = @db AND TABLE_NAME = 'inq_ticket' AND COLUMN_NAME = 'system_exception_step' +); +SET @sql_step := IF( + @exists_step = 0, + 'ALTER TABLE `inq_ticket` ADD COLUMN `system_exception_step` varchar(64) NULL COMMENT ''系统异常步骤:tms_quote/tms_cabin/llm/group'' AFTER `system_exception_reason`', + 'SELECT 1' +); +PREPARE stmt_step FROM @sql_step; +EXECUTE stmt_step; +DEALLOCATE PREPARE stmt_step; + +SET @exists_payload := ( + SELECT COUNT(*) FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = @db AND TABLE_NAME = 'inq_ticket' AND COLUMN_NAME = 'system_exception_payload' +); +SET @sql_payload := IF( + @exists_payload = 0, + 'ALTER TABLE `inq_ticket` ADD COLUMN `system_exception_payload` text NULL COMMENT ''系统异常重试上下文 JSON'' AFTER `system_exception_step`', + 'SELECT 1' +); +PREPARE stmt_payload FROM @sql_payload; +EXECUTE stmt_payload; +DEALLOCATE PREPARE stmt_payload; diff --git a/inquiry-backend/src/views/inquiry/workOrder/index.vue b/inquiry-backend/src/views/inquiry/workOrder/index.vue index 81cdec2..23effa1 100644 --- a/inquiry-backend/src/views/inquiry/workOrder/index.vue +++ b/inquiry-backend/src/views/inquiry/workOrder/index.vue @@ -598,7 +598,7 @@ function refresh() {