Skip to content

RegularTask完整流程

从 TurnInput 路由到 RegularTask、启动预热、模型循环和统一收尾,解释普通 Turn 的真实生命周期。

基于rust-v0.150.0
CodexRustRuntimeRegularTask

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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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
prewarmUnavailable创建普通 ModelClientSession,继续 run_turn
prewarmCancelled记录输入并返回,外层 abort owner 负责 TurnAborted
run_turnTurnAborted统一收尾发送 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

rust
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

rust
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

rust
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. 阅读检查 ​

使用下面的只读搜索可回到当前源码:

bash
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 的 取消、替换和资源清理。