fix(router): report and cool down deployments on fetch_stream() failure

Problem: for providers whose stream opens lazily (Gemini/Vertex),
Router._acompletion calls CustomStreamWrapper.fetch_stream() after
litellm.acompletion() has already returned. If fetch_stream() raises
(e.g. a 429), the exception only reaches router.py's local except
block, which bumps fail_calls and re-raises but never calls
logging_obj.failure_handler / async_failure_handler. No failure
callback fires and the deployment is never cooled down, so the same
broken deployment gets retried immediately on every subsequent
request. A provider whose error surfaces inside acompletion() itself
(e.g. openai/) does not have this problem, since litellm.utils.client
wraps that call and fires the failure handlers normally.

Root cause: CustomStreamWrapper.fetch_stream()
(litellm/litellm_core_utils/streaming_handler.py) had no error
handling around self.make_call() - an exception there propagated
raw to the caller with no logging.

Fix:
- fetch_stream() now wraps make_call() in try/except and routes any
  exception through the existing _log_stream_failure_and_raise()
  helper, which already dispatches failure callbacks
  (dispatch_failure_handlers) the same way in-stream chunk failures
  do. This reuses the router's existing failure-callback ->
  async_deployment_callback_on_failure -> _set_cooldown_deployments
  wiring (confirmed in router.py), so no new cooldown logic was
  needed.
- Added a de-dupe guard in _log_stream_failure_and_raise: callers
  that reach fetch_stream() indirectly through __anext__() (which
  also calls _log_stream_failure_and_raise in its own except block)
  would otherwise report the same exception twice. The method now
  marks the exception after its first dispatch and skips the second.

Testing:
- Full existing suite: tests/test_litellm/litellm_core_utils/test_streaming_handler.py
  and test_max_streaming_duration.py - 144 passed (5 initial failures
  were an unrelated missing `proto-plus` dependency in a fresh clone,
  confirmed by installing it and re-running clean).
- Verified with the self-contained repro from the issue (local aiohttp
  mock servers, no real provider needed): before the fix,
  failure_callbacks=[] on both requests and the failing deployment is
  retried every time; after the fix, failure_callbacks=['gemini-bad']
  on request 1 and request 2 skips the cooled-down deployment
  entirely, matching the issue's expected output exactly.

Fixes #42757
This commit is contained in:
Harshvardhan Kasliwal 2026-09-24 03:51:47 +05:30
parent 49d0ece934
commit 57016b6aec

View file

@ -1961,8 +1961,19 @@ class CustomStreamWrapper:
async def fetch_stream(self):
if self.completion_stream is None and self.make_call is not None:
# Call make_call to get the completion stream
self.completion_stream = await self.make_call(client=litellm.module_level_aclient)
try:
# Call make_call to get the completion stream
self.completion_stream = await self.make_call(client=litellm.module_level_aclient)
except Exception as e:
# make_call() can raise before any chunk is ever pulled (e.g. the
# provider rejects the request when the stream is opened, such as
# a 429 on Gemini/Vertex). Router._acompletion calls fetch_stream()
# directly right after litellm.acompletion() returns - for
# providers whose stream opens lazily - which is outside the
# logging wrapper that normally fires failure callbacks. Without
# this, no failure_handler/async_failure_handler ever runs, so the
# deployment never cools down and gets retried immediately.
self._log_stream_failure_and_raise(e)
self._stream_iter = self.completion_stream.__aiter__()
return self.completion_stream
@ -2197,13 +2208,21 @@ class CustomStreamWrapper:
return processed_chunk
def _log_stream_failure_and_raise(self, e: Exception) -> NoReturn:
traceback_exception: Final = traceback.format_exc()
if self.logging_obj is not None:
# Guard against double-reporting: fetch_stream() may already have
# dispatched failure callbacks for this exact exception before it
# propagated up into __anext__'s except block, which also routes here.
already_logged: Final = getattr(e, "_litellm_stream_failure_logged", False)
if not already_logged and self.logging_obj is not None:
traceback_exception: Final = traceback.format_exc()
self._record_partial_usage_for_failure()
## LOGGING
asyncio.create_task(
self.logging_obj.dispatch_failure_handlers(e, traceback_exception, prefer_async_handlers=True)
)
try:
e._litellm_stream_failure_logged = True # type: ignore[attr-defined]
except Exception:
pass
self._handle_stream_fallback_error(e)
def _record_partial_usage_for_failure(self) -> None: