Skip to content

Commit ea9b6e1

Browse files
committed
core: emit responses api call analytics
1 parent 355fb1c commit ea9b6e1

2 files changed

Lines changed: 149 additions & 2 deletions

File tree

codex-rs/analytics/src/client.rs

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ use tokio::sync::mpsc;
4242
const ANALYTICS_EVENTS_QUEUE_SIZE: usize = 256;
4343
const ANALYTICS_EVENTS_TIMEOUT: Duration = Duration::from_secs(10);
4444
const ANALYTICS_EVENT_DEDUPE_MAX_KEYS: usize = 4096;
45+
const RESPONSES_API_ERROR_MAX_BYTES: usize = 1024;
4546

4647
#[derive(Clone)]
4748
pub(crate) struct AnalyticsEventsQueue {
@@ -230,13 +231,14 @@ impl AnalyticsEventsClient {
230231
&input.output_items,
231232
));
232233
let token_usage = input.token_usage;
234+
let error = input.error.map(truncate_responses_api_error);
233235
let event = CodexResponsesApiCallFact {
234236
thread_id: tracking.thread_id,
235237
turn_id: tracking.turn_id,
236238
responses_id: input.responses_id,
237239
turn_responses_call_index: input.turn_responses_call_index,
238240
status: input.status,
239-
error: input.error,
241+
error,
240242
started_at: input.started_at,
241243
completed_at: input.completed_at,
242244
duration_ms: input.duration_ms,
@@ -339,6 +341,18 @@ impl AnalyticsEventsClient {
339341
}
340342
}
341343

344+
fn truncate_responses_api_error(mut error: String) -> String {
345+
if error.len() <= RESPONSES_API_ERROR_MAX_BYTES {
346+
return error;
347+
}
348+
let mut truncate_at = RESPONSES_API_ERROR_MAX_BYTES;
349+
while !error.is_char_boundary(truncate_at) {
350+
truncate_at -= 1;
351+
}
352+
error.truncate(truncate_at);
353+
error
354+
}
355+
342356
async fn send_track_events(
343357
auth_manager: &Arc<AuthManager>,
344358
base_url: &str,

codex-rs/core/src/session/turn.rs

Lines changed: 134 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
use std::collections::HashMap;
22
use std::collections::HashSet;
33
use std::sync::Arc;
4+
use std::time::Instant;
45

56
use crate::SkillInjections;
67
use crate::SkillLoadOutcome;
@@ -59,11 +60,14 @@ use crate::unavailable_tool::collect_unavailable_called_tools;
5960
use crate::util::backoff;
6061
use crate::util::error_or_panic;
6162
use codex_analytics::AppInvocation;
63+
use codex_analytics::CodexResponsesApiCallInput;
64+
use codex_analytics::CodexResponsesApiCallStatus;
6265
use codex_analytics::CompactionPhase;
6366
use codex_analytics::CompactionReason;
6467
use codex_analytics::InvocationType;
6568
use codex_analytics::TurnResolvedConfigFact;
6669
use codex_analytics::build_track_events_context;
70+
use codex_analytics::now_unix_seconds;
6771
use codex_async_utils::OrCancelExt;
6872
use codex_features::Feature;
6973
use codex_hooks::HookEvent;
@@ -91,6 +95,7 @@ use codex_protocol::protocol::EventMsg;
9195
use codex_protocol::protocol::PlanDeltaEvent;
9296
use codex_protocol::protocol::ReasoningContentDeltaEvent;
9397
use codex_protocol::protocol::ReasoningRawContentDeltaEvent;
98+
use codex_protocol::protocol::TokenUsage;
9499
use codex_protocol::protocol::TurnDiffEvent;
95100
use codex_protocol::protocol::WarningEvent;
96101
use codex_protocol::user_input::UserInput;
@@ -406,6 +411,7 @@ pub(crate) async fn run_turn(
406411
// 1. At the start of a turn, so the fresh user prompt in `input` gets sampled first.
407412
// 2. After auto-compact, when model/tool continuation needs to resume before any steer.
408413
let mut can_drain_pending_input = input.is_empty();
414+
let mut next_turn_responses_call_index: u64 = 0;
409415

410416
loop {
411417
if run_pending_session_start_hooks(&sess, &turn_context).await {
@@ -483,6 +489,7 @@ pub(crate) async fn run_turn(
483489
Arc::clone(&turn_diff_tracker),
484490
&mut client_session,
485491
turn_metadata_header.as_deref(),
492+
&mut next_turn_responses_call_index,
486493
sampling_request_input,
487494
&explicitly_enabled_connectors,
488495
skills_outcome,
@@ -1032,6 +1039,64 @@ fn filter_deferred_dynamic_tool_spec(
10321039
}
10331040
}
10341041

1042+
struct ResponsesApiCallAttempt {
1043+
turn_responses_call_index: u64,
1044+
started_at: u64,
1045+
started_instant: Instant,
1046+
input_items: Vec<ResponseItem>,
1047+
output_items: Vec<ResponseItem>,
1048+
responses_id: Option<String>,
1049+
token_usage: Option<TokenUsage>,
1050+
}
1051+
1052+
impl ResponsesApiCallAttempt {
1053+
fn new(turn_responses_call_index: u64, input_items: Vec<ResponseItem>) -> Self {
1054+
Self {
1055+
turn_responses_call_index,
1056+
started_at: now_unix_seconds(),
1057+
started_instant: Instant::now(),
1058+
input_items,
1059+
output_items: Vec::new(),
1060+
responses_id: None,
1061+
token_usage: None,
1062+
}
1063+
}
1064+
}
1065+
1066+
fn track_responses_api_call_attempt(
1067+
sess: &Session,
1068+
turn_context: &TurnContext,
1069+
attempt: ResponsesApiCallAttempt,
1070+
status: CodexResponsesApiCallStatus,
1071+
error: Option<String>,
1072+
) {
1073+
if !sess.enabled(Feature::GeneralAnalytics) {
1074+
return;
1075+
}
1076+
let input = CodexResponsesApiCallInput {
1077+
responses_id: attempt.responses_id,
1078+
turn_responses_call_index: attempt.turn_responses_call_index,
1079+
status,
1080+
error,
1081+
started_at: attempt.started_at,
1082+
completed_at: Some(now_unix_seconds()),
1083+
duration_ms: Some(attempt.started_instant.elapsed().as_millis() as u64),
1084+
input_items: attempt.input_items,
1085+
output_items: attempt.output_items,
1086+
token_usage: attempt.token_usage,
1087+
};
1088+
sess.services
1089+
.analytics_events_client
1090+
.track_responses_api_call(
1091+
build_track_events_context(
1092+
turn_context.model_info.slug.clone(),
1093+
sess.conversation_id.to_string(),
1094+
turn_context.sub_id.clone(),
1095+
),
1096+
input,
1097+
);
1098+
}
1099+
10351100
#[allow(clippy::too_many_arguments)]
10361101
#[instrument(level = "trace",
10371102
skip_all,
@@ -1047,6 +1112,7 @@ async fn run_sampling_request(
10471112
turn_diff_tracker: SharedTurnDiffTracker,
10481113
client_session: &mut ModelClientSession,
10491114
turn_metadata_header: Option<&str>,
1115+
next_turn_responses_call_index: &mut u64,
10501116
input: Vec<ResponseItem>,
10511117
explicitly_enabled_connectors: &HashSet<String>,
10521118
skills_outcome: Option<&SkillLoadOutcome>,
@@ -1097,6 +1163,10 @@ async fn run_sampling_request(
10971163
turn_context.as_ref(),
10981164
base_instructions.clone(),
10991165
);
1166+
let turn_responses_call_index = *next_turn_responses_call_index;
1167+
*next_turn_responses_call_index = (*next_turn_responses_call_index).saturating_add(1);
1168+
let mut responses_api_call_attempt =
1169+
ResponsesApiCallAttempt::new(turn_responses_call_index, prompt.input.clone());
11001170
let err = match try_run_sampling_request(
11011171
tool_runtime.clone(),
11021172
Arc::clone(&sess),
@@ -1106,28 +1176,63 @@ async fn run_sampling_request(
11061176
Arc::clone(&turn_diff_tracker),
11071177
server_model_warning_emitted_for_turn,
11081178
&prompt,
1179+
&mut responses_api_call_attempt,
11091180
cancellation_token.child_token(),
11101181
)
11111182
.await
11121183
{
11131184
Ok(output) => {
1185+
track_responses_api_call_attempt(
1186+
sess.as_ref(),
1187+
turn_context.as_ref(),
1188+
responses_api_call_attempt,
1189+
CodexResponsesApiCallStatus::Completed,
1190+
/*error*/ None,
1191+
);
11141192
return Ok(output);
11151193
}
11161194
Err(CodexErr::ContextWindowExceeded) => {
11171195
sess.set_total_tokens_full(&turn_context).await;
1196+
track_responses_api_call_attempt(
1197+
sess.as_ref(),
1198+
turn_context.as_ref(),
1199+
responses_api_call_attempt,
1200+
CodexResponsesApiCallStatus::Failed,
1201+
Some("context window exceeded".to_string()),
1202+
);
11181203
return Err(CodexErr::ContextWindowExceeded);
11191204
}
11201205
Err(CodexErr::UsageLimitReached(e)) => {
11211206
let rate_limits = e.rate_limits.clone();
11221207
if let Some(rate_limits) = rate_limits {
11231208
sess.update_rate_limits(&turn_context, *rate_limits).await;
11241209
}
1210+
track_responses_api_call_attempt(
1211+
sess.as_ref(),
1212+
turn_context.as_ref(),
1213+
responses_api_call_attempt,
1214+
CodexResponsesApiCallStatus::Failed,
1215+
Some(e.to_string()),
1216+
);
11251217
return Err(CodexErr::UsageLimitReached(e));
11261218
}
11271219
Err(err) => err,
11281220
};
11291221

11301222
if !err.is_retryable() {
1223+
let status = if matches!(err, CodexErr::TurnAborted) {
1224+
CodexResponsesApiCallStatus::Interrupted
1225+
} else {
1226+
CodexResponsesApiCallStatus::Failed
1227+
};
1228+
let error = Some(format!("{err:#}"));
1229+
track_responses_api_call_attempt(
1230+
sess.as_ref(),
1231+
turn_context.as_ref(),
1232+
responses_api_call_attempt,
1233+
status,
1234+
error,
1235+
);
11311236
return Err(err);
11321237
}
11331238

@@ -1139,6 +1244,14 @@ async fn run_sampling_request(
11391244
&turn_context.model_info,
11401245
)
11411246
{
1247+
let error = Some(format!("{err:#}"));
1248+
track_responses_api_call_attempt(
1249+
sess.as_ref(),
1250+
turn_context.as_ref(),
1251+
responses_api_call_attempt,
1252+
CodexResponsesApiCallStatus::Failed,
1253+
error,
1254+
);
11421255
sess.send_event(
11431256
&turn_context,
11441257
EventMsg::Warning(WarningEvent {
@@ -1160,6 +1273,14 @@ async fn run_sampling_request(
11601273
warn!(
11611274
"stream disconnected - retrying sampling request ({retries}/{max_retries} in {delay:?})...",
11621275
);
1276+
let error = Some(format!("{err:#}"));
1277+
track_responses_api_call_attempt(
1278+
sess.as_ref(),
1279+
turn_context.as_ref(),
1280+
responses_api_call_attempt,
1281+
CodexResponsesApiCallStatus::Failed,
1282+
error,
1283+
);
11631284

11641285
// In release builds, hide the first websocket retry notification to reduce noisy
11651286
// transient reconnect messages. In debug builds, keep full visibility for diagnosis.
@@ -1179,6 +1300,14 @@ async fn run_sampling_request(
11791300
}
11801301
tokio::time::sleep(delay).await;
11811302
} else {
1303+
let error = Some(format!("{err:#}"));
1304+
track_responses_api_call_attempt(
1305+
sess.as_ref(),
1306+
turn_context.as_ref(),
1307+
responses_api_call_attempt,
1308+
CodexResponsesApiCallStatus::Failed,
1309+
error,
1310+
);
11821311
return Err(err);
11831312
}
11841313
}
@@ -1895,6 +2024,7 @@ async fn try_run_sampling_request(
18952024
turn_diff_tracker: SharedTurnDiffTracker,
18962025
server_model_warning_emitted_for_turn: &mut bool,
18972026
prompt: &Prompt,
2027+
responses_api_call_attempt: &mut ResponsesApiCallAttempt,
18982028
cancellation_token: CancellationToken,
18992029
) -> CodexResult<SamplingRequestResult> {
19002030
feedback_tags!(
@@ -1970,6 +2100,7 @@ async fn try_run_sampling_request(
19702100
match event {
19712101
ResponseEvent::Created => {}
19722102
ResponseEvent::OutputItemDone(item) => {
2103+
responses_api_call_attempt.output_items.push(item.clone());
19732104
if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take()
19742105
&& let Some(event) = consumer.flush_on_complete()
19752106
{
@@ -2140,9 +2271,11 @@ async fn try_run_sampling_request(
21402271
sess.services.models_manager.refresh_if_new_etag(etag).await;
21412272
}
21422273
ResponseEvent::Completed {
2143-
response_id: _,
2274+
response_id,
21442275
token_usage,
21452276
} => {
2277+
responses_api_call_attempt.responses_id = Some(response_id);
2278+
responses_api_call_attempt.token_usage = token_usage.clone();
21462279
flush_assistant_text_segments_all(
21472280
&sess,
21482281
&turn_context,

0 commit comments

Comments
 (0)