From a92a5085bd01203737876915d609aac7f64b4638 Mon Sep 17 00:00:00 2001 From: neil-oai Date: Thu, 9 Apr 2026 14:52:37 -0400 Subject: [PATCH] Forward app-server turn clientMetadata to Responses (#16009) ## Summary App-server v2 already receives turn-scoped `clientMetadata`, but the Rust app-server was dropping it before the outbound Responses request. This change keeps the fix lightweight by threading that metadata through the existing turn-metadata path rather than inventing a new transport. ## What we're trying to do and why We want turn-scoped metadata from the app-server protocol layer, especially fields like Hermes/GAAS run IDs, to survive all the way to the actual Responses API request so it is visible in downstream websocket request logging and analytics. The specific bug was: - app-server protocol uses camelCase `clientMetadata` - Responses transport already has an existing turn metadata carrier: `x-codex-turn-metadata` - websocket transport already rewrites that header into `request.request_body.client_metadata["x-codex-turn-metadata"]` - but the Rust app-server never parsed or stored `clientMetadata`, so nothing from the app-server request was making it into that existing path This PR fixes that without adding a new header or a second metadata channel. ## How we did it ### Protocol surface - Add optional `clientMetadata` to v2 `TurnStartParams` and `TurnSteerParams` - Regenerate the JSON schema / TypeScript fixtures - Update app-server docs to describe the field and its behavior ### Runtime plumbing - Add a dedicated core op for app-server user input carrying turn-scoped metadata: `Op::UserInputWithClientMetadata` - Wire `turn/start` and `turn/steer` through that op / signature path instead of dropping the metadata at the message-processor boundary - Store the metadata in `TurnMetadataState` ### Transport behavior - Reuse the existing serialized `x-codex-turn-metadata` payload - Merge the new app-server `clientMetadata` into that JSON additively - Do **not** replace built-in reserved fields already present in the turn metadata payload - Keep websocket behavior unchanged at the outer shape level: it still sends only `client_metadata["x-codex-turn-metadata"]`, but that JSON string now contains the merged fields - Keep HTTP fallback behavior unchanged except that the existing `x-codex-turn-metadata` header now includes the merged fields too ### Request shape before / after Before, a websocket `response.create` looked like: ```json { "type": "response.create", "client_metadata": { "x-codex-turn-metadata": "{\"session_id\":\"...\",\"turn_id\":\"...\"}" } } ``` Even if the app-server caller supplied `clientMetadata`, it was not represented there. After, the same request shape is preserved, but the serialized payload now includes the new turn-scoped fields: ```json { "type": "response.create", "client_metadata": { "x-codex-turn-metadata": "{\"session_id\":\"...\",\"turn_id\":\"...\",\"fiber_run_id\":\"fiber-start-123\",\"origin\":\"gaas\"}" } } ``` ## Validation ### Targeted tests added / updated - protocol round-trip coverage for `clientMetadata` on `turn/start` and `turn/steer` - protocol round-trip coverage for `Op::UserInputWithClientMetadata` - `TurnMetadataState` merge test proving client metadata is added without overwriting reserved built-in fields - websocket request-shape test proving outbound `response.create` contains merged metadata inside `client_metadata["x-codex-turn-metadata"]` - app-server integration tests proving: - `turn/start` forwards `clientMetadata` into the outbound Responses request path - websocket warmup + real turn request both behave correctly - `turn/steer` updates the follow-up request metadata ### Commands run - `just write-app-server-schema` - `cargo test -p codex-app-server-protocol` - `cargo test -p codex-protocol` - `cargo test -p codex-core turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields --lib` - `cargo test -p codex-core --test all responses_websocket_preserves_custom_turn_metadata_fields` - `cargo test -p codex-app-server --test all client_metadata` - `cargo test -p codex-app-server --test all turn_start_forwards_client_metadata_to_responses_websocket_request_body_v2 -- --nocapture` - `just fmt` - `just fix -p codex-core -p codex-protocol -p codex-app-server-protocol -p codex-app-server` - `just fix -p codex-exec -p codex-tui-app-server` - `just argument-comment-lint` ### Full suite note `cargo test` in `codex-rs` still fails in: - `suite::v2::turn_interrupt::turn_interrupt_resolves_pending_command_approval_request` I verified that same failure on a clean detached `HEAD` worktree with an isolated `CARGO_TARGET_DIR`, so it is not caused by this patch. --- .../schema/typescript/v2/TurnSteerParams.ts | 5 +- .../src/protocol/common.rs | 1 + .../app-server-protocol/src/protocol/v2.rs | 13 +- .../app-server/src/codex_message_processor.rs | 7 +- .../src/message_processor/tracing_tests.rs | 1 + .../src/transport/remote_control/tests.rs | 11 +- .../src/transport/remote_control/websocket.rs | 28 +- .../tests/suite/v2/client_metadata.rs | 366 ++++++++++++++++++ codex-rs/app-server/tests/suite/v2/mod.rs | 1 + .../app-server/tests/suite/v2/turn_start.rs | 2 + .../app-server/tests/suite/v2/turn_steer.rs | 3 + codex-rs/core/src/agent/control_tests.rs | 3 + codex-rs/core/src/codex.rs | 32 +- codex-rs/core/src/codex_delegate.rs | 1 + codex-rs/core/src/codex_tests.rs | 27 +- codex-rs/core/src/codex_thread.rs | 6 +- .../src/tools/handlers/multi_agents_tests.rs | 1 + codex-rs/core/src/turn_metadata.rs | 45 ++- codex-rs/core/src/turn_metadata_tests.rs | 27 ++ codex-rs/core/tests/suite/abort_tasks.rs | 5 + codex-rs/core/tests/suite/client.rs | 31 ++ .../core/tests/suite/client_websockets.rs | 69 ++++ .../tests/suite/collaboration_instructions.rs | 14 + codex-rs/core/tests/suite/compact.rs | 29 ++ codex-rs/core/tests/suite/compact_remote.rs | 43 ++ .../core/tests/suite/compact_resume_fork.rs | 1 + codex-rs/core/tests/suite/fork_thread.rs | 1 + codex-rs/core/tests/suite/hooks.rs | 2 + codex-rs/core/tests/suite/items.rs | 9 + .../core/tests/suite/model_visible_layout.rs | 3 + codex-rs/core/tests/suite/otel.rs | 22 ++ codex-rs/core/tests/suite/pending_input.rs | 4 + .../core/tests/suite/permissions_messages.rs | 15 + codex-rs/core/tests/suite/plugins.rs | 3 + codex-rs/core/tests/suite/prompt_caching.rs | 10 + codex-rs/core/tests/suite/quota_exceeded.rs | 1 + .../core/tests/suite/realtime_conversation.rs | 1 + .../core/tests/suite/request_compression.rs | 2 + codex-rs/core/tests/suite/resume.rs | 7 + codex-rs/core/tests/suite/review.rs | 1 + codex-rs/core/tests/suite/search_tool.rs | 1 + codex-rs/core/tests/suite/sqlite_state.rs | 13 +- .../suite/stream_error_allows_next_turn.rs | 2 + .../core/tests/suite/stream_no_completed.rs | 1 + .../core/tests/suite/user_notification.rs | 1 + codex-rs/core/tests/suite/window_headers.rs | 1 + codex-rs/exec/src/lib.rs | 1 + codex-rs/mcp-server/src/codex_tool_runner.rs | 2 + codex-rs/protocol/src/protocol.rs | 34 ++ codex-rs/state/src/runtime/logs.rs | 1 + codex-rs/tui/src/app_server_session.rs | 2 + 51 files changed, 867 insertions(+), 45 deletions(-) create mode 100644 codex-rs/app-server/tests/suite/v2/client_metadata.rs diff --git a/codex-rs/app-server-protocol/schema/typescript/v2/TurnSteerParams.ts b/codex-rs/app-server-protocol/schema/typescript/v2/TurnSteerParams.ts index 2c84f195c..dae166b40 100644 --- a/codex-rs/app-server-protocol/schema/typescript/v2/TurnSteerParams.ts +++ b/codex-rs/app-server-protocol/schema/typescript/v2/TurnSteerParams.ts @@ -3,9 +3,8 @@ // This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. import type { UserInput } from "./UserInput"; -export type TurnSteerParams = { threadId: string, input: Array, -/** +export type TurnSteerParams = {threadId: string, input: Array, /** * Required active turn id precondition. The request fails when it does not * match the currently active turn. */ -expectedTurnId: string, }; +expectedTurnId: string}; diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index 1f3a935a7..731288272 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -388,6 +388,7 @@ client_request_definitions! { }, TurnSteer => "turn/steer" { params: v2::TurnSteerParams, + inspect_params: true, response: v2::TurnSteerResponse, }, TurnInterrupt => "turn/interrupt" { diff --git a/codex-rs/app-server-protocol/src/protocol/v2.rs b/codex-rs/app-server-protocol/src/protocol/v2.rs index 4e4d6264e..de96ef2bf 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2.rs @@ -4034,6 +4034,10 @@ pub enum TurnStatus { pub struct TurnStartParams { pub thread_id: String, pub input: Vec, + /// Optional turn-scoped Responses API client metadata. + #[experimental("turn/start.responsesapiClientMetadata")] + #[ts(optional = nullable)] + pub responsesapi_client_metadata: Option>, /// Override the working directory for this turn and subsequent turns. #[ts(optional = nullable)] pub cwd: Option, @@ -4144,12 +4148,18 @@ pub struct TurnStartResponse { pub turn: Turn, } -#[derive(Serialize, Deserialize, Debug, Default, Clone, PartialEq, JsonSchema, TS)] +#[derive( + Serialize, Deserialize, Debug, Default, Clone, PartialEq, JsonSchema, TS, ExperimentalApi, +)] #[serde(rename_all = "camelCase")] #[ts(export_to = "v2/")] pub struct TurnSteerParams { pub thread_id: String, pub input: Vec, + /// Optional turn-scoped Responses API client metadata. + #[experimental("turn/steer.responsesapiClientMetadata")] + #[ts(optional = nullable)] + pub responsesapi_client_metadata: Option>, /// Required active turn id precondition. The request fails when it does not /// match the currently active turn. pub expected_turn_id: String, @@ -8422,6 +8432,7 @@ mod tests { let without_override = TurnStartParams { thread_id: "thread_123".to_string(), input: vec![], + responsesapi_client_metadata: None, cwd: None, approval_policy: None, approvals_reviewer: None, diff --git a/codex-rs/app-server/src/codex_message_processor.rs b/codex-rs/app-server/src/codex_message_processor.rs index aae240e3f..9cb8403f1 100644 --- a/codex-rs/app-server/src/codex_message_processor.rs +++ b/codex-rs/app-server/src/codex_message_processor.rs @@ -6666,6 +6666,7 @@ impl CodexMessageProcessor { Op::UserInput { items: mapped_items, final_output_json_schema: params.output_schema, + responsesapi_client_metadata: params.responsesapi_client_metadata, }, ) .await; @@ -6746,7 +6747,11 @@ impl CodexMessageProcessor { .collect(); match thread - .steer_input(mapped_items, Some(¶ms.expected_turn_id)) + .steer_input( + mapped_items, + Some(¶ms.expected_turn_id), + params.responsesapi_client_metadata, + ) .await { Ok(turn_id) => { diff --git a/codex-rs/app-server/src/message_processor/tracing_tests.rs b/codex-rs/app-server/src/message_processor/tracing_tests.rs index ef88364a3..d0fe22e9b 100644 --- a/codex-rs/app-server/src/message_processor/tracing_tests.rs +++ b/codex-rs/app-server/src/message_processor/tracing_tests.rs @@ -606,6 +606,7 @@ async fn turn_start_jsonrpc_span_parents_core_turn_spans() -> Result<()> { text: "hello".to_string(), text_elements: Vec::new(), }], + responsesapi_client_metadata: None, cwd: None, approval_policy: None, sandbox_policy: None, diff --git a/codex-rs/app-server/src/transport/remote_control/tests.rs b/codex-rs/app-server/src/transport/remote_control/tests.rs index 9d430ccfb..21808a3a1 100644 --- a/codex-rs/app-server/src/transport/remote_control/tests.rs +++ b/codex-rs/app-server/src/transport/remote_control/tests.rs @@ -106,13 +106,10 @@ async fn remote_control_state_runtime(codex_home: &TempDir) -> Arc } fn remote_control_url_for_listener(listener: &TcpListener) -> String { - format!( - "http://localhost:{}/backend-api/", - listener - .local_addr() - .expect("listener should have a local addr") - .port() - ) + let addr = listener + .local_addr() + .expect("listener should have a local addr"); + format!("http://{addr}/backend-api/") } #[tokio::test] diff --git a/codex-rs/app-server/src/transport/remote_control/websocket.rs b/codex-rs/app-server/src/transport/remote_control/websocket.rs index a0387ef6c..42924c2d3 100644 --- a/codex-rs/app-server/src/transport/remote_control/websocket.rs +++ b/codex-rs/app-server/src/transport/remote_control/websocket.rs @@ -931,13 +931,10 @@ mod tests { } fn remote_control_url_for_listener(listener: &TcpListener) -> String { - format!( - "http://localhost:{}/backend-api/", - listener - .local_addr() - .expect("listener should have a local addr") - .port() - ) + let addr = listener + .local_addr() + .expect("listener should have a local addr"); + format!("http://{addr}/backend-api/") } fn remote_control_auth_dot_json(access_token: &str) -> AuthDotJson { @@ -1041,14 +1038,6 @@ mod tests { let remote_control_url = remote_control_url_for_listener(&listener); let remote_control_target = normalize_remote_control_url(&remote_control_url).expect("target should parse"); - let server_task = tokio::spawn(async move { - let (stream, request_line) = accept_http_request(&listener).await; - assert_eq!( - request_line, - "GET /backend-api/wham/remote/control/server HTTP/1.1" - ); - respond_with_status_and_headers(stream, "401 Unauthorized", &[], "unauthorized").await; - }); let codex_home = TempDir::new().expect("temp dir should create"); save_auth( codex_home.path(), @@ -1076,6 +1065,15 @@ mod tests { ) .expect("fresh auth should save"); + let server_task = tokio::spawn(async move { + let (stream, request_line) = accept_http_request(&listener).await; + assert_eq!( + request_line, + "GET /backend-api/wham/remote/control/server HTTP/1.1" + ); + respond_with_status_and_headers(stream, "401 Unauthorized", &[], "unauthorized").await; + }); + let err = connect_remote_control_websocket( &remote_control_target, Some(state_db.as_ref()), diff --git a/codex-rs/app-server/tests/suite/v2/client_metadata.rs b/codex-rs/app-server/tests/suite/v2/client_metadata.rs new file mode 100644 index 000000000..c85febd7d --- /dev/null +++ b/codex-rs/app-server/tests/suite/v2/client_metadata.rs @@ -0,0 +1,366 @@ +use anyhow::Result; +use app_test_support::McpProcess; +use app_test_support::to_response; +use codex_app_server_protocol::JSONRPCResponse; +use codex_app_server_protocol::RequestId; +use codex_app_server_protocol::ThreadStartParams; +use codex_app_server_protocol::ThreadStartResponse; +use codex_app_server_protocol::TurnStartParams; +use codex_app_server_protocol::TurnStartResponse; +use codex_app_server_protocol::TurnSteerParams; +use codex_app_server_protocol::TurnSteerResponse; +use codex_app_server_protocol::UserInput as V2UserInput; +use core_test_support::responses; +use core_test_support::skip_if_no_network; +use pretty_assertions::assert_eq; +use std::collections::HashMap; +use std::path::Path; +use tempfile::TempDir; +use tokio::time::timeout; + +const DEFAULT_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); + +#[tokio::test] +async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result<()> { + skip_if_no_network!(Ok(())); + + let server = responses::start_mock_server().await; + let response_mock = responses::mount_sse_once( + &server, + responses::sse(vec![ + responses::ev_response_created("resp-1"), + responses::ev_assistant_message("msg-1", "Done"), + responses::ev_completed("resp-1"), + ]), + ) + .await; + + let codex_home = TempDir::new()?; + create_config_toml( + codex_home.path(), + &server.uri(), + /*supports_websockets*/ false, + )?; + + let mut mcp = McpProcess::new(codex_home.path()).await?; + timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; + + let thread_req = mcp + .send_thread_start_request(ThreadStartParams::default()) + .await?; + let thread_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(thread_req)), + ) + .await??; + let ThreadStartResponse { thread, .. } = to_response::(thread_resp)?; + + let client_metadata = HashMap::from([ + ("fiber_run_id".to_string(), "fiber-start-123".to_string()), + ("origin".to_string(), "gaas".to_string()), + ]); + let turn_req = mcp + .send_turn_start_request(TurnStartParams { + thread_id: thread.id, + input: vec![V2UserInput::Text { + text: "Hello".to_string(), + text_elements: Vec::new(), + }], + responsesapi_client_metadata: Some(client_metadata.clone()), + ..Default::default() + }) + .await?; + let turn_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(turn_req)), + ) + .await??; + let TurnStartResponse { turn } = to_response::(turn_resp)?; + + timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_notification_message("turn/completed"), + ) + .await??; + + let request = response_mock.single_request(); + let metadata = request + .header("x-codex-turn-metadata") + .as_deref() + .map(parse_json_header) + .unwrap_or_else(|| panic!("missing x-codex-turn-metadata header")); + assert_eq!(metadata["fiber_run_id"].as_str(), Some("fiber-start-123")); + assert_eq!(metadata["origin"].as_str(), Some("gaas")); + assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str())); + assert!(metadata.get("session_id").is_some()); + + Ok(()) +} + +#[tokio::test] +async fn turn_steer_updates_client_metadata_on_follow_up_responses_request_v2() -> Result<()> { + skip_if_no_network!(Ok(())); + + let codex_home = TempDir::new()?; + + let server = responses::start_mock_server().await; + let first_response = responses::sse_response(responses::sse(vec![ + responses::ev_response_created("resp-1"), + responses::ev_assistant_message("msg-1", "Working"), + responses::ev_completed("resp-1"), + ])) + .set_delay(std::time::Duration::from_secs(2)); + let second_response = responses::sse_response(responses::sse(vec![ + responses::ev_response_created("resp-2"), + responses::ev_assistant_message("msg-2", "Done"), + responses::ev_completed("resp-2"), + ])); + let request_log = + responses::mount_response_sequence(&server, vec![first_response, second_response]).await; + + create_config_toml( + codex_home.path(), + &server.uri(), + /*supports_websockets*/ false, + )?; + + let mut mcp = McpProcess::new(codex_home.path()).await?; + timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; + + let thread_req = mcp + .send_thread_start_request(ThreadStartParams::default()) + .await?; + let thread_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(thread_req)), + ) + .await??; + let ThreadStartResponse { thread, .. } = to_response::(thread_resp)?; + + let start_metadata = + HashMap::from([("fiber_run_id".to_string(), "fiber-start-123".to_string())]); + let turn_req = mcp + .send_turn_start_request(TurnStartParams { + thread_id: thread.id.clone(), + input: vec![V2UserInput::Text { + text: "Run sleep".to_string(), + text_elements: Vec::new(), + }], + responsesapi_client_metadata: Some(start_metadata.clone()), + ..Default::default() + }) + .await?; + let turn_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(turn_req)), + ) + .await??; + let TurnStartResponse { turn } = to_response::(turn_resp)?; + let turn_id = turn.id.clone(); + + timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_notification_message("turn/started"), + ) + .await??; + wait_for_request_count(&request_log, /*expected*/ 1).await?; + + let steer_metadata = HashMap::from([ + ("fiber_run_id".to_string(), "fiber-steer-456".to_string()), + ("origin".to_string(), "gaas".to_string()), + ]); + let steer_req = mcp + .send_turn_steer_request(TurnSteerParams { + thread_id: thread.id.clone(), + input: vec![V2UserInput::Text { + text: "Focus on the failure".to_string(), + text_elements: Vec::new(), + }], + responsesapi_client_metadata: Some(steer_metadata.clone()), + expected_turn_id: turn_id.clone(), + }) + .await?; + let steer_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(steer_req)), + ) + .await??; + let _turn: TurnSteerResponse = to_response::(steer_resp)?; + + timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_notification_message("turn/completed"), + ) + .await??; + + let requests = request_log.requests(); + assert_eq!(requests.len(), 2); + let first_metadata = requests[0] + .header("x-codex-turn-metadata") + .as_deref() + .map(parse_json_header) + .unwrap_or_else(|| panic!("missing first x-codex-turn-metadata header")); + assert_eq!( + first_metadata["fiber_run_id"].as_str(), + Some("fiber-start-123") + ); + assert_eq!(first_metadata["turn_id"].as_str(), Some(turn_id.as_str())); + + let second_metadata = requests[1] + .header("x-codex-turn-metadata") + .as_deref() + .map(parse_json_header) + .unwrap_or_else(|| panic!("missing second x-codex-turn-metadata header")); + assert_eq!( + second_metadata["fiber_run_id"].as_str(), + Some("fiber-steer-456") + ); + assert_eq!(second_metadata["origin"].as_str(), Some("gaas")); + assert_eq!(second_metadata["turn_id"].as_str(), Some(turn_id.as_str())); + + Ok(()) +} + +#[tokio::test] +async fn turn_start_forwards_client_metadata_to_responses_websocket_request_body_v2() -> Result<()> +{ + skip_if_no_network!(Ok(())); + + let websocket_server = responses::start_websocket_server(vec![vec![ + vec![ + responses::ev_response_created("warm-1"), + responses::ev_completed("warm-1"), + ], + vec![ + responses::ev_response_created("resp-1"), + responses::ev_assistant_message("msg-1", "Done"), + responses::ev_completed("resp-1"), + ], + ]]) + .await; + + let codex_home = TempDir::new()?; + create_config_toml( + codex_home.path(), + &websocket_server.uri().replacen("ws://", "http://", 1), + /*supports_websockets*/ true, + )?; + + let mut mcp = McpProcess::new(codex_home.path()).await?; + timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; + + let thread_req = mcp + .send_thread_start_request(ThreadStartParams::default()) + .await?; + let thread_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(thread_req)), + ) + .await??; + let ThreadStartResponse { thread, .. } = to_response::(thread_resp)?; + + let client_metadata = HashMap::from([ + ("fiber_run_id".to_string(), "fiber-start-123".to_string()), + ("origin".to_string(), "gaas".to_string()), + ]); + let turn_req = mcp + .send_turn_start_request(TurnStartParams { + thread_id: thread.id, + input: vec![V2UserInput::Text { + text: "Hello".to_string(), + text_elements: Vec::new(), + }], + responsesapi_client_metadata: Some(client_metadata), + ..Default::default() + }) + .await?; + let turn_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(turn_req)), + ) + .await??; + let TurnStartResponse { turn } = to_response::(turn_resp)?; + + timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_notification_message("turn/completed"), + ) + .await??; + + let warmup = websocket_server + .wait_for_request(/*connection_index*/ 0, /*request_index*/ 0) + .await + .body_json(); + let request = websocket_server + .wait_for_request(/*connection_index*/ 0, /*request_index*/ 1) + .await + .body_json(); + + assert_eq!(warmup["type"].as_str(), Some("response.create")); + assert_eq!(warmup["generate"].as_bool(), Some(false)); + assert_eq!(request["type"].as_str(), Some("response.create")); + assert_eq!(request["previous_response_id"].as_str(), Some("warm-1")); + + let metadata = request["client_metadata"]["x-codex-turn-metadata"] + .as_str() + .map(parse_json_header) + .unwrap_or_else(|| panic!("missing websocket x-codex-turn-metadata client metadata")); + assert_eq!(metadata["fiber_run_id"].as_str(), Some("fiber-start-123")); + assert_eq!(metadata["origin"].as_str(), Some("gaas")); + assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str())); + assert!(metadata.get("session_id").is_some()); + + websocket_server.shutdown().await; + Ok(()) +} + +fn create_config_toml( + codex_home: &Path, + server_uri: &str, + supports_websockets: bool, +) -> std::io::Result<()> { + let config_toml = codex_home.join("config.toml"); + std::fs::write( + config_toml, + format!( + r#" +model = "mock-model" +approval_policy = "never" +sandbox_mode = "read-only" + +model_provider = "mock_provider" + +[model_providers.mock_provider] +name = "Mock provider for test" +base_url = "{server_uri}/v1" +wire_api = "responses" +request_max_retries = 0 +stream_max_retries = 0 +supports_websockets = {supports_websockets} +"# + ), + ) +} + +fn parse_json_header(value: &str) -> serde_json::Value { + match serde_json::from_str(value) { + Ok(value) => value, + Err(err) => panic!("metadata header should be valid json: {err}"), + } +} + +async fn wait_for_request_count( + request_log: &core_test_support::responses::ResponseMock, + expected: usize, +) -> Result<()> { + timeout(DEFAULT_READ_TIMEOUT, async { + loop { + if request_log.requests().len() >= expected { + return; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await?; + Ok(()) +} diff --git a/codex-rs/app-server/tests/suite/v2/mod.rs b/codex-rs/app-server/tests/suite/v2/mod.rs index 6cd67daa5..1bb74ccc6 100644 --- a/codex-rs/app-server/tests/suite/v2/mod.rs +++ b/codex-rs/app-server/tests/suite/v2/mod.rs @@ -1,6 +1,7 @@ mod account; mod analytics; mod app_list; +mod client_metadata; mod collaboration_mode_list; #[cfg(unix)] mod command_exec; diff --git a/codex-rs/app-server/tests/suite/v2/turn_start.rs b/codex-rs/app-server/tests/suite/v2/turn_start.rs index 59591dd0d..8f276e331 100644 --- a/codex-rs/app-server/tests/suite/v2/turn_start.rs +++ b/codex-rs/app-server/tests/suite/v2/turn_start.rs @@ -1377,6 +1377,7 @@ async fn turn_start_updates_sandbox_and_cwd_between_turns_v2() -> Result<()> { text: "first turn".to_string(), text_elements: Vec::new(), }], + responsesapi_client_metadata: None, cwd: Some(first_cwd.clone()), approval_policy: Some(codex_app_server_protocol::AskForApproval::Never), approvals_reviewer: None, @@ -1416,6 +1417,7 @@ async fn turn_start_updates_sandbox_and_cwd_between_turns_v2() -> Result<()> { text: "second turn".to_string(), text_elements: Vec::new(), }], + responsesapi_client_metadata: None, cwd: Some(second_cwd.clone()), approval_policy: Some(codex_app_server_protocol::AskForApproval::Never), approvals_reviewer: None, diff --git a/codex-rs/app-server/tests/suite/v2/turn_steer.rs b/codex-rs/app-server/tests/suite/v2/turn_steer.rs index 5d1b3cc22..a93bf6c6a 100644 --- a/codex-rs/app-server/tests/suite/v2/turn_steer.rs +++ b/codex-rs/app-server/tests/suite/v2/turn_steer.rs @@ -57,6 +57,7 @@ async fn turn_steer_requires_active_turn() -> Result<()> { text: "steer".to_string(), text_elements: Vec::new(), }], + responsesapi_client_metadata: None, expected_turn_id: "turn-does-not-exist".to_string(), }) .await?; @@ -145,6 +146,7 @@ async fn turn_steer_rejects_oversized_text_input() -> Result<()> { text: oversized_input.clone(), text_elements: Vec::new(), }], + responsesapi_client_metadata: None, expected_turn_id: turn.id.clone(), }) .await?; @@ -247,6 +249,7 @@ async fn turn_steer_returns_active_turn_id() -> Result<()> { text: "steer".to_string(), text_elements: Vec::new(), }], + responsesapi_client_metadata: None, expected_turn_id: turn.id.clone(), }) .await?; diff --git a/codex-rs/core/src/agent/control_tests.rs b/codex-rs/core/src/agent/control_tests.rs index 54693228f..a2c2b2279 100644 --- a/codex-rs/core/src/agent/control_tests.rs +++ b/codex-rs/core/src/agent/control_tests.rs @@ -425,6 +425,7 @@ async fn send_input_submits_user_message() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }, ); let captured = harness @@ -571,6 +572,7 @@ async fn spawn_agent_creates_thread_and_sends_prompt() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }, ); let captured = harness @@ -678,6 +680,7 @@ async fn spawn_agent_can_fork_parent_thread_history_with_sanitized_items() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }, ); let captured = harness diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 6375f2a83..01d7e068c 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -763,8 +763,11 @@ impl Codex { &self, input: Vec, expected_turn_id: Option<&str>, + responsesapi_client_metadata: Option>, ) -> Result { - self.session.steer_input(input, expected_turn_id).await + self.session + .steer_input(input, expected_turn_id, responsesapi_client_metadata) + .await } pub(crate) async fn set_app_server_client_info( @@ -2264,6 +2267,7 @@ impl Session { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }, ) .await; @@ -4136,6 +4140,7 @@ impl Session { &self, input: Vec, expected_turn_id: Option<&str>, + responsesapi_client_metadata: Option>, ) -> Result { if input.is_empty() { return Err(SteerInputError::EmptyInput); @@ -4174,6 +4179,15 @@ impl Session { None => return Err(SteerInputError::NoActiveTurn(input)), } + if let Some(responsesapi_client_metadata) = responsesapi_client_metadata + && let Some((_, active_task)) = active_turn.tasks.first() + { + active_task + .turn_context + .turn_metadata_state + .set_responsesapi_client_metadata(responsesapi_client_metadata); + } + let mut turn_state = active_turn.turn_state.lock().await; turn_state.push_pending_input(input.into()); turn_state.accept_mailbox_delivery_for_current_turn(); @@ -4942,7 +4956,7 @@ mod handlers { } pub async fn user_input_or_turn(sess: &Arc, sub_id: String, op: Op) { - let (items, updates) = match op { + let (items, updates, responsesapi_client_metadata) = match op { Op::UserTurn { cwd, approval_policy, @@ -4983,17 +4997,20 @@ mod handlers { app_server_client_name: None, app_server_client_version: None, }, + None, ) } Op::UserInput { items, final_output_json_schema, + responsesapi_client_metadata, } => ( items, SessionSettingsUpdate { final_output_json_schema: Some(final_output_json_schema), ..Default::default() }, + responsesapi_client_metadata, ), _ => unreachable!(), }; @@ -5005,11 +5022,20 @@ mod handlers { sess.maybe_emit_unknown_model_warning_for_turn(current_context.as_ref()) .await; match sess - .steer_input(items.clone(), /*expected_turn_id*/ None) + .steer_input( + items.clone(), + /*expected_turn_id*/ None, + responsesapi_client_metadata.clone(), + ) .await { Ok(_) => current_context.session_telemetry.user_prompt(&items), Err(SteerInputError::NoActiveTurn(items)) => { + if let Some(responsesapi_client_metadata) = responsesapi_client_metadata { + current_context + .turn_metadata_state + .set_responsesapi_client_metadata(responsesapi_client_metadata); + } current_context.session_telemetry.user_prompt(&items); sess.refresh_mcp_servers_if_requested(¤t_context) .await; diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index aef0a92d7..3d74773ba 100644 --- a/codex-rs/core/src/codex_delegate.rs +++ b/codex-rs/core/src/codex_delegate.rs @@ -185,6 +185,7 @@ pub(crate) async fn run_codex_thread_one_shot( io.submit(Op::UserInput { items: input, final_output_json_schema, + responsesapi_client_metadata: None, }) .await?; diff --git a/codex-rs/core/src/codex_tests.rs b/codex-rs/core/src/codex_tests.rs index fdc83f31b..6eb66a519 100644 --- a/codex-rs/core/src/codex_tests.rs +++ b/codex-rs/core/src/codex_tests.rs @@ -1331,6 +1331,7 @@ async fn fork_startup_context_then_first_turn_diff_snapshot() -> anyhow::Result< text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1386,6 +1387,7 @@ async fn fork_startup_context_then_first_turn_diff_snapshot() -> anyhow::Result< text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&forked.thread, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -3279,6 +3281,7 @@ fn op_kind_distinguishes_turn_ops() { Op::UserInput { items: vec![], final_output_json_schema: None, + responsesapi_client_metadata: None, } .kind(), "user_input" @@ -4857,7 +4860,9 @@ async fn steer_input_requires_active_turn() { }]; let err = sess - .steer_input(input, /*expected_turn_id*/ None) + .steer_input( + input, /*expected_turn_id*/ None, /*responsesapi_client_metadata*/ None, + ) .await .expect_err("steering without active turn should fail"); @@ -4886,7 +4891,11 @@ async fn steer_input_enforces_expected_turn_id() { text_elements: Vec::new(), }]; let err = sess - .steer_input(steer_input, Some("different-turn-id")) + .steer_input( + steer_input, + Some("different-turn-id"), + /*responsesapi_client_metadata*/ None, + ) .await .expect_err("mismatched expected turn id should fail"); @@ -4928,7 +4937,11 @@ async fn steer_input_rejects_non_regular_turns() { text_elements: Vec::new(), }]; let err = sess - .steer_input(steer_input, /*expected_turn_id*/ None) + .steer_input( + steer_input, + /*expected_turn_id*/ None, + /*responsesapi_client_metadata*/ None, + ) .await .expect_err("steering a non-regular turn should fail"); @@ -4960,7 +4973,11 @@ async fn steer_input_returns_active_turn_id() { text_elements: Vec::new(), }]; let turn_id = sess - .steer_input(steer_input, Some(&tc.sub_id)) + .steer_input( + steer_input, + Some(&tc.sub_id), + /*responsesapi_client_metadata*/ None, + ) .await .expect("steering with matching expected turn id should succeed"); @@ -5147,6 +5164,7 @@ async fn steered_input_reopens_mailbox_delivery_for_current_turn() { text_elements: Vec::new(), }], Some(&tc.sub_id), + /*responsesapi_client_metadata*/ None, ) .await .expect("steered input should be accepted"); @@ -5191,6 +5209,7 @@ async fn stale_defer_mailbox_delivery_does_not_override_steered_input() { text_elements: Vec::new(), }], Some(&tc.sub_id), + /*responsesapi_client_metadata*/ None, ) .await .expect("steered input should be accepted"); diff --git a/codex-rs/core/src/codex_thread.rs b/codex-rs/core/src/codex_thread.rs index 9727cc208..f95c615a6 100644 --- a/codex-rs/core/src/codex_thread.rs +++ b/codex-rs/core/src/codex_thread.rs @@ -23,6 +23,7 @@ use codex_protocol::protocol::TokenUsage; use codex_protocol::protocol::W3cTraceContext; use codex_protocol::user_input::UserInput; use rmcp::model::ReadResourceRequestParams; +use std::collections::HashMap; use std::path::PathBuf; use tokio::sync::Mutex; use tokio::sync::watch; @@ -97,8 +98,11 @@ impl CodexThread { &self, input: Vec, expected_turn_id: Option<&str>, + responsesapi_client_metadata: Option>, ) -> Result { - self.codex.steer_input(input, expected_turn_id).await + self.codex + .steer_input(input, expected_turn_id, responsesapi_client_metadata) + .await } pub async fn set_app_server_client_info( diff --git a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs index bcb47f3b6..ca74b29d9 100644 --- a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs +++ b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs @@ -2004,6 +2004,7 @@ async fn send_input_accepts_structured_items() { }, ], final_output_json_schema: None, + responsesapi_client_metadata: None, }; let captured = manager .captured_ops() diff --git a/codex-rs/core/src/turn_metadata.rs b/codex-rs/core/src/turn_metadata.rs index b3794d325..e4385096a 100644 --- a/codex-rs/core/src/turn_metadata.rs +++ b/codex-rs/core/src/turn_metadata.rs @@ -1,4 +1,5 @@ use std::collections::BTreeMap; +use std::collections::HashMap; use std::path::Path; use std::path::PathBuf; use std::sync::Arc; @@ -6,6 +7,7 @@ use std::sync::Mutex; use std::sync::RwLock; use serde::Serialize; +use serde_json::Value; use tokio::task::JoinHandle; use crate::sandbox_tags::sandbox_tag; @@ -69,6 +71,20 @@ impl TurnMetadataBag { } } +fn merge_responsesapi_client_metadata( + header: &str, + responsesapi_client_metadata: Option<&HashMap>, +) -> Option { + let responsesapi_client_metadata = responsesapi_client_metadata?; + let mut metadata = serde_json::from_str::>(header).ok()?; + for (key, value) in responsesapi_client_metadata { + metadata + .entry(key.clone()) + .or_insert_with(|| Value::String(value.clone())); + } + serde_json::to_string(&metadata).ok() +} + fn build_turn_metadata_bag( session_id: Option, turn_id: Option, @@ -129,6 +145,7 @@ pub(crate) struct TurnMetadataState { base_metadata: TurnMetadataBag, base_header: String, enriched_header: Arc>>, + responsesapi_client_metadata: Arc>>>, enrichment_task: Arc>>>, } @@ -159,21 +176,30 @@ impl TurnMetadataState { base_metadata, base_header, enriched_header: Arc::new(RwLock::new(None)), + responsesapi_client_metadata: Arc::new(RwLock::new(None)), enrichment_task: Arc::new(Mutex::new(None)), } } pub(crate) fn current_header_value(&self) -> Option { - if let Some(header) = self + let header = if let Some(header) = self .enriched_header .read() .unwrap_or_else(std::sync::PoisonError::into_inner) .as_ref() .cloned() { - return Some(header); - } - Some(self.base_header.clone()) + header + } else { + self.base_header.clone() + }; + let responsesapi_client_metadata = self + .responsesapi_client_metadata + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone(); + merge_responsesapi_client_metadata(&header, responsesapi_client_metadata.as_ref()) + .or(Some(header)) } pub(crate) fn current_meta_value(&self) -> Option { @@ -181,6 +207,17 @@ impl TurnMetadataState { .and_then(|header| serde_json::from_str(&header).ok()) } + pub(crate) fn set_responsesapi_client_metadata( + &self, + responsesapi_client_metadata: HashMap, + ) { + *self + .responsesapi_client_metadata + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner) = + Some(responsesapi_client_metadata); + } + pub(crate) fn spawn_git_enrichment_task(&self) { if self.repo_root.is_none() { return; diff --git a/codex-rs/core/src/turn_metadata_tests.rs b/codex-rs/core/src/turn_metadata_tests.rs index 5da26563f..dfd9322a9 100644 --- a/codex-rs/core/src/turn_metadata_tests.rs +++ b/codex-rs/core/src/turn_metadata_tests.rs @@ -1,6 +1,7 @@ use super::*; use serde_json::Value; +use std::collections::HashMap; use tempfile::TempDir; use tokio::process::Command; @@ -83,3 +84,29 @@ fn turn_metadata_state_uses_platform_sandbox_tag() { assert_eq!(sandbox_name, Some(expected_sandbox)); assert_eq!(session_id, Some("session-a")); } + +#[test] +fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields() { + let temp_dir = TempDir::new().expect("temp dir"); + let cwd = temp_dir.path().to_path_buf(); + let sandbox_policy = SandboxPolicy::new_read_only_policy(); + + let state = TurnMetadataState::new( + "session-a".to_string(), + "turn-a".to_string(), + cwd, + &sandbox_policy, + WindowsSandboxLevel::Disabled, + ); + state.set_responsesapi_client_metadata(HashMap::from([ + ("fiber_run_id".to_string(), "fiber-123".to_string()), + ("session_id".to_string(), "client-supplied".to_string()), + ])); + + let header = state.current_header_value().expect("header"); + let json: Value = serde_json::from_str(&header).expect("json"); + + assert_eq!(json["fiber_run_id"].as_str(), Some("fiber-123")); + assert_eq!(json["session_id"].as_str(), Some("session-a")); + assert_eq!(json["turn_id"].as_str(), Some("turn-a")); +} diff --git a/codex-rs/core/tests/suite/abort_tasks.rs b/codex-rs/core/tests/suite/abort_tasks.rs index af3c70b74..4911d8dd6 100644 --- a/codex-rs/core/tests/suite/abort_tasks.rs +++ b/codex-rs/core/tests/suite/abort_tasks.rs @@ -51,6 +51,7 @@ async fn interrupt_long_running_tool_emits_turn_aborted() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -105,6 +106,7 @@ async fn interrupt_tool_records_history_entries() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -123,6 +125,7 @@ async fn interrupt_tool_records_history_entries() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -203,6 +206,7 @@ async fn interrupt_persists_turn_aborted_marker_in_next_request() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -221,6 +225,7 @@ async fn interrupt_persists_turn_aborted_marker_in_next_request() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index 5200d252c..a2a110b2d 100644 --- a/codex-rs/core/tests/suite/client.rs +++ b/codex-rs/core/tests/suite/client.rs @@ -387,6 +387,7 @@ async fn resume_includes_initial_messages_and_sends_prior_items() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -749,6 +750,7 @@ async fn includes_conversation_id_and_model_headers_in_request() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -947,6 +949,7 @@ async fn includes_base_instructions_override_in_request() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1000,6 +1003,7 @@ async fn chatgpt_auth_sends_correct_request() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1112,6 +1116,7 @@ async fn prefers_apikey_when_config_prefers_apikey_even_with_chatgpt_tokens() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1148,6 +1153,7 @@ async fn includes_user_instructions_message_in_request() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1233,6 +1239,7 @@ async fn includes_apps_guidance_as_developer_message_for_chatgpt_auth() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1293,6 +1300,7 @@ async fn omits_apps_guidance_for_api_key_auth_even_when_feature_enabled() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1349,6 +1357,7 @@ async fn omits_apps_guidance_when_configured_off() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1388,6 +1397,7 @@ async fn omits_environment_context_when_configured_off() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1442,6 +1452,7 @@ async fn skills_append_to_developer_message() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1493,6 +1504,7 @@ async fn includes_configured_effort_in_request() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1535,6 +1547,7 @@ async fn includes_no_effort_in_request() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1575,6 +1588,7 @@ async fn includes_default_reasoning_effort_in_request_when_defined_by_model_info text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1687,6 +1701,7 @@ async fn configured_reasoning_summary_is_sent() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1802,6 +1817,7 @@ async fn reasoning_summary_is_omitted_when_disabled() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1858,6 +1874,7 @@ async fn reasoning_summary_none_overrides_model_catalog_default() -> anyhow::Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1894,6 +1911,7 @@ async fn includes_default_verbosity_in_request() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1939,6 +1957,7 @@ async fn configured_verbosity_not_sent_for_models_without_support() -> anyhow::R text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1983,6 +2002,7 @@ async fn configured_verbosity_is_sent() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2032,6 +2052,7 @@ async fn includes_developer_instructions_message_in_request() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2322,6 +2343,7 @@ async fn token_count_includes_rate_limits_snapshot() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2487,6 +2509,7 @@ async fn usage_limit_error_emits_rate_limit_event() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submission should succeed while emitting usage limit error events"); @@ -2561,6 +2584,7 @@ async fn context_window_error_sets_total_tokens_to_model_window() -> anyhow::Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -2573,6 +2597,7 @@ async fn context_window_error_sets_total_tokens_to_model_window() -> anyhow::Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -2655,6 +2680,7 @@ async fn incomplete_response_emits_content_filter_error_message() -> anyhow::Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -2762,6 +2788,7 @@ async fn azure_overrides_assign_properties_used_for_responses_url() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2847,6 +2874,7 @@ async fn env_var_overrides_loaded_auth() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2908,6 +2936,7 @@ async fn history_dedupes_streamed_and_final_messages_across_turns() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2921,6 +2950,7 @@ async fn history_dedupes_streamed_and_final_messages_across_turns() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2934,6 +2964,7 @@ async fn history_dedupes_streamed_and_final_messages_across_turns() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/client_websockets.rs b/codex-rs/core/tests/suite/client_websockets.rs index 9a76638f6..424f6764e 100755 --- a/codex-rs/core/tests/suite/client_websockets.rs +++ b/codex-rs/core/tests/suite/client_websockets.rs @@ -1003,6 +1003,7 @@ async fn responses_websocket_usage_limit_error_emits_rate_limit_event() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submission should succeed while emitting usage limit error events"); @@ -1088,6 +1089,7 @@ async fn responses_websocket_invalid_request_error_with_status_is_forwarded() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submission should succeed while emitting invalid request events"); @@ -1265,6 +1267,56 @@ async fn responses_websocket_forwards_turn_metadata_on_initial_and_incremental_c server.shutdown().await; } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn responses_websocket_preserves_custom_turn_metadata_fields() { + skip_if_no_network!(); + + let server = start_websocket_server(vec![vec![vec![ + ev_response_created("resp-1"), + ev_completed("resp-1"), + ]]]) + .await; + + let harness = websocket_harness(&server).await; + let mut client_session = harness.client.new_session(); + let prompt = prompt_with_input(vec![message_item("hello")]); + let turn_metadata = json!({ + "turn_id": "turn-123", + "fiber_run_id": "fiber-123", + "origin": "app-server", + }) + .to_string(); + + stream_until_complete_with_turn_metadata( + &mut client_session, + &harness, + &prompt, + /*service_tier*/ None, + Some(&turn_metadata), + ) + .await; + + let body = server + .single_connection() + .first() + .expect("missing request") + .body_json(); + + assert_eq!(body["type"].as_str(), Some("response.create")); + assert_eq!( + body["client_metadata"]["x-codex-turn-metadata"] + .as_str() + .map(|value| serde_json::from_str::(value).expect("valid json")), + Some(json!({ + "turn_id": "turn-123", + "fiber_run_id": "fiber-123", + "origin": "app-server", + })) + ); + + server.shutdown().await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn responses_websocket_uses_previous_response_id_when_prefix_after_completed() { skip_if_no_network!(); @@ -1817,6 +1869,23 @@ async fn stream_until_complete_with_turn_metadata( prompt: &Prompt, service_tier: Option, turn_metadata_header: Option<&str>, +) { + stream_until_complete_with_request_metadata( + client_session, + harness, + prompt, + service_tier, + turn_metadata_header, + ) + .await; +} + +async fn stream_until_complete_with_request_metadata( + client_session: &mut ModelClientSession, + harness: &WebsocketTestHarness, + prompt: &Prompt, + service_tier: Option, + turn_metadata_header: Option<&str>, ) { let mut stream = client_session .stream( diff --git a/codex-rs/core/tests/suite/collaboration_instructions.rs b/codex-rs/core/tests/suite/collaboration_instructions.rs index df8d6e4fe..83c0e1399 100644 --- a/codex-rs/core/tests/suite/collaboration_instructions.rs +++ b/codex-rs/core/tests/suite/collaboration_instructions.rs @@ -84,6 +84,7 @@ async fn no_collaboration_instructions_by_default() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -143,6 +144,7 @@ async fn user_input_includes_collaboration_instructions_after_override() -> Resu text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -241,6 +243,7 @@ async fn override_then_next_turn_uses_updated_collaboration_instructions() -> Re text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -364,6 +367,7 @@ async fn collaboration_mode_update_emits_new_instruction_message() -> Result<()> text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -391,6 +395,7 @@ async fn collaboration_mode_update_emits_new_instruction_message() -> Result<()> text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -447,6 +452,7 @@ async fn collaboration_mode_update_noop_does_not_append() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -474,6 +480,7 @@ async fn collaboration_mode_update_noop_does_not_append() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -532,6 +539,7 @@ async fn collaboration_mode_update_emits_new_instruction_message_when_mode_chang text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -562,6 +570,7 @@ async fn collaboration_mode_update_emits_new_instruction_message_when_mode_chang text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -621,6 +630,7 @@ async fn collaboration_mode_update_noop_does_not_append_when_mode_is_unchanged() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -651,6 +661,7 @@ async fn collaboration_mode_update_noop_does_not_append_when_mode_is_unchanged() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -714,6 +725,7 @@ async fn resume_replays_collaboration_instructions() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -727,6 +739,7 @@ async fn resume_replays_collaboration_instructions() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -783,6 +796,7 @@ async fn empty_collaboration_instructions_are_ignored() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; diff --git a/codex-rs/core/tests/suite/compact.rs b/codex-rs/core/tests/suite/compact.rs index df8fad124..dc5f127df 100644 --- a/codex-rs/core/tests/suite/compact.rs +++ b/codex-rs/core/tests/suite/compact.rs @@ -245,6 +245,7 @@ async fn summarize_context_three_requests_and_instructions() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -267,6 +268,7 @@ async fn summarize_context_three_requests_and_instructions() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -443,6 +445,7 @@ async fn manual_compact_uses_custom_prompt() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit first user turn"); @@ -587,6 +590,7 @@ async fn manual_compact_emits_context_compaction_items() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -750,6 +754,7 @@ async fn multiple_auto_compact_per_task_runs_after_token_limit_hit() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit user input"); @@ -1249,6 +1254,7 @@ async fn auto_compact_runs_after_token_limit_hit() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1262,6 +1268,7 @@ async fn auto_compact_runs_after_token_limit_hit() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1275,6 +1282,7 @@ async fn auto_compact_runs_after_token_limit_hit() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1443,6 +1451,7 @@ async fn auto_compact_emits_context_compaction_items() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1521,6 +1530,7 @@ async fn auto_compact_starts_after_turn_started() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1533,6 +1543,7 @@ async fn auto_compact_starts_after_turn_started() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1545,6 +1556,7 @@ async fn auto_compact_starts_after_turn_started() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2037,6 +2049,7 @@ async fn auto_compact_persists_rollout_entries() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2049,6 +2062,7 @@ async fn auto_compact_persists_rollout_entries() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2061,6 +2075,7 @@ async fn auto_compact_persists_rollout_entries() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2147,6 +2162,7 @@ async fn manual_compact_retries_after_context_window_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2258,6 +2274,7 @@ async fn manual_compact_non_context_failure_retries_then_emits_task_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit user input"); @@ -2350,6 +2367,7 @@ async fn manual_compact_twice_preserves_latest_user_messages() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2365,6 +2383,7 @@ async fn manual_compact_twice_preserves_latest_user_messages() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2380,6 +2399,7 @@ async fn manual_compact_twice_preserves_latest_user_messages() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2541,6 +2561,7 @@ async fn auto_compact_allows_multiple_attempts_when_interleaved_with_other_turn_ text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2643,6 +2664,7 @@ async fn snapshot_request_shape_mid_turn_continuation_compaction() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2839,6 +2861,7 @@ async fn auto_compact_counts_encrypted_reasoning_before_last_user() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -2956,6 +2979,7 @@ async fn auto_compact_runs_when_reasoning_header_clears_between_turns() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -3015,6 +3039,7 @@ async fn snapshot_request_shape_pre_turn_compaction_including_incoming_user_mess text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit user input"); @@ -3050,6 +3075,7 @@ async fn snapshot_request_shape_pre_turn_compaction_including_incoming_user_mess }, ], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit user input"); @@ -3260,6 +3286,7 @@ async fn snapshot_request_shape_pre_turn_compaction_context_window_exceeded() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit first user"); @@ -3272,6 +3299,7 @@ async fn snapshot_request_shape_pre_turn_compaction_context_window_exceeded() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit second user"); @@ -3342,6 +3370,7 @@ async fn snapshot_request_shape_manual_compact_without_previous_user_messages() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit follow-up user input"); diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index d8015812a..a2681423f 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -237,6 +237,7 @@ async fn remote_compact_replaces_history_for_followups() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -251,6 +252,7 @@ async fn remote_compact_replaces_history_for_followups() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -381,6 +383,7 @@ async fn remote_compact_runs_automatically() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -455,6 +458,7 @@ async fn remote_compact_trims_function_call_history_to_fit_context_window() -> R text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -466,6 +470,7 @@ async fn remote_compact_trims_function_call_history_to_fit_context_window() -> R text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -581,6 +586,7 @@ async fn auto_remote_compact_trims_function_call_history_to_fit_context_window() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -592,6 +598,7 @@ async fn auto_remote_compact_trims_function_call_history_to_fit_context_window() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -609,6 +616,7 @@ async fn auto_remote_compact_trims_function_call_history_to_fit_context_window() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -706,6 +714,7 @@ async fn auto_remote_compact_failure_stops_agent_loop() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -717,6 +726,7 @@ async fn auto_remote_compact_failure_stops_agent_loop() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -808,6 +818,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&baseline_codex, |event| { @@ -822,6 +833,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&baseline_codex, |event| { @@ -907,6 +919,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&override_codex, |event| { @@ -921,6 +934,7 @@ async fn remote_compact_trim_estimate_uses_session_base_instructions() -> Result text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&override_codex, |event| { @@ -989,6 +1003,7 @@ async fn remote_manual_compact_emits_context_compaction_items() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -1067,6 +1082,7 @@ async fn remote_manual_compact_failure_emits_task_error_event() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -1149,6 +1165,7 @@ async fn remote_compact_persists_replacement_history_in_rollout() -> Result<()> text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1290,6 +1307,7 @@ async fn remote_compact_and_resume_refresh_stale_developer_instructions() -> Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1305,6 +1323,7 @@ async fn remote_compact_and_resume_refresh_stale_developer_instructions() -> Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1327,6 +1346,7 @@ async fn remote_compact_and_resume_refresh_stale_developer_instructions() -> Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1421,6 +1441,7 @@ async fn remote_compact_refreshes_stale_developer_instructions_without_resume() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1435,6 +1456,7 @@ async fn remote_compact_refreshes_stale_developer_instructions_without_resume() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1504,6 +1526,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_sta text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1515,6 +1538,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_sta text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1579,6 +1603,7 @@ async fn remote_request_uses_custom_experimental_realtime_start_instructions() - text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1637,6 +1662,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_end text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1650,6 +1676,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_restates_realtime_end text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1722,6 +1749,7 @@ async fn snapshot_request_shape_remote_manual_compact_restates_realtime_start() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1736,6 +1764,7 @@ async fn snapshot_request_shape_remote_manual_compact_restates_realtime_start() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1816,6 +1845,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_does_not_restate_real text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1829,6 +1859,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_does_not_restate_real text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1917,6 +1948,7 @@ async fn snapshot_request_shape_remote_compact_resume_restates_realtime_end() -> text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -1944,6 +1976,7 @@ async fn snapshot_request_shape_remote_compact_resume_restates_realtime_end() -> text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2037,6 +2070,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_including_incoming_us text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2121,6 +2155,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_strips_incoming_model text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2147,6 +2182,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_strips_incoming_model text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2263,6 +2299,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_context_window_exceed text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2274,6 +2311,7 @@ async fn snapshot_request_shape_remote_pre_turn_compaction_context_window_exceed text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; let error_message = wait_for_event_match(&codex, |event| match event { @@ -2356,6 +2394,7 @@ async fn snapshot_request_shape_remote_mid_turn_continuation_compaction() -> Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2431,6 +2470,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_summary_only_reinject text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2514,6 +2554,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_multi_summary_reinjec text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2528,6 +2569,7 @@ async fn snapshot_request_shape_remote_mid_turn_compaction_multi_summary_reinjec text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -2607,6 +2649,7 @@ async fn snapshot_request_shape_remote_manual_compact_without_previous_user_mess text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; diff --git a/codex-rs/core/tests/suite/compact_resume_fork.rs b/codex-rs/core/tests/suite/compact_resume_fork.rs index 8f4f2ae75..3b63e53af 100644 --- a/codex-rs/core/tests/suite/compact_resume_fork.rs +++ b/codex-rs/core/tests/suite/compact_resume_fork.rs @@ -808,6 +808,7 @@ async fn user_turn(conversation: &Arc, text: &str) { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .expect("submit user turn"); diff --git a/codex-rs/core/tests/suite/fork_thread.rs b/codex-rs/core/tests/suite/fork_thread.rs index e24cb74d7..dc7f2151f 100644 --- a/codex-rs/core/tests/suite/fork_thread.rs +++ b/codex-rs/core/tests/suite/fork_thread.rs @@ -53,6 +53,7 @@ async fn fork_thread_twice_drops_to_first_message() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/hooks.rs b/codex-rs/core/tests/suite/hooks.rs index c6e8a48d2..47cd251bb 100644 --- a/codex-rs/core/tests/suite/hooks.rs +++ b/codex-rs/core/tests/suite/hooks.rs @@ -904,6 +904,7 @@ async fn blocked_queued_prompt_does_not_strand_earlier_accepted_prompt() -> Resu text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -920,6 +921,7 @@ async fn blocked_queued_prompt_does_not_strand_earlier_accepted_prompt() -> Resu text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; } diff --git a/codex-rs/core/tests/suite/items.rs b/codex-rs/core/tests/suite/items.rs index f949cd4b9..314452cb2 100644 --- a/codex-rs/core/tests/suite/items.rs +++ b/codex-rs/core/tests/suite/items.rs @@ -86,6 +86,7 @@ async fn user_message_item_is_emitted() -> anyhow::Result<()> { .submit(Op::UserInput { items: vec![expected_input.clone()], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -142,6 +143,7 @@ async fn assistant_message_item_is_emitted() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -200,6 +202,7 @@ async fn reasoning_item_is_emitted() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -259,6 +262,7 @@ async fn web_search_item_is_emitted() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -323,6 +327,7 @@ async fn image_generation_call_event_is_emitted() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -385,6 +390,7 @@ async fn image_generation_call_event_is_emitted_when_image_save_fails() -> anyho text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -438,6 +444,7 @@ async fn agent_message_content_delta_has_item_metadata() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -1085,6 +1092,7 @@ async fn reasoning_content_delta_has_item_metadata() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -1144,6 +1152,7 @@ async fn reasoning_raw_content_delta_respects_flag() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; diff --git a/codex-rs/core/tests/suite/model_visible_layout.rs b/codex-rs/core/tests/suite/model_visible_layout.rs index 13e281914..dbf977aae 100644 --- a/codex-rs/core/tests/suite/model_visible_layout.rs +++ b/codex-rs/core/tests/suite/model_visible_layout.rs @@ -322,6 +322,7 @@ async fn snapshot_model_visible_layout_resume_with_personality_change() -> Resul text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -421,6 +422,7 @@ async fn snapshot_model_visible_layout_resume_override_matches_rollout_model() - text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -466,6 +468,7 @@ async fn snapshot_model_visible_layout_resume_override_matches_rollout_model() - text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |event| { diff --git a/codex-rs/core/tests/suite/otel.rs b/codex-rs/core/tests/suite/otel.rs index 96df3fd2f..d771ed19a 100644 --- a/codex-rs/core/tests/suite/otel.rs +++ b/codex-rs/core/tests/suite/otel.rs @@ -107,6 +107,7 @@ async fn responses_api_emits_api_request_event() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -150,6 +151,7 @@ async fn process_sse_emits_tracing_for_output_item() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -193,6 +195,7 @@ async fn process_sse_emits_failed_event_on_parse_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -237,6 +240,7 @@ async fn process_sse_records_failed_event_when_stream_closes_without_completed() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -301,6 +305,7 @@ async fn process_sse_failed_event_records_response_error_message() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -363,6 +368,7 @@ async fn process_sse_failed_event_logs_parse_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -412,6 +418,7 @@ async fn process_sse_failed_event_logs_missing_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -470,6 +477,7 @@ async fn process_sse_failed_event_logs_response_completed_parse_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -522,6 +530,7 @@ async fn process_sse_emits_completed_telemetry() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -594,6 +603,7 @@ async fn handle_responses_span_records_response_kind_and_tool_name() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -678,6 +688,7 @@ async fn record_responses_sets_span_fields_for_response_events() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -762,6 +773,7 @@ async fn handle_response_item_records_tool_result_for_custom_tool_call() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -835,6 +847,7 @@ async fn handle_response_item_records_tool_result_for_function_call() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -918,6 +931,7 @@ async fn handle_response_item_records_tool_result_for_local_shell_missing_ids() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -986,6 +1000,7 @@ async fn handle_response_item_records_tool_result_for_local_shell_call() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1094,6 +1109,7 @@ async fn handle_container_exec_autoapprove_from_config_records_tool_decision() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1145,6 +1161,7 @@ async fn handle_container_exec_user_approved_records_tool_decision() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1211,6 +1228,7 @@ async fn handle_container_exec_user_approved_for_session_records_tool_decision() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1277,6 +1295,7 @@ async fn handle_sandbox_error_user_approves_retry_records_tool_decision() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1343,6 +1362,7 @@ async fn handle_container_exec_user_denies_records_tool_decision() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1409,6 +1429,7 @@ async fn handle_sandbox_error_user_approves_for_session_records_tool_decision() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -1476,6 +1497,7 @@ async fn handle_sandbox_error_user_denies_records_tool_decision() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/pending_input.rs b/codex-rs/core/tests/suite/pending_input.rs index 08596522e..406d6a3d5 100644 --- a/codex-rs/core/tests/suite/pending_input.rs +++ b/codex-rs/core/tests/suite/pending_input.rs @@ -100,6 +100,7 @@ async fn submit_user_input(codex: &CodexThread, text: &str) { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap_or_else(|err| panic!("submit user input: {err}")); @@ -136,6 +137,7 @@ async fn steer_user_input(codex: &CodexThread, text: &str) { text_elements: Vec::new(), }], /*expected_turn_id*/ None, + /*responsesapi_client_metadata*/ None, ) .await .unwrap_or_else(|err| panic!("steer user input: {err:?}")); @@ -267,6 +269,7 @@ async fn injected_user_input_triggers_follow_up_request_with_deltas() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -283,6 +286,7 @@ async fn injected_user_input_triggers_follow_up_request_with_deltas() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/permissions_messages.rs b/codex-rs/core/tests/suite/permissions_messages.rs index d1738db90..ec6b7d67d 100644 --- a/codex-rs/core/tests/suite/permissions_messages.rs +++ b/codex-rs/core/tests/suite/permissions_messages.rs @@ -53,6 +53,7 @@ async fn permissions_message_sent_once_on_start() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -90,6 +91,7 @@ async fn permissions_message_added_on_override_change() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -117,6 +119,7 @@ async fn permissions_message_added_on_override_change() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -160,6 +163,7 @@ async fn permissions_message_not_added_when_no_change() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -171,6 +175,7 @@ async fn permissions_message_not_added_when_no_change() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -214,6 +219,7 @@ async fn permissions_message_omitted_when_disabled() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -241,6 +247,7 @@ async fn permissions_message_omitted_when_disabled() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -297,6 +304,7 @@ async fn resume_replays_permissions_messages() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -326,6 +334,7 @@ async fn resume_replays_permissions_messages() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -339,6 +348,7 @@ async fn resume_replays_permissions_messages() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -396,6 +406,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -425,6 +436,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&initial.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -444,6 +456,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -476,6 +489,7 @@ async fn resume_and_fork_append_permissions_messages() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&forked.thread, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -528,6 +542,7 @@ async fn permissions_message_includes_writable_roots() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&test.codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; diff --git a/codex-rs/core/tests/suite/plugins.rs b/codex-rs/core/tests/suite/plugins.rs index cd0ad330a..ab2c85ac0 100644 --- a/codex-rs/core/tests/suite/plugins.rs +++ b/codex-rs/core/tests/suite/plugins.rs @@ -213,6 +213,7 @@ async fn capability_sections_render_in_developer_message_in_order() -> Result<() text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -288,6 +289,7 @@ async fn explicit_plugin_mentions_inject_plugin_guidance() -> Result<()> { path: format!("plugin://{SAMPLE_PLUGIN_CONFIG_NAME}"), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -360,6 +362,7 @@ async fn explicit_plugin_mentions_track_plugin_used_analytics() -> Result<()> { path: format!("plugin://{SAMPLE_PLUGIN_CONFIG_NAME}"), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; diff --git a/codex-rs/core/tests/suite/prompt_caching.rs b/codex-rs/core/tests/suite/prompt_caching.rs index 4d34e691d..4bcf91417 100644 --- a/codex-rs/core/tests/suite/prompt_caching.rs +++ b/codex-rs/core/tests/suite/prompt_caching.rs @@ -151,6 +151,7 @@ async fn prompt_tools_are_consistent_across_requests() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -162,6 +163,7 @@ async fn prompt_tools_are_consistent_across_requests() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -247,6 +249,7 @@ async fn gpt_5_tools_without_apply_patch_append_apply_patch_instructions() -> an text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -258,6 +261,7 @@ async fn gpt_5_tools_without_apply_patch_append_apply_patch_instructions() -> an text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -320,6 +324,7 @@ async fn prefixes_context_and_instructions_once_and_consistently_across_requests text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -331,6 +336,7 @@ async fn prefixes_context_and_instructions_once_and_consistently_across_requests text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -412,6 +418,7 @@ async fn overrides_turn_context_but_keeps_cached_prefix_and_key_constant() -> an text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -448,6 +455,7 @@ async fn overrides_turn_context_but_keeps_cached_prefix_and_key_constant() -> an text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; @@ -530,6 +538,7 @@ async fn override_before_first_turn_emits_environment_context() -> anyhow::Resul text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -681,6 +690,7 @@ async fn per_turn_overrides_keep_cached_prefix_and_key_constant() -> anyhow::Res text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; diff --git a/codex-rs/core/tests/suite/quota_exceeded.rs b/codex-rs/core/tests/suite/quota_exceeded.rs index 466680e6a..c59f2d86b 100644 --- a/codex-rs/core/tests/suite/quota_exceeded.rs +++ b/codex-rs/core/tests/suite/quota_exceeded.rs @@ -46,6 +46,7 @@ async fn quota_exceeded_emits_single_error_event() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/realtime_conversation.rs b/codex-rs/core/tests/suite/realtime_conversation.rs index 0f68af428..a8f7f10fe 100644 --- a/codex-rs/core/tests/suite/realtime_conversation.rs +++ b/codex-rs/core/tests/suite/realtime_conversation.rs @@ -2650,6 +2650,7 @@ async fn inbound_handoff_request_steers_active_turn() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; diff --git a/codex-rs/core/tests/suite/request_compression.rs b/codex-rs/core/tests/suite/request_compression.rs index fb0125e66..44d37f311 100644 --- a/codex-rs/core/tests/suite/request_compression.rs +++ b/codex-rs/core/tests/suite/request_compression.rs @@ -45,6 +45,7 @@ async fn request_body_is_zstd_compressed_for_codex_backend_when_enabled() -> any text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -92,6 +93,7 @@ async fn request_body_is_not_compressed_for_api_key_auth_even_when_enabled() -> text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; diff --git a/codex-rs/core/tests/suite/resume.rs b/codex-rs/core/tests/suite/resume.rs index a77ad60b5..11b575077 100644 --- a/codex-rs/core/tests/suite/resume.rs +++ b/codex-rs/core/tests/suite/resume.rs @@ -91,6 +91,7 @@ async fn resume_includes_initial_messages_from_rollout_events() -> Result<()> { text_elements: text_elements.clone(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -176,6 +177,7 @@ async fn resume_includes_initial_messages_from_reasoning_events() -> Result<()> text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; @@ -265,6 +267,7 @@ async fn resume_switches_models_preserves_base_instructions() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -305,6 +308,7 @@ async fn resume_switches_models_preserves_base_instructions() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |event| { @@ -320,6 +324,7 @@ async fn resume_switches_models_preserves_base_instructions() -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |event| { @@ -390,6 +395,7 @@ async fn resume_model_switch_is_not_duplicated_after_pre_turn_override() -> Resu text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; @@ -433,6 +439,7 @@ async fn resume_model_switch_is_not_duplicated_after_pre_turn_override() -> Resu text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&resumed.codex, |event| { diff --git a/codex-rs/core/tests/suite/review.rs b/codex-rs/core/tests/suite/review.rs index 503bc8745..4fae0c456 100644 --- a/codex-rs/core/tests/suite/review.rs +++ b/codex-rs/core/tests/suite/review.rs @@ -726,6 +726,7 @@ async fn review_history_surfaces_in_parent_session() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/search_tool.rs b/codex-rs/core/tests/suite/search_tool.rs index d41785bd4..790cea37e 100644 --- a/codex-rs/core/tests/suite/search_tool.rs +++ b/codex-rs/core/tests/suite/search_tool.rs @@ -427,6 +427,7 @@ async fn tool_search_returns_deferred_tools_without_follow_up_tool_injection() - text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; diff --git a/codex-rs/core/tests/suite/sqlite_state.rs b/codex-rs/core/tests/suite/sqlite_state.rs index f35152e18..15f412201 100644 --- a/codex-rs/core/tests/suite/sqlite_state.rs +++ b/codex-rs/core/tests/suite/sqlite_state.rs @@ -463,16 +463,17 @@ async fn tool_call_logs_include_thread_id() -> Result<()> { let db = test.codex.state_db().expect("state db enabled"); let expected_thread_id = test.session_configured.session_id.to_string(); - let subscriber = tracing_subscriber::registry().with(codex_state::log_db::start(db.clone())); - let dispatch = tracing::Dispatch::new(subscriber); - let _guard = tracing::dispatcher::set_default(&dispatch); - test.submit_turn("run a shell command").await?; - { + + let log_db_layer = codex_state::log_db::start(db.clone()); + let subscriber = tracing_subscriber::registry().with(log_db_layer.clone()); + let dispatch = tracing::Dispatch::new(subscriber); + tracing::dispatcher::with_default(&dispatch, || { let span = tracing::info_span!("test_log_span", thread_id = %expected_thread_id); let _entered = span.enter(); tracing::info!("ToolCall: shell_command {{\"command\":\"echo hello\"}}"); - } + }); + log_db_layer.flush().await; let mut found = None; for _ in 0..80 { diff --git a/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs b/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs index 19d9d27cf..254ac7b12 100644 --- a/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs +++ b/codex-rs/core/tests/suite/stream_error_allows_next_turn.rs @@ -98,6 +98,7 @@ async fn continue_after_stream_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); @@ -117,6 +118,7 @@ async fn continue_after_stream_error() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/stream_no_completed.rs b/codex-rs/core/tests/suite/stream_no_completed.rs index 0c4d7cb49..149ba1c53 100644 --- a/codex-rs/core/tests/suite/stream_no_completed.rs +++ b/codex-rs/core/tests/suite/stream_no_completed.rs @@ -82,6 +82,7 @@ async fn retries_on_early_close() { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await .unwrap(); diff --git a/codex-rs/core/tests/suite/user_notification.rs b/codex-rs/core/tests/suite/user_notification.rs index a2303c0ec..3d02c2004 100644 --- a/codex-rs/core/tests/suite/user_notification.rs +++ b/codex-rs/core/tests/suite/user_notification.rs @@ -62,6 +62,7 @@ mv "${tmp_path}" "${payload_path}""#, text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(&codex, |ev| matches!(ev, EventMsg::TurnComplete(_))).await; diff --git a/codex-rs/core/tests/suite/window_headers.rs b/codex-rs/core/tests/suite/window_headers.rs index 76f135d20..bfcdcf25e 100644 --- a/codex-rs/core/tests/suite/window_headers.rs +++ b/codex-rs/core/tests/suite/window_headers.rs @@ -109,6 +109,7 @@ async fn submit_user_turn(codex: &Arc, text: &str) -> Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await?; wait_for_event(codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; diff --git a/codex-rs/exec/src/lib.rs b/codex-rs/exec/src/lib.rs index 10a5fa1da..f3b281a82 100644 --- a/codex-rs/exec/src/lib.rs +++ b/codex-rs/exec/src/lib.rs @@ -717,6 +717,7 @@ async fn run_exec_session(args: ExecRunArgs) -> anyhow::Result<()> { params: TurnStartParams { thread_id: primary_thread_id_for_span.clone(), input: items.into_iter().map(Into::into).collect(), + responsesapi_client_metadata: None, cwd: Some(default_cwd), approval_policy: Some(default_approval_policy.into()), approvals_reviewer: None, diff --git a/codex-rs/mcp-server/src/codex_tool_runner.rs b/codex-rs/mcp-server/src/codex_tool_runner.rs index 2eb0ea496..91fcc8a27 100644 --- a/codex-rs/mcp-server/src/codex_tool_runner.rs +++ b/codex-rs/mcp-server/src/codex_tool_runner.rs @@ -114,6 +114,7 @@ pub async fn run_codex_tool_session( text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }, trace: None, }; @@ -161,6 +162,7 @@ pub async fn run_codex_tool_session_reply( text_elements: Vec::new(), }], final_output_json_schema: None, + responsesapi_client_metadata: None, }) .await { diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index 8083db935..f1bbc0a7f 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -381,6 +381,9 @@ pub enum Op { /// Optional JSON Schema used to constrain the final assistant message for this turn. #[serde(skip_serializing_if = "Option::is_none")] final_output_json_schema: Option, + /// Optional turn-scoped Responses API `client_metadata`. + #[serde(default, skip_serializing_if = "Option::is_none")] + responsesapi_client_metadata: Option>, }, /// Similar to [`Op::UserInput`], but contains additional context required @@ -654,6 +657,7 @@ impl From> for Op { Op::UserInput { items: value, final_output_json_schema: None, + responsesapi_client_metadata: None, } } } @@ -4704,6 +4708,7 @@ mod tests { let op = Op::UserInput { items: Vec::new(), final_output_json_schema: None, + responsesapi_client_metadata: None, }; let json_op = serde_json::to_value(op)?; @@ -4721,6 +4726,7 @@ mod tests { Op::UserInput { items: Vec::new(), final_output_json_schema: None, + responsesapi_client_metadata: None, } ); @@ -4740,6 +4746,7 @@ mod tests { let op = Op::UserInput { items: Vec::new(), final_output_json_schema: Some(schema.clone()), + responsesapi_client_metadata: None, }; let json_op = serde_json::to_value(op)?; @@ -4755,6 +4762,33 @@ mod tests { Ok(()) } + #[test] + fn user_input_with_responsesapi_client_metadata_round_trips() -> Result<()> { + let op = Op::UserInput { + items: Vec::new(), + final_output_json_schema: None, + responsesapi_client_metadata: Some(HashMap::from([( + "fiber_run_id".to_string(), + "fiber-123".to_string(), + )])), + }; + + let json_op = serde_json::to_value(&op)?; + assert_eq!( + json_op, + json!({ + "type": "user_input", + "items": [], + "responsesapi_client_metadata": { + "fiber_run_id": "fiber-123", + } + }) + ); + assert_eq!(serde_json::from_value::(json_op)?, op); + + Ok(()) + } + #[test] fn user_input_text_serializes_empty_text_elements() -> Result<()> { let input = UserInput::Text { diff --git a/codex-rs/state/src/runtime/logs.rs b/codex-rs/state/src/runtime/logs.rs index 30bbe1882..56ef31d5d 100644 --- a/codex-rs/state/src/runtime/logs.rs +++ b/codex-rs/state/src/runtime/logs.rs @@ -638,6 +638,7 @@ mod tests { .await .expect("insert legacy log row"); pool.close().await; + drop(pool); let runtime = StateRuntime::init(codex_home.clone(), "test-provider".to_string()) .await diff --git a/codex-rs/tui/src/app_server_session.rs b/codex-rs/tui/src/app_server_session.rs index 501dbeb01..52d639517 100644 --- a/codex-rs/tui/src/app_server_session.rs +++ b/codex-rs/tui/src/app_server_session.rs @@ -421,6 +421,7 @@ impl AppServerSession { params: TurnStartParams { thread_id: thread_id.to_string(), input: items.into_iter().map(Into::into).collect(), + responsesapi_client_metadata: None, cwd: Some(cwd), approval_policy: Some(approval_policy.into()), approvals_reviewer: Some(approvals_reviewer.into()), @@ -471,6 +472,7 @@ impl AppServerSession { params: TurnSteerParams { thread_id: thread_id.to_string(), input: items.into_iter().map(Into::into).collect(), + responsesapi_client_metadata: None, expected_turn_id: turn_id, }, })