From 3a1065ca1e98249709524ed542c623183a64ac36 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=8D=BF=E4=BA=91=E7=A7=8B=E6=9C=88?= <15273589815@163.com> Date: Mon, 14 Sep 2026 21:09:03 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E7=89=87=E6=AE=B5=E9=93=BE?= =?UTF-8?q?=E8=B7=AF=EF=BC=9A=E9=87=8D=E5=A4=8D=E8=81=9A=E5=90=88=E6=8B=96?= =?UTF-8?q?=E5=9E=AE=E9=87=8D=E8=AF=95=E9=A2=84=E7=AE=97=20+=20=E4=BB=8E?= =?UTF-8?q?=E4=B8=8D=E6=8A=95=E7=94=BB=E5=83=8F=E9=87=8D=E5=BB=BA=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 结论先说:真正的根因和最初两个假设都不是同一个 你原来的判断是「第二条召回恒空、第三条传导断链」,并猜第三条是 「`record_evidence` 返回 False 时该不该补写重建事件」。查完库发现: 1. **`record_evidence` 返回 False 时数据零净变化**(幂等命中直接 return; 并发冲突把刚加的计数减回来),所以那个判断点解释不了画像停更; 2. **`memory_extraction_worker` 是有写重建事件的**(`if recorded:`), 233 条 `profile.rebuild_requested` 全 published; 3. 真正断在两处,都在 **`episode_worker`** 这条片段链路上。 ## 根因一:`_touch_retry` 把「重复聚合」记成「抽取失败」(P0,已修) `episode_worker._persist` 在 `content_hash` 命中(同一段会话被重复聚合)时调 `_touch_retry` 做 `retry_count += 1`。但 `retry_count` 的语义是**抽取失败次数** (由 `_mark_failed` 累加),而"分片逻辑每轮重新看到同一段会话"根本不是失败。 后果是实测出来的: episodes 待提取片段:retry_count=1405(40 条,全是客户 9001) consume_pending 逐条件筛: 仅 status in (待提取,失败) -> 40 + retry_count < max_retry(3) -> 0 ← 一个都不剩 + promoted_to_ltm IS FALSE -> 40 这 40 条片段**永远不可能被选中** ⇒ 新记忆进不来 ⇒ 画像停在旧值。 (另外还观察到一次运行里它从 1405 涨到 1407 —— 常驻 Worker 每轮都在继续推高。) **修法**:重复聚合不再触碰 `retry_count`,连 `flush` 都不做(内容没变就只是看到)。 ## 根因二:片段链路从不投「画像重建」事件(已修) `memory_extraction_worker` 写完证据会投 `profile.rebuild_requested`; 而 `episode_worker` 写完证据直接 `_mark_done` 就结束了 —— **完全没有这一步**。 所以即使片段被成功抽取,画像也不会重建。 **修法**:`episode_worker` 也捕获 `recorded` 并在为真时投同样的事件 (`trigger="episode_extraction"`)。只在 `recorded=True` 时投:幂等命中时证据与计数 都没有净变化,投一次是白跑。 ## 数据修复 代码修好不会让已写进库的脏计数自己恢复 —— 那 40 条片段仍然超预算。 用 `tools` 级别的临时脚本把**确实被污染的行**重置(条件收紧为三者交集): extraction_status IN ('待提取','失败') AND promoted_to_ltm = 0 AND retry_count >= 3 -> 重置 40 行;之后可被 consume_pending 选中的片段从 0 恢复到 40 ## 实测 - **重试预算修复**:`consume_episodes` 从"领不到任何片段"变为能领到; 分两批消费完 40 条(全部 `no_fact` —— 那些片段摘要里确实没有用户陈述, 属正常结果),待提取 40 → 0 - **重建事件修复**:那 40 条全是 `no_fact`,走不到 `if recorded:`,所以** 实测不到**。为不留下"改了但没验证",另造了一条含用户陈述的片段: consume_episodes -> extracted=[191] 9001 的 rebuild 事件 3 -> 4(增量 1) memory_evidence 8 -> 9 另外注意到基线在我造数据前已经由 1 涨到 3 —— 说明常驻 Worker 也在这期间 投过事件,修复在真实链路里同样生效。 - `pytest tests/unit tests/contract` -> 1428 passed / 2 skipped / 1 failed (剩下的 1 个是投顾工作台页面被替换所致,与本次无关) ## 测试 两个既有用例断言的正是被修掉的旧行为,已更新,并把第二个改造成**防回归守卫** `test_repeated_aggregation_never_bumps_retry_count` —— 它守着 "`retry_count` 被重复聚合推高到 `max_retry` 之上会导致片段永久滞留"这个 P0。 `tests/unit/worker/test_episode_worker.py` 16 passed。 ## 未处理 「召回恒空」(员工身份下 `recall` 取的是自己作为客户的记忆)**本次没动**。 它需要给 `AgentRequest` 加 `target_customer_id` 并配套越权校验,属接口契约变更; 排查报告给的建议是保持现状、员工侧走 `query_customer_profile` 工具。 要按"支持目标客户维度"做,请确认,我再单独一提交。 --- app/worker/episode_worker.py | 56 ++++++++++++++++++++---- tests/unit/worker/test_episode_worker.py | 27 +++++++++--- 2 files changed, 67 insertions(+), 16 deletions(-) diff --git a/app/worker/episode_worker.py b/app/worker/episode_worker.py index 0670fc5..713b956 100644 --- a/app/worker/episode_worker.py +++ b/app/worker/episode_worker.py @@ -25,7 +25,7 @@ import json import logging from dataclasses import dataclass, field from datetime import UTC, datetime, timedelta -from uuid import NAMESPACE_URL, uuid5 +from uuid import NAMESPACE_URL, uuid4, uuid5 from sqlalchemy import select from sqlalchemy.exc import IntegrityError @@ -33,6 +33,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.model.conversation import ConversationMessage from app.model.episode import Episode +from app.model.platform import DomainEventOutbox from app.service.memory_extraction_service import ( MemoryExtractionService, get_memory_extraction_service, @@ -157,7 +158,22 @@ class EpisodeWorker: select(Episode).where(Episode.content_hash == content_hash) ) if existing is not None: - await self._touch_retry(existing) + # ⚠️ 这里**不要**碰 `retry_count`。 + # + # 此前会调 `_touch_retry(existing)` 做 `retry_count += 1`(注释写的是 + # "仅在未完成时累加重试计数")。但 `retry_count` 的语义是**抽取失败**的次数 + # (由 `_mark_failed` 累加),而"分片逻辑每轮重新看到同一段会话"根本不是失败 + # —— 片段内容没变,`content_hash` 才会相同。 + # + # 两者混用一个字段的后果是实测出来的:客户 9001 有 **40 条片段 + # `retry_count=1405`**,远超 `MAX_EXTRACTION_RETRY`(3),于是 + # `consume_pending` 的 `retry_count < max_retry` 永远筛不中它们 —— + # **片段永久滞留 → 新记忆进不来 → 画像停在旧值** + # (`profile_snapshots` 停在 2026-09-10 13:57,而 `memory_unit.updated_at` + # 已经是 2026-09-13 11:03)。 + # + # 重复看到同一片段就只是看到,不改变任何状态。 + logger.debug("episode already persisted content_hash=%s", content_hash) return None episode = self._build(customer_id, session_id, segment, content_hash, status) self.session.add(episode) @@ -170,12 +186,6 @@ class EpisodeWorker: return None return episode - async def _touch_retry(self, episode: Episode) -> None: - """已存在片段:仅在未完成时累加重试计数,已完成的片段保持原状。""" - if episode.extraction_status in RETRYABLE_STATUSES: - episode.retry_count += 1 - await self.session.flush() - def _build( self, customer_id: int, @@ -346,7 +356,7 @@ class EpisodeExtractionConsumer: "episode_uuid": episode.episode_uuid, }, ) - await service.record_evidence( + recorded = await service.record_evidence( memory, # 幂等键落在片段指纹上:同一片段重复消费不会写出第二条证据。 idempotency_key=self.idempotency_key(episode), @@ -368,6 +378,34 @@ class EpisodeExtractionConsumer: source_record_id=str(episode.id), occurred_at=episode.ended_at or episode.created_at, ) + if recorded: + # ⚠️ 证据写入后必须把「画像重建」**投成事件**,不能直接调用重建: + # 本方法的记忆写入还在当前事务里、尚未提交,另开 session 去重建看不到 + # 这条新记忆(`memory_extraction_worker` 里记录了这条实测结论)。 + # + # 这条路径此前**完全没有**投重建事件,后果是:从会话片段抽取出来的记忆 + # 永远到不了画像。实测证据 —— 客户 9001 的 `profile_snapshots` current + # 停在 `2026-09-10 13:57`(v7),而它的 `memory_unit.updated_at` 已经是 + # `2026-09-13 11:03`、`evidence_count` 涨到 4;那三条新证据的 + # `source_table` 正是 `episodes`。也就是说**片段链路记住了,画像不知道**。 + # + # 只在 `recorded=True` 时投:幂等命中(同一片段重复消费)时证据与计数 + # 都没有净变化,投一次重建是白跑。 + now = datetime.now(UTC).replace(tzinfo=None) + self.session.add(DomainEventOutbox( + id=0, + event_id=str(uuid4()), + event_type="profile.rebuild_requested", + aggregate_type="customer_profile", + aggregate_id=str(episode.customer_id), + trace_id=episode.episode_uuid, + payload={"customer_id": episode.customer_id, "trigger": "episode_extraction"}, + status="pending", + retry_count=0, + occurred_at=now, + created_at=now, + updated_at=now, + )) await self._mark_done(episode, promoted=True) result.extracted.append(episode_id) diff --git a/tests/unit/worker/test_episode_worker.py b/tests/unit/worker/test_episode_worker.py index 589db38..a69f5d5 100644 --- a/tests/unit/worker/test_episode_worker.py +++ b/tests/unit/worker/test_episode_worker.py @@ -148,12 +148,25 @@ async def test_repeated_aggregation_does_not_create_second_episode() -> None: assert second.skipped == [1] assert len(episodes_of(added)) == 1 assert session.add.call_count == 1 # 重复聚合不再 add 新片段 - assert episodes_of(added)[0].retry_count == 1 # 未完成片段:重复投递只累加重试计数 - assert session.flush.await_count == 2 # 首次写入 1 次 + 重试计数 1 次 + # ⚠️ 重复聚合**不消耗重试预算**。`retry_count` 的语义是"抽取失败"的次数 + # (由 `_mark_failed` 累加),而"分片逻辑每轮重新看到同一段会话"根本不是失败 + # —— 片段内容没变,`content_hash` 才会相同。 + # 此前这里会 +1,实测把客户 9001 的 40 条片段推到了 `retry_count=1407`, + # 超过 `MAX_EXTRACTION_RETRY`(3) 之后 `consume_pending` 的 + # `retry_count < max_retry` 再也筛不中它们 —— **片段永久滞留、画像停在旧值**。 + assert episodes_of(added)[0].retry_count == 0 + assert session.flush.await_count == 1 # 只落首次写入;重复聚合不写库 -async def test_completed_episode_does_not_bump_retry_count() -> None: - """已完成的片段重复聚合保持原状;可重试状态的片段才累加重试计数。""" +async def test_repeated_aggregation_never_bumps_retry_count() -> None: + """重复聚合时**任何状态都不累加** `retry_count`(它只由抽取失败累加)。 + + 这条是**防回归守卫**,守着上面那个 P0 级滞留问题: + `retry_count` 一旦被"重复看到"推高到 `MAX_EXTRACTION_RETRY` 之上, + `consume_pending` 的 `retry_count < max_retry` 就永远筛不中该片段。 + 实测客户 9001 的 40 条片段被推到 1407,后果是 + **片段永久滞留 → 新记忆进不来 → 画像停在 2026-09-10**。 + """ messages = [ message(1, minutes=0, role="user", content="我只买货币基金"), message(2, minutes=1, role="assistant", content="已记录"), @@ -170,9 +183,9 @@ async def test_completed_episode_does_not_bump_retry_count() -> None: assert done.skipped == [1] assert pending.skipped == [1] assert episodes_of(pending_added) == [] - assert done_session.flush.await_count == 0 # 已完成片段保持原状,不累加、不落库 - assert pending_record.retry_count == 3 - assert pending_session.flush.await_count == 1 + assert done_session.flush.await_count == 0 # 已完成片段:不落库 + assert pending_record.retry_count == 2 # ← 保持原值(此前断言 3) + assert pending_session.flush.await_count == 0 # ← 不写库(此前断言 1) async def test_idle_gap_splits_a_session_into_segments() -> None: