diff --git a/litellm/proxy/pass_through_endpoints/pass_through_endpoints.py b/litellm/proxy/pass_through_endpoints/pass_through_endpoints.py index 29d6fdb9c73..883c592d56c 100644 --- a/litellm/proxy/pass_through_endpoints/pass_through_endpoints.py +++ b/litellm/proxy/pass_through_endpoints/pass_through_endpoints.py @@ -920,7 +920,13 @@ class _PreviewReportingStream(httpx.AsyncByteStream): async def abandon(self) -> None: self._abandoned = True self.dispatch() - await self._upstream.aclose() + try: + await self._upstream.aclose() + except Exception as err: # noqa: BLE001 # aclose failures must not mask the exception that abandoned the relay + self._log_warning( + "pass_through_endpoint: closing the abandoned upstream error response failed: %s", + type(err).__name__, + ) async def __aiter__(self) -> AsyncIterator[bytes]: total = 0 # rebind-ok: running byte count against the preview budget diff --git a/tests/test_litellm/proxy/pass_through_endpoints/test_pass_through_endpoints.py b/tests/test_litellm/proxy/pass_through_endpoints/test_pass_through_endpoints.py index a32de6a5d17..c6eff6f6896 100644 --- a/tests/test_litellm/proxy/pass_through_endpoints/test_pass_through_endpoints.py +++ b/tests/test_litellm/proxy/pass_through_endpoints/test_pass_through_endpoints.py @@ -4173,12 +4173,15 @@ class _UpstreamErrorBodyStream(httpx.AsyncByteStream): class _UpstreamErrorBodyStreamCloseTracking(_UpstreamErrorBodyStream): - def __init__(self, body: bytes) -> None: + def __init__(self, body: bytes, close_error: Exception | None = None) -> None: super().__init__(body) + self._close_error: Final = close_error self.closed = False async def aclose(self) -> None: self.closed = True + if self._close_error is not None: + raise self._close_error def _upstream_error_request() -> MagicMock: @@ -5096,6 +5099,38 @@ async def test_preview_reporting_stream_abandon_does_not_log_a_client_disconnect assert all("client disconnected" not in call[0] for call in warnings), warnings +@pytest.mark.asyncio +async def test_preview_reporting_stream_abandon_keeps_the_original_exception_when_aclose_fails(): + """A failing upstream close must not replace the exception that abandoned the + relay: abandon returns without raising, still reports, and logs the close failure.""" + upstream_stream: Final = _UpstreamErrorBodyStreamCloseTracking( + b"body", close_error=httpx.CloseError("close failed") + ) + upstream_response: Final = httpx.Response( + status_code=500, + headers={"content-type": "text/event-stream"}, + stream=upstream_stream, + request=httpx.Request("POST", "http://target-api.com/v1beta/models/claude-nope-9:streamGenerateContent"), + ) + reported: list[bytes] = [] + log_warning: Final = MagicMock() + + async def report(preview: bytes) -> None: + reported.append(preview) + + relay: Final = _PreviewReportingStream( + upstream=upstream_response, + report=report, + log_warning=log_warning, + ) + + await relay.abandon() + await relay.dispatch() + assert reported == [b""], reported + warnings: Final = [call.args[0] for call in log_warning.call_args_list] + assert any("closing the abandoned upstream error response failed" in message for message in warnings), warnings + + @pytest.mark.asyncio async def test_preview_report_collected_runs_without_disconnect_warning_after_clean_end(): upstream_response: Final = httpx.Response(