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。
| 动作 | 作用范围 | TurnAbortReason | Session 是否继续 | ShutdownComplete |
|---|---|---|---|---|
| Interrupt | 当前 RunningTask | Interrupted | 是 | 无 |
| Replaced | 当前 RunningTask | Replaced | 是,新 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 prewarm | SessionState handle | abort task | await join结束 | 忽略 join result |
| Realtime | RealtimeConversationManager | take state、cancel token | await input/fanout task | shutdown结果被外层忽略 |
| active Task | ActiveTurn::RunningTask | cancel→100ms grace→abort→Task hook | 是 | TurnAborted、warning、flush日志 |
| Unified exec | UnifiedExecManager | drain registry、unregister approval、terminate() | 不确认OS进程退出 | 无聚合错误返回 |
| Code Mode | CodeModeService | 标记closing、join init、session shutdown | 是 | warning,继续关闭 |
| MCP prewarm | Session字段 | cancel worker token | await JoinHandle | join异常warning |
| MCP refresh/runtime | McpRefresh/McpRuntime | acquire gate→close→runtime shutdown | 是 | gate关闭日志在后续调用可见 |
| Guardian | GuardianReviewSessionManager | cancel、take sessions、逐个shutdown | 是 | child shutdown_and_wait结果忽略 |
| SessionEnd hooks | Hooks service | root-only执行 | 是 | transcript flush失败warning |
| LiveThread | SessionServices | metadata 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
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
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
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(取消核心)
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)
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,但下面的顺序已经足以形成 待关闭路径:
- Session 仍有 RunningTask 时,先处理
Op::InterAgentCommunication,把trigger_turn=true的消息留在 mailbox;pending-work scheduler 因 active slot 非空而不启动新 Turn。 - 紧接着处理
Op::Shutdown;abort_all_tasks(Interrupted)移除旧 RunningTask,随后再次调用maybe_start_turn_for_pending_work()。 - 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
// :: 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
// :: 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
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
// :: 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
// :: 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
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
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
// 上方完整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 join | abort后忽略结果 | 继续 Realtime |
| Realtime task join | 内部忽略 join error | 继续 Task abort |
| Task 100ms未完成 | warning后强制abort | 执行Task abort hook |
| marker/terminal flush | warning | 仍发送TurnAborted并继续关闭 |
| Unified exec process terminate | 不聚合confirmed-exit结果 | 继续Code Mode |
| Code Mode shutdown | warning | 继续MCP |
| MCP worker join | warning | 继续关闭gate/runtime |
| Guardian child shutdown | child结果忽略 | 继续hooks |
| SessionEnd transcript flush | warning | 仍运行hook |
| LiveThread shutdown | Error Event | 仍deliver ShutdownComplete |
| Event receiver已关闭 | debug丢弃Event | loop仍完成 |
这是一种“尽力清理所有 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
// :: 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
// :: 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
// :: 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. 关闭顺序验证
完成本文后,可以用下面四个问题检验自己是否真正理解关闭链:
- 从
SessionIo::shutdown_and_wait()开始,按顺序说出 submission loop、runtime teardown、thread-stop、 LiveThread和ShutdownComplete,并指出哪个步骤以后不能再追加rollout。 - 解释为什么 Interrupted Task测试期望三次 flush,而正常 TurnComplete只期望两次。
- 给定“客户端没收到ShutdownComplete,但后台submission loop已经退出”,分别检查 event receiver和隐式 channel-close路径,说明这不必然代表资源泄漏。
- 阅读下面的只读搜索结果,确认每个资源是否 cancel、join、abort还是只发termination signal:
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依赖装配,将本文 的回收顺序与资源创建顺序逐项反向对应。
