diff --git a/codex-rs/core/src/compact.rs b/codex-rs/core/src/compact.rs index 4ff35c09e..4dc4c34e9 100644 --- a/codex-rs/core/src/compact.rs +++ b/codex-rs/core/src/compact.rs @@ -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, + active_context_tokens_before: Option, ) { + 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(), diff --git a/codex-rs/core/src/compact_remote.rs b/codex-rs/core/src/compact_remote.rs index cbac86fb9..444b1c3a6 100644 --- a/codex-rs/core/src/compact_remote.rs +++ b/codex-rs/core/src/compact_remote.rs @@ -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, 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 { diff --git a/codex-rs/core/src/compact_remote_v2.rs b/codex-rs/core/src/compact_remote_v2.rs index 77ef1629c..ac64f7f0f 100644 --- a/codex-rs/core/src/compact_remote_v2.rs +++ b/codex-rs/core/src/compact_remote_v2.rs @@ -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, +} + 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 { 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 { 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, + }) + ); } }