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 <noreply@openai.com>
This commit is contained in:
Charles Cunningham
2026-03-05 21:47:11 -08:00
parent 71033e034c
commit d180f88e79

View File

@@ -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<String>,
server_model: Option<String>,
request_id: Option<String>,
telemetry: Option<Arc<dyn WebsocketTelemetry>>,
}
@@ -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(|_| "<telemetry>"))
.finish()
}
@@ -200,7 +197,6 @@ impl ResponsesWebsocketConnection {
server_reasoning_included: bool,
models_etag: Option<String>,
server_model: Option<String>,
request_id: Option<String>,
telemetry: Option<Arc<dyn WebsocketTelemetry>>,
) -> 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<A: AuthProvider> ResponsesWebsocketClient<A> {
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<A: AuthProvider> ResponsesWebsocketClient<A> {
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<Arc<OnceLock<String>>>,
) -> Result<
(
WsStream,
bool,
Option<String>,
Option<String>,
Option<String>,
),
ApiError,
> {
) -> Result<(WsStream, bool, Option<String>, Option<String>), 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,
))
}