300 lines
10 KiB
Python
300 lines
10 KiB
Python
"""
|
|
转人工:缺项连问计数、虚拟工单号、入口/冻结判定。
|
|
|
|
本文件职责:只给流程提供纯函数与轻量序号;不发企微、不改主账六态。
|
|
调用:Worker / 私聊流程;禁止回调线程生成号后同步建群。
|
|
为何独立:海运/陆运/空运共用同一套计数与号码,避免三份 if。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import threading
|
|
from datetime import datetime
|
|
from typing import Any, Optional
|
|
from zoneinfo import ZoneInfo
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_SHANGHAI = ZoneInfo("Asia/Shanghai")
|
|
|
|
# 新询价第一次补问不计数;之后同一批缺项再问满 3 次仍不齐,第 3 次出入口。
|
|
CLARIFY_LIMIT = 3
|
|
VIRTUAL_PREFIX = "VT"
|
|
|
|
# 陆运类型+线路还对不上时,转人工摘要仍列出这一套协同项,全部为「-」。
|
|
_LAND_HANDOFF_COLLAB_FALLBACK = (
|
|
"客户名称",
|
|
"HS编码",
|
|
"贸易条款",
|
|
"包装方式",
|
|
"是否为危险品",
|
|
)
|
|
|
|
_seq_lock = threading.Lock()
|
|
_seq_by_day: dict[str, int] = {}
|
|
|
|
HANDOFF_STATUS = "转人工"
|
|
TERMINAL_STATUSES = frozenset({"已成交", "未成交", "已关闭", HANDOFF_STATUS})
|
|
FROZEN_STATUSES = frozenset({HANDOFF_STATUS})
|
|
|
|
|
|
def missing_signature(keys: list[str] | tuple[str, ...] | None) -> str:
|
|
"""缺项集合签名:顺序无关,用来判断是不是同一批。"""
|
|
return "|".join(sorted({str(k).strip() for k in (keys or []) if str(k).strip()}))
|
|
|
|
|
|
def note_clarify_round(sess: Any, missing_keys: list[str] | tuple[str, ...] | None) -> str:
|
|
"""
|
|
记下本轮缺项。返回 ask(继续补问)或 handoff(该出转人工入口)。
|
|
|
|
新询价(或缺项集合刚变)的第一次补问不计数、不打(1/3)。
|
|
之后同一签名再问:1/3、2/3;第 3 次仍不齐返回 handoff(调用方先发 3/3 再出入口)。
|
|
签名变了:计数清零,并作废「3 次反问」那张入口(handoff_offer=clarify)。
|
|
副作用:改 sess.clarify_* / 可能清 handoff_offer。
|
|
"""
|
|
sig = missing_signature(missing_keys)
|
|
if not sig:
|
|
if getattr(sess, "handoff_offer", "") == "clarify":
|
|
sess.handoff_offer = ""
|
|
sess.clarify_missing_sig = ""
|
|
sess.clarify_retry = 0
|
|
return "ask"
|
|
prev = str(getattr(sess, "clarify_missing_sig", "") or "")
|
|
retry = int(getattr(sess, "clarify_retry", 0) or 0)
|
|
if prev == sig:
|
|
retry = retry + 1
|
|
sess.clarify_retry = retry
|
|
if retry >= CLARIFY_LIMIT:
|
|
return "handoff"
|
|
return "ask"
|
|
if str(getattr(sess, "handoff_offer", "") or "") == "clarify":
|
|
sess.handoff_offer = ""
|
|
sess.clarify_missing_sig = sig
|
|
sess.clarify_retry = 0
|
|
return "ask"
|
|
|
|
|
|
def next_virtual_no(*, now: datetime | None = None) -> str:
|
|
"""
|
|
虚拟工单号:VT + 上海日历日 8 位 + 当日 4 位流水。
|
|
|
|
不占正式 WO。进程内计数;有 Redis 时尽量用 INCR,失败则回退内存。
|
|
"""
|
|
day = (now or datetime.now(_SHANGHAI)).strftime("%Y%m%d")
|
|
seq = _next_seq(day)
|
|
return f"{VIRTUAL_PREFIX}{day}{seq:04d}"
|
|
|
|
|
|
def _next_seq(day: str) -> int:
|
|
try:
|
|
from agent.config import get_settings
|
|
from agent.redis_coord.keys import join_key
|
|
|
|
settings = get_settings()
|
|
raw = getattr(settings, "redis_url", "") or ""
|
|
prefix = getattr(settings, "redis_key_prefix", "") or "inquiry_robot:"
|
|
if raw:
|
|
import redis
|
|
|
|
client = redis.Redis.from_url(raw, socket_timeout=2, socket_connect_timeout=2)
|
|
key = join_key(prefix, "vtseq", day)
|
|
n = int(client.incr(key))
|
|
client.expire(key, 3 * 24 * 3600)
|
|
return n
|
|
except Exception:
|
|
logger.info("虚拟号 Redis 序号不可用,改用进程内计数 day=%s", day)
|
|
with _seq_lock:
|
|
n = _seq_by_day.get(day, 0) + 1
|
|
_seq_by_day[day] = n
|
|
return n
|
|
|
|
|
|
def reset_virtual_seq_for_test() -> None:
|
|
"""单测重置当日流水。"""
|
|
with _seq_lock:
|
|
_seq_by_day.clear()
|
|
|
|
|
|
def is_virtual_no(work_order_no: str) -> bool:
|
|
raw = (work_order_no or "").strip().upper()
|
|
return raw.startswith(VIRTUAL_PREFIX) and len(raw) >= 10
|
|
|
|
|
|
def is_frozen_status(status: str) -> bool:
|
|
return (status or "").strip() in FROZEN_STATUSES
|
|
|
|
|
|
def is_handoff_frozen(sess: Any) -> bool:
|
|
"""这一票已经真正转人工,私聊不应再把它当当前询价。"""
|
|
if sess is None:
|
|
return False
|
|
if bool(getattr(sess, "handoff_committed", False)):
|
|
return True
|
|
return is_frozen_status(str(getattr(sess, "status", "") or ""))
|
|
|
|
|
|
def is_closed_for_handoff(status: str) -> bool:
|
|
"""已成交/未成交/已关闭不再出转人工入口。转人工本身也不再重复入口。"""
|
|
st = (status or "").strip()
|
|
return st in {"已成交", "未成交", "已关闭", HANDOFF_STATUS}
|
|
|
|
|
|
def display_ticket_no(sess: Any) -> str:
|
|
"""群/卡片展示号:正式号优先,否则虚拟号。"""
|
|
return (
|
|
str(getattr(sess, "work_order_no", "") or "").strip()
|
|
or str(getattr(sess, "virtual_work_order_no", "") or "").strip()
|
|
)
|
|
|
|
|
|
def ensure_virtual_no(sess: Any) -> str:
|
|
"""还没有正式号时确保有虚拟号。已有正式号不生成。"""
|
|
real = str(getattr(sess, "work_order_no", "") or "").strip()
|
|
if real:
|
|
return real
|
|
vt = str(getattr(sess, "virtual_work_order_no", "") or "").strip()
|
|
if vt:
|
|
return vt
|
|
vt = next_virtual_no()
|
|
sess.virtual_work_order_no = vt
|
|
return vt
|
|
|
|
|
|
# 这些阶段还没为本票建正式单,转人工入口必须用虚拟号。
|
|
_PRE_CREATE_PHASES = frozenset(
|
|
{
|
|
"",
|
|
"idle",
|
|
"clarify",
|
|
"need_mode",
|
|
"need_land_options",
|
|
"wait_confirm",
|
|
"handoff_offer",
|
|
}
|
|
)
|
|
|
|
|
|
def ticket_no_for_handoff_offer(sess: Any) -> str:
|
|
"""
|
|
转人工入口上给销售看的号码。
|
|
|
|
本票还在补问/核对本、没建正式单:给 VT,丢掉会话里挂着的旧 WO。
|
|
已经出过正式单(已报价/协同中等):沿用 WO。
|
|
"""
|
|
phase = str(getattr(sess, "phase", "") or "").strip()
|
|
real = str(getattr(sess, "work_order_no", "") or "").strip()
|
|
if real and not is_virtual_no(real) and phase not in _PRE_CREATE_PHASES:
|
|
return real
|
|
if real and not is_virtual_no(real):
|
|
sess.work_order_no = ""
|
|
return ensure_virtual_no(sess)
|
|
|
|
|
|
def inherit_pre_create_ticket(sess: Any, new_sess: Any) -> None:
|
|
"""
|
|
同一轮补问可以接着用号;从已出票旧书签开新补问则不带旧 WO/群。
|
|
"""
|
|
prev = str(getattr(sess, "phase", "") or "").strip()
|
|
if prev not in _PRE_CREATE_PHASES:
|
|
return
|
|
if not str(getattr(new_sess, "work_order_no", "") or "").strip():
|
|
new_sess.work_order_no = str(getattr(sess, "work_order_no", "") or "").strip()
|
|
if not str(getattr(new_sess, "collab_chat_id", "") or "").strip():
|
|
new_sess.collab_chat_id = str(getattr(sess, "collab_chat_id", "") or "").strip()
|
|
|
|
|
|
def void_virtual_no(sess: Any) -> None:
|
|
"""正式建单后作废虚拟号。"""
|
|
sess.virtual_work_order_no = ""
|
|
|
|
|
|
def has_collab_group(sess: Any) -> bool:
|
|
return bool(str(getattr(sess, "collab_chat_id", "") or "").strip())
|
|
|
|
|
|
def transport_label(business_line: str) -> str:
|
|
code = (business_line or "").strip().upper()
|
|
return {"SEA": "海运", "AIR": "空运", "LAND": "陆运"}.get(code, "")
|
|
|
|
|
|
def try_commit_air_group_handoff(flow: Any, *, chat_id: str, work_order_no: str) -> str:
|
|
"""
|
|
空运群里出现待转人工的工单号:发转人工摘要并冻结。
|
|
对不上或虚拟号已作废:返回空,交给原激活逻辑。
|
|
"""
|
|
no = (work_order_no or "").strip().upper()
|
|
if not no:
|
|
return ""
|
|
getter = getattr(flow, "session_by_work_order", None)
|
|
sess = getter(no) if callable(getter) else None
|
|
if sess is None:
|
|
return ""
|
|
if not sess.handoff_offer or sess.handoff_committed:
|
|
return ""
|
|
display = display_ticket_no(sess).upper()
|
|
if display != no:
|
|
return ""
|
|
from agent.policy import inquiry_copy as copy
|
|
from agent.policy.sea_group_ops import group_client_of, send_group_text
|
|
|
|
book = getattr(flow, "ledger", None)
|
|
airline_names: list[str] = []
|
|
lister = getattr(book, "list_staff_by_role", None) if book else None
|
|
if callable(lister):
|
|
try:
|
|
for row in lister(role_code="air") or []:
|
|
name = str(row.get("name") or "").strip()
|
|
if name:
|
|
airline_names.append(name)
|
|
except Exception:
|
|
logger.exception("空运转人工取航线名单失败")
|
|
sales = ""
|
|
namer = getattr(flow, "_sales_display_name", None)
|
|
if callable(namer):
|
|
sales = namer(sess.sender_id)
|
|
brief = copy.handoff_group_brief(
|
|
work_order_no=display_ticket_no(sess),
|
|
facts=sess.facts,
|
|
business_line="AIR",
|
|
sales_name=sales,
|
|
airline_names=airline_names,
|
|
)
|
|
send_group_text(group_client_of(flow), chat_id, brief)
|
|
real = (sess.work_order_no or "").strip()
|
|
if real and not is_virtual_no(real):
|
|
trans = getattr(book, "transition", None) if book else None
|
|
if callable(trans):
|
|
trans(
|
|
work_order_no=real,
|
|
to_status=HANDOFF_STATUS,
|
|
remark=copy.HANDOFF_CLARIFY_REASON
|
|
if sess.handoff_offer == "clarify"
|
|
else copy.HANDOFF_INTENT_REASON,
|
|
)
|
|
sess.status = HANDOFF_STATUS
|
|
from agent.policy.event_exception import clear_wait
|
|
|
|
clear_wait(book, work_order_no=real)
|
|
sess.handoff_committed = True
|
|
sess.handoff_offer = ""
|
|
sess.phase = "handoff_done"
|
|
sess.collab_chat_id = chat_id
|
|
saver = getattr(flow, "_save", None)
|
|
if callable(saver):
|
|
saver(sess)
|
|
return "handoff_done"
|
|
|
|
|
|
def collab_keys_for_handoff(facts: dict[str, str] | None, business_line: str) -> tuple[str, ...]:
|
|
"""转人工摘要要列出的协同键;对不上陆运套餐时用兜底五项。"""
|
|
from agent.schema.field_validate import SEA_COLLAB_FIELDS, collab_fields_for_ticket
|
|
|
|
line = (business_line or "").strip().upper()
|
|
keys = collab_fields_for_ticket(facts, line)
|
|
if keys:
|
|
return keys
|
|
if line == "LAND":
|
|
return _LAND_HANDOFF_COLLAB_FALLBACK
|
|
return SEA_COLLAB_FIELDS
|