From 6ad2341250314bb1a9b32ebdf9365bde060e51d9 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 25 Sep 2026 22:07:03 +0900 Subject: [PATCH] fix(gateway): skip immediate same-agent retry on provider 429 --- contextual_orchestrator/orchestrator.py | 5 +++++ tests/test_rate_limit_aware_admission.py | 4 +++- 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/contextual_orchestrator/orchestrator.py b/contextual_orchestrator/orchestrator.py index 6b7a69d80..846c6ae53 100644 --- a/contextual_orchestrator/orchestrator.py +++ b/contextual_orchestrator/orchestrator.py @@ -11016,6 +11016,11 @@ def attempt(agent: ModelAgent) -> EndpointAttempt[Any]: # incidental wording in an upstream error body (e.g. a # 400 that happens to mention "invalid arguments"). decision = classify_provider_transport_failure(exc.retryable) + if exc.provider_status == 429: + # Quota refusal already selected a cooldown above; + # another immediate call to this agent only spends + # the same quota and must not open its health circuit. + decision = replace(downgrade_to_failover(decision), circuit_failure=False) elif isinstance(exc, ProviderResponseError): if allowed_agent_ids is None: raise diff --git a/tests/test_rate_limit_aware_admission.py b/tests/test_rate_limit_aware_admission.py index 97b56f1e3..8796b392a 100644 --- a/tests/test_rate_limit_aware_admission.py +++ b/tests/test_rate_limit_aware_admission.py @@ -23,6 +23,7 @@ import urllib.error from contextlib import contextmanager from copy import deepcopy +from dataclasses import replace from email.message import Message from typing import Any @@ -506,8 +507,9 @@ def _post_chat_completion(port: int, payload: dict[str, Any], token: str): def test_http_route_once_waits_out_storm_and_serves_the_request() -> None: """orchestrator/free over /v1/chat/completions (route_once) waits out a 429 storm.""" orchestrator = TaskOrchestrator( - _free_route_agents(), tool_retry_attempts=0, rate_limit_wait_seconds=5.0 + _free_route_agents(), tool_retry_attempts=1, rate_limit_wait_seconds=5.0 ) + orchestrator.policy = replace(orchestrator.policy, realtime_judge=False) slept: list[float] = [] orchestrator._rate_limit_sleep = slept.append chat_outcomes = QueuedChatOutcomes(