Skip to content

Commit 96d9770

Browse files
committed
fix
1 parent 403d5d2 commit 96d9770

1 file changed

Lines changed: 18 additions & 14 deletions

File tree

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

Lines changed: 18 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -247,9 +247,9 @@ def disconnect_reason(
247247
def _observe_unwind(method: str, chain_task: "asyncio.Future[Optional[str]]") -> None:
248248
"""Consume what a cancelled chain raised while unwinding.
249249
250-
Nobody awaits the chain once the caller has been answered, and the shield that let the
251-
outside cancel through stops watching the chain the moment its own future is cancelled;
252-
without this the loop would report the exception as never retrieved at garbage collection.
250+
The caller has been (or is being) answered without it, so nobody awaits the chain
251+
anymore; unless its exception is retrieved here the loop reports it as never retrieved
252+
when the task is garbage collected.
253253
"""
254254
if chain_task.cancelled():
255255
return
@@ -678,14 +678,15 @@ def _on_deadline() -> None:
678678

679679
deadline = loop.call_later(invocation.response_timeout, _on_deadline)
680680
try:
681-
# shielded: a cancel from outside (the room disconnecting) is raised here at
682-
# once. Awaiting the chain directly would instead forward the cancel to it and
683-
# keep this task parked until the chain finished, so a handler that ignored
684-
# cancellation held up room.disconnect() for as long as the caller's deadline.
685-
return await asyncio.shield(chain_task)
681+
# asyncio.wait never cancels what it waits on, so a cancel from outside (the room
682+
# disconnecting) is raised here at once with the chain untouched. Awaiting the
683+
# chain directly would forward the cancel to it and keep this task parked until
684+
# the chain finished, so a handler that ignored cancellation held up
685+
# room.disconnect() for as long as the caller's deadline. (Not asyncio.shield:
686+
# since Python 3.14 it reports the inner future's exception to the loop's
687+
# exception handler once its outer future was cancelled.)
688+
await asyncio.wait([chain_task])
686689
except asyncio.CancelledError:
687-
if deadline_fired:
688-
raise RpcError._built_in(RpcError.ErrorCode.RESPONSE_TIMEOUT) from None
689690
# cancelled from outside: stop the chain and let it unwind before answering the
690691
# caller, but not for long; this is the path room.disconnect() waits on
691692
chain_task.cancel()
@@ -701,13 +702,16 @@ def _on_deadline() -> None:
701702
else:
702703
_observe_unwind(invocation.method, chain_task)
703704
raise RpcError._built_in(RpcError.ErrorCode.RECIPIENT_DISCONNECTED) from None
704-
except Exception:
705-
if deadline_fired:
706-
raise RpcError._built_in(RpcError.ErrorCode.RESPONSE_TIMEOUT) from None
707-
raise
708705
finally:
709706
deadline.cancel()
710707

708+
if deadline_fired:
709+
# whatever the chain raised while unwinding is not the caller's business
710+
_observe_unwind(invocation.method, chain_task)
711+
raise RpcError._built_in(RpcError.ErrorCode.RESPONSE_TIMEOUT)
712+
# the chain's own outcome: its exception, if any, propagates unchanged
713+
return chain_task.result()
714+
711715
async def _invoke_rpc_handler(self, invocation: RpcInvocationData) -> Optional[str]:
712716
"""Run the registered handler for ``invocation`` (the innermost step of the chain).
713717

0 commit comments

Comments
 (0)