雨天小六

读懂 Codex(9.14):JSONL Writer Task、Flush 和故障注入

· 更新于 2026-08-03 · 专栏:读懂 Codex

#Codex#Agent Runtime#Python#软件架构

具体问题与边界

批量追加第 k 行失败时,pending 应该删除多少,下一次 Flush 又应重试哪一段?

JsonlRollout 拥有 pending 语义记录;JsonlWriter 专门串行阻塞 I/O;JsonlSink 可注入故障;PrefixWriteError 报告已确认前缀。 本节区分“从官方源码得到的生产事实”和“为教学实现作出的 Python 选择”,不会把后一种包装成 Codex 的等价实现。

状态所有权

对象所有者生命周期持久化
pending recordsJsonlRolloutappend 到 confirmed内存
Writer command queueJsonlWriterwriter task 生命周期
immutable batch一次 Flush等待 sink ack
confirmed countJsonlSink/PrefixWriteError一次 batch
JSONL complete lines文件系统跨进程是,恢复权威

正常路径

JSONL Writer Task、Flush 和故障注入正常路径图
图 9.14-1:批量追加第 k 行失败时,pending 应该删除多少,下一次 Flush 又应重试哪一段?
  1. Session 的 append 只把 JSON object 加到 pending,不直接打开文件;锁保护 pending 与 flush 前缀。
  2. flush 在锁内复制 immutable batch,送入容量 16 的 Writer command Queue,并等待 Future。
  3. Writer Task 用 asyncio.to_thread 执行阻塞 sink;FileJsonlSink 逐行写、flush、fsync,每确认一行递增 count。
  4. 成功时只有 confirmed 等于 batch 长度才删除 pending 前缀;静默少确认被当作 RuntimeError。
  5. 部分失败抛 PrefixWriteError(confirmed=k),Rollout 只删除前 k 个,保留 suffix 给下一次 flush。追加前若旧尾无换行,先补 delimiter,避免新 JSON 粘在坏尾。

顺序为什么不能交换

Session 的 append 只把 JSON object 加到 pending,不直接打开文件;锁保护 pending 与 flush 前缀
→ flush 在锁内复制 immutable batch,送入容量 16 的 Writer command Queue,并等待 Future
→ Writer Task 用 `asyncio.to_thread` 执行阻塞 sink;FileJsonlSink 逐行写、flush、fsync,每确认一行递增 count
→ 成功时只有 confirmed 等于 batch 长度才删除 pending 前缀;静默少确认被当作 RuntimeError
→ 部分失败抛 PrefixWriteError(confirmed=k),Rollout 只删除前 k 个,保留 suffix 给下一次 flush

箭头代表可见性与所有权转移,不是松散依赖。Policy、Approval、外部副作用、规范 Item 和 durability 各自有提交点,后一步不能替前一步作更强承诺。

Python 风格伪代码

async def rollout.flush():
    async with lock:
        if not pending: return
        batch = tuple(pending)
        try:
            confirmed = await writer.write(batch)
        except PrefixWriteError as exc:
            del pending[:exc.confirmed]
            raise
        if confirmed != len(batch):
            raise ProtocolError("short confirmation")
        del pending[:confirmed]

async def writer_loop():
    while command := await queue.get():
        if command is STOP: return
        try:
            count = await to_thread(sink.write_batch, path, command.batch)
        except Exception as exc:
            command.future.set_exception(exc)
        else:
            command.future.set_result(count)

def file_sink(batch):
    repair_missing_tail_newline_if_needed()
    for record in batch:
        append_json_line(record); flush(); fsync()
        confirmed += 1
    return confirmed

失败、取消与恢复

JSONL Writer Task、Flush 和故障注入失败路径图
图 9.14-2:每个失败点都列出已经发生的副作用和仍可安全执行的恢复动作。
故障点已留下的状态处理
第 k 行写失败前 k-1 行可能已耐久confirmed 前缀从 pending 删除,后缀重试
sink 少确认但不报错无法知道真实状态协议错误,不删除未确认部分
坏尾无 newline下一行可能粘连追加前补 delimiter,loader 仍报告旧坏行
writer task exceptionFuture 返回异常Turn 进入持久化失败路径
Shutdown 未 drainpending 可能丢失Session 先 close Rollout,再发事件流哨兵

不变量

只有 sink 明确确认的完整前缀可以从 pending 删除;Flush 失败后不得重写已确认前缀,也不得丢弃未确认后缀。

设计思路与限制

逐行 fsync 很保守且慢,适合作为语义演示,不是吞吐最优策略。即使 fsync 也不是跨设备全局事务;Mini 没有跨进程 writer lock、ordinal 与 SQLite projection。

测试与复现

cd examples/mini-codex
uv run pytest -q -k 'test_writer_retries_only_unconfirmed_suffix or test_resume_reports_incomplete_turn_and_repairs_prompt_only'
uv run mypy src
  • test_writer_retries_only_unconfirmed_suffix
  • test_resume_reports_incomplete_turn_and_repairs_prompt_only

官方源码导航

Mini Codex 对照

  • src/mini_codex/persistence/writer.py:Writer Task、Sink 与 PrefixWriteError
  • src/mini_codex/persistence/rollout.py:pending、flush、close 与记录类型
  • src/mini_codex/runtime/session.py:Turn 边界显式 flush/close

本节结论

只有 sink 明确确认的完整前缀可以从 pending 删除;Flush 失败后不得重写已确认前缀,也不得丢弃未确认后缀。

阅读导航

上一节:9.13 · 下一节:9.15

评论


← 返回文章列表