From 61fe23159ed275e444e9b7614bb59a3ac4639f2a Mon Sep 17 00:00:00 2001 From: marksteinbrick-oai Date: Tue, 14 Apr 2026 08:55:12 -0700 Subject: [PATCH] [codex-analytics] add session source to client metadata (#17374) ## Summary Adds `thread_source` field to the existing Codex turn metadata sent to Responses API - Sends `thread_source: "user"` for user-initiated sessions: CLI, VS Code, and Exec - Sends `thread_source: "subagent"` for subagent sessions - Omits `thread_source` for MCP, custom, and unknown session sources - Uses the existing turn metadata transport: - HTTP requests send through the `x-codex-turn-metadata` header - WebSocket `response.create` requests send through `client_metadata["x-codex-turn-metadata"]` ## Testing - `cargo test -p codex-protocol session_source_thread_source_name_classifies_user_and_subagent_sources` - `cargo test -p codex-core turn_metadata_state` - `cargo test -p codex-core --test responses_headers responses_stream_includes_turn_metadata_header_for_git_workspace_e2e -- --nocapture` --- codex-rs/analytics/src/events.rs | 9 ------ codex-rs/analytics/src/reducer.rs | 3 +- codex-rs/core/src/codex.rs | 2 ++ codex-rs/core/src/turn_metadata.rs | 9 ++++++ codex-rs/core/src/turn_metadata_tests.rs | 32 +++++++++++++++++++ codex-rs/core/tests/responses_headers.rs | 18 +++++++++++ .../core/tests/suite/client_websockets.rs | 7 ++-- codex-rs/protocol/src/protocol.rs | 27 ++++++++++++++++ 8 files changed, 94 insertions(+), 13 deletions(-) diff --git a/codex-rs/analytics/src/events.rs b/codex-rs/analytics/src/events.rs index 87f5a4fbe..eb312da14 100644 --- a/codex-rs/analytics/src/events.rs +++ b/codex-rs/analytics/src/events.rs @@ -15,7 +15,6 @@ use codex_plugin::PluginTelemetryMetadata; use codex_protocol::approvals::NetworkApprovalProtocol; use codex_protocol::models::PermissionProfile; use codex_protocol::models::SandboxPermissions; -use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; use serde::Serialize; @@ -530,14 +529,6 @@ pub(crate) fn codex_plugin_used_metadata( } } -pub(crate) fn thread_source_name(thread_source: &SessionSource) -> Option<&'static str> { - match thread_source { - SessionSource::Cli | SessionSource::VSCode | SessionSource::Exec => Some("user"), - SessionSource::SubAgent(_) => Some("subagent"), - SessionSource::Mcp | SessionSource::Custom(_) | SessionSource::Unknown => None, - } -} - pub(crate) fn current_runtime_metadata() -> CodexRuntimeMetadata { let os_info = os_info::get(); CodexRuntimeMetadata { diff --git a/codex-rs/analytics/src/reducer.rs b/codex-rs/analytics/src/reducer.rs index 772a091b1..a5e5165bd 100644 --- a/codex-rs/analytics/src/reducer.rs +++ b/codex-rs/analytics/src/reducer.rs @@ -26,7 +26,6 @@ use crate::events::plugin_state_event_type; use crate::events::subagent_parent_thread_id; use crate::events::subagent_source_name; use crate::events::subagent_thread_started_event_request; -use crate::events::thread_source_name; use crate::facts::AnalyticsFact; use crate::facts::AnalyticsJsonRpcError; use crate::facts::AppMentionedInput; @@ -107,7 +106,7 @@ impl ThreadMetadataState { | SessionSource::Unknown => (None, None), }; Self { - thread_source: thread_source_name(session_source), + thread_source: session_source.thread_source_name(), initialization_mode, subagent_source, parent_thread_id, diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 845c5c085..21725cd33 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -1570,6 +1570,7 @@ impl Session { let per_turn_config = Arc::new(per_turn_config); let turn_metadata_state = Arc::new(TurnMetadataState::new( conversation_id.to_string(), + &session_source, sub_id.clone(), cwd.to_path_buf(), session_configuration.sandbox_policy.get(), @@ -5983,6 +5984,7 @@ async fn spawn_review_thread( let review_turn_id = sub_id.to_string(); let turn_metadata_state = Arc::new(TurnMetadataState::new( sess.conversation_id.to_string(), + &session_source, review_turn_id.clone(), parent_turn_context.cwd.to_path_buf(), parent_turn_context.sandbox_policy.get(), diff --git a/codex-rs/core/src/turn_metadata.rs b/codex-rs/core/src/turn_metadata.rs index e4385096a..9bdffa45d 100644 --- a/codex-rs/core/src/turn_metadata.rs +++ b/codex-rs/core/src/turn_metadata.rs @@ -17,6 +17,7 @@ use codex_git_utils::get_has_changes; use codex_git_utils::get_head_commit_hash; use codex_protocol::config_types::WindowsSandboxLevel; use codex_protocol::protocol::SandboxPolicy; +use codex_protocol::protocol::SessionSource; #[derive(Clone, Debug, Default)] struct WorkspaceGitMetadata { @@ -58,6 +59,8 @@ pub(crate) struct TurnMetadataBag { #[serde(default, skip_serializing_if = "Option::is_none")] session_id: Option, #[serde(default, skip_serializing_if = "Option::is_none")] + thread_source: Option<&'static str>, + #[serde(default, skip_serializing_if = "Option::is_none")] turn_id: Option, #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] workspaces: BTreeMap, @@ -87,6 +90,7 @@ fn merge_responsesapi_client_metadata( fn build_turn_metadata_bag( session_id: Option, + thread_source: Option<&'static str>, turn_id: Option, sandbox: Option, repo_root: Option, @@ -101,6 +105,7 @@ fn build_turn_metadata_bag( TurnMetadataBag { session_id, + thread_source, turn_id, workspaces, sandbox, @@ -126,6 +131,7 @@ pub async fn build_turn_metadata_header(cwd: &Path, sandbox: Option<&str>) -> Op build_turn_metadata_bag( /*session_id*/ None, + /*thread_source*/ None, /*turn_id*/ None, sandbox.map(ToString::to_string), repo_root, @@ -152,6 +158,7 @@ pub(crate) struct TurnMetadataState { impl TurnMetadataState { pub(crate) fn new( session_id: String, + session_source: &SessionSource, turn_id: String, cwd: PathBuf, sandbox_policy: &SandboxPolicy, @@ -161,6 +168,7 @@ impl TurnMetadataState { let sandbox = Some(sandbox_tag(sandbox_policy, windows_sandbox_level).to_string()); let base_metadata = build_turn_metadata_bag( Some(session_id), + session_source.thread_source_name(), Some(turn_id), sandbox, /*repo_root*/ None, @@ -240,6 +248,7 @@ impl TurnMetadataState { let enriched_metadata = build_turn_metadata_bag( state.base_metadata.session_id.clone(), + state.base_metadata.thread_source, state.base_metadata.turn_id.clone(), state.base_metadata.sandbox.clone(), Some(repo_root), diff --git a/codex-rs/core/src/turn_metadata_tests.rs b/codex-rs/core/src/turn_metadata_tests.rs index dfd9322a9..71fd29625 100644 --- a/codex-rs/core/src/turn_metadata_tests.rs +++ b/codex-rs/core/src/turn_metadata_tests.rs @@ -1,5 +1,7 @@ use super::*; +use codex_protocol::protocol::SessionSource; +use codex_protocol::protocol::SubAgentSource; use serde_json::Value; use std::collections::HashMap; use tempfile::TempDir; @@ -69,6 +71,7 @@ fn turn_metadata_state_uses_platform_sandbox_tag() { let state = TurnMetadataState::new( "session-a".to_string(), + &SessionSource::Exec, "turn-a".to_string(), cwd, &sandbox_policy, @@ -79,10 +82,36 @@ fn turn_metadata_state_uses_platform_sandbox_tag() { 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_source = json.get("thread_source").and_then(Value::as_str); let expected_sandbox = sandbox_tag(&sandbox_policy, WindowsSandboxLevel::Disabled); assert_eq!(sandbox_name, Some(expected_sandbox)); assert_eq!(session_id, Some("session-a")); + assert_eq!(thread_source, Some("user")); + assert!(json.get("session_source").is_none()); +} + +#[test] +fn turn_metadata_state_classifies_subagent_thread_source() { + let temp_dir = TempDir::new().expect("temp dir"); + let cwd = temp_dir.path().to_path_buf(); + let sandbox_policy = SandboxPolicy::new_read_only_policy(); + let session_source = SessionSource::SubAgent(SubAgentSource::Review); + + let state = TurnMetadataState::new( + "session-a".to_string(), + &session_source, + "turn-a".to_string(), + cwd, + &sandbox_policy, + WindowsSandboxLevel::Disabled, + ); + + 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] @@ -93,6 +122,7 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( let state = TurnMetadataState::new( "session-a".to_string(), + &SessionSource::Exec, "turn-a".to_string(), cwd, &sandbox_policy, @@ -101,6 +131,7 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( state.set_responsesapi_client_metadata(HashMap::from([ ("fiber_run_id".to_string(), "fiber-123".to_string()), ("session_id".to_string(), "client-supplied".to_string()), + ("thread_source".to_string(), "client-supplied".to_string()), ])); let header = state.current_header_value().expect("header"); @@ -108,5 +139,6 @@ fn turn_metadata_state_merges_client_metadata_without_replacing_reserved_fields( assert_eq!(json["fiber_run_id"].as_str(), Some("fiber-123")); assert_eq!(json["session_id"].as_str(), Some("session-a")); + assert_eq!(json["thread_source"].as_str(), Some("user")); assert_eq!(json["turn_id"].as_str(), Some("turn-a")); } diff --git a/codex-rs/core/tests/responses_headers.rs b/codex-rs/core/tests/responses_headers.rs index 03870beab..849895651 100644 --- a/codex-rs/core/tests/responses_headers.rs +++ b/codex-rs/core/tests/responses_headers.rs @@ -437,6 +437,12 @@ async fn responses_stream_includes_turn_metadata_header_for_git_workspace_e2e() .and_then(serde_json::Value::as_str), Some("none") ); + assert_eq!( + initial_parsed + .get("thread_source") + .and_then(serde_json::Value::as_str), + Some("user") + ); let git_config_global = cwd.join("empty-git-config"); std::fs::write(&git_config_global, "").expect("write empty git config"); @@ -528,6 +534,18 @@ async fn responses_stream_includes_turn_metadata_header_for_git_workspace_e2e() .get("turn_id") .and_then(serde_json::Value::as_str) .expect("second turn_id should be present"); + assert_eq!( + first_parsed + .get("thread_source") + .and_then(serde_json::Value::as_str), + Some("user") + ); + assert_eq!( + second_parsed + .get("thread_source") + .and_then(serde_json::Value::as_str), + Some("user") + ); assert_eq!( first_turn_id, second_turn_id, "requests should share turn_id" diff --git a/codex-rs/core/tests/suite/client_websockets.rs b/codex-rs/core/tests/suite/client_websockets.rs index 424f6764e..eaa8a9610 100755 --- a/codex-rs/core/tests/suite/client_websockets.rs +++ b/codex-rs/core/tests/suite/client_websockets.rs @@ -1209,8 +1209,9 @@ 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","sandbox":"workspace-write"}"#; - let enriched_turn_metadata = r#"{"turn_id":"turn-123","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 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"), @@ -1259,6 +1260,8 @@ async fn responses_websocket_forwards_turn_metadata_on_initial_and_incremental_c 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["workspaces"][0]["has_changes"].as_bool(), Some(true) diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index f56e54fd2..2f9d7021a 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -2633,6 +2633,15 @@ impl SessionSource { }) } + /// Low cardinality thread source label for analytics. + pub fn thread_source_name(&self) -> Option<&'static str> { + match self { + SessionSource::Cli | SessionSource::VSCode | SessionSource::Exec => Some("user"), + SessionSource::SubAgent(_) => Some("subagent"), + SessionSource::Mcp | SessionSource::Custom(_) | SessionSource::Unknown => None, + } + } + pub fn get_nickname(&self) -> Option { match self { SessionSource::SubAgent(SubAgentSource::ThreadSpawn { agent_nickname, .. }) => { @@ -3882,6 +3891,24 @@ mod tests { ); } + #[test] + fn session_source_thread_source_name_classifies_user_and_subagent_sources() { + for (source, expected) in [ + (SessionSource::Cli, Some("user")), + (SessionSource::VSCode, Some("user")), + (SessionSource::Exec, Some("user")), + ( + SessionSource::SubAgent(SubAgentSource::Review), + Some("subagent"), + ), + (SessionSource::Mcp, None), + (SessionSource::Custom("atlas".to_string()), None), + (SessionSource::Unknown, None), + ] { + assert_eq!(source.thread_source_name(), expected); + } + } + #[test] fn session_source_restriction_product_defaults_non_subagent_sources_to_codex() { assert_eq!(