251 lines
9.6 KiB
Python
251 lines
9.6 KiB
Python
"""
|
|
入站分发壳:RouteDecision.intent → 唯一 handler。
|
|
|
|
本文件职责:维护意图→可调用对象表;dispatch_inbound 只查表调用。
|
|
禁止:在本文件写补问/建单/TMS 业务;禁止一个函数处理全部卡片。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from typing import Callable, Optional
|
|
|
|
from agent.channel.wecom.models import InboundMessage
|
|
from agent.handlers import abnormal, card_action, close_deal, continue_thread, new_inquiry
|
|
from agent.handlers.echo_text import handle_echo_text
|
|
from agent.handlers.group_collab import handle_group_collab
|
|
from agent.handlers.image_inquiry import handle_image_inquiry
|
|
from agent.handlers.text_inquiry import handle_text_inquiry, handle_text_other, is_land_confirm_action
|
|
from agent.routing.decision import RouteDecision
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Handler 签名约定:只收 InboundMessage(出站/store 由 handler 自行取依赖)
|
|
HandlerFn = Callable[[InboundMessage], None]
|
|
|
|
|
|
def _noop(message: InboundMessage) -> None:
|
|
"""忽略类意图:不做事。"""
|
|
logger.info("dispatch.noop sender=%s msg=%s", message.sender_id, message.message_id)
|
|
|
|
|
|
def _wrap_echo(message: InboundMessage) -> None:
|
|
"""联调残留:仅 echo_text 意图使用。"""
|
|
from agent.channel.queue import get_message_store
|
|
from agent.config import get_settings
|
|
|
|
settings = get_settings()
|
|
handle_echo_text(
|
|
message,
|
|
get_message_store(),
|
|
ytd_env=settings.ytd_env,
|
|
reply_template=settings.wecom_callback_reply_text,
|
|
)
|
|
|
|
|
|
def _wrap_text_inquiry(message: InboundMessage) -> None:
|
|
"""普通文字询价:空运闭环。"""
|
|
handle_text_inquiry(message)
|
|
|
|
|
|
def _wrap_text_other(message: InboundMessage) -> None:
|
|
"""闲聊:回首次引导;若正等按钮则当按钮。"""
|
|
handle_text_other(message)
|
|
|
|
|
|
def _wrap_new(message: InboundMessage) -> None:
|
|
new_inquiry.handle_new_inquiry(message)
|
|
|
|
|
|
def _wrap_image(message: InboundMessage) -> None:
|
|
"""私聊图片询价:认图后当口播新开,复用文字流程。"""
|
|
handle_image_inquiry(message)
|
|
|
|
|
|
def _wrap_continue(message: InboundMessage, *, thread_id: str = "") -> None:
|
|
continue_thread.handle_continue_thread(message, thread_id=thread_id)
|
|
|
|
|
|
def _wrap_card(message: InboundMessage) -> None:
|
|
"""
|
|
卡片意图:把企微 TaskId / EventKey 原样交给流程,按当时工单号找回会话。
|
|
action 用正文(EventKey 或按钮文案),禁止写死 unknown。
|
|
"""
|
|
from agent.policy import inquiry_copy as copy
|
|
|
|
action = copy.canonical_button(message.content or "") or (message.content or "").strip()
|
|
payload = dict(message.raw or {})
|
|
payload.setdefault("event_key", message.content or "")
|
|
if message.chat_id:
|
|
payload["chat_id"] = message.chat_id
|
|
payload["chat_type"] = message.chat_type
|
|
card_action.handle_card_action(
|
|
sender_id=message.sender_id,
|
|
thread_id="",
|
|
wait_version=0,
|
|
action=action or "unknown",
|
|
payload=payload,
|
|
)
|
|
|
|
|
|
def _wrap_close(message: InboundMessage) -> None:
|
|
close_deal.handle_close_deal(
|
|
sender_id=message.sender_id,
|
|
inquiry_no="",
|
|
outcome="unknown",
|
|
)
|
|
|
|
|
|
def _wrap_abnormal(message: InboundMessage) -> None:
|
|
abnormal.handle_abnormal(sender_id=message.sender_id, inquiry_no="", reason="dispatch")
|
|
|
|
|
|
def _wrap_route_failed(message: InboundMessage) -> None:
|
|
"""模型没给出路由结论:请销售再说一次,不假装听懂。"""
|
|
from agent.channel.queue import get_message_store
|
|
from agent.policy import inquiry_copy as copy
|
|
|
|
get_message_store().enqueue_outbound(
|
|
touser=message.sender_id,
|
|
content=copy.ROUTE_MODEL_FAILED,
|
|
dedupe_key=f"route-fail:{message.message_id}",
|
|
)
|
|
|
|
|
|
# 意图 → handler(空壳或联调 echo);扩展时只改表,不改 ChannelRuntime 分支森林
|
|
HANDLER_REGISTRY: dict[str, HandlerFn] = {
|
|
"echo_text": _wrap_echo,
|
|
"ordinary_text_inquiry": _wrap_text_inquiry,
|
|
"ordinary_text_other": _wrap_text_other,
|
|
"attachment_inquiry": _wrap_new,
|
|
"image_inquiry": _wrap_image,
|
|
"multi_segment_transport": _wrap_new,
|
|
"tms_read_only_quote": _wrap_continue,
|
|
"new_inquiry": _wrap_new,
|
|
"continue_thread": _wrap_continue,
|
|
"card_action": _wrap_card,
|
|
"close_deal": _wrap_close,
|
|
"abnormal": _wrap_abnormal,
|
|
"reject_no_sender": _noop,
|
|
"ignore_empty": _noop,
|
|
"stub_needs_model": _wrap_echo,
|
|
"route_model_failed": _wrap_route_failed,
|
|
"route_parse_failed": _wrap_route_failed,
|
|
}
|
|
|
|
|
|
def _is_air_collab_room(message: InboundMessage) -> bool:
|
|
"""当前群已激活空运工单。查不到或海运群返回 False。"""
|
|
try:
|
|
from agent.handlers.air_group_collab import get_air_text_flow
|
|
from agent.policy.air_group_ops import bound_air_ticket
|
|
|
|
return bound_air_ticket(get_air_text_flow(), message.chat_id or "") is not None
|
|
except Exception:
|
|
logger.exception("判断空运协同群失败 chat=%s", getattr(message, "chat_id", ""))
|
|
return False
|
|
|
|
|
|
def resolve_handler(intent: str) -> Optional[HandlerFn]:
|
|
"""查表;未知意图返回 None。"""
|
|
return HANDLER_REGISTRY.get(intent)
|
|
|
|
|
|
def dispatch_inbound(message: InboundMessage, decision: RouteDecision) -> str:
|
|
"""
|
|
按 RouteDecision 调用唯一 handler。
|
|
|
|
返回:实际调用的 intent 名,或 unknown_intent / skipped。
|
|
副作用:仅 handler 内产生;本函数不做业务。
|
|
"""
|
|
if (message.chat_type or "") == "group":
|
|
from agent.policy import inquiry_copy as copy
|
|
|
|
event = str((message.raw or {}).get("event") or "").lower()
|
|
raw_text = (message.content or "").strip()
|
|
# 只有真点卡才走成交卡;群里打「成交」两个字要当文字成交,不能当私聊点卡丢掉。
|
|
is_card_event = event in {"template_card_event", "click", "templatecardevent"}
|
|
is_card_key = raw_text.startswith(
|
|
(copy.BTN_DEAL_KEY + ":", copy.BTN_LOST_KEY + ":", copy.BTN_NEGOTIATE_KEY + ":")
|
|
)
|
|
if is_card_event or is_card_key:
|
|
_wrap_card(message)
|
|
return "card_action"
|
|
source = str((message.raw or {}).get("source") or "").lower()
|
|
if source in {"aibot", "aibot_bridge"}:
|
|
from agent.handlers.air_group_collab import handle_air_group
|
|
|
|
phase = handle_air_group(message)
|
|
logger.info("dispatch air_group phase=%s chat=%s", phase, message.chat_id)
|
|
return "air_group_collab"
|
|
# 发文件不能 @:空运群附件靠会话存档听见。
|
|
# 文字本应走 BOT;BOT 没推到时,存档里已经带 @询价小助手 的短句按空运补听。
|
|
kind = (message.msg_type or "text").lower()
|
|
air_room = _is_air_collab_room(message)
|
|
if source == "wshoto_archive" and kind in {"text", ""}:
|
|
from agent.routing.air_group_intent import is_missed_air_bot_at
|
|
|
|
if is_missed_air_bot_at(raw_text, room_already_air=air_room):
|
|
from agent.handlers.air_group_collab import handle_air_group
|
|
|
|
phase = handle_air_group(message)
|
|
logger.info("dispatch air_group archive_at phase=%s chat=%s", phase, message.chat_id)
|
|
return "air_group_collab"
|
|
if air_room:
|
|
if kind in {"file", "image", "video"}:
|
|
from agent.handlers.air_group_collab import handle_air_group
|
|
|
|
phase = handle_air_group(message)
|
|
logger.info("dispatch air_group archive_file phase=%s chat=%s", phase, message.chat_id)
|
|
return "air_group_collab"
|
|
logger.info("dispatch 空运群存档文字忽略,须 @ BOT chat=%s", message.chat_id)
|
|
return "air_archive_text_ignore"
|
|
phase = handle_group_collab(message)
|
|
logger.info("dispatch group_collab phase=%s chat=%s", phase, message.chat_id)
|
|
return "group_collab"
|
|
|
|
if not message.has_sender() and decision.intent != "reject_no_sender":
|
|
logger.warning("dispatch 拒绝:无 sender_id")
|
|
HANDLER_REGISTRY["reject_no_sender"](message)
|
|
return "reject_no_sender"
|
|
|
|
kind = (message.msg_type or "text").lower()
|
|
if kind == "image":
|
|
handle_image_inquiry(message)
|
|
return "image_inquiry"
|
|
if kind in {"text", ""}:
|
|
from agent.handlers.image_inquiry import append_wave_text, is_active_image_wave
|
|
|
|
if is_active_image_wave(message.sender_id) and (message.content or "").strip():
|
|
append_wave_text(message.sender_id, message.content or "")
|
|
logger.info("dispatch 图片波次收字 sender=%s", message.sender_id)
|
|
return "image_wave_follow_text"
|
|
|
|
fn = resolve_handler(decision.intent)
|
|
if fn is None:
|
|
logger.warning("dispatch 未知意图 intent=%s → noop", decision.intent)
|
|
_noop(message)
|
|
return "unknown_intent"
|
|
|
|
logger.info(
|
|
"dispatch intent=%s sender=%s msg=%s wo=%s land_confirm=%s",
|
|
decision.intent,
|
|
message.sender_id,
|
|
message.message_id,
|
|
decision.thread_id or "-",
|
|
is_land_confirm_action(message) if message.has_sender() else False,
|
|
)
|
|
if decision.intent in {
|
|
"ordinary_text_other",
|
|
"route_model_failed",
|
|
"route_parse_failed",
|
|
} and is_land_confirm_action(message):
|
|
logger.info("dispatch 陆运核对卡确定:纠正闲聊路由 sender=%s", message.sender_id)
|
|
handle_text_inquiry(message)
|
|
return "ordinary_text_inquiry"
|
|
if decision.intent == "continue_thread":
|
|
continue_thread.handle_continue_thread(message, thread_id=decision.thread_id)
|
|
return decision.intent
|
|
fn(message)
|
|
return decision.intent
|