diff --git a/codex-rs/codex-api/src/api_bridge.rs b/codex-rs/codex-api/src/api_bridge.rs index 6b79b4754a..70da94e3fb 100644 --- a/codex-rs/codex-api/src/api_bridge.rs +++ b/codex-rs/codex-api/src/api_bridge.rs @@ -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::(&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 diff --git a/codex-rs/codex-api/src/api_bridge_tests.rs b/codex-rs/codex-api/src/api_bridge_tests.rs index 51fb5a28b2..de94d21d24 100644 --- a/codex-rs/codex-api/src/api_bridge_tests.rs +++ b/codex-rs/codex-api/src/api_bridge_tests.rs @@ -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] diff --git a/codex-rs/codex-api/src/sse/responses.rs b/codex-rs/codex-api/src/sse/responses.rs index e4f5f04d62..bb401db1bd 100644 --- a/codex-rs/codex-api/src/sse/responses.rs +++ b/codex-rs/codex-api/src/sse/responses.rs @@ -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 { - 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, diff --git a/codex-rs/core/tests/suite/quota_exceeded.rs b/codex-rs/core/tests/suite/quota_exceeded.rs index 7497a0f776..852fe5af59 100644 --- a/codex-rs/core/tests/suite/quota_exceeded.rs +++ b/codex-rs/core/tests/suite/quota_exceeded.rs @@ -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(()) } diff --git a/codex-rs/core/tests/suite/retry_after.rs b/codex-rs/core/tests/suite/retry_after.rs index 991afbe92d..b7704b0b6a 100644 --- a/codex-rs/core/tests/suite/retry_after.rs +++ b/codex-rs/core/tests/suite/retry_after.rs @@ -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![