fix(logging): route realtime success logging through the bounded worker (#31733)

RealTimeStreaming.log_messages dispatched the success handler with a bare
asyncio.create_task, bypassing GLOBAL_LOGGING_WORKER (which gives a per-coroutine
timeout and a concurrency cap). On a long-lived realtime websocket a slow logging
callback left one suspended task per logged turn, each pinning that turn's
assembled response, accumulating without bound (~12-15k in-flight under load in a
repro) until OOM. Route realtime success logging through the bounded worker so
in-flight logging is capped and a hung callback is cancelled at the worker
timeout.

The chat and responses streaming success-logging paths are intentionally left
unchanged: their success callbacks must complete within the call's event-loop run
(the non-streaming path pairs the worker with a synchronous callback; the
streaming path has no such companion), so deferring them through the worker would
drop logs for one-shot SDK calls and breaks test_async_custom_handler_stream.
Bounding those paths needs a load-shedding approach and is left to a follow-up.

(cherry picked from commit d4c33b2b59)
This commit is contained in:
mubashir1osmani 2026-06-30 12:54:47 -07:00 • committed by Yuneng Jiang
parent 13285671f2
commit 9426f7b9dd
No known key found for this signature in database
2 changed files with 28 additions and 2 deletions

View file

@ -5,6 +5,7 @@ from typing import TYPE_CHECKING, Any, Dict, List, Optional, Union, cast
import litellm
from litellm._logging import verbose_logger
from litellm.litellm_core_utils.logging_worker import GLOBAL_LOGGING_WORKER
from litellm.llms.base_llm.realtime.transformation import BaseRealtimeConfig
from litellm.types.llms.openai import (
OpenAIRealtimeEvents,
@ -327,8 +328,10 @@ class RealTimeStreaming:
self.tool_calls
)
## ASYNC LOGGING
# Create an event loop for the new thread
asyncio.create_task(self.logging_obj.async_success_handler(self.messages))
# Route through the bounded logging worker (per-coroutine timeout +
# concurrency cap) instead of a bare create_task, so a slow callback
# can't leave suspended tasks pinning each call's response in memory.
GLOBAL_LOGGING_WORKER.ensure_initialized_and_enqueue(self.logging_obj.async_success_handler(self.messages))
## SYNC LOGGING
executor.submit(self.logging_obj.success_handler(self.messages))

View file

@ -2750,3 +2750,26 @@ def test_non_bidi_setup_left_untouched_for_followup_capable_providers():
assert streaming._maybe_inject_guardrail_auto_response_disable(msg) == msg
finally:
litellm.callbacks = []
@pytest.mark.asyncio
async def test_log_messages_routes_async_logging_through_bounded_worker():
"""Realtime success logging must go through GLOBAL_LOGGING_WORKER (bounded
queue + per-coroutine timeout), not a bare asyncio.create_task. A bare task
has no timeout/concurrency cap, so when a logging callback is slow every
realtime turn leaves a suspended task pinning its response in memory -> an
unbounded leak. Regression for that fix."""
logging_obj = MagicMock()
streaming = RealTimeStreaming(MagicMock(), MagicMock(), logging_obj)
streaming.messages = [{"type": "session.created"}]
with (
patch("litellm.litellm_core_utils.realtime_streaming.GLOBAL_LOGGING_WORKER") as mock_worker,
patch("litellm.litellm_core_utils.realtime_streaming.asyncio.create_task") as mock_create_task,
patch("litellm.litellm_core_utils.realtime_streaming.executor.submit"),
):
await streaming.log_messages()
mock_worker.ensure_initialized_and_enqueue.assert_called_once()
# the bare create_task path must no longer be used for success logging
mock_create_task.assert_not_called()