Fix retry classification for throttling and quota errors (#45602)

## Why

`slow_down` errors were treated as terminal server overloads, while exhausted credit balances and spend limits fell through to retryable stream errors.

## What changed

- Classify `slow_down` as a retryable rate limit in HTTP 503 responses and SSE failures, preserving the server message and parsing retry delays from SSE error messages.
- Map `credit_balance_exhausted`, `organization_spend_limit_exceeded`, and `project_spend_limit_exceeded` SSE errors to quota exhaustion so they terminate without retries.

## Testing

Extend tests to cover HTTP error classification and retryability, `slow_down` stream retry exhaustion and message-provided delays, and a single `UsageLimitExceeded` error with no retries for each quota error code.

GitOrigin-RevId: 7985eba71ba4a6fe897a45957e448fadd42f3963
This commit is contained in:
Steve Coffey
2026-09-15 04:09:41 +00:00
committed by copyberry
parent fc269b66ad
commit 31ffe2bc9a
5 changed files with 83 additions and 38 deletions

View File

@@ -74,15 +74,21 @@ pub fn map_api_error(err: ApiError) -> CodexErr {
if status == http::StatusCode::SERVICE_UNAVAILABLE
&& let Ok(value) = serde_json::from_str::<serde_json::Value>(&body_text)
&& matches!(
value
.get("error")
.and_then(|error| error.get("code"))
.and_then(serde_json::Value::as_str),
Some("server_is_overloaded" | "slow_down")
)
&& let Some(error) = value.get("error")
{
return CodexErr::ServerOverloaded;
match error.get("code").and_then(Value::as_str) {
Some("server_is_overloaded") => return CodexErr::ServerOverloaded,
Some("slow_down") => {
return CodexErr::new(CodexErrorDetails::RateLimitExceeded(
error
.get("message")
.and_then(Value::as_str)
.unwrap_or_default()
.to_owned(),
));
}
_ => {}
}
}
if (status == http::StatusCode::BAD_REQUEST

View File

@@ -52,21 +52,29 @@ fn map_api_error_preserves_retry_delay() {
}
#[test]
fn map_api_error_maps_server_overloaded_from_503_body() {
let body = serde_json::json!({
"error": {
"code": "server_is_overloaded"
}
})
.to_string();
let err = map_api_error(ApiError::Transport(TransportError::Http {
status: http::StatusCode::SERVICE_UNAVAILABLE,
url: Some("http://example.com/v1/responses".to_string()),
headers: None,
body: Some(body),
}));
assert!(matches!(err.details(), CodexErrorDetails::ServerOverloaded));
fn map_api_error_distinguishes_capacity_from_slow_down() {
for (code, expected, retryable) in [
(
"server_is_overloaded",
CodexErrorInfo::ServerOverloaded,
false,
),
("slow_down", CodexErrorInfo::RateLimitExceeded, true),
("unknown_error", CodexErrorInfo::Other, true),
] {
let err = map_api_error(ApiError::Transport(TransportError::Http {
status: http::StatusCode::SERVICE_UNAVAILABLE,
url: None,
headers: None,
body: Some(
serde_json::json!({"error": {"code": code, "message": "retry later"}}).to_string(),
),
}));
assert_eq!(
(err.to_codex_protocol_error(), err.is_retryable()),
(expected, retryable)
);
}
}
#[test]

View File

@@ -455,7 +455,7 @@ pub fn process_responses_event(
let delay = try_parse_retry_after(&error);
let message = error.message.unwrap_or_default();
response_error = match error.code.as_deref() {
Some("rate_limit_exceeded") => {
Some("rate_limit_exceeded" | "slow_down") => {
ApiError::RateLimitExceeded { message, delay }
}
_ => ApiError::Retryable { message, delay },
@@ -683,7 +683,10 @@ async fn process_sse_with_treatment(
}
fn try_parse_retry_after(err: &Error) -> Option<Duration> {
if err.code.as_deref() != Some("rate_limit_exceeded") {
if !matches!(
err.code.as_deref(),
Some("rate_limit_exceeded" | "slow_down")
) {
return None;
}
@@ -713,7 +716,15 @@ fn is_context_window_error(error: &Error) -> bool {
}
fn is_quota_exceeded_error(error: &Error) -> bool {
error.code.as_deref() == Some("insufficient_quota")
matches!(
error.code.as_deref(),
Some(
"insufficient_quota"
| "credit_balance_exhausted"
| "organization_spend_limit_exceeded"
| "project_spend_limit_exceeded"
)
)
}
fn is_usage_not_included(error: &Error) -> bool {
@@ -726,7 +737,6 @@ fn is_cyber_policy_error(error: &Error) -> bool {
fn is_server_overloaded_error(error: &Error) -> bool {
error.code.as_deref() == Some("server_is_overloaded")
|| error.code.as_deref() == Some("slow_down")
}
fn cyber_policy_fallback_message() -> String {
@@ -1118,6 +1128,7 @@ mod tests {
async fn failed_response_classification_uses_error_code() {
for (code, message) in [
("rate_limit_exceeded", "Temporary limit."),
("slow_down", "Temporary limit."),
(
"unknown_error",
"Rate limit reached. Please try again in 1s.",
@@ -1131,7 +1142,7 @@ mod tests {
let events = collect_events(&[sse.as_bytes()]).await;
match (code, events.as_slice()) {
(
"rate_limit_exceeded",
"rate_limit_exceeded" | "slow_down",
[
Err(ApiError::RateLimitExceeded {
message: actual,

View File

@@ -1,5 +1,6 @@
use anyhow::Result;
use codex_core::TurnInputRequest;
use codex_protocol::protocol::CodexErrorInfo;
use codex_protocol::protocol::EventMsg;
use codex_protocol::user_input::UserInput;
use core_test_support::responses::ev_response_created;
@@ -12,8 +13,12 @@ use core_test_support::wait_for_event;
use pretty_assertions::assert_eq;
use serde_json::json;
#[test_case::test_case("insufficient_quota"; "quota")]
#[test_case::test_case("credit_balance_exhausted"; "credit_balance")]
#[test_case::test_case("organization_spend_limit_exceeded"; "organization_spend_limit")]
#[test_case::test_case("project_spend_limit_exceeded"; "project_spend_limit")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn quota_exceeded_emits_single_error_event() -> Result<()> {
async fn quota_exceeded_emits_single_error_event(code: &str) -> Result<()> {
skip_if_no_network!(Ok(()));
let server = start_mock_server().await;
@@ -28,7 +33,7 @@ async fn quota_exceeded_emits_single_error_event() -> Result<()> {
"response": {
"id": "resp-1",
"error": {
"code": "insufficient_quota",
"code": code,
"message": "You exceeded your current quota, please check your plan and billing details."
}
}
@@ -37,15 +42,14 @@ async fn quota_exceeded_emits_single_error_event() -> Result<()> {
)
.await;
let test = builder.build(&server).await?;
let test = builder.build_with_auto_env(&server).await?;
test.codex
.start_or_steer_turn(TurnInputRequest::user_input(vec![UserInput::Text {
text: "quota?".into(),
text_elements: Vec::new(),
}]))
.await
.unwrap();
.await?;
let mut error_events = 0;
@@ -59,6 +63,10 @@ async fn quota_exceeded_emits_single_error_event() -> Result<()> {
err.message,
"Quota exceeded. Check your plan and billing details."
);
assert_eq!(
err.codex_error_info,
Some(CodexErrorInfo::UsageLimitExceeded)
);
}
EventMsg::TurnComplete(_) => break,
_ => {}
@@ -66,6 +74,14 @@ async fn quota_exceeded_emits_single_error_event() -> Result<()> {
}
assert_eq!(error_events, 1, "expected exactly one Codex:Error event");
let requests = server.received_requests().await.expect("recorded requests");
assert_eq!(
requests
.iter()
.filter(|request| request.url.path().ends_with("/responses"))
.count(),
1
);
Ok(())
}

View File

@@ -989,8 +989,10 @@ async fn sse_failure_uses_local_backoff_despite_retry_after() -> Result<()> {
}
/// Headerless sampled stream rate limits exhaust retries before one terminal error.
#[test_case::test_case("rate_limit_exceeded"; "rate_limit")]
#[test_case::test_case("slow_down"; "slow_down")]
#[tokio::test(flavor = "current_thread")]
async fn sse_failure_without_retry_after_exhausts_stream_retries() -> Result<()> {
async fn sse_failure_without_retry_after_exhausts_stream_retries(code: &str) -> Result<()> {
skip_if_no_network!(Ok(()));
let mut telemetry = RetryTelemetryCapture::install();
@@ -1000,12 +1002,12 @@ async fn sse_failure_without_retry_after_exhausts_stream_retries() -> Result<()>
vec![
responses::sse_response(responses::sse_failed(
"rate-limited",
"rate_limit_exceeded",
code,
"Rate limit exceeded.",
)),
responses::sse_response(responses::sse_failed(
"still-rate-limited",
"rate_limit_exceeded",
code,
"Rate limit exceeded.",
)),
],
@@ -1069,8 +1071,10 @@ async fn sse_failure_without_retry_after_exhausts_stream_retries() -> Result<()>
}
/// Rate-limit messages already provide an exact retry delay without an HTTP header.
#[test_case::test_case("rate_limit_exceeded"; "rate_limit")]
#[test_case::test_case("slow_down"; "slow_down")]
#[tokio::test(flavor = "current_thread")]
async fn sse_rate_limit_message_uses_server_advised_retry_delay() -> Result<()> {
async fn sse_rate_limit_message_uses_server_advised_retry_delay(code: &str) -> Result<()> {
skip_if_no_network!(Ok(()));
let mut telemetry = RetryTelemetryCapture::install();
@@ -1080,7 +1084,7 @@ async fn sse_rate_limit_message_uses_server_advised_retry_delay() -> Result<()>
vec![
responses::sse_response(responses::sse_failed(
"rate-limited",
"rate_limit_exceeded",
code,
"Rate limit exceeded. Please try again in 1s.",
)),
responses::sse_response(responses::sse(vec![