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: