JSONL 原子写与崩溃恢复
深入剖析 session-persistence-jsonl 的原子写保证、write-behind 批量策略(默认200ms窗口)、以及 interruptedTurnClosers 崩溃恢复算法。每个事件占一行 JSON、部分行等于未完成事件、重启后丢弃最后不完整行并合成确定性关闭事件。
引言:持久化不只是”保存到文件”
你可能以为 Session 持久化就是调一下 fs.writeFile——数据进文件,完工。但在一个 append-only 事件日志架构中,持久化面临三道根本难题:
- 原子性——部分写入不能留下损坏的日志;
- 性能——每个事件独立 fsync 会杀死吞吐量;
- 恢复——进程崩溃后,日志必须能自愈到一致状态。
下面我们就按真实落盘顺序把这套解法拆开:write-behind 批量策略怎么控吞吐,原子物化与追加怎么守住边界,torn-tail 怎么定义并被丢弃,最后 interruptedTurnClosers 怎么补上确定性收尾。
先看结论:JSONL 的换行符就是提交边界,write-behind 负责把性能从”每事件一次 I/O”变成”一段时间一次 I/O”,崩溃恢复只需要丢弃 torn tail 并补一个确定性收尾。
后面的细节基本都绕着这三件事转:
- 换行符 = commit 边界:没有换行的最后一行就是未提交碎片,读取时自然丢弃。
- write-behind = 工程吞吐:事件先进内存日志,I/O 批量落盘,默认 200ms 窗口是吞吐与延迟的取舍。
- 确定性收尾 = 自愈:崩溃时 turn 可能只写到一半,重启后用
interruptedTurnClosers合成关闭事件,让日志回到一致状态。
带着这三点去读细节,你会发现整套方案没有魔法:它靠的是明确的物理边界(换行)、可控的批量策略(write-behind),以及一个可证明正确的修复函数(closers)。
第一层:JSONL 格式——一行一事件的契约
文件结构
每个 Session 在磁盘上对应一个 .jsonl(或 .jsonl.zstd)文件。文件第一行是带 type: 'session' 标签的 header record,之后每一行是一个 SessionEvent 或一个 packed chunk-row。
export function eventLines(events: readonly SessionEvent[], packChunks: boolean): string {
const records: readonly StorageRecord[] = packChunks ? packChunkRuns(events) : events
return records.map(record => JSON.stringify(record)).join('\n')
}
关键设计:每个换行符是一个原子边界。写入过程中如果断电,最后一行可能不完整(没有换行结尾)——读取时自然被识别为”未提交的 torn tail”并丢弃。你不需要任何额外的事务日志或 WAL。
Header 与事件的分离
Header line 使用 type: 'session' 标签,永远不会与事件类型(如 turn/start、assistant/chunk)冲突。SessionLogScanner 在构造时就解析 header,之后逐块处理事件行。
Scanner 的 write(chunk) 方法在原始字节上做换行搜索——只有完整的行(以 0x0A 结尾)才会被 decode 和 parse。跨 chunk 的片段被安全缓存后拼接。最终 finish() 时,没有换行结尾的残余字节被静默忽略——这就是 torn-tail 的物理定义。
第二层:Write-Behind 批量——200ms 窗口的工程取舍
问题
Session 的 append() 是同步提交到内存日志的(调用方永远不阻塞于 I/O)。但如果每个事件独立 fsync,一次 token stream 的几百个 assistant/chunk 事件就会触发几百次磁盘刷新——性能灾难。
SessionWriteBehind 设计
核心状态机:
idle → enqueue(event) → arm timer (200ms)
↓ deadline expires
startBackground() → write(batch) → durable
↓ success
continueAutomatic() → (if deadline expired during write, start next)
当你调用 session.append('assistant/chunk', {...}) 时,coordinator 的 session/event 监听器将事件的 structuredClone 拷贝入 pending 队列。如果队列从空变为非空,一个 200ms 的定时器开始倒计时。
为什么是 200ms? 这是 DEFAULT_WRITE_BATCH_MAX_DELAY_MS 的值:
200ms 足够在典型 token stream 中累积数十个事件,同时对人类用户不可察觉——一个 token 的平均间隔约 20-50ms,所以一个批量能攒 4-10 个 chunk 事件。
失败重试与背压
private startWrite(background: boolean): Promise<void> {
const batch = this.pending.splice(0)
// ...
const active = operation
.catch((error: unknown) => {
this.pending = batch.concat(this.pending)
// ... pause automatic scheduling
if (background) this.options.reportBackgroundFailure(error)
throw error
})
.finally(() => { this.active = undefined })
this.active = active
return active
}
如果后端写入失败(磁盘满、权限错误),batch 被完整地放回队列头部——后续的 flush() 会重试,事件顺序永远不乱。同时 automatic 调度暂停(automaticPaused = true),直到下一次 enqueue 显式恢复。这确保后台错误不会无限重试、不会丢数据。
显式 Flush 屏障
flush() 是调用者的即时持久化保证:
flush(): Promise<void> {
if (this.barrier !== undefined) return this.barrier
this.cancelTimer()
// ... create barrier
void this.drainBarrier(barrier.resolve, barrier.reject)
return barrier.promise
}
多个并发 flush() 共享同一个 barrier(this.barrier),避免重复刷盘。drainBarrier 先等待可能正在进行的 active write,然后循环写入直到队列清空。Coordinator 在以下时机调用 flush:
session/flush事件(checkpoint policy 触发)- Session dispose 时的 retirement drain
- Service dispose 时的全局 quiescence
第三层:原子物化——首次写入的安全保证
当一个新 Session 首次产生事件时,后端需要原子地创建文件(header + 首批事件一起落盘)。这里的”原子”指的是:要么文件完整存在(header + 事件),要么不存在——没有中间态。
POSIX 路径
步骤:
mkdir(root/project/session, recursive, 0o700)+syncDir每一级- 写入临时文件 (
*.tmp) +handle.sync() link(tmp, finalPath)——如果 finalPath 已存在,link 返回 EEXIST,两个进程不可能互相覆盖syncDir(dir)——确保目录条目 crash-durable- 清理临时文件
为什么是 link() 而不是 rename()? 因为 rename() 会静默覆盖已有文件,而 link() 在目标存在时失败。这是 TOCTOU 安全的关键:即使两个进程同时 materialize 同一个 id,只有一个能成功。
追加路径的原子性
对于已存在文件的追加(appendLines),后端使用了不同但同样安全的策略:
private async appendLines(meta, events): Promise<void> {
const handle = await open(path, 'a')
const { size: before } = await handle.stat()
try {
await handle.writeFile(content)
await handle.sync()
} catch (error) {
await this.rollbackAppend(path, before)
throw error
}
}
如果 writeFile 或 sync 失败,文件被 truncate 回写入前的大小并 fsync——把部分写入的字节彻底清除。Coordinator 的 cursor 没有前进,下次 flush 会重试同一个 batch。
第四层:Torn-Tail 检测与截断
进程在 writeFile 和 sync 之间崩溃——磁盘上留下了部分写入的字节。下次读取时怎么办?
纯文本模式
scanLog 和 SessionLogScanner.finish() 的语义是:没有换行结尾的最后一片字节是 torn tail。committedBytes 记录了最后一个完整行的结束位置。
finish(): SessionLogScan {
this.finished = true
return { meta: this.meta, events: this.events, committedBytes: this.committedBytes }
}
如果 committedBytes < buffer.byteLength,coordinator 知道有 torn tail 需要修复。
Zstandard 模式
Zstd 编码下每次 append 是一个独立 frame。scanZstdFrames 扫描 magic number 边界;一个不完整的最后 frame(tornStart)被尝试 prefix-decode——能恢复多少行就恢复多少,剩余丢弃。
修复执行
async commitRepair(meta, tornMarker, closers): Promise<void> {
if (tornMarker !== undefined) await this.repair(meta, tornMarker.truncateTo)
const repairedEvents = [...(tornMarker?.recoveredEvents ?? []), ...closers]
if (repairedEvents.length > 0) await this.appendLines(meta, repairedEvents)
}
两步操作:
truncate(path, offset)+ fsync——物理删除 torn bytes- append 恢复的完整行 + 合成的 closer 事件
Coordinator 文档明确说这不需要原子(“NOT required to be atomic”)——因为即使第二步失败,下次加载时 truncated 的日志仍然有效(只是缺少 closers,下次再合成)。
第五层:interruptedTurnClosers——确定性恢复算法
截断 torn tail 后,日志可能停在一个 turn 的中间——没有 turn/end,也许有悬空的 tool call。Provider(如 DeepSeek)的 API 会拒绝没有匹配 tool result 的 assistant message。所以恢复必须合成关闭事件。
算法概览
interruptedTurnClosers(events) 纯函数,输入是已提交的事件前缀,输出是需要追加的合成事件列表:
状态机追踪三个游标:
openTurn: number | null——当前开放的 turnopenStep: number | null——当前开放的 steppendingCalls: Map<CallId, ...>——已声明但未收到 result 的 tool calls
事件消费规则:
| 事件 | 效果 |
|---|---|
turn/start | openTurn = turn,清空 step 和 calls |
turn/end | openTurn = null,清空一切 |
step/start | openStep = step |
step/end | openStep = null,清空 calls |
assistant/message | 注册其中的 tool-call blocks 到 pendingCalls |
tool/call | 记录 callSeq(标记为”已开始”) |
tool/result | 从 pendingCalls 删除对应 callId |
平衡日志——无需修复
if (openTurn === null || last === undefined) return []
如果扫描结束时 openTurn === null,说明最后一个 turn 已经正常关闭——不需要任何合成事件。
合成 tool/result
对于每个悬空的 tool call,合成一个 isError: true 的 tool result。合成的错误消息区分两种状态:
未开始(没有 tool/call 事件记录):
“The tool call was interrupted before the Harness recorded it as started. Retry it if it is still needed.”
结果未知(有 tool/call 但没有 tool/result):
“The tool call was interrupted after it was recorded, but no result was durably recorded. Its outcome is unknown. Decide whether to retry…”
这个区分至关重要——已开始的 tool call 可能有副作用(写了文件、执行了命令),模型需要知道不能盲目重试。
合成 step/end 和 turn/end
if (openStep !== null) {
closers.push({ type: 'step/end', seq: seq++, time, data: { turn: openTurn, step: openStep } })
}
closers.push({ type: 'turn/end', seq: seq++, time, data: { turn: openTurn, reason: { kind: 'interrupted' } } })
顺序严格:先关闭 step(不然 turn/end 时 step 仍然开放,违反不变量),再关闭 turn。
确定性保证
所有合成事件复用最后一个真实事件的 time:
const time = last.time
不编造”未来时间”,不使用 Date.now()。这确保了:
- 相同输入永远产生相同输出(纯函数)
- 合成事件的时间戳不会比它们”声称”发生之后的事件更新
- 测试可以精确断言
seq 从 last.seq + 1 开始连续递增,保持全局 seq 连续性。
第六层:Coordinator 的编排全景
加载路径(prepareCore)
loadStored(id)
→ assertVersion + adoptStoredEvents(验证+冻结)
→ interruptedTurnClosers(storedEvents)(计算合成事件)
→ balanced = [...storedEvents, ...closers]
→ Session.fromRestore(id, balanced, meta)(构造 Session 对象)
注意 closers 和 tornMarker 分开保存——commitPrepared 阶段才执行物理修复。这意味着 prepare 是只读的、可缓存的,只有 commit 才写磁盘。
写入路径(installWritePath)
session/event → initFor(session) → writes.enqueue(event)
session/flush → flush(session) → writes.flush()
session/disposed → retire(session) → flush + cleanup
Coordinator 用 per-id promise chain(serialize)保证同一个 session 的写操作永远不交错——即使 flush 和 background write 并发也不会产生乱序。
Live Adoption(HMR 场景)
当热重载发生时,已有的 live Session 不会重放 session/created。Coordinator 的 initFor 对每个存活 Session 执行”live adoption”:
- 读取磁盘上已有的 prefix
- 验证 live Session 的 seed 覆盖了 stored prefix
- 只截断 torn tail(不关闭 open turn——live Session 还在运行!)
- 追加 live seed 中超出 stored 部分的 suffix
第七层:Revision 稳定性与并发读保护
读取路径面临一个微妙问题:如果你先 stat() 获取文件大小,再 readFile() 读取内容,写入方可能在这两步之间追加了字节——你读到的内容比 stat 时更长,包含了部分写入的数据。
JSONL 后端用 readStableFile 解决:
private async readStableFile(path, signal): Promise<{ buffer; revision }> {
for (;;) {
const before = fileRevision(await stat(path, { bigint: true }))
const buffer = await readFile(path, { signal })
const after = fileRevision(await stat(path, { bigint: true }))
if (before === after) return { buffer, revision: after }
}
}
双 stat 包围 readFile:如果前后 revision 不一致(dev/ino/size/mtime/ctime 任一变化),重试。对于一个 append-only 文件,这保证你读到的 buffer 对应一个完整的、不变的物理快照。
Revision 本身是五元组字符串 dev:ino:size:mtimeNs:ctimeNs,用于 Coordinator 的 prepared-source 缓存失效:如果 readStoredRevision 返回的值不等于缓存时的 revision,缓存的 PreparedSource 被丢弃并重建。
第八层:Chunk Packing——不是恢复但影响恢复
Delta chunk 事件(token-level assistant/chunk)在存储时被打包为 text-chunks / reasoning-chunks / tool-call-chunks 行——一个 packed row 代替几十个独立事件行,实测节省约 60% 空间。
读取时 decodeStorageRecord 无条件展开 packed rows 回原始事件。所以 torn-tail 检测在原始字节级工作——一个不完整的 packed row 被丢弃,其包含的所有事件一起丢失。这是安全的:它们都是同一个 flush batch 的子集,要么全在要么全不在。
第九层:Seq 连续性——不变量与防御
整个系统依赖一个核心不变量:event.seq === log.length。Scanner 在解码时验证:
if (event.seq !== this.events.length) {
this.issue = new Error(`seq gap in committed region at line ${this.eventLine}`)
}
一旦发现 seq gap,后续事件不再追加到 prefix(避免在损坏区域上堆积)。但如果 gap 后面出现了 turn/end,则立即抛出错误——因为一个跨越损坏区域的 turn boundary 意味着无法安全恢复。
Coordinator 的 appendCore 同样强制验证:
for (const [i, event] of events.entries()) {
if (event.seq !== state.cursor + i) {
throw new Error(`append seq mismatch...`)
}
}
这形成了端到端的 seq 完整性链:内存 → 序列化 → 磁盘 → 加载 → 验证。
第十层:Dispose 与 Retirement 的有序关闭
当一个 Session 被 dispose 时(scope 卸载、用户关闭),write-behind 队列中可能还有未刷盘的事件。Coordinator 通过 retire() 保证有序收尾:
private retire(session: Session): void {
const retirement = this.retireCore(session)
this.retirements.set(session.id, retirement)
// ...
}
private async retireCore(session: Session): Promise<void> {
await this.flush(session)
await this.serialize(id, () => {
this.live.delete(session)
if (this.states.get(id)?.owner === session) this.states.delete(id)
})
}
关键点:
- Retirement 是一个被跟踪的 promise——
prepare()和load()在尝试操作同一个 id 前会先waitForRetirement - Flush 失败不会阻塞整个系统——错误被记录但 retirement promise 仍然 settle
- Service dispose(进程优雅退出)时,所有 live session 并行 flush,然后等待所有 chain settle
这确保了:即使你在 token stream 进行中关闭应用,已产生的事件也会尽最大努力落盘。如果仍然失败(磁盘物理不可写),下次加载时 torn-tail 机制兜底。
故障模式全景
| 故障场景 | 磁盘状态 | 恢复行为 |
|---|---|---|
| append 中断(write 后 sync 前) | 部分行无换行 | torn tail 截断 |
| sync 后 cursor 更新前崩溃 | 完整行已持久 | 下次 appendLiveBatch 跳过 seq < cursor 的事件 |
| materialize 中断(link 前) | temp 文件存在 | 下次 create → no artifact → 重新 materialize |
| materialize 中断(link 后 syncDir 前) | finalPath 存在但目录未 sync | 极端情况下 finalPath 消失 → 同上 |
| Turn 中间崩溃 | 完整事件 + open turn | interruptedTurnClosers 合成 closers |
| Tool call 已执行但 result 未写 | tool/call 在日志中 | 合成 TOOL_OUTCOME_UNKNOWN result |
| Tool call 未执行 | 只有 assistant/message 中的声明 | 合成 TOOL_NOT_STARTED result |
收口:可靠性从字节边界开始
- JSONL 换行就是事务边界——完整行 = 提交,不完整行 = 丢弃,几乎零成本。
- Write-behind 200ms 是最大延迟,不是最小延迟——deadline 到期或显式 flush,谁先到就先写。
- 恢复是一段确定性纯函数——相同的日志前缀永远合成相同的收尾事件,所以能测试、能推理。
- “未开始”和“结果未知”必须分开——模型需要这点信息来判断能不能重试。
- 修复本身不需要再包一层原子性——修复是幂等的,多跑一次结果也一样。
- Coordinator 的 per-id 操作必须串行化——一把逻辑锁消掉所有并发交错。
所以这里讲的已经不是“保存到文件”,而是一条从字节级物理边界一路撑到语义级事件一致性的可靠性栈。