Files
group_fqcd_jr/app/service/agent_run_application_service.py
张胜宇 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

202 lines
9.3 KiB
Python
Raw Permalink 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 hashlib
import json
from collections.abc import Sequence
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from uuid import uuid4
from sqlalchemy import select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.actor import VISITOR_ACTOR_TYPE, is_visitor
from app.core.contracts import AgentRequest, DomainEvent, RequestContext
from app.core.conversation_privacy import sanitize_customer_service_message
from app.core.customer_service_rules import chitchat_streak
from app.core.errors import (
ForbiddenAgentError,
IdempotencyConflictError,
SessionNotAccessibleError,
)
from app.model.audit import InteractionAudit
from app.model.conversation import ConversationMessage
from app.model.platform import AgentRun, RequestIdempotency
from app.model.session import ConversationSession
from app.repository.conversation_repository import ConversationRepository
from app.repository.outbox_repository import OutboxRepository
from app.service.agent.bootstrap import get_agent_factory
from app.service.agent.factory import AgentFactory
@dataclass(frozen=True)
class RunAccepted:
run_id: str
trace_id: str
status: str = "queued"
def build_outbox_metadata(
request: AgentRequest, prior_user_messages: Sequence[str], *,
clarification_round: int = 0, session_context: Sequence[str] = (),
) -> dict[str, object]:
"""构造 Worker 使用的内部元数据,**不信任外部传入的客服计数**。
客服的三个字段一律由服务端覆写:
- ``chitchat_streak``:由已落库的历史重算(客户端写 0 就能绕过连续闲聊收口);
- ``clarification_round``:来自会话行,是 E5a 澄清的轮次上限依据;
- ``session_context``:当前会话内的已脱敏上下文。
其余 Agent 一直是把 ``request.metadata`` 原样透出,这里保持与它们相同的行为。
上限收紧的原因:``model_copy(update=...)`` **不做校验**,越界值会被静默写入
(``clarification_round`` 契约上是 ``le=2``、``session_context`` 是 ``max_length=6``)。
"""
metadata = request.metadata
if request.agent_type == "customer_service":
metadata = metadata.model_copy(update={
"chitchat_streak": chitchat_streak(prior_user_messages, request.message),
"clarification_round": min(max(clarification_round, 0), 2),
"session_context": tuple(session_context[-6:]),
})
return metadata.model_dump(mode="json")
class AgentRunApplicationService:
def __init__(
self, session: AsyncSession, factory: AgentFactory | None = None,
) -> None:
self.session = session
self.factory = factory if factory is not None else get_agent_factory()
async def _recent_user_messages(
self, request: AgentRequest, user_id: int, *, limit: int = 3
) -> tuple[str, ...]:
"""取本会话最近若干条**已脱敏**的用户消息,供连续闲聊计数使用。
必须走 `ConversationRepository`:它按 `customer_id == user_id` 过滤。一期实现只按
`session_id` 取历史、跨主体可读——那正是 `F-01` 要修的缺陷本体,所以这里不是
"恢复旧实现",而是在**新链路**上用带主体过滤的查询重写(与 Worker 侧的
`runtime._conversation_history` 同源同口径)。
"""
rows = await ConversationRepository(self.session).messages(
request.session_id, user_id, limit
)
messages = tuple(
str(row.content or "").strip() for row in reversed(rows)
if str(row.role) == "user" and str(row.content or "").strip()
)
return messages[-limit:]
async def accept(self, request: AgentRequest, context: RequestContext) -> RunAccepted:
try:
self.factory.authorize(request.agent_type, context)
except ForbiddenAgentError:
async with self.session.begin():
self.session.add(InteractionAudit(
actor_type="user", actor_id=int(context.user_id), portal=context.portal,
action_type="agent.access_denied", session_id=request.session_id,
detail={"agent_type": request.agent_type, "trace_id": context.trace_id},
created_at=datetime.now(UTC).replace(tzinfo=None),
))
raise
# 客服原文不进入会话与异步链路:**先脱敏**,再用脱敏后的形状参与幂等哈希——
# 否则同一请求的两次提交会因「原文 vs 脱敏文本」算出两个哈希而互相冲突。
stored_request = request
if request.agent_type == "customer_service":
stored_request = request.model_copy(update={
"message": sanitize_customer_service_message(request.message)
})
user_id = int(context.user_id)
request_hash = hashlib.sha256(
json.dumps(stored_request.model_dump(mode="json"), sort_keys=True).encode("utf-8")
).hexdigest()
now = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
session_row = await self.session.scalar(
select(ConversationSession).where(
ConversationSession.session_id == request.session_id,
ConversationSession.user_id == user_id,
).with_for_update()
)
if session_row is not None:
if session_row.status != "active" or session_row.agent_type != request.agent_type:
raise ForbiddenAgentError("会话不可用于当前 Agent")
session_row.message_count += 1
session_row.last_active_at = now
owner = await self.session.scalar(
select(ConversationMessage.customer_id)
.where(ConversationMessage.session_id == request.session_id)
.where(ConversationMessage.customer_id.is_not(None))
.limit(1)
)
if owner is not None and owner != user_id:
raise SessionNotAccessibleError("会话不属于当前用户")
existing = await self.session.scalar(
select(RequestIdempotency).where(
RequestIdempotency.user_id == user_id,
RequestIdempotency.agent_type == request.agent_type,
RequestIdempotency.idempotency_key == request.idempotency_key,
)
)
if existing is not None:
if existing.request_hash != request_hash:
raise IdempotencyConflictError("同一幂等键对应不同请求")
run = await self.session.scalar(
select(AgentRun).where(AgentRun.idempotency_id == existing.id)
)
if run is None:
raise RuntimeError("idempotency record has no run")
return RunAccepted(run.run_id, run.trace_id, run.status)
clarification_round = (
session_row.clarification_round if session_row is not None else 0
)
prior_user_messages = (
await self._recent_user_messages(request, user_id)
if request.agent_type == "customer_service"
else ()
)
outbox_metadata = build_outbox_metadata(
stored_request, prior_user_messages,
clarification_round=clarification_round,
)
trace_id = context.trace_id
message = ConversationMessage(
session_id=request.session_id, customer_id=user_id, portal="api",
role="user", content=stored_request.message, trace_id=trace_id, created_at=now,
)
self.session.add(message)
await self.session.flush()
idem = RequestIdempotency(
user_id=user_id, session_id=request.session_id, agent_type=request.agent_type,
idempotency_key=request.idempotency_key, request_hash=request_hash,
trace_id=trace_id, expire_at=now + timedelta(hours=24),
created_at=now, updated_at=now,
)
self.session.add(idem)
try:
await self.session.flush()
except IntegrityError as exc:
raise IdempotencyConflictError("幂等键正在被并发请求占用") from exc
run_id = str(uuid4())
run = AgentRun(
run_id=run_id, idempotency_id=idem.id, session_id=request.session_id,
user_id=user_id, agent_type=request.agent_type, trace_id=trace_id,
request_message_id=message.id, created_at=now, updated_at=now,
)
self.session.add(run)
await self.session.flush()
await OutboxRepository(self.session).append(DomainEvent(
event_id=str(uuid4()), event_type="agent.run_requested", aggregate_type="agent_run",
aggregate_id=run_id, trace_id=trace_id,
payload={
"run_id": run_id,
"actor_type": VISITOR_ACTOR_TYPE if is_visitor(context) else "authenticated",
"metadata": outbox_metadata,
},
occurred_at=now,
))
return RunAccepted(run_id, trace_id)