Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 40 additions & 11 deletions litellm/llms/custom_httpx/aiohttp_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -185,17 +185,7 @@ def _get_valid_client_session(self) -> ClientSession:

# If session is from a different or closed loop, recreate it
if session_loop is None or session_loop != current_loop or session_loop.is_closed():
# Close old session to prevent leaks
old_session = self.client
try:
if not old_session.closed:
try:
asyncio.create_task(old_session.close())
except RuntimeError:
# Different event loop - can't schedule task, rely on GC
verbose_logger.debug("Old session from different loop, relying on GC")
except Exception as e:
verbose_logger.debug(f"Error closing old session: {e}")
self._dispose_stale_session(self.client, session_loop)

# Create a new session in the current event loop
if hasattr(self, "_client_factory") and callable(self._client_factory):
Expand All @@ -212,6 +202,45 @@ def _get_valid_client_session(self) -> ClientSession:

return self.client

@staticmethod
def _dispose_stale_session(session: ClientSession, session_loop: Optional[asyncio.AbstractEventLoop]) -> None:
"""
Release a session that belongs to a different event loop before dropping it.

A ClientSession must be closed on the loop that created it. Scheduling its
async close() on the current loop (or relying on GC) leaves the session
marked open and its connector's sockets unreleased, which surfaces as
aiohttp's "Unclosed client session" warning and leaks file descriptors.

When the owning loop is still running we schedule close() on that loop.
When it is gone (None) or already closed, its selector and sockets are
already torn down, so we detach the connector synchronously to mark the
session closed and drop the dangling reference without touching the dead loop.
"""
if session.closed:
return

if session_loop is not None and not session_loop.is_closed():
try:
session_loop.call_soon_threadsafe(lambda: asyncio.ensure_future(session.close(), loop=session_loop))
return
except RuntimeError as e:
verbose_logger.debug(f"Could not schedule close on owning loop, detaching instead: {e}")

detach = getattr(session, "detach", None)
if callable(detach):
detach()
else: # pragma: no cover - defensive for older aiohttp without detach()
connector = getattr(session, "_connector", None)
if connector is not None:
connector_close = getattr(connector, "_close", None) or getattr(connector, "close", None)
if callable(connector_close):
try:
connector_close()
except Exception as e:
verbose_logger.debug(f"Error closing stale connector: {e}")
setattr(session, "_connector", None)

async def _make_aiohttp_request(
self,
client_session: ClientSession,
Expand Down
58 changes: 58 additions & 0 deletions tests/autofix/test_aiohttp_transport_session_leak.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
import asyncio
import gc
import os
import sys
import warnings

import aiohttp
import pytest

sys.path.insert(0, os.path.abspath("../../.."))

from litellm.llms.custom_httpx.aiohttp_transport import LiteLLMAiohttpTransport


def _create_session_on_separate_closed_loop() -> aiohttp.ClientSession:
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)

async def _mk() -> aiohttp.ClientSession:
return aiohttp.ClientSession()

session = loop.run_until_complete(_mk())
loop.close()
return session


@pytest.mark.asyncio
async def test_stale_loop_session_is_closed_not_leaked():
"""
Regression: "Unclosed client session".

When _get_valid_client_session finds that the cached session belongs to a
different (here: closed) event loop, it recreates the session for the
current loop. The previous session must be released synchronously; the old
behavior scheduled a fire-and-forget close on the wrong loop (or relied on
GC), leaving the session marked open, which is exactly what triggers
aiohttp's "Unclosed client session" warning and leaks its sockets.
"""
stale_session = _create_session_on_separate_closed_loop()
assert not stale_session.closed

transport = LiteLLMAiohttpTransport(client=lambda: aiohttp.ClientSession())
transport.client = stale_session

with warnings.catch_warnings(record=True) as caught:
warnings.simplefilter("always")
new_session = transport._get_valid_client_session()

assert new_session is not stale_session, "expected a fresh session for the running loop"
assert stale_session.closed, "stale session from the old loop leaked (was not closed)"

del stale_session
gc.collect()
unclosed = [w for w in caught if "Unclosed client session" in str(w.message)]
assert not unclosed, f"aiohttp emitted an unclosed-session warning: {[str(w.message) for w in unclosed]}"

if not new_session.closed:
await new_session.close()
Loading