From b8f344ec6bea5a8ea7848e69ee2f3d2430855d5f Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 22:38:48 +0530 Subject: [PATCH 1/4] Release unfinished events when a BroadcastQueue subscriber is removed Add release method to mark unfinished tasks as done. --- livekit-rtc/livekit/rtc/_utils.py | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/livekit-rtc/livekit/rtc/_utils.py b/livekit-rtc/livekit/rtc/_utils.py index 91175815..ab01dec2 100644 --- a/livekit-rtc/livekit/rtc/_utils.py +++ b/livekit-rtc/livekit/rtc/_utils.py @@ -113,6 +113,15 @@ async def wait_for(self, fnc: Callable[[T], bool]) -> T: self.task_done() + def release(self) -> None: + """Mark every unfinished item as done. + + Lets a ``join()`` that is already waiting on this queue return once there + is no consumer left for it. + """ + while self._unfinished_tasks > 0: + self.task_done() + class BroadcastQueue(Generic[T]): """Queue with multiple subscribers.""" @@ -135,6 +144,10 @@ def subscribe(self) -> Queue[T]: def unsubscribe(self, queue: Queue[T]) -> None: self._subscribers.remove(queue) + # A join() may already hold this queue. Events that are still queued, or that + # the subscriber took but never finished, have no consumer anymore, so + # release them instead of leaving the joiner waiting forever. + queue.release() async def join(self) -> None: async with self._lock: From 576b416aa0ba2505847383ec9cc63035402c5a0f Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 22:39:16 +0530 Subject: [PATCH 2/4] Add BroadcastQueue unsubscribe tests Added unit tests for BroadcastQueue subscriber removal behavior, ensuring proper handling of joins and unsubscribes. --- livekit-rtc/tests/test_broadcast_queue.py | 98 +++++++++++++++++++++++ 1 file changed, 98 insertions(+) create mode 100644 livekit-rtc/tests/test_broadcast_queue.py diff --git a/livekit-rtc/tests/test_broadcast_queue.py b/livekit-rtc/tests/test_broadcast_queue.py new file mode 100644 index 00000000..e28dab14 --- /dev/null +++ b/livekit-rtc/tests/test_broadcast_queue.py @@ -0,0 +1,98 @@ +"""Unit tests for BroadcastQueue subscriber removal. + +``Room._listen_task`` calls ``BroadcastQueue.join()`` after every event, so a +subscriber that goes away (for example a cancelled ``publish_track``) must not +leave that join waiting on events nobody will finish. +""" + +from __future__ import annotations + +import asyncio + +import pytest + +from livekit.rtc._utils import BroadcastQueue + + +async def _join_finishes(queue: BroadcastQueue, timeout: float = 1.0) -> bool: + try: + await asyncio.wait_for(queue.join(), timeout) + except asyncio.TimeoutError: + return False + return True + + +@pytest.mark.asyncio +async def test_join_waits_for_a_subscriber_that_has_not_finished(): + broadcast: BroadcastQueue[int] = BroadcastQueue() + subscriber = broadcast.subscribe() + broadcast.put_nowait(1) + + assert not await _join_finishes(broadcast, timeout=0.1) + + await subscriber.get() + subscriber.task_done() + assert await _join_finishes(broadcast) + + +@pytest.mark.asyncio +async def test_unsubscribe_releases_a_join_waiting_on_a_queued_event(): + broadcast: BroadcastQueue[int] = BroadcastQueue() + subscriber = broadcast.subscribe() + broadcast.put_nowait(1) + + joiner = asyncio.create_task(broadcast.join()) + await asyncio.sleep(0) + assert not joiner.done() + + broadcast.unsubscribe(subscriber) + + await asyncio.wait_for(joiner, 1.0) + assert broadcast.len_subscribers() == 0 + + +@pytest.mark.asyncio +async def test_unsubscribe_releases_a_join_waiting_on_an_event_that_was_taken(): + broadcast: BroadcastQueue[int] = BroadcastQueue() + subscriber = broadcast.subscribe() + broadcast.put_nowait(1) + await subscriber.get() # taken, but the subscriber never calls task_done() + + joiner = asyncio.create_task(broadcast.join()) + await asyncio.sleep(0) + assert not joiner.done() + + broadcast.unsubscribe(subscriber) + + await asyncio.wait_for(joiner, 1.0) + + +@pytest.mark.asyncio +async def test_unsubscribe_keeps_other_subscribers_in_the_join(): + broadcast: BroadcastQueue[int] = BroadcastQueue() + gone = broadcast.subscribe() + staying = broadcast.subscribe() + broadcast.put_nowait(1) + + joiner = asyncio.create_task(broadcast.join()) + await asyncio.sleep(0) + broadcast.unsubscribe(gone) + await asyncio.sleep(0.05) + assert not joiner.done() + + await staying.get() + staying.task_done() + await asyncio.wait_for(joiner, 1.0) + + +@pytest.mark.asyncio +async def test_unsubscribe_after_every_event_is_finished_is_a_no_op(): + broadcast: BroadcastQueue[int] = BroadcastQueue() + subscriber = broadcast.subscribe() + broadcast.put_nowait(1) + await subscriber.get() + subscriber.task_done() + + broadcast.unsubscribe(subscriber) + + assert await _join_finishes(broadcast) From 374423a245551b5a02f82f8c8d4cb201fac47df3 Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 22:57:59 +0530 Subject: [PATCH 3/4] Fix type-check: release() without private _unfinished_tasks --- livekit-rtc/livekit/rtc/_utils.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/livekit-rtc/livekit/rtc/_utils.py b/livekit-rtc/livekit/rtc/_utils.py index ab01dec2..088fb0e1 100644 --- a/livekit-rtc/livekit/rtc/_utils.py +++ b/livekit-rtc/livekit/rtc/_utils.py @@ -119,8 +119,12 @@ def release(self) -> None: Lets a ``join()`` that is already waiting on this queue return once there is no consumer left for it. """ - while self._unfinished_tasks > 0: - self.task_done() + while True: + try: + self.task_done() + except ValueError: + # task_done() raises once every item is accounted for + return class BroadcastQueue(Generic[T]): From 20e92466e432e5ec8e9481771aff970ac9a14a26 Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 22:58:24 +0530 Subject: [PATCH 4/4] Fix type-check: annotate test return types --- livekit-rtc/tests/test_broadcast_queue.py | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/livekit-rtc/tests/test_broadcast_queue.py b/livekit-rtc/tests/test_broadcast_queue.py index e28dab14..ff7b3745 100644 --- a/livekit-rtc/tests/test_broadcast_queue.py +++ b/livekit-rtc/tests/test_broadcast_queue.py @@ -23,7 +23,7 @@ async def _join_finishes(queue: BroadcastQueue, timeout: float = 1.0) -> bool: @pytest.mark.asyncio -async def test_join_waits_for_a_subscriber_that_has_not_finished(): +async def test_join_waits_for_a_subscriber_that_has_not_finished() -> None: broadcast: BroadcastQueue[int] = BroadcastQueue() subscriber = broadcast.subscribe() broadcast.put_nowait(1) @@ -36,7 +36,7 @@ async def test_join_waits_for_a_subscriber_that_has_not_finished(): @pytest.mark.asyncio -async def test_unsubscribe_releases_a_join_waiting_on_a_queued_event(): +async def test_unsubscribe_releases_a_join_waiting_on_a_queued_event() -> None: broadcast: BroadcastQueue[int] = BroadcastQueue() subscriber = broadcast.subscribe() broadcast.put_nowait(1) @@ -52,7 +52,7 @@ async def test_unsubscribe_releases_a_join_waiting_on_a_queued_event(): @pytest.mark.asyncio -async def test_unsubscribe_releases_a_join_waiting_on_an_event_that_was_taken(): +async def test_unsubscribe_releases_a_join_waiting_on_an_event_that_was_taken() -> None: broadcast: BroadcastQueue[int] = BroadcastQueue() subscriber = broadcast.subscribe() broadcast.put_nowait(1) @@ -68,7 +68,7 @@ async def test_unsubscribe_releases_a_join_waiting_on_an_event_that_was_taken(): @pytest.mark.asyncio -async def test_unsubscribe_keeps_other_subscribers_in_the_join(): +async def test_unsubscribe_keeps_other_subscribers_in_the_join() -> None: broadcast: BroadcastQueue[int] = BroadcastQueue() gone = broadcast.subscribe() staying = broadcast.subscribe() @@ -86,7 +86,7 @@ async def test_unsubscribe_keeps_other_subscribers_in_the_join(): @pytest.mark.asyncio -async def test_unsubscribe_after_every_event_is_finished_is_a_no_op(): +async def test_unsubscribe_after_every_event_is_finished_is_a_no_op() -> None: broadcast: BroadcastQueue[int] = BroadcastQueue() subscriber = broadcast.subscribe() broadcast.put_nowait(1)