Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 36 additions & 1 deletion contextual_orchestrator/orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand Down
53 changes: 53 additions & 0 deletions tests/test_passthrough_provider_failover.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)})
Expand Down