Skip to content

Commit 6e120fb

Browse files
committed
rtc: close the room when Room.connect is cancelled
The FFI server has no cancel path for an in-flight connect: it answers the connect request and then waits for ReadyForRoomEventRequest, which connect() sends as its last statement. A coroutine cancelled anywhere inside connect() never reaches it, the server times out after 15s and panics, and the panic handler sends SIGTERM to the process. Hand the room to a task that survives the cancellation, answer the pending ready request and disconnect. disconnect() waits for that task so callers can close deterministically, and the room no longer stays joined server-side, which is what evicts a retry using the same identity. Fixes #784 Refs #804
1 parent f0fd85a commit 6e120fb

2 files changed

Lines changed: 185 additions & 5 deletions

File tree

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

Lines changed: 67 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@
3131
from ._proto.room_pb2 import ConnectionState, SimulateScenarioKind
3232
from ._proto.track_pb2 import TrackKind
3333
from ._proto.rpc_pb2 import RpcMethodInvocationEvent
34-
from ._utils import BroadcastQueue
34+
from ._utils import BroadcastQueue, Queue, task_done_logger
3535
from .e2ee import E2EEManager, E2EEOptions
3636
from .log import logger
3737
from .participant import (
@@ -184,6 +184,7 @@ def __init__(
184184
self._room_queue = BroadcastQueue[proto_ffi.FfiEvent]()
185185
self._info = proto_room.RoomInfo()
186186
self._rpc_invocation_tasks: set[asyncio.Task] = set()
187+
self._aborted_connect_tasks: set[asyncio.Task] = set()
187188

188189
self._remote_participants: Dict[str, RemoteParticipant] = {}
189190
self._connection_state = ConnectionState.CONN_DISCONNECTED
@@ -554,13 +555,25 @@ def on_participant_connected(participant):
554555
self._ffi_queue = FfiClient.instance.queue.subscribe(self._loop)
555556

556557
queue = FfiClient.instance.queue.subscribe()
558+
aborted = False
557559
try:
558560
resp = FfiClient.instance.request(req)
559-
cb: proto_ffi.FfiEvent = await queue.wait_for(
560-
lambda e: e.connect.async_id == resp.connect.async_id
561-
)
561+
try:
562+
cb: proto_ffi.FfiEvent = await queue.wait_for(
563+
lambda e: e.connect.async_id == resp.connect.async_id
564+
)
565+
except asyncio.CancelledError:
566+
# the FFI server is already connecting and expects a ReadyForRoomEvent
567+
# once it answers. leaving that unanswered panics it, and the panic
568+
# handler terminates the process, so close the room from a task that
569+
# outlives this cancellation.
570+
aborted = True
571+
FfiClient.instance.queue.unsubscribe(self._ffi_queue)
572+
self._close_aborted_connect(resp.connect.async_id, queue)
573+
raise
562574
finally:
563-
FfiClient.instance.queue.unsubscribe(queue)
575+
if not aborted:
576+
FfiClient.instance.queue.unsubscribe(queue)
564577

565578
if cb.connect.error:
566579
FfiClient.instance.queue.unsubscribe(self._ffi_queue)
@@ -602,6 +615,50 @@ def on_participant_connected(participant):
602615
ready_req.ready_for_room_event.room_handle = self._ffi_handle.handle
603616
FfiClient.instance.request(ready_req)
604617

618+
def _close_aborted_connect(self, async_id: int, queue: Queue[proto_ffi.FfiEvent]) -> None:
619+
"""Close a room that connect() was cancelled before it could own.
620+
621+
Takes ownership of `queue`. The FFI server has no cancel path for an in-flight
622+
connect, so the room has to be created and then disconnected. Without this the
623+
room also stays joined server-side and reconnecting with the same identity
624+
evicts the new session as a duplicate.
625+
"""
626+
627+
async def _close() -> None:
628+
try:
629+
cb: proto_ffi.FfiEvent = await queue.wait_for(
630+
lambda e: e.connect.async_id == async_id
631+
)
632+
finally:
633+
FfiClient.instance.queue.unsubscribe(queue)
634+
635+
if cb.connect.error:
636+
return
637+
638+
ffi_handle = FfiHandle(cb.connect.result.room.handle.id)
639+
640+
ready_req = proto_ffi.FfiRequest()
641+
ready_req.ready_for_room_event.room_handle = ffi_handle.handle
642+
FfiClient.instance.request(ready_req)
643+
644+
close_req = proto_ffi.FfiRequest()
645+
close_req.disconnect.room_handle = ffi_handle.handle
646+
close_req.disconnect.reason = DisconnectReason.CLIENT_INITIATED
647+
close_queue = FfiClient.instance.queue.subscribe()
648+
try:
649+
resp = FfiClient.instance.request(close_req)
650+
await close_queue.wait_for(
651+
lambda e: e.disconnect.async_id == resp.disconnect.async_id
652+
)
653+
finally:
654+
FfiClient.instance.queue.unsubscribe(close_queue)
655+
656+
task = self._loop.create_task(_close())
657+
self._aborted_connect_tasks.add(task)
658+
task.add_done_callback(self._aborted_connect_tasks.discard)
659+
# a failure here still ends in an FFI panic, so it must not be swallowed
660+
task.add_done_callback(task_done_logger)
661+
605662
async def get_rtc_stats(self) -> RtcStats:
606663
if not self.isconnected():
607664
raise RuntimeError("the room isn't connected")
@@ -681,6 +738,11 @@ async def disconnect(
681738
self, *, reason: DisconnectReason.ValueType = DisconnectReason.CLIENT_INITIATED
682739
) -> None:
683740
"""Disconnects from the room."""
741+
if self._aborted_connect_tasks:
742+
# a cancelled connect may still be closing a room the FFI server opened.
743+
# wait for it so disconnect() leaves nothing behind.
744+
await asyncio.gather(*tuple(self._aborted_connect_tasks), return_exceptions=True)
745+
684746
if not self.isconnected():
685747
return
686748

Lines changed: 118 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,118 @@
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+
import asyncio
16+
17+
import pytest
18+
19+
from livekit import rtc
20+
from livekit.rtc import room as room_mod
21+
from livekit.rtc._ffi_client import FfiClient
22+
from livekit.rtc._proto import ffi_pb2 as proto_ffi
23+
from utils import wait_until # type: ignore[import-not-found]
24+
25+
CONNECT_ASYNC_ID = 101
26+
DISCONNECT_ASYNC_ID = 202
27+
ROOM_HANDLE = 7
28+
29+
30+
class _FakeHandle:
31+
"""Stand-in for FfiHandle so a made-up handle id is never dropped natively."""
32+
33+
def __init__(self, handle: int) -> None:
34+
self.handle = handle
35+
36+
37+
def _install_fake_ffi(monkeypatch: pytest.MonkeyPatch) -> list[proto_ffi.FfiRequest]:
38+
"""Record every FfiRequest and answer the ones the cancel path waits on."""
39+
requests: list[proto_ffi.FfiRequest] = []
40+
41+
def fake_request(req: proto_ffi.FfiRequest) -> proto_ffi.FfiResponse:
42+
requests.append(req)
43+
resp = proto_ffi.FfiResponse()
44+
which = req.WhichOneof("message")
45+
if which == "connect":
46+
# the connect callback is delivered by the test, not here
47+
resp.connect.async_id = CONNECT_ASYNC_ID
48+
elif which == "disconnect":
49+
resp.disconnect.async_id = DISCONNECT_ASYNC_ID
50+
event = proto_ffi.FfiEvent()
51+
event.disconnect.async_id = DISCONNECT_ASYNC_ID
52+
FfiClient.instance.queue.put(event)
53+
return resp
54+
55+
monkeypatch.setattr(FfiClient.instance, "request", fake_request)
56+
monkeypatch.setattr(room_mod, "FfiHandle", _FakeHandle)
57+
return requests
58+
59+
60+
def _deliver_connect_callback() -> None:
61+
event = proto_ffi.FfiEvent()
62+
event.connect.async_id = CONNECT_ASYNC_ID
63+
event.connect.result.room.handle.id = ROOM_HANDLE
64+
FfiClient.instance.queue.put(event)
65+
66+
67+
async def test_cancelled_connect_answers_ready_and_closes_the_room(
68+
monkeypatch: pytest.MonkeyPatch,
69+
) -> None:
70+
requests = _install_fake_ffi(monkeypatch)
71+
subscribers_before = len(FfiClient.instance.queue._subscribers)
72+
73+
room = rtc.Room()
74+
task = asyncio.create_task(room.connect("ws://localhost:7880", "token"))
75+
await wait_until(lambda: bool(requests), message="connect request never issued")
76+
77+
task.cancel()
78+
with pytest.raises(asyncio.CancelledError):
79+
await task
80+
81+
# the FFI server does not cancel an in-flight connect: it answers, then waits for
82+
# ReadyForRoomEvent. an unanswered wait panics it and the panic kills the process.
83+
_deliver_connect_callback()
84+
await room.disconnect()
85+
86+
assert [req.WhichOneof("message") for req in requests] == [
87+
"connect",
88+
"ready_for_room_event",
89+
"disconnect",
90+
]
91+
assert requests[1].ready_for_room_event.room_handle == ROOM_HANDLE
92+
assert requests[2].disconnect.room_handle == ROOM_HANDLE
93+
assert len(FfiClient.instance.queue._subscribers) == subscribers_before
94+
95+
96+
async def test_cancelled_connect_leaves_no_pending_work_when_the_server_errors(
97+
monkeypatch: pytest.MonkeyPatch,
98+
) -> None:
99+
requests = _install_fake_ffi(monkeypatch)
100+
subscribers_before = len(FfiClient.instance.queue._subscribers)
101+
102+
room = rtc.Room()
103+
task = asyncio.create_task(room.connect("ws://localhost:7880", "token"))
104+
await wait_until(lambda: bool(requests), message="connect request never issued")
105+
106+
task.cancel()
107+
with pytest.raises(asyncio.CancelledError):
108+
await task
109+
110+
event = proto_ffi.FfiEvent()
111+
event.connect.async_id = CONNECT_ASYNC_ID
112+
event.connect.error = "could not connect"
113+
FfiClient.instance.queue.put(event)
114+
await room.disconnect()
115+
116+
# there is no room to close, so nothing follows the connect
117+
assert [req.WhichOneof("message") for req in requests] == ["connect"]
118+
assert len(FfiClient.instance.queue._subscribers) == subscribers_before

0 commit comments

Comments
 (0)