Skip to content

ThreadStatus计算

从运行事实聚合和等待请求所有权出发,追踪 ThreadStatus 的事件更新、watch 与通知、历史状态修正,以及卸载和优雅退出的不同判据。

基于rust-v0.150.0
CodexRustAppServerConcurrency

ThreadStatus计算 ​

界面显示 Active 时,模型未必正在生成:线程也可能只是等待一个审批结果。反过来,监听任务尚未更新状态时,读取请求已经能看见一个真正的进行中 Turn。如果客户端把状态看作“上一个事件名称的翻译”,这两个窗口都会解释不通。

本文沿 Core 事件进入 App Server 的路径,追踪运行事实如何聚合成状态,再看等待计数、异步释放和两种消费者。读者需要了解 Rust Mutex、Drop、异步任务与 Tokio watch;先读服务端请求与客户端响应理解审批回调,再读ThreadRead历史加载理解历史与当前运行状态的区别。本文不重复完整 Turn 执行流程,目标是让读者能判断某个状态标签来自哪个事实、何时生效,以及哪里可能存在观察延迟。

先区分三个层次:TurnStatus 描述单轮任务;Core 的 AgentStatus 描述执行者;App Server 的 ThreadStatus 面向整个对话的加载、运行和交互等待。Thread 还可能有持久归档标记,但它不属于下面这个枚举。一个 Thread 没有加载,不表示其持久历史已被删除。

1. 四种标签的判据 ​

1.1 线上表示 ​

公开类型只有四个主状态,等待审批和等待用户输入都放在 Active 的附加标记里。

源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs

相关函数/类型:ThreadStatus / ThreadActiveFlag;行号:1624-1645。

rust
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(tag = "type", rename_all = "camelCase")]
#[ts(tag = "type")]
#[ts(export_to = "v2/")]
pub enum ThreadStatus {
    NotLoaded,
    Idle,
    SystemError,
    #[serde(rename_all = "camelCase")]
    #[ts(rename_all = "camelCase")]
    Active {
        active_flags: Vec<ThreadActiveFlag>,
    },
}

#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub enum ThreadActiveFlag {
    WaitingOnApproval,
    WaitingOnUserInput,
}

JSON 使用 type 标签和 camelCase,例如 {"type":"active","activeFlags":["waitingOnApproval"]}。没有单独的 Running、Completed 或 Archived 变体。读完一轮后 Thread 可以保持 Idle,下一轮还能继续,因此不能把 Core 的终态判断原样套到这个公开枚举上。

1.2 运行事实 ​

聚合输入没有直接保存一个 ThreadStatus 枚举,而是保留独立事实。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:RuntimeFacts;行号:430-437。

rust
struct RuntimeFacts {
    is_loaded: bool,
    running: bool,
    pending_permission_requests: u32,
    pending_user_input_requests: u32,
    has_system_error: bool,
}

is_loaded 表示 watch 管理器记录的加载事实;running 表示它观察到的轮次运行;两个 u32 计数表示尚待释放的交互请求;has_system_error 保存错误事实。计数允许同类请求重叠,布尔值不足以表达“两个审批中只完成一个”。这些字段默认均为零或 false,但普通更新入口还会把对象标记为 loaded。

真正的优先级全部集中在一个纯计算函数中。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:loaded_thread_status;行号:438-460。

rust
fn loaded_thread_status(runtime: &RuntimeFacts) -> ThreadStatus {
    if !runtime.is_loaded {
        return ThreadStatus::NotLoaded;
    }

    let mut active_flags = Vec::new();
    if runtime.pending_permission_requests > 0 {
        active_flags.push(ThreadActiveFlag::WaitingOnApproval);
    }
    if runtime.pending_user_input_requests > 0 {
        active_flags.push(ThreadActiveFlag::WaitingOnUserInput);
    }

    if runtime.running || !active_flags.is_empty() {
        return ThreadStatus::Active { active_flags };
    }

    if runtime.has_system_error {
        return ThreadStatus::SystemError;
    }

    ThreadStatus::Idle
}

第一层先处理未加载。随后按固定顺序构造 approval、user-input 标记;只要 running 或任一计数非零,就返回 Active。只有没有运行和等待时,错误事实才显示成 SystemError;最后才是 Idle。因此 error 与 active 不是简单互斥字段,显示结果取决于判定顺序。

以下图对应实际 return 顺序。

例如 loaded=true、running=false、permission=1 时得到 Active,但不是一个正在采样的 Turn。loaded=false 时,即使旧错误事实仍在,也先返回 NotLoaded。后面的消费者正是利用这些区别,而不是统一判断 status != Idle。

2. 事件怎样改变事实 ​

2.1 状态所有者 ​

ThreadWatchManager 可被多个任务克隆,共享一个锁内状态和运行计数 watch。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:ThreadWatchManager / ThreadWatchActiveGuard;行号:18-30。

rust
#[derive(Clone)]
pub(crate) struct ThreadWatchManager {
    state: Arc<Mutex<ThreadWatchState>>,
    outgoing: Option<Arc<OutgoingMessageSender>>,
    running_turn_count_tx: watch::Sender<usize>,
}

pub(crate) struct ThreadWatchActiveGuard {
    manager: ThreadWatchManager,
    thread_id: String,
    guard_type: ThreadWatchActiveGuardType,
    handle: tokio::runtime::Handle,
}

guard 保存 manager、Thread ID、请求类型和 Tokio runtime handle;它没有携带模型工具本体,也不负责执行审批决策。outgoing 可为空,因而本地状态更新不以“存在外部通知发送器”为前提。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:ThreadWatchState;行号:310-314。

rust
#[derive(Default)]
struct ThreadWatchState {
    runtime_by_thread_id: HashMap<String, RuntimeFacts>,
    status_watcher_by_thread_id: HashMap<String, watch::Sender<ThreadStatus>>,
}

runtime_by_thread_id 与 status_watcher_by_thread_id 分开:前者保存事实,后者保存观察通道。删除一个运行事实不意味着必须销毁所有观察者;观察者还需要接收 NotLoaded。

图中的 map 都是当前服务器进程的状态,没有自动从数据库恢复旧等待计数的语义。持久日志可帮助重建历史,不能重建上个进程尚未完成的 oneshot 审批回调。

2.2 listener 顺序 ​

事件先从 CodexThread::next_event 进入 listener,再更新 ThreadState 的当前轮次视图,最后交给 bespoke 转换。

源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs

相关函数/类型:listener 先更新历史再翻译事件;行号:313-393。

rust
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(),
        );
    }

    // Track the event before emitting any typed translations
    // so thread-local state such as raw event opt-in stays
    // synchronized with the conversation.
    let (raw_events_enabled, realtime_effects) = {
        let mut thread_state = thread_state.lock().await;
        thread_state.track_current_turn_event(&event.id, &event.msg);
        let realtime_effects = if realtime_history_enabled
            && thread_state.realtime_history.should_observe(&event.msg)
        {
            let active_turn_id = thread_state.active_turn_snapshot().map(|turn| turn.id);
            thread_state
                .realtime_history
                .observe(&event.msg, active_turn_id.as_deref())
        } else {
            RealtimeEventEffects::default()
        };
        (thread_state.experimental_raw_events, realtime_effects)
    };
    if matches!(
        &event.msg,
        EventMsg::RawResponseItem(_) | EventMsg::RawResponseCompleted(_)
    ) && !raw_events_enabled
    {
        continue;
    }
    let subscribed_connection_ids = thread_state_manager
        .subscribed_connection_ids(conversation_id)
        .await;
    let thread_outgoing = ThreadScopedOutgoingMessageSender::new(
        outgoing_for_task.clone(),
        subscribed_connection_ids,
        conversation_id,
    );

    apply_realtime_event_effects(
        conversation.as_ref(),
        &thread_outgoing,
        conversation_id,
        realtime_effects,
    )
    .await;

    apply_bespoke_event_handling(
        event.clone(),
        conversation_id,
        conversation.clone(),
        thread_manager.clone(),
        thread_outgoing,
        thread_state.clone(),
        thread_watch_manager.clone(),
        thread_list_state_permit.clone(),
        fallback_model_provider.clone(),
    )
    .await;
    if matches!(event.msg, EventMsg::ShutdownComplete)
        && let Some(completion_tx) = thread_state
            .lock()
            .await
            .take_shutdown_drain_waiter()
    {
        let _ = completion_tx.send(());
    }
}

track_current_turn_event 在 typed notification 之前执行;随后才调用状态管理器对应的 note 方法。这个时序中间存在异步工作,所以历史快照与 watch 的观察时刻可以不同。next_event 失败会记录 warning 并退出循环,不能把“读不到下个事件”自行解释成某个终态通知已经发布。

2.3 开始与结束 ​

开始、正常完成先取消旧的服务端待回应请求,再改变事实。

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

相关函数/类型:TurnStarted / TurnComplete 分支;行号:159-205。

rust
EventMsg::TurnStarted(payload) => {
    // While not technically necessary as it was already done on TurnComplete, be extra cautios and abort any pending server requests.
    outgoing.abort_pending_server_requests().await;
    thread_watch_manager
        .note_turn_started(&conversation_id.to_string())
        .await;
    let turn = {
        let state = thread_state.lock().await;
        let mut turn = state.active_turn_snapshot().unwrap_or_else(|| Turn {
            id: payload.turn_id.clone(),
            items: Vec::new(),
            items_view: TurnItemsView::NotLoaded,
            error: None,
            status: TurnStatus::InProgress,
            started_at: payload.started_at,
            completed_at: None,
            duration_ms: None,
        });
        turn.items.clear();
        turn.items_view = TurnItemsView::NotLoaded;
        turn
    };
    let notification = TurnStartedNotification {
        thread_id: conversation_id.to_string(),
        turn,
    };
    outgoing
        .send_server_notification(ServerNotification::TurnStarted(notification))
        .await;
}
EventMsg::TurnComplete(turn_complete_event) => {
    // All per-thread requests are bound to a turn, so abort them.
    outgoing.abort_pending_server_requests().await;
    respond_to_pending_interrupts(&thread_state, &outgoing).await;
    let turn_failed = thread_state.lock().await.turn_summary.last_error.is_some();
    thread_watch_manager
        .note_turn_completed(&conversation_id.to_string(), turn_failed)
        .await;
    handle_turn_complete(
        conversation_id,
        event_turn_id,
        turn_complete_event,
        &outgoing,
        &thread_state,
    )
    .await;
}

TurnStarted 在构造并发送 turn/started 前更新 ThreadStatus;TurnComplete 在完成通知前清理状态。这里的 turn_failed 取自轮次摘要,但是否用于错误显示,要看被调用函数。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:note_turn_started;行号:147-154。

rust
pub(crate) async fn note_turn_started(&self, thread_id: &str) {
    self.update_runtime_for_thread(thread_id, |runtime| {
        runtime.is_loaded = true;
        runtime.running = true;
        runtime.has_system_error = false;
    })
    .await;
}

新 Turn 设置 loaded/running,并清除旧错误标记;它没有在这里直接清除两个等待计数。旧请求取消与其 guard 释放有独立时序,不能只看这一段赋值就推断回调已经全部消失。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:note_turn_completed / note_turn_interrupted;行号:156-162。

rust
pub(crate) async fn note_turn_completed(&self, thread_id: &str, _failed: bool) {
    self.clear_active_state(thread_id).await;
}

pub(crate) async fn note_turn_interrupted(&self, thread_id: &str) {
    self.clear_active_state(thread_id).await;
}

_failed 在实现中未使用。完成或中断都调用同一个清理函数。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:clear_active_state;行号:184-191。

rust
async fn clear_active_state(&self, thread_id: &str) {
    self.update_runtime_for_thread(thread_id, move |runtime| {
        runtime.running = false;
        runtime.pending_permission_requests = 0;
        runtime.pending_user_input_requests = 0;
    })
    .await;
}

这一步清空 running 和两种等待计数,但不修改 has_system_error。因此“Turn 完成且 failed=true 必然变成 SystemError”不是这段代码的语义;错误事实由独立事件入口设置。

2.4 错误与关闭 ​

以下是三个事件分支中与状态相关的摘录,省略号分开不同 match 分支。

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

相关函数/类型:Error / TurnAborted / ShutdownComplete 状态入口;行号:939-942,1117-1124,1240-1245。

rust
EventMsg::Error(ev) => {
    thread_watch_manager
        .note_system_error(&conversation_id.to_string())
        .await;
// ...
EventMsg::TurnAborted(turn_aborted_event) => {
    // All per-thread requests are bound to a turn, so abort them.
    outgoing.abort_pending_server_requests().await;
    respond_to_pending_interrupts(&thread_state, &outgoing).await;

    thread_watch_manager
        .note_turn_interrupted(&conversation_id.to_string())
        .await;
// ...
EventMsg::ShutdownComplete => {
    thread_watch_manager
        .note_thread_shutdown(&conversation_id.to_string())
        .await;
}

EventMsg::Error 在其后续特例处理之前就调用 note_system_error。是否再向用户发送某条错误通知,是另一层判断,不应该反过来决定这里是否已记录错误事实。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:note_system_error;行号:174-182。

rust
pub(crate) async fn note_system_error(&self, thread_id: &str) {
    self.update_runtime_for_thread(thread_id, |runtime| {
        runtime.running = false;
        runtime.pending_permission_requests = 0;
        runtime.pending_user_input_requests = 0;
        runtime.has_system_error = true;
    })
    .await;
}

系统错误清空运行和等待,同时置错。下一次 TurnStarted 清除它,而单纯 complete/interrupt 不清除。排障时要按事件顺序重放字段,而不是只统计收到几条 Error。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:note_thread_shutdown;行号:164-172。

rust
pub(crate) async fn note_thread_shutdown(&self, thread_id: &str) {
    self.update_runtime_for_thread(thread_id, |runtime| {
        runtime.running = false;
        runtime.pending_permission_requests = 0;
        runtime.pending_user_input_requests = 0;
        runtime.is_loaded = false;
    })
    .await;
}

ShutdownComplete 标记未加载并清空活动计数,但保留错误字段。因为投影先检查 is_loaded,外部显示仍是 NotLoaded。单纯 upsert 不清旧错误,真正新一轮开始才清错,这些入口承担不同职责。

3. 等待请求的所有权 ​

3.1 计数入场 ​

申请等待状态时先增加对应计数,再返回持有该等待的 guard。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:note_pending_request;行号:209-221。

rust
async fn note_pending_request(
    &self,
    thread_id: &str,
    guard_type: ThreadWatchActiveGuardType,
) -> ThreadWatchActiveGuard {
    self.update_runtime_for_thread(thread_id, move |runtime| {
        runtime.is_loaded = true;
        let counter = Self::pending_counter(runtime, guard_type);
        *counter = counter.saturating_add(1);
    })
    .await;
    ThreadWatchActiveGuard::new(self.clone(), thread_id.to_string(), guard_type)
}

saturating_add 在 u32 上限饱和,防止溢出变成零。该入口只设置 loaded,没有设置 running;于是一个等待请求可显示 Active,而 running count 仍是零。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:pending_counter;行号:283-291。

rust
fn pending_counter(
    runtime: &mut RuntimeFacts,
    guard_type: ThreadWatchActiveGuardType,
) -> &mut u32 {
    match guard_type {
        ThreadWatchActiveGuardType::Permission => &mut runtime.pending_permission_requests,
        ThreadWatchActiveGuardType::UserInput => &mut runtime.pending_user_input_requests,
    }
}

两种计数不能互相代偿。即使一个审批已经返回,只要仍在等待用户输入,waitingOnUserInput 就应保留。同一种类型的多个 guard 则共用其计数。

3.2 审批任务 ​

文件修改审批提供了完整的生产者与消费者链。

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

相关函数/类型:ApplyPatchApprovalRequest 的 guard 所有权;行号:576-604。

rust
EventMsg::ApplyPatchApprovalRequest(event) => {
    let permission_guard = thread_watch_manager
        .note_permission_requested(&conversation_id.to_string())
        .await;
    let item_id = event.call_id.clone();

    let params = FileChangeRequestApprovalParams {
        thread_id: conversation_id.to_string(),
        turn_id: event.turn_id.clone(),
        item_id: item_id.clone(),
        started_at_ms: event.started_at_ms,
        reason: event.reason.clone(),
        grant_root: event.grant_root.clone(),
    };
    let (pending_request_id, rx) = outgoing
        .send_request(ServerRequestPayload::FileChangeRequestApproval(params))
        .await;
    tokio::spawn(async move {
        on_file_change_request_approval_response(
            item_id,
            pending_request_id,
            rx,
            conversation,
            thread_state.clone(),
            permission_guard,
        )
        .await;
    });
}

guard 在发送 ServerRequest 之前创建,随后随 async move 进入等待回应的任务。这样回调任务的生命周期持有“仍在等待”这一事实,而不是依赖每个成功/失败分支手工记得减一次。

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

相关函数/类型:on_file_change_request_approval_response;行号:1945-1984。

rust
async fn on_file_change_request_approval_response(
    item_id: String,
    pending_request_id: RequestId,
    receiver: oneshot::Receiver<ClientRequestResult>,
    codex: Arc<CodexThread>,
    thread_state: Arc<Mutex<ThreadState>>,
    permission_guard: ThreadWatchActiveGuard,
) {
    let response = receiver.await;
    resolve_server_request_on_thread_listener(&thread_state, pending_request_id).await;
    drop(permission_guard);
    let decision = match response {
        Ok(Ok(value)) => match serde_json::from_value::<FileChangeRequestApprovalResponse>(value) {
            Ok(response) => map_file_change_approval_decision(response.decision),
            Err(err) => {
                error!("failed to deserialize FileChangeRequestApprovalResponse: {err}");
                ReviewDecision::denied("approval request failed")
            }
        },
        Ok(Err(err)) if is_turn_transition_server_request_error(&err) => return,
        Ok(Err(err)) => {
            error!("request failed with client error: {err:?}");
            ReviewDecision::denied("approval request failed")
        }
        Err(err) => {
            error!("request failed: {err:?}");
            ReviewDecision::denied("approval request failed")
        }
    };

    if let Err(err) = codex
        .submit(Op::PatchApproval {
            id: item_id,
            decision,
        })
        .await
    {
        error!("failed to submit PatchApproval: {err}");
    }
}

回应先完成 listener 内的请求结算,再 drop guard,然后解析决策并提交 Op::PatchApproval。客户端错误、通道关闭或错误 JSON 都走拒绝;轮次切换导致的请求取消可直接返回。guard 释放在这些分支前发生,状态等待结束不代表 patch 已获准或已经执行。

用户输入沿同样的所有权方式工作,但错误决策不同。

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

相关函数/类型:on_request_user_input_response;行号:1652-1731。

rust
async fn on_request_user_input_response(
    event_turn_id: String,
    pending_request_id: RequestId,
    receiver: oneshot::Receiver<ClientRequestResult>,
    conversation: Arc<CodexThread>,
    thread_state: Arc<Mutex<ThreadState>>,
    user_input_guard: ThreadWatchActiveGuard,
) {
    let response = receiver.await;
    resolve_server_request_on_thread_listener(&thread_state, pending_request_id).await;
    drop(user_input_guard);
    let value = match response {
        Ok(Ok(value)) => value,
        Ok(Err(err)) if is_turn_transition_server_request_error(&err) => return,
        Ok(Err(err)) => {
            error!("request failed with client error: {err:?}");
            let empty = CoreRequestUserInputResponse {
                answers: HashMap::new(),
            };
            if let Err(err) = conversation
                .submit(Op::UserInputAnswer {
                    id: event_turn_id,
                    response: empty,
                })
                .await
            {
                error!("failed to submit UserInputAnswer: {err}");
            }
            return;
        }
        Err(err) => {
            error!("request failed: {err:?}");
            let empty = CoreRequestUserInputResponse {
                answers: HashMap::new(),
            };
            if let Err(err) = conversation
                .submit(Op::UserInputAnswer {
                    id: event_turn_id,
                    response: empty,
                })
                .await
            {
                error!("failed to submit UserInputAnswer: {err}");
            }
            return;
        }
    };

    let response =
        serde_json::from_value::<ToolRequestUserInputResponse>(value).unwrap_or_else(|err| {
            error!("failed to deserialize ToolRequestUserInputResponse: {err}");
            ToolRequestUserInputResponse {
                answers: HashMap::new(),
            }
        });
    let response = CoreRequestUserInputResponse {
        answers: response
            .answers
            .into_iter()
            .map(|(id, answer)| {
                (
                    id,
                    CoreRequestUserInputAnswer {
                        answers: answer.answers,
                    },
                )
            })
            .collect(),
    };

    if let Err(err) = conversation
        .submit(Op::UserInputAnswer {
            id: event_turn_id,
            response,
        })
        .await
    {
        error!("failed to submit UserInputAnswer: {err}");
    }
}

无法取得或解析回答时提交空 answers;特定轮次切换错误直接返回。审批失败对应拒绝、用户输入失败对应空答案,二者都应释放等待,却不能合并成同一种业务结果。

3.3 异步释放 ​

Drop 不能在同步析构里 await,所以实现保存 runtime handle 并生成清理任务。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:drop;行号:48-57。

rust
fn drop(&mut self) {
    let manager = self.manager.clone();
    let thread_id = self.thread_id.clone();
    let guard_type = self.guard_type;
    self.handle.spawn(async move {
        manager
            .note_active_guard_released(thread_id, guard_type)
            .await;
    });
}

drop(guard) 返回时,只保证已安排异步更新,不保证 runtime map 或客户端状态已改变。测试若立即 assert Idle,会把调度延迟误当成业务错误。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:note_active_guard_released;行号:262-272。

rust
async fn note_active_guard_released(
    &self,
    thread_id: String,
    guard_type: ThreadWatchActiveGuardType,
) {
    self.update_runtime_for_thread(&thread_id, move |runtime| {
        let counter = Self::pending_counter(runtime, guard_type);
        *counter = counter.saturating_sub(1);
    })
    .await;
}

释放用 saturating_sub 防止完成阶段已经归零后再次减少时下溢。它不携带 Turn ID 或 generation;不能据此宣称跨轮次旧 guard 的释放一定与新请求完全隔离。调查切换轮次时的等待标记异常,需要同时追踪旧回调的取消完成和新请求的注册时刻。

按调用顺序,释放与审批执行是两条随后继续的工作。

图中的异步箭头不是“Core 已批准”的因果证明。客户端看到等待标记消失时,工具仍可能继续执行,running 仍可为 true。

上游测试先启动 Turn,然后同时持有审批与用户输入 guard,分两次释放。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:status_updates_track_single_thread;行号:516-561。

rust
let permission_guard = manager
    .note_permission_requested(INTERACTIVE_THREAD_ID)
    .await;
assert_eq!(
    manager
        .loaded_status_for_thread(INTERACTIVE_THREAD_ID)
        .await,
    ThreadStatus::Active {
        active_flags: vec![ThreadActiveFlag::WaitingOnApproval],
    },
);

let user_input_guard = manager
    .note_user_input_requested(INTERACTIVE_THREAD_ID)
    .await;
assert_eq!(
    manager
        .loaded_status_for_thread(INTERACTIVE_THREAD_ID)
        .await,
    ThreadStatus::Active {
        active_flags: vec![
            ThreadActiveFlag::WaitingOnApproval,
            ThreadActiveFlag::WaitingOnUserInput,
        ],
    },
);

drop(permission_guard);
wait_for_status(
    &manager,
    INTERACTIVE_THREAD_ID,
    ThreadStatus::Active {
        active_flags: vec![ThreadActiveFlag::WaitingOnUserInput],
    },
)
.await;

drop(user_input_guard);
wait_for_status(
    &manager,
    INTERACTIVE_THREAD_ID,
    ThreadStatus::Active {
        active_flags: vec![],
    },
)
.await;

第一次释放后只剩 WaitingOnUserInput,第二次释放后仍是没有 flags 的 Active,因为 Turn 还未结束。测试使用 wait_for_status 等待异步释放,之后才 complete 成为 Idle。这个序列直接验证“交互等待”和“轮次运行”可以独立改变。

4. 状态发布 ​

4.1 一次变更 ​

普通运行事实更新会先确保 map 中有条目,并默认将其标记为加载。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:ThreadWatchState::update_runtime;行号:350-367。

rust
fn update_runtime<F>(
    &mut self,
    thread_id: &str,
    mutate: F,
) -> Option<ThreadStatusChangedNotification>
where
    F: FnOnce(&mut RuntimeFacts),
{
    let previous_status = self.status_for(thread_id);
    let runtime = self
        .runtime_by_thread_id
        .entry(thread_id.to_string())
        .or_default();
    runtime.is_loaded = true;
    mutate(runtime);
    self.update_status_watcher_for_thread(thread_id);
    self.status_changed_notification(thread_id.to_string(), previous_status)
}

闭包可以覆盖 loaded,例如 shutdown 将其设回 false。随后先更新本地 watcher,再计算是否需要对外状态通知。某个调用只改变计数却不改变聚合结果时,不需要再通知客户端同一个状态。

外层统一维护 running count,并在释放锁后发送出站通知。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:mutate_and_publish;行号:223-253。

rust
async fn mutate_and_publish<F>(&self, mutate: F)
where
    F: FnOnce(&mut ThreadWatchState) -> Option<ThreadStatusChangedNotification>,
{
    let notification = {
        let mut state = self.state.lock().await;
        let notification = mutate(&mut state);
        let running_turn_count = state
            .runtime_by_thread_id
            .values()
            .filter(|runtime| runtime.running)
            .count();
        self.running_turn_count_tx.send_if_modified(|current| {
            if *current == running_turn_count {
                false
            } else {
                *current = running_turn_count;
                true
            }
        });
        notification
    };

    if let Some(notification) = notification
        && let Some(outgoing) = &self.outgoing
    {
        outgoing
            .send_server_notification(ServerNotification::ThreadStatusChanged(notification))
            .await;
    }
}

计数是对整个 map 的 running 字段求和,单次变更包含一次与已跟踪条目数相关的扫描;不能从函数名误判为常数时间增减。watch 的 send_if_modified 在状态锁内完成,外部发送在锁外 await,避免慢出站队列阻塞状态读取。

这也限定了顺序保证:一个变更中的本地快照先更新,出站随后排队;跨多个并发调用,不能仅凭这把锁声称网络通知已成为全局持久事件日志。

4.2 合并与去重 ​

订阅可以发生在 Thread 尚未跟踪时。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:ThreadWatchState::subscribe;行号:380-387。

rust
fn subscribe(&mut self, thread_id: String) -> watch::Receiver<ThreadStatus> {
    let status = self.loaded_status_for_thread(&thread_id);
    let sender = self
        .status_watcher_by_thread_id
        .entry(thread_id)
        .or_insert_with(|| watch::channel(status.clone()).0);
    sender.subscribe()
}

初始值来自当前状态,未知 Thread 得到 NotLoaded。watch 保存的是最新值,慢消费者不应假定每个瞬时状态都能逐个收到。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:update_status_watcher;行号:394-412。

rust
fn update_status_watcher(&mut self, thread_id: &str, status: &ThreadStatus) {
    let remove_watcher = if let Some(sender) = self.status_watcher_by_thread_id.get(thread_id) {
        let status = status.clone();
        let _ = sender.send_if_modified(|current| {
            if *current == status {
                false
            } else {
                *current = status;
                true
            }
        });
        sender.receiver_count() == 0
    } else {
        false
    };
    if remove_watcher {
        self.status_watcher_by_thread_id.remove(thread_id);
    }
}

状态相等就不修改版本;更新时发现 receiver_count 为零,移除 sender。这是更新时清理,不能说最后一个 receiver drop 后 map 会同步消失。runtime facts 与 watcher map 的存活条件也不同。

对外通知独立比较前后投影。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:status_changed_notification;行号:414-426。

rust
fn status_changed_notification(
    &self,
    thread_id: String,
    previous_status: Option<ThreadStatus>,
) -> Option<ThreadStatusChangedNotification> {
    let status = self.status_for(&thread_id)?;

    if previous_status.as_ref() == Some(&status) {
        return None;
    }

    Some(ThreadStatusChangedNotification { thread_id, status })
}

第二个同类审批把 permission 从 1 改成 2,外部 flags 仍只包含一项,不会因此重复发 Active 通知。客户端不能从通知次数反推出同时存在多少请求。

4.3 静默装配与移除 ​

新建或恢复对象时可以先静默安装运行事实。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:ThreadWatchState::upsert_thread;行号:317-334。

rust
fn upsert_thread(
    &mut self,
    thread_id: String,
    emit_notification: bool,
) -> Option<ThreadStatusChangedNotification> {
    let previous_status = self.status_for(&thread_id);
    let runtime = self
        .runtime_by_thread_id
        .entry(thread_id.clone())
        .or_default();
    runtime.is_loaded = true;
    self.update_status_watcher_for_thread(&thread_id);
    if emit_notification {
        self.status_changed_notification(thread_id, previous_status)
    } else {
        None
    }
}

emit_notification=false 仅省略外部状态通知,本地 watcher 仍更新。这样新对象可先由 thread/started 介绍,避免客户端先收到陌生 ID 的 status/changed。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:silent_upsert_skips_initial_notification;行号:764-796。

rust
async fn silent_upsert_skips_initial_notification() {
    let (outgoing_tx, mut outgoing_rx) = mpsc::channel(8);
    let manager = ThreadWatchManager::new_with_outgoing(Arc::new(OutgoingMessageSender::new(
        outgoing_tx,
        codex_analytics::AnalyticsEventsClient::disabled(),
    )));

    manager.upsert_thread_silently(INTERACTIVE_THREAD_ID).await;

    assert_eq!(
        manager
            .loaded_status_for_thread(INTERACTIVE_THREAD_ID)
            .await,
        ThreadStatus::Idle,
    );
    assert!(
        timeout(Duration::from_millis(100), outgoing_rx.recv())
            .await
            .is_err(),
        "silent upsert should not emit thread/status/changed"
    );

    manager.note_turn_started(INTERACTIVE_THREAD_ID).await;
    assert_eq!(
        recv_status_changed_notification(&mut outgoing_rx).await,
        ThreadStatusChangedNotification {
            thread_id: INTERACTIVE_THREAD_ID.to_string(),
            status: ThreadStatus::Active {
                active_flags: vec![],
            },
        },
    );
}

测试要求静默 upsert 后可读到 Idle,但短时间没有出站通知;随后 TurnStarted 必须仍发 Active。这否定了“静默注册永久禁止通知”的理解。

移除则删除 facts 并通知 NotLoaded。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:ThreadWatchState::remove_thread;行号:336-348。

rust
fn remove_thread(&mut self, thread_id: &str) -> Option<ThreadStatusChangedNotification> {
    let previous_status = self.status_for(thread_id);
    self.runtime_by_thread_id.remove(thread_id);
    self.update_status_watcher(thread_id, &ThreadStatus::NotLoaded);
    if previous_status.is_some() && previous_status != Some(ThreadStatus::NotLoaded) {
        Some(ThreadStatusChangedNotification {
            thread_id: thread_id.to_string(),
            status: ThreadStatus::NotLoaded,
        })
    } else {
        None
    }
}

未知或此前已是 NotLoaded 时不重复发外部通知。本地 watcher 仍收到必要的新值;这与简单删除 HashMap 条目后让消费者一直停在旧 Active 不同。

4.4 连接过滤 ​

即使管理器已发布状态,某条连接也可以选择不收这个方法。

源码文件:codex-rs/app-server/src/transport.rs

相关函数/类型:should_skip_notification_for_connection;行号:100-123。

rust
fn should_skip_notification_for_connection(
    connection_state: &OutboundConnectionState,
    message: &OutgoingMessage,
) -> bool {
    let Ok(opted_out_notification_methods) = connection_state.opted_out_notification_methods.read()
    else {
        warn!("failed to read outbound opted-out notifications");
        return false;
    };
    match message {
        OutgoingMessage::AppServerNotification(envelope) => {
            if envelope.notification.experimental_reason().is_some()
                && !connection_state
                    .experimental_api_enabled
                    .load(Ordering::Acquire)
            {
                return true;
            }
            let method = envelope.notification.to_string();
            opted_out_notification_methods.contains(method.as_str())
        }
        _ => false,
    }
}

过滤在出站连接层按方法名执行,不会回滚内部 facts、watch 或 running count。读取过滤集合失败时返回 false,因此这个错误分支也不能被描述成一律静默丢弃通知。

集成测试在 initialize 里订阅排除名单,随后执行一轮普通 Turn。

源码文件:codex-rs/app-server/tests/suite/v2/thread_status.rs

相关函数/类型:thread_status_changed_can_be_opted_out;行号:150-156,187-210。

rust
Some(InitializeCapabilities {
    experimental_api: true,
    request_attestation: false,
    opt_out_notification_methods: Some(vec!["thread/status/changed".to_string()]),
    mcp_server_openai_form_elicitation: false,
    extensions: None,
}),
// ...
timeout(
    DEFAULT_READ_TIMEOUT,
    mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;

let status_update = timeout(
    std::time::Duration::from_millis(500),
    mcp.read_stream_until_notification_message("thread/status/changed"),
)
.await;
match status_update {
    Err(_) => {}
    Ok(Ok(notification)) => {
        anyhow::bail!(
            "thread/status/changed should be filtered by optOutNotificationMethods; got: {notification:?}"
        );
    }
    Ok(Err(err)) => {
        anyhow::bail!(
            "expected timeout waiting for filtered thread/status/changed, got: {err}"
        );
    }
}

它先确认 turn/completed 到达,再要求等待 status/changed 超时。这个输入与断言组合证明被屏蔽的是某个通知方法,而非整个执行流程;客户端没有收到状态变化时,不能立即判定 Core 卡住。

5. 读时的状态修正 ​

5.1 进行中窗口 ​

历史视图与 watch 有观察时差,读取侧用一个有限的提升规则补偿。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:resolve_thread_status;行号:294-308。

rust
pub(crate) fn resolve_thread_status(
    status: ThreadStatus,
    has_in_progress_turn: bool,
) -> ThreadStatus {
    // Running-turn events can arrive before the watch runtime state is observed by
    // the listener loop. In that window we prefer to reflect a real active turn as
    // `Active` instead of `Idle`/`NotLoaded`.
    if has_in_progress_turn && matches!(status, ThreadStatus::Idle | ThreadStatus::NotLoaded) {
        return ThreadStatus::Active {
            active_flags: Vec::new(),
        };
    }

    status
}

只将 Idle 或 NotLoaded 在 has_in_progress_turn=true 时提升为无 flags 的 Active,SystemError 不会在这里被覆盖,已有 Active flags 也不被清空。该 helper 返回一个用于响应的值,没有回写运行事实或增加 running count。

热恢复路径的输入来自真实执行者与 active Turn 快照。

源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs

相关函数/类型:handle_pending_thread_resume_request 的状态来源;行号:642-646,665-674。

rust
let has_live_in_progress_turn =
    matches!(conversation.agent_status().await, AgentStatus::Running)
        || active_turn
            .as_ref()
            .is_some_and(|turn| matches!(turn.status, TurnStatus::InProgress));
// ...

let thread_status = thread_watch_manager
    .loaded_status_for_thread(&thread.id)
    .await;

set_thread_status_and_interrupt_stale_turns(
    &mut thread,
    thread_status.clone(),
    has_live_in_progress_turn,
);

AgentStatus::Running 或 active snapshot 的 InProgress 支撑 live 判断,不是只看到磁盘日志中残留一个 TurnStarted 就假定进程仍在工作。

源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs

相关函数/类型:set_thread_status_and_interrupt_stale_turns;行号:953-967。

rust
pub(super) fn set_thread_status_and_interrupt_stale_turns(
    thread: &mut Thread,
    loaded_status: ThreadStatus,
    has_live_in_progress_turn: bool,
) {
    let status = resolve_thread_status(loaded_status, has_live_in_progress_turn);
    if !matches!(status, ThreadStatus::Active { .. }) {
        for turn in &mut thread.turns {
            if matches!(turn.status, TurnStatus::InProgress) {
                turn.status = TurnStatus::Interrupted;
            }
        }
    }
    thread.status = status;
}

如果最终 ThreadStatus 不是 Active,旧历史里仍是 InProgress 的 Turn 被投影为 Interrupted。Thread 级 SystemError 和 Turn 级 Interrupted 可以同时出现,两者回答的是不同问题。这个响应处理不应写成“修改源 rollout 以中断旧任务”。

下图区分响应修正与运行事实。

图中没有从响应回写 RuntimeFacts 的箭头。读取响应暂时显示 Active,不足以证明内部运行计数已更新;反过来,历史内容被裁剪也不直接改变 watch。

5.2 子 Agent 补齐 ​

批量状态查询保留“未跟踪”和“明确 NotLoaded”的区别。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:ThreadWatchManager::loaded_statuses_for_threads;行号:117-130。

rust
pub(crate) async fn loaded_statuses_for_threads(
    &self,
    thread_ids: impl IntoIterator<Item = String>,
) -> HashMap<String, ThreadStatus> {
    let state = self.state.lock().await;
    thread_ids
        .into_iter()
        .filter_map(|thread_id| {
            state
                .status_for(&thread_id)
                .map(|status| (thread_id, status))
        })
        .collect()
}

缺少 facts 的 ID 不插入 map;已经记录 shutdown 的 ID 则可以以 NotLoaded 值存在。这个区别会影响子 Agent 的回退。

源码文件:codex-rs/app-server/src/request_processors/thread_enrichment.rs

相关函数/类型:enrich_loaded_threads;行号:12-79。

rust
pub(super) async fn enrich_loaded_threads<T>(
    thread_manager: &ThreadManager,
    thread_watch_manager: &ThreadWatchManager,
    threads: &mut [T],
    mut as_thread: impl FnMut(&mut T) -> &mut Thread,
) {
    let statuses = thread_watch_manager
        .loaded_statuses_for_threads(
            threads
                .iter_mut()
                .map(&mut as_thread)
                .map(|thread| thread.id.clone()),
        )
        .await;

    futures::future::join_all(threads.iter_mut().map(as_thread).map(|thread| {
        let statuses = &statuses;
        async move {
            let watched_status = statuses.get(&thread.id);
            if let Some(status) = watched_status {
                thread.status = status.clone();
            }

            if !matches!(
                &thread.source,
                SessionSource::SubAgent(SubAgentSource::ThreadSpawn { .. })
            ) || matches!(watched_status, Some(ThreadStatus::NotLoaded))
            {
                return;
            }

            let Ok(thread_id) = ThreadId::from_string(&thread.id) else {
                return;
            };
            let Ok(loaded_thread) = thread_manager.get_thread(thread_id).await else {
                return;
            };
            match loaded_thread.agent_status().await {
                AgentStatus::Running => {
                    if watched_status.is_none() {
                        thread.status = resolve_thread_status(
                            ThreadStatus::Idle,
                            /*has_in_progress_turn*/ true,
                        );
                    }
                }
                AgentStatus::PendingInit | AgentStatus::Interrupted | AgentStatus::Completed(_) => {
                    if watched_status.is_none() {
                        thread.status = ThreadStatus::Idle;
                    }
                }
                AgentStatus::Errored(_) => {
                    thread.status = ThreadStatus::SystemError;
                }
                AgentStatus::Shutdown | AgentStatus::NotFound => {
                    thread.status = ThreadStatus::NotLoaded;
                    return;
                }
            }
            let config_snapshot = loaded_thread.config_snapshot().await;
            thread.can_accept_direct_input = Some(can_accept_direct_input(
                loaded_thread.multi_agent_version(),
                &config_snapshot.session_source,
            ));
        }
    }))
    .await;
}

普通线程或明确 watched NotLoaded 不走子 Agent 补齐。对于 ThreadSpawn 来源,Core Running 在缺 watcher 时补 Active,PendingInit/Interrupted/Completed 补 Idle;Errored 可得到 SystemError,Shutdown/NotFound 得到 NotLoaded。不能用一个 AgentStatus 到 ThreadStatus 的无条件转换表替代这段带来源和 watcher 条件的实现。

6. 两种运行消费者 ​

6.1 延迟卸载 ​

无人订阅的线程是否可以卸载,首先看有没有外部连接与是否 Active。

源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs

相关函数/类型:unloading_target;行号:63-70。

rust
fn unloading_target(&self) -> Option<Instant> {
    match (self.has_subscribers, self.is_active) {
        ((false, has_no_subscribers_since), (false, is_inactive_since)) => {
            Some(std::cmp::max(has_no_subscribers_since, is_inactive_since) + self.delay)
        }
        _ => None,
    }
}

两项都为 false 时,deadline 取“最后一次变为无订阅”和“最后一次变为不活动”中较晚的时刻,加卸载延迟。这不是仅从 TurnComplete 开始一个固定倒计时。

源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs

相关函数/类型:sync_receiver_values;行号:72-82。

rust
fn sync_receiver_values(&mut self) {
    let has_subscribers = *self.has_subscribers_rx.borrow();
    if self.has_subscribers.0 != has_subscribers {
        self.has_subscribers = (has_subscribers, Instant::now());
    }

    let is_active = matches!(*self.thread_status_rx.borrow(), ThreadStatus::Active { .. });
    if self.is_active.0 != is_active {
        self.is_active = (is_active, Instant::now());
    }
}

只在布尔状态发生改变时刷新时间。Active 中 approval/user-input 标记互换,仍然是 active=true,不会被当成一次退出活动后重新进入。

listener 真正触发卸载前还核对 Core 状态。

源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs

相关函数/类型:listener 卸载前核对 AgentStatus;行号:398-414。

rust
if !unloading_state.should_unload_now() {
    continue;
}
if matches!(conversation.agent_status().await, AgentStatus::Running) {
    unloading_state.note_thread_activity_observed();
    continue;
}
{
    let mut pending_thread_unloads = pending_thread_unloads.lock().await;
    if pending_thread_unloads.contains(&conversation_id) {
        continue;
    }
    if !unloading_state.should_unload_now() {
        continue;
    }
    pending_thread_unloads.insert(conversation_id);
}

若 Core 仍 Running,记下新观察时刻并继续等待,补偿 watch 与 Core 的竞态窗口。pending unload 集合防止重复安排卸载。具体 shutdown 与资源回收可接着阅读Thread归档与删除,但普通空闲卸载与归档的错误处理不能互相套用。

6.2 运行计数 ​

优雅退出关注真正运行的轮次,因此消费 running count,而非 Active Thread 总数。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:has_running_turns_tracks_runtime_running_flag_only;行号:685-703。

rust
async fn has_running_turns_tracks_runtime_running_flag_only() {
    let manager = ThreadWatchManager::new();
    manager.upsert_thread(INTERACTIVE_THREAD_ID).await;

    assert_eq!(manager.running_turn_count().await, 0);

    let _permission_guard = manager
        .note_permission_requested(INTERACTIVE_THREAD_ID)
        .await;
    assert_eq!(manager.running_turn_count().await, 0);

    manager.note_turn_started(INTERACTIVE_THREAD_ID).await;
    assert_eq!(manager.running_turn_count().await, 1);

    manager
        .note_turn_completed(INTERACTIVE_THREAD_ID, false)
        .await;
    assert_eq!(manager.running_turn_count().await, 0);
}

单独创建 permission guard 后 count 仍为 0;TurnStarted 后为 1;完成后回到 0。测试中的 Thread 在第一阶段已经是等待审批的 Active,这正是两个消费者采取不同判据的理由。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:running_turn_watch_notifies_only_when_count_changes;行号:706-723。

rust
async fn running_turn_watch_notifies_only_when_count_changes() {
    let manager = ThreadWatchManager::new();
    let mut count = manager.subscribe_running_turn_count();

    manager.upsert_thread(INTERACTIVE_THREAD_ID).await;
    manager.note_turn_started(INTERACTIVE_THREAD_ID).await;
    assert!(count.has_changed().expect("watch remains open"));
    assert_eq!(*count.borrow_and_update(), 1);

    let _permission_guard = manager
        .note_permission_requested(INTERACTIVE_THREAD_ID)
        .await;
    assert!(!count.has_changed().expect("watch remains open"));

    manager.note_thread_shutdown(INTERACTIVE_THREAD_ID).await;
    assert!(count.has_changed().expect("watch remains open"));
    assert_eq!(*count.borrow_and_update(), 0);
}

等待标记变化没有改变 count,计数 watch 就不发布新值。一次状态变更可能唤醒 UI 状态观察者,却不唤醒退出 drain;反过来,不应因为计数 watch 没变化就判断状态逻辑没执行。

6.3 优雅退出 ​

运行计数的外部用途必须连同启用条件阅读。

源码文件:codex-rs/app-server/src/lib.rs

相关函数/类型:优雅重启 gate;行号:730-732。

rust
let single_client_mode = matches!(&transport, AppServerTransport::Stdio);
let graceful_signal_restart_enabled =
    runtime_options.install_shutdown_signal_handler && !single_client_mode;

安装关闭信号处理器且不是 stdio 单客户端模式,才启用这条信号驱动的优雅重启路径。不要把它写成所有 App Server 形态都使用相同退出协议。

源码文件:codex-rs/app-server/src/lib.rs

相关函数/类型:update;行号:257-283。

rust
fn update(&mut self, running_turn_count: usize, connection_count: usize) -> ShutdownAction {
    if !self.requested {
        return ShutdownAction::Noop;
    }

    if self.forced || running_turn_count == 0 {
        if self.forced {
            info!(
                "received second shutdown signal; forcing restart with {running_turn_count} running assistant turn(s) and {connection_count} connection(s)"
            );
        } else {
            info!(
                "shutdown signal restart: no assistant turns running; stopping acceptor and disconnecting {connection_count} connection(s)"
            );
        }
        return ShutdownAction::Finish;
    }

    if self.last_logged_running_turn_count != Some(running_turn_count) {
        info!(
            "shutdown signal restart: waiting for {running_turn_count} running assistant turn(s) to finish"
        );
        self.last_logged_running_turn_count = Some(running_turn_count);
    }

    ShutdownAction::Noop
}

未请求关闭时无动作;请求后若强制或 count 为零就 Finish。非零时等待并只在计数变化时更新日志。当前连接数参与日志与断连,但不是“还有客户端就永远不退出”的判据。

源码文件:codex-rs/app-server/src/lib.rs

相关函数/类型:running count 的关闭消费者;行号:943-956,970-974。

rust
let running_turn_count = {
    let running_turn_count = running_turn_count_rx.borrow();
    *running_turn_count
};
if matches!(
    shutdown_state.update(running_turn_count, connections.len()),
    ShutdownAction::Finish
) {
    transport_shutdown_token.cancel();
    let _ = outbound_control_tx
        .send(OutboundControlEvent::DisconnectAll)
        .await;
    break "shutdown_requested";
}
// ...
changed = running_turn_count_rx.changed(), if graceful_signal_restart_enabled && shutdown_state.requested() => {
    if changed.is_err() {
        warn!("running-turn watcher closed during graceful restart drain");
    }
}

主循环读取 watch 最新值,满足 Finish 就取消 transport token 并发送 DisconnectAll;等待期间监听 count 改变。Unix 集成测试构造一个持续约三秒的模型响应,发送 SIGINT 或 SIGTERM,要求 300ms 内进程不能提前退出,之后才成功退出。它证明该平台信号路径的 drain 行为,不证明所有平台上的信号语义一致。

两个消费者对同一组事实使用不同投影。

等待请求会阻止基于 Active 的空闲卸载,但自身不增加运行轮次计数。若修改任一判据,需要分别检查两条链的测试,不能仅凭 UI 标签正确就认为退出时序也正确。

7. 用时序检验状态 ​

7.1 断言的强弱 ​

单元测试可精确控制字段变化;完整集成测试还受到事件排队和卸载时机影响。下面这段真实测试不能只看变量名 saw_idle_after_turn 就概括成“结束后严格为 Idle”。

源码文件:codex-rs/app-server/tests/suite/v2/thread_status.rs

相关函数/类型:thread_status_changed_emits_runtime_updates 的实际断言;行号:82-101,111-124。

rust
match notification.status {
    ThreadStatus::Active { .. } => {
        saw_active_running = true;
    }
    ThreadStatus::Idle => {
        if saw_active_running {
            saw_idle_after_turn = true;
        }
    }
    ThreadStatus::SystemError => {
        if saw_active_running {
            saw_idle_after_turn = true;
        }
    }
    ThreadStatus::NotLoaded => {
        if saw_active_running {
            saw_idle_after_turn = true;
        }
    }
}
// ...
assert!(
    saw_active_running,
    "expected running active flag in thread/status/changed notifications"
);
assert!(
    saw_idle_after_turn,
    "expected idle status after turn completion in thread/status/changed notifications"
);
timeout(
    DEFAULT_READ_TIMEOUT,
    mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;

它接受 Active 之后的 Idle、SystemError 或 NotLoaded,再确认 turn/completed。其证明范围是运行状态变化确实穿过服务器到达客户端,并非所有执行结束路径都只能产生 Idle。精确聚合顺序应由前面的 RuntimeFacts 单元测试约束。

另一个单元测试分别为两个 Thread 订阅 watch,仅启动其中一个;目标 watcher 必须变化,另一 watcher 在短等待内保持 Idle。把按 Thread 订阅误做全局广播,会直接违反这个测试。对外 status 通知广播与进程内按 Thread watch 也必须分开理解。

7.2 异常定位 ​

现象先核对的事实可定位的函数
无模型生成但显示 Active两种 pending count 是否非零loaded_thread_status
审批回复后标记短暂未消失guard drop 的异步任务是否已执行note_active_guard_released
Turn 完成后仍 SystemError此前 Error 是否设置标记,是否已有新 TurnStartednote_system_error、note_turn_started
list 与运行响应短暂不同watched 状态、live Running、active snapshot 的时点resolve_thread_status、enrich_loaded_threads
客户端无 status 通知但 Turn 正常结束optOutNotificationMethods、状态是否实际改变should_skip_notification_for_connection
没有订阅仍暂时不卸载Active 状态、较晚的时间起点、Core 再检查unloading_target
已请求退出仍等待真正 running count 和强制标记ShutdownState::update

以下命令在 Codex 仓库根目录执行。集成测试使用临时配置和模拟模型;信号测试仅适用于其 Unix 测试模块。

bash
just test --locked -p codex-app-server --lib \
  -E 'test(thread_status::tests)'
just test --locked -p codex-app-server --test all \
  -E 'test(suite::v2::thread_status::)'
just test --locked -p codex-app-server --test all \
  -E 'test(websocket_transport_ctrl_c_waits_for_running_turn_before_exit) or test(websocket_transport_sigterm_waits_for_running_turn_before_exit)'
rg -n 'note_turn_started|note_system_error|note_permission_requested' \
  codex-rs/app-server/src/bespoke_event_handling.rs
rg -n 'loaded_thread_status|mutate_and_publish|resolve_thread_status' \
  codex-rs/app-server/src/thread_status.rs

最后按顺序推演:注册 Thread、创建两个审批 guard、启动 Turn、释放一个 guard、收到 Error、收到 TurnComplete、再启动新 Turn。逐步写出五个 RuntimeFacts 字段和公开状态,再分别判断 running count 与卸载资格。对于涉及 Drop 的步骤,要区分“刚安排清理任务”和“清理已经执行”两个时点;只有把这个异步边界画出来,才能用日志和测试判断状态是真的错误,还是观察尚未收敛。