Compare commits

..

1 Commits

Author SHA1 Message Date
briquet ff662a1c8e 🔒️(backend) upgrade base image to python:3.13.15-alpine3.24
Fixes an XML injection attack in libexpat (CVE-2026-93990)
2026-09-24 19:52:33 +02:00
42 changed files with 300 additions and 1378 deletions
+1 -6
View File
@@ -16,7 +16,6 @@ and this project adheres to
### Changed
- 🔥(backend) remove unused API viewset and permission helpers
- 🔊(backend) pin the dockerflow logger level to WARNING
- 🚑️(summary) serve health endpoints with the dockerflow router
- ♻️(backend) serve the dockerflow views early in the middleware stack
@@ -30,7 +29,6 @@ and this project adheres to
- ⬆️(frontend) upgrade humanize-duration from 3.33.2 to 3.34.1
- ⬆️(addons) upgrade i18next from 26.4.0 to 26.4.1
- 🔖(helm) release chart 0.0.28
- ♻️(backend) decouple recording event handling from LiveKit egress statuses
### Fixed
@@ -44,10 +42,7 @@ and this project adheres to
- 🔒️(backend) reject inactive users in resource server backend
- 🐛(frontend) fix file permissions in the Docker image
- 🚸(frontend) inform user that recording waits until a track is published
- 🔒(backend) upgrade base image to python:3.13.5-alpine3.24
- 🐛(backend) handle failed and aborted egresses
- 🩹(frontend) notify participants when a recording fails or is aborted
- 🔒️(frontend) fix HIGH CVE-2026-93990 in libexpat
- 🔒(backend) upgrade base image to python:3.13.15-alpine3.24
## [1.31.0] - 2026-09-08
+1 -2
View File
@@ -57,8 +57,7 @@ RUN npx webpack --mode production
FROM nginxinc/nginx-unprivileged:1.30.4-alpine3.24 AS frontend-production
USER root
RUN apk upgrade --no-cache libexpat && \
apk del curl
RUN apk del curl
USER nginx
USER nginx
-3
View File
@@ -76,9 +76,6 @@ def get_frontend_configuration(request):
"authenticated_users_can_edit_display_name": (
settings.AUTHENTICATED_PARTICIPANTS_CAN_EDIT_DISPLAY_NAME
),
"encryption": {
"is_enabled": settings.ENCRYPTION_ENABLED,
},
}
frontend_configuration.update(settings.FRONTEND_CONFIGURATION)
return Response(frontend_configuration)
+13
View File
@@ -12,6 +12,10 @@ from ..services.participants_management import (
ParticipantsManagementException,
)
ACTION_FOR_METHOD_TO_PERMISSION = {
"versions_detail": {"DELETE": "versions_destroy", "GET": "versions_retrieve"}
}
class IsAuthenticated(permissions.BasePermission):
"""
@@ -23,6 +27,15 @@ class IsAuthenticated(permissions.BasePermission):
return bool(request.auth) or request.user.is_authenticated
class IsAuthenticatedOrSafe(IsAuthenticated):
"""Allows access to authenticated users (or anonymous users but only on safe methods)."""
def has_permission(self, request, view):
if request.method in permissions.SAFE_METHODS:
return True
return super().has_permission(request, view)
class IsSelf(IsAuthenticated):
"""
Allows access only to authenticated users. Alternative method checking the presence
+1 -47
View File
@@ -38,7 +38,6 @@ class UserSerializer(serializers.ModelSerializer):
"short_name",
"timezone",
"language",
"default_encryption_mode",
"default_room_access_level",
"default_room_configuration",
]
@@ -54,14 +53,6 @@ class UserSerializer(serializers.ModelSerializer):
raise serializers.ValidationError(e.errors()) from e
return value
def validate_default_encryption_mode(self, value):
"""Reject a non-none default when the server has encryption disabled."""
if value != models.EncryptionMode.NONE and not settings.ENCRYPTION_ENABLED:
raise serializers.ValidationError(
_("End-to-end encryption is disabled on this server.")
)
return value
class UserLightSerializer(serializers.ModelSerializer):
"""Serialize users with limited fields."""
@@ -103,7 +94,6 @@ class ResourceAccessSerializerMixin:
raise PermissionDenied(
"Only owners of a room can assign other users as owners."
)
return data
def validate_resource(self, resource):
@@ -158,15 +148,7 @@ class RoomSerializer(serializers.ModelSerializer):
class Meta:
model = models.Room
fields = [
"id",
"name",
"slug",
"configuration",
"access_level",
"pin_code",
"encryption_mode",
]
fields = ["id", "name", "slug", "configuration", "access_level", "pin_code"]
read_only_fields = ["id", "slug", "pin_code"]
def validate_configuration(self, value):
@@ -179,32 +161,6 @@ class RoomSerializer(serializers.ModelSerializer):
raise serializers.ValidationError(e.errors()) from e
return value
def validate_encryption_mode(self, value):
"""Encryption mode is part of the link's semantics (the passphrase
lives in the URL hash for `basic` rooms) so it cannot be changed once
the room exists."""
instance = self.instance
if instance and instance.encryption_mode != value:
raise serializers.ValidationError(
"Encryption mode cannot be changed after room creation."
)
return value
def validate_access_level(self, value):
"""Encrypted rooms must stay restricted — the lobby is the only way
to enforce per-participant admission, and basic encryption relies on
the host vetting each joiner before they receive the in-URL key."""
instance = self.instance
if (
instance
and instance.encryption_mode != models.EncryptionMode.NONE
and value != models.RoomAccessLevel.RESTRICTED
):
raise serializers.ValidationError(
"Encrypted rooms require restricted access level."
)
return value
def to_representation(self, instance):
"""
Add users only for administrator users.
@@ -241,14 +197,12 @@ class RoomSerializer(serializers.ModelSerializer):
if should_access_room:
room_id = f"{instance.id!s}"
username = request.query_params.get("username", None)
output["livekit"] = utils.generate_livekit_config(
room_id=room_id,
user=request.user,
username=username,
configuration=output["configuration"],
role=role,
encryption_mode=instance.encryption_mode,
)
else:
del output["pin_code"]
+56 -24
View File
@@ -95,6 +95,60 @@ from .feature_flag import FeatureFlag
logger = getLogger(__name__)
class NestedGenericViewSet(viewsets.GenericViewSet):
"""
A generic Viewset aims to be used in a nested route context.
e.g: `/api/v1.0/resource_1/<resource_1_pk>/resource_2/<resource_2_pk>/`
It allows to define all url kwargs and lookup fields to perform the lookup.
"""
lookup_fields: list[str] = ["pk"]
lookup_url_kwargs: list[str] = []
def __getattribute__(self, file):
"""
This method is overridden to allow to get the last lookup field or lookup url kwarg
when accessing the `lookup_field` or `lookup_url_kwarg` attribute. This is useful
to keep compatibility with all methods used by the parent class `GenericViewSet`.
"""
if file in ["lookup_field", "lookup_url_kwarg"]:
return getattr(self, file + "s", [None])[-1]
return super().__getattribute__(file)
def get_queryset(self):
"""
Get the list of files for this view.
`lookup_fields` attribute is enumerated here to perform the nested lookup.
"""
queryset = super().get_queryset()
# The last lookup field is removed to perform the nested lookup as it corresponds
# to the object pk, it is used within get_object method.
lookup_url_kwargs = (
self.lookup_url_kwargs[:-1]
if self.lookup_url_kwargs
else self.lookup_fields[:-1]
)
filter_kwargs = {}
for index, lookup_url_kwarg in enumerate(lookup_url_kwargs):
if lookup_url_kwarg not in self.kwargs:
raise KeyError(
f"Expected view {self.__class__.__name__} to be called with a URL "
f'keyword argument named "{lookup_url_kwarg}". Fix your URL conf, or '
"set the `.lookup_fields` attribute on the view correctly."
)
filter_kwargs.update(
{self.lookup_fields[index]: self.kwargs[lookup_url_kwarg]}
)
return queryset.filter(**filter_kwargs)
class SerializerPerActionMixin:
"""
A mixin to allow to define serializer classes for each action.
@@ -249,18 +303,6 @@ class RoomViewSet(
Apply the user's default room preferences (access level and configuration)
unless the request explicitly provides its own values.
"""
encryption_mode = serializer.validated_data.get(
"encryption_mode", models.EncryptionMode.NONE
)
if (
encryption_mode != models.EncryptionMode.NONE
and not settings.ENCRYPTION_ENABLED
):
raise drf_exceptions.ValidationError(
{"encryption_mode": "Encryption is not enabled on this server."}
)
user = self.request.user
save_kwargs = {}
@@ -325,21 +367,16 @@ class RoomViewSet(
"""Start recording a room."""
serializer = serializers.StartRecordingSerializer(data=request.data)
if not serializer.is_valid():
return drf_response.Response(
{"detail": "Invalid request."},
status=drf_status.HTTP_400_BAD_REQUEST,
{"detail": "Invalid request."}, status=drf_status.HTTP_400_BAD_REQUEST
)
mode = serializer.validated_data["mode"]
options = serializer.validated_data.get("options")
room = self.get_object()
if room.is_encrypted:
raise drf_exceptions.ValidationError(
{"detail": "Recording is unavailable in encrypted rooms."}
)
try:
with transaction.atomic():
recording = models.Recording.objects.create(
@@ -662,11 +699,6 @@ class RoomViewSet(
room = self.get_object()
if room.is_encrypted:
raise drf_exceptions.ValidationError(
{"detail": "Subtitles are unavailable in encrypted rooms."}
)
try:
SubtitleService().start_subtitle(room)
except SubtitleException:
-12
View File
@@ -5,7 +5,6 @@ Core application enums declaration
import re
from django.conf import global_settings, settings
from django.db import models
from django.utils.translation import gettext_lazy as _
UUID_REGEX = (
@@ -33,14 +32,3 @@ ALL_LANGUAGES = getattr(
"ALL_LANGUAGES",
[(language, _(name)) for language, name in global_settings.LANGUAGES],
)
class EncryptionMode(models.TextChoices):
"""Encryption mode for a room.
Kept as an enum (not a boolean) so future modes — e.g. a vault-managed
per-user key flow — can be added without another schema migration.
"""
NONE = "none", _("No encryption")
BASIC = "basic", _("Passphrase-in-URL encryption")
@@ -1,18 +0,0 @@
# Generated by Django 5.2.16 on 2026-09-23 16:49
from django.db import migrations, models
class Migration(migrations.Migration):
dependencies = [
('core', '0022_user_default_room_access_level_and_more'),
]
operations = [
migrations.AlterField(
model_name='recording',
name='status',
field=models.CharField(choices=[('initiated', 'Initiated'), ('active', 'Active'), ('stopped', 'Stopped'), ('saved', 'Saved'), ('aborted', 'Aborted'), ('failed', 'Failed'), ('failed_to_start', 'Failed to Start'), ('failed_to_stop', 'Failed to Stop'), ('notification_succeeded', 'Notification succeeded'), ('external_process_successful', 'External process successful'), ('external_process_failed', 'External process failed')], default='initiated', max_length=50),
),
]
@@ -1,46 +0,0 @@
"""Add Room.encryption_mode and User.default_encryption_mode (enum-based).
We store the mode as an enum (CharField with choices) rather than a boolean
so a future "advanced" mode (per-user vault keys, etc.) can be added without
a schema migration.
"""
from django.db import migrations, models
class Migration(migrations.Migration):
dependencies = [
("core", "0023_alter_recording_status"),
]
operations = [
migrations.AddField(
model_name="room",
name="encryption_mode",
field=models.CharField(
choices=[
("none", "No encryption"),
("basic", "Passphrase-in-URL encryption"),
],
default="none",
help_text="End-to-end encryption mode for this room.",
max_length=20,
verbose_name="Encryption mode",
),
),
migrations.AddField(
model_name="user",
name="default_encryption_mode",
field=models.CharField(
choices=[
("none", "No encryption"),
("basic", "Passphrase-in-URL encryption"),
],
default="none",
help_text="Encryption mode pre-selected when this user creates a new meeting.",
max_length=20,
verbose_name="Default encryption mode",
),
),
]
+8 -76
View File
@@ -26,7 +26,6 @@ from lasuite.tools.email import get_domain_from_email
from timezone_field import TimeZoneField
from . import fields, utils
from .enums import EncryptionMode
from .recording.enums import FileExtension
from .validators import sub_validator
@@ -59,7 +58,6 @@ class RecordingStatusChoices(models.TextChoices):
STOPPED = "stopped", _("Stopped")
SAVED = "saved", _("Saved")
ABORTED = "aborted", _("Aborted")
FAILED = "failed", _("Failed")
FAILED_TO_START = "failed_to_start", _("Failed to Start")
FAILED_TO_STOP = "failed_to_stop", _("Failed to Stop")
NOTIFICATION_SUCCEEDED = "notification_succeeded", _("Notification succeeded")
@@ -81,13 +79,17 @@ class RecordingStatusChoices(models.TextChoices):
cls.STOPPED,
cls.SAVED,
cls.ABORTED,
cls.FAILED,
cls.EXTERNAL_PROCESS_SUCCESSFUL,
cls.EXTERNAL_PROCESS_FAILED,
cls.FAILED_TO_START,
cls.FAILED_TO_STOP,
}
@classmethod
def is_unsuccessful(cls, status):
"""Determine if the recording status represents an unsuccessful state."""
return status in {cls.ABORTED, cls.FAILED_TO_START, cls.FAILED_TO_STOP}
class RecordingModeChoices(models.TextChoices):
"""Recording mode choices."""
@@ -217,15 +219,6 @@ class User(AbstractBaseUser, BaseModel, auth_models.PermissionsMixin):
"Unselect this instead of deleting accounts."
),
)
default_encryption_mode = models.CharField(
_("Default encryption mode"),
max_length=20,
choices=EncryptionMode.choices,
default=EncryptionMode.NONE,
help_text=_(
"Encryption mode pre-selected when this user creates a new meeting."
),
)
objects = auth_models.UserManager()
@@ -421,13 +414,6 @@ class Room(Resource):
choices=RoomAccessLevel.choices,
default=settings.RESOURCE_DEFAULT_ACCESS_LEVEL,
)
encryption_mode = models.CharField(
max_length=20,
choices=EncryptionMode.choices,
default=EncryptionMode.NONE,
verbose_name=_("Encryption mode"),
help_text=_("End-to-end encryption mode for this room."),
)
# Public configuration exposed to any room participant via the API
configuration = models.JSONField(
blank=True,
@@ -454,68 +440,20 @@ class Room(Resource):
return capfirst(self.name)
def save(self, *args, **kwargs):
"""Restrict new encrypted rooms and allocate PINs for unencrypted rooms.
"""Generate a unique n-digit pin code for new rooms."""
Skip PIN allocation for encrypted rooms — the SIP gateway will
always reject calls to them (no way to derive the key), and the
PIN namespace is finite (10**length): no point burning slots that
can never be dialed.
Also run `clean()` so the encryption invariants are enforced on
every save path (ORM, admin, shell), not only via the DRF
serializer.
"""
# Override both explicit access levels and user defaults on creation.
# Updates remain subject to clean() instead of being silently normalized.
if self._state.adding and self.is_encrypted:
self.access_level = RoomAccessLevel.RESTRICTED
self.clean()
# Roomkit devices also join by PIN, so a PIN is needed as soon as
# either integration is enabled.
if (
(settings.ROOM_TELEPHONY_ENABLED or settings.ROOMKIT_ENABLED)
and not self.pk
and not self.pin_code
and self.encryption_mode == EncryptionMode.NONE
):
self.pin_code = self.generate_unique_pin_code(
length=settings.ROOM_TELEPHONY_PIN_LENGTH
)
super().save(*args, **kwargs)
def clean(self):
"""Enforce encryption-mode invariants outside DRF.
Two rules:
- `encryption_mode` is set at creation and never mutated afterwards
(the URL-hash passphrase encodes assumptions about it).
- An encrypted room must be at the RESTRICTED access level so the
host vets joiners before they ever see the in-URL key.
"""
super().clean()
if self.pk is not None:
previous = Room.objects.filter(pk=self.pk).only("encryption_mode").first()
if (
previous is not None
and previous.encryption_mode != self.encryption_mode
):
raise ValidationError(
{
"encryption_mode": _(
"Encryption mode cannot be changed after room creation."
)
}
)
if (
self.encryption_mode != EncryptionMode.NONE
and self.access_level != RoomAccessLevel.RESTRICTED
):
raise ValidationError(
{
"access_level": _(
"Encrypted rooms must use the 'restricted' access level."
)
}
)
def clean_fields(self, exclude=None):
"""
Automatically generate the slug from the name and make sure it does not look like a UUID.
@@ -538,11 +476,6 @@ class Room(Resource):
"""Check if a room is public"""
return self.access_level == RoomAccessLevel.PUBLIC
@property
def is_encrypted(self):
"""Convenience: any non-none encryption mode counts as encrypted."""
return self.encryption_mode != EncryptionMode.NONE
@staticmethod
def generate_unique_pin_code(length):
"""Generate a unique n-digit PIN code"""
@@ -649,7 +582,6 @@ class Recording(BaseModel):
4. NOTIFICATION_SUCCEEDED: External service has been notified of this recording
Error States:
- FAILED: Egress failed mid-recording
- FAILED_TO_START: Worker failed to initialize recording
- FAILED_TO_STOP: Worker failed during stop operation
- ABORTED: Recording was terminated before completion
-47
View File
@@ -8,50 +8,3 @@ class FileExtension(Enum):
OGG = "ogg"
MP4 = "mp4"
class RecordingWorkerEvent(Enum):
"""Lifecycle events a recording worker reports about a recording.
It is intended to be free of SFU-specific vocabulary.
"""
# The worker accepted the request but is not recording yet.
STARTING = "starting"
# The worker is recording.
STARTED = "started"
# The worker stopped recording and is flushing the media file.
SAVING = "saving"
# The recording ended, its media file is available.
COMPLETED = "completed"
# The recording ended on its configured limit, its media file is available.
LIMIT_REACHED = "limit reached"
# The worker stopped before it ever started recording, there is no media file.
ABORTED = "aborted"
# The worker hit a runtime error once recording had started; its media file
# may be available.
FAILED = "failed"
@classmethod
def is_terminal(cls, event):
"""Determine if the event ends the recording's lifecycle (successful or not)."""
return event in TERMINAL_EVENTS
SUCCESSFUL_EVENTS = frozenset(
{
RecordingWorkerEvent.COMPLETED,
RecordingWorkerEvent.LIMIT_REACHED,
}
)
UNSUCCESSFUL_EVENTS = frozenset(
{
RecordingWorkerEvent.ABORTED,
RecordingWorkerEvent.FAILED,
}
)
TERMINAL_EVENTS = SUCCESSFUL_EVENTS | UNSUCCESSFUL_EVENTS
@@ -1,13 +1,13 @@
"""Recording-related Events Service"""
"""Recording-related LiveKit Events Service"""
# pylint: disable=no-member
from logging import getLogger
from livekit import api
from core import models, utils
from core.models import Recording
from core.recording.enums import (
UNSUCCESSFUL_EVENTS,
RecordingWorkerEvent,
)
from core.recording.event.notification import notification_service
from core.services.room_management import (
RoomManagement,
@@ -26,110 +26,22 @@ class RecordingNotSavableError(Exception):
"""Recording cannot be saved because it is either in an error state or has already been saved"""
# Notification sent to the room's participants, per event and recording mode.
NOTIFICATION_PREFIXES = {
models.RecordingModeChoices.SCREEN_RECORDING: "screenRecording",
models.RecordingModeChoices.TRANSCRIPT: "transcription",
}
NOTIFICATION_SUFFIXES = {
RecordingWorkerEvent.LIMIT_REACHED: "LimitReached",
RecordingWorkerEvent.FAILED: "Failed",
RecordingWorkerEvent.ABORTED: "Aborted",
}
def get_notification_type(recording_mode, event):
"""Generate corresponding notification type string."""
try:
return f"{NOTIFICATION_PREFIXES[recording_mode]}{NOTIFICATION_SUFFIXES[event]}"
except KeyError:
return None
# Recording status in the room's metadata, per event.
ROOM_METADATA_RECORDING_STATUSES = {
RecordingWorkerEvent.STARTED: "started",
RecordingWorkerEvent.SAVING: "saving",
}
class RecordingEventsService:
"""Handles recording-related worker events.
Two entry points: `handle_update` for the events a running recording
reports, and `handle_terminal_event` for the one ending it.
"""
"""Handles recording-related LiveKit webhook events."""
@staticmethod
def log_worker_error(recording, event, error=None, error_code=None):
"""Log FAILED at error level and expected ABORTED outcomes at info level."""
if event == RecordingWorkerEvent.FAILED:
log = logger.error
elif event == RecordingWorkerEvent.ABORTED:
log = logger.info
else:
return
log(
"Recording worker reported %s for recording %s (room=%s, mode=%s): %s (error_code=%s)",
event.value,
recording.id,
recording.room.id,
recording.mode,
error or "no error reported",
error_code or "no error_code reported",
)
@staticmethod
def _notify_participants(recording: Recording, event: RecordingWorkerEvent):
"""Notify the room's participants that a recording ended on the given event."""
recording_mode = recording.options.get("original_mode", None) or recording.mode
notification_type = get_notification_type(recording_mode, event)
if not notification_type:
logger.warning(
"Could not find notification type for: "
"room=%s, recording_id=%s, mode=%s, event=%s",
recording.room.id,
recording.id,
recording_mode,
event.value,
)
return
try:
utils.notify_participants(
room_name=str(recording.room.id),
notification_data={"type": notification_type},
)
except utils.NotificationError as e:
raise RecordingEventsError(
f"Failed to notify participants in room '{recording.room.id}' about "
f"recording {event.value} (recording_id={recording.id})"
) from e
@staticmethod
def _log_notification_failure(recording, event: RecordingWorkerEvent):
"""Log a participant notification error on an unsuccessful recording."""
logger.exception(
"Failed to notify participants that recording %s %s (room=%s)",
recording.id,
event.value,
recording.room.id,
)
@staticmethod
def handle_update(recording: Recording, event: RecordingWorkerEvent):
"""Handle non-terminal worker events and sync recording state to room metadata.
Terminal events are dispatched through `handle_terminal_event` instead.
"""
def handle_update(recording: Recording, egress_status):
"""Handle egress status updates and sync recording state to room metadata."""
room_name = str(recording.room.id)
recording_status = ROOM_METADATA_RECORDING_STATUSES.get(event)
status_mapping = {
api.EgressStatus.EGRESS_ACTIVE: "started",
api.EgressStatus.EGRESS_ENDING: "saving",
api.EgressStatus.EGRESS_ABORTED: "aborted",
}
recording_status = status_mapping.get(egress_status)
if recording_status:
try:
RoomManagement.update_metadata(
@@ -143,113 +55,42 @@ class RecordingEventsService:
except RoomManagementException as e:
logger.exception("Failed to update room's metadata: %s", e)
def handle_terminal_event(self, recording: Recording, event: RecordingWorkerEvent):
"""Run the appropriate handlers for a terminal event, given the recording's state."""
if not RecordingWorkerEvent.is_terminal(event):
logger.warning(
"Ignoring non-terminal event %s dispatched as a terminal event "
"for recording %s.",
event.value,
recording.id,
)
return
if event in UNSUCCESSFUL_EVENTS:
self._flag_unsuccessful_recording(recording, event)
else:
self._save_successful_recording(recording, event)
def _flag_unsuccessful_recording(
self, recording: Recording, event: RecordingWorkerEvent
):
"""Persist the outcome of a recording the worker announced as unsuccessful."""
# Aborted
if event == RecordingWorkerEvent.ABORTED:
if recording.status == models.RecordingStatusChoices.ACTIVE:
self._apply_outcome(recording, event, self._handle_aborted)
return
# Failed
if event == RecordingWorkerEvent.FAILED:
if recording.is_savable():
self._apply_outcome(recording, event, self._handle_failed)
return
logger.error(
"Unsuccessful event %s has no handler; recording %s keeps status '%s'.",
event.value,
recording.id,
recording.status,
)
def _save_successful_recording(
self, recording: Recording, event: RecordingWorkerEvent
):
"""Save a recording whose media file the worker made available."""
# Limit reached
if (
event == RecordingWorkerEvent.LIMIT_REACHED
and recording.status == models.RecordingStatusChoices.ACTIVE
):
self._apply_outcome(recording, event, self._handle_limit_reached)
try:
self._handle_successful(recording)
except RecordingNotSavableError:
logger.warning(
"Recording %s is not savable on a completed recording "
"(already saved or in an error state); ignoring.",
recording.id,
)
def _apply_outcome(
self, recording: Recording, event: RecordingWorkerEvent, handler
):
"""Keep notification failure non-fatal."""
try:
handler(recording)
except RecordingEventsError:
self._log_notification_failure(recording, event)
@classmethod
def _handle_limit_reached(cls, recording: Recording):
@staticmethod
def handle_limit_reached(recording: Recording):
"""Stop recording and notify participants when limit is reached."""
recording.status = models.RecordingStatusChoices.STOPPED
recording.save()
cls._notify_participants(recording, RecordingWorkerEvent.LIMIT_REACHED)
notification_mapping = {
models.RecordingModeChoices.SCREEN_RECORDING: "screenRecordingLimitReached",
models.RecordingModeChoices.TRANSCRIPT: "transcriptionLimitReached",
}
@classmethod
def _handle_failed(cls, recording: Recording):
"""Set recording status to failed, matching the worker event, and notify participants.
notification_type = notification_mapping.get(recording.mode)
if not notification_type:
return
FAILED: used when an actual runtime/pipeline error occurs after the
recording has started
"""
recording.status = models.RecordingStatusChoices.FAILED
recording.save()
cls._notify_participants(recording, RecordingWorkerEvent.FAILED)
@classmethod
def _handle_aborted(cls, recording: Recording):
"""Set recording status to aborted, matching the worker event, and notify participants.
ABORTED: used when the worker stops before it ever became
active/recording
"""
recording.status = models.RecordingStatusChoices.ABORTED
recording.save()
cls._notify_participants(recording, RecordingWorkerEvent.ABORTED)
try:
utils.notify_participants(
room_name=str(recording.room.id),
notification_data={"type": notification_type},
)
except utils.NotificationError as e:
logger.exception(
"Failed to notify participants about recording limit reached: "
"room=%s, recording_id=%s, mode=%s",
recording.room.id,
recording.id,
recording.mode,
)
raise RecordingEventsError(
f"Failed to notify participants in room '{recording.room.id}' about "
f"recording limit reached (recording_id={recording.id})"
) from e
@staticmethod
def _handle_successful(recording: Recording):
def handle_complete(recording: Recording):
"""Notify external services and save recording."""
if not recording.is_savable():
+5 -37
View File
@@ -2,8 +2,6 @@
# pylint: disable=no-member
import logging
from asgiref.sync import async_to_sync
from livekit import api as livekit_api
@@ -12,8 +10,6 @@ from ..enums import FileExtension
from .exceptions import WorkerConnectionError, WorkerResponseError
from .factories import WorkerServiceConfig
logger = logging.getLogger(__name__)
class BaseEgressService:
"""Base egress defining common methods to manage and interact with LiveKit egress processes."""
@@ -53,22 +49,6 @@ class BaseEgressService:
finally:
await lkapi.aclose()
@staticmethod
def _log_egress_error(response, event: str):
"""Log the reason LiveKit reported an unsuccessful egress on stop.
Mirrors the logging done in the 'egress_ended' webhook. The
StopEgress response carries the same error fields.
"""
logger.error(
"Egress %s on stop (egress_id=%s, status=%s): %s (error_code=%s)",
event,
response.egress_id,
livekit_api.EgressStatus.Name(response.status),
response.error or "no error reported",
response.error_code or "no error_code reported",
)
def stop(self, worker_id: str) -> str:
"""Stop an ongoing egress worker.
The StopEgressRequest is shared among all types of egress,
@@ -86,26 +66,14 @@ class BaseEgressService:
"LiveKit response is missing the recording status."
)
# To avoid exposing EgressStatus values and coupling with LiveKit outside of this class,
# the response status is mapped to simpler "ABORTED", "STOPPED" or "FAILED_TO_STOP" strings.
if response.status == livekit_api.EgressStatus.EGRESS_ABORTED:
return "ABORTED"
if response.status == livekit_api.EgressStatus.EGRESS_ENDING:
return "STOPPED"
if response.status == livekit_api.EgressStatus.EGRESS_LIMIT_REACHED:
return "STOPPED"
# Cases below should be very infrequent as status changes should be
# received and processed by `handle_ended`, thus `stop` would not
# be called (unless failure and stop are very close in time).
# We therefore accept not to notify the user in this code branch.
# This could be fixed in a future refactoring.
if response.status == livekit_api.EgressStatus.EGRESS_ABORTED:
self._log_egress_error(response, "aborted")
return "ABORTED"
if response.status == livekit_api.EgressStatus.EGRESS_FAILED:
self._log_egress_error(response, "failed")
return "FAILED"
self._log_egress_error(response, "failed to stop")
return "FAILED_TO_STOP"
def start(self, room_name, recording_id):
+34 -55
View File
@@ -12,12 +12,15 @@ from django.conf import settings
from livekit import api
from core import models
from core.recording.enums import RecordingWorkerEvent
from core.recording.services.metadata_collector import (
MetadataCollectorException,
MetadataCollectorService,
)
from core.recording.services.recording_events import RecordingEventsService
from core.recording.services.recording_events import (
RecordingEventsError,
RecordingEventsService,
RecordingNotSavableError,
)
from .lobby import LobbyService
from .presence import PresenceCache
@@ -81,30 +84,6 @@ class LiveKitWebhookEventType(Enum):
INGRESS_ENDED = "ingress_ended"
# LiveKit egress statuses mapped to recording worker event statuses
EGRESS_STATUS_TO_RECORDING_EVENT = {
api.EgressStatus.EGRESS_STARTING: RecordingWorkerEvent.STARTING,
api.EgressStatus.EGRESS_ACTIVE: RecordingWorkerEvent.STARTED,
api.EgressStatus.EGRESS_ENDING: RecordingWorkerEvent.SAVING,
api.EgressStatus.EGRESS_COMPLETE: RecordingWorkerEvent.COMPLETED,
api.EgressStatus.EGRESS_LIMIT_REACHED: RecordingWorkerEvent.LIMIT_REACHED,
api.EgressStatus.EGRESS_ABORTED: RecordingWorkerEvent.ABORTED,
api.EgressStatus.EGRESS_FAILED: RecordingWorkerEvent.FAILED,
}
def to_recording_event(egress_status):
"""Translate a LiveKit egress status into a recording worker event."""
event = EGRESS_STATUS_TO_RECORDING_EVENT.get(egress_status)
if event is None:
logger.warning(
"Unmapped LiveKit egress status '%s', ignoring the event.",
egress_status,
)
return event
class LiveKitEventsService:
"""Service for processing and handling LiveKit webhook events and notifications."""
@@ -194,20 +173,12 @@ class LiveKitEventsService:
f"Recording with worker ID {egress_id} does not exist"
) from err
event = to_recording_event(data.egress_info.status)
if event is None:
return
self.recording_events.handle_update(recording, event)
egress_status = data.egress_info.status
self.recording_events.handle_update(recording, egress_status)
def _handle_egress_ended(self, data):
"""Handle 'egress_ended' event.
"""Handle 'egress_ended' event."""
Egress ended is sent with one of these statuses:
EGRESS_COMPLETE, EGRESS_FAILED, EGRESS_ABORTED, EGRESS_LIMIT_REACHED
"""
# Fetch recording
try:
recording = models.Recording.objects.select_related("room").get(
worker_id=data.egress_info.egress_id
@@ -217,17 +188,6 @@ class LiveKitEventsService:
f"Recording with worker ID {data.egress_info.egress_id} does not exist"
) from err
event = to_recording_event(data.egress_info.status)
# Log if/why the recording failed
self.recording_events.log_worker_error(
recording,
event,
error=data.egress_info.error,
error_code=data.egress_info.error_code,
)
# Update room
try:
room_name = str(recording.room.id)
RoomManagement.update_metadata(
@@ -241,17 +201,38 @@ class LiveKitEventsService:
except RoomManagementException as e:
logger.exception("Failed to update room's metadata: %s", e)
# Stop metadata collector
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 event is None:
return
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
self.recording_events.handle_terminal_event(recording, event)
# Finalize the recording, the egress has uploaded the file to the storage
if 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
@staticmethod
def _is_connection_test_room(room_name: str) -> bool:
@@ -275,9 +256,7 @@ class LiveKitEventsService:
except models.Room.DoesNotExist as err:
raise ActionFailedError(f"Room with ID {room_id} does not exist") from err
if (
settings.ROOM_TELEPHONY_ENABLED or settings.ROOMKIT_ENABLED
) and not room.is_encrypted:
if settings.ROOM_TELEPHONY_ENABLED or settings.ROOMKIT_ENABLED:
try:
self.sip_management.ensure_dispatch_rule(room)
except SIPException as e:
+5 -22
View File
@@ -48,11 +48,8 @@ class LobbyParticipant:
color: str
id: str
entered_at: str
# Whether the user signed in (e.g. via ProConnect). Surfaced to admins so
# they can decide whether to accept self-declared identities.
is_authenticated: bool = False
def to_dict(self) -> Dict[str, object]:
def to_dict(self) -> Dict[str, str]:
"""Serialize the participant object to a dict representation."""
return {
"status": self.status.value,
@@ -60,7 +57,6 @@ class LobbyParticipant:
"id": self.id,
"color": self.color,
"entered_at": self.entered_at,
"is_authenticated": self.is_authenticated,
}
@classmethod
@@ -76,7 +72,6 @@ class LobbyParticipant:
id=data["id"],
color=data["color"],
entered_at=data["entered_at"],
is_authenticated=bool(data.get("is_authenticated", False)),
)
except (KeyError, ValueError) as e:
logger.exception("Error creating Participant from dict:")
@@ -149,7 +144,7 @@ class LobbyService:
key=settings.LOBBY_COOKIE_NAME,
value=participant_id,
httponly=True,
secure=not settings.DEBUG,
secure=True,
samesite="Lax",
)
@@ -213,7 +208,6 @@ class LobbyService:
id=participant_id,
color=utils.generate_color(participant_id),
entered_at=timezone.now().isoformat(),
is_authenticated=request.user.is_authenticated,
)
else:
participant.status = LobbyParticipantStatus.ACCEPTED
@@ -226,24 +220,19 @@ class LobbyService:
configuration=room.configuration,
participant_id=participant_id,
role=user_role,
encryption_mode=room.encryption_mode,
)
return participant, livekit_config
livekit_config = None
if participant is None:
participant = self.enter(
room.id,
participant_id,
username,
is_authenticated=request.user.is_authenticated,
)
participant = self.enter(room.id, participant_id, username)
elif participant.status == LobbyParticipantStatus.WAITING:
self.refresh_waiting_status(room.id, participant_id)
elif participant.status == LobbyParticipantStatus.ACCEPTED:
# wrongly named, contains access token to join a room
livekit_config = utils.generate_livekit_config(
room_id=room_id,
user=request.user,
@@ -252,7 +241,6 @@ class LobbyService:
configuration=room.configuration,
participant_id=participant_id,
role=user_role,
encryption_mode=room.encryption_mode,
)
return participant, livekit_config
@@ -270,11 +258,7 @@ class LobbyService:
self._index_touch(room_id)
def enter(
self,
room_id: UUID,
participant_id: str,
username: str,
is_authenticated: bool = False,
self, room_id: UUID, participant_id: str, username: str
) -> LobbyParticipant:
"""Add participant to waiting lobby."""
@@ -286,7 +270,6 @@ class LobbyService:
id=participant_id,
color=color,
entered_at=timezone.now().isoformat(),
is_authenticated=is_authenticated,
)
try:
@@ -76,6 +76,10 @@ class RoomManagement:
except TwirpError as e:
if e.code == "not_found":
logger.warning(
"Room %s not found in LiveKit, skipping metadata update",
room_name,
)
raise RoomNotFoundException("Room does not exist") from e
logger.exception(
@@ -2,23 +2,18 @@
Test RecordingEventsService service.
"""
# pylint: disable=redefined-outer-name,protected-access
# pylint: disable=redefined-outer-name
import logging
from unittest import mock
import pytest
from core.factories import RecordingFactory
from core.recording.enums import RecordingWorkerEvent
from core.recording.services.recording_events import (
RecordingEventsError,
RecordingEventsService,
RecordingNotSavableError,
)
from core.services.room_management import (
RoomManagementException,
)
from core.utils import NotificationError
pytestmark = pytest.mark.django_db
@@ -39,10 +34,10 @@ def service():
)
@mock.patch("core.utils.notify_participants")
def test_handle_limit_reached_success(mock_notify, mode, notification_type, service):
"""Test _handle_limit_reached stops recording and notifies participants."""
"""Test handle_limit_reached stops recording and notifies participants."""
recording = RecordingFactory(status="active", mode=mode)
service._handle_limit_reached(recording)
service.handle_limit_reached(recording)
assert recording.status == "stopped"
mock_notify.assert_called_once_with(
@@ -53,69 +48,13 @@ def test_handle_limit_reached_success(mock_notify, mode, notification_type, serv
@pytest.mark.parametrize(
("mode", "notification_type"),
(
("screen_recording", "screenRecordingFailed"),
("transcript", "transcriptionFailed"),
("screen_recording", "screenRecordingLimitReached"),
("transcript", "transcriptionLimitReached"),
),
)
@mock.patch("core.utils.notify_participants")
def test_handle_failed_success(mock_notify, mode, notification_type, service):
"""Test _handle_failed marks recording as failed and notifies participants."""
recording = RecordingFactory(status="active", mode=mode)
service._handle_failed(recording)
assert recording.status == "failed"
mock_notify.assert_called_once_with(
room_name=str(recording.room.id), notification_data={"type": notification_type}
)
@pytest.mark.parametrize(
("mode", "notification_type"),
(
("screen_recording", "screenRecordingAborted"),
("transcript", "transcriptionAborted"),
),
)
@mock.patch("core.utils.notify_participants")
def test_handle_aborted_success(mock_notify, mode, notification_type, service):
"""Test _handle_aborted marks recording as aborted and notifies participants."""
recording = RecordingFactory(status="active", mode=mode)
service._handle_aborted(recording)
assert recording.status == "aborted"
mock_notify.assert_called_once_with(
room_name=str(recording.room.id), notification_data={"type": notification_type}
)
@pytest.mark.parametrize(
("mode", "notification_prefix"),
(("screen_recording", "screenRecording"), ("transcript", "transcription")),
)
@pytest.mark.parametrize(
("handler", "expected_status", "event", "notification_suffix"),
(
("_handle_limit_reached", "stopped", "limit reached", "LimitReached"),
("_handle_failed", "failed", "failed", "Failed"),
("_handle_aborted", "aborted", "aborted", "Aborted"),
),
)
@mock.patch("core.utils.notify_participants")
def test_handle_event_notification_error( # noqa: PLR0913, PLR0917
mock_notify,
handler,
expected_status,
event,
notification_suffix,
mode,
notification_prefix,
service,
): # pylint: disable=too-many-arguments,too-many-positional-arguments
"""Test handlers raise RecordingEventsError when notifying participants fails,
while still applying the recording status of their event.
"""
def test_handle_limit_reached_error(mock_notify, mode, notification_type, service):
"""Test handle_limit_reached raises RecordingEventsError when notification fails."""
mock_notify.side_effect = NotificationError("Error notifying")
@@ -123,15 +62,14 @@ def test_handle_event_notification_error( # noqa: PLR0913, PLR0917
with pytest.raises(
RecordingEventsError,
match=rf"Failed to notify participants in room '.+' "
rf"about recording {event} \(recording_id=.+\)",
match=r"Failed to notify participants in room '.+' "
r"about recording limit reached \(recording_id=.+\)",
):
getattr(service, handler)(recording)
service.handle_limit_reached(recording)
assert recording.status == expected_status
assert recording.status == "stopped"
mock_notify.assert_called_once_with(
room_name=str(recording.room.id),
notification_data={"type": f"{notification_prefix}{notification_suffix}"},
room_name=str(recording.room.id), notification_data={"type": notification_type}
)
@@ -144,19 +82,19 @@ def test_handle_event_notification_error( # noqa: PLR0913, PLR0917
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
def test_handle_successful_saves_recording( # pylint: disable=too-many-arguments, too-many-positional-arguments
def test_handle_complete_saves_recording( # pylint: disable=too-many-arguments, too-many-positional-arguments
mock_notify_external_services,
notify_return_value,
expected_status,
status,
service,
):
"""Test _handle_successful notifies external services and saves a savable recording."""
"""Test handle_complete notifies external services and saves a savable recording."""
mock_notify_external_services.return_value = notify_return_value
recording = RecordingFactory(status=status)
service._handle_successful(recording)
service.handle_complete(recording)
mock_notify_external_services.assert_called_once_with(recording)
@@ -166,293 +104,23 @@ def test_handle_successful_saves_recording( # pylint: disable=too-many-argument
@pytest.mark.parametrize(
"status",
[
"initiated",
"saved",
"notification_succeeded",
"aborted",
"failed",
"failed_to_start",
],
["initiated", "saved", "notification_succeeded", "aborted", "failed_to_start"],
)
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
def test_handle_successful_non_savable_recording(
def test_handle_complete_non_savable_recording(
mock_notify_external_services, status, service
):
"""Test _handle_successful refuses recordings that are already saved or in error."""
"""Test handle_complete refuses recordings that are already saved or in error."""
recording = RecordingFactory(status=status)
with pytest.raises(RecordingNotSavableError):
service._handle_successful(recording)
service.handle_complete(recording)
mock_notify_external_services.assert_not_called()
recording.refresh_from_db()
assert recording.status == status
@pytest.mark.parametrize(
("event", "recording_status"),
(
(RecordingWorkerEvent.STARTED, "started"),
(RecordingWorkerEvent.SAVING, "saving"),
),
)
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_update_syncs_room_metadata(
mock_update_metadata, event, recording_status, service
):
"""Test handle_update updates the room's metadata."""
recording = RecordingFactory(status="active")
service.handle_update(recording, event)
mock_update_metadata.assert_called_once_with(
str(recording.room.id), {"recording_status": recording_status}
)
@pytest.mark.parametrize(
"event",
(
RecordingWorkerEvent.STARTING,
RecordingWorkerEvent.COMPLETED,
RecordingWorkerEvent.LIMIT_REACHED,
RecordingWorkerEvent.ABORTED,
RecordingWorkerEvent.FAILED,
),
)
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_update_ignores_events_without_a_metadata_status(
mock_update_metadata, event, service
):
"""Test handle_update doesn't update metadata for events it doesn't match."""
recording = RecordingFactory(status="active")
service.handle_update(recording, event)
mock_update_metadata.assert_not_called()
@pytest.mark.parametrize(
("event", "initial_status", "expected_status", "notification_type"),
(
(RecordingWorkerEvent.LIMIT_REACHED, "active", "saved", "LimitReached"),
(RecordingWorkerEvent.LIMIT_REACHED, "stopped", "saved", None),
(RecordingWorkerEvent.LIMIT_REACHED, "saved", "saved", None),
(RecordingWorkerEvent.ABORTED, "active", "aborted", "Aborted"),
(RecordingWorkerEvent.ABORTED, "failed_to_stop", "failed_to_stop", None),
(RecordingWorkerEvent.FAILED, "active", "failed", "Failed"),
(RecordingWorkerEvent.FAILED, "stopped", "failed", "Failed"),
(RecordingWorkerEvent.FAILED, "aborted", "aborted", None),
(RecordingWorkerEvent.COMPLETED, "active", "saved", None),
(RecordingWorkerEvent.COMPLETED, "saved", "saved", None),
),
)
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
def test_handle_terminal_event_dispatches_on_event_and_status( # noqa: PLR0913, PLR0917
mock_notify,
mock_notify_external_services,
event,
initial_status,
expected_status,
notification_type,
service,
): # pylint: disable=too-many-arguments,too-many-positional-arguments
"""Test handle_terminal_event chooses the right handler from the event and status."""
mock_notify_external_services.return_value = False
recording = RecordingFactory(status=initial_status, mode="screen_recording")
service.handle_terminal_event(recording, event)
recording.refresh_from_db()
assert recording.status == expected_status
if notification_type is None:
mock_notify.assert_not_called()
else:
mock_notify.assert_called_once_with(
room_name=str(recording.room.id),
notification_data={"type": f"screenRecording{notification_type}"},
)
@pytest.mark.parametrize(
"event",
(
RecordingWorkerEvent.STARTING,
RecordingWorkerEvent.STARTED,
RecordingWorkerEvent.SAVING,
),
)
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
def test_handle_terminal_event_ignores_non_terminal_events(
mock_notify, mock_notify_external_services, event, service, caplog
):
"""Test handle_terminal_event refuses non-terminal events."""
recording = RecordingFactory(status="active")
with caplog.at_level(logging.WARNING):
service.handle_terminal_event(recording, event)
assert f"Ignoring non-terminal event {event.value}" in caplog.text
mock_notify.assert_not_called()
mock_notify_external_services.assert_not_called()
recording.refresh_from_db()
assert recording.status == "active"
@pytest.mark.parametrize(
("event", "expected_status"),
(
(RecordingWorkerEvent.LIMIT_REACHED, "saved"),
(RecordingWorkerEvent.ABORTED, "aborted"),
(RecordingWorkerEvent.FAILED, "failed"),
),
)
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
def test_handle_terminal_event_survives_a_notification_failure( # noqa: PLR0913, PLR0917
mock_notify,
mock_notify_external_services,
event,
expected_status,
service,
caplog,
): # pylint: disable=too-many-arguments,too-many-positional-arguments
"""Test handle_terminal_event logs a notification failure instead of raising.
The recording status must still be persisted: participants missing their
notification should not disturb recording.
"""
mock_notify_external_services.return_value = False
mock_notify.side_effect = NotificationError("Error notifying")
recording = RecordingFactory(status="active")
with caplog.at_level(logging.ERROR):
service.handle_terminal_event(recording, event)
assert f"Failed to notify participants that recording {recording.id}" in caplog.text
recording.refresh_from_db()
assert recording.status == expected_status
@pytest.mark.parametrize(
"status",
["failed_to_start", "aborted", "failed", "failed_to_stop", "saved", "initiated"],
)
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
def test_handle_terminal_event_ignores_a_non_savable_recording(
mock_notify_external_services, status, service, caplog
):
"""Test handle_terminal_event handles a redelivered event idempotently.
A terminal event may be redelivered for an already finalized recording;
this must not raise, otherwise the webhook would 500 and be retried.
"""
recording = RecordingFactory(status=status)
with caplog.at_level(logging.WARNING):
service.handle_terminal_event(recording, RecordingWorkerEvent.COMPLETED)
assert f"Recording {recording.id} is not savable" in caplog.text
mock_notify_external_services.assert_not_called()
recording.refresh_from_db()
assert recording.status == status
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_update_survives_a_metadata_failure(
mock_update_metadata, service, caplog
):
"""Test handle_update logs a metadata failure instead of raising."""
mock_update_metadata.side_effect = RoomManagementException("Error updating")
recording = RecordingFactory(status="active")
with caplog.at_level(logging.ERROR):
service.handle_update(recording, RecordingWorkerEvent.SAVING)
assert "Failed to update room's metadata" in caplog.text
@pytest.mark.parametrize(
("event", "expected_level"),
(
(RecordingWorkerEvent.ABORTED, logging.INFO),
(RecordingWorkerEvent.FAILED, logging.ERROR),
),
)
def test_log_worker_error_reports_an_unsuccessful_event(
event, expected_level, service, caplog
):
"""Test log_worker_error records the reason the recording did not succeed."""
recording = RecordingFactory(status="active", mode="screen_recording")
with caplog.at_level(logging.INFO):
service.log_worker_error(
recording, event, error="could not connect to the room", error_code=500
)
assert (
f"Recording worker reported {event.value} for recording {recording.id}"
in caplog.text
)
assert "could not connect to the room" in caplog.text
assert "error_code=500" in caplog.text
worker_logs = [
record
for record in caplog.records
if record.name == "core.recording.services.recording_events"
]
assert [record.levelno for record in worker_logs] == [expected_level]
@pytest.mark.parametrize(
"event",
(
RecordingWorkerEvent.STARTING,
RecordingWorkerEvent.STARTED,
RecordingWorkerEvent.SAVING,
RecordingWorkerEvent.COMPLETED,
RecordingWorkerEvent.LIMIT_REACHED,
None,
),
)
def test_log_worker_error_stays_quiet_on_anything_else(event, service, caplog):
"""Test log_worker_error ignores events other than FAILED and ABORTED."""
recording = RecordingFactory(status="active")
with caplog.at_level(logging.INFO):
service.log_worker_error(recording, event, error="some error", error_code=500)
assert "Recording worker reported" not in caplog.text
@@ -224,7 +224,6 @@ def test_api_recording_retrieve_expired(settings):
RecordingStatusChoices.INITIATED,
RecordingStatusChoices.ACTIVE,
RecordingStatusChoices.SAVED,
RecordingStatusChoices.FAILED,
RecordingStatusChoices.FAILED_TO_START,
RecordingStatusChoices.FAILED_TO_STOP,
RecordingStatusChoices.ABORTED,
@@ -4,7 +4,6 @@ Test worker service classes.
# pylint: disable=protected-access,redefined-outer-name,unused-argument,no-member
import logging
from unittest.mock import AsyncMock, Mock, patch
import pytest
@@ -155,9 +154,9 @@ def test_base_egress_filepath_construction(service, filename, extension, expecte
"response_status,expected_result",
[
(livekit_api.EgressStatus.EGRESS_ABORTED, "ABORTED"),
(livekit_api.EgressStatus.EGRESS_FAILED, "FAILED"),
(livekit_api.EgressStatus.EGRESS_COMPLETE, "FAILED_TO_STOP"),
(livekit_api.EgressStatus.EGRESS_ENDING, "STOPPED"),
(livekit_api.EgressStatus.EGRESS_FAILED, "FAILED_TO_STOP"),
],
)
def test_base_egress_stop_with_status(service, response_status, expected_result):
@@ -176,32 +175,6 @@ def test_base_egress_stop_with_status(service, response_status, expected_result)
assert result == expected_result
@pytest.mark.parametrize(
"response_status,event",
[
(livekit_api.EgressStatus.EGRESS_ABORTED, "aborted"),
(livekit_api.EgressStatus.EGRESS_FAILED, "failed"),
(livekit_api.EgressStatus.EGRESS_COMPLETE, "failed to stop"),
],
)
def test_base_egress_stop_logs_livekit_error(service, response_status, event, caplog):
"""Should log the reason LiveKit reported for an unsuccessful stop."""
mock_response = Mock(
status=response_status,
egress_id="test_worker_id",
error="could not connect to the room",
error_code=500,
)
service._handle_request = Mock(return_value=mock_response)
with caplog.at_level(logging.ERROR):
service.stop("test_worker_id")
assert f"Egress {event} on stop (egress_id=test_worker_id" in caplog.text
assert "could not connect to the room" in caplog.text
assert "error_code=500" in caplog.text
def test_base_egress_stop_missing_status(service):
"""Test stop method when response is missing status"""
# Mock _handle_request with missing status
@@ -312,34 +312,3 @@ def test_api_rooms_create_authenticated_blank_user_default_access_level():
assert response.status_code == 201
room = Room.objects.get()
assert room.access_level == settings.RESOURCE_DEFAULT_ACCESS_LEVEL
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
@pytest.mark.parametrize(
"user_default", [None, RoomAccessLevel.PUBLIC, RoomAccessLevel.TRUSTED]
)
@pytest.mark.parametrize(
"requested_access", [None, RoomAccessLevel.PUBLIC, RoomAccessLevel.TRUSTED]
)
def test_api_rooms_create_encryption_access_precedence(
settings, encryption_mode, user_default, requested_access
):
"""Encryption overrides request and user access defaults only for encrypted rooms."""
settings.ENCRYPTION_ENABLED = True
user = UserFactory(default_room_access_level=user_default)
client = APIClient()
client.force_login(user)
data = {"name": "New room", "encryption_mode": encryption_mode}
if requested_access is not None:
data["access_level"] = requested_access
response = client.post("/api/v1.0/rooms/", data)
assert response.status_code == 201
expected_access = (
RoomAccessLevel.RESTRICTED
if encryption_mode == "basic"
else requested_access or user_default or settings.RESOURCE_DEFAULT_ACCESS_LEVEL
)
assert response.json()["access_level"] == expected_access
assert Room.objects.get().access_level == expected_access
@@ -62,7 +62,6 @@ def test_request_entry_anonymous(settings):
"status": "waiting",
"color": "mocked-color",
"entered_at": "2025-01-01T10:00:00+00:00",
"is_authenticated": False,
"livekit": None,
}
@@ -76,12 +75,9 @@ def test_request_entry_anonymous(settings):
@freeze_time("2025-01-01 10:00:00")
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
def test_request_entry_authenticated_user(settings, encryption_mode):
def test_request_entry_authenticated_user(settings):
"""Authenticated users should be allowed to request entry."""
room = RoomFactory(
access_level=RoomAccessLevel.RESTRICTED, encryption_mode=encryption_mode
)
room = RoomFactory(access_level=RoomAccessLevel.RESTRICTED)
user = UserFactory()
client = APIClient()
client.force_login(user)
@@ -117,7 +113,6 @@ def test_request_entry_authenticated_user(settings, encryption_mode):
"status": "waiting",
"color": "mocked-color",
"entered_at": "2025-01-01T10:00:00+00:00",
"is_authenticated": True,
"livekit": None,
}
@@ -194,7 +189,6 @@ def test_request_entry_with_existing_participants(settings):
"entered_at": "2025-01-01T10:00:00+00:00",
"status": "waiting",
"color": "mocked-color",
"is_authenticated": False,
"livekit": None,
}
@@ -249,7 +243,6 @@ def test_request_entry_public_room(settings):
"entered_at": "2025-01-01T10:00:00+00:00",
"status": "accepted",
"color": "mocked-color",
"is_authenticated": False,
"livekit": {"token": "test-token"},
}
@@ -304,7 +297,6 @@ def test_request_entry_authenticated_user_public_room(settings):
"entered_at": "2025-01-01T10:00:00+00:00",
"status": "accepted",
"color": "mocked-color",
"is_authenticated": True,
"livekit": {"token": "test-token"},
}
@@ -362,7 +354,6 @@ def test_request_entry_waiting_participant_public_room(settings):
"status": "accepted",
"color": "#123456",
"entered_at": "2025-01-01T10:00:00+00:00",
"is_authenticated": False,
"livekit": {"token": "test-token"},
}
@@ -632,7 +623,6 @@ def test_list_waiting_participants_success(settings):
"username": "user2",
"status": "waiting",
"color": "#654321",
"is_authenticated": False,
"entered_at": "2025-01-01T10:05:00+00:00",
},
{
@@ -640,7 +630,6 @@ def test_list_waiting_participants_success(settings):
"username": "user1",
"status": "waiting",
"color": "#123456",
"is_authenticated": False,
"entered_at": "2025-01-01T10:00:00+00:00",
},
]
@@ -33,7 +33,6 @@ def test_api_rooms_retrieve_anonymous_private_pk():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -53,7 +52,6 @@ def test_api_rooms_retrieve_anonymous_trusted_pk():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -72,7 +70,6 @@ def test_api_rooms_retrieve_anonymous_private_pk_no_dashes():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -89,7 +86,6 @@ def test_api_rooms_retrieve_anonymous_private_slug():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -106,7 +102,6 @@ def test_api_rooms_retrieve_anonymous_private_slug_not_normalized():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -222,7 +217,6 @@ def test_api_rooms_retrieve_anonymous_public(mock_token):
"name": room.name,
"pin_code": room.pin_code,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
mock_token.assert_called_once()
@@ -269,7 +263,6 @@ def test_api_rooms_retrieve_authenticated_public(mock_token):
"name": room.name,
"pin_code": room.pin_code,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
mock_token.assert_called_once_with(
@@ -280,7 +273,6 @@ def test_api_rooms_retrieve_authenticated_public(mock_token):
sources=["camera"],
role=None,
participant_id=None,
encryption_mode="none",
)
@@ -322,7 +314,6 @@ def test_api_rooms_retrieve_authenticated_trusted(mock_token):
"name": room.name,
"pin_code": room.pin_code,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
mock_token.assert_called_once_with(
@@ -333,7 +324,6 @@ def test_api_rooms_retrieve_authenticated_trusted(mock_token):
sources=None,
role=None,
participant_id=None,
encryption_mode="none",
)
@@ -359,7 +349,6 @@ def test_api_rooms_retrieve_authenticated():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -411,7 +400,6 @@ def test_api_rooms_retrieve_members(mock_token, django_assert_num_queries, setti
"name": room.name,
"pin_code": room.pin_code,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
mock_token.assert_called_once_with(
@@ -422,7 +410,6 @@ def test_api_rooms_retrieve_members(mock_token, django_assert_num_queries, setti
sources=["camera"],
role=str(RoleChoices.MEMBER),
participant_id=None,
encryption_mode="none",
)
@@ -474,7 +461,6 @@ def test_api_rooms_retrieve_administrators(
"short_name": other_user_access.user.short_name,
"timezone": "UTC",
"language": other_user_access.user.language,
"default_encryption_mode": "none",
},
"resource": str(room.id),
"role": other_user_access.role,
@@ -490,7 +476,6 @@ def test_api_rooms_retrieve_administrators(
"short_name": user_access.user.short_name,
"timezone": "UTC",
"language": user_access.user.language,
"default_encryption_mode": "none",
},
"resource": str(room.id),
"role": user_access.role,
@@ -511,7 +496,6 @@ def test_api_rooms_retrieve_administrators(
"name": room.name,
"pin_code": room.pin_code,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
mock_token.assert_called_once_with(
@@ -522,25 +506,4 @@ def test_api_rooms_retrieve_administrators(
sources=None,
role=str(user_access.role),
participant_id=None,
encryption_mode="none",
)
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
@mock.patch("core.utils.generate_token", return_value="test-token")
def test_api_rooms_retrieve_custom_username(mock_token, encryption_mode, settings):
"""Encryption does not override the participant's requested display name."""
settings.AUTHENTICATED_PARTICIPANTS_CAN_EDIT_DISPLAY_NAME = True
user = UserFactory(full_name="Profile Name")
room = RoomFactory(
access_level=RoomAccessLevel.RESTRICTED, encryption_mode=encryption_mode
)
UserResourceAccessFactory(resource=room, user=user, role="owner")
client = APIClient()
client.force_login(user)
response = client.get(f"/api/v1.0/rooms/{room.id}/", {"username": "Custom Name"})
assert response.status_code == 200
assert mock_token.call_args.kwargs["username"] == "Custom Name"
assert mock_token.call_args.kwargs["encryption_mode"] == encryption_mode
@@ -410,24 +410,3 @@ def test_api_rooms_update_livekit_sync_failure(mock_update_metadata, exception):
"configuration": {"can_publish_sources": ["camera"]},
},
)
@pytest.mark.parametrize(
"access_level", [RoomAccessLevel.PUBLIC, RoomAccessLevel.TRUSTED]
)
def test_api_rooms_update_encrypted_access_rejected(access_level):
"""API updates cannot change an encrypted room away from restricted access."""
room = RoomFactory(encryption_mode="basic")
user = UserFactory()
room.accesses.create(user=user, role="owner")
client = APIClient()
client.force_login(user)
response = client.patch(
f"/api/v1.0/rooms/{room.id}/", {"access_level": access_level}
)
assert response.status_code == 400
assert "access_level" in response.json()
room.refresh_from_db()
assert room.access_level == RoomAccessLevel.RESTRICTED
@@ -3,7 +3,6 @@ Test LiveKitEvents service.
"""
# pylint: disable=W0621,W0613, W0212, E0611
import logging
import uuid
from unittest import mock
@@ -11,16 +10,13 @@ import pytest
from livekit.api import EgressStatus
from core.factories import RecordingFactory, RoomFactory
from core.recording.enums import RecordingWorkerEvent
from core.recording.services.recording_events import RecordingEventsService
from core.services.livekit_events import (
EGRESS_STATUS_TO_RECORDING_EVENT,
ActionFailedError,
AuthenticationError,
InvalidPayloadError,
LiveKitEventsService,
api,
to_recording_event,
)
from core.services.lobby import LobbyService
from core.services.room_management import RoomManagementException
@@ -81,7 +77,7 @@ def test_initialization(
def test_handle_egress_ended_success( # pylint: disable=too-many-arguments, too-many-positional-arguments
mock_update_metadata, mock_notify, mode, notification_type, service
):
"""Should successfully stop recording and notify all participant."""
"""Should successfully stop recording and notifies all participant."""
recording = RecordingFactory(worker_id="worker-1", mode=mode, status="active")
mock_data = mock.MagicMock()
@@ -108,6 +104,7 @@ def test_handle_egress_ended_success( # pylint: disable=too-many-arguments, too
(
(EgressStatus.EGRESS_ACTIVE, "started"),
(EgressStatus.EGRESS_ENDING, "saving"),
(EgressStatus.EGRESS_ABORTED, "aborted"),
),
)
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
@@ -128,6 +125,29 @@ def test_handle_egress_updated_success(
)
@pytest.mark.parametrize(
"egress_status",
(
EgressStatus.EGRESS_FAILED,
EgressStatus.EGRESS_LIMIT_REACHED,
),
)
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_updated_non_handled(
mock_update_metadata, egress_status, service
):
"""Should ignore certain egress status and don't trigger metadata updates."""
recording = RecordingFactory(worker_id="worker-1", status="initiated")
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = egress_status
service._handle_egress_updated(mock_data)
mock_update_metadata.assert_not_called()
@pytest.mark.parametrize(
("mode", "notification_type"),
(
@@ -160,38 +180,33 @@ def test_handle_egress_ended_metadata_update_fails( # pylint: disable=too-many-
assert recording.status == "saved"
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_ended_notification_fails(
mock_update_metadata, mock_notify, mock_notify_external_services, service
mock_update_metadata, mock_notify, service
):
"""Should still stop and save the recording when notifying participants fails."""
mock_notify_external_services.return_value = False
mock_notify.side_effect = NotificationError("Error notifying")
"""Should raise ActionFailedError when notification fails but still stop recording."""
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 = EgressStatus.EGRESS_LIMIT_REACHED
service._handle_egress_ended(mock_data)
mock_notify.side_effect = NotificationError("Error notifying")
with pytest.raises(
ActionFailedError,
match=r"Failed to process limit reached event for recording .+",
):
service._handle_egress_ended(mock_data)
recording.refresh_from_db()
assert recording.status == "stopped"
mock_notify.assert_called_once_with(
room_name=str(recording.room.id),
notification_data={"type": "screenRecordingLimitReached"},
)
mock_update_metadata.assert_called_once_with(
str(recording.room.id), remove_keys=["recording_mode", "recording_status"]
)
recording.refresh_from_db()
assert recording.status == "saved"
@mock.patch("core.utils.notify_participants")
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
@@ -217,25 +232,17 @@ def test_handle_egress_ended_recording_not_found(
assert recording.status == "active"
@pytest.mark.parametrize(
"egress_status",
(
EgressStatus.EGRESS_FAILED,
EgressStatus.EGRESS_ABORTED,
EgressStatus.EGRESS_LIMIT_REACHED,
),
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_ended_recording_should_not_be_saved(
mock_update_metadata, mock_notify, egress_status, service
def test_handle_egress_ended_recording_not_active(
mock_update_metadata, mock_notify, service
):
"""Don't update status for recordings that must not be saved."""
"""Should ignore non-active recordings."""
recording = RecordingFactory(worker_id="worker-1", status="failed_to_stop")
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = "worker-1"
mock_data.egress_info.status = egress_status
mock_data.egress_info.status = EgressStatus.EGRESS_LIMIT_REACHED
service._handle_egress_ended(mock_data)
@@ -248,24 +255,12 @@ def test_handle_egress_ended_recording_should_not_be_saved(
assert recording.status == "failed_to_stop"
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_ended_complete_does_not_notify_participants(
mock_update_metadata, mock_notify, mock_notify_external_services, service
def test_handle_egress_ended_recording_not_limit_reached(
mock_update_metadata, mock_notify, service
):
"""Shouldn't notify participants on a successful egress.
EGRESS_COMPLETE is the only ended status that notifies no one: limit
reached, aborted and failed egresses each send their own notification.
A stopped recording is simply finalized, which is the nominal flow once
the egress uploaded the file of a user-initiated stop.
"""
mock_notify_external_services.return_value = False
"""Should ignore egress non-limit-reached statuses."""
recording = RecordingFactory(worker_id="worker-1", status="stopped")
mock_data = mock.MagicMock()
@@ -278,10 +273,7 @@ def test_handle_egress_ended_complete_does_not_notify_participants(
mock_update_metadata.assert_called_once_with(
str(recording.room.id), remove_keys=["recording_mode", "recording_status"]
)
mock_notify_external_services.assert_called_once_with(recording)
recording.refresh_from_db()
assert recording.status == "saved"
assert recording.status == "stopped"
@mock.patch("core.services.livekit_events.MetadataCollectorService")
@@ -341,163 +333,96 @@ def test_handle_egress_ended_does_not_call_metadata_collector_stop_when_conditio
mock_collector.stop.assert_not_called()
@pytest.mark.parametrize(
("egress_status", "recording_status", "event", "expected_level"),
(
(EgressStatus.EGRESS_ABORTED, "active", "aborted", logging.INFO),
(EgressStatus.EGRESS_FAILED, "active", "failed", logging.ERROR),
# The synchronous stop may already have persisted the terminal status,
# and the error details exist only in the webhook payload.
(EgressStatus.EGRESS_ABORTED, "aborted", "aborted", logging.INFO),
(EgressStatus.EGRESS_FAILED, "failed", "failed", logging.ERROR),
),
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_ended_logs_livekit_error( # noqa: PLR0913, PLR0917
mock_update_metadata,
mock_notify,
egress_status,
recording_status,
event,
expected_level,
service,
caplog,
): # pylint: disable=too-many-arguments,too-many-positional-arguments
"""Should log the reason LiveKit reported an unsuccessful egress."""
recording = RecordingFactory(worker_id="worker-1", status=recording_status)
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = egress_status
mock_data.egress_info.error = "could not connect to the room"
mock_data.egress_info.error_code = 500
with caplog.at_level(logging.INFO):
service._handle_egress_ended(mock_data)
assert (
f"Recording worker reported {event} for recording {recording.id}" in caplog.text
)
assert "could not connect to the room" in caplog.text
assert "error_code=500" in caplog.text
worker_logs = [
record
for record in caplog.records
if record.name == "core.recording.services.recording_events"
]
assert [record.levelno for record in worker_logs] == [expected_level]
@pytest.mark.parametrize(
"egress_status",
(EgressStatus.EGRESS_COMPLETE, EgressStatus.EGRESS_LIMIT_REACHED),
)
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_ended_does_not_log_error_on_successful_egress( # noqa: PLR0913, PLR0917
@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, PLR0917
mock_update_metadata,
mock_notify,
mock_notify_external_services,
notify_return_value,
recording_status,
egress_status,
service,
caplog,
): # pylint: disable=too-many-arguments,too-many-positional-arguments
"""Shouldn't log an egress error when LiveKit reports a successful egress."""
"""Should notify external services and save the recording on egress completion
(EGRESS_COMPLETE or EGRESS_LIMIT_REACHED).
"""
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
with caplog.at_level(logging.ERROR):
service._handle_egress_ended(mock_data)
service._handle_egress_ended(mock_data)
assert "Recording worker reported" not in caplog.text
mock_notify_external_services.assert_called_once_with(recording)
recording.refresh_from_db()
assert recording.status == recording_status
@mock.patch("core.utils.notify_participants")
@pytest.mark.parametrize(
"egress_status",
[
EgressStatus.EGRESS_STARTING,
EgressStatus.EGRESS_ACTIVE,
EgressStatus.EGRESS_ENDING,
EgressStatus.EGRESS_FAILED,
EgressStatus.EGRESS_ABORTED,
],
)
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_ended_logs_livekit_error_before_cleaning_up(
mock_update_metadata, mock_notify, service, caplog
def test_handle_egress_ended_does_not_save_on_wrong_status(
mock_update_metadata, egress_status, service
):
"""Should log the failure reason even when the cleanup fails afterwards."""
mock_update_metadata.side_effect = RuntimeError("LiveKit is unreachable")
"""Shouldn't save on invalid status."""
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 = EgressStatus.EGRESS_FAILED
mock_data.egress_info.error = "could not connect to the room"
mock_data.egress_info.error_code = 500
mock_data.egress_info.status = egress_status
with caplog.at_level(logging.ERROR), pytest.raises(RuntimeError):
service._handle_egress_ended(mock_data)
service._handle_egress_ended(mock_data)
assert (
f"Recording worker reported failed for recording {recording.id}" in caplog.text
)
assert "could not connect to the room" in caplog.text
def test_egress_status_mapping_covers_every_livekit_status():
"""Every egress status LiveKit can report must translate to a recording event."""
unmapped = [
name
for name in EgressStatus.keys()
if getattr(EgressStatus, name) not in EGRESS_STATUS_TO_RECORDING_EVENT
]
assert not unmapped
recording.refresh_from_db()
assert recording.status == "active"
@pytest.mark.parametrize(
("egress_status", "expected_event"),
(
(EgressStatus.EGRESS_STARTING, RecordingWorkerEvent.STARTING),
(EgressStatus.EGRESS_ACTIVE, RecordingWorkerEvent.STARTED),
(EgressStatus.EGRESS_ENDING, RecordingWorkerEvent.SAVING),
(EgressStatus.EGRESS_COMPLETE, RecordingWorkerEvent.COMPLETED),
(EgressStatus.EGRESS_LIMIT_REACHED, RecordingWorkerEvent.LIMIT_REACHED),
(EgressStatus.EGRESS_ABORTED, RecordingWorkerEvent.ABORTED),
(EgressStatus.EGRESS_FAILED, RecordingWorkerEvent.FAILED),
),
"status", ["failed_to_start", "aborted", "failed_to_stop", "saved", "initiated"]
)
def test_to_recording_event_translates_egress_status(egress_status, expected_event):
"""Should translate a LiveKit egress status into a recording worker event."""
assert to_recording_event(egress_status) == expected_event
def test_to_recording_event_returns_none_on_unmapped_status(caplog):
"""Should warn and return None when LiveKit reports an unknown status."""
with caplog.at_level(logging.WARNING):
event = to_recording_event(999)
assert event is None
assert "Unmapped LiveKit egress status" in caplog.text
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_updated_ignores_unmapped_status(mock_update_metadata, service):
"""Shouldn't touch the room's metadata when the egress status is unknown."""
def test_handle_egress_ended_ignores_non_savable_recording(
mock_update_metadata, status, service
):
"""Should handle non-savable recordings idempotently without raising.
RecordingFactory(worker_id="worker-1", status="active")
'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.
"""
recording = RecordingFactory(worker_id="worker-1", status=status)
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = "worker-1"
mock_data.egress_info.status = 999
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = EgressStatus.EGRESS_COMPLETE
service._handle_egress_updated(mock_data)
service._handle_egress_ended(mock_data)
mock_update_metadata.assert_not_called()
recording.refresh_from_db()
assert recording.status == status
@mock.patch.object(LobbyService, "clear_room_cache")
+1 -11
View File
@@ -303,7 +303,6 @@ def test_request_entry_public_room(
configuration=room.configuration,
participant_id="test-participant-id",
role=None,
encryption_mode="none",
)
lobby_service._get_participant.assert_called_once_with(room.id, participant_id)
@@ -343,7 +342,6 @@ def test_request_entry_trusted_room(
configuration=room.configuration,
participant_id="test-participant-id",
role=None,
encryption_mode="none",
)
lobby_service._get_participant.assert_called_once_with(room.id, participant_id)
@@ -376,12 +374,7 @@ def test_request_entry_new_participant(
assert participant == participant_data
assert livekit_config is None
mock_enter.assert_called_once_with(
room.id,
participant_id,
username,
is_authenticated=request.user.is_authenticated,
)
mock_enter.assert_called_once_with(room.id, participant_id, username)
lobby_service._get_participant.assert_called_once_with(room.id, participant_id)
@@ -449,7 +442,6 @@ def test_request_entry_accepted_participant(
configuration=room.configuration,
participant_id="test-participant-id",
role=None,
encryption_mode="none",
)
lobby_service._get_participant.assert_called_once_with(room.id, participant_id)
@@ -491,7 +483,6 @@ def test_request_entry_participant_with_role(
configuration=room.configuration,
participant_id="test-participant-id",
role="administrator",
encryption_mode="none",
)
lobby_service._get_participant.assert_called_once_with(room.id, participant_id)
@@ -895,7 +886,6 @@ def test_update_participant_status_success(mock_cache, lobby_service, participan
"id": participant_id,
"color": "#123456",
"entered_at": "2025-01-01T10:00:00+00:00",
"is_authenticated": False,
}
mock_cache.set.assert_called_once_with(
"mocked_cache_key", expected_data, timeout=60
-1
View File
@@ -127,7 +127,6 @@ def test_api_users_retrieve_me_authenticated(settings):
"short_name": user.short_name,
"language": user.language,
"timezone": "UTC",
"default_encryption_mode": "none",
}
@@ -157,7 +157,6 @@ def test_models_recording_is_savable_normal():
@pytest.mark.parametrize(
"status",
[
RecordingStatusChoices.FAILED,
RecordingStatusChoices.FAILED_TO_STOP,
RecordingStatusChoices.FAILED_TO_START,
RecordingStatusChoices.ABORTED,
@@ -280,7 +279,6 @@ def test_models_recording_is_saved_false_initiated():
@pytest.mark.parametrize(
"status",
[
RecordingStatusChoices.FAILED,
RecordingStatusChoices.FAILED_TO_STOP,
RecordingStatusChoices.FAILED_TO_START,
RecordingStatusChoices.ABORTED,
@@ -360,35 +360,3 @@ def test_pin_generation_upper_bound(mock_randbelow, settings):
# Assert called with the right exclusive upper bound, 10^5
mock_randbelow.assert_called_with(100000)
@pytest.mark.parametrize("access_level", [None, *RoomAccessLevel.values])
@pytest.mark.parametrize("use_manager", [False, True])
def test_models_encrypted_room_creation_is_restricted(access_level, use_manager):
"""Both save and manager creation normalize encrypted rooms before validation."""
fields = {"name": "Encrypted room", "encryption_mode": "basic"}
if access_level is not None:
fields["access_level"] = access_level
if use_manager:
room = Room.objects.create(**fields)
else:
room = Room(**fields)
# A caller may assign the primary key before the first save.
room.pk = room.id
room.save()
room.refresh_from_db()
assert room.access_level == RoomAccessLevel.RESTRICTED
@pytest.mark.parametrize(
"access_level", [RoomAccessLevel.PUBLIC, RoomAccessLevel.TRUSTED]
)
def test_models_encrypted_room_access_update_rejected(access_level):
"""Existing encrypted rooms reject incompatible access instead of normalizing it."""
room = Room.objects.create(name="Encrypted room", encryption_mode="basic")
room.access_level = access_level
with pytest.raises(ValidationError) as excinfo:
room.save()
assert "access_level" in excinfo.value.message_dict
room.refresh_from_db()
assert room.access_level == RoomAccessLevel.RESTRICTED
+4 -18
View File
@@ -57,35 +57,21 @@ def test_generate_token_authenticated_fallback_user_representation():
assert claims["name"] == str(user)
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
def test_generate_token_explicit_username_overrides_default(encryption_mode):
def test_generate_token_explicit_username_overrides_default():
"""An explicitly provided username should take precedence over the full name."""
user = UserFactory(full_name="Jane Doe")
token = generate_token(
room="my-room",
user=user,
username="Custom Name",
encryption_mode=encryption_mode,
)
token = generate_token(room="my-room", user=user, username="Custom Name")
claims = decode_token(token)
assert claims["name"] == "Custom Name"
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
def test_authenticated_username_ignored_when_editing_disabled(
settings, encryption_mode
):
def test_authenticated_username_ignored_when_editing_disabled(settings):
"""With editing disabled, an authenticated user's username is ignored."""
settings.AUTHENTICATED_PARTICIPANTS_CAN_EDIT_DISPLAY_NAME = False
user = UserFactory(full_name="Jane Doe")
token = generate_token(
room="my-room",
user=user,
username="Custom Name",
encryption_mode=encryption_mode,
)
token = generate_token(room="my-room", user=user, username="Custom Name")
claims = decode_token(token)
assert claims["name"] == "Jane Doe"
+3 -15
View File
@@ -34,9 +34,6 @@ from livekit.api import ( # pylint: disable=E0611
TwirpError,
VideoGrants,
)
from livekit.protocol.room import RoomConfiguration # pylint: disable=E0611
from core.enums import EncryptionMode
logger = logging.getLogger(__name__)
@@ -72,7 +69,6 @@ def generate_token( # noqa: PLR0917
role: Optional[str] = None,
participant_id: Optional[str] = None,
ttl: Optional[timedelta] = None,
encryption_mode: str = "none",
) -> str:
"""Generate a LiveKit access token for a user in a specific room.
@@ -95,8 +91,10 @@ def generate_token( # noqa: PLR0917
"""
is_admin_or_owner = role in ("owner", "administrator")
if is_admin_or_owner:
sources = settings.LIVEKIT_DEFAULT_SOURCES
if is_admin_or_owner or sources is None:
if sources is None:
sources = settings.LIVEKIT_DEFAULT_SOURCES
video_grants = VideoGrants(
@@ -143,14 +141,6 @@ def generate_token( # noqa: PLR0917
if ttl is not None:
token = token.with_ttl(ttl)
if encryption_mode != EncryptionMode.NONE:
token = token.with_room_config(
RoomConfiguration(
name=room,
metadata=json.dumps({"encryption_mode": encryption_mode}),
)
)
return token.to_jwt()
@@ -162,7 +152,6 @@ def generate_livekit_config( # noqa: PLR0917
color: Optional[str] = None,
configuration: Optional[dict] = None,
participant_id: Optional[str] = None,
encryption_mode: str = "none",
) -> dict:
"""Generate LiveKit configuration for room access.
@@ -195,7 +184,6 @@ def generate_livekit_config( # noqa: PLR0917
sources=sources,
role=role,
participant_id=participant_id,
encryption_mode=encryption_mode,
),
}
-4
View File
@@ -978,10 +978,6 @@ class Base(Configuration):
environ_prefix=None,
)
ENCRYPTION_ENABLED = values.BooleanValue(
False, environ_name="ENCRYPTION_ENABLED", environ_prefix=None
)
# External Applications
APPLICATION_ENABLED = values.BooleanValue(
False, environ_name="APPLICATION_ENABLED", environ_prefix=None
+1 -2
View File
@@ -45,8 +45,7 @@ RUN npm run build
FROM nginxinc/nginx-unprivileged:1.30.4-alpine3.24 AS frontend-production
USER root
RUN apk upgrade --no-cache libexpat && \
apk del curl
RUN apk del curl
USER nginx
# Un-privileged user running the application
@@ -96,10 +96,6 @@ export const MainNotificationToast = () => {
case NotificationType.ScreenRecordingStopped:
case NotificationType.TranscriptionLimitReached:
case NotificationType.ScreenRecordingLimitReached:
case NotificationType.TranscriptionFailed:
case NotificationType.ScreenRecordingFailed:
case NotificationType.TranscriptionAborted:
case NotificationType.ScreenRecordingAborted:
toastQueue.add(
{
participant,
@@ -10,15 +10,11 @@ export enum NotificationType {
TranscriptionStarted = 'transcriptionStarted',
TranscriptionStopped = 'transcriptionStopped',
TranscriptionLimitReached = 'transcriptionLimitReached',
TranscriptionFailed = 'transcriptionFailed',
TranscriptionAborted = 'transcriptionAborted',
TranscriptionRequested = 'transcriptionRequested',
ScreenRecordingStarted = 'screenRecordingStarted',
ScreenRecordingStopped = 'screenRecordingStopped',
ScreenRecordingRequested = 'screenRecordingRequested',
ScreenRecordingLimitReached = 'screenRecordingLimitReached',
ScreenRecordingFailed = 'screenRecordingFailed',
ScreenRecordingAborted = 'screenRecordingAborted',
RecordingSaving = 'recordingSaving',
PermissionsRemoved = 'permissionsRemoved',
RoleChanged = 'roleChanged',
@@ -22,20 +22,12 @@ export function ToastAnyRecording({ state, ...props }: Readonly<ToastProps>) {
return 'transcript.stopped'
case NotificationType.TranscriptionLimitReached:
return 'transcript.limitReached'
case NotificationType.TranscriptionFailed:
return 'transcript.failed'
case NotificationType.TranscriptionAborted:
return 'transcript.aborted'
case NotificationType.ScreenRecordingStarted:
return 'screenRecording.started'
case NotificationType.ScreenRecordingStopped:
return 'screenRecording.stopped'
case NotificationType.ScreenRecordingLimitReached:
return 'screenRecording.limitReached'
case NotificationType.ScreenRecordingFailed:
return 'screenRecording.failed'
case NotificationType.ScreenRecordingAborted:
return 'screenRecording.aborted'
default:
return
}
@@ -58,10 +58,6 @@ const renderToast = (
case NotificationType.ScreenRecordingStarted:
case NotificationType.ScreenRecordingStopped:
case NotificationType.ScreenRecordingLimitReached:
case NotificationType.TranscriptionFailed:
case NotificationType.ScreenRecordingFailed:
case NotificationType.TranscriptionAborted:
case NotificationType.ScreenRecordingAborted:
return <ToastAnyRecording key={toast.key} toast={toast} state={state} />
case NotificationType.TranscriptionRequested:
+2 -3
View File
@@ -9,9 +9,8 @@ export enum RecordingStatus {
Stopped = 'stopped',
Saved = 'saved',
Aborted = 'aborted',
Failed = 'failed',
FailedToStart = 'failed_to_start',
FailedToStop = 'failed_to_stop',
FailedToStart = 'failedToStart',
FailedToStop = 'failedToStop',
NotificationSucceed = 'notification_succeeded',
ExternalProcessSuccessful = 'external_process_successful',
ExternalProcessFailed = 'external_process_failed',
@@ -37,8 +37,6 @@
"started": "{{name}} hat die Meeting-Transkription gestartet.",
"stopped": "{{name}} hat die Meeting-Transkription gestoppt.",
"limitReached": "Die Transkription hat die maximal zulässige Dauer überschritten und wird automatisch gespeichert.",
"failed": "Die Transkription wurde unerwartet beendet und konnte nicht gespeichert werden.",
"aborted": "Die Transkription konnte nicht gestartet werden.",
"requested": "{{name}} möchte die Meeting-Transkription starten."
},
"screenRecording": {
@@ -46,8 +44,6 @@
"started": "{{name}} hat die Meeting-Aufzeichnung gestartet.",
"stopped": "{{name}} hat die Meeting-Aufzeichnung gestoppt.",
"limitReached": "Die Aufzeichnung hat die maximal zulässige Dauer überschritten und wird automatisch gespeichert.",
"failed": "Die Aufzeichnung wurde unerwartet beendet und konnte nicht gespeichert werden.",
"aborted": "Die Aufzeichnung konnte nicht gestartet werden.",
"requested": "{{name}} möchte die Meeting-Aufzeichnung starten."
},
"recordingSave": {
@@ -37,8 +37,6 @@
"started": "{{name}} started the meeting transcription.",
"stopped": "{{name}} stopped the meeting transcription.",
"limitReached": "The transcription has exceeded the maximum allowed duration and will be automatically saved.",
"failed": "The transcription stopped unexpectedly and could not be saved.",
"aborted": "The transcription could not be started.",
"requested": "{{name}} wants to start the meeting transcription."
},
"screenRecording": {
@@ -46,8 +44,6 @@
"started": "{{name}} started the meeting recording.",
"stopped": "{{name}} stopped the meeting recording.",
"limitReached": "The recording has exceeded the maximum allowed duration and will be automatically saved.",
"failed": "The recording stopped unexpectedly and could not be saved.",
"aborted": "The recording could not be started.",
"requested": "{{name}} wants to start the meeting recording."
},
"recordingSave": {
@@ -37,8 +37,6 @@
"started": "{{name}} ha iniciado la transcripción de la reunión.",
"stopped": "{{name}} ha detenido la transcripción de la reunión.",
"limitReached": "La transcripción ha superado la duración máxima permitida, se va a guardar automáticamente.",
"failed": "La transcripción se ha interrumpido de forma inesperada y no se ha podido guardar.",
"aborted": "No se ha podido iniciar la transcripción.",
"requested": "{{name}} ha solicitado iniciar la transcripción."
},
"screenRecording": {
@@ -46,8 +44,6 @@
"started": "{{name}} ha iniciado la grabación de la reunión.",
"stopped": "{{name}} ha detenido la grabación de la reunión.",
"limitReached": "La grabación ha superado la duración máxima permitida, se va a guardar automáticamente.",
"failed": "La grabación se ha interrumpido de forma inesperada y no se ha podido guardar.",
"aborted": "No se ha podido iniciar la grabación.",
"requested": "{{name}} ha solicitado iniciar la grabación."
},
"recordingSave": {
@@ -37,8 +37,6 @@
"started": "{{name}} a démarré la transcription de la réunion.",
"stopped": "{{name}} a arrêté la transcription de la réunion.",
"limitReached": "La transcription a dépassé la durée maximale autorisée, elle va être automatiquement sauvegardée.",
"failed": "La transcription s'est interrompue de façon inattendue et n'a pas pu être sauvegardée.",
"aborted": "La transcription n'a pas pu être démarrée.",
"requested": "{{name}} a demandé à démarrer la transcription."
},
"screenRecording": {
@@ -46,8 +44,6 @@
"started": "{{name}} a démarré l'enregistrement de la réunion.",
"stopped": "{{name}} a arrêté l'enregistrement de la réunion.",
"limitReached": "L'enregistrement a dépassé la durée maximale autorisée, il va être automatiquement sauvegardé.",
"failed": "L'enregistrement s'est interrompu de façon inattendue et n'a pas pu être sauvegardé.",
"aborted": "L'enregistrement n'a pas pu être démarré.",
"requested": "{{name}} a demandé à démarrer l'enregistrement."
},
"recordingSave": {
@@ -37,8 +37,6 @@
"started": "{{name}} is de transcriptie van de vergadering gestart.",
"stopped": "{{name}} heeft de transcriptie van de vergadering gestopt.",
"limitReached": "De transcriptie heeft de maximaal toegestane duur overschreden en wordt automatisch opgeslagen.",
"failed": "De transcriptie is onverwacht gestopt en kon niet worden opgeslagen.",
"aborted": "De transcriptie kon niet worden gestart.",
"requested": "{{name}} wil graag de transcriptie starten."
},
"screenRecording": {
@@ -46,8 +44,6 @@
"started": "{{name}} is begonnen met het opnemen van de vergadering.",
"stopped": "{{name}} is gestopt met het opnemen van de vergadering.",
"limitReached": "De opname heeft de maximaal toegestane duur overschreden en wordt automatisch opgeslagen.",
"failed": "De opname is onverwacht gestopt en kon niet worden opgeslagen.",
"aborted": "De opname kon niet worden gestart.",
"requested": "{{name}} wil graag de opname starten."
},
"recordingSave": {