diff --git a/app/service/customer_profile_candidate_service.py b/app/service/customer_profile_candidate_service.py index b4abe60..65f37a8 100644 --- a/app/service/customer_profile_candidate_service.py +++ b/app/service/customer_profile_candidate_service.py @@ -18,6 +18,7 @@ from app.core.errors import GenericResourceNotFoundError, InvalidStateError from app.infrastructure.db import SessionFactory from app.model.audit import InteractionAudit from app.model.memory import MemoryConflict, MemorySyncOutbox, MemoryUnit +from app.model.platform import DomainEventOutbox from app.model.profile import ProfileSnapshot from app.service.agent.bootstrap import get_memory_cache_adapter from app.service.authorization_service import AuthorizationService @@ -147,6 +148,34 @@ class CustomerProfileCandidateService: candidate.promoted_at = now candidate.version += 1 await self._write_profile_snapshot(session, candidate, reviewer_id, now) + # 追加一条**画像重建**事件:`_write_profile_snapshot` 只写"画像快照", + # 而 `user_facts`(记忆 → 画像字段的中间层)与 `fin_customer_profile.risk_tags` + # 是由 `ProfileAssemblyService.rebuild()` 里的 `promote_facts` 写的。 + # + # 不投这条事件的后果(2026-09-14 实测):批准候选后 `memory_unit` 立刻变 active、 + # 快照版本 +1、投影也投出去了,但**画像字段仍停在旧值**(客户说"稳健型", + # 画像里还是"进取型"),要等一次无关的重建(或手动 `tools/rebuild_profile.py`) + # 才收敛 —— 演示"说完 → 批准 → 画像变了"会当场翻车。 + # + # 走事件而不是在这里直接重建:本方法所在事务**还没提交**, + # 另开 session 去重建看不到刚写下的记忆(这是 `runtime.dispatch_profile_rebuild` + # 注释里记录的同一个坑)。事件只可能在提交之后被消费。 + session.add(DomainEventOutbox( + # `id=0`:该 ORM 类没声明 autoincrement,不显式给主键会让 INSERT 报错 + # (MySQL 把 0 当作"取下一个自增值")。与 `episode_worker` 的写法一致。 + id=0, + event_id=str(uuid4()), + event_type="profile.rebuild_requested", + aggregate_type="customer_profile", + aggregate_id=str(candidate.customer_id), + trace_id=str(candidate.memory_uuid), + payload={"customer_id": candidate.customer_id, "trigger": "candidate_promoted"}, + status="pending", + retry_count=0, + occurred_at=now, + created_at=now, + updated_at=now, + )) await MemoryService(session, cache=get_memory_cache_adapter()).invalidate_recall_cache( int(candidate.customer_id) ) diff --git a/app/service/profile_assembly_service.py b/app/service/profile_assembly_service.py index ce68ecc..7d0aeb4 100644 --- a/app/service/profile_assembly_service.py +++ b/app/service/profile_assembly_service.py @@ -226,16 +226,32 @@ class ProfileAssemblyService: ) -> None: """写入新版本快照并把旧版本置为非当前。 - 唯一键 `uk_profile_snapshot_current` 建立在生成列 `current_customer_id` 上, - 保证「每个客户最多一条 current」;因此必须先清旧再写新,顺序不能反。 + 唯一键 `uk_profile_snapshot_current` 建在列 `current_customer_id` 上(**不是** + `is_current`),保证「每个客户最多一条 current」;因此**换当前版本时必须 + 同时清掉旧行的那一列**,只改 `is_current` 是不够的。 + + ⚠️ 2026-09-14 修:这里此前两件事都没做对 —— 旧行只置 `is_current=False`、 + 新行**不写** `current_customer_id`。后果是**这个唯一键从来没起作用** + (实测 13 行该列全 NULL),而且埋了一颗雷:候选批准路径( + `CustomerProfileCandidateService._write_profile_snapshot`)是**会写**这一列的, + 它能正常工作的前提是"这一列当前没人占"。只要这个客户先被重建过一次 + (旧 current 行仍占着 `current_customer_id=<客户号>`),下一次批准候选就会 + 撞 `Duplicate entry '<客户号>' for key 'uk_profile_snapshot_current'` → 整次批准 500。 + 演示链路"客户说完 → 管理员批准"会在这里断掉,所以必须按唯一键的真实语义写。 """ previous = list(await self.session.scalars( select(ProfileSnapshot).where( - ProfileSnapshot.customer_id == customer_id, ProfileSnapshot.is_current.is_(True) + ProfileSnapshot.customer_id == customer_id, + or_( + ProfileSnapshot.is_current.is_(True), + ProfileSnapshot.current_customer_id == customer_id, + ), ) )) for row in previous: row.is_current = False + # 归还唯一键的占用(旧行变历史版本后该列必须为空)。 + row.current_customer_id = None row.updated_at = now await self.session.flush() @@ -252,6 +268,8 @@ class ProfileAssemblyService: generation_basis=basis, snapshot_hash=sha256(payload.encode("utf-8")).hexdigest(), is_current=True, + # 当前版本必须显式写客户号(`app/model/profile.py` 的模块 docstring 第 2 条)。 + current_customer_id=customer_id, generated_at=now, created_at=now, updated_at=now, diff --git a/docs/44-演示流程.md b/docs/44-演示流程.md index 8704644..168a5d7 100644 --- a/docs/44-演示流程.md +++ b/docs/44-演示流程.md @@ -9,8 +9,9 @@ > 已开放自助审核发布、管理员工作台为八个标签页、角色权限数与已知空数据)。数字类结论都标了 > 实测日期——**换环境或换库后要重查**。 > -> **配套文档**(想看"为什么"而不是"怎么演"就看这两份): +> **配套文档**(想看"为什么"而不是"怎么演"就看这几份): > - `docs/演示用/全功能流程-大白话版.md` —— 从登入讲到六个门户的全部流程,不用术语; +> - `docs/演示用/记忆系统演示文档-2026-09-14.md` —— **记忆系统怎么演**(含三条可复现命令与实测输出); > - `docs/演示用/记忆架构与Worker工作流程-大白话版.md` —— 记忆链路与后台 Worker 的详细工作流程; > - `docs/演示用/今日盈亏实现说明-2026-09-14.md` —— 「今日盈亏」的口径、数据坑与复算命令。 @@ -326,7 +327,7 @@ Invoke-RestMethod -Uri http://127.0.0.1:8000/api/v1/users/me/orders -Method Post | 五角色权限对比 | 用五个账号分别登录,看入口守卫各跳哪里 | 前端只是体验层,真正的拦截在服务端权限码与数据范围 | | 知识入库 | `python tools/seed_knowledge_demo.py` 后复测场景 1 | 知识是**可运营**的,不是写死的 | | 平台体检 | `python tools/e2e_smoke_test.py` | 40 项全绿,交付质量可自证 | -| 记忆链路 | 客户连问几轮 → 看 `memory_unit` 里新增长期记忆 → Worker 投影到 Neo4j/Milvus | "记住客户"是可查的,不是模型自己说记住了(见 `docs/演示用/记忆架构与Worker工作流程-大白话版.md`) | +| 记忆链路 | 客户连问几轮 → 看 `memory_unit` 里新增长期记忆 → Worker 投影到 Neo4j/Milvus | "记住客户"是可查的,不是模型自己说记住了。**完整演法(含一键命令与实测输出)见 `docs/演示用/记忆系统演示文档-2026-09-14.md`**;快速跑:`python tools/memory_demo_chain.py` | | 后台排障 | 停掉 Worker → 客服变"繁忙" → 查 `agent_run.status='queued'` → 起回 Worker 自动恢复 | 异步三段式的可观测性 | --- diff --git a/docs/演示用/记忆系统演示文档-2026-09-14.md b/docs/演示用/记忆系统演示文档-2026-09-14.md new file mode 100644 index 0000000..0a7f2c4 --- /dev/null +++ b/docs/演示用/记忆系统演示文档-2026-09-14.md @@ -0,0 +1,321 @@ +# 记忆系统演示文档(照着演) + +> **读者**:负责演示的人。本文按"**操作 → 看到什么 → 这体现什么**"三段式写,可以直接照着念。 +> **核心思路**:记忆系统最怕"讲得很玄、看不到东西"。所以这里每个场景都配一条**可复现的命令**, +> 把"库里真有什么"摆出来 —— **能自证,才叫演示**。 +> **配套**:原理看 `记忆架构与Worker工作流程-大白话版.md`;整体流程看 `全功能流程-大白话版.md`。 +> **实测日期**:2026-09-14(数字会随数据变化,**以现场命令输出为准**)。 + +--- + +## 0. 演示前 5 分钟准备 + +### 0.1 必须满足的两件事 + +| 前提 | 为什么 | 怎么确认 | +|---|---|---| +| **API 在跑** | 所有操作走接口 | 浏览器能打开 | +| **Worker 在跑** | 记忆是**异步**提炼的:受理 202 → Worker 领单 → 抽候选 → 重建画像 | `python -m app.worker` 的窗口还在;**它停了记忆链路会静默不动**,而页面只显示"客服繁忙" | + +> ⚠️ 只跑 `uvicorn` 不跑 Worker 是**最常见的翻车点**:你说完话,什么都不发生,也没有任何报错。 +> 一句话排查:`agent_run` 最新那行是不是 `status='queued'` 且 `worker_id` 为空。 + +### 0.2 三条命令(演示时要用到的全部工具) + +```powershell +# ① 一眼看清"记忆链路现在是什么样"(只读,不写数据) +python tools/probe_memory_state.py 9001 + +# ② 一键走完整条链路:客户说 → 候选 → 客户确认 → 管理员批准 → 记忆与画像收敛 +# (会真的写入数据,这正是它的意义) +python tools/memory_demo_chain.py + +# ③ 演示"谁能读、谁读不到"(只读) +python tools/memory_recall_demo.py +``` + +### 0.3 演示账号 + +| 角色 | 账号 | 在记忆这条链上干什么 | +|---|---|---| +| 客户 | `cust_t` / `123456` | 说那句话;**确认**候选(第一把钥匙) | +| 管理员 | `admin_t` / `88888888` | **批准**候选(第二把钥匙);在「画像候选」页操作 | +| 风控专员 | `risk_t` / `666666` | 演示"有权限 + 有归属 → 能读到客户记忆" | +| 运营 | `offsite_t` / `offsite123` | 演示"没有权限 → 读不到(失败关闭)" | + +--- + +## 1. 场景一:先看"记忆到底存在哪"(1.5 min,只读) + +**操作**:跑探针。 + +```powershell +python tools/probe_memory_state.py 9001 +``` + +**看到什么**(2026-09-14 实测): + +``` +memory_unit 3 ← 长期记忆本体(客户 9001 有 3 条) +memory_evidence ... ← 每条记忆的证据(能追到"客户哪句话说的") +memory_conflict ... ← 同一键内容变化时的留痕 +memory_sync_outbox 已 processed ← 投影到 Milvus / Neo4j 的投递台账 +user_facts 3 ← 事实层(记忆 → 画像字段的中间层) +profile_snapshots 30 版 ← 画像的历史版本(每次变更只新增,不覆盖) +``` + +**这体现什么**(照着念): + +> 记忆不是一个黑盒,而是**一个账本加两份副本**: +> **账本在 MySQL**(记忆本体、证据、冲突留痕、画像版本), +> **关系副本在 Neo4j**、**向量副本在 Milvus**(都是为了检索快)。 +> 副本丢了能从账本重建;**账本与副本不一致时,一律以账本为准。** + +**被问到"证据在哪"**:`memory_evidence` 里每条记录都带 `source_table` / `source_record_id` / +`evidence_excerpt` —— 可以当场指出"这条记忆来自 `conversation_message` 第 4668 行,原话是……"。 +这是"**可追溯**",不是"模型说它记住了"。 + +--- + +## 2. 场景二 ★:让平台真的记住一件事(3 min,核心场景) + +这是整套演示里最有说服力的一段:**客户说一句 → 平台抽成候选 → 客户自己确认 → 管理员批准 → +记忆与画像都变**。全程走真实接口与真实权限。 + +### 2.1 快速版(30 秒,一条命令) + +```powershell +python tools/memory_demo_chain.py +``` + +**实测输出**(节选,2026-09-14): + +``` +--- 演示前(客户 9001) --- + active 长期记忆: [('preference:horizon', '投资期限十年以上,长期持有'), ...] + 画像字段 : [('C5', '投资期限十年以上,长期持有', '"自述:preference:risk_level=保守型"')] +【1】客户在客服会话里说:我的投资期限是三年以内。 + 受理 202,run_id=…(等 Worker 处理…) 运行终态:succeeded +【2】等它被抽成画像候选 候选 #56x:preference:horizon = 三年以内(置信度 0.9) +【3】第一把钥匙:客户本人确认 状态:verified(candidate → verified) +【4】第二把钥匙:管理员批准 状态:active(verified → active) +【5】等 Worker 自动收敛 第 1 次轮询(2s):user_facts 已更新 +--- 演示后(客户 9001) --- + active 长期记忆: [('preference:horizon', '三年以内'), ...] + 画像字段 : [('C5', '三年以内', ...)] ← 画像字段真的变了 + 画像快照 : 共 30 版,最高 v30 ← 每变一次多一版,历史不丢 +``` + +### 2.2 界面版(推荐现场用,能点的地方就点) + +| 步 | 在哪点 | 说什么 | +|---|---|---| +| 1 | 客户门户 → 右下角客服浮窗 → 输入"**我的投资期限是三年以内**"(任何命中记忆信号的话都行) | "客户随口说了一句自己的偏好" | +| 2 | 管理员工作台 →「**画像候选**」标签页 | "它没有直接进档案,而是先变成一条**候选**" | +| 3 | (客户确认这一步目前**只有接口**,见下) | "客户本人得点头" | +| 4 | 回到「画像候选」→ 点「**批准**」 | "管理员再核一次,才进正式记忆" | +| 5 | 跑 `python tools/probe_memory_state.py 9001` | "现在看库里:记忆 active、画像多了一版" | + +**第 3 步的接口兜底**(客户确认目前没有页面,如实说明即可): + +```powershell +$c = Invoke-RestMethod -Uri http://127.0.0.1:8000/api/v1/auth/tokens -Method Post ` + -Body '{"username":"cust_t","password":"123456"}' -ContentType application/json +$h = @{ Authorization = "Bearer $($c.data.access_token)" } +# 看自己名下待确认的候选 +Invoke-RestMethod -Uri http://127.0.0.1:8000/api/v1/users/me/memory-candidates -Headers $h | + ConvertTo-Json -Depth 4 +# 确认(把 <候选ID> 换成上一步拿到的 candidate_id) +Invoke-RestMethod -Uri "http://127.0.0.1:8000/api/v1/users/me/memory-candidates/<候选ID>/decisions" ` + -Method Post -Headers ($h + @{ "Idempotency-Key" = [guid]::NewGuid().ToString("N") }) ` + -Body '{"decision":"confirmed"}' -ContentType application/json +``` + +### 2.3 这体现什么(这段一定要讲) + +> 平台**不会**因为客户说了一句话就改档案。它分三步: +> ① 从对话里**抽**一条候选(带置信度、带证据); +> ② **客户本人确认**(这是他的数据); +> ③ **管理员批准**(这是合规要求)。 +> 两把钥匙都插上,才写进正式记忆,并顺着 **记忆 → 事实 → 画像字段** 一层层更新。 +> +> 为什么这么做:客服场景里客户是**随口说**的。让它直接进风控与投顾要读的画像, +> 等于用一个没核过的数字去支撑决策。 + +**被追问时的三个加分点**: + +1. **同一键只留一条生效**:客户后来说"三年以内",旧值会被置为 `invalidated`, + 并在 `memory_conflict` 留一条"谁替换了谁",不是悄悄覆盖。 +2. **风险等级只认问卷**:`investor_type`(C1–C5)**只**来自风险测评, + 记忆里哪怕有"风险偏好"也只进 `risk_tags` 并标注"**自述**"。 + 所以会出现"问卷 C5 / 自述保守 / 行为买 R4"这种**三方不一致**——那本身就是风控信号。 +3. **证据门槛**:记忆要升成"事实"需要**证据 ≥2 条**或**置信度 ≥0.90**,一句话不够格就先躺着。 + +--- + +## 3. 场景三:谁能读到、谁读不到(2 min,只读,最容易被问) + +**操作**: + +```powershell +python tools/memory_recall_demo.py +``` + +**实测输出**(节选): + +| 身份 | 库里的角色/权限/归属 | 可读客户范围 | 召回 | +|---|---|---|---| +| 客户本人 9001 | `customer`,`self` | `[9001]` | **3 条**(稳健型 / 低亏损容忍 / 三年以内) | +| 另一个客户 12001 | `customer`,`self` | `[12001]` | 0 条(读的是**他自己**名下) | +| 风控 9002 | `risk_operator`,有 `memory:read:customer`,归属 `[9001]` | `[9001]` | **3 条** | +| 投顾 9020 | `advisor`,有权限码,归属 `[9001, 9101-9104]` | 同上 | **3 条** | +| 运营 9006 | `operator`,**没有** `memory:read:customer` | **(空)** | 0 条 → 日志点名"缺能力码,失败关闭" | +| 管理员 9003 | `admin`,**有**权限码,但**没有归属行** | **(空)** | 0 条 → 日志点名"没有生效的归属客户" | + +**这体现什么**: + +> 读谁的记忆由**一个文件**说了算(`app/core/memory_scope.py`): +> **客户只读自己;员工必须同时满足"有能力码"和"客户分配给了我"**,缺一条就是**空**,并且日志点名原因。 +> +> 为什么这么严:**员工号和客户号是同一号段**(演示数据里客户 9001-9020、员工 9002/9020 并存)。 +> 早年三处各自判断口径,结果一边是"员工永远读不到"(把员工号当客户号查), +> 另一边更危险——**按号查会读到陌生客户的记忆并塞进提示词,还不报错**。 +> 金融场景最不能接受的就是这种"看起来正常"的越权。 + +> 💡 现场最漂亮的一句:**"管理员权限最大,但他也读不到这位客户的记忆 —— +> 因为权限码解决的是'能不能读他人客户数据',归属关系解决的是'谁负责谁',两者是**与**关系。"** + +**已知边界(被问到就照实说)**: + +- **客服 Agent 不召回客户记忆**(`recalls_customer_memory=False`)。客服那条线走的是场景二的 + "候选"链路。所以"客服怎么不记得我上次说的话"——那是**特意关掉的**。 +- 目前**真正消费召回结果的只有风控 Agent**(把记忆拼进系统提示词)。 +- **风控的定时扫描上下文(`system` 身份)读不到记忆**:它没有归属客户, + 根因是"召回发生在拿到目标客户之前",属已登记的遗留(两条出路待定)。 + +--- + +## 4. 场景四(备选):画像的历史版本与投影(1 min,只读) + +```powershell +python tools/probe_memory_state.py 9001 +``` + +**看到什么**:`profile_snapshots` 里同一个客户有**很多版**(实测 30 版),只有一版 +`is_current=1`,其余是历史。 + +**这体现什么**: + +> 画像**只新增版本、不原地覆盖** —— 因为风控与投顾要能回答"当时是凭什么给的结论"。 +> 每次变更多一版,`generation_basis` 还逐字段记了来源(这个字段是问卷来的、那个是记忆来的)。 + +另外讲一句可观测性:记忆写入后平台会**异步**把它投影到两个派生存储 +(关系进 Neo4j、向量进 Milvus),投递台账在 `memory_sync_outbox`。 +**投影失败不会假装成功**:`status` 会停在 `failed`/`dead` 并带 `last_error`, +Milvus 没配时事件**保持 pending**(而不是标成"已同步")。 + +--- + +## 5. 场景五(备选):把 Worker 停掉会怎样(1.5 min,风险低但很能说明问题) + +| 步 | 操作 | 看到什么 | +|---|---|---| +| 1 | 关掉 Worker 窗口 | 页面一切正常(**这就是坑**) | +| 2 | 客户再问一句客服 | 一直"客服繁忙/超时";**服务端没有任何报错** | +| 3 | 查 `agent_run` 最新一行 | `status='queued'`、`worker_id` 为空 | +| 4 | 重新起 Worker(`python -m app.worker`) | 几秒内自己追上,队列清空 | + +**这体现什么**:Agent 与记忆都是"**受理 → 排队 → Worker 执行**"的异步三段式。 +**停了不会有报错,只会静默不动** —— 所以"看板正常但什么都没发生"时,第一个要查的就是 Worker。 + +--- + +## 6. 建议的演示顺序与时长 + +| 序 | 场景 | 时长 | 目的 | +|---|---|---|---| +| 1 | 场景一:探针看"存在哪" | 1.5 min | 先把"不玄"立住 | +| 2 | **场景二:记住一件事**(界面版 + 命令兜底) | 3 min | **主线**,最能打 | +| 3 | 场景三:谁能读谁读不到 | 2 min | 讲清权限与安全 | +| 4 | 场景四:画像版本与投影 | 1 min | 讲可追溯与可观测 | +| 5 | 场景五(时间够再演):停 Worker | 1.5 min | 讲异步与排障 | + +合计约 **9 分钟**(含答疑 15 分钟)。 + +--- + +## 7. 可能被问到的问题(建议话术) + +**Q:记忆存在哪里?会不会丢?** +A:**值在 MySQL**(记忆本体、证据、画像版本,唯一权威);**关系在 Neo4j、向量在 Milvus**, +这两份是**可以重建的副本**。副本不一致时以 MySQL 为准。 + +**Q:模型会不会自己乱记?** +A:不会。① 只抽**受控键**(风险偏好、投资期限、流动性约束……词表在 +`app/service/memory_taxonomy.py`);② 每条记忆**带证据**(能追回原话);③ 要升成"事实" +需**证据 ≥2 条或置信度 ≥0.90**;④ 客服这类"随口说"的必须**客户确认 + 管理员批准**才进正式档案。 + +**Q:为什么客服聊天不写长期记忆?** +A:客服链路走的是**候选画像**(两把钥匙),不直接写。要让客服也参与召回, +得先解决"召回发生在目标客户确定之前"这个已登记的遗留。 + +**Q:员工能看到客户的记忆吗?** +A:能看到的必须**同时**满足两条:有 `memory:read:customer` 权限码 **且** 这位客户分配给他。 +缺一条读到的就是空,而且日志会点名原因("缺能力码"还是"归属未维护")。 +**管理员权限最大也不例外** —— 归属表里没有就不是他的客户。 + +**Q:记忆会过期吗?** +A:会。记忆有 `valid_until`,召回时只取未过期的,并按 `置信度 × exp(-距今天数/365)` 衰减排序。 + +**Q:客户要求删掉记忆怎么办?** +A:业务侧只写一条 `memory.deletion_requested` 事件,由 Worker 级联:记忆失效/删除 → 删证据 → +**按记忆 id 通知两个派生存储清理** → 写审计 → 清召回热缓存。 +**清不掉就如实留痕**(审计里记 `skipped` + 原因),不假装删了。 +⚠️ 当前**单条记忆失效**(非客户级联)不清理投影,属已知不对称。 + +**Q:改了记忆,画像会立刻变吗?** +A:会,但路径是**异步**的:记忆 → `user_facts`(事实层)→ 画像字段 → 新画像版本, +由 Worker 消费 `profile.rebuild_requested` 事件完成,实测**约 2 秒**。 +(这条是 2026-09-14 补的:此前**候选批准这条链**不发重建事件,画像字段会停在旧值。) + +**Q:演示会不会改坏数据?** +A:会**真的写入**(记忆、事实、画像版本各加一条)。想隔离就用另一个客户号重跑; +演示完不用清理——历史版本本来就是设计的一部分(画像只新增、不覆盖)。 + +--- + +## 8. 演示前自检(30 秒) + +```powershell +python tools/probe_memory_state.py 9001 # 表行数、事件是否堆在 pending +python tools/memory_recall_demo.py # 六个身份的召回结果是否如表格所示 +``` + +两条都正常再看一遍这张清单: + +- [ ] API **和 Worker** 都在跑(两个窗口都在) +- [ ] `domain_event_outbox` 里没有大堆 `pending`(有 = Worker 没在消费) +- [ ] 客户 `cust_t` 能登录,客户门户能打开浮窗 +- [ ] 管理员能登录,能打开「画像候选」标签页 +- [ ] 演示用的那句话**命中记忆信号**(含"投资期限 / 风险偏好 / 流动性 / 养老"等词,见 §9) + +--- + +## 9. 附:哪些话会被记住(受控信号词表摘录) + +来自 `app/service/memory_taxonomy.py`(**不是随便一句都会被记**): + +| 记忆键 | 命中示例词 | +|---|---| +| `preference:risk_level` | 风险偏好、风险承受、稳健型、保守型、激进型 | +| `preference:horizon` | 投资期限、投资周期、长期持有、短期投资 | +| `preference:product_type` | 偏好股票/债券/货币/指数、只买、只投 | +| `preference:communication` | 沟通方式、联系我、推送、短信通知 | +| `constraint:liquidity` | 流动性、随时赎回、急用钱、不能锁定 | +| `constraint:loss_tolerance` | 不能亏、怕亏、最大回撤、不能接受亏损 | +| `constraint:exclusion` | 不买、不投、别推荐、禁止 | +| `profile:occupation` / `profile:family` | 我的职业、已婚、有孩子、赡养 | +| `goal:target` / `goal:retirement` | 我的目标是、攒够、买房、养老、退休 | + +> 演示时挑一句**客户口吻**的话(如"我的投资期限是三年以内,不着急用钱"), +> 比念术语自然,也更容易命中。 diff --git a/tests/integration/test_profile_snapshot_current_invariant_mysql.py b/tests/integration/test_profile_snapshot_current_invariant_mysql.py new file mode 100644 index 0000000..457b329 --- /dev/null +++ b/tests/integration/test_profile_snapshot_current_invariant_mysql.py @@ -0,0 +1,172 @@ +"""画像快照"当前版本"不变式的真机回归。 + +## 守的是什么 + +`profile_snapshots` 的唯一键 `uk_profile_snapshot_current` 建在列 `current_customer_id` +上(**不是** `is_current`),语义是「每个客户最多一条当前版本」。有两个写入方: + +1. `ProfileGenerationService` / `CustomerProfileCandidateService` —— 清旧 + 显式写新(正确); +2. `ProfileAssemblyService._write_snapshot` —— **2026-09-14 之前两件事都没做**: + 旧行只置 `is_current=False`(不清 `current_customer_id`)、新行不写该列。 + 后果有两条,第二条会让演示直接断掉: + + - 这个唯一键从来没起作用(实测 13 行该列全 NULL); + - 只要客户**先被重建过一次**(旧 current 行仍占着 `current_customer_id=<客户号>`), + 下一次「批准画像候选」就会撞 + `Duplicate entry '<客户号>' for key 'uk_profile_snapshot_current'` → 整次批准 **500**。 + 而"客户说完 → 管理员批准"正是记忆系统演示的主线。 + +所以这个用例按**真实顺序**跑一遍:先重建画像(会写下 current 行),再批准一条候选, +断言两件事都没炸、且不变式成立。 +""" + +import asyncio +from datetime import UTC, datetime +from uuid import uuid4 + +import pytest +from sqlalchemy import select, text, update + +from app.infrastructure.db import SessionFactory +from app.model.memory import MemoryUnit +from app.model.profile import ProfileSnapshot +from app.service.profile_assembly_service import ProfileAssemblyService + +CUSTOMER_ID = 9001 +#: 用例自己造的受控键(故意不用 preference:horizon,避免碰到演示数据里那条)。 +MEMORY_KEY = "preference:communication" + + +async def _snapshot_rows() -> list[ProfileSnapshot]: + async with SessionFactory() as session: + return list( + await session.scalars( + select(ProfileSnapshot) + .where(ProfileSnapshot.customer_id == CUSTOMER_ID) + .order_by(ProfileSnapshot.id) + ) + ) + + +async def _rebuild() -> dict[str, object]: + async with SessionFactory() as session, session.begin(): + return await ProfileAssemblyService(session).rebuild(CUSTOMER_ID) + + +async def _create_candidate(memory_key: str, content: str) -> int: + """直接造一条 candidate(不经对话,避免用例依赖模型与 Worker)。""" + now = datetime.now(UTC).replace(tzinfo=None) + async with SessionFactory() as session, session.begin(): + memory = MemoryUnit( + # 不给 `id`:单列整型主键默认就是自增(给 0 反而会被当成显式值)。 + memory_uuid=str(uuid4()), + customer_id=CUSTOMER_ID, + memory_key=memory_key, + content=content, + structured_value={"value": content}, + memory_type="preference", + source_type="AI对话提取", + source_confidence=0.95, + confidence=0.95, + evidence_count=1, + conflict_count=0, + recall_count=0, + status="candidate", + valid_from=now, + valid_until=None, + promoted_at=None, + version=1, + created_at=now, + updated_at=now, + ) + session.add(memory) + await session.flush() + return int(memory.id) + + +async def _approve(candidate_id: int) -> None: + """走服务层真实入口(与管理员端点同一条路径)。""" + from app.core.contracts import RequestContext + from app.service.customer_profile_candidate_service import ( + CustomerProfileCandidateService, + ) + + context = RequestContext( + user_id="9003", + trace_id="snapshot-invariant-test", + roles=("admin",), + permissions=("memory:candidate:review",), + data_scope="all", + ) + await CustomerProfileCandidateService().review_by_admin( + candidate_id, "approved", context, comment="回归用例" + ) + + +async def _cleanup(candidate_id: int) -> None: + """删掉本轮造的候选行;并把同键残留置为 invalidated(不污染演示数据)。""" + async with SessionFactory() as session, session.begin(): + await session.execute( + text("DELETE FROM memory_unit WHERE id = :cid"), {"cid": candidate_id} + ) + await session.execute( + update(MemoryUnit) + .where( + MemoryUnit.customer_id == CUSTOMER_ID, + MemoryUnit.memory_key == MEMORY_KEY, + MemoryUnit.status.in_(("candidate", "verified", "active")), + ) + .values(status="invalidated") + ) + + +async def _pending_rebuild_events() -> int: + async with SessionFactory() as session: + return int( + ( + await session.execute( + text( + "SELECT COUNT(*) FROM domain_event_outbox " + "WHERE event_type = 'profile.rebuild_requested' " + "AND aggregate_id = :cid" + ), + {"cid": str(CUSTOMER_ID)}, + ) + ).all()[0][0] + ) + + +@pytest.mark.integration +def test_profile_rebuild_then_candidate_approval_keeps_current_snapshot_invariant() -> None: + """重建过画像之后,再批准候选**不能**因为唯一键冲突而 500。""" + candidate_id = asyncio.run(_create_candidate(MEMORY_KEY, "短信通知优先")) + try: + # ① 重建一次(这就是"埋雷"的那一步:会写下一行 current 快照) + asyncio.run(_rebuild()) + + rows = asyncio.run(_snapshot_rows()) + current = [row for row in rows if row.is_current] + assert len(current) == 1, "每个客户只能有一条 is_current" + assert current[0].current_customer_id == CUSTOMER_ID, ( + "当前版本必须显式写 current_customer_id(唯一键建在这一列上)" + ) + holders = [row for row in rows if row.current_customer_id == CUSTOMER_ID] + assert len(holders) == 1, ( + f"唯一键的占用方只能有一个,实得 {len(holders)} 个 —— " + "历史版本必须归还这一列(否则下一次写当前版本会撞 Duplicate entry)" + ) + + # ② 再批准一条候选:修复前这里会以 IntegrityError(1062) → 500 结束 + asyncio.run(_approve(candidate_id)) + + rows = asyncio.run(_snapshot_rows()) + current = [row for row in rows if row.is_current] + assert len(current) == 1 + assert current[0].current_customer_id == CUSTOMER_ID + assert len([row for row in rows if row.current_customer_id == CUSTOMER_ID]) == 1 + # 批准后还会追加一条"画像重建"事件,让 user_facts / 画像字段自动收敛。 + assert asyncio.run(_pending_rebuild_events()) >= 1, ( + "候选批准必须投一条 profile.rebuild_requested,否则画像字段会停在旧值" + ) + finally: + asyncio.run(_cleanup(candidate_id)) diff --git a/tools/memory_demo_chain.py b/tools/memory_demo_chain.py new file mode 100644 index 0000000..74b20dc --- /dev/null +++ b/tools/memory_demo_chain.py @@ -0,0 +1,258 @@ +"""记忆系统一键演示(写数据!):客户说一句话 → 候选 → 两把钥匙 → 记忆与画像收敛。 + +**它做什么**(全程走真实 HTTP,不绕过权限与幂等): + +1. 打印"演示前"的记忆与画像快照; +2. 用客户账号在客服会话里说一句带**记忆信号**的话(默认:投资期限); +3. 等 Worker 把它抽成**画像候选**(`memory_unit.status='candidate'`); +4. 客户确认(第一把钥匙,`memory:candidate:confirm`); +5. 管理员批准(第二把钥匙,`memory:candidate:review` + admin); +6. 等 Worker 自动收敛(`profile.rebuild_requested`),打印"演示后"快照。 + +**前置**:API 与 **Worker 都要在跑**(`python -m app.worker`),否则第 3/6 步永远停在 +`queued`。演示会**真的写入数据**(这正是它的意义),要用 `--customer` 换一个客户来隔离。 + +用法:: + + python tools/memory_demo_chain.py # 客户 9001,默认那句话 + python tools/memory_demo_chain.py --customer 10001 + python tools/memory_demo_chain.py --message "我的投资期限是十年以上" +""" + +from __future__ import annotations + +import argparse +import asyncio +import sys +import time +import uuid +from typing import Any + +import httpx +from sqlalchemy import text + +from app.infrastructure.db import SessionFactory + +if hasattr(sys.stdout, "reconfigure"): + sys.stdout.reconfigure(errors="replace") + +BASE = "http://127.0.0.1:8000" +DEFAULT_MESSAGE = "我的投资期限是五年以上,长期持有,不着急用钱。" +#: 演示账号(与演示数据一致;客户号 9001 对应 cust_t)。 +CUSTOMER_ACCOUNT = ("cust_t", "123456") +ADMIN_ACCOUNT = ("admin_t", "88888888") + + +def _hdr(token: str, *, idempotent: bool = False) -> dict[str, str]: + headers = {"Authorization": f"Bearer {token}"} + if idempotent: + headers["Idempotency-Key"] = uuid.uuid4().hex + return headers + + +def _login(client: httpx.Client, account: tuple[str, str]) -> tuple[str, int]: + """返回 `(令牌, 用户号)` —— 用户号从登录响应里取(`data.user_id`),别写死。""" + response = client.post( + "/api/v1/auth/tokens", json={"username": account[0], "password": account[1]} + ) + if response.status_code != 200: + raise SystemExit(f"登录 {account[0]} 失败:HTTP {response.status_code} {response.text[:200]}") + data = response.json()["data"] + return str(data["access_token"]), int(data["user_id"]) + + +async def _snapshot(customer_id: int) -> dict[str, Any]: + async with SessionFactory() as session: + memories = [ + tuple(row) + for row in ( + await session.execute( + text( + "SELECT memory_key, content, status, version FROM memory_unit " + "WHERE customer_id = :cid AND status = 'active' ORDER BY memory_key" + ), + {"cid": customer_id}, + ) + ).all() + ] + profile = [ + tuple(row) + for row in ( + await session.execute( + text( + "SELECT investor_type, investment_horizon, LEFT(risk_tags, 90) " + "FROM fin_customer_profile WHERE customer_id = :cid" + ), + {"cid": customer_id}, + ) + ).all() + ] + facts = [ + tuple(row) + for row in ( + await session.execute( + text( + "SELECT fact_key, fact_value FROM user_facts " + "WHERE customer_id = :cid ORDER BY fact_key" + ), + {"cid": customer_id}, + ) + ).all() + ] + count, top = ( + await session.execute( + text( + "SELECT COUNT(*), COALESCE(MAX(version), 0) FROM profile_snapshots " + "WHERE customer_id = :cid" + ), + {"cid": customer_id}, + ) + ).all()[0] + pending = ( + await session.execute( + text( + "SELECT COUNT(*) FROM domain_event_outbox " + "WHERE event_type = 'profile.rebuild_requested' AND status = 'pending'" + ) + ) + ).all()[0][0] + return { + "memories": memories, + "profile": profile, + "facts": facts, + "snapshots": f"共 {count} 版,最高 v{top}", + "pending_rebuild": pending, + } + + +def _dump(title: str, data: dict[str, Any]) -> None: + print(f"\n--- {title} ---") + memories = data["memories"] + print(" active 长期记忆:", memories or "(无)") + print(" 画像字段 :", data["profile"] or "(无画像行)") + print(" user_facts :", data["facts"] or "(无)") + print(" 画像快照 :", data["snapshots"]) + print(" 待消费的重建事件:", data["pending_rebuild"]) + + +def _wait_run(client: httpx.Client, token: str, run_id: str) -> str: + for _ in range(40): + time.sleep(1.5) + payload = client.get(f"/api/v1/agent-runs/{run_id}", headers=_hdr(token)).json()["data"] + if payload.get("status") in ("succeeded", "failed", "cancelled"): + return str(payload["status"]) + return "timeout" + + +def run(message: str) -> int: + with httpx.Client(base_url=BASE, timeout=60) as client: + # 客户号从**登录响应**里取,不写死:否则"看的是 A、写的是 B",演示会前后矛盾。 + customer, customer_id = _login(client, CUSTOMER_ACCOUNT) + admin, _admin_id = _login(client, ADMIN_ACCOUNT) + existing_ids = { + row["candidate_id"] + for row in ( + client.get( + "/api/v1/users/me/memory-candidates", headers=_hdr(customer) + ).json().get("data") + or [] + ) + } + + before = asyncio.run(_snapshot(customer_id)) + _dump(f"演示前(客户 {customer_id})", before) + + print(f"\n【1】客户在客服会话里说:{message}") + conversation = client.post( + "/api/v1/conversations", headers=_hdr(customer, idempotent=True), + json={"agent_type": "customer_service"}, + ).json()["data"] + run_response = client.post( + "/api/v1/agent-runs", headers=_hdr(customer, idempotent=True), + json={ + "agent_type": "customer_service", + "session_id": conversation["session_id"], + "message": message, + "idempotency_key": uuid.uuid4().hex, + }, + ) + run_id = run_response.json()["data"]["run_id"] + print(f" 受理 202,run_id={run_id}(等 Worker 处理…)") + status = _wait_run(client, customer, run_id) + print(f" 运行终态:{status}") + if status != "succeeded": + print(" ⚠️ 运行没有成功,后续演示会断在这里(先确认 Worker 在跑)") + return 1 + + print("\n【2】等它被抽成画像候选(memory_unit.status='candidate')") + candidate = None + for _ in range(20): + rows = client.get( + "/api/v1/users/me/memory-candidates", headers=_hdr(customer) + ).json().get("data") or [] + # 只认**这次新产生**的候选:库里可能躺着以前演示留下的 candidate, + # 按 id 差集挑,避免"批准了一条旧候选"这种假演示。 + fresh = [ + row for row in rows + if row.get("status") == "candidate" and row["candidate_id"] not in existing_ids + ] + if fresh: + candidate = fresh[0] + break + time.sleep(1.5) + if candidate is None: + print(" ⚠️ 没有产生候选:这句话没命中记忆信号词表(换一句试试)") + return 1 + print(f" 候选 #{candidate['candidate_id']}:{candidate['memory_key']} = " + f"{candidate['value']}(置信度 {candidate['confidence']})") + + print("\n【3】第一把钥匙:客户本人确认") + decided = client.post( + f"/api/v1/users/me/memory-candidates/{candidate['candidate_id']}/decisions", + headers=_hdr(customer, idempotent=True), json={"decision": "confirmed"}, + ).json()["data"] + print(f" 状态:{decided['status']}(candidate → verified)") + + print("\n【4】第二把钥匙:管理员批准(这一步才进正式记忆)") + review = client.post( + f"/api/v1/admin/customer-profile-candidates/{candidate['candidate_id']}/reviews", + headers=_hdr(admin, idempotent=True), + json={"decision": "approved", "comment": "演示:客户已确认"}, + ) + if review.status_code != 200: + print(f" ⚠️ 批准失败 HTTP {review.status_code}:{review.text[:200]}") + return 1 + print(f" 状态:{review.json()['data']['status']}(verified → active)") + + print("\n【5】等 Worker 自动收敛(记忆 → 事实 → 画像字段 → 投影)") + converged = False + for attempt in range(15): + time.sleep(2) + current = asyncio.run(_snapshot(customer_id)) + # 比**内容**而不是条数:同一个记忆键的值变了(约三年 → 十年以上)条数不变, + # 只比条数会误判成"没收敛"(踩过)。 + if current["facts"] != before["facts"]: + print(f" 第 {attempt + 1} 次轮询({(attempt + 1) * 2}s):user_facts 已更新") + converged = True + break + if not converged: + print(" ⚠️ 没等到收敛:确认 Worker 在跑,或看 memory_sync_outbox.status") + + _dump(f"演示后(客户 {customer_id})", asyncio.run(_snapshot(customer_id))) + print( + "\n提示:想核对投影落没落库,跑 " + f"`python tools/probe_memory_state.py {customer_id}`;" + "想演「谁能读、谁读不到」跑 `python tools/memory_recall_demo.py`。" + ) + return 0 + + +def main() -> int: + parser = argparse.ArgumentParser(description="记忆系统一键演示(会写数据)") + parser.add_argument("--message", default=DEFAULT_MESSAGE, help="客户说的那句话") + args = parser.parse_args() + return run(args.message) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tools/memory_recall_demo.py b/tools/memory_recall_demo.py new file mode 100644 index 0000000..afe1ded --- /dev/null +++ b/tools/memory_recall_demo.py @@ -0,0 +1,102 @@ +"""记忆召回演示(只读):**同一个客户,换不同身份去读,看谁能读到、谁读不到。** + +## 为什么需要它 + +"记忆系统的权限"这套东西光讲没用 —— 演示时最好当场看到: +用**客户本人**的身份能读到自己,用**别的客户**身份读到的是空的, +用**有权限且有归属**的风控身份能读到,用**没权限**的运营身份读不到(而且日志会点名原因)。 + +这个脚本就是干这个的:按平台**真实**的身份解析(`IdentityService`,读的是库里的 +角色/权限/客户归属)构造上下文,再走**产线同一条召回链路** +(`PlatformGovernance` + `build_memory_recall_service`,含 Redis 热缓存与 Milvus 语义通道), +逐个身份打印召回结果与降级原因。全部只读,不写任何数据。 + +用法:: + + python tools/memory_recall_demo.py # 默认看客户 9001 + python tools/memory_recall_demo.py --customer 10001 + python tools/memory_recall_demo.py --identity 9002 # 只看某个身份 +""" + +from __future__ import annotations + +import argparse +import asyncio +import sys + +from app.core.contracts import RequestContext +from app.core.memory_scope import REQUIRED_EMPLOYEE_PERMISSION, customer_memory_scope +from app.service.agent.bootstrap import build_memory_recall_service +from app.service.agent.governance import PlatformGovernance +from app.service.identity_service import IdentityService + +if hasattr(sys.stdout, "reconfigure"): + sys.stdout.reconfigure(errors="replace") + +#: 演示身份:(显示名, user_id, 说明)。真实的角色/权限/归属都从库里查,不在这里编。 +DEMO_IDENTITIES: tuple[tuple[str, str, str], ...] = ( + ("客户本人(9001 / cust_t)", "9001", "客户身份:只读自己"), + ("另一个客户(12001)", "12001", "客户身份:读的是**他自己**名下,不是 9001"), + ("风控专员(9002 / risk_t)", "9002", "员工:有 memory:read:customer 且 9001 已分配给他"), + ("投顾(9020 / advisor_t)", "9020", "员工:同上,看归属表里有没有"), + ("运营(9006 / offsite_t)", "9006", "员工:**没有** memory:read:customer → 失败关闭"), + ("管理员(9003 / admin_t)", "9003", "员工:有权限码,但要看归属表"), +) + + +def _reason(context: RequestContext, has_scope: bool) -> str: + """把"为什么 0 条"讲成人话(判定本身仍由 `customer_memory_scope` 做,这里只解释)。""" + if not has_scope: + if REQUIRED_EMPLOYEE_PERMISSION not in context.permissions: + return (f"该身份缺少 {REQUIRED_EMPLOYEE_PERMISSION} 能力码 → **失败关闭**" + "(归属关系只解决'谁负责谁',不构成读他人客户数据的授权)") + return ("该身份在 `sys_customer_assignment` 里**没有生效的归属客户** → **失败关闭**" + "(这不是记忆坏了,是没授权;要读指定客户可先维护分配关系)") + return "可读范围内这些客户确实没有 active 的长期记忆(或都过期了)" + + +async def show(customer_id: int, only: str | None) -> None: + governance = PlatformGovernance(recall_factory=build_memory_recall_service) + identity = IdentityService() + + for label, user_id, note in DEMO_IDENTITIES: + if only and only != user_id: + continue + context = await identity.resolve( + RequestContext(user_id=user_id, trace_id=f"memory-recall-demo-{user_id}") + ) + items = await governance.recall(context) + scope = customer_memory_scope(context) + print("=" * 78) + print(f"{label} ({note})") + print(f" 角色={list(context.roles)} 数据范围={context.data_scope} " + f"归属客户={list(context.customer_ids)}") + print(f" memory:read:customer = " + f"{REQUIRED_EMPLOYEE_PERMISSION in context.permissions}") + print(f" 可读客户范围 = {list(scope) or '(空)'}") + if not items: + print(f" 召回:**0 条** —— 原因:{_reason(context, bool(scope))}") + continue + print(f" 召回:{len(items)} 条") + # `RecalledMemory` 就是**送进模型上下文**的那一份(只有 uuid/客户号/正文), + # 置信度与来源在更下层的 `RecallItem` 里,这里刻意不展开 —— 演示时讲清这一点即可。 + for item in items: + print(f" · 客户 {item.customer_id}:{item.content[:60]} " + f"(uuid={item.memory_uuid[:8]}…)") + print("=" * 78) + print(f"提示:这里看的是**召回**(把记忆塞进模型上下文)。" + f"客户 {customer_id} 的记忆本体在 memory_unit 表里,用 " + f"`python tools/probe_memory_state.py {customer_id}` 看。") + + +def main() -> int: + parser = argparse.ArgumentParser(description="记忆召回按身份的演示(只读)") + parser.add_argument("--customer", type=int, default=9001, help="被观察的客户号(默认 9001)") + parser.add_argument("--identity", default=None, help="只看某一个身份(user_id)") + args = parser.parse_args() + asyncio.run(show(args.customer, args.identity)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main())