diff --git a/codex-rs/analytics/src/analytics_client.rs b/codex-rs/analytics/src/analytics_client.rs index 48dadc7b93..8431ffba5f 100644 --- a/codex-rs/analytics/src/analytics_client.rs +++ b/codex-rs/analytics/src/analytics_client.rs @@ -18,6 +18,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; @@ -48,19 +49,11 @@ 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, - 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, @@ -72,6 +65,24 @@ pub struct CodexThreadInitializedEvent { #[derive(Clone, Copy)] pub struct CodexTurnSteerEvent; +#[derive(Clone)] +pub struct CodexThreadInitializedInput { + pub thread_id: String, + pub model: String, + pub product_client_id: 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 { @@ -113,9 +124,72 @@ pub struct AppInvocation { pub invocation_type: Option, } +pub enum AnalyticsInput { + ThreadInitialized(CodexThreadInitializedInput), + TurnEvent(TurnEventInput), + TurnSteer(TurnSteerInput), + SkillInvoked(SkillInvokedInput), + AppMentioned(AppMentionedInput), + AppUsed(AppUsedInput), + PluginUsed(PluginUsedInput), + PluginStateChanged(PluginStateChangedInput), +} + +pub struct TurnEventInput { + pub tracking: TrackEventsContext, + pub turn_event: CodexTurnEvent, +} + +pub struct TurnSteerInput { + pub tracking: TrackEventsContext, + pub turn_steer: CodexTurnSteerEvent, +} + +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>>, } @@ -130,42 +204,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::TurnSteer(job) => { - send_track_turn_steer(&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 { @@ -175,8 +218,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"); } @@ -229,156 +272,102 @@ 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); + pub fn track_thread_initialized(&self, input: CodexThreadInitializedInput) { + self.record(AnalyticsInput::ThreadInitialized(input)); } 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_turn_steer(&self, tracking: TrackEventsContext, turn_steer: CodexTurnSteerEvent) { - track_turn_steer( - &self.queue, - self.analytics_enabled, - Some(tracking), + self.record(AnalyticsInput::TurnSteer(TurnSteerInput { + tracking, turn_steer, - ); + })); } 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), - TurnSteer(TrackTurnSteerJob), - 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 TrackTurnSteerJob { - analytics_enabled: Option, - tracking: TrackEventsContext, - turn_steer: CodexTurnSteerEvent, -} - -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; @@ -429,16 +418,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, @@ -548,388 +527,186 @@ 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::TurnSteer(input) => { + self.ingest_turn_steer(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_turn_steer( - queue: &AnalyticsEventsQueue, - analytics_enabled: Option, - tracking: Option, - turn_steer: CodexTurnSteerEvent, -) { - if analytics_enabled == Some(false) { - return; - } - let Some(tracking) = tracking else { - return; - }; - let job = TrackEventsJob::TurnSteer(TrackTurnSteerJob { - analytics_enabled, - tracking, - turn_steer, - }); - 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", + fn ingest_thread_initialized( + &mut self, + 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, + 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, }; - 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(), - ); - 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()), - }, + self.threads.insert( + event.thread_id.clone(), + ThreadState { + _initialized_event: event.clone(), }, + ); + out.push(TrackEventRequest::ThreadInitialized( + codex_thread_initialized_event_request(product_client_id, 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), - }, - )]; + fn ingest_turn_steer(&mut self, input: TurnSteerInput, out: &mut Vec) { + let TurnSteerInput { + tracking, + turn_steer, + } = input; + out.push(TrackEventRequest::TurnSteer(CodexTurnSteerEventRequest { + event_type: "codex_turn_event", + event_params: codex_turn_steer_event_params(&tracking, turn_steer), + })); + } - 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| { + 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()), + }, + }, + )); + } + } + 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_turn_steer(auth_manager: &AuthManager, base_url: &str, job: TrackTurnSteerJob) { - let TrackTurnSteerJob { - analytics_enabled, - tracking, - turn_steer, - } = job; - let events = vec![TrackEventRequest::TurnSteer(CodexTurnSteerEventRequest { - event_type: "codex_turn_event", - event_params: codex_turn_steer_event_params(&tracking, turn_steer), - })]; - - 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 { @@ -1012,29 +789,27 @@ fn personality_mode(personality: Option) -> Option { Some(personality) => Some(personality.to_string()), } } +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( +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, - 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, @@ -1099,13 +874,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/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index 7c7a4f4f0a..c0d45c840c 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::CodexTurnSteerEvent; @@ -17,7 +16,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::codex_turn_steer_event_params; use super::normalize_path_for_skill_id; @@ -292,29 +291,19 @@ fn turn_steer_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(), - 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, 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"); @@ -326,16 +315,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", @@ -349,29 +328,19 @@ 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(), - 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, 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"); @@ -381,29 +350,19 @@ 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(), - 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, 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!( diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index c848470c04..cfbe7a258b 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -1,12 +1,22 @@ 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::CodexThreadContext; pub use analytics_client::CodexThreadInitializedEvent; +pub use analytics_client::CodexThreadInitializedInput; pub use analytics_client::CodexTurnEvent; pub use analytics_client::CodexTurnSteerEvent; 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::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 f221a7fdf9..06604357f5 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -351,7 +351,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::CodexTurnEvent; use codex_analytics::CodexTurnSteerEvent; use codex_analytics::InitializationMode; @@ -673,43 +674,26 @@ 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(), - 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, - 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, + product_client_id: crate::default_client::originator().value, 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. diff --git a/codex-rs/core/tests/suite/items.rs b/codex-rs/core/tests/suite/items.rs index 7f755d6174..a163c61871 100644 --- a/codex-rs/core/tests/suite/items.rs +++ b/codex-rs/core/tests/suite/items.rs @@ -34,7 +34,6 @@ use core_test_support::responses::ev_response_created; use core_test_support::responses::ev_web_search_call_added_partial; use core_test_support::responses::ev_web_search_call_done; use core_test_support::responses::mount_sse_once; -use core_test_support::responses::mount_sse_sequence; use core_test_support::responses::sse; use core_test_support::responses::start_mock_server; use core_test_support::skip_if_no_network; @@ -48,7 +47,6 @@ use std::path::Path; use std::path::PathBuf; use std::time::Duration; use std::time::Instant; -use tempfile::tempdir; fn image_generation_artifact_path(codex_home: &Path, session_id: &str, call_id: &str) -> PathBuf { fn sanitize(value: &str) -> String { @@ -74,36 +72,6 @@ fn image_generation_artifact_path(codex_home: &Path, session_id: &str, call_id: .join(format!("{}.png", sanitize(call_id))) } -async fn wait_for_analytics_event( - server: &wiremock::MockServer, - event_type: &str, - event_match: impl Fn(&Value) -> bool, -) -> Value { - let deadline = Instant::now() + Duration::from_secs(10); - loop { - let requests = server.received_requests().await.unwrap_or_default(); - if let Some(event) = requests - .into_iter() - .filter(|request| request.url.path() == "/codex/analytics-events/events") - .find_map(|request| { - let payload: Value = serde_json::from_slice(&request.body).ok()?; - payload["events"].as_array().and_then(|events| { - events - .iter() - .find(|event| event["event_type"] == event_type && event_match(event)) - .cloned() - }) - }) - { - return event; - } - if Instant::now() >= deadline { - panic!("timed out waiting for analytics event {event_type}"); - } - tokio::time::sleep(Duration::from_millis(50)).await; - } -} - #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn user_message_item_is_emitted() -> anyhow::Result<()> { skip_if_no_network!(Ok(())); @@ -194,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().to_path_buf(), + cwd: config.cwd.to_path_buf(), approval_policy: AskForApproval::Never, approvals_reviewer: None, sandbox_policy: SandboxPolicy::new_read_only_policy(), @@ -217,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"); @@ -272,129 +240,6 @@ async fn user_turn_tracks_turn_metadata_analytics() -> anyhow::Result<()> { Ok(()) } -#[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn user_turn_tracks_turn_steer_analytics() -> anyhow::Result<()> { - skip_if_no_network!(Ok(())); - - let server = start_mock_server().await; - let temp = tempdir()?; - let unblock_path = temp.path().join("unblock-steering"); - let command = format!( - "while [ ! -f \"{}\" ]; do sleep 0.01; done; echo done", - unblock_path.display() - ); - let call_id = "shell-steering-call"; - - mount_sse_sequence( - &server, - vec![ - sse(vec![ - ev_response_created("resp-1"), - core_test_support::responses::ev_function_call( - call_id, - "shell", - &serde_json::to_string(&serde_json::json!({ - "command": ["/bin/sh", "-c", command], - }))?, - ), - ev_completed("resp-1"), - ]), - sse(vec![ - ev_assistant_message("msg-2", "done"), - ev_completed("resp-2"), - ]), - ], - ) - .await; - - let chatgpt_base_url = server.uri(); - let test = test_codex() - .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) - .with_model("gpt-5") - .with_config(move |config| { - config.chatgpt_base_url = chatgpt_base_url; - }) - .build(&server) - .await?; - let codex = test.codex.clone(); - let turn_model = test.session_configured.model.clone(); - - codex - .submit(Op::UserTurn { - items: vec![UserInput::Text { - text: "start steering flow".into(), - text_elements: Vec::new(), - }], - final_output_json_schema: None, - cwd: test.cwd_path().to_path_buf(), - approval_policy: AskForApproval::Never, - approvals_reviewer: None, - sandbox_policy: SandboxPolicy::DangerFullAccess, - model: turn_model, - effort: None, - summary: None, - service_tier: None, - collaboration_mode: None, - personality: None, - }) - .await?; - - let turn_id = wait_for_event_match(&codex, |ev| match ev { - EventMsg::TurnStarted(event) => Some(event.turn_id.clone()), - _ => None, - }) - .await; - - wait_for_event_match(&codex, |ev| match ev { - EventMsg::ExecCommandBegin(event) if event.call_id == call_id => Some(()), - _ => None, - }) - .await; - - let steered_turn_id = codex - .steer_input( - vec![UserInput::Text { - text: "steering metadata check".into(), - text_elements: Vec::new(), - }], - Some(turn_id.as_str()), - ) - .await - .expect("steer should succeed on active turn"); - assert_eq!(steered_turn_id, turn_id); - - std::fs::write(&unblock_path, "go")?; - - wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; - - let event = wait_for_analytics_event(&server, "codex_turn_event", |event| { - event["event_params"].get("model").is_some() - && event["event_params"].get("sandbox_policy").is_none() - }) - .await; - let event_params = &event["event_params"]; - - assert_eq!( - event_params["product_client_id"], - serde_json::json!(codex_core::default_client::originator().value) - ); - assert_eq!(event_params["model"], "gpt-5"); - assert!(event_params["thread_id"].as_str().is_some()); - assert!(event_params["turn_id"].as_str().is_some()); - assert!(event_params.get("sandbox_policy").is_none()); - assert!(event_params.get("reasoning_effort").is_none()); - assert!(event_params.get("reasoning_summary").is_none()); - assert!(event_params.get("service_tier").is_none()); - assert!(event_params.get("approval_policy").is_none()); - assert!(event_params.get("approvals_reviewer").is_none()); - assert!(event_params.get("sandbox_network_access").is_none()); - assert!(event_params.get("collaboration_mode").is_none()); - assert!(event_params.get("personality").is_none()); - assert!(event_params.get("num_input_images").is_none()); - - Ok(()) -} - #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn resumed_thread_turn_tracks_is_not_first_turn_analytics() -> anyhow::Result<()> { skip_if_no_network!(Ok(())); @@ -480,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"); diff --git a/codex-rs/core/tests/suite/thread_metadata.rs b/codex-rs/core/tests/suite/thread_metadata.rs index a307d25059..3cf0e22ecd 100644 --- a/codex-rs/core/tests/suite/thread_metadata.rs +++ b/codex-rs/core/tests/suite/thread_metadata.rs @@ -1,13 +1,6 @@ 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; @@ -23,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) @@ -69,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"); @@ -101,108 +76,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(()) -}