Skip to content

Commit b804793

Browse files
committed
core: emit responses api call analytics
1 parent 8e29b28 commit b804793

2 files changed

Lines changed: 168 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: 153 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,9 @@
11
use std::collections::HashMap;
22
use std::collections::HashSet;
33
use std::sync::Arc;
4+
use std::time::Instant;
5+
use std::time::SystemTime;
6+
use std::time::UNIX_EPOCH;
47

58
use crate::SkillInjections;
69
use crate::SkillLoadOutcome;
@@ -58,6 +61,8 @@ use crate::unavailable_tool::collect_unavailable_called_tools;
5861
use crate::util::backoff;
5962
use crate::util::error_or_panic;
6063
use codex_analytics::AppInvocation;
64+
use codex_analytics::CodexResponsesApiCallInput;
65+
use codex_analytics::CodexResponsesApiCallStatus;
6166
use codex_analytics::CompactionPhase;
6267
use codex_analytics::CompactionReason;
6368
use codex_analytics::InvocationType;
@@ -90,6 +95,7 @@ use codex_protocol::protocol::EventMsg;
9095
use codex_protocol::protocol::PlanDeltaEvent;
9196
use codex_protocol::protocol::ReasoningContentDeltaEvent;
9297
use codex_protocol::protocol::ReasoningRawContentDeltaEvent;
98+
use codex_protocol::protocol::TokenUsage;
9399
use codex_protocol::protocol::TurnDiffEvent;
94100
use codex_protocol::protocol::WarningEvent;
95101
use codex_protocol::user_input::UserInput;
@@ -394,6 +400,7 @@ pub(crate) async fn run_turn(
394400
// 1. At the start of a turn, so the fresh user prompt in `input` gets sampled first.
395401
// 2. After auto-compact, when model/tool continuation needs to resume before any steer.
396402
let mut can_drain_pending_input = input.is_empty();
403+
let mut next_turn_responses_call_index: u64 = 0;
397404

398405
loop {
399406
if run_pending_session_start_hooks(&sess, &turn_context).await {
@@ -465,12 +472,15 @@ pub(crate) async fn run_turn(
465472
.map(|user_message| user_message.message())
466473
.collect::<Vec<String>>();
467474
let turn_metadata_header = turn_context.turn_metadata_state.current_header_value();
475+
let turn_responses_call_index = next_turn_responses_call_index;
476+
next_turn_responses_call_index = next_turn_responses_call_index.saturating_add(1);
468477
match run_sampling_request(
469478
Arc::clone(&sess),
470479
Arc::clone(&turn_context),
471480
Arc::clone(&turn_diff_tracker),
472481
&mut client_session,
473482
turn_metadata_header.as_deref(),
483+
turn_responses_call_index,
474484
sampling_request_input,
475485
&explicitly_enabled_connectors,
476486
skills_outcome,
@@ -991,6 +1001,102 @@ pub(crate) fn build_prompt(
9911001
}
9921002
}
9931003

1004+
struct ResponsesApiCallObservation {
1005+
turn_responses_call_index: u64,
1006+
started_at: u64,
1007+
started_instant: Instant,
1008+
input_items: Vec<ResponseItem>,
1009+
output_items: Vec<ResponseItem>,
1010+
responses_id: Option<String>,
1011+
token_usage: Option<TokenUsage>,
1012+
}
1013+
1014+
impl ResponsesApiCallObservation {
1015+
fn new(turn_responses_call_index: u64) -> Self {
1016+
Self {
1017+
turn_responses_call_index,
1018+
started_at: current_unix_seconds(),
1019+
started_instant: Instant::now(),
1020+
input_items: Vec::new(),
1021+
output_items: Vec::new(),
1022+
responses_id: None,
1023+
token_usage: None,
1024+
}
1025+
}
1026+
1027+
fn record_prompt(&mut self, prompt: &Prompt) {
1028+
self.input_items = prompt.input.clone();
1029+
self.output_items.clear();
1030+
self.responses_id = None;
1031+
self.token_usage = None;
1032+
}
1033+
1034+
fn record_output_item_done(&mut self, item: &ResponseItem) {
1035+
self.output_items.push(item.clone());
1036+
}
1037+
1038+
fn record_completed(&mut self, response_id: String, token_usage: Option<TokenUsage>) {
1039+
self.responses_id = Some(response_id);
1040+
self.token_usage = token_usage;
1041+
}
1042+
1043+
fn into_input(
1044+
self,
1045+
status: CodexResponsesApiCallStatus,
1046+
error: Option<String>,
1047+
) -> CodexResponsesApiCallInput {
1048+
let completed_at = current_unix_seconds();
1049+
CodexResponsesApiCallInput {
1050+
responses_id: self.responses_id,
1051+
turn_responses_call_index: self.turn_responses_call_index,
1052+
status,
1053+
error,
1054+
started_at: self.started_at,
1055+
completed_at: Some(completed_at),
1056+
duration_ms: Some(self.started_instant.elapsed().as_millis() as u64),
1057+
input_items: self.input_items,
1058+
output_items: self.output_items,
1059+
token_usage: self.token_usage,
1060+
}
1061+
}
1062+
}
1063+
1064+
fn current_unix_seconds() -> u64 {
1065+
SystemTime::now()
1066+
.duration_since(UNIX_EPOCH)
1067+
.unwrap_or_default()
1068+
.as_secs()
1069+
}
1070+
1071+
fn responses_api_call_error_status(err: &CodexErr) -> CodexResponsesApiCallStatus {
1072+
match err {
1073+
CodexErr::TurnAborted => CodexResponsesApiCallStatus::Interrupted,
1074+
_ => CodexResponsesApiCallStatus::Failed,
1075+
}
1076+
}
1077+
1078+
fn track_responses_api_call_observation(
1079+
sess: &Session,
1080+
turn_context: &TurnContext,
1081+
observation: ResponsesApiCallObservation,
1082+
status: CodexResponsesApiCallStatus,
1083+
error: Option<String>,
1084+
) {
1085+
if !sess.enabled(Feature::GeneralAnalytics) {
1086+
return;
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+
observation.into_input(status, error),
1097+
);
1098+
}
1099+
9941100
#[allow(clippy::too_many_arguments)]
9951101
#[instrument(level = "trace",
9961102
skip_all,
@@ -1006,6 +1112,7 @@ async fn run_sampling_request(
10061112
turn_diff_tracker: SharedTurnDiffTracker,
10071113
client_session: &mut ModelClientSession,
10081114
turn_metadata_header: Option<&str>,
1115+
turn_responses_call_index: u64,
10091116
input: Vec<ResponseItem>,
10101117
explicitly_enabled_connectors: &HashSet<String>,
10111118
skills_outcome: Option<&SkillLoadOutcome>,
@@ -1042,6 +1149,8 @@ async fn run_sampling_request(
10421149
.await;
10431150
let mut retries = 0;
10441151
let mut initial_input = Some(input);
1152+
let mut responses_api_call_observation =
1153+
ResponsesApiCallObservation::new(turn_responses_call_index);
10451154
loop {
10461155
let prompt_input = if let Some(input) = initial_input.take() {
10471156
input
@@ -1056,6 +1165,7 @@ async fn run_sampling_request(
10561165
turn_context.as_ref(),
10571166
base_instructions.clone(),
10581167
);
1168+
responses_api_call_observation.record_prompt(&prompt);
10591169
let err = match try_run_sampling_request(
10601170
tool_runtime.clone(),
10611171
Arc::clone(&sess),
@@ -1065,28 +1175,59 @@ async fn run_sampling_request(
10651175
Arc::clone(&turn_diff_tracker),
10661176
server_model_warning_emitted_for_turn,
10671177
&prompt,
1178+
&mut responses_api_call_observation,
10681179
cancellation_token.child_token(),
10691180
)
10701181
.await
10711182
{
10721183
Ok(output) => {
1184+
track_responses_api_call_observation(
1185+
sess.as_ref(),
1186+
turn_context.as_ref(),
1187+
responses_api_call_observation,
1188+
CodexResponsesApiCallStatus::Completed,
1189+
/*error*/ None,
1190+
);
10731191
return Ok(output);
10741192
}
10751193
Err(CodexErr::ContextWindowExceeded) => {
10761194
sess.set_total_tokens_full(&turn_context).await;
1195+
track_responses_api_call_observation(
1196+
sess.as_ref(),
1197+
turn_context.as_ref(),
1198+
responses_api_call_observation,
1199+
CodexResponsesApiCallStatus::Failed,
1200+
Some("context window exceeded".to_string()),
1201+
);
10771202
return Err(CodexErr::ContextWindowExceeded);
10781203
}
10791204
Err(CodexErr::UsageLimitReached(e)) => {
10801205
let rate_limits = e.rate_limits.clone();
10811206
if let Some(rate_limits) = rate_limits {
10821207
sess.update_rate_limits(&turn_context, *rate_limits).await;
10831208
}
1209+
track_responses_api_call_observation(
1210+
sess.as_ref(),
1211+
turn_context.as_ref(),
1212+
responses_api_call_observation,
1213+
CodexResponsesApiCallStatus::Failed,
1214+
Some(e.to_string()),
1215+
);
10841216
return Err(CodexErr::UsageLimitReached(e));
10851217
}
10861218
Err(err) => err,
10871219
};
10881220

10891221
if !err.is_retryable() {
1222+
let status = responses_api_call_error_status(&err);
1223+
let error = Some(format!("{err:#}"));
1224+
track_responses_api_call_observation(
1225+
sess.as_ref(),
1226+
turn_context.as_ref(),
1227+
responses_api_call_observation,
1228+
status,
1229+
error,
1230+
);
10901231
return Err(err);
10911232
}
10921233

@@ -1138,6 +1279,14 @@ async fn run_sampling_request(
11381279
}
11391280
tokio::time::sleep(delay).await;
11401281
} else {
1282+
let error = Some(format!("{err:#}"));
1283+
track_responses_api_call_observation(
1284+
sess.as_ref(),
1285+
turn_context.as_ref(),
1286+
responses_api_call_observation,
1287+
CodexResponsesApiCallStatus::Failed,
1288+
error,
1289+
);
11411290
return Err(err);
11421291
}
11431292
}
@@ -1850,6 +1999,7 @@ async fn try_run_sampling_request(
18501999
turn_diff_tracker: SharedTurnDiffTracker,
18512000
server_model_warning_emitted_for_turn: &mut bool,
18522001
prompt: &Prompt,
2002+
responses_api_call_observation: &mut ResponsesApiCallObservation,
18532003
cancellation_token: CancellationToken,
18542004
) -> CodexResult<SamplingRequestResult> {
18552005
feedback_tags!(
@@ -1925,6 +2075,7 @@ async fn try_run_sampling_request(
19252075
match event {
19262076
ResponseEvent::Created => {}
19272077
ResponseEvent::OutputItemDone(item) => {
2078+
responses_api_call_observation.record_output_item_done(&item);
19282079
if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take()
19292080
&& let Some(event) = consumer.flush_on_complete()
19302081
{
@@ -2095,9 +2246,10 @@ async fn try_run_sampling_request(
20952246
sess.services.models_manager.refresh_if_new_etag(etag).await;
20962247
}
20972248
ResponseEvent::Completed {
2098-
response_id: _,
2249+
response_id,
20992250
token_usage,
21002251
} => {
2252+
responses_api_call_observation.record_completed(response_id, token_usage.clone());
21012253
flush_assistant_text_segments_all(
21022254
&sess,
21032255
&turn_context,

0 commit comments

Comments
 (0)