diff --git a/codex-rs/app-server-protocol/src/protocol/v2/turn.rs b/codex-rs/app-server-protocol/src/protocol/v2/turn.rs index 84567fb4c..836c6d4e9 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2/turn.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2/turn.rs @@ -68,7 +68,13 @@ pub struct TurnStartParams { #[ts(optional = nullable)] pub client_user_message_id: Option, pub input: Vec, - /// Optional turn-scoped Responses API client metadata. + /// Optional metadata to enrich Codex's ResponsesAPI turn metadata. + /// + /// Entries are flattened into the JSON string sent as + /// `client_metadata["x-codex-turn-metadata"]` on ResponsesAPI HTTP and websocket requests. + /// + /// They are not sent as top-level ResponsesAPI `client_metadata` keys, and reserved keys + /// such as `session_id`, `thread_id`, `turn_id`, and `window_id` cannot be overridden. #[experimental("turn/start.responsesapiClientMetadata")] #[ts(optional = nullable)] pub responsesapi_client_metadata: Option>, @@ -161,7 +167,13 @@ pub struct TurnSteerParams { #[ts(optional = nullable)] pub client_user_message_id: Option, pub input: Vec, - /// Optional turn-scoped Responses API client metadata. + /// Optional metadata to enrich Codex's ResponsesAPI turn metadata. + /// + /// Entries are flattened into the JSON string sent as + /// `client_metadata["x-codex-turn-metadata"]` on ResponsesAPI HTTP and websocket requests. + /// + /// They are not sent as top-level ResponsesAPI `client_metadata` keys, and reserved keys + /// such as `session_id`, `thread_id`, `turn_id`, and `window_id` cannot be overridden. #[experimental("turn/steer.responsesapiClientMetadata")] #[ts(optional = nullable)] pub responsesapi_client_metadata: Option>, diff --git a/codex-rs/app-server/tests/suite/v2/client_metadata.rs b/codex-rs/app-server/tests/suite/v2/client_metadata.rs index 873a1cc04..681e11ae3 100644 --- a/codex-rs/app-server/tests/suite/v2/client_metadata.rs +++ b/codex-rs/app-server/tests/suite/v2/client_metadata.rs @@ -112,6 +112,7 @@ async fn turn_start_forwards_client_metadata_to_responses_request_v2() -> Result assert_eq!(metadata["origin"].as_str(), Some("gaas")); assert_eq!(metadata["thread_source"].as_str(), Some("client-supplied")); assert_eq!(metadata["turn_id"].as_str(), Some(turn.id.as_str())); + assert!(metadata.get("installation_id").is_some()); assert!(metadata.get("session_id").is_some()); assert_eq!( metadata["window_id"].as_str(), diff --git a/codex-rs/core/src/client.rs b/codex-rs/core/src/client.rs index a7220d20d..75ed2039a 100644 --- a/codex-rs/core/src/client.rs +++ b/codex-rs/core/src/client.rs @@ -70,7 +70,6 @@ use codex_login::default_client::build_reqwest_client; use codex_otel::SessionTelemetry; use codex_otel::current_span_w3c_trace_context; -use codex_protocol::SessionId; use codex_protocol::ThreadId; use codex_protocol::config_types::ReasoningSummary as ReasoningSummaryConfig; use codex_protocol::config_types::Verbosity as VerbosityConfig; @@ -79,7 +78,6 @@ use codex_protocol::openai_models::ModelInfo; use codex_protocol::openai_models::ReasoningEffort as ReasoningEffortConfig; use codex_protocol::protocol::InternalSessionSource; use codex_protocol::protocol::SessionSource; -use codex_protocol::protocol::SubAgentSource; use codex_protocol::protocol::W3cTraceContext; use codex_rollout_trace::CompactionTraceContext; use codex_rollout_trace::InferenceTraceAttempt; @@ -111,6 +109,8 @@ use crate::client_common::Prompt; use crate::client_common::ResponseEvent; use crate::client_common::ResponseStream; use crate::feedback_tags; +use crate::responses_metadata::CodexResponsesMetadata; +use crate::responses_metadata::subagent_header_value; use crate::util::emit_feedback_auth_recovery_tags; use codex_api::map_api_error; use codex_feedback::FeedbackRequestTags; @@ -169,13 +169,10 @@ pub(crate) struct CompactConversationRequestSettings { /// configuration is per turn and is passed explicitly to streaming/unary methods. #[derive(Debug)] struct ModelClientState { - session_id: SessionId, thread_id: ThreadId, - installation_id: String, provider: SharedModelProvider, auth_env_telemetry: AuthEnvTelemetry, session_source: SessionSource, - parent_thread_id: Option, model_verbosity: Option, enable_request_compression: bool, include_timing_metrics: bool, @@ -319,12 +316,9 @@ impl ModelClient { /// are passed to [`ModelClientSession::stream`] (and other turn-scoped methods) explicitly. pub fn new( auth_manager: Option>, - session_id: SessionId, thread_id: ThreadId, - installation_id: String, provider_info: ModelProviderInfo, session_source: SessionSource, - parent_thread_id: Option, model_verbosity: Option, enable_request_compression: bool, include_timing_metrics: bool, @@ -341,13 +335,10 @@ impl ModelClient { let include_attestation = model_provider.supports_attestation(); Self { state: Arc::new(ModelClientState { - session_id, thread_id, - installation_id, provider: model_provider, auth_env_telemetry, session_source, - parent_thread_id, model_verbosity, enable_request_compression, include_timing_metrics, @@ -444,8 +435,7 @@ impl ModelClient { settings: CompactConversationRequestSettings, session_telemetry: &SessionTelemetry, compaction_trace: &CompactionTraceContext, - window_id: &str, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, ) -> Result> { if prompt.input.is_empty() { return Ok(Vec::new()); @@ -469,7 +459,7 @@ impl ModelClient { settings.effort, settings.summary, settings.service_tier, - window_id, + responses_metadata, )?; let ResponsesApiRequest { model, @@ -496,18 +486,17 @@ impl ModelClient { }; let mut extra_headers = ApiHeaderMap::new(); - if let Ok(header_value) = HeaderValue::from_str(&self.state.installation_id) { + if let Ok(header_value) = HeaderValue::from_str(&responses_metadata.installation_id) { extra_headers.insert(X_CODEX_INSTALLATION_ID_HEADER, header_value); } extra_headers.extend(build_responses_headers( self.state.beta_features_header.as_deref(), /*turn_state*/ None, - parse_turn_metadata_header(turn_metadata_header).as_ref(), )); - extra_headers.extend(self.build_responses_identity_headers(Some(window_id))); + extra_headers.extend(self.build_responses_compatibility_headers(responses_metadata)); extra_headers.extend(build_session_headers( - Some(self.state.session_id.to_string()), - Some(self.state.thread_id.to_string()), + Some(responses_metadata.session_id.to_string()), + Some(responses_metadata.thread_id.to_string()), )); if let Some(header_value) = self.generate_attestation_header_for().await { extra_headers.insert(X_OAI_ATTESTATION_HEADER, header_value); @@ -626,50 +615,29 @@ impl ModelClient { extra_headers } - fn build_responses_identity_headers(&self, window_id: Option<&str>) -> ApiHeaderMap { - let mut extra_headers = self.build_subagent_headers(); - if let Some(parent_thread_id) = parent_thread_id_header_value(self.state.parent_thread_id) - && let Ok(val) = HeaderValue::from_str(&parent_thread_id) - { - extra_headers.insert(X_CODEX_PARENT_THREAD_ID_HEADER, val); - } - if let Some(window_id) = window_id - && let Ok(val) = HeaderValue::from_str(window_id) - { - extra_headers.insert(X_CODEX_WINDOW_ID_HEADER, val); + fn build_responses_compatibility_headers( + &self, + responses_metadata: &CodexResponsesMetadata, + ) -> ApiHeaderMap { + let mut extra_headers = responses_metadata.compatibility_headers(); + if matches!( + self.state.session_source, + SessionSource::Internal(InternalSessionSource::MemoryConsolidation) + ) { + extra_headers.insert( + X_OPENAI_MEMGEN_REQUEST_HEADER, + HeaderValue::from_static("true"), + ); } extra_headers } fn build_ws_client_metadata( &self, - window_id: &str, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, use_responses_lite: bool, ) -> HashMap { - let mut client_metadata = HashMap::new(); - client_metadata.insert( - X_CODEX_INSTALLATION_ID_HEADER.to_string(), - self.state.installation_id.clone(), - ); - client_metadata.insert(X_CODEX_WINDOW_ID_HEADER.to_string(), window_id.to_string()); - if let Some(subagent) = subagent_header_value(&self.state.session_source) { - client_metadata.insert(X_OPENAI_SUBAGENT_HEADER.to_string(), subagent); - } - if let Some(parent_thread_id) = parent_thread_id_header_value(self.state.parent_thread_id) { - client_metadata.insert( - X_CODEX_PARENT_THREAD_ID_HEADER.to_string(), - parent_thread_id, - ); - } - if let Some(turn_metadata_header) = parse_turn_metadata_header(turn_metadata_header) - && let Ok(turn_metadata) = turn_metadata_header.to_str() - { - client_metadata.insert( - X_CODEX_TURN_METADATA_HEADER.to_string(), - turn_metadata.to_string(), - ); - } + let mut client_metadata = responses_metadata.client_metadata(); if use_responses_lite { client_metadata.insert( WS_REQUEST_HEADER_RESPONSES_LITE_CLIENT_METADATA_KEY.to_string(), @@ -743,7 +711,7 @@ impl ModelClient { effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, - window_id: &str, + responses_metadata: &CodexResponsesMetadata, ) -> Result { let instructions = &prompt.base_instructions.text; let input = prompt.get_formatted_input_for_request(model_info.use_responses_lite); @@ -786,13 +754,7 @@ impl ModelClient { service_tier, prompt_cache_key, text, - client_metadata: Some(HashMap::from([ - ( - X_CODEX_INSTALLATION_ID_HEADER.to_string(), - self.state.installation_id.clone(), - ), - (X_CODEX_WINDOW_ID_HEADER.to_string(), window_id.to_string()), - ])), + client_metadata: Some(responses_metadata.client_metadata()), }; Ok(request) } @@ -835,13 +797,13 @@ impl ModelClient { session_telemetry: &SessionTelemetry, api_provider: codex_api::Provider, api_auth: SharedAuthProvider, + responses_metadata: &CodexResponsesMetadata, turn_state: Option>>, - turn_metadata_header: Option<&str>, auth_context: AuthRequestTelemetryContext, request_route_telemetry: RequestRouteTelemetry, ) -> std::result::Result { let headers = self - .build_websocket_headers(turn_state.as_ref(), turn_metadata_header) + .build_websocket_headers(responses_metadata, turn_state.as_ref()) .await; let websocket_telemetry = ModelClientSession::build_websocket_telemetry( session_telemetry, @@ -921,22 +883,19 @@ impl ModelClient { /// replayed on reconnect within the same turn. async fn build_websocket_headers( &self, + responses_metadata: &CodexResponsesMetadata, turn_state: Option<&Arc>>, - turn_metadata_header: Option<&str>, ) -> ApiHeaderMap { - let turn_metadata_header = parse_turn_metadata_header(turn_metadata_header); - let session_id = self.state.session_id.to_string(); - let thread_id = self.state.thread_id.to_string(); - let mut headers = build_responses_headers( - self.state.beta_features_header.as_deref(), - turn_state, - turn_metadata_header.as_ref(), - ); - if let Ok(header_value) = HeaderValue::from_str(&thread_id) { + let mut headers = + build_responses_headers(self.state.beta_features_header.as_deref(), turn_state); + if let Ok(header_value) = HeaderValue::from_str(&responses_metadata.thread_id) { headers.insert("x-client-request-id", header_value); } - headers.extend(build_session_headers(Some(session_id), Some(thread_id))); - headers.extend(self.build_responses_identity_headers(/*window_id*/ None)); + headers.extend(build_session_headers( + Some(responses_metadata.session_id.to_string()), + Some(responses_metadata.thread_id.to_string()), + )); + headers.extend(self.build_responses_compatibility_headers(responses_metadata)); if let Some(header_value) = self.generate_attestation_header_for().await { headers.insert(X_OAI_ATTESTATION_HEADER, header_value); } @@ -979,27 +938,22 @@ impl ModelClientSession { /// regardless of transport choice. async fn build_responses_options( &self, - window_id: &str, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, compression: Compression, use_responses_lite: bool, ) -> ApiResponsesOptions { - let turn_metadata_header = parse_turn_metadata_header(turn_metadata_header); - let session_id = self.client.state.session_id.to_string(); - let thread_id = self.client.state.thread_id.to_string(); ApiResponsesOptions { - session_id: Some(session_id), - thread_id: Some(thread_id), + session_id: Some(responses_metadata.session_id.to_string()), + thread_id: Some(responses_metadata.thread_id.to_string()), session_source: Some(self.client.state.session_source.clone()), extra_headers: { let mut headers = build_responses_headers( self.client.state.beta_features_header.as_deref(), Some(&self.turn_state), - turn_metadata_header.as_ref(), ); headers.extend( self.client - .build_responses_identity_headers(Some(window_id)), + .build_responses_compatibility_headers(responses_metadata), ); if let Some(header_value) = self.client.generate_attestation_header_for().await { headers.insert(X_OAI_ATTESTATION_HEADER, header_value); @@ -1026,12 +980,12 @@ impl ModelClientSession { let previous_request = self.websocket_session.last_request.as_ref()?; let mut previous_without_input = previous_request.clone(); previous_without_input.input.clear(); + previous_without_input.client_metadata = None; let mut request_without_input = request.clone(); request_without_input.input.clear(); + request_without_input.client_metadata = None; if previous_without_input != request_without_input { - trace!( - "incremental request failed, properties didn't match {previous_without_input:?} != {request_without_input:?}" - ); + trace!("incremental request failed, websocket reuse properties didn't match"); return None; } @@ -1100,6 +1054,7 @@ impl ModelClientSession { pub async fn preconnect_websocket( &mut self, session_telemetry: &SessionTelemetry, + responses_metadata: &CodexResponsesMetadata, ) -> std::result::Result<(), ApiError> { if !self.client.responses_websocket_enabled() { return Ok(()); @@ -1124,8 +1079,8 @@ impl ModelClientSession { session_telemetry, client_setup.api_provider, client_setup.api_auth, + responses_metadata, Some(Arc::clone(&self.turn_state)), - /*turn_metadata_header*/ None, auth_context, RequestRouteTelemetry::for_endpoint(RESPONSES_ENDPOINT), ) @@ -1145,7 +1100,7 @@ impl ModelClientSession { wire_api = %self.client.state.provider.info().wire_api, transport = "responses_websocket", api.path = "responses", - turn.has_metadata_header = params.turn_metadata_header.is_some() + turn.has_metadata_header = params.responses_metadata.has_turn_metadata() ) )] async fn websocket_connection( @@ -1156,7 +1111,7 @@ impl ModelClientSession { session_telemetry, api_provider, api_auth, - turn_metadata_header, + responses_metadata, options, auth_context, request_route_telemetry, @@ -1180,8 +1135,8 @@ impl ModelClientSession { session_telemetry, api_provider, api_auth, + responses_metadata, Some(turn_state), - turn_metadata_header, auth_context, request_route_telemetry, ) @@ -1236,19 +1191,18 @@ impl ModelClientSession { transport = "responses_http", http.method = "POST", api.path = "responses", - turn.has_metadata_header = turn_metadata_header.is_some() + turn.has_metadata_header = responses_metadata.has_turn_metadata() ) )] async fn stream_responses_api( &self, - window_id: &str, prompt: &Prompt, model_info: &ModelInfo, session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, inference_trace: &InferenceTraceContext, ) -> Result { let auth_manager = self.client.state.provider.auth_manager(); @@ -1273,8 +1227,7 @@ impl ModelClientSession { let compression = self.responses_request_compression(client_setup.auth.as_ref()); let mut options = self .build_responses_options( - window_id, - turn_metadata_header, + responses_metadata, compression, model_info.use_responses_lite, ) @@ -1287,7 +1240,7 @@ impl ModelClientSession { effort.clone(), summary, service_tier.clone(), - window_id, + responses_metadata, )?; let inference_trace_attempt = inference_trace.start_attempt(); inference_trace_attempt.add_request_headers(&mut options.extra_headers); @@ -1355,20 +1308,19 @@ impl ModelClientSession { wire_api = %self.client.state.provider.info().wire_api, transport = "responses_websocket", api.path = "responses", - turn.has_metadata_header = turn_metadata_header.is_some(), + turn.has_metadata_header = responses_metadata.has_turn_metadata(), websocket.warmup = warmup ) )] async fn stream_responses_websocket( &mut self, - window_id: &str, prompt: &Prompt, model_info: &ModelInfo, session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, warmup: bool, request_trace: Option, inference_trace: &InferenceTraceContext, @@ -1390,8 +1342,7 @@ impl ModelClientSession { let options = self .build_responses_options( - window_id, - turn_metadata_header, + responses_metadata, compression, model_info.use_responses_lite, ) @@ -1403,13 +1354,12 @@ impl ModelClientSession { effort.clone(), summary, service_tier.clone(), - window_id, + responses_metadata, )?; let mut ws_payload = ResponseCreateWsRequest { client_metadata: response_create_client_metadata( Some(self.client.build_ws_client_metadata( - window_id, - turn_metadata_header, + responses_metadata, model_info.use_responses_lite, )), request_trace.as_ref(), @@ -1425,7 +1375,7 @@ impl ModelClientSession { session_telemetry, api_provider: client_setup.api_provider, api_auth: client_setup.api_auth, - turn_metadata_header, + responses_metadata, options: &options, auth_context: request_auth_context, request_route_telemetry: RequestRouteTelemetry::for_endpoint( @@ -1544,14 +1494,13 @@ impl ModelClientSession { #[allow(clippy::too_many_arguments)] pub async fn prewarm_websocket( &mut self, - window_id: &str, prompt: &Prompt, model_info: &ModelInfo, session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, ) -> Result<()> { if !self.client.responses_websocket_enabled() { return Ok(()); @@ -1563,14 +1512,13 @@ impl ModelClientSession { let disabled_trace = InferenceTraceContext::disabled(); match self .stream_responses_websocket( - window_id, prompt, model_info, session_telemetry, effort, summary, service_tier, - turn_metadata_header, + responses_metadata, /*warmup*/ true, current_span_w3c_trace_context(), &disabled_trace, @@ -1607,14 +1555,13 @@ impl ModelClientSession { /// branches. pub async fn stream( &mut self, - window_id: &str, prompt: &Prompt, model_info: &ModelInfo, session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, inference_trace: &InferenceTraceContext, ) -> Result { let wire_api = self.client.state.provider.info().wire_api; @@ -1624,14 +1571,13 @@ impl ModelClientSession { let request_trace = current_span_w3c_trace_context(); match self .stream_responses_websocket( - window_id, prompt, model_info, session_telemetry, effort.clone(), summary, service_tier.clone(), - turn_metadata_header, + responses_metadata, /*warmup*/ false, request_trace, inference_trace, @@ -1646,14 +1592,13 @@ impl ModelClientSession { } self.stream_responses_api( - window_id, prompt, model_info, session_telemetry, effort, summary, service_tier, - turn_metadata_header, + responses_metadata, inference_trace, ) .await @@ -1680,14 +1625,6 @@ impl ModelClientSession { } } -/// Parses per-turn metadata into an HTTP header value. -/// -/// Invalid values are treated as absent so callers can compare and propagate -/// metadata with the same sanitization path used when constructing headers. -fn parse_turn_metadata_header(turn_metadata_header: Option<&str>) -> Option { - turn_metadata_header.and_then(|value| HeaderValue::from_str(value).ok()) -} - /// Stamp a ResponsesWsRequest with the current time. /// /// Meant to be called just before sending the request over the socket, to capture realistic @@ -1709,11 +1646,9 @@ fn stamp_ws_stream_request_start_ms(request: &mut ResponsesWsRequest) { /// /// - `x-codex-beta-features`: comma-separated beta feature keys enabled for the session. /// - `x-codex-turn-state`: sticky routing token captured earlier in the turn. -/// - `x-codex-turn-metadata`: optional per-turn metadata for observability. fn build_responses_headers( beta_features_header: Option<&str>, turn_state: Option<&Arc>>, - turn_metadata_header: Option<&HeaderValue>, ) -> ApiHeaderMap { let mut headers = ApiHeaderMap::new(); if let Some(value) = beta_features_header @@ -1728,9 +1663,6 @@ fn build_responses_headers( { headers.insert(X_CODEX_TURN_STATE_HEADER, header_value); } - if let Some(header_value) = turn_metadata_header { - headers.insert(X_CODEX_TURN_METADATA_HEADER, header_value.clone()); - } headers } @@ -1743,31 +1675,6 @@ fn add_responses_lite_header(headers: &mut ApiHeaderMap, use_responses_lite: boo } } -fn subagent_header_value(session_source: &SessionSource) -> Option { - match session_source { - SessionSource::SubAgent(subagent_source) => match subagent_source { - SubAgentSource::Review => Some("review".to_string()), - SubAgentSource::Compact => Some("compact".to_string()), - SubAgentSource::MemoryConsolidation => Some("memory_consolidation".to_string()), - SubAgentSource::ThreadSpawn { .. } => Some("collab_spawn".to_string()), - SubAgentSource::Other(label) => Some(label.clone()), - }, - SessionSource::Internal(InternalSessionSource::MemoryConsolidation) => { - Some("memory_consolidation".to_string()) - } - SessionSource::Cli - | SessionSource::VSCode - | SessionSource::Exec - | SessionSource::Mcp - | SessionSource::Custom(_) - | SessionSource::Unknown => None, - } -} - -fn parent_thread_id_header_value(parent_thread_id: Option) -> Option { - parent_thread_id.map(|parent_thread_id| parent_thread_id.to_string()) -} - const RESPONSE_STREAM_CHANNEL_CAPACITY: usize = 1600; const STREAM_DROPPED_REASON: &str = "response stream dropped before provider terminal event"; @@ -2004,7 +1911,7 @@ struct WebsocketConnectParams<'a> { session_telemetry: &'a SessionTelemetry, api_provider: codex_api::Provider, api_auth: SharedAuthProvider, - turn_metadata_header: Option<&'a str>, + responses_metadata: &'a CodexResponsesMetadata, options: &'a ApiResponsesOptions, auth_context: AuthRequestTelemetryContext, request_route_telemetry: RequestRouteTelemetry, diff --git a/codex-rs/core/src/client_tests.rs b/codex-rs/core/src/client_tests.rs index d93bd8bf7..18cd20223 100644 --- a/codex-rs/core/src/client_tests.rs +++ b/codex-rs/core/src/client_tests.rs @@ -10,6 +10,9 @@ use super::X_OPENAI_SUBAGENT_HEADER; use crate::AttestationContext; use crate::AttestationProvider; use crate::GenerateAttestationFuture; +use crate::responses_metadata::CodexResponsesMetadata; +use crate::test_support::TestCodexResponsesRequestKind; +use crate::test_support::responses_metadata as test_responses_metadata; use codex_api::ApiError; use codex_api::ResponseEvent; use codex_app_server_protocol::AuthMode; @@ -21,7 +24,6 @@ use codex_model_provider_info::ModelProviderInfo; use codex_model_provider_info::WireApi; use codex_model_provider_info::create_oss_provider_with_base_url; use codex_otel::SessionTelemetry; -use codex_protocol::SessionId; use codex_protocol::ThreadId; use codex_protocol::models::ContentItem; use codex_protocol::models::ResponseItem; @@ -60,24 +62,16 @@ use tracing_subscriber::layer::SubscriberExt; use tracing_subscriber::registry::LookupSpan; use tracing_subscriber::util::SubscriberInitExt; -fn test_model_client(session_source: SessionSource) -> ModelClient { - test_model_client_with_parent(session_source, /*parent_thread_id*/ None) -} +const TEST_INSTALLATION_ID: &str = "11111111-1111-4111-8111-111111111111"; -fn test_model_client_with_parent( - session_source: SessionSource, - parent_thread_id: Option, -) -> ModelClient { +fn test_model_client(session_source: SessionSource) -> ModelClient { let provider = create_oss_provider_with_base_url("https://example.com/v1", WireApi::Responses); let thread_id = ThreadId::new(); ModelClient::new( /*auth_manager*/ None, - thread_id.into(), thread_id, - /*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(), provider, session_source, - parent_thread_id, /*model_verbosity*/ None, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, @@ -86,6 +80,26 @@ fn test_model_client_with_parent( ) } +fn test_responses_metadata_for_client( + client: &ModelClient, + turn_id: Option<&str>, + window_id: String, + parent_thread_id: Option, + request_kind: TestCodexResponsesRequestKind, +) -> CodexResponsesMetadata { + let thread_id = client.state.thread_id.to_string(); + test_responses_metadata( + TEST_INSTALLATION_ID, + &thread_id, + &thread_id, + turn_id, + window_id, + &client.state.session_source, + parent_thread_id, + request_kind, + ) +} + fn test_model_info() -> ModelInfo { serde_json::from_value(json!({ "slug": "gpt-test", @@ -280,48 +294,63 @@ fn build_subagent_headers_sets_internal_memory_consolidation_label() { #[test] fn build_ws_client_metadata_includes_window_lineage_and_turn_metadata() { let parent_thread_id = ThreadId::new(); - let client = test_model_client_with_parent( - SessionSource::SubAgent(SubAgentSource::ThreadSpawn { - parent_thread_id, - depth: 2, - agent_path: None, - agent_nickname: None, - agent_role: None, - }), - Some(parent_thread_id), - ); + let client = test_model_client(SessionSource::SubAgent(SubAgentSource::ThreadSpawn { + parent_thread_id, + depth: 2, + agent_path: None, + agent_nickname: None, + agent_role: None, + })); - let thread_id = client.state.thread_id; - let window_id = format!("{thread_id}:1"); - let client_metadata = client.build_ws_client_metadata( - &window_id, - Some(r#"{"turn_id":"turn-123"}"#), - /*use_responses_lite*/ false, + let thread_id = client.state.thread_id.to_string(); + let expected_window_id = format!("{thread_id}:1"); + let responses_metadata = test_responses_metadata_for_client( + &client, + Some("turn-123"), + expected_window_id.clone(), + Some(parent_thread_id), + TestCodexResponsesRequestKind::Turn, ); + let client_metadata = + client.build_ws_client_metadata(&responses_metadata, /*use_responses_lite*/ false); + let parent_thread_id = parent_thread_id.to_string(); + let turn_metadata: serde_json::Value = serde_json::from_str( + client_metadata + .get(X_CODEX_TURN_METADATA_HEADER) + .expect("turn metadata"), + ) + .expect("valid turn metadata"); + for (client_key, metadata_key, expected) in [ + ( + X_CODEX_INSTALLATION_ID_HEADER, + "installation_id", + "11111111-1111-4111-8111-111111111111", + ), + ("session_id", "session_id", thread_id.as_str()), + ("thread_id", "thread_id", thread_id.as_str()), + ("turn_id", "turn_id", "turn-123"), + ( + X_CODEX_WINDOW_ID_HEADER, + "window_id", + expected_window_id.as_str(), + ), + ( + X_CODEX_PARENT_THREAD_ID_HEADER, + "parent_thread_id", + parent_thread_id.as_str(), + ), + ] { + assert_eq!( + client_metadata.get(client_key).map(String::as_str), + Some(expected) + ); + assert_eq!(turn_metadata[metadata_key].as_str(), Some(expected)); + } assert_eq!( - client_metadata, - std::collections::HashMap::from([ - ( - X_CODEX_INSTALLATION_ID_HEADER.to_string(), - "11111111-1111-4111-8111-111111111111".to_string(), - ), - ( - X_CODEX_WINDOW_ID_HEADER.to_string(), - format!("{thread_id}:1"), - ), - ( - X_OPENAI_SUBAGENT_HEADER.to_string(), - "collab_spawn".to_string(), - ), - ( - X_CODEX_PARENT_THREAD_ID_HEADER.to_string(), - parent_thread_id.to_string(), - ), - ( - X_CODEX_TURN_METADATA_HEADER.to_string(), - r#"{"turn_id":"turn-123"}"#.to_string(), - ), - ]) + client_metadata + .get(X_OPENAI_SUBAGENT_HEADER) + .map(String::as_str), + Some("collab_spawn") ); } @@ -529,12 +558,9 @@ fn model_client_with_counting_attestation( }; let model_client = ModelClient::new( auth_manager, - SessionId::new(), ThreadId::new(), - /*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(), provider, SessionSource::Exec, - /*parent_thread_id*/ None, /*model_verbosity*/ None, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, @@ -550,9 +576,16 @@ fn model_client_with_counting_attestation( async fn websocket_handshake_includes_attestation_for_chatgpt_codex_responses() { let (model_client, attestation_calls) = model_client_with_counting_attestation(/*include_attestation*/ true); + let responses_metadata = test_responses_metadata_for_client( + &model_client, + /*turn_id*/ None, + format!("{}:0", model_client.state.thread_id), + /*parent_thread_id*/ None, + TestCodexResponsesRequestKind::WebsocketConnection, + ); let headers = model_client - .build_websocket_headers(/*turn_state*/ None, /*turn_metadata_header*/ None) + .build_websocket_headers(&responses_metadata, /*turn_state*/ None) .await; assert_eq!( diff --git a/codex-rs/core/src/compact.rs b/codex-rs/core/src/compact.rs index 26b5e33bf..6ffb779d3 100644 --- a/codex-rs/core/src/compact.rs +++ b/codex-rs/core/src/compact.rs @@ -8,12 +8,14 @@ use crate::hook_runtime::PostCompactHookOutcome; use crate::hook_runtime::PreCompactHookOutcome; use crate::hook_runtime::run_post_compact_hooks; use crate::hook_runtime::run_pre_compact_hooks; +use crate::responses_metadata::CodexResponsesMetadata; +use crate::responses_metadata::CodexResponsesRequestKind; +use crate::responses_metadata::CompactionTurnMetadata; #[cfg(test)] use crate::session::PreviousTurnSettings; use crate::session::session::Session; use crate::session::turn::get_last_assistant_message_from_turn; use crate::session::turn_context::TurnContext; -use crate::turn_metadata::CompactionTurnMetadata; use crate::util::backoff; use codex_analytics::CodexCompactionEvent; use codex_analytics::CompactionImplementation; @@ -215,6 +217,12 @@ async fn run_compact_task_inner_impl( // Reuse one client session so turn-scoped state (sticky routing, websocket incremental // request tracking) // survives retries within this compact turn. + let window_id = sess.current_window_id().await; + let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata( + sess.installation_id.clone(), + window_id, + CodexResponsesRequestKind::Compaction(compaction_metadata), + ); loop { // Clone is required because of the loop @@ -228,16 +236,11 @@ async fn run_compact_task_inner_impl( personality: turn_context.personality, ..Default::default() }; - let window_id = sess.current_window_id().await; - let turn_metadata_header = turn_context - .turn_metadata_state - .current_header_value_for_compaction(&window_id, compaction_metadata); let attempt_result = drain_to_completed( &sess, turn_context.as_ref(), &mut client_session, - &window_id, - turn_metadata_header.as_deref(), + &responses_metadata, &prompt, ) .await; @@ -587,20 +590,18 @@ async fn drain_to_completed( sess: &Session, turn_context: &TurnContext, client_session: &mut ModelClientSession, - window_id: &str, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, prompt: &Prompt, ) -> CodexResult<()> { let mut stream = client_session .stream( - window_id, prompt, &turn_context.model_info, &turn_context.session_telemetry, turn_context.reasoning_effort.clone(), turn_context.reasoning_summary, turn_context.config.service_tier.clone(), - turn_metadata_header, + responses_metadata, // Rollout tracing currently models remote compaction only; local compaction streams // are left untraced until the reducer has a first-class local compaction lifecycle. &InferenceTraceContext::disabled(), diff --git a/codex-rs/core/src/compact_remote.rs b/codex-rs/core/src/compact_remote.rs index 93457053f..0822d0f0a 100644 --- a/codex-rs/core/src/compact_remote.rs +++ b/codex-rs/core/src/compact_remote.rs @@ -12,10 +12,11 @@ use crate::hook_runtime::PostCompactHookOutcome; use crate::hook_runtime::PreCompactHookOutcome; use crate::hook_runtime::run_post_compact_hooks; use crate::hook_runtime::run_pre_compact_hooks; +use crate::responses_metadata::CodexResponsesRequestKind; +use crate::responses_metadata::CompactionTurnMetadata; use crate::session::session::Session; use crate::session::turn::built_tools; use crate::session::turn_context::TurnContext; -use crate::turn_metadata::CompactionTurnMetadata; use codex_analytics::CompactionImplementation; use codex_analytics::CompactionPhase; use codex_analytics::CompactionReason; @@ -225,9 +226,11 @@ async fn run_remote_compact_task_inner_impl( output_schema_strict: true, }; let window_id = sess.current_window_id().await; - let turn_metadata_header = turn_context - .turn_metadata_state - .current_header_value_for_compaction(&window_id, compaction_metadata); + let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata( + sess.installation_id.clone(), + window_id, + CodexResponsesRequestKind::Compaction(compaction_metadata), + ); let mut new_history = sess .services .model_client @@ -245,8 +248,7 @@ async fn run_remote_compact_task_inner_impl( }, &turn_context.session_telemetry, &compaction_trace, - &window_id, - turn_metadata_header.as_deref(), + &responses_metadata, ) .await?; let new_window_id = sess.advance_auto_compact_window_id().await; diff --git a/codex-rs/core/src/compact_remote_v2.rs b/codex-rs/core/src/compact_remote_v2.rs index 5ab7a68a5..047c2527c 100644 --- a/codex-rs/core/src/compact_remote_v2.rs +++ b/codex-rs/core/src/compact_remote_v2.rs @@ -15,12 +15,14 @@ use crate::hook_runtime::PostCompactHookOutcome; use crate::hook_runtime::PreCompactHookOutcome; use crate::hook_runtime::run_post_compact_hooks; use crate::hook_runtime::run_pre_compact_hooks; +use crate::responses_metadata::CodexResponsesMetadata; +use crate::responses_metadata::CodexResponsesRequestKind; +use crate::responses_metadata::CompactionTurnMetadata; use crate::responses_retry::ResponsesStreamRequest; use crate::responses_retry::handle_retryable_response_stream_error; use crate::session::session::Session; use crate::session::turn::built_tools; use crate::session::turn_context::TurnContext; -use crate::turn_metadata::CompactionTurnMetadata; use codex_analytics::CompactionImplementation; use codex_analytics::CompactionPhase; use codex_analytics::CompactionReason; @@ -241,9 +243,11 @@ async fn run_remote_compact_task_inner_impl( }; let window_id = sess.current_window_id().await; - let turn_metadata_header = turn_context - .turn_metadata_state - .current_header_value_for_compaction(&window_id, compaction_metadata); + let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata( + sess.installation_id.clone(), + window_id, + CodexResponsesRequestKind::Compaction(compaction_metadata), + ); let trace_attempt = compaction_trace.start_attempt(&serde_json::json!({ "model": turn_context.model_info.slug.as_str(), "instructions": prompt.base_instructions.text.as_str(), @@ -264,8 +268,7 @@ async fn run_remote_compact_task_inner_impl( turn_context, client_session, &prompt, - &window_id, - turn_metadata_header.as_deref(), + &responses_metadata, ) .await; @@ -327,8 +330,7 @@ async fn run_remote_compaction_request_v2( turn_context: &TurnContext, client_session: &mut ModelClientSession, prompt: &Prompt, - window_id: &str, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, ) -> CodexResult { let max_retries = turn_context .provider @@ -339,14 +341,13 @@ async fn run_remote_compaction_request_v2( loop { let result = match client_session .stream( - window_id, prompt, &turn_context.model_info, &turn_context.session_telemetry, turn_context.reasoning_effort.clone(), turn_context.reasoning_summary, turn_context.config.service_tier.clone(), - turn_metadata_header, + responses_metadata, &InferenceTraceContext::disabled(), ) .await diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 775556e47..073c06ced 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -12,9 +12,12 @@ mod client_common; mod realtime_context; mod realtime_conversation; mod realtime_prompt; +mod responses_metadata; mod responses_retry; pub(crate) mod session; +pub use responses_metadata::CodexResponsesMetadata; pub use session::SteerInputError; +pub use turn_metadata::detached_memory_responses_metadata; mod codex_thread; mod compact_remote; mod compact_remote_v2; @@ -190,7 +193,6 @@ pub use exec_policy::check_execpolicy_for_warnings; pub use exec_policy::format_exec_policy_error_with_source; pub use exec_policy::load_exec_policy; pub use installation_id::resolve_installation_id; -pub use turn_metadata::build_turn_metadata_header; pub mod compact; mod memory_usage; pub mod otel_init; diff --git a/codex-rs/core/src/responses_metadata.rs b/codex-rs/core/src/responses_metadata.rs new file mode 100644 index 000000000..6e86ed8eb --- /dev/null +++ b/codex-rs/core/src/responses_metadata.rs @@ -0,0 +1,366 @@ +use std::collections::BTreeMap; +use std::collections::HashMap; + +use codex_analytics::CompactionImplementation; +use codex_analytics::CompactionPhase; +use codex_analytics::CompactionReason; +use codex_analytics::CompactionStrategy; +use codex_analytics::CompactionTrigger; +use codex_protocol::ThreadId; +use codex_protocol::protocol::InternalSessionSource; +use codex_protocol::protocol::SessionSource; +use codex_protocol::protocol::SubAgentSource; +use codex_utils_string::to_ascii_json_string; +use http::HeaderMap as ApiHeaderMap; +use http::HeaderValue; +use serde::Serialize; +use serde_json::Value; + +use crate::client::X_CODEX_INSTALLATION_ID_HEADER; +use crate::client::X_CODEX_PARENT_THREAD_ID_HEADER; +use crate::client::X_CODEX_TURN_METADATA_HEADER; +use crate::client::X_CODEX_WINDOW_ID_HEADER; +use crate::client::X_OPENAI_SUBAGENT_HEADER; + +pub(crate) const INSTALLATION_ID_KEY: &str = "installation_id"; +pub(crate) const SESSION_ID_KEY: &str = "session_id"; +pub(crate) const THREAD_ID_KEY: &str = "thread_id"; +pub(crate) const TURN_ID_KEY: &str = "turn_id"; +pub(crate) const WINDOW_ID_KEY: &str = "window_id"; +pub(crate) const REQUEST_KIND_KEY: &str = "request_kind"; +pub(crate) const COMPACTION_KEY: &str = "compaction"; +pub(crate) const TURN_STARTED_AT_UNIX_MS_KEY: &str = "turn_started_at_unix_ms"; + +pub(crate) const FORKED_FROM_THREAD_ID_KEY: &str = "forked_from_thread_id"; +pub(crate) const PARENT_THREAD_ID_KEY: &str = "parent_thread_id"; +pub(crate) const SUBAGENT_KIND_KEY: &str = "subagent_kind"; +pub(crate) const SANDBOX_KEY: &str = "sandbox"; +pub(crate) const WORKSPACES_KEY: &str = "workspaces"; + +// App-server clients can specify additional metadata in the `responsesapi_client_metadata` param +// when submitting a turn, but they must not override fields owned by core. +const RESERVED_METADATA_KEYS: &[&str] = &[ + INSTALLATION_ID_KEY, + X_CODEX_INSTALLATION_ID_HEADER, + SESSION_ID_KEY, + THREAD_ID_KEY, + TURN_ID_KEY, + WINDOW_ID_KEY, + X_CODEX_WINDOW_ID_HEADER, + X_CODEX_TURN_METADATA_HEADER, + X_CODEX_PARENT_THREAD_ID_HEADER, + X_OPENAI_SUBAGENT_HEADER, + REQUEST_KIND_KEY, + COMPACTION_KEY, + TURN_STARTED_AT_UNIX_MS_KEY, + FORKED_FROM_THREAD_ID_KEY, + PARENT_THREAD_ID_KEY, + SUBAGENT_KIND_KEY, + SANDBOX_KEY, + WORKSPACES_KEY, +]; + +/// Metadata attached to model requests whose purpose is conversation compaction. +/// +/// This covers both local compaction requests sent through the normal `/responses` path and remote +/// compaction requests sent through `/responses/compact`. These fields describe the operation at +/// dispatch time. Post-response outcomes such as status, error, duration, and token deltas remain +/// in compaction analytics events. +#[derive(Clone, Copy, Debug, Serialize)] +pub(crate) struct CompactionTurnMetadata { + trigger: CompactionTrigger, + reason: CompactionReason, + implementation: CompactionImplementation, + phase: CompactionPhase, + strategy: CompactionStrategy, +} + +impl CompactionTurnMetadata { + pub(crate) fn new( + trigger: CompactionTrigger, + reason: CompactionReason, + implementation: CompactionImplementation, + phase: CompactionPhase, + ) -> Self { + Self { + trigger, + reason, + implementation, + phase, + strategy: CompactionStrategy::Memento, + } + } +} + +#[derive(Clone, Copy, Debug)] +pub(crate) enum CodexResponsesRequestKind { + Turn, + Prewarm, + Compaction(CompactionTurnMetadata), + Memory, +} + +impl CodexResponsesRequestKind { + fn metadata(self) -> (&'static str, Option) { + match self { + CodexResponsesRequestKind::Turn => ("turn", None), + CodexResponsesRequestKind::Prewarm => ("prewarm", None), + CodexResponsesRequestKind::Compaction(metadata) => ("compaction", Some(metadata)), + CodexResponsesRequestKind::Memory => ("memory", None), + } + } + + fn has_turn_identity(self) -> bool { + !matches!(self, CodexResponsesRequestKind::Memory) + } +} + +#[derive(Clone, Debug, Serialize, Default)] +pub(crate) struct TurnMetadataWorkspace { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) associated_remote_urls: Option>, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) latest_git_commit_hash: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) has_changes: Option, +} + +/// Caller-owned snapshot of Codex metadata sent to ResponsesAPI. +/// +/// The full Codex turn metadata blob is transported canonically as +/// `client_metadata["x-codex-turn-metadata"]`. Flat `client_metadata` keys and direct HTTP/ws +/// headers are generated compatibility projections of this snapshot, not separate sources of +/// truth. +#[derive(Clone, Debug)] +pub struct CodexResponsesMetadata { + pub(crate) installation_id: String, + pub(crate) session_id: String, + pub(crate) thread_id: String, + pub(crate) turn_id: Option, + pub(crate) window_id: String, + pub(crate) request_kind: Option, + pub(crate) forked_from_thread_id: Option, + pub(crate) parent_thread_id: Option, + pub(crate) subagent_header: Option, + pub(crate) subagent_kind: Option, + pub(crate) sandbox: Option, + pub(crate) workspaces: BTreeMap, + pub(crate) turn_started_at_unix_ms: Option, + pub(crate) extra: BTreeMap, +} + +impl CodexResponsesMetadata { + pub(crate) fn new( + installation_id: String, + session_id: String, + thread_id: String, + window_id: String, + ) -> Self { + Self { + installation_id, + session_id, + thread_id, + turn_id: None, + window_id, + request_kind: None, + forked_from_thread_id: None, + parent_thread_id: None, + subagent_header: None, + subagent_kind: None, + sandbox: None, + workspaces: BTreeMap::new(), + turn_started_at_unix_ms: None, + extra: BTreeMap::new(), + } + } + + pub(crate) fn has_turn_metadata(&self) -> bool { + self.request_kind.is_some() + } + + pub(crate) fn turn_metadata_json(&self) -> Option { + to_ascii_json_string(&self.turn_metadata_payload()).ok() + } + + pub(crate) fn turn_metadata_value(&self) -> Option { + serde_json::to_value(self.turn_metadata_payload()).ok() + } + + pub(crate) fn client_metadata(&self) -> HashMap { + let mut client_metadata = HashMap::from([ + ( + X_CODEX_INSTALLATION_ID_HEADER.to_string(), + self.installation_id.clone(), + ), + (SESSION_ID_KEY.to_string(), self.session_id.clone()), + (THREAD_ID_KEY.to_string(), self.thread_id.clone()), + (X_CODEX_WINDOW_ID_HEADER.to_string(), self.window_id.clone()), + ]); + if let Some(turn_id) = &self.turn_id { + client_metadata.insert(TURN_ID_KEY.to_string(), turn_id.clone()); + } + if let Some(subagent_header) = &self.subagent_header { + client_metadata.insert( + X_OPENAI_SUBAGENT_HEADER.to_string(), + subagent_header.clone(), + ); + } + if let Some(parent_thread_id) = self.parent_thread_id { + client_metadata.insert( + X_CODEX_PARENT_THREAD_ID_HEADER.to_string(), + parent_thread_id.to_string(), + ); + } + if self.has_turn_metadata() + && let Some(turn_metadata_json) = self.turn_metadata_json() + { + client_metadata.insert(X_CODEX_TURN_METADATA_HEADER.to_string(), turn_metadata_json); + } + client_metadata + } + + pub(crate) fn compatibility_headers(&self) -> ApiHeaderMap { + let mut headers = ApiHeaderMap::new(); + insert_header(&mut headers, X_CODEX_WINDOW_ID_HEADER, &self.window_id); + // Direct x-codex-turn-metadata is compatibility output. New per-request consumers should + // prefer client_metadata["x-codex-turn-metadata"], which is rendered from this same object. + if self.has_turn_metadata() + && let Some(turn_metadata_json) = self.turn_metadata_json() + { + insert_header( + &mut headers, + X_CODEX_TURN_METADATA_HEADER, + &turn_metadata_json, + ); + } + if let Some(parent_thread_id) = self.parent_thread_id { + insert_header( + &mut headers, + X_CODEX_PARENT_THREAD_ID_HEADER, + &parent_thread_id.to_string(), + ); + } + if let Some(subagent_header) = &self.subagent_header { + insert_header(&mut headers, X_OPENAI_SUBAGENT_HEADER, subagent_header); + } + headers + } + + fn turn_metadata_payload(&self) -> CodexTurnMetadataPayload<'_> { + let request_kind = self.request_kind; + let (request_kind_value, compaction) = request_kind.map_or((None, None), |request_kind| { + let (request_kind, compaction) = request_kind.metadata(); + (Some(request_kind), compaction) + }); + let has_turn_identity = + request_kind.is_none_or(CodexResponsesRequestKind::has_turn_identity); + let has_request_identity = + request_kind.is_some_and(CodexResponsesRequestKind::has_turn_identity); + CodexTurnMetadataPayload { + installation_id: has_request_identity.then_some(self.installation_id.as_str()), + session_id: has_turn_identity.then_some(self.session_id.as_str()), + thread_id: has_turn_identity.then_some(self.thread_id.as_str()), + turn_id: has_turn_identity + .then_some(self.turn_id.as_deref()) + .flatten(), + window_id: has_request_identity.then_some(self.window_id.as_str()), + request_kind: request_kind_value, + forked_from_thread_id: self.forked_from_thread_id, + parent_thread_id: self.parent_thread_id, + subagent_kind: self.subagent_kind.as_deref(), + sandbox: self.sandbox.as_deref(), + workspaces: non_empty_workspaces(&self.workspaces), + turn_started_at_unix_ms: self.turn_started_at_unix_ms, + compaction, + // responsesapi_client_metadata enriches the Codex turn metadata blob, not literal + // top-level Responses client_metadata. Reserved Codex-owned keys are filtered when + // these extras enter turn state. + extra: &self.extra, + } + } +} + +pub(crate) fn subagent_header_value(session_source: &SessionSource) -> Option { + match session_source { + SessionSource::SubAgent(subagent_source) => match subagent_source { + SubAgentSource::Review => Some("review".to_string()), + SubAgentSource::Compact => Some("compact".to_string()), + SubAgentSource::MemoryConsolidation => Some("memory_consolidation".to_string()), + SubAgentSource::ThreadSpawn { .. } => Some("collab_spawn".to_string()), + SubAgentSource::Other(label) => Some(label.clone()), + }, + SessionSource::Internal(InternalSessionSource::MemoryConsolidation) => { + Some("memory_consolidation".to_string()) + } + SessionSource::Cli + | SessionSource::VSCode + | SessionSource::Exec + | SessionSource::Mcp + | SessionSource::Custom(_) + | SessionSource::Unknown => None, + } +} + +pub(crate) fn subagent_metadata_kind(session_source: &SessionSource) -> Option { + match session_source { + SessionSource::SubAgent(subagent_source) => Some(subagent_source.kind().to_string()), + SessionSource::Cli + | SessionSource::VSCode + | SessionSource::Exec + | SessionSource::Mcp + | SessionSource::Custom(_) + | SessionSource::Internal(_) + | SessionSource::Unknown => None, + } +} + +fn insert_header(headers: &mut ApiHeaderMap, name: &'static str, value: &str) { + if let Ok(header_value) = HeaderValue::from_str(value) { + headers.insert(name, header_value); + } +} + +pub(crate) fn filter_extra_metadata(extra: HashMap) -> BTreeMap { + extra + .into_iter() + .filter(|(key, _)| !RESERVED_METADATA_KEYS.contains(&key.as_str())) + .collect() +} + +fn non_empty_workspaces( + workspaces: &BTreeMap, +) -> Option<&BTreeMap> { + (!workspaces.is_empty()).then_some(workspaces) +} + +#[derive(Serialize)] +struct CodexTurnMetadataPayload<'a> { + #[serde(default, skip_serializing_if = "Option::is_none")] + installation_id: Option<&'a str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + session_id: Option<&'a str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + thread_id: Option<&'a str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + turn_id: Option<&'a str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + window_id: Option<&'a str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + request_kind: Option<&'static str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + forked_from_thread_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + parent_thread_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + subagent_kind: Option<&'a str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + sandbox: Option<&'a str>, + #[serde(default, skip_serializing_if = "Option::is_none")] + workspaces: Option<&'a BTreeMap>, + #[serde(default, skip_serializing_if = "Option::is_none")] + turn_started_at_unix_ms: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + compaction: Option, + #[serde(flatten)] + extra: &'a BTreeMap, +} diff --git a/codex-rs/core/src/session/review.rs b/codex-rs/core/src/session/review.rs index 60c028036..9a7dfc829 100644 --- a/codex-rs/core/src/session/review.rs +++ b/codex-rs/core/src/session/review.rs @@ -91,7 +91,6 @@ pub(super) async fn spawn_review_thread( forked_from_thread_id, parent_turn_context.parent_thread_id, &session_source, - parent_turn_context.thread_source.clone(), review_turn_id.clone(), #[allow(deprecated)] parent_turn_context.cwd.clone(), diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index cce3e10c7..2d14ceb5f 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -1021,12 +1021,9 @@ impl Session { attestation_provider: attestation_provider.clone(), model_client: ModelClient::new( Some(Arc::clone(&auth_manager)), - session_id, thread_id, - installation_id.clone(), session_configuration.provider.clone(), session_configuration.session_source.clone(), - session_configuration.parent_thread_id, config.model_verbosity, config.features.enabled(Feature::EnableRequestCompression), config.features.enabled(Feature::RuntimeMetrics), diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index f44f9a74f..0dbb13fcd 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -423,12 +423,9 @@ fn test_model_client_session() -> crate::client::ModelClientSession { .expect("test thread id should be valid"); crate::client::ModelClient::new( /*auth_manager*/ None, - thread_id.into(), thread_id, - /*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(), ModelProviderInfo::create_openai_provider(/* base_url */ /*base_url*/ None), codex_protocol::protocol::SessionSource::Exec, - /*parent_thread_id*/ None, /*model_verbosity*/ None, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, @@ -4996,12 +4993,9 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { attestation_provider: None, model_client: ModelClient::new( Some(auth_manager.clone()), - thread_id.into(), thread_id, - /*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(), session_configuration.provider.clone(), session_configuration.session_source.clone(), - session_configuration.parent_thread_id, config.model_verbosity, config.features.enabled(Feature::EnableRequestCompression), config.features.enabled(Feature::RuntimeMetrics), @@ -7074,12 +7068,9 @@ where attestation_provider: None, model_client: ModelClient::new( Some(Arc::clone(&auth_manager)), - thread_id.into(), thread_id, - /*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(), session_configuration.provider.clone(), session_configuration.session_source.clone(), - session_configuration.parent_thread_id, config.model_verbosity, config.features.enabled(Feature::EnableRequestCompression), config.features.enabled(Feature::RuntimeMetrics), diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 7e429dcac..7bb34aec8 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -35,6 +35,8 @@ use crate::mentions::collect_explicit_app_ids; use crate::mentions::collect_explicit_plugin_mentions; use crate::mentions::collect_tool_mentions_from_messages; use crate::plugins::build_plugin_injections; +use crate::responses_metadata::CodexResponsesMetadata; +use crate::responses_metadata::CodexResponsesRequestKind; use crate::responses_retry::ResponsesStreamRequest; use crate::responses_retry::handle_retryable_response_stream_error; use crate::session::PreviousTurnSettings; @@ -223,9 +225,11 @@ pub(crate) async fn run_turn( .await; let window_id = sess.current_window_id().await; - let turn_metadata_header = turn_context - .turn_metadata_state - .current_header_value_for_model_request(&window_id); + let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata( + sess.installation_id.clone(), + window_id, + CodexResponsesRequestKind::Turn, + ); let tokens_before_sampling = sess.get_total_token_usage().await; match run_sampling_request( Arc::clone(&sess), @@ -233,8 +237,7 @@ pub(crate) async fn run_turn( Arc::clone(&turn_extension_data), Arc::clone(&turn_diff_tracker), &mut client_session, - &window_id, - turn_metadata_header.as_deref(), + &responses_metadata, sampling_request_input.clone(), cancellation_token.child_token(), ) @@ -1029,8 +1032,7 @@ async fn run_sampling_request( turn_store: Arc, turn_diff_tracker: SharedTurnDiffTracker, client_session: &mut ModelClientSession, - window_id: &str, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, input: Vec, cancellation_token: CancellationToken, ) -> CodexResult { @@ -1073,8 +1075,7 @@ async fn run_sampling_request( Arc::clone(&turn_context), Arc::clone(&turn_store), client_session, - window_id, - turn_metadata_header, + responses_metadata, Arc::clone(&turn_diff_tracker), &prompt, cancellation_token.child_token(), @@ -1803,8 +1804,7 @@ async fn try_run_sampling_request( turn_context: Arc, turn_store: Arc, client_session: &mut ModelClientSession, - window_id: &str, - turn_metadata_header: Option<&str>, + responses_metadata: &CodexResponsesMetadata, turn_diff_tracker: SharedTurnDiffTracker, prompt: &Prompt, cancellation_token: CancellationToken, @@ -1825,14 +1825,13 @@ async fn try_run_sampling_request( let sampling_timing_guard = turn_context.turn_timing_state.begin_sampling(); let mut stream = client_session .stream( - window_id, prompt, &turn_context.model_info, &turn_context.session_telemetry, turn_context.reasoning_effort.clone(), turn_context.reasoning_summary, turn_context.config.service_tier.clone(), - turn_metadata_header, + responses_metadata, &inference_trace, ) .instrument(trace_span!("stream_request")) diff --git a/codex-rs/core/src/session/turn_context.rs b/codex-rs/core/src/session/turn_context.rs index 1ee432370..5e5d5b760 100644 --- a/codex-rs/core/src/session/turn_context.rs +++ b/codex-rs/core/src/session/turn_context.rs @@ -517,7 +517,6 @@ impl Session { session_configuration.forked_from_thread_id, session_configuration.parent_thread_id, &session_configuration.session_source, - session_configuration.thread_source.clone(), sub_id.clone(), cwd.clone(), &session_configuration.permission_profile(), diff --git a/codex-rs/core/src/session_startup_prewarm.rs b/codex-rs/core/src/session_startup_prewarm.rs index fda498254..4fee5e774 100644 --- a/codex-rs/core/src/session_startup_prewarm.rs +++ b/codex-rs/core/src/session_startup_prewarm.rs @@ -9,6 +9,7 @@ use tracing::info; use tracing::warn; use crate::client::ModelClientSession; +use crate::responses_metadata::CodexResponsesRequestKind; use crate::session::INITIAL_SUBMIT_ID; use crate::session::session::Session; use crate::session::turn::build_prompt; @@ -267,21 +268,24 @@ async fn schedule_startup_prewarm_inner( /*status*/ None, ); let window_id = session.current_window_id().await; - let startup_turn_metadata_header = startup_turn_context + let responses_metadata = startup_turn_context .turn_metadata_state - .current_header_value_for_prewarm(&window_id); + .to_responses_metadata( + session.installation_id.clone(), + window_id, + CodexResponsesRequestKind::Prewarm, + ); let mut client_session = session.services.model_client.new_session(); let websocket_warmup_started_at = Instant::now(); client_session .prewarm_websocket( - &window_id, &startup_prompt, &startup_turn_context.model_info, &startup_turn_context.session_telemetry, startup_turn_context.reasoning_effort.clone(), startup_turn_context.reasoning_summary, startup_turn_context.config.service_tier.clone(), - startup_turn_metadata_header.as_deref(), + &responses_metadata, ) .await?; startup_turn_context.session_telemetry.record_startup_phase( diff --git a/codex-rs/core/src/test_support.rs b/codex-rs/core/src/test_support.rs index 7c084838f..60e80ca29 100644 --- a/codex-rs/core/src/test_support.rs +++ b/codex-rs/core/src/test_support.rs @@ -20,13 +20,19 @@ use codex_models_manager::collaboration_mode_presets; use codex_models_manager::manager::SharedModelsManager; use codex_models_manager::test_support::construct_model_info_offline_for_tests; use codex_models_manager::test_support::get_model_offline_for_tests; +use codex_protocol::ThreadId; use codex_protocol::config_types::CollaborationModeMask; use codex_protocol::openai_models::ModelInfo; use codex_protocol::openai_models::ModelPreset; +use codex_protocol::protocol::SessionSource; use once_cell::sync::Lazy; use crate::ThreadManager; use crate::config::Config; +use crate::responses_metadata::CodexResponsesMetadata; +use crate::responses_metadata::CodexResponsesRequestKind; +use crate::responses_metadata::subagent_header_value; +use crate::responses_metadata::subagent_metadata_kind; use crate::thread_manager; use crate::unified_exec; @@ -146,6 +152,44 @@ pub fn construct_model_info_offline(model: &str, config: &Config) -> ModelInfo { construct_model_info_offline_for_tests(model, &config.to_models_manager_config()) } +#[derive(Clone, Copy)] +pub enum TestCodexResponsesRequestKind { + Turn, + Prewarm, + WebsocketConnection, +} + +#[allow(clippy::too_many_arguments)] +pub fn responses_metadata( + installation_id: &str, + session_id: &str, + thread_id: &str, + turn_id: Option<&str>, + window_id: String, + session_source: &SessionSource, + parent_thread_id: Option, + request_kind: TestCodexResponsesRequestKind, +) -> CodexResponsesMetadata { + let request_kind = match request_kind { + TestCodexResponsesRequestKind::Turn => Some(CodexResponsesRequestKind::Turn), + TestCodexResponsesRequestKind::Prewarm => Some(CodexResponsesRequestKind::Prewarm), + TestCodexResponsesRequestKind::WebsocketConnection => None, + }; + CodexResponsesMetadata { + turn_id: request_kind.and(turn_id.map(ToString::to_string)), + request_kind, + parent_thread_id, + subagent_header: subagent_header_value(session_source), + subagent_kind: request_kind.and_then(|_| subagent_metadata_kind(session_source)), + ..CodexResponsesMetadata::new( + installation_id.to_string(), + session_id.to_string(), + thread_id.to_string(), + window_id, + ) + } +} + pub fn all_model_presets() -> &'static Vec { &TEST_MODEL_PRESETS } diff --git a/codex-rs/core/src/turn_metadata.rs b/codex-rs/core/src/turn_metadata.rs index e66a1fbb5..d13c08181 100644 --- a/codex-rs/core/src/turn_metadata.rs +++ b/codex-rs/core/src/turn_metadata.rs @@ -6,16 +6,15 @@ use std::sync::RwLock; use std::sync::atomic::AtomicBool; use std::sync::atomic::Ordering; -use codex_analytics::CompactionImplementation; -use codex_analytics::CompactionPhase; -use codex_analytics::CompactionReason; -use codex_analytics::CompactionStrategy; -use codex_analytics::CompactionTrigger; -use codex_utils_string::to_ascii_json_string; -use serde::Serialize; use serde_json::Value; use tokio::task::JoinHandle; +use crate::responses_metadata::CodexResponsesMetadata; +use crate::responses_metadata::CodexResponsesRequestKind; +use crate::responses_metadata::TurnMetadataWorkspace; +use crate::responses_metadata::filter_extra_metadata; +use crate::responses_metadata::subagent_header_value; +use crate::responses_metadata::subagent_metadata_kind; use crate::sandbox_tags::permission_profile_sandbox_tag; use codex_git_utils::get_git_remote_urls_assume_git_repo; use codex_git_utils::get_git_repo_root; @@ -26,62 +25,18 @@ use codex_protocol::config_types::WindowsSandboxLevel; use codex_protocol::models::PermissionProfile; use codex_protocol::openai_models::ReasoningEffort as ReasoningEffortConfig; use codex_protocol::protocol::SessionSource; -use codex_protocol::protocol::ThreadSource; use codex_utils_absolute_path::AbsolutePathBuf; const MODEL_KEY: &str = "model"; const REASONING_EFFORT_KEY: &str = "reasoning_effort"; -const TURN_STARTED_AT_UNIX_MS_KEY: &str = "turn_started_at_unix_ms"; const USER_INPUT_REQUESTED_DURING_TURN_KEY: &str = "user_input_requested_during_turn"; const WORKSPACE_KIND_KEY: &str = "workspace_kind"; -const REQUEST_KIND_KEY: &str = "request_kind"; -const COMPACTION_KEY: &str = "compaction"; -const WINDOW_ID_KEY: &str = "window_id"; pub(crate) struct McpTurnMetadataContext<'a> { pub(crate) model: &'a str, pub(crate) reasoning_effort: Option, } -/// Metadata present only on outbound model requests that perform compaction. -/// -/// These fields describe the operation at dispatch time. Post-response outcomes such as status, -/// error, duration, and token deltas remain in compaction analytics events. -#[derive(Clone, Copy, Debug, Serialize)] -pub(crate) struct CompactionTurnMetadata { - trigger: CompactionTrigger, - reason: CompactionReason, - implementation: CompactionImplementation, - phase: CompactionPhase, - strategy: CompactionStrategy, -} - -impl CompactionTurnMetadata { - pub(crate) fn new( - trigger: CompactionTrigger, - reason: CompactionReason, - implementation: CompactionImplementation, - phase: CompactionPhase, - ) -> Self { - Self { - trigger, - reason, - implementation, - phase, - strategy: CompactionStrategy::Memento, - } - } -} - -#[derive(Clone, Copy, Debug, Serialize)] -#[serde(rename_all = "snake_case")] -enum TurnMetadataRequestKind { - Turn, - Prewarm, - Compaction, - Memory, -} - #[derive(Clone, Debug, Default)] struct WorkspaceGitMetadata { associated_remote_urls: Option>, @@ -97,16 +52,6 @@ impl WorkspaceGitMetadata { } } -#[derive(Clone, Debug, Serialize, Default)] -struct TurnMetadataWorkspace { - #[serde(default, skip_serializing_if = "Option::is_none")] - associated_remote_urls: Option>, - #[serde(default, skip_serializing_if = "Option::is_none")] - latest_git_commit_hash: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - has_changes: Option, -} - impl From for TurnMetadataWorkspace { fn from(value: WorkspaceGitMetadata) -> Self { Self { @@ -117,141 +62,40 @@ impl From for TurnMetadataWorkspace { } } -/// Base payload for the outbound model request `x-codex-turn-metadata` header. -/// -/// Turn-owned state populates identity fields, including optional fork and subagent lineage. A -/// concrete request kind is added at outbound model dispatch so turns, startup prewarm, and -/// compaction remain distinguishable. Detached memory requests are constructed as `memory` -/// directly. -#[derive(Clone, Debug, Serialize)] -pub(crate) struct TurnMetadataBag { - #[serde(default, skip_serializing_if = "Option::is_none")] - request_kind: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - session_id: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - thread_id: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - forked_from_thread_id: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - parent_thread_id: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - subagent_kind: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - thread_source: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - turn_id: Option, - #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] - workspaces: BTreeMap, - #[serde(default, skip_serializing_if = "Option::is_none")] - sandbox: Option, -} - -impl TurnMetadataBag { - fn with_workspace_git_metadata( - mut self, - repo_root: Option, - workspace_git_metadata: Option, - ) -> Self { - if let (Some(repo_root), Some(workspace_git_metadata)) = (repo_root, workspace_git_metadata) - && !workspace_git_metadata.is_empty() - { - self.workspaces - .insert(repo_root, workspace_git_metadata.into()); - } - self - } - - fn to_header_value(&self) -> Option { - to_ascii_json_string(self).ok() - } -} - -fn merge_turn_metadata( - header: &str, - turn_started_at_unix_ms: Option, - responsesapi_client_metadata: Option<&HashMap>, -) -> Option { - if turn_started_at_unix_ms.is_none() && responsesapi_client_metadata.is_none() { - return None; - } - - let mut metadata = serde_json::from_str::>(header).ok()?; - if let Some(turn_started_at_unix_ms) = turn_started_at_unix_ms { - metadata.insert( - TURN_STARTED_AT_UNIX_MS_KEY.to_string(), - Value::Number(turn_started_at_unix_ms.into()), - ); - } - if let Some(responsesapi_client_metadata) = responsesapi_client_metadata { - for (key, value) in responsesapi_client_metadata { - if matches!( - key.as_str(), - "session_id" - | "thread_id" - | "turn_id" - | TURN_STARTED_AT_UNIX_MS_KEY - | "forked_from_thread_id" - | "parent_thread_id" - | "subagent_kind" - | REQUEST_KIND_KEY - | COMPACTION_KEY - | WINDOW_ID_KEY - ) { - continue; - } - metadata - .entry(key.clone()) - .or_insert_with(|| Value::String(value.clone())); - } - } - to_ascii_json_string(&metadata).ok() -} - -pub async fn build_turn_metadata_header( +#[allow(clippy::too_many_arguments)] +pub async fn detached_memory_responses_metadata( + installation_id: String, + session_id: String, + thread_id: String, + window_id: String, + session_source: &SessionSource, cwd: &AbsolutePathBuf, sandbox: Option<&str>, -) -> Option { - let repo_root = get_git_repo_root(cwd).map(|root| root.to_string_lossy().into_owned()); - - let (head_commit_hash, associated_remote_urls, has_changes) = tokio::join!( - get_head_commit_hash(cwd), - get_git_remote_urls_assume_git_repo(cwd), - get_has_changes(cwd), - ); - let latest_git_commit_hash = head_commit_hash.map(|sha| sha.0); - TurnMetadataBag { - request_kind: Some(TurnMetadataRequestKind::Memory), - session_id: None, - thread_id: None, - forked_from_thread_id: None, - parent_thread_id: None, - subagent_kind: None, - thread_source: None, - turn_id: None, - workspaces: BTreeMap::new(), +) -> CodexResponsesMetadata { + CodexResponsesMetadata { + request_kind: Some(CodexResponsesRequestKind::Memory), + subagent_header: subagent_header_value(session_source), sandbox: sandbox.map(ToString::to_string), + workspaces: memory_workspaces(cwd).await, + ..CodexResponsesMetadata::new(installation_id, session_id, thread_id, window_id) } - .with_workspace_git_metadata( - repo_root, - Some(WorkspaceGitMetadata { - associated_remote_urls, - latest_git_commit_hash, - has_changes, - }), - ) - .to_header_value() } #[derive(Clone, Debug)] pub(crate) struct TurnMetadataState { cwd: AbsolutePathBuf, repo_root: Option, - base_metadata: TurnMetadataBag, - base_header: Option, - enriched_header: Arc>>, + session_id: String, + thread_id: String, + forked_from_thread_id: Option, + parent_thread_id: Option, + subagent_header: Option, + subagent_kind: Option, + turn_id: String, + sandbox: Option, + enriched_workspaces: Arc>>>, turn_started_at_unix_ms: Arc>>, - responsesapi_client_metadata: Arc>>>, + responsesapi_client_metadata: Arc>>, user_input_requested_during_turn: Arc, enrichment_task: Arc>>>, } @@ -264,7 +108,6 @@ impl TurnMetadataState { forked_from_thread_id: Option, parent_thread_id: Option, session_source: &SessionSource, - thread_source: Option, turn_id: String, cwd: AbsolutePathBuf, permission_profile: &PermissionProfile, @@ -280,79 +123,34 @@ impl TurnMetadataState { ) .to_string(), ); - let subagent_kind = match session_source { - SessionSource::SubAgent(subagent_source) => Some(subagent_source.kind().to_string()), - SessionSource::Cli - | SessionSource::VSCode - | SessionSource::Exec - | SessionSource::Mcp - | SessionSource::Custom(_) - | SessionSource::Internal(_) - | SessionSource::Unknown => None, - }; - let base_metadata = TurnMetadataBag { - request_kind: None, - session_id: Some(session_id), - thread_id: Some(thread_id), - forked_from_thread_id, - parent_thread_id, - subagent_kind, - thread_source, - turn_id: Some(turn_id), - workspaces: BTreeMap::new(), - sandbox, - }; - let base_header = base_metadata.to_header_value(); - Self { cwd, repo_root, - base_metadata, - base_header, - enriched_header: Arc::new(RwLock::new(None)), + session_id, + thread_id, + forked_from_thread_id, + parent_thread_id, + subagent_header: subagent_header_value(session_source), + subagent_kind: subagent_metadata_kind(session_source), + turn_id, + sandbox, + enriched_workspaces: Arc::new(RwLock::new(None)), turn_started_at_unix_ms: Arc::new(RwLock::new(None)), - responsesapi_client_metadata: Arc::new(RwLock::new(None)), + responsesapi_client_metadata: Arc::new(RwLock::new(BTreeMap::new())), user_input_requested_during_turn: Arc::new(AtomicBool::new(false)), enrichment_task: Arc::new(Mutex::new(None)), } } - pub(crate) fn current_header_value(&self) -> Option { - let header = if let Some(header) = self - .enriched_header - .read() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .as_ref() - .cloned() - { - header - } else { - self.base_header.clone()? - }; - let turn_started_at_unix_ms = *self - .turn_started_at_unix_ms - .read() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let responsesapi_client_metadata = self - .responsesapi_client_metadata - .read() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(); - merge_turn_metadata( - &header, - turn_started_at_unix_ms, - responsesapi_client_metadata.as_ref(), - ) - .or(Some(header)) - } - pub(crate) fn current_meta_value_for_mcp_request( &self, context: McpTurnMetadataContext<'_>, ) -> Option { - let header = self.current_header_value()?; - let mut metadata = serde_json::from_str::>(&header).ok()?; - metadata.remove(REQUEST_KIND_KEY); + let Value::Object(mut metadata) = + self.responses_metadata_template().turn_metadata_value()? + else { + return None; + }; metadata.insert( MODEL_KEY.to_string(), Value::String(context.model.to_string()), @@ -382,50 +180,18 @@ impl TurnMetadataState { Some(Value::Object(metadata)) } - fn current_header_value_for_model_request_kind( + pub(crate) fn to_responses_metadata( &self, - window_id: &str, - request_kind: TurnMetadataRequestKind, - ) -> Option { - let header = self.current_header_value()?; - let mut metadata = serde_json::from_str::>(&header).ok()?; - metadata.insert( - REQUEST_KIND_KEY.to_string(), - serde_json::to_value(request_kind).ok()?, - ); - metadata.insert( - WINDOW_ID_KEY.to_string(), - Value::String(window_id.to_string()), - ); - to_ascii_json_string(&metadata).ok() - } - - pub(crate) fn current_header_value_for_model_request(&self, window_id: &str) -> Option { - self.current_header_value_for_model_request_kind(window_id, TurnMetadataRequestKind::Turn) - } - - pub(crate) fn current_header_value_for_prewarm(&self, window_id: &str) -> Option { - self.current_header_value_for_model_request_kind( + installation_id: String, + window_id: String, + request_kind: CodexResponsesRequestKind, + ) -> CodexResponsesMetadata { + CodexResponsesMetadata { + installation_id, window_id, - TurnMetadataRequestKind::Prewarm, - ) - } - - pub(crate) fn current_header_value_for_compaction( - &self, - window_id: &str, - compaction: CompactionTurnMetadata, - ) -> Option { - let header = self.current_header_value_for_model_request_kind( - window_id, - TurnMetadataRequestKind::Compaction, - )?; - let mut metadata = serde_json::from_str::>(&header).ok()?; - metadata.insert( - COMPACTION_KEY.to_string(), - serde_json::to_value(compaction).ok()?, - ); - to_ascii_json_string(&metadata).ok() + request_kind: Some(request_kind), + ..self.responses_metadata_template() + } } pub(crate) fn mark_user_input_requested_during_turn(&self) { @@ -441,15 +207,54 @@ impl TurnMetadataState { .responsesapi_client_metadata .write() .unwrap_or_else(std::sync::PoisonError::into_inner) = - Some(responsesapi_client_metadata); + filter_extra_metadata(responsesapi_client_metadata); } pub(crate) fn workspace_kind(&self) -> Option { self.responsesapi_client_metadata .read() .unwrap_or_else(std::sync::PoisonError::into_inner) - .as_ref() - .and_then(|metadata| metadata.get(WORKSPACE_KIND_KEY).cloned()) + .get(WORKSPACE_KIND_KEY) + .cloned() + } + + fn responses_metadata_template(&self) -> CodexResponsesMetadata { + CodexResponsesMetadata { + turn_id: Some(self.turn_id.clone()), + forked_from_thread_id: self.forked_from_thread_id, + parent_thread_id: self.parent_thread_id, + subagent_header: self.subagent_header.clone(), + subagent_kind: self.subagent_kind.clone(), + sandbox: self.sandbox.clone(), + workspaces: self.current_workspaces(), + turn_started_at_unix_ms: self.current_turn_started_at_unix_ms(), + extra: self + .responsesapi_client_metadata + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone(), + ..CodexResponsesMetadata::new( + String::new(), + self.session_id.clone(), + self.thread_id.clone(), + String::new(), + ) + } + } + + fn current_workspaces(&self) -> BTreeMap { + self.enriched_workspaces + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + .unwrap_or_default() + } + + fn current_turn_started_at_unix_ms(&self) -> Option { + *self + .turn_started_at_unix_ms + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) } pub(crate) fn set_turn_started_at_unix_ms(&self, turn_started_at_unix_ms: i64) { @@ -479,20 +284,16 @@ impl TurnMetadataState { return; }; - let enriched_metadata = state - .base_metadata - .clone() - .with_workspace_git_metadata(Some(repo_root), Some(workspace_git_metadata)); - if enriched_metadata.workspaces.is_empty() { + if workspace_git_metadata.is_empty() { return; } - if let Some(header_value) = enriched_metadata.to_header_value() { - *state - .enriched_header - .write() - .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(header_value); - } + let mut workspaces = BTreeMap::new(); + workspaces.insert(repo_root, workspace_git_metadata.into()); + *state + .enriched_workspaces + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(workspaces); })); } @@ -522,6 +323,27 @@ impl TurnMetadataState { } } +async fn memory_workspaces(cwd: &AbsolutePathBuf) -> BTreeMap { + let repo_root = get_git_repo_root(cwd).map(|root| root.to_string_lossy().into_owned()); + let (head_commit_hash, associated_remote_urls, has_changes) = tokio::join!( + get_head_commit_hash(cwd), + get_git_remote_urls_assume_git_repo(cwd), + get_has_changes(cwd), + ); + let workspace_git_metadata = WorkspaceGitMetadata { + associated_remote_urls, + latest_git_commit_hash: head_commit_hash.map(|sha| sha.0), + has_changes, + }; + let mut workspaces = BTreeMap::new(); + if let Some(repo_root) = repo_root + && !workspace_git_metadata.is_empty() + { + workspaces.insert(repo_root, workspace_git_metadata.into()); + } + workspaces +} + #[cfg(test)] #[path = "turn_metadata_tests.rs"] mod tests; diff --git a/codex-rs/core/src/turn_metadata_tests.rs b/codex-rs/core/src/turn_metadata_tests.rs index 32a559848..5e6d1f32b 100644 --- a/codex-rs/core/src/turn_metadata_tests.rs +++ b/codex-rs/core/src/turn_metadata_tests.rs @@ -1,11 +1,18 @@ use super::*; +use crate::responses_metadata::CodexResponsesRequestKind; +use crate::responses_metadata::CompactionTurnMetadata; +use crate::responses_metadata::INSTALLATION_ID_KEY; +use crate::responses_metadata::WINDOW_ID_KEY; use crate::sandbox_tags::permission_profile_sandbox_tag; +use codex_analytics::CompactionImplementation; +use codex_analytics::CompactionPhase; +use codex_analytics::CompactionReason; +use codex_analytics::CompactionTrigger; use codex_protocol::models::PermissionProfile; use codex_protocol::openai_models::ReasoningEffort as ReasoningEffortConfig; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; -use codex_protocol::protocol::ThreadSource; use codex_utils_absolute_path::AbsolutePathBuf; use core_test_support::PathBufExt; use core_test_support::PathExt; @@ -23,6 +30,44 @@ fn test_mcp_turn_metadata_context() -> McpTurnMetadataContext<'static> { } } +fn test_responses_metadata_json( + state: &TurnMetadataState, + window_id: &str, + request_kind: CodexResponsesRequestKind, +) -> String { + state + .to_responses_metadata( + "installation-a".to_string(), + window_id.to_string(), + request_kind, + ) + .turn_metadata_json() + .expect("turn metadata json") +} + +fn test_turn_responses_metadata_json(state: &TurnMetadataState, window_id: &str) -> String { + test_responses_metadata_json(state, window_id, CodexResponsesRequestKind::Turn) +} + +fn test_compaction_responses_metadata_json( + state: &TurnMetadataState, + window_id: &str, + compaction: CompactionTurnMetadata, +) -> String { + test_responses_metadata_json( + state, + window_id, + CodexResponsesRequestKind::Compaction(compaction), + ) +} + +fn test_turn_metadata_header(state: &TurnMetadataState) -> String { + state + .responses_metadata_template() + .turn_metadata_json() + .expect("header") +} + async fn create_clean_git_repo(repo_name: &str) -> (TempDir, AbsolutePathBuf) { let temp_dir = TempDir::new().expect("temp dir"); let repo_path = temp_dir.path().join(repo_name).abs(); @@ -64,12 +109,21 @@ async fn create_clean_git_repo(repo_name: &str) -> (TempDir, AbsolutePathBuf) { } #[tokio::test] -async fn build_turn_metadata_header_marks_detached_memory_without_turn_identity() { +async fn detached_memory_responses_metadata_omits_turn_identity() { let (_temp_dir, repo_path) = create_clean_git_repo("repo-東京").await; - let header = build_turn_metadata_header(&repo_path, Some("none")) - .await - .expect("header"); + let header = detached_memory_responses_metadata( + String::new(), + String::new(), + String::new(), + String::new(), + &SessionSource::Unknown, + &repo_path, + Some("none"), + ) + .await + .turn_metadata_json() + .expect("header"); assert!(header.is_ascii()); assert!(!header.contains("東京")); let parsed: Value = serde_json::from_str(&header).expect("valid json"); @@ -100,13 +154,22 @@ async fn build_turn_metadata_header_marks_detached_memory_without_turn_identity( } #[tokio::test] -async fn build_turn_metadata_header_marks_memory_without_workspace_metadata() { +async fn detached_memory_responses_metadata_omits_empty_workspace_metadata() { let temp_dir = TempDir::new().expect("temp dir"); let cwd = temp_dir.path().abs(); - let header = build_turn_metadata_header(&cwd, /*sandbox*/ None) - .await - .expect("detached memory should emit its request kind"); + let header = detached_memory_responses_metadata( + String::new(), + String::new(), + String::new(), + String::new(), + &SessionSource::Unknown, + &cwd, + /*sandbox*/ None, + ) + .await + .turn_metadata_json() + .expect("detached memory should emit its request kind"); let parsed: Value = serde_json::from_str(&header).expect("valid json"); assert_eq!(parsed, serde_json::json!({"request_kind": "memory"})); @@ -124,7 +187,6 @@ fn turn_metadata_state_uses_platform_sandbox_tag() { /*forked_from_thread_id*/ None, /*parent_thread_id*/ None, &SessionSource::Exec, - Some(ThreadSource::User), "turn-a".to_string(), cwd, &permission_profile, @@ -132,12 +194,11 @@ fn turn_metadata_state_uses_platform_sandbox_tag() { /*enforce_managed_network*/ false, ); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); let sandbox_name = json.get("sandbox").and_then(Value::as_str); let session_id = json.get("session_id").and_then(Value::as_str); let thread_id = json.get("thread_id").and_then(Value::as_str); - let thread_source = json.get("thread_source").and_then(Value::as_str); assert!(json.get("request_kind").is_none()); let expected_sandbox = permission_profile_sandbox_tag( @@ -148,39 +209,12 @@ fn turn_metadata_state_uses_platform_sandbox_tag() { assert_eq!(sandbox_name, Some(expected_sandbox)); assert_eq!(session_id, Some("session-a")); assert_eq!(thread_id, Some("thread-a")); - assert_eq!(thread_source, Some("user")); assert!(json.get("forked_from_thread_id").is_none()); assert!(json.get("parent_thread_id").is_none()); assert!(json.get("subagent_kind").is_none()); assert!(json.get("session_source").is_none()); } -#[test] -fn turn_metadata_state_uses_explicit_subagent_thread_source() { - let temp_dir = TempDir::new().expect("temp dir"); - let cwd = temp_dir.path().abs(); - let permission_profile = PermissionProfile::read_only(); - let state = TurnMetadataState::new( - "session-a".to_string(), - "thread-a".to_string(), - /*forked_from_thread_id*/ None, - /*parent_thread_id*/ None, - &SessionSource::Exec, - Some(ThreadSource::Subagent), - "turn-a".to_string(), - cwd, - &permission_profile, - WindowsSandboxLevel::Disabled, - /*enforce_managed_network*/ false, - ); - - let header = state.current_header_value().expect("header"); - let json: Value = serde_json::from_str(&header).expect("json"); - - assert_eq!(json["thread_source"].as_str(), Some("subagent")); - assert!(json.get("session_source").is_none()); -} - #[test] fn turn_metadata_state_includes_root_fork_lineage() { let temp_dir = TempDir::new().expect("temp dir"); @@ -195,7 +229,6 @@ fn turn_metadata_state_includes_root_fork_lineage() { Some(source_thread_id), /*parent_thread_id*/ None, &SessionSource::Exec, - Some(ThreadSource::User), "turn-a".to_string(), cwd, &permission_profile, @@ -203,7 +236,7 @@ fn turn_metadata_state_includes_root_fork_lineage() { /*enforce_managed_network*/ false, ); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); assert_eq!( @@ -234,7 +267,6 @@ fn turn_metadata_state_includes_thread_spawn_subagent_parent_without_fork() { agent_nickname: None, agent_role: None, }), - Some(ThreadSource::Subagent), "turn-a".to_string(), cwd, &permission_profile, @@ -242,7 +274,7 @@ fn turn_metadata_state_includes_thread_spawn_subagent_parent_without_fork() { /*enforce_managed_network*/ false, ); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); assert!(json.get("forked_from_thread_id").is_none()); @@ -273,7 +305,6 @@ fn turn_metadata_state_includes_forked_thread_spawn_subagent_lineage() { agent_nickname: None, agent_role: None, }), - Some(ThreadSource::Subagent), "turn-a".to_string(), cwd, &permission_profile, @@ -281,7 +312,7 @@ fn turn_metadata_state_includes_forked_thread_spawn_subagent_lineage() { /*enforce_managed_network*/ false, ); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); assert_eq!( @@ -318,7 +349,6 @@ fn turn_metadata_state_includes_known_parent_for_non_thread_spawn_subagents_with /*forked_from_thread_id*/ None, Some(parent_thread_id), &SessionSource::SubAgent(subagent_source), - Some(ThreadSource::Subagent), "turn-a".to_string(), cwd.clone(), &permission_profile, @@ -326,7 +356,7 @@ fn turn_metadata_state_includes_known_parent_for_non_thread_spawn_subagents_with /*enforce_managed_network*/ false, ); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); assert!(json.get("forked_from_thread_id").is_none()); @@ -350,7 +380,6 @@ fn turn_metadata_state_includes_turn_started_at_unix_ms_after_start() { /*forked_from_thread_id*/ None, /*parent_thread_id*/ None, &SessionSource::Exec, - Some(ThreadSource::User), "turn-a".to_string(), cwd, &permission_profile, @@ -359,7 +388,7 @@ fn turn_metadata_state_includes_turn_started_at_unix_ms_after_start() { ); state.set_turn_started_at_unix_ms(/*turn_started_at_unix_ms*/ 1_700_000_000_123); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); assert_eq!( @@ -380,7 +409,6 @@ fn turn_metadata_state_includes_model_and_reasoning_effort_only_in_request_meta( /*forked_from_thread_id*/ None, /*parent_thread_id*/ None, &SessionSource::Exec, - /*thread_source*/ None, "turn-a".to_string(), cwd, &permission_profile, @@ -388,7 +416,7 @@ fn turn_metadata_state_includes_model_and_reasoning_effort_only_in_request_meta( /*enforce_managed_network*/ false, ); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let header_json: Value = serde_json::from_str(&header).expect("json"); assert!(header_json.get("model").is_none()); assert!(header_json.get("reasoning_effort").is_none()); @@ -429,7 +457,6 @@ fn turn_metadata_state_marks_user_input_requested_during_turn_only_for_mcp_reque /*forked_from_thread_id*/ None, /*parent_thread_id*/ None, &SessionSource::Exec, - /*thread_source*/ None, "turn-a".to_string(), cwd, &permission_profile, @@ -437,7 +464,7 @@ fn turn_metadata_state_marks_user_input_requested_during_turn_only_for_mcp_reque /*enforce_managed_network*/ false, ); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let header_json: Value = serde_json::from_str(&header).expect("json"); assert!( header_json @@ -452,7 +479,7 @@ fn turn_metadata_state_marks_user_input_requested_during_turn_only_for_mcp_reque state.mark_user_input_requested_during_turn(); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let header_json: Value = serde_json::from_str(&header).expect("json"); assert!( header_json @@ -482,7 +509,6 @@ fn turn_metadata_state_ignores_client_reserved_metadata_before_start() { /*forked_from_thread_id*/ None, /*parent_thread_id*/ None, &SessionSource::Exec, - Some(ThreadSource::User), "turn-a".to_string(), cwd, &permission_profile, @@ -505,7 +531,7 @@ fn turn_metadata_state_ignores_client_reserved_metadata_before_start() { ("subagent_kind".to_string(), "client-supplied".to_string()), ])); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); assert!(json.get("turn_started_at_unix_ms").is_none()); @@ -536,7 +562,6 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( agent_nickname: None, agent_role: None, }), - Some(ThreadSource::User), "turn-a".to_string(), cwd, &permission_profile, @@ -554,6 +579,19 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( ), ("session_id".to_string(), "client-supplied".to_string()), ("thread_id".to_string(), "client-supplied".to_string()), + ("installation_id".to_string(), "client-supplied".to_string()), + ( + "x-codex-installation-id".to_string(), + "client-supplied".to_string(), + ), + ( + "x-codex-parent-thread-id".to_string(), + "client-supplied".to_string(), + ), + ( + "x-openai-subagent".to_string(), + "client-supplied".to_string(), + ), ( "forked_from_thread_id".to_string(), "client-supplied".to_string(), @@ -574,7 +612,7 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( ])); state.set_turn_started_at_unix_ms(/*turn_started_at_unix_ms*/ 1_700_000_000_123); - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); assert!(header.is_ascii()); assert!(!header.contains("東京")); let json: Value = serde_json::from_str(&header).expect("json"); @@ -586,6 +624,10 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( assert_eq!(json["reasoning_effort"].as_str(), Some("client-supplied")); assert_eq!(json["session_id"].as_str(), Some("session-a")); assert_eq!(json["thread_id"].as_str(), Some("thread-a")); + assert!(json.get(INSTALLATION_ID_KEY).is_none()); + assert!(json.get("x-codex-installation-id").is_none()); + assert!(json.get("x-codex-parent-thread-id").is_none()); + assert!(json.get("x-openai-subagent").is_none()); assert_eq!( json["forked_from_thread_id"].as_str(), Some("44444444-4444-4444-8444-444444444444") @@ -595,7 +637,7 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( Some("55555555-5555-4555-8555-555555555555") ); assert_eq!(json["subagent_kind"].as_str(), Some("thread_spawn")); - assert_eq!(json["thread_source"].as_str(), Some("user")); + assert_eq!(json["thread_source"].as_str(), Some("client-supplied")); assert_eq!(json["turn_id"].as_str(), Some("turn-a")); assert!(json.get("request_kind").is_none()); assert!(json.get(WINDOW_ID_KEY).is_none()); @@ -604,12 +646,14 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( Some(1_700_000_000_123) ); - let model_request_header = state - .current_header_value_for_model_request("thread-a:1") - .expect("model request header"); + let model_request_header = test_turn_responses_metadata_json(&state, "thread-a:1"); let model_request_json: Value = serde_json::from_str(&model_request_header).expect("model request json"); assert_eq!(model_request_json["request_kind"].as_str(), Some("turn")); + assert_eq!( + model_request_json[INSTALLATION_ID_KEY].as_str(), + Some("installation-a") + ); assert_eq!( model_request_json[WINDOW_ID_KEY].as_str(), Some("thread-a:1") @@ -635,7 +679,6 @@ fn turn_metadata_state_overlays_compaction_only_on_compaction_requests() { /*forked_from_thread_id*/ None, /*parent_thread_id*/ None, &SessionSource::Exec, - Some(ThreadSource::User), "turn-a".to_string(), cwd, &permission_profile, @@ -647,17 +690,16 @@ fn turn_metadata_state_overlays_compaction_only_on_compaction_requests() { "client-supplied".to_string(), )])); - let compact_header = state - .current_header_value_for_compaction( - "thread-a:2", - CompactionTurnMetadata::new( - CompactionTrigger::Auto, - CompactionReason::ContextLimit, - CompactionImplementation::ResponsesCompactionV2, - CompactionPhase::MidTurn, - ), - ) - .expect("compact header"); + let compact_header = test_compaction_responses_metadata_json( + &state, + "thread-a:2", + CompactionTurnMetadata::new( + CompactionTrigger::Auto, + CompactionReason::ContextLimit, + CompactionImplementation::ResponsesCompactionV2, + CompactionPhase::MidTurn, + ), + ); let compact_json: Value = serde_json::from_str(&compact_header).expect("json"); assert_eq!(compact_json["request_kind"].as_str(), Some("compaction")); assert_eq!(compact_json["turn_id"].as_str(), Some("turn-a")); @@ -673,9 +715,7 @@ fn turn_metadata_state_overlays_compaction_only_on_compaction_requests() { }) ); - let regular_header = state - .current_header_value_for_model_request("thread-a:3") - .expect("regular header"); + let regular_header = test_turn_responses_metadata_json(&state, "thread-a:3"); let regular_json: Value = serde_json::from_str(®ular_header).expect("json"); assert_eq!(regular_json["request_kind"].as_str(), Some("turn")); assert_eq!(regular_json[WINDOW_ID_KEY].as_str(), Some("thread-a:3")); @@ -701,7 +741,6 @@ async fn turn_metadata_state_preserves_lineage_after_git_enrichment() { agent_nickname: None, agent_role: None, }), - Some(ThreadSource::Subagent), "turn-a".to_string(), repo_path, &permission_profile, @@ -713,7 +752,7 @@ async fn turn_metadata_state_preserves_lineage_after_git_enrichment() { let json = tokio::time::timeout(Duration::from_secs(2), async { loop { - let header = state.current_header_value().expect("header"); + let header = test_turn_metadata_header(&state); let json: Value = serde_json::from_str(&header).expect("json"); if json .get("workspaces") diff --git a/codex-rs/core/tests/common/context_snapshot.rs b/codex-rs/core/tests/common/context_snapshot.rs index 8aaefbbf1..7a165214c 100644 --- a/codex-rs/core/tests/common/context_snapshot.rs +++ b/codex-rs/core/tests/common/context_snapshot.rs @@ -297,9 +297,11 @@ fn canonicalize_json_snapshot_value(value: &mut Value, options: &ContextSnapshot fn format_snapshot_json_string(text: &str, options: &ContextSnapshotOptions) -> String { let normalized = match options.render_mode { ContextSnapshotRenderMode::RedactedText - | ContextSnapshotRenderMode::KindWithTextPrefix { .. } => normalize_snapshot_uuids( - &normalize_snapshot_line_endings(&canonicalize_snapshot_text(text)), - ), + | ContextSnapshotRenderMode::KindWithTextPrefix { .. } => { + normalize_snapshot_dynamic_values(&normalize_snapshot_line_endings( + &canonicalize_snapshot_text(text), + )) + } ContextSnapshotRenderMode::FullText => normalize_snapshot_line_endings(text), ContextSnapshotRenderMode::KindOnly => unreachable!(), }; @@ -440,7 +442,7 @@ fn normalize_dynamic_snapshot_paths(text: &str) -> String { .into_owned() } -fn normalize_snapshot_uuids(text: &str) -> String { +fn normalize_snapshot_dynamic_values(text: &str) -> String { static UUID_RE: OnceLock = OnceLock::new(); let uuid_re = UUID_RE.get_or_init(|| { Regex::new( @@ -448,7 +450,20 @@ fn normalize_snapshot_uuids(text: &str) -> String { ) .expect("uuid regex should compile") }); - uuid_re.replace_all(text, "").into_owned() + static TURN_STARTED_AT_UNIX_MS_RE: OnceLock = OnceLock::new(); + let turn_started_at_unix_ms_re = TURN_STARTED_AT_UNIX_MS_RE.get_or_init(|| { + Regex::new(r#""turn_started_at_unix_ms":\d+"#) + .expect("turn_started_at_unix_ms regex should compile") + }); + static SANDBOX_RE: OnceLock = OnceLock::new(); + let sandbox_re = SANDBOX_RE + .get_or_init(|| Regex::new(r#""sandbox":"[^"]+""#).expect("sandbox regex should compile")); + let text = uuid_re.replace_all(text, ""); + let text = + turn_started_at_unix_ms_re.replace_all(&text, r#""turn_started_at_unix_ms":"#); + sandbox_re + .replace_all(&text, r#""sandbox":"""#) + .into_owned() } #[cfg(test)] @@ -456,6 +471,7 @@ mod tests { use super::ContextSnapshotOptions; use super::ContextSnapshotRenderMode; use super::format_response_items_snapshot; + use super::format_snapshot_json_string; use pretty_assertions::assert_eq; use serde_json::json; @@ -708,4 +724,17 @@ mod tests { "00:message/developer:## Skills\\n- openai-docs: helper (file: /openai-docs/SKILL.md)" ); } + + #[test] + fn redacted_text_mode_normalizes_turn_metadata_dynamic_json_strings() { + let rendered = format_snapshot_json_string( + r#"{"turn_id":"019eaded-ba5c-7d40-8a81-a4dcebc4679e","sandbox":"seccomp","turn_started_at_unix_ms":1781035793002}"#, + &ContextSnapshotOptions::default(), + ); + + assert_eq!( + rendered, + r#"{"turn_id":"","sandbox":"","turn_started_at_unix_ms":}"# + ); + } } diff --git a/codex-rs/core/tests/common/lib.rs b/codex-rs/core/tests/common/lib.rs index 94d72d619..f9732b1b5 100644 --- a/codex-rs/core/tests/common/lib.rs +++ b/codex-rs/core/tests/common/lib.rs @@ -15,6 +15,8 @@ use codex_core::CodexThread; use codex_core::config::Config; use codex_core::config::ConfigBuilder; use codex_core::config::ConfigOverrides; +pub use codex_core::test_support::TestCodexResponsesRequestKind; +pub use codex_core::test_support::responses_metadata; use codex_utils_absolute_path::AbsolutePathBuf; pub use codex_utils_absolute_path::test_support::PathBufExt; pub use codex_utils_absolute_path::test_support::PathExt; diff --git a/codex-rs/core/tests/responses_headers.rs b/codex-rs/core/tests/responses_headers.rs index af87762c0..cda5ffcdb 100644 --- a/codex-rs/core/tests/responses_headers.rs +++ b/codex-rs/core/tests/responses_headers.rs @@ -15,8 +15,10 @@ use codex_protocol::models::ContentItem; use codex_protocol::models::ResponseItem; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; +use core_test_support::TestCodexResponsesRequestKind; use core_test_support::load_default_config_for_test; use core_test_support::responses; +use core_test_support::responses_metadata as test_responses_metadata; use core_test_support::test_codex::test_codex; use futures::StreamExt; use pretty_assertions::assert_eq; @@ -32,6 +34,23 @@ fn normalize_git_remote_url(url: &str) -> String { } const TEST_INSTALLATION_ID: &str = "11111111-1111-4111-8111-111111111111"; +fn test_turn_responses_metadata( + _client: &ModelClient, + thread_id: ThreadId, + session_source: &SessionSource, +) -> codex_core::CodexResponsesMetadata { + let thread_id = thread_id.to_string(); + test_responses_metadata( + TEST_INSTALLATION_ID, + &thread_id, + &thread_id, + /*turn_id*/ None, + format!("{thread_id}:0"), + session_source, + /*parent_thread_id*/ None, + TestCodexResponsesRequestKind::Turn, + ) +} #[tokio::test] async fn responses_stream_includes_subagent_header_on_review() { @@ -101,18 +120,16 @@ async fn responses_stream_includes_subagent_header_on_review() { let client = ModelClient::new( /*auth_manager*/ None, - thread_id.into(), thread_id, - /*installation_id*/ TEST_INSTALLATION_ID.to_string(), provider.clone(), - session_source, - /*parent_thread_id*/ None, + session_source.clone(), config.model_verbosity, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, /*beta_features_header*/ None, /*attestation_provider*/ None, ); + let responses_metadata = test_turn_responses_metadata(&client, thread_id, &session_source); let mut client_session = client.new_session(); let mut prompt = Prompt::default(); @@ -127,14 +144,13 @@ async fn responses_stream_includes_subagent_header_on_review() { let mut stream = client_session .stream( - &expected_window_id, &prompt, &model_info, &session_telemetry, effort, summary.unwrap_or(model_info.default_reasoning_summary), /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -218,7 +234,6 @@ async fn responses_stream_includes_subagent_header_on_other() { let session_source = SessionSource::SubAgent(SubAgentSource::Other("my-task".to_string())); let model_info = codex_core::test_support::construct_model_info_offline(model.as_str(), &config); - let window_id = format!("{thread_id}:0"); let session_telemetry = SessionTelemetry::new( thread_id, @@ -235,18 +250,16 @@ async fn responses_stream_includes_subagent_header_on_other() { let client = ModelClient::new( /*auth_manager*/ None, - thread_id.into(), thread_id, - /*installation_id*/ TEST_INSTALLATION_ID.to_string(), provider.clone(), - session_source, - /*parent_thread_id*/ None, + session_source.clone(), config.model_verbosity, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, /*beta_features_header*/ None, /*attestation_provider*/ None, ); + let responses_metadata = test_turn_responses_metadata(&client, thread_id, &session_source); let mut client_session = client.new_session(); let mut prompt = Prompt::default(); @@ -261,14 +274,13 @@ async fn responses_stream_includes_subagent_header_on_other() { let mut stream = client_session .stream( - &window_id, &prompt, &model_info, &session_telemetry, effort, summary.unwrap_or(model_info.default_reasoning_summary), /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -339,7 +351,6 @@ async fn responses_respects_model_info_overrides_from_config() { SessionSource::SubAgent(SubAgentSource::Other("override-check".to_string())); let model_info = codex_core::test_support::construct_model_info_offline(model.as_str(), &config); - let window_id = format!("{thread_id}:0"); let session_telemetry = SessionTelemetry::new( thread_id, model.as_str(), @@ -355,18 +366,16 @@ async fn responses_respects_model_info_overrides_from_config() { let client = ModelClient::new( /*auth_manager*/ None, - thread_id.into(), thread_id, - /*installation_id*/ TEST_INSTALLATION_ID.to_string(), provider.clone(), - session_source, - /*parent_thread_id*/ None, + session_source.clone(), config.model_verbosity, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, /*beta_features_header*/ None, /*attestation_provider*/ None, ); + let responses_metadata = test_turn_responses_metadata(&client, thread_id, &session_source); let mut client_session = client.new_session(); let mut prompt = Prompt::default(); @@ -381,14 +390,13 @@ async fn responses_respects_model_info_overrides_from_config() { let mut stream = client_session .stream( - &window_id, &prompt, &model_info, &session_telemetry, effort, summary.unwrap_or(model_info.default_reasoning_summary), /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index b7cd636d5..cf8d2ec17 100644 --- a/codex-rs/core/tests/suite/client.rs +++ b/codex-rs/core/tests/suite/client.rs @@ -49,6 +49,7 @@ use codex_protocol::protocol::SessionMetaLine; use codex_protocol::protocol::SessionSource; use codex_protocol::user_input::UserInput; use core_test_support::PathBufExt; +use core_test_support::TestCodexResponsesRequestKind; use core_test_support::apps_test_server::AppsTestServer; use core_test_support::load_default_config_for_test; use core_test_support::responses::ResponsesRequest; @@ -63,6 +64,7 @@ use core_test_support::responses::mount_sse_once_match; use core_test_support::responses::mount_sse_sequence; use core_test_support::responses::sse; use core_test_support::responses::sse_failed; +use core_test_support::responses_metadata as test_responses_metadata; use core_test_support::skip_if_no_network; use core_test_support::test_codex::TestCodex; use core_test_support::test_codex::local_selections; @@ -90,6 +92,24 @@ use wiremock::matchers::query_param; const INSTALLATION_ID_FILENAME: &str = "installation_id"; const TEST_WINDOW_ID: &str = "test-thread:0"; +const TEST_INSTALLATION_ID: &str = "11111111-1111-4111-8111-111111111111"; + +fn test_turn_responses_metadata( + _client: &ModelClient, + thread_id: ThreadId, +) -> codex_core::CodexResponsesMetadata { + let thread_id = thread_id.to_string(); + test_responses_metadata( + TEST_INSTALLATION_ID, + &thread_id, + &thread_id, + /*turn_id*/ None, + TEST_WINDOW_ID.to_string(), + &SessionSource::Exec, + /*parent_thread_id*/ None, + TestCodexResponsesRequestKind::Turn, + ) +} #[expect(clippy::unwrap_used)] fn assert_message_role(request_body: &serde_json::Value, role: &str) { @@ -113,6 +133,41 @@ fn message_input_text_contains(request: &ResponsesRequest, role: &str, needle: & .any(|text| text.contains(needle)) } +fn assert_codex_client_metadata( + request_body: &serde_json::Value, + installation_id: &str, + session_id: &str, + thread_id: &str, +) { + let client_metadata = &request_body["client_metadata"]; + assert_eq!( + client_metadata["x-codex-installation-id"].as_str(), + Some(installation_id) + ); + assert_eq!(client_metadata["session_id"].as_str(), Some(session_id)); + assert_eq!(client_metadata["thread_id"].as_str(), Some(thread_id)); + let Some(turn_metadata_str) = client_metadata["x-codex-turn-metadata"].as_str() else { + panic!("missing x-codex-turn-metadata client metadata"); + }; + let Ok(turn_metadata) = serde_json::from_str::(turn_metadata_str) else { + panic!("invalid x-codex-turn-metadata json"); + }; + assert_eq!( + turn_metadata["installation_id"].as_str(), + Some(installation_id) + ); + assert_eq!(turn_metadata["session_id"].as_str(), Some(session_id)); + assert_eq!(turn_metadata["thread_id"].as_str(), Some(thread_id)); + assert_eq!( + client_metadata["turn_id"].as_str(), + turn_metadata["turn_id"].as_str() + ); + assert_eq!( + client_metadata["x-codex-window-id"].as_str(), + turn_metadata["window_id"].as_str() + ); +} + /// Writes an `auth.json` into the provided `codex_home` with the specified parameters. /// Returns the fake JWT string written to `tokens.id_token`. #[expect(clippy::unwrap_used)] @@ -780,9 +835,10 @@ async fn includes_session_id_thread_id_and_model_headers_in_request() { let installation_id = std::fs::read_to_string(test.codex_home_path().join(INSTALLATION_ID_FILENAME)) .expect("read installation id"); + let session_id_string = expected_session_id.to_string(); let thread_id_string = expected_thread_id.to_string(); - assert_eq!(request_session_id, expected_session_id.to_string()); + assert_eq!(request_session_id, session_id_string.as_str()); assert_eq!(request_thread_id, thread_id_string.as_str()); assert_eq!(request_originator, originator().value); assert_eq!(request_authorization, "Bearer Test API Key"); @@ -790,9 +846,11 @@ async fn includes_session_id_thread_id_and_model_headers_in_request() { request_body["prompt_cache_key"].as_str(), Some(thread_id_string.as_str()) ); - assert_eq!( - request_body["client_metadata"]["x-codex-installation-id"].as_str(), - Some(installation_id.as_str()) + assert_codex_client_metadata( + &request_body, + installation_id.as_str(), + session_id_string.as_str(), + thread_id_string.as_str(), ); } @@ -898,18 +956,16 @@ async fn send_provider_auth_request(server: &MockServer, auth: ModelProviderAuth Some(AuthManager::from_auth_for_testing(CodexAuth::from_api_key( "unused-api-key", ))), - thread_id.into(), thread_id, - /*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(), provider, SessionSource::Exec, - /*parent_thread_id*/ None, config.model_verbosity, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, /*beta_features_header*/ None, /*attestation_provider*/ None, ); + let responses_metadata = test_turn_responses_metadata(&client, thread_id); let mut client_session = client.new_session(); let mut prompt = Prompt::default(); prompt.input.push(ResponseItem::Message { @@ -923,14 +979,13 @@ async fn send_provider_auth_request(server: &MockServer, auth: ModelProviderAuth let mut stream = client_session .stream( - TEST_WINDOW_ID, &prompt, &model_info, &session_telemetry, effort, summary.unwrap_or(ReasoningSummary::Auto), /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -1054,15 +1109,19 @@ async fn chatgpt_auth_sends_correct_request() { let installation_id = std::fs::read_to_string(test.codex_home_path().join(INSTALLATION_ID_FILENAME)) .expect("read installation id"); - assert_eq!(request_session_id, expected_session_id.to_string()); - assert_eq!(request_thread_id, expected_thread_id.to_string()); + let session_id_string = expected_session_id.to_string(); + let thread_id_string = expected_thread_id.to_string(); + assert_eq!(request_session_id, session_id_string.as_str()); + assert_eq!(request_thread_id, thread_id_string.as_str()); assert_eq!(request_originator, originator().value); assert_eq!(request_authorization, "Bearer Access Token"); assert_eq!(request_chatgpt_account_id, "account_id"); - assert_eq!( - request_body["client_metadata"]["x-codex-installation-id"].as_str(), - Some(installation_id.as_str()) + assert_codex_client_metadata( + &request_body, + installation_id.as_str(), + session_id_string.as_str(), + thread_id_string.as_str(), ); assert!(request_body["stream"].as_bool().unwrap()); assert_eq!( @@ -2391,18 +2450,16 @@ async fn azure_responses_request_includes_store_and_reasoning_ids() { let client = ModelClient::new( /*auth_manager*/ None, - thread_id.into(), thread_id, - /*installation_id*/ "11111111-1111-4111-8111-111111111111".to_string(), provider.clone(), SessionSource::Exec, - /*parent_thread_id*/ None, config.model_verbosity, /*enable_request_compression*/ false, /*include_timing_metrics*/ false, /*beta_features_header*/ None, /*attestation_provider*/ None, ); + let responses_metadata = test_turn_responses_metadata(&client, thread_id); let mut client_session = client.new_session(); let mut prompt = Prompt::default(); @@ -2470,14 +2527,13 @@ async fn azure_responses_request_includes_store_and_reasoning_ids() { let mut stream = client_session .stream( - TEST_WINDOW_ID, &prompt, &model_info, &session_telemetry, effort, summary.unwrap_or(ReasoningSummary::Auto), /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await diff --git a/codex-rs/core/tests/suite/client_websockets.rs b/codex-rs/core/tests/suite/client_websockets.rs index 8b135b6b7..4aac3dcbd 100755 --- a/codex-rs/core/tests/suite/client_websockets.rs +++ b/codex-rs/core/tests/suite/client_websockets.rs @@ -1,6 +1,7 @@ #![allow(clippy::expect_used, clippy::unwrap_used)] use codex_api::WS_REQUEST_HEADER_TRACEPARENT_CLIENT_METADATA_KEY; use codex_api::WS_REQUEST_HEADER_TRACESTATE_CLIENT_METADATA_KEY; +use codex_core::CodexResponsesMetadata; use codex_core::ModelClient; use codex_core::ModelClientSession; use codex_core::Prompt; @@ -35,6 +36,7 @@ use codex_rollout_trace::InferenceTraceContext; use codex_rollout_trace::RawTraceEventPayload; use codex_rollout_trace::TraceWriter; use codex_rollout_trace::replay_bundle; +use core_test_support::TestCodexResponsesRequestKind; use core_test_support::load_default_config_for_test; use core_test_support::responses::WebSocketConnectionConfig; use core_test_support::responses::WebSocketTestServer; @@ -43,6 +45,7 @@ use core_test_support::responses::ev_completed; use core_test_support::responses::ev_response_created; use core_test_support::responses::start_websocket_server; use core_test_support::responses::start_websocket_server_with_headers; +use core_test_support::responses_metadata as test_responses_metadata; use core_test_support::skip_if_no_network; use core_test_support::test_codex::test_codex; use core_test_support::tracing::install_test_tracing; @@ -106,6 +109,42 @@ struct WebsocketTestHarness { session_telemetry: SessionTelemetry, } +fn responses_metadata( + harness: &WebsocketTestHarness, + turn_id: Option<&str>, + request_kind: TestCodexResponsesRequestKind, +) -> CodexResponsesMetadata { + test_responses_metadata( + TEST_INSTALLATION_ID, + &harness.session_id.to_string(), + &harness.thread_id.to_string(), + turn_id, + TEST_WINDOW_ID.to_string(), + &SessionSource::Exec, + /*parent_thread_id*/ None, + request_kind, + ) +} + +fn turn_metadata(harness: &WebsocketTestHarness, turn_id: Option<&str>) -> CodexResponsesMetadata { + responses_metadata(harness, turn_id, TestCodexResponsesRequestKind::Turn) +} + +fn prewarm_metadata( + harness: &WebsocketTestHarness, + turn_id: Option<&str>, +) -> CodexResponsesMetadata { + responses_metadata(harness, turn_id, TestCodexResponsesRequestKind::Prewarm) +} + +fn websocket_connection_metadata(harness: &WebsocketTestHarness) -> CodexResponsesMetadata { + responses_metadata( + harness, + /*turn_id*/ None, + TestCodexResponsesRequestKind::WebsocketConnection, + ) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn responses_websocket_streams_request() { skip_if_no_network!(); @@ -274,8 +313,9 @@ async fn responses_websocket_preconnect_does_not_replace_turn_trace_payload() { let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); + let responses_metadata = websocket_connection_metadata(&harness); client_session - .preconnect_websocket(&harness.session_telemetry) + .preconnect_websocket(&harness.session_telemetry, &responses_metadata) .await .expect("websocket preconnect failed"); let prompt = prompt_with_input(vec![message_item("hello")]); @@ -310,8 +350,9 @@ async fn responses_websocket_preconnect_reuses_connection() { let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); + let responses_metadata = websocket_connection_metadata(&harness); client_session - .preconnect_websocket(&harness.session_telemetry) + .preconnect_websocket(&harness.session_telemetry, &responses_metadata) .await .expect("websocket preconnect failed"); let prompt = prompt_with_input(vec![message_item("hello")]); @@ -322,7 +363,10 @@ async fn responses_websocket_preconnect_reuses_connection() { server.single_handshake().header(USER_AGENT_HEADER), Some(codex_login::default_client::get_codex_user_agent()) ); - assert_eq!(server.single_handshake().header("x-codex-window-id"), None); + assert_eq!( + server.single_handshake().header("x-codex-window-id"), + Some(TEST_WINDOW_ID.to_string()) + ); let connection = server.single_connection(); assert_eq!(connection.len(), 1); @@ -342,16 +386,16 @@ async fn responses_websocket_request_prewarm_reuses_connection() { let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ true).await; let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); + let responses_metadata = prewarm_metadata(&harness, /*turn_id*/ None); client_session .prewarm_websocket( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, ) .await .expect("websocket prewarm failed"); @@ -376,6 +420,16 @@ async fn responses_websocket_request_prewarm_reuses_connection() { assert_eq!(warmup["type"].as_str(), Some("response.create")); assert_eq!(warmup["generate"].as_bool(), Some(false)); assert_eq!(warmup["tools"], serde_json::json!([])); + let warmup_turn_metadata: serde_json::Value = serde_json::from_str( + warmup["client_metadata"]["x-codex-turn-metadata"] + .as_str() + .expect("warmup turn metadata"), + ) + .expect("valid warmup turn metadata"); + assert_eq!( + warmup_turn_metadata["request_kind"].as_str(), + Some("prewarm") + ); assert_eq!(follow_up["type"].as_str(), Some("response.create")); assert_eq!(follow_up["previous_response_id"].as_str(), Some("warm-1")); assert_eq!(follow_up["input"], serde_json::json!([])); @@ -383,6 +437,49 @@ async fn responses_websocket_request_prewarm_reuses_connection() { server.shutdown().await; } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn responses_websocket_request_prewarm_uses_caller_supplied_metadata() { + skip_if_no_network!(); + + let server = start_websocket_server(vec![vec![vec![ + ev_response_created("warm-1"), + ev_completed("warm-1"), + ]]]) + .await; + + let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ true).await; + let mut client_session = harness.client.new_session(); + let prompt = prompt_with_input(vec![message_item("hello")]); + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); + client_session + .prewarm_websocket( + &prompt, + &harness.model_info, + &harness.session_telemetry, + harness.effort.clone(), + harness.summary, + /*service_tier*/ None, + &responses_metadata, + ) + .await + .expect("websocket prewarm failed"); + + let warmup = server + .single_connection() + .first() + .expect("missing warmup request") + .body_json(); + let warmup_turn_metadata: serde_json::Value = serde_json::from_str( + warmup["client_metadata"]["x-codex-turn-metadata"] + .as_str() + .expect("warmup turn metadata"), + ) + .expect("valid warmup turn metadata"); + assert_eq!(warmup_turn_metadata["request_kind"].as_str(), Some("turn")); + + server.shutdown().await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn responses_websocket_request_prewarm_traces_logical_request() { skip_if_no_network!(); @@ -396,17 +493,17 @@ async fn responses_websocket_request_prewarm_traces_logical_request() { let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ true).await; let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); + let prewarm_responses_metadata = prewarm_metadata(&harness, /*turn_id*/ None); client_session .prewarm_websocket( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &prewarm_responses_metadata, ) .await .expect("websocket prewarm failed"); @@ -443,16 +540,16 @@ async fn responses_websocket_request_prewarm_traces_logical_request() { "test-provider".to_string(), ); + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); let mut stream = client_session .stream( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &inference_trace, ) .await @@ -610,21 +707,22 @@ async fn responses_websocket_preconnect_is_reused_even_with_header_changes() { let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); + let preconnect_metadata = websocket_connection_metadata(&harness); client_session - .preconnect_websocket(&harness.session_telemetry) + .preconnect_websocket(&harness.session_telemetry, &preconnect_metadata) .await .expect("websocket preconnect failed"); let prompt = prompt_with_input(vec![message_item("hello")]); + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); let mut stream = client_session .stream( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -655,29 +753,29 @@ async fn responses_websocket_request_prewarm_is_reused_even_with_header_changes( let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ true).await; let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); + let prewarm_responses_metadata = prewarm_metadata(&harness, /*turn_id*/ None); client_session .prewarm_websocket( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &prewarm_responses_metadata, ) .await .expect("websocket prewarm failed"); + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); let mut stream = client_session .stream( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -723,16 +821,16 @@ async fn responses_websocket_prewarm_uses_v2_when_provider_supports_websockets() let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ false).await; let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); + let responses_metadata = prewarm_metadata(&harness, /*turn_id*/ None); client_session .prewarm_websocket( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, ) .await .expect("websocket prewarm failed"); @@ -780,13 +878,18 @@ async fn responses_websocket_preconnect_runs_when_only_v2_feature_enabled() { let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ true).await; let mut client_session = harness.client.new_session(); + let responses_metadata = websocket_connection_metadata(&harness); client_session - .preconnect_websocket(&harness.session_telemetry) + .preconnect_websocket(&harness.session_telemetry, &responses_metadata) .await .expect("websocket preconnect failed"); assert_eq!(server.handshakes().len(), 1); assert_eq!(server.single_connection().len(), 0); + assert_eq!( + server.single_handshake().header("x-codex-turn-metadata"), + None + ); let prompt = prompt_with_input(vec![message_item("hello")]); stream_until_complete(&mut client_session, &harness, &prompt).await; @@ -1072,17 +1175,17 @@ async fn responses_websocket_emits_reasoning_included_event() { let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); let mut stream = client_session .stream( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -1147,17 +1250,17 @@ async fn responses_websocket_emits_rate_limit_events() { let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); let mut stream = client_session .stream( - TEST_WINDOW_ID, &prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -1471,30 +1574,29 @@ async fn responses_websocket_forwards_turn_metadata_on_initial_and_incremental_c let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); - let first_turn_metadata = - r#"{"turn_id":"turn-123","thread_source":"user","sandbox":"workspace-write"}"#; - let enriched_turn_metadata = r#"{"turn_id":"turn-123","thread_source":"user","sandbox":"workspace-write","workspaces":[{"root_path":"/tmp/repo","latest_git_commit_hash":"abc123","associated_remote_urls":["git@github.com:openai/codex.git"],"has_changes":true}]}"#; let prompt_one = prompt_with_input(vec![message_item("hello")]); let prompt_two = prompt_with_input(vec![ message_item("hello"), assistant_message_item("msg-1", "assistant output"), message_item("second"), ]); + let first_responses_metadata = turn_metadata(&harness, Some("turn-123")); + let second_responses_metadata = turn_metadata(&harness, Some("turn-456")); - stream_until_complete_with_turn_metadata( + stream_until_complete_with_metadata( &mut client_session, &harness, &prompt_one, /*service_tier*/ None, - Some(first_turn_metadata), + &first_responses_metadata, ) .await; - stream_until_complete_with_turn_metadata( + stream_until_complete_with_metadata( &mut client_session, &harness, &prompt_two, /*service_tier*/ None, - Some(enriched_turn_metadata), + &second_responses_metadata, ) .await; @@ -1504,36 +1606,37 @@ async fn responses_websocket_forwards_turn_metadata_on_initial_and_incremental_c let second = connection.get(1).expect("missing request").body_json(); assert_eq!(first["type"].as_str(), Some("response.create")); - assert_eq!( - first["client_metadata"]["x-codex-turn-metadata"].as_str(), - Some(first_turn_metadata) - ); assert_eq!(second["type"].as_str(), Some("response.create")); assert_eq!(second["previous_response_id"].as_str(), Some("resp-1")); - assert_eq!( - second["client_metadata"]["x-codex-turn-metadata"].as_str(), - Some(enriched_turn_metadata) - ); - - let first_metadata: serde_json::Value = - serde_json::from_str(first_turn_metadata).expect("first metadata should be valid json"); - let second_metadata: serde_json::Value = serde_json::from_str(enriched_turn_metadata) - .expect("enriched metadata should be valid json"); + let first_metadata: serde_json::Value = serde_json::from_str( + first["client_metadata"]["x-codex-turn-metadata"] + .as_str() + .expect("first turn metadata"), + ) + .expect("first metadata should be valid json"); + let second_metadata: serde_json::Value = serde_json::from_str( + second["client_metadata"]["x-codex-turn-metadata"] + .as_str() + .expect("second turn metadata"), + ) + .expect("second metadata should be valid json"); assert_eq!(first_metadata["turn_id"].as_str(), Some("turn-123")); - assert_eq!(second_metadata["turn_id"].as_str(), Some("turn-123")); - assert_eq!(first_metadata["thread_source"].as_str(), Some("user")); - assert_eq!(second_metadata["thread_source"].as_str(), Some("user")); + assert_eq!(second_metadata["turn_id"].as_str(), Some("turn-456")); assert_eq!( - second_metadata["workspaces"][0]["has_changes"].as_bool(), - Some(true) + first["client_metadata"]["turn_id"].as_str(), + first_metadata["turn_id"].as_str() + ); + assert_eq!( + second["client_metadata"]["turn_id"].as_str(), + second_metadata["turn_id"].as_str() ); server.shutdown().await; } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn responses_websocket_preserves_custom_turn_metadata_fields() { +async fn responses_websocket_sends_canonical_turn_metadata() { skip_if_no_network!(); let server = start_websocket_server(vec![vec![vec![ @@ -1545,19 +1648,14 @@ async fn responses_websocket_preserves_custom_turn_metadata_fields() { 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(); + let responses_metadata = turn_metadata(&harness, Some("turn-123")); - stream_until_complete_with_turn_metadata( + stream_until_complete_with_metadata( &mut client_session, &harness, &prompt, /*service_tier*/ None, - Some(&turn_metadata), + &responses_metadata, ) .await; @@ -1568,15 +1666,16 @@ async fn responses_websocket_preserves_custom_turn_metadata_fields() { .body_json(); assert_eq!(body["type"].as_str(), Some("response.create")); - assert_eq!( + let turn_metadata: serde_json::Value = serde_json::from_str( 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", - })) + .expect("turn metadata"), + ) + .expect("valid turn metadata"); + assert_eq!(turn_metadata["turn_id"].as_str(), Some("turn-123")); + assert_eq!( + body["client_metadata"]["turn_id"].as_str(), + turn_metadata["turn_id"].as_str() ); server.shutdown().await; @@ -1803,16 +1902,16 @@ async fn responses_websocket_v2_after_error_uses_full_create_without_previous_re stream_until_complete(&mut session, &harness, &prompt_one).await; + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); let mut second_stream = session .stream( - TEST_WINDOW_ID, &prompt_two, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -1892,16 +1991,16 @@ async fn responses_websocket_v2_surfaces_terminal_error_without_close_handshake( stream_until_complete(&mut session, &harness, &prompt_one).await; + let responses_metadata = turn_metadata(&harness, /*turn_id*/ None); let mut second_stream = session .stream( - TEST_WINDOW_ID, &prompt_two, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -2081,12 +2180,9 @@ async fn websocket_harness_with_provider_options( let summary = ReasoningSummary::Auto; let client = ModelClient::new( /*auth_manager*/ None, - session_id, thread_id, - /*installation_id*/ TEST_INSTALLATION_ID.to_string(), provider.clone(), SessionSource::Exec, - /*parent_thread_id*/ None, config.model_verbosity, /*enable_request_compression*/ false, runtime_metrics_enabled, @@ -2127,16 +2223,16 @@ async fn stream_until_complete_with_model_info( model_info: &ModelInfo, expected_response_id: &str, ) { + let responses_metadata = turn_metadata(harness, /*turn_id*/ None); let mut stream = client_session .stream( - TEST_WINDOW_ID, prompt, model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, /*service_tier*/ None, - /*turn_metadata_header*/ None, + &responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await @@ -2161,50 +2257,33 @@ async fn stream_until_complete_with_service_tier( prompt: &Prompt, service_tier: Option, ) { - stream_until_complete_with_turn_metadata( + let responses_metadata = turn_metadata(harness, /*turn_id*/ None); + stream_until_complete_with_metadata( client_session, harness, prompt, service_tier, - /*turn_metadata_header*/ None, + &responses_metadata, ) .await; } -async fn stream_until_complete_with_turn_metadata( +async fn stream_until_complete_with_metadata( client_session: &mut ModelClientSession, harness: &WebsocketTestHarness, 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>, + responses_metadata: &CodexResponsesMetadata, ) { let mut stream = client_session .stream( - TEST_WINDOW_ID, prompt, &harness.model_info, &harness.session_telemetry, harness.effort.clone(), harness.summary, service_tier.map(|service_tier| service_tier.request_value().to_string()), - turn_metadata_header, + responses_metadata, &codex_rollout_trace::InferenceTraceContext::disabled(), ) .await diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index 2a1d459f9..a63397c55 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -386,6 +386,10 @@ async fn remote_compact_replaces_history_for_followups() -> Result<()> { .expect("remote compact request should include turn metadata"), ) .expect("remote compact turn metadata should be valid json"); + assert_eq!( + compact_request.header("x-codex-installation-id").as_deref(), + compact_metadata["installation_id"].as_str() + ); assert!( compact_metadata["turn_id"] .as_str() diff --git a/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_api_auth_prompt_cache_key_request_diff.snap b/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_api_auth_prompt_cache_key_request_diff.snap index 89796cb91..0df96a7c0 100644 --- a/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_api_auth_prompt_cache_key_request_diff.snap +++ b/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_api_auth_prompt_cache_key_request_diff.snap @@ -7,7 +7,11 @@ Scenario: After five varied API-key-auth turns, remote manual compaction omits s --- Last Normal /responses Request +++ Remote /responses/compact Request - "client_metadata": { +- "session_id": "", +- "thread_id": "", +- "turn_id": "", - "x-codex-installation-id": "", +- "x-codex-turn-metadata": "{\"installation_id\":\"\",\"session_id\":\"\",\"thread_id\":\"\",\"turn_id\":\"\",\"window_id\":\":0\",\"request_kind\":\"turn\",\"sandbox\":\"\",\"turn_started_at_unix_ms\":}", - "x-codex-window-id": ":0" - }, - "include": [ diff --git a/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_chatgpt_auth_service_tier_prompt_cache_key_request_diff.snap b/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_chatgpt_auth_service_tier_prompt_cache_key_request_diff.snap index ab6c9d1b3..d959d526c 100644 --- a/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_chatgpt_auth_service_tier_prompt_cache_key_request_diff.snap +++ b/codex-rs/core/tests/suite/snapshots/all__suite__compact_remote__remote_manual_compact_chatgpt_auth_service_tier_prompt_cache_key_request_diff.snap @@ -7,7 +7,11 @@ Scenario: After five varied ChatGPT-auth turns, remote manual compaction reuses --- Last Normal /responses Request +++ Remote /responses/compact Request - "client_metadata": { +- "session_id": "", +- "thread_id": "", +- "turn_id": "", - "x-codex-installation-id": "", +- "x-codex-turn-metadata": "{\"installation_id\":\"\",\"session_id\":\"\",\"thread_id\":\"\",\"turn_id\":\"\",\"window_id\":\":0\",\"request_kind\":\"turn\",\"sandbox\":\"\",\"turn_started_at_unix_ms\":}", - "x-codex-window-id": ":0" - }, - "include": [ diff --git a/codex-rs/memories/write/src/runtime.rs b/codex-rs/memories/write/src/runtime.rs index 0b6718f01..8dae6ca61 100644 --- a/codex-rs/memories/write/src/runtime.rs +++ b/codex-rs/memories/write/src/runtime.rs @@ -7,6 +7,7 @@ use codex_core::StartThreadOptions; use codex_core::ThreadManager; use codex_core::config::Config; use codex_core::content_items_to_text; +use codex_core::detached_memory_responses_metadata; use codex_core::resolve_installation_id; use codex_features::Feature; use codex_login::AuthManager; @@ -49,7 +50,6 @@ pub(crate) struct StageOneRequestContext { pub(crate) reasoning_effort: Option, pub(crate) reasoning_summary: ReasoningSummary, pub(crate) service_tier: Option, - pub(crate) turn_metadata_header: Option, } impl StageOneRequestContext { @@ -199,8 +199,6 @@ impl MemoryStartupContext { .get_models_manager() .get_model_info(model_name, &config.to_models_manager_config()) .await; - let turn_metadata_header = - codex_core::build_turn_metadata_header(&config.cwd, /*sandbox*/ None).await; let reasoning_summary = config .model_reasoning_summary .unwrap_or(model_info.default_reasoning_summary); @@ -214,7 +212,6 @@ impl MemoryStartupContext { reasoning_effort: Some(reasoning_effort), reasoning_summary, service_tier: config_snapshot.service_tier, - turn_metadata_header, } } @@ -227,14 +224,13 @@ impl MemoryStartupContext { let installation_id = resolve_installation_id(&config.codex_home).await?; let config_snapshot = self.thread.config_snapshot().await; let session_source = config_snapshot.session_source; + let session_id = SessionId::from(self.thread_id); + let session_id_string = session_id.to_string(); let model_client = ModelClient::new( Some(Arc::clone(&self.auth_manager)), - SessionId::from(self.thread_id), // We use thread_id to detach this query from the foreground user session. self.thread_id, - installation_id, config.model_provider.clone(), - session_source, - config_snapshot.parent_thread_id, + session_source.clone(), config.model_verbosity, config.features.enabled(Feature::EnableRequestCompression), config.features.enabled(Feature::RuntimeMetrics), @@ -244,16 +240,25 @@ impl MemoryStartupContext { let mut client_session = model_client.new_session(); let window_id = format!("{}:0", self.thread_id); + let responses_metadata = detached_memory_responses_metadata( + installation_id, + session_id_string, + self.thread_id.to_string(), + window_id, + &session_source, + &config.cwd, + /*sandbox*/ None, + ) + .await; let mut stream = client_session .stream( - &window_id, prompt, &context.model_info, &context.session_telemetry, context.reasoning_effort.clone(), context.reasoning_summary, context.service_tier.clone(), - context.turn_metadata_header.as_deref(), + &responses_metadata, &InferenceTraceContext::disabled(), ) .await?;