mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
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) <noreply@anthropic.com>
This commit is contained in:
parent
fe080a86b2
commit
d725179b95
2 changed files with 102 additions and 0 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue