同一 session 的 transcript 在高并发 mirror 写入下,append 必须串行化才能保住 entry 顺序;而 foldSessionSummary 则设计为无副作用的纯函数,把摘要折叠这一计算密集型步骤从并发锁中剥离出去。两者一守一放、各司其职,是 SessionStore 在保证顺序的同时仍能横向扩展的关键。下文拆解这套协作机制,并给出可落地的并发模型与最小代码示例。
一、为什么 append 要序列化
Agent SDK 在长会话里会持续把 transcript_mirror 帧写入外部存储。即便本地磁盘已是事实源,mirror 这一副本仍要被多份消费者(resume、列出 subagent、审计)同时读取。一旦同一 session 的 append 出现并发交错,恢复出来的 transcript 就会丢掉消息时序,进而破坏工具调用与结果的对应关系。
官方契约因此写得相当直接:实现自定义 SessionStore 时,append 必须对同一 SessionKey 原子。InMemorySessionStore 的实现就是用 ConcurrentHashMap.compute 持有分段锁完成”读—改—写”,避免 lost update。
| 存储后端 | 原子写方式 | 注意事项 |
|---|---|---|
| 内存 | ConcurrentHashMap.compute | 进程重启即丢 |
| Redis | MULTI 事务 + 哈希标签 | Cluster 下需 {session} 标签 |
| S3/对象存储 | 顺序 part 文件 + List + Merge | 客户端时钟偏差会带来排序问题 |
| PostgreSQL | jsonb 列 + BIGSERIAL | jsonb 会重排 key,仅依赖序号即可 |
二、为什么 foldSessionSummary 必须是纯函数
把”摘要更新”从锁里剥出来是关键工程取舍。摘要(SummaryEntry)的维护是计算密集型:要把追加的 entries 与上一份 SummaryEntry 合并、统计 token、提取 metadata。这些操作如果也持锁,会让 append 的尾延迟被摘要拖累。
FoldSessionSummary(prev, key, entries) 的契约因此被设计为:
- 不读存储、不写存储;
MTime不由 fold 设置,由 adapter 在落盘后用存储时钟单独打戳;Data是不透明字典,adapter 原样透传;- 对
key.Subpath != ""不调用,避免 subagent 写入污染主会话摘要。
折叠结果出来后再由 append 在锁内一次性写回主存储 + 摘要存储。锁里只做”读最新 prev → fold → 写”,计算与 I/O 被拆分到锁外。
三、协作流程:把”加锁段”压到最小
import asyncio, json, pathlib
from typing import Any
class FileSessionStore:
"""最小可运行的 SessionStore 演示:单文件 + 摘要侧车"""
def __init__(self, root: pathlib.Path):
self.root = root
# 每 session 一把锁
self._locks: dict[str, asyncio.Lock] = {}
def _lock_for(self, session_id: str) -> asyncio.Lock:
if session_id not in self._locks:
self._locks[session_id] = asyncio.Lock()
return self._locks[session_id]
async def append(self, session_id: str, entries: list[dict]) -> None:
lock = self._lock_for(session_id)
async with lock:
main = self.root / f"{session_id}.jsonl"
sidecar = self.root / f"{session_id}.summary.json"
# 1. 计算部分:读 prev、折叠摘要(纯函数)
prev: dict[str, Any] | None = None
if sidecar.exists():
prev = json.loads(sidecar.read_text())
new_summary = fold_session_summary(prev, session_id, entries)
# 2. 写入部分:append 主体 + 覆盖侧车
with main.open("a", encoding="utf-8") as f:
for e in entries:
f.write(json.dumps(e, ensure_ascii=False) + "\n")
sidecar.write_text(json.dumps(new_summary, ensure_ascii=False))
def fold_session_summary(prev, session_id, entries):
"""与 SDK 契约对齐的纯函数版本"""
summary = {"sessionId": session_id, "mtime": 0, "data": {}}
if prev:
summary["data"] = dict(prev.get("data", {}))
# 示例:累计 user 消息条数
summary["data"]["userCount"] = summary["data"].get("userCount", 0)
summary["data"]["userCount"] += sum(
1 for e in entries if e.get("type") == "user"
)
return summary
这段代码体现了”锁里只做不可分割的两件事”:算摘要、写存储。读 prev、算新摘要都是纯计算,逻辑上完全可被外移;目前把它们放在 async with lock 内只是为了避免”读 prev 与写主文件”之间被另一线程插队。
理解这种”分而不散”的关键是看清 I/O 与计算的边界。读 prev 摘要与写主 transcript 之间存在”读—改—写”窗口,必须串行;计算新摘要与写摘要侧车之间也存在同样的窗口,但摘要计算本身没有副作用,所以可以放在锁内、也可以放在锁外的 worker 线程里用 channel 把结果回流。生产实现里通常会把”读 prev + fold”前置到一个无锁读取阶段,把锁内收缩到”原子写主文件 + 原子写侧车”两件事,进一步压低持锁时长。读者如果是从零开始实现,最稳的路径是先用上面的 asyncio.Lock 版本跑通功能,再用 benchmark 决定是否值得引入 worker 池。
四、四种典型并发策略对比
| 策略 | 一致性 | 吞吐 | 实现成本 | 适用场景 |
|---|---|---|---|---|
| 单 key 全局锁 | 强 | 低 | 低 | 单进程、低并发 |
compute / MULTI |
强 | 中 | 中 | Redis / 内存 |
| CAS + 重试 | 强 | 中 | 中 | 分布式 KV |
| 文件原子写 + 顺序 part | 最终一致 | 高 | 高 | 对象存储 |
CAS 方案常用”读取 prev 摘要 + 带上版本号 → fold → CAS 写回”,版本号失配就重读 prev 重算。Redis Cluster 下还要用 hash tag 把同一 session 的 key 锁在同一个 slot 内,否则 MULTI 跨节点会失败。
五、append 与 foldSessionSummary 的协作边界
| 步骤 | 归属 | 是否在锁内 | 是否纯函数 |
|---|---|---|---|
| 读取 prev 摘要 | adapter | 是 | 否(I/O) |
| 计算新摘要 | FoldSessionSummary | 视实现而定 | 是 |
| 写入主 transcript | adapter | 是 | 否(I/O) |
| 写入侧车摘要 | adapter | 是 | 否(I/O) |
| 设置 MTime | adapter | 是 | 否 |
原则是:把 IO 与”读—改—写”放锁内,把纯计算尽可能外移或置于锁内不持 I/O 的极小段里。落库时序的契约是 MTime 必须与同一 session 的 ListEntry.MTime 共享时钟,否则后续 ListSessionSummaries 的排序会出现偏差。
六、落地时的两条易错点
第一条:把 fold 写成了有副作用的函数,比如顺手”修正 prev 里某个字段”。这会让重放与重算结果不一致,也让 CAS 重试时拿到错误的 prev。原则是 fold 只能基于 (prev, key, entries) 三个入参推导出新摘要。
第二条:用单文件无锁追加 + 文件锁。文件锁在容器与跨主机场景下不可靠,应优先选用 backend 自带的原子原语(Redis MULTI、PostgreSQL 事务、CAS)。当 backend 没有原子原语时,退而用顺序 part 文件加一次性合并。
到此,SessionStore 的一致性链路就完整了:append 在锁内做 IO,foldSessionSummary 在锁外做纯计算,两者通过 SummaryEntry 侧车解耦,既保住 entry 顺序又避免长尾延迟。
常见问题(FAQ)
Q1:append 一定要加锁吗?
同一 SessionKey 必须串行,否则多端并发追加会破坏 entry 顺序。
Q2:foldSessionSummary 为何不直接返回 MTime?
MTime 必须是存储时钟,fold 不接触存储,由 adapter 写入后单独打戳。
Q3:subagent 摘要怎么处理?
对 key.Subpath != "" 不要调用 fold,避免子会话污染主会话摘要。