fix(http_handler): dispose aiohttp session when AsyncHTTPHandler is finalized without a running loop (#36670)

* fix(http_handler): dispose aiohttp session when finalized without a running loop

AsyncHTTPHandler.__del__ can only schedule an async close when a running
event loop exists at finalization time; in any other context (worker
threads whose loop has closed, sync contexts, interpreter shutdown) the
RuntimeError from get_running_loop() is swallowed and the underlying
aiohttp ClientSession is abandoned to GC, emitting 'Unclosed client
session' / 'Unclosed connector' warnings.

This is the disposal gap left after the recycle-time fix: clients created
for short-lived event loops (the loop-id-keyed LLM client cache mints one
handler per loop) are never recycled - they live and die with their loop,
and their finalization is precisely the loop-less case.

Fix:
- no running loop: fall back to the connector's synchronous teardown via
  LiteLLMAiohttpTransport._mark_connector_closed - the same finalizer-safe
  path used for dead-loop recycles - honoring _owns_session so a shared
  session is never closed.
- running loop: keep the async close, but hold a strong reference to the
  scheduled task until it completes (a bare create_task() result may be
  collected before running), mirroring _background_close_tasks.

Tests: loop-less finalization closes a dead-loop session; running-loop
finalization registers and drains the close task; the sync fallback
respects session ownership. All three fail without the fix.

* lint: conform new finalizer code to the type-discipline budget

Final on the five never-rebound locals (LIT010); the class-level task
registry keeps its mutable set with the sanctioned mutable-ok reason,
mirroring the aiohttp transport's registry (LIT001).

* lint: reasoned pyright ignore on the cross-class teardown call

The handler deliberately reuses the transport's finalizer-safe connector
teardown; no public seam exists and an async close can never run at
loop-less finalization. Clears the net-new reportPrivateUsage the
basedpyright budget gate flagged once the LIT stage passed.

* fix(http_handler): retrieve exceptions from finalizer close tasks

A bare discard done-callback dropped the task without consuming its
exception, so a failing aclose() emitted "Task exception was never
retrieved" at GC, the same noise class this path exists to remove.
Mirror the transport's _on_close_task_done: discard, early-return on
cancellation, retrieve and debug-log the exception.

* fix(http_handler): dispose foreign-loop sessions instead of scheduling aclose on the live loop

GC on a live loop (e.g. the app's) of a handler whose session belongs to
another, possibly dead, loop scheduled aclose() on the current loop, the
cross-loop path the transport refuses. Route both that case and the
loop-less case through the transport's lifecycle-aware
_close_recycled_session, which picks async close on the session's own
loop, threadsafe handoff, or the synchronous connector teardown.

Regression test: a dead-loop session collected while another loop runs
is disposed without scheduling anything on that loop.

* chore: retrigger CI (test_mcp_logging payload-order flake, also failed on litellm_spendlogs_fallback_metadata minutes earlier)

* test(mcp): select the MCP tool-call payload instead of the last-delivered one

TestMCPLogger kept a single last-writer slot; an async success event from
another call (a mocked acompletion whose log task lands late) races the
MCP event for it, so the cost assertions intermittently read the wrong
payload. This PR's finalizer change shifts task interleaving on the loop
and tips that latent race over (also seen on an unrelated PR minutes
earlier). Collect call_type=call_mcp_tool payloads in their own list and
assert on those.

* test(mcp): MCPLoggerHook inherits the order-independent payload capture

It duplicated TestMCPLogger's init and success handler verbatim; the
hook test reads the same MCP payload selection, so subclass instead.
This commit is contained in:
Anmol Jaiswal 2026-08-25 20:42:10 +05:30 • committed by GitHub
parent 31a67561ab
commit bb27bfd9a7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 246 additions and 39 deletions

View file

@ -9,7 +9,7 @@ import threading
import time
from collections.abc import AsyncIterable, Callable, Iterable, Mapping
from http.cookiejar import CookieJar, DefaultCookiePolicy
from typing import TYPE_CHECKING, Any, Final, Optional, TypeAlias, TypedDict
from typing import TYPE_CHECKING, Any, ClassVar, Final, Optional, TypeAlias, TypedDict
import certifi
import httpx
@ -933,11 +933,83 @@ class AsyncHTTPHandler:
response.raise_for_status()
return response
# Strong references to finalizer-scheduled client-close tasks. A bare
# create_task() result may be garbage-collected before it runs, leaving
# the underlying aiohttp session unclosed ("Unclosed client session").
# Mirrors LiteLLMAiohttpTransport._background_close_tasks.
_finalizer_close_tasks: ClassVar[set["asyncio.Task[None]"]] = set() # mutable-ok: strong refs for pending closes
@classmethod
def _on_finalizer_close_done(cls, task: "asyncio.Task[None]") -> None:
cls._finalizer_close_tasks.discard(task)
if task.cancelled():
return
exc: Final = task.exception()
if exc is not None:
verbose_logger.debug("Error closing client at finalization: %s", exc)
def _aiohttp_session_bound_elsewhere(self, loop: asyncio.AbstractEventLoop) -> bool:
"""True when the wrapped aiohttp session is bound to a loop other than
``loop`` — awaiting ``aclose()`` here would touch that loop's internals."""
from litellm.llms.custom_httpx.aiohttp_transport import (
LiteLLMAiohttpTransport,
)
transport: Final = getattr(self._client, "_transport", None)
if not isinstance(transport, LiteLLMAiohttpTransport):
return False
session: Final = transport.client
if not isinstance(session, ClientSession) or session.closed:
return False
return getattr(session, "_loop", None) is not loop
def _dispose_wrapped_aiohttp_session(self) -> None:
"""Dispose the wrapped aiohttp session when ``aclose()`` cannot run here.
Finalization either has no running loop, or a loop the session is not
bound to. Delegating to the transport's lifecycle-aware disposal picks
the safe path per session state (async close on its own loop, threadsafe
handoff to a loop running elsewhere, or the synchronous connector
teardown that flips the flags ``ClientSession.__del__`` checks), so no
"Unclosed client session" / "Unclosed connector" warnings fire at
garbage collection.
"""
from litellm.llms.custom_httpx.aiohttp_transport import (
LiteLLMAiohttpTransport,
)
transport: Final = getattr(self._client, "_transport", None)
if not isinstance(transport, LiteLLMAiohttpTransport):
return
# A shared session (e.g. the proxy's) is never this handler's to close.
if not getattr(transport, "_owns_session", False):
return
session: Final = transport.client
if isinstance(session, ClientSession) and not session.closed:
transport._close_recycled_session(session) # pyright: ignore[reportPrivateUsage] # deliberate reuse of the transport's lifecycle-aware disposal; an async close can never run in this context
def __del__(self) -> None:
try:
if not _handler_may_close_client(sys.getrefcount(self._client), self._owns_client):
return
asyncio.get_running_loop().create_task(self._client.aclose())
try:
loop: Final = asyncio.get_running_loop()
except RuntimeError:
# No running loop at finalization time (worker threads after
# their loop closed, interpreter/worker shutdown, GC in a
# sync context). An async close can never run here.
self._dispose_wrapped_aiohttp_session()
return
if self._aiohttp_session_bound_elsewhere(loop):
# GC ran on a live loop (e.g. the app's) but the session
# belongs to another, possibly dead, loop — awaiting aclose()
# here is the cross-loop path the transport refuses.
self._dispose_wrapped_aiohttp_session()
return
task: Final = loop.create_task(self._client.aclose())
cls: Final = type(self)
cls._finalizer_close_tasks.add(task)
task.add_done_callback(cls._on_finalizer_close_done)
except Exception:
pass

View file

@ -24,12 +24,20 @@ from mcp.types import Tool as MCPTool, CallToolResult, TextContent
class TestMCPLogger(CustomLogger):
def __init__(self):
self.standard_logging_payload = None
self.mcp_tool_call_payloads = []
super().__init__()
async def async_log_success_event(self, kwargs, response_obj, start_time, end_time):
print("success event")
self.standard_logging_payload = kwargs.get("standard_logging_object", None)
print(f"Captured standard_logging_payload: {self.standard_logging_payload}")
payload = kwargs.get("standard_logging_object", None)
self.standard_logging_payload = payload
# Async success events from other calls (e.g. a mocked acompletion whose
# log task is delivered late) race with the MCP event for the single
# last-writer slot; keep MCP tool calls in their own list so assertions
# are order-independent.
if payload is not None and payload.get("call_type") == "call_mcp_tool":
self.mcp_tool_call_payloads.append(payload)
print(f"Captured standard_logging_payload: {payload}")
def _set_authorized_user(server_ids):
@ -138,7 +146,11 @@ async def test_mcp_cost_tracking():
# wait 1-2 seconds for logging to be processed
await asyncio.sleep(2)
logged_standard_logging_payload = test_logger.standard_logging_payload
logged_standard_logging_payload = (
test_logger.mcp_tool_call_payloads[-1]
if test_logger.mcp_tool_call_payloads
else None
)
print("logged_standard_logging_payload", logged_standard_logging_payload)
# Add assertions
@ -277,7 +289,11 @@ async def test_mcp_cost_tracking_per_tool():
# wait for logging to be processed
await asyncio.sleep(2)
logged_standard_logging_payload_1 = test_logger.standard_logging_payload
logged_standard_logging_payload_1 = (
test_logger.mcp_tool_call_payloads[-1]
if test_logger.mcp_tool_call_payloads
else None
)
print(
"logged_standard_logging_payload_1", logged_standard_logging_payload_1
)
@ -290,6 +306,7 @@ async def test_mcp_cost_tracking_per_tool():
# Reset logger for second test
test_logger.standard_logging_payload = None
test_logger.mcp_tool_call_payloads.clear()
# Test 2: Call cheap_tool - should cost 0.1
response2 = await mcp_server_tool_call(
@ -300,7 +317,11 @@ async def test_mcp_cost_tracking_per_tool():
# wait for logging to be processed
await asyncio.sleep(2)
logged_standard_logging_payload_2 = test_logger.standard_logging_payload
logged_standard_logging_payload_2 = (
test_logger.mcp_tool_call_payloads[-1]
if test_logger.mcp_tool_call_payloads
else None
)
print(
"logged_standard_logging_payload_2", logged_standard_logging_payload_2
)
@ -329,16 +350,7 @@ async def test_mcp_cost_tracking_per_tool():
assert mock_client.call_tool.call_count == 2
class MCPLoggerHook(CustomLogger):
def __init__(self):
self.standard_logging_payload = None
super().__init__()
async def async_log_success_event(self, kwargs, response_obj, start_time, end_time):
print("success event")
self.standard_logging_payload = kwargs.get("standard_logging_object", None)
print(f"Captured standard_logging_payload: {self.standard_logging_payload}")
class MCPLoggerHook(TestMCPLogger):
async def async_post_mcp_tool_call_hook(
self, kwargs, response_obj: MCPPostCallResponseObject, start_time, end_time
) -> Optional[MCPPostCallResponseObject]:
@ -436,7 +448,11 @@ async def test_mcp_tool_call_hook():
await asyncio.sleep(2)
# check logged standard logging payload
logged_standard_logging_payload = test_logger.standard_logging_payload
logged_standard_logging_payload = (
test_logger.mcp_tool_call_payloads[-1]
if test_logger.mcp_tool_call_payloads
else None
)
print("logged_standard_logging_payload", logged_standard_logging_payload)
assert (
logged_standard_logging_payload is not None

View file

@ -56,9 +56,7 @@ async def test_async_post_streaming_status_error_should_not_wait_forever_for_bod
litellm_handler = AsyncHTTPHandler()
await litellm_handler.client.aclose()
litellm_handler.client = httpx.AsyncClient(
transport=httpx.MockTransport(mock_handler)
)
litellm_handler.client = httpx.AsyncClient(transport=httpx.MockTransport(mock_handler))
try:
with pytest.raises(MaskedHTTPStatusError) as exc_info:
await asyncio.wait_for(
@ -202,9 +200,7 @@ async def test_ssl_verification_with_aiohttp_transport(monkeypatch: pytest.Monke
transport_connector = transport._get_valid_client_session().connector
assert isinstance(transport_connector, TCPConnector)
aiohttp_session = aiohttp.ClientSession(
connector=aiohttp.TCPConnector(ssl=False)
)
aiohttp_session = aiohttp.ClientSession(connector=aiohttp.TCPConnector(ssl=False))
try:
aiohttp_connector = aiohttp_session.connector
assert isinstance(aiohttp_connector, aiohttp.TCPConnector)
@ -378,7 +374,8 @@ async def test_get_async_httpx_client_with_shared_session():
# Test with shared session
client = get_async_httpx_client(
llm_provider=LlmProviders.ANTHROPIC, shared_session=mock_session # type: ignore
llm_provider=LlmProviders.ANTHROPIC,
shared_session=mock_session, # type: ignore
)
# Verify the client was created successfully
@ -397,9 +394,7 @@ async def test_get_async_httpx_client_without_shared_session():
from litellm.types.utils import LlmProviders
# Test without shared session
client = get_async_httpx_client(
llm_provider=LlmProviders.ANTHROPIC, shared_session=None
)
client = get_async_httpx_client(llm_provider=LlmProviders.ANTHROPIC, shared_session=None)
# Verify the client was created successfully
assert client is not None
@ -476,11 +471,13 @@ async def test_session_reuse_integration():
# Create two clients with the same session
client1 = get_async_httpx_client(
llm_provider=LlmProviders.ANTHROPIC, shared_session=mock_session # type: ignore
llm_provider=LlmProviders.ANTHROPIC,
shared_session=mock_session, # type: ignore
)
client2 = get_async_httpx_client(
llm_provider=LlmProviders.OPENAI, shared_session=mock_session # type: ignore
llm_provider=LlmProviders.OPENAI,
shared_session=mock_session, # type: ignore
)
# Both clients should be created successfully
@ -512,9 +509,7 @@ async def test_session_reuse_integration():
(None, None, None, False), # None value - skip configuration
],
)
def test_ssl_ecdh_curve(
env_curve, litellm_curve, expected_curve, should_call, monkeypatch
):
def test_ssl_ecdh_curve(env_curve, litellm_curve, expected_curve, should_call, monkeypatch):
"""Test SSL ECDH curve configuration with valid curves and precedence"""
from litellm.llms.custom_httpx.http_handler import _ssl_context_cache
@ -717,9 +712,7 @@ class TestDefaultCachedClientTimeoutHonorsRequestTimeout:
_default_cached_client_timeout,
)
monkeypatch.setattr(
litellm, "request_timeout", litellm.constants.DEFAULT_REQUEST_TIMEOUT_SECONDS
)
monkeypatch.setattr(litellm, "request_timeout", litellm.constants.DEFAULT_REQUEST_TIMEOUT_SECONDS)
monkeypatch.setattr(litellm, "request_timeout_explicitly_set", False)
assert _default_cached_client_timeout() is _DEFAULT_TIMEOUT
@ -734,9 +727,7 @@ class TestDefaultCachedClientTimeoutHonorsRequestTimeout:
assert resolved.read == 300.0
assert resolved.connect == 5.0
def test_cached_async_client_built_with_explicit_request_timeout(
self, monkeypatch: pytest.MonkeyPatch
):
def test_cached_async_client_built_with_explicit_request_timeout(self, monkeypatch: pytest.MonkeyPatch):
from litellm.caching.llm_caching_handler import LLMClientCache
from litellm.llms.custom_httpx.http_handler import get_async_httpx_client
from litellm.types.utils import LlmProviders
@ -1195,3 +1186,131 @@ async def test_aiohttp_session_never_replays_one_upstreams_cookie_to_another():
assert len(jar) == 0
assert dict(jar.filter_cookies(URL("https://upstream-a.example.com"))) == {}
await session.close()
def _mint_session_on_dead_loop(handler: AsyncHTTPHandler) -> ClientSession:
"""Create the transport's real ClientSession on a loop that then closes.
This is the lifecycle of every client minted for a short-lived event loop
(the loop-id-keyed LLM client cache creates one handler per loop): the
session outlives its loop and can only ever be disposed loop-lessly.
"""
transport = handler.client._transport
assert isinstance(transport, LiteLLMAiohttpTransport)
loop = asyncio.new_event_loop()
async def _create() -> ClientSession:
return transport._get_valid_client_session()
session = loop.run_until_complete(_create())
loop.close()
return session
def test_finalizer_without_running_loop_closes_dead_loop_session():
"""A handler finalized with no running event loop must still dispose its
aiohttp session.
The async close can never run in that context; without the synchronous
fallback the session and its connector are abandoned to GC and emit
"Unclosed client session" / "Unclosed connector" warnings."""
handler = AsyncHTTPHandler(timeout=61.0)
session = _mint_session_on_dead_loop(handler)
assert not session.closed
del handler
gc.collect()
assert session.closed
@pytest.mark.asyncio
async def test_finalizer_with_running_loop_schedules_close_and_holds_task_ref():
"""With a running loop, finalization schedules an async close and must keep
a strong reference to the task until it completes — a bare create_task()
result may be collected before it runs, leaving the session unclosed."""
handler = AsyncHTTPHandler(timeout=61.0)
transport = handler.client._transport
assert isinstance(transport, LiteLLMAiohttpTransport)
session = transport._get_valid_client_session()
assert not session.closed
del transport
baseline_tasks = set(AsyncHTTPHandler._finalizer_close_tasks)
del handler
gc.collect()
scheduled = AsyncHTTPHandler._finalizer_close_tasks - baseline_tasks
assert len(scheduled) == 1
await asyncio.gather(*scheduled)
assert session.closed
assert not (AsyncHTTPHandler._finalizer_close_tasks & scheduled)
@pytest.mark.asyncio
async def test_sync_close_helper_respects_session_ownership():
"""The loop-less fallback closes only sessions the transport owns; a
shared session (e.g. the proxy's) must never be closed by a handler."""
owned_handler = AsyncHTTPHandler(timeout=61.0)
owned_transport = owned_handler.client._transport
assert isinstance(owned_transport, LiteLLMAiohttpTransport)
owned_session = owned_transport._get_valid_client_session()
baseline = set(LiteLLMAiohttpTransport._background_close_tasks)
owned_handler._dispose_wrapped_aiohttp_session()
scheduled = LiteLLMAiohttpTransport._background_close_tasks - baseline
await asyncio.gather(*scheduled)
assert owned_session.closed
shared_session = ClientSession()
shared_handler = AsyncHTTPHandler(timeout=61.0, shared_session=shared_session)
shared_transport = shared_handler.client._transport
assert isinstance(shared_transport, LiteLLMAiohttpTransport)
assert shared_transport._owns_session is False
shared_handler._dispose_wrapped_aiohttp_session()
assert not shared_session.closed
await shared_session.close()
await shared_handler.close()
await owned_handler.close()
@pytest.mark.asyncio
async def test_finalizer_close_done_consumes_exception():
"""A failing finalizer close must have its exception retrieved by the done
callback, or asyncio emits "Task exception was never retrieved" at GC —
the same log noise the finalizer path exists to eliminate."""
async def failing_close() -> None:
raise RuntimeError("close failed")
task = asyncio.get_running_loop().create_task(failing_close())
AsyncHTTPHandler._finalizer_close_tasks.add(task)
await asyncio.sleep(0)
AsyncHTTPHandler._on_finalizer_close_done(task)
assert task not in AsyncHTTPHandler._finalizer_close_tasks
cancelled = asyncio.get_running_loop().create_task(asyncio.sleep(30))
cancelled.cancel()
await asyncio.sleep(0)
AsyncHTTPHandler._on_finalizer_close_done(cancelled)
@pytest.mark.asyncio
async def test_finalizer_on_live_loop_disposes_foreign_loop_session_without_scheduling():
"""GC on a live loop (e.g. the app's) of a handler whose session belongs to
another, dead loop must not schedule aclose() here — that is the cross-loop
path the transport refuses — and must still dispose the session."""
handler = AsyncHTTPHandler(timeout=61.0)
session = await asyncio.to_thread(_mint_session_on_dead_loop, handler)
assert not session.closed
baseline_tasks = set(AsyncHTTPHandler._finalizer_close_tasks)
del handler
gc.collect()
assert AsyncHTTPHandler._finalizer_close_tasks == baseline_tasks
assert session.closed