From bf8fb10a8dd1b2e9922ccdbdfdecce58ab7d5b95 Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Tue, 21 Apr 2026 10:14:03 -0700 Subject: [PATCH] core: emit responses api call analytics --- codex-rs/analytics/src/client.rs | 16 +++- codex-rs/core/src/session/turn.rs | 135 +++++++++++++++++++++++++++++- 2 files changed, 149 insertions(+), 2 deletions(-) diff --git a/codex-rs/analytics/src/client.rs b/codex-rs/analytics/src/client.rs index e0de1975dd..f1449e0ec0 100644 --- a/codex-rs/analytics/src/client.rs +++ b/codex-rs/analytics/src/client.rs @@ -42,6 +42,7 @@ use tokio::sync::mpsc; const ANALYTICS_EVENTS_QUEUE_SIZE: usize = 256; const ANALYTICS_EVENTS_TIMEOUT: Duration = Duration::from_secs(10); const ANALYTICS_EVENT_DEDUPE_MAX_KEYS: usize = 4096; +const RESPONSES_API_ERROR_MAX_BYTES: usize = 1024; #[derive(Clone)] pub(crate) struct AnalyticsEventsQueue { @@ -230,13 +231,14 @@ impl AnalyticsEventsClient { &input.output_items, )); let token_usage = input.token_usage; + let error = input.error.map(truncate_responses_api_error); let event = CodexResponsesApiCallFact { thread_id: tracking.thread_id, turn_id: tracking.turn_id, responses_id: input.responses_id, turn_responses_call_index: input.turn_responses_call_index, status: input.status, - error: input.error, + error, started_at: input.started_at, completed_at: input.completed_at, duration_ms: input.duration_ms, @@ -339,6 +341,18 @@ impl AnalyticsEventsClient { } } +fn truncate_responses_api_error(mut error: String) -> String { + if error.len() <= RESPONSES_API_ERROR_MAX_BYTES { + return error; + } + let mut truncate_at = RESPONSES_API_ERROR_MAX_BYTES; + while !error.is_char_boundary(truncate_at) { + truncate_at -= 1; + } + error.truncate(truncate_at); + error +} + async fn send_track_events( auth_manager: &Arc, base_url: &str, diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 7b27c8080a..153c8e3b08 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -1,6 +1,7 @@ use std::collections::HashMap; use std::collections::HashSet; use std::sync::Arc; +use std::time::Instant; use crate::SkillInjections; use crate::SkillLoadOutcome; @@ -59,11 +60,14 @@ use crate::unavailable_tool::collect_unavailable_called_tools; use crate::util::backoff; use crate::util::error_or_panic; use codex_analytics::AppInvocation; +use codex_analytics::CodexResponsesApiCallInput; +use codex_analytics::CodexResponsesApiCallStatus; use codex_analytics::CompactionPhase; use codex_analytics::CompactionReason; use codex_analytics::InvocationType; use codex_analytics::TurnResolvedConfigFact; use codex_analytics::build_track_events_context; +use codex_analytics::now_unix_seconds; use codex_async_utils::OrCancelExt; use codex_features::Feature; use codex_hooks::HookEvent; @@ -91,6 +95,7 @@ use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::PlanDeltaEvent; use codex_protocol::protocol::ReasoningContentDeltaEvent; use codex_protocol::protocol::ReasoningRawContentDeltaEvent; +use codex_protocol::protocol::TokenUsage; use codex_protocol::protocol::TurnDiffEvent; use codex_protocol::protocol::WarningEvent; use codex_protocol::user_input::UserInput; @@ -406,6 +411,7 @@ pub(crate) async fn run_turn( // 1. At the start of a turn, so the fresh user prompt in `input` gets sampled first. // 2. After auto-compact, when model/tool continuation needs to resume before any steer. let mut can_drain_pending_input = input.is_empty(); + let mut next_turn_responses_call_index: u64 = 0; loop { if run_pending_session_start_hooks(&sess, &turn_context).await { @@ -483,6 +489,7 @@ pub(crate) async fn run_turn( Arc::clone(&turn_diff_tracker), &mut client_session, turn_metadata_header.as_deref(), + &mut next_turn_responses_call_index, sampling_request_input, &explicitly_enabled_connectors, skills_outcome, @@ -1032,6 +1039,64 @@ fn filter_deferred_dynamic_tool_spec( } } +struct ResponsesApiCallAttempt { + turn_responses_call_index: u64, + started_at: u64, + started_instant: Instant, + input_items: Vec, + output_items: Vec, + responses_id: Option, + token_usage: Option, +} + +impl ResponsesApiCallAttempt { + fn new(turn_responses_call_index: u64, input_items: Vec) -> Self { + Self { + turn_responses_call_index, + started_at: now_unix_seconds(), + started_instant: Instant::now(), + input_items, + output_items: Vec::new(), + responses_id: None, + token_usage: None, + } + } +} + +fn emit_responses_api_call_attempt( + sess: &Session, + turn_context: &TurnContext, + attempt: ResponsesApiCallAttempt, + status: CodexResponsesApiCallStatus, + error: Option, +) { + if !sess.enabled(Feature::GeneralAnalytics) { + return; + } + let input = CodexResponsesApiCallInput { + responses_id: attempt.responses_id, + turn_responses_call_index: attempt.turn_responses_call_index, + status, + error, + started_at: attempt.started_at, + completed_at: Some(now_unix_seconds()), + duration_ms: Some(attempt.started_instant.elapsed().as_millis() as u64), + input_items: attempt.input_items, + output_items: attempt.output_items, + token_usage: attempt.token_usage, + }; + sess.services + .analytics_events_client + .track_responses_api_call( + build_track_events_context( + turn_context.model_info.slug.clone(), + sess.conversation_id.to_string(), + turn_context.sub_id.clone(), + ), + input, + ); +} + #[allow(clippy::too_many_arguments)] #[instrument(level = "trace", skip_all, @@ -1047,6 +1112,7 @@ async fn run_sampling_request( turn_diff_tracker: SharedTurnDiffTracker, client_session: &mut ModelClientSession, turn_metadata_header: Option<&str>, + next_turn_responses_call_index: &mut u64, input: Vec, explicitly_enabled_connectors: &HashSet, skills_outcome: Option<&SkillLoadOutcome>, @@ -1097,6 +1163,10 @@ async fn run_sampling_request( turn_context.as_ref(), base_instructions.clone(), ); + let turn_responses_call_index = *next_turn_responses_call_index; + *next_turn_responses_call_index = (*next_turn_responses_call_index).saturating_add(1); + let mut responses_api_call_attempt = + ResponsesApiCallAttempt::new(turn_responses_call_index, prompt.input.clone()); let err = match try_run_sampling_request( tool_runtime.clone(), Arc::clone(&sess), @@ -1106,15 +1176,30 @@ async fn run_sampling_request( Arc::clone(&turn_diff_tracker), server_model_warning_emitted_for_turn, &prompt, + &mut responses_api_call_attempt, cancellation_token.child_token(), ) .await { Ok(output) => { + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + CodexResponsesApiCallStatus::Completed, + /*error*/ None, + ); return Ok(output); } Err(CodexErr::ContextWindowExceeded) => { sess.set_total_tokens_full(&turn_context).await; + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + CodexResponsesApiCallStatus::Failed, + Some("context window exceeded".to_string()), + ); return Err(CodexErr::ContextWindowExceeded); } Err(CodexErr::UsageLimitReached(e)) => { @@ -1122,12 +1207,32 @@ async fn run_sampling_request( if let Some(rate_limits) = rate_limits { sess.update_rate_limits(&turn_context, *rate_limits).await; } + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + CodexResponsesApiCallStatus::Failed, + Some(e.to_string()), + ); return Err(CodexErr::UsageLimitReached(e)); } Err(err) => err, }; if !err.is_retryable() { + let status = if matches!(err, CodexErr::TurnAborted) { + CodexResponsesApiCallStatus::Interrupted + } else { + CodexResponsesApiCallStatus::Failed + }; + let error = Some(format!("{err:#}")); + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + status, + error, + ); return Err(err); } @@ -1139,6 +1244,14 @@ async fn run_sampling_request( &turn_context.model_info, ) { + let error = Some(format!("{err:#}")); + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + CodexResponsesApiCallStatus::Failed, + error, + ); sess.send_event( &turn_context, EventMsg::Warning(WarningEvent { @@ -1160,6 +1273,14 @@ async fn run_sampling_request( warn!( "stream disconnected - retrying sampling request ({retries}/{max_retries} in {delay:?})...", ); + let error = Some(format!("{err:#}")); + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + CodexResponsesApiCallStatus::Failed, + error, + ); // In release builds, hide the first websocket retry notification to reduce noisy // transient reconnect messages. In debug builds, keep full visibility for diagnosis. @@ -1179,6 +1300,14 @@ async fn run_sampling_request( } tokio::time::sleep(delay).await; } else { + let error = Some(format!("{err:#}")); + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + CodexResponsesApiCallStatus::Failed, + error, + ); return Err(err); } } @@ -1895,6 +2024,7 @@ async fn try_run_sampling_request( turn_diff_tracker: SharedTurnDiffTracker, server_model_warning_emitted_for_turn: &mut bool, prompt: &Prompt, + responses_api_call_attempt: &mut ResponsesApiCallAttempt, cancellation_token: CancellationToken, ) -> CodexResult { feedback_tags!( @@ -1970,6 +2100,7 @@ async fn try_run_sampling_request( match event { ResponseEvent::Created => {} ResponseEvent::OutputItemDone(item) => { + responses_api_call_attempt.output_items.push(item.clone()); if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take() && let Some(event) = consumer.flush_on_complete() { @@ -2140,9 +2271,11 @@ async fn try_run_sampling_request( sess.services.models_manager.refresh_if_new_etag(etag).await; } ResponseEvent::Completed { - response_id: _, + response_id, token_usage, } => { + responses_api_call_attempt.responses_id = Some(response_id); + responses_api_call_attempt.token_usage = token_usage.clone(); flush_assistant_text_segments_all( &sess, &turn_context,