RegularTask完整流程
普通用户输入不会直接调用模型,而是先进入 TurnInputRequest 路由。Core 根据当前是否存在可 steer 的 Regular Turn,返回 Started、Steered 或带 NotSubmittedReason 的结果;只有 Started 分支才会 创建 RegularTask。本文追踪 Started 之后的生命周期,并解释 TurnStarted、startup prewarm、run_turn、 pending input 和 TurnComplete/TurnAborted 不是同一个时刻。
阅读本文前,建议先读 Task抽象与生命周期 和 Turn主循环与退出条件。前者解释 SessionTask 的统一 owner, 后者解释 run_turn 内部的 sampling 与工具回灌;本文只负责输入路由、RegularTask 外层循环和统一收尾。
1. 输入路由
start_or_steer_turn 使用 TurnInputMode::StartOrSteer。Session 先尝试把用户输入注入当前 Regular Turn;只有 NoActiveTurn 才创建新的 TurnContext、合并 additional context 并启动 RegularTask。Steered 不会创建第二个 task。
源码位置:codex-rs/core/src/session/turn_input.rs :: start_or_steer
match session
.steer_input(
&mut items,
additional_context.clone(),
/*expected_turn_id*/ None,
settings.required_active_final_output_json_schema(),
client_id.clone(),
responsesapi_client_metadata.clone(),
incoming_root_turn_id,
)
.await
{
Ok(turn_id) => {
settings.apply_steered(session, submission_id).await?;
Ok(TurnInputSubmission::Steered { turn_id })
}
Err(NotSubmittedReason::NoActiveTurn) => {
let turn_context = settings
.apply_started(session, submission_id.clone())
.await?;
let mut task_input = merge_additional_context_input(session, additional_context).await;
if !items.is_empty() {
task_input.push(TurnInput::UserInput {
content: items,
client_id,
});
}
session
.spawn_task(turn_context, task_input, RegularTask::new())
.await;
Ok(TurnInputSubmission::Started {
turn_id: submission_id,
})
}
Err(reason) => Ok(TurnInputSubmission::NotSubmitted { reason }),
}TurnInputRequest 还可以携带 parent_turn_id、root_turn_id 和 responsesapi_client_metadata。新 Turn 会把这些字段写入 TurnContext;Steered 输入只追加到当前 Turn,不重写已有的 root lineage。
Started 路径先应用约束后的 settings,再把 additional context 放在用户输入之前;Steered 路径把输入 放进当前 Turn 的 pending queue。Review、Compact、expected Turn ID 不匹配、空输入和 Plan mode 都可能 得到稳定的 NotSubmittedReason。
2. 任务启动
Session::spawn_task 先以 Replaced 原因中止旧 Task,再清理 connector selection,最后调用 start_task。后者读取 pending mailbox、设置 ActiveTurn、记录 timing,创建 cancellation token 和 Tokio task。Task 安装完成后,start_or_steer_turn 才返回 Started。
源码位置:codex-rs/core/src/tasks/mod.rs :: Session::spawn_task, Session::start_task
pub async fn spawn_task<T: SessionTask>(
self: &Arc<Self>,
turn_context: Arc<TurnContext>,
input: Vec<TurnInput>,
task: T,
) {
self.abort_all_tasks(TurnAbortReason::Replaced).await;
self.clear_connector_selection().await;
self.start_task(
turn_context,
input,
task,
MailboxParentProvenance::Ignore,
)
.await;
}start_task 会先把仍属于当前 turn state 的 pending input 合并进 task 输入,然后发出 turn-start lifecycle;Tokio closure 完成后统一调用 Session::on_task_finished。因此 Started 只表示 新 task 已安装,不表示模型请求已发送。
3. TurnStarted
RegularTask 是零字段 SessionTask,kind 为 Regular,span 名为 session_task.turn。它先发布 TurnStarted,再消费一次性的 startup prewarm handle;客户端不需要等待 websocket 预热才能看到 Turn 生命周期开始。
源码位置:codex-rs/core/src/tasks/regular.rs :: RegularTask::run
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,
});
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
};ctx 是已经安装到 RunningTask 的正式 TurnContext。发送事件后,任务才决定使用预热 session、创建 普通 ModelClientSession,或在取消时提前结束。
4. 预热分支
SessionStartupPrewarmHandle::resolve 按 handle 启动时间计算剩余 timeout。Ready 只交给第一次 run_turn;Unavailable 包括未调度、setup/join 失败和 timeout;Cancelled 会 abort 预热 task 并停止 进入模型循环。
源码位置:codex-rs/core/src/session_startup_prewarm.rs :: SessionStartupPrewarmHandle::resolve
let age_at_first_turn = started_at.elapsed();
let remaining = timeout.saturating_sub(age_at_first_turn);
match tokio::select! {
_ = cancellation_token.cancelled() => None,
result = tokio::time::timeout(remaining, &mut task) => Some(result),
} {
Some(Ok(result)) => Self::resolution_from_join_result(result, started_at),
Some(Err(_elapsed)) => {
task.abort();
SessionStartupPrewarmResolution::Unavailable {
status: "timed_out",
prewarm_duration: Some(started_at.elapsed()),
}
}
None => {
task.abort();
SessionStartupPrewarmResolution::Cancelled
}
}timeout 从 prewarm 启动时开始计算,首个 Turn 到达较晚时不会重新获得完整等待时间。Unavailable 只损失 预热收益,RegularTask 随后让 run_turn 创建普通 client session。
5. 外层循环
一次 run_turn 内部会在模型 sampling、工具执行和结果回灌之间循环;RegularTask 还维护一层更外部 的 pending-input 循环。当 run_turn 返回后,如果 InputQueue 仍有 pending input,RegularTask 用空 next_input 再次进入 run_turn,避免把已经进入当前 Turn 的输入重新作为首轮输入。
源码位置:codex-rs/core/src/tasks/regular.rs :: RegularTask::run
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,
prewarmed_client_session.take(),
cancellation_token.child_token(),
)
.await?;
if !sess.input_queue.has_pending_input(&sess.active_turn).await {
return Ok(last_agent_message);
}
next_input = Vec::new();
}Option::take() 保证预热 session 最多交给第一次 run_turn。同一次 run_turn 内的多次 sampling 复用 该 session;外层补跑时传入 None,会创建新的 ModelClientSession。
6. 统一收尾
所有 SessionTask 都从 start_task 创建的 Tokio closure 进入 Session::on_task_finished。该方法 先区分正常返回、TurnAborted 和意外错误,再接管 ActiveTurn 中剩余 pending input,运行 hooks,计算 token usage,最后发出 TurnComplete 或 TurnAborted。
源码位置:codex-rs/core/src/tasks/mod.rs :: Session::on_task_finished
let (last_agent_message, abort_reason) = match task_result {
Ok(last_agent_message) => (last_agent_message, None),
Err(err) if matches!(err.details(), CodexErrorDetails::TurnAborted) => {
(None, Some(TurnAbortReason::Interrupted))
}
Err(err) => {
self.emit_turn_error_lifecycle(
turn_context.as_ref(),
err.to_codex_protocol_error(),
)
.await;
self.send_event(
turn_context.as_ref(),
EventMsg::Error(err.to_error_event(/*message_prefix*/ None)),
)
.await;
(None, None)
}
};
let pending_input = self
.input_queue
.take_pending_input_for_turn_state(turn_state.as_ref())
.await;
run_hooks_and_record_inputs(
self,
&turn_context,
&pending_input,
PersistContext::Standard,
)
.await;pending input 的接管发生在清理 ActiveTurn task 之后、发送终态事件之前。正常返回根据 terminal_error 选择 Completed 或 Failed idle cause,最终发送带 timing、usage 和可选 error 的 TurnComplete;取消分支发送 TurnAborted。
7. 失败边界
| 阶段 | 结果 | 后续动作 |
|---|---|---|
| 输入路由 | NotSubmittedReason | 不创建 RegularTask,不发送 TurnStarted |
| prewarm | Unavailable | 创建普通 ModelClientSession,继续 run_turn |
| prewarm | Cancelled | 记录输入并返回,外层 abort owner 负责 TurnAborted |
| run_turn | TurnAborted | 统一收尾发送 TurnAborted |
| run_turn | 普通错误 | 写入 terminal error,发送 Error 和 TurnComplete |
| 收尾 | late pending input | 先记录输入,再清理 ActiveTurn |
prewarm unavailable 不是模型错误;它只表示优化路径未能提供可复用 client session。相反,普通 run_turn 错误会进入 TurnContext 的 terminal error,并由统一收尾转换为可观察事件。
8. 测试路径
第一个测试让 startup prewarm 保持 pending,再启动 RegularTask;即使预热未完成,也必须先收到 TurnStarted。
源码位置:codex-rs/core/src/session/tests.rs :: regular_turn_emits_turn_started_with_trace_id_without_waiting_for_startup_prewarm
let first = tokio::time::timeout(
std::time::Duration::from_millis(200),
rx.recv(),
)
.await
.expect("expected turn started event without waiting for startup prewarm")
.expect("channel open");
let EventMsg::TurnStarted(turn_started) = first.msg else {
panic!("expected turn started event");
};
assert_eq!(turn_started.turn_id, tc.sub_id);
assert_eq!(turn_started.trace_id, tc.trace_id);第二个测试验证 active Turn 下 idle-only 输入不会创建第二个 RegularTask。
源码位置:codex-rs/core/src/session/turn_input_tests.rs :: start_only_rejects_active_turn_without_injecting
let submission = submit_start_only(
&session,
SubmittedTurnInput::ResponseItem(user_message("synthetic idle input")),
)
.await;
assert_eq!(
submission,
TurnInputSubmission::NotSubmitted {
reason: NotSubmittedReason::NotIdle,
}
);
assert_eq!(
session.input_queue.get_pending_input(&session.active_turn).await,
(Vec::<TurnInput>::new(), None, None),
);第三个测试让任务结束时仍有 pending user input,确认 finish owner 会把它写入 history 和事件流,而不是 在清理 active slot 时丢弃。
源码位置:codex-rs/core/src/session/tests.rs :: task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input
let submission = submit_steer_only(
&session,
pending_user_input.clone(),
&turn_context.sub_id,
)
.await;
assert!(matches!(submission, TurnInputSubmission::Steered { .. }));
session
.on_task_finished(Arc::clone(&turn_context), Ok(None))
.await;
let history = session.clone_history().await;
assert!(history_contains_user_message(&history, "late pending input"));这些测试分别覆盖首事件顺序、active/idleness 路由和 finish owner 的输入接管;它们不证明所有并发调度交错, 也不把 RegularTask 的职责扩展到 Compact、Review 或 UserShell task。
补充任务生命周期状态图,强调 TurnStarted、预热解析、循环和统一收尾的边界。
9. 阅读检查
使用下面的只读搜索可回到当前源码:
rg -n "start_or_steer|TurnInputSubmission|RegularTask|consume_startup_prewarm|on_task_finished" \
codex-rs/core/src需要继续深入时,阅读 Turn主循环与退出条件 查看 run_turn 的 sampling/tool loop,或阅读 Task抽象与生命周期 查看 RunningTask 的 取消、替换和资源清理。
