fix(router): retry /v1/messages stream drops that happen before first content

On the aanthropic_messages route, a transport drop before the first
content chunk (e.g. aiohttp 'Connection closed' / TransferEncodingError
from a Databricks AI Gateway) is raised from the stream iterator —
outside the retried call boundary — so num_retries / retry_policy never
applied and the request surfaced as a 500 after a single upstream
attempt (issue #44238; 9 dropped streams in one week across three
Databricks-served models in the report).

When nothing has reached the client and no fallbacks are configured,
retry the same group once from the stream recovery path; after any
frame reached the client (or with fallbacks configured, whose recovery
owns the failure) the existing behavior is unchanged. The retry is
structurally bounded: a retry that drops again comes back from a raw
stream outside the recovery gate.
This commit is contained in:
JingHao-Leon 2026-10-03 05:16:12 +08:00
parent d76c279c95
commit 36f0a0fe6d
2 changed files with 673 additions and 805 deletions

View file

@ -5551,10 +5551,12 @@ class Router:
model, initial_kwargs
)
buffered_lifecycle_chunks: tuple[bytes, ...] = () # rebind-ok: flushed once committed or on decline
chunks_sent_to_client = 0 # any frame that reached the client blocks a same-group retry
try:
async for chunk in source_iterator:
if _anthropic_stream_forwards_ping_live(chunk, has_generated_content):
yield chunk
chunks_sent_to_client += 1
continue
if _anthropic_stream_commits_now(chunk, has_generated_content, len(buffered_lifecycle_chunks)):
has_generated_content = True
@ -5609,9 +5611,36 @@ class Router:
yield buffered_chunk
buffered_lifecycle_chunks = ()
yield chunk
chunks_sent_to_client += 1
for buffered_chunk in buffered_lifecycle_chunks:
yield buffered_chunk
except Exception as stream_error: # noqa: BLE001 # any raised provider error must reach the fallback gate
# A transport drop before ANY frame reached the client never
# reaches the retry machinery — the error is raised from the
# stream iterator, outside the retried call boundary — and
# surfaced as a 500 after a single upstream attempt (issue
# #44238). With no fallbacks configured there is nothing to
# hand the failure to, so retry the same group once; once any
# frame reached the client (or fallbacks are configured, whose
# recovery owns the failure) keep the existing behavior.
if (
chunks_sent_to_client == 0
and not fallbacks_disabled_for_request(initial_kwargs)
and not initial_kwargs.get("fallbacks", self.fallbacks)
):
verbose_router_logger.info(
"Anthropic messages stream dropped before first content; retrying the same group once"
)
retry_kwargs: Final = {
**initial_kwargs,
"original_function": self._ageneric_api_call_with_fallbacks_anthropic_messages_attempt,
}
retry_stream = await self._ageneric_api_call_with_fallbacks_anthropic_messages_attempt(
**retry_kwargs
)
async for item in retry_stream:
yield item
return
async for item in self._aanthropic_messages_recover_stream_error(
stream_error,
has_generated_content,

File diff suppressed because it is too large Load diff