From d180f88e79c892e4e92ae072abf8e506aa71d91f Mon Sep 17 00:00:00 2001 From: Charles Cunningham Date: Thu, 5 Mar 2026 21:47:11 -0800 Subject: [PATCH] codex: address PR review feedback (#13681) Do not reuse the websocket upgrade request ID as per-turn rollout metadata for later stream requests. Co-authored-by: Codex --- .../src/endpoint/responses_websocket.rs | 27 +++++-------------- 1 file changed, 6 insertions(+), 21 deletions(-) diff --git a/codex-rs/codex-api/src/endpoint/responses_websocket.rs b/codex-rs/codex-api/src/endpoint/responses_websocket.rs index c0643219c0..9a169dea82 100644 --- a/codex-rs/codex-api/src/endpoint/responses_websocket.rs +++ b/codex-rs/codex-api/src/endpoint/responses_websocket.rs @@ -3,7 +3,6 @@ use crate::auth::add_auth_headers_to_header_map; use crate::common::ResponseEvent; use crate::common::ResponseStream; use crate::common::ResponsesWsRequest; -use crate::common::extract_request_id; use crate::error::ApiError; use crate::provider::Provider; use crate::rate_limits::parse_rate_limit_event; @@ -175,7 +174,6 @@ pub struct ResponsesWebsocketConnection { server_reasoning_included: bool, models_etag: Option, server_model: Option, - request_id: Option, telemetry: Option>, } @@ -187,7 +185,6 @@ impl std::fmt::Debug for ResponsesWebsocketConnection { .field("server_reasoning_included", &self.server_reasoning_included) .field("models_etag", &self.models_etag) .field("server_model", &self.server_model) - .field("request_id", &self.request_id) .field("telemetry", &self.telemetry.as_ref().map(|_| "")) .finish() } @@ -200,7 +197,6 @@ impl ResponsesWebsocketConnection { server_reasoning_included: bool, models_etag: Option, server_model: Option, - request_id: Option, telemetry: Option>, ) -> Self { Self { @@ -209,7 +205,6 @@ impl ResponsesWebsocketConnection { server_reasoning_included, models_etag, server_model, - request_id, telemetry, } } @@ -229,7 +224,6 @@ impl ResponsesWebsocketConnection { let server_reasoning_included = self.server_reasoning_included; let models_etag = self.models_etag.clone(); let server_model = self.server_model.clone(); - let request_id = self.request_id.clone(); let telemetry = self.telemetry.clone(); let request_body = serde_json::to_value(&request).map_err(|err| { ApiError::Stream(format!("failed to encode websocket request: {err}")) @@ -274,7 +268,10 @@ impl ResponsesWebsocketConnection { Ok(ResponseStream { rx_event, - initial_request_id: request_id, + // Websocket upgrade response headers are scoped to the connection, not + // individual `stream_request` calls, so they cannot be used as per-turn + // request IDs in rollout metadata. + initial_request_id: None, }) } } @@ -305,7 +302,7 @@ impl ResponsesWebsocketClient { merge_request_headers(&self.provider.headers, extra_headers, default_headers); add_auth_headers_to_header_map(&self.auth, &mut headers); - let (stream, server_reasoning_included, models_etag, server_model, request_id) = + let (stream, server_reasoning_included, models_etag, server_model) = connect_websocket(ws_url, headers, turn_state.clone()).await?; Ok(ResponsesWebsocketConnection::new( stream, @@ -313,7 +310,6 @@ impl ResponsesWebsocketClient { server_reasoning_included, models_etag, server_model, - request_id, telemetry, )) } @@ -338,16 +334,7 @@ async fn connect_websocket( url: Url, headers: HeaderMap, turn_state: Option>>, -) -> Result< - ( - WsStream, - bool, - Option, - Option, - Option, - ), - ApiError, -> { +) -> Result<(WsStream, bool, Option, Option), ApiError> { ensure_rustls_crypto_provider(); info!("connecting to websocket: {url}"); @@ -389,7 +376,6 @@ async fn connect_websocket( .get(OPENAI_MODEL_HEADER) .and_then(|value| value.to_str().ok()) .map(ToString::to_string); - let request_id = extract_request_id(response.headers()); if let Some(turn_state) = turn_state && let Some(header_value) = response .headers() @@ -403,7 +389,6 @@ async fn connect_websocket( reasoning_included, models_etag, server_model, - request_id, )) }