Files
inquiry_robot/inquiry-agent/agent/channel/runtime.py
T
jillion886andCursor 1a2e1d6cb2 补齐成交跟进:按线路推成交通知、收下未成交原因、协商中按工作日再问。
成交走询价消息服务;陆运/海运私聊下一句能记下原因;协商中回执写明顺延三个工作日和下次 9:00。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-19 16:08:48 +08:00

167 lines
6.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
通道消费循环:inbox → RouteDecision → dispatch → handler;outbox → 企微发送。
本文件职责:在 HTTP 进程内独立线程跑快路径,不与 LibreOffice 串队。
禁止:在回调请求线程调用本循环的阻塞逻辑;禁止 Redis 全局单例锁。
"""
from __future__ import annotations
import logging
import threading
from typing import Optional
from agent.channel.outbox.sender import PERMANENT_SEND_ERRORS, WeComAppClient
from agent.channel.queue import get_message_store
from agent.config import Settings
from agent.routing import IntentRouter
from agent.routing.dispatch import dispatch_inbound
logger = logging.getLogger(__name__)
# 避免 httpx 把 access_token / secret 打进 INFO
logging.getLogger("httpx").setLevel(logging.WARNING)
class ChannelRuntime:
"""
企微通道运行时:inbox/outbox 各一条守护线程。
启动时机:FastAPI lifespan;停止:进程退出。
inbox:私聊自由文字走 DeepSeek A 认意图(本线程,非回调);群消息和点卡不进模型。
"""
def __init__(self, settings: Settings) -> None:
self._settings = settings
self._store = get_message_store()
self._stop = threading.Event()
self._threads: list[threading.Thread] = []
# 开网后销售原话交给模型认「续办/新询价」;未开网才走 stub。
self._router = IntentRouter(allow_network=bool(settings.llm_allow_network))
agent_id = int(settings.wecom_agent_id or "0") or 0
self._client = WeComAppClient(
corp_id=settings.wecom_corp_id,
secret=settings.wecom_secret,
agent_id=agent_id,
)
notify_id = int(getattr(settings, "wecom_notify_agent_id", "") or "0") or 0
notify_secret = str(getattr(settings, "wecom_notify_secret", "") or "").strip()
self._notify_client = None
if notify_id and notify_secret:
self._notify_client = WeComAppClient(
corp_id=settings.wecom_corp_id,
secret=notify_secret,
agent_id=notify_id,
)
@property
def store(self):
return self._store
def start(self) -> None:
"""启动 inbox/outbox 消费线程。"""
self._stop.clear()
self._threads = [
threading.Thread(target=self._inbox_loop, name="wecom-inbox", daemon=True),
threading.Thread(target=self._outbox_loop, name="wecom-outbox", daemon=True),
]
for t in self._threads:
t.start()
logger.info("ChannelRuntime 已启动(inbox/outbox 独立线程;dispatch 分发)")
def stop(self) -> None:
self._stop.set()
def _inbox_loop(self) -> None:
while not self._stop.is_set():
item = self._store.claim_inbound()
if item is None:
self._store.wait_inbox(1.0)
continue
try:
message = item.message
# 群消息由 dispatch 按 BOT/存档分流,不消耗 DeepSeek A,避免空运群先卡 2~3 秒。
if (message.chat_type or "") == "group":
from agent.routing.decision import RouteDecision
decision = RouteDecision(intent="ordinary_text_other")
elif (message.msg_type or "text").lower() == "image":
from agent.routing.decision import RouteDecision
# 私聊图片不进 DeepSeek A(正文为空会被判闲聊)。
decision = RouteDecision(intent="image_inquiry", confidence=1.0)
else:
decision = self._router.route_inbound(
text=message.content or "",
sender_id=message.sender_id or "",
event=str((message.raw or {}).get("event") or ""),
)
dispatch_inbound(message, decision)
self._store.complete_inbound(item.id)
except Exception as exc: # noqa: BLE001
logger.exception("inbox 处理失败 id=%s", item.id)
self._store.complete_inbound(item.id, error=str(exc))
def _outbox_loop(self) -> None:
while not self._stop.is_set():
item = self._store.claim_outbound()
if item is None:
self._store.wait_outbox(1.0)
continue
# 过期丢弃:超过 30 分钟的出站不再发(对齐「过期不计业务成功」)
import time
if time.time() - item.created_at > 1800:
logger.warning("outbox 过期丢弃 id=%s", item.id)
self._store.complete_outbound(item.id)
continue
payload = getattr(item, "payload", None) or {}
client = self._client
if payload.get("wecom_app") == "notify":
if self._notify_client is None:
logger.error("询价消息服务未配置,成交通知丢弃 id=%s", item.id)
self._store.complete_outbound(
item.id, error="notify_app_missing", permanent=True
)
continue
client = self._notify_client
result = client.send_outbound(
touser=item.touser,
content=item.content,
payload=payload,
)
if result.ok:
self._store.complete_outbound(item.id)
elif result.errcode in PERMANENT_SEND_ERRORS:
logger.error(
"outbox 永久失败不再重试 id=%s err=%s:%s",
item.id,
result.errcode,
result.errmsg,
)
self._store.complete_outbound(
item.id,
error=f"permanent:{result.errcode}:{result.errmsg}",
permanent=True,
)
else:
self._store.complete_outbound(
item.id,
error=f"{result.errcode}:{result.errmsg}",
)
_RUNTIME: Optional[ChannelRuntime] = None
def get_channel_runtime() -> Optional[ChannelRuntime]:
return _RUNTIME
def start_channel_runtime(settings: Settings) -> ChannelRuntime:
"""幂等启动通道运行时。"""
global _RUNTIME
if _RUNTIME is None:
_RUNTIME = ChannelRuntime(settings)
_RUNTIME.start()
return _RUNTIME