Skip to content

Commit f4a3cf2

Browse files
authored
fix(rtc): filter media events from the room FFI subscription (#821)
1 parent fb265ac commit f4a3cf2

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)