[codex-analytics] report compaction request token counts (#25946)

## Why

Compaction analytics need token counts that better represent the request
being compacted. The existing session snapshot can diverge from the
actual remote compaction request after output rewriting, and remote v2
can use server-side Responses usage when available.

## What changed

- Add an optional `active_context_tokens_before` override to
`CompactionAnalyticsAttempt::track(...)` for remote compaction when it
has a better before-token value than the begin-time session snapshot.
The local `/compact` path passes no override.
- For remote v1 `responses_compact`, subtract the estimated token delta
from pre-compaction output rewriting from the session snapshot, capped
by locally-added tokens since the last successful API response.
- For remote v2 `responses_compaction_v2`, use the same bounded
output-rewrite fallback as remote v1, then overwrite
`active_context_tokens_before` with server `token_usage.input_tokens`
from the `response.completed` event when present.
- Keep the existing v2 compaction-output validation while carrying the
completed response token usage through `collect_compaction_output`.

## Verification

- `just fmt`
- `just test -p codex-core
collect_compaction_output_accepts_additional_output_items`
- `git diff --check`
This commit is contained in:
rhan-oai
2026-06-03 20:00:44 -07:00
committed by GitHub
Unverified
parent 6bcccb0ee6
commit c143a86de8
3 changed files with 159 additions and 35 deletions
+26 -4
View File
@@ -146,7 +146,12 @@ async fn run_compact_task_inner(
PreCompactHookOutcome::Stopped { reason } => {
let error = reason.unwrap_or_else(|| "PreCompact hook stopped execution".to_string());
attempt
.track(sess.as_ref(), CompactionStatus::Interrupted, Some(error))
.track(
sess.as_ref(),
CompactionStatus::Interrupted,
Some(error),
/*active_context_tokens_before*/ None,
)
.await;
return Err(CodexErr::TurnAborted);
}
@@ -164,11 +169,25 @@ async fn run_compact_task_inner(
if result.is_ok() {
let post_compact_outcome = run_post_compact_hooks(&sess, &turn_context, trigger).await;
if let PostCompactHookOutcome::Stopped = post_compact_outcome {
attempt.track(sess.as_ref(), status, error).await;
attempt
.track(
sess.as_ref(),
status,
error,
/*active_context_tokens_before*/ None,
)
.await;
return Err(CodexErr::TurnAborted);
}
}
attempt.track(sess.as_ref(), status, error).await;
attempt
.track(
sess.as_ref(),
status,
error,
/*active_context_tokens_before*/ None,
)
.await;
result.map(|_| ())
}
@@ -344,7 +363,10 @@ impl CompactionAnalyticsAttempt {
sess: &Session,
status: CompactionStatus,
error: Option<String>,
active_context_tokens_before: Option<i64>,
) {
let active_context_tokens_before =
active_context_tokens_before.unwrap_or(self.active_context_tokens_before);
let active_context_tokens_after = sess.get_total_token_usage().await;
sess.services
.analytics_events_client
@@ -358,7 +380,7 @@ impl CompactionAnalyticsAttempt {
strategy: CompactionStrategy::Memento,
status,
error,
active_context_tokens_before: self.active_context_tokens_before,
active_context_tokens_before,
active_context_tokens_after,
started_at: self.started_at,
completed_at: now_unix_seconds(),
+49 -14
View File
@@ -100,6 +100,7 @@ async fn run_remote_compact_task_inner(
CompactionImplementation::ResponsesCompact,
phase,
);
let mut active_context_tokens_before = sess.get_total_token_usage().await;
let attempt = CompactionAnalyticsAttempt::begin(
sess.as_ref(),
turn_context.as_ref(),
@@ -119,6 +120,7 @@ async fn run_remote_compact_task_inner(
sess.as_ref(),
codex_analytics::CompactionStatus::Interrupted,
Some(error),
Some(active_context_tokens_before),
)
.await;
return Err(CodexErr::TurnAborted);
@@ -129,6 +131,7 @@ async fn run_remote_compact_task_inner(
turn_context,
initial_context_injection,
compaction_metadata,
&mut active_context_tokens_before,
)
.await;
let status = compaction_status_from_result(&result);
@@ -136,11 +139,25 @@ async fn run_remote_compact_task_inner(
if result.is_ok() {
let post_compact_outcome = run_post_compact_hooks(sess, turn_context, trigger).await;
if let PostCompactHookOutcome::Stopped = post_compact_outcome {
attempt.track(sess.as_ref(), status, error).await;
attempt
.track(
sess.as_ref(),
status,
error,
Some(active_context_tokens_before),
)
.await;
return Err(CodexErr::TurnAborted);
}
}
attempt.track(sess.as_ref(), status, error.clone()).await;
attempt
.track(
sess.as_ref(),
status,
error.clone(),
Some(active_context_tokens_before),
)
.await;
if let Err(err) = result {
sess.track_turn_codex_error(turn_context, &err);
let event = EventMsg::Error(
@@ -157,6 +174,7 @@ async fn run_remote_compact_task_inner_impl(
turn_context: &Arc<TurnContext>,
initial_context_injection: InitialContextInjection,
compaction_metadata: CompactionTurnMetadata,
active_context_tokens_before: &mut i64,
) -> CodexResult<()> {
let context_compaction_item = ContextCompactionItem::new();
// Use the UI compaction item ID as the trace compaction ID so protocol lifecycle events,
@@ -172,11 +190,12 @@ async fn run_remote_compact_task_inner_impl(
.await;
let mut history = sess.clone_history().await;
let base_instructions = sess.get_base_instructions().await;
let rewritten_outputs = trim_function_call_history_to_fit_context_window(
&mut history,
turn_context.as_ref(),
&base_instructions,
);
let (rewritten_outputs, estimated_deleted_tokens) =
trim_function_call_history_to_fit_context_window(
&mut history,
turn_context.as_ref(),
&base_instructions,
);
if rewritten_outputs > 0 {
info!(
turn_id = %turn_context.sub_id,
@@ -184,6 +203,14 @@ async fn run_remote_compact_task_inner_impl(
"rewrote history outputs before remote compaction"
);
}
if estimated_deleted_tokens > 0 {
let max_local_deleted_tokens = sess
.get_total_token_usage_breakdown()
.await
.estimated_tokens_of_items_added_since_last_successful_api_response;
*active_context_tokens_before = (*active_context_tokens_before)
.saturating_sub(estimated_deleted_tokens.min(max_local_deleted_tokens));
}
// This is the history selected for remote compaction, after any output rewriting required to
// fit the compact endpoint. The checkpoint below records it separately from the next sampling
// request, whose prompt will repeat current developer/context prefix items.
@@ -382,18 +409,21 @@ pub(crate) fn trim_function_call_history_to_fit_context_window(
history: &mut ContextManager,
turn_context: &TurnContext,
base_instructions: &BaseInstructions,
) -> usize {
) -> (usize, i64) {
let Some(context_window) = turn_context.model_context_window() else {
return 0;
return (0, 0);
};
let mut rewritten_outputs = 0usize;
let mut estimated_deleted_tokens = 0i64;
let item_count = history.raw_items().len();
for index in (0..item_count).rev() {
if history
.estimate_token_count_with_base_instructions(base_instructions)
.is_none_or(|estimated_tokens| estimated_tokens <= context_window)
{
let Some(estimated_tokens_before) =
history.estimate_token_count_with_base_instructions(base_instructions)
else {
break;
};
if estimated_tokens_before <= context_window {
break;
}
let Some(rewritten_item) = history
@@ -406,10 +436,15 @@ pub(crate) fn trim_function_call_history_to_fit_context_window(
let mut items = history.raw_items().to_vec();
items[index] = rewritten_item;
history.replace(items);
let estimated_tokens_after = history
.estimate_token_count_with_base_instructions(base_instructions)
.unwrap_or_default();
rewritten_outputs += 1;
estimated_deleted_tokens = estimated_deleted_tokens
.saturating_add(estimated_tokens_before.saturating_sub(estimated_tokens_after));
}
rewritten_outputs
(rewritten_outputs, estimated_deleted_tokens)
}
fn rewritten_output_for_context_window(item: &ResponseItem) -> Option<ResponseItem> {
+84 -17
View File
@@ -35,6 +35,7 @@ use codex_protocol::models::ContentItem;
use codex_protocol::models::ResponseItem;
use codex_protocol::protocol::CompactedItem;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::TokenUsage;
use codex_protocol::protocol::TruncationPolicy;
use codex_protocol::protocol::TurnStartedEvent;
use codex_rollout_trace::CompactionCheckpointTracePayload;
@@ -112,6 +113,7 @@ async fn run_remote_compact_task_inner(
CompactionImplementation::ResponsesCompactionV2,
phase,
);
let mut active_context_tokens_before = sess.get_total_token_usage().await;
let attempt = CompactionAnalyticsAttempt::begin(
sess.as_ref(),
turn_context.as_ref(),
@@ -131,6 +133,7 @@ async fn run_remote_compact_task_inner(
sess.as_ref(),
codex_analytics::CompactionStatus::Interrupted,
Some(error),
Some(active_context_tokens_before),
)
.await;
return Err(CodexErr::TurnAborted);
@@ -142,6 +145,7 @@ async fn run_remote_compact_task_inner(
client_session,
initial_context_injection,
compaction_metadata,
&mut active_context_tokens_before,
)
.await;
let status = compaction_status_from_result(&result);
@@ -149,11 +153,25 @@ async fn run_remote_compact_task_inner(
if result.is_ok() {
let post_compact_outcome = run_post_compact_hooks(sess, turn_context, trigger).await;
if let PostCompactHookOutcome::Stopped = post_compact_outcome {
attempt.track(sess.as_ref(), status, error).await;
attempt
.track(
sess.as_ref(),
status,
error,
Some(active_context_tokens_before),
)
.await;
return Err(CodexErr::TurnAborted);
}
}
attempt.track(sess.as_ref(), status, error.clone()).await;
attempt
.track(
sess.as_ref(),
status,
error.clone(),
Some(active_context_tokens_before),
)
.await;
if let Err(err) = result {
sess.track_turn_codex_error(turn_context, &err);
let event = EventMsg::Error(
@@ -171,6 +189,7 @@ async fn run_remote_compact_task_inner_impl(
client_session: Option<&mut ModelClientSession>,
initial_context_injection: InitialContextInjection,
compaction_metadata: CompactionTurnMetadata,
active_context_tokens_before: &mut i64,
) -> CodexResult<()> {
let context_compaction_item = ContextCompactionItem::new();
let compaction_trace = sess.services.rollout_thread_trace.compaction_trace_context(
@@ -185,11 +204,12 @@ async fn run_remote_compact_task_inner_impl(
let mut history = sess.clone_history().await;
let base_instructions = sess.get_base_instructions().await;
let rewritten_outputs = trim_function_call_history_to_fit_context_window(
&mut history,
turn_context.as_ref(),
&base_instructions,
);
let (rewritten_outputs, estimated_deleted_tokens) =
trim_function_call_history_to_fit_context_window(
&mut history,
turn_context.as_ref(),
&base_instructions,
);
if rewritten_outputs > 0 {
info!(
turn_id = %turn_context.sub_id,
@@ -197,6 +217,14 @@ async fn run_remote_compact_task_inner_impl(
"rewrote history outputs before remote compaction v2"
);
}
if estimated_deleted_tokens > 0 {
let max_local_deleted_tokens = sess
.get_total_token_usage_breakdown()
.await
.estimated_tokens_of_items_added_since_last_successful_api_response;
*active_context_tokens_before = (*active_context_tokens_before)
.saturating_sub(estimated_deleted_tokens.min(max_local_deleted_tokens));
}
let trace_input_history = history.raw_items().to_vec();
let prompt_input = history.for_prompt(&turn_context.model_info.input_modalities);
@@ -249,9 +277,16 @@ async fn run_remote_compact_task_inner_impl(
trace_attempt.record_result(
compaction_output_result
.as_ref()
.map(|(item, _)| std::slice::from_ref(item)),
.map(|output| std::slice::from_ref(&output.compaction_output)),
);
let (compaction_output, response_id) = compaction_output_result?;
let RemoteCompactionV2Output {
compaction_output,
response_id,
token_usage,
} = compaction_output_result?;
if let Some(token_usage) = token_usage {
*active_context_tokens_before = token_usage.input_tokens;
}
let compacted_history = build_v2_compacted_history(&prompt_input, compaction_output);
let new_history = process_compacted_history(
sess.as_ref(),
@@ -288,13 +323,19 @@ async fn run_remote_compact_task_inner_impl(
Ok(())
}
struct RemoteCompactionV2Output {
compaction_output: ResponseItem,
response_id: String,
token_usage: Option<TokenUsage>,
}
async fn run_remote_compaction_request_v2(
sess: &Session,
turn_context: &TurnContext,
client_session: &mut ModelClientSession,
prompt: &Prompt,
turn_metadata_header: Option<&str>,
) -> CodexResult<(ResponseItem, String)> {
) -> CodexResult<RemoteCompactionV2Output> {
let max_retries = turn_context
.provider
.info()
@@ -364,11 +405,12 @@ async fn log_remote_compaction_request_failure(
async fn collect_compaction_output(
mut stream: ResponseStream,
) -> CodexResult<(ResponseItem, String)> {
) -> CodexResult<RemoteCompactionV2Output> {
let mut output_item_count = 0usize;
let mut compaction_count = 0usize;
let mut compaction_output = None;
let mut completed_response_id = None;
let mut completed_token_usage = None;
while let Some(event) = stream.next().await {
match event? {
ResponseEvent::OutputItemDone(item) => {
@@ -380,8 +422,13 @@ async fn collect_compaction_output(
}
}
}
ResponseEvent::Completed { response_id, .. } => {
ResponseEvent::Completed {
response_id,
token_usage,
..
} => {
completed_response_id = Some(response_id);
completed_token_usage = token_usage;
break;
}
_ => {}
@@ -404,7 +451,11 @@ async fn collect_compaction_output(
let Some(compaction_output) = compaction_output else {
unreachable!("compaction output must exist when count is exactly one");
};
Ok((compaction_output, response_id))
Ok(RemoteCompactionV2Output {
compaction_output,
response_id,
token_usage: completed_token_usage,
})
}
fn build_v2_compacted_history(
@@ -737,16 +788,32 @@ mod tests {
Ok(ResponseEvent::OutputItemDone(compaction.clone())),
Ok(ResponseEvent::Completed {
response_id: "resp-compact".to_string(),
token_usage: None,
token_usage: Some(TokenUsage {
input_tokens: 123_456,
cached_input_tokens: 7_890,
output_tokens: 42,
reasoning_output_tokens: 5,
total_tokens: 123_498,
}),
end_turn: Some(true),
}),
]);
let (output, response_id) = collect_compaction_output(stream)
let output = collect_compaction_output(stream)
.await
.expect("compaction should be collected");
assert_eq!(output, compaction);
assert_eq!(response_id, "resp-compact");
assert_eq!(output.compaction_output, compaction);
assert_eq!(output.response_id, "resp-compact");
assert_eq!(
output.token_usage,
Some(TokenUsage {
input_tokens: 123_456,
cached_input_tokens: 7_890,
output_tokens: 42,
reasoning_output_tokens: 5,
total_tokens: 123_498,
})
);
}
}