RolloutRecorder 可以被克隆,也可能由多个异步调用点同时使用,但一个 Rollout 文件仍要求确定的追加顺序。若每个调用者都直接持有文件句柄,就必须共同协调文件偏移、Ordinal、延迟创建、失败后的待写后缀以及关闭时机。锁可以保护每次写操作,却很难自然表达“刷新此前所有命令”这样的顺序屏障。
Codex 采用的是单写者结构:Recorder 克隆只持有一个有界 MPSC Sender;文件句柄和全部可变写入状态归一个 Tokio Writer Task 所有。调用者通过四种命令表达追加、物化、刷新和关闭,由单一 Receiver 确定解释顺序。
本节关心的是这套并发协议。JSONL 单行格式、失败前缀重试和 Persist、Flush、Shutdown 的业务调用场景,后续章节还会分别展开。
Recorder 本身不是文件对象
RolloutRecorder 只有三个主要字段:
RolloutRecorder
├── tx: Sender<RolloutCmd>
├── writer_task: Arc<RolloutWriterTask>
└── rollout_path: PathBuf
writer_task 保存任务句柄和最后一次终止性失败,便于 Sender 失效后给出更具体的错误。真正的文件句柄位于 RolloutWriterState 中,而这个 State 只被 rollout_writer 任务移动并持有:
RolloutWriterState
├── writer: Option<JsonlWriter>
├── deferred_log_file_info
├── pending_items
├── meta
├── ordinal_state
├── rollout_path / cwd
└── last_logged_error
这种所有权分配把并发问题缩小成队列排序问题。Recorder 克隆之间无需共享可变文件对象;所有 Ordinal 推进、待写项删除和句柄重开都发生在同一个任务中。
Create 与 Resume 从不同状态开始
构造 Recorder 时,两种模式的 Writer State 不同:
| 构造模式 | writer | meta | 初始物理文件 |
|---|---|---|---|
Create | None | 待写的 SessionMeta | 不立即创建 |
Resume | 已打开为追加模式的 JsonlWriter | None | 已存在并完成尾部检查 |
Create 模式先计算路径和 Session 元数据,但延迟创建文件。这允许一个新建 Thread 在没有任何需要持久化的工作时消失,不留下只有元数据的空 Rollout。Resume 模式则必须立刻打开现有文件,并从已有内容恢复下一 Ordinal。
无论哪种模式,构造函数最后都创建容量为 256 的 tokio::sync::mpsc 通道,并启动同一个 Writer Loop。这个数字限制的是尚未被 Receiver 取走的命令数量,不是 Rollout 行数,也不是一个 Thread 最多能保存的事件数。
四种命令代表四种完成语义
RolloutCmd 的四个变体并不等价:
| 命令 | 是否带 ack | Writer 行为 | 调用成功能证明什么 |
|---|---|---|---|
AddItems | 否 | 追加到 pending_items;已物化时立即尝试写出 | 命令已进入队列,不能单独证明 I/O 成功 |
Persist | 是 | 确保文件打开,写 SessionMeta 和 pending,并刷新 | Writer 已完成本次物化/写入尝试并返回结果 |
Flush | 是 | 处理此前命令;有内容时写出并刷新 | 此屏障前的待写项已成功处理,或收到明确错误 |
Shutdown | 是 | 尝试排空;成功才确认并退出循环 | 此前待写项已处理,Writer Loop 将停止 |
record_canonical_items 会先忽略空切片,然后等待 tx.send(AddItems) 完成。这里的 await 只可能因有界队列反压而等待;命令一旦入队,函数就返回,没有第二个 oneshot 等待磁盘结果。
另外三种命令都携带 oneshot::Sender<Result<()>>。调用者先等待命令进入 MPSC,再等待 Writer Task 回传实际执行结果。这层两阶段等待把“无法排队”和“排队后 I/O 失败”区分开来。
async flush():
ack_tx, ack_rx = oneshot()
await command_queue.send(Flush(ack_tx))
return await ack_rx
AddItems 在两种状态下行为不同
Writer Loop 收到 AddItems 后,先把对象追加到 pending_items,然后调用 flush_if_materialized。
如果 Recorder 仍处于延迟创建状态,这个函数立即返回,pending 只留在内存中。后续 Persist、Flush 或 Shutdown 才会打开文件,先写 SessionMeta,再按顺序写 pending。
如果文件已经物化,AddItems 会立即尝试写出。这次 I/O 失败不会通过 AddItems 返回给原调用者,因为命令没有 ack;State 会丢弃失效的 writer 句柄、保留未成功写完的后缀,等待后续屏障重新打开文件并重试。
因此,下列两个表述必须区分:
record_canonical_items(...).await成功:AddItems 已入队;flush().await成功:Writer 已确认屏障前的待写项成功越过当前实现的刷新边界。
本地 ThreadStore 的常规追加路径会先调用 record_canonical_items,紧接着调用 flush,正是为了把一次 Store append 从“入队成功”提升到“得到 Writer 确认”。
MPSC 队列定义顺序和反压
多个 Sender 可以并发调用,但 Receiver 一次只处理一条命令。屏障的含义来自命令在队列中的位置:Writer 收到某个 Flush 时,排在它之前的 AddItems 已经先被解释;排在它之后的命令不属于这次确认范围。
这里的“之前”指 MPSC 实际接收顺序,不是不同任务开始调用函数的墙上时钟顺序。两个任务同时发送时,谁先完成入队,谁才位于队列前面。若业务需要跨调用者的更强顺序,必须在 Sender 之前建立额外协调。
容量 256 提供反压。Receiver 长时间落后时,后续 send Future 会等待空位,从而把压力传回生产者,而不是让未处理命令无限占用内存。但这个上限也不能直接换算为“最多 256 个待写 Item”:一个 AddItems 可以携带多个 Item,Writer State 还可能因为 I/O 故障保留 pending 后缀。
一个屏障内部会尝试两次
Persist、Flush 和 Shutdown 都通过 write_pending_with_recovery 写出状态:
write_pending_with_recovery(operation):
result = write_pending_once()
if result is success:
clear_last_error()
return success
drop_writer_but_keep_unwritten_suffix()
reopen_and_try_write_pending_once()
if retry succeeds:
clear_last_error()
return success
else:
drop_writer_but_keep_unwritten_suffix()
return final_error
write_pending_once 依次完成四件事:打开或重新打开文件;如有需要先写 SessionMeta;逐项写 pending 并推进 Ordinal;最后调用文件 flush。成功写完的前缀会从 pending_items 删除,失败点之后的后缀保留。因此重试不会从内存队列的第一个已确认 Item 重新开始。
这里的恢复模式针对可重新打开文件的普通 I/O 故障。第一次失败会丢弃当前句柄,第二次使用新句柄重试;第二次仍失败才通过 ack 返回错误。Writer Loop 本身不会因此退出,后续 Flush 或 Shutdown 仍可再次尝试。
Persist、Flush 和 Shutdown 的边界并不相同
三者共享写入函数,但对延迟创建状态和生命周期的处理不同。
Persist 即使没有 pending,也会打开文件并写入 SessionMeta。它的职责就是显式物化,重复调用是幂等的。
Flush 在“仍延迟创建且 pending 为空”时直接成功,不创建文件;一旦 pending 非空,它会物化并写出。它确认已有工作,不主动制造一个只有元数据的新历史。
Shutdown 对延迟且空的 Recorder 同样可以直接成功并退出。若有 pending,则必须先排空。更重要的是,Shutdown 写入失败时只把错误发给当前调用者,Writer Loop 不执行 break,仍保持存活以便调用者修复条件后重试;只有成功确认才退出循环。
如果所有 Recorder Sender 被直接丢弃,Receiver 循环会因通道关闭而结束,它不会自动把这件事等同于一个成功的 Shutdown 命令。对于需要确认排空的生命周期,调用者必须显式等待 shutdown(),不能把 Rust 对象析构当作持久化屏障。
Flush 不是断电级 fsync 承诺
JsonlWriter::write_line 使用 write_all 写入一行及其换行符,然后调用 Tokio 文件的 flush;write_pending_once 在批次结尾还会再次调用 flush。源码中这一链路没有调用 sync_all 或 sync_data。
所以,本节所说的“刷新成功”应严格理解为:异步文件写入已经由 Writer Task 完成并通过 flush 返回,pending 与 Ordinal 状态也已相应推进。它不能仅凭这些代码被扩大解释为操作系统崩溃或突然断电后仍必然存在。若系统需要这种保证,还要明确文件同步、目录同步和文件系统语义。
取消和任务消失留下什么状态
| 故障或取消点 | 可能留下的状态 | 后续含义 |
|---|---|---|
send 尚未完成时调用者取消 | 命令可能未入队 | 不能假定 Writer 看见过该操作 |
AddItems 已入队,实际写入失败 | 未写后缀仍在 pending_items | 后续 Flush/Persist/Shutdown 可重试 |
| 屏障命令已入队,调用者取消 ack 等待 | Writer 仍可能完成命令,ack 接收端已消失 | 本次调用者不知道结果,应以新屏障重新确认 |
| Receiver/Writer Task 已消失 | 新命令无法入队,或 ack 通道关闭 | API 返回队列/任务错误,不得伪装成功 |
Shutdown 写入失败 | Writer Loop 保持存活,pending 保留 | 修复故障后可以再次 Shutdown |
Writer 发送 ack 时会忽略“接收者已经取消”的发送错误。这是合理的:I/O 状态已经发生,无法因等待者离开而回滚。取消只会改变谁知道结果,不会撤销已经排队的命令。
设计取舍:串行化的是状态变更,不是所有上层工作
单写者会让同一 Rollout 的文件变更串行执行,牺牲了对同一文件的写并行度,换来清晰的 Ordinal、pending 和屏障语义。模型请求、工具执行和其他 Thread 的 Recorder 仍可以并发;被串行化的只是一个 Recorder 内的命令流。
源码还有 append_rollout_item_to_path,用于给未加载 Thread 追加已经过滤的元数据项。它的注释明确要求 Live Session 使用 RolloutRecorder::record_canonical_items 来保持与会话流的顺序。因此,“单一 Sender”是活跃 Thread 写入协议,不应被误写成仓库中任何 Rollout 文件访问都只能经过这个对象。
可用于核对实现的不变量
- 只有 Writer Task 修改
RolloutWriterState、文件句柄和 Ordinal; AddItems的 API 成功只确认入队,不能替代带 ack 的屏障;- 一个
Flush只覆盖在 MPSC 接收顺序中位于它之前的命令; - 延迟且空的
Flush/Shutdown不应凭空创建 Rollout; Persist即使无 pending,也应能物化SessionMeta;- I/O 失败后只删除成功前缀,未写后缀必须保留以供重试;
Shutdown失败不能终止 Writer Loop,成功才结束;- 通道关闭或任务消失后,API 不能返回虚假成功;
flush的实现保证不能被表述为源码中不存在的fsync。
源码与测试证据
| 位置 | 可以核对的结论 |
|---|---|
codex-rs/rollout/src/recorder.rs | Recorder 字段、四种命令、256 容量、Writer State、重试和 ack 协议 |
codex-rs/thread-store/src/local/live_writer.rs | Store 追加在 record_canonical_items 后显式调用 flush |
codex-rs/rollout/src/recorder_tests.rs::recorder_materializes_on_flush_with_pending_items | 延迟文件由带 pending 的 Flush 物化,顺序和 Persist 幂等性成立 |
codex-rs/rollout/src/recorder_tests.rs::persist_reports_filesystem_error_and_retries_buffered_items | 物化失败保留缓存,修复路径后 Flush 可重试 |
codex-rs/rollout/src/recorder_tests.rs::writer_state_retries_write_error_before_reporting_flush_success | Writer 句柄失败后会重开并在成功后确认 Flush |
codex-rs/rollout/src/recorder_tests.rs::paginated_ordinal_overflow_fails_without_appending | Ordinal 无法推进时屏障返回错误且不追加 Item |
本节边界
本节建立了 Recorder 内部的命令完成语义,但没有完整讨论延迟创建时 SessionMeta 与第一批 Item 的次序,也没有穷举不同生命周期 API 对上层 Store 的保证。8.5 将先分析首次物化,8.6 和 8.7 再分别讨论前缀重试与生命周期屏障。
源码导航
codex-rs/rollout/src/recorder.rscodex-rs/rollout/src/recorder_tests.rscodex-rs/thread-store/src/local/live_writer.rs
评论
登录后即可评论