diff --git a/docs/演示用/记忆架构与Worker工作流程-大白话版.md b/docs/演示用/记忆架构与Worker工作流程-大白话版.md new file mode 100644 index 0000000..bd9849a --- /dev/null +++ b/docs/演示用/记忆架构与Worker工作流程-大白话版.md @@ -0,0 +1,409 @@ +# 记忆架构与底座 Worker · 大白话版(2026-09-14) + +> **读者**:想知道"客户说过的话,平台到底记在哪儿、怎么记住、后台那个 Worker 一天到晚在忙什么"的人。 +> **写法**:不用术语;每一句都能在代码或数据库里查到,关键处标了文件路径。 +> **配套**:全平台流程看 `全功能流程-大白话版.md`;演示照读看 `docs/44-演示流程.md`。 +> **数字口径**:文中"当前库里……"之类的数字都是 **2026-09-14 实测**,换库/换环境后要重查。 + +--- + +## 0. 先把这两件事串起来(一句话) + +``` +客户说话 → API 受理(落库、立刻回 202)→ 请求进队列 + ↓ + 【底座 Worker】把活儿领走:跑 Agent、投事件、写记忆、投影向量 + ↓ + 记忆落在 MySQL(账本:记忆本体 + 证据 + 画像版本) + ↓ + 派生出两份副本:Neo4j 存关系、Milvus 存向量(都是为了下次检索) + ↓ + 下次对话时召回相关记忆 → 拼进给模型的上下文 +``` + +**所以第一部分讲"记忆存在哪、怎么写、怎么写错、怎么忘",第二部分讲"这些慢活是谁在干、干到哪一步、怎么判断它死了"。** +两者是同一件事的两端:记忆架构里所有的"↓ 由 Worker 做"都对应第二部分的某条工作线。 + +--- + +## 第一部分:记忆架构(客户说过的话,平台到底记在哪) + +### 1.1 四个存储,只有一个算"账本" + +| 存储 | 存什么 | 地位 | +|---|---|---| +| **MySQL** | **长期记忆本体**:`memory_unit`(每条记忆:键、正文、结构化值、置信度、状态、有效期、版本、证据条数;带一个**生成列 + 唯一键**,数据库层面保证"同一客户同一记忆键只有一条生效")、`memory_evidence`(每条记忆的证据与幂等键)、`memory_conflict`(同一个键内容变了就留痕)、`user_facts`(事实清单)、`fin_customer_profile`(画像当前值)、`profile_snapshots`(画像的历史版本) | **唯一权威(账本)** | +| **Neo4j** | 客户与标签/产品之间的**关系**:`(客户)-[:PREFERS / :HAS_GOAL]->(标签)`、`(客户)-[:HOLDS]->(产品)` | 派生视图,可重建;**不一致时一律以 MySQL 为准** | +| **Milvus** | 长期记忆的**向量**,集合名 `user_long_term_memory_v1`(是**代码里的常量**,不是配置项;写/读/删三处共用,有单测守着) | 派生视图,可重建 | +| **Redis** | 只做**加速**:召回热缓存(键 `mem:recall:...`,**300 秒**过期)+ 客服当前会话的短期上下文(滑动 30 分钟、绝对上限 24 小时) | 可随时丢,丢了不影响正确性 | + +一句话:**MySQL 是账本,另外三个都是"能从账本重新算出来的副本"**。 +一条记忆只有在 MySQL 里存在,才算真的"记住了"。 + +> Milvus 里只放 `preference:` / `goal:` 这两类前缀的记忆(偏好与目标), +> `constraint:` / `profile:` 之类不进向量库,并会在日志里留痕说明为什么跳过。 + +### 1.2 一条记忆的一生(走"非客服"这条主流链路) + +| 步 | 谁做 | 落什么 | +|---|---|---| +| ① | API 侧 Agent 跑完 | **同一个事务**里写:结果消息、审计、`agent.run_completed` 事件、`memory.extraction_requested` 事件 | +| ② | **Worker** 消费该事件 | 回查原对话 → 让模型按**严格 JSON + 受控词表**提炼(提炼失败就失败关闭,不硬猜)→ 写 `memory_unit` + `memory_evidence`(幂等键 `memory.extraction_requested:{事件id}`) | +| ③ | 同上,仅当**确实新增了证据** | 投一条 `profile.rebuild_requested` 事件 | +| ④ | **Worker** 消费该事件 | 记忆 → `user_facts`(要**证据 ≥2 条或置信度 ≥0.90** 才提升)→ 写 `fin_customer_profile`(**只映射白名单字段**)→ 新增一版 `profile_snapshots`(版本号 +1、`is_current` 切换) | +| ⑤ | 同一步里顺带做 | 图投影(全部 `MERGE`,幂等)+ 往 `memory_sync_outbox` 写一行 `milvus` | +| ⑥ | **Worker** 消费投影事件 | 写向量(按 `memory_uuid` + `version` 幂等:版本更高才覆盖) | +| ⑦ | 下次对话 | 召回(见 1.5),把相关记忆拼进给模型的上下文 | + +三个"为什么这么设计"的点: + +- **为什么画像重建要绕一层事件**:记忆写入时那条数据还在**没提交的事务**里, + 另开连接去重建画像根本看不到它(实测过:快照加了、事实没进、图里也没多关系)。 + 事件只可能在提交之后被消费,那时数据一定可见。 +- **`investor_type`(投资者类型)永远只来自风险测评**,记忆碰不到它。 + 记忆能影响的只有偏好、期限、风险标签这类字段——这是合规边界,不是实现偷懒。 +- **证据门槛**:一句话就改画像太危险,所以要么有两条以上证据,要么单条置信度 ≥0.90。 + +### 1.3 另一条路:客服对话走"候选画像",需要两把钥匙 + +**客服 Agent 不写长期记忆**(`recalls_customer_memory=False`,且明确不抽记忆)。客服链路是: + +``` +客服对话命中信号 → 写一条 memory_unit(status='candidate') ← 只是"候选",不进画像 + ↓ 客户本人在"我的画像候选"里确认(权限 memory:candidate:confirm) + status='verified'(已核实) + ↓ 管理员批准(权限 memory:candidate:review,且必须 admin) + status='active' → 同时写冲突留痕 + 新画像版本 + 投两条投影事件 +``` + +为什么多这一道:客服场景里客户是**随口说**的("我比较喜欢稳一点的"), +直接进正式画像风险太大;让**客户自己确认 + 管理员批准**,才当事实用。 + +### 1.4 什么才算"值得记住" + +不抽的情况(刻意的): + +- **客服对话**不抽长期记忆(走 1.3 的候选链路); +- **访客**不抽(没有身份,记住谁?)。 + +其余情况要满足三者之一才抽:**工具有成功产出**(拿到了权威数据)、 +**发生了业务事件**(风险测评完成、适当性评估、成交等)、**用户消息命中受控信号词表** +(默认 13 个键:风险偏好、投资期限、家庭状况、目标之类)。 +抽不出来就标记"已完成但未提升",**不硬凑一条记忆**。 + +### 1.5 召回:三路合流,但没有"图"这一路 + +| 路 | 规则 | +|---|---| +| ① MySQL 结构化(主力) | 取 `status='active'` 且未过期的记忆,按 `置信度 × exp(-距今天数/365)` **衰减排序** | +| ② Milvus 语义 | 需同时满足"Milvus 可达 + 有 embedding 端点",否则**显式降级**并记原因,①照常返回;语义置信度 = `相似度 × 0.7`(**刻意压低**:不能让"像"压过"是") | +| ③ Redis 热缓存 | 键命中就**直接返回、短路①②**,TTL 300 秒;记忆一写就主动失效(枚举常见参数组合的键来删) | + +合并口径:以 `memory_uuid` 为身份合并,结构化优先、向量命中追加, +同一个 uuid 取较高置信度,最后按置信度降序截断。 +**没有分数阈值,只有条数上限**:单客户默认 10 条(上限 100),跨客户合并后再取 10 条。 + +> ⚠️ **召回没有"图"这一路**:Neo4j 只服务画像展示、投顾推荐与风控工具, +> 召回代码里根本不引用图适配器。谁说"图数据库参与记忆召回",那是记错了。 +> +> 另一个当前的事实:**真正消费召回结果的只有风控 Agent**(它把记忆拼进系统提示词); +> 客服 Agent 明确**不召回客户记忆**(`recalls_customer_memory=False`),其余 Agent 默认为"召回"。 +> 所以演示时如果问"客服怎么不记得我上次说的话"——那是**特意关掉的**, +> 客服的记忆走的是 1.3 的候选画像链路,不是召回。 + +### 1.6 谁能读到谁的记忆:一个文件说了算 + +| 身份 | 能读谁 | +|---|---| +| **客户** | **只有自己**,别家客户一律读不到 | +| **员工**(风控/投顾/运营/管理员/system) | **必须同时满足两条**:① 有 `memory:read:customer` 能力码;② 客户已在 `sys_customer_assignment` 里分配给自己。**缺一条就是空** | +| 访客 | 无(上游已拦) | + +细节:单次最多 **10 个**归属客户(按客户号**升序**截断——为了"同一输入结果永远一致"), +合并后最多 **10 条**记忆;解析不出数字的客户号一律判**越界**(宁可拒掉一条来路不明的记忆)。 +演示数据里客户 9001 已分配给投顾 9020 与风控 9002,所以这两个账号能读到 9001 的记忆, +而运营账号(只有 NL2SQL 与场外的权限)**读不到**——它不是"被系统忘了",是**没有授权**。 + +**为什么要单独抽一个文件**(`app/core/memory_scope.py`):这里踩过两个坑—— +以前三处各自判断口径不一,结果是 **① 员工身份召回恒空**(把员工号当客户号查,永远查不到), +**② 越权陷阱**:员工号和客户号在同一号段(演示数据里客户 9001-9020、员工 9002/9020 并存), +按号查会**读到陌生客户的记忆并塞进提示词,而且不报错**。 +金融场景最不能接受的就是这种"看起来正常"的越权。 + +> **"归属没维护"不是故障,是"没有授权"**,所以日志会点名原因 +> (是缺能力码,还是没分配客户),让运维知道该去维护分配表, +> 而不是跑去查"记忆是不是坏了"。 + +### 1.7 怎么把记忆忘掉(失效与删除) + +``` +业务侧只写一条事件 memory.deletion_requested(mode = invalidate 或 delete),不直接动表 + ↓ Worker 消费 +① memory_unit 置为失效(或按 id 物理删除) +② memory_evidence 物理删除 +③ 按 memory_uuid 投投影清理事件(memory.invalidated / memory.deleted) +④ 写审计(action_type = memory.customer_lifecycle) +⑤ 删该客户的召回热缓存 + ↓ 投影侧 +⑥ 删对应 user_facts → 用剩余证据重建画像 → 图对账(以画像为准:缺的补、多的删,删完再核对一次) +⑦ Milvus 按 memory_uuid 删向量 +``` + +**清不掉就如实留痕**,绝不假装删了:审计里记 `status ∈ {cleaned, skipped, skipped_no_client}` +加原因串(比如 `graph_client_unavailable`、`vector_delete_failed`、`vector_collection_absent`)。 + +> ⚠️ 一处**已知的不对称**:客户级联失效会清理投影, +> 但**单条记忆失效**(`MemoryService.invalidate`)只把状态改掉,**不写事件、不清理投影**。 +> 也就是说那条记忆的向量还留在 Milvus 里(虽然召回时会按 `active` 过滤掉)。 +> 演示时不用讲,但如果被追问"删掉一条记忆向量删不删",如实说"当前只有客户级联会删投影"。 + +### 1.8 当前库里实际是什么样(2026-09-14 实测) + +| 表 | 行数 | 说明 | +|---|---|---| +| `memory_unit` | **2** | 都是客户 9001 的:`preference:risk_level`(激进型,v5)、`preference:horizon`(约三年,v4) | +| `memory_evidence` | 9 | 每条记忆的证据 | +| `memory_conflict` | 4 | 同一键内容变化的历史留痕 | +| `episodes` | 191 | 会话片段;其中 **185 行是"已完成但未提升"**(处理过了,只是没抽出值得长期记住的事实),6 行提升成了记忆 | +| `user_facts` | 3 | 事实清单 | +| `profile_snapshots` | 13 | 画像版本(9001 有 11 版、10001/10002 各 1 版) | +| `fin_customer_profile` | 6 | 画像当前值 | +| `memory_sync_outbox` | 4 | 全是 10001/10002 开户测评的投影,**都已 processed** | +| Milvus `user_long_term_memory_v1` | **2 条** | 正是 9001 的那两条记忆(version 5 / 4) | +| Redis `mem:recall*` | 0 个键 | TTL 300 秒早就过了,属正常 | + +> ⚠️ **一个要如实说明的观察**:9001 的向量**确实在 Milvus 里**, +> 但 `memory_sync_outbox` 现在**没有任何 9001 的行**(那 4 行都是 10001/10002)。 +> 这说明它的投递来源不在当前这张表里(行被清理过,或由更早的路径/探针写入), +> **具体时点未查证**。所以**不能**拿"Milvus 里有向量"来断言"投影链路当前工作正常"—— +> 要证明这条链路,得让 Worker 跑起来、再触发一次记忆写入,然后看 outbox 里有没有新行。 + +### 1.9 已知遗留(逐条核实"现在是否仍然存在") + +| # | 遗留 | 现状 | 影响 | +|---|---|---|---| +| ① | **片段抽取的幂等哈希不覆盖"抽出来的内容"** | **仍在**:`content_hash` 只由客户+会话+每条消息的 id/角色/正文决定,抽取结果不在键里 | 换模型或改提示词后,同一个片段重复消费**不会重写记忆、也不留痕** | +| ② | **投顾线两处生产者不发 `memory_sources`** | **仍在**:靠 Worker 侧回查兜底,**每次都会 warning** | 功能可用,但契约没根治;warning 就是提醒 | +| ③ | **风控的 system 上下文召回恒空** | **仍在**:扫描用 `user_id="0"`、`roles=("system",)`、没有归属行 | 根因是"召回发生在 `handle()` 之前、上下文里没有本次目标客户",而 `RequestContext` 至今没有"目标客户"字段。两条出路(补分配行 / 引入显式目标客户)**待你决策** | +| ④ | 投影清理客户端未装配时的降级 | 机制仍在,但**生产入口已经装配**了真实清理服务 | 正常启动的 Worker 走真实删除 | +| ⑤ | **(本次新发现)`profile_snapshots.current_customer_id` 从不被写入** | **仍在**:13 行该列**全部为 NULL**,而唯一键正建在这一列上 | "每个客户最多一条 current"目前**只由代码逻辑保证,数据库没有约束**;一旦代码出错,会出现两条 `is_current=1` 而不被数据库拦下 | +| ⑥ | 文档/注释过期 | `docs/37` 称 `memory_unit`/`user_facts`/`episodes` "目前都是 0 行"(现为 2/3/191);`tools/rebuild_profile.py` 注释称"画像组装没有自动触发"(现在有事件驱动) | 会误导接手的人,建议按实测更新 | +| ⑦ | **(本次一并核实)`user_facts` 没有 `(客户, 事实键)` 唯一键** | 靠服务层"先查后写"保证不重复 | 并发写入时理论上可能写出两条同键事实;当前单 Worker 场景触发不到 | +| ⑧ | 两个"备用/未接线"的组件 | `app/worker/graph_projection_worker.py` 的 `GraphProjectionWorker` **全仓零实例化点**;`ProjectionReconciliationService` 也**没有生产调用点** | 看代码时别以为它们在跑;实际图投影走的是 `ProfileGraphProjectionService` | + +### 1.10 给业务方的一句话 + +> 平台记住一个客户,靠的不是"模型自己记得",而是: +> **在 MySQL 里落一条可追溯的记忆(带证据、带置信度、带有效期)→ 够格的才升成画像 → +> 再从账本派生出关系和向量两份副本用于检索**。 +> 谁能读、读几条、什么时候忘掉,全部有明确规则和留痕; +> 派生副本丢了可以从账本重建,账本不一致时一律以账本为准。 + +--- + +## 第二部分:底座 Worker 到底在忙什么 + +### 2.1 为什么必须有它:三段式 + +平台的 Agent 请求**不是**"一问一答的同步调用",而是**三段式**: + +``` +① 受理(API):收到请求 → 落库 → 立刻回 202「收到了」 ← 几十毫秒 +② 排队:请求变成 agent_run 表里的一行,状态 queued +③ 执行(Worker):Worker 把这一行领走、真正跑 Agent、把结果落库 ← 几秒 +``` + +页面拿到 `202` 之后靠**轮询**问"好了吗"(客服浮窗每 0.5 秒一次、最多 40 次;NL2SQL 每 1.2 秒一次)。 +所以 **Worker 不跑 = 请求永远停在 `queued` = 页面只会说"客服繁忙/超时"**, +而服务端**没有任何报错**。这是这套架构最容易误判的一环。 + +除了 Agent,还有一大批"慢慢做就行"的活也归它(见 2.3)。 + +### 2.2 怎么起、起了什么 + +| | | +|---|---| +| 启动 | `python -m app.worker`(`start.ps1` 会单独开一个窗口起它) | +| 只跑一轮 | `python -m app.worker --once`(诊断用:跑一轮就退出,异常直接抛出来) | +| 轮询节奏 | 默认**每秒一轮**(`worker_poll_seconds`,`app/core/config.py`) | +| 装配在哪 | `app/worker/__main__.py`:显式注入 `relationships`(图库)、`projection_cleaner`(记忆删除时的投影清理)、`settings`,并构造 `OffsiteMailWorker` | + +**装配是刻意的"在入口处一眼可见"**:向量库/图库的客户端不在这里硬造, +由 `WorkerRuntime` 惰性去 `bootstrap` 里取——取不到就**显式降级**(见 2.8), +而不是"假装配好了、跑起来再说"。 + +### 2.3 一轮 `run_once()` 的四件事(顺序固定) + +| 顺序 | 做什么 | 单轮上限 | 频率 | +|---|---|---|---| +| ① | **领域事件投递**:把 `domain_event_outbox` 里的待办事件一条条消费掉 | 10 条 | 每轮 | +| ② | **记忆/画像投影**:消费 `memory_sync_outbox`,把记忆投到 Milvus 与 Neo4j | 20 条 | 每轮 | +| ③ | **会话片段**:先把若干轮对话聚合成片段,再把片段抽成长期记忆 | 50 个客户 / 20 个片段 | **每 30 轮才做一次** | +| ④ | **Agent 运行**:领**最老的那一条** `queued/running/cancel_requested` 请求并执行 | 1 条 | 每轮 | + +失败隔离做得比较细:②③ 出错只**告警**、不影响其它步骤; +①④ 出错会冒到入口,常驻模式**记堆栈 + 退避重试(不退出)**,`--once` 模式才抛出来。 + +> 为什么 ③ 要节流成 1/30 轮:片段聚合要按客户维度做 distinct 查询, +> 每轮都做就是白白烧数据库。 + +### 2.4 工作线一:Agent 运行(最要紧的一条) + +**状态机**(`agent_run.status`): + +``` +queued ──领取──> running ──成功──> succeeded + │ └─失败(可重试且未超限)──> queued(退避后重来) + ├────────失败(不可重试/超限)────> failed + └────────收到取消──────────────> cancelled +(另有 cancel_requested:已请求取消、还没落地) +``` + +**领取与"不会被两个 Worker 抢同一单"**:领取用一条带条件的 `UPDATE` +(`status in (queued,running)` 且 `locked_until` 为空或已过期)→ 置 `running` + +`worker_id` + `locked_until = now + 60 秒`。后到的 Worker 抢不到就是抢不到。 + +**心跳**:执行期间每 **20 秒**(= 租约 60 秒 ÷ 3)续租一次;续租失败说明"活被别人接管了", +立刻取消本地任务,并把这次运行记成 `RUN_LEASE_LOST`(可重试)。 + +**重试**:可重试的失败会回到 `queued`,退避 `min(60, 2^已试次数)`; +超过 `worker_retry_limit`(默认 3 次)才落 `failed`。 + +**取消**:HTTP 侧只把状态改成 `cancel_requested`;Worker 在两个点上收口—— +执行前看到就直接落 `cancelled`,执行中则由"心跳续租失败"触发取消。 +(**代码里没有"执行中主动轮询取消标志"的逻辑**,这是刻意的:靠租约而不是靠轮询。) + +**执行前会重新解析身份与权限**(`restore_context`):哪怕这条请求是几分钟前受理的, +执行时也按**当时的**权限与数据范围判定——权限被收回了,活儿就不干。 + +**结果由谁落库**:`AgentPersistenceService.complete_run` 在**同一个事务**里写 +结果消息、`run.status=succeeded`、审计、幂等记录(`request_idempotency`),并发出一条 +`agent.run_completed` 事件。中途还会再校验一次租约——对不上就抛错,避免"过期的人写结果"。 + +### 2.5 工作线二:领域事件(`domain_event_outbox`) + +**这是平台的"可靠待办清单"**:先落库、再干活,干了什么都要认账。 + +| 项 | 实际情况 | +|---|---| +| 领什么 | `status in (pending, failed)` **且** `event_type` 在白名单里 **且** 重试时间已到;按 `id` 升序、`FOR UPDATE SKIP LOCKED`(抢锁不阻塞别人) | +| 白名单是什么 | 就是 `WorkerRuntime` 里那张硬编码的处理器表;**知识向量那两个事件只有在 Milvus 写适配器就绪时才会被并进去** | +| 事件类型(当前) | `agent.run_requested`、`memory.extraction_requested`、`customer_profile.candidate_requested`、`agent.run_completed`、`config.cache_invalidate_requested`、`memory.deletion_requested`、`memory.invalidated`、`memory.deleted`、`profile.rebuild_requested`、`conversation.transfer_requested`,加上知识的 `knowledge.vector_sync_requested` / `knowledge.vector_delete_requested` | +| 幂等在哪 | 两层:投递层用 `(event_id, 消费者名)` 唯一键先查后跳过;业务层各 handler 自己再幂等(知识按主键覆盖、图投影用 `MERGE`) | +| 失败怎么办 | 重试次数 +1、退避 `min(300, 2^次数)`;累计到 **5 次**转 **`dead`(死信)** | +| 状态取值 | `pending` / `failed` / `published` / `dead` | +| 错误写什么 | 只写**代码里写死的文案或异常类名**,不写原始异常字符串——防止数据库连接串、密钥之类混进日志 | + +### 2.6 工作线三:记忆与画像投影(`memory_sync_outbox`) + +领取条件、退避公式、死信阈值(5 次)与上一节一致,两个"目标存储"各有一个处理器: + +| 目标 | 干什么 | +|---|---| +| `milvus` | 把该客户的长期记忆写进向量集合(`MilvusProfileProjection.upsert`,集合名见第一部分) | +| `neo4j` | 复用主干的图投影服务重建该客户的关系;**图库不可用时抛"可恢复错误"走重试**,绝不记成"已投递" | + +两个"知道就好"的细节: + +- 事件里少了 `memory_sources`(要投影哪些记忆)时,**会回查该客户当前有效的记忆兜底**, + 并且**每次留一条 warning**——目的是让"谁没按契约发字段"始终看得见,而不是悄悄兼容; +- **没配 Milvus 客户端时,这一批事件根本不消费、保持 `pending`**,只在日志里告警一次。 + 这跟"注册一个空处理器然后判死信"是**刻意区分**的:前者是"能力没装配",后者是"投递失败"。 + +### 2.7 工作线四:会话片段与知识向量 + +**会话片段(`episodes`)**:平台不会把每一句话都当记忆,而是先把对话**切成片段**再提炼。 + +- 切法:**相邻两条消息间隔超过 30 分钟就算新片段**(静默窗口 `DEFAULT_GAP_MINUTES=30`); +- 幂等:每个片段算一个 `content_hash`(客户 + 会话 + 每条消息的 id/角色/正文), + 靠表的唯一键去重,重复聚合直接跳过; +- 抽取:一轮最多消费 20 个片段,只挑"待提取/失败且重试未超限"的; + 抽成功 → 标记"已完成 + 已提升为长期记忆",并发一条 `profile.rebuild_requested`; + 抽失败 → 记"失败"、重试次数 +1,**一行记忆都不写**(不会半截写入)。 + +**知识向量**:`knowledge.vector_sync_requested` 与 `knowledge.vector_delete_requested` 两个事件, +写进哪个集合是**按知识那一行自己的字段**取的(白名单三个集合:faq / product / policy), +只有 `status='active'` 的知识才写,向量维度必须等于 **1024** 否则失败关闭。 + +### 2.8 工作线五:场外收件(和 Agent 同一个进程) + +`OffsiteMailWorker` 就在**同一个入口、同一个循环**里跑(历史上它曾因一次文件覆盖被删掉接线, +导致邮箱完全没人收——现在源码里专门留了注释记这件事)。 + +它每一轮要过**四道闸门**,任何一道不满足就什么都不做(只写日志): + +1. 场外总开关开着; +2. IMAP 收件开着; +3. **真实识别就绪**:OCR 与模型两项健康检查都 `ok`(没就绪时**明确拒绝**用 Mock 结果糊过去); +4. 配了"场外操作身份"(`offsite_worker_user_id`),并且这个身份能解析出角色与权限。 + +然后:领**游标租约**(`SELECT ... FOR UPDATE` + `lease_until`,心跳每 20 秒续租; +续租时影响行数不为 1 就取消本轮)→ 拉一批邮件 → 落库。 +单封邮件失败重试到上限(默认 3 次)就置 `blocked` 并写审计 `offsite.mail_cursor_blocked`。 +另外每轮开头还会做一次"救火":把**卡在"发送中"超过 900 秒**的通知改判为"发送失败"并留痕。 + +### 2.9 它什么时候"故意不做事"(显式降级清单) + +这些情况都是**刻意**的:宁可说"没做成",也不写一条假的"已完成"。 + +| 情况 | 行为 | +|---|---| +| 没配 Milvus(知识向量) | 不注册处理器,事件保持 `pending`,日志告警一次 | +| 没配 Milvus(记忆向量) | 不消费,事件保持 `pending` | +| 没注入投影清理客户端 | 记 warning + 审计(`action_type=memory.projection_cleanup`、`status=skipped_no_client`),**事件仍标记已消费**(权威库已正确,投影是可重建的派生数据) | +| 图投影降级 | 抛可恢复错误 → 走重试 | +| Redis 删缓存失败 / 会话短期记忆写失败 | 只 warning,不影响主流程 | +| 场外识别未就绪 / 操作身份未配 | 记 error 后**不处理**,不会用假数据蒙过去 | + +### 2.10 怎么判断它活着、卡住了、还是积压了 + +**看表**(比看日志快): + +| 症状 | 看哪里 | +|---|---| +| 客服"超时" | `agent_run` 里最新那条是不是 `status='queued'` 且 `worker_id` 为空 | +| 有一单卡住不动 | 该行 `locked_until` 是否早已过期、`attempt_count` 是否不再增长 | +| 事件堆积 | `domain_event_outbox` 的 `status/retry_count/next_retry_at` 分布 | +| 记忆没进向量库 | `memory_sync_outbox` 的 `status/last_error` | +| 记忆没被提炼 | `episodes.extraction_status/retry_count` | +| 场外不收信 | `offsite_mail_cursor.status='blocked'` 与 `next_retry_at` | + +**看日志关键词**:`worker round failed; retrying after backoff`、`run failed run_id=`、 +`lease renewal failed`、`profile projection consumption failed`、`episode aggregation failed`、 +`knowledge vector handlers not registered`、`profile projection disabled`、`projection cleanup degraded`、 +`场外邮件处理失败`。 + +**当前死信实况(2026-09-14 实测,供对照)**:`domain_event_outbox` 里 +`published` 1415 条、**`dead` 570 条**,全都是**重试满 5 次**留下的(另有 1 条是早期"没有处理器"); +原因分布是 `ValueError` 372 条(集中在 09-09/09-10,早期链路)、 +`OutboxHandlerError: run not found` 190 条(最近到 09-14)、 +知识向量 `RecoverableAgentError` 7 条(09-13)。 +**这不是"现在的故障",而是可靠队列的正常留痕**——它证明"失败了不会悄悄丢"。 + +### 2.11 三个常被搞错的边界 + +1. **风控定时扫描不在这个 Worker 里**。`app/worker/risk_scan_scheduler.py` 是**独立进程** + (`python -m app.worker.risk_scan_scheduler`),而且**默认关闭**; + 页面上那个「扫描」按钮走的是 API(`POST /api/v1/risk/alerts/scan`),同步执行、最长 60 秒。 + 两者**共用同一把数据库咨询锁**互斥:手工扫描抢不到锁会报"忙",定时扫描抢不到就跳过本轮。 +2. **`app/worker/agent_run_worker.py` 里的 `AgentRunWorker` 生产入口没有用**(只有测试引用)。 + 生产上领取与执行走的是 `WorkerRuntime.execute` / `AgentRunRepository.claim`。 + 看代码时别被这个类带偏。 +3. **`app/worker/graph_projection_worker.py` 与 `ProjectionReconciliationService` 都是"没接线"的**: + 前者全仓零实例化点,后者没有生产调用点(只有工具里提到)。 + 真正在跑的图投影是 `ProfileGraphProjectionService`(见第一部分 1.2 的 ⑤ 步)。 + +### 2.12 一页速查 + +| 想看什么 | 看哪里 | +|---|---| +| 请求有没有被处理 | `agent_run.status` / `worker_id` / `locked_until` / `attempt_count` | +| 事件有没有堆积 | `domain_event_outbox.status`(`pending`/`failed`/`published`/`dead`)+ `retry_count` + `next_retry_at` | +| 记忆有没有进向量库 | `memory_sync_outbox`(`target_store ∈ {milvus, neo4j}`,**小写**) | +| 片段有没有被提炼 | `episodes.extraction_status` + `promoted_to_ltm` | +| 场外有没有在收信 | `offsite_mail_cursor.status`(`blocked` 就是卡住了) | +| 记忆删除有没有清干净 | `interaction_audit` 里 `action_type='memory.projection_cleanup'` 的 `detail.status` | + +**重启 Worker 的三种方式**:`start.ps1`(连 API 一起)、`python -m app.worker`(只起它)、 +`python -m app.worker --once`(跑一轮看有没有报错)。 +停掉它不会让已有数据坏掉,只会让"慢活"暂停——重新起来会接着做(队列是数据库里的)。