core: emit responses api call analytics

This commit is contained in:
rhan-oai
2026-04-21 10:14:03 -07:00
parent 355fb1c285
commit bf8fb10a8d
2 changed files with 149 additions and 2 deletions

View File

@@ -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<AuthManager>,
base_url: &str,

View File

@@ -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<ResponseItem>,
output_items: Vec<ResponseItem>,
responses_id: Option<String>,
token_usage: Option<TokenUsage>,
}
impl ResponsesApiCallAttempt {
fn new(turn_responses_call_index: u64, input_items: Vec<ResponseItem>) -> 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<String>,
) {
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<ResponseItem>,
explicitly_enabled_connectors: &HashSet<String>,
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<SamplingRequestResult> {
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,