Skip to content

Commit 50f1867

Browse files
committed
fix
1 parent 96d9770 commit 50f1867

2 files changed

Lines changed: 45 additions & 2 deletions

File tree

‎livekit-rtc/livekit/rtc/participant.py‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -673,8 +673,10 @@ async def _run_incoming_chain(self, invocation: RpcInvocationData) -> Optional[s
673673

674674
def _on_deadline() -> None:
675675
nonlocal deadline_fired
676-
deadline_fired = True
677-
chain_task.cancel()
676+
# only a cancel the chain accepted counts: cancel() is False when the chain has
677+
# already finished, which can happen in the same loop iteration the timer fires
678+
# while this task has not resumed yet; that result is the caller's, not a timeout
679+
deadline_fired = chain_task.cancel()
678680

679681
deadline = loop.call_later(invocation.response_timeout, _on_deadline)
680682
try:
@@ -709,6 +711,16 @@ def _on_deadline() -> None:
709711
# whatever the chain raised while unwinding is not the caller's business
710712
_observe_unwind(invocation.method, chain_task)
711713
raise RpcError._built_in(RpcError.ErrorCode.RESPONSE_TIMEOUT)
714+
if chain_task.cancelled():
715+
# nothing here cancelled it: an interceptor or the handler let a CancelledError
716+
# of its own escape. That is a failure of the handler, and it must still be
717+
# answered; result() would re-raise the CancelledError past the response code
718+
# and leave the caller waiting for its timeout.
719+
logger.warning(
720+
"RPC handler for %s raised CancelledError; returning APPLICATION_ERROR",
721+
invocation.method,
722+
)
723+
raise RpcError._built_in(RpcError.ErrorCode.APPLICATION_ERROR)
712724
# the chain's own outcome: its exception, if any, propagates unchanged
713725
return chain_task.result()
714726

‎livekit-rtc/tests/test_rpc_interceptors.py‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -302,6 +302,37 @@ async def slow(data: RpcInvocationData) -> str:
302302
assert info.value.code == rtc.RpcError.ErrorCode.RESPONSE_TIMEOUT
303303

304304

305+
async def test_chain_finishing_as_the_deadline_fires_keeps_its_result() -> None:
306+
"""The deadline timer can run in the same loop iteration the chain completes in, before
307+
the waiting task resumes: its cancel is rejected (the chain is done) and the result is
308+
the caller's. A zero deadline with an immediate handler pins that ordering."""
309+
lp = _participant()
310+
lp._rpc_handlers["fast"] = lambda data: "ok"
311+
312+
result = await lp._run_incoming_chain(
313+
RpcInvocationData("r1", "alice", "{}", 0.0, method="fast")
314+
)
315+
assert result == "ok"
316+
317+
318+
async def test_cancelled_error_escaping_the_chain_is_an_application_error(
319+
caplog: pytest.LogCaptureFixture,
320+
) -> None:
321+
"""A CancelledError that an interceptor or the handler lets escape on its own is not a
322+
deadline nor a disconnect; the caller must still get an answer rather than wait out its
323+
timeout because the cancellation slipped past the response code."""
324+
lp = _participant()
325+
326+
async def leaks_cancel(data: RpcInvocationData) -> str:
327+
raise asyncio.CancelledError()
328+
329+
lp._rpc_handlers["leaky"] = leaks_cancel
330+
with pytest.raises(rtc.RpcError) as info:
331+
await lp._run_incoming_chain(RpcInvocationData("r1", "alice", "{}", 5.0, method="leaky"))
332+
assert info.value.code == rtc.RpcError.ErrorCode.APPLICATION_ERROR
333+
assert any("leaky" in r.getMessage() for r in caplog.records)
334+
335+
305336
async def test_outside_cancellation_maps_to_recipient_disconnected() -> None:
306337
lp = _participant()
307338
started = asyncio.Event()

0 commit comments

Comments
 (0)