From 0e9f42fb46b52f00a5072acfe2630dafe95a1d99 Mon Sep 17 00:00:00 2001 From: rhan-oai Date: Tue, 21 Apr 2026 09:19:04 -0700 Subject: [PATCH] core: emit responses api call analytics --- codex-rs/analytics/src/client.rs | 16 +++- codex-rs/core/src/session/turn.rs | 154 +++++++++++++++++++++++++++++- 2 files changed, 168 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..0b61ae871c 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -1,6 +1,9 @@ use std::collections::HashMap; use std::collections::HashSet; use std::sync::Arc; +use std::time::Instant; +use std::time::SystemTime; +use std::time::UNIX_EPOCH; use crate::SkillInjections; use crate::SkillLoadOutcome; @@ -59,6 +62,8 @@ 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; @@ -91,6 +96,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 +412,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 { @@ -477,12 +484,15 @@ pub(crate) async fn run_turn( .map(|user_message| user_message.message()) .collect::>(); let turn_metadata_header = turn_context.turn_metadata_state.current_header_value(); + let turn_responses_call_index = next_turn_responses_call_index; + next_turn_responses_call_index = next_turn_responses_call_index.saturating_add(1); match run_sampling_request( Arc::clone(&sess), Arc::clone(&turn_context), Arc::clone(&turn_diff_tracker), &mut client_session, turn_metadata_header.as_deref(), + turn_responses_call_index, sampling_request_input, &explicitly_enabled_connectors, skills_outcome, @@ -1032,6 +1042,102 @@ fn filter_deferred_dynamic_tool_spec( } } +struct ResponsesApiCallObservation { + turn_responses_call_index: u64, + started_at: u64, + started_instant: Instant, + input_items: Vec, + output_items: Vec, + responses_id: Option, + token_usage: Option, +} + +impl ResponsesApiCallObservation { + fn new(turn_responses_call_index: u64) -> Self { + Self { + turn_responses_call_index, + started_at: current_unix_seconds(), + started_instant: Instant::now(), + input_items: Vec::new(), + output_items: Vec::new(), + responses_id: None, + token_usage: None, + } + } + + fn record_prompt(&mut self, prompt: &Prompt) { + self.input_items = prompt.input.clone(); + self.output_items.clear(); + self.responses_id = None; + self.token_usage = None; + } + + fn record_output_item_done(&mut self, item: &ResponseItem) { + self.output_items.push(item.clone()); + } + + fn record_completed(&mut self, response_id: String, token_usage: Option) { + self.responses_id = Some(response_id); + self.token_usage = token_usage; + } + + fn into_input( + self, + status: CodexResponsesApiCallStatus, + error: Option, + ) -> CodexResponsesApiCallInput { + let completed_at = current_unix_seconds(); + CodexResponsesApiCallInput { + responses_id: self.responses_id, + turn_responses_call_index: self.turn_responses_call_index, + status, + error, + started_at: self.started_at, + completed_at: Some(completed_at), + duration_ms: Some(self.started_instant.elapsed().as_millis() as u64), + input_items: self.input_items, + output_items: self.output_items, + token_usage: self.token_usage, + } + } +} + +fn current_unix_seconds() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs() +} + +fn responses_api_call_error_status(err: &CodexErr) -> CodexResponsesApiCallStatus { + match err { + CodexErr::TurnAborted => CodexResponsesApiCallStatus::Interrupted, + _ => CodexResponsesApiCallStatus::Failed, + } +} + +fn track_responses_api_call_observation( + sess: &Session, + turn_context: &TurnContext, + observation: ResponsesApiCallObservation, + status: CodexResponsesApiCallStatus, + error: Option, +) { + if !sess.enabled(Feature::GeneralAnalytics) { + return; + } + 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(), + ), + observation.into_input(status, error), + ); +} + #[allow(clippy::too_many_arguments)] #[instrument(level = "trace", skip_all, @@ -1047,6 +1153,7 @@ async fn run_sampling_request( turn_diff_tracker: SharedTurnDiffTracker, client_session: &mut ModelClientSession, turn_metadata_header: Option<&str>, + turn_responses_call_index: u64, input: Vec, explicitly_enabled_connectors: &HashSet, skills_outcome: Option<&SkillLoadOutcome>, @@ -1083,6 +1190,8 @@ async fn run_sampling_request( .await; let mut retries = 0; let mut initial_input = Some(input); + let mut responses_api_call_observation = + ResponsesApiCallObservation::new(turn_responses_call_index); loop { let prompt_input = if let Some(input) = initial_input.take() { input @@ -1097,6 +1206,7 @@ async fn run_sampling_request( turn_context.as_ref(), base_instructions.clone(), ); + responses_api_call_observation.record_prompt(&prompt); let err = match try_run_sampling_request( tool_runtime.clone(), Arc::clone(&sess), @@ -1106,15 +1216,30 @@ async fn run_sampling_request( Arc::clone(&turn_diff_tracker), server_model_warning_emitted_for_turn, &prompt, + &mut responses_api_call_observation, cancellation_token.child_token(), ) .await { Ok(output) => { + track_responses_api_call_observation( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_observation, + CodexResponsesApiCallStatus::Completed, + /*error*/ None, + ); return Ok(output); } Err(CodexErr::ContextWindowExceeded) => { sess.set_total_tokens_full(&turn_context).await; + track_responses_api_call_observation( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_observation, + CodexResponsesApiCallStatus::Failed, + Some("context window exceeded".to_string()), + ); return Err(CodexErr::ContextWindowExceeded); } Err(CodexErr::UsageLimitReached(e)) => { @@ -1122,12 +1247,28 @@ async fn run_sampling_request( if let Some(rate_limits) = rate_limits { sess.update_rate_limits(&turn_context, *rate_limits).await; } + track_responses_api_call_observation( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_observation, + CodexResponsesApiCallStatus::Failed, + Some(e.to_string()), + ); return Err(CodexErr::UsageLimitReached(e)); } Err(err) => err, }; if !err.is_retryable() { + let status = responses_api_call_error_status(&err); + let error = Some(format!("{err:#}")); + track_responses_api_call_observation( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_observation, + status, + error, + ); return Err(err); } @@ -1179,6 +1320,14 @@ async fn run_sampling_request( } tokio::time::sleep(delay).await; } else { + let error = Some(format!("{err:#}")); + track_responses_api_call_observation( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_observation, + CodexResponsesApiCallStatus::Failed, + error, + ); return Err(err); } } @@ -1895,6 +2044,7 @@ async fn try_run_sampling_request( turn_diff_tracker: SharedTurnDiffTracker, server_model_warning_emitted_for_turn: &mut bool, prompt: &Prompt, + responses_api_call_observation: &mut ResponsesApiCallObservation, cancellation_token: CancellationToken, ) -> CodexResult { feedback_tags!( @@ -1970,6 +2120,7 @@ async fn try_run_sampling_request( match event { ResponseEvent::Created => {} ResponseEvent::OutputItemDone(item) => { + responses_api_call_observation.record_output_item_done(&item); if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take() && let Some(event) = consumer.flush_on_complete() { @@ -2140,9 +2291,10 @@ 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_observation.record_completed(response_id, token_usage.clone()); flush_assistant_text_segments_all( &sess, &turn_context,