Skip to content

Commit 591b8ce

Browse files
committed
fix(rtc): filter media events from the room FFI subscription
Exclude audio and video stream events before scheduling room callbacks, while preserving publication and other request callbacks. Add regression coverage for dispatch filtering, publish/unpublish completion, and subscription cleanup.
1 parent 64a34b6 commit 591b8ce

2 files changed

Lines changed: 166 additions & 1 deletion

File tree

‎livekit-rtc/livekit/rtc/room.py‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -551,7 +551,14 @@ def on_participant_connected(participant):
551551
req.connect.options.rtc_config.ice_servers.extend(options.rtc_config.ice_servers)
552552

553553
# subscribe before connecting so we don't miss any events
554-
self._ffi_queue = FfiClient.instance.queue.subscribe(self._loop)
554+
# Media streams have their own subscriptions. Keep other events, including
555+
# the publish/unpublish callbacks forwarded through _room_queue.
556+
self._ffi_queue = FfiClient.instance.queue.subscribe(
557+
self._loop,
558+
filter_fn=lambda e: (
559+
e.WhichOneof("message") not in ("audio_stream_event", "video_stream_event")
560+
),
561+
)
555562

556563
queue = FfiClient.instance.queue.subscribe()
557564
try:

‎tests/rtc/test_room_ffi_filter.py‎

Lines changed: 158 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,158 @@
1+
# Copyright 2026 LiveKit, Inc.
2+
#
3+
# Licensed under the Apache License, Version 2.0 (the "License");
4+
# you may not use this file except in compliance with the License.
5+
# You may obtain a copy of the License at
6+
#
7+
# http://www.apache.org/licenses/LICENSE-2.0
8+
#
9+
# Unless required by applicable law or agreed to in writing, software
10+
# distributed under the License is distributed on an "AS IS" BASIS,
11+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
# See the License for the specific language governing permissions and
13+
# limitations under the License.
14+
15+
"""Exercise Room.connect's real subscription without a native FFI or server."""
16+
17+
import asyncio
18+
from collections.abc import AsyncIterator
19+
from types import SimpleNamespace
20+
from unittest.mock import Mock, patch
21+
22+
import pytest
23+
24+
from livekit import rtc
25+
from livekit.rtc._ffi_client import FfiClient, FfiQueue
26+
from livekit.rtc._proto import ffi_pb2 as proto_ffi
27+
from livekit.rtc._proto import track_pb2 as proto_track
28+
29+
30+
@pytest.fixture
31+
async def connected_room(
32+
monkeypatch: pytest.MonkeyPatch,
33+
) -> AsyncIterator[tuple[rtc.Room, FfiQueue[proto_ffi.FfiEvent]]]:
34+
queue = FfiQueue[proto_ffi.FfiEvent]()
35+
36+
def request(req: proto_ffi.FfiRequest) -> proto_ffi.FfiResponse:
37+
response = proto_ffi.FfiResponse()
38+
event = proto_ffi.FfiEvent()
39+
which = req.WhichOneof("message")
40+
if which == "connect":
41+
response.connect.async_id = event.connect.async_id = 1
42+
event.connect.result.room.handle.id = 1
43+
event.connect.result.room.info.sid = "RM_test"
44+
elif which == "ready_for_room_event":
45+
return response
46+
elif which == "publish_track":
47+
response.publish_track.async_id = event.publish_track.async_id = 2
48+
event.publish_track.publication.info.sid = "TR_test"
49+
elif which == "unpublish_track":
50+
response.unpublish_track.async_id = event.unpublish_track.async_id = 3
51+
elif which == "disconnect":
52+
response.disconnect.async_id = event.disconnect.async_id = 4
53+
eos = proto_ffi.FfiEvent()
54+
eos.room_event.room_handle = req.disconnect.room_handle
55+
eos.room_event.eos.SetInParent()
56+
queue.put(eos)
57+
else:
58+
raise AssertionError(f"unexpected FFI request: {which}")
59+
queue.put(event)
60+
return response
61+
62+
# A different PID prevents synthetic handles from being dropped natively.
63+
monkeypatch.setattr(
64+
FfiClient, "_instance", SimpleNamespace(queue=queue, request=request, _pid=-1)
65+
)
66+
room = rtc.Room()
67+
await asyncio.wait_for(room.connect("wss://example.invalid", "test-token"), 1)
68+
assert room._ffi_handle is not None
69+
room._ffi_handle.mark_consumed()
70+
try:
71+
yield room, queue
72+
finally:
73+
try:
74+
await asyncio.wait_for(room.disconnect(), 1)
75+
finally:
76+
# Do not access the restored FFI singleton from Room.__del__ later.
77+
room._ffi_handle = None
78+
assert not queue._subscribers
79+
80+
81+
@pytest.mark.parametrize("event_type", ["audio_stream_event", "video_stream_event"])
82+
async def test_media_events_do_not_schedule_room_callbacks(
83+
connected_room: tuple[rtc.Room, FfiQueue[proto_ffi.FfiEvent]], event_type: str
84+
) -> None:
85+
room, queue = connected_room
86+
loop = asyncio.get_running_loop()
87+
media_queue = queue.subscribe(loop, filter_fn=lambda e: e.WhichOneof("message") == event_type)
88+
event = proto_ffi.FfiEvent()
89+
getattr(event, event_type).SetInParent()
90+
try:
91+
with patch.object(
92+
loop, "call_soon_threadsafe", wraps=loop.call_soon_threadsafe
93+
) as schedule:
94+
for _ in range(100):
95+
queue.put(event)
96+
97+
# Only the independent media subscriber should incur an event-loop wakeup.
98+
assert schedule.call_count == 100
99+
assert all(call.args[0] == media_queue.put_nowait for call in schedule.call_args_list)
100+
await asyncio.sleep(0)
101+
assert media_queue.qsize() == 100
102+
assert room._ffi_queue.empty()
103+
finally:
104+
queue.unsubscribe(media_queue)
105+
106+
107+
@pytest.mark.parametrize(
108+
"event_type",
109+
[
110+
"room_event",
111+
"rpc_method_invocation",
112+
"publish_track",
113+
"unpublish_track",
114+
"capture_audio_frame",
115+
],
116+
)
117+
async def test_room_keeps_control_events_and_request_callbacks(
118+
connected_room: tuple[rtc.Room, FfiQueue[proto_ffi.FfiEvent]],
119+
monkeypatch: pytest.MonkeyPatch,
120+
event_type: str,
121+
) -> None:
122+
room, queue = connected_room
123+
room_handler = Mock()
124+
rpc_handler = Mock()
125+
monkeypatch.setattr(room, "_on_room_event", room_handler)
126+
monkeypatch.setattr(room, "_on_rpc_method_invocation", rpc_handler)
127+
subscriber = room._room_queue.subscribe()
128+
event = proto_ffi.FfiEvent()
129+
getattr(event, event_type).SetInParent()
130+
if event_type == "room_event":
131+
event.room_event.room_handle = 1
132+
try:
133+
queue.put(event)
134+
received = await asyncio.wait_for(subscriber.get(), 1)
135+
subscriber.task_done()
136+
assert received is event
137+
if event_type == "room_event":
138+
room_handler.assert_called_once_with(event.room_event)
139+
elif event_type == "rpc_method_invocation":
140+
rpc_handler.assert_called_once_with(event.rpc_method_invocation)
141+
finally:
142+
room._room_queue.unsubscribe(subscriber)
143+
144+
145+
async def test_publish_and_unpublish_complete_through_room_subscription(
146+
connected_room: tuple[rtc.Room, FfiQueue[proto_ffi.FfiEvent]],
147+
) -> None:
148+
room, _ = connected_room
149+
track = rtc.LocalAudioTrack(proto_track.OwnedTrack())
150+
participant = room.local_participant
151+
152+
publication = await asyncio.wait_for(participant.publish_track(track), 1)
153+
assert participant.track_publications["TR_test"] is publication
154+
assert publication.track is track
155+
156+
await asyncio.wait_for(participant.unpublish_track(publication.sid), 1)
157+
assert not participant.track_publications
158+
assert publication.track is None

0 commit comments

Comments
 (0)