相对第一版 46fc976 的完整变更。组员迁移对照表见 docs/20。
一、对外契约对齐 docs/05(破坏性,共 4 处,组员需按 docs/20 调整)
1) 配置发布端点改为文档规定的复数资源名:submit→validations、
approve→reviews(需 body decision)、activate→activations、
rollback→rollbacks;第一版这 4 个动词式路径 docs/05 从未定义过。
2) 错误码由 8 个笼统码改为 15 个具体语义码(FORBIDDEN→AGENT_PERMISSION_DENIED、
UNAUTHORIZED→AUTHENTICATION_REQUIRED、CONFLICT→RESOURCE_VERSION_CONFLICT、
RESOURCE_NOT_FOUND→RUN_NOT_FOUND/SESSION_NOT_FOUND 等),
输入类错误状态码 400→422。
3) POST /api/v1/agent-runs 与 GET /api/v1/agent-runs/{run_id} 统一为
{data, meta} 信封(data 内字段名与语义未变)。
4) 错误响应体统一为 {error:{code,message,retryable,field_errors}, meta:{trace_id}},
不再返回 FastAPI 默认的 {"detail": ...}。
二、数据库基线与约束
新增 39 张表的基线迁移(链根)与联合唯一键纠偏(4 张表、删 8 增 4,幂等收敛);
撤下 config_release 的双人复核 CHECK(应用层已允许自审,审核节点保留,
自审如实写入 reviewer_id);记忆 active key 生成列与唯一键;
activate 开始记录 supersedes_release_id 使版本链可追溯。
docs/00 基线未修改,未重命名或删除任何表与字段。
三、修复会静默出错或无报错的缺陷
- 跑完集成测试后平台会静默失去生效配置:清理只删自己创建的版本,却没有恢复被它
顶成 superseded 的原生效版本,且审计一并删除因而完全无痕,表现为所有工具被拒
但没有任何报错。已修清理逻辑并加恢复。
- Worker 单轮异常导致进程退出;记忆抽取调用方的“事务已开始”异常;
召回缓存丢失 degraded 标记;连接时区未生效导致 created_at/updated_at 差 8 小时;
.env 与 os.getenv 密钥来源分裂导致“没有可用的已批准模型端点”。
- 记忆信号识别漏判与跨键误命中;SSE 未带 Accept 的协商行为。
四、功能补齐
记忆链路 P1/P2/P3(抽取、受控词表、召回与缓存、生命周期级联及投影事件)、
fin_* 场内交易只读 ORM 层、agent_intent_config 状态流转并在运行期真正生效、
限流(Redis 固定窗口、故障一律放行)、游标校验、trace_id 中间件、
示例业务 Agent fund_query_demo 与一键端到端验证脚本,以及审计/指纹/迁移状态工具。
五、文档与验证
新增 docs/19(业务 Agent 接入实操)、docs/20(第一版迁移指南)与 docs/evidence 证据;
docs/01/02/06/08/09/17 同步实现现状。
验证结果:ruff 通过、mypy 103 文件无错、unit+contract 447 passed、
integration 29 passed、acceptance_check --production 7 PASS、
demo_agent_e2e 9/9 PASS(含失败关闭反证)。
329 lines
14 KiB
Python
329 lines
14 KiB
Python
"""记忆链路端到端探针:受理 → Worker 执行 → complete_run → 事件消费 → 抽取 → memory_unit。
|
||
|
||
为什么需要它:记忆链路上的缺陷(Worker 未注册消费、事件不携带正文而消费者期望正文、
|
||
抽取语义与幂等错误、抽取被静默跳过)在业务 Agent 接入前**不会通过现有数据自然暴露**——
|
||
`jr` 库里 `memory_unit` 一直是 0 行。探针使用生产装配(`get_agent_factory()`、真实 MySQL、
|
||
真实 `WorkerRuntime` 租约与治理链、真实 `OutboxWorker` 消费)跑完整链路并逐项断言。
|
||
|
||
判据(见 `docs/evidence/20260909-memory-baseline-before.md` 与 `-memory-chain-acceptance.md`):
|
||
1. 运行成功后存在 `memory.extraction_requested` 事件,payload **不含正文**,只有定位字段;
|
||
2. 消费后记忆由**抽取结果**产生:受控语义键 + 结构化值,而不是用户原文整句;
|
||
3. 同一 `event_id` 重复消费不产生第二条记忆(幂等边界)。
|
||
|
||
关于装配:生产链路注入了 `IntentClassifier`,而意图分类需要已配置的模型端点;
|
||
本地 `model_endpoint_config` 为 0 行时按设计失败关闭,run 无法进入 `complete_run`。
|
||
因此探针用最小装配工厂(只注入治理)跳过意图分类,并**注入确定性替身模型服务**
|
||
(`WorkerRuntime(model_service=...)`)驱动抽取,使本探针既不依赖真实模型端点,
|
||
又能覆盖"抽取 → 记忆"这一段。生产工厂的装配完整性单独断言。
|
||
|
||
用法:
|
||
python tools/memory_chain_probe.py # 执行、断言并清理
|
||
python tools/memory_chain_probe.py --keep # 保留测试数据以便排查
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import sys
|
||
from uuid import uuid4
|
||
|
||
from sqlalchemy import and_, delete, select
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.core.contracts import AgentDefinition, AgentRequest, CoreResult, RequestContext
|
||
from app.core.errors import AgentTypeNotFoundError
|
||
from app.infrastructure.db import SessionFactory
|
||
from app.model.audit import InteractionAudit
|
||
from app.model.conversation import ConversationMessage
|
||
from app.model.memory import MemoryEvidence, MemoryUnit
|
||
from app.model.platform import AgentRun, DomainEventOutbox, OutboxDelivery, RequestIdempotency
|
||
from app.service.agent.base import BaseAgent
|
||
from app.service.agent.bootstrap import get_agent_factory
|
||
from app.service.agent.factory import AgentFactory
|
||
from app.service.agent.governance import AgentGovernance, PlatformGovernance
|
||
from app.service.agent_run_application_service import AgentRunApplicationService
|
||
from app.service.identity_service import IdentityService
|
||
from app.worker.memory_extraction_worker import MemoryExtractionWorker
|
||
from app.worker.runtime import WorkerRuntime
|
||
|
||
CUSTOMER_ID = 9001
|
||
AGENT_TYPE = "memory_chain_probe"
|
||
MESSAGE = "请记住:我的风险偏好是稳健型,后续建议请按稳健型说明。"
|
||
SESSION_PREFIX = "probe-"
|
||
|
||
# 抽取结果标识:P2 之后记忆键来自受控词表、内容是结构化值,不再以会话前缀命名,
|
||
# 因此清理与计数按"受控键 + 值"精确识别探针产物。
|
||
PROBE_MEMORY_KEY = "preference:risk_level"
|
||
PROBE_MEMORY_VALUE = "稳健型"
|
||
PROBE_MEMORY_TYPE = "preference"
|
||
EXTRACTION_JSON = (
|
||
'{"memory_key": "preference:risk_level", "value": "稳健型", '
|
||
'"memory_type": "preference", "confidence": 0.9}'
|
||
)
|
||
|
||
failures: list[str] = []
|
||
|
||
|
||
def check(condition: bool, description: str) -> None:
|
||
print((" PASS " if condition else " FAIL ") + description)
|
||
if not condition:
|
||
failures.append(description)
|
||
|
||
|
||
class MemoryProbeAgent(BaseAgent):
|
||
"""最小探针 Agent:只声明意图,不调用任何工具。"""
|
||
|
||
definition = AgentDefinition(
|
||
agent_type=AGENT_TYPE,
|
||
version="probe-1",
|
||
allowed_roles=("customer",),
|
||
allowed_portals=("api",),
|
||
allowed_tools=(),
|
||
supported_intents=("general",),
|
||
)
|
||
|
||
async def handle(self, request: AgentRequest, context: RequestContext) -> CoreResult:
|
||
del request, context
|
||
return CoreResult(text="已记录您的偏好说明,后续将按稳健型为您解释。")
|
||
|
||
|
||
class _StubExecution:
|
||
def __init__(self, text: str) -> None:
|
||
self.text = text
|
||
|
||
|
||
class StubModelService:
|
||
"""确定性替身模型:只回放抽取用的严格 JSON,不访问网络。"""
|
||
|
||
def __init__(self, payload: str = EXTRACTION_JSON) -> None:
|
||
self.payload = payload
|
||
self.calls = 0
|
||
|
||
async def generate(self, endpoints: object, prompt: str, **kwargs: object) -> _StubExecution:
|
||
del endpoints, prompt, kwargs
|
||
self.calls += 1
|
||
return _StubExecution(self.payload)
|
||
|
||
|
||
class StubEndpointResolver:
|
||
"""确定性替身端点解析器:返回一个占位端点。
|
||
|
||
`MemoryExtractionService` 在调用模型前必须先解析出可用端点,端点缺失时按设计
|
||
失败关闭。本地 `model_endpoint_config` 为 0 行,因此探针注入本替身,使抽取路径
|
||
可被完整验证,而不需要往库里塞一条假的模型端点配置。
|
||
"""
|
||
|
||
def __init__(self) -> None:
|
||
self.calls = 0
|
||
|
||
async def resolve(self, *, agent_type: str, task_type: str) -> list[object]:
|
||
del agent_type, task_type
|
||
self.calls += 1
|
||
return [object()]
|
||
|
||
|
||
def build_probe_factory() -> AgentFactory:
|
||
"""最小装配:只注入治理,避开需要模型端点的意图分类。"""
|
||
governance: AgentGovernance = PlatformGovernance()
|
||
factory = AgentFactory(governance=governance)
|
||
try:
|
||
factory.definition(AGENT_TYPE)
|
||
except AgentTypeNotFoundError:
|
||
factory.register(
|
||
MemoryProbeAgent.definition,
|
||
lambda _context: MemoryProbeAgent(MemoryProbeAgent.definition),
|
||
)
|
||
return factory
|
||
|
||
|
||
def _probe_memory_filter():
|
||
return and_(
|
||
MemoryUnit.customer_id == CUSTOMER_ID,
|
||
MemoryUnit.memory_key == PROBE_MEMORY_KEY,
|
||
)
|
||
|
||
|
||
async def purge_probe_residue() -> None:
|
||
"""清理历史探针残留,使本探针可重复运行。"""
|
||
async with SessionFactory() as session, session.begin():
|
||
memories = (await session.scalars(select(MemoryUnit).where(_probe_memory_filter()))).all()
|
||
for memory in memories:
|
||
await session.execute(
|
||
delete(MemoryEvidence).where(MemoryEvidence.memory_id == memory.id)
|
||
)
|
||
await session.delete(memory)
|
||
probe_runs = (
|
||
await session.scalars(
|
||
select(AgentRun).where(AgentRun.session_id.like(f"{SESSION_PREFIX}%"))
|
||
)
|
||
).all()
|
||
for run in probe_runs:
|
||
events = (
|
||
await session.scalars(
|
||
select(DomainEventOutbox.event_id).where(
|
||
DomainEventOutbox.aggregate_id == run.run_id
|
||
)
|
||
)
|
||
).all()
|
||
if events:
|
||
await session.execute(
|
||
delete(OutboxDelivery).where(OutboxDelivery.event_id.in_(list(events)))
|
||
)
|
||
await session.execute(
|
||
delete(DomainEventOutbox).where(DomainEventOutbox.aggregate_id == run.run_id)
|
||
)
|
||
await session.execute(
|
||
delete(RequestIdempotency).where(RequestIdempotency.id == run.idempotency_id)
|
||
)
|
||
await session.execute(
|
||
delete(ConversationMessage).where(
|
||
ConversationMessage.session_id == run.session_id
|
||
)
|
||
)
|
||
await session.execute(
|
||
delete(InteractionAudit).where(InteractionAudit.session_id == run.session_id)
|
||
)
|
||
await session.delete(run)
|
||
|
||
|
||
async def count_memories(session: AsyncSession) -> int:
|
||
memories = (await session.scalars(select(MemoryUnit).where(_probe_memory_filter()))).all()
|
||
return len(memories)
|
||
|
||
|
||
async def main() -> int:
|
||
keep = "--keep" in sys.argv
|
||
session_id = f"{SESSION_PREFIX}{uuid4()}"
|
||
idempotency_key = f"probe{uuid4().hex}"
|
||
context = RequestContext(
|
||
user_id=str(CUSTOMER_ID),
|
||
trace_id=str(uuid4()),
|
||
roles=("customer",),
|
||
permissions=("agent:run",),
|
||
portal="api",
|
||
)
|
||
|
||
print("1) 前置检查与生产装配核对")
|
||
await purge_probe_residue()
|
||
check(True, "历史探针残留已清理")
|
||
production = get_agent_factory()
|
||
check(production._model_service is not None, "生产工厂已注入模型服务")
|
||
check(production._tool_executor is not None, "生产工厂已注入工具执行器")
|
||
check(production._intent_classifier is not None, "生产工厂已注入意图分类器")
|
||
resolved = await IdentityService().resolve(context)
|
||
if "agent:run" not in resolved.permissions:
|
||
print(" FAIL 客户 9001 缺少 agent:run 权限;请先运行 tools/seed_test_rbac.py")
|
||
return 1
|
||
check(True, f"客户 {CUSTOMER_ID} 实时权限已加载(角色 {resolved.roles})")
|
||
|
||
factory = build_probe_factory()
|
||
stub = StubModelService()
|
||
resolver = StubEndpointResolver()
|
||
check(True, f"探针工厂已装配(最小治理装配:{AGENT_TYPE})")
|
||
|
||
print("2) 受理运行")
|
||
request = AgentRequest(
|
||
agent_type=AGENT_TYPE,
|
||
message=MESSAGE,
|
||
session_id=session_id,
|
||
idempotency_key=idempotency_key,
|
||
)
|
||
async with SessionFactory() as session:
|
||
accepted = await AgentRunApplicationService(session, factory).accept(request, resolved)
|
||
run_id = accepted.run_id
|
||
check(bool(run_id), f"受理成功 run_id={run_id}")
|
||
|
||
print("3) Worker 执行(真实租约、治理链与落库)")
|
||
runtime = WorkerRuntime(factory=factory, model_service=stub, endpoint_resolver=resolver)
|
||
executed = await runtime.execute(run_id)
|
||
check(executed, "Worker 领取并执行完成")
|
||
async with SessionFactory() as session:
|
||
run = await session.scalar(select(AgentRun).where(AgentRun.run_id == run_id))
|
||
check(
|
||
run is not None and run.status == "succeeded",
|
||
f"运行终态 = {run.status if run is not None else None}"
|
||
+ (f"(error_code={run.error_code})" if run is not None and run.error_code else ""),
|
||
)
|
||
|
||
print("4) 事件契约检查")
|
||
async with SessionFactory() as session:
|
||
event = await session.scalar(
|
||
select(DomainEventOutbox).where(
|
||
DomainEventOutbox.aggregate_id == run_id,
|
||
DomainEventOutbox.event_type == "memory.extraction_requested",
|
||
)
|
||
)
|
||
check(event is not None, "complete_run 在同一事务写入了 memory.extraction_requested")
|
||
if event is None:
|
||
print(" 说明:运行未成功时不会产生记忆事件,后续断言一并失败属预期")
|
||
payload: dict[str, object] = {}
|
||
event_id = run_id
|
||
else:
|
||
payload = dict(event.payload)
|
||
event_id = event.event_id
|
||
check(bool(payload) and "content" not in payload, "事件 payload 不携带正文(只带定位信息)")
|
||
check(
|
||
{"message_id", "customer_id"} <= set(payload),
|
||
f"事件定位字段完整:{sorted(payload)}",
|
||
)
|
||
|
||
print("5) 消费事件并核对抽取结果")
|
||
rounds = 0
|
||
while rounds < 10:
|
||
if not await runtime.dispatch_one(run_id=run_id):
|
||
break
|
||
rounds += 1
|
||
check(rounds >= 1, f"Worker 轮询 {rounds} 轮后该 run 的事件队列清空")
|
||
check(stub.calls >= 1, f"抽取模型被真实调用({stub.calls} 次)")
|
||
check(resolver.calls >= 1, f"抽取端点经解析器解析({resolver.calls} 次)")
|
||
async with SessionFactory() as session:
|
||
before = await count_memories(session)
|
||
memory = await session.scalar(select(MemoryUnit).where(_probe_memory_filter()))
|
||
check(before >= 1, f"memory_unit 出现探针记忆({before} 行)")
|
||
check(
|
||
memory is not None and memory.memory_key == PROBE_MEMORY_KEY,
|
||
f"记忆键来自受控词表:{memory.memory_key if memory is not None else None}",
|
||
)
|
||
check(
|
||
memory is not None and (memory.content or "") == PROBE_MEMORY_VALUE,
|
||
f"记忆内容是抽取的结构化值而非用户原文:{memory.content if memory is not None else None}",
|
||
)
|
||
check(
|
||
memory is not None and memory.memory_type == PROBE_MEMORY_TYPE,
|
||
f"记忆类型与键前缀同构:{memory.memory_type if memory is not None else None}",
|
||
)
|
||
check(memory is not None and memory.customer_id == CUSTOMER_ID, "记忆归属为发起运行的客户")
|
||
|
||
print("6) 幂等:同一 event_id 重复消费")
|
||
async with SessionFactory() as session:
|
||
# 复用同一抽取器:这样重复消费若未被幂等拦截就会真的走抽取并写入,
|
||
# 断言才有意义(而不是因为抽取不可用而"恰好"没写)。
|
||
worker = MemoryExtractionWorker(session, extractor=runtime.memory_extraction)
|
||
again = await worker.handle(payload, event_id=event_id)
|
||
await session.commit()
|
||
after = await count_memories(session)
|
||
check(again is False, "重复消费被幂等边界拦截")
|
||
check(after == before, f"重复消费未新增记忆({before} → {after})")
|
||
|
||
if keep:
|
||
print(f"\n--keep 已启用:保留测试数据 session_id={session_id} run_id={run_id}")
|
||
else:
|
||
print("7) 清理测试数据")
|
||
await purge_probe_residue()
|
||
async with SessionFactory() as session:
|
||
remaining = await count_memories(session)
|
||
check(remaining == 0, f"探针记忆已清理(剩余 {remaining} 行)")
|
||
|
||
print()
|
||
if failures:
|
||
print(f"FAILED: {len(failures)} 项未通过")
|
||
for item in failures:
|
||
print(f" - {item}")
|
||
return 1
|
||
print("PASSED: 记忆链路端到端连通(受理 → 执行 → 事件 → 消费 → 抽取 → 记忆 → 幂等)")
|
||
return 0
|
||
|
||
|
||
if __name__ == "__main__":
|
||
raise SystemExit(asyncio.run(main()))
|