feature: advanced memory 服务化,支持 redis/sql - #326
Conversation
12ca615 to
3efda2f
Compare
AI Code Review审查结论通过 审查范围:f05797d..12ca615(feature: advanced memory 服务化,支持 redis/sql),48 个文件、+3210/−239,含 10 个测试文件。计划符合性:Redis/SQL 双后端(配置、scoped runtime、存储层、TTL 组、cleanup、示例与测试)整体落实了服务化意图,local 路径保持兼容。主要风险集中在三处由本变更引入的实现缺口:1) SQL backend 下 tool-result 替换路径错位( 发现的问题中等
问题: 触发条件: 配置 实际影响: 传给模型的 修正方向: 让 SQL scoped runtime 也生成 中等
问题: 触发条件: 同一 session 的两个并发写入者(例如 Runner 的 post-turn 后台处理与主循环并发提交 transcript,这是本变更中已文档化的并发模型)同时带上相同 实际影响: 其中一个写入者得到 修正方向: 对 较低
问题: 触发条件: 写入耗时超过 30s 时,锁自动过期,另一节点/进程可用不同 token 抢先获取同一锁。 实际影响: 两个写入者并发进入同一 scope 的写路径,出现重复写、交叉写或后写覆盖前写;原本等待锁的一方在 修正方向: 在写入循环中定期续期锁 TTL(如每 10s 较低
问题: 触发条件: 配置 实际影响: 同一会话的元数据存本地磁盘、内容存数据库,多节点部署下元数据不可共享且与内容不一致; 修正方向: 与 redis 一样在 |
| persisted_path = ( | ||
| Path( | ||
| f"advanced-memory://{self._runtime.config.redis_key_prefix}/" | ||
| f"{self._runtime.scope.app_name}/{self._runtime.scope.user_id}/{session_id}/" | ||
| f"tool/{candidate.result_id}" | ||
| ) | ||
| if hasattr(self._runtime, "scope") and self._runtime.config.storage_backend == "redis" | ||
| else self._runtime.paths.tool_result_path(session_id, candidate.result_id) | ||
| ) |
There was a problem hiding this comment.
问题: _build_replacement 只在 storage_backend == "redis" 且 runtime 带 scope 时构造 advanced-memory://... URL,其余情况(包括 SQL backend 的 scoped runtime)一律走 self._runtime.paths.tool_result_path(...) 生成本地磁盘路径。但 SQL backend 下写入方 SqlToolResultStore.write()(_sql_stores.py:368)实际返回 advanced-memory://sql/{app}/{user}/{session}/tool/{result_id},读取出口并不存在本地文件。
触发条件: 配置 storage_backend="sql" 且启用 AdvancedMemoryToolResultBudget,运行时任何 tool result 超过上下文预算被替换时。
实际影响: 传给模型的 persisted_output.path 指向不存在的本地文件;模型/下游诊断按该路径读取或拼接时无法定位真实数据,工具结果替换功能在 SQL 后端下失效,且与 Redis 分支的行为不一致。
修正方向: 让 SQL scoped runtime 也生成 advanced-memory://sql/{app}/{user}/{session}/tool/{result_id} URL,例如将判断改为按 storage_backend 三路分支(redis/sql 分别生成各自 URL,local 走磁盘路径),或直接调用与存储层一致的路径构造逻辑。
| async with self._storage.create_db_session() as db: | ||
| dedupe_id = self._dedupe_id(session_id, unique_key, value) | ||
| seen_key = (self._app_name, self._user_id, session_id, unique_key, value) | ||
| seen = await self._storage.get( | ||
| db, | ||
| SqlKey(key=(dedupe_id,), storage_cls=SqlTranscriptSeen), | ||
| ) | ||
| if seen is not None and not self._expired(seen.expires_at): | ||
| await self._refresh_session_scope(db, session_id) | ||
| await self._storage.commit(db) | ||
| return Path(f"advanced-memory://sql/{self._app_name}/{self._user_id}/{session_id}/transcript"), False | ||
| if seen is not None: | ||
| await self._storage.delete( | ||
| db, | ||
| SqlKey(key=(dedupe_id,), storage_cls=SqlTranscriptSeen), | ||
| SqlCondition(filters=[ | ||
| SqlTranscriptSeen.dedupe_id == dedupe_id, | ||
| ]), | ||
| ) | ||
| payload.setdefault("recorded_at", self._now().replace(tzinfo=timezone.utc).isoformat()) | ||
| await self._storage.add(db, SqlTranscriptSeen( | ||
| dedupe_id=dedupe_id, | ||
| app_name=seen_key[0], | ||
| user_id=seen_key[1], | ||
| session_id=seen_key[2], | ||
| unique_key=seen_key[3], | ||
| unique_value=seen_key[4], | ||
| expires_at=self._expiry(self._config.session_ttl_seconds), | ||
| )) | ||
| await self._storage.add(db, SqlTranscript( | ||
| app_name=self._app_name, | ||
| user_id=self._user_id, | ||
| session_id=session_id, | ||
| record_id=uuid.uuid4().hex, | ||
| payload=json.dumps(payload, ensure_ascii=False), | ||
| expires_at=self._expiry(self._config.session_ttl_seconds), | ||
| )) | ||
| await self._refresh_session_scope(db, session_id) | ||
| await self._storage.commit(db) | ||
| return Path(f"advanced-memory://sql/{self._app_name}/{self._user_id}/{session_id}/transcript"), True |
There was a problem hiding this comment.
问题: SqlTranscriptStore.append_unique 是“先查 SqlTranscriptSeen 再 INSERT”的无锁 check-then-act:seen is None 与真正 INSERT 之间没有 SELECT ... FOR UPDATE(get_for_update 已存在但此处未用)、锁或重试。与 Redis 端 _APPEND_UNIQUE_SCRIPT 的原子 Lua 及本地存储的 threading.Lock 相比,缺少并发保护。
触发条件: 同一 session 的两个并发写入者(例如 Runner 的 post-turn 后台处理与主循环并发提交 transcript,这是本变更中已文档化的并发模型)同时带上相同 unique_key+value,两者都读到 seen is None 后各自 INSERT 同一 dedupe_id 主键。
实际影响: 其中一个写入者得到 IntegrityError,该 transcript 记录写入失败,可能中断 post-turn 处理流程;同时两者都会返回“已追加”语义,去重约定被破坏。
修正方向: 对 SqlTranscriptSeen 行使用 get_for_update 做行级锁,或在 IntegrityError 上捕获并回退为“已存在”返回 (path, False),与 Redis 原子语义保持一致。
| @asynccontextmanager | ||
| async def _memory_write_lock(self): | ||
| """Serialize long-term memory writes across processes and nodes.""" | ||
| token = uuid4().hex | ||
| key = self._memory_lock_key() | ||
| deadline = asyncio.get_running_loop().time() + self._config.memory_lock_acquire_timeout_seconds | ||
| acquired = False | ||
| while asyncio.get_running_loop().time() < deadline: | ||
| result = await self._command( | ||
| "set", | ||
| key, | ||
| token, | ||
| nx=True, | ||
| ex=self._config.memory_lock_ttl_seconds, | ||
| _command_expire=RedisExpire( | ||
| key=key, | ||
| ttl=Ttl(ttl_seconds=self._config.memory_lock_ttl_seconds), | ||
| ), | ||
| ) | ||
| if result is True or result in (b"OK", "OK"): | ||
| acquired = True | ||
| break | ||
| await asyncio.sleep(min(0.05, max(0.0, deadline - asyncio.get_running_loop().time()))) | ||
| 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 用 SET key token NX EX=30 获取锁,锁的 TTL 固定 30 秒且没有续期(keepalive)机制;持有锁期间的长时间写入可能超过 30 秒(例如大规模 topic 内容或慢网络下的持久化)。
触发条件: 写入耗时超过 30s 时,锁自动过期,另一节点/进程可用不同 token 抢先获取同一锁。
实际影响: 两个写入者并发进入同一 scope 的写路径,出现重复写、交叉写或后写覆盖前写;原本等待锁的一方在 memory_lock_acquire_timeout_seconds(默认 10s)内拿不到锁还会抛 TimeoutError,长写入场景下写操作失败率上升。
修正方向: 在写入循环中定期续期锁 TTL(如每 10s PEXPIRE 一次),或按写入数据量自适应放宽锁 TTL;释放时继续用 token 校验防止误删他人锁。
| if self._runtime.config.storage_backend == "redis": | ||
| raise ValueError( | ||
| "AdvancedMemorySessionService is file-backed; use RedisSessionService with " | ||
| "AdvancedMemoryService when AdvancedMemoryConfig.storage_backend='redis'" | ||
| ) |
There was a problem hiding this comment.
问题: AdvancedMemorySessionService.__init__ 只拒绝 storage_backend == "redis",SQL backend 被放行,但该 service 仍以本地文件承载 session 元数据(session.json、_state.json 位于 tenants/... 目录),transcript 与内存数据却写入 SQL 数据库,形成碎片化的状态层。
触发条件: 配置 storage_backend="sql" 并调用 AdvancedMemorySessionService.bind(...) 使用该 service 恢复/持久化会话事件。
实际影响: 同一会话的元数据存本地磁盘、内容存数据库,多节点部署下元数据不可共享且与内容不一致;emit_transcript_records 的清理 glob(*/*/{SESSION}/*/session.json)也只覆盖本地布局,数据生命周期管理分裂。
修正方向: 与 redis 一样在 __init__ 中拒绝 SQL backend(引导用户改用 SQL 专属的 SqlSessionService),或为 SQL backend 实现基于数据库的 session 元数据存储,保证状态层单一。
AI Code Review审查结论不通过 审查范围: 发现的问题严重
问题: 触发条件: 使用 redis 后端时,每次 实际影响: transcript 流(会话历史)会被固定 24h 过期(重写会话时历史丢失、 修正方向: 不要在 严重
问题: 触发条件: 配置 实际影响: 长航时进程中的长期记忆永不按 TTL 过期(已到期内容继续被注入上下文,发酵为陈旧记忆); 修正方向: 读操作不应刷新整个组。改为仅在写入时 中等
问题: 触发条件: 任何在 service 关闭后继续使用同一 runtime(或示例中两个进程共享同一 Redis)的调用,如 runner 的 实际影响: 对已 修正方向: 让 中等
问题: 触发条件: 实际影响: 无 修正方向: 校验 中等
问题: 触发条件: 索引首行(如超长 memory 条目)长度超过 25KB。 实际影响: 记忆索引在超大条目下被完全丢弃,模型无法看到任何索引条目,违背“读取受限前缀”的设计意图(本地 修正方向: 与本地实现一致,先按行加载再截断字节——例如遍历行并只跳过超限的行而保留可读前缀( 较低
问题: 触发条件: 按 README 默认配置运行( 实际影响: 示例无法演示会话 TTL 清理,残留数据长期占用 Redis;与 README 声称的“自动过期”演示不符。 修正方向: 在未配置 TTL 时保留默认 |
| await self._command("xadd", stream, {"data": json.dumps(payload)}) | ||
| await self._refresh_session_ttl(session_id, stream) | ||
| return Path(f"advanced-memory://{stream}") | ||
|
|
||
| async def append_unique(self, session_id: str, record: Mapping[str, Any], *, unique_key: str) -> tuple[Path, bool]: | ||
| payload = dict(record) | ||
| value = payload.get(unique_key) | ||
| if not isinstance(value, str) or not value: | ||
| raise ValueError(f"Transcript unique key {unique_key!r} must be a non-empty string") | ||
| payload.setdefault("recorded_at", datetime.now(timezone.utc).isoformat()) | ||
| stream = f"{self._session_base(session_id)}:transcript" | ||
| seen = f"{stream}:seen:{unique_key}" | ||
| async with self._storage.create_db_session() as connection: | ||
| added = await self._storage.execute_command( | ||
| connection, | ||
| RedisCommand( | ||
| method="eval", | ||
| args=(_APPEND_UNIQUE_SCRIPT, 2, stream, seen, value, json.dumps(payload, ensure_ascii=False)), | ||
| )) | ||
| await self._refresh_session_ttl(session_id, stream, seen) | ||
| return Path(f"advanced-memory://{stream}"), bool(added) | ||
|
|
||
| async def read_all(self, session_id: str) -> list[dict[str, Any]]: | ||
| stream = f"{self._session_base(session_id)}:transcript" | ||
| entries = await self._command("xrange", stream, "-", "+") | ||
| await self._refresh_session_ttl(session_id, stream) | ||
| records: list[dict[str, Any]] = [] |
There was a problem hiding this comment.
问题: RedisTranscriptStore 的 append/append_unique 用 XADD 写 transcript 流,execute_command 的 EXPIRE_METHOD: list[str] 与 int 返回值的 lower_method in EXPIRE_METHOD 判定永远为假(XADD 返回字节/字符串 ID),且 RedisCommand 缺省 RedisExpire() 的 Ttl 默认 enable=True/ttl_seconds=86400,于是每次 xadd 都会附带一个 24h 的固定过期时间——流会莫名过期丢失,而与 session_ttl_seconds 无关。同时 _refresh_ttl_group 的 if ttl is None: return 在未配置 TTL 时跳过全部注册,:keys 注册表与 :seen: 去重键从不刷新。
触发条件: 使用 redis 后端时,每次 transcripts.append/append_unique(每次事件写入)即触发;session_ttl_seconds 配置与否都会发生。
实际影响: transcript 流(会话历史)会被固定 24h 过期(重写会话时历史丢失、AutoCompact/HistorySnip 的状态恢复基于 transcript 而失真),未配置 TTL 时又永不清理、数据无限增长。
修正方向: 不要在 append 的 RedisCommand 中携带 RedisExpire(_command 已默认 RedisExpire()),或为流单独构建无过期命令;同时让 _refresh_ttl_group 在 ttl is None 时仍执行注册(或在 append_unique 中对 :seen: 键单独设置与流一致的 TTL)。
| def _read_topic_sync(self, path: Path) -> str | None: | ||
| if _expire_memory_dir(self._paths.memory_dir, self._config) or not path.exists(): | ||
| return None | ||
| _refresh_memory_dir(self._paths.memory_dir) | ||
| return path.read_text(encoding=self._config.encoding) | ||
|
|
||
| async def read_topic_frontmatter(self, topic_name: str) -> str | None: | ||
| """Read only the frontmatter of a detail memory topic.""" | ||
| path = self._paths.memory_topic_path(topic_name) | ||
| return await asyncio.to_thread(self._read_frontmatter, path) | ||
|
|
||
| def _read_optional_text(self, path: Path) -> str | None: | ||
| """Synchronously read an optional text file.""" | ||
| if not path.exists(): | ||
| return None | ||
| return path.read_text(encoding=self._config.encoding) | ||
| return await asyncio.to_thread(self._read_frontmatter_sync, path) | ||
|
|
||
| def _read_frontmatter(self, path: Path) -> str | None: | ||
| def _read_frontmatter_sync(self, path: Path) -> str | None: | ||
| """Synchronously read a topic's bounded frontmatter block.""" |
There was a problem hiding this comment.
问题: LongTermMemoryStore._read_index_sync/_read_topic_sync/_list_topics_sync 在每次读操作中调用 _refresh_memory_dir,后者对 memory 目录下所有 .md 文件执行 touch。MEMORY.md 的 mtime 因此跟随每次模型请求的 read_index 不断刷新,_is_expired 永远不成立,_expire_memory_dir 永不触发——TTL 到期完全失效。
触发条件: 配置 memory_ttl_seconds 后,任何保持运行的服务(每次模型调用都注入 read_index)都会持续延长 MEMORY 组生命周期。
实际影响: 长航时进程中的长期记忆永不按 TTL 过期(已到期内容继续被注入上下文,发酵为陈旧记忆);LocalAdvancedMemoryCleanup 周期任务退化为空转。
修正方向: 读操作不应刷新整个组。改为仅在写入时 _refresh_memory_dir,或让每次读只更新组内该文件的 activity,避免一旦“读即保活”导致无法到期。
| @@ -43,11 +84,165 @@ def create(cls, config: AdvancedMemoryConfig | None = None) -> "AdvancedMemoryRu | |||
| session_memory=SessionMemoryStore(resolved_config, paths), | |||
| tool_results=ToolResultStore(resolved_config, paths), | |||
| transcripts=TranscriptStore(resolved_config, paths), | |||
| _redis_storage=redis_storage, | |||
| _sql_storage=sql_storage, | |||
| _sql_cleanup=sql_cleanup, | |||
| _local_cleanup=local_cleanup, | |||
| ) | |||
|
|
|||
There was a problem hiding this comment.
问题: AdvancedMemoryService.close() 与 AdvancedMemorySessionService.close()(_backend.close())都会调用 self._runtime.close(),进而关闭共享的 RedisStorage 连接池。bind() 返回的 TranscriptSessionService 包装层 close() 只关闭委托(旧版),Runner.close() 同时关闭 session service 与 memory service,双路径重复关闭同一池。
触发条件: 任何在 service 关闭后继续使用同一 runtime(或示例中两个进程共享同一 Redis)的调用,如 runner 的 close 顺序在异步任务/后处理线程仍引用池时触发。
实际影响: 对已 disconnect() 的 AsyncConnectionPool 发起命令会抛出 AbortError,导致运行中的 transcript/会话读写在关闭阶段失败。
修正方向: 让 close() 幂等(close 后标记已关闭并短路),或在 TranscriptSessionService.close() 转发到 memory_runtime.close() 时避免重复关闭,保证单一所有权。
| def _cleanup_expired_sessions(self) -> None: | ||
| """Delete session directories idle longer than the configured TTL.""" | ||
| cutoff = time.time() - self.session_config.ttl.ttl_seconds | ||
| root = self._runtime.paths.session_root_dir | ||
| if not root.exists(): | ||
| tenants_root = self._runtime.config.root_dir / "tenants" | ||
| if not tenants_root.exists(): | ||
| return | ||
| for metadata_path in root.glob("*/session.json"): | ||
| for metadata_path in tenants_root.glob(f"*/*/{self._runtime.config.session_dir_name}/*/session.json"): | ||
| try: | ||
| if metadata_path.stat().st_mtime < cutoff: | ||
| shutil.rmtree(metadata_path.parent, ignore_errors=True) |
There was a problem hiding this comment.
问题: _cleanup_expired_sessions 与 list_sessions 将租户根目录硬编码为 root_dir/tenants,并使用 */*/SESSION/*/session.json 的 glob 扫描。app_name/user_id 经 _collision_safe_component 仅对特殊字符追加哈希,普通合法字符(如 -、.、_)直接落盘;若某一用户的 user_id 包含 /(MemoryScope.__post_init__ 的 _safe_component 只拒绝前后空白与控制字符,不拒绝 /),生成的目录组件会被替换为 _,且 tenant_root_dir 的解析使 glob 跨越目录层级。
触发条件: user_id 或 app_name 含 /(或两次不同合法值碰撞)时,list_sessions(user_id=...) 或清理任务遍历到的路径可能跨越租户边界。
实际影响: 无 user_id 的 list_sessions 可列出其它租户的会话元数据(信息泄露),或清理任务误删其它租户会话。
修正方向: 校验 user_id/app_name 拒绝 / 与路径分隔符,并改用遍历 tenants/<app>/<user>/SESSION/* 时按真实列出的目录解析(禁止 glob 越层)。
| 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=SqlMemoryIndex)) | ||
| if row is None or self._expired(row.expires_at): | ||
| return "" | ||
| await self._refresh_memory_scope(db) | ||
| await self._storage.commit(db) | ||
| content = row.content | ||
| lines, used_bytes = [], 0 | ||
| for line in content.splitlines(keepends=True)[:self._config.memory_index_max_lines]: | ||
| size = len(line.encode(self._config.encoding)) | ||
| if used_bytes + size > self._config.memory_index_max_bytes: | ||
| break | ||
| lines.append(line) | ||
| used_bytes += size |
There was a problem hiding this comment.
问题: SqlLongTermMemoryStore.read_index 的限长逻辑为 content.splitlines(keepends=True)[:max_lines] 后累加字节;当第一行(含换行)就超过 memory_index_max_bytes(默认 25000)时,used_bytes + size > max 在第一行即成立并 break,lines 为空,导致整段内容被截断为 ""。
触发条件: 索引首行(如超长 memory 条目)长度超过 25KB。
实际影响: 记忆索引在超大条目下被完全丢弃,模型无法看到任何索引条目,违背“读取受限前缀”的设计意图(本地 _read_index_sync 对同条件返回的是读取到空串而非丢弃)。
修正方向: 与本地实现一致,先按行加载再截断字节——例如遍历行并只跳过超限的行而保留可读前缀(if used_bytes + size > max: break 前先 append 首行)。
|
|
||
| def create_redis_session_service(redis_url: str) -> RedisSessionService: | ||
| """Create session storage with the Advanced Memory session TTL.""" | ||
| session_ttl = os.getenv("SESSION_TTL") | ||
| ttl_seconds = int(session_ttl) if session_ttl else 0 | ||
| return RedisSessionService( | ||
| db_url=redis_url, | ||
| is_async=True, | ||
| session_config=SessionServiceConfig(ttl=SessionServiceConfig.create_ttl_config( | ||
| enable=bool(session_ttl), | ||
| ttl_seconds=ttl_seconds, | ||
| cleanup_interval_seconds=ttl_seconds, | ||
| ), ), | ||
| ) | ||
|
|
There was a problem hiding this comment.
问题: create_redis_session_service 在 SESSION_TTL 未设置(.env 默认空)时构造 SessionServiceConfig.create_ttl_config(enable=False, ttl_seconds=0, cleanup_interval_seconds=0);need_ttl_expire() 要求 enable、ttl_seconds>0、cleanup_interval_seconds>0 全部成立,因此会话级 TTL 完全关闭,而 Advanced Memory 的 session_ttl_seconds 分支又只在 true 时启用。
触发条件: 按 README 默认配置运行(.env 中 SESSION_TTL 留空)。
实际影响: 示例无法演示会话 TTL 清理,残留数据长期占用 Redis;与 README 声称的“自动过期”演示不符。
修正方向: 在未配置 TTL 时保留默认 ttl_seconds(如 DEFAULT_TTL_SECONDS)并 enable=False,或将 README 明确标注需设置 SESSION_TTL。
Please enter the commit message for your changes. Lines starting
3efda2f to
36ee472
Compare
AI Code Review审查结论通过 ReviewCodeArtifact: Advanced Memory 服务化变更审查结果 发现的问题中等
问题: 中等
问题: 中等
问题: 较低
问题: 较低
问题: SQL 后端各表的 较低
问题: |
| persisted_path_text = str(persisted_path).replace( | ||
| "advanced-memory:/", | ||
| "advanced-memory://", | ||
| 1, | ||
| ) | ||
| replacement.replacement_response["persisted_output"]["path"] = persisted_path_text |
There was a problem hiding this comment.
问题: _persist_replacement(trpc_agent_sdk/advanced_memory/_tool_result_budget.py:297-302)对 persisted_path 无条件执行 str(path).replace("advanced-memory:/", "advanced-memory://", 1)。该替换的意图是兼容旧本地路径格式,但 RedisToolResultStore/SqlToolResultStore 返回的 Path 本身已经是 advanced-memory://{...} 双斜杠格式,替换后变成 advanced-memory:///sql/... 三重斜杠;而本地 ToolResultStore 返回真实文件路径,replace 不生效。触发条件: 使用 storage_backend="sql" 或 storage_backend="redis" 时,任一超过 tool_result_max_chars 的工具结果被持久化替换(ToolResultBudget.apply → _persist_replacement)。实际影响: 模型可见的 persisted_output.path 与 store 返回的规范路径(advanced-memory://sql/app/user/session/tool/id)不一致,变成 advanced-memory:///sql/...,且本次提交自带的新测试 test_sql_replacement_reports_sql_storage_path 断言 startswith("advanced-memory://sql/") 必然失败;该 path 还以 persisted_path 字段写入 transcript,造成跨端路径解析不一致。修正方向: 删除盲目的 replace,或仅当字符串以 advanced-memory:/(单斜杠,旧格式)开头时才执行替换;redis/sql store 直接透传其返回的 Path 文本。
| if value != value.strip() or any(character.isspace() and character not in {" "} for character in value): | ||
| raise ValueError(f"{field_name} must not contain leading/trailing or control whitespace") | ||
| if any(ord(character) < 32 or ord(character) == 127 for character in value): | ||
| raise ValueError(f"{field_name} must not contain control characters") |
There was a problem hiding this comment.
问题: _safe_component(trpc_agent_sdk/advanced_memory/_paths.py:22-23)新增了对前导/尾随空格、空白字符和控制字符的 ValueError 校验,此前这些值只会被静默替换为 _。MemoryScope.__post_init__ 在 for_scope/for_session 时对该限制生效,且 for_session 由 Advanced Memory 管线(LongTermMemoryContext.apply、TranscriptSessionService、memory tools 等)在每次模型请求时调用。触发条件: 传入的 app_name 或 user_id 含前导/尾随空格(例:" alice ")、\t 等内部空白或控制字符——其他 SessionService(RedisSessionService/SqlSessionService/内存版)对相同输入不做任何限制,可正常持久化。实际影响: 这些此前可正常工作的用户标识会从静默 sanitize 变为在会话运行中途抛 ValueError,导致带 Advanced Memory 的 Runner 调用直接失败;属于向后兼容性破坏。修正方向: 保持 sanitize(替换为 _)的旧行为,或仅在确实需要时校验并在文档中明确声明该限制,避免与其他 SessionService 的输入约束不一致。
| raise ValueError("runtime and config must describe the same Advanced Memory configuration") | ||
| self._runtime = runtime or AdvancedMemoryRuntime.create(config) | ||
| if self._runtime.config.storage_backend == "redis": | ||
| raise ValueError("AdvancedMemorySessionService is file-backed; use RedisSessionService with " | ||
| "AdvancedMemoryService when AdvancedMemoryConfig.storage_backend='redis'") |
There was a problem hiding this comment.
问题: AdvancedMemorySessionService.__init__ 只对 storage_backend == "redis" 抛出拒绝消息(trpc_agent_sdk/sessions/_advanced_memory_session_service.py:355-356),并未对 storage_backend == "sql" 做同样检查,但 README 明确说明 SQL 组合必须使用 SqlSessionService + AdvancedMemoryService。触发条件: 用户按 redis 同款写法构造 AdvancedMemorySessionService(config=AdvancedMemoryConfig(storage_backend="sql", sql_url=...))。实际影响: Session 元数据写入本地文件(tenants/.../session.json),而 Advanced Memory 数据写入 SQL 表,形成双存储、双 TTL 计时器并存的分裂状态;session 元数据的本地清理路径(_cleanup_expired_sessions 仅扫描 tenants/*/*/SESSION)与 SQL 清理任务各自独立执行,两个清理器可能对同一 session 的数据状态不一致,且用户不会收到任何配置错误提示。修正方向: 与 redis 分支对称,对 storage_backend == "sql" 也 raise,提示改用 SqlSessionService + AdvancedMemoryService(或为文件型 session 后端明确支持 SQL 存储并统一清理逻辑)。
| f"Memory directory: " | ||
| f"{runtime.paths.memory_dir if config.storage_backend == 'local' else 'Redis'}\n" | ||
| f"Index file: " | ||
| f"{runtime.paths.memory_index_path if config.storage_backend == 'local' else 'Redis memory index'}\n" | ||
| f"<index>\n{index.rstrip()}\n</index>\n" |
There was a problem hiding this comment.
问题: LongTermMemoryContext.apply 在构造模型指令时,storage_backend 非 local 一律输出字面量 Redis/Redis memory index(trpc_agent_sdk/advanced_memory/_memory_context.py:82-84),未区分 sql 后端。触发条件: AdvancedMemoryConfig(storage_backend="sql") 且 long_term_memory_injection_enabled=True,模型请求注入 memory index 时。实际影响: 模型被告知记忆存储在 "Redis",与实际 SQL 存储不符,可能误导模型对记忆的持久性、位置或可达性的判断;同时该展示文本对 redis 后端也有误导(真实键带 redis_key_prefix)。修正方向: 按后端分别输出(如 SQL/Redis/本地目录),或直接使用 store 提供的规范路径/标识。
| app_name: Mapped[str] = mapped_column(String(DEFAULT_MAX_KEY_LENGTH), primary_key=True) | ||
| user_id: Mapped[str] = mapped_column(String(DEFAULT_MAX_KEY_LENGTH), primary_key=True) |
There was a problem hiding this comment.
问题: SQL 后端各表的 app_name/user_id/session_id/record_id/result_id 主键均为 String(128)/String(64)(trpc_agent_sdk/advanced_memory/_sql_stores.py:39-98),写入前没有长度校验,而 Redis/SQL 会话服务对同样的输入不做长度限制(uuid 或用户自定义 id 均可)。触发条件: 在 MySQL 上使用超过 128 字符(或 64 字符的 record_id)的 app_name/user_id/session_id/result_id 调用 save_memory/transcripts.append/tool_results.write 等写入路径。实际影响: MySQL 对超长主键插入抛 DataError/OperationalError,写入中途失败,同事务内后续记录可能丢失;错误信息对用户不友好且难以定位到具体字段。修正方向: 写入前对键字段执行与 _paths 一致的显式校验(限制长度或拒绝超长值),或把主键列改为 Text/自增替代键(配合唯一索引)。
3f45e1f to
6cc5c29
Compare
6cc5c29 to
e955a59
Compare
AI Code Review审查结论不通过 审查范围:base f05797d..head 3f45e1f(36ee472+3f45e1f 两个提交,49 个文件,+3476/-246),核心为 advanced memory 服务化(新增 Redis/SQL 存储后端、租户化路径、TTL 清理、可选保留 transcript)及文件存储路径优化。已通读全部 SDK 变更文件、全部 11 个测试文件及 3 个示例目录,并对关键跨文件契约(工具 tool_context 注入、Runner close 流程、SQLAlchemy AsyncSession 语义、Redis TTL 组机制、协同锁键)做了交叉验证。 计划符合性:需求"advanced memory 服务化,支持 redis/sql"基本实现,功能覆盖完整(主题、索引、session memory、transcript、tool result 在三种后端间契约一致),文件路径租户化(tenants//)与"可选保留原始记录"(session_ttl_delete_transcripts 语义贯穿 local/redis/sql 三端)也已落地。 主要风险(按严重度):
测试充分性:新增 11 个测试文件覆盖了 SQL/Redis 存储基本读写、TTL 组刷新、transcript 保留语义和租户隔离,但 (a) 未覆盖 SQL 首插 + refresh 的失败路径,(b) Redis delete_session 前缀键删除用 MagicMock 未断言误删。 门禁结论:FAILED——存在一个高置信 SEVERE(SQL 会话服务运行时崩溃)和多个中等风险,建议修复后再合入。详细的 5 个可操作问题见 comments。 发现的问题严重
问题: 触发条件: SQL 后端( 实际影响: 新会话创建/读取直接抛 修正方向: 在 中等
问题: 触发条件: Redis 多租户/多会话场景下调用 实际影响: 删除一个会话时误删其他会话的全部记忆数据(transcript、session memory、tool result、dedupe 状态),数据永久丢失且不可恢复。 修正方向: 删除前先对每个候选键用 中等
问题: 本次变更把会话级互斥键统一改为 触发条件: 同一会话的 post-turn 提取(写 session memory + checkpoint)与下一次请求前的 autocompact(读 session memory + 读 transcript)交错执行;以及两个不同租户使用相同 session_id 时,autocompact 层发生锁碰撞。 实际影响: 同一会话内 session-memory 提取与基于其 checkpoint 的模型无关压缩不再被同一把锁串行化,可能出现压缩读到未落盘的旧 checkpoint、或提取盖掉正在压缩期间读取的记忆;跨租户则使本应独立的两个用户因相同 session_id 在同一锁上互相阻塞。 修正方向: 将 中等
问题: 触发条件: redis 后端启用( 实际影响: 锁键 TTL 被错误地覆盖为 24 小时,替代本应在 30 秒后自动失效的安全网;进程崩溃后其余写入者持续抛 修正方向: 为 较低
问题: 触发条件: 应用层传入长度 >128 的 user_id/session_id(或冗长 app_name)且使用 MySQL+utf8mb4(示例默认组合);同步 SQLite 不强制 VARCHAR 长度,掩盖问题。 实际影响: MySQL 下 修正方向: 在 |
| await self._sql_storage.commit(sql_session) | ||
| # Assigning a SQL expression expires the server-generated timestamp | ||
| # even when expire_on_commit=False. Refresh it before callers access | ||
| # update_time outside SQLAlchemy's async greenlet. | ||
| await self._sql_storage.refresh(sql_session, storage_session) |
There was a problem hiding this comment.
问题: SqlSessionService 在本次变更中新增了多处 commit(sql_session) 后立即调用 refresh(sql_session, ...)(_sql_session_service.py:444-445、616-617、653-654、741-745;__init__ 中 kwargs.setdefault("expire_on_commit", False) 是配套改动)。当对应 ORM 行是本次事务中新建的(create_session 首次创建 Session/_update_app_state/_update_user_state 首次创建 app/user state 行)时,commit 已把该行 flush 出事务边界,AsyncSession.refresh() 会因无法在该会话内定位对象而抛出 InvalidRequestError(SQLAlchemy async 下对已过期/已提交的新对象 refresh 直接报错)。
触发条件: SQL 后端(is_async=True,示例默认 MySQL+aiomysql),首次调用 get_session/create_session/append_event 处理尚不存在的 session 或 app/user state;create_session 几乎必现(所有新会话都走新建分支)。同步 sqlite 后端可能掩盖该问题,因此测试 test_sql_stores.py 未发现。
实际影响: 新会话创建/读取直接抛 InvalidRequestError,SQL 部署(示例默认配方 mysql+aiomysql)完全不可用;已有行场景下的重复刷新还多消耗一轮 DB 往返。
修正方向: 在 create_session/append_event/_get_app_state/_get_user_state 中仅对『本次事务从库中读到的既有行』执行 refresh,对新建行改为直接在对象上赋值服务器时间戳即可(MySQL 下可对返回 lastrowid/时间戳做一次轻量重查,或使用 RETURNING 语义从插入结果中取 update_time),避免对已提交的新对象调用 refresh。
| async def delete_session(self, session_id: str) -> None: | ||
| """Delete all Advanced Memory keys for one session.""" | ||
| session_base = self._session_base(session_id) | ||
| registry = self._session_registry(session_id) | ||
| keys: set[str] = {registry} | ||
| tracked = await self._command("smembers", registry) or [] | ||
| keys.update(value for value in (self._text(item) for item in tracked) if value) | ||
|
|
||
| cursor: Any = 0 | ||
| pattern = f"{session_base}:*" | ||
| while True: | ||
| cursor, scanned = await self._command( | ||
| "scan", | ||
| cursor, | ||
| match=pattern, | ||
| count=100, | ||
| ) | ||
| keys.update(value for value in (self._text(item) for item in scanned) if value) | ||
| if int(cursor) == 0: | ||
| break | ||
| if keys: | ||
| await self._command("delete", *keys) |
There was a problem hiding this comment.
问题: _RedisStore.delete_session(_redis_stores.py:144-164)用 scan(match=f"{session_base}:*") 收集待删除键,而 session_base 由 redis_key_prefix:{app:user:session_id} 构成;当两个会话 ID 互为字符串前缀(如 sess 与 sess1)时,较短的 session_base 前缀会匹配到较长会话的 :summary/:transcript/:tool: 等全部键,scan 收集后统一 delete,跨会话误删。
触发条件: Redis 多租户/多会话场景下调用 delete_session(session_id),且存在另一个 session_id 以前者为前缀(例如 "a" 与 "ab"),或不同 user 的 _user_base 组件恰为前缀关系。
实际影响: 删除一个会话时误删其他会话的全部记忆数据(transcript、session memory、tool result、dedupe 状态),数据永久丢失且不可恢复。
修正方向: 删除前先对每个候选键用 type/exists 校验其所属的租户-会话关系(例如只删除 registry 成员 + 精确 {session_base}: 后紧跟已知子资源名(summary/transcript/keys/tool/seen)的键),或改用 registry 精确枚举而非前缀 scan。
| key = self._runtime.session_key(session_id) if hasattr(self._runtime, "session_key") else session_id | ||
| lock = self._session_locks.get(key) |
There was a problem hiding this comment.
问题: 本次变更把会话级互斥键统一改为 self._runtime.session_key(session_id)(本行 _session_lock 与下方 _load_state 的 state_key 同时引入该约定),但 AutoCompact._latest_session_memory_record 仍用裸 session_id 取协调锁(_autocompact.py:381 的 coordination.guard(session_id, ...)),而 SessionMemoryExtractor.extract_if_needed 已改用 session_key(_session_memory.py:667-668)。两者在同一个 SessionOperationCoordinator 上使用不同的字符串键,互不互斥。
触发条件: 同一会话的 post-turn 提取(写 session memory + checkpoint)与下一次请求前的 autocompact(读 session memory + 读 transcript)交错执行;以及两个不同租户使用相同 session_id 时,autocompact 层发生锁碰撞。
实际影响: 同一会话内 session-memory 提取与基于其 checkpoint 的模型无关压缩不再被同一把锁串行化,可能出现压缩读到未落盘的旧 checkpoint、或提取盖掉正在压缩期间读取的记忆;跨租户则使本应独立的两个用户因相同 session_id 在同一锁上互相阻塞。
修正方向: 将 _autocompact.py:381 的 guard(session_id) 统一改为与本行约定一致的 session_key(self._runtime.session_key(session_id)),并补充多租户同 ID 并发测试。
| result = await self._command( | ||
| "set", | ||
| key, | ||
| token, | ||
| nx=True, | ||
| ex=self._config.memory_lock_ttl_seconds, | ||
| _command_expire=RedisExpire( | ||
| key=key, | ||
| ttl=Ttl(ttl_seconds=self._config.memory_lock_ttl_seconds), | ||
| ), | ||
| ) | ||
| if result is True or result in (b"OK", "OK"): | ||
| acquired = True |
There was a problem hiding this comment.
问题: _memory_write_lock(_redis_stores.py:71-102)在 set nx ex=memory_lock_ttl_seconds 之外又传了 _command_expire=RedisExpire(key=key, ttl=Ttl(ttl_seconds=memory_lock_ttl_seconds));Storage.execute_command(storage/_redis.py)的 EXPIRE_METHOD 分支会统一再执行一次 expire(),而 RedisExpire.ttl 是 Ttl 模型,未显式设置的 cleanup_interval_seconds 取默认 3600、enable=True,导致 need_ttl_expire() 恒真——即使 _command_expire 显式给了 30s,expire() 也会按默认 Ttl 重发一次 EXPIRE key 86400(默认 DEFAULT_TTL_SECONDS),覆盖 30s 锁 TTL。
触发条件: redis 后端启用(storage_backend="redis")且多个进程并发写长期记忆时任一进程在持锁期间崩溃或卡死。
实际影响: 锁键 TTL 被错误地覆盖为 24 小时,替代本应在 30 秒后自动失效的安全网;进程崩溃后其余写入者持续抛 TimeoutError(memory_lock_acquire_timeout_seconds=10s)最长可达 24 小时,直到人工干预。
修正方向: 为 RedisExpire 的 Ttl 显式设置 enable=False(该锁已用 EX 自管理 TTL,不需要 execute_command 的二次 expire),或在 execute_command 的 EXPIRE_METHOD 分支中跳过已带 EX 的 set(判断 kwargs 含 ex),避免任何二次 EXPIRE。
| app_name: str | ||
| user_id: str | ||
|
|
||
| def __post_init__(self) -> None: | ||
| _safe_component(self.app_name, field_name="app_name") | ||
| _safe_component(self.user_id, field_name="user_id") | ||
|
|
||
| @property | ||
| def storage_key(self) -> str: | ||
| """Return a stable process-local key for locks and caches.""" |
There was a problem hiding this comment.
问题: MemoryScope.__post_init__(_paths.py:49-51)与 _safe_component 只校验空白/控制字符,不限制长度;但 SQL 后端所有主键列均声明为 String(DEFAULT_MAX_KEY_LENGTH)(_sql_stores.py 各模型,常量值为 128)。超过 128 字节的 app_name/user_id/session_id 会作为完整原值写入主键列。
触发条件: 应用层传入长度 >128 的 user_id/session_id(或冗长 app_name)且使用 MySQL+utf8mb4(示例默认组合);同步 SQLite 不强制 VARCHAR 长度,掩盖问题。
实际影响: MySQL 下 _sql_stores 的写入/查询直接抛 DataError(键过长),SQL 后端对超长标识符完全不可用;Redis 侧键名不加长度上限(虽然 Redis 512MB 上限远不会触发,但键会异常膨胀)。
修正方向: 在 MemoryScope.__post_init__/AdvancedMemoryConfig 校验处对 app_name/user_id/session_id 增加长度上限(如 ≤64/128 字符),或在 SQL 存储层对超长值截断+哈希(与 _collision_safe_component 同样的思路)后再入库,并补充超长标识符的 SQL/Redis 测试。
AI Code Review审查结论不通过 审查范围:base_commit f05797d..head_commit 6cc5c29,共 49 个文件(+3478/−246),“advanced memory 服务化,支持 redis/sql”。已按 max 深度审查全部非生成变更文件:新增 Redis/SQL 存储后端( 发现的问题严重
问题: 触发条件: SQL 或 Redis 后端下,同一 app/user 租户在两个及以上的进程/节点上运行,且模型在相近时间先后调用 实际影响: 先写入的长期记忆条目从索引中被整体覆盖, 修正方向: 将“读索引→写索引”放到与 中等
问题: 触发条件: SQL 或 Redis 存储后端,Runner 默认 实际影响: 进程退出阶段的重复关闭异常( 修正方向: 在 中等
问题: 触发条件: 用户直接调用 实际影响: 无作用域调用下跨租户数据串扰与静默覆盖,未启用“服务化”语义时行为未变,但启用多租户后这部分入口形成旁路。 修正方向: 在无 中等
问题: Advanced Memory 的长事务一致性依赖事务内的 触发条件: 启用 实际影响: 偶发的刚写入即被删除、或多进程下过期判定漂移导致的记忆/会话数据丢失或残留,属于数据完整性隐患。 修正方向: 为清理器与写事务引入互斥或原子的“过期+删除”判定(如清理进程持有租户级锁或在删除前用 较低
问题: 触发条件: local 后端下服务正常关闭(Runner close)或 实际影响: 关闭阶段清理与正在进行的记忆写入竞争时可能删除目标文件,或在周期循环已覆盖该轮扫描后产生无意义的重复磁盘 I/O;低频率场景下影响有限。 修正方向: 让 较低
问题: 触发条件: SQL 后端部署在数据库节点与应用节点时钟不同步(NTP 偏差、容器时区差异)或数据库默认时区与 UTC 应用时区不匹配的环境。 实际影响: 记忆索引与主题被提前整体过期删除(数据丢失)或延迟过期(残留),影响长期记忆的可靠性。 修正方向: 统一以单一时钟为准(例如全部以数据库 |
| async with self._index_lock(runtime): | ||
| path = await runtime.long_term_memory.write_topic( | ||
| filename, | ||
| document, | ||
| ) | ||
| entries = _parse_index(await self._runtime.long_term_memory.read_index()) | ||
| entries = _parse_index(await runtime.long_term_memory.read_index()) | ||
| new_entry = MemoryIndexEntry( | ||
| name=name, | ||
| filename=path.name, | ||
| summary=summary, | ||
| ) | ||
| entries = [entry for entry in entries if entry.filename != new_entry.filename] | ||
| entries.insert(0, new_entry) | ||
| await self._runtime.long_term_memory.write_index(entries) | ||
| updated_at = parse_memory_updated_at(await self._runtime.long_term_memory.read_topic(filename) or "") | ||
| await runtime.long_term_memory.write_index(entries) |
There was a problem hiding this comment.
问题: AdvancedMemoryTools.save_memory 的“写主题→读索引→写索引”三步由 _index_lock(进程内 asyncio.Lock)串联(见 trpc_agent_sdk/tools/_advanced_memory_tool.py:114-127),但该锁只约束同一进程内的调用。SQL 后端中 SqlLongTermMemoryStore.write_index 虽然在事务内对 tenant 行 get_for_update 持锁(_sql_stores.py:263-266),可 read_index(_sql_stores.py:240-255)运行在单独的短事务里,位于 DB 行锁保护范围之外:保存流程是“提交 topic 事务后,另开事务读 index,再开事务持行锁写 index”。两个进程并发保存同一租户时,进程 A 与进程 B 先后读到同一旧索引,各自插入自己的条目,后提交者整体覆盖先提交者的条目,先提交者的记忆丢失。Redis 后端同样只有进程内 asyncio.Lock 与短持有期的 SET NX 锁,_redis_stores.py:194-199 的索引写与读不在同一锁作用域内,跨进程并发时索引仍会互相覆盖。
触发条件: SQL 或 Redis 后端下,同一 app/user 租户在两个及以上的进程/节点上运行,且模型在相近时间先后调用 save_memory(如多副本部署、横向扩容的 HTTP 服务)。
实际影响: 先写入的长期记忆条目从索引中被整体覆盖,MEMORY.md 索引与已落盘的 topic 文件不一致,模型按索引无法找到已保存的记忆,出现静默丢数据;这是该 PR 引入的 redis/sql 新后端在“服务化”目标下的核心正确性缺陷。
修正方向: 将“读索引→写索引”放到与 write_index 相同的 DB 事务与行锁作用域内,例如在 SqlLongTermMemoryStore 提供 read-merge-write 原子方法(SELECT ... FOR UPDATE 持有锁后读、合并、写一个事务内完成),save_memory 直接调用它;Redis 端相应地把索引读改写包进同一次分布式锁(或 Lua 脚本)作用域,不能只锁写步骤。
| def _runtime_for_context(self, tool_context: Any | None) -> Any: | ||
| """Resolve storage from the authenticated session, never tool arguments.""" | ||
| if tool_context is None: | ||
| return self._runtime | ||
| session = getattr(tool_context, "session", None) | ||
| return self._runtime.for_session(session) |
There was a problem hiding this comment.
问题: AdvancedMemoryTools._runtime_for_context(trpc_agent_sdk/tools/_advanced_memory_tool.py:74-79)在 tool_context is None 或无 session 时返回未作用域的根 runtime,save_memory/read_memory/list_memory_index 随即读/写全局 MEMORY 目录;LongTermMemoryContext.apply(_memory_context.py:35-46)也在 ctx is None 时直接用根 runtime 注入索引。关闭“服务化”后所有租户(for_scope/for_session 之外)的调用都会落到一个共享的全局命名空间,多个 app/user 的记忆互相读写、互相覆盖,与 MemoryScope 隔离设计相矛盾。
触发条件: 用户直接调用 create_advanced_memory_tools(runtime) 构造工具且不带 tool_context 调用;或任何未转交 ctx 的模型前置回调路径。
实际影响: 无作用域调用下跨租户数据串扰与静默覆盖,未启用“服务化”语义时行为未变,但启用多租户后这部分入口形成旁路。
修正方向: 在无 ctx/无 session 时抛出明确错误(要求必须携带 tool_context 且包含 session),或为未作用域调用拒绝执行并提示使用 for_scope/for_session,从根上消除全局命名空间旁路。
| async def _refresh_session_scope(self, db: Any, session_id: str) -> None: | ||
| expiry = self._expiry(self._config.session_ttl_seconds) | ||
| if expiry is None: | ||
| return | ||
| tables = ( | ||
| (SqlSessionMemory, (self._app_name, self._user_id, session_id)), | ||
| (SqlToolResult, (self._app_name, self._user_id, session_id)), | ||
| ) | ||
| if self._config.session_ttl_delete_transcripts: | ||
| tables = ( | ||
| (SqlTranscript, (self._app_name, self._user_id, session_id)), | ||
| (SqlTranscriptSeen, (self._app_name, self._user_id, session_id)), | ||
| *tables, | ||
| ) | ||
| for model, key in tables: | ||
| rows = await self._storage.query( | ||
| db, | ||
| SqlKey(key=key, storage_cls=model), | ||
| SqlCondition(filters=[ | ||
| getattr(model, "app_name") == self._app_name, | ||
| getattr(model, "user_id") == self._user_id, | ||
| getattr(model, "session_id") == session_id, | ||
| getattr(model, "expires_at").is_(None) | (getattr(model, "expires_at") > self._now()), | ||
| ]), | ||
| ) | ||
| for row in rows: | ||
| row.expires_at = expiry |
There was a problem hiding this comment.
问题: Advanced Memory 的长事务一致性依赖事务内的 _refresh_memory_scope/_refresh_session_scope 刷新 TTL(_sql_stores.py:132-180),但同步的过期清理器 SqlAdvancedMemoryCleanup 使用同一时间基准偏置:write_index/read_topic 用 self._now()(_sql_stores.py:116-124)推进过期时间,而读写线程间没有原子性保证,若清理器恰好在写入事务提交前发现某行 expires_at <= now 并删除,该行可能刚被写入就被清理(或相反:清理器读到的游标状态陈旧)。由于每个写事务都在提交前批量刷新整个 scope/session 的 TTL,清理器必须与刷新保持互斥才能正确;当前实现只依赖“刷新的行不会被清理器删”的先后时序,无锁或事务级协调。此外 _expiry(None) 返回 None、读路径 _expired(None)=False 的组合在 TTL 关闭时虽正确,但清理器在所有 TTL 均为 None 时仍会启动并按固定间隔空转扫描。
触发条件: 启用 memory_ttl_seconds/session_ttl_seconds 且 sql_cleanup_interval_seconds 较短时,清理任务与高频写入并发;或 TTL/清理间隔配置下多进程共享同一 SQL 库。
实际影响: 偶发的刚写入即被删除、或多进程下过期判定漂移导致的记忆/会话数据丢失或残留,属于数据完整性隐患。
修正方向: 为清理器与写事务引入互斥或原子的“过期+删除”判定(如清理进程持有租户级锁或在删除前用 expires_at <= now AND expires_at >= now - interval 双边界约束),并让清理间隔只在至少一个 TTL 启用时启动,避免无 TTL 时空转。
| async def close(self) -> None: | ||
| if self._task is not None: | ||
| await self.cleanup_once() | ||
| if self._stop_event is not None: | ||
| self._stop_event.set() | ||
| if self._task is not None and not self._task.done(): | ||
| self._task.cancel() | ||
| await asyncio.gather(self._task, return_exceptions=True) | ||
| self._task = None | ||
| self._stop_event = None |
There was a problem hiding this comment.
问题: LocalAdvancedMemoryCleanup.close()(trpc_agent_sdk/advanced_memory/_storage.py:490-499)在 self._task is not None 时执行 await self.cleanup_once(),_run 循环退出后再次全盘扫描 root/MEMORY、root/SESSION 与全部 tenants/*/*/ 目录(_cleanup_sync,_storage.py:450-480);LocalAdvancedMemoryCleanup.start() 又会在每次 initialize() 时对同一 root 重复启动(runtime 与每个 scope 后端共享)。这是关闭阶段的冗余清理,与后台循环的清理重叠执行,可能对仍在慢速读取的文件做 unlink/rmtree(_expire_memory_dir 直接 unlink,_expire_session_dir 直接 rmtree),且 close 后 star 重复调用可能重建已释放的 _stop_event。
触发条件: local 后端下服务正常关闭(Runner close)或 AdvancedMemoryRuntime.close() 被调用,且存在一个已启动的清理任务。
实际影响: 关闭阶段清理与正在进行的记忆写入竞争时可能删除目标文件,或在周期循环已覆盖该轮扫描后产生无意义的重复磁盘 I/O;低频率场景下影响有限。
修正方向: 让 close() 只取消任务并仅当循环尚未执行过清理时补一次 cleanup_once(),或为关闭清理结果加幂等/互斥保护,避免与循环的清理并发执行。
| async def write_index(self, entries: list[MemoryIndexEntry]) -> None: | ||
| content = "\n".join(entry.to_markdown() for entry in entries) | ||
| if content: | ||
| content += "\n" | ||
| async with self._storage.create_db_session() as db: | ||
| # Keep the tenant's lock row locked until this transaction commits. | ||
| await self._storage.get_for_update( | ||
| db, | ||
| SqlKey(key=(self._app_name, self._user_id), storage_cls=SqlMemoryIndex), | ||
| ) | ||
| key = SqlKey(key=(self._app_name, self._user_id), storage_cls=SqlMemoryIndex) | ||
| row = await self._storage.get(db, key) | ||
| if row is None: | ||
| row = SqlMemoryIndex(app_name=self._app_name, user_id=self._user_id) | ||
| await self._storage.add(db, row) | ||
| row.content = content | ||
| row.expires_at = self._expiry(self._config.memory_ttl_seconds) | ||
| await self._refresh_memory_scope(db) |
There was a problem hiding this comment.
问题: SqlLongTermMemoryStore.write_index 与 read_index 交替刷新整个租户 scope 的 expires_at(_sql_stores.py:245,274),所有记忆文件与索引作为一个整体在同一时钟推进 TTL;但 TypeDecorator=PreciseTimestamp(DateTime(fsp=6))依赖数据库时钟与 func.now(),而应用侧 _expiry 使用 datetime.now(timezone.utc) 计算(_sql_stores.py:116-124),双时钟在跨时区部署或数据库与应用时钟偏差时,expires_at <= now 的过期判定可能提前或滞后 TTL 的整个周期,导致记忆在 TTL 未到即被清理;本地后端 _is_expired(_storage.py:49-52)的 mtime 在同一进程内是自洽的,但 SQL 后端引入的是数据库与应用两个时钟源。
触发条件: SQL 后端部署在数据库节点与应用节点时钟不同步(NTP 偏差、容器时区差异)或数据库默认时区与 UTC 应用时区不匹配的环境。
实际影响: 记忆索引与主题被提前整体过期删除(数据丢失)或延迟过期(残留),影响长期记忆的可靠性。
修正方向: 统一以单一时钟为准(例如全部以数据库 expires_at = func.now() + ttl 计算,或应用侧统一 TIMESTAMP WITHOUT TIME ZONE 且强制 UTC),并避免跨时钟的比较。
AI Code Review审查结论不通过 审查范围: 发现的问题严重
问题: 触发条件: 两个不同的应用/用户(如 实际影响: 跨调用方的会话数据可被静默覆盖/读取:后创建的会话会覆盖同名目录内容, 修正方向: 恢复所有权检查: 中等
问题: 触发条件: 同 实际影响: 调用方通过 修正方向: 恢复 中等
问题: 触发条件: 同一 实际影响: 相同事件被重复追加到 transcript, 修正方向: 删除过期 seen 时不只删 seen,应将同 scope 的过期 transcript 行一并清理(或在 delete seen 的条件中同时覆盖对应 中等
问题: 触发条件: 多进程/多节点部署下同一 topic 或 index 并发 实际影响: 修正方向: 使用有锁行的固定哨兵行(例如按 scope 预先 upsert 一个专用 lock 行再对其 较低
问题: 触发条件: 实际影响: 单次请求的 Redis RTT 与组内 key 数线性增长,长会话下刷新耗时和负载明显放大; 修正方向: 改为多种小命令组合(如 |
| await self._read_session(app_name, user_id, resolved_id) | ||
| global_state = await self._read_global_state(app_name, user_id) | ||
| global_state["app"].update(state_delta.app_state_delta) | ||
| global_state["user"].update(state_delta.user_state_delta) | ||
| await self._write_global_state(app_name, user_id, global_state) | ||
| await self._write_session(session) |
There was a problem hiding this comment.
问题: create_session 中 await self._read_session(app_name, user_id, resolved_id) 的返回值被直接丢弃,base 版本对已有会话的所有权校验(existing.app_name != app_name or existing.user_id != user_id 时 raise ValueError)在新的 scoped 路径下被一并移除;新 get_session 也只检查 session is None,不再校验 app_name/user_id 归属。
触发条件: 两个不同的应用/用户(如 app-A/user-1 与 app-B/user-2)在同一 AdvancedMemoryRuntime 根下先后用相同的显式 session_id 调用 create_session,或 get_session 携带与存储记录不匹配的 app_name/user_id。_metadata_path 返回 tenants/<app>/<user>/sessions/<id>/session.json,不同租户的路径互不重叠,旧守卫正是为防止路径同名误写而存在。
实际影响: 跨调用方的会话数据可被静默覆盖/读取:后创建的会话会覆盖同名目录内容,get_session 对错误租户组合也返回会话数据,破坏多租户隔离,且此变更前该守卫是明确的公开行为。
修正方向: 恢复所有权检查:existing = await self._read_session(app_name, user_id, resolved_id),若 existing is not None and (existing.app_name != app_name or existing.user_id != user_id) 则 raise ValueError;同时在 get_session 中恢复 session.app_name != app_name or session.user_id != user_id 时返回 None 的守卫。
| global_state["app"].update(state_delta.app_state_delta) | ||
| global_state["user"].update(state_delta.user_state_delta) | ||
| await self._write_global_state(app_name, user_id, global_state) |
There was a problem hiding this comment.
问题: create_session 从「按 key 增量合并」回归为「整块覆盖」:global_state["app"].update(...) 与 global_state["user"].update(...) 直接将读入的整个 scope 字典原地更新并写回(_write_global_state 每次都把指针指向的同一字典整体序列化覆盖文件)。base 版本用 setdefault(app_name, {}).update(...) 保留同 scope 内其他键。
触发条件: 同 app_name/user_id 下:1) 先 create_session 写入 app:xxx/user:yyy 状态键,再创建第二个会话并传入新状态键(或调用方并发创建会话),前一键被覆盖删除;2) append_event 的 state_delta 路径沿用相同 read-modify-write 覆盖语义。
实际影响: 调用方通过 state 持久化的应用级/用户级状态键静默丢失,跨会话持久状态不再累积,属于数据完整性回归。
修正方向: 恢复 setdefault 合并语义(global_state["app"].setdefault(...).update(...))或先与已存在键显式合并再写回,避免整字典覆盖。
| if not isinstance(value, str) or not value: | ||
| raise ValueError(f"Transcript unique key {unique_key!r} must be a non-empty string") | ||
| async with self._storage.create_db_session() as db: | ||
| dedupe_id = self._dedupe_id(session_id, unique_key, value) | ||
| seen_key = (self._app_name, self._user_id, session_id, unique_key, value) | ||
| seen = await self._storage.get( | ||
| db, | ||
| SqlKey(key=(dedupe_id, ), storage_cls=SqlTranscriptSeen), | ||
| ) | ||
| if seen is not None and not self._expired(seen.expires_at): | ||
| await self._refresh_session_scope(db, session_id) | ||
| await self._storage.commit(db) | ||
| return Path(f"advanced-memory://sql/{self._app_name}/{self._user_id}/{session_id}/transcript"), False | ||
| if seen is not None: | ||
| await self._storage.delete( | ||
| db, | ||
| SqlKey(key=(dedupe_id, ), storage_cls=SqlTranscriptSeen), | ||
| SqlCondition(filters=[ | ||
| SqlTranscriptSeen.dedupe_id == dedupe_id, | ||
| ]), |
There was a problem hiding this comment.
问题: SqlTranscriptStore.append_unique 在 seen 行已过期(超过 session_ttl)时删除该 seen 记录,但仍向 SqlTranscript 追加新的 transcript 行。seen 与 transcript 的过期策略在 session_ttl_delete_transcripts=True 下使用同一 TTL,两者先于 SQL 清理任务(sql_cleanup_interval_seconds 默认 60 秒)分别独立过期。
触发条件: 同一 session_id 跨 TTL 周期继续追加相同 unique_key(如事件 event_id)的记录:过期后的 seen 行被删、旧 transcript 行仍保留,append_unique 按全新事件放行并再写一条。
实际影响: 相同事件被重复追加到 transcript,read_all/_restore_events 恢复出重复事件,破坏去重保证,与 find_last_event_id 的父链拼接和清理测试意图矛盾。
修正方向: 删除过期 seen 时不只删 seen,应将同 scope 的过期 transcript 行一并清理(或在 delete seen 的条件中同时覆盖对应 SqlTranscript 行),使去重状态与数据状态一致。
| async def write_topic(self, topic_name: str, document: MemoryDocument) -> Path: | ||
| name = self._paths.memory_topic_path(topic_name).name | ||
| async with self._storage.create_db_session() as db: | ||
| # Serialize all long-term writes for this app/user scope. | ||
| await self._storage.get_for_update( | ||
| db, | ||
| SqlKey(key=(self._app_name, self._user_id), storage_cls=SqlMemoryIndex), | ||
| ) | ||
| key = self._topic_key(name) | ||
| row = await self._storage.get(db, SqlKey(key=key, storage_cls=SqlMemoryTopic)) | ||
| if row is None: |
There was a problem hiding this comment.
问题: write_topic/write_index 的并发串行化依赖 get_for_update(SqlMemoryIndex),但 SqlMemoryIndex 行由本事务首次 add 创建——get_for_update 只对已提交/已存在的行加行锁,首写场景锁不到任何行;同时 SQLite(示例默认 sqlite:///advanced-memory.db,_set_sqlite_pragma 无锁支持)对 FOR UPDATE 静默忽略。两事务对同一 key 读未命中后分别 insert,commit 时触发主键冲突(IntegrityError 由 commit 的 rollback+重新 raise 抛出)或丢失更新。
触发条件: 多进程/多节点部署下同一 topic 或 index 并发 save_memory,或分布式锁(AdvancedMemoryTools._index_lock)仅覆盖进程内 asyncio 锁时跨进程并发。
实际影响: save_memory 报告未捕获的 IntegrityError(工具调用失败),或索引写入相互覆盖导致 MEMORY.md 与 topic 文件不一致。
修正方向: 使用有锁行的固定哨兵行(例如按 scope 预先 upsert 一个专用 lock 行再对其 get_for_update),或对 (app_name, user_id) 采用数据库层唯一约束+upsert,替代「先锁再查再插」的 TOCTOU 模式。
| async def _refresh_ttl_group( | ||
| self, | ||
| registry: str, | ||
| keys: list[str], | ||
| ttl: int | None, | ||
| skip_prefixes: tuple[str, ...] = (), | ||
| ) -> None: | ||
| """Track and refresh every key in one logical memory group.""" | ||
| if ttl is None: | ||
| return | ||
| if keys: | ||
| await self._command("sadd", registry, *keys) | ||
| tracked = await self._command("smembers", registry) or [] | ||
| tracked_keys = {self._text(value) for value in tracked} | ||
| tracked_keys.update(keys) | ||
| for key in tracked_keys: | ||
| if key and not key.startswith(skip_prefixes): | ||
| await self._command("expire", key, ttl) |
There was a problem hiding this comment.
问题: _refresh_ttl_group 对每次请求(读/写 transcript、session memory、memory index 等)都执行 SADD + SMEMBERS(全量拉取组内所有 key)+ 对每个 tracked key 逐个发 EXPIRE 命令,即 O(N) 次 Redis 往返;且 skip 判断用 startswith(前缀),可能命中前缀未被实际跟踪的 key。
触发条件: session_ttl_seconds/memory_ttl_seconds 已配置的 redis 后端,每个模型请求(append_event、read_index、transcript restore 等)触发 TTL 组刷新,组内 key 数量随会话累计增长。
实际影响: 单次请求的 Redis RTT 与组内 key 数线性增长,长会话下刷新耗时和负载明显放大;startswith 前缀匹配也可能误跳过本应刷新的 key(如 :transcript 子键 :transcript:seen:...,当前实现恰好依赖该前缀相同的 skip 语义,语义脆弱)。
修正方向: 改为多种小命令组合(如 PEXPIRE 批量/每条命令带 EX 在写时一并设置),或利用 Redis 7+ EXPIRE 批量命令/keys 精度限制,将 TTL 刷新收敛为每写操作一次;skip 判断改为精确的 key 集合比较而非前缀匹配。
|
|
||
| def create_session_service() -> AdvancedMemorySessionService: | ||
| """Create the persistent Advanced Memory session service.""" | ||
| memory_ttl = os.getenv("M_TTL") |
|
|
||
| # Canonical Session-oriented names. The implementation classes retain their | ||
| # historical internal names until the long-term-memory runtime is fully split. | ||
| SessionCompactConfig = AdvancedMemoryConfig |
There was a problem hiding this comment.
这个是当前这类压缩使用的配置,就直接取名AdvancedCompactConfig
| from trpc_agent_sdk.context import InvocationContext | ||
| from trpc_agent_sdk.sessions import Session | ||
|
|
||
| from ._runtime import AdvancedMemoryRuntime |
There was a problem hiding this comment.
TYPE_CHECKING是为解决循环依赖的,这里的
from ._runtime import AdvancedMemoryRuntime
from ._session_memory import SessionMemoryExtractor
存在循环依赖吗?
| from ._session_memory import SessionMemoryExtractor | ||
|
|
||
|
|
||
| class SessionCompactManager: |
There was a problem hiding this comment.
新建一个文件 _base_manager.py 里面抽象基类
BaseSessionCompactManager
然后SessionCompactManager换成AdvancedSessionCompactManager来继承这个基类
后续我们之前的默认的压缩也可以继承这个基类,然后使用的地方都直接传入这个基类即可
AI Code Review审查结论不通过 本审查覆盖 base 发现的问题严重
问题: 新增的 AutoCompact 流水线在每次压缩成功后调用 触发条件: 使用 实际影响: 1) 整个会话生命周期内事件表写放大呈 O(n²)(每次压缩重写全部历史行含长文本负载),长会话下 SQL 写吞吐与响应延迟显著劣化,多个 writer 并发时行删除/重插互相干扰;2) 修正方向: 为压缩持久化提供专用原子路径:像 Redis 的 |
| changed = compact_events( | ||
| summary_event, | ||
| boundary_event_id, | ||
| compaction_id=compaction_id, | ||
| ) | ||
| if changed: | ||
| await ctx.session_service.update_session(ctx.session) |
There was a problem hiding this comment.
问题: 新增的 AutoCompact 流水线在每次压缩成功后调用 _persist_session_compaction → ctx.session_service.update_session(ctx.session)(_autocompact.py:588),而 SQL 后端的 update_session(_sql_session_service.py:647-677)实现为:先 delete 该 session 的全部事件行,再逐条 add 重插 session.events 与 historical_events。该实现未被本次变更修正:本次仅在 append_event 引入 get_for_update 行锁,update_session 仍是无锁 get + 全量删除重插;而本次变更让它在每次压缩时反复触发。
触发条件: 使用 SqlSessionService + session_compact_config 且会话持续增长的场景:每次 AutoCompact 边界(字符/token 阈值)触发一次全量重写;historical_events 随会话单调增长,因此重写规模随历史累积放大;若同一会话存在并发请求(服务器/A2A/ag_ui 场景),另一请求的 append_event 数据可能落在本路径的 delete 与重插之间。
实际影响: 1) 整个会话生命周期内事件表写放大呈 O(n²)(每次压缩重写全部历史行含长文本负载),长会话下 SQL 写吞吐与响应延迟显著劣化,多个 writer 并发时行删除/重插互相干扰;2) update_session 无行锁的删除窗口与并发的 append_event(带 get_for_update)提交存在竞态,可能静默丢失刚追加的事件。
修正方向: 为压缩持久化提供专用原子路径:像 Redis 的 patch_session_state 一样新增只持久化 compact_events 结果的 SQL 方法(在 get_for_update 行锁内重写该次压缩影响的队列即可),或在 update_session 内先 get_for_update 再基于存储版本合并事件行,避免全量 delete+重插;并为 SQL 后端补充压缩后重读校验与并发 append 的测试。
b5d5108 to
2219dbf
Compare
AI Code Review审查结论不通过 审查范围: 发现的问题严重
问题: 新增的 触发条件: (1)同一 session 的 post-turn autocompact(写入 compact 结果)与其后/同时进行的另一轮 run 的 实际影响: 后提交方会整组删除并重插事件行,把另一路已持久化的 Event 从存储中抹除,同时用旧内存 修正方向: 让 严重
问题: 新增的 触发条件: 一次 run 内 autocompact 成功后紧随其后的 session-memory 提取提交 checkpoint( 实际影响: checkpoint(含 修正方向: 将 state 写入与事件写入统一到同一行锁/同一事务内完成( 中等
问题: 触发条件: (1)post-turn 延迟模式下 worker 线程与主循环同时为同一 session 调用 实际影响: 该 session 的所有 session-memory 提取及依赖它的 autocompact( 修正方向: 为 中等
问题: 触发条件: 任意一次 实际影响: 该 session 的指定事件/记录(transcript event、session-memory checkpoint、tool result)永久缺失,后续 修正方向: 先 较低
问题: 触发条件: 同一 Runner/agent 上同时配置 实际影响: 启动时 修正方向: 在两个 setup 路径间共享同一个 runtime 实例(如以 较低
问题: 新增的 触发条件: (1) 实际影响: 轻则 修正方向: 为 |
| "session_compaction_source": record.source, | ||
| "session_compaction_boundary_signature": record.boundary_signature, | ||
| "session_compaction_boundary_occurrence": record.boundary_occurrence, | ||
| }, | ||
| ) | ||
| active_before = list(ctx.session.events) | ||
| historical_before = list(ctx.session.historical_events) | ||
| last_update_before = ctx.session.last_update_time | ||
| try: | ||
| changed = compact_events( | ||
| summary_event, | ||
| boundary_event_id, | ||
| compaction_id=compaction_id, | ||
| ) | ||
| if changed: | ||
| await ctx.session_service.update_session(ctx.session) | ||
| except Exception: | ||
| ctx.session.events = active_before |
There was a problem hiding this comment.
问题: 新增的 _persist_session_compaction(本提交新增,autocompact 成功路径)通过 ctx.session_service.update_session(ctx.session) 持久化 compact_events 的结果,而 SQL 后端的 update_session 是无行锁的 delete-all-then-insert,且用 storage_session.state = session.state 整体覆盖 state。该调用与 append_event(get_for_update 行锁)及新引入的 patch_session_state 之间没有任何互斥,也没有会话级串行化,使这个在单进程内不存在的竞态在本提交内变为可稳定触发:跨副本部署时完全无锁。
触发条件: (1)同一 session 的 post-turn autocompact(写入 compact 结果)与其后/同时进行的另一轮 run 的 append_event 或另一 autocompact 并发执行;(2)两个 worker/副本并发处理同一 session,一方持有的 session 对象 events 列表不包含对方刚持久化的 event,另一方在 update_session 落盘。
实际影响: 后提交方会整组删除并重插事件行,把另一路已持久化的 Event 从存储中抹除,同时用旧内存 session.state 覆盖新的 state 写入(含 SESSION_MEMORY_STATE_KEY),造成事件丢失、state 回退、记忆/摘要与存储不一致,属于数据损坏。
修正方向: 让 update_session 与 append_event/patch_session_state 使用同一行级互斥(get_for_update)并在锁内基于当前行做增量更新而非 delete-all-insert;同时在 _persist_session_compaction 中对 update_session 与 append_event 应用同一会话级串行化;若需保留 delete-all-insert 语义,至少在删除/插入前校验行版本。
| @override | ||
| async def patch_session_state( | ||
| self, | ||
| session: Session, | ||
| state_delta: dict[str, Any], | ||
| ) -> None: | ||
| """Merge state under a row lock without touching persisted Events.""" | ||
| key = SqlKey( | ||
| key=(session.app_name, session.user_id, session.id), | ||
| storage_cls=StorageSession, | ||
| ) | ||
| async with self._sql_storage.create_db_session() as sql_session: | ||
| storage_session: Optional[StorageSession] = (await self._sql_storage.get_for_update(sql_session, key)) | ||
| if storage_session is None: | ||
| raise ValueError(f"Session {session.id} was not found") | ||
| merged_state = dict(storage_session.state or {}) | ||
| merged_state.update(state_delta) | ||
| storage_session.state = merged_state # type: ignore | ||
| await self._sql_storage.commit(sql_session) | ||
| await self._sql_storage.refresh(sql_session, storage_session) | ||
| session.state.update(state_delta) | ||
| session.last_update_time = storage_session.update_timestamp_tz | ||
|
|
There was a problem hiding this comment.
问题: 新增的 patch_session_state 在 SQL 后端用 get_for_update 行锁内合并 state,能保护 append_event 同键写入;但同一 run 内 _persist_checkpoint(patch_session_state)与 _persist_session_compaction(update_session)先后执行,两者锁的种类不同(行锁 vs 无锁 get),无法互相串行化;若 update_session 在 patch_session_state 之后落盘,会用其持有的旧 session.state 整体覆盖新写入的 checkpoint,且两操作各提交各的事务,交错时行锁只保护了 state 列而不保护事件删除/重插。
触发条件: 一次 run 内 autocompact 成功后紧随其后的 session-memory 提取提交 checkpoint(extract_if_needed(force=True)),与另一 run 的 update_session 并发落盘;或 update_session 与 patch_session_state 在同一 session 上交错提交。
实际影响: checkpoint(含 SESSION_MEMORY_STATE_KEY)或 compaction 结果可能被覆盖丢失,导致下次提取从过期 checkpoint 重放整段历史,或 summary 锚点与存储事件不一致;轻则记忆重复/陈旧,重则数据损坏。
修正方向: 将 state 写入与事件写入统一到同一行锁/同一事务内完成(update_session 改为读当前行的 state 合并后更新,而不是整对象覆盖),或为 update_session 增加行版本校验后重试;至少保证一次 run 内先 update_session 再 patch_session_state,杜绝覆盖顺序反转。
| async with self._runtime.coordination.guard(session_key) as acquired: | ||
| if not acquired: | ||
| return SessionMemoryExtractionResult(False, "coordination-timeout") | ||
| records = await self._runtime.transcripts.read_all(session.id) | ||
| checkpoint = self._last_checkpoint(records) | ||
| if self.uses_session_state: |
There was a problem hiding this comment.
问题: extract_if_needed 的 guard(session_key) 未设置超时(timeout=None),CrossLoopLock.acquire 以 10ms 轮询等待;而 autocompact 的 _apply_scoped 持有的是另一把锁(_session_lock),两把锁不能互斥;若临界区内 _persist_checkpoint 抛异常或协程被 CancelledError 打断而锁未释放(此处无 try/finally 保证),后续所有 extract_if_needed 会永久阻塞,post-turn 主循环与其共享该 coordinator 锁。
触发条件: (1)post-turn 延迟模式下 worker 线程与主循环同时为同一 session 调用 extract_if_needed;(2)持锁协程因取消或异常未释放锁;(3)guard 无超时,等待方无限轮询。
实际影响: 该 session 的所有 session-memory 提取及依赖它的 autocompact(_apply_scoped 会调用 extract_if_needed)永久挂起,post-turn 队列堆积,run 卡死。
修正方向: 为 extract_if_needed 的 guard 提供与 autocompact 一致的 session_memory_wait_timeout_seconds 超时;为临界区加 try/finally 确保锁必释放;并让 _apply_scoped 内对 extract_if_needed 的调用与 extractor 使用同一把锁(或复用 _session_lock),避免两套锁互相死锁。
| from ._paths import AdvancedMemoryPaths | ||
|
|
||
| _APPEND_UNIQUE_SCRIPT = """ | ||
| if redis.call('SADD', KEYS[2], ARGV[1]) == 0 then return 0 end | ||
| redis.call('XADD', KEYS[1], '*', 'data', ARGV[2]) | ||
| return 1 | ||
| """ |
There was a problem hiding this comment.
问题: _APPEND_UNIQUE_SCRIPT 的 SADD 与 XADD 两步之间没有原子性:SADD 成功、XADD 因异常返回而中止时,seen 键已包含该 unique_key 值,后续同一 dedupe 值再调 append_unique 会直接命中 SADD == 0 返回 False,记录被永久丢弃且无日志;脚本对 XADD 的失败语义没有任何恢复处理。
触发条件: 任意一次 append_unique 在 SADD 之后 XADD 失败(流被误删、内存超限、连接瞬断、脚本执行超时)后重试相同记录。
实际影响: 该 session 的指定事件/记录(transcript event、session-memory checkpoint、tool result)永久缺失,后续 _event_records_after_checkpoint 可能因缺记录跳过后续事件,session memory 与 autocompact 状态静默不一致,属于记录级数据丢失。
修正方向: 先 XADD 再 SADD、失败时回滚已插入的流条目(XDEL),或在 SADD 返回 0 时校验流中是否确有该记录、缺失则重建 seen;至少为失败路径记录日志并允许重试。
| compact_config = getattr(session_service, "session_compact_config", None) | ||
| from trpc_agent_sdk.sessions.compact import BaseSessionCompactConfig | ||
| if isinstance(compact_config, BaseSessionCompactConfig): | ||
| compact_config.setup(agent, session_service) |
There was a problem hiding this comment.
问题: Runner.__init__ 新增的 config-based setup 路径仅当 session_service.session_compact_config 存在时调用 compact_config.setup(agent, session_service),该 setup 会 AdvancedMemoryRuntime.create 一个新的 runtime 并安装全部 before-model 回调;而 AdvancedMemoryService.bind(另一条 setup 路径)也会创建独立 runtime。同一 Runner 上同时使用两者,或在 session_compact_config 之外再调用 setup_advanced_session_compact/setup_context_compression,会实例化两个互不共享锁与存储的 runtime,install_staged_callback 会因 runtime 不同抛 ValueError 或在幂等保护下静默失败。
触发条件: 同一 Runner/agent 上同时配置 session_compact_config 与 AdvancedMemoryService(AdvancedMemoryService(config=...) 构造尤其容易误触发,bind 只检查自身 integration,不感知 config-based setup)。
实际影响: 启动时 ValueError 中断;或双 runtime 并发写同一 session,各自持有独立的 CrossLoopLock/store,锁不互通导致 transcript 缺失/重复、记忆写入竞争。
修正方向: 在两个 setup 路径间共享同一个 runtime 实例(如以 session_service 为键缓存 AdvancedMemoryRuntime.create 结果,或在 bind 中检测到已由 config-based setup 安装时直接复用已建 runtime/回调)。
| async def patch_session_state( | ||
| self, | ||
| session: Session, | ||
| state_delta: dict[str, Any], | ||
| ) -> None: | ||
| """Atomically merge state while preserving concurrently written Events.""" | ||
| script = """ | ||
| local raw = redis.call('GET', KEYS[1]) | ||
| if not raw then | ||
| return false | ||
| end | ||
| local value = cjson.decode(raw) | ||
| local delta = cjson.decode(ARGV[1]) | ||
| if not value.state then | ||
| value.state = {} | ||
| end | ||
| for key, item in pairs(delta) do | ||
| value.state[key] = item | ||
| end | ||
| if type(value.events) == 'table' and next(value.events) == nil then | ||
| value.events = cjson.empty_array | ||
| end | ||
| if type(value.historical_events) == 'table' and next(value.historical_events) == nil then | ||
| value.historical_events = cjson.empty_array | ||
| end | ||
| if type(value.historicalEvents) == 'table' and next(value.historicalEvents) == nil then | ||
| value.historicalEvents = cjson.empty_array | ||
| end | ||
| local timestamp = tonumber(ARGV[2]) | ||
| if value.last_update_time ~= nil then | ||
| value.last_update_time = timestamp | ||
| end | ||
| if value.lastUpdateTime ~= nil then | ||
| value.lastUpdateTime = timestamp | ||
| end | ||
| local encoded = cjson.encode(value) | ||
| local ttl = tonumber(ARGV[3]) | ||
| if ttl > 0 then | ||
| redis.call('SET', KEYS[1], encoded, 'EX', ttl) | ||
| else | ||
| redis.call('SET', KEYS[1], encoded) | ||
| end | ||
| return encoded | ||
| """ | ||
| timestamp = time.time() | ||
| ttl = (int(self._session_config.ttl.ttl_seconds) if self._session_config.ttl.need_ttl_expire() else 0) | ||
| key = session_key(session.app_name, session.user_id, session.id) | ||
| async with self._redis_storage.create_db_session() as redis_session: | ||
| result = await self._redis_storage.execute_command( | ||
| redis_session, | ||
| RedisCommand( | ||
| method="eval", | ||
| args=( | ||
| script, | ||
| 1, | ||
| key, | ||
| json.dumps(state_delta, default=str), | ||
| timestamp, | ||
| ttl, | ||
| ), | ||
| ), | ||
| ) | ||
| if not result: | ||
| raise ValueError(f"Session {session.id} was not found") | ||
| stored_session = _session_from_storage_json(result) | ||
| session.state.update(state_delta) | ||
| session.last_update_time = stored_session.last_update_time | ||
|
|
There was a problem hiding this comment.
问题: 新增的 patch_session_state(Redis)用 Lua 内 GET 读取当前值,session 键不存在时脚本返回 false,Python 侧 if not result: raise ValueError 直接抛错;且该脚本原子地合并 state 并整体 SET,与 update_session 的 _get_session + _set_session(JSON 整体覆盖)之间没有重试或版本校验,交错时后写方会用旧内存快照覆盖新 state。
触发条件: (1)extract_if_needed 的 checkpoint 经 patch_session_state 提交与另一轮 run 的 update_session 并发;(2)session 恰好在 patch_session_state 前被 delete_session 删除。
实际影响: 轻则 ValueError 中断 post-turn 处理(_run_post_turn_processing 捕获后状态写入静默丢失),重则 checkpoint/state 被覆盖回退;Redis 端没有重试窗口,多副本并发写也未被互斥。
修正方向: 为 patch_session_state 的失败结果区分「键不存在」与「键存在但写失败」并提供一次重试;将 state 与事件写入统一为同一 Lua 脚本/同一事务(或对 update_session 增加乐观版本号校验),避免 update_session 用旧内存快照覆盖 patch_session_state 的原子合并结果。
feature: advanced memory 服务化,支持 redis/sql
多用户实现
组件支持 redis 存储
组件支持 sql 存储
组件支持简单自定义记忆