From aee1847303d03f1c8ea726ee9d275db547016a18 Mon Sep 17 00:00:00 2001 From: leo <260626284+cameledev@users.noreply.github.com> Date: Wed, 24 Jun 2026 17:44:04 +0200 Subject: [PATCH] =?UTF-8?q?=E2=9C=A8(backend)=20add=20LiveKit=20egress=5Fe?= =?UTF-8?q?nded=20fallback=20for=20saving=20recordings?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- CHANGELOG.md | 1 + docs/features/recording.md | 2 +- docs/installation/kubernetes.md | 2 +- src/backend/core/api/viewsets.py | 26 ++-- src/backend/core/recording/event/parsers.py | 2 +- .../recording/services/recording_events.py | 25 +++ src/backend/core/services/livekit_events.py | 37 ++++- .../test_api_recordings_storage_hook.py | 41 +++++ .../tests/services/test_livekit_events.py | 147 +++++++++++++++++- 9 files changed, 254 insertions(+), 29 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6b1249a4..6baafde4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,7 @@ and this project adheres to - ✨(backend) add command to clean pending and deleted files - 🧱(helm) run clean files command as cronjob +- ✨(backend) add fallback to save recordings without S3/MinIO webhooks ### Changed diff --git a/docs/features/recording.md b/docs/features/recording.md index ec566475..13cc06cb 100644 --- a/docs/features/recording.md +++ b/docs/features/recording.md @@ -96,7 +96,7 @@ sequenceDiagram | **RECORDING_WORKER_CLASSES** | Dict | `{ "screen_recording": "core.recording.worker.services.VideoCompositeEgressService", "transcript": "core.recording.worker.services.AudioCompositeEgressService" }` | Maps recording types to their worker service classes. | | **RECORDING_EVENT_PARSER_CLASS** | String | `"core.recording.event.parsers.MinioParser"` | Class responsible for parsing storage events and updating the backend. | | **RECORDING_ENABLE_STORAGE_EVENT_AUTH** | Boolean | `True` | Enable authentication for storage event webhook requests. | -| **RECORDING_STORAGE_EVENT_ENABLE** | Boolean | `False` | Enable handling of storage events (must configure webhook in storage). | +| **RECORDING_STORAGE_EVENT_ENABLE** | Boolean | `False` | Enable handling of storage events (must configure webhook in storage). If `False`, fallback to LiveKit egress complete webhook. | | **RECORDING_STORAGE_EVENT_TOKEN** | Secret/File | `None` | Token used to authenticate storage webhook requests, if `RECORDING_ENABLE_STORAGE_EVENT_AUTH` is enabled. | | **RECORDING_EXPIRATION_DAYS** | Integer | `None` | Number of days before recordings expire. Should match bucket lifecycle policy. Set to `None` for no expiration. | | **RECORDING_MAX_DURATION** | Integer | `None` | Maximum duration of a recording in milliseconds. Must be synced with the LiveKit Egress configuration. Set to None for unlimited duration. When the maximum duration is reached, the recording is automatically stopped and saved, and the user is prompted in the frontend with an alert message. | diff --git a/docs/installation/kubernetes.md b/docs/installation/kubernetes.md index bf582bea..746df027 100644 --- a/docs/installation/kubernetes.md +++ b/docs/installation/kubernetes.md @@ -344,7 +344,7 @@ These are the environmental options available on meet backend. | RECORDING_WORKER_CLASSES | Worker classes for recording | {"screen_recording": "core.recording.worker.services.VideoCompositeEgressService","transcript": "core.recording.worker.services.AudioCompositeEgressService"} | | RECORDING_EVENT_PARSER_CLASS | Storage event engine for recording | core.recording.event.parsers.MinioParser | | RECORDING_ENABLE_STORAGE_EVENT_AUTH | Enable storage event authorization | true | -| RECORDING_STORAGE_EVENT_ENABLE | Enable recording storage events | false | +| RECORDING_STORAGE_EVENT_ENABLE | Enable recording storage events. If false, fallback to egress webhook. | false | | RECORDING_STORAGE_EVENT_TOKEN | Recording storage event token | | | RECORDING_EXPIRATION_DAYS | Recording expiration in days | | | RECORDING_MAX_DURATION | Maximum recording duration in milliseconds. Must match LiveKit Egress configuration exactly. | | diff --git a/src/backend/core/api/viewsets.py b/src/backend/core/api/viewsets.py index e731cf42..e1777e62 100644 --- a/src/backend/core/api/viewsets.py +++ b/src/backend/core/api/viewsets.py @@ -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."}, diff --git a/src/backend/core/recording/event/parsers.py b/src/backend/core/recording/event/parsers.py index 02f6db10..9e6a71ea 100644 --- a/src/backend/core/recording/event/parsers.py +++ b/src/backend/core/recording/event/parsers.py @@ -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: diff --git a/src/backend/core/recording/services/recording_events.py b/src/backend/core/recording/services/recording_events.py index 8fac30aa..3bebac0d 100644 --- a/src/backend/core/recording/services/recording_events.py +++ b/src/backend/core/recording/services/recording_events.py @@ -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() diff --git a/src/backend/core/services/livekit_events.py b/src/backend/core/services/livekit_events.py index e8a8822a..5852f6d4 100644 --- a/src/backend/core/services/livekit_events.py +++ b/src/backend/core/services/livekit_events.py @@ -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.""" diff --git a/src/backend/core/tests/recording/test_api_recordings_storage_hook.py b/src/backend/core/tests/recording/test_api_recordings_storage_hook.py index 2ef88798..5603523c 100644 --- a/src/backend/core/tests/recording/test_api_recordings_storage_hook.py +++ b/src/backend/core/tests/recording/test_api_recordings_storage_hook.py @@ -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 + ) diff --git a/src/backend/core/tests/services/test_livekit_events.py b/src/backend/core/tests/services/test_livekit_events.py index 6ace6d7a..56d0e0e5 100644 --- a/src/backend/core/tests/services/test_livekit_events.py +++ b/src/backend/core/tests/services/test_livekit_events.py @@ -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(