From 28611de7b47e9f97c0cdca09be925dcefd400d23 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 | 30 ++-- codex-rs/analytics/src/facts.rs | 6 +- codex-rs/analytics/src/lib.rs | 19 +++ codex-rs/analytics/src/response_items.rs | 4 +- codex-rs/core/src/session/turn.rs | 170 ++++++++++++++++++++++- 5 files changed, 211 insertions(+), 18 deletions(-) diff --git a/codex-rs/analytics/src/client.rs b/codex-rs/analytics/src/client.rs index e0de1975dd..79c77770ba 100644 --- a/codex-rs/analytics/src/client.rs +++ b/codex-rs/analytics/src/client.rs @@ -10,7 +10,6 @@ use crate::facts::AppMentionedInput; use crate::facts::AppUsedInput; use crate::facts::CodexResponsesApiCallFact; use crate::facts::CodexResponsesApiCallInput; -use crate::facts::CodexResponsesApiItemPhase; use crate::facts::CustomAnalyticsFact; use crate::facts::HookRunFact; use crate::facts::HookRunInput; @@ -23,7 +22,6 @@ use crate::facts::TrackEventsContext; use crate::facts::TurnResolvedConfigFact; use crate::facts::TurnTokenUsageFact; use crate::reducer::AnalyticsReducer; -use crate::response_items::response_items_metadata; use codex_app_server_protocol::ClientRequest; use codex_app_server_protocol::ClientResponse; use codex_app_server_protocol::InitializeParams; @@ -42,6 +40,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 { @@ -223,25 +222,20 @@ impl AnalyticsEventsClient { tracking: TrackEventsContext, input: CodexResponsesApiCallInput, ) { - let mut items = - response_items_metadata(CodexResponsesApiItemPhase::Input, &input.input_items); - items.extend(response_items_metadata( - CodexResponsesApiItemPhase::Output, - &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, - input_item_count: input.input_items.len(), - output_item_count: input.output_items.len(), + input_item_count: input.input_item_count, + output_item_count: input.output_item_count, input_tokens: token_usage.as_ref().map(|usage| usage.input_tokens), cached_input_tokens: token_usage.as_ref().map(|usage| usage.cached_input_tokens), output_tokens: token_usage.as_ref().map(|usage| usage.output_tokens), @@ -249,7 +243,7 @@ impl AnalyticsEventsClient { .as_ref() .map(|usage| usage.reasoning_output_tokens), total_tokens: token_usage.as_ref().map(|usage| usage.total_tokens), - items, + items: input.items, }; self.record_fact(AnalyticsFact::Custom( @@ -339,6 +333,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/analytics/src/facts.rs b/codex-rs/analytics/src/facts.rs index 5f8fb34317..c53becb7ba 100644 --- a/codex-rs/analytics/src/facts.rs +++ b/codex-rs/analytics/src/facts.rs @@ -14,7 +14,6 @@ use codex_protocol::config_types::Personality; use codex_protocol::config_types::ReasoningSummary; use codex_protocol::config_types::ServiceTier; use codex_protocol::models::MessagePhase; -use codex_protocol::models::ResponseItem; use codex_protocol::openai_models::ReasoningEffort; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::HookEventName; @@ -322,8 +321,9 @@ pub struct CodexResponsesApiCallInput { pub started_at: u64, pub completed_at: Option, pub duration_ms: Option, - pub input_items: Vec, - pub output_items: Vec, + pub input_item_count: usize, + pub output_item_count: usize, + pub items: Vec, pub token_usage: Option, } diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index cc95fa8587..0e19a2a248 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -4,6 +4,7 @@ mod facts; mod reducer; mod response_items; +use codex_protocol::models::ResponseItem; use serde::Serialize; use std::time::SystemTime; use std::time::UNIX_EPOCH; @@ -22,6 +23,7 @@ pub use facts::AppInvocation; pub use facts::CodexCompactionEvent; pub use facts::CodexResponsesApiCallInput; pub use facts::CodexResponsesApiCallStatus; +pub use facts::CodexResponsesApiItemMetadata; pub use facts::CodexTurnSteerEvent; pub use facts::CompactionImplementation; pub use facts::CompactionPhase; @@ -54,6 +56,23 @@ pub fn now_unix_seconds() -> u64 { .as_secs() } +pub fn responses_api_input_items_metadata( + items: &[ResponseItem], +) -> Vec { + response_items::response_items_metadata(facts::CodexResponsesApiItemPhase::Input, items) +} + +pub fn responses_api_output_item_metadata( + item_index: usize, + item: &ResponseItem, +) -> CodexResponsesApiItemMetadata { + response_items::response_item_metadata( + facts::CodexResponsesApiItemPhase::Output, + item_index, + item, + ) +} + pub(crate) fn serialized_string(value: &T) -> Option { match serde_json::to_value(value).ok()? { serde_json::Value::String(value) => Some(value), diff --git a/codex-rs/analytics/src/response_items.rs b/codex-rs/analytics/src/response_items.rs index 91c927120c..d45c209124 100644 --- a/codex-rs/analytics/src/response_items.rs +++ b/codex-rs/analytics/src/response_items.rs @@ -13,7 +13,7 @@ use codex_protocol::models::ReasoningItemContent; use codex_protocol::models::ReasoningItemReasoningSummary; use codex_protocol::models::ResponseItem; -pub fn response_items_metadata( +pub(crate) fn response_items_metadata( phase: CodexResponsesApiItemPhase, items: &[ResponseItem], ) -> Vec { @@ -24,7 +24,7 @@ pub fn response_items_metadata( .collect() } -fn response_item_metadata( +pub(crate) fn response_item_metadata( item_phase: CodexResponsesApiItemPhase, item_index: usize, item: &ResponseItem, diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 7b27c8080a..022e6f3446 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,17 @@ 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::CodexResponsesApiItemMetadata; 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_analytics::responses_api_input_items_metadata; +use codex_analytics::responses_api_output_item_metadata; use codex_async_utils::OrCancelExt; use codex_features::Feature; use codex_hooks::HookEvent; @@ -91,6 +98,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 +414,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 +492,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, @@ -495,6 +505,7 @@ pub(crate) async fn run_turn( let SamplingRequestResult { needs_follow_up: model_needs_follow_up, last_agent_message: sampling_request_last_agent_message, + .. } = sampling_request_output; can_drain_pending_input = true; let has_pending_input = sess.has_pending_input().await; @@ -1032,6 +1043,81 @@ fn filter_deferred_dynamic_tool_spec( } } +struct ResponsesApiCallAttempt { + turn_responses_call_index: u64, + started_at: u64, + started_instant: Instant, + stream_started: bool, + input_item_count: usize, + output_item_count: usize, + items: Vec, + responses_id: Option, + token_usage: Option, +} + +impl ResponsesApiCallAttempt { + fn new(turn_responses_call_index: u64, input_items: &[ResponseItem]) -> Self { + let input_item_count = input_items.len(); + let items = responses_api_input_items_metadata(input_items); + Self { + turn_responses_call_index, + started_at: now_unix_seconds(), + started_instant: Instant::now(), + stream_started: false, + input_item_count, + output_item_count: 0, + items, + responses_id: None, + token_usage: None, + } + } + + fn record_output_item(&mut self, item: &ResponseItem) { + let item_index = self.output_item_count; + self.items + .push(responses_api_output_item_metadata(item_index, item)); + self.output_item_count += 1; + } +} + +fn emit_responses_api_call_attempt( + sess: &Session, + turn_context: &TurnContext, + attempt: ResponsesApiCallAttempt, + status: CodexResponsesApiCallStatus, + error: Option, +) { + if !sess.enabled(Feature::GeneralAnalytics) { + return; + } + if status == CodexResponsesApiCallStatus::Interrupted && !attempt.stream_started { + 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_item_count: attempt.input_item_count, + output_item_count: attempt.output_item_count, + items: attempt.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 +1133,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 +1184,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); let err = match try_run_sampling_request( tool_runtime.clone(), Arc::clone(&sess), @@ -1106,15 +1197,40 @@ 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) => { + let status = if output.response_completed { + CodexResponsesApiCallStatus::Completed + } else { + CodexResponsesApiCallStatus::Interrupted + }; + let error = if output.response_completed { + None + } else { + Some("stream preempted by pending input".to_string()) + }; + emit_responses_api_call_attempt( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_attempt, + status, + error, + ); 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 +1238,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 +1275,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 +1304,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 +1331,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); } } @@ -1330,6 +1490,7 @@ pub(crate) async fn built_tools( struct SamplingRequestResult { needs_follow_up: bool, last_agent_message: Option, + response_completed: bool, } /// Ephemeral per-response state for streaming a single proposed plan. @@ -1895,6 +2056,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!( @@ -1918,6 +2080,7 @@ async fn try_run_sampling_request( .instrument(trace_span!("stream_request")) .or_cancel(&cancellation_token) .await??; + responses_api_call_attempt.stream_started = true; let mut in_flight: FuturesOrdered>> = FuturesOrdered::new(); let mut needs_follow_up = false; @@ -1970,6 +2133,7 @@ async fn try_run_sampling_request( match event { ResponseEvent::Created => {} ResponseEvent::OutputItemDone(item) => { + responses_api_call_attempt.record_output_item(&item); if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take() && let Some(event) = consumer.flush_on_complete() { @@ -2049,6 +2213,7 @@ async fn try_run_sampling_request( break Ok(SamplingRequestResult { needs_follow_up: true, last_agent_message, + response_completed: false, }); } } @@ -2140,9 +2305,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, @@ -2157,6 +2324,7 @@ async fn try_run_sampling_request( break Ok(SamplingRequestResult { needs_follow_up, last_agent_message, + response_completed: true, }); } ResponseEvent::OutputTextDelta(delta) => {