mirror of
https://github.com/openai/codex.git
synced 2026-09-20 12:47:38 +00:00
## What changed Add `voice_session_id` to turn, app-use, skill-invocation, and MCP tool-call analytics events. Track realtime session lifecycle and handoff markers to preserve attribution when a handoff starts a turn after the voice session closes, or steers an active turn. Keep that attribution from carrying over to subsequent text turns. Record handoff markers using only the thread ID, without retaining transcript content in analytics facts. ## Testing Add coverage for attribution after realtime closure, active-turn steering without tagging the next text turn, and handoff markers that exclude transcripts. Extend event serialization and MCP tool-call tests for `voice_session_id`. GitOrigin-RevId: 024b0e7e4c6fbaf8325574de6c24c048bb4ca032
974 lines
31 KiB
Rust
974 lines
31 KiB
Rust
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);
|
|
// Covers two sequential POSTs plus queue/barrier scheduling; additional queued sends remain best-effort.
|
|
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()
|
|
{
|
|
//TODO: add a metric for this
|
|
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),
|
|
));
|
|
}
|
|
|
|
/// Logs Guardian turn and tool events that bypass the app-server listener.
|
|
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,
|
|
/*thread_originator*/ 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,
|
|
});
|
|
}
|
|
|
|
/// Records analytics-relevant notifications without cloning ignored variants.
|
|
pub fn track_notification(&self, notification: &ServerNotification) {
|
|
if let ServerNotification::ThreadRealtimeItemAdded(handoff) = notification {
|
|
if handoff.item.get("type").and_then(serde_json::Value::as_str)
|
|
== Some("handoff_request")
|
|
{
|
|
self.record_fact(AnalyticsFact::RealtimeHandoffRequested {
|
|
thread_id: handoff.thread_id.clone(),
|
|
});
|
|
}
|
|
return;
|
|
}
|
|
if !matches!(
|
|
notification,
|
|
ServerNotification::ThreadArchived(_)
|
|
| ServerNotification::ThreadClosed(_)
|
|
| ServerNotification::ThreadUnarchived(_)
|
|
| ServerNotification::ThreadRealtimeStarted(_)
|
|
| ServerNotification::ThreadRealtimeClosed(_)
|
|
| 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,
|
|
)
|
|
},
|
|
})
|
|
}
|
|
// Legacy tool events accompany canonical items. Messages, reasoning, and review
|
|
// content must not enter the analytics queue.
|
|
_ => 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;
|