后台任务生命周期
客户端已经收到响应,服务器为什么还在读取配置、连接 MCP 服务或发送费用遥测?关闭了一个对象, 为什么它发起的 HTTP 请求仍可能结束?这些现象都与“后台任务”有关,但不能用一套定时循环解释。 有的工作在 App Server 启动时创建,有的跟随一个 Thread,有的只在一次导入请求之后继续运行。
本文面向理解 Rust Arc、Drop、async/await 和 channel 的读者。可以先读 AppServer启动与依赖装配,确认 MessageProcessor 与 ThreadManager 的关系;断连清理与资源回收 解释连接如何退出。 这里继续追踪那些跨越一次 RPC 调用的工作:谁唤醒它,谁拥有它,谁读取结果,以及取消究竟停在哪一个 await。 遥测事件的字段和交付另见 请求追踪与 Analytics。
先固定对象层级:Thread 是对外可寻址的会话;Core 的 Session 持有这个会话的运行状态;Turn 是一次 输入处理轮次;Tokio task 是调度器上的异步执行单元。一个 Thread 可以先后执行多个 Turn,一个 Session 也可以拥有不属于当前 Turn 的预热 task。下文的“任务”指异步执行单元时,不等同于 Core 的用户任务。
读完后,应能从一个“配置已返回、效果未出现”或“关闭后仍有工作”的现象,找到负责的 owner、唤醒条件、 状态字段和完成消费者。本文展开生命周期与调度关系,模型目录内容、MCP 协议握手和导入文件格式只追到必要边界。
1. 任务所有权
1.1 装配入口
MessageProcessor::new 先取得 ThreadManager 的模型管理器,再创建模型刷新、可选费用收集和技能监听。 三个成员的 Rust 类型已经给出了不同的退出线索。
源码文件:codex-rs/app-server/src/message_processor.rs
相关函数/类型:MessageProcessor / MessageProcessor::new(L135–L139、L366–L377,局部节选)
// 作者注:三个字段使用不同所有权类型;费用 worker 可因配置条件完全不创建。
// ...
pub(crate) struct MessageProcessor {
outgoing: Arc<OutgoingMessageSender>,
models_refresh_worker: ModelsRefreshWorker,
turn_cost_worker: Option<TurnCostWorker>,
skills_watcher: Arc<SkillsWatcher>,
// ...
let models_refresh_worker =
crate::models_refresh_worker::spawn(&models_manager, config.http_client_factory());
let turn_cost_worker =
TurnCostWorker::spawn(Arc::clone(&config), Arc::clone(&auth_manager));
thread_manager
.plugins_manager()
.set_analytics_events_client(analytics_events_client.clone());
let skills_watcher = SkillsWatcher::new(
thread_manager.skills_service(),
&config.codex_home,
outgoing.clone(),
);
// ...ModelsRefreshWorker 是直接成员;费用 worker 用 Option 表示可能根本没有启动;SkillsWatcher 被多个 processor 共享,释放一份 Arc 不等于最后一个 owner 已离开。即使字段都叫 worker,也不能推导出同样的关闭方式。
下面只画任务相关的持有关系。实线表示 owner 或传入 task 的强引用,虚线表示后台循环使用的弱引用; 图中的 task 状态指闭包实际持有的数据,不是额外定义的一套 worker 框架。
最有用的阅读方法是把“发出停止信号”和“确认 task 已退出”分开。CancellationToken::cancel() 设置共享 取消状态,只有执行流检查它或等待 cancelled() 才会响应;JoinHandle.await 才等待任务结束。 Tokio 的 JoinHandle 被丢弃时会分离任务,本身不会自动 abort 它。
1.2 生命周期分组
| 工作 | 触发源 | 状态的主要 owner | 结果消费者 | 结束方式 |
|---|---|---|---|---|
| 模型刷新 | 启动即刷新,随后 sleep | worker 闭包,manager 用 Weak | 模型目录读取者 | cancel,或 Weak 升级失败 |
| 费用补报 | Core Event、认证变化、150 秒 tick | 单个 WorkerRuntime | SessionTelemetry | cancel、channel 关闭、条目完成或重试淘汰 |
| 技能监听 | 文件变更与节流接收 | SkillsWatcher、注册守卫 | 技能缓存和客户端通知 | token、订阅关闭、注册 Drop |
| OTel 重载 | AuthManager watch | 外层 server 的 reloader task | tracing 层和 exporter | 外层 cancel 后 join |
| MCP 预热 | 配置变脏、认证变化 | 每个 Core Session | runtime 与实际工具调用 | Session 取消并 join |
| 外部迁移尾声 | 一次导入产生 pending 项 | 分离的导入 task | progress/completed 通知、历史记录 | 两个导入分支结束后汇总 |
后两项尤其容易混淆。App Server 的 MCP 配置函数是被请求处理器等待的普通 async 函数,Core 内部另有 预热 worker;迁移尾声则在请求路径中直接 tokio::spawn,没有自动加入通用的后台任务追踪器。
2. 模型刷新节拍
2.1 弱引用循环
模型 worker 的关闭接口很短,但它决定了之后所有退出推理的边界。
源码文件:codex-rs/app-server/src/models_refresh_worker.rs
相关函数/类型:ModelsRefreshWorker / shutdown / Drop(L10–L28,摘录)
// 作者注:持有 JoinHandle 不等于等待它;这里的显式关闭和 Drop 都只发送取消。
const MODELS_REFRESH_INTERVAL: Duration = Duration::from_secs(4 * 60 + 30);
#[derive(Debug)]
pub(crate) struct ModelsRefreshWorker {
shutdown: CancellationToken,
_task: JoinHandle<()>,
}
impl ModelsRefreshWorker {
pub(crate) fn shutdown(&self) {
self.shutdown.cancel();
}
}
impl Drop for ModelsRefreshWorker {
fn drop(&mut self) {
self.shutdown();
}
}shutdown() 与 Drop 做同一件事:取消 token。_task 的下划线表示 owner 不主动读取其结果,不能把 这个字段解释成“析构时自动等待”。真正何时结束,必须继续看循环。
源码文件:codex-rs/app-server/src/models_refresh_worker.rs
相关函数/类型:spawn_with_interval(L37–L68,摘录)
// 作者注:弱引用只在一次刷新期间升级;刷新完成后才开始 sleep,网络 await 本身未放入取消 select。
fn spawn_with_interval(
models_manager: &SharedModelsManager,
http_client_factory: HttpClientFactory,
refresh_interval: Duration,
) -> ModelsRefreshWorker {
let models_manager = Arc::downgrade(models_manager);
let shutdown = CancellationToken::new();
let worker_shutdown = shutdown.clone();
let task = tokio::spawn(async move {
loop {
if worker_shutdown.is_cancelled() {
break;
}
let Some(models_manager) = models_manager.upgrade() else {
break;
};
models_manager
.list_models(RefreshStrategy::Online, http_client_factory.clone())
.await;
// 作者注:在 270 秒等待前释放强引用,避免循环仅为保活而拥有 manager。
drop(models_manager);
tokio::select! {
_ = worker_shutdown.cancelled() => break,
_ = tokio::time::sleep(refresh_interval) => {}
}
}
});
ModelsRefreshWorker {
shutdown,
_task: task,
}
}循环在开始时检查取消,再把 Weak 升级为一次刷新所需的强引用。如果 manager 已经被释放,就退出。 list_models(Online) 在第一次循环立即执行;完成后显式 drop(models_manager),再等待 270 秒。
因此,相邻两次刷新开始的间隔大致是“上一轮刷新耗时 + 270 秒 + 调度延迟”。这里没有固定周期的 tokio::time::interval,也不会因为一次慢请求错过多个 tick 而补跑几次刷新。循环自身逐次 await, 所以这个 worker 不会同时启动两轮刷新;其他调用方是否也在刷新,是 manager 的另一个并发问题。
弱引用也没有使当前刷新变成可取消。升级成功后,这一轮局部变量会持有 manager,直到 await 返回。 取消发生在 sleep 中时可以较快退出;发生在网络读取中时,当前调用仍可能完成、更新缓存,然后才退出。
2.2 目录消费者
worker 没有使用 list_models 返回的列表。它调用这个 API 的目的,是触发 manager 内部的目录更新。
源码文件:codex-rs/models-manager/src/manager.rs
相关函数/类型:ModelsManager::list_models / OpenAiModelsManager::raw_model_catalog(L88–L105、L340–L354,摘录)
// 作者注:worker 忽略返回列表,但 manager 内部刷新缓存;失败被记录后仍读取已有目录。
fn list_models(
&self,
refresh_strategy: RefreshStrategy,
http_client_factory: HttpClientFactory,
) -> ModelsManagerFuture<'_, Vec<ModelPreset>> {
Box::pin(
async move {
let catalog = self
.raw_model_catalog(refresh_strategy, http_client_factory)
.await;
self.build_available_models(catalog.models)
}
.instrument(tracing::info_span!(
"list_models",
refresh_strategy = %refresh_strategy
)),
)
}
// ...
async fn raw_model_catalog(
&self,
refresh_strategy: RefreshStrategy,
http_client_factory: HttpClientFactory,
) -> ModelsResponse {
if let Err(err) = self
.refresh_available_models(refresh_strategy, &http_client_factory)
.await
{
error!("failed to refresh available models: {err}");
}
ModelsResponse {
models: self.get_remote_models().await,
}
}默认 trait 实现先读 raw catalog,再构造可见模型列表。OpenAI manager 捕获刷新错误后仍返回当前内存目录, 因此 worker 不需要 match Result,也不能依据“future 正常返回”认定远端读取成功。用户可能继续看到旧目录, 日志里同时出现 failed to refresh available models。
Online 的含义还受 manager 能力检查约束。
源码文件:codex-rs/models-manager/src/manager.rs
相关函数/类型:refresh_available_models / should_refresh_models(L375–L388、L405–L409、L437–L439,局部节选)
// 作者注:Online 分支仍位于能力 gate 之后;一次循环不保证发生一次网络请求。
// ...
async fn refresh_available_models(
&self,
refresh_strategy: RefreshStrategy,
http_client_factory: &HttpClientFactory,
) -> CoreResult<()> {
if !self.should_refresh_models().await {
if matches!(
refresh_strategy,
RefreshStrategy::Offline | RefreshStrategy::OnlineIfUncached
) {
self.try_load_cache().await;
}
return Ok(());
}
// ...
RefreshStrategy::Online => {
// Always fetch from network
self.fetch_and_update_models(http_client_factory).await
}
}
// ...
async fn should_refresh_models(&self) -> bool {
self.endpoint_client.uses_codex_backend().await || self.endpoint_client.has_command_auth()
}
// ...should_refresh_models() 先判断是否使用 Codex backend 或拥有 command auth。通过后,Online 才进入 远端获取分支;不能用 worker 循环次数作为 HTTP 请求数。如果选择的是静态目录实现,也应沿其 trait 实现判断。
源码文件:codex-rs/models-manager/src/manager.rs
相关函数/类型:fetch_and_update_models(L412–L435,摘录)
// 作者注:远端读取成功后先更新内存目录和 ETag;落盘失败仅记录日志,不撤销内存更新。
async fn fetch_and_update_models(
&self,
http_client_factory: &HttpClientFactory,
) -> CoreResult<()> {
let client_version = crate::client_version_to_whole();
let (models, etag) = self
.endpoint_client
.list_models(&client_version, http_client_factory.clone())
.await?;
self.apply_remote_models(models.clone()).await;
*self.etag.write().await = etag.clone();
if let Some(cache) = self.cache.as_ref() {
let entry = ModelsCacheEntry {
fetched_at: Utc::now(),
etag,
client_version: Some(client_version),
models,
};
if let Err(err) = cache.store(&entry).await {
error!("failed to write models cache: {err}");
}
}
Ok(())
}远端读取失败时不会越过 ? 更新目录。成功时先写内存模型与 ETag,再尝试保存缓存;缓存保存失败不会 回滚内存更新。这解释了“当前进程已看到新模型,但下次启动仍可能读旧磁盘缓存”的条件路径。 目录筛选与缓存内容详见 模型目录加载。
2.3 关闭回归
测试特意把第一次刷新做成失败、第二次做成可控阻塞,然后在阻塞期间释放 worker。
源码文件:codex-rs/app-server/src/models_refresh_worker_tests.rs
相关函数/类型:TestModelsEndpoint::list_models / refreshes_immediately_periodically_and_stops_when_dropped(L62–L72、L76–L97,局部节选)
// 作者注:fixture 第一次失败、第二次阻塞;drop worker 后显式放行第二次,最后断言没有第三次。
// ...
Box::pin(async move {
let fetch_index = self.fetch_count.fetch_add(1, Ordering::SeqCst);
self.fetched.notify_one();
if fetch_index == 0 {
return Err(CodexErr::Io(std::io::Error::other("test failure")));
}
if fetch_index == 1 {
self.release_second_fetch.notified().await;
}
Ok((Vec::new(), None))
})
// ...
#[tokio::test]
async fn refreshes_immediately_periodically_and_stops_when_dropped() {
let codex_home = tempdir().expect("temp dir");
let endpoint = TestModelsEndpoint::new();
let models_manager: SharedModelsManager = Arc::new(OpenAiModelsManager::new(
codex_home.path().to_path_buf(),
endpoint.clone(),
/*auth_manager*/ None,
));
let worker = spawn_with_interval(
&models_manager,
HttpClientFactory::new(OutboundProxyPolicy::ReqwestDefault),
Duration::from_millis(10),
);
endpoint.wait_for_fetch_count(/*expected*/ 2).await;
drop(worker);
endpoint.release_second_fetch.notify_one();
tokio::time::sleep(Duration::from_millis(30)).await;
assert_eq!(endpoint.fetch_count.load(Ordering::SeqCst), 2);
}
// ...wait_for_fetch_count(2) 说明第一次错误没有杀死循环。drop(worker) 之后还要 release_second_fetch.notify_one(),表明测试允许在途读取结束。最终计数仍为 2,证明停止信号阻止了后续轮次。 它没有断言网络 future 在 Drop 时立即被取消,也没有测量生产环境每轮恰好间隔 270 秒。
3. 费用补报状态
3.1 启动条件
费用补报处理的是“Turn 已结束,后端稍后才能返回完整计费结果”。它不在当前 Turn 的关键路径里等待结算, 而是先记住必要状态,再轮询后端。下面几个限制分别约束时间、队列、状态表和请求大小。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:POLL_INTERVAL / TurnCostEntry / BackendAvailability(L24–L29、L61–L74、L89–L95,摘录)
// 作者注:队列容量、追踪上限、批次和重试次数约束不同资源;费用状态与后端可用性也彼此独立。
const POLL_INTERVAL: Duration = Duration::from_secs(150);
const REQUEST_TIMEOUT: Duration = Duration::from_secs(15);
const OBSERVATION_CHANNEL_CAPACITY: usize = 16_384;
const MAX_TRACKED_TURNS: usize = 4_096;
const MAX_QUERY_TURNS: usize = 100;
const MAX_STALLED_POLL_ATTEMPTS: u8 = 5;
// ...
enum TurnCostStatus {
Running,
Completed,
Interrupted,
}
struct TurnCostEntry {
thread_id: ThreadId,
session_telemetry: SessionTelemetry,
expected_response_count: u64,
status: TurnCostStatus,
next_poll_at: Instant,
attempt_count: u8,
}
// ...
enum BackendAvailability {
AwaitingAuthChange,
RetryProbe,
Ready,
Disabled,
}TurnCostStatus 描述被观察 Turn 的状态,BackendAvailability 描述是否值得向费用后端请求。 后端可用不代表某个 Turn 的价格已就绪,Turn 完成也不代表后台一定能拿到完整价格。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:TurnCostWorker::spawn(L97–L140,摘录)
// 作者注:日志或 metrics exporter 任一启用即可;Bedrock 排除,OpenAI 与其他 provider 选不同后端。
pub(crate) fn spawn(config: Arc<Config>, auth_manager: Arc<AuthManager>) -> Option<Self> {
let has_otel_log_exporter = matches!(
config.otel.exporter,
OtelExporterKind::OtlpHttp { .. } | OtelExporterKind::OtlpGrpc { .. }
);
let has_otel_metrics_exporter = matches!(
config.otel.metrics_exporter,
OtelExporterKind::OtlpHttp { .. } | OtelExporterKind::OtlpGrpc { .. }
);
if !(has_otel_log_exporter || has_otel_metrics_exporter)
|| config.model_provider.is_amazon_bedrock()
{
return None;
}
let is_openai = config.model_provider.is_openai();
let backend = if is_openai {
TurnCostBackend::OpenAiApiKey(Arc::clone(&auth_manager))
} else {
TurnCostBackend::ModelProvider(create_model_provider(
config.model_provider.clone(),
Some(Arc::clone(&auth_manager)),
))
};
let (sender, receiver) = mpsc::channel(OBSERVATION_CHANNEL_CAPACITY);
let shutdown = CancellationToken::new();
let runtime = WorkerRuntime {
config: Arc::clone(&config),
backend: backend.clone(),
turns: HashMap::new(),
};
let worker_shutdown = shutdown.clone();
let task = tokio::spawn(async move {
runtime.run(receiver, worker_shutdown).await;
});
Some(Self {
handle: TurnCostWorkerHandle {
sender,
backend,
config,
},
shutdown,
_task: task,
})
}启用 OTLP 日志或 metrics exporter 即满足第一道条件;只配置 trace exporter 不满足这两个判断。 Bedrock 在这里直接排除。OpenAI provider 使用 AuthManager 驱动的 API-key backend,其他 provider 通过 create_model_provider 获取接口与认证。
启动配置以 Arc<Config> 留在 handle 和 runtime 中。它不是每次观察都重新加载的全局最新配置; 后续线程的 provider 与它不匹配时,会在入口被过滤。这也是费用缺失不能只排查网络的原因。
3.2 观察入口
费用观察发生在 Thread listener 读取原始 Core 事件之后、对外 typed 通知翻译之前。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:ensure_listener_task_running(L313–L329,局部节选)
// 作者注:原始 Core Event 在 typed 翻译前被费用观察器读取,Started 才按需获取 SessionTelemetry。
// ...
event = conversation.next_event() => {
let event = match event {
Ok(event) => event,
Err(err) => {
tracing::warn!("thread.next_event() failed with: {err}");
break;
}
};
if let Some(worker) = &turn_cost_worker {
worker.observe_event(
conversation_id,
config.as_ref(),
&event,
|| conversation.session_telemetry(),
);
}
// ...传入的是 listener 当时使用的配置和一个懒执行的 telemetry 闭包。只有开始事件需要保存 SessionTelemetry,其他事件不必重新查 Thread 或构造它。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:TurnCostWorkerHandle::observe_event(L158–L190,摘录)
// 作者注:provider 匹配与 API-key 门控先执行;try_send 不等待容量,丢失观察不会阻塞主事件链。
pub(crate) fn observe_event(
&self,
thread_id: ThreadId,
thread_config: &Config,
event: &Event,
session_telemetry: impl FnOnce() -> SessionTelemetry,
) {
if thread_config.model_provider != self.config.model_provider {
return;
}
if let TurnCostBackend::OpenAiApiKey(auth_manager) = &self.backend {
let Some(auth) = auth_manager.auth_cached() else {
return;
};
if !auth.is_api_key_auth() {
return;
}
}
let kind = match &event.msg {
EventMsg::TurnStarted(_) => TurnCostObservationKind::Started {
session_telemetry: Box::new(session_telemetry()),
},
EventMsg::RawResponseCompleted(_) => TurnCostObservationKind::ResponseCompleted,
EventMsg::TurnComplete(_) => TurnCostObservationKind::Finished { interrupted: false },
EventMsg::TurnAborted(_) => TurnCostObservationKind::Finished { interrupted: true },
_ => return,
};
let _ = self.sender.try_send(TurnCostObservation {
thread_id,
turn_id: event.id.clone(),
kind,
});
}这里有三道信息缩减:不匹配的 provider 不观察,OpenAI 的非 API-key 认证不观察,其余事件只保留 Started、RawResponseCompleted、TurnComplete 和 TurnAborted。入队使用 try_send,容量满或接收端关闭时 错误被忽略;费用观测不会为了补报阻塞主事件流。
代价也很具体:如果 Started 丢失,后续完成事件找不到条目;如果完成事件丢失,已有条目可能一直保持 Running,本文件没有额外的 Running 超时淘汰逻辑。4096 的追踪上限限制内存增长,但不是每个 Turn 都会被计费的保证。
3.3 后端探测
后台启动先发一次探测,再进入事件、认证和定时器的竞争。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:WorkerRuntime::run / run_with_backend_availability(L194–L256,摘录)
// 作者注:先探测后循环;认证变化清空旧 Turn 和排队观察;select 只在分支边界优先处理取消。
async fn run(self, receiver: mpsc::Receiver<TurnCostObservation>, shutdown: CancellationToken) {
let auth_changes = match &self.backend {
TurnCostBackend::OpenAiApiKey(auth_manager) => {
Some(auth_manager.auth_change_receiver())
}
TurnCostBackend::ModelProvider(_) => None,
};
let backend_availability = self.probe_backend().await;
self.run_with_backend_availability(receiver, shutdown, auth_changes, backend_availability)
.await;
}
async fn run_with_backend_availability(
mut self,
mut receiver: mpsc::Receiver<TurnCostObservation>,
shutdown: CancellationToken,
mut auth_changes: Option<tokio::sync::watch::Receiver<u64>>,
mut backend_availability: BackendAvailability,
) {
let mut ticker = tokio::time::interval(POLL_INTERVAL);
// 作者注:消费 interval 的立即 tick,第一次定时动作在后续周期。
ticker.tick().await;
loop {
tokio::select! {
biased;
_ = shutdown.cancelled() => break,
changed = async {
match auth_changes.as_mut() {
Some(auth_changes) => auth_changes.changed().await,
None => std::future::pending().await,
}
} => {
if changed.is_err() {
break;
}
// 作者注:旧账户的未补报费用不迁移到新账户。
self.turns.clear();
while receiver.try_recv().is_ok() {}
backend_availability = self.probe_backend().await;
}
observation = receiver.recv() => {
let Some(observation) = observation else {
break;
};
if !matches!(
backend_availability,
BackendAvailability::Ready | BackendAvailability::RetryProbe
) {
continue;
}
self.record_observation(observation);
}
_ = ticker.tick() => {
match backend_availability {
BackendAvailability::Ready => self.poll_due().await,
BackendAvailability::RetryProbe => {
backend_availability = self.probe_backend().await;
}
BackendAvailability::AwaitingAuthChange
| BackendAvailability::Disabled => {}
}
}
}
}
}探测使用随机 Turn ID,只判断 endpoint 是否可用,不查询某个真实 Turn 的账单。成功返回一个空结果集 仍可以表示 Ready。下面是探测结果到四个枚举值的完整分类。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:WorkerRuntime::probe_backend(L258–L294,摘录)
// 作者注:同为 401/403,OpenAI 等认证变化,自定义 provider 定时重探;其他 4xx 可禁用。
async fn probe_backend(&self) -> BackendAvailability {
let probe_turn_ids = [uuid::Uuid::new_v4().to_string()];
match tokio::time::timeout(REQUEST_TIMEOUT, self.query_turn_costs(&probe_turn_ids)).await {
Ok(Ok(Some(_))) => BackendAvailability::Ready,
Ok(Ok(None)) => match self.backend {
TurnCostBackend::OpenAiApiKey(_) => BackendAvailability::AwaitingAuthChange,
TurnCostBackend::ModelProvider(_) => BackendAvailability::Disabled,
},
Ok(Err(error)) => match error.status().map(|status| status.as_u16()) {
Some(401 | 403) if matches!(self.backend, TurnCostBackend::OpenAiApiKey(_)) => {
tracing::debug!(
"turn cost worker waiting for auth change after backend availability check: {error}"
);
BackendAvailability::AwaitingAuthChange
}
Some(401 | 403 | 429) => BackendAvailability::RetryProbe,
Some(400..=499) => {
tracing::debug!(
"turn cost worker disabled by backend availability check: {error}"
);
BackendAvailability::Disabled
}
_ => {
tracing::debug!(
"turn cost worker backend availability check failed temporarily: {error}"
);
BackendAvailability::RetryProbe
}
},
Err(_) => {
tracing::debug!(
"turn cost worker backend availability check timed out; will retry"
);
BackendAvailability::RetryProbe
}
}
}这张状态图画 RetryProbe 在定时重新探测后的去向。四个节点都是 App Server 状态,因此使用同一模块颜色; 初次探测和认证变化后的探测也使用下方的分类表。
| 探测结果 | OpenAI API-key backend | 自定义 provider backend |
|---|---|---|
Some(costs),包括空列表 | Ready | Ready |
None | AwaitingAuthChange | Disabled |
| 401 / 403 | AwaitingAuthChange | RetryProbe |
| 429 | RetryProbe | RetryProbe |
| 其他 4xx | Disabled | Disabled |
| 其他错误或 15 秒超时 | RetryProbe | RetryProbe |
OpenAI 在任何当前状态下收到 auth change 都先清空 Turn 表, 排空观察队列,再用同一张表重新选择状态。自定义 provider 分支没有订阅这个 auth watch,因而 401/403 需要靠 RetryProbe 的定时路径重新读取 provider 认证,不能永久停在“等 AuthManager 唤醒”。
Ready 在 tick 时查询到期 Turn,RetryProbe 只重新探测;两者都接收观察。AwaitingAuthChange 和 Disabled 会消费但忽略观察,不会为之后补发保存一份历史。
biased 让同时就绪的取消分支优先,但探测和 poll_due().await 已被选中后,内部 HTTP 等待不再与该 token 竞争。它们有自己的 15 秒请求预算;“优先取消”不能解释成随时抢占任意指令。
3.4 结算条件
观察进入单消费者后,不需要为 turns 再加 Mutex。状态表只由这个 runtime 的执行流修改。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:WorkerRuntime::record_observation(L296–L334,摘录)
// 作者注:只给 Running 累加 response 数;完成/中断冻结预期数并允许下次 tick 查询。
fn record_observation(&mut self, observation: TurnCostObservation) {
match observation.kind {
TurnCostObservationKind::Started { session_telemetry } => {
if self.turns.len() < MAX_TRACKED_TURNS {
self.turns
.entry(observation.turn_id)
.or_insert(TurnCostEntry {
thread_id: observation.thread_id,
session_telemetry: *session_telemetry,
expected_response_count: 0,
status: TurnCostStatus::Running,
next_poll_at: Instant::now(),
attempt_count: 0,
});
}
}
TurnCostObservationKind::ResponseCompleted => {
if let Some(entry) = self.turns.get_mut(&observation.turn_id)
&& entry.status == TurnCostStatus::Running
{
entry.expected_response_count = entry.expected_response_count.saturating_add(1);
}
}
TurnCostObservationKind::Finished { interrupted } => {
let Some(entry) = self.turns.get_mut(&observation.turn_id) else {
return;
};
if entry.status != TurnCostStatus::Running {
return;
}
entry.status = if interrupted {
TurnCostStatus::Interrupted
} else {
TurnCostStatus::Completed
};
entry.next_poll_at = Instant::now();
}
}
}Started 保存 telemetry、Thread ID、响应预期数和下一查询时间。每个 RawResponseCompleted 只在 Running 期间增加预期数;终止事件冻结计数,并把查询时间设为现在。真正查询仍要等 worker 的 tick。 键使用 event.id 对应的 Turn ID 字符串,Thread ID 留在条目里用于归因和诊断。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:WorkerRuntime::poll_due / poll_api_key_entries(L336–L379,摘录)
// 作者注:只从已结束且到期的 Turn 选最多 100 项;Some 列表缺项会重试,但 Ok(None) 直接返回。
async fn poll_due(&mut self) {
let now = Instant::now();
let due_turn_ids: Vec<String> = self
.turns
.iter()
.filter(|(_, entry)| {
entry.status != TurnCostStatus::Running && entry.next_poll_at <= now
})
.take(MAX_QUERY_TURNS)
.map(|(turn_id, _)| turn_id.clone())
.collect();
if !due_turn_ids.is_empty() {
self.poll_api_key_entries(&due_turn_ids).await;
}
}
async fn poll_api_key_entries(&mut self, turn_ids: &[String]) {
let costs =
match tokio::time::timeout(REQUEST_TIMEOUT, self.query_turn_costs(turn_ids)).await {
Ok(Ok(Some(costs))) => costs,
Ok(Ok(None)) => return,
Ok(Err(error)) => {
warn!("failed to query API-key turn costs: {error}");
self.retry_entries(turn_ids);
return;
}
Err(_) => {
warn!("timed out querying API-key turn costs");
self.retry_entries(turn_ids);
return;
}
};
let costs_by_turn: HashMap<String, ApiKeyTurnCost> = costs
.into_iter()
.map(|cost| (cost.turn_id.clone(), cost))
.collect();
for turn_id in turn_ids {
let Some(cost) = costs_by_turn.get(turn_id) else {
self.retry_entry(turn_id);
continue;
};
self.process_api_key_cost(turn_id, cost);
}
}最多选取 100 个已结束且到期的条目,集合来自 HashMap 迭代,不能承诺按结束时间公平排序。 HTTP 错误或超时会让本批条目安排重试;成功列表缺少某个 Turn 也会重试。注意 Ok(None) 直接返回, 没有递增失败次数,不能把所有“没拿到价格”都归为同一种失败计数。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:WorkerRuntime::process_api_key_cost(L439–L475,摘录)
// 作者注:Priced、金额和足够的 response 数同时满足才移除条目并写遥测;responses 优先于 event_count。
fn process_api_key_cost(&mut self, turn_id: &str, cost: &ApiKeyTurnCost) {
if cost.status != ApiKeyTurnCostStatus::Priced {
self.retry_entry(turn_id);
return;
}
let response_count = cost
.responses
.as_ref()
.map(|responses| responses.len() as u64)
.or(cost.event_count);
let (Some(total_usd), Some(response_count)) = (cost.total_usd.as_deref(), response_count)
else {
self.retry_entry(turn_id);
return;
};
let Some(entry) = self.turns.get(turn_id) else {
return;
};
if response_count < entry.expected_response_count {
self.retry_entry(turn_id);
return;
}
let mut session_telemetry = entry.session_telemetry.clone();
if let Some(model) = cost.model.as_deref() {
session_telemetry = session_telemetry.with_model(model, model);
}
let Some(entry) = self.turns.remove(turn_id) else {
return;
};
session_telemetry.record_turn_cost(
turn_id,
total_usd,
entry.status == TurnCostStatus::Interrupted,
cost.speed.as_deref(),
cost.reasoning_effort.as_deref(),
);
}Priced 只是必要条件。金额、响应数量也要齐全,并且后端响应数不能少于本地观察数。 如果后端同时提供 responses 明细与 event_count,代码先用明细长度,避免一个较大的汇总数掩盖未结算响应。 成功后先从表中移除条目,再调用捕获的 telemetry;不要求 Thread 仍可从 ThreadManager 查到。 这一步表示提交遥测,exporter 的交付结果仍有独立边界。
源码文件:codex-rs/app-server/src/turn_cost_worker.rs
相关函数/类型:WorkerRuntime::retry_entry(L483–L499,摘录)
// 作者注:第 5 次无进展就移除;此前把下一次截止时间推迟 150 秒。
fn retry_entry(&mut self, turn_id: &str) {
let Some(entry) = self.turns.get_mut(turn_id) else {
return;
};
entry.attempt_count = entry.attempt_count.saturating_add(1);
if entry.attempt_count >= MAX_STALLED_POLL_ATTEMPTS {
warn!(
thread_id = %entry.thread_id,
turn_id,
attempts = MAX_STALLED_POLL_ATTEMPTS,
"dropping turn cost event after repeated unsuccessful polls"
);
self.turns.remove(turn_id);
return;
}
entry.next_poll_at = Instant::now() + POLL_INTERVAL;
}失败次数累计到 5 时删除条目,之前每次将下一查询时刻推迟 150 秒。这里没有指数退避,也没有重试到成功为止。 观测队列、认证切换和追踪容量都可能更早丢失条目,因此这套机制不能用于要求逐 Turn 完整对账的账务系统。
3.5 测试窗口
下面的认证测试先确认未登录时任务没有退出,再写入测试 API key 并让 AuthManager reload。 测试准备的 mock endpoint 期望收到一次请求。
源码文件:codex-rs/app-server/src/turn_cost_worker_tests.rs
相关函数/类型:worker_waits_for_late_api_key_login(L126–L151,局部节选)
// 作者注:先证明等待登录时 task 仍存活,再 reload API-key 并等待一次实际 mock HTTP 请求。
// ...
let auth_home = TempDir::new().expect("temporary auth home");
let auth_manager = auth_manager_at(auth_home.path()).await;
let runtime = test_runtime(&server, Arc::clone(&auth_manager)).await;
let (_sender, receiver) = mpsc::channel(OBSERVATION_CHANNEL_CAPACITY);
let shutdown = CancellationToken::new();
let mut task = tokio::spawn(runtime.run(receiver, shutdown.clone()));
assert!(
timeout(Duration::from_millis(/*millis*/ 50), &mut task)
.await
.is_err(),
"worker exited while waiting for login"
);
login_with_api_key(
auth_home.path(),
"sk-test",
AuthCredentialsStoreMode::File,
AuthKeyringBackendKind::default(),
)
.expect("write API-key auth");
assert!(auth_manager.reload().await);
wait_for_request_count(&server, /*expected*/ 1).await;
shutdown.cancel();
task.await.expect("worker task");
server.verify().await;
// ...50 毫秒 timeout 是“任务尚未完成”的断言,不是生产重试间隔。真正的恢复证据是 reload 返回变化,随后 mock 收到请求;测试最终显式 cancel 并等待任务结束。它不声称 custom provider 也依赖同一 watch。
价格完整性测试则直接构造一个预期两次响应的 Completed 条目。
源码文件:codex-rs/app-server/src/turn_cost_worker_tests.rs
相关函数/类型:priced_cost_waits_for_every_response_when_response_costs_are_available(L344–L384,局部节选)
// 作者注:预期两次 response,明细只有一项时即使 event_count=2 也保留;补第二项后才删除。
// ...
runtime.turns.insert(
turn_id.to_string(),
TurnCostEntry {
thread_id,
session_telemetry: test_session_telemetry(thread_id),
expected_response_count: 2,
status: TurnCostStatus::Completed,
next_poll_at: Instant::now(),
attempt_count: 0,
},
);
let mut cost = ApiKeyTurnCost {
turn_id: turn_id.to_string(),
status: ApiKeyTurnCostStatus::Priced,
total_usd: Some("1.25".to_string()),
event_count: Some(2),
responses: Some(vec![ApiKeyResponseCost {
response_id: "resp-one".to_string(),
total_usd: "1.25".to_string(),
}]),
model: Some("gpt-5.6".to_string()),
speed: Some("fast".to_string()),
reasoning_effort: Some("high".to_string()),
};
runtime.process_api_key_cost(turn_id, &cost);
let entry = runtime.turns.get(turn_id).expect("turn remains tracked");
assert_eq!(entry.attempt_count, 1);
cost.event_count = None;
cost.responses
.as_mut()
.expect("response costs")
.push(ApiKeyResponseCost {
response_id: "resp-two".to_string(),
total_usd: "0.50".to_string(),
});
runtime.process_api_key_cost(turn_id, &cost);
assert_eq!(runtime.turns.len(), 0);
// ...第一次输入虽然 event_count=2,但明细只有一条,attempt_count 增为 1 且条目保留;补入第二条后,即使 清掉 event_count,条目也被移除。它验证的是本地结算条件,没有断言某个真实后端何时完成计价。
4. 技能变更通知
4.1 注册守卫
技能监听首先建立订阅、取消关系和节流间隔。
源码文件:codex-rs/app-server/src/skills_watcher.rs
相关函数/类型:WATCHER_THROTTLE_INTERVAL / SkillsWatcher / SkillsWatcher::new(L25–L67,摘录)
// 作者注:生产节流 10 秒,单元编译 50 毫秒;watcher 初始化失败降为 noop,DropGuard 负责最后取消。
#[cfg(not(test))]
const WATCHER_THROTTLE_INTERVAL: Duration = Duration::from_secs(10);
#[cfg(test)]
const WATCHER_THROTTLE_INTERVAL: Duration = Duration::from_millis(50);
pub(crate) struct SkillsWatcher {
subscriber: FileWatcherSubscriber,
runtime_extra_roots_registration: Mutex<WatchRegistration>,
shutdown_token: CancellationToken,
_shutdown_drop_guard: DropGuard,
}
impl SkillsWatcher {
pub(crate) fn new(
skills_service: Arc<HostSkillsService>,
codex_home: &AbsolutePathBuf,
outgoing: Arc<OutgoingMessageSender>,
) -> Arc<Self> {
let file_watcher = match FileWatcher::new() {
Ok(file_watcher) => Arc::new(file_watcher),
Err(err) => {
warn!("failed to initialize skills file watcher: {err}");
Arc::new(FileWatcher::noop())
}
};
let (subscriber, rx) = file_watcher.add_subscriber();
let shutdown_token = CancellationToken::new();
let shutdown_drop_guard = shutdown_token.clone().drop_guard();
let system_skills_root = system_cache_root_dir(codex_home);
Self::spawn_event_loop(
rx,
skills_service,
system_skills_root,
outgoing,
shutdown_token.child_token(),
);
Arc::new(Self {
subscriber,
runtime_extra_roots_registration: Mutex::new(WatchRegistration::default()),
shutdown_token,
_shutdown_drop_guard: shutdown_drop_guard,
})
}初始化 OS watcher 失败时仍返回 SkillsWatcher,但底层是 noop;“对象存在”不能证明正在监听磁盘。 DropGuard 在最后一个 watcher owner 释放时取消 token。生产二进制使用 10 秒节流,单元测试编译使用 50 毫秒,不能把测试中的时间常数当作用户配置项。
每个 Thread 实际监听哪些根,还要检查它选中的执行环境。
源码文件:codex-rs/app-server/src/skills_watcher.rs
相关函数/类型:SkillsWatcher::register_thread_config(L89–L131,摘录)
// 作者注:只从第一个环境取根;缺环境、未知环境或远程环境返回空注册,不假装监听远端磁盘。
pub(crate) async fn register_thread_config(
&self,
config: &Config,
thread_manager: &ThreadManager,
environments: &[TurnEnvironmentSelection],
) -> WatchRegistration {
let Some(environment_selection) = environments.first() else {
return WatchRegistration::default();
};
let Some(environment) = thread_manager
.environment_manager()
.get_environment(&environment_selection.environment_id)
else {
warn!(
"failed to register skills watcher for unknown environment `{}`",
environment_selection.environment_id
);
return WatchRegistration::default();
};
if environment.is_remote() {
return WatchRegistration::default();
}
let plugins_input = config.plugins_config_input();
let plugins_manager = thread_manager.plugins_manager();
let plugin_outcome = 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(),
);
let roots = thread_manager
.skills_service()
.watchable_skill_root_paths(&skills_input, environment.get_filesystem())
.await
.into_iter()
.map(|path| WatchPath {
path: path.into_path_buf(),
recursive: true,
})
.collect();
self.subscriber.register_paths(roots)
}只检查第一个环境选择;没有环境、环境 ID 无法解析或环境是远程时,返回空注册。正常路径把有效插件 技能根、cwd 和配置层交给 skills service,获取可监听根后递归注册。它不会因为模型运行在远端,就把 主机 notify 自动变成远端文件事件桥。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:ensure_listener_task_running(L239–L248、L253–L263,局部节选)
// 作者注:注册结果移交给 ThreadState 的 listener;提前返回时临时 registration 自动释放。
// ...
let config = conversation.config().await;
let environments = conversation.environment_selections().await;
let watch_registration = listener_task_context
.skills_watcher
.register_thread_config(
config.as_ref(),
listener_task_context.thread_manager.as_ref(),
&environments,
)
.await;
// ...
let (mut listener_command_rx, listener_generation) = {
let mut thread_state = thread_state.lock().await;
if thread_state.listener_matches(&conversation) {
return Ok(());
}
let (listener_command_rx, listener_generation) = thread_state.set_listener(
cancel_tx,
&conversation,
watch_registration,
thread_settings_baseline,
);
// ...Thread listener 建立时接收注册守卫,并将其交给 ThreadState。如果发现相同 listener 已存在而提前返回, 这次临时得到的注册会自动 Drop;已有 listener 的注册仍由已有 state 持有。
源码文件:codex-rs/app-server/src/thread_state.rs
相关函数/类型:ThreadState::clear_listener(L148–L157,摘录)
// 作者注:将 registration 替换为空值会 Drop 旧注册;取消监听与注销 OS 路径引用是两件相接的事。
pub(crate) fn clear_listener(&mut self) {
if let Some(cancel_tx) = self.cancel_tx.take() {
let _ = cancel_tx.send(());
}
self.shutdown_drain_waiter = None;
self.listener_command_tx = None;
self.current_turn_history.reset();
self.listener_thread = None;
self.watch_registration = WatchRegistration::default();
}清 listener 会用空注册替换旧值,从而执行旧守卫的析构。底层守卫如何注销,不能只凭 RAII 这个名字猜测。
源码文件:codex-rs/file-watcher/src/lib.rs
相关函数/类型:WatchRegistration / Drop(L345–L350、L362–L368,摘录)
// 作者注:注册守卫对 watcher 只有 Weak;最后 Drop 时注销本订阅的路径引用。
/// RAII guard for a set of active path registrations.
pub struct WatchRegistration {
file_watcher: std::sync::Weak<FileWatcher>,
subscriber_id: SubscriberId,
watched_paths: Vec<SubscriberWatchKey>,
}
// ...
impl Drop for WatchRegistration {
fn drop(&mut self) {
if let Some(file_watcher) = self.file_watcher.upgrade() {
file_watcher.unregister_paths(self.subscriber_id, &self.watched_paths);
}
}
}守卫保存 subscriber ID 与注册的路径集合,对 FileWatcher 只留 Weak。释放时若 watcher 还在,就注销 本注册拥有的路径引用;同一路径可能仍被其他注册持有,并不一定马上取消整个 OS watch。 另一个 runtime_extra_roots_registration 用 Mutex 持有运行时额外根,替换注册采用相同释放规则。
4.2 合并与节流
文件系统可以对一次保存产生多次通知。底层先把变化路径合入集合,再唤醒消费者。
源码文件:codex-rs/file-watcher/src/lib.rs
相关函数/类型:ReceiverInner / Receiver::recv / WatchSender::add_changed_paths(L107–L147,摘录)
// 作者注:路径存入 BTreeSet 合并去重;Notify 只负责唤醒,真正数据保存在集合里。
struct ReceiverInner {
changed_paths: AsyncMutex<BTreeSet<PathBuf>>,
notify: Notify,
sender_count: AtomicUsize,
}
impl Receiver {
/// Waits for the next batch of changed paths, or returns `None` once the
/// corresponding subscriber has been removed and no more events can arrive.
pub async fn recv(&mut self) -> Option<FileWatcherEvent> {
loop {
let notified = self.inner.notify.notified();
{
let mut changed_paths = self.inner.changed_paths.lock().await;
if !changed_paths.is_empty() {
return Some(FileWatcherEvent {
paths: std::mem::take(&mut *changed_paths).into_iter().collect(),
});
}
if self.inner.sender_count.load(Ordering::Acquire) == 0 {
return None;
}
}
notified.await;
}
}
}
impl WatchSender {
async fn add_changed_paths(&self, paths: &[PathBuf]) {
if paths.is_empty() {
return;
}
let mut changed_paths = self.inner.changed_paths.lock().await;
let previous_len = changed_paths.len();
changed_paths.extend(paths.iter().cloned());
if changed_paths.len() != previous_len {
self.inner.notify.notify_one();
}
}Notify 不承载事件内容,BTreeSet<PathBuf> 才承载待处理路径。接收时用 mem::take 原子取走当前集合; 发送者数变为零且集合为空时才返回 None。所以停止发送者后,底层接收器仍可能先交出最后一批路径。
源码文件:codex-rs/file-watcher/src/lib.rs
相关函数/类型:ThrottledWatchReceiver::recv(L241–L253,摘录)
// 作者注:限速截止时间取自上一次发出;新文件事件不重置窗口,首次接收没有额外等待。
/// Receives the next event, enforcing the configured minimum delay after
/// the previous emission.
pub async fn recv(&mut self) -> Option<FileWatcherEvent> {
if let Some(next_allowed) = self.next_allowed {
sleep_until(next_allowed).await;
}
let event = self.rx.recv().await;
if event.is_some() {
self.next_allowed = Some(Instant::now() + self.interval);
}
event
}这是 throttle:从上一次交出事件开始,至少间隔一个窗口后再交下一批。首次没有 next_allowed,收到变化 就可交出;新事件只在集合中合并,不把截止时间不断往后延。它与“最后一次变化后静默一段时间才发出”的 debounce 不同,也不保留每一次底层文件操作的顺序。
4.3 缓存先失效
App Server 消费合并事件时,还有系统技能过滤和通知副作用。
源码文件:codex-rs/app-server/src/skills_watcher.rs
相关函数/类型:SkillsWatcher::spawn_event_loop(L133–L170,摘录)
// 作者注:先清缓存再广播通知;仅当整批路径都在 system root 下才忽略。
fn spawn_event_loop(
rx: Receiver,
skills_service: Arc<HostSkillsService>,
system_skills_root: AbsolutePathBuf,
outgoing: Arc<OutgoingMessageSender>,
shutdown_token: CancellationToken,
) {
let mut rx = ThrottledWatchReceiver::new(rx, WATCHER_THROTTLE_INTERVAL);
let Ok(handle) = tokio::runtime::Handle::try_current() else {
warn!("skills watcher listener skipped: no Tokio runtime available");
return;
};
handle.spawn(async move {
loop {
// 作者注:取消仅竞争 recv;已开始的通知发送不在这一 select 内。
let event = tokio::select! {
_ = shutdown_token.cancelled() => break,
event = rx.recv() => event,
};
let Some(event) = event else {
break;
};
// The legacy user-skills root contains `.system` and is watched recursively.
if event
.paths
.iter()
.all(|path| path.starts_with(system_skills_root.as_path()))
{
continue;
}
skills_service.clear_cache();
outgoing
.send_server_notification(ServerNotification::SkillsChanged(
SkillsChangedNotification {},
))
.await;
}
});
}过滤条件是 all:整批路径都在 system skills root 下才跳过。只要混合批次还包含一个普通技能路径, 就清空技能缓存并发送 SkillsChanged。该通知不携带完整技能列表,客户端仍需通过列表 API 获取新值。
下图展示一次普通本地技能变更。第一次通知不强制等待 10 秒;后续接收受上一批发出时间限制。
列表请求由 catalog processor 处理,再通过 outgoing 返回;缓存先失效使通知后的 force_reload=false 读取也能看到重新加载的数据。通知入队不保证客户端已经收到,队列细节可接着读 OutgoingMessage写队列。
取消只与 rx.recv() 竞争,清缓存后正在等待的通知发送不在同一个 select 里。底层 receiver 的“可交出 尾批”也不意味着上层 SkillsWatcher 取消时必定把尾批发给客户端;上层可以先选中取消分支退出。
4.4 两层断言
节流测试直接使用内存 watch channel,分离算法语义与 OS 事件的不确定性。
源码文件:codex-rs/file-watcher/src/file_watcher_tests.rs
相关函数/类型:throttled_receiver_coalesces_within_interval(L23–L52,摘录)
// 作者注:首批 a 立即接收,半个窗口内第二次 recv 超时,随后得到合并的 b/c。
#[tokio::test]
async fn throttled_receiver_coalesces_within_interval() {
let (tx, rx) = watch_channel();
let mut throttled = ThrottledWatchReceiver::new(rx, TEST_THROTTLE_INTERVAL);
tx.add_changed_paths(&[path("a")]).await;
let first = timeout(Duration::from_secs(1), throttled.recv())
.await
.expect("first emit timeout");
assert_eq!(
first,
Some(FileWatcherEvent {
paths: vec![path("a")],
})
);
tx.add_changed_paths(&[path("b"), path("c")]).await;
let blocked = timeout(TEST_THROTTLE_INTERVAL / 2, throttled.recv()).await;
assert_eq!(blocked.is_err(), true);
let second = timeout(TEST_THROTTLE_INTERVAL * 2, throttled.recv())
.await
.expect("second emit timeout");
assert_eq!(
second,
Some(FileWatcherEvent {
paths: vec![path("b"), path("c")],
})
);
}它验证首批 a、半窗口内的 timeout 和之后合并的 b/c,不会证明所有操作系统对某种编辑器保存策略 产生同样事件。注册 Drop 测试使用 noop watcher 检查路径计数清除,也不能替代真实文件监听测试。
真实 App Server 测试先建立 Thread、读取初始技能,再写入新描述。
源码文件:codex-rs/app-server/tests/suite/v2/skills_list.rs
相关函数/类型:skills_changed_notification_is_emitted_after_skill_change(L1232–L1238、L1303–L1332,局部节选)
// 作者注:fixture 已建立 Thread 和初始技能缓存;修改文件后等通知,force_reload=false 仍读到新描述。
// ...
#[tokio::test]
async fn skills_changed_notification_is_emitted_after_skill_change() -> Result<()> {
// TODO(anp): Remove after skill watching can bridge host-local storage into remote exec.
skip_if_remote!(
Ok(()),
"host-local skill changes are not visible to remote executors"
);
// ...
let skill_path = codex_home
.path()
.join("skills")
.join("demo")
.join("SKILL.md");
std::fs::write(
&skill_path,
"---\nname: demo\ndescription: updated\n---\n\n# Updated\n",
)?;
expect_skills_changed_notification(&mut mcp, WATCHER_TIMEOUT).await?;
let updated_skills_request_id = mcp
.send_skills_list_request(SkillsListParams {
cwds: vec![codex_home.path().to_path_buf()],
force_reload: false,
})
.await?;
let SkillsListResponse { data } = timeout(
DEFAULT_TIMEOUT,
mcp.read_response(updated_skills_request_id),
)
.await??;
assert_eq!(data.len(), 1);
assert!(
data[0]
.skills
.iter()
.any(|skill| skill.name == "demo" && skill.description == "updated")
);
Ok(())
// ...后一次请求明确设置 force_reload: false,最终仍断言描述变成 updated,把“收到通知”和“缓存已经失效” 连起来。skip_if_remote! 说明远程执行器不能观察主机本地技能变化的分支不属于这个测试的覆盖范围。
5. 账户触发重载
5.1 唤醒来源
OTel reloader 的 owner 位于 App Server 外层启动函数,关闭 token 与 transport 共享。
源码文件:codex-rs/app-server/src/lib.rs
相关函数/类型:run_main_with_transport_options(L839–L846,局部节选)
// 作者注:外层 server 持有 reloader handle,并把 transport shutdown token 传给它。
// ...
let otel_reloader_handle = otel_reloader::spawn(
otel,
otel_logger_reload_handle,
config_manager.clone(),
Arc::clone(&auth_manager),
default_analytics_enabled,
transport_shutdown_token.clone(),
);
// ...这个任务返回 JoinHandle<()> 供外层退出时等待,不是 MessageProcessor 的普通成员 worker。 它也没有监听配置文件修改,而是等待认证状态变化。
源码文件:codex-rs/app-server/src/otel_reloader.rs
相关函数/类型:spawn(L55–L75,局部节选)
// 作者注:唤醒源是 AuthManager 的 watch;50 毫秒延迟给账户处理器安装 cloud loader 留出机会。
// ...
let mut auth_changes = auth_manager.auth_change_receiver();
tokio::spawn(async move {
loop {
tokio::select! {
_ = shutdown_token.cancelled() => break,
changed = auth_changes.changed() => {
if changed.is_err() {
break;
}
// Account handlers install the new cloud loader after publishing auth changes.
tokio::time::sleep(Duration::from_millis(/*millis*/ 50)).await;
let config = match config_manager.load_latest_config(/*fallback_cwd*/ None).await {
Ok(config) => config,
Err(error) => {
warn!(%error, "failed to reload telemetry config after account change");
continue;
}
};
// ...配置文件可以在此之前发生变化,但读取最新配置的触发是 auth_changes.changed()。代码注释解释了 50 毫秒延迟的意图:账户处理器先发布认证变化,再安装新的 cloud loader。这个固定 sleep 给后一步留时间, 本身不是 cloud loader 已安装的确认消息,也不能用来证明任意调度延迟下的强同步关系。
5.2 安装与回收
读取配置成功后,task 构造新 provider,替换可重载 logger 层,再交接 provider 所有权。
源码文件:codex-rs/app-server/src/otel_reloader.rs
相关函数/类型:spawn(L76–L111,局部节选)
// 作者注:成功安装 logger 层后替换 provider;旧 provider 的 blocking shutdown 被分离,最后一个则在退出时等待。
// ...
let next_provider = match codex_core::otel_init::build_provider(
&config,
env!("CARGO_PKG_VERSION"),
Some(OTEL_SERVICE_NAME),
default_analytics_enabled,
) {
Ok(provider) => provider,
Err(error) => {
warn!(%error, "failed to rebuild telemetry exporters after account change");
continue;
}
};
if let Err(error) = logger_reload_handle.reload(
next_provider
.as_ref()
.and_then(OtelProvider::logger_export_layer)
.map(Layer::boxed),
) {
warn!(%error, "failed to install telemetry exporters after account change");
continue;
}
if let Some(previous_provider) = std::mem::replace(&mut provider, next_provider) {
drop(tokio::task::spawn_blocking(move || previous_provider.shutdown()));
}
info!(
event.name = "codex.app_server.otel_reloaded",
"reloaded telemetry exporters after account change"
);
}
}
}
if let Some(provider) = provider {
let _ = tokio::task::spawn_blocking(move || provider.shutdown()).await;
}
})
// ...失败应按阶段理解:配置读取失败不会进入构造;provider 构造失败不会进入 logger reload;reload 失败不会 执行本地 provider 字段的 mem::replace。这些顺序只说明本函数没有走到哪一步,不能推出构造 provider 期间涉及的所有全局设施都具备事务回滚。
成功替换后,旧 provider 的 shutdown() 放到 blocking pool,立即丢弃 handle,使热切换不等待它关闭。 退出循环时最后一个 provider 则通过 spawn_blocking(...).await 等待。因而 reloader 的 handle 完成, 只显式覆盖最后持有的 provider,不是一份“所有历史替换任务均已 join”的清单。
另外,取消只在循环 select 的选择点检查。当前 50 毫秒 sleep、配置加载或一次安装流程已开始时, 不会被外层 token 自动抢占。
5.3 切换断言
集成测试建立 /initial 与 /next 两组 mock collector 路径,在初始账户创建 Thread 后写入下一组 endpoint,再通过账户登录 RPC 触发认证变化。
源码文件:codex-rs/app-server/tests/suite/v2/otel.rs
相关函数/类型:account_switch_reloads_telemetry_collectors_and_preserves_trace_context(L79–L92、L112–L131,局部节选)
// 作者注:先等登录响应和重载日志,再产生 Thread 事件;最终检查新旧 collector 的身份和 trace 及新 metrics。
// ...
let response: LoginAccountResponse =
timeout(TEST_TIMEOUT, app_server.read_response(request_id)).await??;
assert_eq!(response, LoginAccountResponse::ChatgptAuthTokens {});
timeout(
TEST_TIMEOUT,
app_server.wait_for_json_log_event("codex.app_server.otel_reloaded"),
)
.await??;
app_server
.start_thread(ThreadStartParams::default())
.await?;
let status = timeout(TEST_TIMEOUT, app_server.shutdown_gracefully()).await??;
assert!(status.success(), "app-server did not shut down cleanly");
// ...
assert!(
initial_logs.contains(INITIAL_EMAIL),
"the initial account's logs were not exported: {initial_logs}"
);
assert!(
initial_traces.contains(PARENT_TRACE_ID),
"the initial account's trace context was not propagated: {initial_traces}"
);
assert!(
next_logs.contains(NEXT_EMAIL),
"the next account's logs did not reach its collector: {next_logs}"
);
assert!(
next_traces.contains(PARENT_TRACE_ID),
"the next account's trace context was not propagated: {next_traces}"
);
assert!(
next_metrics.contains("codex.thread.started"),
"the next account's metrics did not reach its collector: {next_metrics}"
);
// ...测试先等待 codex.app_server.otel_reloaded,再创建下一次 Thread,最后正常关闭 server。断言检查旧日志 包含旧身份、新日志包含新身份,两组 trace 都包含给定父 Trace ID,新 metrics 包含线程启动指标。 它验证了这个账户切换场景中 exporter 与父链接力,不能证明“仅修改 TOML 就会立即重载”,也不覆盖所有 provider 构造失败后的全局状态。
6. MCP 刷新接力
6.1 配置预检
MCP 是模型调用外部工具的协议。这里需分清三层:线程保存的配置、已发布的 runtime 快照、某个 MCP server 的连接就绪状态。App Server 刷新配置首先处理第一层。
源码文件:codex-rs/app-server/src/mcp_refresh.rs
相关函数/类型:reload_mcp_config(L9–L29,摘录)
// 作者注:先收集所有 thread/config,预检失败时尚未进入应用循环;这不是带回滚的跨线程事务。
pub(crate) async fn reload_mcp_config(
thread_manager: &Arc<ThreadManager>,
config_manager: &ConfigManager,
) -> io::Result<()> {
config_manager
.load_latest_config(/*fallback_cwd*/ None)
.await?;
let mut refreshes = Vec::new();
for thread_id in thread_manager.list_thread_ids().await {
let thread = thread_manager
.get_thread(thread_id)
.await
.map_err(|err| io::Error::other(format!("failed to load thread {thread_id}: {err}")))?;
let config = load_refresh_config(thread.as_ref(), config_manager).await?;
refreshes.push((thread, config));
}
for (thread, config) in refreshes {
thread.refresh_mcp_config(config).await;
}
Ok(())
}strict 路径先加载全局配置,再收集各 Thread 与其新配置,全部成功后才执行应用循环。预检中途失败时, 尚未调用任何 thread.refresh_mcp_config;这是“先计划后应用”的保证,不是应用阶段失败后回滚所有 Thread。
源码文件:codex-rs/app-server/src/mcp_refresh.rs
相关函数/类型:reload_mcp_config_best_effort / load_refresh_config(L31–L62,摘录)
// 作者注:逐线程加载、逐线程更新;失败记录后继续,配置由该线程现有上下文导出。
pub(crate) async fn reload_mcp_config_best_effort(
thread_manager: &Arc<ThreadManager>,
config_manager: &ConfigManager,
) {
for thread_id in thread_manager.list_thread_ids().await {
let thread = match thread_manager.get_thread(thread_id).await {
Ok(thread) => thread,
Err(err) => {
warn!(%thread_id, %err, "failed to load thread for MCP configuration refresh");
continue;
}
};
let config = match load_refresh_config(thread.as_ref(), config_manager).await {
Ok(config) => config,
Err(err) => {
warn!(%thread_id, %err, "failed to load thread MCP configuration");
continue;
}
};
thread.refresh_mcp_config(config).await;
}
}
async fn load_refresh_config(
thread: &CodexThread,
config_manager: &ConfigManager,
) -> io::Result<Config> {
let thread_config = thread.config().await;
config_manager
.load_latest_config_for_thread(thread_config.as_ref())
.await
}best effort 路径则每拿到一个健康线程就更新它,某线程消失或配置加载失败只记日志继续。 load_latest_config_for_thread 以该线程当前配置为上下文,不能把它替换成没有 cwd、覆盖项的单次全局加载。
6.2 唤醒与快照
CodexThread::refresh_mcp_config 在 core/src/codex_thread.rs 转交给 Session。真正的更新范围如下。
源码文件:codex-rs/core/src/session/mod.rs
相关函数/类型:Session::refresh_mcp_config(L1804–L1822,摘录)
// 作者注:锁内只更新 MCP 相关输入并标脏,释放锁后安排预热;返回不等待服务连接。
pub(crate) async fn refresh_mcp_config(&self, next_config: Config) {
let mut state = self.state.lock().await;
let mut config = (*state.session_configuration.original_config_do_not_use).clone();
config.config_layer_stack = next_config
.config_layer_stack
.with_user_layer_from(&config.config_layer_stack);
config.mcp_servers = next_config.mcp_servers;
config.mcp_oauth_credentials_store_mode = next_config.mcp_oauth_credentials_store_mode;
if let Err(err) = config.features.set_enabled(
Feature::SecretAuthStorage,
next_config.features.enabled(Feature::SecretAuthStorage),
) {
warn!("failed to refresh MCP auth storage config: {err}");
}
state.session_configuration.original_config_do_not_use = Arc::new(config);
self.mark_mcp_runtime_dirty();
drop(state);
self.schedule_mcp_prewarm();
}它克隆当前 Session 配置,只替换 MCP server、OAuth storage、相关配置层与 SecretAuthStorage 输入, 保留当前 user layer 语义;标记 runtime dirty 后释放状态锁,再安排预热。返回时不等待服务端握手,更不会 顺便把新配置中的模型等无关字段全部应用到已有会话。
每个 Session 的唤醒 channel 只有一个位置。
源码文件:codex-rs/core/src/session/session.rs
相关函数/类型:Session::new(L1471–L1487,局部节选)
// 作者注:每个 Session 只放一个待处理唤醒;配置本体另存在 Session state,不在 channel 传快照。
// ...
let (mcp_prewarm_tx, mcp_prewarm_rx) = async_channel::bounded(1);
let sess = Arc::new(Session {
thread_id,
installation_id,
tx_event: tx_event.clone(),
agent_status,
state: Mutex::new(state),
managed_network_proxy_refresh_lock: Semaphore::new(/*permits*/ 1),
features: config.features.clone(),
windows_sandbox_proxy_settings_mode,
multi_agent_version,
mcp_refresh: McpRefresh::new(),
mcp_elicitation_reviewer_handle: OnceLock::new(),
mcp_elicitation_lifecycle_handle: OnceLock::new(),
mcp_prewarm_tx,
mcp_prewarm_shutdown: CancellationToken::new(),
mcp_prewarm_task: std::sync::Mutex::new(None),
// ...channel 中传 (),不是完整配置快照。多次变化可以合并成“需要检查最新状态”这一件事;真实状态保存在 Session 中。如果改为给每次变化排入旧配置快照,就可能让慢 worker 按旧版本依次连接,延长新状态生效时间。
源码文件:codex-rs/core/src/session/mcp_prewarm.rs
相关函数/类型:start_mcp_prewarm_worker / schedule_mcp_prewarm(L14–L60,摘录)
// 作者注:空闲时仅持 Weak;预热内部 select 可直接取消 refresh future,pending 的恢复由 guard 保障。
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
},
};
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);
}
pub(super) fn schedule_mcp_prewarm(&self) {
let _ = self.mcp_prewarm_tx.try_send(());
}worker 空闲时只持 Weak,收到请求或 auth change 后再升级 Session。认证唤醒会额外标脏,普通唤醒则依赖 调用方已标脏。try_send(()) 在 channel 满时可以合并信号,因为后面的 refresh 会重新读取最新状态。
这一循环和模型刷新不同:它在 session.refresh_mcp_if_dirty() 外面又放了一层取消 select。停止预热时, 可以直接丢弃正在刷新的 future,因此被 claim 的 dirty 标记必须有补偿机制。
6.3 发布门闩
McpRefresh 同时保存 pending 位和单许可 semaphore。
源码文件:codex-rs/core/src/session/mcp_refresh.rs
相关函数/类型:McpRefresh / McpRefreshInvalidationGuard(L7–L55,摘录)
// 作者注:一个 semaphore 串行发布;claim 清 pending,发布前 future 被丢弃时 guard 把 pending 置回。
/// Owns MCP invalidation and the single gate used to publish runtime updates.
pub(super) struct McpRefresh {
pending: AtomicBool,
gate: Semaphore,
}
impl McpRefresh {
pub(super) fn new() -> Self {
Self {
pending: AtomicBool::new(false),
gate: Semaphore::new(/*permits*/ 1),
}
}
pub(super) fn invalidate(&self) {
self.pending.store(true, Ordering::Release);
}
pub(super) fn is_pending(&self) -> bool {
self.pending.load(Ordering::Acquire)
}
pub(super) fn claim(&self) -> bool {
self.pending.swap(false, Ordering::AcqRel)
}
#[tracing::instrument(name = "mcp.runtime.refresh_wait", skip_all)]
pub(super) async fn acquire(&self) -> Result<SemaphorePermit<'_>, AcquireError> {
self.gate.acquire().await
}
pub(super) fn close(&self) {
self.gate.close();
}
}
/// Restores a claimed refresh when its task is cancelled before publication.
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 {
self.refresh.invalidate();
}
}
}semaphore 使多个调用方不能并发发布;pending 位表示仍需刷新。claim() 把 true 换成 false,不能靠它 单独证明工作已完成,因为后续还有多个 await。McpRefreshInvalidationGuard 在 published=false 时 恢复 pending,弥补 future 被取消而不走正常返回的路径。
源码文件:codex-rs/core/src/session/mcp.rs
相关函数/类型:Session::refresh_mcp_if_dirty(L173–L198、L227–L239,局部节选)
// 作者注:门闩内重读最新认证和 desired state;省略的是能力与配置投影,成功发布后检查是否又被标脏。
// ...
/// Publishes changed MCP state, waiting for any refresh already in progress.
#[tracing::instrument(name = "mcp.runtime.refresh_if_dirty", skip_all)]
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
.mcp_runtime
.current_auth_matches(auth.as_ref())
{
self.mark_mcp_runtime_dirty();
}
if !self.mcp_refresh.claim() {
return;
}
let mut refresh_invalidation = McpRefreshInvalidationGuard {
refresh: &self.mcp_refresh,
published: false,
};
let auth = self.services.auth_manager.auth().await;
let desired = self.latest_mcp_desired_state(auth).await;
// ...
self.publish_mcp_runtime(
&desired,
mcp_projection,
&ready_selected_capability_roots,
Some(self.mcp_elicitation_reviewer()),
)
.await;
refresh_invalidation.published = true;
if !self.mcp_refresh.is_pending() {
return;
}
}
}
// ...持有门闩后先检查当前认证是否匹配 runtime,再 claim。接着重新加载认证和 desired state,经过能力及 配置投影,发布 runtime;成功后才将 guard 标为 published。如果发布过程中又被 invalidate,就在同一门闩 内继续一轮,而不是把旧配置已发布误认为所有变化已处理。
下面给出一次配置变化与工具调用竞争的时序。预热提高准备速度,实际调用仍会走刷新检查。
图中没有把“配置函数返回”连接为“server ready”:中间仍有预热、runtime 发布和所选服务的启动等待。 如果工具调用先取得门闩,它可以承担刷新;后台任务不是唯一的正确性承担者。
6.4 实际消费者
发布函数将 desired state 与选中的能力根投影成 runtime input,再调用 runtime 替换接口。
源码文件:codex-rs/core/src/session/mcp_runtime.rs
相关函数/类型:Session::publish_mcp_runtime(L279–L303,摘录)
// 作者注:投影成 runtime input 后调用 replace,再发布所选插件资料;这里的发布不等于所有 client ready。
pub(super) async fn publish_mcp_runtime(
&self,
desired: &McpDesiredState,
mcp_projection: McpRuntimeProjection,
ready_selected_capability_roots: &[SelectedCapabilityRoot],
elicitation_reviewer: Option<ElicitationReviewerHandle>,
) {
let mcp_projection = self
.project_selected_environment_mcp_servers(
&desired.session_source,
&desired.config,
&desired.environments,
mcp_projection,
)
.await;
let selected_plugins = mcp_projection.selected_plugins.clone();
let input = self.build_mcp_runtime_input(
desired,
mcp_projection,
ready_selected_capability_roots,
elicitation_reviewer,
);
self.services.mcp_runtime.replace(input).await;
self.services.thread_extension_data.insert(selected_plugins);
}McpRuntime::replace 位于 codex-rs/codex-mcp/src/runtime.rs。它协调连接集合并发布不可变快照; publish 最后 current.store(...),随后释放 publication gate。发布的配置可以被读取,但具体连接可能 仍在启动。要判断工具能否调用,还要看消费者。
源码文件:codex-rs/core/src/session/mcp_runtime.rs
相关函数/类型:Session::prepare_mcp_call(L61–L73,摘录)
// 作者注:真正调用前仍执行 refresh_if_dirty;预热丢唤醒或被取消不会替代这条正确性路径。
/// Captures this session's current MCP client and catalog for one tool call.
pub(crate) async fn prepare_mcp_call(
self: &Arc<Self>,
server: &str,
tool: &str,
) -> Option<PreparedMcpCall> {
self.refresh_mcp_if_dirty().await;
self.services
.mcp_runtime
.current_binding_for_call(server)
.await?
.prepare_call(server, tool)
}调用路径先自己等待 dirty refresh,再获取与目标 server 对应的 binding。这里保留的是一次调用所需的 具体 client 与工具目录关系,不能拿一份“最新工具名字列表”代替实际绑定。
源码文件:codex-rs/codex-mcp/src/runtime.rs
相关函数/类型:McpRuntime::current_binding_for_call(L416–L429,摘录)
// 作者注:调用绑定等待选中 server startup;读取 current_config 则明确不等待 clients。
/// Captures the current runtime after its selected server has finished startup.
pub async fn current_binding_for_call(&self, server: &str) -> Option<Arc<McpBinding>> {
let current = self.current.load_full();
current.config.as_ref()?;
if !current.connections.wait_for_server_startup(server).await {
return None;
}
Self::binding_from_published_runtime(current, /*required_servers*/ &[]).await
}
/// Returns the latest published configuration without waiting for clients.
pub fn current_config(&self) -> Option<Arc<McpConfig>> {
self.current.load().config.clone()
}current_binding_for_call 等待所选 server startup,失败返回 None;current_config 明确不等待 clients。 两者都是读取当前 runtime,却提供完全不同的可用性保证。“已读到新配置”不能用来判断一次连接已经成功。
6.5 取消恢复
取消测试不依赖真实 MCP 服务,而是持住 Session state 锁,让刷新停在 claim 之后、发布之前。
源码文件:codex-rs/core/src/session/tests.rs
相关函数/类型:cancelled_mcp_refresh_remains_pending(L8682–L8710,摘录)
// 作者注:持住 state 锁让刷新停在发布前;drop future 恢复 pending,下一次刷新才清掉它。
#[tokio::test]
async fn cancelled_mcp_refresh_remains_pending() {
let (session, _turn_context) = make_session_and_context().await;
let session = Arc::new(session);
{
let _state = session.state.lock().await;
{
let mut refresh = Box::pin(session.refresh_mcp_if_dirty());
let mut context = std::task::Context::from_waker(futures::task::noop_waker_ref());
assert!(std::future::Future::poll(refresh.as_mut(), &mut context).is_pending());
assert!(
!session.mcp_refresh.is_pending(),
"the refresh should have claimed its pending invalidation"
);
}
}
assert!(
session.mcp_refresh.is_pending(),
"a cancelled refresh must leave the runtime dirty"
);
session.refresh_mcp_if_dirty().await;
assert!(
!session.mcp_refresh.is_pending(),
"the next refresh should publish the pending runtime"
);
}第一次手动 poll 得到 Pending,并确认 dirty 已被 claim。离开作用域丢弃 future 后,pending 必须重新为 true; 释放锁再次刷新后才为 false。这个测试直接约束了 guard 的作用,不需要用“取消安全”一词替代字段推演。
App Server 的 best effort 测试关注另一层:局部配置加载失败是否阻止其他 Thread。
源码文件:codex-rs/app-server/src/mcp_refresh.rs
相关函数/类型:best_effort_refresh_updates_healthy_threads(L124–L147,摘录)
// 作者注:good/bad loader 各尝试一次,只让健康线程切换到 Secrets,坏线程保留 Direct。
#[tokio::test]
async fn best_effort_refresh_updates_healthy_threads() -> anyhow::Result<()> {
let (temp_dir, thread_manager, config_manager, loader) = refresh_test_state().await?;
std::fs::write(
temp_dir.path().join(codex_config::CONFIG_TOML_FILE),
"[features]\nsecret_auth_storage = true\n",
)?;
reload_mcp_config_best_effort(&thread_manager, &config_manager).await;
assert_eq!(loader.good_loads.load(Ordering::Relaxed), 1);
assert_eq!(loader.bad_loads.load(Ordering::Relaxed), 1);
for thread_id in thread_manager.list_thread_ids().await {
let thread = thread_manager.get_thread(thread_id).await?;
let config = thread.config().await;
let expected = if config.cwd.ends_with("good") {
AuthKeyringBackendKind::Secrets
} else {
AuthKeyringBackendKind::Direct
};
assert_eq!(config.auth_keyring_backend_kind(), expected);
}
Ok(())
}fixture 的 good/bad loader 各被调用一次;good 线程切换到 Secrets,bad 保留 Direct。 它没有等待所有 MCP endpoint 初始化成功,也没有断言整个 Config 被热替换。
7. 迁移后台尾声
7.1 响应与接纳
外部配置迁移由 externalAgentConfig/import 进入。前置代码先检查 Memory 导入开关与选择项,之后生成 独立 import_id。这个 ID 关联后续进度,不能与 RPC request ID 混用。
源码文件:codex-rs/app-server/src/external_agent_migration/processor.rs
相关函数/类型:ExternalAgentConfigRequestProcessor::import(L197–L230,局部节选)
// 作者注:RPC response 前已经执行同步导入和必要的配置变更处理;返回体只携带 import_id。
// ...
let import_id = Uuid::new_v4().to_string();
let analytics_source = params.source.clone().unwrap_or_default();
let provider_id = params.provider_id.clone();
let migration_service = self
.migration_service
.with_migration_source(params.migration_source.as_deref());
let needs_runtime_refresh = migration_items_need_runtime_refresh(¶ms.migration_items);
let has_migration_items = !params.migration_items.is_empty();
let has_plugin_imports = params.migration_items.iter().any(|item| {
matches!(
item.item_type,
ExternalAgentConfigMigrationItemType::Plugins
)
});
let (pending_session_imports, session_validation_result) =
self.validate_pending_session_imports(¶ms, &migration_service);
let import_outcome = self
.import_external_agent_config(params, &migration_service)
.await;
if needs_runtime_refresh {
self.config_processor.handle_config_mutation().await;
}
self.outgoing
.send_response(
request_id,
ExternalAgentConfigImportResponse {
import_id: import_id.clone(),
},
)
.await;
if !has_migration_items {
return Ok(());
}
// ...import_external_agent_config 在 RPC 响应前已经执行。Config、Skills、McpServerConfig、Hooks、Commands、 Plugins 等类型还会触发 handle_config_mutation(),之后才返回 import_id。因此这里既不是“响应前零工作”, 也不是“收到响应就全部完成”。空 migration_items 返回后直接结束,没有后面的完成通知路径。
源码文件:codex-rs/app-server/src/external_agent_migration/processor.rs
相关函数/类型:ExternalAgentConfigRequestProcessor::import(L232–L256,局部节选)
// 作者注:即使没有后台项,非空导入仍发 progress/completed;空 migrationItems 走前一段的提前返回。
// ...
let mut completed_item_results = Vec::new();
if let Some(session_validation_result) = session_validation_result {
send_import_progress(&self.outgoing, &import_id, &session_validation_result).await;
completed_item_results.push(session_validation_result);
}
for item_result in import_outcome.item_results {
send_import_progress(&self.outgoing, &import_id, &item_result).await;
completed_item_results.push(item_result);
}
let has_background_imports = !import_outcome.pending_plugin_imports.is_empty()
|| !pending_session_imports.is_empty();
if !has_background_imports {
send_completed_import_notification(
&self.outgoing,
self.state_db.as_ref(),
&self.analytics_events_client,
import_id,
analytics_source,
provider_id,
&completed_item_results,
)
.await;
return Ok(());
}
// ...非空导入先报告已知的校验/同步结果。如果没有 pending session/plugin,就立即组成完成通知并返回。 客户端不能根据“是否创建后台 task”决定是否等待 completed,而应根据本次请求类型与 import ID 追踪协议结果。
7.2 两条导入分支
有 pending 工作时,方法把后续需要的依赖移进一个独立 Tokio task。
源码文件:codex-rs/app-server/src/external_agent_migration/processor.rs
相关函数/类型:ExternalAgentConfigRequestProcessor::import(L258–L291,局部节选)
// 作者注:依赖被 clone/move 到独立 task;这里没有把 JoinHandle 放进 processor 的 TaskTracker。
// ...
let session_importer = self.session_importer.clone();
let outgoing = Arc::clone(&self.outgoing);
let state_db = self.state_db.clone();
let analytics_events_client = self.analytics_events_client.clone();
let thread_manager = Arc::clone(&self.thread_manager);
let session_metadata_mode = migration_service.session_metadata_mode();
let plugin_migration_service = migration_service;
let session_import_result = (!pending_session_imports.is_empty()).then(|| {
CoreImportItemResult::new(
CoreMigrationItemType::Sessions,
"Import sessions".to_string(),
/*cwd*/ None,
)
});
let pending_plugin_imports = import_outcome.pending_plugin_imports;
tokio::spawn(async move {
let connector_names_by_source_path =
detected_session_connectors(&plugin_migration_service, &pending_session_imports).0;
let session_progress_outgoing = Arc::clone(&outgoing);
let session_import_id = import_id.clone();
let session_imports = async move {
let session_import_result = session_import_result?;
let item_result = session_importer
.import_sessions(
pending_session_imports,
session_import_result,
session_metadata_mode,
connector_names_by_source_path,
)
.await;
send_import_progress(&session_progress_outgoing, &session_import_id, &item_result)
.await;
Some(item_result)
};
// ...该 task 持有 outgoing、state DB、analytics、ThreadManager、迁移 service 等引用。源码没有保存这次 spawn 返回的 handle,也没有给它连接取消 token。因此仅关闭原 RPC 的连接,不能由这段代码推导出迁移已取消; 持有的强引用也可能比请求函数活得更久。
会话导入是一个 future,插件导入是另一个;插件分支内部仍然逐项处理。
源码文件:codex-rs/app-server/src/external_agent_migration/processor.rs
相关函数/类型:ExternalAgentConfigRequestProcessor::import(L292–L331,局部节选)
// 作者注:插件分支内部逐项 await;单项错误进结果集合,然后继续下一项。
// ...
let plugin_progress_outgoing = Arc::clone(&outgoing);
let plugin_import_id = import_id.clone();
let plugin_imports = async move {
let mut item_results = Vec::new();
for pending_plugin_import in pending_plugin_imports {
let mut item_result = CoreImportItemResult::new(
CoreMigrationItemType::Plugins,
pending_plugin_import.description.clone(),
pending_plugin_import.cwd.clone(),
);
match plugin_migration_service
.import_plugins(
pending_plugin_import.cwd.as_deref(),
Some(pending_plugin_import.details),
)
.await
{
Ok(plugin_outcome) => {
apply_plugin_outcome_to_item_result(&mut item_result, plugin_outcome);
}
Err(error) => {
record_import_error(
&mut item_result,
"plugin_import",
/*sub_error_type*/ None,
error.to_string(),
/*source*/ None,
);
}
}
send_import_progress(
&plugin_progress_outgoing,
&plugin_import_id,
&item_result,
)
.await;
item_results.push(item_result);
}
item_results
};
// ...每项插件导入的成功或失败都会写入结果并发送 progress,失败不会直接 ? 退出整个分支。 这提供了项目粒度的结果,但不提供“任意项目失败就撤销此前写入”的事务语义。
源码文件:codex-rs/app-server/src/external_agent_migration/processor.rs
相关函数/类型:ExternalAgentConfigRequestProcessor::import(L332–L355,局部节选)
// 作者注:会话与插件两个 future 并发等待后汇总;清插件/技能缓存,然后发送完成通知。
// ...
let (session_result, plugin_results) = tokio::join!(session_imports, plugin_imports);
let mut background_item_results = Vec::new();
if let Some(session_result) = session_result {
background_item_results.push(session_result);
}
background_item_results.extend(plugin_results);
completed_item_results.extend(background_item_results);
if has_plugin_imports {
thread_manager.plugins_manager().clear_cache();
thread_manager.skills_service().clear_cache();
}
send_completed_import_notification(
&outgoing,
state_db.as_ref(),
&analytics_events_client,
import_id,
analytics_source,
provider_id,
&completed_item_results,
)
.await;
});
Ok(())
// ...tokio::join! 让两个 future 并发推进,并等待两者都返回;它没有为每个插件再生成 task,也不是遇到首个 错误就停止的 try_join!。汇合后合并结果,必要时清除插件与技能缓存,最后发送 completed。
下图画的是“存在 pending 项”的路径。同步结果、后台进度和最终结果沿同一个 import ID 汇合。
这里两个分支的通知先后受实际执行时间影响;最终 completed 在 join 之后。图中的客户端箭头都经 outgoing 队列,不能进一步证明断连客户端已经接收,或 server 进程被终止后该 task 仍有持久化执行保证。
7.3 完成记录
完成 helper 把聚合结果投向日志、Analytics、状态库与通知四个消费者。
源码文件:codex-rs/app-server/src/external_agent_migration/processor.rs
相关函数/类型:send_completed_import_notification(L547–L580,摘录)
// 作者注:先组成终态和分析事实,再尝试保存历史;DB 写失败记录 warning 后仍发通知。
async fn send_completed_import_notification(
outgoing: &OutgoingMessageSender,
state_db: Option<&StateDbHandle>,
analytics_events_client: &AnalyticsEventsClient,
import_id: String,
analytics_source: String,
provider_id: Option<String>,
item_results: &[CoreImportItemResult],
) {
let notification = completed_notification(import_id, item_results);
log_completed_import_failures(¬ification);
track_completed_import_notification(
analytics_events_client,
&analytics_source,
provider_id.as_deref().unwrap_or_default(),
¬ification,
);
if let Some(state_db) = state_db
&& let Err(err) =
record_completed_import_notification(state_db, provider_id.as_deref(), ¬ification)
.await
{
tracing::warn!(
import_id = %notification.import_id,
error = %err,
"failed to record external agent config import completion"
);
}
outgoing
.send_server_notification(ServerNotification::ExternalAgentConfigImportCompleted(
notification,
))
.await;
}state DB 缺席时仍可发送完成通知;写入出错时记录 warning 后继续发送。因此“看到了 completed”和 “历史库成功保存了本次结果”不是逻辑等价。反过来,成功写历史之后连接断开,也可能使客户端看不到通知。
真实失败测试在已有 config.toml 放入非法 TOML,再发起 Config 与 Commands 导入。
源码文件:codex-rs/app-server/tests/suite/v2/external_agent_config.rs
相关函数/类型:external_agent_config_import_reports_failed_sync_import_in_completion(L1225–L1229、L1268–L1296,局部节选)
// 作者注:fixture 写入非法既有 TOML;RPC 仍返回 import_id,完成结果中 Config 含失败类型与阶段。
// ...
std::fs::write(
source_home.join("settings.json"),
r#"{"env":{"FOO":"bar"}}"#,
)?;
std::fs::write(codex_home.path().join("config.toml"), "invalid = [")?;
// ...
let response: ExternalAgentConfigImportResponse =
timeout(DEFAULT_TIMEOUT, mcp.read_response(request_id)).await??;
let import_id = assert_import_response(response);
let completed: ExternalAgentConfigImportCompletedNotification = timeout(
DEFAULT_TIMEOUT,
mcp.read_notification("externalAgentConfig/import/completed"),
)
.await??;
assert_eq!(completed.import_id, import_id);
let config_result = completed
.item_type_results
.iter()
.find(|result| result.item_type == ExternalAgentConfigMigrationItemType::Config)
.expect("config result");
assert!(config_result.successes.is_empty());
assert_eq!(config_result.failures.len(), 1);
let config_failure = &config_result.failures[0];
assert_eq!(
config_failure.error_type.as_deref(),
Some("invalid_existing_config")
);
assert_eq!(config_failure.failure_stage, "import_request_failed");
assert!(
config_failure
.message
.contains("invalid existing config.toml"),
"unexpected failure: {config_failure:?}"
);
// ...RPC 响应仍给出 import ID,Config 项在 completed 中报告 invalid_existing_config 和 import_request_failed,successes 为空。客户端应检查每种 item type 的 successes/failures,不能把 HTTP/RPC 成功或 completed 这个事件名当作“所有导入成功”。
另一个测试专门走 pending plugin 路径。
源码文件:codex-rs/app-server/tests/suite/v2/external_agent_config.rs
相关函数/类型:external_agent_config_import_sends_completion_notification_after_pending_plugins_finish(L1760–L1777、L1806–L1816,局部节选)
// 作者注:无效 marketplace 让 pending 分支无需真实网络 clone;断言只确认 completed 的 import_id。
// ...
let codex_home = TempDir::new()?;
let source_home = external_agent_home(codex_home.path());
std::fs::create_dir_all(&source_home)?;
// This test only needs a pending non-local plugin import. Use an invalid
// source so the background completion path cannot make a real network clone.
std::fs::write(
source_home.join("settings.json"),
r#"{
"enabledPlugins": {
"formatter@acme-tools": true
},
"extraKnownMarketplaces": {
"acme-tools": {
"source": "not a valid marketplace source"
}
}
}"#,
)?;
// ...
let response: ExternalAgentConfigImportResponse =
timeout(DEFAULT_TIMEOUT, mcp.read_response(request_id)).await??;
let import_id = assert_import_response(response);
let completed: ExternalAgentConfigImportCompletedNotification = timeout(
DEFAULT_TIMEOUT,
mcp.read_notification("externalAgentConfig/import/completed"),
)
.await??;
assert_eq!(completed.import_id, import_id);
Ok(())
// ...它用非法 marketplace source 避免真实网络 clone,最终只断言 completed 与响应的 import ID 相同。 这个测试确认后台路径能结束并通知,不证明插件安装成功;插件是否可用应由有效来源场景中的结果和目录读取断言确认。
8. 退出现场
8.1 分阶段停止
回到最初的问题:关闭 App Server 时,不能把所有工作交给一个名字叫 drain_background_tasks 的函数想象完成。
源码文件:codex-rs/app-server/src/message_processor.rs
相关函数/类型:clear_runtime_references / drain_background_tasks(L589–L594、L768–L775,摘录)
// 作者注:两个清理入口覆盖的 owner 不完全相同;drain 中停止模型/费用后才等待线程后台任务。
pub(crate) fn clear_runtime_references(&self) {
self.account_processor.clear_external_auth();
self.apps_processor.shutdown();
self.models_refresh_worker.shutdown();
self.skills_watcher.shutdown();
}
// ...
pub(crate) async fn drain_background_tasks(&self) {
self.models_refresh_worker.shutdown();
if let Some(worker) = &self.turn_cost_worker {
worker.shutdown();
}
self.thread_processor.drain_background_tasks().await;
}clear_runtime_references 清外部认证引用、停止 Apps、模型和技能监听;drain_background_tasks 则停止 模型与费用,然后等待 Thread processor 的后台任务。两个入口职责有交集但不相等,调用方必须按真实退出路径阅读。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:drain_background_tasks / shutdown_threads(L1226–L1234、L1240–L1251,摘录)
// 作者注:TaskTracker close+wait 有 10 秒预算,超时只继续;随后 Core thread 的 bounded shutdown 是另一阶段。
pub(crate) async fn drain_background_tasks(&self) {
self.background_tasks.close();
if tokio::time::timeout(Duration::from_secs(10), self.background_tasks.wait())
.await
.is_err()
{
warn!("timed out waiting for background tasks to shut down; proceeding");
}
}
// ...
pub(crate) async fn shutdown_threads(&self) {
let report = self
.thread_manager
.shutdown_all_threads_bounded(Duration::from_secs(10))
.await;
for thread_id in report.submit_failed {
warn!("failed to submit Shutdown to thread {thread_id}");
}
for thread_id in report.timed_out {
warn!("timed out waiting for thread {thread_id} to shut down");
}
}这里 TaskTracker::close() 配合 wait(),只等待纳入该 tracker 的任务。10 秒 timeout 发生后记录日志并 继续;代码没有遍历 handle 再统一 abort。随后 shutdown_all_threads_bounded 有另一份 10 秒预算,分别报告 提交 Shutdown 失败与等待超时的 Thread。
Core MCP 预热的收尾提供了一个明确“发信号并等待”的对照。
源码文件:codex-rs/core/src/session/mcp_prewarm.rs
相关函数/类型:stop_mcp_prewarm_worker(L62–L74,摘录)
// 作者注:与仅 cancel 的 worker 不同,此处取走并 await 保存的 handle。
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)
.take();
if let Some(worker) = worker
&& let Err(error) = worker.await
{
warn!(%error, "MCP prewarm worker stopped unexpectedly");
}
}take handle 后释放同步 Mutex,再 await task,避免持锁等待。worker 的内层取消 select 会丢弃当前 refresh, 由前面的 guard 恢复未完成状态;接下来 Session 的关闭路径取得同一发布门闩。
源码文件:codex-rs/core/src/session/handlers.rs
相关函数/类型:shutdown_session_runtime(L412–L417,局部节选)
// 作者注:先停预热,再取得发布门闩并关闭它,最后关闭 runtime,避免与后台发布并发。
// ...
sess.stop_mcp_prewarm_worker().await;
{
let _refresh = sess.mcp_refresh.acquire().await;
sess.mcp_refresh.close();
sess.services.mcp_runtime.shutdown().await;
}
// ...先停预热,再关刷新 gate 和 runtime,避免后台仍在发布新连接时并发拆掉 runtime。 这是 Session 层的顺序,不能拿它替代外层 transport 或迁移 task 的退出分析。
源码文件:codex-rs/app-server/src/lib.rs
相关函数/类型:run_main_with_transport_options(L1180–L1192、L1202–L1211,局部节选)
// 作者注:graceful 与 forced 的处理器收尾不同;外层等待 router 后取消 transport token 并 join OTel reloader。
// ...
if !shutdown_state.forced() {
futures::future::join_all(connections.iter().map(
|(&connection_id, connection_state)| {
processor.connection_closed(connection_id, &connection_state.session)
},
))
.await;
connection_cleanup_tasks.drain().await;
processor.drain_background_tasks().await;
processor.shutdown_threads().await;
} else {
connection_cleanup_tasks.abort();
}
// ...
drop(transport_event_tx);
let _ = processor_handle.await;
let _ = outbound_handle.await;
transport_shutdown_token.cancel();
let _ = otel_reloader_handle.await;
for handle in transport_accept_handles {
let _ = handle.await;
}
// ...graceful 路径先对连接做清理并 drain connection cleanup,再处理后台任务和 Core Thread;forced 分支改为 abort 连接清理任务,跳过前面这段有序等待。外层之后等待 processor/router,取消 transport token,再等待 OTel reloader 与 acceptor。不同阶段的预算和等待对象不能简单合并成“整个 server 最多 10 秒退出”。
8.2 现象到状态
| 现象 | 先定位的状态或入口 | 可以得到的判断 |
|---|---|---|
| 模型 worker 已 drop,仍出现一次远端结果 | spawn_with_interval 的当前 await | 在途刷新没有被 token select 包围,完成一次不等于继续循环 |
| Turn 完成但没有费用 | spawn 条件、provider 匹配、认证、观察队列、availability、条目计数 | 只有排除这些条件后,才进入后端缺项与重试分析 |
| 技能修改没有通知 | 环境选择、noop fallback、registration、system-root 全批过滤 | 文件存在不等于本地 watcher 已注册该根 |
| 通知出现但界面未更新 | skills/list 与客户端状态 | 通知本身没有新技能正文,客户端仍需读取 |
| 修改 OTel 配置后无重载 | auth watch 与重载日志 | 单独文件修改不是此任务的唤醒源 |
| MCP 配置返回但调用还在等 | pending、刷新 gate、选中 server startup | 配置更新、发布和连接就绪是三个阶段 |
| 导入 RPC 成功但某项失败 | import ID 对应的 completed.item_type_results | 每项失败与传输响应分开表达 |
| graceful drain 超时仍有工作 | TaskTracker 的登记范围和各类 handle owner | timeout 不会自动 abort 未登记或仍运行的任务 |
可以从下面三组真实测试开始练习定位,在 Codex 仓库根目录执行:
just test --locked -p codex-app-server --lib -E 'test(models_refresh_worker::tests::) | test(turn_cost_worker::tests::) | test(mcp_refresh::tests::)'
just test --locked -p codex-file-watcher --lib -E 'test(throttled_receiver_) | test(watch_registration_drop_unregisters_paths)'
just test --locked -p codex-core --lib -E 'test(cancelled_mcp_refresh_remains_pending) | test(mcp_refresh_detects_shared_auth_manager_changes)'再做两个源码推演:沿一次技能变更复述“路径注册 → 合并 → 节流 → 缓存失效 → 通知 → 列表读取”,指出每个 对象的 owner;把 MCP 预热的 refresh future 设想为在任意 await 处被丢弃,说明 pending 何时应为 true, 以及实际工具调用为什么仍要检查刷新。能定位这两条链,就能把“后台还在工作”拆成具体可验证的状态问题。
