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
#[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
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
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
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
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
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
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_after | HTTP mock server 返回 x-codex-turn-state,执行两个 turn | 同 turn 后续请求携带 token,下一 turn 不携带 | HTTP header 的 turn 隔离 |
websocket_turn_state_persists_within_turn_and_resets_after | WebSocket metadata 返回 token,执行后续请求 | 物理连接可复用,但新 turn 仍从空 state 开始 | WebSocket body metadata 的 turn 边界 |
websocket_turn_state_is_stable_within_turn | 后续 metadata 尝试返回新值 | 同 turn 后续请求仍使用第一个 token | OnceLock 只保留首个值 |
responses_websocket_streams_request | mock WebSocket 收到一个 Responses 请求 | handshake 头、model、stream 和 input 存在 | 一条正常 WebSocket 路径 |
skip_if_no_network! 会让这些集成测试在无网络环境跳过;因此测试证明的是 mock transport 下的 header、body 和状态迁移,不是真实 provider 网络故障下的重试成功率。
这些测试覆盖 mock transport 下的 turn state 隔离、首值稳定性和正常 WebSocket 请求;它们不覆盖真实 provider 的网络故障、认证刷新或服务端重试行为。
10. 复现与练习
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 复用。
