Files
group_fqcd_jr/app/service/agent_persistence_service.py
T
qyqy 928d0bcea3 chore: 按评审恢复 5 份文档、审计补 agent_type、修订环境口径
评审意见落地(§3.3 驳回删除 / §3.1 审计补充 / §1.4 解释器 / §2 编号):
- docs/04、06、10、13、99 全部恢复(评审:删除收益为零、保留成本同样为零);
  AGENTS.md 改为"保留但仅作历史参考"并列入 D 类,ARCHIVE 归档说明加作废声明
- 治理审计补留痕:interaction_audit.detail 增加 agent_type 与 governance_rewrite
  (治理层会改写对外输出,事后必须能追溯到是哪个 Agent 触发的;不改表结构,detail 是 JSON 列)
  新增 tests/unit/service/test_agent_persistence_audit.py 锁住该契约
- AGENTS.md 修订环境口径:解释器各用本机可用的那个(.venv 被 gitignore、不进仓库);
  config_release 与 Milvus schema 均属环境数据、不随代码合并,相关结论必须带环境限定;
  测试基线 1034;mypy 数字双方不可比(本机未装 sqlalchemy2-stubs,报错集中在模型层)
- docs/26 JWT 文档:因 21 已被风控迁移清单占用而改名,PR 描述里会单独说明
2026-09-11 18:42:46 +08:00

123 lines
6.4 KiB
Python

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={"calls": [call.model_dump(mode="json")
for call in result.result.tool_calls]},
)
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