Skip to content

Commit 233d677

Browse files
committed
core: emit responses api call analytics
1 parent e0fb689 commit 233d677

2 files changed

Lines changed: 133 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: 118 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 {
@@ -477,12 +483,15 @@ pub(crate) async fn run_turn(
477483
.map(|user_message| user_message.message())
478484
.collect::<Vec<String>>();
479485
let turn_metadata_header = turn_context.turn_metadata_state.current_header_value();
486+
let turn_responses_call_index = next_turn_responses_call_index;
487+
next_turn_responses_call_index = next_turn_responses_call_index.saturating_add(1);
480488
match run_sampling_request(
481489
Arc::clone(&sess),
482490
Arc::clone(&turn_context),
483491
Arc::clone(&turn_diff_tracker),
484492
&mut client_session,
485493
turn_metadata_header.as_deref(),
494+
turn_responses_call_index,
486495
sampling_request_input,
487496
&explicitly_enabled_connectors,
488497
skills_outcome,
@@ -1032,6 +1041,64 @@ fn filter_deferred_dynamic_tool_spec(
10321041
}
10331042
}
10341043

1044+
struct ResponsesApiCallAttempt {
1045+
turn_responses_call_index: u64,
1046+
started_at: u64,
1047+
started_instant: Instant,
1048+
input_items: Vec<ResponseItem>,
1049+
output_items: Vec<ResponseItem>,
1050+
responses_id: Option<String>,
1051+
token_usage: Option<TokenUsage>,
1052+
}
1053+
1054+
impl ResponsesApiCallAttempt {
1055+
fn new(turn_responses_call_index: u64, input_items: Vec<ResponseItem>) -> Self {
1056+
Self {
1057+
turn_responses_call_index,
1058+
started_at: now_unix_seconds(),
1059+
started_instant: Instant::now(),
1060+
input_items,
1061+
output_items: Vec::new(),
1062+
responses_id: None,
1063+
token_usage: None,
1064+
}
1065+
}
1066+
}
1067+
1068+
fn track_responses_api_call_attempt(
1069+
sess: &Session,
1070+
turn_context: &TurnContext,
1071+
attempt: ResponsesApiCallAttempt,
1072+
status: CodexResponsesApiCallStatus,
1073+
error: Option<String>,
1074+
) {
1075+
if !sess.enabled(Feature::GeneralAnalytics) {
1076+
return;
1077+
}
1078+
let input = CodexResponsesApiCallInput {
1079+
responses_id: attempt.responses_id,
1080+
turn_responses_call_index: attempt.turn_responses_call_index,
1081+
status,
1082+
error,
1083+
started_at: attempt.started_at,
1084+
completed_at: Some(now_unix_seconds()),
1085+
duration_ms: Some(attempt.started_instant.elapsed().as_millis() as u64),
1086+
input_items: attempt.input_items,
1087+
output_items: attempt.output_items,
1088+
token_usage: attempt.token_usage,
1089+
};
1090+
sess.services
1091+
.analytics_events_client
1092+
.track_responses_api_call(
1093+
build_track_events_context(
1094+
turn_context.model_info.slug.clone(),
1095+
sess.conversation_id.to_string(),
1096+
turn_context.sub_id.clone(),
1097+
),
1098+
input,
1099+
);
1100+
}
1101+
10351102
#[allow(clippy::too_many_arguments)]
10361103
#[instrument(level = "trace",
10371104
skip_all,
@@ -1047,6 +1114,7 @@ async fn run_sampling_request(
10471114
turn_diff_tracker: SharedTurnDiffTracker,
10481115
client_session: &mut ModelClientSession,
10491116
turn_metadata_header: Option<&str>,
1117+
turn_responses_call_index: u64,
10501118
input: Vec<ResponseItem>,
10511119
explicitly_enabled_connectors: &HashSet<String>,
10521120
skills_outcome: Option<&SkillLoadOutcome>,
@@ -1097,6 +1165,8 @@ async fn run_sampling_request(
10971165
turn_context.as_ref(),
10981166
base_instructions.clone(),
10991167
);
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

@@ -1179,6 +1284,14 @@ async fn run_sampling_request(
11791284
}
11801285
tokio::time::sleep(delay).await;
11811286
} else {
1287+
let error = Some(format!("{err:#}"));
1288+
track_responses_api_call_attempt(
1289+
sess.as_ref(),
1290+
turn_context.as_ref(),
1291+
responses_api_call_attempt,
1292+
CodexResponsesApiCallStatus::Failed,
1293+
error,
1294+
);
11821295
return Err(err);
11831296
}
11841297
}
@@ -1895,6 +2008,7 @@ async fn try_run_sampling_request(
18952008
turn_diff_tracker: SharedTurnDiffTracker,
18962009
server_model_warning_emitted_for_turn: &mut bool,
18972010
prompt: &Prompt,
2011+
responses_api_call_attempt: &mut ResponsesApiCallAttempt,
18982012
cancellation_token: CancellationToken,
18992013
) -> CodexResult<SamplingRequestResult> {
19002014
feedback_tags!(
@@ -1970,6 +2084,7 @@ async fn try_run_sampling_request(
19702084
match event {
19712085
ResponseEvent::Created => {}
19722086
ResponseEvent::OutputItemDone(item) => {
2087+
responses_api_call_attempt.output_items.push(item.clone());
19732088
if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take()
19742089
&& let Some(event) = consumer.flush_on_complete()
19752090
{
@@ -2140,9 +2255,11 @@ async fn try_run_sampling_request(
21402255
sess.services.models_manager.refresh_if_new_etag(etag).await;
21412256
}
21422257
ResponseEvent::Completed {
2143-
response_id: _,
2258+
response_id,
21442259
token_usage,
21452260
} => {
2261+
responses_api_call_attempt.responses_id = Some(response_id);
2262+
responses_api_call_attempt.token_usage = token_usage.clone();
21462263
flush_assistant_text_segments_all(
21472264
&sess,
21482265
&turn_context,

0 commit comments

Comments
 (0)