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 refresh | Core Session | Session::spawn() 内 | 当前加载结果为空也可继续 | Turn 使用 manager 中已加载结果 |
| Plugin/Skill warmup | Core Session 使用共享 service | Session::spawn() 内 | 记录 Skill error,继续构造 | Turn 再获取配置对应的 Skill snapshot |
| Model startup prewarm | Core Session | SessionConfigured 后 | 超时或失败后首轮新建 client session | run_turn() 正常连接与 fallback |
| MCP prewarm worker | Core Session | 初始 MCP runtime 安装后 | 保持 dirty,后续刷新重试 | 每个 model step/tool call 主动 refresh_mcp_if_dirty() |
| Skill file watcher | App Server | App 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节选)
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
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
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(发布与预热顺序节选)
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
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
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
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
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(预热消费节选)
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
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
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
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
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(节选)
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
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
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
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
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
// :: 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
// :: 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 主线落回源码:
rg -n "prewarm|SessionConfigured|TurnStarted|CancellationToken" codex-rs/core/src