From 29ddeb71ad6317b41cbb591ec1941ccd971e2e66 Mon Sep 17 00:00:00 2001 From: Cooper Gamble Date: Mon, 9 Mar 2026 02:03:06 +0000 Subject: [PATCH] [codex-core] Fix inline compaction event ordering and token accounting [ci changed_files] Co-authored-by: Codex --- codex-rs/core/src/codex.rs | 11 +- codex-rs/core/src/stream_events_utils.rs | 14 +- codex-rs/core/tests/suite/compact_remote.rs | 194 ++++++++++++++++++++ 3 files changed, 208 insertions(+), 11 deletions(-) diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index fe998edca6..bac18f35c6 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -7368,14 +7368,12 @@ async fn try_run_sampling_request( &mut assistant_message_stream_parsers, ) .await; + let mut applied_server_side_compaction_checkpoint = false; if let Some(PendingServerSideCompactionCheckpoint { history_at_checkpoint, item, - turn_item, }) = pending_server_side_compaction_checkpoint.take() { - let turn_item = TurnItem::ContextCompaction(turn_item); - sess.emit_turn_item_started(&turn_context, &turn_item).await; sess.apply_server_side_compaction_checkpoint( turn_context.as_ref(), item, @@ -7385,11 +7383,12 @@ async fn try_run_sampling_request( history_at_checkpoint.as_slice(), ) .await; - sess.emit_turn_item_completed(&turn_context, turn_item) + applied_server_side_compaction_checkpoint = true; + } + if !applied_server_side_compaction_checkpoint { + sess.update_token_usage_info(&turn_context, token_usage.as_ref()) .await; } - sess.update_token_usage_info(&turn_context, token_usage.as_ref()) - .await; should_emit_turn_diff = true; needs_follow_up |= sess.has_pending_input().await; diff --git a/codex-rs/core/src/stream_events_utils.rs b/codex-rs/core/src/stream_events_utils.rs index 78dc7b997a..b53904ab77 100644 --- a/codex-rs/core/src/stream_events_utils.rs +++ b/codex-rs/core/src/stream_events_utils.rs @@ -152,7 +152,6 @@ pub(crate) struct OutputItemResult { pub(crate) struct PendingServerSideCompactionCheckpoint { pub history_at_checkpoint: Vec, pub item: ResponseItem, - pub turn_item: ContextCompactionItem, } pub(crate) struct HandleOutputCtx { @@ -172,18 +171,23 @@ pub(crate) async fn handle_output_item_done( let plan_mode = ctx.turn_context.collaboration_mode.mode == ModeKind::Plan; if matches!(item, ResponseItem::Compaction { .. }) { - let compaction_item = match previously_active_item { + let turn_item = TurnItem::ContextCompaction(match previously_active_item { Some(TurnItem::ContextCompaction(item)) => item, _ => ContextCompactionItem::new(), - }; + }); debug!( turn_id = %ctx.turn_context.sub_id, - "buffering streamed server-side compaction item until response.completed" + "emitting streamed server-side compaction item and buffering history rewrite until response.completed" ); + ctx.sess + .emit_turn_item_started(&ctx.turn_context, &turn_item) + .await; + ctx.sess + .emit_turn_item_completed(&ctx.turn_context, turn_item) + .await; output.pending_server_side_compaction = Some(PendingServerSideCompactionCheckpoint { history_at_checkpoint: ctx.sess.clone_history().await.raw_items().to_vec(), item, - turn_item: compaction_item, }); return Ok(output); } diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index ae68c37bd9..a2129b29d1 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -546,6 +546,96 @@ async fn auto_server_side_compaction_uses_inline_context_management() -> Result< Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn auto_server_side_compaction_emits_events_before_later_streamed_items() -> Result<()> { + skip_if_no_network!(Ok(())); + + let compact_threshold = 120; + let harness = TestCodexHarness::with_builder( + test_codex() + .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) + .with_config(move |config| { + config + .features + .enable(Feature::ServerSideCompaction) + .expect("enable server-side compaction"); + config.model_auto_compact_token_limit = Some(compact_threshold); + }), + ) + .await?; + let codex = harness.test().codex.clone(); + + let _responses_mock = responses::mount_sse_sequence( + harness.server(), + vec![ + responses::sse(vec![ + responses::ev_assistant_message("m1", "FIRST_REMOTE_REPLY"), + responses::ev_completed_with_tokens("resp-1", 500), + ]), + responses::sse(vec![ + responses::ev_compaction(&summary_with_prefix("INLINE_SERVER_SUMMARY")), + responses::ev_assistant_message("m2", "AFTER_INLINE_REPLY"), + responses::ev_completed_with_tokens("resp-2", 80), + ]), + ], + ) + .await; + + submit_text_turn_and_wait(&codex, "inline compact turn one").await?; + + codex + .submit(Op::UserInput { + items: vec![UserInput::Text { + text: "inline compact turn two".to_string(), + text_elements: Vec::new(), + }], + final_output_json_schema: None, + }) + .await?; + + let mut event_index = 0usize; + let mut context_compacted_index = None; + let mut assistant_completed_index = None; + let mut saw_turn_complete = false; + while !saw_turn_complete + || context_compacted_index.is_none() + || assistant_completed_index.is_none() + { + let event = codex.next_event().await.expect("event"); + match event.msg { + EventMsg::ContextCompacted(_) => { + context_compacted_index = Some(event_index); + } + EventMsg::ItemCompleted(ItemCompletedEvent { + item: TurnItem::AgentMessage(item), + .. + }) if item.content.iter().any(|entry| { + matches!( + entry, + codex_protocol::items::AgentMessageContent::Text { text } + if text == "AFTER_INLINE_REPLY" + ) + }) => + { + assistant_completed_index = Some(event_index); + } + EventMsg::TurnComplete(_) => { + saw_turn_complete = true; + } + _ => {} + } + event_index += 1; + } + + assert!( + context_compacted_index.expect("context compacted event") + < assistant_completed_index.expect("assistant completed event"), + "expected the inline compaction event to arrive before later streamed assistant items" + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn auto_server_side_compaction_keeps_current_turn_inputs_for_follow_ups() -> Result<()> { skip_if_no_network!(Ok(())); @@ -645,6 +735,110 @@ async fn auto_server_side_compaction_keeps_current_turn_inputs_for_follow_ups() Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn auto_server_side_compaction_preserves_recomputed_token_estimate() -> Result<()> { + skip_if_no_network!(Ok(())); + + let compact_threshold = 120; + let inline_summary = summary_with_prefix("INLINE_SERVER_SUMMARY"); + + let harness = TestCodexHarness::with_builder( + test_codex() + .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) + .with_config(move |config| { + config + .features + .enable(Feature::ServerSideCompaction) + .expect("enable server-side compaction"); + config.model_auto_compact_token_limit = Some(compact_threshold); + }), + ) + .await?; + let codex = harness.test().codex.clone(); + + let responses_mock = responses::mount_sse_sequence( + harness.server(), + vec![ + responses::sse(vec![ + responses::ev_assistant_message("m1", "FIRST_REMOTE_REPLY"), + responses::ev_completed_with_tokens("resp-1", 500), + ]), + responses::sse(vec![ + responses::ev_compaction(&inline_summary), + responses::ev_function_call( + "call-inline-stale-token-usage", + DUMMY_FUNCTION_NAME, + "{}", + ), + responses::ev_completed_with_tokens("resp-2", 500), + ]), + responses::sse(vec![ + responses::ev_assistant_message("m3", "AFTER_INLINE_TOOL_REPLY"), + responses::ev_completed("resp-3"), + ]), + ], + ) + .await; + + codex + .submit(Op::UserInput { + items: vec![UserInput::Text { + text: "inline compact turn one".to_string(), + text_elements: Vec::new(), + }], + final_output_json_schema: None, + }) + .await?; + wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; + + codex + .submit(Op::UserInput { + items: vec![UserInput::Text { + text: "inline compact turn two".to_string(), + text_elements: Vec::new(), + }], + final_output_json_schema: None, + }) + .await?; + + let mut saw_turn_complete = false; + let mut token_usage_events = Vec::new(); + while !saw_turn_complete { + let event = codex.next_event().await.expect("event"); + match event.msg { + EventMsg::TokenCount(token_count) => { + if let Some(last_token_usage) = token_count + .info + .as_ref() + .map(|info| info.last_token_usage.total_tokens) + { + token_usage_events.push(last_token_usage); + } + } + EventMsg::TurnComplete(_) => { + saw_turn_complete = true; + } + _ => {} + } + } + + let requests = responses_mock.requests(); + assert_eq!( + requests.len(), + 3, + "expected initial request, inline-compacted tool call, and same-turn follow-up" + ); + assert!( + token_usage_events + .iter() + .copied() + .any(|tokens| tokens > 500), + "expected a post-compaction token event to keep the recomputed local estimate instead of reverting to the provider-reported total: {token_usage_events:?}" + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn auto_server_side_compaction_follow_up_preserves_model_switch_updates() -> Result<()> { skip_if_no_network!(Ok(()));