ARTICLE · 1121436
Codex 源码-Agent Loop / Turn 执行循环
Codex 源码解析系列
第 10 讲:Agent Loop / Turn 执行循环
基于 OpenAI Codex 源码 · 2026-10-04
💡 本讲一句话:Codex 的 agent loop 不是"while 循环里调模型"这么简单——它是一个三层嵌套的状态机:外层 run_turn 管"要不要再来一轮",中层 run_sampling_request 管"这次请求失败能不能重试",内层 SSE 事件循环管"流里每个 item 怎么分发"。工具调用被做成挂起的 future按序收集,用户中途插话(steer)有专门的排空时机控制。读完你会明白一个生产级 agent loop 的骨架长什么样。
一、先建立地图:一个 turn 里到底有几层循环
前几讲我们把 Session / TurnContext / StepContext(第 8 讲)和系统提示词组装(第 9 讲)拆完了。这一讲进入 agent loop 本体:用户输入进来之后,模型、工具、harness 是怎么一轮轮"对话"下去的?
源码里这条链路横跨两个文件:主循环在 core/src/session/turn.rs(3039 行),流事件分发在 core/src/stream_events_utils.rs。骨架是三层嵌套:
🔹 外层 run_turn(L163):turn 级循环——每轮采样后决定"继续 / 压缩后再来 / Stop hook 拦截 / 结束"🔹 中层 run_sampling_request(L1538):请求级重试状态机——流断了、可恢复错误,按 provider 配置的重试上限恢复🔹 内层 try_run_sampling_request(L2409):SSE 事件消费循环——逐事件分发,工具调用挂成 future,直到 response.completed
先看外层入口。run_turn 的文档注释把整个 agent loop 的契约写得很直白:
📄 codex-rs/core/src/session/turn.rs (第 149-170 行)
/// Takes initial turn input and runs a loop where, at each sampling request, // 文档注释:run_turn 拿到本轮输入后进入循环,每次采样请求模型要么要工具、要么回消息/// the model replies with either: // ——回复只有两种形态/// - requested function calls // 形态一:函数调用(harness 执行完把结果喂回去)/// - an assistant message // 形态二:助手消息(记入历史,本轮结束)pub(crate) async fn run_turn( // Turn 的执行入口:一个 turn = "用户输入 → 模型/工具多轮往返 → 最终回复" sess: Arc<Session>, // 会话本体(共享状态、事件总线、input_queue),Arc 因为要跨任务传递 turn_context: Arc<TurnContext>, // 本轮上下文快照:配置、审批策略、provider、sub_id mut input: Vec<TurnInput>, // 用户输入列表(文本/图片/@提及等);mut 因为 guardian 路径会改写它 mcp_startup_requirements: &mut McpStartupRequirements, // MCP 启动需求跨 turn 保留:显式 @ 的 server/plugin 必须拉起 prewarmed_client_session: Option<ModelClientSession>, // 预热好的模型会话(缓存 WebSocket/粘性路由),没有就新建 cancellation_token: CancellationToken, // 取消令牌:用户中断时所有 await 点都能快速退出) -> CodexResult<Option<String>> { // 返回本轮最后一条助手消息;None = 被 hook 拦截或出错提前结束为什么这样设计:注意返回值是 CodexResult<Option<String>>——"出错"和"正常结束但没有最终消息"是两个不同的信号。hook 拦截、压缩失败这类情况返回 Ok(None):turn 体面地结束了,用户还能继续对话;而取消(TurnAborted)才走 Err。这个区分贯穿整个函数——后面每个失败分支你都会看到"先记输入、再广播错误、最后 Ok(None)"的固定套路。
二、进主循环之前:turn 预处理管线
run_turn 开头有约 250 行"进场检查",顺序是:guardian 待审输入 → 异步 hook 结果落账 → pre-sampling compact → MCP 需求解析 → StepContext 捕获 → skills/plugins 注入 → session start hooks → 用户输入过 hook 接受。其中最有设计含量的是压缩失败的处理——因为压缩发生在"新输入落库之前",失败时第一件事不是报错,而是先把用户输入记下来:
📄 codex-rs/core/src/session/turn.rs (第 177-222 行,节选)
let mut client_session = // 复用预热会话,否则新建——WebSocket + sticky routing 状态是 turn 级缓存 prewarmed_client_session.unwrap_or_else(|| sess.services.model_client.new_session()); // 有预热会话就复用,没有才新建——省掉 WebSocket 握手开销if let Err(err) = run_pre_sampling_compact( // 采样前先压缩:历史超预算就先 compact,避免第一轮请求直接爆窗口 &sess, // 传会话本体:compact 要读写历史和事件总线 &turn_context, // 传本轮上下文快照:读配置里的压缩阈值/模型信息 &mut client_session, // 可变借用:remote compact 可能切换底层连接状态 &cancellation_token, // 取消令牌:用户中断时 compact 立即停止).await{ // ——压缩失败分支开始 // Compaction runs before the new input is recorded, so preserve it on every failure. // 压缩发生在"新输入落库之前",失败时必须先把用户输入记下来,否则丢消息 run_hooks_and_record_inputs(&sess, &turn_context, &turn_context.capture_current_model_info(), &input, PersistContext::Standard).await; if matches!(err.details(), CodexErrorDetails::TurnAborted) { // 取消类错误直接上抛,不当普通错误广播 return Err(err); } let error = err.to_codex_protocol_error(); // 其余失败:先走 turn 错误生命周期(遥测/状态机) sess.emit_turn_error_lifecycle(turn_context.as_ref(), error.clone()).await; sess.send_event(turn_context.as_ref(), EventMsg::Error(err.to_error_event(message_prefix))).await; // 再广播 Error 事件——顺序保证 hook 先跑完,防止客户端抢着往"还在保 prompt 的 turn"里塞新输入 return Ok(None); // 本轮到此为止:无最终消息,但用户可继续对话}为什么这样设计:源码注释点破了关键约束——"Publish the failure only after prompt hooks finish, so clients cannot react to an error by steering follow-up input into a turn still preserving its prompt."(错误事件必须在 hook 跑完之后才发布,防止客户端抢着塞新输入)。这是一个竞态防护:如果先广播错误、hook 还没把用户输入落库,客户端立刻发下一条消息,两条输入的顺序就可能错乱。固定套路"记输入 → 生命周期 → 广播 → Ok(None)"就是为此存在的。
整条预处理管线可以浓缩成一张表(后面几讲会逐个展开):
三、主循环骨架:pending input 的排空时机控制
预处理通过后进入 agent loop。这里有一个非常微妙的状态变量 can_drain_pending_input——它控制"用户中途塞进来的新消息(steer)什么时候允许进入历史"。源码注释解释了两个必须延迟排空的时机:
📄 codex-rs/core/src/session/turn.rs (第 406-434 行)
let mut last_agent_message: Option<String> = None; // 记录本轮最后一条助手消息,作为 run_turn 的返回值let mut stop_hook_active: bool = false; // Stop hook 是否已生效过——防止"hook 要求继续 → 再触发 hook"死循环let turn_diff_tracker = Arc::new(tokio::sync::Mutex::new( // TurnDiffTracker:跟踪本轮文件改动供 UI 展示 diff;用户视角一个 turn 内共享 TurnDiffTracker::with_environment_display_roots(display_roots), // 用各环境的展示根目录初始化 diff 跟踪器)); // ——Mutex 包裹:工具执行与 UI 读取并发访问同一 tracker// `ModelClientSession` is turn-scoped and caches WebSocket + sticky routing state, so we reuse // 客户端会话是 turn 级缓存,重试间复用同一实例// one instance across retries within this turn. // ——所以整个 loop 里 client_session 不重建let mut next_step_context = Some(first_step_context); // 下一步的 StepContext(工具/权限/MCP 快照);Some=复用,None=重新捕获let mut guardian_budget_compacted = false; // Guardian 预算压缩只允许一次——防止无效压缩无限循环loop { // ===== agent loop 主循环:采样 → 分发输出 → 决策是否再来一轮 ===== let pending_input = if can_drain_pending_input { // 只有"允许排空"时才取队列里用户中途塞进来的新输入(steer) sess.input_queue.get_pending_input(&sess.active_turn).await.0 // 取出并清空该 turn 的待处理输入 } else { Vec::new() // 否则留队:比如 compact 之后要先让模型把工具续跑完,不能插话打断 }; if run_hooks_and_record_inputs(&sess, &turn_context, &turn_context.capture_current_model_info(), &pending_input, PersistContext::Standard).await { // 排空出来的输入同样要过 hook:拒绝则 break 结束本轮 break; }为什么这样设计:注释里写得很清楚,延迟排空有两个场景:① turn 刚开始时,新鲜的用户输入要先被采样(不能先被队列里的旧 steer 淹没);② auto-compact 之后,模型/工具的续跑必须先恢复,再处理插话。本质上是消息顺序的语义保证:steer 是"对当前进展的补充指令",如果它抢在工具结果之前进入历史,模型的上下文就乱了。一个 bool 变量 + 每轮循环顶部的条件排空,就把这个时序问题收敛了。
四、中层:采样请求与重试状态机
run_sampling_request(L1538)包在 try_run_sampling_request 外面,职责是可重试错误的恢复循环。它维护三个关键状态:重试上限、一次性消费的初始输入、以及"已执行工具调用 → output"的映射:
📄 codex-rs/core/src/session/turn.rs (第 1561-1600 行,节选)
let max_retries = turn_context.provider.info().stream_max_retries(); // 重试上限来自 provider 配置——不同后端容忍度不同let mut retry_state = ResponsesStreamRetryState::default(); // 重试状态机:记录已试过的退避/切换策略,避免重复同一动作let mut initial_input = Some(input); // 首次请求用调用方准备好的输入;Some→None 一次性消费let mut original_input = None; // 记住"原始 prompt 输入",采样结束后回写历史(保证记录的是真实发过的内容)let mut executed_tool_calls_by_output = HashMap::new(); // output → 已执行工具调用的映射:重试时把旧结果重新挂到新请求上loop { // ===== 流式请求 + 可恢复错误恢复循环 ===== turn_context.extension_data.remove::<codex_api::ResponseId>(); // 清掉上一响应的 ResponseId——重试后新工具调用不能归属到旧的 response,否则审计/回放错乱 let prompt_input = if let Some(input) = initial_input.take() { // 第一次:用现成输入 input } else { // 重试:从历史重新拼 prompt(历史里已含上一轮的工具结果) sess.clone_history().await.for_prompt(&step_context.settings.model_info.input_modalities) }; let mut prompt_input = prompt_input; sess.services.executed_tool_calls.attach_to_prompt(&mut prompt_input, &mut executed_tool_calls_by_output); // 把已执行的工具调用结果重新附着到重试请求上——模型看到的上下文不丢为什么这样设计:两个细节值得注意。第一,remove::<codex_api::ResponseId>() 在每次重试前执行——注释说"A retry must not attribute the next tool call to the previous response"(重试不能把下一个工具调用归属到上一个响应)。这是审计正确性:rollout/回放里每个工具调用必须挂在它真正所属的 response 下。第二,executed_tool_calls_by_output 这个 HashMap 解决的是"流断在工具执行之后、结果回传之前"的尴尬窗口——重试时把已经跑完的工具结果重新挂回去,避免重复执行(shell 命令重跑可能产生副作用)。
五、内层:SSE 事件消费循环
try_run_sampling_request(L2409)是真正读流的地方。它先建好"在途工具调用队列",然后进入事件循环——注意 FuturesOrdered 的选择:
📄 codex-rs/core/src/session/turn.rs (第 2458-2523 行,节选)
let mut in_flight: FuturesOrdered<InFlightFuture<'static>> = FuturesOrdered::new(); // 在途工具调用队列:FuturesOrdered 保证按发出顺序收集结果,历史顺序不乱let mut needs_follow_up = false; // 本轮采样是否还需要"下一轮请求"(有工具调用/服务端要求续跑)let outcome: CodexResult<SamplingRequestResult> = loop { // ===== SSE 事件消费主循环:一直读到 response.completed ===== let event = match stream.next().instrument(trace_span!(parent: &receiving_span, "receiving")).or_cancel(&cancellation_token).await { // or_cancel:取消令牌一触发,整个流立即中断 Ok(event) => event, Err(codex_async_utils::CancelErr::Cancelled) => { break Err(CodexErr::TurnAborted); } // 被取消 → TurnAborted,上层识别后直接返回(不当普通错误广播) }; let event = match event { // 流结束三态:正常事件 / 传输错误 / 提前断流 Some(Ok(event)) => event, Some(Err(err)) => break Err(err), // 网络/协议错误上抛,交给外层重试状态机判断是否可恢复 None => { break Err(CodexErr::Stream("stream closed before response.completed".into())); } // 服务端没发 completed 就断流 = 不完整响应,必须报错而不是静默收尾 }; sess.services.session_telemetry.record_responses(&handle_responses, &event); // 每个事件都过遥测:OTEL span 上挂 token/工具字段(第 44 讲展开)为什么这样设计:FuturesOrdered(而不是 FuturesUnordered)是刻意的:模型一次响应里可能发多个工具调用,它们并发执行但按序收集结果——因为历史里的 FunctionCallOutput 必须和 FunctionCall 一一对应、顺序一致,否则下一轮 prompt 的语义就变了。另一个细节是"提前断流"被显式建模成错误(stream closed before response.completed):SSE 协议里 response.completed 是"本轮采样完整结束"的契约,没收到它就不能假装成功——这是流式系统里典型的 fail-loud 设计。
六、OutputItemDone:工具分发的心脏
事件循环里最重要的分支是 ResponseEvent::OutputItemDone(L2538)——每个输出项"完成"时,交给 handle_output_item_done(stream_events_utils.rs L300)做三分支分发:
📄 codex-rs/core/src/stream_events_utils.rs (第 308-345 行)
match ToolRouter::build_tool_call(item.clone()) { // 每个完成的输出项先问路由:这是工具调用吗? Ok(Some(call)) => { // ——是:记日志、落库、把执行 future 挂进在途队列(不阻塞流消费) call_trace::received(ctx.sess.thread_id, &call.tool_name, &call.call_id, call_trace::Receipt::ModelTurn(&ctx.step_context.turn.sub_id)); // 调用链追踪:记录"模型第 N 轮发起了这个工具",供审计回放 ctx.sess.input_queue.accept_mailbox_delivery_for_current_turn(&ctx.sess.active_turn, &ctx.step_context.turn.sub_id).await; // 有工具要跑 = 本轮还没结束,放行 mailbox(多智能体邮件)投递 record_completed_response_item(ctx.sess.as_ref(), ctx.step_context.as_ref(), &item).await; // 立即持久化工具调用项——即使 turn 随后被取消,历史/rollout 也保持完整 let cancellation_token = ctx.cancellation_token.child_token(); // 派生子令牌:取消 turn 时正在跑的工具一起停 let tool_future: InFlightFuture<'static> = Box::pin( // 工具执行是异步 future,先挂起不 await——流里还有后续事件要处理 ctx.tool_runtime.clone().handle_tool_call(call, cancellation_token), ); output.needs_follow_up = true; // 有工具调用 ⇒ 采样结束后必须再发一轮请求把结果喂回模型 output.tool_future = Some(tool_future); } Ok(None) => { // ——不是工具:消息/推理等,走"非工具项"收尾(finalize → emit_turn_item_completed → 记历史)为什么这样设计:这个分支的精髓是"先落库、再挂 future"的顺序。注释说得很直白:"This records items immediately so history and rollout stay in sync even if the turn is later cancelled."(立即记录,保证即使 turn 随后被取消,历史和 rollout 也保持同步)。工具执行本身被 Box::pin 成 future 挂进 in_flight 队列——流消费线程不被阻塞,后续事件(比如模型继续输出下一条消息)照常处理;等 response.completed 之后才统一 await 收集结果。这就是"并发执行、顺序落账"的完整闭环。
把内层循环里所有事件的处理方式汇总一下:
七、外层决策:采样之后,继续还是收工?
SamplingRequestResult(L1788)只有两个字段:needs_follow_up + last_agent_message。外层 run_turn 拿到它之后,结合 token 状态做"窗口滚动"决策——这是 agent loop 里最容易被低估的一段:
📄 codex-rs/core/src/session/turn.rs (第 599-640 行,节选)
let should_roll_over = needs_follow_up && (sess.take_new_context_window_request().await || token_limit_reached); // 需要续跑 且(服务端要求开新窗口 或 本地 token 到限)⇒ 先压缩再续let allow_auto_compact_fallback = !should_roll_over && !token_limit_reached; // 不满足 rollover 但还有余量 ⇒ 允许"自动压缩兜底"路径if should_roll_over { // ===== 上下文窗口滚动:MidTurn 压缩,把历史压到预算内 ===== if let Err(err) = run_auto_compact(&sess, Arc::clone(&step_context), /*fallback_step_context*/ None, &mut client_session, InitialContextInjection::BeforeLastUserMessage { world_state: Arc::clone(&world_state), step_context: Arc::clone(&step_context) }, CompactionReason::ContextLimit, CompactionPhase::MidTurn).await { // 压缩失败:取消直接上抛;其余走错误生命周期后结束本轮 if matches!(err.details(), CodexErrorDetails::TurnAborted) { return Err(err); } let error = err.to_codex_protocol_error(); sess.emit_turn_error_lifecycle(turn_context.as_ref(), error.clone()).await; return Ok(None); } can_drain_pending_input = !model_needs_follow_up; // 压缩后若模型还要续跑,先让它把活干完再排空用户插话——顺序敏感 continue; // 回到 loop 顶部:用压缩后的历史重新发起采样请求}为什么这样设计:注意 should_roll_over 要求两个条件同时成立:既需要续跑、又确实到限了。如果模型已经说完了(needs_follow_up=false),哪怕 token 到限也不压缩——没有下一轮请求,压缩就是白烧一次 LLM 调用。而 can_drain_pending_input = !model_needs_follow_up 这一行再次体现了第三节的时序哲学:压缩之后,模型的续跑优先于用户的插话。源码里还有一句乐观注释:"as long as compaction works well in getting us way below the token limit, we shouldn't worry about being in an infinite loop."(只要压缩能把 token 压到远低于上限,就不必担心死循环)——这是用"压缩质量"换"循环安全"的赌注。
把整个 turn 的所有退出/继续路径汇总成一张表:
八、全景:一个 turn 的完整生命周期
一个 turn 的完整生命周期(run_turn L163-783)
① Pre-sampling compact:run_pre_sampling_compact
历史超预算先压缩(PreTurn);失败时先记用户输入再广播错误,Ok(None) 收场。
▼
② 输入落账 + hook 接受:run_hooks_and_record_inputs(TurnStart)
用户输入过 hook 检查,拒绝的进不了模型上下文;同时预热 shell snapshot。
▼
③ 采样请求:run_sampling_request → try_run_sampling_request
历史 + base instructions 拼 prompt,开 SSE 流;可恢复错误进重试状态机(≤max_retries)。
▼
④ 流事件分发:OutputItemDone → handle_output_item_done
🔹 工具调用:先落库,再挂 future 进 in_flight(FuturesOrdered 按序收集)🔹 消息/推理:finalize → emit_turn_item_completed → 记历史
▼
⑤ 采样后决策:needs_follow_up + token_status
🔹 需续跑且窗口满:run_auto_compact(MidTurn) → continue🔹 无需续跑:Stop hook 裁决 break(收工)或 continue(带 prompt 继续)
▼
⑥ 返回:Ok(last_agent_message)
最后一条助手消息交给上层(TUI/app-server)展示;None = hook 拦截或出错提前结束。
九、总结:agent loop 的三个设计内核
把 run_turn 的 780 行读完,真正值得带走的不是某个分支,而是三个反复出现的设计内核:
🔹 三层循环各管一件事:turn 级(要不要再来一轮)/ request 级(这次失败能不能重试)/ stream 级(流里每个 item 怎么分发)。每层只处理自己粒度的错误——取消在 stream 层变 TurnAborted,可恢复错误在 request 层被重试吃掉,窗口到限才升级到 turn 层触发压缩。错误不跨层乱跑🔹 "先落账、再执行"的顺序纪律:工具调用先持久化再挂 future;压缩失败先把用户输入记下来;错误事件等 hook 跑完才广播。所有顺序保证都服务于同一个目标——任何时刻崩溃/取消,历史和 rollout 都是自洽的🔹 时序敏感的 bool 状态机:can_drain_pending_input、stop_hook_active、guardian_budget_compacted 三个标志位分别管"插话时机 / hook 死循环防护 / 压缩重试上限"——agent loop 的复杂度不在算法,在这些顺序约束
下一讲我们跳出单个 turn,看 CodexThread 与 ThreadManager:多个 turn、多个会话是怎么被组织成"线程"的——run_turn 只是它调用的众多协程之一。
📚 系列导航
← 第 9 讲:System Prompt 组装与注入
→ 第 11 讲:CodexThread 与 ThreadManager
关注公众号「AI技术推荐官」获取更多源码解析内容