Learning AI Quality 返回 KuthorX Blog II博客首页

第 14 章

第 3 周:Codex 的 agent loop 比我们的多了什么?

对照 openai/codex 源码读生产级 agent loop:三层循环、needs_follow_up 停止条件、流式解析时工具一完整就开跑、读写锁控制并行、按调用顺序回填、错误分类与指数退避和 WebSocket→HTTP 降级、用户中途插话,以及为什么它没有 max_turns。一段讲解视频,一个工具时间轴和重试模拟器。

第 2 周手写的 agent loop 只有二三十行:调模型,stop_reason 是 tool_use 就执行工具,把结果塞回去,再调模型。OpenAI 的 Codex 是一个每天被大量开发者使用的编码 agent,它的 loop 骨架和我们的一模一样,但多了很多东西:模型还在吐字,工具就已经开跑了;有的工具能一起跑,有的必须独占;断流了会退避重试,重试用完还会换一条传输通道;用户可以在 agent 干活时插一句话。这一章对照源码,把这些"多出来的东西"一段一段读清楚,再从测开的角度看每一段怎么测。

本章所有源码引用都基于 openai/codex commit 7993248(799324821d36a822923cee7814d3b80f7ec3cf99,2026-10-01),路径省略 codex-rs/ 前缀,行号可以点进去看。不会 Rust 也没关系,代码下面都有逐行解释。

讲解视频

互动演示

两个模拟器,都不调真实 API。第一个是工具执行时间轴:调模型输出速度、每个工具调用的参数长度、每个工具的耗时和类型(可并行的 exec_command 还是独占的 apply_patch),对比"等整个响应结束再执行"和"工具一完整就执行"的总耗时,回填顺序和读写锁排队都按 Codex 的规则计算。第二个是重试模拟器:注入断流、429、500、503、连接失败、上下文超窗等故障,看每次请求之后等多久、什么时候降级、什么时候放弃,默认参数取自源码。页面底部有自动判分的练习。

互动演示:工具时间轴与重试模拟 在新标签页打开

三层循环

第 2 周的 loop 只有一层 for _ in range(max_turns)。Codex 的一个 turn(用户发一条消息,到 agent 交还控制权)里套了三层:

run_turn                                        core/src/session/turn.rs:163
└─ loop {                 // ① 外层:一个 turn 里的多次采样(step)     :426
     取出用户中途插话 pending_input                                  :430
     run_sampling_request
     └─ loop {            // ② 采样重试                                 :1653
          try_run_sampling_request
          └─ loop { stream.next() }   // ③ 流处理                      :2618
               output_item.done 是工具调用 → 立刻开跑,future 放进 in_flight
               response.completed → 记 token、看 end_turn,退出流循环
          drain_in_flight:按调用顺序把工具结果写进历史                 :3162
          出错 → handle_response_stream_error:分类、退避、降级
        }
     needs_follow_up = 模型要继续 || 有待处理的插话                     :566
     !needs_follow_up → Stop hook → break                             :653
   }

run_turn 的文档注释( core/src/session/turn.rs:149-162 )说得很直白:每次采样,模型要么回工具调用,要么回一条 assistant 消息;回的是工具调用,就执行,把输出送进下一次采样;只回了消息,turn 就结束。这就是第 2 周的 while 循环。

第 2 周的 loopCodex差别
stop_reason == "tool_use" 就继续needs_follow_up:本次采样出现过工具调用、end_turn == false、或者有用户插话多了"用户插话"这个来源
等整个响应回来,再执行工具流式解析,一个工具调用完整就开跑工具和模型生成重叠
工具串行执行声明可并行的拿读锁同时跑,其余拿写锁独占按工具声明决定能否并行
请求失败就抛异常错误分成"终止 / 只听服务端 / 可重试",指数退避、Retry-After、WebSocket→HTTP 降级失败是常态,按类型处理
max_turns + token 预算 + 超时 + 重复检测没有步数上限;靠上下文压缩、用户中断、终止类错误兜底交互式和无人值守的取舍不同

停止条件:needs_follow_up 由什么决定

Codex 调的是 OpenAI Responses API,没有 Claude 那样的 stop_reason: "tool_use",它看输出里有没有工具调用 item。三个来源:

// core/src/stream_events_utils.rs:350-357(流里收到一个完整的工具调用)
let tool_future: InFlightFuture<'static> = Box::pin(
    ctx.tool_runtime
        .clone()
        .handle_tool_call(call, cancellation_token),
);

output.needs_follow_up = true;
output.tool_future = Some(tool_future);

// core/src/session/turn.rs:2994-2996(response.completed 里的 end_turn)
if let Some(false) = end_turn {
    needs_follow_up = true;
}

// core/src/session/turn.rs:566(采样结束后,外层循环再看一眼插话队列)
let needs_follow_up = model_needs_follow_up || has_pending_input;

逐段解释:

  • 第一段:handle_tool_call(...) 返回一个"等工具结束"的 future,Box::pin 把它装箱,方便放进队列。紧接着把 needs_follow_up 置为 true:模型要看到工具结果才能继续,所以必须再采样一次。源码在 stream_events_utils.rs:325-357 。工具名不认识也走这条路:ToolRouter::build_tool_call 不检查工具名( tools/router.rs:248-301 ),调用照常派发,到注册表里找不到才返回 unsupported call: <工具名>( tools/registry.rs:551-570 ),这句话作为失败的工具输出按序回填,模型下一轮看到后自己改。build_tool_call 本身返回 RespondToModel 时(目前只有 tool_search 参数解析失败一处, router.rs:274 ),走的是另一个分支:不派发,直接写一条工具输出回历史并置 needs_follow_up = true( stream_events_utils.rs:400-424 )。
  • 第二段:if let Some(false) = end_turn 是 Rust 的模式匹配,意思是"end_turn 这个可选字段有值、而且值是 false"。服务端明说"还没完",就继续( turn.rs:2994-2996 )。response.incomplete 且原因是 interrupted 时,SSE 解析层也会把它转成 end_turn: Some(false)( codex-api/src/sse/responses.rs:446-455 )。
  • 第三段:|| 是"或"。模型这一步说完了,但用户在它说话时又发了消息,那就再采样一次,让模型看到新消息( turn.rs:566 )。

三者都不成立,进入 Stop hook( turn.rs:653-757 )。hook 返回 should_block 并附带一段续写提示时,提示作为新消息写进历史,continue 回到外层循环;should_stop 时直接 break。否则看要不要做 turn 结束时的压缩,然后 break。

对照第 2 周"只有 end_turn 才算完成":Codex 也不把截断当成功。response.incomplete 的原因不是 interrupted 或 content_filter 时,会被当成流错误 Incomplete response returned, reason: …( sse/responses.rs:418-435 ),走下面的重试分支。

流处理:工具一完整就开跑,结果按顺序回填

模型在一次响应里调 4 个工具,第 1 个调用的参数在第 2 秒就生成完了,整个响应第 5 秒才结束。第 2 周的写法要等到第 5 秒才开始执行第 1 个工具。Codex 在收到 response.output_item.done、发现它是工具调用的那一刻就开跑:

// core/src/session/turn.rs:2589
let mut in_flight: FuturesOrdered<InFlightFuture<'static>> = FuturesOrdered::new();
// ... 流循环里,每收到一个 OutputItemDone:
// core/src/session/turn.rs:2787-2793
if let Some(tool_future) = output_result.tool_future {
    in_flight.push_back(tool_future);
}
// ...
needs_follow_up |= output_result.needs_follow_up;
// ... 流循环结束后(不管成功还是出错):
// core/src/session/turn.rs:3162-3171
if !in_flight.is_empty() {
    // ...
    drain_in_flight(&mut in_flight, sess.clone(), &step_context).await?;
}
  • FuturesOrdered 是 futures 库里"按放入顺序吐结果"的 future 队列。里面的 future 可以乱序完成,但取结果时一定先拿第 1 个、再拿第 2 个。所以工具结果写进历史的顺序 = 模型发起调用的顺序,和谁先跑完无关。
  • needs_follow_up |= ... 是"只要有一个工具调用,就置为 true"。
  • drain_in_flight( turn.rs:2472-2500 )就是 while let Some(res) = in_flight.next().await,逐个把结果记进历史。它在流循环之后执行( turn.rs:3162-3171 ),流中途出错也会执行,已经开跑的工具结果不会丢。

一个容易看漏的细节:FuturesOrdered 在 push_back 时并不运行 future。futures-util 0.3 源码里的文档写的是 “the future will not be polled at this point”,要等有人调 poll_next。而 drain 要等流结束才调用。那工具怎么提前开跑?答案在 core/src/tools/parallel.rs:196 :handle_tool_call 是一个普通函数(不是 async),它在返回 future 之前就 tokio::spawn 了真正干活的任务。tokio 文档说 spawn 出去的任务 “will start running in the background immediately when spawn is called”。放进 in_flight 的只是一个"等这个任务结束"的句柄。源码注释也写明了这个分工:“The sampling loop collects results in order only after its stream ends."( parallel.rs:229 )

用 Python 改写时也一样:只把协程对象存进列表,它不会跑;要用 asyncio.create_task。

断流时已经开跑的工具怎么办:照样跑完、结果进历史。采样重试时不复用原来的请求体,而是从历史重新组装( turn.rs:1656-1662 ),新请求里已经带着这些工具调用和输出,所以不会再执行一遍。

能省多少

演示的默认场景:首 token 1 秒,模型每秒吐完一个调用(50 token ÷ 50 token/s),4 个工具耗时 3、1、2、1 秒,都可并行。

写法这一步总耗时
第 2 周:等响应结束,串行执行5 + 3 + 1 + 2 + 1 = 12 秒
等响应结束,并行执行5 + max(3, 1, 2, 1) = 8 秒
Codex:一完整就执行,并行max(第 5 秒 completed,2+3, 3+1, 4+2, 5+1) = 6 秒

省的是"模型还在生成后面调用的参数"和"前面工具在执行"的重叠时间。这个收益有条件:模型输出越慢、一次调用越多、前面的工具越慢,省得越多。只有一个调用,或者模型慢、工具快(演示里的"模型慢、工具快"预设,每个工具在下一个调用吐完之前就结束了,最后一个工具仍然要在 completed 之后才跑),提前执行就几乎省不了。不是"总能省一半”。

测开视角:怎么证明"真的提前开跑了"

Codex 有一个专门的集成测试 shell_tools_start_before_response_completed_when_stream_delayed( core/tests/suite/tool_parallelism.rs:304-435 ),做法值得照搬:

  1. mock 服务器先只发 4 个 exec_command 调用,response.completed 被一个 oneshot 闸门卡住不发;
  2. 每个命令往临时文件里写一个毫秒时间戳;
  3. 等文件里攒够 4 个时间戳,才打开闸门发 completed;断言 4 个时间戳都不晚于 completed 的发送时刻。

另一个测试 tool_results_grouped(同文件 :226-302 )断言第二次请求里所有 function_call 都排在所有 function_call_output 之前,而且输出的 call_id 顺序和调用一一对应。断言的是 agent 发给模型的请求体,不是模型说了什么。这套 mock 服务器的写法在「不调真实模型,怎么测一个 Agent?」那一章展开。

并行工具:一把读写锁

组装 Prompt 时 parallel_tool_calls 固定写成 true( turn.rs:1598 ),真正发出去的值是 prompt.parallel_tool_calls && !model_info.use_responses_lite( core/src/client.rs:1001 ):除了走 Responses Lite 的模型,请求里都是 true,模型可以一次发多个调用。但不是所有工具都能同时跑。每次 run_sampling_request 会创建一个 ToolCallRuntime( turn.rs:1638 ),里面有一把 tokio::sync::RwLock<()>( parallel.rs:45-64 )。锁的类型参数是 (),不保护任何数据,只用来排队:

// core/src/tools/parallel.rs:205-209(在 tokio::spawn 出去的任务里)
let guard = if supports_parallel {
    Either::Left(lock.read().await)
} else {
    Either::Right(lock.write().await)
};
  • supports_parallel 来自工具的 supports_parallel_tool_calls() 声明。为 true 拿读锁,多个读锁可以同时持有,这些工具一起跑。
  • 否则拿写锁:要等前面的工具都结束才开始,它跑的时候别人都得等。
  • Either::Left / Right 只是把两种锁守卫装进同一个变量;guard 离开作用域(工具跑完后 drop(guard))时锁释放。
  • 这个方法的默认实现返回 false( tools/src/tool_executor.rs:122-124 ),不声明就独占,默认是保守的。
声明可并行(读锁)没声明(写锁,独占)
exec_command( exec_command.rs:142-144 )、write_stdin、view_image、tool_search、三个 MCP resource 工具;MCP 工具在 server 声明支持并行、或带 read_only_hint 时( handlers/mcp.rs:148-159 )apply_patch、update_plan 等没有覆盖默认实现的工具

两点要注意:

  1. 读锁不等于只读。exec_command 跑 shell,完全可以写文件,但它拿的是读锁。读/写只是锁的模式,意思是"能不能和其他调用同时跑"。shell 命令的安全靠审批和沙箱(见「Agent 要执行 rm -rf,谁来拦?」那一章),不是靠这把锁。
  2. tokio 的 RwLock 是公平的。tokio 文档:等锁的任务按先进先出排队,“if a task that wishes to acquire the write lock is at the head of the queue, read locks will not be given out until the write lock has been released”。所以"读、读、写、读"的顺序里,第 4 个读者要等第 3 个写者跑完。演示里把 c3 换成 apply_patch,总耗时从 6 秒变成 8 秒:c3 在第 4 秒完整,要等 c1 在第 5 秒结束才拿到写锁,跑到第 7 秒;c4 在第 5 秒完整,排在 c3 后面,第 7 秒才开始,第 8 秒结束。

重试:分类、退避、降级

第 2 周的 loop 里,请求失败就是一个异常。Codex 把失败当常态,分两层重试,另外对"连不上"单独处理:

层默认次数重试什么等多久
HTTP 请求层( codex-client/src/retry.rs:22-52 )request_max_retries = 4,上限 1005xx、超时、网络错误;不重试 429(retry_429: false, model-provider-info/src/lib.rs:447-453 )有 Retry-After 用它;否则 200ms × 2n−1,乘 [0.9, 1.1) 的随机抖动
采样层( core/src/responses_retry.rs:57-176 )stream_max_retries = 5,上限 100按 CodexErr::retry_delay 的分类有服务端建议用它;否则 200ms × 2n−1 × [0.9, 1.1)
连接失败(同一文件 :93-118 )不限次数,不占上面的 5 次ConnectionFailed,且是普通会话、非 Bedrock。它只来自 HTTP 传输:连不上时 HTTP 请求层先按第一行重试 4 次(等 0.2、0.4、0.8、1.6 秒左右,retry_transport: true),仍失败才映射成 ConnectionFailed( codex-api/src/api_bridge.rs:263-264 )5 秒起,每次翻倍,封顶 60 秒;每轮之前还有 HTTP 层那约 3 秒

常量出处:

  • DEFAULT_STREAM_IDLE_TIMEOUT_MS = 300_000、DEFAULT_STREAM_MAX_RETRIES = 5、DEFAULT_REQUEST_MAX_RETRIES = 4,上限都是 100( model-provider-info/src/lib.rs:63-72 )。流 300 秒收不到事件,按 idle timeout waiting for SSE 的流错误处理。
  • 采样层退避 INITIAL_DELAY_MS = 200、BACKOFF_FACTOR = 2.0,抖动 random_range(0.9..1.1)( async-utils/src/backoff.rs:7-17 )。注意它和 HTTP 层用的是两个不同的 backoff 函数,参数恰好一样。
  • INITIAL_CONNECTION_RETRY_DELAY = 5s、MAX_CONNECTION_RETRY_DELAY = 60s( responses_retry.rs:23-24 )。连接无限重试由 feature flag unbounded_connection_retries 控制,Stable、默认开( features/src/lib.rs:1330-1334 )。
  • WebSocket 握手失败不会变成 ConnectionFailed:map_ws_error 把 IO 错误映射成 TransportError::Network( codex-api/src/endpoint/responses_websocket.rs:566-593 ),再变成普通的流错误,按采样层计次、用完后降级到 HTTP;握手超过 15 秒(DEFAULT_WEBSOCKET_CONNECT_TIMEOUT_MS, model-provider-info/src/lib.rs:68 )算超时,同样计次。降级到 HTTP 之后还连不上,才进入上表第三行的无限重试。

采样层 5 次重试的名义等待是 0.2 + 0.4 + 0.8 + 1.6 + 3.2 = 6.2 秒,加上抖动在 [5.58, 6.82) 秒之间。退避没有封顶:把 stream_max_retries 调到 10,第 10 次要等 200ms × 29 ≈ 102 秒。

错误分类

CodexErr::retry_delay( protocol/src/error.rs:389-437 )返回 None 表示终止,返回 Some(时长) 表示可以重试:

返回错误含义
None(终止)ContextWindowExceeded、UsageLimitReached、QuotaExceeded、InvalidRequest、Sandbox、Interrupted、TurnAborted、Fatal 等重试也不会好,直接报错
只听服务端ServerOverloaded、RetryLimit服务端给了 Retry-After 才重试,没给就终止
服务端建议或退避Stream(断流、空闲超时)、RateLimitExceeded、Timeout、InternalServerError、ConnectionFailed、UnexpectedStatus、ContentFilter 等可重试

429 不是一个错误,是好几个( codex-api/src/api_bridge.rs:181-240 ):

  • HTTP 429,响应体类型是 usage_limit_reached → UsageLimitReached,终止;
  • 响应体是 insufficient_quota 等额度类 → QuotaExceeded,终止;
  • 其他 429 → RetryLimit,只有带 Retry-After 时才重试;
  • 流里的 response.failed 事件,错误码 rate_limit_exceeded → RateLimitExceeded,可重试,等待时间从错误消息里的 “Please try again in 11.054s” 解析( sse/responses_error.rs:82-88 )。

降级:WebSocket → HTTP

内置的 OpenAI provider 默认走 WebSocket(supports_websockets: true, model-provider-info/src/lib.rs:555 )。采样层的处理顺序是:

// core/src/responses_retry.rs:87-90
let retry_count = retry_state.retries.saturating_add(1);
let Some(delay) = err.retry_delay(retry_count) else {
    return Err(err);
};
// ... 连接失败的无限重试分支 ...
// core/src/responses_retry.rs:120-125, :137-138
if retry_state.retries >= max_retries
    && client_session.try_switch_fallback_transport(
        &turn_context.session_telemetry,
        turn_context.model_info(),
    )
{
    // ...
    retry_state.retries = 0;
    return Ok(());
}
  • 先问 retry_delay:let Some(delay) = ... else { return Err(err) } 的意思是"拿不到延迟就直接把错误返回",终止类错误在这里就出去了。
  • 重试次数用完、而且还能降级时,切到 HTTP SSE,计数清零,相当于 HTTP 再给一轮完整的重试次数。切换本身不等退避;但服务端给过 Retry-After 时仍要等到那个时刻,源码注释:“Changing transport must not bypass the server’s retry deadline.”
  • 降级是粘性的:force_http_fallback 用 swap(true) 置位( core/src/client.rs:654-673 ),整个会话只发生一次。

所以默认配置下断流一直不恢复:WebSocket 1 次 + 重试 5 次,第 6 次失败时降级,HTTP 再 1 次 + 重试 5 次,一共 12 次请求后 turn 以错误结束。演示 2 把"连续失败次数"拉到 12 以上就能看到。

对应的测试

  • core/tests/suite/websocket_fallback.rs:85-128 :stream_max_retries = 2 时,WebSocket 尝试 1 + 2 次后降级,HTTP 请求 1 次成功(断言里 WebSocket 是 4 次,多出的 1 次是启动时的预热连接,测试注释写明了)。
  • core/tests/suite/retry_after.rs:344-400 :429 带 Retry-After: 1,断言实际等了至少 1 秒、一共 2 次请求、turn 正常完成。
  • retry_after.rs:502-506 :503 server_is_overloaded、request_max_retries = 2。没有 Retry-After 时 3 次请求就结束;带 Retry-After: 0 且 stream_max_retries = 2 时是 3 × 3 = 9 次。两层重试是相乘的。把演示 2 的参数调成一样(演示里服务端建议等待 0 表示"没给",复现第二个数字时调到 1 秒),能复现这两个数字。

用户中途插话

模型正在跑,用户又发了一句"别改测试文件"。Codex 不开新 turn,而是把这句话塞进当前 turn。 core/src/session/turn_input.rs 顶部注释说,这里决定一条输入是 “starts a turn, steers an active turn, or is rejected”:

  1. steer_input 把消息放进当前 turn 的 pending input 队列( turn_input.rs:734-740 )。Review 和 Compact 任务不接受插话( :683-695 )。
  2. 外层循环每一圈开头取出 pending input,作为用户消息写进历史( turn.rs:426-447 )。
  3. 采样结束后 needs_follow_up = model_needs_follow_up || has_pending_input:哪怕模型已经说完了,只要有插话,就再采样一次。

默认行为是:不取消当前流,模型要到下一次请求才看到插话;但 wait_agent、sleep 这类等待型工具会被新输入提前结束(见下面的测试和同文件的 any_new_input_interrupts_sleep, pending_input.rs:657-769 ),普通工具照常跑完。想更快打断,有一个还在开发中的 instant_interrupt(Under Development、默认关, features/src/lib.rs:1120-1124 ):开启后每次采样会监听插话队列,一有用户输入就取消当前流( session/input_queue.rs:253-273 )。用户在 TUI 里按 Esc( tui/src/keymap.rs:1662 把 Esc 绑定到 interrupt_turn)是另一条路:CancellationToken 一路传下去,turn 以 TurnAborted 结束。

测开视角:插话是一个经典的竞态场景,要测"插话在哪个时刻到达":采样中、工具执行中、Stop hook 判定前后。Codex 的 core/tests/suite/pending_input.rs:593 steer_interrupts_wait_agent_and_is_sent_in_follow_up_request 就是一例:模型调了一个等待子 agent 的工具,等待期间插话,断言工具输出是 “Wait interrupted by new input.",并且第二次请求里依次有原始提问和插话。

没有 max_turns:取舍和无人值守

在 core/src/session 和 core/src/tasks 里 grep max_turns|max_steps|max_iterations,没有结果,turn 内的步数没有上限。token 预算类的开关 token_budget 和 rollout_budget 都是 Under Development、默认关( features/src/lib.rs:1760-1764 、 :1772-1776 )。本机 codex-cli 0.159.2 跑 codex features list(不调模型),输出里这几行是:

instant_interrupt                        under development  false
rollout_budget                           under development  false
token_budget                             under development  false
unbounded_connection_retries             stable             true

本机 CLI 版本和这个 commit 不一定一致,但这几个开关的阶段和默认值和源码相同。

它靠什么防失控?

  • 上下文压缩:token 快满时在 turn 中途压缩然后继续。源码注释直接写了这个判断:“as long as compaction works well in getting us way below the token limit, we shouldn’t worry about being in an infinite loop."( turn.rs:612 )
  • 人在旁边:交互式编码 agent,用户随时可以 Esc 或者插话纠偏。
  • 终止类错误:上下文超窗、额度用完等直接结束 turn;采样重试有次数上限(连接失败除外)。

把第 2 周的五个停止条件拿来对照:

第 2 周的停止条件Codex 默认配置下无人值守(CI、批量评测)要不要
模型 end_turn有:needs_follow_up == false要
max_turns没有要:没人按 Esc,模型绕圈会一直跑
token / 成本预算没有(开关默认关)要:批量跑 N 个任务,成本上限要可控
超时turn 级没有墙钟超时,只有流空闲超时 300 秒、WebSocket 握手 15 秒和工具自己的超时等局部超时;连接失败无限重试要:加墙钟超时,否则断网时会一直等
重复调用检测没有建议要,至少在 trace 里统计

结论:Codex 的取舍对"人坐在终端前"的场景是合理的,压缩让长任务能继续,人负责喊停。BugHunt-Bench 这种无人值守的批量评测没有人喊停,要在自己的 harness 里保留 max_turns、墙钟超时和 token 预算;用 codex exec 跑批量任务,也要在外面套一层超时。源码里也有类似的意识:Stop hook 拒绝时,记忆整理这种内部无人值守任务直接报错退出,注释是 “Do not feed managed rejections back into an unattended memory loop."( turn.rs:667 )

动手实验:把第 2 周的 loop 改成"流式 + 提前执行 + 按序回填”

下面是一个纯 Python 标准库的模拟,不调任何 API。模型是一个按脚本吐事件的假流,时间用 asyncio.sleep 缩放。它做四件事:流式读事件,output_item.done 就用 asyncio.create_task 开跑;用一个公平读写锁控制并行;流结束后按调用顺序 await,结果按序回填;采样失败按 Codex 的规则分类、退避、降级。

"""第 3 周 · 生产级 agent loop 动手实验(不调任何 API)。"""
import asyncio
import collections
import random
import time

SCALE = 0.5           # 0.5 真实秒 = 1 模拟秒
TTFT = 1.0            # 首个事件前的等待(模拟秒)
SECS_PER_CALL = 1.0   # 每个工具调用的参数生成耗时:50 token ÷ 50 token/s
PARALLEL_SAFE = {"exec_command": True, "apply_patch": False}  # 对应各 handler 的 supports_parallel_tool_calls()


class Clock:
    def __init__(self):
        self.t0 = time.monotonic()

    def now(self):
        return (time.monotonic() - self.t0) / SCALE


async def sleep(sim_secs):
    await asyncio.sleep(sim_secs * SCALE)


async def fake_stream(calls):
    """假模型:先等 TTFT,然后每 SECS_PER_CALL 吐完一个工具调用,最后 response.completed。"""
    await sleep(TTFT)
    for call_id, name, secs in calls:
        await sleep(SECS_PER_CALL)
        yield {"type": "response.output_item.done",
               "item": {"type": "function_call", "call_id": call_id, "name": name, "secs": secs}}
    yield {"type": "response.completed", "end_turn": None}


class FairRWLock:
    """先来先得的读写锁:队头是写者时,排在它后面的读者也要等(tokio RwLock 的公平策略)。"""

    def __init__(self):
        self.readers, self.writer, self.queue = 0, False, collections.deque()

    def _wake(self):
        while self.queue:
            fut, write = self.queue[0]
            if write:
                if self.readers == 0 and not self.writer:
                    self.queue.popleft()
                    self.writer = True
                    fut.set_result(None)
                return
            if self.writer:
                return
            self.queue.popleft()
            self.readers += 1
            fut.set_result(None)

    async def acquire(self, write):
        fut = asyncio.get_running_loop().create_future()
        self.queue.append((fut, write))
        self._wake()
        await fut

    def release(self, write):
        if write:
            self.writer = False
        else:
            self.readers -= 1
        self._wake()


async def run_tool(item, lock, clock, log):
    write = not PARALLEL_SAFE[item["name"]]
    await lock.acquire(write)
    start = clock.now()
    await sleep(item["secs"])          # 工具真正干活
    end = clock.now()
    lock.release(write)
    log[item["call_id"]].update(start=start, end=end)
    return {"type": "function_call_output", "call_id": item["call_id"]}


async def codex_style(calls):
    """Codex 写法:output_item.done 就 create_task 开跑;流结束后按顺序 await,结果按调用顺序回填。"""
    clock, lock, history = Clock(), FairRWLock(), []
    log = collections.defaultdict(dict)
    in_flight = []                     # 相当于 FuturesOrdered:先进先出
    needs_follow_up = False
    async for ev in fake_stream(calls):
        if ev["type"] == "response.output_item.done":
            item = ev["item"]
            log[item["call_id"]]["done"] = clock.now()
            history.append(item)       # 调用本身先进历史
            in_flight.append(asyncio.create_task(run_tool(item, lock, clock, log)))
            needs_follow_up = True     # 出现过工具调用,就还要再采样一次
        elif ev["type"] == "response.completed":
            if ev["end_turn"] is False:
                needs_follow_up = True
    completed = clock.now()
    for task in in_flight:             # drain_in_flight:按 push 顺序取结果
        out = await task
        history.append(out)
        log[out["call_id"]]["recorded"] = clock.now()
    return completed, clock.now(), log, history, needs_follow_up


# ---------------- 重试:错误分类 + 退避 + 传输降级(常量取自 Codex 源码) ----------------
INITIAL_DELAY_MS = 200                 # async-utils/src/backoff.rs:7
BACKOFF_FACTOR = 2.0                   # async-utils/src/backoff.rs:8
DEFAULT_STREAM_MAX_RETRIES = 5         # model-provider-info/src/lib.rs:64
CONN_DELAY0, CONN_DELAY_MAX = 5.0, 60.0  # core/src/responses_retry.rs:23-24
TERMINAL = {"context_window_exceeded", "usage_limit_reached", "invalid_request", "quota_exceeded"}
SERVER_DELAY_ONLY = {"http_429", "server_overloaded"}   # 只有服务端给了 Retry-After 才重试


def backoff_secs(attempt, rng):
    base = int(INITIAL_DELAY_MS * BACKOFF_FACTOR ** max(attempt - 1, 0))
    return int(base * rng.uniform(0.9, 1.1)) / 1000


def retry_delay(kind, retry_after, attempt, rng):
    """对应 CodexErr::retry_delay:None 表示终止,不再重试。"""
    if kind in TERMINAL:
        return None
    if kind in SERVER_DELAY_ONLY:
        return retry_after
    return retry_after if retry_after is not None else backoff_secs(attempt, rng)

完整文件还有第 2 周写法的对照函数 week2_style、虚拟时钟上的重试循环 sample_with_retry 和打印逻辑,共 230 行左右,放在学习目录的 week03_Codex源码/code/stream_loop_sim.py。本机 python3 stream_loop_sim.py 的输出(节选,6 个重试场景选了 3 个;单位是模拟秒;真实 sleep 有几毫秒调度开销,偶尔会比理论值多 0.1):

A. 第 2 周写法:等响应结束,串行执行
  调用  工具           收到   开始   结束   回填
  c1    exec_command    2.0    5.0    8.0    8.0
  c2    exec_command    3.0    8.0    9.0    9.0
  c3    exec_command    4.0    9.0   11.0   11.0
  c4    exec_command    5.0   11.0   12.0   12.0
  response.completed 在 5.0 秒,这一步总耗时 12.0 秒

B. Codex 写法:工具一完整就开跑,全部可并行
  调用  工具           收到   开始   结束   回填
  c1    exec_command    2.0    2.0    5.0    5.0
  c2    exec_command    3.0    3.0    4.0    5.0
  c3    exec_command    4.0    4.0    6.0    6.0
  c4    exec_command    5.0    5.0    6.0    6.0
  response.completed 在 5.0 秒,这一步总耗时 6.0 秒
  回填顺序 ['c1', 'c2', 'c3', 'c4'],needs_follow_up = True

C. Codex 写法:c3 换成 apply_patch(写锁,独占)
  调用  工具           收到   开始   结束   回填
  c1    exec_command    2.0    2.0    5.0    5.0
  c2    exec_command    3.0    3.0    4.0    5.0
  c3    apply_patch     4.0    5.0    7.0    7.0
  c4    exec_command    5.0    7.0    8.0    8.0
  response.completed 在 5.0 秒,这一步总耗时 8.0 秒

重试模拟(虚拟时钟,抖动用固定随机种子)

[断流 7 次:5 次重试用完后降级到 HTTP]
  t=  0.00s  第 1 次请求(WebSocket)失败:stream_closed
             重试 1/5,等 0.192s
  t=  0.19s  第 2 次请求(WebSocket)失败:stream_closed
             重试 2/5,等 0.372s
  t=  0.56s  第 3 次请求(WebSocket)失败:stream_closed
             重试 3/5,等 0.824s
  t=  1.39s  第 4 次请求(WebSocket)失败:stream_closed
             重试 4/5,等 1.463s
  t=  2.85s  第 5 次请求(WebSocket)失败:stream_closed
             重试 5/5,等 3.222s
  t=  6.07s  第 6 次请求(WebSocket)失败:stream_closed
             重试用完,WebSocket 降级到 HTTP,计数清零
  t=  6.07s  第 7 次请求(HTTP)失败:stream_closed
             重试 1/5,等 0.182s
  t=  6.26s  第 8 次请求(HTTP)成功

[HTTP 429,没有 Retry-After]
  t=  0.00s  第 1 次请求(HTTP)失败:http_429
             retry_delay 返回 None:终止,不重试,turn 以错误结束

[HTTP 连不上服务器:3 轮采样都失败]
  t=  0.00s  第 1 次请求(HTTP)失败:connection_failed
             HTTP 层先重试 4 次(等 0.186 + 0.412 + 0.731 + 1.611 = 2.94s),仍连不上
             采样层收到 ConnectionFailed:等 5s,不占重试次数
  t=  7.94s  第 2 次请求(HTTP)失败:connection_failed
             HTTP 层先重试 4 次(等 0.182 + 0.400 + 0.725 + 1.578 = 2.88s),仍连不上
             采样层收到 ConnectionFailed:等 10s,不占重试次数
  t= 20.82s  第 3 次请求(HTTP)失败:connection_failed
             HTTP 层先重试 4 次(等 0.183 + 0.393 + 0.852 + 1.479 = 2.91s),仍连不上
             采样层收到 ConnectionFailed:等 20s,不占重试次数
  t= 43.73s  第 4 次请求(HTTP)成功

B 里的 c2 第 4 秒就跑完了,但要等 c1 在第 5 秒回填之后才轮到它,这就是按序回填。C 里的 c4 是读者,却要等排在它前面的写者 c3 跑完,这就是公平读写锁。三组数字和演示 1 的读数一致:A 是默认参数下「第 2 周写法」那一格,B、C 分别是「默认」和「c3 换成 apply_patch」两个预设下的 Codex 写法。

这个模拟有几处简化:HTTP 层只模拟了连接错误的重试,5xx 重试没有放进 sample_with_retry(互动演示 2 里有);连不上只演示了纯 HTTP,WebSocket 握手失败按断流计次的那段没放进来;asyncio 是单线程的,工具在这里是 sleep,真实的阻塞调用要放到线程池里;读写锁只实现了最基本的排队,没有处理取消。

测开视角:这一章能落到哪些用例

机制用例断言
工具提前执行mock 流把 completed 卡住,工具写时间戳工具开始时间早于 completed
按序回填让后发起的工具先结束第二次请求里输出顺序和调用顺序一致
读写锁一个独占工具夹在两个可并行工具中间独占工具执行期间没有其他工具在跑
断流重试流在第 2 个工具调用后断开已开跑的工具只执行一次;重试请求里带着它的输出
退避连续断流 k 次请求次数、每次间隔落在 [0.9, 1.1) × 名义值内
Retry-After429 带 / 不带 Retry-After带:等待 ≥ 指定秒数后成功;不带:不重试
降级WebSocket 一直断第 stream_max_retries + 1 次失败后改走 HTTP,计数清零
插话在采样中、工具执行中各插一次插话出现在下一次请求里,turn 没有提前结束
无人值守兜底假模型永远调工具你自己的 harness 在 max_turns 处停下

这些用例的共同点:都不依赖真实模型,用一个可编排的假 Responses 服务器就能确定地跑出来。

常见错误说法

  • “Codex 的停止条件就是没有工具调用”:还有 end_turn == false 和用户插话两个来源,没有工具调用但有插话,turn 会继续。
  • “把 future 放进 FuturesOrdered,工具就开始跑了”:FuturesOrdered 不会在 push_back 时运行 future。Codex 能提前执行,是因为 handle_tool_call 里先 tokio::spawn 了任务。
  • “工具提前执行总能把时间省一半”:收益取决于模型输出速度、调用个数和工具耗时。只有一个调用,或者模型生成很快、工具很慢(重叠窗口很短)时,提前执行几乎不省;工具独占只会让工具之间不能并行,仍能和模型生成重叠。演示 1 默认参数下把 4 个工具都换成 apply_patch,等响应结束再执行要 12.0 秒,提前执行 9.0 秒,省 3.0 秒,比全部可并行时省的 2.0 秒(8.0→6.0)还多。
  • “拿读锁的工具都是只读的”:exec_command 拿读锁,但它跑 shell,照样能写文件。读/写只是"能不能和别人同时跑”。
  • “遇到 429 就退避重试”:HTTP 429 没有 Retry-After 时 Codex 不重试,额度类 429 直接终止;只有流里的 rate_limit_exceeded 才走退避。
  • “Codex 最多重试 5 次”:采样层默认 5 次,用完后从 WebSocket 降级到 HTTP 还有一轮;HTTP 层对 5xx 和连接错误另有 4 次,两层相乘;HTTP 连接失败在这 4 次之后进入不占次数、不限次数的连接重试。
  • “生产级 agent 都没有 max_turns,所以我们也不需要”:Codex 是交互式的,人负责喊停。无人值守的批量评测要保留 max_turns、墙钟超时和 token 预算。

下一章:工具系统与 MCP client。RespondToModel 和 Fatal 的边界在哪,apply_patch 为什么用语法约束而不是 JSON 参数,MCP 工具名冲突怎么处理。