Skip to content

ModelClient结构

从 ModelClient 与 ModelClientSession 的分层出发,追踪认证、传输回退、turn state 和 WebSocket 连接缓存的生命周期。

基于rust-v0.150.0
CodexRustModelClient

ModelClient结构 ​

本文只研究 codex-rs/core/src/client.rs 中的客户端生命周期,不重复 ModelInfo能力模型 对模型字段的解释,也不重复 推理强度与服务层 对 reasoning 和 service tier 的解释,更不把实际 Responses body 的字段逐一展开。问题是:为什么一个 ModelClient 可以跨 turn 复用,而 ModelClientSession 必须每个 turn 新建?当 WebSocket 失败时,什么状态被清理,什么状态会让后续 turn 直接使用 HTTP?

读者需要能阅读 Rust 的 Arc、AtomicBool、OnceLock 和 Drop。本文的覆盖范围是客户端的状态和传输选择;不覆盖 provider 服务端的路由质量,也不覆盖所有 retry 分支的无网络复现。

1. 两个生命周期 ​

ModelClient 是 session 范围的 Clone 类型,其 state 用 Arc 共享;ModelClientSession 是 turn 范围的可变对象,保存当前 WebSocket 会话和 sticky routing token。这个分层把“跨 turn 稳定配置”与“本 turn 状态”分开。

2. 共享状态 ​

先看 owner:ModelClientState 是跨 turn 共享的状态包,但它不保存完整 Config。源码注释明确说明,大多数 turn 配置在调用时显式传入。因此不能把 ModelClientState 误读成“全局配置容器”。provider 的 auth/recovery 策略仍由 ModelProvider 拥有,ModelClient 只协调调用。

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

rust
#[derive(Debug)]
struct ModelClientState {
    thread_id: ThreadId,
    provider: SharedModelProvider,
    auth_env_telemetry: AuthEnvTelemetry,
    session_source: SessionSource,
    originator: String,
    model_verbosity: Option<VerbosityConfig>,
    enable_request_compression: bool,
    include_timing_metrics: bool,
    beta_features_header: Option<String>,
    concurrent_reasoning_summaries_enabled: bool,
    include_attestation: bool,
    attestation_provider: Option<Arc<dyn AttestationProvider>>,
    disable_websockets: AtomicBool,
    agent_identity_session_fallback: AgentIdentitySessionFallback,
    cached_websocket_session: StdMutex<WebsocketSession>,
}

provider 和 thread_id 是请求路由所需的稳定身份;beta_features_header 、压缩和 attestation 是请求选项;disable_websockets 和 cached_websocket_session 则是传输生命周期的状态。这些字段的所有者是 session client,而不是某一个 Prompt。

3. 构造与注入 ​

ModelClient::new 做两件事:通过 create_model_provider 建立 provider 抽象,再从 provider 能力和 auth 状态计算遥测与 attestation 开关。它不进行网络请求;网络设置在后续 turn 方法中解析。

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

rust
let model_provider = create_model_provider(provider_info, auth_manager);
let codex_api_key_env_enabled = model_provider
    .auth_manager()
    .as_ref()
    .is_some_and(|manager| manager.codex_api_key_env_enabled());
let auth_env_telemetry =
    collect_auth_env_telemetry(model_provider.info(), codex_api_key_env_enabled);
let include_attestation = model_provider.supports_attestation();
Self {
    state: Arc::new(ModelClientState {
        thread_id,
        provider: model_provider,
        auth_env_telemetry,
        session_source,
        originator,
        model_verbosity,
        enable_request_compression,
        include_timing_metrics,
        beta_features_header,
        concurrent_reasoning_summaries_enabled,
        include_attestation,
        attestation_provider,
        disable_websockets: AtomicBool::new(false),
        agent_identity_session_fallback: AgentIdentitySessionFallback::default(),
        cached_websocket_session: StdMutex::new(WebsocketSession::default()),
    }),
    agent_identity_policy,
    prompt_cache_key_override: None,
    http_client_factory,
}

http_client_factory 和 agent_identity_policy 没有放进 ModelClientState,说明它们属于 client 实例的请求调用依赖;而 state 中的字段需要随 clone 共享。这是一个可用于调试的所有者边界。

4. Turn会话 ​

new_session 只创建 turn 对象,不做 network I/O。它从 session 级 cache 中取出 WebsocketSession,并为当前 turn 新建一个空的 OnceLock<String>。

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

rust
pub fn new_session(&self) -> ModelClientSession {
    ModelClientSession {
        client: self.clone(),
        websocket_session: self.take_cached_websocket_session(),
        turn_state: Arc::new(OnceLock::new()),
    }
}

ModelClientSession 不应跨 turn 复用:源码文档指出,复用会把上一个 turn 的 sticky-routing token 带到下一个 turn。这不是性能偏好,而是客户端与服务端合同的边界。物理 WebSocket 连接可以通过 WebsocketSession 缓存复用,但 turn_state 每次 new_session 都重新创建。

5. Turn State ​

turn_state 是一个 turn 范围的 OnceLock:初始请求不携带它,服务端在 turn 开始后返回值,后续请求复用它。下一个 turn 会重新建立 OnceLock,因此即使物理 WebSocket 连接被缓存,turn state 也不会跨 turn 泄漏。

build_responses_headers 只会在 OnceLock 已有值时添加 x-codex-turn-state;这个头与请求 body 中的 client metadata 并非同一个数据源,实际传输路径会按 transport 调整。

6. 回退边界 ​

WebSocket 是否可用由两个条件决定:provider 声明支持,且 session 级 disable_websockets 还没有被设为 true。一旦 force_http_fallback 激活回退,它会写入原子标志、清空缓存的 WebSocket session,并记录 telemetry。

源码位置:codex-rs/core/src/client.rs :: ModelClient::responses_websocket_enabled 与 ModelClient::force_http_fallback

rust
pub fn responses_websocket_enabled(&self) -> bool {
    if !self.state.provider.info().supports_websockets
        || self.state.disable_websockets.load(Ordering::Relaxed)
    {
        return false;
    }

    true
}

pub(crate) fn force_http_fallback(
    &self,
    session_telemetry: &SessionTelemetry,
    _model_info: &ModelInfo,
) -> bool {
    let websocket_enabled = self.responses_websocket_enabled();
    let activated =
        websocket_enabled && !self.state.disable_websockets.swap(true, Ordering::Relaxed);
    if activated {
        warn!("falling back to HTTP");
        session_telemetry.counter(
            "codex.transport.fallback_to_http",
            /*inc*/ 1,
            &[("from_wire_api", "responses_websocket")],
        );
    }

    self.store_cached_websocket_session(WebsocketSession::default());
    activated
}

这里有两个容易误读的点。第一,provider 不支持 WebSocket 时只是返回 false,不会触发 telemetry fallback;第二,真正的失败回退会清空缓存并永久影响当前 Codex session 的后续 turn,而不是只影响当前请求。

7. Stream路由 ​

ModelClientSession::stream 在 Responses wire API 下先检查 WebSocket;WebSocket 路径返回 FallbackToHttp 时才调用 try_switch_fallback_transport,然后转入 HTTP Responses API。这是传输选择的 owner:request builder 负责 body,session 负责选择与回退。

WebsocketSession 还保存上一次请求、响应 receiver 和 warmup 标志,用于判断后续请求是否只是当前 turn 的 incremental extension。请求比较会显式忽略 metadata,但新增 request 字段必须进入穷举匹配,不能默认允许复用。

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

rust
struct WebsocketSession {
    connection: Option<ApiWebSocketConnection>,
    last_request: Option<ResponsesApiRequest>,
    last_response_rx: Option<oneshot::Receiver<LastResponse>>,
    last_response_from_untraced_warmup: bool,
    connection_reused: StdMutex<bool>,
}

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

rust
pub async fn stream(
    &mut self,
    prompt: &Prompt,
    model_info: &ModelInfo,
    session_telemetry: &SessionTelemetry,
    effort: Option<ReasoningEffortConfig>,
    summary: ReasoningSummaryConfig,
    service_tier: Option<String>,
    responses_metadata: &CodexResponsesMetadata,
    inference_trace: &InferenceTraceContext,
) -> Result<ResponseStream> {
    let wire_api = self.client.state.provider.info().wire_api;
    match wire_api {
        WireApi::Responses => {
            if self.client.responses_websocket_enabled() {
                let request_trace = current_span_w3c_trace_context();
                match self
                    .stream_responses_websocket(
                        prompt,
                        model_info,
                        session_telemetry,
                        effort.clone(),
                        summary,
                        service_tier.clone(),
                        responses_metadata,
                        /*warmup*/ false,
                        request_trace,
                        inference_trace,
                    )
                    .await?
                {
                    WebsocketStreamOutcome::Stream(stream) => return Ok(stream),
                    WebsocketStreamOutcome::FallbackToHttp => {
                        self.try_switch_fallback_transport(session_telemetry, model_info);
                    }
                }
            }

            self.stream_responses_api(
                prompt,
                model_info,
                session_telemetry,
                effort,
                summary,
                service_tier,
                responses_metadata,
                inference_trace,
            )
            .await
        }
    }
}

当前实现的路由分支是 Responses 路径;阅读重点是它先尝试 WebSocket,失败后才转 HTTP。WebSocket 已经被 session fallback 禁用时,当前 turn 和后续 turn 都会跳过第一段。

8. 回收与恢复 ​

turn 结束时,Drop for ModelClientSession 会把当前 WebsocketSession 移回 session 缓存。但 fallback 前已经调用 store_cached_websocket_session(WebsocketSession::default()),所以失败路径不会把旧连接重新恢复。“恢复”在这里是使用 HTTP 继续完成 turn,不是自动重新启用 WebSocket。

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

rust
impl Drop for ModelClientSession {
    fn drop(&mut self) {
        let websocket_session = std::mem::take(&mut self.websocket_session);
        self.client
            .store_cached_websocket_session(websocket_session);
    }
}

9. Client测试 ​

测试输入与动作断言覆盖范围
responses_turn_state_persists_within_turn_and_resets_afterHTTP mock server 返回 x-codex-turn-state,执行两个 turn同 turn 后续请求携带 token,下一 turn 不携带HTTP header 的 turn 隔离
websocket_turn_state_persists_within_turn_and_resets_afterWebSocket metadata 返回 token,执行后续请求物理连接可复用,但新 turn 仍从空 state 开始WebSocket body metadata 的 turn 边界
websocket_turn_state_is_stable_within_turn后续 metadata 尝试返回新值同 turn 后续请求仍使用第一个 tokenOnceLock 只保留首个值
responses_websocket_streams_requestmock WebSocket 收到一个 Responses 请求handshake 头、model、stream 和 input 存在一条正常 WebSocket 路径

skip_if_no_network! 会让这些集成测试在无网络环境跳过;因此测试证明的是 mock transport 下的 header、body 和状态迁移,不是真实 provider 网络故障下的重试成功率。

这些测试覆盖 mock transport 下的 turn state 隔离、首值稳定性和正常 WebSocket 请求;它们不覆盖真实 provider 的网络故障、认证刷新或服务端重试行为。

10. 复现与练习 ​

bash
RUST_MIN_STACK=16777216 cargo test -p codex-core responses_turn_state_persists_within_turn_and_resets_after
RUST_MIN_STACK=16777216 cargo test -p codex-core websocket_turn_state_persists_within_turn_and_resets_after
RUST_MIN_STACK=16777216 cargo test -p codex-core websocket_turn_state_is_stable_within_turn
RUST_MIN_STACK=16777216 cargo test -p codex-core responses_websocket_streams_request

源码练习:先把 ModelClientSession::new_session 中的 OnceLock::new() 改成跨 session 共用,再运行前三个 turn-state 测试,观察哪一个断言最先暴露跨 turn 泄漏;然后只移除 force_http_fallback 的 store_cached_websocket_session,思考为什么失败连接可能被 Drop 再次放回 cache。不要只看测试是否通过,还要检查状态 owner 和清理时机。

11. 源码导航 ​

先读本文的 ModelClientState、new 与 new_session,建立 session/turn 所有者边界;再读 responses_websocket_enabled、force_http_fallback 和 stream,追踪正常与降级路径;最后对照 turn-state 测试,验证传输连接可以跨 turn 缓存,而 sticky state 不能跨 turn 复用。