Skip to content

Session启动预热

区分 Session 构造期加载、首轮模型预热、MCP 后台刷新与 Skill 文件监听的生命周期和正确性边界。

基于rust-v0.150.0
CodexRustSessionPrewarm

Session启动预热 ​

Codex 启动一个 Session 时,Plugin、Skill、AGENTS.md、MCP 和模型连接都会出现“提前准备”的行为, 但它们并不属于同一种预热。Plugin/Skill 与 AGENTS.md 是 Session::spawn() 返回前必须完成的构造工作; 模型 WebSocket 是首轮可以复用、失败后允许降级的单次优化;MCP prewarm 则是贯穿 Session 生命周期的 合并刷新 worker;Skill 文件监听甚至不由 Core Session 持有,而是 App Server 的宿主能力。

如果只看到多个 tokio::spawn()、tokio::join!() 就把它们统称为“并行预热”,会误判三个关键问题: Session 何时可以对外发布、首轮为什么可能短暂等待,以及后台任务失败后谁负责兜底。

本文从 Session依赖装配 的发布屏障继续,默认读者已经知道 SessionServices 与 LiveThread 的 owner。它只区分构造期加载、模型 startup prewarm、MCP worker 和 宿主 watcher,不重复 Session核心数据结构 的字段分组。

1. 四类准备工作 ​

下面按“是否阻塞 Session 构造”和“是否位于正确性路径”拆分启动准备。图中的虚线表达失效通知, 不是构造调用关系。

四类机制的承诺不同:

机制所有者启动时机失败结果正确性兜底
AGENTS.md refreshCore SessionSession::spawn() 内当前加载结果为空也可继续Turn 使用 manager 中已加载结果
Plugin/Skill warmupCore Session 使用共享 serviceSession::spawn() 内记录 Skill error,继续构造Turn 再获取配置对应的 Skill snapshot
Model startup prewarmCore SessionSessionConfigured 后超时或失败后首轮新建 client sessionrun_turn() 正常连接与 fallback
MCP prewarm workerCore Session初始 MCP runtime 安装后保持 dirty,后续刷新重试每个 model step/tool call 主动 refresh_mcp_if_dirty()
Skill file watcherApp ServerApp Server 生命周期内watcher 退化为 noop后续显式加载仍可读取文件系统

所以“构造期并行加载”缩短的是 Session::spawn() 的总等待时间;“发布后 prewarm”隐藏的是客户端看到 Session 到提交首轮之间的空档;文件 watcher 则解决长生命周期进程中的缓存失效。它们不能互相替代。

2. 构造期加载 ​

Session::spawn() 的内部装配先解析环境,再并行刷新 AGENTS.md、加载 Plugin/Skill 快照并查询 thread name。 这里使用 tokio::join!,意味着三个 future 都结束后才能继续,不是脱离 Session 生命周期的后台任务。

源码位置:codex-rs/core/src/session/session.rs :: spawn_internal(构造期warmup节选)

rust
let resolved_environments = turn_environments.snapshot().await;
let agents_md_manager = Arc::new(AgentsMdManager::new(user_instructions));

let plugin_skill_warmup = warm_plugins_and_skills_for_session_init(
    Arc::clone(&config),
    Arc::clone(&plugins_manager),
    Arc::clone(&skills_service),
    // 必须使用已解析环境的filesystem,不能默认假设所有文件都在宿主本地。
    &resolved_environments,
)
.instrument(info_span!(
    "session_init.plugin_skill_warmup",
    otel.name = "session_init.plugin_skill_warmup",
));
let thread_name_lookup =
    thread_title_from_thread_store(live_thread_init.as_ref(), &thread_store, thread_id)
        .instrument(info_span!(
            "session_init.thread_name_lookup",
            otel.name = "session_init.thread_name_lookup",
        ));

let ((), plugin_skill_errors, thread_name) = tokio::join!(
    // join并发轮询三项,但Session构造必须等待三项全部结束。
    agents_md_manager.refresh(config.as_ref(), &resolved_environments),
    plugin_skill_warmup,
    thread_name_lookup,
);
for err in &plugin_skill_errors {
    // Skill加载错误降级为日志,不把整个Thread变成不可用。
    error!("failed to load skill {}: {}", err.path.display(), err.message);
}
session_configuration.thread_name = thread_name.clone();

Plugin 与 Skill 并不是两条独立 future。Core 先解析当前配置启用的 Plugin,再把 Plugin 暴露的 Skill root 和预制 snapshot 交给 HostSkillsService;二者具有明确的数据依赖,只在这条 future 与另外两条 初始化工作之间并行。

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

rust
async fn warm_plugins_and_skills_for_session_init(
    config: Arc<Config>,
    plugins_manager: Arc<PluginsManager>,
    skills_service: Arc<HostSkillsService>,
    turn_environments: &TurnEnvironmentSnapshot,
) -> Vec<SkillError> {
    let fs = turn_environments.primary_filesystem();
    let plugins_input = config.plugins_config_input();
    // 先求出有效Plugin集合,Skill root不能脱离Plugin配置单独推导。
    let plugin_outcome = plugins_manager.plugins_for_config(&plugins_input).await;
    let effective_skill_roots = plugin_outcome.effective_plugin_skill_roots();
    let plugin_skill_snapshots =
        plugins_manager.plugin_skill_snapshots_for_config(&plugins_input);
    let skills_input = skills_load_input_from_config(config.as_ref(), effective_skill_roots)
        // Plugin自带快照与文件系统root被合并成一次Skill加载输入。
        .with_plugin_skill_snapshots(plugin_skill_snapshots);

    skills_service
        .snapshot_for_config(&skills_input, fs)
        .await
        .outcome()
        .errors
        .clone()
}

AGENTS.md 也不是无条件重复扫描。AgentsMdManager 以 environment selections 作为 cache key;选择没有 变化时直接返回,变化时在锁外读取项目指令,最后再替换缓存,避免持锁执行文件系统 I/O。

源码位置:codex-rs/core/src/agents_md_manager.rs :: AgentsMdManager::refresh

rust
pub(crate) async fn refresh(
    &self,
    config: &Config,
    environments: &TurnEnvironmentSnapshot,
) {
    let selections = environments.to_selections();
    if self.cache.lock().await.selections.as_ref() == Some(&selections) {
        // 环境选择未变化时复用Session级结果,不重新发现AGENTS.md。
        return;
    }

    // 文件读取发生在cache锁之外,其他只读调用不会被慢I/O长期阻塞。
    let loaded = load_project_instructions(
        config,
        self.user_instructions.clone(),
        environments,
    )
    .await
    .map(Arc::new);
    let mut cache = self.cache.lock().await;
    cache.selections = Some(selections);
    cache.loaded = loaded;
}

3. 配置完成事件 ​

Session Arc 创建后,Core 先发送 SessionConfigured 和启动期 warning,随后才安装正式 MCP runtime、 启动 MCP worker 并调度模型预热。这个顺序使客户端先获得 canonical Session 身份和配置,后续 MCP 连接事件、预热失败日志或首轮事件才不会出现在一个“尚未配置”的 Thread 上。

对应的发布顺序直接写在 spawn_internal() 的 Session 装配尾部:

源码位置:codex-rs/core/src/session/session.rs :: spawn_internal(发布与预热顺序节选)

rust
for event in events {
    // events的第一项固定是SessionConfigured,随后才是构造期warning。
    sess.send_event_raw(event).await;
}
turn_environments.start_connection_event_forwarding(tx_event.clone());

// ... 若启动期间认证发生变化,此处重新计算MCP projection。
sess.install_initial_mcp_runtime(
    &session_configuration,
    latest_auth,
    mcp_projection,
    &resolved_environments,
    mcp_runtime_cwd,
)
.await?;
// worker只能在初始binding安装后启动,避免后台刷新与首次发布争用顺序。
sess.start_mcp_prewarm_worker(mcp_prewarm_rx, mcp_auth_changes);
// 模型预热在SessionConfigured之后调度,不延迟客户端接收配置事件。
sess.schedule_startup_prewarm(session_configuration.base_instructions.clone())
    .await;

// initial history可能发事件,因此同样被放在SessionConfigured之后。
Box::pin(sess.record_initial_history(initial_history)).await;

这里的 schedule_startup_prewarm() 虽然是 async 方法,但它只完成分支判断、spawn 和 handle 入库; 真正的网络准备发生在 spawned task 中。相反,install_initial_mcp_runtime(...).await? 仍是 Session 构造的 可失败屏障。

当前测试还明确约束了首轮行为:RegularTask 发出 TurnStarted 不需要等待 startup prewarm 完成; 如果首轮等待期间收到 interrupt,则会产生 TurnAborted,而不是把 prewarm 当成不可取消的构造屏障。 startup prewarm 因而是可消费的一次性优化句柄,不是 Session 对外发布前必须成功的第二个初始化阶段。

4. 模型预热 ​

4.1 WebSocket 关闭 ​

模型预热先检查 Responses WebSocket。功能关闭时不创建 SessionStartupPrewarmHandle,只在后台调用 prewarm_auth(),提前完成当前 client setup 与 Agent Identity bearer fallback 的准备;失败只记录 warning。

源码位置:codex-rs/core/src/session_startup_prewarm.rs :: Session::schedule_startup_prewarm

rust
pub(crate) async fn schedule_startup_prewarm(
    self: &Arc<Self>,
    base_instructions: String,
) {
    if !self.services.model_client.responses_websocket_enabled() {
        let model_client = self.services.model_client.clone();
        tokio::spawn(async move {
            // 非WebSocket路径只解析认证/client setup,不发送模型推理请求。
            if let Err(err) = model_client.prewarm_auth().await {
                warn!("startup auth prewarm failed: {err:#}");
            }
        });
        return;
    }

    let session_telemetry = self.services.session_telemetry.clone();
    let websocket_connect_timeout = self.provider().await.websocket_connect_timeout();
    let started_at = Instant::now();
    let startup_prewarm_session = Arc::clone(self);
    let startup_prewarm = tokio::spawn(
        async move {
            // task拥有Session强引用,外层handle负责消费、超时或关闭时abort。
            let result = schedule_startup_prewarm_inner(
                startup_prewarm_session,
                base_instructions,
            )
            .await;
            let status = if result.is_ok() { "ready" } else { "failed" };
            session_telemetry.record_startup_phase(
                "startup_prewarm_total",
                started_at.elapsed(),
                Some(status),
            );
            session_telemetry.record_duration(
                STARTUP_PREWARM_DURATION_METRIC,
                started_at.elapsed(),
                &[("status", status)],
            );
            result
        }
        .instrument(trace_span!(
            "startup_prewarm",
            otel.name = "startup_prewarm",
            thread.id = %self.thread_id(),
        )),
    );
    self.set_session_startup_prewarm(SessionStartupPrewarmHandle::new(
        startup_prewarm,
        started_at,
        websocket_connect_timeout,
    ))
    .await;
}

prewarm_auth() 本身很薄:它复用普通模型请求的 current_client_setup(),因此不是维护第二套认证逻辑。

源码位置:codex-rs/core/src/client.rs :: ModelClient::prewarm_auth

rust
pub(crate) async fn prewarm_auth(&self) -> Result<()> {
    // 丢弃setup值但保留解析、注册与缓存产生的副作用。
    self.current_client_setup().await.map(|_| ())
}

4.2 WebSocket预热 ​

WebSocket 分支创建专用 startup TurnContext,捕获 StepContext 和 tool router,再用空用户输入构造 Prompt。这不是一次普通 run_turn():它跳过 Git enrichment,使用 Prewarm request kind,并用独立 CancellationToken 捕获当时的工具快照。

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

rust
let startup_turn_context = session
    .new_startup_prewarm_turn_with_sub_id(INITIAL_SUBMIT_ID.to_owned())
    .await;
// Guardian预热与阶段耗时记录不改变下面的tool/prompt数据流,此处省略。
let startup_cancellation_token = CancellationToken::new();

// prewarm早于run_turn,必须主动捕获一次真实StepContext来获得tool router。
let step_context = session
    .capture_step_context(
        Arc::clone(&startup_turn_context),
        &startup_cancellation_token,
    )
    .await?;
let startup_router = Arc::clone(&step_context.tool_router);
let startup_prompt = build_prompt(
    // 没有用户输入;这里只建立与首轮兼容的instructions和tools请求外壳。
    Vec::new(),
    startup_router.as_ref(),
    startup_turn_context.as_ref(),
    BaseInstructions {
        text: base_instructions,
    },
);
let window_id = session.current_window_id().await;
let responses_metadata = startup_turn_context
    .turn_metadata_state
    .to_responses_metadata(
        session.installation_id.clone(),
        window_id,
        // 服务端和遥测可据此区分预热与真实turn。
        CodexResponsesRequestKind::Prewarm,
    );
let mut client_session = session.services.model_client.new_session();
client_session
    .prewarm_websocket(
        &startup_prompt,
        &startup_turn_context.model_info,
        &startup_turn_context.session_telemetry,
        startup_turn_context.reasoning_effort.clone(),
        startup_turn_context.reasoning_summary,
        startup_turn_context.config.service_tier.clone(),
        &responses_metadata,
    )
    .await?;
Ok(client_session)

对象所有权决定了“复用”发生在哪一层:ModelClient 属于 Session,ModelClientSession 属于一个 Turn; 预热 task 先持有后者,首个 RegularTask 再一次性取走它。

真正发到 WebSocket 的 warmup 是 V2 response.create,generate=false。它等待 Completed 后才把 ModelClientSession 标记为 ready,使真实请求复用同一条连接和 warmup response id;它不是一轮隐藏的 模型回答,也不会作为一次 inference attempt 写入 rollout trace。

源码位置:codex-rs/core/src/client.rs :: ModelClientSession::prewarm_websocket

rust
pub async fn prewarm_websocket(
    &mut self,
    prompt: &Prompt,
    model_info: &ModelInfo,
    session_telemetry: &SessionTelemetry,
    effort: Option<ReasoningEffortConfig>,
    summary: ReasoningSummaryConfig,
    service_tier: Option<String>,
    responses_metadata: &CodexResponsesMetadata,
) -> Result<()> {
    if !self.client.responses_websocket_enabled() {
        return Ok(());
    }
    if self.websocket_session.last_request.is_some() {
        // 已有请求状态时不覆盖现有WebSocket会话链。
        return Ok(());
    }

    let disabled_trace = InferenceTraceContext::disabled();
    match self
        .stream_responses_websocket(
            prompt,
            model_info,
            session_telemetry,
            effort,
            summary,
            service_tier,
            responses_metadata,
            /*warmup*/ true,
            current_span_w3c_trace_context(),
            &disabled_trace,
        )
        .await
    {
        Ok(WebsocketStreamOutcome::Stream(mut stream)) => {
            while let Some(event) = stream.next().await {
                match event {
                    // 只有服务端确认Completed后,这个session才可交给首轮复用。
                    Ok(ResponseEvent::Completed { .. }) => break,
                    Err(err) => return Err(err),
                    _ => {}
                }
            }
            Ok(())
        }
        Ok(WebsocketStreamOutcome::FallbackToHttp) => {
            // 例如握手返回426时切换HTTP,预热失败不阻断真实turn。
            self.try_switch_fallback_transport(session_telemetry, model_info);
            Ok(())
        }
        Err(err) => Err(err),
    }
}

5. 首轮Regular Turn ​

5.1 TurnStarted ​

首个普通 Turn 到达时,RegularTask 先发送 TurnStarted,再等待预热 handle。这样 UI 能立刻进入 running 状态,并拿到 trace id;网络握手慢不会把 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,
    });
    // 先发布TurnStarted;消费预热即使等待,也不隐藏首轮已启动的事实。
    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
}
.instrument(trace_span!("regular_task.prepare_run_turn"))
.await;

let prewarmed_client_session = match prewarmed_client_session {
    SessionStartupPrewarmResolution::Cancelled => {
        // 中断期间仍记录输入与hook结果,随后按Task取消语义收尾。
        run_hooks_and_record_inputs(&sess, &ctx, &input).await;
        return Ok(None);
    }
    // failed、timed_out、join_failed与not_scheduled都回退为普通新session。
    SessionStartupPrewarmResolution::Unavailable { .. } => None,
    SessionStartupPrewarmResolution::Ready(prewarmed) => Some(*prewarmed),
};

5.2 Handle ​

SessionState::take_session_startup_prewarm() 使用 Option::take()。因此只有第一个 Regular Turn 能获得 handle;同一 Turn 因 steer/inject 继续循环时,prewarmed_client_session.take() 也只在首次 run_turn() 传入,后续循环改用正常新建路径。

源码位置:codex-rs/core/src/state/session.rs :: SessionState::take_session_startup_prewarm

rust
pub(crate) fn take_session_startup_prewarm(
    &mut self,
) -> Option<SessionStartupPrewarmHandle> {
    // Option::take把状态原位改成None,所有权不会被第二个Turn重复取得。
    self.startup_prewarm.take()
}

let mut prewarmed_client_session = prewarmed_client_session;
loop {
    let last_agent_message = run_turn(
        Arc::clone(&sess),
        Arc::clone(&ctx),
        next_input,
        // 第一次调用后这里变成None,后续pending 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();
}

Handle 的解析状态比“成功/失败”更细。超时预算从 Session 调度预热时开始,而不是首轮开始等待时重新 计时;因此 remaining = timeout - age_at_first_turn,首轮不会额外再等一个完整连接超时。

状态机中的终态在实现里被折叠成三个调用方分支:可复用、不可用但可降级,以及本轮已取消。

源码位置:codex-rs/core/src/session_startup_prewarm.rs :: SessionStartupPrewarmHandle::resolve

rust
let age_at_first_turn = started_at.elapsed();
// timeout从预热调度时计时;首轮只等待尚未消耗的剩余预算。
let remaining = timeout.saturating_sub(age_at_first_turn);

let resolution = if task.is_finished() {
    Self::resolution_from_join_result(task.await, started_at)
} else {
    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)) => {
            // 超时后主动abort,防止失去消费者的连接task继续占用Session资源。
            task.abort();
            SessionStartupPrewarmResolution::Unavailable {
                status: "timed_out",
                prewarm_duration: Some(started_at.elapsed()),
            }
        }
        None => {
            // Turn取消与预热不可用不同:调用方直接结束这次RegularTask。
            task.abort();
            return SessionStartupPrewarmResolution::Cancelled;
        }
    }
};

6. MCP Prewarm ​

模型预热只服务首个普通 Turn,MCP worker 则持续监听两类信号:显式 refresh request 和认证版本变化。 它保存 Weak<Session>,使 worker 不会仅凭自身强引用延长 Session 生命周期。

6.1 单容量Channel ​

request_mcp_runtime_refresh() 先设置原子 dirty 标记,再向容量为一的 channel 执行 try_send(())。 channel 满时丢弃新的 unit message 是安全的,因为待处理事实保存在 McpRefresh.pending,消息只负责唤醒。

源码位置:codex-rs/core/src/session/mcp_prewarm.rs :: Session MCP prewarm scheduling

rust
pub(crate) fn request_mcp_runtime_refresh(&self) {
    // 正确状态存入dirty flag,不能只依赖可能被合并的channel消息。
    self.mark_mcp_runtime_dirty();
    self.schedule_mcp_prewarm();
}

pub(super) fn schedule_mcp_prewarm(&self) {
    // bounded(1)满时忽略发送失败,把一串变化合并成一次“读取最新状态”。
    let _ = self.mcp_prewarm_tx.try_send(());
}

pub(super) fn start_mcp_prewarm_worker(
    self: &Arc<Self>,
    requests: async_channel::Receiver<()>,
    mut auth_changes: tokio::sync::watch::Receiver<u64>,
) {
    let session = Arc::downgrade(self);
    let shutdown = self.mcp_prewarm_shutdown.clone();
    let worker = self.services.runtime_handle.spawn(async move {
        loop {
            let auth_changed = tokio::select! {
                biased;
                _ = shutdown.cancelled() => break,
                request = requests.recv() => {
                    if request.is_err() { break; }
                    false
                },
                auth_change = auth_changes.changed() => {
                    if auth_change.is_err() { break; }
                    true
                },
            };
            // Session已释放时Weak升级失败,worker自然退出。
            let Some(session) = session.upgrade() else { break; };
            if auth_changed {
                session.mark_mcp_runtime_dirty();
            }
            tokio::select! {
                biased;
                _ = shutdown.cancelled() => break,
                _ = session.refresh_mcp_if_dirty() => {},
            }
        }
    });
    *self.mcp_prewarm_task.lock()
        .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(worker);
}

6.2 后台刷新边界 ​

Core 文件顶部直接把 MCP prewarm 定义为 best effort:worker 只尽早准备最新 Thread 状态,每个确切 model step 仍走 refresh_mcp_if_dirty()。刷新由单 permit semaphore 串行化;claim 后若 future 被取消, RAII guard 在 Drop 时重新 invalidate,避免丢失已认领但尚未发布的变化。

源码位置:codex-rs/core/src/session/mcp_refresh.rs :: McpRefresh与InvalidationGuard

rust
pub(super) struct McpRefresh {
    pending: AtomicBool,
    // 单permit保证同一Session只有一条runtime发布链。
    gate: Semaphore,
}

impl McpRefresh {
    pub(super) fn invalidate(&self) {
        self.pending.store(true, Ordering::Release);
    }

    pub(super) fn claim(&self) -> bool {
        // 原子取走当前dirty;刷新期间的新invalidate会再次把它设为true。
        self.pending.swap(false, Ordering::AcqRel)
    }
}

pub(super) struct McpRefreshInvalidationGuard<'a> {
    pub(super) refresh: &'a McpRefresh,
    pub(super) published: bool,
}

impl Drop for McpRefreshInvalidationGuard<'_> {
    fn drop(&mut self) {
        if !self.published {
            // task取消或panic发生在publish前时,恢复dirty供下一条路径重试。
            self.refresh.invalidate();
        }
    }
}

这个 guard 保护的是“失效信号不丢失”;真正的重建与发布仍由持有 gate 的刷新循环完成。

源码位置:codex-rs/core/src/session/mcp.rs :: Session::refresh_mcp_if_dirty(节选)

rust
pub(crate) async fn refresh_mcp_if_dirty(self: &Arc<Self>) {
    let Ok(_refresh) = self.mcp_refresh.acquire().await else {
        error!("MCP runtime refresh semaphore closed");
        return;
    };
    loop {
        let auth = self.services.auth_manager.auth_cached();
        if self
            .services
            .plugins_manager
            .set_auth_mode(auth.as_ref().map(CodexAuth::api_auth_mode))
            || !self
                .services
                .mcp_runtime
                .current_auth_matches(auth.as_ref())
        {
            // cached auth与当前binding不一致时,即使没有外部消息也重新标脏。
            self.mark_mcp_runtime_dirty();
        }

        if !self.mcp_refresh.claim() {
            // 没有dirty时立即返回,正确性检查成本不等于每步重建runtime。
            return;
        }
        let mut refresh_invalidation = McpRefreshInvalidationGuard {
            refresh: &self.mcp_refresh,
            published: false,
        };
        let auth = self.services.auth_manager.auth().await;
        self.services
            .plugins_manager
            .set_auth_mode(auth.as_ref().map(CodexAuth::api_auth_mode));
        let desired = self.latest_mcp_desired_state(auth).await;
        let selected_capability_roots = self
            .resolve_selected_capability_roots_for_step(&desired.environments)
            .await;
        let ready_selected_capability_roots =
            Self::ready_selected_capability_roots(&selected_capability_roots);
        let executor_capability_discovery = self
            .executor_capability_discovery_for_step(
                &desired.config,
                &ready_selected_capability_roots,
                &desired.environments,
                desired.windows_sandbox_level,
            )
            .await;
        let mcp_projection = self
            .services
            .mcp_manager
            .runtime_config_for_step(
                &desired.config,
                &self.services.mcp_thread_init,
                &self.services.thread_extension_data,
                McpThreadIdentity {
                    session_source: &desired.session_source,
                    originator: &desired.originator,
                },
                &ready_selected_capability_roots,
                // executor侧发现结果也属于这次runtime projection的输入。
                executor_capability_discovery.as_deref(),
            )
            .await;
        self.publish_mcp_runtime(
            &desired,
            mcp_projection,
            &ready_selected_capability_roots,
            Some(self.mcp_elicitation_reviewer()),
        )
        .await;
        refresh_invalidation.published = true;
        // 发布期间若又失效,就在持有同一gate时继续收敛到最新状态。
        if !self.mcp_refresh.is_pending() {
            return;
        }
    }
}

7. Skill 文件监听 ​

Core 的 Session 构造只加载 Skill 快照,并没有创建文件 watcher。长生命周期的 App Server 创建一个 SkillsWatcher,为本地 Thread 的 Skill roots 注册递归监听;远程 environment 直接跳过,因为宿主 进程的本地 watcher 无法代表远端文件系统。

监听路径还刻意排除两类 root:Plugin root 有独立的生命周期失效机制,System Skill 在 watcher 启动前 已经安装。这样 watcher 不会把所有来源重复纳入同一套文件事件策略。

源码位置:codex-rs/app-server/src/skills_watcher.rs :: SkillsWatcher::register_thread_config

rust
pub(crate) async fn register_thread_config(
    &self,
    config: &Config,
    thread_manager: &ThreadManager,
    environments: &[TurnEnvironmentSelection],
) -> WatchRegistration {
    let Some(selection) = environments.first() else {
        return WatchRegistration::default();
    };
    let Some(environment) = thread_manager
        .environment_manager()
        .get_environment(&selection.environment_id)
    else {
        // 未知environment无法安全猜测watch path,返回空registration。
        return WatchRegistration::default();
    };
    if environment.is_remote() {
        // App Server本地FileWatcher不监视远端executor文件系统。
        return WatchRegistration::default();
    }

    let plugins_input = config.plugins_config_input();
    let plugin_outcome = thread_manager
        .plugins_manager()
        .plugins_for_config(&plugins_input)
        .await;
    let skills_input = HostSkillsLoadInput::new(
        config.cwd.clone(),
        plugin_outcome.effective_plugin_skill_roots(),
        config.config_layer_stack.clone(),
        config.bundled_skills_enabled(),
    );
    let roots = thread_manager.skills_service()
        .skill_roots_for_config(&skills_input, Some(environment.get_filesystem()))
        .await
        .into_iter()
        // Plugin和System root各有自己的失效/安装生命周期,不在这里重复监听。
        .filter(|root| root.plugin_identity.is_none() && root.scope != SkillScope::System)
        .map(|root| WatchPath {
            path: root.path.into_path_buf(),
            recursive: true,
        })
        .collect();
    self.subscriber.register_paths(roots)
}

事件循环采用节流 receiver。一次文件变化会清除共享 Skill cache,并向客户端发送 SkillsChanged; 它不会直接修改某个正在运行的 TurnContext。已捕获的 Turn 快照保持稳定,后续 list/turn 才读取新快照。

源码位置:codex-rs/app-server/src/skills_watcher.rs :: SkillsWatcher::spawn_event_loop

rust
let mut rx = ThrottledWatchReceiver::new(rx, WATCHER_THROTTLE_INTERVAL);
handle.spawn(async move {
    loop {
        let event = tokio::select! {
            _ = shutdown_token.cancelled() => break,
            event = rx.recv() => event,
        };
        let Some(event) = event else { break; };
        if event.paths.iter().all(|path| {
            path.starts_with(system_skills_root.as_path())
        }) {
            // legacy user root递归覆盖.system;系统Skill事件在此去重。
            continue;
        }
        // watcher只做cache invalidation,不就地改写任何正在运行的TurnContext。
        skills_service.clear_cache();
        outgoing
            .send_server_notification(ServerNotification::SkillsChanged(
                SkillsChangedNotification {},
            ))
            .await;
    }
});

这也解释了不同宿主的行为差异:App Server 客户端能收到 SkillsChanged,但直接嵌入 codex-core 的宿主 不会凭空拥有这项监听能力;它需要自行决定何时清缓存或重新加载。

8. 取消和关闭回收 ​

模型 prewarm handle 包装 AbortOnDropHandle,首轮消费、超时和 Session shutdown 都有明确终点; MCP worker 则用 Session 级 CancellationToken 停止。关闭顺序先处理 startup prewarm,再中止 active tasks, 最后停 MCP worker并关闭 runtime,避免刷新 task 在 runtime shutdown 期间再次发布 binding。

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

rust
async fn shutdown_session_runtime(sess: &Arc<Session>) {
    if let Some(startup_prewarm) = sess.take_session_startup_prewarm().await {
        // 首轮尚未消费时先abort并await,释放其Session强引用和网络准备任务。
        startup_prewarm.abort().await;
    }
    let _ = sess.conversation.shutdown().await;
    sess.abort_all_tasks(TurnAbortReason::Interrupted).await;
    sess.services.unified_exec_manager.terminate_all_processes().await;
    if let Err(err) = sess.services.code_mode_service.shutdown().await {
        warn!("failed to shutdown code mode session: {err}");
    }

    // 先停止可能发起refresh的worker,再关闭refresh gate和MCP runtime。
    sess.stop_mcp_prewarm_worker().await;
    {
        let _refresh = sess.mcp_refresh.acquire().await;
        sess.mcp_refresh.close();
        sess.services.mcp_runtime.shutdown().await;
    }
    sess.guardian_review_session.shutdown().await;
    crate::hook_runtime::run_session_end_hooks(sess).await;
}

stop_mcp_prewarm_worker() 先 cancel,再从同步 Mutex 中取出 JoinHandle,释放锁后才 await。若持锁 await, 其他清理路径可能无法访问 task slot;当前实现避免了这种锁跨 await 的问题。

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

rust
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)
        // 在同步锁内只转移所有权,不在锁内等待task结束。
        .take();
    if let Some(worker) = worker
        && let Err(error) = worker.await
    {
        warn!(%error, "MCP prewarm worker stopped unexpectedly");
    }
}

9. 测试预热是否越界 ​

启动预热最容易出现的回归不是“完全不能运行”,而是事件顺序、复用次数或取消语义悄悄变化。当前测试 从三层锁定边界:Core 单元测试验证 TurnStarted 不等待;WebSocket 集成测试验证一次握手承载 warmup 和真实 turn;MCP 集成测试验证 shutdown 不等待卡住的 MCP startup。

源码位置:codex-rs/core/tests/suite/agent_websocket.rs

rust
// :: websocket_first_turn_uses_startup_prewarm_and_create(断言节选)
test.submit_turn_with_policy("hello", test.config.legacy_sandbox_policy())
    .await?;

// warmup与真实turn共用一次WebSocket handshake。
assert_eq!(server.handshakes().len(), 1);
let connection = server.single_connection();
assert_eq!(connection.len(), 2);
let warmup = connection.first().expect("missing warmup request").body_json();
let turn = connection.get(1).expect("missing turn request").body_json();
assert_eq!(warmup["type"].as_str(), Some("response.create"));
// generate=false证明这是连接准备,不生成用户可见回答。
assert_eq!(warmup["generate"].as_bool(), Some(false));
assert_eq!(turn["type"].as_str(), Some("response.create"));

连接复用测试验证传输层结果,下面的单元测试则专门锁定对外事件顺序。

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

rust
// :: regular_turn_emits_turn_started_with_trace_id_without_waiting_for_startup_prewarm(节选)
let (_tx, startup_prewarm_rx) = tokio::sync::oneshot::channel::<()>();
let handle = tokio::spawn(async move {
    // sender被保留且不发送,模拟永远尚未ready的startup prewarm。
    let _ = startup_prewarm_rx.await;
    Ok(test_model_client_session())
});
sess.set_session_startup_prewarm(
    SessionStartupPrewarmHandle::new(
        handle,
        std::time::Instant::now(),
        crate::client::WEBSOCKET_CONNECT_TIMEOUT,
    ),
).await;

sess.spawn_task(Arc::clone(&tc), Vec::new(), RegularTask::new()).await;
let first = tokio::time::timeout(Duration::from_millis(200), rx.recv())
    .await
    // 若RegularTask先等待prewarm,这个200ms断言就会失败。
    .expect("expected turn started event without waiting for startup prewarm")
    .expect("channel open");
assert!(matches!(first.msg, EventMsg::TurnStarted(_)));

诊断启动变慢时,应先按生命周期定位,而不是笼统搜索“prewarm”:

  • SessionConfigured 迟迟不出现:检查构造期 environment、AGENTS、Plugin/Skill、thread store 和初始 MCP 安装;模型 startup task 此时还不是阻塞点。
  • TurnStarted 已出现但首个请求延迟:检查 startup_prewarm.resolve、WebSocket connect timeout 和认证恢复。
  • MCP 工具在配置变化后短暂缺失:检查 dirty 是否设置、worker 是否收到唤醒,以及 step 侧 refresh_mcp_if_dirty() 是否执行;不能只看后台 worker。
  • Skill 文件变化没有通知:先确认宿主是否为 App Server、environment 是否本地、监听 root 是否被过滤; 这不是 Core Session model prewarm 的职责。
  • shutdown 卡住:检查 startup handle 是否 abort、MCP worker 是否先于 refresh gate/runtime 关闭。

这些边界共同形成一个原则:预热可以隐藏延迟,但不能成为唯一的正确性来源。凡是允许后台失败的机制, 都必须在首个消费者或每次关键消费点保留同步兜底;凡是构造不变量,则必须留在 Session 发布屏障之前。

首轮真正注入哪些模型上下文见 Session启动上下文;预热 task 与 MCP worker 如何在关闭时取消见 Session关闭流程。

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

bash
rg -n "prewarm|SessionConfigured|TurnStarted|CancellationToken" codex-rs/core/src