Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions docs/development/07-transports-and-production.md
Original file line number Diff line number Diff line change
Expand Up @@ -227,6 +227,8 @@ WebRTC has two layers:
stateless/per-server route handling, admission, and per-offer transport
construction. Each accepted offer also installs one named runtime cleanup
task that waits for peer closure and releases its Session and capacity slot.
Shutdown cancels those waiters, runs forced finalizers in a separate named
cohort, and reports unexpected results with their task identities.

Shared config/stats/health/CORS/root behavior lives once in
[`server/_webrtc_handlers.py`](../../src/easycat/server/_webrtc_handlers.py).
Expand Down
78 changes: 69 additions & 9 deletions src/easycat/server/webrtc_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,8 @@

_OFFER_CLEANUP_TASK = "webrtc_offer_cleanup"
_OFFER_CLEANUP_COHORT = "webrtc-offer-cleanup"
_FORCE_CLEANUP_TASK = "webrtc_force_cleanup"
_FORCE_CLEANUP_COHORT = "webrtc-force-cleanup"

# Per-connection factory seam (NO ``ConnectionContext`` type): a per-transport
# ``Callable[[WebRTCTransport], EasyConfig | Session]``.
Expand Down Expand Up @@ -187,6 +189,16 @@ def __init__(
logger=logger,
failure_message="WebRTC offer cleanup task failed",
drop_if_closed=False,
release_standalone_when_idle=True,
)
self._force_cleanup_task_scope = RuntimeTaskScope(
owner_label="webrtc-routes-force-cleanup",
member_name=_FORCE_CLEANUP_TASK,
cohort=_FORCE_CLEANUP_COHORT,
logger=logger,
failure_message="WebRTC forced cleanup task failed",
drop_if_closed=False,
release_standalone_when_idle=True,
)
self._cleanup_task_keys: dict[asyncio.Task[Any], int] = {}
self._released_cleanup_keys: set[int] = set()
Expand Down Expand Up @@ -561,21 +573,69 @@ async def cancel_cleanup_tasks(self, *, timeout_s: float | None = None) -> None:
for task, _key in pending:
task.cancel()
if pending:
finalizers = [
asyncio.create_task(self._finalize_session_cleanup(key, force=True))
for _task, key in pending
if key is not None
]
finalizers: list[asyncio.Task[Any]] = []
for _task, key in pending:
if key is None:
continue
finalizer = self._force_cleanup_task_scope.create_task(
self._finalize_session_cleanup(key, force=True),
task_name=f"easycat-webrtc-force-cleanup-{key}",
)
assert finalizer is not None
finalizers.append(finalizer)
cleanup_tasks = [*(task for task, _key in pending), *finalizers]
cleanup = asyncio.gather(
*(task for task, _key in pending),
*finalizers,
*cleanup_tasks,
return_exceptions=True,
)
results: list[Any | BaseException] | None
if timeout_s is None:
await cleanup
results = await cleanup
else:
await _await_with_hard_timeout(cleanup, timeout_s=timeout_s)
completed = await _await_with_hard_timeout(cleanup, timeout_s=timeout_s)
results = cleanup.result() if completed else None
Comment thread
yisding marked this conversation as resolved.
if results is not None:
self._report_cleanup_results(
cleanup_tasks,
results,
explicitly_cancelled={task for task, _key in pending},
)
else:
for finalizer in finalizers:
finalizer.add_done_callback(self._report_late_cleanup_result)
await self._cleanup_task_scope.release_standalone_if_empty()
await self._force_cleanup_task_scope.release_standalone_if_empty()

def _report_cleanup_results(
self,
tasks: list[asyncio.Task[Any]],
results: list[Any],
*,
explicitly_cancelled: set[asyncio.Task[Any]],
) -> None:
"""Log unexpected cleanup failures with the owning task's identity."""
for task, result in zip(tasks, results, strict=True):
if not isinstance(result, BaseException):
continue
if isinstance(result, asyncio.CancelledError) and task in explicitly_cancelled:
continue
logger.error(
"WebRTC cleanup task %s failed",
task.get_name(),
exc_info=result,
)

def _report_late_cleanup_result(self, task: asyncio.Task[Any]) -> None:
"""Report a forced finalizer that settles after the hard deadline."""
if task.cancelled():
return
error = task.exception()
if error is not None:
logger.error(
"WebRTC cleanup task %s failed",
task.get_name(),
exc_info=error,
)

async def _stop_managed_session(self, key: int, force: bool) -> None:
"""Route drain teardown through the manager's keyed stop ownership."""
Expand Down
7 changes: 3 additions & 4 deletions tests/ratchets/source-baseline.json
Original file line number Diff line number Diff line change
@@ -1,12 +1,12 @@
{
"version": 1,
"rationale": "WebRTC per-offer cleanup workers now use named RuntimeTaskScope ownership.",
"rationale": "WebRTC forced cleanup finalizers now use named ownership and report unexpected results.",
"counts": {
"cancelled_error_handler": 100,
"epoch_field": 17,
"gather_return_exceptions": 22,
"module_task_set": 1,
"raw_task_spawn": 34,
"raw_task_spawn": 33,
"shield_loop": 4,
"task_cancelling": 51
},
Expand Down Expand Up @@ -136,7 +136,7 @@
"gather_return_exceptions\tserver/transports.py\tclose_websocket_connections\tasyncio.gather(return_exceptions=True)\t768ec4176f5a8ad9\t0",
"gather_return_exceptions\tserver/voice_server.py\tVoiceServer._cancel_ws_handler_tasks\tasyncio.gather(return_exceptions=True)\t77ded9771f58812f\t0",
"gather_return_exceptions\tserver/voice_server.py\tVoiceServer._close_active_ws_connections\tasyncio.gather(return_exceptions=True)\t4d4f60904d8d66b3\t0",
"gather_return_exceptions\tserver/webrtc_routes.py\tWebRTCRoutes.cancel_cleanup_tasks\tasyncio.gather(return_exceptions=True)\t69a335c95c0827de\t0",
"gather_return_exceptions\tserver/webrtc_routes.py\tWebRTCRoutes.cancel_cleanup_tasks\tasyncio.gather(return_exceptions=True)\t189be02028bb19e4\t0",
"gather_return_exceptions\tsession/_tts_scheduler.py\tTTSScheduler.begin_synthesis_with_bot_start\tasyncio.gather(return_exceptions=True)\t42eb62b100557e55\t0",
"gather_return_exceptions\tsession/_tts_scheduler.py\tTTSScheduler.begin_synthesis_with_bot_start\tasyncio.gather(return_exceptions=True)\t42eb62b100557e55\t1",
"gather_return_exceptions\tsession/_tts_scheduler.py\tTTSScheduler.begin_synthesis_with_bot_start\tasyncio.gather(return_exceptions=True)\t42eb62b100557e55\t2",
Expand Down Expand Up @@ -168,7 +168,6 @@
"raw_task_spawn\tserver/voice_server.py\tVoiceServer._close_active_ws_connections\tasyncio.create_task\tbcc4bc77941655eb\t0",
"raw_task_spawn\tserver/voice_server.py\tVoiceServer._prepare_listener_cleanup_task\tasyncio.ensure_future\t5c7a78f6e7637ffa\t0",
"raw_task_spawn\tserver/voice_server.py\tVoiceServer._stop_unlocked\tasyncio.create_task\t17a66ec385918dba\t0",
"raw_task_spawn\tserver/webrtc_routes.py\tWebRTCRoutes.cancel_cleanup_tasks\tasyncio.create_task\tf891ff9a3484306c\t0",
"raw_task_spawn\tserver/webrtc_routes.py\t_shutdown_standalone_webrtc\tasyncio.create_task\tb2e5ccf8c0d214f5\t0",
"raw_task_spawn\tsession/_journal_sink.py\tSessionJournalSink.append_record_async\tasyncio.create_task\t933bcec6345d252c\t0",
"raw_task_spawn\tsession/_session.py\tSession._cut_off_turn_playback\tasyncio.create_task\t002fb7c8f4960eaf\t0",
Expand Down
84 changes: 84 additions & 0 deletions tests/server/test_capacity_gate_drain.py
Original file line number Diff line number Diff line change
Expand Up @@ -737,6 +737,90 @@ async def test_webrtc_cleanup_and_drain_share_one_force_escalatable_stop() -> No
assert routes._cleanup_task_scope.tasks() == ()


@pytest.mark.parametrize("finalizer_fails", [False, True])
async def test_webrtc_cleanup_reports_only_unexpected_results_with_task_name(
caplog: pytest.LogCaptureFixture,
monkeypatch: pytest.MonkeyPatch,
*,
finalizer_fails: bool,
) -> None:
routes = WebRTCRoutes(
WebRTCTransportConfig(static_dir=None),
auth=None,
config_factory=lambda _transport: _GracefulSession(), # type: ignore[arg-type]
gate=CapacityGate(max_sessions=1),
manager=SessionManager(),
runtime_feedback=False,
)

async def finalize(_key: int, *, force: bool) -> None:
if force and finalizer_fails:
raise RuntimeError("forced finalizer failed")

monkeypatch.setattr(routes, "_finalize_session_cleanup", finalize)
cleanup = routes._start_cleanup_task(
7,
_NeverClosedTransport(), # type: ignore[arg-type]
)

with caplog.at_level("ERROR", logger="easycat.server.webrtc_routes"):
await routes.cancel_cleanup_tasks()

messages = [record.getMessage() for record in caplog.records]
expected = "WebRTC cleanup task easycat-webrtc-force-cleanup-7 failed"
if finalizer_fails:
assert expected in messages
else:
assert not any(message.startswith("WebRTC cleanup task ") for message in messages)
assert cleanup.cancelled()
assert routes._cleanup_task_scope.tasks() == ()
assert routes._force_cleanup_task_scope.tasks() == ()


async def test_webrtc_cleanup_reports_forced_finalizer_failure_after_timeout(
caplog: pytest.LogCaptureFixture,
monkeypatch: pytest.MonkeyPatch,
) -> None:
routes = WebRTCRoutes(
WebRTCTransportConfig(static_dir=None),
auth=None,
config_factory=lambda _transport: _GracefulSession(), # type: ignore[arg-type]
gate=CapacityGate(max_sessions=1),
manager=SessionManager(),
runtime_feedback=False,
)
cancellation_seen = asyncio.Event()
release = asyncio.Event()

async def finalize(_key: int, *, force: bool) -> None:
assert force is True
try:
await release.wait()
except asyncio.CancelledError:
cancellation_seen.set()
await release.wait()
raise RuntimeError("late forced finalizer failed")

monkeypatch.setattr(routes, "_finalize_session_cleanup", finalize)
routes._start_cleanup_task(
7,
_NeverClosedTransport(), # type: ignore[arg-type]
)

with caplog.at_level("ERROR", logger="easycat.server.webrtc_routes"):
await routes.cancel_cleanup_tasks(timeout_s=0.0)
await cancellation_seen.wait()
release.set()
for _ in range(20):
if routes._force_cleanup_task_scope.tasks() == ():
break
await asyncio.sleep(0)

assert "WebRTC cleanup task easycat-webrtc-force-cleanup-7 failed" in caplog.text
assert "late forced finalizer failed" in caplog.text
assert routes._force_cleanup_task_scope.tasks() == ()


async def test_drain_with_no_active_sessions_is_a_noop() -> None:
gate: CapacityGate[int] = CapacityGate(max_sessions=4)
await gate.drain(list, drain_timeout_s=1.0, force_after=True)
Expand Down