From 7daa98793d00d14aaab57390a76178694a622584 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 14:02:09 -0700 Subject: [PATCH 01/15] remove flaky test --- codex-rs/core/tests/suite/thread_metadata.rs | 107 ------------------- 1 file changed, 107 deletions(-) diff --git a/codex-rs/core/tests/suite/thread_metadata.rs b/codex-rs/core/tests/suite/thread_metadata.rs index a307d25059..23e1437582 100644 --- a/codex-rs/core/tests/suite/thread_metadata.rs +++ b/codex-rs/core/tests/suite/thread_metadata.rs @@ -1,13 +1,11 @@ use anyhow::Result; use codex_core::CodexAuth; -use codex_core::RolloutRecorder; use codex_core::config::Constrained; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::config_types::Personality; use codex_protocol::config_types::ReasoningSummary; use codex_protocol::config_types::ServiceTier; use codex_protocol::openai_models::ReasoningEffort; -use codex_protocol::protocol::InitialHistory; use core_test_support::responses::start_mock_server; use core_test_support::test_codex::test_codex; use std::time::Duration; @@ -101,108 +99,3 @@ async fn thread_initialization_tracks_thread_initialized_analytics() -> Result<( Ok(()) } - -#[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn resumed_thread_emits_thread_initialized_analytics() -> Result<()> { - let server = start_mock_server().await; - let chatgpt_base_url = server.uri(); - - let initial = test_codex() - .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) - .with_config({ - let chatgpt_base_url = chatgpt_base_url.clone(); - move |config| { - config.chatgpt_base_url = chatgpt_base_url; - config.model = Some("gpt-5".to_string()); - config.model_reasoning_effort = Some(ReasoningEffort::High); - config.model_reasoning_summary = Some(ReasoningSummary::Detailed); - config.service_tier = Some(ServiceTier::Flex); - config.approvals_reviewer = ApprovalsReviewer::GuardianSubagent; - config.permissions.sandbox_policy = Constrained::allow_any( - codex_protocol::protocol::SandboxPolicy::new_workspace_write_policy(), - ); - config.personality = Some(Personality::Friendly); - } - }) - .build(&server) - .await?; - - let rollout_path = initial - .codex - .rollout_path() - .expect("rollout path for initial thread"); - let resume_deadline = Instant::now() + Duration::from_secs(10); - loop { - if matches!( - RolloutRecorder::get_rollout_history(&rollout_path).await?, - InitialHistory::Resumed(_) - ) { - break; - } - if Instant::now() >= resume_deadline { - panic!("timed out waiting for rollout to become resumable"); - } - tokio::time::sleep(Duration::from_millis(50)).await; - } - - let _resumed = test_codex() - .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) - .with_config(move |config| { - config.chatgpt_base_url = chatgpt_base_url; - config.model = Some("gpt-5".to_string()); - config.model_reasoning_effort = Some(ReasoningEffort::High); - config.model_reasoning_summary = Some(ReasoningSummary::Detailed); - config.service_tier = Some(ServiceTier::Flex); - config.approvals_reviewer = ApprovalsReviewer::GuardianSubagent; - config.permissions.sandbox_policy = Constrained::allow_any( - codex_protocol::protocol::SandboxPolicy::new_workspace_write_policy(), - ); - config.personality = Some(Personality::Friendly); - }) - .resume(&server, initial.home.clone(), rollout_path) - .await?; - - let deadline = Instant::now() + Duration::from_secs(10); - let analytics_request = loop { - let requests = server.received_requests().await.unwrap_or_default(); - if let Some(request) = requests.into_iter().find(|request| { - if request.url.path() != "/codex/analytics-events/events" { - return false; - } - let Ok(payload) = serde_json::from_slice::(&request.body) else { - return false; - }; - payload["events"] - .as_array() - .into_iter() - .flatten() - .any(|event| { - event["event_type"] == "codex_thread_initialized" - && event["event_params"]["initialization_mode"] == "resumed" - }) - }) { - break request; - } - if Instant::now() >= deadline { - panic!("timed out waiting for resumed thread analytics request"); - } - tokio::time::sleep(Duration::from_millis(50)).await; - }; - - let payload: serde_json::Value = - serde_json::from_slice(&analytics_request.body).expect("analytics payload"); - let event = payload["events"] - .as_array() - .expect("events array") - .iter() - .find(|event| { - event["event_type"] == "codex_thread_initialized" - && event["event_params"]["initialization_mode"] == "resumed" - }) - .expect("codex_thread_initialized resumed event should be present"); - - assert_eq!(event["event_params"]["session_source"], "user"); - assert_eq!(event["event_params"]["initialization_mode"], "resumed"); - - Ok(()) -} From 6c6d49a49a829181a57209b404bffaa67ad14bf0 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 14:24:15 -0700 Subject: [PATCH 02/15] lint --- codex-rs/core/tests/suite/items.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/codex-rs/core/tests/suite/items.rs b/codex-rs/core/tests/suite/items.rs index 81fca14e2f..a0f80f29bb 100644 --- a/codex-rs/core/tests/suite/items.rs +++ b/codex-rs/core/tests/suite/items.rs @@ -162,7 +162,7 @@ async fn user_turn_tracks_turn_metadata_analytics() -> anyhow::Result<()> { text_elements: Vec::new(), }], final_output_json_schema: None, - cwd: config.cwd.clone(), + cwd: config.cwd.to_path_buf(), approval_policy: AskForApproval::Never, approvals_reviewer: None, sandbox_policy: SandboxPolicy::new_read_only_policy(), From 4a163ee2a935ab5760eee79dca09040336c3191a Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 14:36:12 -0700 Subject: [PATCH 03/15] remove turn context --- codex-rs/analytics/src/analytics_client.rs | 61 ------------------- .../analytics/src/analytics_client_tests.rs | 48 --------------- codex-rs/core/src/codex.rs | 20 ------ codex-rs/core/tests/suite/thread_metadata.rs | 23 ------- 4 files changed, 152 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index 8d3e7ea72b..be45543274 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -4,14 +4,6 @@ use codex_login::AuthManager; use codex_login::default_client::create_client; use codex_login::default_client::originator; use codex_plugin::PluginTelemetryMetadata; -use codex_protocol::config_types::ApprovalsReviewer; -use codex_protocol::config_types::ModeKind; -use codex_protocol::config_types::Personality; -use codex_protocol::config_types::ReasoningSummary; -use codex_protocol::config_types::ServiceTier; -use codex_protocol::openai_models::ReasoningEffort; -use codex_protocol::protocol::AskForApproval; -use codex_protocol::protocol::SandboxPolicy; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SkillScope; use codex_protocol::protocol::SubAgentSource; @@ -37,16 +29,6 @@ pub struct TrackEventsContext { pub struct CodexThreadInitializedEvent { pub thread_id: String, pub model: String, - pub model_provider: String, - pub reasoning_effort: Option, - pub reasoning_summary: Option, - pub service_tier: Option, - pub approval_policy: AskForApproval, - pub approvals_reviewer: ApprovalsReviewer, - pub sandbox_policy: SandboxPolicy, - pub sandbox_network_access: bool, - pub collaboration_mode: ModeKind, - pub personality: Option, pub ephemeral: bool, pub session_source: SessionSource, pub initialization_mode: InitializationMode, @@ -372,16 +354,6 @@ struct CodexThreadInitializedEventParams { thread_id: String, product_client_id: String, model: String, - model_provider: String, - reasoning_effort: Option, - reasoning_summary: Option, - service_tier: String, - approval_policy: String, - approvals_reviewer: String, - sandbox_policy: &'static str, - sandbox_network_access: bool, - collaboration_mode: &'static str, - personality: Option, ephemeral: bool, session_source: Option<&'static str>, initialization_mode: InitializationMode, @@ -786,21 +758,6 @@ fn codex_thread_initialized_event_params( thread_id: thread_event.thread_id, product_client_id: originator().value, model: thread_event.model, - model_provider: thread_event.model_provider, - reasoning_effort: thread_event.reasoning_effort.map(|value| value.to_string()), - reasoning_summary: thread_event - .reasoning_summary - .map(|value| value.to_string()), - service_tier: thread_event - .service_tier - .map(|value| value.to_string()) - .unwrap_or_else(|| "default".to_string()), - approval_policy: thread_event.approval_policy.to_string(), - approvals_reviewer: thread_event.approvals_reviewer.to_string(), - sandbox_policy: sandbox_policy_mode(&thread_event.sandbox_policy), - sandbox_network_access: thread_event.sandbox_network_access, - collaboration_mode: collaboration_mode_mode(thread_event.collaboration_mode), - personality: thread_event.personality.map(|value| value.to_string()), ephemeral: thread_event.ephemeral, session_source: session_source_name(&thread_event.session_source), initialization_mode: thread_event.initialization_mode, @@ -845,24 +802,6 @@ fn codex_plugin_used_metadata( } } -fn sandbox_policy_mode(sandbox_policy: &SandboxPolicy) -> &'static str { - match sandbox_policy { - SandboxPolicy::DangerFullAccess => "full_access", - SandboxPolicy::ReadOnly { .. } => "read_only", - SandboxPolicy::WorkspaceWrite { .. } => "workspace_write", - SandboxPolicy::ExternalSandbox { .. } => "external_sandbox", - } -} - -fn collaboration_mode_mode(mode: ModeKind) -> &'static str { - match mode { - ModeKind::Plan => "plan", - ModeKind::Default => "default", - ModeKind::PairProgramming => "pair_programming", - ModeKind::Execute => "execute", - } -} - fn session_source_name(session_source: &SessionSource) -> Option<&'static str> { match session_source { SessionSource::Cli | SessionSource::VSCode | SessionSource::Exec => Some("user"), diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index 964428beee..eaac08430f 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -20,14 +20,6 @@ use codex_plugin::AppConnectorId; use codex_plugin::PluginCapabilitySummary; use codex_plugin::PluginId; use codex_plugin::PluginTelemetryMetadata; -use codex_protocol::config_types::ApprovalsReviewer; -use codex_protocol::config_types::ModeKind; -use codex_protocol::config_types::Personality; -use codex_protocol::config_types::ReasoningSummary; -use codex_protocol::config_types::ServiceTier; -use codex_protocol::openai_models::ReasoningEffort; -use codex_protocol::protocol::AskForApproval; -use codex_protocol::protocol::SandboxPolicy; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; use pretty_assertions::assert_eq; @@ -207,16 +199,6 @@ fn thread_initialized_event_serializes_expected_shape() { event_params: codex_thread_initialized_event_params(CodexThreadInitializedEvent { thread_id: "thread-0".to_string(), model: "gpt-5".to_string(), - model_provider: "openai".to_string(), - reasoning_effort: Some(ReasoningEffort::High), - reasoning_summary: Some(ReasoningSummary::Detailed), - service_tier: Some(ServiceTier::Flex), - approval_policy: AskForApproval::OnRequest, - approvals_reviewer: ApprovalsReviewer::GuardianSubagent, - sandbox_policy: SandboxPolicy::new_read_only_policy(), - sandbox_network_access: false, - collaboration_mode: ModeKind::Plan, - personality: Some(Personality::Friendly), ephemeral: true, session_source: SessionSource::Exec, initialization_mode: InitializationMode::New, @@ -236,16 +218,6 @@ fn thread_initialized_event_serializes_expected_shape() { "thread_id": "thread-0", "product_client_id": originator().value, "model": "gpt-5", - "model_provider": "openai", - "reasoning_effort": "high", - "reasoning_summary": "detailed", - "service_tier": "flex", - "approval_policy": "on-request", - "approvals_reviewer": "guardian_subagent", - "sandbox_policy": "read_only", - "sandbox_network_access": false, - "collaboration_mode": "plan", - "personality": "friendly", "ephemeral": true, "session_source": "user", "initialization_mode": "new", @@ -264,16 +236,6 @@ fn thread_initialized_event_serializes_subagent_source() { event_params: codex_thread_initialized_event_params(CodexThreadInitializedEvent { thread_id: "thread-1".to_string(), model: "gpt-5".to_string(), - model_provider: "openai".to_string(), - reasoning_effort: None, - reasoning_summary: None, - service_tier: None, - approval_policy: AskForApproval::OnRequest, - approvals_reviewer: ApprovalsReviewer::User, - sandbox_policy: SandboxPolicy::new_read_only_policy(), - sandbox_network_access: false, - collaboration_mode: ModeKind::Default, - personality: None, ephemeral: false, session_source: SessionSource::SubAgent(SubAgentSource::Review), initialization_mode: InitializationMode::New, @@ -296,16 +258,6 @@ fn thread_initialized_event_omits_non_user_non_subagent_session_source() { event_params: codex_thread_initialized_event_params(CodexThreadInitializedEvent { thread_id: "thread-2".to_string(), model: "gpt-5".to_string(), - model_provider: "openai".to_string(), - reasoning_effort: None, - reasoning_summary: None, - service_tier: None, - approval_policy: AskForApproval::OnRequest, - approvals_reviewer: ApprovalsReviewer::User, - sandbox_policy: SandboxPolicy::new_read_only_policy(), - sandbox_network_access: false, - collaboration_mode: ModeKind::Default, - personality: None, ephemeral: false, session_source: SessionSource::Mcp, initialization_mode: InitializationMode::New, diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 7fb1e80d75..e3bf4049ff 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -676,26 +676,6 @@ impl Codex { .collaboration_mode .model() .to_string(), - model_provider: thread_initialized_configuration - .original_config_do_not_use - .model_provider_id - .clone(), - reasoning_effort: thread_initialized_configuration - .collaboration_mode - .reasoning_effort(), - reasoning_summary: thread_initialized_configuration.model_reasoning_summary, - service_tier: thread_initialized_configuration.service_tier, - approval_policy: thread_initialized_configuration.approval_policy.value(), - approvals_reviewer: thread_initialized_configuration.approvals_reviewer, - sandbox_policy: thread_initialized_configuration - .sandbox_policy - .get() - .clone(), - sandbox_network_access: thread_initialized_configuration - .network_sandbox_policy - .is_enabled(), - collaboration_mode: thread_initialized_configuration.collaboration_mode.mode, - personality: thread_initialized_configuration.personality, ephemeral: thread_initialized_configuration .original_config_do_not_use .ephemeral, diff --git a/codex-rs/core/tests/suite/thread_metadata.rs b/codex-rs/core/tests/suite/thread_metadata.rs index 23e1437582..3cf0e22ecd 100644 --- a/codex-rs/core/tests/suite/thread_metadata.rs +++ b/codex-rs/core/tests/suite/thread_metadata.rs @@ -1,11 +1,6 @@ use anyhow::Result; use codex_core::CodexAuth; use codex_core::config::Constrained; -use codex_protocol::config_types::ApprovalsReviewer; -use codex_protocol::config_types::Personality; -use codex_protocol::config_types::ReasoningSummary; -use codex_protocol::config_types::ServiceTier; -use codex_protocol::openai_models::ReasoningEffort; use core_test_support::responses::start_mock_server; use core_test_support::test_codex::test_codex; use std::time::Duration; @@ -21,14 +16,9 @@ async fn thread_initialization_tracks_thread_initialized_analytics() -> Result<( .with_config(move |config| { config.chatgpt_base_url = chatgpt_base_url; config.model = Some("gpt-5".to_string()); - config.model_reasoning_effort = Some(ReasoningEffort::High); - config.model_reasoning_summary = Some(ReasoningSummary::Detailed); - config.service_tier = Some(ServiceTier::Flex); - config.approvals_reviewer = ApprovalsReviewer::GuardianSubagent; config.permissions.sandbox_policy = Constrained::allow_any( codex_protocol::protocol::SandboxPolicy::new_workspace_write_policy(), ); - config.personality = Some(Personality::Friendly); config.ephemeral = true; }) .build(&server) @@ -67,19 +57,6 @@ async fn thread_initialization_tracks_thread_initialized_analytics() -> Result<( serde_json::json!(codex_core::default_client::originator().value) ); assert_eq!(event["event_params"]["model"], "gpt-5"); - assert_eq!(event["event_params"]["model_provider"], "openai"); - assert_eq!(event["event_params"]["reasoning_effort"], "high"); - assert_eq!(event["event_params"]["reasoning_summary"], "detailed"); - assert_eq!(event["event_params"]["service_tier"], "flex"); - assert_eq!(event["event_params"]["approval_policy"], "on-request"); - assert_eq!( - event["event_params"]["approvals_reviewer"], - "guardian_subagent" - ); - assert_eq!(event["event_params"]["sandbox_policy"], "workspace_write"); - assert_eq!(event["event_params"]["sandbox_network_access"], false); - assert_eq!(event["event_params"]["collaboration_mode"], "default"); - assert_eq!(event["event_params"]["personality"], "friendly"); assert_eq!(event["event_params"]["ephemeral"], true); assert_eq!(event["event_params"]["session_source"], "user"); assert_eq!(event["event_params"]["initialization_mode"], "new"); From f4790429c1b09fb45f4186eb021d2ec97bcd134a Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 14:54:33 -0700 Subject: [PATCH 04/15] fix test --- codex-rs/core/tests/suite/items.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/codex-rs/core/tests/suite/items.rs b/codex-rs/core/tests/suite/items.rs index a0f80f29bb..5ecdb879fa 100644 --- a/codex-rs/core/tests/suite/items.rs +++ b/codex-rs/core/tests/suite/items.rs @@ -325,10 +325,10 @@ async fn resumed_thread_turn_tracks_is_not_first_turn_analytics() -> anyhow::Res }) .filter(|event| event["event_type"] == "codex_turn_event") .collect::>(); - if let Some(event) = turn_events.last().cloned() { - if turn_events.len() >= 2 { - break event; - } + if let Some(event) = turn_events.last().cloned() + && turn_events.len() >= 2 + { + break event; } if Instant::now() >= deadline { panic!("timed out waiting for resumed turn analytics event"); From b41bc22e4535fb63755a9c19d994cb3da8f7cdd3 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:05:08 -0700 Subject: [PATCH 05/15] establish reducer structure --- codex-rs/analytics/src/analytics_client.rs | 672 ++++++++------------- codex-rs/analytics/src/lib.rs | 9 + 2 files changed, 257 insertions(+), 424 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index be45543274..1bccea915c 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -10,6 +10,7 @@ use codex_protocol::protocol::SubAgentSource; use serde::Serialize; use sha1::Digest; use sha1::Sha1; +use std::collections::HashMap; use std::collections::HashSet; use std::path::Path; use std::path::PathBuf; @@ -78,9 +79,64 @@ pub struct AppInvocation { pub invocation_type: Option, } +pub enum AnalyticsInput { + ThreadInitialized(ThreadInitializedInput), + SkillInvoked(SkillInvokedInput), + AppMentioned(AppMentionedInput), + AppUsed(AppUsedInput), + PluginUsed(PluginUsedInput), + PluginStateChanged(PluginStateChangedInput), +} + +pub struct ThreadInitializedInput { + pub thread_event: CodexThreadInitializedEvent, +} + +pub struct SkillInvokedInput { + pub tracking: TrackEventsContext, + pub invocations: Vec, +} + +pub struct AppMentionedInput { + pub tracking: TrackEventsContext, + pub mentions: Vec, +} + +pub struct AppUsedInput { + pub tracking: TrackEventsContext, + pub app: AppInvocation, +} + +pub struct PluginUsedInput { + pub tracking: TrackEventsContext, + pub plugin: PluginTelemetryMetadata, +} + +pub struct PluginStateChangedInput { + pub plugin: PluginTelemetryMetadata, + pub state: PluginState, +} + +#[derive(Clone, Copy)] +pub enum PluginState { + Installed, + Uninstalled, + Enabled, + Disabled, +} + +#[derive(Default)] +pub struct AnalyticsReducer { + threads: HashMap, +} + +struct ThreadState { + _initialized_event: CodexThreadInitializedEvent, +} + #[derive(Clone)] pub(crate) struct AnalyticsEventsQueue { - sender: mpsc::Sender, + sender: mpsc::Sender, app_used_emitted_keys: Arc>>, plugin_used_emitted_keys: Arc>>, } @@ -95,36 +151,11 @@ impl AnalyticsEventsQueue { pub(crate) fn new(auth_manager: Arc, base_url: String) -> Self { let (sender, mut receiver) = mpsc::channel(ANALYTICS_EVENTS_QUEUE_SIZE); tokio::spawn(async move { + let mut reducer = AnalyticsReducer::default(); while let Some(job) = receiver.recv().await { - match job { - TrackEventsJob::SkillInvocations(job) => { - send_track_skill_invocations(&auth_manager, &base_url, job).await; - } - TrackEventsJob::ThreadInitialized(job) => { - send_track_thread_initialized(&auth_manager, &base_url, job).await; - } - TrackEventsJob::AppMentioned(job) => { - send_track_app_mentioned(&auth_manager, &base_url, job).await; - } - TrackEventsJob::AppUsed(job) => { - send_track_app_used(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginUsed(job) => { - send_track_plugin_used(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginInstalled(job) => { - send_track_plugin_installed(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginUninstalled(job) => { - send_track_plugin_uninstalled(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginEnabled(job) => { - send_track_plugin_enabled(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginDisabled(job) => { - send_track_plugin_disabled(&auth_manager, &base_url, job).await; - } - } + let mut events = Vec::new(); + reducer.ingest(job, &mut events).await; + send_track_events(&auth_manager, &base_url, events).await; } }); Self { @@ -134,8 +165,8 @@ impl AnalyticsEventsQueue { } } - fn try_send(&self, job: TrackEventsJob) { - if self.sender.try_send(job).is_err() { + fn try_send(&self, input: AnalyticsInput) { + if self.sender.try_send(input).is_err() { //TODO: add a metric for this tracing::warn!("dropping analytics events: queue is full"); } @@ -188,124 +219,90 @@ impl AnalyticsEventsClient { tracking: TrackEventsContext, invocations: Vec, ) { - track_skill_invocations( - &self.queue, - self.analytics_enabled, - Some(tracking), + if invocations.is_empty() { + return; + } + self.record(AnalyticsInput::SkillInvoked(SkillInvokedInput { + tracking, invocations, - ); + })); } pub fn track_thread_initialized(&self, thread_event: CodexThreadInitializedEvent) { - track_thread_initialized(&self.queue, self.analytics_enabled, thread_event); + self.record(AnalyticsInput::ThreadInitialized(ThreadInitializedInput { + thread_event, + })); } pub fn track_app_mentioned(&self, tracking: TrackEventsContext, mentions: Vec) { - track_app_mentioned( - &self.queue, - self.analytics_enabled, - Some(tracking), + if mentions.is_empty() { + return; + } + self.record(AnalyticsInput::AppMentioned(AppMentionedInput { + tracking, mentions, - ); + })); } pub fn track_app_used(&self, tracking: TrackEventsContext, app: AppInvocation) { - track_app_used(&self.queue, self.analytics_enabled, Some(tracking), app); + if !self.queue.should_enqueue_app_used(&tracking, &app) { + return; + } + self.record(AnalyticsInput::AppUsed(AppUsedInput { tracking, app })); } pub fn track_plugin_used(&self, tracking: TrackEventsContext, plugin: PluginTelemetryMetadata) { - track_plugin_used(&self.queue, self.analytics_enabled, Some(tracking), plugin); + if !self.queue.should_enqueue_plugin_used(&tracking, &plugin) { + return; + } + self.record(AnalyticsInput::PluginUsed(PluginUsedInput { + tracking, + plugin, + })); } pub fn track_plugin_installed(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Installed, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Installed, + }, + )); } pub fn track_plugin_uninstalled(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Uninstalled, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Uninstalled, + }, + )); } pub fn track_plugin_enabled(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Enabled, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Enabled, + }, + )); } pub fn track_plugin_disabled(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Disabled, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Disabled, + }, + )); } -} -enum TrackEventsJob { - SkillInvocations(TrackSkillInvocationsJob), - ThreadInitialized(TrackThreadInitializedJob), - AppMentioned(TrackAppMentionedJob), - AppUsed(TrackAppUsedJob), - PluginUsed(TrackPluginUsedJob), - PluginInstalled(TrackPluginManagementJob), - PluginUninstalled(TrackPluginManagementJob), - PluginEnabled(TrackPluginManagementJob), - PluginDisabled(TrackPluginManagementJob), -} - -struct TrackSkillInvocationsJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - invocations: Vec, -} - -struct TrackThreadInitializedJob { - analytics_enabled: Option, - thread_event: CodexThreadInitializedEvent, -} - -struct TrackAppMentionedJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - mentions: Vec, -} - -struct TrackAppUsedJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - app: AppInvocation, -} - -struct TrackPluginUsedJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - plugin: PluginTelemetryMetadata, -} - -struct TrackPluginManagementJob { - analytics_enabled: Option, - plugin: PluginTelemetryMetadata, -} - -#[derive(Clone, Copy)] -enum PluginManagementEventType { - Installed, - Uninstalled, - Enabled, - Disabled, + pub fn record(&self, input: AnalyticsInput) { + if self.analytics_enabled == Some(false) { + return; + } + self.queue.try_send(input); + } } const ANALYTICS_EVENTS_QUEUE_SIZE: usize = 256; @@ -423,320 +420,151 @@ struct CodexPluginUsedEventRequest { event_params: CodexPluginUsedMetadata, } -pub(crate) fn track_skill_invocations( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - invocations: Vec, -) { - if analytics_enabled == Some(false) { - return; +impl AnalyticsReducer { + async fn ingest(&mut self, input: AnalyticsInput, out: &mut Vec) { + match input { + AnalyticsInput::ThreadInitialized(input) => { + self.ingest_thread_initialized(input, out); + } + AnalyticsInput::SkillInvoked(input) => { + self.ingest_skill_invoked(input, out).await; + } + AnalyticsInput::AppMentioned(input) => { + self.ingest_app_mentioned(input, out); + } + AnalyticsInput::AppUsed(input) => { + self.ingest_app_used(input, out); + } + AnalyticsInput::PluginUsed(input) => { + self.ingest_plugin_used(input, out); + } + AnalyticsInput::PluginStateChanged(input) => { + self.ingest_plugin_state_changed(input, out); + } + } } - let Some(tracking) = tracking else { - return; - }; - if invocations.is_empty() { - return; - } - let job = TrackEventsJob::SkillInvocations(TrackSkillInvocationsJob { - analytics_enabled, - tracking, - invocations, - }); - queue.try_send(job); -} -pub(crate) fn track_thread_initialized( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - thread_event: CodexThreadInitializedEvent, -) { - if analytics_enabled == Some(false) { - return; - } - let job = TrackEventsJob::ThreadInitialized(TrackThreadInitializedJob { - analytics_enabled, - thread_event, - }); - queue.try_send(job); -} - -pub(crate) fn track_app_mentioned( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - mentions: Vec, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - if mentions.is_empty() { - return; - } - let job = TrackEventsJob::AppMentioned(TrackAppMentionedJob { - analytics_enabled, - tracking, - mentions, - }); - queue.try_send(job); -} - -pub(crate) fn track_app_used( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - app: AppInvocation, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - if !queue.should_enqueue_app_used(&tracking, &app) { - return; - } - let job = TrackEventsJob::AppUsed(TrackAppUsedJob { - analytics_enabled, - tracking, - app, - }); - queue.try_send(job); -} - -pub(crate) fn track_plugin_used( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - plugin: PluginTelemetryMetadata, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - if !queue.should_enqueue_plugin_used(&tracking, &plugin) { - return; - } - let job = TrackEventsJob::PluginUsed(TrackPluginUsedJob { - analytics_enabled, - tracking, - plugin, - }); - queue.try_send(job); -} - -fn track_plugin_management( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - event_type: PluginManagementEventType, - plugin: PluginTelemetryMetadata, -) { - if analytics_enabled == Some(false) { - return; - } - let job = TrackPluginManagementJob { - analytics_enabled, - plugin, - }; - let job = match event_type { - PluginManagementEventType::Installed => TrackEventsJob::PluginInstalled(job), - PluginManagementEventType::Uninstalled => TrackEventsJob::PluginUninstalled(job), - PluginManagementEventType::Enabled => TrackEventsJob::PluginEnabled(job), - PluginManagementEventType::Disabled => TrackEventsJob::PluginDisabled(job), - }; - queue.try_send(job); -} - -async fn send_track_skill_invocations( - auth_manager: &AuthManager, - base_url: &str, - job: TrackSkillInvocationsJob, -) { - let TrackSkillInvocationsJob { - analytics_enabled, - tracking, - invocations, - } = job; - let mut events = Vec::with_capacity(invocations.len()); - for invocation in invocations { - let skill_scope = match invocation.skill_scope { - SkillScope::User => "user", - SkillScope::Repo => "repo", - SkillScope::System => "system", - SkillScope::Admin => "admin", - }; - let repo_root = get_git_repo_root(invocation.skill_path.as_path()); - let repo_url = if let Some(root) = repo_root.as_ref() { - collect_git_info(root) - .await - .and_then(|info| info.repository_url) - } else { - None - }; - let skill_id = skill_id_for_local_skill( - repo_url.as_deref(), - repo_root.as_deref(), - invocation.skill_path.as_path(), - invocation.skill_name.as_str(), + fn ingest_thread_initialized( + &mut self, + input: ThreadInitializedInput, + out: &mut Vec, + ) { + self.threads.insert( + input.thread_event.thread_id.clone(), + ThreadState { + _initialized_event: input.thread_event.clone(), + }, ); - events.push(TrackEventRequest::SkillInvocation( - SkillInvocationEventRequest { - event_type: "skill_invocation", - skill_id, - skill_name: invocation.skill_name.clone(), - event_params: SkillInvocationEventParams { - thread_id: Some(tracking.thread_id.clone()), - invoke_type: Some(invocation.invocation_type), - model_slug: Some(tracking.model_slug.clone()), - product_client_id: Some(originator().value), - repo_url, - skill_scope: Some(skill_scope.to_string()), - }, + out.push(TrackEventRequest::ThreadInitialized( + CodexThreadInitializedEventRequest { + event_type: "codex_thread_initialized", + event_params: codex_thread_initialized_event_params(input.thread_event), }, )); } - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} + async fn ingest_skill_invoked( + &mut self, + input: SkillInvokedInput, + out: &mut Vec, + ) { + let SkillInvokedInput { + tracking, + invocations, + } = input; + for invocation in invocations { + let skill_scope = match invocation.skill_scope { + SkillScope::User => "user", + SkillScope::Repo => "repo", + SkillScope::System => "system", + SkillScope::Admin => "admin", + }; + let repo_root = get_git_repo_root(invocation.skill_path.as_path()); + let repo_url = if let Some(root) = repo_root.as_ref() { + collect_git_info(root) + .await + .and_then(|info| info.repository_url) + } else { + None + }; + let skill_id = skill_id_for_local_skill( + repo_url.as_deref(), + repo_root.as_deref(), + invocation.skill_path.as_path(), + invocation.skill_name.as_str(), + ); + out.push(TrackEventRequest::SkillInvocation( + SkillInvocationEventRequest { + event_type: "skill_invocation", + skill_id, + skill_name: invocation.skill_name.clone(), + event_params: SkillInvocationEventParams { + thread_id: Some(tracking.thread_id.clone()), + invoke_type: Some(invocation.invocation_type), + model_slug: Some(tracking.model_slug.clone()), + product_client_id: Some(originator().value), + repo_url, + skill_scope: Some(skill_scope.to_string()), + }, + }, + )); + } + } -async fn send_track_thread_initialized( - auth_manager: &AuthManager, - base_url: &str, - job: TrackThreadInitializedJob, -) { - let TrackThreadInitializedJob { - analytics_enabled, - thread_event, - } = job; - let events = vec![TrackEventRequest::ThreadInitialized( - CodexThreadInitializedEventRequest { - event_type: "codex_thread_initialized", - event_params: codex_thread_initialized_event_params(thread_event), - }, - )]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} - -async fn send_track_app_mentioned( - auth_manager: &AuthManager, - base_url: &str, - job: TrackAppMentionedJob, -) { - let TrackAppMentionedJob { - analytics_enabled, - tracking, - mentions, - } = job; - let events = mentions - .into_iter() - .map(|mention| { + fn ingest_app_mentioned(&mut self, input: AppMentionedInput, out: &mut Vec) { + let AppMentionedInput { tracking, mentions } = input; + out.extend(mentions.into_iter().map(|mention| { let event_params = codex_app_metadata(&tracking, mention); TrackEventRequest::AppMentioned(CodexAppMentionedEventRequest { event_type: "codex_app_mentioned", event_params, }) - }) - .collect::>(); + })); + } - send_track_events(auth_manager, analytics_enabled, base_url, events).await; + fn ingest_app_used(&mut self, input: AppUsedInput, out: &mut Vec) { + let AppUsedInput { tracking, app } = input; + let event_params = codex_app_metadata(&tracking, app); + out.push(TrackEventRequest::AppUsed(CodexAppUsedEventRequest { + event_type: "codex_app_used", + event_params, + })); + } + + fn ingest_plugin_used(&mut self, input: PluginUsedInput, out: &mut Vec) { + let PluginUsedInput { tracking, plugin } = input; + out.push(TrackEventRequest::PluginUsed(CodexPluginUsedEventRequest { + event_type: "codex_plugin_used", + event_params: codex_plugin_used_metadata(&tracking, plugin), + })); + } + + fn ingest_plugin_state_changed( + &mut self, + input: PluginStateChangedInput, + out: &mut Vec, + ) { + let PluginStateChangedInput { plugin, state } = input; + let event = CodexPluginEventRequest { + event_type: plugin_state_event_type(state), + event_params: codex_plugin_metadata(plugin), + }; + out.push(match state { + PluginState::Installed => TrackEventRequest::PluginInstalled(event), + PluginState::Uninstalled => TrackEventRequest::PluginUninstalled(event), + PluginState::Enabled => TrackEventRequest::PluginEnabled(event), + PluginState::Disabled => TrackEventRequest::PluginDisabled(event), + }); + } } -async fn send_track_app_used(auth_manager: &AuthManager, base_url: &str, job: TrackAppUsedJob) { - let TrackAppUsedJob { - analytics_enabled, - tracking, - app, - } = job; - let event_params = codex_app_metadata(&tracking, app); - let events = vec![TrackEventRequest::AppUsed(CodexAppUsedEventRequest { - event_type: "codex_app_used", - event_params, - })]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} - -async fn send_track_plugin_used( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginUsedJob, -) { - let TrackPluginUsedJob { - analytics_enabled, - tracking, - plugin, - } = job; - let events = vec![TrackEventRequest::PluginUsed(CodexPluginUsedEventRequest { - event_type: "codex_plugin_used", - event_params: codex_plugin_used_metadata(&tracking, plugin), - })]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} - -async fn send_track_plugin_installed( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_installed").await; -} - -async fn send_track_plugin_uninstalled( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_uninstalled") - .await; -} - -async fn send_track_plugin_enabled( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_enabled").await; -} - -async fn send_track_plugin_disabled( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_disabled").await; -} - -async fn send_track_plugin_management_event( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, - event_type: &'static str, -) { - let TrackPluginManagementJob { - analytics_enabled, - plugin, - } = job; - let event_params = codex_plugin_metadata(plugin); - let event = CodexPluginEventRequest { - event_type, - event_params, - }; - let events = vec![match event_type { - "codex_plugin_installed" => TrackEventRequest::PluginInstalled(event), - "codex_plugin_uninstalled" => TrackEventRequest::PluginUninstalled(event), - "codex_plugin_enabled" => TrackEventRequest::PluginEnabled(event), - "codex_plugin_disabled" => TrackEventRequest::PluginDisabled(event), - _ => unreachable!("unknown plugin management event type"), - }]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; +fn plugin_state_event_type(state: PluginState) -> &'static str { + match state { + PluginState::Installed => "codex_plugin_installed", + PluginState::Uninstalled => "codex_plugin_uninstalled", + PluginState::Enabled => "codex_plugin_enabled", + PluginState::Disabled => "codex_plugin_disabled", + } } fn codex_app_metadata(tracking: &TrackEventsContext, app: AppInvocation) -> CodexAppMetadata { @@ -822,13 +650,9 @@ fn subagent_source_name(subagent_source: SubAgentSource) -> String { async fn send_track_events( auth_manager: &AuthManager, - analytics_enabled: Option, base_url: &str, events: Vec, ) { - if analytics_enabled == Some(false) { - return; - } if events.is_empty() { return; } diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index bf9fab83a4..63455c0616 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -1,10 +1,19 @@ mod analytics_client; pub use analytics_client::AnalyticsEventsClient; +pub use analytics_client::AnalyticsInput; +pub use analytics_client::AnalyticsReducer; pub use analytics_client::AppInvocation; +pub use analytics_client::AppMentionedInput; +pub use analytics_client::AppUsedInput; pub use analytics_client::CodexThreadInitializedEvent; pub use analytics_client::InitializationMode; pub use analytics_client::InvocationType; +pub use analytics_client::PluginState; +pub use analytics_client::PluginStateChangedInput; +pub use analytics_client::PluginUsedInput; pub use analytics_client::SkillInvocation; +pub use analytics_client::SkillInvokedInput; +pub use analytics_client::ThreadInitializedInput; pub use analytics_client::TrackEventsContext; pub use analytics_client::build_track_events_context; From f0bc8dec3ea62ef67925b2b4366afec03b17873a Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:28:52 -0700 Subject: [PATCH 06/15] consolidate thread context --- codex-rs/analytics/src/analytics_client.rs | 47 ++++++++++++++++------ codex-rs/analytics/src/lib.rs | 3 +- codex-rs/core/src/codex.rs | 21 +++++----- 3 files changed, 48 insertions(+), 23 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index 1bccea915c..f42059a622 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -38,6 +38,23 @@ pub struct CodexThreadInitializedEvent { pub created_at: u64, } +#[derive(Clone)] +pub struct CodexThreadInitializedInput { + pub thread_id: String, + pub model: String, + pub created_at: u64, + pub thread_context: CodexThreadContext, +} + +#[derive(Clone)] +pub struct CodexThreadContext { + pub ephemeral: bool, + pub session_source: SessionSource, + pub initialization_mode: InitializationMode, + pub subagent_source: Option, + pub parent_thread_id: Option, +} + #[derive(Clone, Copy, Debug, Serialize)] #[serde(rename_all = "snake_case")] pub enum InitializationMode { @@ -80,7 +97,7 @@ pub struct AppInvocation { } pub enum AnalyticsInput { - ThreadInitialized(ThreadInitializedInput), + ThreadInitialized(CodexThreadInitializedInput), SkillInvoked(SkillInvokedInput), AppMentioned(AppMentionedInput), AppUsed(AppUsedInput), @@ -88,10 +105,6 @@ pub enum AnalyticsInput { PluginStateChanged(PluginStateChangedInput), } -pub struct ThreadInitializedInput { - pub thread_event: CodexThreadInitializedEvent, -} - pub struct SkillInvokedInput { pub tracking: TrackEventsContext, pub invocations: Vec, @@ -228,10 +241,8 @@ impl AnalyticsEventsClient { })); } - pub fn track_thread_initialized(&self, thread_event: CodexThreadInitializedEvent) { - self.record(AnalyticsInput::ThreadInitialized(ThreadInitializedInput { - thread_event, - })); + pub fn track_thread_initialized(&self, input: CodexThreadInitializedInput) { + self.record(AnalyticsInput::ThreadInitialized(input)); } pub fn track_app_mentioned(&self, tracking: TrackEventsContext, mentions: Vec) { @@ -446,19 +457,29 @@ impl AnalyticsReducer { fn ingest_thread_initialized( &mut self, - input: ThreadInitializedInput, + input: CodexThreadInitializedInput, out: &mut Vec, ) { + let event = CodexThreadInitializedEvent { + thread_id: input.thread_id, + model: input.model, + ephemeral: input.thread_context.ephemeral, + session_source: input.thread_context.session_source, + initialization_mode: input.thread_context.initialization_mode, + subagent_source: input.thread_context.subagent_source, + parent_thread_id: input.thread_context.parent_thread_id, + created_at: input.created_at, + }; self.threads.insert( - input.thread_event.thread_id.clone(), + event.thread_id.clone(), ThreadState { - _initialized_event: input.thread_event.clone(), + _initialized_event: event.clone(), }, ); out.push(TrackEventRequest::ThreadInitialized( CodexThreadInitializedEventRequest { event_type: "codex_thread_initialized", - event_params: codex_thread_initialized_event_params(input.thread_event), + event_params: codex_thread_initialized_event_params(event), }, )); } diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index 63455c0616..4ce20a0079 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -6,7 +6,9 @@ pub use analytics_client::AnalyticsReducer; pub use analytics_client::AppInvocation; pub use analytics_client::AppMentionedInput; pub use analytics_client::AppUsedInput; +pub use analytics_client::CodexThreadContext; pub use analytics_client::CodexThreadInitializedEvent; +pub use analytics_client::CodexThreadInitializedInput; pub use analytics_client::InitializationMode; pub use analytics_client::InvocationType; pub use analytics_client::PluginState; @@ -14,6 +16,5 @@ pub use analytics_client::PluginStateChangedInput; pub use analytics_client::PluginUsedInput; pub use analytics_client::SkillInvocation; pub use analytics_client::SkillInvokedInput; -pub use analytics_client::ThreadInitializedInput; pub use analytics_client::TrackEventsContext; pub use analytics_client::build_track_events_context; diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index e3bf4049ff..d333ab7cc9 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -350,7 +350,8 @@ use crate::util::backoff; use crate::windows_sandbox::WindowsSandboxLevelExt; use codex_analytics::AnalyticsEventsClient; use codex_analytics::AppInvocation; -use codex_analytics::CodexThreadInitializedEvent; +use codex_analytics::CodexThreadContext; +use codex_analytics::CodexThreadInitializedInput; use codex_analytics::InitializationMode; use codex_analytics::InvocationType; use codex_analytics::build_track_events_context; @@ -670,23 +671,25 @@ impl Codex { session .services .analytics_events_client - .track_thread_initialized(CodexThreadInitializedEvent { + .track_thread_initialized(CodexThreadInitializedInput { thread_id: thread_id.to_string(), model: thread_initialized_configuration .collaboration_mode .model() .to_string(), - ephemeral: thread_initialized_configuration - .original_config_do_not_use - .ephemeral, - initialization_mode, - subagent_source: session_source_subagent_source(&thread_session_source), - parent_thread_id: session_source_parent_thread_id(&thread_session_source), - session_source: thread_session_source, created_at: SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_secs(), + thread_context: CodexThreadContext { + ephemeral: thread_initialized_configuration + .original_config_do_not_use + .ephemeral, + session_source: thread_session_source, + initialization_mode, + subagent_source: session_source_subagent_source(&thread_session_source), + parent_thread_id: session_source_parent_thread_id(&thread_session_source), + }, }); // This task will run until Op::Shutdown is received. From 2737227745c2eb646b5e5e93030c84272fe0bca7 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:32:12 -0700 Subject: [PATCH 07/15] add product client id --- codex-rs/analytics/src/analytics_client.rs | 29 ++++++++++++++++++---- codex-rs/core/src/codex.rs | 1 + 2 files changed, 25 insertions(+), 5 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index f42059a622..43675ca0ab 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -42,6 +42,7 @@ pub struct CodexThreadInitializedEvent { pub struct CodexThreadInitializedInput { pub thread_id: String, pub model: String, + pub product_client_id: String, pub created_at: u64, pub thread_context: CodexThreadContext, } @@ -460,6 +461,7 @@ impl AnalyticsReducer { input: CodexThreadInitializedInput, out: &mut Vec, ) { + let product_client_id = input.product_client_id.clone(); let event = CodexThreadInitializedEvent { thread_id: input.thread_id, model: input.model, @@ -477,10 +479,7 @@ impl AnalyticsReducer { }, ); out.push(TrackEventRequest::ThreadInitialized( - CodexThreadInitializedEventRequest { - event_type: "codex_thread_initialized", - event_params: codex_thread_initialized_event_params(event), - }, + codex_thread_initialized_event_request(product_client_id, event), )); } @@ -602,10 +601,30 @@ fn codex_app_metadata(tracking: &TrackEventsContext, app: AppInvocation) -> Code fn codex_thread_initialized_event_params( thread_event: CodexThreadInitializedEvent, +) -> CodexThreadInitializedEventParams { + codex_thread_initialized_event_params_with_product_client_id(originator().value, thread_event) +} + +fn codex_thread_initialized_event_request( + product_client_id: String, + thread_event: CodexThreadInitializedEvent, +) -> CodexThreadInitializedEventRequest { + CodexThreadInitializedEventRequest { + event_type: "codex_thread_initialized", + event_params: codex_thread_initialized_event_params_with_product_client_id( + product_client_id, + thread_event, + ), + } +} + +fn codex_thread_initialized_event_params_with_product_client_id( + product_client_id: String, + thread_event: CodexThreadInitializedEvent, ) -> CodexThreadInitializedEventParams { CodexThreadInitializedEventParams { thread_id: thread_event.thread_id, - product_client_id: originator().value, + product_client_id, model: thread_event.model, ephemeral: thread_event.ephemeral, session_source: session_source_name(&thread_event.session_source), diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index d333ab7cc9..68b6a296e6 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -677,6 +677,7 @@ impl Codex { .collaboration_mode .model() .to_string(), + product_client_id: crate::default_client::originator().value, created_at: SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() From e9202dea95b940f94f6872bf8201a470d5f32a07 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:05:08 -0700 Subject: [PATCH 08/15] establish reducer structure --- codex-rs/analytics/src/analytics_client.rs | 752 ++++++++------------- codex-rs/analytics/src/lib.rs | 9 + 2 files changed, 289 insertions(+), 472 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index d47f9950f4..39436b3d4c 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -4,12 +4,21 @@ use codex_login::AuthManager; use codex_login::default_client::create_client; use codex_login::default_client::originator; use codex_plugin::PluginTelemetryMetadata; +use codex_protocol::config_types::ApprovalsReviewer; +use codex_protocol::config_types::ModeKind; +use codex_protocol::config_types::Personality; +use codex_protocol::config_types::ReasoningSummary; +use codex_protocol::config_types::ServiceTier; +use codex_protocol::openai_models::ReasoningEffort; +use codex_protocol::protocol::AskForApproval; +use codex_protocol::protocol::SandboxPolicy; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SkillScope; use codex_protocol::protocol::SubAgentSource; use serde::Serialize; use sha1::Digest; use sha1::Sha1; +use std::collections::HashMap; use std::collections::HashSet; use std::path::Path; use std::path::PathBuf; @@ -40,6 +49,8 @@ pub struct CodexTurnEvent { pub num_input_images: usize, pub is_first_turn: bool, } + +#[derive(Clone)] pub struct CodexThreadInitializedEvent { pub thread_id: String, pub model: String, @@ -92,9 +103,70 @@ pub struct AppInvocation { pub invocation_type: Option, } +pub enum AnalyticsInput { + ThreadInitialized(ThreadInitializedInput), + TurnEvent(TurnEventInput), + SkillInvoked(SkillInvokedInput), + AppMentioned(AppMentionedInput), + AppUsed(AppUsedInput), + PluginUsed(PluginUsedInput), + PluginStateChanged(PluginStateChangedInput), +} + +pub struct ThreadInitializedInput { + pub thread_event: CodexThreadInitializedEvent, +} + +pub struct TurnEventInput { + pub tracking: TrackEventsContext, + pub turn_event: CodexTurnEvent, +} + +pub struct SkillInvokedInput { + pub tracking: TrackEventsContext, + pub invocations: Vec, +} + +pub struct AppMentionedInput { + pub tracking: TrackEventsContext, + pub mentions: Vec, +} + +pub struct AppUsedInput { + pub tracking: TrackEventsContext, + pub app: AppInvocation, +} + +pub struct PluginUsedInput { + pub tracking: TrackEventsContext, + pub plugin: PluginTelemetryMetadata, +} + +pub struct PluginStateChangedInput { + pub plugin: PluginTelemetryMetadata, + pub state: PluginState, +} + +#[derive(Clone, Copy)] +pub enum PluginState { + Installed, + Uninstalled, + Enabled, + Disabled, +} + +#[derive(Default)] +pub struct AnalyticsReducer { + threads: HashMap, +} + +struct ThreadState { + _initialized_event: CodexThreadInitializedEvent, +} + #[derive(Clone)] pub(crate) struct AnalyticsEventsQueue { - sender: mpsc::Sender, + sender: mpsc::Sender, app_used_emitted_keys: Arc>>, plugin_used_emitted_keys: Arc>>, } @@ -109,39 +181,11 @@ impl AnalyticsEventsQueue { pub(crate) fn new(auth_manager: Arc, base_url: String) -> Self { let (sender, mut receiver) = mpsc::channel(ANALYTICS_EVENTS_QUEUE_SIZE); tokio::spawn(async move { + let mut reducer = AnalyticsReducer::default(); while let Some(job) = receiver.recv().await { - match job { - TrackEventsJob::SkillInvocations(job) => { - send_track_skill_invocations(&auth_manager, &base_url, job).await; - } - TrackEventsJob::ThreadInitialized(job) => { - send_track_thread_initialized(&auth_manager, &base_url, job).await; - } - TrackEventsJob::AppMentioned(job) => { - send_track_app_mentioned(&auth_manager, &base_url, job).await; - } - TrackEventsJob::AppUsed(job) => { - send_track_app_used(&auth_manager, &base_url, job).await; - } - TrackEventsJob::TurnEvent(job) => { - send_track_turn_event(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginUsed(job) => { - send_track_plugin_used(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginInstalled(job) => { - send_track_plugin_installed(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginUninstalled(job) => { - send_track_plugin_uninstalled(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginEnabled(job) => { - send_track_plugin_enabled(&auth_manager, &base_url, job).await; - } - TrackEventsJob::PluginDisabled(job) => { - send_track_plugin_disabled(&auth_manager, &base_url, job).await; - } - } + let mut events = Vec::new(); + reducer.ingest(job, &mut events).await; + send_track_events(&auth_manager, &base_url, events).await; } }); Self { @@ -151,8 +195,8 @@ impl AnalyticsEventsQueue { } } - fn try_send(&self, job: TrackEventsJob) { - if self.sender.try_send(job).is_err() { + fn try_send(&self, input: AnalyticsInput) { + if self.sender.try_send(input).is_err() { //TODO: add a metric for this tracing::warn!("dropping analytics events: queue is full"); } @@ -205,140 +249,97 @@ impl AnalyticsEventsClient { tracking: TrackEventsContext, invocations: Vec, ) { - track_skill_invocations( - &self.queue, - self.analytics_enabled, - Some(tracking), + if invocations.is_empty() { + return; + } + self.record(AnalyticsInput::SkillInvoked(SkillInvokedInput { + tracking, invocations, - ); + })); } pub fn track_thread_initialized(&self, thread_event: CodexThreadInitializedEvent) { - track_thread_initialized(&self.queue, self.analytics_enabled, thread_event); + self.record(AnalyticsInput::ThreadInitialized(ThreadInitializedInput { + thread_event, + })); } pub fn track_app_mentioned(&self, tracking: TrackEventsContext, mentions: Vec) { - track_app_mentioned( - &self.queue, - self.analytics_enabled, - Some(tracking), + if mentions.is_empty() { + return; + } + self.record(AnalyticsInput::AppMentioned(AppMentionedInput { + tracking, mentions, - ); + })); } pub fn track_app_used(&self, tracking: TrackEventsContext, app: AppInvocation) { - track_app_used(&self.queue, self.analytics_enabled, Some(tracking), app); + if !self.queue.should_enqueue_app_used(&tracking, &app) { + return; + } + self.record(AnalyticsInput::AppUsed(AppUsedInput { tracking, app })); } pub fn track_plugin_used(&self, tracking: TrackEventsContext, plugin: PluginTelemetryMetadata) { - track_plugin_used(&self.queue, self.analytics_enabled, Some(tracking), plugin); + if !self.queue.should_enqueue_plugin_used(&tracking, &plugin) { + return; + } + self.record(AnalyticsInput::PluginUsed(PluginUsedInput { + tracking, + plugin, + })); } pub fn track_turn_event(&self, tracking: TrackEventsContext, turn_event: CodexTurnEvent) { - track_turn_event( - &self.queue, - self.analytics_enabled, - Some(tracking), + self.record(AnalyticsInput::TurnEvent(TurnEventInput { + tracking, turn_event, - ); + })); } pub fn track_plugin_installed(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Installed, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Installed, + }, + )); } pub fn track_plugin_uninstalled(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Uninstalled, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Uninstalled, + }, + )); } pub fn track_plugin_enabled(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Enabled, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Enabled, + }, + )); } pub fn track_plugin_disabled(&self, plugin: PluginTelemetryMetadata) { - track_plugin_management( - &self.queue, - self.analytics_enabled, - PluginManagementEventType::Disabled, - plugin, - ); + self.record(AnalyticsInput::PluginStateChanged( + PluginStateChangedInput { + plugin, + state: PluginState::Disabled, + }, + )); } -} -enum TrackEventsJob { - SkillInvocations(TrackSkillInvocationsJob), - ThreadInitialized(TrackThreadInitializedJob), - AppMentioned(TrackAppMentionedJob), - AppUsed(TrackAppUsedJob), - TurnEvent(TrackTurnEventJob), - PluginUsed(TrackPluginUsedJob), - PluginInstalled(TrackPluginManagementJob), - PluginUninstalled(TrackPluginManagementJob), - PluginEnabled(TrackPluginManagementJob), - PluginDisabled(TrackPluginManagementJob), -} - -struct TrackSkillInvocationsJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - invocations: Vec, -} - -struct TrackThreadInitializedJob { - analytics_enabled: Option, - thread_event: CodexThreadInitializedEvent, -} - -struct TrackAppMentionedJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - mentions: Vec, -} - -struct TrackAppUsedJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - app: AppInvocation, -} - -struct TrackTurnEventJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - turn_event: CodexTurnEvent, -} - -struct TrackPluginUsedJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - plugin: PluginTelemetryMetadata, -} - -struct TrackPluginManagementJob { - analytics_enabled: Option, - plugin: PluginTelemetryMetadata, -} - -#[derive(Clone, Copy)] -enum PluginManagementEventType { - Installed, - Uninstalled, - Enabled, - Disabled, + pub fn record(&self, input: AnalyticsInput) { + if self.analytics_enabled == Some(false) { + return; + } + self.queue.try_send(input); + } } const ANALYTICS_EVENTS_QUEUE_SIZE: usize = 256; @@ -483,354 +484,165 @@ struct CodexPluginUsedEventRequest { event_params: CodexPluginUsedMetadata, } -pub(crate) fn track_skill_invocations( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - invocations: Vec, -) { - if analytics_enabled == Some(false) { - return; +impl AnalyticsReducer { + async fn ingest(&mut self, input: AnalyticsInput, out: &mut Vec) { + match input { + AnalyticsInput::ThreadInitialized(input) => { + self.ingest_thread_initialized(input, out); + } + AnalyticsInput::TurnEvent(input) => { + self.ingest_turn_event(input, out); + } + AnalyticsInput::SkillInvoked(input) => { + self.ingest_skill_invoked(input, out).await; + } + AnalyticsInput::AppMentioned(input) => { + self.ingest_app_mentioned(input, out); + } + AnalyticsInput::AppUsed(input) => { + self.ingest_app_used(input, out); + } + AnalyticsInput::PluginUsed(input) => { + self.ingest_plugin_used(input, out); + } + AnalyticsInput::PluginStateChanged(input) => { + self.ingest_plugin_state_changed(input, out); + } + } } - let Some(tracking) = tracking else { - return; - }; - if invocations.is_empty() { - return; - } - let job = TrackEventsJob::SkillInvocations(TrackSkillInvocationsJob { - analytics_enabled, - tracking, - invocations, - }); - queue.try_send(job); -} -pub(crate) fn track_thread_initialized( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - thread_event: CodexThreadInitializedEvent, -) { - if analytics_enabled == Some(false) { - return; - } - let job = TrackEventsJob::ThreadInitialized(TrackThreadInitializedJob { - analytics_enabled, - thread_event, - }); - queue.try_send(job); -} - -pub(crate) fn track_app_mentioned( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - mentions: Vec, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - if mentions.is_empty() { - return; - } - let job = TrackEventsJob::AppMentioned(TrackAppMentionedJob { - analytics_enabled, - tracking, - mentions, - }); - queue.try_send(job); -} - -pub(crate) fn track_app_used( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - app: AppInvocation, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - if !queue.should_enqueue_app_used(&tracking, &app) { - return; - } - let job = TrackEventsJob::AppUsed(TrackAppUsedJob { - analytics_enabled, - tracking, - app, - }); - queue.try_send(job); -} - -pub(crate) fn track_turn_event( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - turn_event: CodexTurnEvent, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - let job = TrackEventsJob::TurnEvent(TrackTurnEventJob { - analytics_enabled, - tracking, - turn_event, - }); - queue.try_send(job); -} - -pub(crate) fn track_plugin_used( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - plugin: PluginTelemetryMetadata, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - if !queue.should_enqueue_plugin_used(&tracking, &plugin) { - return; - } - let job = TrackEventsJob::PluginUsed(TrackPluginUsedJob { - analytics_enabled, - tracking, - plugin, - }); - queue.try_send(job); -} - -fn track_plugin_management( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - event_type: PluginManagementEventType, - plugin: PluginTelemetryMetadata, -) { - if analytics_enabled == Some(false) { - return; - } - let job = TrackPluginManagementJob { - analytics_enabled, - plugin, - }; - let job = match event_type { - PluginManagementEventType::Installed => TrackEventsJob::PluginInstalled(job), - PluginManagementEventType::Uninstalled => TrackEventsJob::PluginUninstalled(job), - PluginManagementEventType::Enabled => TrackEventsJob::PluginEnabled(job), - PluginManagementEventType::Disabled => TrackEventsJob::PluginDisabled(job), - }; - queue.try_send(job); -} - -async fn send_track_skill_invocations( - auth_manager: &AuthManager, - base_url: &str, - job: TrackSkillInvocationsJob, -) { - let TrackSkillInvocationsJob { - analytics_enabled, - tracking, - invocations, - } = job; - let mut events = Vec::with_capacity(invocations.len()); - for invocation in invocations { - let skill_scope = match invocation.skill_scope { - SkillScope::User => "user", - SkillScope::Repo => "repo", - SkillScope::System => "system", - SkillScope::Admin => "admin", - }; - let repo_root = get_git_repo_root(invocation.skill_path.as_path()); - let repo_url = if let Some(root) = repo_root.as_ref() { - collect_git_info(root) - .await - .and_then(|info| info.repository_url) - } else { - None - }; - let skill_id = skill_id_for_local_skill( - repo_url.as_deref(), - repo_root.as_deref(), - invocation.skill_path.as_path(), - invocation.skill_name.as_str(), + fn ingest_thread_initialized( + &mut self, + input: ThreadInitializedInput, + out: &mut Vec, + ) { + self.threads.insert( + input.thread_event.thread_id.clone(), + ThreadState { + _initialized_event: input.thread_event.clone(), + }, ); - events.push(TrackEventRequest::SkillInvocation( - SkillInvocationEventRequest { - event_type: "skill_invocation", - skill_id, - skill_name: invocation.skill_name.clone(), - event_params: SkillInvocationEventParams { - thread_id: Some(tracking.thread_id.clone()), - invoke_type: Some(invocation.invocation_type), - model_slug: Some(tracking.model_slug.clone()), - product_client_id: Some(originator().value), - repo_url, - skill_scope: Some(skill_scope.to_string()), - }, + out.push(TrackEventRequest::ThreadInitialized( + CodexThreadInitializedEventRequest { + event_type: "codex_thread_initialized", + event_params: codex_thread_initialized_event_params(input.thread_event), }, )); } - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} + fn ingest_turn_event(&mut self, input: TurnEventInput, out: &mut Vec) { + let TurnEventInput { + tracking, + turn_event, + } = input; + out.push(TrackEventRequest::TurnEvent(CodexTurnEventRequest { + event_type: "codex_turn_event", + event_params: codex_turn_event_params(&tracking, turn_event), + })); + } -async fn send_track_thread_initialized( - auth_manager: &AuthManager, - base_url: &str, - job: TrackThreadInitializedJob, -) { - let TrackThreadInitializedJob { - analytics_enabled, - thread_event, - } = job; - let events = vec![TrackEventRequest::ThreadInitialized( - CodexThreadInitializedEventRequest { - event_type: "codex_thread_initialized", - event_params: codex_thread_initialized_event_params(thread_event), - }, - )]; + async fn ingest_skill_invoked( + &mut self, + input: SkillInvokedInput, + out: &mut Vec, + ) { + let SkillInvokedInput { + tracking, + invocations, + } = input; + for invocation in invocations { + let skill_scope = match invocation.skill_scope { + SkillScope::User => "user", + SkillScope::Repo => "repo", + SkillScope::System => "system", + SkillScope::Admin => "admin", + }; + let repo_root = get_git_repo_root(invocation.skill_path.as_path()); + let repo_url = if let Some(root) = repo_root.as_ref() { + collect_git_info(root) + .await + .and_then(|info| info.repository_url) + } else { + None + }; + let skill_id = skill_id_for_local_skill( + repo_url.as_deref(), + repo_root.as_deref(), + invocation.skill_path.as_path(), + invocation.skill_name.as_str(), + ); + out.push(TrackEventRequest::SkillInvocation( + SkillInvocationEventRequest { + event_type: "skill_invocation", + skill_id, + skill_name: invocation.skill_name.clone(), + event_params: SkillInvocationEventParams { + thread_id: Some(tracking.thread_id.clone()), + invoke_type: Some(invocation.invocation_type), + model_slug: Some(tracking.model_slug.clone()), + product_client_id: Some(originator().value), + repo_url, + skill_scope: Some(skill_scope.to_string()), + }, + }, + )); + } + } - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} - -async fn send_track_app_mentioned( - auth_manager: &AuthManager, - base_url: &str, - job: TrackAppMentionedJob, -) { - let TrackAppMentionedJob { - analytics_enabled, - tracking, - mentions, - } = job; - let events = mentions - .into_iter() - .map(|mention| { + fn ingest_app_mentioned(&mut self, input: AppMentionedInput, out: &mut Vec) { + let AppMentionedInput { tracking, mentions } = input; + out.extend(mentions.into_iter().map(|mention| { let event_params = codex_app_metadata(&tracking, mention); TrackEventRequest::AppMentioned(CodexAppMentionedEventRequest { event_type: "codex_app_mentioned", event_params, }) - }) - .collect::>(); + })); + } - send_track_events(auth_manager, analytics_enabled, base_url, events).await; + fn ingest_app_used(&mut self, input: AppUsedInput, out: &mut Vec) { + let AppUsedInput { tracking, app } = input; + let event_params = codex_app_metadata(&tracking, app); + out.push(TrackEventRequest::AppUsed(CodexAppUsedEventRequest { + event_type: "codex_app_used", + event_params, + })); + } + + fn ingest_plugin_used(&mut self, input: PluginUsedInput, out: &mut Vec) { + let PluginUsedInput { tracking, plugin } = input; + out.push(TrackEventRequest::PluginUsed(CodexPluginUsedEventRequest { + event_type: "codex_plugin_used", + event_params: codex_plugin_used_metadata(&tracking, plugin), + })); + } + + fn ingest_plugin_state_changed( + &mut self, + input: PluginStateChangedInput, + out: &mut Vec, + ) { + let PluginStateChangedInput { plugin, state } = input; + let event = CodexPluginEventRequest { + event_type: plugin_state_event_type(state), + event_params: codex_plugin_metadata(plugin), + }; + out.push(match state { + PluginState::Installed => TrackEventRequest::PluginInstalled(event), + PluginState::Uninstalled => TrackEventRequest::PluginUninstalled(event), + PluginState::Enabled => TrackEventRequest::PluginEnabled(event), + PluginState::Disabled => TrackEventRequest::PluginDisabled(event), + }); + } } -async fn send_track_app_used(auth_manager: &AuthManager, base_url: &str, job: TrackAppUsedJob) { - let TrackAppUsedJob { - analytics_enabled, - tracking, - app, - } = job; - let event_params = codex_app_metadata(&tracking, app); - let events = vec![TrackEventRequest::AppUsed(CodexAppUsedEventRequest { - event_type: "codex_app_used", - event_params, - })]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} - -async fn send_track_turn_event(auth_manager: &AuthManager, base_url: &str, job: TrackTurnEventJob) { - let TrackTurnEventJob { - analytics_enabled, - tracking, - turn_event, - } = job; - let events = vec![TrackEventRequest::TurnEvent(CodexTurnEventRequest { - event_type: "codex_turn_event", - event_params: codex_turn_event_params(&tracking, turn_event), - })]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} - -async fn send_track_plugin_used( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginUsedJob, -) { - let TrackPluginUsedJob { - analytics_enabled, - tracking, - plugin, - } = job; - let events = vec![TrackEventRequest::PluginUsed(CodexPluginUsedEventRequest { - event_type: "codex_plugin_used", - event_params: codex_plugin_used_metadata(&tracking, plugin), - })]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; -} - -async fn send_track_plugin_installed( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_installed").await; -} - -async fn send_track_plugin_uninstalled( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_uninstalled") - .await; -} - -async fn send_track_plugin_enabled( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_enabled").await; -} - -async fn send_track_plugin_disabled( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, -) { - send_track_plugin_management_event(auth_manager, base_url, job, "codex_plugin_disabled").await; -} - -async fn send_track_plugin_management_event( - auth_manager: &AuthManager, - base_url: &str, - job: TrackPluginManagementJob, - event_type: &'static str, -) { - let TrackPluginManagementJob { - analytics_enabled, - plugin, - } = job; - let event_params = codex_plugin_metadata(plugin); - let event = CodexPluginEventRequest { - event_type, - event_params, - }; - let events = vec![match event_type { - "codex_plugin_installed" => TrackEventRequest::PluginInstalled(event), - "codex_plugin_uninstalled" => TrackEventRequest::PluginUninstalled(event), - "codex_plugin_enabled" => TrackEventRequest::PluginEnabled(event), - "codex_plugin_disabled" => TrackEventRequest::PluginDisabled(event), - _ => unreachable!("unknown plugin management event type"), - }]; - - send_track_events(auth_manager, analytics_enabled, base_url, events).await; +fn plugin_state_event_type(state: PluginState) -> &'static str { + match state { + PluginState::Installed => "codex_plugin_installed", + PluginState::Uninstalled => "codex_plugin_uninstalled", + PluginState::Enabled => "codex_plugin_enabled", + PluginState::Disabled => "codex_plugin_disabled", + } } fn codex_app_metadata(tracking: &TrackEventsContext, app: AppInvocation) -> CodexAppMetadata { @@ -973,13 +785,9 @@ fn subagent_source_name(subagent_source: SubAgentSource) -> String { async fn send_track_events( auth_manager: &AuthManager, - analytics_enabled: Option, base_url: &str, events: Vec, ) { - if analytics_enabled == Some(false) { - return; - } if events.is_empty() { return; } diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index b76546753a..57029a890a 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -1,11 +1,20 @@ mod analytics_client; pub use analytics_client::AnalyticsEventsClient; +pub use analytics_client::AnalyticsInput; +pub use analytics_client::AnalyticsReducer; pub use analytics_client::AppInvocation; +pub use analytics_client::AppMentionedInput; +pub use analytics_client::AppUsedInput; pub use analytics_client::CodexThreadInitializedEvent; pub use analytics_client::CodexTurnEvent; pub use analytics_client::InitializationMode; pub use analytics_client::InvocationType; +pub use analytics_client::PluginState; +pub use analytics_client::PluginStateChangedInput; +pub use analytics_client::PluginUsedInput; pub use analytics_client::SkillInvocation; +pub use analytics_client::SkillInvokedInput; +pub use analytics_client::ThreadInitializedInput; pub use analytics_client::TrackEventsContext; pub use analytics_client::build_track_events_context; From 6aa5461d30d9cc5ba2480e51e3d54c718d8a05d7 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:36:23 -0700 Subject: [PATCH 09/15] int --- codex-rs/analytics/src/analytics_client.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index 43675ca0ab..30d679ef4a 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -599,6 +599,7 @@ fn codex_app_metadata(tracking: &TrackEventsContext, app: AppInvocation) -> Code } } +#[cfg(test)] fn codex_thread_initialized_event_params( thread_event: CodexThreadInitializedEvent, ) -> CodexThreadInitializedEventParams { From fdd3eaa93d9e2d3086ada4e53a4837fd90585461 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:43:36 -0700 Subject: [PATCH 10/15] diff cleanup --- codex-rs/analytics/src/analytics_client.rs | 1 - codex-rs/analytics/src/lib.rs | 10 +--------- codex-rs/core/src/codex.rs | 4 ---- 3 files changed, 1 insertion(+), 14 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index b9236771b4..ab5f5cb447 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -542,7 +542,6 @@ impl AnalyticsReducer { event.thread_id.clone(), ThreadState { _initialized_event: event.clone(), ->>>>>>> origin/rhan/oob-thread-metadata }, ); out.push(TrackEventRequest::ThreadInitialized( diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index 0c19744c13..2a486cdc67 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -6,14 +6,10 @@ pub use analytics_client::AnalyticsReducer; pub use analytics_client::AppInvocation; pub use analytics_client::AppMentionedInput; pub use analytics_client::AppUsedInput; -<<<<<<< HEAD -pub use analytics_client::CodexThreadInitializedEvent; -pub use analytics_client::CodexTurnEvent; -======= pub use analytics_client::CodexThreadContext; pub use analytics_client::CodexThreadInitializedEvent; pub use analytics_client::CodexThreadInitializedInput; ->>>>>>> origin/rhan/oob-thread-metadata +pub use analytics_client::CodexTurnEvent; pub use analytics_client::InitializationMode; pub use analytics_client::InvocationType; pub use analytics_client::PluginState; @@ -21,9 +17,5 @@ pub use analytics_client::PluginStateChangedInput; pub use analytics_client::PluginUsedInput; pub use analytics_client::SkillInvocation; pub use analytics_client::SkillInvokedInput; -<<<<<<< HEAD -pub use analytics_client::ThreadInitializedInput; -======= ->>>>>>> origin/rhan/oob-thread-metadata pub use analytics_client::TrackEventsContext; pub use analytics_client::build_track_events_context; diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 47ec1dda79..0b35b91f67 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -351,13 +351,9 @@ use crate::util::backoff; use crate::windows_sandbox::WindowsSandboxLevelExt; use codex_analytics::AnalyticsEventsClient; use codex_analytics::AppInvocation; -<<<<<<< HEAD -use codex_analytics::CodexThreadInitializedEvent; use codex_analytics::CodexTurnEvent; -======= use codex_analytics::CodexThreadContext; use codex_analytics::CodexThreadInitializedInput; ->>>>>>> origin/rhan/oob-thread-metadata use codex_analytics::InitializationMode; use codex_analytics::InvocationType; use codex_analytics::build_track_events_context; From a779c8f4d3f8c499a76fb0cbe2ff48b9fba38d3f Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:44:17 -0700 Subject: [PATCH 11/15] fmt --- codex-rs/core/src/codex.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 0b35b91f67..846fefd32d 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -351,9 +351,9 @@ use crate::util::backoff; use crate::windows_sandbox::WindowsSandboxLevelExt; use codex_analytics::AnalyticsEventsClient; use codex_analytics::AppInvocation; -use codex_analytics::CodexTurnEvent; use codex_analytics::CodexThreadContext; use codex_analytics::CodexThreadInitializedInput; +use codex_analytics::CodexTurnEvent; use codex_analytics::InitializationMode; use codex_analytics::InvocationType; use codex_analytics::build_track_events_context; From 80367abdf10e3d97e2ece5eef0f4263158ab88fc Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 15:59:23 -0700 Subject: [PATCH 12/15] imports --- codex-rs/analytics/src/analytics_client_tests.rs | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index e26a027d6d..0d43b2dc80 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -23,6 +23,14 @@ use codex_plugin::AppConnectorId; use codex_plugin::PluginCapabilitySummary; use codex_plugin::PluginId; use codex_plugin::PluginTelemetryMetadata; +use codex_protocol::config_types::ApprovalsReviewer; +use codex_protocol::config_types::ModeKind; +use codex_protocol::config_types::Personality; +use codex_protocol::config_types::ReasoningSummary; +use codex_protocol::config_types::ServiceTier; +use codex_protocol::openai_models::ReasoningEffort; +use codex_protocol::protocol::AskForApproval; +use codex_protocol::protocol::SandboxPolicy; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; use pretty_assertions::assert_eq; From 6f6f5064b248e00e8255af893be6a6fd03a0e2c1 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 16:13:10 -0700 Subject: [PATCH 13/15] remove now dead code --- codex-rs/analytics/src/analytics_client.rs | 6 ---- .../analytics/src/analytics_client_tests.rs | 30 +++++++++---------- 2 files changed, 15 insertions(+), 21 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index ab5f5cb447..acd01d6033 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -731,12 +731,6 @@ fn personality_mode(personality: Option) -> Option { Some(personality) => Some(personality.to_string()), } } -fn codex_thread_initialized_event_params( - thread_event: CodexThreadInitializedEvent, -) -> CodexThreadInitializedEventParams { - codex_thread_initialized_event_params_with_product_client_id(originator().value, thread_event) -} - fn codex_thread_initialized_event_request( product_client_id: String, thread_event: CodexThreadInitializedEvent, diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index 0d43b2dc80..f909bc43fe 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -261,9 +261,9 @@ fn turn_event_serializes_expected_shape() { #[test] fn thread_initialized_event_serializes_expected_shape() { - let event = TrackEventRequest::ThreadInitialized(CodexThreadInitializedEventRequest { - event_type: "codex_thread_initialized", - event_params: codex_thread_initialized_event_params(CodexThreadInitializedEvent { + let event = TrackEventRequest::ThreadInitialized(codex_thread_initialized_event_request( + originator().value, + CodexThreadInitializedEvent { thread_id: "thread-0".to_string(), model: "gpt-5".to_string(), ephemeral: true, @@ -272,8 +272,8 @@ fn thread_initialized_event_serializes_expected_shape() { subagent_source: None, parent_thread_id: None, created_at: 1_716_000_000, - }), - }); + }, + )); let payload = serde_json::to_value(&event).expect("serialize thread initialized event"); @@ -298,9 +298,9 @@ fn thread_initialized_event_serializes_expected_shape() { #[test] fn thread_initialized_event_serializes_subagent_source() { - let event = TrackEventRequest::ThreadInitialized(CodexThreadInitializedEventRequest { - event_type: "codex_thread_initialized", - event_params: codex_thread_initialized_event_params(CodexThreadInitializedEvent { + let event = TrackEventRequest::ThreadInitialized(codex_thread_initialized_event_request( + originator().value, + CodexThreadInitializedEvent { thread_id: "thread-1".to_string(), model: "gpt-5".to_string(), ephemeral: false, @@ -309,8 +309,8 @@ fn thread_initialized_event_serializes_subagent_source() { subagent_source: Some(SubAgentSource::Review), parent_thread_id: None, created_at: 1, - }), - }); + }, + )); let payload = serde_json::to_value(&event).expect("serialize subagent thread initialized event"); @@ -320,9 +320,9 @@ fn thread_initialized_event_serializes_subagent_source() { #[test] fn thread_initialized_event_omits_non_user_non_subagent_session_source() { - let event = TrackEventRequest::ThreadInitialized(CodexThreadInitializedEventRequest { - event_type: "codex_thread_initialized", - event_params: codex_thread_initialized_event_params(CodexThreadInitializedEvent { + let event = TrackEventRequest::ThreadInitialized(codex_thread_initialized_event_request( + originator().value, + CodexThreadInitializedEvent { thread_id: "thread-2".to_string(), model: "gpt-5".to_string(), ephemeral: false, @@ -331,8 +331,8 @@ fn thread_initialized_event_omits_non_user_non_subagent_session_source() { subagent_source: None, parent_thread_id: None, created_at: 1, - }), - }); + }, + )); let payload = serde_json::to_value(&event).expect("serialize mcp thread initialized event"); assert_eq!( From e04ceaa2eb89e0b48bd4d48b96d0543ae65a4cac Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 16:38:06 -0700 Subject: [PATCH 14/15] update thread event imports --- codex-rs/analytics/src/analytics_client_tests.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index f909bc43fe..5d0954ca9e 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -5,7 +5,6 @@ use super::CodexAppUsedEventRequest; use super::CodexPluginEventRequest; use super::CodexPluginUsedEventRequest; use super::CodexThreadInitializedEvent; -use super::CodexThreadInitializedEventRequest; use super::CodexTurnEvent; use super::CodexTurnEventRequest; use super::InitializationMode; @@ -15,7 +14,7 @@ use super::TrackEventsContext; use super::codex_app_metadata; use super::codex_plugin_metadata; use super::codex_plugin_used_metadata; -use super::codex_thread_initialized_event_params; +use super::codex_thread_initialized_event_request; use super::codex_turn_event_params; use super::normalize_path_for_skill_id; use codex_login::default_client::originator; From ff52cac307b22b9b7ca9205ad5a8479ffebb75b9 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Wed, 25 Mar 2026 16:56:37 -0700 Subject: [PATCH 15/15] test: wait for turn event --- codex-rs/core/tests/suite/items.rs | 32 +++++++++++++++--------------- 1 file changed, 16 insertions(+), 16 deletions(-) diff --git a/codex-rs/core/tests/suite/items.rs b/codex-rs/core/tests/suite/items.rs index 5ecdb879fa..a163c61871 100644 --- a/codex-rs/core/tests/suite/items.rs +++ b/codex-rs/core/tests/suite/items.rs @@ -185,30 +185,30 @@ async fn user_turn_tracks_turn_metadata_analytics() -> anyhow::Result<()> { wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; let deadline = Instant::now() + Duration::from_secs(10); - let analytics_request = loop { + let event = loop { let requests = server.received_requests().await.unwrap_or_default(); - if let Some(request) = requests + let turn_events = requests .into_iter() - .find(|request| request.url.path() == "/codex/analytics-events/events") - { - break request; + .filter(|request| request.url.path() == "/codex/analytics-events/events") + .filter_map(|request| serde_json::from_slice::(&request.body).ok()) + .flat_map(|payload| { + payload["events"] + .as_array() + .cloned() + .unwrap_or_default() + .into_iter() + }) + .filter(|event| event["event_type"] == "codex_turn_event") + .collect::>(); + if let Some(event) = turn_events.last().cloned() { + break event; } if Instant::now() >= deadline { - panic!("timed out waiting for turn analytics request"); + panic!("timed out waiting for turn analytics event"); } tokio::time::sleep(Duration::from_millis(50)).await; }; - let payload: Value = serde_json::from_slice(&analytics_request.body)?; - let event = payload["events"] - .as_array() - .and_then(|events| { - events - .iter() - .find(|event| event["event_type"] == "codex_turn_event") - }) - .expect("codex_turn_event should be present"); - let event_params = &event["event_params"]; assert_eq!(event_params["sandbox_policy"], "read_only");