Skip to content

Session关闭流程

追踪 Session 如何停止任务与后台服务、建立持久化屏障,并区分中断、替换、显式关闭和通道断开。

基于rust-v0.150.0
CodexRustSessionShutdown

Session关闭流程 ​

本文默认读者了解 Rust 的 Arc、CancellationToken 和 async task,但不要求熟悉 Codex 的关闭实现。建议 先阅读 Session启动预热、 Session输入队列 和 Session运行时处理:它们分别解释本文要停止的预热任务、 Turn/mailbox 所有权和 Op::Shutdown 分派入口。

读完本文,应能从 SessionIo::shutdown_and_wait() 追踪到每种资源的 owner 和停止动作,解释为什么 TurnAborted 需要多次 rollout flush,以及区分“关闭请求已经入队”“submission loop 已退出”“客户端收到 ShutdownComplete”这三个不同事实。

本文只讨论 Session/Thread 关闭编排,不展开每种 OS 进程后端如何发送信号,也不重复解释 MCP refresh、 mailbox 和 Op 分派的完整实现;这些机制分别由前置文章承担,本文只追踪它们在 teardown 中的消费方式。

1. 四种停止状态 ​

Codex 中容易被统称为 shutdown 的动作实际有四种。Interrupt 停止当前 Task但保留 Session;Replaced 是新 Task 替换旧 Task;显式 Shutdown 关闭整个 Session 并尝试返回完成事件;所有 submission sender 被丢弃则让 loop 隐式退出并清理,但没有请求方可关联 ShutdownComplete。

动作作用范围TurnAbortReasonSession 是否继续ShutdownComplete
Interrupt当前 RunningTaskInterrupted是无
Replaced当前 RunningTaskReplaced是,新 Task 接管无
显式 Shutdown整个 Session/Thread runtime当前 Task 使用 Interrupted否尝试发送
submission channel close整个 Session/Thread runtime当前 Task 使用 Interrupted否不发送

这四条路径复用了部分函数,但语义不同。尤其是 Interrupted 可以写入模型可见的中断 marker,并允许 mailbox 触发下一 Turn;Replaced 不写中断 marker,也不启动 pending mailbox。

2. 关闭资源清单 ​

Session 没有一个能自动递归 await 全部子资源的 owner。SessionServices、Session 自身字段和持久化 LiveThread 分别持有不同资源,关闭函数必须显式按依赖顺序操作。

资源的停止协议并不统一:

资源Owner停止动作是否等待失败可见性
startup model prewarmSessionState handleabort taskawait join结束忽略 join result
RealtimeRealtimeConversationManagertake state、cancel tokenawait input/fanout taskshutdown结果被外层忽略
active TaskActiveTurn::RunningTaskcancel→100ms grace→abort→Task hook是TurnAborted、warning、flush日志
Unified execUnifiedExecManagerdrain registry、unregister approval、terminate()不确认OS进程退出无聚合错误返回
Code ModeCodeModeService标记closing、join init、session shutdown是warning,继续关闭
MCP prewarmSession字段cancel worker tokenawait JoinHandlejoin异常warning
MCP refresh/runtimeMcpRefresh/McpRuntimeacquire gate→close→runtime shutdown是gate关闭日志在后续调用可见
GuardianGuardianReviewSessionManagercancel、take sessions、逐个shutdown是child shutdown_and_wait结果忽略
SessionEnd hooksHooks serviceroot-only执行是transcript flush失败warning
LiveThreadSessionServicesmetadata flush→store shutdown是Error Event后仍完成关闭

“是否等待”必须具体到停止协议。Unified exec 等待的是 registry drain 和 termination signal 调用,不是每个 OS 子进程的 confirmed exit;Realtime、MCP worker 和 Guardian 则确实等待其 task/session 完成。

3. Shutdown时序 ​

SessionIo::shutdown_and_wait() 先 clone 一个共享 termination future,再提交 Op::Shutdown。即使 send 因 loop 已经死亡返回 InternalAgentDied,它仍等待 termination future;这让多个调用方可以安全等待同一个 loop 结束,而不是把“第二次提交失败”误判为资源尚未回收。

源码位置:codex-rs/core/src/session/mod.rs :: SessionIo::shutdown_and_wait

rust
pub(crate) async fn shutdown_and_wait(&self) -> CodexResult<()> {
    // 先clone共享future,避免submit后loop立刻退出导致失去等待句柄。
    let session_loop_termination = self.session_loop_termination.clone();
    match self.submit(Op::Shutdown).await {
        Ok(_) => {}
        // loop已关闭等价于“不需要再成功入队”,但仍必须等待它真正终止。
        Err(err) if matches!(err.details(), CodexErrorDetails::InternalAgentDied) => {}
        Err(err) => return Err(err),
    }
    session_loop_termination.await;
    Ok(())
}

SessionLoopTermination 是 Shared<BoxFuture<'static, ()>>,底层 JoinHandle 只 await 一次,多个 waiter 共享 完成结果。显式 handler 的外层顺序如下。

当前测试还覆盖“Shutdown 已经进行中”的重复等待:第二个 shutdown_and_wait() 不会再次依赖一条新的 Shutdown submission,而是等待同一个 loop termination。Guardian review session 也分为 cached 与 tracked 两类,teardown 必须分别取出并等待,否则 Session 主 loop 虽已结束,受信任 reviewer 仍可能存活。

源码位置:codex-rs/core/src/session/handlers.rs :: shutdown

rust
pub async fn shutdown(sess: &Arc<Session>, sub_id: String) -> bool {
    shutdown_session_runtime(sess).await;
    let history = sess.clone_history().await;
    let turn_count = history
        .raw_items()
        .iter()
        // contextual user wrappers不算真实会话Turn。
        .filter(|item| is_user_turn_boundary(item))
        .count();
    sess.services.session_telemetry.counter(
        "codex.conversation.turn.count",
        i64::try_from(turn_count).unwrap_or(0),
        &[],
    );

    emit_thread_stop_lifecycle(sess.as_ref()).await;
    if let Some(live_thread) = sess.live_thread()
        && let Err(e) = live_thread.shutdown().await
    {
        warn!("failed to shutdown thread persistence: {e}");
        // Store失败不阻止协议结束;Error发送本身也是best effort。
        sess.send_event_raw(Event {
            id: sub_id.clone(),
            msg: EventMsg::Error(ErrorEvent {
                message: "Failed to shutdown thread persistence".to_string(),
                codex_error_info: Some(CodexErrorInfo::Other),
            }),
        })
        .await;
    }

    let event = Event { id: sub_id, msg: EventMsg::ShutdownComplete };
    // writer已关闭,不能再用会追加rollout的send_event_raw。
    sess.services.rollout_thread_trace.record_protocol_event(&event.msg);
    sess.deliver_event_raw(event).await;
    sess.services
        .rollout_thread_trace
        .record_ended(codex_rollout_trace::RolloutStatus::Completed);
    true
}

ShutdownComplete 表示显式 handler 已走到末尾,不保证客户端一定收到:若 event receiver 已关闭, deliver_event_raw() 只记录 debug 并丢弃事件。调用方需要等待 session_loop_termination 判断本地关闭完成, 不能只等待 event stream。

4. RunningTask ​

4.1 ActiveTurn ​

abort_all_tasks() 先 take_active_turn(),让后续查询立即看不到旧 active slot;再取出 RunningTask 并执行 handle_task_abort()。Task abort 完成后才运行 extension turn-abort lifecycle,最后清 TurnState waiter 和 pending input。

源码位置:codex-rs/core/src/tasks/mod.rs :: Session::abort_all_tasks

rust
pub async fn abort_all_tasks(self: &Arc<Self>, reason: TurnAbortReason) {
    let mut aborted_turn = false;
    let mut active_turn_to_clear = None;
    let mut turn_context = None;
    if let Some(mut active_turn) = self.take_active_turn().await {
        let task = active_turn.task.take();
        aborted_turn = task.is_some();
        turn_context = task.as_ref().map(|task| Arc::clone(&task.turn_context));
        if let Some(task) = task {
            // RunningTask的终端事件由handle_task_abort负责。
            self.handle_task_abort(task, reason.clone()).await;
        }
        if aborted_turn {
            active_turn_to_clear = Some(active_turn);
        }
    }
    if let Some(turn_context) = turn_context.as_deref() {
        self.emit_turn_abort_lifecycle(
            reason.clone(),
            turn_context.extension_data.as_ref(),
        )
        .await;
    }
    if let Some(active_turn) = active_turn_to_clear {
        // 先让Task观察cancel,再drop approval/input waiters,避免生成伪拒绝结果。
        self.input_queue.clear_pending(&active_turn).await;
    }
    if reason == TurnAbortReason::Interrupted && aborted_turn {
        self.maybe_start_turn_for_pending_work().await;
    }
}

ActiveTurn 可能处于 reservation 状态:slot 存在但 task=None。这种情况下 aborted_turn=false,原 TurnState pending input 不会被 clear_pending() 删除。测试 abort_empty_active_turn_preserves_pending_input 专门锁定这一点。

4.2 Grace Period ​

handle_task_abort() 首先取消根 token并取消 Git enrichment;最多等待 100ms 的 done notification;随后 无条件 abort Tokio handle,再调用具体 SessionTask::abort()。ReviewTask 利用 abort hook写回退出 Review 模式;普通 Task 默认 no-op。

源码位置:codex-rs/core/src/tasks/mod.rs :: Session::handle_task_abort(取消核心)

rust
if task.cancellation_token.is_cancelled() {
    // 重复abort不再发送第二个TurnAborted。
    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,
        );
    }
}

// 即使done已到达,abort已完成的handle也是安全的;owner随后执行类型化清理。
task.handle.abort();
session_task
    .abort(Arc::clone(self), Arc::clone(&task.turn_context))
    .await;

grace window 只等待 SessionTask wrapper 的 done,并不承诺每个工具自行关闭。外层 shutdown 还要显式 终止 Unified exec、Code Mode、MCP 和 Guardian,这正是资源账本不能省略的原因。

4.3 Interrupted ​

只有 Interrupted 可能写 <turn_aborted> 模型可见 marker;marker role由 MultiAgentVersion 决定,V2 使用 developer message,旧版本使用 contextual user message。marker 必须在 TurnAborted 之前 flush,因为 客户端收到 abort event 后可能同步重读 rollout。

Task body自身结束时 wrapper会做一次 flush;abort path再做 marker flush和 terminal-event flush。因此一个 正常观察 cancellation 的 interrupted Task测试期望三次 flush,而正常完成只期望两次。

源码位置:codex-rs/core/src/tasks/mod.rs :: Session::handle_task_abort(marker与terminal barrier)

rust
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在TurnAborted可见前建立durability barrier。
    if let Err(err) = self.flush_rollout().await {
        warn!("failed to flush interrupted-turn marker before emitting TurnAborted: {err}");
    }
}

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;
// terminal event追加后再flush,供恢复路径观察完整终态。
if let Err(err) = self.flush_rollout().await {
    warn!("failed to flush rollout after emitting terminal turn event: {err}");
}

4.4 Shutdown ​

shutdown_session_runtime() 调用 abort_all_tasks(Interrupted)。该函数为普通 Interrupt 设计:中断后如果 Session mailbox 有 trigger work,会调用 maybe_start_turn_for_pending_work()。当前 shutdown 函数没有在这 一步前设置 closing flag,也没有随后第二次 abort ActiveTurn。

这里不要求两个 submission handler 真正并发。submission_loop 串行处理 Op,但下面的顺序已经足以形成 待关闭路径:

  1. Session 仍有 RunningTask 时,先处理 Op::InterAgentCommunication,把 trigger_turn=true 的消息留在 mailbox;pending-work scheduler 因 active slot 非空而不启动新 Turn。
  2. 紧接着处理 Op::Shutdown;abort_all_tasks(Interrupted) 移除旧 RunningTask,随后再次调用 maybe_start_turn_for_pending_work()。
  3. scheduler 此时看到 Session idle 且 trigger mail 仍在队列中,于是可以创建新的 RegularTask;外层 teardown 随后继续关闭 Unified exec、Code Mode 和 MCP runtime,却没有再次终止这个新 Task。

尚未由专门回归测试覆盖

以上结论是该版本控制流能够支持的风险推断,不等同于已经由回归测试复现的产品故障。该版本没有 覆盖“pending trigger mail 先入队,再 Shutdown”的专门测试,也没有 Session 级 closing gate 或 teardown 末尾的第二次 abort。因此不能宣称关闭开始后绝不会创建新 Turn;这个行为仍需源码修复或专门测试锁定。

5. Realtime等待 ​

5.1 Realtime任务 ​

Realtime shutdown 在 state Mutex 内只执行 take(),随后释放锁并停止 ConversationState。它设置 realtime_active=false、取消 stop token、await input task,再按 shutdown 模式 await fanout task。这样新调用 立即看到 conversation 不再运行,慢 task join也不会长期持有 state lock。

源码位置:codex-rs/core/src/realtime_conversation.rs

rust
// :: RealtimeConversationManager::shutdown, stop_conversation_state
pub(crate) async fn shutdown(&self) -> CodexResult<()> {
    let state = {
        let mut guard = self.state.lock().await;
        // 先从共享状态移除,阻止新的audio/text继续投递旧会话。
        guard.take()
    };
    if let Some(state) = state {
        stop_conversation_state(state, RealtimeFanoutTaskStop::Await).await;
    }
    Ok(())
}

async fn stop_conversation_state(
    mut state: ConversationState,
    fanout_task_stop: RealtimeFanoutTaskStop,
) {
    state.realtime_active.store(false, Ordering::Relaxed);
    state.stop_token.cancel();
    let _ = state.input_task.await;
    if let Some(fanout_task) = state.fanout_task.take() {
        match fanout_task_stop {
            RealtimeFanoutTaskStop::Await => { let _ = fanout_task.await; }
            RealtimeFanoutTaskStop::Detach => {}
        }
    }
}

外层 shutdown_session_runtime() 忽略 Realtime 的 CodexResult,不过当前 shutdown() 实现始终返回 Ok; 具体 task join error也被丢弃。这里的保证是“等待 task结束”,不是“向客户端报告每个 Realtime 清理错误”。

5.2 Unified ​

Unified exec manager 在锁内 drain全部 ProcessEntry并清 reserved IDs,释放锁后逐项注销 network approval, 再调用同步 process.terminate()。函数不调用 terminate_confirmed().await,因此返回时 registry 已清空且终止 请求已发出,但 OS 进程退出 watcher可能仍在收尾。

源码位置:codex-rs/core/src/unified_exec/process_manager.rs

rust
// :: UnifiedExecManager::terminate_all_processes
pub(crate) async fn terminate_all_processes(&self) {
    let entries: Vec<ProcessEntry> = {
        let mut processes = self.process_store.lock().await;
        // 先转移owner并清reservation,后续查询不会再把旧进程视为managed active。
        let entries = processes.processes.drain().map(|(_, entry)| entry).collect();
        processes.reserved_process_ids.clear();
        entries
    };

    for entry in entries {
        unregister_network_approval_for_entry(&entry).await;
        // 这里只发终止,不等待confirmed exit;与terminate_process语义不同。
        entry.process.terminate();
    }
}

这解释了为什么“Session teardown完成”不应被解释为“所有系统进程已被 wait/reap”。Manager已经放弃运行时 所有权并发送终止,具体 Process实现和 watcher负责最终 OS 状态。

5.3 Code Mode 防止 ​

CodeModeService 使用 AtomicBool shutting_down 加 Tokio OnceCell。shutdown先发布 closing,再用 get_or_try_init 等待正在进行的初始化;若从未初始化,closure返回错误而不会为了关闭创建一个无用 host。 即使 create_session 在 closing前开始、closing后返回,session() 也会立刻关闭新 session并拒绝交付。

源码位置:codex-rs/core/src/tools/code_mode/mod.rs :: CodeModeService::shutdown, session

rust
pub(crate) async fn shutdown(&self) -> Result<(), String> {
    self.shutting_down.store(true, Ordering::Release);
    match self
        .session
        .get_or_try_init(|| async {
            // 未初始化时用错误占位,不启动仅为shutdown服务的host。
            Err::<Arc<dyn CodeModeSession>, String>(
                "code mode session is shutting down".to_string(),
            )
        })
        .await
    {
        Ok(session) => session.shutdown().await,
        Err(_) => Ok(()),
    }
}

async fn session(&self) -> Result<Arc<dyn CodeModeSession>, String> {
    if self.shutting_down.load(Ordering::Acquire) {
        return Err("code mode session is shutting down".to_string());
    }
    self.session
        .get_or_try_init(|| async {
            if self.shutting_down.load(Ordering::Acquire) {
                return Err("code mode session is shutting down".to_string());
            }
            let session = self.session_provider
                .create_session(self.dispatch_broker.clone()).await?;
            if self.shutting_down.load(Ordering::Acquire) {
                // 初始化与shutdown竞态时不发布晚到session。
                let _ = session.shutdown().await;
                return Err("code mode session is shutting down".to_string());
            }
            Ok(session)
        })
        .await
        .map(Arc::clone)
}

外层对 Code Mode shutdown error只 warning并继续,以免辅助执行服务阻止 MCP、hooks 和 Store回收。

6. MCP 关闭 ​

MCP prewarm worker可以因 refresh request或auth watch启动 refresh_mcp_if_dirty();MCP step/tool consumer也可 直接走同一个 refresh gate。关闭时先 cancel并join worker,随后 acquire单 permit gate等待已在进行的刷新, 再 close semaphore并 shutdown runtime。

源码位置:codex-rs/core/src/session/mcp_prewarm.rs

rust
// :: Session::stop_mcp_prewarm_worker
pub(super) async fn stop_mcp_prewarm_worker(&self) {
    self.mcp_prewarm_shutdown.cancel();
    let worker = self
        .mcp_prewarm_task
        .lock()
        .unwrap_or_else(std::sync::PoisonError::into_inner)
        // JoinHandle在同步锁内take,await发生在锁外。
        .take();
    if let Some(worker) = worker
        && let Err(error) = worker.await
    {
        warn!(%error, "MCP prewarm worker stopped unexpectedly");
    }
}

如果先 shutdown runtime再停 worker,worker可能在关闭期间重新计算并发布 binding;如果先 close gate但不 等待旧 permit,in-flight refresh仍可能持有发布权。当前顺序同时关闭了这两个窗口。

7. Guardian资源 ​

7.1 Guardian关闭 ​

Guardian manager先 cancel自己的初始化 token,在锁内 take trunk和全部 ephemeral reviews,释放锁后逐个 shutdown。每个 GuardianReviewSession再 cancel review token并调用 child SessionIo::shutdown_and_wait()。

源码位置:codex-rs/core/src/guardian/review_session.rs

rust
// :: GuardianReviewSessionManager::shutdown
pub(crate) async fn shutdown(&self) {
    self.cancellation_token.cancel();
    let (review_session, ephemeral_reviews) = {
        let mut state = self.state.lock().await;
        // 先移出manager状态,避免持锁递归关闭child Session。
        (
            state.trunk.take(),
            std::mem::take(&mut state.ephemeral_reviews),
        )
    };
    if let Some(review_session) = review_session {
        review_session.shutdown().await;
    }
    for review_session in ephemeral_reviews {
        review_session.shutdown().await;
    }
}

manager按顺序等待 ephemeral reviews,而不是并行 join;当前实现优先简单、确定的所有权释放顺序。child shutdown错误在 GuardianReviewSession::shutdown() 中被忽略,因此父 Session不会因 reviewer关闭错误卡死。

7.2 SessionEnd hook ​

SessionEnd hooks位于 Guardian/MCP关闭之后、thread-stop extension之前。它只对 root Session运行;执行前 获取/物化 transcript path并 flush rollout。flush失败 warning,但 hook仍运行;hook control output不再改变 已关闭中的 Session。

源码位置:codex-rs/core/src/hook_runtime.rs :: run_session_end_hooks

rust
let hooks = sess.hooks();
let preview_runs = hooks.preview_session_end();
if preview_runs.is_empty() {
    return;
}
let turn_context = sess.new_default_turn().await;
if matches!(&turn_context.session_source, SessionSource::SubAgent(_)) {
    // subagent使用SubagentStop等生命周期,不运行用户SessionEnd hook。
    return;
}
let request = SessionEndRequest {
    session_id: sess.session_id().into(),
    turn_id: turn_context.sub_id.clone(),
    #[allow(deprecated)]
    cwd: turn_context.cwd.clone(),
    transcript_path: sess.hook_transcript_path().await,
};
if let Err(err) = sess.flush_rollout().await {
    warn!("failed to flush transcript before SessionEnd hook: {err}");
}
emit_hook_started_events(sess, &turn_context, preview_runs).await;
let outcome = hooks.run_session_end(request).await;
emit_hook_completed_events(sess, &turn_context, outcome.hook_events).await;

thread-stop extensions随后仍可读取 session_extension_data 和 thread_extension_data;此时 runtime资源已经 停止,但 LiveThread writer尚未关闭。最后才执行 Store shutdown。

7.3 LiveThread关闭 ​

LiveThread::shutdown() 先 flush pending metadata update,再调用 store-specific shutdown。Local paginated store会 shutdown recorder、尝试把 durable rollout materialize到 SQLite,再同步 rollout path;其中 SQLite projection失败只 warning,但 recorder shutdown或path同步失败会向上返回错误。

源码位置:codex-rs/thread-store/src/live_thread.rs :: LiveThread::shutdown

rust
pub async fn shutdown(&self) -> ThreadStoreResult<()> {
    // metadata必须先写入仍然存活的store writer。
    self.flush_pending_metadata_update_for_existing_history().await?;
    self.thread_store.shutdown_thread(self.thread_id).await
}

Store关闭后不能再追加 ShutdownComplete。显式 handler手动记录 rollout trace并调用 deliver_event_raw(); 测试通过 InMemoryThreadStore call counter确认只有一次 shutdown_thread,完成事件没有触发 append。

8. Submission ​

submission_loop 使用 while let Ok(sub) = rx_sub.recv().await。显式 Shutdown返回 true并设置 shutdown_received;sender全部被丢弃时 recv返回Err,loop在尾部执行同一 runtime teardown、thread-stop和 LiveThread shutdown。

隐式路径与显式路径有三个差异:不统计 conversation turn count、不发送 ShutdownComplete、不调用 rollout trace record_ended;但资源回收、turn abort、thread-stop和Store shutdown仍然执行。真实 loop 在完整 match sub.op 之后使用下面的尾部收敛两种退出原因:

源码位置:codex-rs/core/src/session/handlers.rs :: submission_loop channel-close tail

rust
// 上方完整match结束后,should_exit只在Op::Shutdown时为true。
if should_exit {
    shutdown_received = true;
    break;
}
// 这个右括号结束while let;channel close也会直接到达后续分支。
}
if !shutdown_received {
    shutdown_session_runtime(&sess).await;
    emit_thread_stop_lifecycle(sess.as_ref()).await;
    if let Some(live_thread) = sess.live_thread()
        && let Err(err) = live_thread.shutdown().await
    {
        warn!(
            "failed to shutdown thread persistence after submission channel closed: {err}"
        );
    }
}
debug!("Agent loop exited");

9. 失败后续清理 ​

关闭的目标是最大限度回收其余资源,所以多个错误被降级为 warning或best effort。只有仍可等待的 owner 协议被顺序 await;某一步失败通常不阻止下一步。

失败点当前处理后续动作
startup prewarm joinabort后忽略结果继续 Realtime
Realtime task join内部忽略 join error继续 Task abort
Task 100ms未完成warning后强制abort执行Task abort hook
marker/terminal flushwarning仍发送TurnAborted并继续关闭
Unified exec process terminate不聚合confirmed-exit结果继续Code Mode
Code Mode shutdownwarning继续MCP
MCP worker joinwarning继续关闭gate/runtime
Guardian child shutdownchild结果忽略继续hooks
SessionEnd transcript flushwarning仍运行hook
LiveThread shutdownError Event仍deliver ShutdownComplete
Event receiver已关闭debug丢弃Eventloop仍完成

这是一种“尽力清理所有 owner”的策略,不是事务回滚。对初学者最重要的判断是:Error/Warning Event能否 送达取决于 event receiver,而资源回收不能依赖客户端继续读取。

10. 关闭顺序测试 ​

10.1 Channel close ​

测试先安装同时实现 TurnAbort 和 ThreadStop lifecycle 的 recorder,再启动永不结束但监听 cancellation 的 RegularTask,最后直接 drop submission sender。关键断言是记录顺序严格等于 ["turn_abort", "thread_stop"]。

源码位置:codex-rs/core/src/session/tests.rs

rust
// :: submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle(断言节选)
session
    .spawn_task(
        Arc::new(turn_context),
        Vec::new(),
        NeverEndingTask {
            kind: TaskKind::Regular,
            listen_to_cancellation_token: true,
        },
    )
    .await;

let (tx_sub, rx_sub) = async_channel::bounded(1);
drop(tx_sub);
submission_loop(
    Arc::clone(&session),
    session.get_config().await,
    rx_sub,
)
.await;

// 检查隐式退出包含完整Turn abort,并且extension顺序没有颠倒。
assert_eq!(
    vec!["turn_abort", "thread_stop"],
    *calls
        .lock()
        .unwrap_or_else(std::sync::PoisonError::into_inner),
);

它没有覆盖 ShutdownComplete,因为隐式路径本来不发送该事件;也没有覆盖 pending trigger mail 先入队、 随后 channel close 触发 teardown 的顺序。

10.2 Interrupted ​

测试用 InMemoryThreadStore记录调用次数:Task观察 cancellation后 wrapper先 flush;abort path在 marker后 flush;TurnAborted追加后再次 flush。断言恰好为3,直接锁定本文的两个额外durability barrier。

源码位置:codex-rs/core/src/session/tests.rs

rust
// :: turn_aborted_flushes_terminal_event_after_delivery(断言节选)
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;
assert!(matches!(
    event.msg,
    EventMsg::TurnAborted(TurnAbortedEvent {
        reason: TurnAbortReason::Interrupted,
        ..
    })
));
abort_task.await.expect("abort task should finish");

// task-body、marker、terminal event三个barrier缺一都会让计数变化。
let calls = wait_for_flush_count(&store, /*expected_flushes*/ 3).await;
assert_eq!(3, calls.flush_thread);

它覆盖调用顺序和数量,不覆盖真实磁盘断电语义;后者属于 ThreadStore/rollout 实现测试。

10.3 Store Shutdown ​

测试给 Session换入 InMemoryThreadStore,调用显式 handler,然后断言只有 create和shutdown各一次,没有 append/flush 等额外 Store call。这说明完成事件走 direct deliver,而不是重新写已关闭 Store。

源码位置:codex-rs/core/src/session/tests.rs

rust
// :: shutdown_complete_does_not_append_to_thread_store_after_shutdown(断言节选)
assert!(handlers::shutdown(&session, "sub-1".to_string()).await);

// 若ShutdownComplete误走send_event_raw,这里会出现shutdown后的append调用。
assert_eq!(
    InMemoryThreadStoreCalls {
        create_thread: 1,
        shutdown_thread: 1,
        ..Default::default()
    },
    store.calls().await,
);

10.4 SessionEnd hook ​

集成测试让模型产生真实回答,配置一个读取 transcript 的 SessionEnd hook,再执行 shutdown_and_wait()。断言 hook只运行一次、transcript_exists=true,并同时包含用户输入和模型回答;另一个 测试覆盖 Review/ThreadSpawn subagent 不运行 SessionEnd。

这组测试覆盖 hook 的 transcript barrier 和 root-only 条件,但不覆盖 hook 失败会停止关闭的路径;源码明确选择继续。

11. 关闭顺序验证 ​

完成本文后,可以用下面四个问题检验自己是否真正理解关闭链:

  1. 从 SessionIo::shutdown_and_wait() 开始,按顺序说出 submission loop、runtime teardown、thread-stop、 LiveThread和ShutdownComplete,并指出哪个步骤以后不能再追加rollout。
  2. 解释为什么 Interrupted Task测试期望三次 flush,而正常 TurnComplete只期望两次。
  3. 给定“客户端没收到ShutdownComplete,但后台submission loop已经退出”,分别检查 event receiver和隐式 channel-close路径,说明这不必然代表资源泄漏。
  4. 阅读下面的只读搜索结果,确认每个资源是否 cancel、join、abort还是只发termination signal:
bash
rg -n "shutdown_session_runtime|handle_task_abort|stop_mcp_prewarm_worker|terminate_all_processes|run_session_end_hooks" \
  codex-rs/core/src

还应保留一个未解决问题:当前 shutdown 复用 abort_all_tasks(Interrupted),而该函数可能为已经排队的 trigger mail 启动新 Turn。除非新增 closing gate、teardown 后的第二次 abort,或专门测试证明该顺序无害, 否则不能把这个关闭期调度缺口从设计说明中省略。

继续学习时,可以阅读 CodexThread背压,理解 event receiver 关闭为何不会反向取消 Session;再回到 Session依赖装配,将本文 的回收顺序与资源创建顺序逐项反向对应。