Codex 在终态附近不只 flush 一次。正常完成至少需要“普通 items 完成后”和“terminal event 追加后”两道屏障;用户中断还在两者之间增加“中断历史标记已持久化”屏障。它们分别保证不同观察者不会读到半套事实。
send_event 的顺序:先持久化,再投递
send_event_raw_with_persistence 先把 EventMsg 转成 RolloutItem 提交给 store,再记录 trace,最后向 Event channel 投递。这里的“提交”可能进入 buffering writer,不代表已经 fsync/flush;因此终态事件发给客户端后还需要显式 flush_rollout()。
async def normal_task_wrapper():
result = await task.run()
await flush_rollout() # barrier 1: items/tool outputs
await on_task_finished(result)
async def on_task_finished(result):
await send_event(TurnComplete(...)) # persist enqueue -> deliver
await flush_rollout() # barrier 2: terminal event
测试 turn_complete_flushes_terminal_event_after_delivery 使用内存 store 计数,明确要求两次 flush。注释还解释原因:buffering thread writer不会因为再无事件就自动把 terminal item推到底层。
Interrupted 为什么需要第三道屏障
Interrupt 会向模型历史写 <turn_aborted> 类标记,让后续恢复或下一轮模型知道上一轮不是自然结束。某些客户端收到 TurnAborted 后会同步重读 Rollout;因此 marker 必须在发事件前 flush。
async def interrupt_barriers(task):
task.token.cancel()
await wait_or_abort(task)
# wrapper若协作返回,会先做普通 item flush
marker = build_interrupted_turn_marker(task.config)
if marker:
await record_conversation_item(marker)
await flush_rollout() # marker-before-event barrier
await send_event(TurnAborted(...))
await flush_rollout() # terminal-event barrier
对应测试要求中断路径共观察到三次 flush:task runner、marker、terminal event,并验证事件接收顺序是 RawResponseItem marker 后 TurnAborted,之后没有额外事件。
Flush 失败为什么不撤销已发终态
任务 wrapper 的第一道 flush 失败会发 Warning,但 Runtime继续完成;终态后的 flush 失败只 warning。因为 event 可能已投递,试图回滚终态会制造更严重的双终态或永不结束。LiveThread/Store 保留重试能力,客户端则得到“转录保存失败、Runtime 仍继续”的明确告警。
因此 flush 是可见性屏障,不是跨 Event channel 与存储的两阶段事务。系统选择单调前进:已产生的事实不撤回,持久化失败被显式暴露并重试。
评论
登录后即可评论