From d725179b9586306dc816947fc4b9888f51eac20d Mon Sep 17 00:00:00 2001 From: JF Date: Sat, 28 Mar 2026 07:22:15 -0400 Subject: [PATCH] fix(logging_worker): cancel orphaned tasks on event loop change to prevent warnings When the event loop changes (e.g. between asyncio.run() calls or in test suites), _ensure_queue() drops references to the old worker task without cancelling it. The orphaned task, still PENDING on the now-closed loop, triggers "Task was destroyed but it is pending!" warnings and RuntimeError during garbage collection. Similarly, _flush_on_exit() creates a new event loop to drain remaining queue items but never cleans up the old worker task. Fix: add _discard_orphaned_tasks() helper that cancels old tasks (catching RuntimeError if the loop is already closed) and sets _log_destroy_pending = False to suppress the asyncio destroy warning. Called from both _ensure_queue() and _flush_on_exit(). Co-Authored-By: Claude Opus 4.6 (1M context) --- litellm/litellm_core_utils/logging_worker.py | 34 ++++++++++ .../litellm_core_utils/test_logging_worker.py | 68 +++++++++++++++++++ 2 files changed, 102 insertions(+) diff --git a/litellm/litellm_core_utils/logging_worker.py b/litellm/litellm_core_utils/logging_worker.py index 7f00c47c1ff..7aa51068099 100644 --- a/litellm/litellm_core_utils/logging_worker.py +++ b/litellm/litellm_core_utils/logging_worker.py @@ -71,6 +71,12 @@ class LoggingWorker: verbose_logger.debug( "LoggingWorker: Event loop changed, reinitializing queue and worker" ) + # Cancel orphaned tasks bound to the old (likely closed) event + # loop and suppress "Task was destroyed but it is pending!" + # warnings. We cannot await these tasks because their event loop + # is no longer running, so we mark them to skip the __del__ + # warning instead. + self._discard_orphaned_tasks() # Clear old state - these are bound to the old loop self._queue = None self._sem = None @@ -411,6 +417,29 @@ class LoggingWorker: except asyncio.QueueEmpty: break + def _discard_orphaned_tasks(self) -> None: + """Cancel orphaned tasks and suppress their destroy warnings. + + When the event loop changes or the process is exiting, pending tasks + bound to the old (likely closed) loop cannot be properly awaited. + We attempt to cancel them and, regardless of whether cancel() + succeeds (it may raise ``RuntimeError`` if the loop is already + closed), suppress the ``"Task was destroyed but it is pending!"`` + warning by setting ``_log_destroy_pending = False``. + """ + all_tasks = list(self._running_tasks) + if self._worker_task is not None: + all_tasks.append(self._worker_task) + for t in all_tasks: + if not t.done(): + try: + t.cancel() + except RuntimeError: + pass # Event loop is already closed + t._log_destroy_pending = False + self._worker_task = None + self._running_tasks.clear() + def _safe_log(self, level: str, message: str) -> None: """ Safely log a message during shutdown, suppressing errors if logging is closed. @@ -466,6 +495,11 @@ class LoggingWorker: Note: All logging in this method is wrapped to handle cases where logging handlers are closed during shutdown. """ + # Cancel the old worker task — its event loop is already closed. + # Suppress "Task was destroyed but it is pending!" warnings since + # the closed loop cannot process cancellation. + self._discard_orphaned_tasks() + if self._queue is None: self._safe_log("debug", "[LoggingWorker] atexit: No queue initialized") return diff --git a/tests/test_litellm/litellm_core_utils/test_logging_worker.py b/tests/test_litellm/litellm_core_utils/test_logging_worker.py index a44b821db87..db71aa4c03b 100644 --- a/tests/test_litellm/litellm_core_utils/test_logging_worker.py +++ b/tests/test_litellm/litellm_core_utils/test_logging_worker.py @@ -360,3 +360,71 @@ class TestLoggingWorker: assert worker2._bound_loop is not None await worker2.stop() + + @pytest.mark.asyncio + async def test_event_loop_change_cancels_orphaned_tasks(self): + """Test that switching event loops cancels old tasks and suppresses warnings. + + When the event loop changes (e.g. between asyncio.run() calls or in + test suites), _ensure_queue() must cancel the old worker task and set + _log_destroy_pending = False so that garbage-collecting the orphaned + task does not emit "Task was destroyed but it is pending!" warnings. + """ + worker = LoggingWorker(timeout=1.0, max_queue_size=10) + + # Start the worker — creates _worker_task on the current loop + worker.start() + await asyncio.sleep(0.05) + + old_task = worker._worker_task + assert old_task is not None and not old_task.done() + + # Simulate an event loop change: bind worker to a different loop + # so that the next _ensure_queue() call detects the mismatch. + worker._bound_loop = asyncio.new_event_loop() + + # Re-start triggers _ensure_queue() which should cancel the old task + worker.start() + await asyncio.sleep(0.05) + + # The old task should have been cancelled and marked to suppress + # the "Task was destroyed but it is pending!" warning. + assert old_task.cancelled() or old_task.done() + assert old_task._log_destroy_pending is False + + # The worker should be running on the current loop now + assert worker._worker_task is not None + assert worker._worker_task is not old_task + + await worker.stop() + + def test_flush_on_exit_cancels_worker_task(self): + """Test that _flush_on_exit cancels the worker task to avoid warnings. + + When the atexit handler fires, the original event loop is closed. + _flush_on_exit must cancel the orphaned worker task and suppress + its destroy warning before creating a new loop to drain the queue. + """ + loop = asyncio.new_event_loop() + worker = LoggingWorker(timeout=1.0, max_queue_size=10) + + # Start the worker on a loop, then close it (simulating process exit) + loop.run_until_complete(self._start_worker(worker)) + old_task = worker._worker_task + assert old_task is not None + + # Close the loop — this is what happens before atexit fires + loop.close() + + # Now _flush_on_exit should clean up the orphaned task without + # raising RuntimeError from the closed event loop. + worker._flush_on_exit() + + assert old_task._log_destroy_pending is False + assert worker._worker_task is None + assert len(worker._running_tasks) == 0 + + @staticmethod + async def _start_worker(worker: LoggingWorker): + worker.start() + await asyncio.sleep(0.05)