fix(logging_worker): make flush() survive an event loop change

flush() awaited join() on whatever queue the worker held, even one bound to
an event loop that has since closed. Its unfinished counter is never
decremented on the new loop, so the first flush() after a loop change hung
until pytest-timeout killed it and every later one raised "is bound to a
different event loop" from the queue's Event. The CircleCI unit job has
been red on every branch since the first tests that flush without
enqueueing landed, and an SDK script that flushes from a second
asyncio.run() hangs the same way.

flush() now goes through start() first, which carries the tasks stranded
on the previous loop onto the current one and guarantees a worker there to
drain them, the same loop-change handling every other entry point already
had.
This commit is contained in:
mateo-berri 2026-09-21 15:44:26 -07:00
parent 3d26a29a1a
commit 212ab630b8
2 changed files with 40 additions and 0 deletions

View file

@ -484,9 +484,15 @@ class LoggingWorker:
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.
``start()`` runs first so that, after an event loop change, the tasks
still on the previous loop's queue move onto this loop and a worker
here drains them; joining the old queue directly would wait on a
counter nothing on this loop ever decrements.
"""
if self._queue is None:
return
self.start()
await self._queue.join()
async def clear_queue(self):

View file

@ -180,6 +180,40 @@ class TestLoggingWorker:
assert sorted(fired) == ["first", "second"]
@pytest.mark.parametrize("stranded", ["still_queued", "dequeued_never_started"])
def test_flush_on_new_loop_drains_tasks_stranded_on_previous_loop(self, stranded):
"""
Regression: ``flush()`` from a new event loop used to ``join()`` the queue bound to the
previous loop, whose unfinished counter nothing on the new loop ever decrements. The first
such flush hung until pytest-timeout killed it and every later one raised
``RuntimeError: ... is bound to a different event loop`` from the queue's Event.
"""
worker = LoggingWorker(timeout=1.0, max_queue_size=10)
fired = []
async def marker():
fired.append(True)
async def enqueue_on_first_loop():
if stranded == "still_queued":
worker._ensure_queue()
worker.enqueue(marker())
return
worker.ensure_initialized_and_enqueue(marker())
asyncio.run(enqueue_on_first_loop())
assert worker._queue is not None
expected_shape = (1, 0) if stranded == "still_queued" else (0, 1)
assert (worker._queue.qsize(), len(worker._unstarted_dequeued_tasks())) == expected_shape
assert fired == [], "precondition: the callback never ran before the first loop closed"
async def flush_on_second_loop():
await asyncio.wait_for(worker.flush(), timeout=5)
asyncio.run(flush_on_second_loop())
assert fired == [True]
def test_flush_on_exit_swallows_cancellation_and_drains_remaining(self):
"""A callback raising CancelledError must not abort the atexit flush of later events."""
worker = LoggingWorker(timeout=1.0, max_queue_size=10)