From 57016b6aec17e5607947a7d857104852412e565c Mon Sep 17 00:00:00 2001 From: Harshvardhan Kasliwal Date: Thu, 24 Sep 2026 03:51:47 +0530 Subject: [PATCH] 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 --- .../litellm_core_utils/streaming_handler.py | 27 ++++++++++++++++--- 1 file changed, 23 insertions(+), 4 deletions(-) diff --git a/litellm/litellm_core_utils/streaming_handler.py b/litellm/litellm_core_utils/streaming_handler.py index f97a274708f..90f45828de0 100644 --- a/litellm/litellm_core_utils/streaming_handler.py +++ b/litellm/litellm_core_utils/streaming_handler.py @@ -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: