Files
inquiry_robot/inquiry-agent/agent/channel/aibot/__init__.py
T
jillion886andCursor cfbce2fd69 落地空运单聊询价字段与线路选择。
必填核对后再建单;多条线路走 H5 点选;费用只展示该线路 TMS 回包,不套用整单空运费。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-22 15:02:32 +08:00

136 lines
4.9 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.
"""
AIBOT 入站入口:桥进程把长连接消息转到本路由。
本文件职责:校验 bridge token、归一 InboundMessage、入 inbox 快回。
禁止:在本请求线程调 LLM / Graph / 出站企微。
路径:/internal/aibot/inbound(对齐 prompt/13)。
"""
from __future__ import annotations
import logging
from typing import Any, Optional
from fastapi import APIRouter, Header, HTTPException
from pydantic import BaseModel, Field
from agent.channel.queue import get_message_store
from agent.channel.wecom.models import InboundMessage
from agent.config import get_settings
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/internal/aibot", tags=["aibot"])
class AibotInboundBody(BaseModel):
"""桥转发的归一字段(桥侧已解析,智能体不再调企微换 ID)。"""
sender_id: str = Field(default="", description="企微 userid,必须来自 AIBOT 报文")
message_id: str = ""
content: str = ""
msg_type: str = "text"
chat_id: str = ""
chat_type: str = "group"
agent_id: str = ""
create_time: int = 0
media: dict[str, Any] = Field(default_factory=dict)
raw: dict[str, Any] = Field(default_factory=dict)
def normalize_aibot_payload(body: dict[str, Any]) -> AibotInboundBody:
"""
同时吃两套入站:智能体骨架字段,以及现网官方桥转发的企微 body
(msgid/chatid/from.userid/text.content)。
"""
data = dict(body or {})
from_user = data.get("from") if isinstance(data.get("from"), dict) else {}
text_obj = data.get("text") if isinstance(data.get("text"), dict) else {}
sender = str(data.get("sender_id") or from_user.get("userid") or "").strip()
mid = str(data.get("message_id") or data.get("msgid") or "").strip()
content = str(data.get("content") or text_obj.get("content") or "")
chat_id = str(data.get("chat_id") or data.get("chatid") or "").strip()
chattype = str(data.get("chat_type") or data.get("chattype") or "group").strip().lower()
chat_type = "c2c" if chattype in {"single", "c2c", "private"} else "group"
raw = dict(data.get("raw") or {})
raw.setdefault("source", "aibot")
if data.get("response_url"):
raw.setdefault("response_url", data.get("response_url"))
req_id = str(data.get("aibot_req_id") or raw.get("aibot_req_id") or "").strip()
if req_id:
raw["aibot_req_id"] = req_id
return AibotInboundBody(
sender_id=sender,
message_id=mid,
content=content,
msg_type=str(data.get("msg_type") or data.get("msgtype") or "text"),
chat_id=chat_id,
chat_type=chat_type,
agent_id=str(data.get("agent_id") or data.get("aibotid") or ""),
create_time=int(data.get("create_time") or 0),
media=dict(data.get("media") or {}),
raw=raw,
)
def _check_bridge_token(x_token: Optional[str]) -> None:
s = get_settings()
expected = (getattr(s, "wecom_aibot_bridge_token", None) or "").strip()
if not expected:
if (s.ytd_env or "").lower() == "prod":
raise HTTPException(status_code=503, detail="AIBOT bridge token 未配置")
logger.warning("aibot inbound:未配置 bridge token,test 跳过校验")
return
if (x_token or "").strip() != expected:
raise HTTPException(status_code=401, detail="unauthorized")
@router.get("/ping")
def aibot_ping() -> dict[str, str]:
return {"status": "ok", "channel": "aibot"}
@router.post("/inbound")
def aibot_inbound(
body: dict[str, Any],
x_aibot_bridge_token: Optional[str] = Header(
default=None, alias="X-Aibot-Bridge-Token"
),
x_ytd_aibot_bridge_token: Optional[str] = Header(
default=None, alias="X-Ytd-Aibot-Bridge-Token"
),
) -> dict[str, Any]:
"""
接收桥转发入站:无 sender_id 拒绝;其余入 inbox。
"""
_check_bridge_token(x_aibot_bridge_token or x_ytd_aibot_bridge_token)
parsed = normalize_aibot_payload(body)
raw = dict(parsed.raw or {})
raw.setdefault("source", "aibot")
msg = InboundMessage(
sender_id=parsed.sender_id.strip(),
message_id=parsed.message_id.strip(),
content=parsed.content or "",
msg_type=parsed.msg_type or "text",
chat_id=parsed.chat_id or "",
chat_type=parsed.chat_type or "group",
agent_id=parsed.agent_id or "",
create_time=int(parsed.create_time or 0),
media=parsed.media or {},
raw=raw,
)
if not msg.has_sender():
raise HTTPException(status_code=400, detail="missing_sender_id")
if not msg.message_id:
raise HTTPException(status_code=400, detail="missing_message_id")
ok, info = get_message_store().enqueue_inbound(msg)
logger.info(
"aibot.inbound ok=%s info=%s sender=%s msg_id=%s",
ok,
info,
msg.sender_id,
msg.message_id,
)
return {"status": "accepted" if ok else "duplicate", "detail": info}