成交走询价消息服务;陆运/海运私聊下一句能记下原因;协商中回执写明顺延三个工作日和下次 9:00。 Co-authored-by: Cursor <cursoragent@cursor.com>
167 lines
6.4 KiB
Python
167 lines
6.4 KiB
Python
"""
|
||
通道消费循环: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
|