Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Serialize Responses routing fields before large inputs (#49675)
## Why

Gateways that inspect request bodies incrementally need routing fields before potentially multi-megabyte prompts so they can route requests without buffering the input first.

## What changed

Serialize `model`, `stream`, and optional `service_tier` before `input` in both HTTP Responses requests and WebSocket `response.create` messages. Preserve field values and omission of unset `service_tier`.

## Testing

Add an HTTP regression test with a 2 MiB input that checks routing field order and unchanged JSON contents. Extend WebSocket serialization coverage to assert the routing fields appear first, and check that HTTP requests omit unset `service_tier`.

GitOrigin-RevId: dbc174e1fe7ca30d0c00f075bc37dd5f97b27b37
  • Loading branch information
jbeckwith-oai authored and copyberry committed Sep 30, 2026
commit ed0cc1a4ab30e1e83f1214e7a368c55b83fd089d
16 changes: 10 additions & 6 deletions codex-rs/codex-api/src/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -277,7 +277,12 @@ impl Serialize for ResponsesApiTools {

#[derive(Debug, Serialize, Clone, PartialEq)]
pub struct ResponsesApiRequest {
// Keep routing fields first: serde serializes struct fields in declaration order, and
// gateways may inspect request bodies incrementally before potentially multi-megabyte input.
pub model: String,
pub stream: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub service_tier: Option<String>,
#[serde(skip_serializing_if = "String::is_empty")]
pub instructions: String,
pub input: Vec<ResponseItem>,
Expand All @@ -287,13 +292,10 @@ pub struct ResponsesApiRequest {
pub parallel_tool_calls: bool,
pub reasoning: Option<Reasoning>,
pub store: bool,
pub stream: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub stream_options: Option<StreamOptions>,
pub include: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub service_tier: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub prompt_cache_key: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub text: Option<TextControls>,
Expand Down Expand Up @@ -330,7 +332,12 @@ impl<'a> From<&'a ResponsesApiRequest> for ResponseCreateWsRequest<'a> {

#[derive(Debug, Serialize)]
pub struct ResponseCreateWsRequest<'a> {
// Keep routing fields first: serde serializes struct fields in declaration order, and
// gateways may inspect request bodies incrementally before potentially multi-megabyte input.
pub model: &'a str,
pub stream: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub service_tier: Option<&'a str>,
#[serde(skip_serializing_if = "str::is_empty")]
pub instructions: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
Expand All @@ -342,13 +349,10 @@ pub struct ResponseCreateWsRequest<'a> {
pub parallel_tool_calls: bool,
pub reasoning: Option<&'a Reasoning>,
pub store: bool,
pub stream: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub stream_options: Option<&'a StreamOptions>,
pub include: &'a [String],
#[serde(skip_serializing_if = "Option::is_none")]
pub service_tier: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
pub prompt_cache_key: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
pub text: Option<&'a TextControls>,
Expand Down
3 changes: 3 additions & 0 deletions codex-rs/codex-api/src/endpoint/responses_websocket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -997,6 +997,9 @@ mod tests {
expected_payload["generate"] = json!(false);
let request_text =
serialize_websocket_request(&request).expect("serialize websocket request");
assert!(request_text.starts_with(
r#"{"type":"response.create","model":"gpt-test","stream":true,"service_tier":"priority","instructions":"Use the available tools.","previous_response_id":"resp-1","input":"#
));
let wire_payload =
serde_json::from_str::<Value>(&request_text).expect("parse websocket request");

Expand Down
58 changes: 58 additions & 0 deletions codex-rs/codex-api/tests/clients.rs
Original file line number Diff line number Diff line change
Expand Up @@ -376,13 +376,71 @@ async fn responses_client_stream_request_preserves_item_ids() -> Result<()> {
serde_json::from_slice(prepared.body.as_deref().expect("body should be JSON"))?;
assert_eq!(body, expected);
assert_eq!(body["input"][0]["id"], "msg_1");
assert_eq!(body.get("service_tier"), None);
assert_eq!(
prepared.headers.get(http::header::CONTENT_TYPE),
Some(&HeaderValue::from_static("application/json"))
);
Ok(())
}

#[tokio::test]
async fn responses_client_stream_request_sends_routing_fields_ahead_of_large_input() -> Result<()> {
let state = RecordingState::default();
let transport = RecordingTransport::new(state.clone());
let client = ResponsesClient::new(transport, provider("openai"), Arc::new(NoAuth));
// Exercise the customer gateway case where routing fields must be available without
// buffering a potentially multi-megabyte prompt first.
let large_input = "x".repeat(2 * 1024 * 1024);
let request = ResponsesApiRequest {
model: "gpt-test".into(),
instructions: "Say hi".into(),
input: vec![ResponseItem::Message {
id: None,
role: "user".into(),
content: vec![ContentItem::InputText { text: large_input }],
phase: None,
internal_chat_message_metadata_passthrough: None,
}],
tools: None,
tool_choice: "auto".into(),
parallel_tool_calls: false,
reasoning: None,
store: false,
stream: true,
stream_options: None,
include: Vec::new(),
service_tier: Some("priority".into()),
prompt_cache_key: None,
text: None,
client_metadata: None,
access_programs: None,
};
let expected = serde_json::to_value(&request)?;

let _stream = client
.stream_request(request, ResponsesOptions::default())
.await?;

let requests = state.take_stream_requests();
assert_eq!(requests.len(), 1);
let body = std::str::from_utf8(request_body_bytes(&requests[0]))?;
let input_position = body.find(r#""input":"#).expect("input should be present");
for routing_field in [r#""model":"#, r#""stream":"#, r#""service_tier":"#] {
assert!(
body.find(routing_field)
.is_some_and(|position| position < input_position),
"{routing_field} should precede input"
);
}
assert!(body.starts_with(
r#"{"model":"gpt-test","stream":true,"service_tier":"priority","instructions":"Say hi","input":[{"type":"message""#
));
assert!(body.len() > 2 * 1024 * 1024);
assert_eq!(serde_json::from_str::<serde_json::Value>(body)?, expected);
Ok(())
}

#[tokio::test]
async fn streaming_client_adds_auth_headers() -> Result<()> {
let state = RecordingState::default();
Expand Down
Loading