fix(codex): rotate the credential pool on Responses HTTP-200 soft failures
The Codex Responses API reports quota exhaustion as HTTP 200 with
`response.status == "failed"` and `error.code == "usage_limit_reached"`.
The SDK never raises, so the exception path's credential-pool rotation
(`recover_after_classification` -> `_recover_with_credential_pool`) never
saw it: `retry_invalid_response` retried the same exhausted account
`max_retries` times and then went straight to cross-provider fallback,
leaving a healthy sibling pool entry unused and the dead entry unmarked.
Fix: `turn_recovery.classify_codex_soft_failure` reshapes `response.error`
as an SDK-style error body and runs it through the existing
`classify_api_error` / `extract_api_error_context` (no second pattern
list); `retry_invalid_response` then tries same-provider pool recovery
FIRST for rate_limit / billing / auth verdicts and falls through to the
existing eager fallback + retry path otherwise. Content-policy and other
non-quota failures never rotate. The `ResponseError` repr that leaked into
the retry trace now shows the provider's message.
Live probe (fake Responses SSE server, 2-entry pool, real AIAgent loop):
before -> 3 requests on tok-A, pool untouched, turn failed;
after -> 1 soft failure on tok-A, cred-0 exhausted (usage_limit_reached),
rotated to tok-B, turn completed; content_policy control -> no
rotation on either side.
Fixes #24159
Salvages #24173 (@jmmaloney4) — recovery order and insertion point; the
hand-rolled pattern classifier is replaced by the shared error classifier.
Co-authored-by: Jack Maloney <jmmaloney4@gmail.com>
This commit is contained in:
@@ -19,7 +19,7 @@ from typing import Any, Dict, List, Optional, Tuple
|
||||
from agent.conversation_compression import COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE
|
||||
from agent.model_metadata import is_output_cap_error, parse_available_output_tokens_from_error
|
||||
from agent.retry_utils import is_zai_coding_overload_error, zai_coding_overload_retry_ceiling
|
||||
from agent.error_classifier import FailoverReason
|
||||
from agent.error_classifier import FailoverReason, classify_api_error
|
||||
from agent.message_sanitization import (
|
||||
_looks_like_image_content_rejection, _sanitize_messages_non_ascii,
|
||||
_sanitize_messages_surrogates, _sanitize_structure_non_ascii, _sanitize_structure_surrogates,
|
||||
@@ -1239,6 +1239,46 @@ def compute_error_backoff(
|
||||
return wait_time
|
||||
|
||||
|
||||
def _codex_soft_failure_error(response: Any) -> Dict[str, Any]:
|
||||
"""``response.error`` of a Codex ``failed``/``cancelled`` Response as ``{"code", "message"}``
|
||||
(the SDK types it as ``ResponseError``; the raw-SSE assembler keeps the dict); ``{}`` when absent."""
|
||||
error_obj = getattr(response, "error", None)
|
||||
if not error_obj:
|
||||
return {}
|
||||
if isinstance(error_obj, dict):
|
||||
fields = error_obj
|
||||
elif hasattr(error_obj, "code") or hasattr(error_obj, "message"):
|
||||
fields = {"code": getattr(error_obj, "code", None), "message": getattr(error_obj, "message", None)}
|
||||
else:
|
||||
fields = {"message": str(error_obj)}
|
||||
return {k: v for k, v in fields.items() if isinstance(v, str) and v.strip()}
|
||||
|
||||
|
||||
class _CodexSoftFailure(Exception):
|
||||
"""A Codex HTTP-200 ``status=failed`` Response reshaped so ``classify_api_error`` and
|
||||
``extract_api_error_context`` read ``response.error`` exactly like an SDK error body."""
|
||||
|
||||
def __init__(self, error: Dict[str, Any]) -> None:
|
||||
super().__init__(error.get("message") or "")
|
||||
self.body = {"error": error}
|
||||
|
||||
|
||||
def classify_codex_soft_failure(agent: Any, response: Any) -> Tuple[Any, Dict[str, Any]]:
|
||||
"""``(classified, error_context)`` for a Codex ``failed``/``cancelled`` Response, or
|
||||
``(None, {})`` when it is not one. The SDK never raises on these HTTP-200 soft failures,
|
||||
so this is the only place their quota/billing/auth semantics reach the credential pool."""
|
||||
if agent.api_mode != "codex_responses":
|
||||
return None, {}
|
||||
if str(getattr(response, "status", "") or "").strip().lower() not in {"failed", "cancelled"}:
|
||||
return None, {}
|
||||
exc = _CodexSoftFailure(_codex_soft_failure_error(response))
|
||||
classified = classify_api_error(
|
||||
exc, provider=getattr(agent, "provider", "") or "", model=getattr(agent, "model", "") or "",
|
||||
base_url=str(getattr(agent, "base_url", "") or ""), api_key=getattr(agent, "api_key", None),
|
||||
)
|
||||
return classified, agent._extract_api_error_context(exc)
|
||||
|
||||
|
||||
def validate_response_shape(agent: Any, response: Any) -> Tuple[bool, List[str]]:
|
||||
"""Validate the raw provider response via the transport; ``(response_invalid,
|
||||
error_details)``. A Codex ``failed``/``cancelled`` status (e.g. quota exhaustion) is
|
||||
@@ -1251,11 +1291,9 @@ def validate_response_shape(agent: Any, response: Any) -> Tuple[bool, List[str]]
|
||||
if agent.api_mode == "codex_responses":
|
||||
_codex_resp_status = str(getattr(response, "status", "") or "").strip().lower()
|
||||
if _codex_resp_status in {"failed", "cancelled"}:
|
||||
_codex_error_obj = getattr(response, "error", None)
|
||||
_codex_error_msg = (
|
||||
_codex_error_obj.get("message") if isinstance(_codex_error_obj, dict)
|
||||
else str(_codex_error_obj) if _codex_error_obj
|
||||
else f"Responses API returned status '{_codex_resp_status}'"
|
||||
_codex_soft_failure_error(response).get("message")
|
||||
or f"Responses API returned status '{_codex_resp_status}'"
|
||||
)
|
||||
logger.warning(
|
||||
"Codex response status='%s' (error=%s). Routing to fallback. %s",
|
||||
@@ -1301,7 +1339,8 @@ def describe_invalid_response(agent: Any, response: Any, api_duration: float) ->
|
||||
provider_name = "Unknown"
|
||||
_has_error = bool(response and hasattr(response, 'error') and response.error)
|
||||
if _has_error:
|
||||
error_msg = str(response.error)
|
||||
# A typed ``ResponseError`` stringifies as its repr; show the provider's message.
|
||||
error_msg = _codex_soft_failure_error(response).get("message") or str(response.error)
|
||||
if hasattr(response.error, 'metadata') and response.error.metadata:
|
||||
provider_name = response.error.metadata.get('provider_name', 'Unknown')
|
||||
elif response and hasattr(response, 'message') and response.message:
|
||||
|
||||
@@ -13,6 +13,7 @@ import logging
|
||||
import time
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
from agent.error_classifier import FailoverReason
|
||||
from agent.turn_api_call import stop_thinking_spinner
|
||||
from agent.turn_failure_copy import invalid_response_failure_reason, provider_label_for, site_copy, stamp_failure
|
||||
from agent.turn_truncation import handle_content_policy_refusal, recover_from_truncation
|
||||
@@ -236,7 +237,9 @@ def retry_invalid_response(
|
||||
else jittered backoff that preserves a pending redirect."""
|
||||
from agent.conversation_loop import _arm_fallback_restart
|
||||
from agent.retry_utils import jittered_backoff
|
||||
from agent.turn_recovery import describe_invalid_response, interruptible_backoff_sleep
|
||||
from agent.turn_recovery import (
|
||||
classify_codex_soft_failure, describe_invalid_response, interruptible_backoff_sleep,
|
||||
)
|
||||
|
||||
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> InvalidResponseVerdict:
|
||||
return InvalidResponseVerdict(
|
||||
@@ -255,6 +258,20 @@ def retry_invalid_response(
|
||||
)
|
||||
# Retry status is buffered and only surfaced if every retry+fallback exhausts.
|
||||
thinking_spinner = stop_thinking_spinner(agent, thinking_spinner)
|
||||
|
||||
# Codex reports quota exhaustion as HTTP 200 ``status=failed`` — the SDK never raises, so the
|
||||
# exception path's credential-pool rotation never sees it. Same-provider recovery for the
|
||||
# pool-recoverable reasons FIRST (a healthy sibling account beats burning cross-provider
|
||||
# fallback); content-policy and other failures keep the fallback/retry path (#24159).
|
||||
_soft, _soft_ctx = classify_codex_soft_failure(agent, response)
|
||||
if _soft is not None and (_soft.reason in (FailoverReason.rate_limit, FailoverReason.billing) or _soft.is_auth):
|
||||
_recovered, _retry.has_retried_429 = agent._recover_with_credential_pool(
|
||||
status_code=None, has_retried_429=_retry.has_retried_429, classified_reason=_soft.reason,
|
||||
error_context=_soft_ctx, billing_unverified=_soft.billing_unverified,
|
||||
)
|
||||
if _recovered:
|
||||
agent._buffer_diagnostic_status(f"🔄 Codex soft failure ({_soft.reason.value}) — switched to the next pool credential, retrying...")
|
||||
return _verdict("continue")
|
||||
retry_count += 1
|
||||
|
||||
# Eager fallback: empty/malformed responses often mean rate limiting.
|
||||
|
||||
104
tests/agent/test_codex_soft_failure_pool_rotation.py
Normal file
104
tests/agent/test_codex_soft_failure_pool_rotation.py
Normal file
@@ -0,0 +1,104 @@
|
||||
"""Codex Responses HTTP-200 soft failures (``response.status == "failed"``) must reach the
|
||||
same-provider credential pool before cross-provider fallback (#24159).
|
||||
|
||||
The SDK never raises on these, so the exception path's ``_recover_with_credential_pool`` never
|
||||
sees them; ``retry_invalid_response`` has to classify ``response.error`` itself.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
from agent.agent_runtime_helpers import recover_with_credential_pool
|
||||
from agent.credential_pool import STATUS_EXHAUSTED, CredentialPool, PooledCredential
|
||||
from agent.turn_response_check import retry_invalid_response
|
||||
from agent.turn_retry_state import TurnRetryState
|
||||
|
||||
_BASE_URL = "https://chatgpt.com/backend-api/codex"
|
||||
|
||||
|
||||
def _entry(i: int) -> PooledCredential:
|
||||
return PooledCredential(
|
||||
provider="openai-codex", id=f"cred-{i}", label=f"acct-{i}", auth_type="api_key", priority=i,
|
||||
source="manual", access_token=f"tok-{i}-1234567890", base_url=_BASE_URL,
|
||||
)
|
||||
|
||||
|
||||
class _Agent:
|
||||
log_prefix = ""
|
||||
quiet_mode = True
|
||||
api_mode = "codex_responses"
|
||||
provider = "openai-codex"
|
||||
model = "gpt-5.1-codex"
|
||||
base_url = _BASE_URL
|
||||
_fallback_chain = ()
|
||||
_fallback_index = 0
|
||||
_credential_pool_revert_id = None
|
||||
|
||||
def __init__(self, pool: CredentialPool) -> None:
|
||||
self._credential_pool = pool
|
||||
self.api_key = pool.select().access_token
|
||||
self.swapped_to: list = []
|
||||
self._try_activate_fallback = MagicMock(return_value=False)
|
||||
|
||||
def _recover_with_credential_pool(self, **kwargs):
|
||||
return recover_with_credential_pool(self, **kwargs)
|
||||
|
||||
def _extract_api_error_context(self, error):
|
||||
from agent.agent_runtime_helpers import extract_api_error_context
|
||||
|
||||
return extract_api_error_context(error)
|
||||
|
||||
def _swap_credential(self, entry):
|
||||
self.swapped_to.append(entry.id)
|
||||
self.api_key = entry.access_token
|
||||
return True
|
||||
|
||||
def _has_pending_fallback(self):
|
||||
return False
|
||||
|
||||
def _clean_error_message(self, msg):
|
||||
return msg
|
||||
|
||||
def __getattr__(self, name):
|
||||
return lambda *args, **kwargs: None
|
||||
|
||||
|
||||
def _soft_failure(code: str, message: str) -> SimpleNamespace:
|
||||
# The SDK types ``response.error`` as ``ResponseError(code=..., message=...)``, not a dict.
|
||||
return SimpleNamespace(status="failed", output=[], output_text="", error=SimpleNamespace(code=code, message=message))
|
||||
|
||||
|
||||
def _run(agent: _Agent, response: SimpleNamespace):
|
||||
return retry_invalid_response(
|
||||
agent, response=response, error_details=["response.status=failed"], _retry=TurnRetryState(),
|
||||
thinking_spinner=None, messages=[], api_messages=[], api_kwargs=None, active_system_prompt=None,
|
||||
conversation_history=None, retry_count=0, max_retries=3, compression_attempts=0, api_call_count=1,
|
||||
api_request_id="r", api_start_time=0.0, api_duration=0.4, effective_task_id="t", turn_id="turn",
|
||||
)
|
||||
|
||||
|
||||
def test_quota_soft_failure_rotates_pool_before_provider_fallback():
|
||||
pool = CredentialPool("openai-codex", [_entry(0), _entry(1)])
|
||||
agent = _Agent(pool)
|
||||
|
||||
verdict = _run(agent, _soft_failure("usage_limit_reached", "You've hit your usage limit. Try again at 3:00 PM."))
|
||||
|
||||
assert verdict.action == "continue"
|
||||
assert agent.swapped_to == ["cred-1"] and agent.api_key == "tok-1-1234567890"
|
||||
benched = next(e for e in pool.entries() if e.id == "cred-0")
|
||||
assert benched.last_status == STATUS_EXHAUSTED and benched.last_error_reason == "usage_limit_reached"
|
||||
agent._try_activate_fallback.assert_not_called()
|
||||
|
||||
|
||||
def test_content_policy_soft_failure_leaves_pool_alone():
|
||||
pool = CredentialPool("openai-codex", [_entry(0), _entry(1)])
|
||||
agent = _Agent(pool)
|
||||
|
||||
verdict = _run(agent, _soft_failure("content_policy_violation", "Your request was rejected by our safety system."))
|
||||
|
||||
assert verdict.action == "continue" # ordinary invalid-response retry path
|
||||
assert agent.swapped_to == [] and agent.api_key == "tok-0-1234567890"
|
||||
assert all(e.last_status is None for e in pool.entries())
|
||||
agent._try_activate_fallback.assert_called()
|
||||
@@ -37,6 +37,10 @@ Your request
|
||||
→ 401 auth expired?
|
||||
→ Try refreshing the token (OAuth)
|
||||
→ Refresh failed → rotate to next pool key
|
||||
→ HTTP 200 but `response.status: failed` (ChatGPT/Codex reports usage limits this way)?
|
||||
→ Same rules as above, keyed on the embedded error code/message:
|
||||
quota/billing/auth → pool rotation first, provider fallback only once the pool is exhausted;
|
||||
content-policy and other failures → no rotation
|
||||
→ Success → continue normally
|
||||
```
|
||||
|
||||
|
||||
Reference in New Issue
Block a user