Files
group_fqcd_jr/app/service/agent_persistence_service.py
T
qyqy a7e2d3eac0 fix(customer-service): 统一客服热线来源;补 transfer_required 出参(docs/05 §6.3 未兑现的一半)
两个都是"代码里存在但没接对"的缺陷,都不是新功能。

## A1 客服热线在代码里有两个值(一个出口给假号码)

- `app/core/customer_service_rules.py:35` `CONTACT_PHONE = "15936583816"` ← 真号码,安全路由 6 处在用
- `app/service/agent/implementations/customer_service.py:155` `HOTLINE = "400-XXX-XXXX"` ← 占位符,兜底出口在用

后果:**同一个客服给客户两个不同的电话号码**。问"风险等级怎么划分"被安全路由处理时给真号码;
问一个知识库答不了的问题走兜底时给 `400-XXX-XXXX` —— 客户按这个号码永远打不通。

修法:`HOTLINE` / `SERVICE_HOURS` 改为**转发** `customer_service_rules` 的两个常量
(不是"改成相同的值",而是引用同一对象,避免日后再次漂移);工作时间也随之从
"每日 7:00-22:00" 统一为 "工作日 09:00-18:00"(与安全路由出口一致)。
新增守卫测试用 `is` 断言对象同一性 —— 值相等挡不住"两边各写一份恰好相同"的漂移。

## A2 `transfer_required` 既没落库也没出参

`docs/05` §6.3 一直规定 `GET /agent-runs/{run_id}` 的 `result` 里有
`transfer_required` / `transfer_reason`,但实现里两个都没有:前端判断"这轮要不要转人工"
只能靠**猜正文里有没有兜底话术的开头**(`docs/24` 自己把这称为权宜之计)。

- 写入侧:`conversation_message` **没有** `transfer_required` 列,加列要迁移且规则 4 禁止改既有
  字段定义 ⇒ 放进 `tool_calls` 这个现成 JSON 列,作为 `calls` 的兄弟键
  (`{"calls": [...], "transfer_required": bool, "transfer_reason": str|None}`)
- 读取侧:`RunQueryService.get` 取出来放进 `result`;**兼容历史行**(`calls` 裸列表 / None →
  按 False/None 处理,不抛异常、也不凭正文猜)

刻意**没做**的一半:`docs/05` §6.3 的 `result` 里还有 `degraded` / `degradation_reason`,
但 `CoreResult` 里根本没有这两个字段(降级信息目前只在工具出参里)—— 补它要改
`CoreResult` 并让各 Agent 传递降级状态,属另一个改动范围。**已在交付说明里注明这一半仍缺。**

## 真机验证

| 问题 | transfer_required | transfer_reason | 正文电话 |
|---|---|---|---|
「请介绍一下量子纠缠在基金估值中的应用」 | **True** | 置信度不足:score=0.571 gap=0.004 | 15936583816 ✅ |
「请帮我计算一下三体问题的数值解」 | **True** | 置信度不足:score=0.499 gap=0.011 | 15936583816 ✅ |
「基金申购后多久确认」(正常知识直返) | False | — | 无(正确) |
「你们公司明天会下雪吗」(闲聊出口) | False | — | 无(正确) |

(第一次我用"下雪"当兜底用例,结果它被闲聊出口正确接住了 —— 是我的期望值写错,不是代码问题。)

门禁:测试 1223 passed(新增 4 个用例)/ 3 failed(均为已知非代码缺陷)/ mypy 0 错。
2026-09-11 21:17:47 +08:00

136 lines
7.4 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.
from datetime import UTC, datetime
from decimal import Decimal
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.errors import RunLeaseLostError
from app.model.audit import InteractionAudit
from app.model.conversation import ConversationMessage
from app.model.platform import AgentRun, DomainEventOutbox, RequestIdempotency
#: 治理层追加免责声明时使用的分隔形状(`app/service/agent/governance.py` 里定义)。
#: 这里只用于**审计留痕**,不参与任何判定:判据是"末尾是否出现这个形状"。
_GOVERNANCE_APPEND_MARKERS: tuple[str, ...] = ("\n\n本内容仅为投资分析参考",)
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,
) -> 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("不能覆盖运行终态")
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={
"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,
},
)
self.session.add(message)
await self.session.flush()
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
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,
))
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,
))
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