青雲的博客
深入浅出 DeepSeek Harness 第四部:事件账本——Session是append-only日志 第 24 章

JSONL 原子写与崩溃恢复

深入剖析 session-persistence-jsonl 的原子写保证、write-behind 批量策略(默认200ms窗口)、以及 interruptedTurnClosers 崩溃恢复算法。每个事件占一行 JSON、部分行等于未完成事件、重启后丢弃最后不完整行并合成确定性关闭事件。

源码版本
47f943859bef60e4160492346772ded9b24f765a
验证日期
Commit
47f943859bef60e4160492346772ded9b24f765a

引言:持久化不只是”保存到文件”

你可能以为 Session 持久化就是调一下 fs.writeFile——数据进文件,完工。但在一个 append-only 事件日志架构中,持久化面临三道根本难题:

  1. 原子性——部分写入不能留下损坏的日志;
  2. 性能——每个事件独立 fsync 会杀死吞吐量;
  3. 恢复——进程崩溃后,日志必须能自愈到一致状态。

下面我们就按真实落盘顺序把这套解法拆开:write-behind 批量策略怎么控吞吐,原子物化与追加怎么守住边界,torn-tail 怎么定义并被丢弃,最后 interruptedTurnClosers 怎么补上确定性收尾。

先看结论:JSONL 的换行符就是提交边界,write-behind 负责把性能从”每事件一次 I/O”变成”一段时间一次 I/O”,崩溃恢复只需要丢弃 torn tail 并补一个确定性收尾。

后面的细节基本都绕着这三件事转:

  1. 换行符 = commit 边界:没有换行的最后一行就是未提交碎片,读取时自然丢弃。
  2. write-behind = 工程吞吐:事件先进内存日志,I/O 批量落盘,默认 200ms 窗口是吞吐与延迟的取舍。
  3. 确定性收尾 = 自愈:崩溃时 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/startassistant/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 路径

步骤:

  1. mkdir(root/project/session, recursive, 0o700) + syncDir 每一级
  2. 写入临时文件 (*.tmp) + handle.sync()
  3. link(tmp, finalPath) ——如果 finalPath 已存在,link 返回 EEXIST,两个进程不可能互相覆盖
  4. syncDir(dir) ——确保目录条目 crash-durable
  5. 清理临时文件

为什么是 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
  }
}

如果 writeFilesync 失败,文件被 truncate 回写入前的大小并 fsync——把部分写入的字节彻底清除。Coordinator 的 cursor 没有前进,下次 flush 会重试同一个 batch。


第四层:Torn-Tail 检测与截断

进程在 writeFilesync 之间崩溃——磁盘上留下了部分写入的字节。下次读取时怎么办?

纯文本模式

scanLogSessionLogScanner.finish() 的语义是:没有换行结尾的最后一片字节是 torn tailcommittedBytes 记录了最后一个完整行的结束位置。

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)
}

两步操作:

  1. truncate(path, offset) + fsync——物理删除 torn bytes
  2. 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——当前开放的 turn
  • openStep: number | null——当前开放的 step
  • pendingCalls: Map<CallId, ...>——已声明但未收到 result 的 tool calls

事件消费规则:

事件效果
turn/startopenTurn = turn,清空 step 和 calls
turn/endopenTurn = null,清空一切
step/startopenStep = step
step/endopenStep = 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 对象)

注意 closerstornMarker 分开保存——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”:

  1. 读取磁盘上已有的 prefix
  2. 验证 live Session 的 seed 覆盖了 stored prefix
  3. 只截断 torn tail(关闭 open turn——live Session 还在运行!)
  4. 追加 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 turninterruptedTurnClosers 合成 closers
Tool call 已执行但 result 未写tool/call 在日志中合成 TOOL_OUTCOME_UNKNOWN result
Tool call 未执行只有 assistant/message 中的声明合成 TOOL_NOT_STARTED result

收口:可靠性从字节边界开始

  1. JSONL 换行就是事务边界——完整行 = 提交,不完整行 = 丢弃,几乎零成本。
  2. Write-behind 200ms 是最大延迟,不是最小延迟——deadline 到期或显式 flush,谁先到就先写。
  3. 恢复是一段确定性纯函数——相同的日志前缀永远合成相同的收尾事件,所以能测试、能推理。
  4. “未开始”和“结果未知”必须分开——模型需要这点信息来判断能不能重试。
  5. 修复本身不需要再包一层原子性——修复是幂等的,多跑一次结果也一样。
  6. Coordinator 的 per-id 操作必须串行化——一把逻辑锁消掉所有并发交错。

所以这里讲的已经不是“保存到文件”,而是一条从字节级物理边界一路撑到语义级事件一致性的可靠性栈。