"""LiveKit Events Service""" # pylint: disable=no-member import re import uuid from enum import Enum from logging import getLogger from django.conf import settings from livekit import api from core import models, utils from core.recording.services.metadata_collector import ( MetadataCollectorException, MetadataCollectorService, ) from core.recording.services.recording_events import ( RecordingEventsError, RecordingEventsService, RecordingNotSavableError, ) from .lobby import LobbyService from .telephony import TelephonyException, TelephonyService logger = getLogger(__name__) class LiveKitWebhookError(Exception): """Base exception for LiveKit webhook processing errors.""" status_code = 500 class AuthenticationError(LiveKitWebhookError): """Authentication failed.""" status_code = 401 class InvalidPayloadError(LiveKitWebhookError): """Invalid webhook payload.""" status_code = 400 class UnsupportedEventTypeError(LiveKitWebhookError): """Unsupported event type.""" status_code = 422 class ActionFailedError(LiveKitWebhookError): """Webhook action fails to process or complete.""" status_code = 500 class LiveKitWebhookEventType(Enum): """LiveKit webhook event types.""" # Room events ROOM_STARTED = "room_started" ROOM_FINISHED = "room_finished" # Participant events PARTICIPANT_JOINED = "participant_joined" PARTICIPANT_LEFT = "participant_left" # Track events TRACK_PUBLISHED = "track_published" TRACK_UNPUBLISHED = "track_unpublished" # Egress events EGRESS_STARTED = "egress_started" EGRESS_UPDATED = "egress_updated" EGRESS_ENDED = "egress_ended" # Ingress events INGRESS_STARTED = "ingress_started" INGRESS_ENDED = "ingress_ended" class LiveKitEventsService: """Service for processing and handling LiveKit webhook events and notifications.""" 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"], ) self.webhook_receiver = api.WebhookReceiver(token_verifier) self.lobby_service = LobbyService() self.telephony_service = TelephonyService() self.recording_events = RecordingEventsService() self._filter_regex = None if settings.LIVEKIT_WEBHOOK_EVENTS_FILTER_REGEX: try: self._filter_regex = re.compile( settings.LIVEKIT_WEBHOOK_EVENTS_FILTER_REGEX ) except re.error: logger.exception( "Invalid LIVEKIT_WEBHOOK_EVENTS_FILTER_REGEX. Webhook filtering disabled." ) def receive(self, request): """Process webhook and route to appropriate handler.""" auth_token = request.headers.get("Authorization") if not auth_token: raise AuthenticationError("Authorization header missing") try: data = self.webhook_receiver.receive( request.body.decode("utf-8"), auth_token ) except Exception as e: raise InvalidPayloadError("Invalid webhook payload") from e room_name = data.room.name or data.egress_info.room_name if self._filter_regex and not self._filter_regex.search(room_name): logger.info("Filtered webhook event for room '%s'", room_name) return try: webhook_type = LiveKitWebhookEventType(data.event) except ValueError as e: raise UnsupportedEventTypeError( f"Unknown webhook type: {data.event}" ) from e # Handle according to received webhook type handler = self._webhook_handlers.get(webhook_type.value) if handler is not None: handler(data) def _handle_egress_updated(self, data): """Handle 'egress_updated' event.""" egress_id = data.egress_info.egress_id try: recording = models.Recording.objects.get(worker_id=egress_id) except models.Recording.DoesNotExist as err: raise ActionFailedError( f"Recording with worker ID {egress_id} does not exist" ) from err egress_status = data.egress_info.status self.recording_events.handle_update(recording, egress_status) def _handle_egress_ended(self, data): """Handle 'egress_ended' event.""" try: recording = models.Recording.objects.select_related("room").get( worker_id=data.egress_info.egress_id ) except models.Recording.DoesNotExist as err: raise ActionFailedError( f"Recording with worker ID {data.egress_info.egress_id} does not exist" ) from err try: room_name = str(recording.room.id) utils.update_room_metadata( room_name, {}, ["recording_mode", "recording_status"] ) except utils.MetadataUpdateException as e: logger.exception("Failed to update room's metadata: %s", e) if recording.options.get("metadata_collector_dispatch_id", None) is not None: try: MetadataCollectorService().stop(recording) except MetadataCollectorException: logger.warning("Failed to stop the MetadataCollectorService") if ( data.egress_info.status == api.EgressStatus.EGRESS_LIMIT_REACHED and recording.status == models.RecordingStatusChoices.ACTIVE ): try: self.recording_events.handle_limit_reached(recording) except RecordingEventsError as e: raise ActionFailedError( 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.""" try: room_id = uuid.UUID(data.room.name) except ValueError as e: logger.warning( "Ignoring room event: room name '%s' is not a valid UUID format.", data.room.name, ) raise ActionFailedError("Failed to process room started event") from e try: room = models.Room.objects.get(id=room_id) except models.Room.DoesNotExist as err: raise ActionFailedError(f"Room with ID {room_id} does not exist") from err if settings.ROOM_TELEPHONY_ENABLED: try: self.telephony_service.create_dispatch_rule(room) except TelephonyException as e: raise ActionFailedError( f"Failed to create telephony dispatch rule for room {room_id}" ) from e def _handle_room_finished(self, data): """Handle 'room_finished' event.""" try: room_id = uuid.UUID(data.room.name) except ValueError as e: logger.warning( "Ignoring room event: room name '%s' is not a valid UUID format.", data.room.name, ) raise ActionFailedError("Failed to process room finished event") from e if settings.ROOM_TELEPHONY_ENABLED: try: self.telephony_service.delete_dispatch_rule(room_id) except TelephonyException as e: raise ActionFailedError( f"Failed to delete telephony dispatch rule for room {room_id}" ) from e try: self.lobby_service.clear_room_cache(room_id) except Exception as e: raise ActionFailedError( f"Failed to clear room cache for room {room_id}" ) from e