mirror of
https://github.com/suitenumerique/meet.git
synced 2026-08-10 02:39:41 +00:00
f4569c64e5
Updating room access rewrote the entire metadata payload, removing information about active recordings. This caused the frontend to lose track of ongoing recordings and could trigger 409 errors when attempting to start a new recording. Consolidate the duplicated metadata update logic into `RoomManagement.update_metadata()` and preserve merge behavior instead of overwriting the full metadata object.
282 lines
8.9 KiB
Python
282 lines
8.9 KiB
Python
"""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
|
|
from core.recording.services.metadata_collector import (
|
|
MetadataCollectorException,
|
|
MetadataCollectorService,
|
|
)
|
|
from core.recording.services.recording_events import (
|
|
RecordingEventsError,
|
|
RecordingEventsService,
|
|
RecordingNotSavableError,
|
|
)
|
|
|
|
from .lobby import LobbyService
|
|
from .room_management import (
|
|
RoomManagement,
|
|
RoomManagementException,
|
|
RoomNotFoundException,
|
|
)
|
|
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)
|
|
RoomManagement().update_metadata(
|
|
room_name, remove_keys=["recording_mode", "recording_status"]
|
|
)
|
|
except RoomNotFoundException:
|
|
logger.info(
|
|
"LiveKit room %s no longer exists, skipping metadata update",
|
|
room_name,
|
|
)
|
|
except RoomManagementException 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
|