夜雨聆风学习资料网

ARTICLE · 1121436

Codex 源码-Agent Loop / Turn 执行循环

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)"就是为此存在的。

整条预处理管线可以浓缩成一张表(后面几讲会逐个展开):

预处理步骤
关键函数(turn.rs)
失败时的行为
Pre-sampling compact(L183)
run_pre_sampling_compact → run_auto_compact(PreTurn)
先记用户输入防丢失,广播 Error,Ok(None) 结束
MCP 需求解析(L234)
required_mcp_servers_for_input + capture_step_context_with_required_mcp_servers
TurnAborted:记输入后上抛;其余错误直接 Err
Skills/插件注入(L308)
build_skills_and_plugins → record_conversation_items
返回 None(取消/失败)→ Ok(None) 结束本轮
Hook 接受检查(L366)
run_hooks_and_record_inputs(PersistContext::TurnStart)
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 收集结果。这就是"并发执行、顺序落账"的完整闭环。

把内层循环里所有事件的处理方式汇总一下:

ResponseEvent
核心动作(turn.rs)
对 needs_follow_up 的影响
OutputItemAdded(L2650)
开始流式预览:emit item-started,为工具参数建 diff consumer
无直接影响(等 Done)
OutputTextDelta / Reasoning*Delta(L2833+)
增量文本/推理 → send_event 推给客户端实时渲染
无
OutputItemDone(L2538)
handle_output_item_done:按"工具调用 vs 消息"三分支分发
工具调用 ⇒ true;RespondToModel ⇒ true
Completed { end_turn }(L2783)
记 token usage、flush 流式段、record response_id,break 出循环
end_turn=false ⇒ 强制 true(服务端要求续跑)

七、外层决策:采样之后,继续还是收工?

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 的所有退出/继续路径汇总成一张表:

决策条件
触发来源
loop 行为
模型回纯消息且无 pending input(L642)
needs_follow_up=false
跑 Stop hook → break,返回 last_agent_message
Stop hook should_block 且带 prompt(L661)
hook 引擎
记录 hook_prompt,stop_hook_active=true,continue
token 到限 + 需要续跑(L612)
context_window_token_status / 服务端 end_turn=false
run_auto_compact(MidTurn) → continue
ContextWindowExceeded(Guardian 预算,L706)
采样错误
压缩一次重试;guardian_budget_compacted 保证不循环
TurnAborted(取消令牌)
用户中断
return Err,不当普通错误广播

八、全景:一个 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技术推荐官」获取更多源码解析内容

相关学习资料