diff --git a/codex-rs/cli/src/doctor.rs b/codex-rs/cli/src/doctor.rs index e3bdc65778..8dd0792f3d 100644 --- a/codex-rs/cli/src/doctor.rs +++ b/codex-rs/cli/src/doctor.rs @@ -2409,10 +2409,6 @@ async fn websocket_reachability_check( Ok(Ok(probe)) => { details.push(format!("handshake result: HTTP {}", probe.status)); details.push(format!("reasoning header: {}", probe.reasoning_included)); - details.push(format!( - "models etag present: {}", - probe.models_etag_present - )); details.push(format!( "server model present: {}", probe.server_model_present diff --git a/codex-rs/codex-api/src/endpoint/responses_websocket.rs b/codex-rs/codex-api/src/endpoint/responses_websocket.rs index 7e13c11953..733658fb39 100644 --- a/codex-rs/codex-api/src/endpoint/responses_websocket.rs +++ b/codex-rs/codex-api/src/endpoint/responses_websocket.rs @@ -184,7 +184,6 @@ pub struct ResponsesWebsocketConnection { // TODO (pakrym): is this the right place for timeout? idle_timeout: Duration, server_reasoning_included: bool, - models_etag: Option, server_model: Option, telemetry: Option>, } @@ -195,7 +194,6 @@ impl std::fmt::Debug for ResponsesWebsocketConnection { .field("stream", &"") .field("idle_timeout", &self.idle_timeout) .field("server_reasoning_included", &self.server_reasoning_included) - .field("models_etag", &self.models_etag) .field("server_model", &self.server_model) .field("telemetry", &self.telemetry.as_ref().map(|_| "")) .finish() @@ -207,7 +205,6 @@ impl ResponsesWebsocketConnection { stream: WsStream, idle_timeout: Duration, server_reasoning_included: bool, - models_etag: Option, server_model: Option, telemetry: Option>, ) -> Self { @@ -215,7 +212,6 @@ impl ResponsesWebsocketConnection { stream: Arc::new(Mutex::new(Some(stream))), idle_timeout, server_reasoning_included, - models_etag, server_model, telemetry, } @@ -242,7 +238,6 @@ impl ResponsesWebsocketConnection { let stream = Arc::clone(&self.stream); let idle_timeout = self.idle_timeout; let server_reasoning_included = self.server_reasoning_included; - let models_etag = self.models_etag.clone(); let server_model = self.server_model.clone(); let telemetry = self.telemetry.clone(); let ResponsesWsRequest::ResponseCreate(ws_request) = &request; @@ -282,9 +277,6 @@ impl ResponsesWebsocketConnection { if let Some(model) = server_model { let _ = tx_event.send(Ok(ResponseEvent::ServerModel(model))).await; } - if let Some(etag) = models_etag { - let _ = tx_event.send(Ok(ResponseEvent::ModelsEtag(etag))).await; - } if server_reasoning_included { let _ = tx_event .send(Ok(ResponseEvent::ServerReasoningIncluded(true))) @@ -356,8 +348,6 @@ pub struct ResponsesWebsocketProbe { pub status: StatusCode, /// Whether the server reported reasoning support in the upgrade response. pub reasoning_included: bool, - /// Whether the server returned a model catalog ETag in the upgrade response. - pub models_etag_present: bool, /// Whether the server returned a server-selected model in the upgrade response. pub server_model_present: bool, /// Close frame received immediately after upgrade, when one arrives quickly. @@ -393,13 +383,12 @@ impl ResponsesWebsocketClient { merge_request_headers(&self.provider.headers, extra_headers, default_headers); self.auth.add_auth_headers(&mut headers); - let (stream, _status, server_reasoning_included, models_etag, server_model) = + let (stream, _status, server_reasoning_included, server_model) = connect_websocket(ws_url, headers, http_client_factory, turn_state.clone()).await?; Ok(ResponsesWebsocketConnection::new( stream, self.provider.stream_idle_timeout, server_reasoning_included, - models_etag, server_model, telemetry, )) @@ -428,14 +417,13 @@ impl ResponsesWebsocketClient { merge_request_headers(&self.provider.headers, extra_headers, default_headers); self.auth.add_auth_headers(&mut headers); - let (mut stream, status, reasoning_included, models_etag, server_model) = - connect_websocket( - ws_url.clone(), - headers, - http_client_factory, - /*turn_state*/ None, - ) - .await?; + let (mut stream, status, reasoning_included, server_model) = connect_websocket( + ws_url.clone(), + headers, + http_client_factory, + /*turn_state*/ None, + ) + .await?; let immediate_close = tokio::time::timeout(immediate_close_timeout, stream.next()) .await .ok() @@ -450,7 +438,6 @@ impl ResponsesWebsocketClient { url: ws_url.to_string(), status, reasoning_included, - models_etag_present: models_etag.is_some(), server_model_present: server_model.is_some(), immediate_close, }) @@ -491,7 +478,7 @@ async fn connect_websocket( headers: HeaderMap, http_client_factory: &HttpClientFactory, turn_state: Option>>, -) -> Result<(WsStream, StatusCode, bool, Option, Option), ApiError> { +) -> Result<(WsStream, StatusCode, bool, Option), ApiError> { info!("connecting to websocket: {url}"); let mut request = url @@ -519,11 +506,6 @@ async fn connect_websocket( }; let reasoning_included = response.headers().contains_key(X_REASONING_INCLUDED_HEADER); - let models_etag = response - .headers() - .get(X_MODELS_ETAG_HEADER) - .and_then(|value| value.to_str().ok()) - .map(ToString::to_string); let server_model = response .headers() .get(OPENAI_MODEL_HEADER) @@ -541,7 +523,6 @@ async fn connect_websocket( WsStream::new(stream), response.status(), reasoning_included, - models_etag, server_model, )) } @@ -736,6 +717,21 @@ async fn run_websocket_response_stream( text.as_str(), timing_log_context, ); + if event.kind() == "codex.response.metadata" + && let Some(etag) = + event + .headers + .as_ref() + .and_then(Value::as_object) + .and_then(|headers| { + json_headers_to_http_headers(headers) + .get(X_MODELS_ETAG_HEADER) + .and_then(|value| value.to_str().ok()) + .map(str::to_string) + }) + { + let _ = tx_event.send(Ok(ResponseEvent::ModelsEtag(etag))).await; + } if let Some(response_turn_state) = event.turn_state() && let Some(turn_state) = turn_state { diff --git a/codex-rs/core/tests/suite/client_websockets.rs b/codex-rs/core/tests/suite/client_websockets.rs index 400282e476..1770da1973 100755 --- a/codex-rs/core/tests/suite/client_websockets.rs +++ b/codex-rs/core/tests/suite/client_websockets.rs @@ -1496,14 +1496,15 @@ async fn responses_websocket_emits_rate_limit_events() { let server = start_websocket_server_with_headers(vec![WebSocketConnectionConfig { requests: vec![vec![ + json!({ + "type": "codex.response.metadata", + "headers": {"x-models-etag": "etag-123"}, + }), rate_limit_event, ev_response_created("resp-1"), ev_completed("resp-1"), ]], - response_headers: vec![ - ("X-Models-Etag".to_string(), "etag-123".to_string()), - ("X-Reasoning-Included".to_string(), "true".to_string()), - ], + response_headers: vec![("X-Reasoning-Included".to_string(), "true".to_string())], accept_delay: None, close_after_requests: true, }])