diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index efbcca6716..aad57dfab2 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; @@ -58,11 +61,15 @@ 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::CodexResponsesApiCallFact; +use codex_analytics::CodexResponsesApiCallStatus; +use codex_analytics::CodexResponsesApiItemPhase; 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::response_items_metadata; use codex_async_utils::OrCancelExt; use codex_features::Feature; use codex_hooks::HookEvent; @@ -90,6 +97,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; @@ -394,6 +402,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 { @@ -465,12 +474,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, @@ -991,6 +1003,127 @@ pub(crate) fn build_prompt( } } +const RESPONSES_API_ANALYTICS_ERROR_MAX_BYTES: usize = 1024; + +struct ResponsesApiCallAnalyticsState { + turn_responses_call_index: u64, + started_at: u64, + started_instant: Instant, + input_items: Vec, + output_items: Vec, + responses_id: Option, + token_usage: Option, +} + +impl ResponsesApiCallAnalyticsState { + 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_fact( + self, + sess: &Session, + turn_context: &TurnContext, + status: CodexResponsesApiCallStatus, + error: Option, + ) -> CodexResponsesApiCallFact { + let completed_at = current_unix_seconds(); + let mut items = + response_items_metadata(CodexResponsesApiItemPhase::Input, &self.input_items); + items.extend(response_items_metadata( + CodexResponsesApiItemPhase::Output, + &self.output_items, + )); + let token_usage = self.token_usage; + CodexResponsesApiCallFact { + thread_id: sess.conversation_id.to_string(), + turn_id: turn_context.sub_id.clone(), + 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_item_count: self.input_items.len(), + output_item_count: self.output_items.len(), + 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), + reasoning_output_tokens: token_usage + .as_ref() + .map(|usage| usage.reasoning_output_tokens), + total_tokens: token_usage.as_ref().map(|usage| usage.total_tokens), + items, + } + } +} + +fn current_unix_seconds() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs() +} + +fn truncate_responses_api_error(mut error: String) -> String { + if error.len() <= RESPONSES_API_ANALYTICS_ERROR_MAX_BYTES { + return error; + } + let mut truncate_at = RESPONSES_API_ANALYTICS_ERROR_MAX_BYTES; + while !error.is_char_boundary(truncate_at) { + truncate_at -= 1; + } + error.truncate(truncate_at); + error +} + +fn responses_api_call_error_status(err: &CodexErr) -> CodexResponsesApiCallStatus { + match err { + CodexErr::TurnAborted => CodexResponsesApiCallStatus::Interrupted, + _ => CodexResponsesApiCallStatus::Failed, + } +} + +fn track_responses_api_call_analytics( + sess: &Session, + turn_context: &TurnContext, + analytics: ResponsesApiCallAnalyticsState, + status: CodexResponsesApiCallStatus, + error: Option, +) { + if !sess.enabled(Feature::GeneralAnalytics) { + return; + } + sess.services + .analytics_events_client + .track_responses_api_call(analytics.into_fact(sess, turn_context, status, error)); +} + #[allow(clippy::too_many_arguments)] #[instrument(level = "trace", skip_all, @@ -1006,6 +1139,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>, @@ -1042,6 +1176,8 @@ async fn run_sampling_request( .await; let mut retries = 0; let mut initial_input = Some(input); + let mut responses_api_call_analytics = + ResponsesApiCallAnalyticsState::new(turn_responses_call_index); loop { let prompt_input = if let Some(input) = initial_input.take() { input @@ -1056,6 +1192,7 @@ async fn run_sampling_request( turn_context.as_ref(), base_instructions.clone(), ); + responses_api_call_analytics.record_prompt(&prompt); let err = match try_run_sampling_request( tool_runtime.clone(), Arc::clone(&sess), @@ -1065,15 +1202,30 @@ async fn run_sampling_request( Arc::clone(&turn_diff_tracker), server_model_warning_emitted_for_turn, &prompt, + &mut responses_api_call_analytics, cancellation_token.child_token(), ) .await { Ok(output) => { + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Completed, + /*error*/ None, + ); return Ok(output); } Err(CodexErr::ContextWindowExceeded) => { sess.set_total_tokens_full(&turn_context).await; + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Failed, + Some("context window exceeded".to_string()), + ); return Err(CodexErr::ContextWindowExceeded); } Err(CodexErr::UsageLimitReached(e)) => { @@ -1081,12 +1233,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_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Failed, + Some(truncate_responses_api_error(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(truncate_responses_api_error(format!("{err:#}"))); + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + status, + error, + ); return Err(err); } @@ -1138,6 +1306,14 @@ async fn run_sampling_request( } tokio::time::sleep(delay).await; } else { + let error = Some(truncate_responses_api_error(format!("{err:#}"))); + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Failed, + error, + ); return Err(err); } } @@ -1850,6 +2026,7 @@ async fn try_run_sampling_request( turn_diff_tracker: SharedTurnDiffTracker, server_model_warning_emitted_for_turn: &mut bool, prompt: &Prompt, + responses_api_call_analytics: &mut ResponsesApiCallAnalyticsState, cancellation_token: CancellationToken, ) -> CodexResult { feedback_tags!( @@ -1925,6 +2102,7 @@ async fn try_run_sampling_request( match event { ResponseEvent::Created => {} ResponseEvent::OutputItemDone(item) => { + responses_api_call_analytics.record_output_item_done(&item); if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take() && let Some(event) = consumer.flush_on_complete() { @@ -2095,9 +2273,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_analytics.record_completed(response_id, token_usage.clone()); flush_assistant_text_segments_all( &sess, &turn_context,