fix(logging_worker): restore queue.join() in flush() to wait for in-flight callbacks

This commit is contained in:
Ishaan Jaffer 2026-04-24 16:33:46 -07:00
parent 8090b29a4b
commit 8b92525c91
No known key found for this signature in database

View file

@ -370,11 +370,17 @@ class LoggingWorker:
self._running_tasks.clear()
async def flush(self) -> None:
"""Flush the logging queue."""
"""Flush the logging queue.
Waits until every enqueued task has completed. ``queue.join()`` blocks
on the queue's unfinished-task counter (decremented by ``task_done()``),
so it correctly handles items that have been dequeued but whose
callback hasn't finished yet — ``queue.empty()`` would return True in
that window and cause us to skip the wait.
"""
if self._queue is None:
return
while not self._queue.empty():
await self._queue.join()
await self._queue.join()
async def clear_queue(self):
"""