Turn主循环与退出条件
Codex 的一次 Turn 可以只有一条 assistant 回复,也可以经历“模型请求 → 多个工具 → 模型请求 → 压缩 → 模型请求”多轮循环。源码中的 run_turn() 负责这条主链,但它不直接发送最终 TurnComplete;取消、普通错误 和正常结束也不会走同一种返回路径。
阅读本文前,建议先掌握 StepContext模型请求 的请求快照边界,以及 Session输入队列 的 pending/steer/mailbox 语义。若需要回看端到端层次,可从 Turn端到端链路 进入;本文只深入 Core 内部循环, 不覆盖跨进程总览。
读完后,应能从任意 ResponseEvent 判断它是否会产生 follow-up,区分“当前 sampling 重试”“当前 run_turn 继续”和“RegularTask 再次调用 run_turn”,并解释为什么模型错误最终可能仍对应 TurnComplete { error: Some(...) } 而不是 TurnAborted。
1. 循环层次
一次 RegularTask 至少涉及四个嵌套边界:Task 输入清空循环、Turn sampling 循环、单次请求的 stream retry 循环,以及 response event 循环。它们复用同一个 TurnContext,但重建 StepContext、Prompt 或 stream 的频率 不同。
| 边界 | 重新创建什么 | 保留什么 | 退出信号 |
|---|---|---|---|
| RegularTask 外层 | 新的 run_turn 调用和空 next_input | 同一 TurnContext、Task token | input queue 为空 |
run_turn loop | 新 StepContext、WorldState diff、Prompt | 同一 ModelClientSession、TurnContext | 无 follow-up 且 hook 不阻止 |
run_sampling_request retry loop | 新 stream,必要时重取 history input | 同一 StepContext、ToolRouter | 成功或不可重试错误 |
try_run_sampling_request event loop | 逐个 event/item/tool future | 当前 Prompt 与 ToolCallRuntime | response.completed、流关闭、取消或错误 |
这个分层解释了一个常见误判:源码中的某个 break 通常只离开当前层,并不一定代表客户端已经看到 TurnComplete。当前实现还把 next_step_context 作为首轮快照复用点:同一 sampling 的 stream retry 保持 StepContext 不变,只有工具结果、pending input 或所需 MCP server 改变时才进入下一次捕获。
2. RegularTask
RegularTask 先发送 TurnStarted 并消费启动预热,再调用 run_turn。如果 run_turn 结束瞬间 input queue 仍有 pending input,它用空 next_input 再调用一次;这处理的是 sampling 边界外最后到达的输入。
源码位置:codex-rs/core/src/tasks/regular.rs,符号 RegularTask::run。
async fn run(
self: Arc<Self>,
sess: Arc<Session>,
ctx: Arc<TurnContext>,
input: Vec<TurnInput>,
cancellation_token: CancellationToken,
) -> SessionTaskResult {
let run_turn_span = trace_span!("run_turn");
let prewarmed_client_session = async {
let event = EventMsg::TurnStarted(TurnStartedEvent {
turn_id: ctx.sub_id.clone(),
trace_id: ctx.trace_id.clone(),
started_at: ctx.turn_timing_state.started_at_unix_secs().await,
model_context_window: ctx.model_context_window(),
collaboration_mode_kind: ctx.mode,
});
// TurnStarted在等待startup prewarm前发送,客户端不会把预热等待误判为Turn未启动。
sess.send_event(ctx.as_ref(), event).await;
sess.set_server_reasoning_included(/*included*/ false).await;
sess.consume_startup_prewarm_for_regular_turn(&cancellation_token)
.await
}
.instrument(trace_span!("regular_task.prepare_run_turn"))
.await;
let prewarmed_client_session = match prewarmed_client_session {
SessionStartupPrewarmResolution::Cancelled => {
// 预热等待期取消仍记录输入/hooks,但不会进入run_turn。
run_hooks_and_record_inputs(&sess, &ctx, &input).await;
return Ok(None);
}
SessionStartupPrewarmResolution::Unavailable { .. } => None,
SessionStartupPrewarmResolution::Ready(prewarmed_client_session) => {
Some(*prewarmed_client_session)
}
};
let mut next_input = input;
let mut prewarmed_client_session = prewarmed_client_session;
loop {
let last_agent_message = run_turn(
Arc::clone(&sess),
Arc::clone(&ctx),
next_input,
// 预热session只交给第一次run_turn;后续复用Turn内既有状态。
prewarmed_client_session.take(),
cancellation_token.child_token(),
)
.instrument(run_turn_span.clone())
.await?;
if !sess.input_queue.has_pending_input(&sess.active_turn).await {
return Ok(last_agent_message);
}
// 原始TurnInput已经记录;外层补跑只消费queue,避免重复写入首轮输入。
next_input = Vec::new();
}
}SessionStartupPrewarmResolution::Unavailable 是性能降级而不是正确性失败:run_turn 会创建正常 client session。 Cancelled 不进入主循环,RegularTask::run 自身返回 Ok(None);若取消来自外部 Interrupt,持有 RunningTask 的 abort 路径负责取消、移出 active task 并发送 TurnAborted,不能仅凭这里的返回值判断终态。
3. 主循环前置门
进入 loop 之前,run_turn 依次处理预采样压缩、显式 MCP/plugin/skill mention、首个 StepContext、初始 WorldState 与输入/hooks。这些步骤任何一个提前结束,都不会进入 sampling 循环。
源码位置:codex-rs/core/src/session/turn.rs,符号 run_turn 的前置阶段。
let mut client_session =
prewarmed_client_session.unwrap_or_else(|| sess.services.model_client.new_session());
if let Err(err) = run_pre_sampling_compact(
&sess,
&turn_context,
&mut client_session,
&cancellation_token,
)
.await
{
if matches!(err.details(), CodexErrorDetails::TurnAborted) {
run_hooks_and_record_inputs(&sess, &turn_context, &input).await;
return Err(err);
}
if matches!(err.details(), CodexErrorDetails::ToolCollision(_)) {
// ToolCollision保持结构化错误,由Task层统一收尾,不转成普通完成。
return Err(err);
}
let error = err.to_codex_protocol_error();
sess.emit_turn_error_lifecycle(turn_context.as_ref(), error.clone())
.await;
error!("Failed to run pre-sampling compact");
// 其他压缩失败已发error lifecycle,返回Ok(None)让Task生成带error的TurnComplete。
return Ok(None);
}
let user_input = turn_user_input(&input);
let (required_servers, mentioned_plugins) =
match required_mcp_servers_for_input(&sess, turn_context.as_ref(), &user_input)
.or_cancel(&cancellation_token)
.await
{
Ok(requirements) => requirements,
Err(err) => {
run_hooks_and_record_inputs(&sess, &turn_context, &input).await;
return Err(err.into());
}
};
let first_step_context = match sess
.capture_step_context_with_required_mcp_servers(
Arc::clone(&turn_context),
&cancellation_token,
&required_servers,
)
.await
{
Ok(step_context) => step_context,
Err(err) if matches!(err.details(), CodexErrorDetails::TurnAborted) => {
// 首个Step取消仍先记录原始输入和hooks,再把TurnAborted交给Task层。
run_hooks_and_record_inputs(&sess, &turn_context, &input).await;
return Err(err);
}
Err(err) => return Err(err),
};
// ...建立WorldState,注入skills/plugins,运行session-start与user-prompt hooks| 门 | 失败或停止 | 对外语义 |
|---|---|---|
| pre-sampling compact | abort/collision/其他失败分流 | Err 或带 error 的正常收尾 |
| required server 解析 | cancellation | 记录输入后 TurnAborted |
| StepContext 捕获 | cancellation、ToolCollision、构建错误 | 不发布部分 Step |
| hooks 与 skill/plugin injection | hook stop 或依赖安装停止 | Ok(None),终态由 Task 层发送 |
4. needs_follow_up
外层循环不直接数工具调用,而是消费 SamplingRequestResult。这个结果由每个 completed output item 和最终 ResponseEvent::Completed 累积而来;之后还要与 input queue 合并,才得到真正的 needs_follow_up。
4.1 工具后续轮次
源码位置:codex-rs/core/src/stream_events_utils.rs,符号 handle_output_item_done。
pub(crate) async fn handle_output_item_done(
ctx: &mut HandleOutputCtx,
item: ResponseItem,
previously_active_item: Option<TurnItem>,
) -> Result<OutputItemResult> {
let mut output = OutputItemResult::default();
let plan_mode = ctx.turn_context.mode == ModeKind::Plan;
match ToolRouter::build_tool_call(item.clone()) {
Ok(Some(call)) => {
ctx.sess
.input_queue
.accept_mailbox_delivery_for_current_turn(
&ctx.sess.active_turn,
&ctx.turn_context.sub_id,
)
.await;
record_completed_response_item(
ctx.sess.as_ref(),
ctx.turn_context.as_ref(),
&item,
)
.await;
let cancellation_token = ctx.cancellation_token.child_token();
let tool_future: InFlightFuture<'static> = Box::pin(
ctx.tool_runtime
.clone()
.handle_tool_call(call, cancellation_token),
);
// 工具结果必须回灌模型,所以创建future的同时把follow-up置true。
output.needs_follow_up = true;
output.tool_future = Some(tool_future);
}
Ok(None) => {
// assistant/reasoning等非工具item被记录,但自身不要求再请求模型。
let finalized_turn_item = finalize_non_tool_response_item(
ctx.sess.as_ref(),
TurnItemContributorPolicy::Run(ctx.turn_store.as_ref()),
&item,
plan_mode,
)
.await;
// ...发送item lifecycle并写入history/rollout
output.last_agent_message =
finalized_facts.and_then(|facts| facts.last_agent_message);
}
Err(FunctionCallError::RespondToModel(message)) => {
// 工具请求被拒绝或需直接回复时,合成output回灌,仍需要follow-up。
let response = ResponseInputItem::FunctionCallOutput {
call_id: String::new(),
output: FunctionCallOutputPayload {
body: FunctionCallOutputBody::Text(message),
..Default::default()
},
};
record_completed_response_item(
ctx.sess.as_ref(),
ctx.turn_context.as_ref(),
&item,
)
.await;
if let Some(response_item) = response_input_to_response_item(&response) {
ctx.sess
.record_conversation_items(
&ctx.turn_context,
std::slice::from_ref(&response_item),
)
.await;
}
output.needs_follow_up = true;
}
Err(FunctionCallError::Fatal(message)) => {
return Err(CodexErr::Fatal(message));
}
}
Ok(output)
}这里的 RespondToModel 不是 Turn 失败:Core 已经把拒绝原因变成 tool output,模型仍有机会调整。只有 Fatal 作为 CodexErr 离开当前 sampling。
4.2 完成事件
源码位置:codex-rs/core/src/session/turn.rs,符号 try_run_sampling_request 的 ResponseEvent::Completed 分支。
ResponseEvent::Completed {
response_id,
token_usage,
end_turn,
} => {
// ...发送RawResponseCompleted并记录token usage
let budget_result = sess
.record_token_usage_info(&turn_context, token_usage.as_ref())
.await;
should_emit_token_count = true;
should_emit_turn_diff = true;
if let Err(err) = budget_result {
break Err(err);
}
if let Some(false) = end_turn {
// 服务端明确声明尚未结束时,即使没有工具item也必须再次sampling。
needs_follow_up = true;
}
break Ok(SamplingRequestResult {
needs_follow_up,
last_agent_message,
});
}此外,commentary/reasoning item 完成时若 mailbox 已有消息,stream loop 可以提前返回 needs_follow_up=true。这让 agent mail 在安全边界抢占后续输出,但 final answer 不会无条件被重新打开。
5. 工具 Future
一个 response 可以产生多个工具 future。它们在 item 完成时入 FuturesOrdered,可由 ToolCallRuntime 的 读写 gate 并发执行;stream loop 退出后统一 drain,工具 output 写入 history,然后外层循环才构造下一份 Prompt。
源码位置:codex-rs/core/src/session/turn.rs,符号 try_run_sampling_request 与 drain_in_flight。
let mut in_flight: FuturesOrdered<
BoxFuture<'static, CodexResult<ResponseInputItem>>
> = FuturesOrdered::new();
let mut needs_follow_up = false;
// ...在ResponseEvent::OutputItemDone中
if let Some(tool_future) = output_result.tool_future {
// FuturesOrdered允许future并行推进,但next按插入顺序产出结果。
in_flight.push_back(tool_future);
}
if let Some(agent_message) = output_result.last_agent_message {
last_agent_message = Some(agent_message);
}
needs_follow_up |= output_result.needs_follow_up;
// ...response event loop结束后
let tool_blocking_timing_guard = if in_flight.is_empty() {
None
} else {
Some(turn_context.turn_timing_state.begin_tool_blocking())
};
drain_in_flight(&mut in_flight, sess.clone(), turn_context.clone()).await?;
drop(tool_blocking_timing_guard);
if cancellation_token.is_cancelled() {
// 即使response.completed已到达,工具等待期间取消仍以TurnAborted覆盖成功outcome。
return Err(CodexErr::TurnAborted);
}
outcome源码位置:codex-rs/core/src/session/turn.rs,符号 drain_in_flight。
async fn drain_in_flight(
in_flight: &mut FuturesOrdered<BoxFuture<'static, CodexResult<ResponseInputItem>>>,
sess: Arc<Session>,
turn_context: Arc<TurnContext>,
) -> CodexResult<()> {
while let Some(res) = in_flight.next().await {
match res {
Ok(response_input) => {
let response_item = response_input.into();
// 下一轮clone_history之前,tool output已经进入conversation history。
sess.record_conversation_items(
&turn_context,
std::slice::from_ref(&response_item),
)
.await;
mark_thread_memory_mode_polluted_if_external_context(
sess.as_ref(),
turn_context.as_ref(),
&response_item,
)
.await;
}
Err(err) => {
error_or_panic(format!(
"in-flight tool future failed during drain: {err}"
));
}
}
}
Ok(())
}非 fatal 工具失败通常已被 ToolCallRuntime 转成可回灌的 failure output;这里的 Err 表示 future 层异常, 源码记录 invariant violation 后继续 drain,避免一个异常 future 丢掉其他已启动工具的结果。
6. 循环控制
sampling 成功后,Core 将模型的 follow-up 与 input queue 合并:
needs_follow_up = model_needs_follow_up || has_pending_input
should_roll_over = needs_follow_up && (new_window_requested || token_limit_reached)只有“还要继续”并且达到换窗条件时才做 mid-turn compaction。若模型已经完成且没有 pending input,即使 token 状态触顶也不会为了一个不会发生的下一请求额外压缩。
源码位置:codex-rs/core/src/session/turn.rs,符号 run_turn 的 post-sampling 分支。
let SamplingRequestResult {
needs_follow_up: model_needs_follow_up,
last_agent_message: sampling_request_last_agent_message,
} = sampling_request_output;
if model_needs_follow_up {
sess.input_queue
.accept_mailbox_delivery_for_current_turn(
&sess.active_turn,
&turn_context.sub_id,
)
.await;
}
can_drain_pending_input = true;
let (has_pending_input, token_status) = async {
let has_pending_input =
sess.input_queue.has_pending_input(&sess.active_turn).await;
let token_status = super::context_window::context_window_token_status(
sess.as_ref(),
turn_context.as_ref(),
)
.await;
(has_pending_input, token_status)
}
.await;
let needs_follow_up = model_needs_follow_up || has_pending_input;
let token_limit_reached = token_status.token_limit_reached;
let should_roll_over = needs_follow_up
&& (sess.take_new_context_window_request().await || token_limit_reached);
if should_roll_over {
// 当前Step交给inline compaction,压缩完成后continue重新捕获下一Step。
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;
}
if !needs_follow_up {
// ...进入stop hooks;未阻止时break
break;
}
continue;这是真实控制状态而非独立 Rust enum。源码中的权威状态仍是局部变量、input queue、token status 和返回值。
7. 六类路径逐条展开
| 路径 | 触发输入/响应 | 关键分支 | 最终结果 |
|---|---|---|---|
| 无工具 | assistant/reasoning + completed | model_needs_follow_up=false | stop hook允许后退出 |
| 单工具 | 一个 tool item + completed | future drain,下一轮带一个 output | 第二次 sampling 后退出或继续 |
| 多工具/多轮 | 多个 tool item,或多次 response 继续 | FuturesOrdered + 多次外层 loop | outputs成组回灌,直到无 follow-up |
| 取消 | stream/工具/压缩期间 token cancel | CodexErr::TurnAborted | Task层发送 TurnAborted |
| 超限 | 成功响应后预算触顶,或API返回 context exceeded | mid-turn compact;或特殊错误分类 | 压缩后继续;或 Error + TurnComplete |
| 模型错误 | retry耗尽、invalid image、usage limit等 | 类型化错误分流 | Error + TurnComplete,Session可继续 |
7.1 无工具完成
assistant message 只更新 last_agent_message,真正退出还要等待 response.completed、tool future drain、pending input 检查和 stop hook。服务端若给 end_turn=false,即使没有工具也继续。
7.2 sampling
第一次 response 记录 function call、执行工具并将 output 写入 history;needs_follow_up=true 令外层继续;第二次 Prompt 才包含 tool output。第二次若只有最终 assistant message 且无 pending input,进入 stop hook 后结束。
7.3 工具数量与轮次
同一 response 内多个工具可并行执行,但 FuturesOrdered 按调用顺序产出并成组写回;模型在下一 response 再次调用工具,则是另一轮 StepContext 与 sampling。不要把“并行三个工具”误画成“三次模型请求”。
7.4 取消与模型响应
stream 的 .or_cancel() 会直接返回 TurnAborted;即使 response.completed 已到达,只要工具 drain 期间 token 被取消,末尾的 is_cancelled() 仍覆盖成功 outcome。这样不会在工具尚未稳定写回时错误报告正常完成。
7.5 超限处理路径
成功 sampling 后的 token_limit_reached 只有在 needs_follow_up=true 时触发 mid-turn compact。若请求本身被 provider 以 ContextWindowExceeded 拒绝,内部 retry loop 不重试,而是把 total token 标记为有效模型窗口, 交给外层错误分支发送 Error 并结束本次 Turn。
7.6 普通模型错误
retryable stream error先按 provider budget 重试;不可重试或重试耗尽后,run_turn 记录 analytics/error lifecycle并发送 Error,然后 break。Task 收尾读取 terminal_error,生成带 error 的 TurnComplete,active turn 清除后下一次用户 Turn 仍可运行。
8. Stop hook
模型与 input queue 都认为无需继续时,Core 才运行 stop hooks。hook 若 should_block 且提供 continuation fragment,Core 把它作为新的 response item 写入 history,重新开放当前 Turn 的 mailbox,然后 continue。
源码位置:codex-rs/core/src/session/turn.rs,符号 run_turn 的 stop hook 分支。
if !needs_follow_up {
last_agent_message = sampling_request_last_agent_message;
let stop_outcome = run_turn_stop_hooks(
&sess,
&turn_context,
stop_hook_active,
last_agent_message.clone(),
)
.await;
if stop_outcome.should_block {
if let Some(hook_prompt_message) =
build_hook_prompt_message(&stop_outcome.continuation_fragments)
{
// continuation进入history,所以后续请求能看到前几次hook反馈。
sess.record_response_item_and_emit_turn_item(
&turn_context,
hook_prompt_message,
)
.await;
sess.input_queue
.accept_mailbox_delivery_for_current_turn(
&sess.active_turn,
&turn_context.sub_id,
)
.await;
stop_hook_active = true;
continue;
} else {
// block却没有prompt无法驱动模型继续,只warning并忽略该block。
sess.send_event(
&turn_context,
EventMsg::Warning(WarningEvent {
message: "Stop hook requested continuation without a prompt; ignoring the block."
.to_string(),
}),
)
.await;
}
}
if stop_outcome.should_stop {
break;
}
if run_legacy_after_agent_hook(
&sess,
&turn_context,
&sampling_request_input,
last_agent_message.clone(),
)
.await
{
return Ok(None);
}
break;
}stop_hook_active 会传给后续 stop hook,避免把第一次完成与 hook 驱动后的再次完成混为同一上下文。最终 last_agent_message 取最后一次允许停止的 assistant message,而不是第一次 draft。
9. 错误退出
源码位置:codex-rs/core/src/session/turn.rs,符号 run_turn 的错误分支。
match sampling_request_result {
Ok((sampling_request_output, sampling_request_input)) => {
// ...正常post-sampling判断
}
Err(err) if matches!(err.details(), CodexErrorDetails::TurnAborted) => {
// 取消保持Err身份,不能降级成带Error的TurnComplete。
return Err(err);
}
Err(codex_error)
if matches!(
codex_error.details(),
CodexErrorDetails::InvalidImageRequest()
) =>
{
sess.track_turn_codex_error(turn_context.as_ref(), &codex_error);
let error = CodexErrorInfo::BadRequest;
sess.emit_turn_error_lifecycle(turn_context.as_ref(), error.clone())
.await;
let event = EventMsg::Error(ErrorEvent {
message: "Invalid image in your last message. Please remove it and try again."
.to_string(),
codex_error_info: Some(error),
});
sess.send_event(&turn_context, event).await;
break;
}
Err(e) => {
info!("Turn error: {e:#}");
let error = e.to_codex_protocol_error();
sess.emit_turn_error_lifecycle(turn_context.as_ref(), error.clone())
.await;
sess.track_turn_codex_error(turn_context.as_ref(), &e);
let event = EventMsg::Error(e.to_error_event(/*message_prefix*/ None));
sess.send_event(&turn_context, event).await;
// break只结束run_turn loop;Session仍存活,Task层稍后发TurnComplete。
break;
}
}10. 终态由Task收尾层
run_turn 正常路径只返回 Option<String>;Task runner 持有 abort reason、TurnState 和 timing,才能发送最终 envelope。普通 ErrorEvent 若 affects_turn_status 会先写入 turn_context.terminal_error,随后成为 TurnComplete.error。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::on_task_finished 的终态构造。
let event = if let Some(reason) = abort_reason {
self.emit_turn_abort_lifecycle(
reason.clone(),
turn_context.extension_data.as_ref(),
)
.await;
// 只有Task确立abort reason才发送TurnAborted。
EventMsg::TurnAborted(TurnAbortedEvent {
turn_id: Some(turn_context.sub_id.clone()),
reason,
started_at,
completed_at,
duration_ms,
})
} else {
let time_to_first_token_ms = turn_context
.turn_timing_state
.time_to_first_token_ms()
.await;
// run_turn中发送的终态Error在这里被带入TurnComplete。
let error = turn_context.terminal_error.lock().await.clone();
self.emit_turn_stop_lifecycle(turn_context.extension_data.as_ref())
.await;
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_context.sub_id.clone(),
last_agent_message,
error,
started_at,
completed_at,
duration_ms,
time_to_first_token_ms,
})
};
self.send_event(turn_context.as_ref(), event).await;这也划清本文与 Task 生命周期专题的边界:本文解释 run_turn 为什么返回;Task 如何注册、替换和清理 RunningTask 由后续专题负责。
11. 主循环反向测试
11.1 多工具并发但输出
源码位置:codex-rs/core/tests/suite/tool_parallelism.rs,测试 tool_results_grouped。
mount_sse_once(
&server,
sse(vec![
json!({"type": "response.created", "response": {"id": "resp-1"}}),
ev_function_call("call-1", "shell_command", &shell_args),
ev_function_call("call-2", "shell_command", &shell_args),
ev_function_call("call-3", "shell_command", &shell_args),
ev_completed("resp-1"),
]),
)
.await;
let tool_output_request = mount_sse_once(
&server,
sse(vec![
ev_assistant_message("msg-1", "done"),
ev_completed("resp-2"),
]),
)
.await;
run_turn(&test, "run shell three times").await?;
let input = tool_output_request.single_request().input();
// ...分别筛出function_call与function_call_output
// 三个call与三个output都进入第二次sampling。
assert_eq!(function_calls.len(), 3);
assert_eq!(function_call_outputs.len(), 3);
for (index, _) in &function_calls {
for (output_index, _) in &function_call_outputs {
assert!(*index < *output_index,
"all function calls must come before outputs");
}
}
let zipped = function_calls
.iter()
.zip(function_call_outputs.iter())
.collect::<Vec<_>>();
for (call, output) in zipped {
// FuturesOrdered保证output顺序与call顺序对应。
assert_eq!(
call.1.get("call_id").and_then(Value::as_str),
output.1.get("call_id").and_then(Value::as_str)
);
}11.2 Stop hook
源码位置:codex-rs/core/tests/suite/hooks.rs,测试 stop_hook_can_block_multiple_times_in_same_turn。
let responses = mount_sse_sequence(
&server,
vec![
sse(vec![
ev_response_created("resp-1"),
ev_assistant_message("msg-1", "draft one"),
ev_completed("resp-1"),
]),
sse(vec![
ev_response_created("resp-2"),
ev_assistant_message("msg-2", "draft two"),
ev_completed("resp-2"),
]),
sse(vec![
ev_response_created("resp-3"),
ev_assistant_message("msg-3", "final draft"),
ev_completed("resp-3"),
]),
],
)
.await;
// ...fixture依次返回两条stop continuation prompt
test.submit_turn("hello from the sea").await?;
let requests = responses.requests();
// 两次block把一次无工具Turn扩展成三次sampling。
assert_eq!(requests.len(), 3);
assert_eq!(
request_hook_prompt_texts(&requests[2]),
vec![
FIRST_CONTINUATION_PROMPT.to_string(),
SECOND_CONTINUATION_PROMPT.to_string(),
],
"third request should retain hook prompts in user history",
);11.3 工具执行中断产生
源码位置:codex-rs/core/tests/suite/abort_tasks.rs,测试 interrupt_long_running_tool_emits_turn_aborted。
let body = sse(vec![
ev_function_call("call_sleep", "shell_command", &args),
ev_completed("done"),
]);
mount_sse_once(&server, body).await;
codex.submit(Op::UserInput {
items: vec![UserInput::Text {
text: "start sleep".into(),
text_elements: Vec::new(),
}],
final_output_json_schema: None,
responsesapi_client_metadata: None,
additional_context: Default::default(),
thread_settings: Default::default(),
})
.await
.unwrap();
// 等待工具真正开始,排除Interrupt发生在sampling前的竞态。
wait_for_event(&codex, |ev| matches!(ev, EventMsg::ExecCommandBegin(_))).await;
codex.submit(Op::Interrupt).await.unwrap();
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnAborted(_))).await;11.4 上下文窗口超限
源码位置:codex-rs/core/tests/suite/client.rs,测试 context_window_error_sets_total_tokens_to_model_window。
const EFFECTIVE_CONTEXT_WINDOW: i64 = (272_000 * 95) / 100;
mount_sse_once_match(
&server,
body_string_contains("trigger context window"),
sse_failed(
"resp_context_window",
"context_length_exceeded",
"Your input exceeds the context window of this model. Please adjust your input and try again.",
),
)
.await;
// ...先seed成功Turn,再提交触发超限的Turn
let token_event = wait_for_event(&codex, |event| {
matches!(
event,
EventMsg::TokenCount(payload)
if payload.info.as_ref().is_some_and(|info| {
info.model_context_window == Some(info.total_token_usage.total_tokens)
&& info.total_token_usage.total_tokens > 0
})
)
}).await;
let EventMsg::TokenCount(token_payload) = token_event else {
unreachable!("wait_for_event returned unexpected event");
};
let info = token_payload
.info
.expect("token usage info present when context window is exceeded");
// 被动超限把usage标记为完整有效窗口,供后续压缩/诊断使用。
assert_eq!(info.model_context_window, Some(EFFECTIVE_CONTEXT_WINDOW));
assert_eq!(info.total_token_usage.total_tokens, EFFECTIVE_CONTEXT_WINDOW);
let error_event = wait_for_event(&codex, |ev| matches!(ev, EventMsg::Error(_))).await;
let expected_context_window_message = CodexErr::ContextWindowExceeded.to_string();
assert!(matches!(
error_event,
EventMsg::Error(ref err) if err.message == expected_context_window_message
));
// 普通模型错误之后仍有TurnComplete,而不是TurnAborted。
wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await;无工具路径由 additional_context_is_model_visible_but_not_a_user_message_item 验证一次 completed response 即 TurnComplete;单工具路径由 shell_command_tool_executes_command_and_streams_output 验证第二次请求包含对应 function output;模型错误后的 Session 复用由 continue_after_stream_error 验证失败 Turn 完成后下一 Turn 仍能 成功。这三项与上面四组测试共同覆盖模块计划要求的六类路径。
上述测试没有覆盖所有 provider 流错误的恢复,也没有覆盖任意 stop hook 组合都会在有限次数内收敛;它们只 覆盖当前版本列出的 follow-up、工具回灌、普通错误和 context-window 边界。
12. Turn不结束
遇到 Turn 持续请求模型时,不要先猜“loop 写死了”。按下面顺序读取事实:
- 查看本次
SamplingRequestResult.needs_follow_up是 tool item、end_turn=false还是 mailbox preempt 设置的; - 检查
has_pending_input是否让模型已完成的 response 重新进入下一轮; - 检查 token status 是否令 follow-up 先经过 mid-turn compaction;
- 检查 stop hook 是否不断返回 continuation fragment;
- 区分 request retry 次数与真正的 sampling 请求数;前者不会重建 StepContext;
- 最后查看 Task 收尾收到的是
Ok还是 TurnAborted,以及terminal_error是否已被 ErrorEvent 设置。
可以选 tool_results_grouped 的三工具 fixture,手工标出四个时间点:tool item 写 history、future 入队、output 写 history、下一 Prompt clone history。若能解释为什么第二次请求不可能早于第三个 output 的 drain,就已经掌握 主循环的关键一致性。Task 如何持有 abort handle、替换旧 Task 并发送终态事件,由 Task抽象与生命周期 继续解释。
可以用下面的只读搜索把本文的 turn loop 主线落回源码:
rg -n "run_turn|model_needs_follow_up|can_drain_pending_input|TurnCompleted" codex-rs/core/src