第 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 周的 loop | Codex | 差别 |
|---|---|---|
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
),做法值得照搬:
- mock 服务器先只发 4 个
exec_command调用,response.completed被一个 oneshot 闸门卡住不发; - 每个命令往临时文件里写一个毫秒时间戳;
- 等文件里攒够 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 等没有覆盖默认实现的工具 |
两点要注意:
- 读锁不等于只读。
exec_command跑 shell,完全可以写文件,但它拿的是读锁。读/写只是锁的模式,意思是"能不能和其他调用同时跑"。shell 命令的安全靠审批和沙箱(见「Agent 要执行 rm -rf,谁来拦?」那一章),不是靠这把锁。 - 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,上限 100 | 5xx、超时、网络错误;不重试 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 flagunbounded_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:503server_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”:
steer_input把消息放进当前 turn 的 pending input 队列(turn_input.rs:734-740)。Review 和 Compact 任务不接受插话(:683-695)。- 外层循环每一圈开头取出 pending input,作为用户消息写进历史(
turn.rs:426-447)。 - 采样结束后
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-After | 429 带 / 不带 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 工具名冲突怎么处理。