diff --git a/litellm/llms/custom_httpx/aiohttp_transport.py b/litellm/llms/custom_httpx/aiohttp_transport.py index 62f707b3622..b97a59a93a6 100644 --- a/litellm/llms/custom_httpx/aiohttp_transport.py +++ b/litellm/llms/custom_httpx/aiohttp_transport.py @@ -116,6 +116,16 @@ class AiohttpResponseStream(httpx.AsyncByteStream): # For other exceptions, use the normal mapping with map_aiohttp_exceptions(): raise + finally: + # Release the aiohttp connection when iteration ends for any + # reason (read timeout, cancellation from a client disconnect, + # GeneratorExit). Without this, abnormally terminated streams + # permanently hold a slot in the TCPConnector pool; once the + # pool is exhausted every request to that host times out (408) + # until the proxy is restarted, even after the backend recovers. + # On a fully-read response the connection was already released + # at EOF and close() is a no-op. + self._aiohttp_response.close() async def aclose(self) -> None: with map_aiohttp_exceptions(): diff --git a/tests/test_litellm/llms/custom_httpx/test_aiohttp_transport.py b/tests/test_litellm/llms/custom_httpx/test_aiohttp_transport.py index 474ffee3304..0c5e386c438 100644 --- a/tests/test_litellm/llms/custom_httpx/test_aiohttp_transport.py +++ b/tests/test_litellm/llms/custom_httpx/test_aiohttp_transport.py @@ -63,10 +63,14 @@ class MockAiohttpResponse: ): self.status = status self.headers = headers or {} + self.closed = False self.content = MockContent( content_chunks, exception_to_raise, exception_at_chunk ) + def close(self): + self.closed = True + async def __aexit__(self, exc_type, exc_val, exc_tb): pass @@ -613,3 +617,64 @@ async def test_handle_session_closed_during_request(): assert counts["requests"] == 2 # First request failed, second succeeded assert counts["sessions"] == 2 # Created 2 sessions for retry assert response.status_code == 200 + + +@pytest.mark.asyncio +async def test_response_stream_closes_response_on_error(): + """ + Regression test for #30192: when body iteration ends with an error, the + underlying aiohttp response must be closed so its connector slot is + released. Leaked slots exhaust the pool and every later request times + out (408) until the proxy restarts, even after the backend recovers. + """ + mock_response = MockAiohttpResponse( + content_chunks=[b"chunk1", b"chunk2"], + exception_to_raise=aiohttp.ServerTimeoutError("read timeout"), + exception_at_chunk=1, + ) + + stream = AiohttpResponseStream(mock_response) # type: ignore + with pytest.raises(httpx.TimeoutException): + async for _ in stream: + pass + + assert mock_response.closed is True + + +@pytest.mark.asyncio +async def test_response_stream_closes_response_on_cancellation(): + """ + Regression test for #30192: a task cancelled mid-stream (e.g. the caller + disconnects during a traffic spike) must not leak its aiohttp connection. + """ + mock_response = MockAiohttpResponse( + content_chunks=[b"chunk1", b"chunk2", b"chunk3"], + exception_to_raise=asyncio.CancelledError(), + exception_at_chunk=1, + ) + + stream = AiohttpResponseStream(mock_response) # type: ignore + with pytest.raises(asyncio.CancelledError): + async for _ in stream: + pass + + assert mock_response.closed is True + + +@pytest.mark.asyncio +async def test_response_stream_closes_response_on_generator_exit(): + """ + Regression test for #30192: when the consumer stops iterating early and the + stream generator is closed (GeneratorExit), the underlying aiohttp response + must still be closed so its connector slot is released. + """ + mock_response = MockAiohttpResponse( + content_chunks=[b"chunk1", b"chunk2", b"chunk3"], + ) + + stream = AiohttpResponseStream(mock_response) # type: ignore + iterator = stream.__aiter__() + assert await iterator.__anext__() == b"chunk1" + await iterator.aclose() + + assert mock_response.closed is True