retry remote compaction v2 requests (#23951)

## Why

Remote compaction v2 sends a normal `/responses` request with a
compaction trigger. It should follow the retry semantics used by normal
Responses streaming calls for transient stream/request failures, while
keeping a smaller per-transport retry budget because compact attempts
can run much longer than normal turns.

## What changed

- Add a v2 compaction retry loop that uses `stream_max_retries`,
matching normal Responses turn retry mechanics.
- Cap the compact v2 retry budget at 2 retries per transport with
`min(stream_max_retries, 2)`.
- Retry retryable request-open and post-open stream collection failures
through the same loop.
- Use the existing 200ms exponential backoff and requested retry delay
handling used by normal turn retries.
- Emit the same `Reconnecting... n/max` stream-error notification
pattern.
- Fall back from WebSockets to HTTPS after the compact v2 stream retry
budget is exhausted, then reset the retry counter for HTTPS.
- Keep final remote-compaction failure logging after retries/fallback
are exhausted.
- Treat compact stream EOF before `response.completed` as a retryable
stream failure.
- Add compact v2 regression coverage with `request_max_retries = 0` and
`stream_max_retries = 2`, covering both request-open failure and
opened-stream EOF in one end-to-end test.

## Tests

- `just fmt`
- `cargo test -p codex-core remote_compact_v2`
- `just fix -p codex-core`
This commit is contained in:
rhan-oai
2026-05-22 10:14:14 -07:00
committed by GitHub
Unverified
parent d53e68954a
commit dac98cb635
2 changed files with 219 additions and 26 deletions
+107 -26
View File
@@ -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<Session>,
@@ -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,
));
};
+112
View File
@@ -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(()));