Skip to content

Task抽象与生命周期

解释 SessionTask 如何被类型擦除、注册为 RunningTask,并在正常完成、取消和替换时由唯一 owner 产生终态。

基于rust-v0.150.0
CodexRustRuntimeTask

Task抽象与生命周期 ​

Codex Core 中的 Task 不是一次模型请求。RegularTask 可以包含多轮 sampling 与工具调用,CompactTask 运行一次 压缩工作流,ReviewTask 管理子会话,UserShellCommandTask 执行独立命令。它们共同实现 SessionTask,由 Session 注册为一个 RunningTask,并共享终态、取消、持久化屏障和指标收尾。

本文接续 Turn主循环与退出条件:前文解释 RegularTask 内部为什么继续或 返回,本文解释返回值由谁接收、active slot 如何释放、外部 Interrupt/Replaced 如何抢占 owner。还建议先了解 Session核心数据结构,避免把 ActiveTurn、TurnState 与 TurnContext 当成同一个对象。

本文不展开四种具体 Task 的业务算法,也不重复 Session 整体关闭顺序。目标是让读者能从 spawn_task() 追踪到 Tokio task、RunningTask、on_task_finished() 或 handle_task_abort(),并判断一个 terminal event 究竟由正常完成 owner 还是 abort owner 发出。

1. Task层级 ​

TurnContext 固定本 Turn 的配置;SessionTask 定义工作流;RunningTask 是 Session 保存的运行句柄; TurnState 保存审批、权限请求、pending input 等可变状态。ActiveTurn 只是把“运行句柄”和“Turn 内状态” 放进同一 active slot,它本身不是后台任务。

源码位置:codex-rs/core/src/state/turn.rs,符号 ActiveTurn、TaskKind 与 RunningTask。

rust
/// Metadata about the currently running turn.
pub(crate) struct ActiveTurn {
    // None既可能表示尚未安装task的reservation,也可能表示task已被finish owner取走。
    pub(crate) task: Option<RunningTask>,
    // TurnState独立放在Arc中,使输入队列与收尾逻辑能在RunningTask被取走后继续处理。
    pub(crate) turn_state: Arc<Mutex<TurnState>>,
}

#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum TaskKind {
    Regular,
    Review,
    Compact,
}

pub(crate) struct RunningTask {
    // done用于abort grace等待;它不是Task完成结果的传输channel。
    pub(crate) done: Arc<Notify>,
    pub(crate) kind: TaskKind,
    // 类型擦除后的实现保留到abort hook完成,不能只保存JoinHandle。
    pub(crate) task: Arc<dyn AnySessionTask>,
    pub(crate) cancellation_token: CancellationToken,
    // RunningTask被意外drop时AbortOnDropHandle仍会终止Tokio task。
    pub(crate) handle: AbortOnDropHandle<()>,
    pub(crate) turn_context: Arc<TurnContext>,
    pub(crate) _agent_execution_guard: Option<AgentExecutionGuard>,
    // Timer recorded when the task drops to capture the full turn duration.
    pub(crate) _timer: Option<codex_otel::Timer>,
}

TaskKind 只有三类,但 SessionTask 有四个当前实现:UserShellCommandTask 的 kind 是 Regular。因此 kind 是调度/遥测分类,不是 Rust 具体类型的完整判别枚举;调试时还要结合 span_name()。

2. SessionTask ​

SessionTask::run 使用 RPITIT 风格返回 future,不适合直接作为 trait object 方法调用;而 RunningTask 又必须 在同一个字段中保存不同实现。AnySessionTask 将 future 装箱为 BoxFuture,完成类型擦除,同时保留可动态 调用的 abort()。

源码位置:codex-rs/core/src/tasks/mod.rs,符号 SessionTask、AnySessionTask 及 blanket impl。

rust
pub(crate) type SessionTaskResult = CodexResult<Option<String>>;

pub(crate) trait SessionTask: Send + Sync + 'static {
    fn kind(&self) -> TaskKind;
    fn span_name(&self) -> &'static str;

    fn run(
        self: Arc<Self>,
        session: Arc<Session>,
        ctx: Arc<TurnContext>,
        input: Vec<TurnInput>,
        cancellation_token: CancellationToken,
    ) -> impl std::future::Future<Output = SessionTaskResult> + Send;

    fn abort(
        &self,
        session: Arc<Session>,
        ctx: Arc<TurnContext>,
    ) -> impl std::future::Future<Output = ()> + Send {
        // 默认abort为空;只有需要补充业务清理的Task覆盖它。
        async move {
            let _ = (session, ctx);
        }
    }
}

pub(crate) trait AnySessionTask: Send + Sync + 'static {
    fn kind(&self) -> TaskKind;
    fn span_name(&self) -> &'static str;
    fn run(
        self: Arc<Self>,
        session: Arc<Session>,
        ctx: Arc<TurnContext>,
        input: Vec<TurnInput>,
        cancellation_token: CancellationToken,
    ) -> BoxFuture<'static, SessionTaskResult>;
    fn abort<'a>(
        &'a self,
        session: Arc<Session>,
        ctx: Arc<TurnContext>,
    ) -> BoxFuture<'a, ()>;
}

impl<T> AnySessionTask for T
where
    T: SessionTask,
{
    fn kind(&self) -> TaskKind {
        SessionTask::kind(self)
    }

    fn span_name(&self) -> &'static str {
        SessionTask::span_name(self)
    }

    fn run(
        self: Arc<Self>,
        session: Arc<Session>,
        ctx: Arc<TurnContext>,
        input: Vec<TurnInput>,
        cancellation_token: CancellationToken,
    ) -> BoxFuture<'static, SessionTaskResult> {
        // blanket impl只做future装箱,不改变具体Task的返回和取消语义。
        Box::pin(SessionTask::run(
            self,
            session,
            ctx,
            input,
            cancellation_token,
        ))
    }

    fn abort<'a>(
        &'a self,
        session: Arc<Session>,
        ctx: Arc<TurnContext>,
    ) -> BoxFuture<'a, ()> {
        Box::pin(SessionTask::abort(self, session, ctx))
    }
}

SessionTaskResult 的三种语义是:Ok(Some(message)) 正常完成并带最终文本;Ok(None) 正常完成但没有最终 文本;Err(CodexErr::TurnAborted) 请求 aborted lifecycle。其他 Err 会在统一收尾中转换成 ErrorEvent, 随后仍使用 TurnComplete envelope。

3. 实现与TaskKind ​

实现TaskKindspan_namerun()核心职责自定义abort()
RegularTaskRegularsession_task.turn多轮模型与工具主循环否
CompactTaskCompactsession_task.compactlocal/remote/token-budget压缩否
ReviewTaskReviewsession_task.review启动review子会话并转发结果是,退出review mode
UserShellCommandTaskRegularsession_task.user_shellstandalone shell turn否

源码位置:codex-rs/core/src/tasks/review.rs,符号 ReviewTask 的 run 与 abort。

rust
impl SessionTask for ReviewTask {
    fn kind(&self) -> TaskKind {
        TaskKind::Review
    }

    fn span_name(&self) -> &'static str {
        "session_task.review"
    }

    async fn run(
        self: Arc<Self>,
        session: Arc<Session>,
        ctx: Arc<TurnContext>,
        input: Vec<TurnInput>,
        cancellation_token: CancellationToken,
    ) -> SessionTaskResult {
        // ...收集UserInput并启动review子会话
        let output = match start_review_conversation(
            session.clone(),
            ctx.clone(),
            user_input,
            cancellation_token.clone(),
        )
        .await
        {
            Some(receiver) => {
                process_review_events(session.clone(), ctx.clone(), receiver).await
            }
            None => None,
        };
        if !cancellation_token.is_cancelled() {
            // 正常路径退出review mode;取消路径留给abort hook做同一清理,避免重复事件。
            exit_review_mode(Arc::clone(&session), output.clone(), ctx.clone()).await;
        }
        Ok(None)
    }

    async fn abort(&self, session: Arc<Session>, ctx: Arc<TurnContext>) {
        // review是当前唯一覆盖abort的实现,保证被替换/中断时仍发退出review状态。
        exit_review_mode(session, /*review_output*/ None, ctx).await;
    }
}

CompactTask 接收 cancellation token 却不直接读取它;它调用的压缩实现负责传播 TurnAborted,超出 grace 时 外层仍会 abort Tokio handle。UserShellCommandTask 则把 token 传入执行函数。这个差异说明 trait 只规定取消 能力与最终强制手段,不要求每个实现用同一种轮询写法。

4. spawn_task ​

spawn_task() 是普通替换入口:先用 Replaced 中止旧 Task,清除 connector selection,再调用 start_task()。旧 Task 的 terminal event 和清理由 abort owner 完成后,新 Task 才开始安装。

当前 start_task() 还会在安装 RunningTask 前把 pending input 合并进 TurnState,记录 turn-start token baseline,并创建带 turn tracing span 的 Tokio task。RunningTask 中的取消 token、AbortOnDropHandle、计时器 和 diagnostics guard 因而覆盖完整 Turn 生命周期,而不只是模型 stream 的局部执行。

源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::spawn_task。

rust
pub async fn spawn_task<T: SessionTask>(
    self: &Arc<Self>,
    turn_context: Arc<TurnContext>,
    input: Vec<TurnInput>,
    task: T,
) {
    // 新任务不会与旧RunningTask并存;Replaced与用户Interrupt是不同abort reason。
    self.abort_all_tasks(TurnAbortReason::Replaced).await;
    self.clear_connector_selection().await;
    self.start_task(
        turn_context,
        input,
        task,
        MailboxParentProvenance::Ignore,
    )
    .await;
}

合成 mailbox Turn 不调用这个固定 provenance,而是用 MailboxParentProvenance::Attribute 将触发消息的父 Turn 写入 metadata。Task lifecycle 因而也承担 lineage 安装点,但不会把 mailbox 内容复制进 TurnContext 字段。

5. start_task事务 ​

start_task() 先创建类型擦除对象、独立 cancellation token 和 done Notify,然后取得或建立 ActiveTurn 的 TurnState,将 Session 级 pending input 转移进去并发送 extension turn-start lifecycle。第二次取得 active lock 后,才 spawn Tokio future 并把 RunningTask 写入 slot。

源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::start_task 的准备阶段。

rust
let task: Arc<dyn AnySessionTask> = Arc::new(task);
let task_kind = task.kind();
let span_name = task.span_name();
let started_at = Instant::now();
let turn_started_at_unix_ms = turn_context
    .turn_timing_state
    .mark_turn_started(started_at)
    .await;
turn_context
    .turn_metadata_state
    .set_turn_started_at_unix_ms(turn_started_at_unix_ms);
let token_usage_at_turn_start = self.total_token_usage().await.unwrap_or_default();

// root token归RunningTask所有;传给实现的是child,外部abort只需cancel root。
let cancellation_token = CancellationToken::new();
let done = Arc::new(Notify::new());

let (pending_items, parent_turn_id) =
    self.input_queue.get_pending_input(&self.active_turn).await;
if let (MailboxParentProvenance::Attribute, Some(id)) =
    (mailbox_parent_provenance, parent_turn_id)
{
    turn_context.turn_metadata_state.set_parent_turn_id(id);
}
let turn_state = {
    let mut active = self.active_turn.lock().await;
    let turn = active.get_or_insert_with(ActiveTurn::default);
    // ActiveTurn可先作为reservation存在,但同一slot不能已有RunningTask。
    debug_assert!(turn.task.is_none());
    Arc::clone(&turn.turn_state)
};
turn_state.lock().await.token_usage_at_turn_start = token_usage_at_turn_start.clone();
self.input_queue
    .extend_pending_input_for_turn_state(turn_state.as_ref(), pending_items)
    .await;
self.emit_turn_start_lifecycle(turn_context.as_ref(), &token_usage_at_turn_start)
    .await;

这里的 turn-start lifecycle 是 extension 回调,不等同于协议 TurnStarted。RegularTask 和 standalone UserShell 在各自 run() 中发送协议事件;Compact 有意不在压缩前注入普通 Turn 上下文。两层事件不能用同一 名字推断顺序。

6. Tokio闭包 ​

spawned future 先运行具体 Task,随后 flush rollout;只有 root token 未取消时才进入 on_task_finished()。 无论正常、错误还是已被取消,闭包最后都 notify_waiters(),供 abort owner 的 grace wait 观察。

源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::start_task 的 spawn 与注册阶段。

rust
let done_clone = Arc::clone(&done);
let session = Arc::clone(self);
let ctx = Arc::clone(&turn_context);
let task_for_run = Arc::clone(&task);
let task_input = input;
let task_cancellation_token = cancellation_token.child_token();

let handle = tokio::spawn(
    async move {
        let ctx_for_finish = Arc::clone(&ctx);
        let task_result = task_for_run
            .run(
                Arc::clone(&session),
                ctx,
                task_input,
                // 再派生child,具体Task不能取消RunningTask root token。
                task_cancellation_token.child_token(),
            )
            .instrument(trace_span!("session_task.run"))
            .await;
        let sess = Arc::clone(&session);
        // 普通items先flush,terminal event稍后还会有第二个flush barrier。
        if let Err(err) = sess.flush_rollout().await {
            warn!("failed to flush rollout before completing turn: {err}");
            sess.send_event(
                ctx_for_finish.as_ref(),
                EventMsg::Warning(WarningEvent {
                    message: format!(
                        "Failed to save the conversation transcript; Codex will continue retrying. Error: {err}"
                    ),
                }),
            )
            .await;
        }
        if !task_cancellation_token.is_cancelled() {
            // 外部abort取消root后由abort owner发终态;此处必须跳过,防止双terminal event。
            sess.on_task_finished(Arc::clone(&ctx_for_finish), task_result)
                .await;
        }
        done_clone.notify_waiters();
    }
    .instrument(task_span),
);

let timer = turn_context
    .session_telemetry
    .start_timer(TURN_E2E_DURATION_METRIC, &[])
    .ok();
let running_task = RunningTask {
    done,
    handle: AbortOnDropHandle::new(handle),
    kind: task_kind,
    task,
    cancellation_token,
    turn_context: Arc::clone(&turn_context),
    _agent_execution_guard: agent_execution_guard,
    _timer: timer,
};
// RunningTask写入active slot后,Session才拥有取消与终止该future的完整句柄集合。
turn.task = Some(running_task);

第一道屏障是 task body 后的 rollout flush,保证普通历史尽量先持久化;第二道是 terminal event 发送后的 flush,保证客户端收到终态后重读 rollout 能看到对应事件。done 不是持久化屏障,它只表示 spawn closure 已经走到最后。

7. 正常完成 Owner ​

on_task_finished() 先把结果归一化,再从 active_turn.task 中 take() RunningTask 并 detach handle。若 task 已经被外部 abort owner 取走,这里得到 None 并立即返回,不再发送任何终态。

源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::on_task_finished 的 owner 获取。

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) => {
        // Task自发返回TurnAborted时,在正常finish owner中归约为Interrupted。
        (None, Some(TurnAbortReason::Interrupted))
    }
    Err(err) => {
        warn!(%err, "session task returned an unexpected error");
        self.emit_turn_error_lifecycle(
            turn_context.as_ref(),
            err.to_codex_protocol_error(),
        )
        .await;
        self.track_turn_codex_error(turn_context.as_ref(), &err);
        self.send_event(
            turn_context.as_ref(),
            EventMsg::Error(err.to_error_event(/*message_prefix*/ None)),
        )
        .await;
        // 普通错误没有abort reason;ErrorEvent进入terminal_error后仍生成TurnComplete。
        (None, None)
    }
};
turn_context
    .turn_metadata_state
    .cancel_git_enrichment_task();

let turn_state = {
    let mut active = self.active_turn.lock().await;
    active.as_mut().and_then(|active_turn| {
        // take成功者成为finish owner;abort路径若先take,这里返回None。
        let task = active_turn.task.take()?;
        task.handle.detach();
        Some(Arc::clone(&active_turn.turn_state))
    })
};
let Some(turn_state) = turn_state else {
    return;
};

detach 很重要:正常完成后 RunningTask 即将 drop,如果不 detach,AbortOnDropHandle 会对已经完成或正在收尾的 Tokio task 再调用 abort。owner 获取在 active lock 内完成,但持久化、metrics 和 extension hooks 都在锁外, 避免长时间占用 Session 的 active slot mutex。

8. Active slot ​

正常 owner 会取出 TurnState 中尚未消费的 input,经 hook 检查后写入 history;然后生成 metrics 与 terminal event。最后只有 active slot 仍指向同一 turn_state 且 task=None 时才清空 ActiveTurn。

源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::on_task_finished 的终态与清理阶段。

rust
let pending_input = self
    .input_queue
    .take_pending_input_for_turn_state(turn_state.as_ref())
    .await;
if !pending_input.is_empty() {
    for pending_input_item in pending_input {
        let hook_outcome =
            inspect_pending_input(self, &turn_context, &pending_input_item).await;
        if hook_outcome.should_stop {
            record_additional_contexts(
                self,
                &turn_context,
                hook_outcome.additional_contexts,
            )
            .await;
        } else {
            // Task返回边界遗留的输入仍进入history,不会因RunningTask结束而静默丢失。
            record_pending_input(
                self,
                &turn_context,
                pending_input_item,
                hook_outcome.additional_contexts,
            )
            .await;
        }
    }
}

// ...计算turn token/tool/memory metrics与timing

let event = if let Some(reason) = abort_reason {
    self.emit_turn_abort_lifecycle(
        reason.clone(),
        turn_context.extension_data.as_ref(),
    )
    .await;
    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;
    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;

let cleared_active_turn = {
    let mut active = self.active_turn.lock().await;
    if let Some(active_turn) = active.as_ref()
        && active_turn.task.is_none()
        // 指针身份防止旧finish逻辑清掉后来安装、但恰好也task=None的新ActiveTurn。
        && Arc::ptr_eq(&active_turn.turn_state, &turn_state)
    {
        *active = None;
        true
    } else {
        false
    }
};
if cleared_active_turn {
    self.emit_thread_idle_lifecycle_if_idle().await;
}
if let Err(err) = self.flush_rollout().await {
    warn!("failed to flush rollout after emitting terminal turn event: {err}");
}
if cleared_active_turn {
    // 只有本owner真正清空slot后,trigger mailbox才有资格启动下一Task。
    self.maybe_start_turn_for_pending_work().await;
}

Arc::ptr_eq 比较的是 TurnState owner,而不是 Turn ID 字符串。即使未来存在 ID 重用或新 reservation,旧收尾也 不能仅凭相等字符串删除新 slot。

9. 外部 Abort Owner ​

abort_all_tasks() 先把整个 ActiveTurn 从 mutex 中 take(),再从中取 RunningTask。这样 spawn closure 即使 同时结束,on_task_finished() 也无法再次取得 task owner。随后先 cancel token,最多等待 100ms done;无论 是否协作完成,都 abort JoinHandle,再调用具体 Task 的 abort hook。

源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::handle_task_abort。

rust
async fn handle_task_abort(
    self: &Arc<Self>,
    task: RunningTask,
    reason: TurnAbortReason,
) {
    let sub_id = task.turn_context.sub_id.clone();
    if task.cancellation_token.is_cancelled() {
        // root已取消表示已有abort owner,不重复发终态。
        return;
    }

    task.cancellation_token.cancel();
    task.turn_context
        .turn_metadata_state
        .cancel_git_enrichment_task();
    let session_task = task.task;

    select! {
        _ = task.done.notified() => {},
        _ = tokio::time::sleep(
            Duration::from_millis(GRACEFULL_INTERRUPTION_TIMEOUT_MS)
        ) => {
            warn!(
                "task {sub_id} didn't complete gracefully after {}ms",
                GRACEFULL_INTERRUPTION_TIMEOUT_MS
            );
        }
    }

    // 即使body已协作返回也abort handle,封住closure尚未结束的flush/finish窗口。
    task.handle.abort();
    session_task
        .abort(Arc::clone(self), Arc::clone(&task.turn_context))
        .await;

    if reason == TurnAbortReason::Interrupted
        && let Some(marker) = interrupted_turn_history_marker(
            InterruptedTurnHistoryMarker::from_config_and_version(
                task.turn_context.config.as_ref(),
                task.turn_context.multi_agent_version,
            ),
        )
    {
        self.record_conversation_items(
            task.turn_context.as_ref(),
            std::slice::from_ref(&marker),
        )
        .await;
        // marker先flush,客户端收到TurnAborted后同步重读即可看到中断事实。
        if let Err(err) = self.flush_rollout().await {
            warn!(
                "failed to flush interrupted-turn marker before emitting TurnAborted: {err}"
            );
        }
    }

    let started_at = task
        .turn_context
        .turn_timing_state
        .started_at_unix_secs()
        .await;
    let (completed_at, duration_ms, profile) = task
        .turn_context
        .turn_timing_state
        .complete_profile_and_duration_ms()
        .await;
    self.services
        .analytics_events_client
        .track_turn_profile(TurnProfileFact {
            turn_id: task.turn_context.sub_id.clone(),
            profile,
        });
    let event = EventMsg::TurnAborted(TurnAbortedEvent {
        turn_id: Some(task.turn_context.sub_id.clone()),
        reason,
        started_at,
        completed_at,
        duration_ms,
    });
    self.send_event(task.turn_context.as_ref(), event).await;
    // ...清guardian circuit breaker
    if let Err(err) = self.flush_rollout().await {
        warn!("failed to flush rollout after emitting terminal turn event: {err}");
    }
}

Interrupted 才写模型可见 marker;Replaced 不写,因为替换不是用户要求模型理解的中断上下文。两者都发送 TurnAborted,只是 reason 不同。abort_all_tasks() 在 handler 返回后才清 pending waiters,避免正在等待的 审批先被 channel drop 转成模型可见 rejection,抢在 TurnAborted 之前。

10. 生命周期状态 ​

这张状态图是多个 owner 字段与返回路径的组合模型,不是源码中的 enum。真正的不变量是:同一时刻只有一个 路径能从 active slot 取走 RunningTask;正常 owner 发送 TurnComplete/自发 abort,外部 abort owner 发送 Interrupt/Replaced 的 TurnAborted。

11. 生命周期测试 ​

11.1 正常完成flush ​

源码位置:codex-rs/core/src/session/tests.rs,测试 turn_complete_flushes_terminal_event_after_delivery。

rust
let (mut sess, tc, rx) = make_session_and_context_with_rx().await;
let store = attach_in_memory_thread_store(
    Arc::get_mut(&mut sess).expect("session should be uniquely owned"),
)
.await;

sess.spawn_task(Arc::clone(&tc), input, CompletingTask).await;
let event = recv_terminal_event(
    &rx,
    TerminalEventKind::TurnComplete,
)
.await;
assert!(matches!(event.msg, EventMsg::TurnComplete(_)));

// 第一次flush位于task body之后,第二次位于TurnComplete发送之后。
let calls = wait_for_flush_count(&store, /*expected_flushes*/ 2).await;
assert_eq!(2, calls.flush_thread);

这个测试使用立即返回 Ok(None) 的 CompletingTask,排除了模型和工具噪声,直接证明统一 spawn/finish plumbing 建立两道持久化屏障。

该测试不等于覆盖真实模型、工具调用和并发输入路径;它只证明正常完成分支中的终态事件与两次 flush 顺序。

11.2 Interrupted ​

源码位置:codex-rs/core/src/session/tests.rs,测试 turn_aborted_flushes_terminal_event_after_delivery。

rust
sess.spawn_task(
    Arc::clone(&tc),
    input,
    NeverEndingTask {
        kind: TaskKind::Regular,
        listen_to_cancellation_token: true,
    },
)
.await;

let abort_task = tokio::spawn({
    let sess = Arc::clone(&sess);
    async move {
        sess.abort_all_tasks(TurnAbortReason::Interrupted).await;
    }
});

let event = recv_terminal_event(&rx, TerminalEventKind::TurnAborted).await;
match event.msg {
    EventMsg::TurnAborted(e) => {
        assert_eq!(TurnAbortReason::Interrupted, e.reason)
    }
    other => panic!("unexpected event: {other:?}"),
}
abort_task.await.expect("abort task should finish");

// task body flush、interrupted marker flush、terminal event flush各一次。
let calls = wait_for_flush_count(&store, /*expected_flushes*/ 3).await;
assert_eq!(3, calls.flush_thread);

它同时覆盖协作取消不会让 spawn closure 再发 TurnComplete:测试只接受 TurnAborted,并精确观察三个 barrier。

11.3 不协作的Task ​

源码位置:codex-rs/core/src/session/tests.rs,测试 abort_regular_task_emits_marker_before_turn_aborted。

rust
sess.spawn_task(
    Arc::clone(&tc),
    input,
    NeverEndingTask {
        kind: TaskKind::Regular,
        // fixture故意不监听token,迫使abort路径走100ms grace再abort handle。
        listen_to_cancellation_token: false,
    },
)
.await;

sess.abort_all_tasks(TurnAbortReason::Interrupted).await;

let marker_evt = tokio::time::timeout(
    std::time::Duration::from_secs(2),
    rx.recv(),
)
.await
.expect("timeout waiting for marker event")
.expect("event");
assert!(matches!(marker_evt.msg, EventMsg::RawResponseItem(_)));

let evt = tokio::time::timeout(
    std::time::Duration::from_secs(2),
    rx.recv(),
)
.await
.expect("timeout waiting for event")
.expect("event");
match evt.msg {
    EventMsg::TurnAborted(e) => {
        assert_eq!(TurnAbortReason::Interrupted, e.reason)
    }
    other => panic!("unexpected event: {other:?}"),
}
// abort owner只发一个terminal event,超时强杀不会制造尾随完成事件。
assert!(rx.try_recv().is_err());

11.4 空槽清理 ​

源码位置:codex-rs/core/src/session/tests.rs,测试 abort_empty_active_turn_preserves_pending_input。

rust
let turn_state = {
    let mut active = sess.active_turn.lock().await;
    let active_turn = active.get_or_insert_with(ActiveTurn::default);
    Arc::clone(&active_turn.turn_state)
};
sess.input_queue
    .extend_pending_input_for_turn_state(
        turn_state.as_ref(),
        vec![TurnInput::ResponseItem(pending_item.clone())],
    )
    .await;

// ActiveTurn存在但task=None,Replaced不应把它误当作正在运行的Task清理输入。
sess.abort_all_tasks(TurnAbortReason::Replaced).await;

assert!(sess.active_turn.lock().await.is_none());
assert_eq!(
    sess.input_queue
        .take_pending_input_for_turn_state(turn_state.as_ref())
        .await,
    vec![TurnInput::ResponseItem(pending_item)]
);

这个测试区分了 ActiveTurn reservation 与 RunningTask:abort_all_tasks 可以移除空 reservation,但只有确实 取得 task 时才执行 pending clear 与 abort lifecycle。

12. Owner定位 ​

遇到 Task 生命周期问题时,可按下面顺序定位:

  1. 先看 active_turn.task 是否仍存在;ActiveTurn 存在不代表后台 future 已安装。
  2. 对照 TaskKind 与 span_name;UserShell 的 kind 也是 Regular,不能只按枚举判断实现。
  3. 正常完成缺 terminal event 时,检查 task body 后 flush、root token 是否已取消、on_task_finished 是否成功 take() task。
  4. 出现双 terminal event 时,检查外部 abort 是否先从 slot 取走 owner,以及 spawn closure 是否在 token 取消后 仍调用正常 finish。
  5. Interrupt 后 marker 缺失时,区分 reason 是否为 Replaced,并检查 marker flush 是否失败。
  6. 旧 Task 收尾影响新 Task 时,检查 active slot 清理是否使用 Arc::ptr_eq(turn_state) 身份门禁。

可以用 NeverEndingTask 的两个模式做只读推演:监听 token 时应在 grace 内 notify;不监听时应超时并 abort handle。两条路径之后都只能看到 marker(若启用)和一个 TurnAborted。普通用户输入如何进入 RegularTask 并 消费 startup prewarm,由 RegularTask完整流程 继续解释;Compact、Review 和 UserShell 的业务差异留给各自专题。

可以用下面的只读搜索把本文的 task lifecycle 主线落回源码:

bash
rg -n "SessionTask|AnySessionTask|RunningTask|abort|join" codex-rs/core/src