diff --git a/codex-rs/core/src/compact_remote_v2.rs b/codex-rs/core/src/compact_remote_v2.rs index 101317a0e..251d37b35 100644 --- a/codex-rs/core/src/compact_remote_v2.rs +++ b/codex-rs/core/src/compact_remote_v2.rs @@ -19,6 +19,7 @@ use crate::hook_runtime::run_pre_compact_hooks; use crate::session::session::Session; use crate::session::turn::built_tools; use crate::session::turn_context::TurnContext; +use crate::util::backoff; use codex_analytics::CompactionImplementation; use codex_analytics::CompactionPhase; use codex_analytics::CompactionReason; @@ -34,18 +35,22 @@ use codex_protocol::protocol::CompactedItem; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::TruncationPolicy; use codex_protocol::protocol::TurnStartedEvent; +use codex_protocol::protocol::WarningEvent; use codex_rollout_trace::CompactionCheckpointTracePayload; use codex_rollout_trace::InferenceTraceContext; use codex_utils_output_truncation::approx_token_count; use codex_utils_output_truncation::truncate_text; use futures::StreamExt; -use futures::TryFutureExt; use tokio_util::sync::CancellationToken; use tracing::info; +use tracing::warn; // Mirror the current /responses/compact retained-message default while the // server-side path remains the reference implementation. const RETAINED_MESSAGE_TOKEN_BUDGET: usize = 64_000; +// Compact attempts can run much longer than normal turns, so keep the per-transport +// retry budget smaller than the general Responses stream retry budget. +const MAX_REMOTE_COMPACTION_V2_STREAM_RETRIES: u64 = 2; pub(crate) async fn run_inline_remote_auto_compact_task( sess: Arc, @@ -277,31 +282,106 @@ async fn run_remote_compaction_request_v2( prompt: &Prompt, turn_metadata_header: Option<&str>, ) -> CodexResult<(ResponseItem, String)> { - let stream = client_session - .stream( - prompt, - &turn_context.model_info, - &turn_context.session_telemetry, - turn_context.reasoning_effort, - turn_context.reasoning_summary, - turn_context.config.service_tier.clone(), - turn_metadata_header, - &InferenceTraceContext::disabled(), - ) - .or_else(|err| async { - let total_usage_breakdown = sess.get_total_token_usage_breakdown().await; - let compact_request_log_data = - build_compact_request_log_data(&prompt.input, &prompt.base_instructions.text); - log_remote_compact_failure( - turn_context, - &compact_request_log_data, - total_usage_breakdown, - &err, - ); + let max_retries = turn_context + .provider + .info() + .stream_max_retries() + .min(MAX_REMOTE_COMPACTION_V2_STREAM_RETRIES); + let mut retries = 0; + loop { + let result = match client_session + .stream( + prompt, + &turn_context.model_info, + &turn_context.session_telemetry, + turn_context.reasoning_effort, + turn_context.reasoning_summary, + turn_context.config.service_tier.clone(), + turn_metadata_header, + &InferenceTraceContext::disabled(), + ) + .await + { + Ok(stream) => collect_compaction_output(stream).await, + Err(err) => Err(err), + }; + + match result { + Ok(compaction_output) => return Ok(compaction_output), + Err(err) if !err.is_retryable() => { + log_remote_compaction_request_failure(sess, turn_context, prompt, &err).await; + return Err(err); + } Err(err) - }) - .await?; - collect_compaction_output(stream).await + if retries >= max_retries + && client_session.try_switch_fallback_transport( + &turn_context.session_telemetry, + &turn_context.model_info, + ) => + { + sess.send_event( + turn_context, + EventMsg::Warning(WarningEvent { + message: format!( + "Falling back from WebSockets to HTTPS transport. {err:#}" + ), + }), + ) + .await; + retries = 0; + } + Err(err) if retries < max_retries => { + retries += 1; + let delay = match &err { + CodexErr::Stream(_, requested_delay) => { + requested_delay.unwrap_or_else(|| backoff(retries)) + } + _ => backoff(retries), + }; + warn!( + turn_id = %turn_context.sub_id, + retries, + max_retries, + compact_error = %err, + "remote compaction v2 stream failed; retrying request after delay" + ); + + let report_error = retries > 1 + || cfg!(debug_assertions) + || !sess.services.model_client.responses_websocket_enabled(); + if report_error { + sess.notify_stream_error( + turn_context, + format!("Reconnecting... {retries}/{max_retries}"), + err, + ) + .await; + } + tokio::time::sleep(delay).await; + } + Err(err) => { + log_remote_compaction_request_failure(sess, turn_context, prompt, &err).await; + return Err(err); + } + } + } +} + +async fn log_remote_compaction_request_failure( + sess: &Session, + turn_context: &TurnContext, + prompt: &Prompt, + err: &CodexErr, +) { + let total_usage_breakdown = sess.get_total_token_usage_breakdown().await; + let compact_request_log_data = + build_compact_request_log_data(&prompt.input, &prompt.base_instructions.text); + log_remote_compact_failure( + turn_context, + &compact_request_log_data, + total_usage_breakdown, + err, + ); } async fn collect_compaction_output( @@ -331,8 +411,9 @@ async fn collect_compaction_output( } let Some(response_id) = completed_response_id else { - return Err(CodexErr::Fatal( + return Err(CodexErr::Stream( "remote compaction v2 stream closed before response.completed".to_string(), + None, )); }; diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index d868abbfe..b614c09d3 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -816,6 +816,118 @@ async fn remote_compact_v2_reuses_compaction_trigger_for_followups() -> Result<( Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn remote_compact_v2_retries_failures_with_stream_retry_budget() -> Result<()> { + skip_if_no_network!(Ok(())); + + let harness = TestCodexHarness::with_builder( + test_codex() + .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) + .with_config(|config| { + let _ = config.features.enable(Feature::RemoteCompactionV2); + config.model_provider.request_max_retries = Some(0); + config.model_provider.stream_max_retries = Some(2); + }), + ) + .await?; + let codex = harness.test().codex.clone(); + + let responses_mock = responses::mount_response_sequence( + harness.server(), + vec![ + responses::sse_response(responses::sse(vec![ + responses::ev_assistant_message("m1", "FIRST_REMOTE_REPLY"), + responses::ev_completed("resp-1"), + ])), + ResponseTemplate::new(500).set_body_string("first compact open failed"), + responses::sse_response(responses::sse(vec![serde_json::json!({ + "type": "response.output_item.done", + "item": { + "type": "compaction", + "encrypted_content": "FAILED_COMPACT_SUMMARY", + } + })])), + responses::sse_response(responses::sse(vec![ + serde_json::json!({ + "type": "response.output_item.done", + "item": { + "type": "compaction", + "encrypted_content": "RETRIED_COMPACT_SUMMARY", + } + }), + responses::ev_completed("resp-compact-retry"), + ])), + responses::sse_response(responses::sse(vec![ + responses::ev_assistant_message("m2", "AFTER_COMPACT_REPLY"), + responses::ev_completed("resp-2"), + ])), + ], + ) + .await; + + codex + .submit(Op::UserInput { + environments: None, + items: vec![UserInput::Text { + text: "hello remote compact".into(), + text_elements: Vec::new(), + }], + final_output_json_schema: None, + responsesapi_client_metadata: None, + thread_settings: Default::default(), + }) + .await?; + wait_for_turn_complete(&codex).await; + + codex.submit(Op::Compact).await?; + wait_for_turn_complete(&codex).await; + + codex + .submit(Op::UserInput { + environments: None, + items: vec![UserInput::Text { + text: "after compact".into(), + text_elements: Vec::new(), + }], + final_output_json_schema: None, + responsesapi_client_metadata: None, + thread_settings: Default::default(), + }) + .await?; + wait_for_turn_complete(&codex).await; + + let response_requests = responses_mock.requests(); + assert_eq!( + 5, + response_requests.len(), + "expected initial turn, failed open, failed stream, compact retry, and follow-up turn" + ); + + for compact_request in &response_requests[1..=3] { + assert_eq!("/v1/responses", compact_request.path()); + assert!( + compact_request + .body_json() + .to_string() + .contains("\"type\":\"compaction_trigger\""), + "expected v2 compaction request to include the compaction_trigger item" + ); + } + + let follow_up_request = response_requests.last().expect("follow-up request missing"); + let follow_up_body = follow_up_request.body_json().to_string(); + assert!( + follow_up_body.contains("RETRIED_COMPACT_SUMMARY"), + "expected follow-up request to include the retried compaction payload" + ); + assert!( + !follow_up_body.contains("FAILED_COMPACT_SUMMARY"), + "expected failed compaction attempt output to be discarded" + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn remote_compact_v2_accepts_additional_output_items_before_compaction() -> Result<()> { skip_if_no_network!(Ok(()));