(backend) add LiveKit egress_ended fallback for saving recordings

Recording lifecycle previously relied exclusively on the MinIO storage-hook
endpoint to transition from STOPPED to SAVED, which made the system dependent
on MinIO bucket notifications and lifecycle configuration. Introduce support
for LiveKit egress_ended webhook (EGRESS_COMPLETE, EGRESS_LIMIT_REACHED) as
 an alternative finalization mechanism for self-hosted deployments that do
 not use MinIO / S3 hooks.

Change behavior of existing configuation RECORDING_STORAGE_EVENT_ENABLE.
When the LiveKit mechanism is enabled (RECORDING_STORAGE_EVENT_ENABLE=False),
RecordingEventsService.handle_complete is triggered from
LiveKitEventsService._handle_egress_ended.
This commit is contained in:
leo
2026-06-24 17:44:04 +02:00
committed by aleb_the_flash
parent 6b7cd8ab2e
commit aee1847303
9 changed files with 254 additions and 29 deletions
+10 -16
View File
@@ -46,12 +46,15 @@ from core.recording.event.exceptions import (
InvalidFileTypeError,
ParsingEventDataError,
)
from core.recording.event.notification import notification_service
from core.recording.event.parsers import get_parser
from core.recording.services.metadata_collector import (
MetadataCollectorException,
MetadataCollectorService,
)
from core.recording.services.recording_events import (
RecordingEventsService,
RecordingNotSavableError,
)
from core.recording.worker.exceptions import (
RecordingStartError,
RecordingStopError,
@@ -972,24 +975,15 @@ class RecordingViewSet(
except models.Recording.DoesNotExist as e:
raise drf_exceptions.NotFound("No recording found for this event.") from e
if not recording.is_savable():
# Save recording
recording_events_service = RecordingEventsService()
try:
recording_events_service.handle_complete(recording)
except RecordingNotSavableError:
raise drf_exceptions.PermissionDenied(
f"Recording with ID {recording_id} cannot be saved because it is either,"
" in an error state or has already been saved."
)
# Attempt to notify external services about the recording
# This is a non-blocking operation - failures are logged but don't interrupt the flow
notification_succeeded = notification_service.notify_external_services(
recording
)
recording.status = (
models.RecordingStatusChoices.NOTIFICATION_SUCCEEDED
if notification_succeeded
else models.RecordingStatusChoices.SAVED
)
recording.save()
) from None
return drf_response.Response(
{"message": "Event processed."},
+1 -1
View File
@@ -58,7 +58,7 @@ class EventParser(Protocol):
def parse(self, data: Dict) -> StorageEvent:
"""Extract storage event data from raw dictionary input."""
def validate(self, data: StorageEvent) -> None:
def validate(self, data: StorageEvent) -> str:
"""Verify storage event data meets all requirements."""
def get_recording_id(self, data: Dict) -> str:
@@ -8,6 +8,7 @@ from livekit import api
from core import models, utils
from core.models import Recording
from core.recording.event.notification import notification_service
logger = getLogger(__name__)
@@ -16,6 +17,10 @@ class RecordingEventsError(Exception):
"""Recording event handling fails."""
class RecordingNotSavableError(Exception):
"""Recording cannot be saved because it is either in an error state or has already been saved"""
class RecordingEventsService:
"""Handles recording-related LiveKit webhook events."""
@@ -73,3 +78,23 @@ class RecordingEventsService:
f"Failed to notify participants in room '{recording.room.id}' about "
f"recording limit reached (recording_id={recording.id})"
) from e
@staticmethod
def handle_complete(recording: Recording):
"""Notify external services and save recording."""
if not recording.is_savable():
raise RecordingNotSavableError
# Attempt to notify external services about the recording
# This is a non-blocking operation - failures are logged but don't interrupt the flow
notification_succeeded = notification_service.notify_external_services(
recording
)
recording.status = (
models.RecordingStatusChoices.NOTIFICATION_SUCCEEDED
if notification_succeeded
else models.RecordingStatusChoices.SAVED
)
recording.save()
+30 -7
View File
@@ -19,6 +19,7 @@ from core.recording.services.metadata_collector import (
from core.recording.services.recording_events import (
RecordingEventsError,
RecordingEventsService,
RecordingNotSavableError,
)
from .lobby import LobbyService
@@ -88,6 +89,13 @@ class LiveKitEventsService:
def __init__(self):
"""Initialize with required services."""
self._webhook_handlers = {
"egress_updated": self._handle_egress_updated,
"egress_ended": self._handle_egress_ended,
"room_started": self._handle_room_started,
"room_finished": self._handle_room_finished,
}
token_verifier = api.TokenVerifier(
settings.LIVEKIT_CONFIGURATION["api_key"],
settings.LIVEKIT_CONFIGURATION["api_secret"],
@@ -135,14 +143,11 @@ class LiveKitEventsService:
f"Unknown webhook type: {data.event}"
) from e
handler_name = f"_handle_{webhook_type.value}"
handler = getattr(self, handler_name, None)
# Handle according to received webhook type
handler = self._webhook_handlers.get(webhook_type.value)
if not handler or not callable(handler):
return
# pylint: disable=not-callable
handler(data)
if handler is not None:
handler(data)
def _handle_egress_updated(self, data):
"""Handle 'egress_updated' event."""
@@ -195,6 +200,24 @@ class LiveKitEventsService:
f"Failed to process limit reached event for recording {recording}"
) from e
# Fallback for completion when no MinIO/S3 webhooks are configured
if (
not settings.RECORDING_STORAGE_EVENT_ENABLE
) and data.egress_info.status in [
api.EgressStatus.EGRESS_COMPLETE,
api.EgressStatus.EGRESS_LIMIT_REACHED,
]:
try:
self.recording_events.handle_complete(recording)
except RecordingNotSavableError:
logger.warning(
"Recording %s is not savable on egress complete "
"(already saved or in an error state); ignoring.",
recording.id,
)
# Silently ignoring EGRESS_ABORTED, EGRESS_FAILED
def _handle_room_started(self, data):
"""Handle 'room_started' event."""
@@ -224,3 +224,44 @@ def test_save_recording_success(recording_settings, mock_get_parser, client, sta
recording.refresh_from_db()
assert recording.status == RecordingStatusChoices.SAVED
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@pytest.mark.parametrize("notification_succeeded", [True, False])
def test_save_recording_notifies_external_services(
mock_notify_external_services,
recording_settings,
mock_get_parser,
client,
notification_succeeded,
):
"""External services should be notified when a recording is saved."""
recording = RecordingFactory(status="active")
mock_parser = mock.Mock()
mock_parser.get_recording_id.return_value = recording.id
mock_get_parser.return_value = mock_parser
mock_notify_external_services.return_value = notification_succeeded
response = client.post(
"/api/v1.0/recordings/storage-hook/",
{"recording_data": "valid-data"},
HTTP_AUTHORIZATION="Bearer testAuthToken",
)
assert response.status_code == 200
assert response.json() == {"message": "Event processed."}
mock_notify_external_services.assert_called_once_with(recording)
recording.refresh_from_db()
assert recording.status == (
RecordingStatusChoices.NOTIFICATION_SUCCEEDED
if notification_succeeded
else RecordingStatusChoices.SAVED
)
@@ -91,7 +91,9 @@ def test_handle_egress_ended_success(
)
recording.refresh_from_db()
assert recording.status == "stopped"
# NB: notify_external_services will return False, so status is "saved"
assert recording.status == "saved"
@pytest.mark.parametrize(
@@ -155,7 +157,7 @@ def test_handle_egress_updated_non_handled(
def test_handle_egress_ended_metadata_update_fails(
mock_update_room_metadata, mock_notify, mode, notification_type, service
):
"""Should successfully stop recording when metadata's update fails."""
"""Should successfully stop and save recording when metadata's update fails."""
recording = RecordingFactory(worker_id="worker-1", mode=mode, status="active")
mock_data = mock.MagicMock()
@@ -170,7 +172,9 @@ def test_handle_egress_ended_metadata_update_fails(
room_name=str(recording.room.id), notification_data={"type": notification_type}
)
recording.refresh_from_db()
assert recording.status == "stopped"
# NB: notify_external_services will return False, so status is "saved"
assert recording.status == "saved"
@mock.patch("core.utils.notify_participants")
@@ -326,6 +330,143 @@ def test_handle_egress_ended_does_not_call_metadata_collector_stop_when_conditio
mock_collector.stop.assert_not_called()
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.utils.update_room_metadata")
@pytest.mark.parametrize(
"egress_status",
[EgressStatus.EGRESS_COMPLETE, EgressStatus.EGRESS_LIMIT_REACHED],
)
@pytest.mark.parametrize(
"notify_return_value, recording_status",
[(True, "notification_succeeded"), (False, "saved")],
)
def test_handle_egress_ended_finalizes_recording( # noqa: PLR0913
mock_update_room_metadata,
mock_notify,
mock_notify_external_services,
notify_return_value,
recording_status,
egress_status,
service,
settings,
): # pylint: disable=too-many-arguments,too-many-positional-arguments
"""Should notify external services and save the recording on egress completion
(EGRESS_COMPLETE or EGRESS_LIMIT_REACHED) when RECORDING_STORAGE_EVENT_ENABLE is False.
"""
settings.RECORDING_STORAGE_EVENT_ENABLE = False
mock_notify_external_services.return_value = notify_return_value
recording = RecordingFactory(worker_id="worker-1", status="active")
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = egress_status
service._handle_egress_ended(mock_data)
mock_notify_external_services.assert_called_once_with(recording)
recording.refresh_from_db()
assert recording.status == recording_status
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.utils.update_room_metadata")
@pytest.mark.parametrize(
"egress_status, expected_status",
[
(EgressStatus.EGRESS_COMPLETE, "active"),
(EgressStatus.EGRESS_LIMIT_REACHED, "stopped"),
],
)
def test_handle_egress_ended_does_not_finalize_when_webhooks_enabled( # noqa: PLR0913
mock_update_room_metadata,
mock_notify,
mock_notify_external_services,
egress_status,
expected_status,
service,
settings,
): # pylint: disable=too-many-arguments,too-many-positional-arguments
"""When storage event webhooks are enabled, egress_ended must not finalize the
recording: external services are never notified. EGRESS_LIMIT_REACHED still stops
the recording, EGRESS_COMPLETE leaves it active.
"""
settings.RECORDING_STORAGE_EVENT_ENABLE = True
recording = RecordingFactory(worker_id="worker-1", status="active")
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = egress_status
service._handle_egress_ended(mock_data)
mock_notify_external_services.assert_not_called()
recording.refresh_from_db()
assert recording.status == expected_status
@pytest.mark.parametrize(
"egress_status",
[
EgressStatus.EGRESS_STARTING,
EgressStatus.EGRESS_ACTIVE,
EgressStatus.EGRESS_ENDING,
EgressStatus.EGRESS_FAILED,
EgressStatus.EGRESS_ABORTED,
],
)
@mock.patch("core.utils.update_room_metadata")
def test_handle_egress_ended_does_not_save_on_wrong_status(
mock_update_room_metadata, egress_status, service, settings
):
"""Shouldn't save on invalid status."""
settings.RECORDING_STORAGE_EVENT_ENABLE = False
recording = RecordingFactory(worker_id="worker-1", status="active")
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = egress_status
service._handle_egress_ended(mock_data)
recording.refresh_from_db()
assert recording.status == "active"
@pytest.mark.parametrize(
"status", ["failed_to_start", "aborted", "failed_to_stop", "saved", "initiated"]
)
@mock.patch("core.utils.update_room_metadata")
def test_handle_egress_ended_ignores_non_savable_recording(
mock_update_room_metadata, status, service, settings
):
"""Should handle non-savable recordings idempotently without raising.
'egress_ended' may be redelivered (e.g. for an already-saved recording);
this must not raise, otherwise the webhook would 500 and LiveKit would retry.
"""
settings.RECORDING_STORAGE_EVENT_ENABLE = False
recording = RecordingFactory(worker_id="worker-1", status=status)
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = EgressStatus.EGRESS_COMPLETE
service._handle_egress_ended(mock_data)
recording.refresh_from_db()
assert recording.status == status
@mock.patch.object(LobbyService, "clear_room_cache")
@mock.patch.object(TelephonyService, "delete_dispatch_rule")
def test_handle_room_finished_clears_cache_and_deletes_dispatch_rule(