Files
inquiry_robot/inquiry-agent/agent/jobs/event_exception_poll.py
T

329 lines
11 KiB
Python

"""
事件异常巡检:Worker HTTP 槽轮询到点工单并告警。
本文件职责:读主账等待、按工作日对照超时小时、会话喊人、可选应用卡片、再记摘要。
禁止:回调线程调用;占 LibreOffice 槽;改六态;反复催已告警的等待。
"""
from __future__ import annotations
import logging
from datetime import datetime
from typing import Any, Callable, Optional
from agent.config import Settings, get_settings
from agent.policy import inquiry_copy as copy
from agent.policy.deal_followup import SHANGHAI, workday_elapsed_hours
from agent.policy.event_exception import (
card_title,
conversation_of,
exception_type,
WAIT_ADOPT,
WAIT_AIR_QUOTE,
WAIT_DEAL,
WAIT_LOST_REASON,
mentions_for_role,
pick_app_targets,
pick_mentions,
timeout_code,
)
logger = logging.getLogger(__name__)
# 采用、成交跟进、未成交原因都是在等销售。点名只找本单销售。
_SALES_WAIT_KINDS = {WAIT_ADOPT, WAIT_DEAL, WAIT_LOST_REASON}
SendGroupFn = Callable[..., bool]
EnqueueFn = Callable[..., tuple[bool, str]]
TimeoutLookupFn = Callable[..., dict[str, Any]]
def _parse_started(raw: Any) -> Optional[datetime]:
text = str(raw or "").strip()
if not text:
return None
text = text.replace("T", " ")[:19]
try:
naive = datetime.strptime(text, "%Y-%m-%d %H:%M:%S")
return naive.replace(tzinfo=SHANGHAI)
except ValueError:
return None
def _hours_of(row: dict[str, Any]) -> float:
try:
return float(row.get("timeoutHours") if row.get("timeoutHours") is not None else row.get("timeout_hours") or 0)
except (TypeError, ValueError):
return 0.0
def poll_event_exceptions(
*,
ledger: Any = None,
settings: Optional[Settings] = None,
now: Optional[datetime] = None,
enqueue: Optional[EnqueueFn] = None,
send_group: Optional[SendGroupFn] = None,
timeout_lookup: Optional[TimeoutLookupFn] = None,
notify_ready: Optional[bool] = None,
) -> int:
"""
巡检未告警等待。返回本轮成功记下异常的工单数。
会话成功才 record;卡片失败不影响会话。禁止进企微回调线程。
"""
cfg = settings or get_settings()
book = ledger
if book is None:
if str(getattr(cfg, "ledger_backend", "http") or "http").lower() == "memory":
return 0
from agent.ledger.http_ledger import HttpLedger
book = HttpLedger()
lister = getattr(book, "list_event_waits", None)
if not callable(lister):
return 0
stamp = now or datetime.now(SHANGHAI)
if stamp.tzinfo is None:
stamp = stamp.replace(tzinfo=SHANGHAI)
# 测服会话/卡片也不加 [test],与销售询价正文一致。
test_prefix = False
lookup = timeout_lookup or getattr(book, "get_timeout_by_code", None)
ready = notify_app_flag(notify_ready)
sent = 0
for row in lister() or []:
if not isinstance(row, dict):
continue
try:
if _handle_one(
row,
book=book,
now=stamp,
test_prefix=test_prefix,
enqueue=enqueue,
send_group=send_group,
timeout_lookup=lookup,
notify_ready=ready,
settings=cfg,
):
sent += 1
except Exception:
logger.exception(
"事件异常巡检失败 wo=%s",
row.get("work_order_no") or row.get("workOrderNo"),
)
return sent
def notify_app_flag(override: Optional[bool]) -> bool:
if override is not None:
return bool(override)
from agent.policy.deal_outcome import notify_app_ready
return notify_app_ready()
def _handle_one(
row: dict[str, Any],
*,
book: Any,
now: datetime,
test_prefix: bool,
enqueue: Optional[EnqueueFn],
send_group: Optional[SendGroupFn],
timeout_lookup: Any,
notify_ready: bool,
settings: Settings,
) -> bool:
no = str(row.get("work_order_no") or row.get("workOrderNo") or "").strip()
kind = str(row.get("event_wait_kind") or row.get("eventWaitKind") or "").strip()
if not no or not kind:
return False
line = str(row.get("business_line") or row.get("businessLine") or "").strip()
facts = dict(row.get("facts") or {})
klass = str(facts.get("运输分类") or facts.get("transportType") or "").strip()
code = timeout_code(wait_kind=kind, business_line=line, transport_class=klass)
if not code or not callable(timeout_lookup):
return False
timeout_row = timeout_lookup(timeout_code=code) or {}
if not timeout_row.get("found"):
return False
hours = _hours_of(timeout_row)
if hours <= 0:
return False
started = _parse_started(row.get("event_wait_started_at") or row.get("eventWaitStartedAt"))
if started is None:
return False
if workday_elapsed_hours(started, now) < hours:
return False
watchers = [dict(x) for x in (row.get("watchers") or []) if isinstance(x, dict)]
sales = _ticket_sales(row, watchers)
people = pick_mentions(wait_kind=kind, sales=sales, members=watchers, business_line=line)
# 空运群往往只绑了群号,成员列表是空的,群里就找不到航线岗,话术会落成「相关人员」。
# 只补航线超时。销售超时不能改去 @ 航线,只 @ 本单销售。
if kind == WAIT_AIR_QUOTE and not people:
people = _airline_staff_mentions(book)
names = [str(p.get("name") or "").strip() for p in people if str(p.get("name") or "").strip()]
chat_id = str(row.get("collab_chat_id") or row.get("collabChatId") or "").strip()
where = conversation_of(wait_kind=kind, collab_chat_id=chat_id)
# 群里用 @。销售超时即使发到私聊,也写成 @销售姓名,不写「您」或「相关人员」。
mention = where == "group" or kind in _SALES_WAIT_KINDS
text = copy.event_exception_session_text(
exception_type=exception_type(kind),
work_order_no=no,
names=names,
hours=hours,
mention=mention,
test_prefix=test_prefix,
)
ok_session = False
if where == "group":
if not chat_id:
logger.warning("事件异常无协同群 wo=%s kind=%s", no, kind)
return False
sender = send_group or _default_send_group
ok_session = bool(sender(chat_id=chat_id, content=text, business_line=line))
else:
uid = str((people[0].get("wecomId") if people else "") or sales.get("wecomId") or "").strip()
if not uid:
logger.warning("事件异常无私聊对象 wo=%s kind=%s", no, kind)
return False
enq = enqueue or _default_enqueue
ok_session, info = enq(
touser=uid,
content=text,
dedupe_key=f"event-ex:{no}:{kind}",
payload={},
)
if not ok_session:
logger.warning("事件异常私聊入队失败 wo=%s info=%s", no, info)
if not ok_session:
return False
if notify_ready:
_send_app_cards(
people=people,
kind=kind,
work_order_no=no,
hours=hours,
test_prefix=test_prefix,
enqueue=enqueue or _default_enqueue,
)
recorder = getattr(book, "record_event_exception", None)
if callable(recorder):
recorder(work_order_no=no, event_exception=exception_type(kind))
return True
def _ticket_sales(row: dict[str, Any], watchers: list[dict[str, Any]]) -> dict[str, Any]:
"""
这张工单的销售。销售超时只 @ 这个人。
先按工单上的销售企微对上关注人,对不上再用岗位是销售的那条。
姓名用工单上的销售名,避免点成「销售」两个字或点到航线。
"""
sid = str(row.get("sales_wecom_id") or row.get("salesWecomId") or "").strip()
sname = str(row.get("sales_name") or row.get("salesName") or "").strip()
matched: dict[str, Any] | None = None
for raw in watchers:
uid = str(raw.get("wecomId") or raw.get("wecom_id") or "").strip()
role = str(raw.get("roleCode") or raw.get("role_code") or "").strip().lower()
if sid and uid == sid:
matched = dict(raw)
break
if matched is None and role == "sales":
matched = dict(raw)
if matched is None:
matched = {
"wecomId": sid,
"roleCode": "sales",
"notifyEventException": 0,
}
name = str(matched.get("name") or "").strip()
if sname and (not name or name in {sid, "销售"}):
matched["name"] = sname
elif not name:
matched["name"] = sname or "销售"
if sid and not str(matched.get("wecomId") or "").strip():
matched["wecomId"] = sid
matched.setdefault("roleCode", "sales")
return matched
def _airline_staff_mentions(book: Any) -> list[dict[str, Any]]:
"""
主账按岗位列出启用的航线同事。
只在群成员里没有航线岗时调用。主账失败或没有这个接口时返回空,话术仍用「相关人员」。
不写工单。可与巡检线程并发,每次只读这一单的员工名单。
"""
lister = getattr(book, "list_staff_by_role", None)
if not callable(lister):
return []
try:
rows = list(lister(role_code="air") or [])
except Exception:
logger.exception("事件异常列航线岗失败")
return []
return mentions_for_role(rows, role_code="air")
def _send_app_cards(
*,
people: list[dict[str, Any]],
kind: str,
work_order_no: str,
hours: float,
test_prefix: bool,
enqueue: EnqueueFn,
) -> None:
payload = copy.event_exception_alert_payload(
card_title=card_title(kind),
work_order_no=work_order_no,
hours=hours,
test_prefix=test_prefix,
)
fallback = copy.event_exception_alert_text(
card_title=card_title(kind),
work_order_no=work_order_no,
hours=hours,
test_prefix=test_prefix,
)
for person in pick_app_targets(people):
uid = str(person.get("wecomId") or "").strip()
if not uid:
continue
ok, info = enqueue(
touser=uid,
content=fallback,
dedupe_key=f"event-ex-app:{work_order_no}:{kind}:{uid}",
payload=payload,
)
if not ok:
logger.warning("事件异常应用卡片入队失败 uid=%s info=%s", uid, info)
def _default_enqueue(**kwargs: Any) -> tuple[bool, str]:
from agent.channel.queue import get_message_store
return get_message_store().enqueue_outbound(**kwargs)
def _default_send_group(*, chat_id: str, content: str, business_line: str = "") -> bool:
line = (business_line or "").strip().upper()
if line == "AIR":
from agent.channel.aibot.reply import AibotReplyClient
out = AibotReplyClient().send_group(chat_id=chat_id, content=content)
return bool(out.get("ok"))
from agent.channel.outbox.sender import WeComAppClient
from agent.config import get_settings
cfg = get_settings()
client = WeComAppClient(
corp_id=cfg.wecom_corp_id,
secret=cfg.wecom_secret,
agent_id=int(cfg.wecom_agent_id or "0") or 0,
)
out = client.send_group(chat_id=chat_id, content=content)
return bool(out.get("ok"))