mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
fix(router): release aiohttp connection when stream iteration ends abnormally (#30271)
* fix(router): release aiohttp connection when stream iteration ends abnormally A streaming response that terminates with a mid-stream read timeout, a task cancellation (client disconnect), or GeneratorExit never closed the underlying aiohttp ClientResponse. aiohttp only auto-releases the connector slot at body EOF, so each abnormally terminated stream permanently leaked one slot from the shared TCPConnector pool. During a backend traffic spike the pool drains; once exhausted every subsequent request to that host waits for a slot, times out and surfaces as a 408, indefinitely, even after the backend recovers. Only a proxy restart cleared the in-memory sessions, which matched the reported symptom of a router stuck returning 408 for a healthy vLLM backend. Close the response in a finally clause when iteration ends. On a fully read response the connection was already released at EOF and close() is a no-op, so keep-alive reuse for normal requests is unchanged. Fixes #30192 * test(aiohttp): cover GeneratorExit path with a mock instead of a live socket The previous slot-release test started a real aiohttp TCP server, which can flake in offline CI and does not exercise this fix's code path directly. Replace it with a dependency-injected mock that closes the stream generator (GeneratorExit) and asserts the response is closed, covering the third abnormal-exit path the finally block handles
This commit is contained in:
parent
9b2edc5a24
commit
22ecb4cca7
2 changed files with 75 additions and 0 deletions
|
|
@ -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():
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue