From 4fd5c35c4f6f51048f47c8680ed0f6a26c608f68 Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Sat, 4 Apr 2026 11:06:43 -0700 Subject: [PATCH] [codex-analytics] subagent analytics (#15915) - creates custom event that emits subagent thread analytics from core - wires client metadata (`product_client_id, client_name, client_version`), through from app-server - creates `created_at `timestamp in core - subagent analytics are behind `FeatureFlag::GeneralAnalytics` PR stack - [[telemetry] thread events #15690](https://github.com/openai/codex/pull/15690) - --> [[telemetry] subagent events #15915](https://github.com/openai/codex/pull/15915) - [[telemetry] turn events #15591](https://github.com/openai/codex/pull/15591) - [[telemetry] steer events #15697](https://github.com/openai/codex/pull/15697) - [[telemetry] queued prompt data #15804](https://github.com/openai/codex/pull/15804) Notes: - core does not spawn a subagent thread for compact, but represented in mapping for consistency `INFO | 2026-04-01 13:08:12 | codex_backend.routers.analytics_events | analytics_events.track_analytics_events:399 | Tracked codex_thread_initialized event params={'thread_id': '019d4aa9-233b-70f2-a958-c3dbae1e30fa', 'product_surface': 'codex', 'app_server_client': {'product_client_id': 'CODEX_CLI', 'client_name': 'codex-tui', 'client_version': '0.0.0', 'rpc_transport': 'in_process', 'experimental_api_enabled': None}, 'runtime': {'codex_rs_version': '0.0.0', 'runtime_os': 'macos', 'runtime_os_version': '26.4.0', 'runtime_arch': 'aarch64'}, 'model': 'gpt-5.3-codex', 'ephemeral': False, 'initialization_mode': 'new', 'created_at': 1775074091, 'thread_source': 'subagent', 'subagent_source': 'thread_spawn', 'parent_thread_id': '019d4aa8-51ec-77e3-bafb-2c1b8e29e385'} | ` `INFO | 2026-04-01 13:08:41 | codex_backend.routers.analytics_events | analytics_events.track_analytics_events:399 | Tracked codex_thread_initialized event params={'thread_id': '019d4aa9-94e3-75f1-8864-ff8ad0e55e1e', 'product_surface': 'codex', 'app_server_client': {'product_client_id': 'CODEX_CLI', 'client_name': 'codex-tui', 'client_version': '0.0.0', 'rpc_transport': 'in_process', 'experimental_api_enabled': None}, 'runtime': {'codex_rs_version': '0.0.0', 'runtime_os': 'macos', 'runtime_os_version': '26.4.0', 'runtime_arch': 'aarch64'}, 'model': 'gpt-5.3-codex', 'ephemeral': False, 'initialization_mode': 'new', 'created_at': 1775074120, 'thread_source': 'subagent', 'subagent_source': 'review', 'parent_thread_id': None} | ` --------- Co-authored-by: jif-oai Co-authored-by: Michael Bolin --- .../analytics/src/analytics_client_tests.rs | 152 ++++++++++++++++++ codex-rs/analytics/src/client.rs | 7 + codex-rs/analytics/src/events.rs | 48 ++++++ codex-rs/analytics/src/facts.rs | 14 ++ codex-rs/analytics/src/lib.rs | 1 + codex-rs/analytics/src/reducer.rs | 15 ++ .../app-server/src/codex_message_processor.rs | 39 ++++- codex-rs/app-server/src/message_processor.rs | 1 + .../app-server/tests/suite/v2/thread_fork.rs | 4 +- .../app-server/tests/suite/v2/thread_start.rs | 4 +- codex-rs/core/src/agent/control.rs | 43 +++++ codex-rs/core/src/codex.rs | 61 ++++++- codex-rs/core/src/codex_delegate.rs | 14 +- codex-rs/core/src/codex_tests.rs | 6 + codex-rs/core/src/codex_thread.rs | 5 +- codex-rs/core/src/memories/phase2.rs | 21 +++ 16 files changed, 424 insertions(+), 11 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index d51112a8b..73ea42d76 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -13,6 +13,7 @@ use crate::events::TrackEventRequest; use crate::events::codex_app_metadata; use crate::events::codex_plugin_metadata; use crate::events::codex_plugin_used_metadata; +use crate::events::subagent_thread_started_event_request; use crate::facts::AnalyticsFact; use crate::facts::AppInvocation; use crate::facts::AppMentionedInput; @@ -24,6 +25,7 @@ use crate::facts::PluginStateChangedInput; use crate::facts::PluginUsedInput; use crate::facts::SkillInvocation; use crate::facts::SkillInvokedInput; +use crate::facts::SubAgentThreadStartedInput; use crate::facts::TrackEventsContext; use crate::reducer::AnalyticsReducer; use crate::reducer::normalize_path_for_skill_id; @@ -47,6 +49,7 @@ use codex_plugin::AppConnectorId; use codex_plugin::PluginCapabilitySummary; use codex_plugin::PluginId; use codex_plugin::PluginTelemetryMetadata; +use codex_protocol::protocol::SubAgentSource; use pretty_assertions::assert_eq; use serde_json::json; use std::collections::HashSet; @@ -446,6 +449,155 @@ async fn initialize_caches_client_and_thread_lifecycle_publishes_once_initialize assert_eq!(payload[0]["event_params"]["parent_thread_id"], json!(null)); } +#[test] +fn subagent_thread_started_review_serializes_expected_shape() { + let event = TrackEventRequest::ThreadInitialized(subagent_thread_started_event_request( + SubAgentThreadStartedInput { + thread_id: "thread-review".to_string(), + product_client_id: "codex-tui".to_string(), + client_name: "codex-tui".to_string(), + client_version: "1.0.0".to_string(), + model: "gpt-5".to_string(), + ephemeral: false, + subagent_source: SubAgentSource::Review, + created_at: 123, + }, + )); + + let payload = serde_json::to_value(&event).expect("serialize review subagent event"); + assert_eq!(payload["event_params"]["thread_source"], "subagent"); + assert_eq!( + payload["event_params"]["app_server_client"]["product_client_id"], + "codex-tui" + ); + assert_eq!( + payload["event_params"]["app_server_client"]["client_name"], + "codex-tui" + ); + assert_eq!( + payload["event_params"]["app_server_client"]["client_version"], + "1.0.0" + ); + assert_eq!( + payload["event_params"]["app_server_client"]["rpc_transport"], + "in_process" + ); + assert_eq!(payload["event_params"]["created_at"], 123); + assert_eq!(payload["event_params"]["initialization_mode"], "new"); + assert_eq!(payload["event_params"]["subagent_source"], "review"); + assert_eq!(payload["event_params"]["parent_thread_id"], json!(null)); +} + +#[test] +fn subagent_thread_started_thread_spawn_serializes_parent_thread_id() { + let parent_thread_id = + codex_protocol::ThreadId::from_string("11111111-1111-1111-1111-111111111111") + .expect("valid thread id"); + let event = TrackEventRequest::ThreadInitialized(subagent_thread_started_event_request( + SubAgentThreadStartedInput { + thread_id: "thread-spawn".to_string(), + product_client_id: "codex-tui".to_string(), + client_name: "codex-tui".to_string(), + client_version: "1.0.0".to_string(), + model: "gpt-5".to_string(), + ephemeral: true, + subagent_source: SubAgentSource::ThreadSpawn { + parent_thread_id, + depth: 1, + agent_path: None, + agent_nickname: None, + agent_role: None, + }, + created_at: 124, + }, + )); + + let payload = serde_json::to_value(&event).expect("serialize thread spawn subagent event"); + assert_eq!(payload["event_params"]["thread_source"], "subagent"); + assert_eq!(payload["event_params"]["subagent_source"], "thread_spawn"); + assert_eq!( + payload["event_params"]["parent_thread_id"], + "11111111-1111-1111-1111-111111111111" + ); +} + +#[test] +fn subagent_thread_started_memory_consolidation_serializes_expected_shape() { + let event = TrackEventRequest::ThreadInitialized(subagent_thread_started_event_request( + SubAgentThreadStartedInput { + thread_id: "thread-memory".to_string(), + product_client_id: "codex-tui".to_string(), + client_name: "codex-tui".to_string(), + client_version: "1.0.0".to_string(), + model: "gpt-5".to_string(), + ephemeral: false, + subagent_source: SubAgentSource::MemoryConsolidation, + created_at: 125, + }, + )); + + let payload = + serde_json::to_value(&event).expect("serialize memory consolidation subagent event"); + assert_eq!( + payload["event_params"]["subagent_source"], + "memory_consolidation" + ); + assert_eq!(payload["event_params"]["parent_thread_id"], json!(null)); +} + +#[test] +fn subagent_thread_started_other_serializes_expected_shape() { + let event = TrackEventRequest::ThreadInitialized(subagent_thread_started_event_request( + SubAgentThreadStartedInput { + thread_id: "thread-guardian".to_string(), + product_client_id: "codex-tui".to_string(), + client_name: "codex-tui".to_string(), + client_version: "1.0.0".to_string(), + model: "gpt-5".to_string(), + ephemeral: false, + subagent_source: SubAgentSource::Other("guardian".to_string()), + created_at: 126, + }, + )); + + let payload = serde_json::to_value(&event).expect("serialize other subagent event"); + assert_eq!(payload["event_params"]["subagent_source"], "guardian"); +} + +#[tokio::test] +async fn subagent_thread_started_publishes_without_initialize() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + reducer + .ingest( + AnalyticsFact::Custom(CustomAnalyticsFact::SubAgentThreadStarted( + SubAgentThreadStartedInput { + thread_id: "thread-review".to_string(), + product_client_id: "codex-tui".to_string(), + client_name: "codex-tui".to_string(), + client_version: "1.0.0".to_string(), + model: "gpt-5".to_string(), + ephemeral: false, + subagent_source: SubAgentSource::Review, + created_at: 127, + }, + )), + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events).expect("serialize events"); + assert_eq!(payload.as_array().expect("events array").len(), 1); + assert_eq!(payload[0]["event_type"], "codex_thread_initialized"); + assert_eq!( + payload[0]["event_params"]["app_server_client"]["product_client_id"], + "codex-tui" + ); + assert_eq!(payload[0]["event_params"]["thread_source"], "subagent"); + assert_eq!(payload[0]["event_params"]["subagent_source"], "review"); +} + #[test] fn plugin_used_event_serializes_expected_shape() { let tracking = TrackEventsContext { diff --git a/codex-rs/analytics/src/client.rs b/codex-rs/analytics/src/client.rs index b8a543a3b..dd300a7bd 100644 --- a/codex-rs/analytics/src/client.rs +++ b/codex-rs/analytics/src/client.rs @@ -11,6 +11,7 @@ use crate::facts::PluginState; use crate::facts::PluginStateChangedInput; use crate::facts::SkillInvocation; use crate::facts::SkillInvokedInput; +use crate::facts::SubAgentThreadStartedInput; use crate::facts::TrackEventsContext; use crate::reducer::AnalyticsReducer; use codex_app_server_protocol::ClientResponse; @@ -144,6 +145,12 @@ impl AnalyticsEventsClient { }); } + pub fn track_subagent_thread_started(&self, input: SubAgentThreadStartedInput) { + self.record_fact(AnalyticsFact::Custom( + CustomAnalyticsFact::SubAgentThreadStarted(input), + )); + } + pub fn track_app_mentioned(&self, tracking: TrackEventsContext, mentions: Vec) { if mentions.is_empty() { return; diff --git a/codex-rs/analytics/src/events.rs b/codex-rs/analytics/src/events.rs index 36efb01a6..885e93bbb 100644 --- a/codex-rs/analytics/src/events.rs +++ b/codex-rs/analytics/src/events.rs @@ -1,10 +1,12 @@ use crate::facts::AppInvocation; use crate::facts::InvocationType; use crate::facts::PluginState; +use crate::facts::SubAgentThreadStartedInput; use crate::facts::TrackEventsContext; use codex_login::default_client::originator; use codex_plugin::PluginTelemetryMetadata; use codex_protocol::protocol::SessionSource; +use codex_protocol::protocol::SubAgentSource; use serde::Serialize; #[derive(Clone, Copy, Debug, Serialize)] @@ -228,3 +230,49 @@ pub(crate) fn current_runtime_metadata() -> CodexRuntimeMetadata { runtime_arch: std::env::consts::ARCH.to_string(), } } + +pub(crate) fn subagent_thread_started_event_request( + input: SubAgentThreadStartedInput, +) -> ThreadInitializedEvent { + let event_params = ThreadInitializedEventParams { + thread_id: input.thread_id, + app_server_client: CodexAppServerClientMetadata { + product_client_id: input.product_client_id, + client_name: Some(input.client_name), + client_version: Some(input.client_version), + rpc_transport: AppServerRpcTransport::InProcess, + experimental_api_enabled: None, + }, + runtime: current_runtime_metadata(), + model: input.model, + ephemeral: input.ephemeral, + thread_source: Some("subagent"), + initialization_mode: ThreadInitializationMode::New, + subagent_source: Some(subagent_source_name(&input.subagent_source)), + parent_thread_id: subagent_parent_thread_id(&input.subagent_source), + created_at: input.created_at, + }; + ThreadInitializedEvent { + event_type: "codex_thread_initialized", + event_params, + } +} + +fn subagent_source_name(subagent_source: &SubAgentSource) -> String { + match subagent_source { + SubAgentSource::Review => "review".to_string(), + SubAgentSource::Compact => "compact".to_string(), + SubAgentSource::ThreadSpawn { .. } => "thread_spawn".to_string(), + SubAgentSource::MemoryConsolidation => "memory_consolidation".to_string(), + SubAgentSource::Other(other) => other.clone(), + } +} + +fn subagent_parent_thread_id(subagent_source: &SubAgentSource) -> Option { + match subagent_source { + SubAgentSource::ThreadSpawn { + parent_thread_id, .. + } => Some(parent_thread_id.to_string()), + _ => None, + } +} diff --git a/codex-rs/analytics/src/facts.rs b/codex-rs/analytics/src/facts.rs index 31b8516e8..e19d15d84 100644 --- a/codex-rs/analytics/src/facts.rs +++ b/codex-rs/analytics/src/facts.rs @@ -7,6 +7,7 @@ use codex_app_server_protocol::RequestId; use codex_app_server_protocol::ServerNotification; use codex_plugin::PluginTelemetryMetadata; use codex_protocol::protocol::SkillScope; +use codex_protocol::protocol::SubAgentSource; use serde::Serialize; use std::path::PathBuf; @@ -50,6 +51,18 @@ pub struct AppInvocation { pub invocation_type: Option, } +#[derive(Clone)] +pub struct SubAgentThreadStartedInput { + pub thread_id: String, + pub product_client_id: String, + pub client_name: String, + pub client_version: String, + pub model: String, + pub ephemeral: bool, + pub subagent_source: SubAgentSource, + pub created_at: u64, +} + #[allow(dead_code)] pub(crate) enum AnalyticsFact { Initialize { @@ -75,6 +88,7 @@ pub(crate) enum AnalyticsFact { } pub(crate) enum CustomAnalyticsFact { + SubAgentThreadStarted(SubAgentThreadStartedInput), SkillInvoked(SkillInvokedInput), AppMentioned(AppMentionedInput), AppUsed(AppUsedInput), diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index 6f927d09c..f2f76ca8c 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -8,6 +8,7 @@ pub use events::AppServerRpcTransport; pub use facts::AppInvocation; pub use facts::InvocationType; pub use facts::SkillInvocation; +pub use facts::SubAgentThreadStartedInput; pub use facts::TrackEventsContext; pub use facts::build_track_events_context; diff --git a/codex-rs/analytics/src/reducer.rs b/codex-rs/analytics/src/reducer.rs index d83cddcbb..63b9c3d5b 100644 --- a/codex-rs/analytics/src/reducer.rs +++ b/codex-rs/analytics/src/reducer.rs @@ -15,6 +15,7 @@ use crate::events::codex_app_metadata; use crate::events::codex_plugin_metadata; use crate::events::codex_plugin_used_metadata; use crate::events::plugin_state_event_type; +use crate::events::subagent_thread_started_event_request; use crate::events::thread_source_name; use crate::facts::AnalyticsFact; use crate::facts::AppMentionedInput; @@ -24,6 +25,7 @@ use crate::facts::PluginState; use crate::facts::PluginStateChangedInput; use crate::facts::PluginUsedInput; use crate::facts::SkillInvokedInput; +use crate::facts::SubAgentThreadStartedInput; use codex_app_server_protocol::ClientResponse; use codex_app_server_protocol::InitializeParams; use codex_git_utils::collect_git_info; @@ -76,6 +78,9 @@ impl AnalyticsReducer { } AnalyticsFact::Notification(_notification) => {} AnalyticsFact::Custom(input) => match input { + CustomAnalyticsFact::SubAgentThreadStarted(input) => { + self.ingest_subagent_thread_started(input, out); + } CustomAnalyticsFact::SkillInvoked(input) => { self.ingest_skill_invoked(input, out).await; } @@ -120,6 +125,16 @@ impl AnalyticsReducer { ); } + fn ingest_subagent_thread_started( + &mut self, + input: SubAgentThreadStartedInput, + out: &mut Vec, + ) { + out.push(TrackEventRequest::ThreadInitialized( + subagent_thread_started_event_request(input), + )); + } + async fn ingest_skill_invoked( &mut self, input: SkillInvokedInput, diff --git a/codex-rs/app-server/src/codex_message_processor.rs b/codex-rs/app-server/src/codex_message_processor.rs index fc9ebf746..6475531c0 100644 --- a/codex-rs/app-server/src/codex_message_processor.rs +++ b/codex-rs/app-server/src/codex_message_processor.rs @@ -691,6 +691,7 @@ impl CodexMessageProcessor { connection_id: ConnectionId, request: ClientRequest, app_server_client_name: Option, + app_server_client_version: Option, request_context: RequestContext, ) { let to_connection_request_id = |request_id| ConnectionRequestId { @@ -707,6 +708,8 @@ impl CodexMessageProcessor { self.thread_start( to_connection_request_id(request_id), params, + app_server_client_name.clone(), + app_server_client_version.clone(), request_context, ) .await; @@ -811,6 +814,7 @@ impl CodexMessageProcessor { to_connection_request_id(request_id), params, app_server_client_name.clone(), + app_server_client_version.clone(), ) .await; } @@ -2057,6 +2061,8 @@ impl CodexMessageProcessor { &self, request_id: ConnectionRequestId, params: ThreadStartParams, + app_server_client_name: Option, + app_server_client_version: Option, request_context: RequestContext, ) { let ThreadStartParams { @@ -2112,6 +2118,8 @@ impl CodexMessageProcessor { runtime_feature_enablement, cloud_requirements, request_id, + app_server_client_name, + app_server_client_version, config, typesafe_overrides, dynamic_tools, @@ -2185,6 +2193,8 @@ impl CodexMessageProcessor { runtime_feature_enablement: BTreeMap, cloud_requirements: CloudRequirementsLoader, request_id: ConnectionRequestId, + app_server_client_name: Option, + app_server_client_version: Option, config_overrides: Option>, typesafe_overrides: ConfigOverrides, dynamic_tools: Option>, @@ -2331,6 +2341,19 @@ impl CodexMessageProcessor { session_configured, .. } = new_conv; + if let Err(error) = Self::set_app_server_client_info( + thread.as_ref(), + app_server_client_name, + app_server_client_version, + ) + .await + { + listener_task_context + .outgoing + .send_error(request_id, error) + .await; + return; + } let config_snapshot = thread .config_snapshot() .instrument(tracing::info_span!( @@ -6474,6 +6497,7 @@ impl CodexMessageProcessor { request_id: ConnectionRequestId, params: TurnStartParams, app_server_client_name: Option, + app_server_client_version: Option, ) { if let Err(error) = Self::validate_v2_input_limit(¶ms.input) { self.outgoing.send_error(request_id, error).await; @@ -6486,8 +6510,12 @@ impl CodexMessageProcessor { return; } }; - if let Err(error) = - Self::set_app_server_client_name(thread.as_ref(), app_server_client_name).await + if let Err(error) = Self::set_app_server_client_info( + thread.as_ref(), + app_server_client_name, + app_server_client_version, + ) + .await { self.outgoing.send_error(request_id, error).await; return; @@ -6581,16 +6609,17 @@ impl CodexMessageProcessor { } } - async fn set_app_server_client_name( + async fn set_app_server_client_info( thread: &CodexThread, app_server_client_name: Option, + app_server_client_version: Option, ) -> Result<(), JSONRPCErrorError> { thread - .set_app_server_client_name(app_server_client_name) + .set_app_server_client_info(app_server_client_name, app_server_client_version) .await .map_err(|err| JSONRPCErrorError { code: INTERNAL_ERROR_CODE, - message: format!("failed to set app server client name: {err}"), + message: format!("failed to set app server client info: {err}"), data: None, }) } diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index 9773b0c90..02b948bc1 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -850,6 +850,7 @@ impl MessageProcessor { connection_id, other, session.app_server_client_name.clone(), + session.client_version.clone(), request_context, ) .boxed() diff --git a/codex-rs/app-server/tests/suite/v2/thread_fork.rs b/codex-rs/app-server/tests/suite/v2/thread_fork.rs index fd65a681a..0849fe9b3 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_fork.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_fork.rs @@ -563,7 +563,7 @@ fn create_config_toml_with_chatgpt_base_url( general_analytics_enabled: bool, ) -> std::io::Result<()> { let general_analytics_toml = if general_analytics_enabled { - "\n[features]\ngeneral_analytics = true\n".to_string() + "\ngeneral_analytics = true".to_string() } else { String::new() }; @@ -578,6 +578,8 @@ sandbox_mode = "read-only" chatgpt_base_url = "{chatgpt_base_url}" model_provider = "mock_provider" + +[features] {general_analytics_toml} [model_providers.mock_provider] diff --git a/codex-rs/app-server/tests/suite/v2/thread_start.rs b/codex-rs/app-server/tests/suite/v2/thread_start.rs index 409e05404..7907e621b 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_start.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_start.rs @@ -827,7 +827,7 @@ fn create_config_toml_with_chatgpt_base_url( general_analytics_enabled: bool, ) -> std::io::Result<()> { let general_analytics_toml = if general_analytics_enabled { - "\n[features]\ngeneral_analytics = true\n".to_string() + "\ngeneral_analytics = true".to_string() } else { String::new() }; @@ -842,6 +842,8 @@ sandbox_mode = "read-only" chatgpt_base_url = "{chatgpt_base_url}" model_provider = "mock_provider" + +[features] {general_analytics_toml} [model_providers.mock_provider] diff --git a/codex-rs/core/src/agent/control.rs b/codex-rs/core/src/agent/control.rs index bec349918..255c6bc2c 100644 --- a/codex-rs/core/src/agent/control.rs +++ b/codex-rs/core/src/agent/control.rs @@ -4,6 +4,7 @@ use crate::agent::registry::AgentRegistry; use crate::agent::role::DEFAULT_ROLE_NAME; use crate::agent::role::resolve_role_config; use crate::agent::status::is_final; +use crate::codex::emit_subagent_session_started; use crate::codex_thread::ThreadConfigSnapshot; use crate::find_archived_thread_path_by_id_str; use crate::find_thread_path_by_id_str; @@ -245,6 +246,48 @@ impl AgentControl { agent_metadata.agent_id = Some(new_thread.thread_id); reservation.commit(agent_metadata.clone()); + if let Some(SessionSource::SubAgent( + subagent_source @ SubAgentSource::ThreadSpawn { + parent_thread_id, .. + }, + )) = notification_source.as_ref() + && new_thread.thread.enabled(Feature::GeneralAnalytics) + { + let client_metadata = match state.get_thread(*parent_thread_id).await { + Ok(parent_thread) => { + parent_thread + .codex + .session + .app_server_client_metadata() + .await + } + Err(error) => { + tracing::warn!( + error = %error, + parent_thread_id = %parent_thread_id, + "skipping subagent thread analytics: failed to load parent thread metadata" + ); + crate::codex::AppServerClientMetadata { + client_name: None, + client_version: None, + } + } + }; + let thread_config = new_thread.thread.codex.thread_config_snapshot().await; + emit_subagent_session_started( + &new_thread + .thread + .codex + .session + .services + .analytics_events_client, + client_metadata, + new_thread.thread_id, + thread_config, + subagent_source.clone(), + ); + } + // Notify a new thread has been created. This notification will be processed by clients // to subscribe or drain this newly created thread. // TODO(jif) add helper for drain diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 504420217..758b59204 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -5,6 +5,8 @@ use std::path::Path; use std::path::PathBuf; use std::sync::Arc; use std::sync::atomic::AtomicU64; +use std::time::SystemTime; +use std::time::UNIX_EPOCH; use crate::agent::AgentControl; use crate::agent::AgentStatus; @@ -48,6 +50,7 @@ use chrono::Utc; use codex_analytics::AnalyticsEventsClient; use codex_analytics::AppInvocation; use codex_analytics::InvocationType; +use codex_analytics::SubAgentThreadStartedInput; use codex_analytics::build_track_events_context; use codex_app_server_protocol::McpServerElicitationRequest; use codex_app_server_protocol::McpServerElicitationRequestParams; @@ -639,6 +642,7 @@ impl Codex { original_config_do_not_use: Arc::clone(&config), metrics_service_name, app_server_client_name: None, + app_server_client_version: None, session_source, dynamic_tools, persist_extended_history, @@ -758,13 +762,15 @@ impl Codex { self.session.steer_input(input, expected_turn_id).await } - pub(crate) async fn set_app_server_client_name( + pub(crate) async fn set_app_server_client_info( &self, app_server_client_name: Option, + app_server_client_version: Option, ) -> ConstraintResult<()> { self.session .update_settings(SessionSettingsUpdate { app_server_client_name, + app_server_client_version, ..Default::default() }) .await @@ -1119,6 +1125,7 @@ pub(crate) struct SessionConfiguration { /// Optional service name tag for session metrics. metrics_service_name: Option, app_server_client_name: Option, + app_server_client_version: Option, /// Source of the session (cli, vscode, exec, mcp, ...) session_source: SessionSource, dynamic_tools: Vec, @@ -1212,6 +1219,9 @@ impl SessionConfiguration { if let Some(app_server_client_name) = updates.app_server_client_name.clone() { next_configuration.app_server_client_name = Some(app_server_client_name); } + if let Some(app_server_client_version) = updates.app_server_client_version.clone() { + next_configuration.app_server_client_version = Some(app_server_client_version); + } Ok(next_configuration) } } @@ -1229,9 +1239,26 @@ pub(crate) struct SessionSettingsUpdate { pub(crate) final_output_json_schema: Option>, pub(crate) personality: Option, pub(crate) app_server_client_name: Option, + pub(crate) app_server_client_version: Option, +} + +pub(crate) struct AppServerClientMetadata { + pub(crate) client_name: Option, + pub(crate) client_version: Option, } impl Session { + pub(crate) async fn app_server_client_metadata(&self) -> AppServerClientMetadata { + let state = self.state.lock().await; + AppServerClientMetadata { + client_name: state.session_configuration.app_server_client_name.clone(), + client_version: state + .session_configuration + .app_server_client_version + .clone(), + } + } + /// Builds the `x-codex-beta-features` header value for this session. /// /// `ModelClient` is session-scoped and intentionally does not depend on the full `Config`, so @@ -4438,6 +4465,37 @@ impl Session { } } +pub(crate) fn emit_subagent_session_started( + analytics_events_client: &AnalyticsEventsClient, + client_metadata: AppServerClientMetadata, + thread_id: ThreadId, + thread_config: ThreadConfigSnapshot, + subagent_source: SubAgentSource, +) { + let AppServerClientMetadata { + client_name, + client_version, + } = client_metadata; + let (Some(client_name), Some(client_version)) = (client_name, client_version) else { + tracing::warn!("skipping subagent thread analytics: missing inherited client metadata"); + return; + }; + let created_at = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs(); + analytics_events_client.track_subagent_thread_started(SubAgentThreadStartedInput { + thread_id: thread_id.to_string(), + product_client_id: client_name.clone(), + client_name, + client_version, + model: thread_config.model, + ephemeral: thread_config.ephemeral, + subagent_source, + created_at, + }); +} + async fn submission_loop(sess: Arc, config: Arc, rx_sub: Receiver) { // To break out of this loop, send Op::Shutdown. while let Ok(sub) = rx_sub.recv().await { @@ -4805,6 +4863,7 @@ mod handlers { final_output_json_schema: Some(final_output_json_schema), personality, app_server_client_name: None, + app_server_client_version: None, }, ) } diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index bc5eb8f64..93d3eba02 100644 --- a/codex-rs/core/src/codex_delegate.rs +++ b/codex-rs/core/src/codex_delegate.rs @@ -36,6 +36,7 @@ use crate::codex::CodexSpawnOk; use crate::codex::SUBMISSION_CHANNEL_CAPACITY; use crate::codex::Session; use crate::codex::TurnContext; +use crate::codex::emit_subagent_session_started; use crate::config::Config; use crate::guardian::GuardianApprovalRequest; use crate::guardian::review_approval_request_with_cancel; @@ -85,7 +86,7 @@ pub(crate) async fn run_codex_thread_interactive( mcp_manager: Arc::clone(&parent_session.services.mcp_manager), skills_watcher: Arc::clone(&parent_session.services.skills_watcher), conversation_history: initial_history.unwrap_or(InitialHistory::New), - session_source: SessionSource::SubAgent(subagent_source), + session_source: SessionSource::SubAgent(subagent_source.clone()), agent_control: parent_session.services.agent_control.clone(), dynamic_tools: Vec::new(), persist_extended_history: false, @@ -96,6 +97,17 @@ pub(crate) async fn run_codex_thread_interactive( parent_trace: None, }) .await?; + if parent_session.enabled(codex_features::Feature::GeneralAnalytics) { + let thread_config = codex.thread_config_snapshot().await; + let client_metadata = parent_session.app_server_client_metadata().await; + emit_subagent_session_started( + &parent_session.services.analytics_events_client, + client_metadata, + codex.session.conversation_id, + thread_config, + subagent_source, + ); + } let codex = Arc::new(codex); // Use a child token so parent cancel cascades but we can scope it to this task diff --git a/codex-rs/core/src/codex_tests.rs b/codex-rs/core/src/codex_tests.rs index 9681987fe..aeb6cb48b 100644 --- a/codex-rs/core/src/codex_tests.rs +++ b/codex-rs/core/src/codex_tests.rs @@ -1843,6 +1843,7 @@ async fn set_rate_limits_retains_previous_credits() { original_config_do_not_use: Arc::clone(&config), metrics_service_name: None, app_server_client_name: None, + app_server_client_version: None, session_source: SessionSource::Exec, dynamic_tools: Vec::new(), persist_extended_history: false, @@ -1944,6 +1945,7 @@ async fn set_rate_limits_updates_plan_type_when_present() { original_config_do_not_use: Arc::clone(&config), metrics_service_name: None, app_server_client_name: None, + app_server_client_version: None, session_source: SessionSource::Exec, dynamic_tools: Vec::new(), persist_extended_history: false, @@ -2292,6 +2294,7 @@ pub(crate) async fn make_session_configuration_for_tests() -> SessionConfigurati original_config_do_not_use: Arc::clone(&config), metrics_service_name: None, app_server_client_name: None, + app_server_client_version: None, session_source: SessionSource::Exec, dynamic_tools: Vec::new(), persist_extended_history: false, @@ -2557,6 +2560,7 @@ async fn session_new_fails_when_zsh_fork_enabled_without_zsh_path() { original_config_do_not_use: Arc::clone(&config), metrics_service_name: None, app_server_client_name: None, + app_server_client_version: None, session_source: SessionSource::Exec, dynamic_tools: Vec::new(), persist_extended_history: false, @@ -2657,6 +2661,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { original_config_do_not_use: Arc::clone(&config), metrics_service_name: None, app_server_client_name: None, + app_server_client_version: None, session_source: SessionSource::Exec, dynamic_tools: Vec::new(), persist_extended_history: false, @@ -3496,6 +3501,7 @@ pub(crate) async fn make_session_and_context_with_dynamic_tools_and_rx( original_config_do_not_use: Arc::clone(&config), metrics_service_name: None, app_server_client_name: None, + app_server_client_version: None, session_source: SessionSource::Exec, dynamic_tools, persist_extended_history: false, diff --git a/codex-rs/core/src/codex_thread.rs b/codex-rs/core/src/codex_thread.rs index d274b8b03..252b1b6ae 100644 --- a/codex-rs/core/src/codex_thread.rs +++ b/codex-rs/core/src/codex_thread.rs @@ -100,12 +100,13 @@ impl CodexThread { self.codex.steer_input(input, expected_turn_id).await } - pub async fn set_app_server_client_name( + pub async fn set_app_server_client_info( &self, app_server_client_name: Option, + app_server_client_version: Option, ) -> ConstraintResult<()> { self.codex - .set_app_server_client_name(app_server_client_name) + .set_app_server_client_info(app_server_client_name, app_server_client_version) .await } diff --git a/codex-rs/core/src/memories/phase2.rs b/codex-rs/core/src/memories/phase2.rs index 78f762ed0..fde135d5e 100644 --- a/codex-rs/core/src/memories/phase2.rs +++ b/codex-rs/core/src/memories/phase2.rs @@ -1,6 +1,7 @@ use crate::agent::AgentStatus; use crate::agent::status::is_final as is_final_agent_status; use crate::codex::Session; +use crate::codex::emit_subagent_session_started; use crate::config::Config; use crate::memories::memory_root; use crate::memories::metrics; @@ -143,6 +144,26 @@ pub(super) async fn run(session: &Arc, config: Arc) { } }; + if let Some(thread_config) = session + .services + .agent_control + .get_agent_config_snapshot(thread_id) + .await + { + if session.enabled(Feature::GeneralAnalytics) { + let client_metadata = session.app_server_client_metadata().await; + emit_subagent_session_started( + &session.services.analytics_events_client, + client_metadata, + thread_id, + thread_config, + SubAgentSource::MemoryConsolidation, + ); + } + } else { + warn!("failed to load memory consolidation thread config for analytics: {thread_id}"); + } + // 6. Spawn the agent handler. agent::handle( session,