diff --git a/codex-rs/core/src/codex/compact.rs b/codex-rs/core/src/codex/compact.rs index b6bca43e10..d1e12cc6bc 100644 --- a/codex-rs/core/src/codex/compact.rs +++ b/codex-rs/core/src/codex/compact.rs @@ -78,7 +78,8 @@ async fn run_compact_task_inner( let mut trimmed_tails: Vec> = Vec::new(); let max_retries = turn_context.client.get_provider().stream_max_retries(); - let mut retries = 0; + let mut context_retries = 0u64; + let mut stream_retries = 0u64; let rollout_item = RolloutItem::TurnContext(TurnContextItem { cwd: turn_context.cwd.clone(), @@ -124,8 +125,7 @@ async fn run_compact_task_inner( if !trimmed.is_empty() { truncated_count += trimmed.len(); trimmed_tails.push(trimmed); - retries += 1; - if retries >= max_retries { + if context_retries >= max_retries { sess.set_total_tokens_full(&sub_id, turn_context.as_ref()) .await; let event = Event { @@ -137,6 +137,9 @@ async fn run_compact_task_inner( sess.send_event(event).await; return; } + context_retries += 1; + stream_retries = 0; + // Keep stream retry budget untouched; we trimmed context successfully. continue; } } @@ -152,12 +155,12 @@ async fn run_compact_task_inner( return; } Err(e) => { - if retries < max_retries { - retries += 1; - let delay = backoff(retries); + if stream_retries < max_retries { + stream_retries += 1; + let delay = backoff(stream_retries); sess.notify_stream_error( &sub_id, - format!("Re-connecting... {retries}/{max_retries}"), + format!("Re-connecting... {stream_retries}/{max_retries}"), ) .await; tokio::time::sleep(delay).await; diff --git a/codex-rs/core/tests/suite/compact.rs b/codex-rs/core/tests/suite/compact.rs index 518888f81a..d1e0bb2d6a 100644 --- a/codex-rs/core/tests/suite/compact.rs +++ b/codex-rs/core/tests/suite/compact.rs @@ -32,6 +32,7 @@ use serde_json::Value; pub(super) const FIRST_REPLY: &str = "FIRST_REPLY"; pub(super) const SUMMARY_TEXT: &str = "SUMMARY_ONLY_CONTEXT"; const THIRD_USER_MSG: &str = "next turn"; +const THIRD_ASSISTANT_MSG: &str = "post compact assistant"; const AUTO_SUMMARY_TEXT: &str = "AUTO_SUMMARY"; const FIRST_AUTO_MSG: &str = "token limit start"; const SECOND_AUTO_MSG: &str = "token limit push"; @@ -646,6 +647,10 @@ async fn manual_compact_retries_after_context_window_error() { ev_assistant_message("m2", SUMMARY_TEXT), ev_completed("r2"), ]); + let third_turn = sse(vec![ + ev_assistant_message("m3", THIRD_ASSISTANT_MSG), + ev_completed("r3"), + ]); let request_log = mount_sse_sequence( &server, @@ -653,6 +658,7 @@ async fn manual_compact_retries_after_context_window_error() { user_turn.clone(), compact_failed.clone(), compact_succeeds.clone(), + third_turn, ], ) .await; @@ -711,8 +717,8 @@ async fn manual_compact_retries_after_context_window_error() { let requests = request_log.requests(); assert_eq!( requests.len(), - 3, - "expected user turn and two compact attempts" + 4, + "expected user turn, two compact attempts, and one follow-up turn" ); let compact_attempt = requests[1].body_json();