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

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

362 lines
11 KiB
Python

"""
意图:空运协同群入站(BOT @ 之后的文字或附件)。
本文件职责:激活/切单后交给 air_group_ops。海运应用群不进这里。
禁止:存档文字当业务;锁舱用当前工单去猜;套 sea_group_ops 建群。
"""
from __future__ import annotations
import logging
from typing import Any, Optional
from agent.channel.wecom.models import InboundMessage
from agent.policy import inquiry_copy as copy
from agent.policy.air_group_ops import (
activate_ticket,
apply_air_confirm_tms,
apply_air_deal,
apply_air_fields,
apply_air_file,
apply_air_lock,
apply_air_quote,
apply_air_release,
apply_air_tms_choice,
current_ticket,
peek_air_pending,
resolve_cabin_ticket,
)
from agent.policy.air_text_flow import get_air_text_flow
from agent.policy.sea_group_ops import (
capture_group_lost_reason,
group_client_of,
is_handoff_ticket,
send_group_text,
ticket_view,
)
from agent.routing.air_group_intent import (
INTENT_ACTIVATE,
INTENT_ADJUST,
INTENT_CONFIRM_TMS,
INTENT_DEAL,
INTENT_DROP_TMS,
INTENT_FIELDS,
INTENT_KEEP_TMS,
INTENT_LOCK,
INTENT_LOST,
INTENT_NEGOTIATE,
INTENT_OTHER,
INTENT_PICK_OPTION,
INTENT_QUOTE,
INTENT_RELEASE,
classify_air_group_text,
)
logger = logging.getLogger(__name__)
def _attach_bot_client(engine, message: InboundMessage) -> None:
"""
每次入站换一次 BOT 出站客户端。
测试注入的 FakeBotGroup 不动。
有 response_url 也不走 stream;空运群正文一律走桥。
把这次入站的 aibot_req_id 带给桥,群里才能引用原 @。
不能复用上一次 @ 的过期 url。
"""
from agent.channel.aibot.reply import AibotReplyClient
existing = engine.group_client()
if existing is not None and not isinstance(existing, AibotReplyClient):
return
engine._group_client = AibotReplyClient(
str((message.raw or {}).get("response_url") or ""),
callback_req_id=str((message.raw or {}).get("aibot_req_id") or ""),
)
def handle_air_group(
message: InboundMessage,
*,
flow=None,
ledger=None,
injected_intent: str = "",
injected_quote: Optional[dict[str, Any]] = None,
) -> str:
"""
空运群入站:文字必须 @ BOT;附件走会话存档,不能 @。
航线 Excel/PDF 当报价;别人的只存档。无 sender_id 拒绝。
BOT 与存档 msgid 对不上:同一句只办一遍,存档后到就跳过,再 @ 仍回。
"""
if (message.chat_type or "") != "group":
return "group_ignore"
if not message.has_sender():
return "reject_no_sender"
engine = flow or get_air_text_flow()
if ledger is not None:
engine.ledger = ledger
_attach_bot_client(engine, message)
chat_id = (message.chat_id or "").strip()
if not chat_id:
return "group_ignore"
kind = (message.msg_type or "text").lower()
if kind in {"file", "image", "video"}:
return _handle_file(engine, message)
text = message.content or ""
from agent.redis_coord.air_inbound import mark_air_text, should_skip_air_text
if should_skip_air_text(
chat_id=chat_id,
sender_id=message.sender_id,
text=text,
raw=message.raw,
):
return "duplicate_inbound"
mark_air_text(chat_id=chat_id, sender_id=message.sender_id, text=text)
wo = copy.extract_work_order_no(text)
if wo:
from agent.policy.handoff_human import try_commit_air_group_handoff
handed = try_commit_air_group_handoff(engine, chat_id=chat_id, work_order_no=wo)
if handed:
return handed
intent = classify_air_group_text(text, injected=injected_intent)
logger.info(
"air_group 文字 chat=%s sender=%s intent=%s wo=%s",
chat_id,
message.sender_id,
intent,
wo or "-",
)
if intent in {INTENT_LOCK, INTENT_RELEASE}:
return _handle_cabin(
engine,
message,
intent=intent,
work_order_no=wo,
injected_quote=injected_quote,
)
ticket = None
if wo:
ticket, phase = activate_ticket(
flow=engine,
chat_id=chat_id,
work_order_no=wo,
speak=intent == INTENT_ACTIVATE or _should_speak_activate(intent),
)
if ticket is None:
return phase
if is_handoff_ticket(ticket):
logger.info("air_group 转人工不再跟 chat=%s", chat_id)
return "handoff_ignore"
if intent == INTENT_ACTIVATE:
return phase
else:
ticket = current_ticket(engine, chat_id)
if ticket is None:
send_group_text(group_client_of(engine), chat_id, copy.air_need_work_order())
return "need_work_order"
if is_handoff_ticket(ticket):
logger.info("air_group 转人工不再跟 chat=%s", chat_id)
return "handoff_ignore"
from agent.policy.system_exception import reply_paused_ticket, ticket_is_paused
if ticket_is_paused(ticket):
reply_paused_ticket(
ticket=ticket,
reply=lambda text: send_group_text(group_client_of(engine), chat_id, text),
)
return "system_exception_paused"
return _handle_bound(
engine,
message,
ticket=ticket,
intent=intent,
injected_quote=injected_quote,
)
def _handle_file(engine, message: InboundMessage) -> str:
"""
存档听到的附件(发文件不能 @)。句中/文件名有工单号先切单;没有就跟当前单。
"""
chat_id = message.chat_id
media = dict(message.media or {})
filename = str(
media.get("filename") or media.get("FileName") or (message.content or "")
).strip()
wo = copy.extract_work_order_no(
" ".join([message.content or "", filename])
)
if wo:
ticket, phase = activate_ticket(
flow=engine, chat_id=chat_id, work_order_no=wo, speak=False
)
if ticket is None:
return phase
else:
ticket = current_ticket(engine, chat_id)
if ticket is None:
send_group_text(group_client_of(engine), chat_id, copy.air_need_work_order())
return "need_work_order"
if is_handoff_ticket(ticket):
logger.info("air_group 转人工不再跟 chat=%s", chat_id)
return "handoff_ignore"
return apply_air_file(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
filename=filename,
media=media,
chat_id=chat_id,
)
def _should_speak_activate(intent: str) -> bool:
"""带工单号又办事时,先切单再办事;切单摘要只在纯激活时说,避免刷屏。"""
return intent in {INTENT_ACTIVATE, INTENT_OTHER}
def _handle_cabin(
engine,
message: InboundMessage,
*,
intent: str,
work_order_no: str,
injected_quote: Optional[dict[str, Any]],
) -> str:
_ = injected_quote
chat_id = message.chat_id
if not work_order_no:
send_group_text(group_client_of(engine), chat_id, copy.air_lock_need_work_order())
return "cabin_need_work_order"
ticket, phase = resolve_cabin_ticket(
flow=engine,
chat_id=chat_id,
work_order_no=work_order_no,
allow_terminal=intent == INTENT_RELEASE,
)
if ticket is None:
return phase
if intent == INTENT_RELEASE:
return apply_air_release(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
chat_id=chat_id,
instruction=message.content or "",
)
return apply_air_lock(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
chat_id=chat_id,
instruction=message.content or "",
)
def _handle_bound(
engine,
message: InboundMessage,
*,
ticket: object,
intent: str,
injected_quote: Optional[dict[str, Any]],
) -> str:
chat_id = message.chat_id
lost = capture_group_lost_reason(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
text=message.content or "",
)
if lost:
return lost
facts = dict(ticket_view(ticket).get("collab_facts") or {})
if intent == INTENT_PICK_OPTION or (
intent == INTENT_OTHER and str(facts.get("__pending_air_lock") or "") == "1"
):
option = _option_from_text(message.content or "")
if option:
return apply_air_lock(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
chat_id=chat_id,
instruction=message.content or "",
option_no=option,
)
if intent in {INTENT_DEAL, INTENT_LOST, INTENT_NEGOTIATE}:
outcome = {
INTENT_DEAL: "已成交",
INTENT_LOST: "未成交",
INTENT_NEGOTIATE: "协商中",
}[intent]
return apply_air_deal(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
outcome=outcome,
chat_id=chat_id,
)
if intent == INTENT_CONFIRM_TMS:
return apply_air_confirm_tms(
flow=engine, ticket=ticket, sender_id=message.sender_id, chat_id=chat_id
)
if intent in {INTENT_KEEP_TMS, INTENT_DROP_TMS}:
return apply_air_tms_choice(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
chat_id=chat_id,
keep_tms=intent == INTENT_KEEP_TMS,
)
if intent == INTENT_FIELDS:
return apply_air_fields(
flow=engine,
ticket=ticket,
text=message.content or "",
sender_id=message.sender_id,
chat_id=chat_id,
)
if intent == INTENT_ADJUST:
from agent.policy.quote_adjust_ops import apply_group_adjust
return apply_group_adjust(
flow=engine,
ticket=ticket,
sender_id=message.sender_id,
text=message.content or "",
sess=None,
chat_id=chat_id,
)
if intent == INTENT_QUOTE:
return apply_air_quote(
flow=engine,
ticket=ticket,
text=message.content or "",
sender_id=message.sender_id,
chat_id=chat_id,
injected_quote=injected_quote,
)
if peek_air_pending(ticket):
send_group_text(group_client_of(engine), chat_id, copy.group_tms_choice_unclear())
return "wait_tms_choice"
logger.info("air_group 其它意图 chat=%s intent=%s", chat_id, intent)
return "group_ignore"
def _option_from_text(text: str) -> str:
blob = (text or "").replace(" ", "").upper()
if "AIR-OPT" in blob or "AIROPT" in blob:
start = blob.find("AIR")
token = []
for ch in blob[start:]:
if ch.isalnum() or ch == "-":
token.append(ch)
else:
break
return "".join(token)
return ""