Feature/advanced memory service - #336
Conversation
Please enter the commit message for your changes. Lines starting
AI Code Review审查结论不通过 审查范围:base f05797d..head fdc52ae(112 文件,+6182/-6678),主题为 advanced memory 服务重构:新增 发现的问题严重
问题: 中等
问题: 中等
问题: 中等
问题: 较低
问题: |
| async with self._session_lock(session_id): | ||
| request.contents = [content.model_copy(deep=True) for content in request.contents] | ||
| state = await self._load_state(session_id) | ||
| reapplied = False | ||
| if state.latest_compaction is not None: | ||
| reapplied = self._apply_record(request, state.latest_compaction) | ||
|
|
||
| request_chars_before = estimate_request_chars(request) | ||
| token_budget_before = tracker.budget(request, ctx) | ||
| token_mode = token_budget_before.token_mode_enabled | ||
| request_tokens_before = token_budget_before.estimate.tokens | ||
| comparison_tokens_before = (tracker.estimate_request_tokens(request) if token_mode else None) | ||
| blocking_reached = (request_tokens_before >= token_budget_before.blocking_threshold_tokens | ||
| if token_mode else request_chars_before >= self._auto_compact_config.blocking_chars) | ||
| if state.consecutive_failures >= auto_compact_config.max_failures and blocking_reached: | ||
| return AdvancedAutoCompactResult( | ||
| compacted=False, | ||
| reapplied=reapplied, |
There was a problem hiding this comment.
问题: _apply_scoped 中的压缩回放(replay)永远失效:_persist_session_compaction(407-454 行)每次成功压缩后通过 Session.compact_events 把 ctx.session.events 替换为「summary 事件 + recent 事件」并调用 update_session 持久化,旧 boundary 事件被移入 historical_events。下一次请求的 request.contents 由新的 active events 重建(BEFORE_MODEL 经 filter 传入的 LlmRequest、AFTER_TURN 经 _request_from_events 均不再包含旧 boundary),_apply_record(305-320 行)用 state.latest_compaction.boundary_signature 在 _find_signature_index 中无法命中,reapplied 恒为 False,每次都落入 735-744 行的 legacy 分支重新调用 _legacy_summary 生成新摘要。触发条件: 上下文持续增长、多次超过 trigger_chars/auto_compact_threshold_tokens 的场景(如多轮工具调用、长上下文会话),每次超阈值都会再次压缩。「回放旧压缩」的设计意图(避免重复压缩)在该状态下无法生效。实际影响: 每个 epoch 都产生一次额外的 LLM 摘要调用并将新的 summary 前缀插入 session.events,模型上下文累积多层「compacted context」摘要导致上下文失真、token 预算失真、延迟与成本上升;historical_events 中 summary 边界不断前移,压缩数据完整性受损。修正方向: 压缩后从 compact_events 落库的 summary 事件中读回持久化状态(如按 session_compaction_id 定位最新 summary 事件),使状态可以直接从 Session/historical 恢复当前的 boundary_signature/occurrence;或将 latest_compaction 与持久化摘要事件绑定(_persist_session_compaction 后回读确认),确保再次超阈值时 _apply_record 能命中最近一次压缩边界而不是重新发起 legacy 全量压缩。
| ctx: InvocationContext, | ||
| ) -> None: | ||
| """Stage the final request fingerprint for persistence on the response Event.""" | ||
| session = getattr(ctx, "session", None) if ctx is not None else None | ||
| state = getattr(session, "state", None) | ||
| session = ctx.session | ||
| state = session.state | ||
| if isinstance(state, dict): | ||
| state["advanced_memory_pending_request_context_fingerprint"] = _request_static_fingerprint(request) | ||
|
|
||
| def budget( |
There was a problem hiding this comment.
问题: record_request_context(_token_budget.py:245-253,由 _auto_compact.py:581 在每次非阻断压缩后调用)把 advanced_memory_pending_request_context_fingerprint 写入 session.state,但本次变更删除了唯一把它转移到事件 custom_metadata 的代码(旧 trpc_agent_sdk/advanced_memory/_session_service.py:160-166 的 append_event 转移逻辑),新原生 InMemorySessionService/SqlSessionService/RedisSessionService.append_event 均无等效写入,仓库内无任何消费者。_latest_usage_baseline(168-198 行)读取 advanced_memory_request_context_fingerprint 的失效守卫(如果 recorded_fingerprint 非 None 且与当前静态指纹不同则跳过该 usage 事件)因此永不触发。触发条件: 任何启用 token 模式(配置了 model_context_window_tokens 或有 resolver)的会话:会话中发生过带 usage_metadata 的事件之后,应用修改了 system_instruction/tools 配置(静态指纹变化),再次发送请求时 budget() 被调用。实际影响: usage 事件指纹失效守卫形同虚设——_latest_usage_baseline 按内容指纹匹配到最后一次 usage 事件并直接复用其过期的 usage_tokens 作为当前请求的估算(source="usage"/"hybrid"),即使当前请求因配置缩小而远未到窗口上限,也会被判定为「窗口已满」提前触发压缩甚至阻断;反过来请求变大使估算偏小也可能跳过必要的压缩。该 state 键还永久滞留 session.state 并在每次后端全量往返(SQL/Redis save_session)中被反复重写。修正方向: 在会话服务 append_event 路径(或 runner 追加事件处)把 session.state 中的 pending fingerprint 迁移到带 usage_metadata 事件的 custom_metadata(键名 advanced_memory_request_context_fingerprint),与 base 行为对齐;若该机制不再需要则同时删除 state 写入与失效守卫。
| async def _load_state(self, session_id: str) -> AdvancedAutoCompactState: | ||
| """Restore process-local compaction state.""" | ||
| state_key = self._runtime.session_key(session_id) | ||
| state = self._states.get(state_key) | ||
| if state is not None: | ||
| return state | ||
| state = AdvancedAutoCompactState(latest_compaction=None, consecutive_failures=0) | ||
| self._states[state_key] = state | ||
| return state |
There was a problem hiding this comment.
问题: _load_state 无条件返回 AdvancedAutoCompactState(latest_compaction=None, consecutive_failures=0),_history_snip._load_state(106-117 行)同样返回空的 snipped_ids/result_hashes,_micro_compact/_tool_result_budget 同款。base 提交中这些状态由 transcript 持久化并在重启后恢复(autocompact-success/autocompact-failure、history-snip、content-replacement 记录),本次变更删除了所有 _persist_success/_persist_failure/transcript 写入而没有替代的持久化。触发条件: 进程重启(部署、崩溃、runner 重建)后同一 session 继续使用 AdvancedAutoCompactSummarizer。实际影响: 熔断计数器归零——即使上次已连续失败接近 max_failures,重启后会再次尝试并把接近硬限制的超大请求放行到模型(应在 blocked 时直接拒绝的请求不再被拦住);已经 snip/clear 的 tool result 会被再次替换处理,result_hashes 冲突检测失效,重复处理与状态漂移。修正方向: 将熔断计数与已处理结果哈希持久化到 session.state(经 update_session_state,与 _SESSION_MEMORY_STATE_KEY 同机制)或专用 transcript,_load_state 从其中恢复而非固定清零。
| async def _refresh_memory_scope(self, db: Any) -> None: | ||
| expiry = self._expiry(self._config.memory_ttl_seconds) | ||
| if expiry is None: | ||
| return | ||
| index = await self._storage.get( | ||
| db, SqlKey( | ||
| key=(self._app_name, self._user_id), | ||
| storage_cls=AdvancedMemorySqlIndex, | ||
| )) | ||
| if index is not None: | ||
| index.expires_at = expiry | ||
| topics = await self._storage.query( | ||
| db, | ||
| SqlKey(key=(self._app_name, self._user_id), storage_cls=AdvancedMemorySqlTopic), | ||
| SqlCondition(filters=[ | ||
| AdvancedMemorySqlTopic.app_name == self._app_name, | ||
| AdvancedMemorySqlTopic.user_id == self._user_id, | ||
| AdvancedMemorySqlTopic.expires_at.is_(None) | ||
| | (AdvancedMemorySqlTopic.expires_at > self._now()), | ||
| ]), | ||
| ) | ||
| for topic in topics: | ||
| topic.expires_at = expiry |
There was a problem hiding this comment.
问题: _refresh_memory_scope 在每次读取(read_index 141-151 行、read_topic 203-205 行、list_topics 246-247 行都调用它)时把 index 行和该租户下所有未过期 topic 行的 expires_at 重新置为 now + memory_ttl_seconds。只要租户被持续读取,expires_at 就不断被推远,_expired 永远判不过期;Redis 后端 _refresh_ttl_group(_redis_storage.py:105-125 行)在每次读后对 registry 内所有键 EXPIRE,行为相同。触发条件: 配置了 memory_ttl_seconds(非 None)并使用 SQL/Redis 后端,且租户有任一常规读取流量(如每请求注入 <advanced-memory-index> 的 read_index)。实际影响: memory_ttl_seconds 对活跃租户永不生效——过期记忆永远不会被 cleanup 清理,数据无限累积,与配置意图和本地后端的读路径行为(_file_storage.py 读不刷新文件 mtime)不一致。修正方向: 读路径不应刷新 expires_at(TTL 应基于写入时间),或仅在写入/初始化时刷新;如需滑动语义,应只刷新被访问的单条记录而非整个租户范围。
| result = await summarizer_manager.create_session_summary_before_model( | ||
| req, | ||
| invocation_ctx, | ||
| ) | ||
| if not result: | ||
| return None | ||
| invocation_ctx.end_invocation = True | ||
| rsp.rsp = result | ||
| rsp.is_continue = False | ||
| rsp.error = None | ||
| session_memory_extractor = summarizer_manager.get_session_memory_extractor() | ||
| if session_memory_extractor is not None: | ||
| await session_memory_extractor.extract_if_needed( | ||
| invocation_ctx, | ||
| force=False, |
There was a problem hiding this comment.
问题: _before 在命令压缩阻塞/完成后设置 invocation_ctx.end_invocation = True,但仓库内(含 filter 框架)没有任何代码读取该字段——end_invocation 只在 agents/_callback.py 的 legacy 回调路径被同一属性自赋值使用,框架控制流仅依赖 rsp.is_continue。触发条件: 挂载 AdvancedAutoCompactSummarizerFilter 且命令压缩返回非 None(产生拦截响应)时。实际影响: 该赋值是无作用的死代码,不会造成功能错误,但给维护者错误的信号(以为设置了「结束本次调用」语义,实际没有),且无法支撑后续基于 end_invocation 的终止判断。修正方向: 删除该赋值,或若需要「终止本次模型调用」语义,改由 rsp.is_continue = False 的既有机制表达(当前已设置),并在 filter 框架中显式消费该字段。
fdc52ae to
9c37094
Compare
AI Code Review审查结论不通过 审查范围: 发现的问题严重
问题: 本次新增的 触发条件: 任何使用 实际影响: 会话事件追加与状态更新在 SQLite 上全部抛 修正方向: 在 中等
问题: usage-baseline fingerprint 的写入管线在本次重构中被删除:base 版本 触发条件: 配置了 实际影响: token 估计永远回退到 修正方向: 在 中等
问题: token 模式下 "request token estimate did not reduce" 的成功检查使用了与异常分支不一致的两套估算:压缩前 触发条件: token 模式下,usage 基线( 实际影响: 真实有效的压缩被判定失败, 修正方向: 统一两处使用同一估算(全部改用 较低
问题: 触发条件: 任何会话事件进入 实际影响: 现版本无功能影响;一旦 修正方向: 简化为 |
| async def get_for_update(self, db: SqlSession, key: SqlKey) -> Any: | ||
| """Get one row while holding a database row lock until commit.""" | ||
| stmt = select(key.storage_cls) | ||
| for column, value in zip(inspect(key.storage_cls).primary_key, key.key): | ||
| stmt = stmt.where(column == value) | ||
| stmt = stmt.with_for_update() | ||
| if isinstance(db, AsyncSession): | ||
| result = await db.execute(stmt) | ||
| else: | ||
| result = db.execute(stmt) | ||
| return result.scalars().first() |
There was a problem hiding this comment.
问题: 本次新增的 SqlStorage.get_for_update 无条件构造 select(...).with_for_update(),编译出 SELECT ... FOR UPDATE;SQLite 不支持该子句(已实测 sqlite3 报 OperationalError: near "FOR": syntax error),且 SQLAlchemy 对 SQLite 本就没有行锁语义。本仓库 SQL 后端会话服务的主要后端即 SQLite(tests/sessions/test_sql_session_service.py 使用 sqlite:///:memory:),因此该路径是必然触发的崩溃而非理论问题。
触发条件: 任何使用 SqlSessionService 且后端为 SQLite 的部署,执行任意一次 append_event(_sql_session_service.py:549)或 update_session_state(_sql_session_service.py:641)即触发;MySQL/PostgreSQL 不受语法影响但同样失去本应提供的锁语义前提。
实际影响: 会话事件追加与状态更新在 SQLite 上全部抛 OperationalError,会话记录无法落盘、状态无法持久化;同时 with_for_update 行锁无法提供跨进程串行化,多进程并发写入时会话状态与事件窗口可能互相覆盖。
修正方向: 在 get_for_update 中检测方言(dialect.supports_for_update),SQLite 下回退为普通 select 并将并发保护降级/改为基于版本号或 updated_at 条件更新的乐观并发方案;同时补一个 SQLite 后端的回归用例覆盖 append_event 与 update_session_state。
| @@ -200,10 +192,10 @@ def _latest_usage_baseline( | |||
| ) | |||
| return None | |||
There was a problem hiding this comment.
问题: usage-baseline fingerprint 的写入管线在本次重构中被删除:base 版本 trpc_agent_sdk/advanced_memory/_session_service.py(append_event 中)会把 advanced_memory_pending_request_context_fingerprint 写入响应事件的 custom_metadata[advanced_memory_request_context_fingerprint],而新架构删除该服务后,全仓库只有 _token_budget.py:242 的 record_request_context 仍在写 session.state,没有任何代码将其转写到事件元数据(仓库内无消费者,唯一测试也是手工设置该字段)。
触发条件: 配置了 model_context_window_tokens(token 模式)且模型返回 usage_metadata 的任何会话;每次 _latest_usage_baseline(_token_budget.py:160-193)扫描事件时都取不到 fingerprint,recorded_fingerprint is None 直接通过。
实际影响: token 估计永远回退到 estimate_request 的 heuristic JSON 估算,真实 usage 计费数据被忽略:大缓存上下文可能被过早压缩(token 低估),或将越过真实模型窗口的请求放行(token 高估),影响会话上下文正确性与稳定性。
修正方向: 在 append_event 或等价持久化点恢复写入:当事件携带 usage_metadata 且 session.state 中存在 pending fingerprint 时,将 advanced_memory_request_context_fingerprint 写入事件 custom_metadata 并随事件持久化(可在 update_session_state 或 Session.compact_events 边界处理)。
| request_chars_after = estimate_request_chars(request) | ||
| if token_mode: | ||
| comparison_tokens_after = tracker.estimate_request_tokens(request) | ||
| if (comparison_tokens_after >= comparison_tokens_before | ||
| and request_chars_after >= request_chars_before): | ||
| raise ValueError("Advanced Auto Compact did not reduce request token estimate") | ||
| elif request_chars_after >= request_chars_before: | ||
| raise ValueError("Advanced Auto Compact did not reduce request size") |
There was a problem hiding this comment.
问题: token 模式下 "request token estimate did not reduce" 的成功检查使用了与异常分支不一致的两套估算:压缩前 comparison_tokens_before(_auto_compact.py:654)与压缩后 comparison_tokens_after(_auto_compact.py:753)都是 estimate_request_tokens 的 heuristic 全量估算,而失败路径(_auto_compact.py:786-788)报告并用于失败递增的却是 request_tokens_before = token_budget_before.estimate.tokens(usage/hybrid 基线)。
触发条件: token 模式下,usage 基线(source=usage/hybrid)大于 heuristic 全量估算,且压缩只削减了部分 token(真实减少但低于阈值差)的场景;此时 comparison_tokens_after >= comparison_tokens_before 成立,误报 "did not reduce"。
实际影响: 真实有效的压缩被判定失败,consecutive_failures 持续累加,达到 max_failures 后进入硬阻塞(返回 blocked 响应终止推理),会话被锁死;同时该异常分支的 error 三元表达式 str(exc) if exc else None if blocked else None 优先级错误,永远无法按设计区分异常与阻塞。
修正方向: 统一两处使用同一估算(全部改用 estimate_request_tokens heuristic 或全部使用 usage 基线),并修正 error 表达式为 (str(exc) if exc else None) if blocked else None 以正确传达失败原因。
| if event.is_summary_event and event.is_summary_event(): | ||
| continue |
There was a problem hiding this comment.
问题: _session_event_records 中 event.is_summary_event and event.is_summary_event() 是调用一个 bound method 后对返回值做布尔判断,再对同一个 method 调用一次;若未来 Event 的 is_summary_event 改为属性或取得返回值与调用结果不一致,该条件会静默失效。当前 Event.is_summary_event 是方法,实际行为正确,但该写法把代码正确性绑定到了方法对象本身非空这一偶然特性。
触发条件: 任何会话事件进入 _session_event_records(会话记忆抽取路径)时;当前不会出错,属于脆弱写法而非即刻缺陷。
实际影响: 现版本无功能影响;一旦 Event.is_summary_event 演化(改属性/改名),summary 事件会被误当成待抽取事件纳入会话记忆文档,导致重复信息甚至边界错位。
修正方向: 简化为 if event.is_summary_event():,消除对方法对象真值性的依赖。
9c37094 to
1bbcc4f
Compare
1bbcc4f to
21633b4
Compare
AI Code Review审查结论通过 审查范围:base_commit f05797d .. head_commit 21633b4(计划:Feature/advanced memory service),共 116 个文件(+6255/-6718),并把 Advanced Memory 从 trpc_agent_sdk/advanced_memory/ 迁移为 trpc_agent_sdk/tools/advanced_memory/(local/Redis/SQL 三种后端),新增 trpc_agent_sdk/sessions/compact/(advanced 压缩管道与 default 摘要器重构),重写 sessions/_session.py 的 compact_events、abc/_compact.py 与三套 SessionService 的 update_session_state 等框架层。 计划符合性:功能主体(三种存储后端、TTL/清理、分布式锁、Session Memory 提取、ToolResultBudget/HistorySnip/MicroCompact/AutoCompact 管道、before_model 过滤器、compact_events 双列表)均已实现且大部分可运行;未发现数据损坏级或必现崩溃级缺陷,因此无 SEVERE 评论。 主要风险(已逐项经代码与 git 历史验证):
被否决的候选(有过硬反证):seen_ids 仅做集合成员运算而非 dict 键;_before blocked 响应经 filter 适配器正常流出(无 AttributeError);migrate_legacy 的父目录由 mkdir(parents=True)+ensure_base_directories 保证存在(无 FileNotFoundError);get_session_summary 返回 cache 对象与 conversation_count 重置行为与 base 一致(非本变更引入);SQL write_index 用全新 entries 覆盖过期行内容(无"过期内容复活"路径)。无兼容性/导入级崩溃。测试覆盖:tests/advanced_memory/ 与 tests/sessions/compact/ 覆盖了主要路径,但对死代码指纹、并发 update_session_state、config=None 构造与 Redis 锁超时释放等场景无测试。 门禁结论:PASSED(无 SEVERE,5 个 MODERATE / 2 个 LOW,均为该变更引入或使其可达的真实问题)。 发现的问题中等
问题: 触发条件: 任何调用方以无参方式构造公开导出的 实际影响: 进程启动即抛 AttributeError 崩溃;目前只是恰好所有示例与测试都显式传了 config 才掩盖了该缺陷,一旦有调用方省略配置参数即不可用。 修正方向: 使用局部变量承接默认值,例如 中等
问题: 触发条件: 启用 token 模式( 实际影响: usage 基线(来自上一个模型响应 Event 的 修正方向: 在响应 Event 追加/持久化路径中把暂存的指纹转入 Event 中等
问题: 新增的 触发条件: Redis/内存后端 + 同一会话在主 loop 继续追加事件/更新 state 的同时,后台 post-turn worker( 实际影响: 快照覆盖导致 state(如 修正方向: 对 Redis 采用 WATCH/MULTI-EXEC 或 Lua 脚本化的条件比较写(如仅当 key 版本未变才写回),或至少让 较低
问题: 触发条件: 长会话中压缩器连续 3 次失败(如摘要 LLM 反复报错、后端存储异常)且请求接近上下文硬限制。 实际影响: 用户侧看到的是模型“正常回复”一段拒绝文本(“Automatic context compaction has failed repeatedly…”),没有错误码、异常或终止信号,应用层无从区分压缩失败与模型主动拒答;触发后仅靠该文本提示用户,无法编程化处理。 修正方向: 让 blocked 路径产出带 较低
问题: 触发条件: 锁竞争激烈且某次获取超时( 实际影响: 在极端时间窗内可能删除并发持有者(后续进程)的锁,破坏跨进程互斥,使两个进程同时进入 修正方向: 仅当 较低
问题: 触发条件: 常驻多租户部署(服务化场景,本 PR 目标)中 app/user/session 数随时间增长。 实际影响: 进程内存随唯一 app/user/session 数量线性增长且不回收(每个 session 的 修正方向: 为缓存引入容量上限/淘汰(如按最后访问时间剔除过期 session 的状态与锁,scope 级 processor 按 LRU 限制),或在会话关闭/删除时显式清理对应键。 较低
问题: 触发条件: 每轮对话的每次模型请求(before_model 管道)都会执行多次 实际影响: 读路径 SQL 放大抬高延迟与连接池压力(高并发下明显),TTL 语义变成仅对完全冷数据生效,与配置预期(按 memory_ttl_seconds 过期删除)不符;文件/Redis 后端同样在每次读时续期。 修正方向: 将 TTL 续期收敛到写路径(写入时设置 |
| def __init__( | ||
| self, | ||
| config: AdvancedAutoCompactSummarizerConfig | None = None, | ||
| *, | ||
| model: LLMModel | None = None, | ||
| session_memory_extractor: SessionMemoryExtractor | None = None, | ||
| ) -> None: | ||
| """Initialize the compressor, summary generator, and session locks.""" | ||
| self._model = model | ||
| self._runtime = AdvancedAutoCompactSummarizerRuntime(config=config or AdvancedAutoCompactSummarizerConfig()) | ||
| self._auto_compact_config = config.auto_compact | ||
| self._session_memory_extractor = self._create_extractor(session_memory_extractor) |
There was a problem hiding this comment.
问题: AdvancedAutoCompactSummarizer.__init__ 中 self._runtime = AdvancedAutoCompactSummarizerRuntime(config=config or AdvancedAutoCompactSummarizerConfig()) 之后,第 128 行 self._auto_compact_config = config.auto_compact 直接解引用形参 config;当调用方省略 config(构造器签名显式允许 config: AdvancedAutoCompactSummarizerConfig | None = None,且上一行刚通过 or 回退生成了默认配置)时,抛 AttributeError: 'NoneType' object has no attribute 'auto_compact',默认配置回退形同虚设。
触发条件: 任何调用方以无参方式构造公开导出的 AdvancedAutoCompactSummarizer()(例如从 YAML/配置缺失块构建 manager,或遵循 runtime config 自带 default_factory 的默认值语义)。
实际影响: 进程启动即抛 AttributeError 崩溃;目前只是恰好所有示例与测试都显式传了 config 才掩盖了该缺陷,一旦有调用方省略配置参数即不可用。
修正方向: 使用局部变量承接默认值,例如 resolved = config or AdvancedAutoCompactSummarizerConfig(),然后 self._runtime = AdvancedAutoCompactSummarizerRuntime(config=resolved)、self._auto_compact_config = resolved.auto_compact,并补一个无参构造的单元测试。
| @classmethod | ||
| def record_request_context( | ||
| self, | ||
| request: "LlmRequest", | ||
| ctx: "InvocationContext | None", | ||
| cls, | ||
| request: LlmRequest, | ||
| ctx: InvocationContext, | ||
| ) -> None: | ||
| """Stage the final request fingerprint for persistence on the response Event.""" | ||
| session = getattr(ctx, "session", None) if ctx is not None else None | ||
| state = getattr(session, "state", None) | ||
| session = ctx.session | ||
| state = session.state | ||
| if isinstance(state, dict): | ||
| state["advanced_memory_pending_request_context_fingerprint"] = _request_static_fingerprint(request) | ||
|
|
There was a problem hiding this comment.
问题: record_request_context 把请求指纹写入 session.state["advanced_memory_pending_request_context_fingerprint"],但全仓库没有任何代码读取该键;而 _latest_usage_baseline(第 178 行)从 Event 的 custom_metadata["advanced_memory_request_context_fingerprint"] 读取指纹,生产代码中没有任何地方把指纹写入 Event metadata(仅测试手动构造)。指纹守卫是死代码。
触发条件: 启用 token 模式(model_context_window_tokens 显式配置或注入 context_window_resolver)时,任一请求都会进入 _latest_usage_baseline;写入 application/system 指令、更换工具集、或同一 Session 上并发/切换请求时,recorded_fingerprint 恒为 None,守卫永不触发。
实际影响: usage 基线(来自上一个模型响应 Event 的 usage_metadata token 数)会被无条件当作当前请求的基线,叠加后缀估算得到失真的预算(warning/auto_compact/blocking 阈值比较),可能提前或推迟压缩,甚至把另一条请求链的用量误算进当前请求,导致错误触发 blocking 或预算判断;作者设计该守卫的意图(跨上下文不连续时回退到全量估算)完全失效。
修正方向: 在响应 Event 追加/持久化路径中把暂存的指纹转入 Event custom_metadata(或在 _latest_usage_baseline 中读取 state 暂存值并校验),使指纹写入与读取在同一生命周期内闭环;并补充“指纹不同→回退 estimated”的集成测试。
| async def update_session_state( | ||
| self, | ||
| session: Session, | ||
| state_delta: dict[str, Any], | ||
| ) -> None: | ||
| """Persist session-scoped state without replacing the caller's Event window.""" | ||
| if not state_delta: | ||
| return | ||
| session.state.update(state_delta) | ||
|
|
||
| async with self._redis_storage.create_db_session() as redis_session: | ||
| key = session_key(session.app_name, session.user_id, session.id) | ||
| storage_session = await self._get_session(redis_session, key) | ||
| if not storage_session: | ||
| logger.warning( | ||
| "Session %s not found in Redis while updating state", | ||
| session.id, | ||
| ) | ||
| return | ||
| storage_session.state.update(state_delta) | ||
| await self._set_session(redis_session, storage_session) | ||
|
|
||
| @override |
There was a problem hiding this comment.
问题: 新增的 update_session_state 对 Redis 会话做无锁读-改-写:先 _get_session 读出完整 Session,storage_session.state.update(state_delta) 后经 _set_session 把整个 Session JSON 写回;无 WATCH/事务/锁。调用方 SessionMemoryExtractor._persist_checkpoint 在每次 Session Memory 提取后写 state,可能发生在 post-turn worker 线程与主 loop 并发执行期间。
触发条件: Redis/内存后端 + 同一会话在主 loop 继续追加事件/更新 state 的同时,后台 post-turn worker(_defer_post_turn_processing=True 时)或另一请求并发执行 update_session_state/update_session;读到的快照与写回之间存在另一处写入时发生覆盖。
实际影响: 快照覆盖导致 state(如 _trpc_agent:summary/session memory checkpoint)或事件窗口丢失:一方刚写入的增量被另一方基于旧快照的整体写回抹掉,Session Memory checkpoint 丢失会触发重复提取或边界错位,用户数据在并发下静默丢失。
修正方向: 对 Redis 采用 WATCH/MULTI-EXEC 或 Lua 脚本化的条件比较写(如仅当 key 版本未变才写回),或至少让 update_session_state 与 append_event/update_session 走同一把分布式锁并按字段合并而非整对象 SET;并为并发写场景补充测试。
| if not result: | ||
| return None | ||
| invocation_ctx.end_invocation = True | ||
| rsp.rsp = result | ||
| rsp.is_continue = False | ||
| rsp.error = None | ||
| session_memory_extractor = summarizer_manager.get_session_memory_extractor() | ||
| if session_memory_extractor is not None: | ||
| await session_memory_extractor.extract_if_needed( | ||
| invocation_ctx, | ||
| force=False, | ||
| ) | ||
| return |
There was a problem hiding this comment.
问题: AdvancedAutoCompactSummarizerFilter._before 在压缩连续失败达到 max_failures 且命中 blocking 阈值时,把 AdvancedAutoCompactResult.blocked 对应的固定提示文本(ADVANCED_AUTOCOMPACT_BLOCKED_MESSAGE)作为 LlmResponse 塞进 rsp.rsp 并 is_continue=False 结束过滤器链,同时置 end_invocation=True。经 filter 适配器(_base_filter._handle_co 流入 run_stream_filters)该响应会作为正常模型回复从 generate_async 流出,而 end_invocation 在全部代码中没有任何读取方。
触发条件: 长会话中压缩器连续 3 次失败(如摘要 LLM 反复报错、后端存储异常)且请求接近上下文硬限制。
实际影响: 用户侧看到的是模型“正常回复”一段拒绝文本(“Automatic context compaction has failed repeatedly…”),没有错误码、异常或终止信号,应用层无从区分压缩失败与模型主动拒答;触发后仅靠该文本提示用户,无法编程化处理。
修正方向: 让 blocked 路径产出带 error_code 的错误响应(或在响应中携带结构化元数据),或在 end_invocation 处落一个真实消费者(如 runner 终止循环并上报错误),而不是把错误伪装成普通模型文本。
| if not acquired: | ||
| raise TimeoutError(f"Timed out acquiring Advanced Memory lock for {self._paths.scope.storage_key}") | ||
| try: | ||
| yield | ||
| finally: | ||
| await self._command( | ||
| "eval", | ||
| _RELEASE_LOCK_SCRIPT, | ||
| 1, | ||
| key, | ||
| token, | ||
| ) |
There was a problem hiding this comment.
问题: _memory_write_lock 的 finally 无条件执行 eval(_RELEASE_LOCK_SCRIPT) 释放锁;当 while 循环在期限内未获得锁(acquired=False)抛出 TimeoutError 时仍会执行该释放。Lua 脚本会用本进程的 token 与 key 现值比较,但若在超时瞬间锁 key 恰好过期并被另一进程重新取得,本进程的释放会失败(比较不相等——脚本返回 0),而另一种情形下释放可能命中别人刚持有的锁。
触发条件: 锁竞争激烈且某次获取超时(memory_lock_acquire_timeout_seconds 到期)、或锁 TTL 在超时与释放之间恰好翻转。
实际影响: 在极端时间窗内可能删除并发持有者(后续进程)的锁,破坏跨进程互斥,使两个进程同时进入 write_index/write_topic 临界区,导致索引与 topic 内容相互覆盖、丢失;多数超时场景是安全 no-op,故概率低。
修正方向: 仅当 acquired 为真时才执行释放(if acquired: await self._command("eval", ...)),并在 finally 中持有 acquired 状态进行判断。
| async def apply( | ||
| self, | ||
| request: LlmRequest, | ||
| *, | ||
| ctx: InvocationContext, | ||
| force: bool = False, | ||
| ) -> AdvancedAutoCompactResult: | ||
| """Run compaction against the current session's tenant namespace.""" | ||
| session_id = ctx.session_id | ||
| if self._runtime.scope: | ||
| return await self._apply_scoped(request, session_id=session_id, ctx=ctx, force=force) | ||
| runtime = self._runtime.for_session(ctx.session) | ||
| processor = self._scoped_processors.get(runtime.scope) | ||
| if processor is None: | ||
| processor = copy.copy(self) | ||
| processor._runtime = runtime | ||
| processor._states = {} | ||
| processor._session_locks = {} | ||
| self._scoped_processors[runtime.scope] = processor | ||
| return await processor.apply(request, ctx=ctx, force=force) |
There was a problem hiding this comment.
问题: _apply_scoped/apply 的 _scoped_processors 缓存与各处理器的 _states/_session_locks 字典没有任何剔除/淘汰逻辑:每个 app/user scope(for_session 产生的新 runtime scope)都会在进程内永久保留一个 processor 实例,每个 session 的 state 与 asyncio.Lock 也永久留存。
触发条件: 常驻多租户部署(服务化场景,本 PR 目标)中 app/user/session 数随时间增长。
实际影响: 进程内存随唯一 app/user/session 数量线性增长且不回收(每个 session 的 AdvancedAutoCompactState 含 latest_compaction.summary 大文本、lock 对象),长时间运行的服务内存持续膨胀,最终可能 OOM;对已删除/过期会话也会重复保留状态。
修正方向: 为缓存引入容量上限/淘汰(如按最后访问时间剔除过期 session 的状态与锁,scope 级 processor 按 LRU 限制),或在会话关闭/删除时显式清理对应键。
| content="", | ||
| expires_at=self._expiry(self._config.memory_ttl_seconds), | ||
| )) | ||
| await self._storage.commit(db) | ||
|
|
||
| async def read_index(self) -> str: | ||
| async with self._storage.create_db_session() as db: | ||
| row = await self._storage.get( | ||
| db, | ||
| SqlKey(key=(self._app_name, self._user_id), storage_cls=AdvancedMemorySqlIndex), | ||
| ) | ||
| if row is None or self._expired(row.expires_at): | ||
| return "" | ||
| await self._refresh_memory_scope(db) | ||
| content = row.content | ||
| valid_topics = await self._storage.query( | ||
| db, | ||
| SqlKey(key=(self._app_name, self._user_id), storage_cls=AdvancedMemorySqlTopic), | ||
| SqlCondition(filters=[ | ||
| AdvancedMemorySqlTopic.app_name == self._app_name, | ||
| AdvancedMemorySqlTopic.user_id == self._user_id, | ||
| AdvancedMemorySqlTopic.expires_at.is_(None) | ||
| | (AdvancedMemorySqlTopic.expires_at > self._now()), | ||
| ]), | ||
| ) | ||
| valid_filenames = {topic.topic_name for topic in valid_topics} | ||
| pruned_content = prune_memory_index(content, valid_filenames) | ||
| if pruned_content != content: | ||
| row.content = pruned_content | ||
| content = pruned_content | ||
| await self._storage.commit(db) | ||
| return limit_memory_index( | ||
| content, | ||
| max_lines=self._config.memory_index_max_lines, | ||
| max_bytes=self._config.memory_index_max_bytes, | ||
| encoding=self._config.encoding, | ||
| ) |
There was a problem hiding this comment.
问题: read_index(第 141-172 行)与 read_topic/list_topics/write_index/write_topic 每次进入都调用 _refresh_memory_scope(第 98-120 行):即使 memory_ttl_seconds is None 也先 get index 行 + 全量 query 该租户所有 topic,并逐行重置 expires_at,随后必定 commit;read_topic_frontmatter 又经 read_topic 逐 topic 触发一遍。预加载路径(AdvancedMemoryPreloader._candidates)按 topic 数放大。
触发条件: 每轮对话的每次模型请求(before_model 管道)都会执行多次 read_index/read_topic;多租户、topic 较多(几十个)且 TTL 启用时,读路径产生“1+2N”条 SQL 与多次 commit,并把连接池占用在读操作上;同时活跃租户由于每次读都把 expires_at 重置为 now+TTL,TTL 永不生效。
实际影响: 读路径 SQL 放大抬高延迟与连接池压力(高并发下明显),TTL 语义变成仅对完全冷数据生效,与配置预期(按 memory_ttl_seconds 过期删除)不符;文件/Redis 后端同样在每次读时续期。
修正方向: 将 TTL 续期收敛到写路径(写入时设置 expires_at),读路径只读不刷新,并去掉 memory_ttl_seconds is None 时无谓的 scope 刷新与 commit;如需读刷新,改为仅对真正过期的行批量延期一次,避免每次读全量重写。
AI Code Review审查结论不通过 审查范围审查 计划符合性整体结构与计划一致:memory 与上下文解耦、前后端分离(storage/runtime/toolset/filters)、BEFORE_MODEL 流式过滤器 + AFTER_TURN post-turn 双触发路径、跨事件循环协调与幂等 compaction 均已落地;旧 API( 主要风险
测试充分性新增/更新测试主要覆盖 门禁结论FAILED:存在 1 项 SEVERE(local 后端 TTL 触发整体清空用户记忆的确定性数据丢失)、3 项 MODERATE 并发正确性缺陷及 3 项 LOW 问题,建议修复后合入;至少应优先修复 SEVERE 清空逻辑并为 TTL 路径补充测试。 发现的问题严重
问题: local 文件后端以 index 文件的 触发条件: 配置 实际影响: 用户长期记忆被静默整体清空且不可恢复(文件被 修正方向: 进入清理前按各 topic 文件自身 mtime 而非 index mtime 判定过期(仅删除自身过期的文件),删除后不再返回空索引而是保留仍有效的条目;并将清理逻辑与 中等
问题: 触发条件: 同一 runner 内不同事件循环(主 loop 与 实际影响: 该 session 的 memory 提取/补提取永久停止, 修正方向: 在 中等
问题: 新增的 触发条件: 多进程/多实例部署,或本次变更新增的 post-turn worker( 实际影响: 并发 append 的 events、state patch 被静默丢弃,会话历史/记忆出现缺失且不可察觉。 修正方向: 对 state 更新与事件写入使用原子路径:以事务( 中等
问题: 触发条件: 实际影响: 并发期间新产生的会话事件被静默丢弃(delete-all 删掉)或触发唯一键冲突异常,compaction 与正常对话交错时数据完整性受损。 修正方向: 在 较低
问题: 触发条件: 熔断条件( 实际影响: 强制熔断时无法区分「成功熔断」与「失败但阻塞」,日志/调用方拿不到任何根因信息,且死代码表达误导后续维护。 修正方向: 显式改写为 较低
问题: 预算统计 触发条件: 选中的 memory 条数较多且正文接近预算上限时(截断判断 实际影响: 每次请求的 system 输入超过配置的字符预算,长上下文请求下 token 消耗超出预期,可能触发下游 token 限制或性能劣化。 修正方向: 在 较低
问题: 触发条件: 同一 session 的 实际影响: 跨 loop 并发时锁不互斥(或抛出 修正方向: 与 extractor 一致改用 |
| def _read_index_sync(self) -> str: | ||
| if _is_expired(self.index_path, self._config.memory_ttl_seconds): | ||
| for path in self._paths.memory_dir.glob("*.md"): | ||
| path.unlink(missing_ok=True) | ||
| return "" | ||
| if not self.index_path.exists(): | ||
| return "" | ||
| with self.index_path.open(encoding=self._config.encoding) as source: |
There was a problem hiding this comment.
问题: local 文件后端以 index 文件的 mtime 判定 TTL 过期,过期即删除 memory_dir 下所有 *.md(含比 index 更新、根本未过期的 topic 文件),随后返回空索引。_is_expired 仅看 index mtime,而读取路径(read_index)不会刷新它,清理任务(约 195-197 行)使用同样模式;这与 Redis/SQL 后端按 key 滑动刷新 TTL 的语义不一致。
触发条件: 配置 memory_ttl_seconds 后,任何一次 read_index 调用(LongTermMemoryContext.apply 在每个 model 请求前都会执行)只要发生在距上次 index 写入超过 TTL 之后——即使期间多次读取、甚至刚写过新的 topic 文件——都会先删除全部主题文件再返回 "";save_memory 的写 topic → read_index → 写 index 流程同样会先被清空。
实际影响: 用户长期记忆被静默整体清空且不可恢复(文件被 unlink),后续索引只剩新写入条目;属确定性数据丢失,且与 Redis/SQL 后端行为不一致,构成跨后端兼容性差异。
修正方向: 进入清理前按各 topic 文件自身 mtime 而非 index mtime 判定过期(仅删除自身过期的文件),删除后不再返回空索引而是保留仍有效的条目;并将清理逻辑与 read_index/save_memory 解耦,或改为仅在显式写路径维护 TTL;同时补充 TTL 过期场景的测试。
|
|
||
| def for_session(self, session: Session) -> AdvancedAutoCompactSummarizerRuntime: | ||
| return AdvancedAutoCompactSummarizerRuntime(config=self.config, | ||
| coordination=self.coordination, |
There was a problem hiding this comment.
问题: guard 中的锁泄漏根因位于 _coordination.py:asyncio.wait_for(lock.acquire(), timeout=...) 超时取消的是外层任务,而 acquire 内部的轮询循环在 asyncio.to_thread(self._lock.acquire, False) 成功返回后、将结果交回协程之前即获得并持有底层 threading.Lock,此时外层已被取消,acquired 保持 False,finally 不会 release,该锁永久遗留。
触发条件: 同一 runner 内不同事件循环(主 loop 与 _PostTurnWorkerThread)或用户手动触发时,两个 guard 并发作用于同一 session_key,且先持有者在 wait_timeout_seconds(默认 15s)内未释放——例如 compaction memory extractor 在锁内执行长 LLM 调用时;此后每次 guard 均 coordination-timeout。
实际影响: 该 session 的 memory 提取/补提取永久停止,extract_if_needed 持续返回非阻塞失败,advanced memory 静默失效,且无任何恢复路径。
修正方向: 在 acquire 内检测取消(轮询循环中加入 asyncio.current_task().cancelled() 检查,获锁后若已取消则立即 release() 并返回 False),或将锁获取改为不可取消的临界区并在 finally 中统一释放;补充超时/取消并发测试。
| async def update_session_state( | ||
| self, | ||
| session: Session, | ||
| state_delta: dict[str, Any], | ||
| ) -> None: | ||
| """Persist session-scoped state without replacing the caller's Event window.""" | ||
| if not state_delta: | ||
| return | ||
| session.state.update(state_delta) | ||
|
|
||
| async with self._redis_storage.create_db_session() as redis_session: | ||
| key = session_key(session.app_name, session.user_id, session.id) | ||
| storage_session = await self._get_session(redis_session, key) | ||
| if not storage_session: | ||
| logger.warning( | ||
| "Session %s not found in Redis while updating state", | ||
| session.id, | ||
| ) | ||
| return | ||
| storage_session.state.update(state_delta) | ||
| await self._set_session(redis_session, storage_session) |
There was a problem hiding this comment.
问题: 新增的 update_session_state 采用「_get_session 读快照 → storage_session.state.update → _set_session 全量回写且 _set_session 串行化整个 storage session」模式,无版本校验或原子更新;在两处读-写之间发生的并发写入会被 stale 快照整体覆盖(lost update)。
触发条件: 多进程/多实例部署,或本次变更新增的 post-turn worker(runners._PostTurnWorkerThread、defer_post_turn_processing)与主 loop 并发操作同一 session——例如 compaction 持久化(_persist_session_compaction → update_session)进行全量写入的时刻,用户主流程恰好 append 新 event。
实际影响: 并发 append 的 events、state patch 被静默丢弃,会话历史/记忆出现缺失且不可察觉。
修正方向: 对 state 更新与事件写入使用原子路径:以事务(MULTI/WATCH)包裹读-写,或对事件写入走 append-only key(如 RPUSH)而仅对 state 做字段级 patch,避免全量 JSON 覆盖;至少对 update_session/update_session_state 增加并发交错测试。
| @override | ||
| async def update_session_state( | ||
| self, | ||
| session: Session, | ||
| state_delta: dict[str, Any], | ||
| ) -> None: | ||
| """Persist session-scoped state without rewriting Event rows.""" | ||
| if not state_delta: | ||
| return | ||
| session.state.update(state_delta) | ||
|
|
||
| async with self._sql_storage.create_db_session() as sql_session: | ||
| session_key = SqlKey( | ||
| key=(session.app_name, session.user_id, session.id), | ||
| storage_cls=StorageSession, | ||
| ) | ||
| storage_session: Optional[StorageSession] = await self._sql_storage.get_for_update( | ||
| sql_session, | ||
| session_key, | ||
| ) | ||
| if storage_session is None: | ||
| logger.warning( | ||
| "Session %s not found in storage while updating state", | ||
| session.id, | ||
| ) | ||
| return | ||
|
|
||
| persisted_state = dict(storage_session.state or {}) | ||
| persisted_state.update(state_delta) | ||
| storage_session.state = persisted_state # type: ignore | ||
| await self._sql_storage.commit(sql_session) | ||
| await self._sql_storage.refresh(sql_session, storage_session) | ||
| session.last_update_time = storage_session.update_timestamp_tz | ||
|
|
There was a problem hiding this comment.
问题: update_session 使用无行锁的 get 读取 storage_session,随后 delete 该 session 全部 event 行并重新 insert 内存快照中的 events;本变更在同一文件内给 append_event(549-552 行)与新增的 update_session_state(625-658 行)均换用 get_for_update 行锁,唯独 update_session 未对齐,形成不一致的并发模型;而该路径正是本变更引入的 compaction 持久化(_auto_compact.py:449 的 ctx.session_service.update_session(ctx.session))的主写入通道。
触发条件: _persist_session_compaction(compaction 完成后)或任何调用 update_session 的路径与并发 append_event 交错:delete-all 发生在其间而未获锁时,并发 append 的行先落库、随后被 delete 删除,或 insert 与外部 append 的键冲突。
实际影响: 并发期间新产生的会话事件被静默丢弃(delete-all 删掉)或触发唯一键冲突异常,compaction 与正常对话交错时数据完整性受损。
修正方向: 在 update_session 中也使用 get_for_update(与 append_event/update_session_state 一致)将 delete+re-add 包裹在行锁事务内,或改为增量差异写入(仅插入快照中不存在的事件);补充与并发 append 交错的测试。
| request_chars_after=request_chars_before, | ||
| consecutive_failures=state.consecutive_failures, | ||
| error=str(exc) if exc else None if blocked else None, | ||
| request_tokens_before=request_tokens_before if token_mode else None, | ||
| request_tokens_after=comparison_tokens_before if token_mode else None, | ||
| token_source=token_budget_before.estimate.source if token_mode else None, |
There was a problem hiding this comment.
问题: error=str(exc) if exc else None if blocked else None 的嵌套三目实际解析为 str(exc) if exc else (None if blocked else None),两条 None 分支等价,blocked 分支是死代码;同时 658-665 行的硬熔断 early-return 路径(blocked=True, error=None)本就把异常信息丢弃。
触发条件: 熔断条件(consecutive_failures >= max_failures + blocking_reached)满足时的强制阻塞,或未触发异常时误入 blocked=True 分支;should_block 为 True 时结果中 error 恒为 None。
实际影响: 强制熔断时无法区分「成功熔断」与「失败但阻塞」,日志/调用方拿不到任何根因信息,且死代码表达误导后续维护。
修正方向: 显式改写为 error = str(exc) if exc is not None else ("blocked by fallback policy" if blocked else None),并让熔断 early-return 携带真实原因。
| def _session_lock(self, session_id: str) -> asyncio.Lock: | ||
| """Return the unique compaction lock for a session.""" | ||
| key = self._runtime.session_key(session_id) | ||
| lock = self._session_locks.get(key) | ||
| if lock is None: | ||
| lock = asyncio.Lock() | ||
| self._session_locks[key] = lock | ||
| return lock |
There was a problem hiding this comment.
问题: _session_lock 返回按 session_key 惰性创建的 asyncio.Lock,其绑定创建时的事件循环;本仓库已存在第二个事件循环(runners._PostTurnWorkerThread 专用 loop),而 compaction memory extractor 的协调特意改用 CrossLoopLock 正是为跨 loop 安全,此处却未对齐。
触发条件: 同一 session 的 AdvancedAutoCompactSummarizer 处理器在两个不同事件循环上并发执行——如同时配置多种触发方式(BEFORE_MODEL 与 AFTER_TURN 混合、或用户从自定义 loop 手动调用 create_session_summary_by_events),_micro_compact/_history_snip/_tool_result_budget 中的同类 asyncio.Lock 亦同。
实际影响: 跨 loop 并发时锁不互斥(或抛出 RuntimeError: ... bound to a different event loop),处理器并发修改共享内存状态,导致历史裁剪/工具结果预算统计不一致。
修正方向: 与 extractor 一致改用 SessionOperationCoordinator/CrossLoopLock(或使锁与 loop 无关),并为 mixed-trigger 场景补充测试;若明确不支持跨 loop 并发,应在文档与入口处显式约束。
Advance memory 服务化
基于 Session Service 的Advanced Compact
基于 tool 的 Advanced Long-term Memory
修改框架中错误引入summary至system instructions