diff --git a/litellm/llms/custom_httpx/aiohttp_transport.py b/litellm/llms/custom_httpx/aiohttp_transport.py index adac6a1b276e..7f2c286882b9 100644 --- a/litellm/llms/custom_httpx/aiohttp_transport.py +++ b/litellm/llms/custom_httpx/aiohttp_transport.py @@ -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): @@ -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, diff --git a/tests/autofix/test_aiohttp_transport_session_leak.py b/tests/autofix/test_aiohttp_transport_session_leak.py new file mode 100644 index 000000000000..a07137de4834 --- /dev/null +++ b/tests/autofix/test_aiohttp_transport_session_leak.py @@ -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()