From 394a60fc75b0b1c978a789b98a1db18c06acc385 Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Tue, 29 Sep 2026 09:53:56 +0000 Subject: [PATCH 1/2] fix(rtc): deliver buffered frames and end iteration at end of stream --- livekit-rtc/livekit/rtc/_utils.py | 3 + livekit-rtc/livekit/rtc/audio_stream.py | 3 +- livekit-rtc/livekit/rtc/video_stream.py | 4 +- tests/rtc/test_stream_end.py | 134 ++++++++++++++++++++++++ 4 files changed, 142 insertions(+), 2 deletions(-) create mode 100644 tests/rtc/test_stream_end.py diff --git a/livekit-rtc/livekit/rtc/_utils.py b/livekit-rtc/livekit/rtc/_utils.py index 91175815..098d1ed2 100644 --- a/livekit-rtc/livekit/rtc/_utils.py +++ b/livekit-rtc/livekit/rtc/_utils.py @@ -93,6 +93,9 @@ async def get(self) -> T: self._event.clear() return self._queue.popleft() + def empty(self) -> bool: + return len(self._queue) == 0 + class Queue(asyncio.Queue[T]): """asyncio.Queue with utility functions.""" diff --git a/livekit-rtc/livekit/rtc/audio_stream.py b/livekit-rtc/livekit/rtc/audio_stream.py index 97772f4f..7be481be 100644 --- a/livekit-rtc/livekit/rtc/audio_stream.py +++ b/livekit-rtc/livekit/rtc/audio_stream.py @@ -335,7 +335,8 @@ def __aiter__(self) -> AsyncIterator[AudioFrameEvent]: return self async def __anext__(self) -> AudioFrameEvent: - if self._task.done(): + # frames queued before the stream ended are still delivered + if self._task.done() and self._queue.empty(): raise StopAsyncIteration item = await self._queue.get() diff --git a/livekit-rtc/livekit/rtc/video_stream.py b/livekit-rtc/livekit/rtc/video_stream.py index eeced550..31243087 100644 --- a/livekit-rtc/livekit/rtc/video_stream.py +++ b/livekit-rtc/livekit/rtc/video_stream.py @@ -160,6 +160,7 @@ async def _run(self) -> None: self._queue.put(event) elif video_event.HasField("eos"): + self._queue.put(None) break FfiClient.instance.queue.unsubscribe(self._ffi_queue) @@ -175,7 +176,8 @@ def __aiter__(self) -> AsyncIterator[VideoFrameEvent]: return self async def __anext__(self) -> VideoFrameEvent: - if self._task.done(): + # frames queued before the stream ended are still delivered + if self._task.done() and self._queue.empty(): raise StopAsyncIteration item = await self._queue.get() diff --git a/tests/rtc/test_stream_end.py b/tests/rtc/test_stream_end.py new file mode 100644 index 00000000..045f6d63 --- /dev/null +++ b/tests/rtc/test_stream_end.py @@ -0,0 +1,134 @@ +# Copyright 2026 LiveKit, Inc. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Audio and video stream iteration at end of stream, without a native FFI.""" + +from __future__ import annotations + +import asyncio +import gc +from collections.abc import Iterator +from types import SimpleNamespace +from unittest.mock import MagicMock, patch + +import pytest +from livekit import rtc +from livekit.rtc import audio_stream, video_stream +from livekit.rtc._ffi_client import FfiClient, FfiQueue +from livekit.rtc._proto import ffi_pb2 as proto_ffi + +STREAM_HANDLE = 7 + + +@pytest.fixture(autouse=True) +def fake_ffi(monkeypatch: pytest.MonkeyPatch) -> Iterator[MagicMock]: + ffi = MagicMock() + ffi.queue = FfiQueue[proto_ffi.FfiEvent]() + monkeypatch.setattr(FfiClient, "_instance", ffi) + yield ffi + # stream finalizers unsubscribe from the FFI queue, so run them while it is faked + gc.collect() + + +def _ffi_handle(handle: int) -> SimpleNamespace: + return SimpleNamespace(handle=handle, disposed=False, dispose=lambda: None) + + +def _track() -> MagicMock: + track = MagicMock() + track._ffi_handle.handle = 1 + return track + + +def _response(kind: str) -> proto_ffi.FfiResponse: + response = proto_ffi.FfiResponse() + getattr(response, kind).stream.handle.id = STREAM_HANDLE + return response + + +def _audio_frame_event() -> proto_ffi.FfiEvent: + event = proto_ffi.FfiEvent() + event.audio_stream_event.stream_handle = STREAM_HANDLE + event.audio_stream_event.frame_received.frame.handle.id = 1 + return event + + +def _audio_eos_event() -> proto_ffi.FfiEvent: + event = proto_ffi.FfiEvent() + event.audio_stream_event.stream_handle = STREAM_HANDLE + event.audio_stream_event.eos.SetInParent() + return event + + +def _video_frame_event() -> proto_ffi.FfiEvent: + event = proto_ffi.FfiEvent() + event.video_stream_event.stream_handle = STREAM_HANDLE + event.video_stream_event.frame_received.buffer.handle.id = 1 + return event + + +def _video_eos_event() -> proto_ffi.FfiEvent: + event = proto_ffi.FfiEvent() + event.video_stream_event.stream_handle = STREAM_HANDLE + event.video_stream_event.eos.SetInParent() + return event + + +async def _collect(stream: rtc.AudioStream | rtc.VideoStream) -> list: + async def read() -> list: + return [event async for event in stream] + + return await asyncio.wait_for(read(), timeout=2) + + +async def test_audio_stream_delivers_frames_queued_before_eos() -> None: + with ( + patch.object(FfiClient.instance, "request", return_value=_response("new_audio_stream")), + patch.object(audio_stream, "FfiHandle", side_effect=_ffi_handle), + patch.object(audio_stream.AudioFrame, "_from_owned_info", side_effect=lambda _: object()), + ): + stream = rtc.AudioStream(_track()) + for event in (_audio_frame_event(), _audio_frame_event(), _audio_eos_event()): + FfiClient.instance.queue.put(event) + # let the stream's task consume every event before the reader starts + await asyncio.wait_for(stream._task, timeout=2) + + assert len(await _collect(stream)) == 2 + + +async def test_video_stream_delivers_frames_queued_before_eos() -> None: + with ( + patch.object(FfiClient.instance, "request", return_value=_response("new_video_stream")), + patch.object(video_stream, "FfiHandle", side_effect=_ffi_handle), + patch.object(video_stream.VideoFrame, "_from_owned_info", side_effect=lambda _: object()), + ): + stream = rtc.VideoStream(_track()) + for event in (_video_frame_event(), _video_frame_event(), _video_eos_event()): + FfiClient.instance.queue.put(event) + await asyncio.wait_for(stream._task, timeout=2) + + assert len(await _collect(stream)) == 2 + + +async def test_video_stream_ends_when_eos_arrives_during_read() -> None: + with ( + patch.object(FfiClient.instance, "request", return_value=_response("new_video_stream")), + patch.object(video_stream, "FfiHandle", side_effect=_ffi_handle), + ): + stream = rtc.VideoStream(_track()) + reader = asyncio.ensure_future(_collect(stream)) + await asyncio.sleep(0.05) # the reader is now waiting for a frame + FfiClient.instance.queue.put(_video_eos_event()) + + assert await asyncio.wait_for(reader, timeout=2) == [] From a122059b12fa1a0850f6050e568719149432e601 Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Tue, 29 Sep 2026 11:13:08 +0000 Subject: [PATCH 2/2] fix(rtc): keep the last frame of a full bounded stream queue at end of stream --- livekit-rtc/livekit/rtc/_utils.py | 5 +++++ livekit-rtc/livekit/rtc/audio_stream.py | 2 +- livekit-rtc/livekit/rtc/video_stream.py | 2 +- tests/rtc/test_stream_end.py | 28 +++++++++++++++++++++++++ 4 files changed, 35 insertions(+), 2 deletions(-) diff --git a/livekit-rtc/livekit/rtc/_utils.py b/livekit-rtc/livekit/rtc/_utils.py index 098d1ed2..9f9b89c4 100644 --- a/livekit-rtc/livekit/rtc/_utils.py +++ b/livekit-rtc/livekit/rtc/_utils.py @@ -87,6 +87,11 @@ def put(self, item: T) -> None: self._queue.append(item) self._event.set() + def put_end(self, item: T) -> None: + """Appends an end-of-stream marker without evicting a queued item, even when full.""" + self._queue.append(item) + self._event.set() + async def get(self) -> T: while len(self._queue) == 0: await self._event.wait() diff --git a/livekit-rtc/livekit/rtc/audio_stream.py b/livekit-rtc/livekit/rtc/audio_stream.py index 7be481be..e088110c 100644 --- a/livekit-rtc/livekit/rtc/audio_stream.py +++ b/livekit-rtc/livekit/rtc/audio_stream.py @@ -310,7 +310,7 @@ async def _run(self) -> None: event = AudioFrameEvent(frame) self._queue.put(event) elif audio_event.HasField("eos"): - self._queue.put(None) + self._queue.put_end(None) break FfiClient.instance.queue.unsubscribe(self._ffi_queue) diff --git a/livekit-rtc/livekit/rtc/video_stream.py b/livekit-rtc/livekit/rtc/video_stream.py index 31243087..1134d1ed 100644 --- a/livekit-rtc/livekit/rtc/video_stream.py +++ b/livekit-rtc/livekit/rtc/video_stream.py @@ -160,7 +160,7 @@ async def _run(self) -> None: self._queue.put(event) elif video_event.HasField("eos"): - self._queue.put(None) + self._queue.put_end(None) break FfiClient.instance.queue.unsubscribe(self._ffi_queue) diff --git a/tests/rtc/test_stream_end.py b/tests/rtc/test_stream_end.py index 045f6d63..713f8029 100644 --- a/tests/rtc/test_stream_end.py +++ b/tests/rtc/test_stream_end.py @@ -132,3 +132,31 @@ async def test_video_stream_ends_when_eos_arrives_during_read() -> None: FfiClient.instance.queue.put(_video_eos_event()) assert await asyncio.wait_for(reader, timeout=2) == [] + + +async def test_audio_stream_keeps_last_frame_when_bounded_queue_is_full_at_eos() -> None: + with ( + patch.object(FfiClient.instance, "request", return_value=_response("new_audio_stream")), + patch.object(audio_stream, "FfiHandle", side_effect=_ffi_handle), + patch.object(audio_stream.AudioFrame, "_from_owned_info", side_effect=lambda _: object()), + ): + stream = rtc.AudioStream(_track(), capacity=1) + for event in (_audio_frame_event(), _audio_eos_event()): + FfiClient.instance.queue.put(event) + await asyncio.wait_for(stream._task, timeout=2) + + assert len(await _collect(stream)) == 1 + + +async def test_video_stream_keeps_last_frame_when_bounded_queue_is_full_at_eos() -> None: + with ( + patch.object(FfiClient.instance, "request", return_value=_response("new_video_stream")), + patch.object(video_stream, "FfiHandle", side_effect=_ffi_handle), + patch.object(video_stream.VideoFrame, "_from_owned_info", side_effect=lambda _: object()), + ): + stream = rtc.VideoStream(_track(), capacity=1) + for event in (_video_frame_event(), _video_eos_event()): + FfiClient.instance.queue.put(event) + await asyncio.wait_for(stream._task, timeout=2) + + assert len(await _collect(stream)) == 1