Files
group_fqcd_jr/app/service/agent_persistence_service.py
T
张胜宇 5d0becb67d 客服 Agent 重构收口:五出口决策链 + 知识库档位隔离 + 前端入参边界(答辩演示版本)
一、客服 Agent 智能增强(正面回应"不智能、动不动就转人工")
- 决策链由 2 个出口扩到 5 个:E1 澄清 / E2 计算型 / E3 知识直返 / E4 证据约束生成 / E5 分级回退
- 转人工从"默认动作"降为最后一档 E5c,只保留 4 类白名单:
  P0 反诈 / P1 账户与个人数据 / P2 写操作与争议 / 用户明确要求人工
- 46 条金标实测(修复前 → 修复后):
  转人工率 43.5% → 10.9%;出口准确率 45.7% → 100%;事实正确率 69.6% → 100%
  禁忌违反 1 → 0;档位越权 / 无出处数字 / 误拒 四项零容忍全 0
- 安全不变量 INV-1~INV-5;零容忍规则未删,改的是挂载点
  (输出侧字面黑名单 → 检索层档位隔离 + 判定层合规词表 + 输出守护)

二、知识库:档位单点化与物理隔离
- 新增 app/core/knowledge_tier.py 作为档位规则唯一落点(G-03),
  knowledge_contracts.py 原定义块改为显式再导出(X as X,非副本)
- 档位过滤由 bool 默认值(fail-open)改为 tiers 必填集合(缺参即 TypeError)
- Milvus 侧四集合按 visibility 分区键物理隔离;双 schema 收敛为一套
- 新增 app/core/actor.py:访客三元组与匿名判定的唯一构造/判定点(G-01/G-01b)
- 新增 app/core/fund_fee_rules.py:费率计算纯函数

三、前端入参边界对齐(本轮 W11 新修,4 处"校验宽于存储")
- message 加 max_length=8000(与浮窗 widget.js 的 maxlength 一致)
- session_id 加 1—64;idempotency_key 上限 128 → 64(对齐列宽 String(64))
- feedback_type 加 max_length=32(对齐列宽 String(32))
- 8 条路径参数补 min_length=1 + max_length=64 + 字符集正则
  ({session_id} / {run_id} / {handover_id})
- 改前超限值会落到 MySQL 才失败(500);改后一律 422 AGENT_INPUT_INVALID + 字段级定位
- 新增 tests/unit/api/test_frontend_boundaries.py(33 例),含"端点表 ↔ OpenAPI 全量对照"

四、投顾模块整体清除(D4.4 / D4.5)
- 删除投顾相关 controller / schema / model / repository / service 及门户页面
- tools/portal_api_check.py 同步作废 AD003/AD005/AD011/A047 四条用例与 advisor_t 登录
  (端点与账号均已不存在,此前稳定报 3 条假红)

五、验证(提交前实测)
- pytest -q:1856 passed / 2 skipped / 0 failed
- ruff check app tools tests:19(= 基线);mypy app:2(= 基线)
- 前端接口契约体检 portal_api_check.py:38 项,通过 34,失败 0,跳过 4
- 全链路冒烟 e2e_smoke_test.py --read-only:31/31
- HTTP 全链路探针 http_probe.py:11/11 succeeded
- 跨文档一致性 _consistency.py:GATE PASS
- 真机边界复验 12 条:12/12 符合预期

六、纪律与文档
- 可改文件白名单 A-09(docs/46)与底座会签申请单 A-10(docs/47,组 1—组 4 全部受理)
- 零 DDL:未新增/修改任何表结构,89 张业务表与基线一致
- 证据留痕:docs/evidence/**(含 46 条金标 score、快照、清除与重建记录)
- 未提交(刻意排除,见提交说明):仓库内 客服agent/ 与 开发文档/ 是 2026-09-16 前的
  过期副本(Todolist 440 行 vs 权威 D2.1 1167 行),权威正本在仓库外;
  _chunks_report.txt 是 tools/build_knowledge_chunks.py 生成的本地产物
2026-09-20 14:33:30 +08:00

326 lines
18 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.
import logging
from datetime import UTC, datetime
from decimal import Decimal
from typing import Any
from uuid import uuid4
from sqlalchemy import select, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.contracts import AgentResult, DomainEvent
from app.core.customer_service_rules import (
ALL_TRANSFER_REASONS,
normalize_transfer_reason,
transfer_priority,
)
from app.core.errors import RunLeaseLostError
from app.model.audit import InteractionAudit
from app.model.conversation import ConversationMessage
from app.model.platform import AgentRun, DomainEventOutbox, HandoverTicket, RequestIdempotency
from app.model.risk import RiskUser
from app.model.session import ConversationSession
from app.service.customer_service_handover_context import (
MAX_SUMMARY_MESSAGES,
CustomerServiceHandoverContext,
build_customer_service_handover_context,
)
logger = logging.getLogger(__name__)
#: 建单白名单只对客服 Agent 生效(`E-01` ①)。这里硬编码而不导入实现层模块,
#: 避免底座(持久化服务)反向依赖业务 Agent 实现。
CUSTOMER_SERVICE_AGENT_TYPE = "customer_service"
#: 治理层追加免责声明时使用的分隔形状(`app/service/agent/governance.py` 里定义)。
#: 这里只用于**审计留痕**,不参与任何判定:判据是"末尾是否出现这个形状"。
_GOVERNANCE_APPEND_MARKERS: tuple[str, ...] = ("\n\n本内容仅为投资分析参考",)
def _suitability_record(result: AgentResult) -> dict[str, Any] | None:
"""取出本轮客服答复附带的适当性裁决留痕(`E-04`);没有则返回 None。
只认 `data['suitability']` 这个形状:其它 Agent 的 `data` 结构各不相同,
这里不做任何猜测式提取,避免把无关字段误记成合规留痕。
"""
data = result.result.data
if not isinstance(data, dict):
return None
candidate = data.get("suitability")
if not isinstance(candidate, dict):
return None
return dict(candidate)
def _governance_rewrote(result: AgentResult) -> bool:
"""治理层是否改写过这次输出(用于审计)。
两条可观测痕迹(都不改协议、只读结果本身):
1. **追加了固定免责声明**:正文末尾出现治理层使用的分隔形状;
2. **拦截并替换**:命中禁用词/硬规则时治理层会把回复换成安全话术并置 `transfer_required`。
保守取值:任一条成立即记 True。它只是审计信息,判错方向的代价是"多标了一次",
不会影响业务行为——因此宁可宽一点,也不为了精确而改动治理协议。
"""
text = result.result.text or ""
appended = any(text.endswith(marker) or marker in text
for marker in _GOVERNANCE_APPEND_MARKERS)
return appended or bool(result.result.transfer_required)
class AgentPersistenceService:
def __init__(self, session: AsyncSession) -> None:
self.session = session
async def complete_run(
self, run_id: str, result: AgentResult, memory_extraction_requested: bool = True,
*, worker_id: str | None = None, profile_candidate_requested: bool = False,
) -> int:
now = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
run = await self.session.scalar(
select(AgentRun).where(AgentRun.run_id == run_id).with_for_update()
)
if run is None:
raise ValueError("run not found")
if result.run_id != run_id:
raise ValueError("result belongs to another run")
if run.status == "succeeded" and run.result_message_id is not None:
return run.result_message_id
if worker_id is not None and (
run.status != "running" or run.worker_id != worker_id
or run.locked_until is None or run.locked_until <= now
):
raise RunLeaseLostError("运行租约失效或已取消")
if run.status not in {"queued", "running"}:
raise RunLeaseLostError("不能覆盖运行终态")
clarification_round = 0
if run.agent_type == "customer_service":
session_row = await self.session.scalar(select(ConversationSession).where(
ConversationSession.session_id == run.session_id,
ConversationSession.user_id == run.user_id,
).with_for_update())
if session_row is not None:
clarification_round = session_row.clarification_round
if result.result.clarification_required:
session_row.clarification_round = min(
session_row.clarification_round + 1, 2
)
else:
session_row.clarification_round = 0
stored_tool_calls: dict[str, Any] = {
"calls": [call.model_dump(mode="json")
for call in result.result.tool_calls],
"transfer_required": bool(result.result.transfer_required),
"transfer_reason": result.result.transfer_reason,
# `E-05`:出口**声明**的主语(产品 / 类目 / 条款)随消息落库,
# 下一轮由 `runtime._conversation_history` 读成
# `ConversationTurn.subject` —— 读侧不再反解回答文本。
"topic": result.result.topic,
}
if result.result.data:
stored_tool_calls["data"] = result.result.data
if result.result.sql:
stored_tool_calls["sql"] = result.result.sql
message = ConversationMessage(
session_id=run.session_id, customer_id=run.user_id, portal="agent",
role="assistant", content=result.result.text,
trace_id=run.trace_id, created_at=now,
intent=result.result.intent.intent if result.result.intent else None,
confidence=(Decimal(str(result.result.intent.confidence))
if result.result.intent else None),
source_references=[ref.model_dump(mode="json")
for ref in result.result.source_references],
# `tool_calls` 是**唯一能承载附加信息的现成 JSON 列**(`conversation_message`
# 没有 `transfer_required` 列,加列要迁移,而规则 4 禁止改既有字段定义)。
# 因此把「转人工标记」作为 `calls` 的**兄弟键**放进来:
# {"calls": [...], "transfer_required": bool, "transfer_reason": str|None}
# 之所以必须落库:`docs/05` §6.3 规定 `GET /agent-runs/{run_id}` 的
# `result` 里要有 `transfer_required` / `transfer_reason`,而它此前
# **既没落库也没出参** —— 前端只能靠"回答里是否含兜底话术开头"来猜要不要转人工
# (`docs/24` 自己把这称为权宜之计)。落库后读写两侧才有同一份真相。
# 读侧允许 `calls` 是裸列表(历史行),见 `RunQueryService.get`。
tool_calls=stored_tool_calls,
)
self.session.add(message)
await self.session.flush()
handover_ticket: HandoverTicket | None = None
handover_context: CustomerServiceHandoverContext | None = None
# `E-04`:适当性裁决要**同时**落消息表(上面 `stored_tool_calls['data']`)
# 与审计表。只落消息表的话,合规要按「谁在什么时候对哪个客户做过风险
# 揭示」去查,就得翻聊天正文 —— 审计表里一行结构化记录才是可核对的凭据。
# 零 DDL:`interaction_audit` 表与 `action_type` 列都已存在。
suitability_record = _suitability_record(result)
if result.result.transfer_required:
# 访客 subject 不是正式用户主键,先按 RiskUser 查询,查不到则保留空归属。
ticket_customer_id = await self.session.scalar(
select(RiskUser.id).where(RiskUser.id == run.user_id)
)
recent_messages = list(await self.session.scalars(
select(ConversationMessage)
.where(ConversationMessage.session_id == run.session_id)
.order_by(ConversationMessage.id.desc())
.limit(MAX_SUMMARY_MESSAGES)
))
recent_messages.reverse()
confidence = (
Decimal(str(result.result.intent.confidence))
if result.result.intent else None
)
# `E-01` ①:**建单白名单只对客服 Agent 生效** —— 其它 Agent 的
# 建单行为一字不改。客服侧的白名单由 `H-04` 收在 `_exit_transfer`
# (`ALL_TRANSFER_REASONS` = 四类白名单 ∪ 治理层合规码,见规则文件)
# (越界即抛错),这里再留一道**响亮告警**:将来若有人又把
# 「答不上来就转人工」塞进客服侧,日志会立刻出现 ERROR。
#
# 为什么**仍然建单**:走到这里说明已经对客户说过「为您转接人工」了,
# 把工单丢掉等于把客户半路扔下 —— 那比原因码不够精确严重得多。
# 因此这里的收敛动作是「纠原因码 + 告警」,不是「拒绝建单」。
if (run.agent_type == CUSTOMER_SERVICE_AGENT_TYPE
and result.result.transfer_reason not in ALL_TRANSFER_REASONS):
logger.error(
"客服转人工原因不在白名单内(回归信号):agent_type=%s "
"transfer_reason=%r",
run.agent_type, result.result.transfer_reason,
)
# `E-01` ③:`reason_code` 是**枚举列**,自由文本一律收敛后再落库。
raw_reason = result.result.transfer_reason
reason_code = normalize_transfer_reason(raw_reason)
if reason_code != raw_reason:
logger.warning(
"转人工原因非枚举码,已收敛:agent_type=%s raw=%r -> %s",
run.agent_type, raw_reason, reason_code,
)
handover_context = build_customer_service_handover_context(
reason_code=reason_code,
clarification_round=clarification_round,
confidence=confidence,
source_references=result.result.source_references,
messages=recent_messages,
)
handover_ticket = HandoverTicket(
ticket_no=f"ticket-{uuid4().hex[:24]}",
session_id=run.session_id,
customer_id=ticket_customer_id,
source_agent=run.agent_type,
source_message_id=message.id,
intent=(result.result.intent.intent if result.result.intent else None),
confidence=confidence,
# `E-01` ②:优先级由原因码映射而来
# (此前从未赋值,恒为模型默认值 `P1`,队列索引白建)。
priority=transfer_priority(reason_code),
reason_code=reason_code,
reason_detail=handover_context.reason_detail,
conversation_summary=handover_context.conversation_summary,
source_references=handover_context.source_references,
status="pending", created_at=now, updated_at=now,
)
self.session.add(handover_ticket)
if suitability_record is not None:
self.session.add(InteractionAudit(
actor_type="agent", actor_id=run.user_id,
target_customer_id=run.user_id, session_id=run.session_id,
portal="agent",
action_type="agent.suitability_disclosed",
detail={
"agent_type": run.agent_type,
"trace_id": run.trace_id,
**suitability_record,
},
created_at=now,
))
run.result_message_id = message.id
run.status = "succeeded"
run.result_version = 1
run.completed_at = now
run.updated_at = now
run.locked_until = None
session_row = await self.session.scalar(select(ConversationSession).where(
ConversationSession.session_id == run.session_id,
ConversationSession.user_id == run.user_id,
).with_for_update())
if session_row is not None and result.result.intent is not None:
session_row.last_intent = result.result.intent.intent
session_row.updated_at = now
self.session.add(InteractionAudit(
actor_type="agent", actor_id=run.user_id, target_customer_id=run.user_id,
session_id=run.session_id, portal="agent", action_type="agent.run_completed",
# `agent_type` 必须落进审计:治理层(`PlatformGovernance.review`)会**改写对外
# 输出**(追加固定免责声明、命中禁用词时整条替换成安全话术),事后要能回答
# "这次改写是哪个 Agent 触发的、改写到了什么程度"。
# `governance_rewrite` 记录治理是否动过输出:正文里出现固定话术的追加形状,
# 或该次运行被标记为需转人工(拦截分支会置 `transfer_required`)。
# 不改表结构:`detail` 是 JSON 列,加键不需要迁移(AGENTS.md 规则 4)。
detail={
"run_id": run_id,
"result_message_id": message.id,
"agent_type": run.agent_type,
"governance_rewrite": _governance_rewrote(result),
},
created_at=now,
))
if handover_ticket is not None:
assert handover_context is not None
self.session.add(InteractionAudit(
actor_type="agent", actor_id=run.user_id,
target_customer_id=handover_ticket.customer_id,
session_id=run.session_id, portal="agent",
action_type="agent.handover_requested",
detail={
"run_id": run_id, "ticket_no": handover_ticket.ticket_no,
"reason_code": handover_ticket.reason_code,
"clarification_round": clarification_round,
"source_reference_count": len(handover_context.source_references),
},
created_at=now,
))
await self.session.execute(
update(RequestIdempotency)
.where(RequestIdempotency.id == run.idempotency_id)
.values(status="completed", result_message_id=message.id, updated_at=now)
)
events = [DomainEvent(
event_id=str(uuid4()), event_type="agent.run_completed", aggregate_type="agent_run",
aggregate_id=run_id, trace_id=run.trace_id,
payload={"run_id": run_id}, occurred_at=now,
)]
if memory_extraction_requested:
events.append(DomainEvent(
event_id=str(uuid4()), event_type="memory.extraction_requested",
aggregate_type="agent_run", aggregate_id=run_id, trace_id=run.trace_id,
payload={"run_id": run_id, "message_id": message.id,
"customer_id": run.user_id}, occurred_at=now,
))
if profile_candidate_requested:
# 候选画像只允许由已登录客服会话触发;Worker 会再次校验身份标记。
events.append(DomainEvent(
event_id=str(uuid4()),
event_type="customer_profile.candidate_requested",
aggregate_type="agent_run", aggregate_id=run_id, trace_id=run.trace_id,
payload={
"run_id": run_id, "message_id": message.id,
"customer_id": run.user_id,
"actor_type": "authenticated_customer",
},
occurred_at=now,
))
if handover_ticket is not None:
assert handover_context is not None
events.append(DomainEvent(
event_id=str(uuid4()), event_type="conversation.transfer_requested",
aggregate_type="conversation", aggregate_id=run.session_id,
trace_id=run.trace_id,
payload={
"ticket_no": handover_ticket.ticket_no,
"handover_context": handover_context.event_metadata,
},
occurred_at=now,
))
for event in events:
self.session.add(DomainEventOutbox(
event_id=event.event_id, event_type=event.event_type,
aggregate_type=event.aggregate_type, aggregate_id=event.aggregate_id,
trace_id=event.trace_id, payload=event.payload, occurred_at=now,
created_at=now, updated_at=now,
))
return message.id