Files
group_fqcd_jr/tests/unit/service/test_memory_service.py
T
lzf_0626 6516ccb385 feat: 第二版——接口契约对齐 docs/05,修复静默故障与数据库基线
相对第一版 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(含失败关闭反证)。
2026-09-10 15:55:54 +08:00

210 lines
7.8 KiB
Python

from datetime import UTC, datetime
from typing import Any
from unittest.mock import AsyncMock, Mock
import pytest
from sqlalchemy.ext.asyncio import AsyncSession
from app.model.memory import MemoryUnit
from app.service.memory_recall_service import CACHE_KEY_PREFIX, MemoryRecallService
from app.service.memory_service import MemoryService
NOW = datetime(2026, 9, 9, 0, 0, 0, tzinfo=UTC).replace(tzinfo=None)
class RecordingCache:
"""缓存替身:记录被删除的键,可注入删除故障(故障必须不阻塞写入)。"""
def __init__(self, *, fail: bool = False) -> None:
self.fail = fail
self.deleted: list[tuple[str, ...]] = []
async def delete(self, *keys: str) -> int:
self.deleted.append(tuple(keys))
if self.fail:
raise ConnectionError("redis unavailable")
return len(keys)
def memory(
memory_id: int, *, customer_id: int = 7, key: str = "conversation.session-1",
content: str = "旧内容", version: int = 1,
) -> MemoryUnit:
return MemoryUnit(
id=memory_id, memory_uuid=f"uuid-{memory_id}", customer_id=customer_id, memory_key=key,
content=content, memory_type="事实候选", source_type="用户自述", source_confidence=0.65,
confidence=0.65, evidence_count=0, conflict_count=0, recall_count=0, status="active",
valid_from=NOW, version=version, created_at=NOW, updated_at=NOW,
)
def fake_session() -> Any:
session = AsyncMock(spec=AsyncSession)
nested = Mock()
nested.__aenter__ = AsyncMock(return_value=None)
nested.__aexit__ = AsyncMock(return_value=False)
session.begin_nested = Mock(return_value=nested)
session.add = Mock()
return session
@pytest.mark.asyncio
async def test_recall_filters_by_customer() -> None:
"""召回必须带 customer_id 过滤,不能跨客户命中。"""
session = fake_session()
session.scalars.return_value = [memory(1)]
await MemoryService(session).recall(7)
statement = str(session.scalars.await_args.args[0])
assert "memory_unit.customer_id" in statement
assert "memory_unit.status" in statement
@pytest.mark.asyncio
async def test_upsert_conflict_does_not_reference_itself() -> None:
existing = memory(5, content="旧内容")
session = fake_session()
session.scalar.side_effect = [existing, None]
added: list[Any] = []
session.add = Mock(side_effect=added.append)
updated = await MemoryService(session).upsert(7, "conversation.session-1", "新内容")
conflict = next(item for item in added if type(item).__name__ == "MemoryConflict")
assert conflict.left_memory_id == 5
assert conflict.right_memory_id != conflict.left_memory_id
assert conflict.status == "auto_resolved"
assert conflict.severity == "low"
assert conflict.resolution is not None
assert conflict.winner_memory_id == 5
assert updated is existing
assert updated.content == "新内容"
assert updated.conflict_count == 1
assert updated.version == 2
@pytest.mark.asyncio
async def test_upsert_leaves_primary_key_to_database() -> None:
"""新建记忆不得显式写 id=0,主键交给自增列。"""
session = fake_session()
session.scalar.side_effect = [None, None]
added: list[Any] = []
session.add = Mock(side_effect=added.append)
created = await MemoryService(session).upsert(7, "conversation.session-1", "新内容")
assert created.id is None
assert created.status == "active"
assert created.customer_id == 7
@pytest.mark.asyncio
async def test_upsert_persists_structured_value_when_created() -> None:
"""新建记忆时结构化值随受控键一起落库,调用方不再需要落原文。"""
session = fake_session()
session.scalar.side_effect = [None, None]
added: list[Any] = []
session.add = Mock(side_effect=added.append)
created = await MemoryService(session).upsert(
7, "preference:risk_level", "稳健型", memory_type="preference", confidence=0.9,
structured_value={"memory_key": "preference:risk_level", "value": "稳健型"},
)
assert created.structured_value == {
"memory_key": "preference:risk_level", "value": "稳健型"}
assert created.memory_type == "preference"
assert created.memory_key == "preference:risk_level"
@pytest.mark.asyncio
async def test_upsert_refreshes_structured_value_on_existing_memory() -> None:
existing = memory(5, key="preference:risk_level", content="保守型")
session = fake_session()
session.scalar.side_effect = [existing, 6]
session.add = Mock()
updated = await MemoryService(session).upsert(
7, "preference:risk_level", "稳健型", memory_type="preference",
structured_value={"value": "稳健型"},
)
assert updated is existing
assert updated.content == "稳健型"
assert updated.structured_value == {"value": "稳健型"}
@pytest.mark.asyncio
async def test_record_evidence_is_idempotent_by_key() -> None:
"""同一 idempotency_key 重复写入不再新增证据、不再增加证据计数。"""
target = memory(5)
session = fake_session()
session.scalar.side_effect = [77]
added: list[Any] = []
session.add = Mock(side_effect=added.append)
recorded = await MemoryService(session).record_evidence(
target, idempotency_key="memory.extraction_requested:event-1", evidence_type="对话",
excerpt="我偏好低风险", snapshot=None, weight=0.05)
assert recorded is False
assert added == []
assert target.evidence_count == 0
# A1:写入路径此前完全不失效召回热缓存,新记忆在 TTL(300 秒)内召回不到。
@pytest.mark.asyncio
async def test_upsert_invalidates_customer_recall_cache() -> None:
"""新建记忆后必须删除该客户的召回热缓存键,且键集与召回服务完全一致。"""
session = fake_session()
session.scalar.side_effect = [None, None]
session.add = Mock()
cache = RecordingCache()
created = await MemoryService(session, cache=cache).upsert(
7, "preference:risk_level", "稳健型")
assert created.memory_key == "preference:risk_level"
assert cache.deleted == [tuple(MemoryRecallService.cache_keys(7))]
assert all(key.startswith(f"{CACHE_KEY_PREFIX}:7:") for key in cache.deleted[0])
@pytest.mark.asyncio
async def test_update_path_invalidates_customer_recall_cache() -> None:
"""更新既有记忆同样改变召回结果,必须一并失效缓存。"""
existing = memory(5, key="preference:risk_level", content="保守型")
session = fake_session()
session.scalar.side_effect = [existing, 6]
session.add = Mock()
cache = RecordingCache()
await MemoryService(session, cache=cache).upsert(7, "preference:risk_level", "稳健型")
assert cache.deleted == [tuple(MemoryRecallService.cache_keys(7))]
@pytest.mark.asyncio
async def test_invalidate_memory_drops_recall_cache() -> None:
"""单条失效也会让已失效记忆在 TTL 内继续被召回,因此同样需要失效缓存。"""
session = fake_session()
session.scalar.side_effect = [memory(5)]
cache = RecordingCache()
assert await MemoryService(session, cache=cache).invalidate("uuid-5", 7) is True
assert cache.deleted == [tuple(MemoryRecallService.cache_keys(7))]
@pytest.mark.asyncio
async def test_cache_failure_never_blocks_memory_write() -> None:
"""缓存只是可重建的加速层:删除失败不得阻塞写入主流程。"""
session = fake_session()
session.scalar.side_effect = [None, None]
session.add = Mock()
cache = RecordingCache(fail=True)
created = await MemoryService(session, cache=cache).upsert(
7, "preference:risk_level", "稳健型")
assert created.content == "稳健型"
assert cache.deleted # 失效动作尝试过,但异常被吞
assert await MemoryService(session, cache=None).invalidate_recall_cache(7) == 0