From 004f92a3c526bddd1b10b267f434efe022351c03 Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Tue, 29 Sep 2026 09:54:48 +0000 Subject: [PATCH 1/2] fix(rtc): keep AudioMixer streams that miss the read timeout --- livekit-rtc/livekit/rtc/audio_mixer.py | 24 +++++++++++++++---- tests/rtc/test_mixer.py | 32 +++++++++++++++++++++++++- 2 files changed, 50 insertions(+), 6 deletions(-) diff --git a/livekit-rtc/livekit/rtc/audio_mixer.py b/livekit-rtc/livekit/rtc/audio_mixer.py index e2f28c6b..0c2986a6 100644 --- a/livekit-rtc/livekit/rtc/audio_mixer.py +++ b/livekit-rtc/livekit/rtc/audio_mixer.py @@ -52,6 +52,10 @@ def __init__( """ self._streams: set[_Stream] = set() self._buffers: dict[_Stream, np.ndarray] = {} + # A read that timed out keeps running here and is awaited again on the + # next round. Cancelling it would also cancel the stream itself when the + # stream is an async generator, which then ends and drops all its audio. + self._pending: dict[_Stream, asyncio.Future[AudioFrame]] = {} self._sample_rate: int = sample_rate self._num_channels: int = num_channels self._chunk_size: int = blocksize if blocksize > 0 else int(sample_rate // 10) @@ -90,6 +94,9 @@ def remove_stream(self, stream: AsyncIterator[AudioFrame]) -> None: """ self._streams.discard(stream) self._buffers.pop(stream, None) + pending = self._pending.pop(stream, None) + if pending is not None: + pending.cancel() def __aiter__(self) -> "AudioMixer": return self @@ -110,6 +117,9 @@ async def aclose(self) -> None: self._mixer_task.cancel() with contextlib.suppress(asyncio.CancelledError): await self._mixer_task + for pending in self._pending.values(): + pending.cancel() + self._pending.clear() def end_input(self) -> None: """ @@ -174,13 +184,17 @@ async def _get_contribution( had_data = buf.shape[0] > 0 exhausted = False while buf.shape[0] < self._chunk_size and not exhausted: - try: - frame = await asyncio.wait_for( - stream.__anext__(), timeout=self._stream_timeout_ms / 1000 - ) - except asyncio.TimeoutError: + pending = self._pending.get(stream) + if pending is None: + pending = asyncio.ensure_future(stream.__anext__()) + self._pending[stream] = pending + done, _ = await asyncio.wait({pending}, timeout=self._stream_timeout_ms / 1000) + if not done: logger.warning(f"AudioMixer: stream {stream} timeout, ignoring") break + self._pending.pop(stream, None) + try: + frame = pending.result() except StopAsyncIteration: exhausted = True break diff --git a/tests/rtc/test_mixer.py b/tests/rtc/test_mixer.py index 1ab75f4b..8a697b90 100644 --- a/tests/rtc/test_mixer.py +++ b/tests/rtc/test_mixer.py @@ -1,9 +1,11 @@ # type: ignore +import asyncio + import numpy as np import pytest -from livekit.rtc import AudioMixer +from livekit.rtc import AudioFrame, AudioMixer from livekit.rtc.utils import sine_wave_generator SAMPLE_RATE = 48000 @@ -56,3 +58,31 @@ async def test_mixer_two_sine_waves(): # Assert that the peaks include 440Hz and 880Hz (with a tolerance of ±5 Hz) assert any(np.isclose(peak_freqs, 440, atol=5)), f"Expected 440 Hz in peaks, got: {peak_freqs}" assert any(np.isclose(peak_freqs, 880, atol=5)), f"Expected 880 Hz in peaks, got: {peak_freqs}" + + +@pytest.mark.asyncio +async def test_mixer_keeps_generator_stream_that_is_slower_than_the_timeout(): + """ + A stream that misses the timeout once must not be lost: its audio is mixed + once it arrives. + """ + + async def slow_stream(): + await asyncio.sleep(0.3) # longer than stream_timeout_ms + for _ in range(5): + yield AudioFrame( + np.ones(BLOCKSIZE, dtype=np.int16).tobytes(), SAMPLE_RATE, 1, BLOCKSIZE + ) + + mixer = AudioMixer(sample_rate=SAMPLE_RATE, num_channels=1, stream_timeout_ms=100) + mixer.add_stream(slow_stream()) + mixer.end_input() + + async def read_all(): + return [frame async for frame in mixer] + + frames = await asyncio.wait_for(read_all(), timeout=5) + await mixer.aclose() + + frames_with_audio = [f for f in frames if np.any(np.frombuffer(f.data.tobytes(), np.int16))] + assert len(frames_with_audio) == 5 From 8a4db8bea90dff18bede9ff9b2b2ad921674262b Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Tue, 29 Sep 2026 11:14:09 +0000 Subject: [PATCH 2/2] fix(rtc): wait for cancelled stream reads to finish when closing AudioMixer --- livekit-rtc/livekit/rtc/audio_mixer.py | 14 ++++++- tests/rtc/test_mixer.py | 55 ++++++++++++++++++++++++++ 2 files changed, 67 insertions(+), 2 deletions(-) diff --git a/livekit-rtc/livekit/rtc/audio_mixer.py b/livekit-rtc/livekit/rtc/audio_mixer.py index 0c2986a6..76a407a7 100644 --- a/livekit-rtc/livekit/rtc/audio_mixer.py +++ b/livekit-rtc/livekit/rtc/audio_mixer.py @@ -56,6 +56,8 @@ def __init__( # next round. Cancelling it would also cancel the stream itself when the # stream is an async generator, which then ends and drops all its audio. self._pending: dict[_Stream, asyncio.Future[AudioFrame]] = {} + # Reads cancelled by remove_stream that may still be finishing their cancellation. + self._cancelling: set[asyncio.Future[AudioFrame]] = set() self._sample_rate: int = sample_rate self._num_channels: int = num_channels self._chunk_size: int = blocksize if blocksize > 0 else int(sample_rate // 10) @@ -97,6 +99,9 @@ def remove_stream(self, stream: AsyncIterator[AudioFrame]) -> None: pending = self._pending.pop(stream, None) if pending is not None: pending.cancel() + # aclose() waits for the cancellation to finish, which can outlive this call. + self._cancelling.add(pending) + pending.add_done_callback(self._cancelling.discard) def __aiter__(self) -> "AudioMixer": return self @@ -117,9 +122,14 @@ async def aclose(self) -> None: self._mixer_task.cancel() with contextlib.suppress(asyncio.CancelledError): await self._mixer_task - for pending in self._pending.values(): - pending.cancel() + pending_reads = list(self._pending.values()) self._pending.clear() + for pending in pending_reads: + pending.cancel() + # Cancelling only requests cancellation; wait so stream cleanup has finished on return. + # Reads already cancelled by remove_stream are awaited but not cancelled a second time, + # which would interrupt their cleanup. + await asyncio.gather(*pending_reads, *self._cancelling, return_exceptions=True) def end_input(self) -> None: """ diff --git a/tests/rtc/test_mixer.py b/tests/rtc/test_mixer.py index 8a697b90..9b736de9 100644 --- a/tests/rtc/test_mixer.py +++ b/tests/rtc/test_mixer.py @@ -86,3 +86,58 @@ async def read_all(): frames_with_audio = [f for f in frames if np.any(np.frombuffer(f.data.tobytes(), np.int16))] assert len(frames_with_audio) == 5 + + +@pytest.mark.asyncio +async def test_mixer_aclose_waits_for_pending_stream_cleanup(): + """ + Closing the mixer must finish cancelling an in-flight stream read, so the + stream's cleanup has completed by the time `aclose` returns. + """ + cleaned_up = False + + async def stream_with_slow_cleanup(): + nonlocal cleaned_up + try: + await asyncio.sleep(10) + yield AudioFrame( + np.ones(BLOCKSIZE, dtype=np.int16).tobytes(), SAMPLE_RATE, 1, BLOCKSIZE + ) + finally: + await asyncio.sleep(0.1) + cleaned_up = True + + mixer = AudioMixer(sample_rate=SAMPLE_RATE, num_channels=1, stream_timeout_ms=50) + mixer.add_stream(stream_with_slow_cleanup()) + await asyncio.sleep(0.2) # the read is pending and has missed the timeout + + await mixer.aclose() + + assert cleaned_up + + +@pytest.mark.asyncio +async def test_mixer_aclose_waits_for_cleanup_of_a_removed_stream(): + """A stream removed while its read is pending is also finished by `aclose`.""" + cleaned_up = False + + async def stream_with_slow_cleanup(): + nonlocal cleaned_up + try: + await asyncio.sleep(10) + yield AudioFrame( + np.ones(BLOCKSIZE, dtype=np.int16).tobytes(), SAMPLE_RATE, 1, BLOCKSIZE + ) + finally: + await asyncio.sleep(0.1) + cleaned_up = True + + stream = stream_with_slow_cleanup() + mixer = AudioMixer(sample_rate=SAMPLE_RATE, num_channels=1, stream_timeout_ms=50) + mixer.add_stream(stream) + await asyncio.sleep(0.2) + mixer.remove_stream(stream) + + await mixer.aclose() + + assert cleaned_up