diff --git a/contextual_orchestrator/orchestrator.py b/contextual_orchestrator/orchestrator.py index a29afe8a7..cc6375f9d 100644 --- a/contextual_orchestrator/orchestrator.py +++ b/contextual_orchestrator/orchestrator.py @@ -11067,6 +11067,18 @@ def attempt(agent: ModelAgent) -> EndpointAttempt[Any]: else self.client.chat(agent, messages) ) except Exception as exc: + if review_no_replay and "review" in agent.tags and isinstance(exc, urllib.error.HTTPError): + response_error = exc + try: + exc = classify_provider_failure( + response_error, agent_id=agent.id, + model=agent.model, transport="chat", + ) + finally: + try: + response_error.close() + except Exception: # noqa: BLE001 - preserve the classified failure + pass if _is_request_too_large_error(exc): break every_failure_was_request_too_large = False @@ -11122,7 +11134,7 @@ def attempt(agent: ModelAgent) -> EndpointAttempt[Any]: # exception below proves no send took place. if exc.provider_status != 429: self._record_failure(agent.id) - raise + raise exc from None # The primary chat call is a bounded, side-effect-free # model request, not a tool invocation: classify from # the provider's own already-computed retryability @@ -11147,6 +11159,29 @@ def attempt(agent: ModelAgent) -> EndpointAttempt[Any]: ) else: decision = classify_tool_failure(exc) + if review_no_replay and "review" in agent.tags and decision.kind in { + ToolFailureKind.UNKNOWN, ToolFailureKind.RATE_LIMITED, + }: + rate_limit_signal = self._rate_limited_provider_signal(exc) + if rate_limit_signal is not None: + signal_status, signal_http_error = rate_limit_signal + self._record_rate_limit( + agent.id, + resolve_retry_after_seconds(signal_http_error) + if signal_http_error is not None else None, + status=signal_status, + ) + if decision.kind is ToolFailureKind.UNKNOWN: + self._record_failure(agent.id) + raise ProviderUpstreamError( + agent_id=agent.id, + model=agent.model, + error_code=PROVIDER_OUTCOME_UNKNOWN_CODE, + message="the provider request outcome is unknown; automatic replay is unsafe", + client_status=502, + retryable=False, + transport="chat", + ) from None action = decision.action # A failed attempt is one Bernoulli stability observation # for measured group routing regardless of what happens next. diff --git a/tests/test_passthrough_provider_failover.py b/tests/test_passthrough_provider_failover.py index 688de79e9..6ea4b7de3 100644 --- a/tests/test_passthrough_provider_failover.py +++ b/tests/test_passthrough_provider_failover.py @@ -818,6 +818,59 @@ def chat(self, agent: ModelAgent, _messages: list[dict[str, Any]]) -> str: assert client.calls == [first.id] +@pytest.mark.parametrize("case, expected", [ + ("plain", (502, "provider_outcome_unknown", False)), + ("wrapped", (502, "provider_outcome_unknown", False)), + ("raw_404", (404, "model_not_found", False)), + ("raw_429", (429, "rate_limit_exceeded", True)), +]) +def test_review_free_invoke_unknown_failure_does_not_replay( + case: str, expected: tuple[int, str, bool] +) -> None: + """An unclassified review chat failure cannot authorize another send.""" + first = ModelAgent("first_agent", "first-model", priority=10, + tags=("cost:free", "review")) + second = ModelAgent("second_agent", "second-model", priority=1, + tags=("cost:free", "review")) + orchestrator = TaskOrchestrator([first, second], tool_retry_attempts=0) + responses: list[urllib.error.HTTPError] = [] + if case.startswith("raw_"): + response = _http_error(int(case[4:])) + failure: BaseException = response + responses.append(response) + else: + failure = RuntimeError("synthetic current send outcome unknown") + if case == "wrapped": + cause = _http_error(404, {"error": {"code": "model_not_found"}}) + responses.append(cause) + failure.__cause__ = cause + calls: list[str] = [] + + def chat(agent: ModelAgent, _messages: list[dict[str, Any]]) -> str: + calls.append(agent.id) + if agent.id == first.id: + raise failure + return "synthetic fallback" + + orchestrator.client.chat = chat + try: + with pytest.raises(ProviderUpstreamError) as caught: + orchestrator._invoke_with_rate_limit_recovery( + first, [{"role": "user", "content": "synthetic"}], + text="synthetic", role="worker", + allowed_agent_ids={first.id, second.id}, + review_no_replay=True, virtual_selector=True, + ) + assert (caught.value.client_status, caught.value.error_code, caught.value.retryable) == expected + assert calls == [first.id] + if case.startswith("raw_"): + assert responses[0].closed + finally: + for response in responses: + response.close() + orchestrator.close() + + def test_explicit_structured_synthesis_normalizes_413() -> None: """A sticky explicit model still exposes request-size rejection as 413 authority.""" client = SequencedProxyClient({"primary_agent": _http_error(413)})