| use crate::events::AppServerRpcTransport; |
| use crate::events::GuardianReviewAnalyticsResult; |
| use crate::events::GuardianReviewTrackContext; |
| use crate::events::TrackEventRequest; |
| use crate::events::TrackEventsRequest; |
| use crate::events::current_runtime_metadata; |
| use crate::facts::AnalyticsFact; |
| use crate::facts::AnalyticsJsonRpcError; |
| use crate::facts::AppInvocation; |
| use crate::facts::AppMentionedInput; |
| use crate::facts::AppUsedInput; |
| use crate::facts::ArtifactOperation; |
| use crate::facts::ArtifactOperationInput; |
| use crate::facts::CodexGoalEvent; |
| use crate::facts::CustomAnalyticsFact; |
| use crate::facts::ElicitationType; |
| use crate::facts::ExternalAgentConfigImportCompletedInput; |
| use crate::facts::ExternalAgentConfigImportFailureInput; |
| use crate::facts::HookRunFact; |
| use crate::facts::HookRunInput; |
| use crate::facts::ImagePreparationFact; |
| use crate::facts::McpToolCallElicitation; |
| use crate::facts::PluginInstallFailedInput; |
| use crate::facts::PluginInstallRequested; |
| use crate::facts::PluginInstallRequestedInput; |
| use crate::facts::PluginInstallSource; |
| use crate::facts::PluginMeasurementsInput; |
| use crate::facts::PluginState; |
| use crate::facts::PluginStateChangedInput; |
| use crate::facts::SkillInvocation; |
| use crate::facts::SkillInvokedInput; |
| use crate::facts::SubAgentThreadStartedInput; |
| use crate::facts::TrackEventsContext; |
| use crate::facts::TurnCodexErrorFact; |
| use crate::facts::TurnProfileFact; |
| use crate::facts::TurnResolvedConfigFact; |
| use crate::facts::TurnTokenUsageFact; |
| use crate::guardian_v2::GuardianV2Event; |
| use crate::now_unix_millis; |
| use crate::reducer::AnalyticsReducer; |
| use crate::reducer::MAX_PLUGIN_MEASUREMENTS_PER_BATCH; |
| use crate::reducer::tracked_tool_item_id; |
| use crate::reducer::valid_plugin_measurement_identifier; |
| use crate::reducer::valid_plugin_measurement_row; |
| use codex_app_server_protocol::ClientRequest; |
| use codex_app_server_protocol::ClientResponsePayload; |
| use codex_app_server_protocol::InitializeParams; |
| use codex_app_server_protocol::ItemCompletedNotification; |
| use codex_app_server_protocol::ItemStartedNotification; |
| use codex_app_server_protocol::JSONRPCErrorError; |
| use codex_app_server_protocol::RequestId; |
| use codex_app_server_protocol::ServerNotification; |
| use codex_app_server_protocol::ServerRequest; |
| use codex_app_server_protocol::ServerResponse; |
| use codex_app_server_protocol::Turn; |
| use codex_app_server_protocol::TurnCompletedNotification; |
| use codex_app_server_protocol::TurnError; |
| use codex_app_server_protocol::TurnItemsView; |
| use codex_app_server_protocol::TurnStartedNotification; |
| use codex_app_server_protocol::TurnStatus; |
| use codex_app_server_protocol::item_event_to_server_notification; |
| use codex_login::AuthManager; |
| use codex_login::CodexAuth; |
| use codex_login::default_client::create_client; |
| use codex_plugin::PluginId; |
| use codex_plugin::PluginTelemetryMetadata; |
| use codex_protocol::ThreadId; |
| use codex_protocol::items::CollabAgentToolCallItem; |
| use codex_protocol::items::CollabAgentToolCallStatus; |
| use codex_protocol::items::TurnItem; |
| use codex_protocol::protocol::Event; |
| use codex_protocol::protocol::EventMsg; |
| use codex_protocol::request_permissions::RequestPermissionsResponse; |
| use std::collections::HashSet; |
| use std::path::PathBuf; |
| use std::sync::Arc; |
| use std::sync::Mutex; |
| use std::time::Duration; |
| use tokio::sync::mpsc; |
| use tokio::sync::oneshot; |
|
|
| const ANALYTICS_EVENTS_QUEUE_SIZE: usize = 256; |
| const ANALYTICS_EVENTS_TIMEOUT: Duration = Duration::from_secs(10); |
| |
| const ANALYTICS_EVENTS_FLUSH_TIMEOUT: Duration = Duration::from_secs(25); |
| const ANALYTICS_EVENT_DEDUPE_MAX_KEYS: usize = 4096; |
|
|
| pub(crate) enum AnalyticsEventsQueueMessage { |
| Fact(Box<AnalyticsFact>), |
| Flush(oneshot::Sender<()>), |
| } |
|
|
| #[derive(Clone)] |
| pub(crate) struct AnalyticsEventsQueue { |
| pub(crate) sender: mpsc::Sender<AnalyticsEventsQueueMessage>, |
| pub(crate) app_used_emitted_keys: Arc<Mutex<HashSet<(String, String)>>>, |
| pub(crate) plugin_used_emitted_keys: Arc<Mutex<HashSet<(String, String)>>>, |
| } |
|
|
| #[derive(Clone)] |
| pub struct AnalyticsEventsClient { |
| queue: Option<AnalyticsEventsQueue>, |
| } |
|
|
| #[derive(Clone, Debug, Eq, PartialEq)] |
| enum AnalyticsEventsDestination { |
| Http { |
| url: String, |
| }, |
| #[cfg(debug_assertions)] |
| CaptureFile { |
| path: PathBuf, |
| }, |
| } |
|
|
| impl AnalyticsEventsDestination { |
| fn from_base_url(base_url: String) -> Self { |
| let capture_file = analytics_capture_file_from_env(); |
| Self::from_base_url_and_capture_file(base_url, capture_file) |
| } |
|
|
| fn from_base_url_and_capture_file(base_url: String, capture_file: Option<PathBuf>) -> Self { |
| #[cfg(debug_assertions)] |
| if let Some(path) = capture_file { |
| if let Err(err) = crate::analytics_capture::initialize(&path) { |
| tracing::error!( |
| path = %path.display(), |
| "failed to initialize analytics event capture; network delivery remains disabled: {err}" |
| ); |
| } |
| tracing::warn!( |
| path = %path.display(), |
| "analytics event capture enabled; network delivery is disabled" |
| ); |
| return Self::CaptureFile { path }; |
| } |
|
|
| #[cfg(not(debug_assertions))] |
| let _ = capture_file; |
|
|
| let base_url = base_url.trim_end_matches('/'); |
| Self::Http { |
| url: format!("{base_url}/codex/analytics-events/events"), |
| } |
| } |
| } |
|
|
| fn analytics_capture_file_from_env() -> Option<PathBuf> { |
| #[cfg(debug_assertions)] |
| { |
| std::env::var_os(crate::analytics_capture::ANALYTICS_EVENTS_CAPTURE_FILE_ENV_VAR) |
| .filter(|value| !value.is_empty()) |
| .map(PathBuf::from) |
| } |
|
|
| #[cfg(not(debug_assertions))] |
| None |
| } |
|
|
| impl AnalyticsEventsQueue { |
| fn new(auth_manager: Arc<AuthManager>, destination: AnalyticsEventsDestination) -> Self { |
| let (sender, mut receiver) = mpsc::channel(ANALYTICS_EVENTS_QUEUE_SIZE); |
| tokio::spawn(async move { |
| let mut reducer = AnalyticsReducer::default(); |
| while let Some(input) = receiver.recv().await { |
| let input = match input { |
| AnalyticsEventsQueueMessage::Fact(input) => *input, |
| AnalyticsEventsQueueMessage::Flush(done_tx) => { |
| let mut events = Vec::new(); |
| reducer.flush(&mut events); |
| send_track_events(&auth_manager, &destination, events).await; |
| let _ = done_tx.send(()); |
| continue; |
| } |
| }; |
| let mut events = Vec::new(); |
| reducer.ingest(input, &mut events).await; |
| send_track_events(&auth_manager, &destination, events).await; |
| } |
| }); |
| Self { |
| sender, |
| app_used_emitted_keys: Arc::new(Mutex::new(HashSet::new())), |
| plugin_used_emitted_keys: Arc::new(Mutex::new(HashSet::new())), |
| } |
| } |
|
|
| fn try_send(&self, input: AnalyticsFact) { |
| if self |
| .sender |
| .try_send(AnalyticsEventsQueueMessage::Fact(Box::new(input))) |
| .is_err() |
| { |
| |
| tracing::warn!("dropping analytics events: queue is full"); |
| } |
| } |
|
|
| pub(crate) fn should_enqueue_app_used( |
| &self, |
| tracking: &TrackEventsContext, |
| app: &AppInvocation, |
| ) -> bool { |
| let Some(connector_id) = app.connector_id.as_ref() else { |
| return true; |
| }; |
| let mut emitted = self |
| .app_used_emitted_keys |
| .lock() |
| .unwrap_or_else(std::sync::PoisonError::into_inner); |
| if emitted.len() >= ANALYTICS_EVENT_DEDUPE_MAX_KEYS { |
| emitted.clear(); |
| } |
| emitted.insert((tracking.turn_id.clone(), connector_id.clone())) |
| } |
|
|
| pub(crate) fn should_enqueue_plugin_used( |
| &self, |
| tracking: &TrackEventsContext, |
| plugin: &PluginTelemetryMetadata, |
| ) -> bool { |
| let mut emitted = self |
| .plugin_used_emitted_keys |
| .lock() |
| .unwrap_or_else(std::sync::PoisonError::into_inner); |
| if emitted.len() >= ANALYTICS_EVENT_DEDUPE_MAX_KEYS { |
| emitted.clear(); |
| } |
| let Some(plugin_id) = plugin |
| .plugin_id |
| .as_ref() |
| .map(PluginId::as_key) |
| .or_else(|| plugin.remote_plugin_id.clone()) |
| else { |
| return true; |
| }; |
| emitted.insert((tracking.turn_id.clone(), plugin_id)) |
| } |
| } |
|
|
| impl AnalyticsEventsClient { |
| pub fn new( |
| auth_manager: Arc<AuthManager>, |
| base_url: String, |
| analytics_enabled: Option<bool>, |
| ) -> Self { |
| let destination = AnalyticsEventsDestination::from_base_url(base_url); |
| Self { |
| queue: (analytics_enabled != Some(false)) |
| .then(|| AnalyticsEventsQueue::new(Arc::clone(&auth_manager), destination)), |
| } |
| } |
|
|
| pub fn disabled() -> Self { |
| Self { queue: None } |
| } |
|
|
| pub async fn flush(&self) { |
| let Some(queue) = self.queue.as_ref() else { |
| return; |
| }; |
| let (done_tx, done_rx) = oneshot::channel(); |
| let flushed = tokio::time::timeout(ANALYTICS_EVENTS_FLUSH_TIMEOUT, async { |
| if queue |
| .sender |
| .send(AnalyticsEventsQueueMessage::Flush(done_tx)) |
| .await |
| .is_err() |
| { |
| return false; |
| } |
| done_rx.await.is_ok() |
| }) |
| .await; |
|
|
| if !matches!(flushed, Ok(true)) { |
| tracing::warn!("timed out or failed while flushing analytics events"); |
| } |
| } |
|
|
| pub fn is_enabled(&self) -> bool { |
| self.queue.is_some() |
| } |
|
|
| pub fn track_plugin_measurements(&self, mut input: PluginMeasurementsInput) { |
| if input.rows.is_empty() |
| || input.rows.len() > MAX_PLUGIN_MEASUREMENTS_PER_BATCH |
| || !valid_plugin_measurement_identifier(&input.operation) |
| { |
| return; |
| } |
| input.rows.retain(valid_plugin_measurement_row); |
| if input.rows.is_empty() { |
| return; |
| } |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::PluginMeasurements(input), |
| )); |
| } |
|
|
| pub fn track_skill_invocations( |
| &self, |
| tracking: TrackEventsContext, |
| invocations: Vec<SkillInvocation>, |
| ) { |
| if invocations.is_empty() { |
| return; |
| } |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::SkillInvoked( |
| SkillInvokedInput { |
| tracking, |
| invocations, |
| }, |
| ))); |
| } |
|
|
| pub fn track_artifact_operation( |
| &self, |
| tracking: TrackEventsContext, |
| operation: ArtifactOperation, |
| ) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::ArtifactOperation(ArtifactOperationInput { |
| tracking, |
| operation, |
| }), |
| )); |
| } |
|
|
| pub fn track_initialize( |
| &self, |
| connection_id: u64, |
| params: InitializeParams, |
| product_client_id: String, |
| rpc_transport: AppServerRpcTransport, |
| ) { |
| self.record_fact(AnalyticsFact::Initialize { |
| connection_id, |
| params, |
| product_client_id, |
| runtime: current_runtime_metadata(), |
| rpc_transport, |
| }); |
| } |
|
|
| pub fn track_subagent_thread_started(&self, input: SubAgentThreadStartedInput) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::SubAgentThreadStarted(input), |
| )); |
| } |
|
|
| |
| pub fn track_guardian_session_event(&self, thread_id: ThreadId, event: &Event) { |
| if let Some(notification) = session_event_to_analytics_notification(thread_id, event) { |
| self.track_notification(¬ification); |
| } |
| } |
|
|
| pub fn track_collab_tool_call( |
| &self, |
| turn_id: String, |
| mut item: CollabAgentToolCallItem, |
| started_at_ms: i64, |
| completed_at_ms: i64, |
| ) { |
| let thread_id = item.sender_thread_id.to_string(); |
| let completed_item = TurnItem::CollabAgentToolCall(item.clone()).into(); |
| item.status = CollabAgentToolCallStatus::InProgress; |
|
|
| self.track_notification(&ServerNotification::ItemStarted(ItemStartedNotification { |
| item: TurnItem::CollabAgentToolCall(item).into(), |
| thread_id: thread_id.clone(), |
| turn_id: turn_id.clone(), |
| started_at_ms, |
| })); |
| self.track_notification(&ServerNotification::ItemCompleted( |
| ItemCompletedNotification { |
| item: completed_item, |
| thread_id, |
| turn_id, |
| completed_at_ms, |
| }, |
| )); |
| } |
|
|
| pub fn track_code_mode_tool_call(&self, input: crate::facts::CodeModeToolCallFact) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::CodeModeToolCall(input), |
| )); |
| } |
|
|
| pub fn track_control_tool_call(&self, input: crate::facts::ControlToolCallFact) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::ControlToolCall( |
| input, |
| ))); |
| } |
|
|
| pub fn track_guardian_review( |
| &self, |
| tracking: &GuardianReviewTrackContext, |
| result: GuardianReviewAnalyticsResult, |
| completed_at_ms: u64, |
| ) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::GuardianReview( |
| Box::new(tracking.event_params(result, completed_at_ms)), |
| ))); |
| } |
|
|
| pub fn track_app_mentioned(&self, tracking: TrackEventsContext, mentions: Vec<AppInvocation>) { |
| if mentions.is_empty() { |
| return; |
| } |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::AppMentioned( |
| AppMentionedInput { tracking, mentions }, |
| ))); |
| } |
|
|
| pub fn track_request( |
| &self, |
| connection_id: u64, |
| request_id: RequestId, |
| request: &ClientRequest, |
| ) { |
| if let ClientRequest::TurnInterrupt { params, .. } = request { |
| if params.turn_id.is_empty() { |
| return; |
| } |
| self.record_fact(AnalyticsFact::ExplicitClientInterruptRequest { |
| connection_id, |
| request_id, |
| turn_id: params.turn_id.clone(), |
| requested_at_ms: now_unix_millis(), |
| }); |
| return; |
| } |
| if !matches!( |
| request, |
| ClientRequest::TurnStart { .. } | ClientRequest::TurnSteer { .. } |
| ) { |
| return; |
| } |
| self.record_fact(AnalyticsFact::ClientRequest { |
| connection_id, |
| request_id, |
| request: Box::new(request.clone()), |
| }); |
| } |
|
|
| pub fn track_app_used( |
| &self, |
| tracking: TrackEventsContext, |
| app: AppInvocation, |
| elicitation_type: Option<ElicitationType>, |
| ) { |
| let Some(queue) = self.queue.as_ref() else { |
| return; |
| }; |
| if !queue.should_enqueue_app_used(&tracking, &app) { |
| return; |
| } |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::AppUsed( |
| AppUsedInput { |
| tracking, |
| app, |
| elicitation_type, |
| }, |
| ))); |
| } |
|
|
| pub fn track_mcp_tool_call_elicitation(&self, input: McpToolCallElicitation) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::McpToolCallElicitation(input), |
| )); |
| } |
|
|
| pub fn track_hook_run(&self, tracking: TrackEventsContext, hook: HookRunFact) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::HookRun( |
| HookRunInput { tracking, hook }, |
| ))); |
| } |
|
|
| pub fn track_plugin_used(&self, tracking: TrackEventsContext, plugin: PluginTelemetryMetadata) { |
| let Some(queue) = self.queue.as_ref() else { |
| return; |
| }; |
| if !queue.should_enqueue_plugin_used(&tracking, &plugin) { |
| return; |
| } |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::PluginUsed( |
| crate::facts::PluginUsedInput { tracking, plugin }, |
| ))); |
| } |
|
|
| pub fn track_plugin_install_requested( |
| &self, |
| tracking: TrackEventsContext, |
| request: PluginInstallRequested, |
| ) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::PluginInstallRequested(PluginInstallRequestedInput { |
| tracking, |
| request, |
| }), |
| )); |
| } |
|
|
| pub fn track_compaction(&self, event: crate::facts::CodexCompactionEvent) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::Compaction( |
| Box::new(event), |
| ))); |
| } |
|
|
| pub fn track_guardian_v2_event(&self, event: GuardianV2Event) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::GuardianV2( |
| Box::new(event), |
| ))); |
| } |
|
|
| pub fn track_goal_event(&self, event: CodexGoalEvent) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::Goal(Box::new( |
| event, |
| )))); |
| } |
|
|
| pub fn track_thread_hint_status(&self, event: crate::thread_hint::ThreadHintStatusEvent) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::ThreadHintStatus(Box::new(event)), |
| )); |
| } |
|
|
| pub fn track_image_preparation(&self, fact: ImagePreparationFact) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::ImagePreparation(Box::new(fact)), |
| )); |
| } |
|
|
| pub fn track_turn_resolved_config(&self, fact: TurnResolvedConfigFact) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::TurnResolvedConfig(Box::new(fact)), |
| )); |
| } |
|
|
| pub fn track_turn_token_usage(&self, fact: TurnTokenUsageFact) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::TurnTokenUsage( |
| Box::new(fact), |
| ))); |
| } |
|
|
| pub fn track_turn_profile(&self, fact: TurnProfileFact) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::TurnProfile( |
| Box::new(fact), |
| ))); |
| } |
|
|
| pub fn track_turn_codex_error(&self, fact: TurnCodexErrorFact) { |
| self.record_fact(AnalyticsFact::Custom(CustomAnalyticsFact::TurnCodexError( |
| Box::new(fact), |
| ))); |
| } |
|
|
| pub fn track_plugin_installed(&self, plugin: PluginTelemetryMetadata) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::PluginStateChanged(PluginStateChangedInput { |
| plugin, |
| state: PluginState::Installed, |
| }), |
| )); |
| } |
|
|
| pub fn track_plugin_install_failed( |
| &self, |
| plugin: PluginTelemetryMetadata, |
| source: PluginInstallSource, |
| error_type: String, |
| sub_error_type: Option<String>, |
| ) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::PluginInstallFailed(PluginInstallFailedInput { |
| plugin, |
| source, |
| error_type, |
| sub_error_type, |
| }), |
| )); |
| } |
|
|
| pub fn track_external_agent_config_import_completed( |
| &self, |
| input: ExternalAgentConfigImportCompletedInput, |
| ) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::ExternalAgentConfigImportCompleted(input), |
| )); |
| } |
|
|
| pub fn track_external_agent_config_import_failure( |
| &self, |
| input: ExternalAgentConfigImportFailureInput, |
| ) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::ExternalAgentConfigImportFailure(input), |
| )); |
| } |
|
|
| pub fn track_plugin_uninstalled(&self, plugin: PluginTelemetryMetadata) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::PluginStateChanged(PluginStateChangedInput { |
| plugin, |
| state: PluginState::Uninstalled, |
| }), |
| )); |
| } |
|
|
| pub fn track_plugin_enabled(&self, plugin: PluginTelemetryMetadata) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::PluginStateChanged(PluginStateChangedInput { |
| plugin, |
| state: PluginState::Enabled, |
| }), |
| )); |
| } |
|
|
| pub fn track_plugin_disabled(&self, plugin: PluginTelemetryMetadata) { |
| self.record_fact(AnalyticsFact::Custom( |
| CustomAnalyticsFact::PluginStateChanged(PluginStateChangedInput { |
| plugin, |
| state: PluginState::Disabled, |
| }), |
| )); |
| } |
|
|
| pub(crate) fn record_fact(&self, input: AnalyticsFact) { |
| if let Some(queue) = self.queue.as_ref() { |
| queue.try_send(input); |
| } |
| } |
|
|
| pub fn track_response( |
| &self, |
| connection_id: u64, |
| request_id: RequestId, |
| response: &ClientResponsePayload, |
| ) { |
| self.track_response_inner( |
| connection_id, |
| request_id, |
| response, |
| None, |
| ); |
| } |
|
|
| pub fn track_response_with_thread_originator( |
| &self, |
| connection_id: u64, |
| request_id: RequestId, |
| response: &ClientResponsePayload, |
| thread_originator: String, |
| ) { |
| self.track_response_inner(connection_id, request_id, response, Some(thread_originator)); |
| } |
|
|
| fn track_response_inner( |
| &self, |
| connection_id: u64, |
| request_id: RequestId, |
| response: &ClientResponsePayload, |
| thread_originator: Option<String>, |
| ) { |
| if !matches!( |
| response, |
| ClientResponsePayload::ThreadStart(_) |
| | ClientResponsePayload::ThreadResume(_) |
| | ClientResponsePayload::ThreadFork(_) |
| | ClientResponsePayload::TurnStart(_) |
| | ClientResponsePayload::TurnSteer(_) |
| | ClientResponsePayload::TurnInterrupt(_) |
| ) { |
| return; |
| } |
| if serde_json::to_writer(std::io::sink(), response).is_err() { |
| return; |
| } |
| self.record_fact(AnalyticsFact::ClientResponse { |
| connection_id, |
| request_id, |
| response: Box::new(response.clone()), |
| thread_originator, |
| }); |
| } |
|
|
| pub fn track_error_response( |
| &self, |
| connection_id: u64, |
| request_id: RequestId, |
| error: JSONRPCErrorError, |
| error_type: Option<AnalyticsJsonRpcError>, |
| ) { |
| self.record_fact(AnalyticsFact::ErrorResponse { |
| connection_id, |
| request_id, |
| error, |
| error_type, |
| }); |
| } |
|
|
| pub fn track_server_request(&self, connection_id: u64, request: ServerRequest) { |
| self.record_fact(AnalyticsFact::ServerRequest { |
| connection_id, |
| request: Box::new(request), |
| }); |
| } |
|
|
| pub fn track_server_response(&self, completed_at_ms: u64, response: ServerResponse) { |
| self.record_fact(AnalyticsFact::ServerResponse { |
| completed_at_ms, |
| response: Box::new(response), |
| }); |
| } |
|
|
| pub fn track_effective_permissions_approval_response( |
| &self, |
| completed_at_ms: u64, |
| request_id: RequestId, |
| response: RequestPermissionsResponse, |
| ) { |
| self.record_fact(AnalyticsFact::EffectivePermissionsApprovalResponse { |
| completed_at_ms, |
| request_id, |
| response: Box::new(response), |
| }); |
| } |
|
|
| pub fn track_server_request_aborted(&self, completed_at_ms: u64, request_id: RequestId) { |
| self.record_fact(AnalyticsFact::ServerRequestAborted { |
| completed_at_ms, |
| request_id, |
| }); |
| } |
|
|
| |
| pub fn track_notification(&self, notification: &ServerNotification) { |
| if !matches!( |
| notification, |
| ServerNotification::ThreadArchived(_) |
| | ServerNotification::ThreadClosed(_) |
| | ServerNotification::ThreadUnarchived(_) |
| | ServerNotification::TurnStarted(_) |
| | ServerNotification::TurnCompleted(_) |
| | ServerNotification::TurnDiffUpdated(_) |
| | ServerNotification::ItemStarted(_) |
| | ServerNotification::ItemCompleted(_) |
| | ServerNotification::ItemGuardianApprovalReviewStarted(_) |
| | ServerNotification::ItemGuardianApprovalReviewCompleted(_) |
| ) { |
| return; |
| } |
| self.record_fact(AnalyticsFact::Notification(Box::new(notification.clone()))); |
| } |
| } |
|
|
| fn session_event_to_analytics_notification( |
| thread_id: ThreadId, |
| event: &Event, |
| ) -> Option<ServerNotification> { |
| let notification = match &event.msg { |
| EventMsg::ItemStarted(started) => item_event_to_server_notification( |
| event.msg.clone(), |
| &started.thread_id.to_string(), |
| &started.turn_id, |
| ), |
| EventMsg::ItemCompleted(completed) => item_event_to_server_notification( |
| event.msg.clone(), |
| &completed.thread_id.to_string(), |
| &completed.turn_id, |
| ), |
| EventMsg::TurnStarted(started) => { |
| ServerNotification::TurnStarted(TurnStartedNotification { |
| thread_id: thread_id.to_string(), |
| turn: Turn { |
| started_at: started.started_at, |
| ..analytics_turn(&started.turn_id, TurnStatus::InProgress) |
| }, |
| }) |
| } |
| EventMsg::TurnComplete(completed) => { |
| let error = completed.error.as_ref().map(|error| TurnError { |
| message: String::new(), |
| codex_error_info: error.codex_error_info.clone().map(Into::into), |
| additional_details: None, |
| misalignment: None, |
| }); |
| let status = if error.is_some() { |
| TurnStatus::Failed |
| } else { |
| TurnStatus::Completed |
| }; |
| ServerNotification::TurnCompleted(TurnCompletedNotification { |
| thread_id: thread_id.to_string(), |
| turn: Turn { |
| error, |
| started_at: completed.started_at, |
| completed_at: completed.completed_at, |
| duration_ms: completed.duration_ms, |
| ..analytics_turn(&completed.turn_id, status) |
| }, |
| }) |
| } |
| EventMsg::TurnAborted(aborted) => { |
| ServerNotification::TurnCompleted(TurnCompletedNotification { |
| thread_id: thread_id.to_string(), |
| turn: Turn { |
| started_at: aborted.started_at, |
| completed_at: aborted.completed_at, |
| duration_ms: aborted.duration_ms, |
| ..analytics_turn( |
| aborted.turn_id.as_deref().unwrap_or(&event.id), |
| TurnStatus::Interrupted, |
| ) |
| }, |
| }) |
| } |
| |
| |
| _ => return None, |
| }; |
| match ¬ification { |
| ServerNotification::ItemStarted(ItemStartedNotification { item, .. }) |
| | ServerNotification::ItemCompleted(ItemCompletedNotification { item, .. }) => { |
| tracked_tool_item_id(item)?; |
| } |
| _ => {} |
| } |
| Some(notification) |
| } |
|
|
| fn analytics_turn(turn_id: &str, status: TurnStatus) -> Turn { |
| Turn { |
| id: turn_id.to_string(), |
| items: Vec::new(), |
| items_view: TurnItemsView::NotLoaded, |
| status, |
| error: None, |
| started_at: None, |
| completed_at: None, |
| duration_ms: None, |
| } |
| } |
|
|
| async fn send_track_events( |
| auth_manager: &AuthManager, |
| destination: &AnalyticsEventsDestination, |
| mut events: Vec<TrackEventRequest>, |
| ) { |
| if events.is_empty() { |
| return; |
| } |
|
|
| let Some(auth) = auth_manager.auth().await else { |
| return; |
| }; |
| if auth.is_api_key_auth() { |
| events.retain(TrackEventRequest::can_send_with_api_key_auth); |
| } else if !auth.uses_codex_backend() { |
| return; |
| } |
| if events.is_empty() { |
| return; |
| } |
|
|
| for events in track_event_request_batches(events) { |
| send_track_events_request(&auth, destination, events).await; |
| } |
| } |
|
|
| fn track_event_request_batches(events: Vec<TrackEventRequest>) -> Vec<Vec<TrackEventRequest>> { |
| let mut batches = Vec::new(); |
| let mut current_batch = Vec::new(); |
|
|
| for event in events { |
| if event.should_send_in_isolated_request() { |
| if !current_batch.is_empty() { |
| batches.push(current_batch); |
| current_batch = Vec::new(); |
| } |
| batches.push(vec![event]); |
| } else { |
| current_batch.push(event); |
| } |
| } |
|
|
| if !current_batch.is_empty() { |
| batches.push(current_batch); |
| } |
|
|
| batches |
| } |
|
|
| async fn send_track_events_request( |
| auth: &CodexAuth, |
| destination: &AnalyticsEventsDestination, |
| events: Vec<TrackEventRequest>, |
| ) { |
| if events.is_empty() { |
| return; |
| } |
|
|
| let payload = TrackEventsRequest { events }; |
|
|
| #[cfg(debug_assertions)] |
| if capture_track_events_request(destination, &payload) { |
| return; |
| } |
|
|
| let url = match destination { |
| AnalyticsEventsDestination::Http { url } => url, |
| #[cfg(debug_assertions)] |
| AnalyticsEventsDestination::CaptureFile { .. } => return, |
| }; |
| let response = create_client() |
| .post(url) |
| .timeout(ANALYTICS_EVENTS_TIMEOUT) |
| .headers(codex_model_provider::auth_provider_from_auth(auth).to_auth_headers()) |
| .header("Content-Type", "application/json") |
| .json(&payload) |
| .send() |
| .await; |
|
|
| match response { |
| Ok(response) if response.status().is_success() => {} |
| Ok(response) => { |
| let status = response.status(); |
| let body = response.text().await.unwrap_or_default(); |
| tracing::warn!("events failed with status {status}: {body}"); |
| } |
| Err(err) => { |
| tracing::warn!("failed to send events request: {err}"); |
| } |
| } |
| } |
|
|
| #[cfg(debug_assertions)] |
| fn capture_track_events_request( |
| destination: &AnalyticsEventsDestination, |
| payload: &TrackEventsRequest, |
| ) -> bool { |
| let AnalyticsEventsDestination::CaptureFile { path } = destination else { |
| return false; |
| }; |
|
|
| if let Err(err) = crate::analytics_capture::append_payload(path, payload) { |
| tracing::error!( |
| path = %path.display(), |
| "failed to capture analytics events; network delivery remains disabled: {err}" |
| ); |
| } |
| true |
| } |
|
|
| #[cfg(test)] |
| #[path = "client_tests.rs"] |
| mod tests; |
|
|