Compare commits

...

6 Commits

Author SHA1 Message Date
lebaudantoine 387269d1d2 wip allow setting encryption on a room 2026-09-25 00:39:28 +02:00
Lebaud Antoine 9a2ad63524 🔥(backend) remove unused API viewset and permission helpers
Dead code elements spotted by @briquet.
2026-09-24 23:08:00 +02:00
lebaudantoine 67f9e54784 🔒️(frontend) fix HIGH CVE-2026-93990 in libexpat
Bump `libexpat` from 2.8.4-r0 to 2.8.5-r0 to address the following
HIGH severity CVE, reported by Trivy on the frontend image
(alpine 3.24.1):

* CVE-2026-93990 — expat: XML injection via malformed UTF-16
  input.

  https://avd.aquasec.com/nvd/cve-2026-93990
2026-09-24 18:05:32 +02:00
lebaudantoine 5d3255ddbc 🔇(backend) drop warning log when room metadata is updated
Most call sites of `update_metadata` already wrap the call in a
try/except that logs the failure at info level.

Remove the warning log inside `update_metadata` itself to avoid
redundant logs, without losing any information.
2026-09-24 18:05:32 +02:00
leo d89b01b681 🐛(recording) handle FAILED and ABORTED LiveKit egresses
`EGRESS_ABORTED` and `EGRESS_FAILED` events were previously ignored, leaving
recordings indefinitely in `ACTIVE` state and potentially blocking subsequent
recordings with 409 errors. Add handling and logging for failed and aborted
egresses, discarding failed recordings while preserving the existing behavior
for savable recordings.

Rename `handle_complete` to `handle_savable` to reflect that it handles both
`EGRESS_COMPLETE` and `EGRESS_LIMIT_REACHED`.

Slight refactor to separate LiveKit event handling from recording concerns as
part of a general separation concern to allow for future SFU swapping.

NB:
- FAILED recordings are currently discarded although exploitable media files
may exist
- There is a theoretical hole: if stop observes EGRESS_FAILED before the
egress_ended webhook is processed, the recording is immediately marked as
FAILED. Since only ACTIVE and STOPPED recordings are savable, the webhook
then skips the failure notification and LiveKit error log. In that rare race
condition, participants may therefore not see the failure toast. We accept
this trade-off for now, as this should be very infrequent.
- Another theoretical hole: There is a short race window where the user
clicks stop while the limit-reached status is being processed. Since the user
explicitly requested the stop, we consider skipping the limit notification
acceptable and do not handle this case.

fix(recording): log aborted worker events at info level
2026-09-24 18:05:32 +02:00
briquet 660c0ed684 🔒️(backend) upgrade base image to python:3.13.15-alpine3.24
Fixes an XML injection attack in libexpat (CVE-2026-93990)
2026-09-24 17:59:24 +02:00
43 changed files with 1379 additions and 300 deletions
+6
View File
@@ -16,6 +16,7 @@ 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
@@ -29,6 +30,7 @@ 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
@@ -42,6 +44,10 @@ 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
## [1.31.0] - 2026-09-08
+1 -1
View File
@@ -1,7 +1,7 @@
# Django Meet
# ---- base image to inherit from ----
FROM python:3.13.5-alpine3.21 AS base
FROM python:3.13.15-alpine3.24 AS base
# Upgrade pip to its latest release to speed up dependencies installation
RUN python -m pip install --upgrade pip
+2 -1
View File
@@ -57,7 +57,8 @@ RUN npx webpack --mode production
FROM nginxinc/nginx-unprivileged:1.30.4-alpine3.24 AS frontend-production
USER root
RUN apk del curl
RUN apk upgrade --no-cache libexpat && \
apk del curl
USER nginx
USER nginx
+3
View File
@@ -76,6 +76,9 @@ 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,10 +12,6 @@ from ..services.participants_management import (
ParticipantsManagementException,
)
ACTION_FOR_METHOD_TO_PERMISSION = {
"versions_detail": {"DELETE": "versions_destroy", "GET": "versions_retrieve"}
}
class IsAuthenticated(permissions.BasePermission):
"""
@@ -27,15 +23,6 @@ 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
+47 -1
View File
@@ -38,6 +38,7 @@ class UserSerializer(serializers.ModelSerializer):
"short_name",
"timezone",
"language",
"default_encryption_mode",
"default_room_access_level",
"default_room_configuration",
]
@@ -53,6 +54,14 @@ 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."""
@@ -94,6 +103,7 @@ class ResourceAccessSerializerMixin:
raise PermissionDenied(
"Only owners of a room can assign other users as owners."
)
return data
def validate_resource(self, resource):
@@ -148,7 +158,15 @@ class RoomSerializer(serializers.ModelSerializer):
class Meta:
model = models.Room
fields = ["id", "name", "slug", "configuration", "access_level", "pin_code"]
fields = [
"id",
"name",
"slug",
"configuration",
"access_level",
"pin_code",
"encryption_mode",
]
read_only_fields = ["id", "slug", "pin_code"]
def validate_configuration(self, value):
@@ -161,6 +179,32 @@ 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.
@@ -197,12 +241,14 @@ 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"]
+24 -56
View File
@@ -95,60 +95,6 @@ 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.
@@ -303,6 +249,18 @@ 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 = {}
@@ -367,16 +325,21 @@ 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(
@@ -699,6 +662,11 @@ 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,6 +5,7 @@ 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 = (
@@ -32,3 +33,14 @@ 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")
@@ -0,0 +1,18 @@
# 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),
),
]
@@ -0,0 +1,46 @@
"""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",
),
),
]
+76 -8
View File
@@ -26,6 +26,7 @@ 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
@@ -58,6 +59,7 @@ 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")
@@ -79,17 +81,13 @@ 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."""
@@ -219,6 +217,15 @@ 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()
@@ -414,6 +421,13 @@ 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,
@@ -440,20 +454,68 @@ class Room(Resource):
return capfirst(self.name)
def save(self, *args, **kwargs):
"""Generate a unique n-digit pin code for new rooms."""
"""Restrict new encrypted rooms and allocate PINs for unencrypted rooms.
# Roomkit devices also join by PIN, so a PIN is needed as soon as
# either integration is enabled.
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()
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.
@@ -476,6 +538,11 @@ 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"""
@@ -582,6 +649,7 @@ 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,3 +8,50 @@ 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 LiveKit Events Service"""
# pylint: disable=no-member
"""Recording-related Events Service"""
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,22 +26,110 @@ 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 LiveKit webhook events."""
"""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.
"""
@staticmethod
def handle_update(recording: Recording, egress_status):
"""Handle egress status updates and sync recording state to room metadata."""
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.
"""
room_name = str(recording.room.id)
status_mapping = {
api.EgressStatus.EGRESS_ACTIVE: "started",
api.EgressStatus.EGRESS_ENDING: "saving",
api.EgressStatus.EGRESS_ABORTED: "aborted",
}
recording_status = status_mapping.get(egress_status)
recording_status = ROOM_METADATA_RECORDING_STATUSES.get(event)
if recording_status:
try:
RoomManagement.update_metadata(
@@ -55,42 +143,113 @@ class RecordingEventsService:
except RoomManagementException as e:
logger.exception("Failed to update room's metadata: %s", e)
@staticmethod
def handle_limit_reached(recording: Recording):
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):
"""Stop recording and notify participants when limit is reached."""
recording.status = models.RecordingStatusChoices.STOPPED
recording.save()
notification_mapping = {
models.RecordingModeChoices.SCREEN_RECORDING: "screenRecordingLimitReached",
models.RecordingModeChoices.TRANSCRIPT: "transcriptionLimitReached",
}
cls._notify_participants(recording, RecordingWorkerEvent.LIMIT_REACHED)
notification_type = notification_mapping.get(recording.mode)
if not notification_type:
return
@classmethod
def _handle_failed(cls, recording: Recording):
"""Set recording status to failed, matching the worker event, and notify participants.
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
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)
@staticmethod
def handle_complete(recording: Recording):
def _handle_successful(recording: Recording):
"""Notify external services and save recording."""
if not recording.is_savable():
+37 -5
View File
@@ -2,6 +2,8 @@
# pylint: disable=no-member
import logging
from asgiref.sync import async_to_sync
from livekit import api as livekit_api
@@ -10,6 +12,8 @@ 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."""
@@ -49,6 +53,22 @@ 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,
@@ -66,14 +86,26 @@ 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):
+55 -34
View File
@@ -12,15 +12,12 @@ 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 (
RecordingEventsError,
RecordingEventsService,
RecordingNotSavableError,
)
from core.recording.services.recording_events import RecordingEventsService
from .lobby import LobbyService
from .presence import PresenceCache
@@ -84,6 +81,30 @@ 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."""
@@ -173,12 +194,20 @@ class LiveKitEventsService:
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)
event = to_recording_event(data.egress_info.status)
if event is None:
return
self.recording_events.handle_update(recording, event)
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
@@ -188,6 +217,17 @@ 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(
@@ -201,38 +241,17 @@ 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 (
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
if event is None:
return
# 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
self.recording_events.handle_terminal_event(recording, event)
@staticmethod
def _is_connection_test_room(room_name: str) -> bool:
@@ -256,7 +275,9 @@ 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:
if (
settings.ROOM_TELEPHONY_ENABLED or settings.ROOMKIT_ENABLED
) and not room.is_encrypted:
try:
self.sip_management.ensure_dispatch_rule(room)
except SIPException as e:
+22 -5
View File
@@ -48,8 +48,11 @@ 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, str]:
def to_dict(self) -> Dict[str, object]:
"""Serialize the participant object to a dict representation."""
return {
"status": self.status.value,
@@ -57,6 +60,7 @@ class LobbyParticipant:
"id": self.id,
"color": self.color,
"entered_at": self.entered_at,
"is_authenticated": self.is_authenticated,
}
@classmethod
@@ -72,6 +76,7 @@ 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:")
@@ -144,7 +149,7 @@ class LobbyService:
key=settings.LOBBY_COOKIE_NAME,
value=participant_id,
httponly=True,
secure=True,
secure=not settings.DEBUG,
samesite="Lax",
)
@@ -208,6 +213,7 @@ 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
@@ -220,19 +226,24 @@ 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)
participant = self.enter(
room.id,
participant_id,
username,
is_authenticated=request.user.is_authenticated,
)
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,
@@ -241,6 +252,7 @@ class LobbyService:
configuration=room.configuration,
participant_id=participant_id,
role=user_role,
encryption_mode=room.encryption_mode,
)
return participant, livekit_config
@@ -258,7 +270,11 @@ class LobbyService:
self._index_touch(room_id)
def enter(
self, room_id: UUID, participant_id: str, username: str
self,
room_id: UUID,
participant_id: str,
username: str,
is_authenticated: bool = False,
) -> LobbyParticipant:
"""Add participant to waiting lobby."""
@@ -270,6 +286,7 @@ class LobbyService:
id=participant_id,
color=color,
entered_at=timezone.now().isoformat(),
is_authenticated=is_authenticated,
)
try:
@@ -76,10 +76,6 @@ 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,18 +2,23 @@
Test RecordingEventsService service.
"""
# pylint: disable=redefined-outer-name
# pylint: disable=redefined-outer-name,protected-access
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
@@ -34,10 +39,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(
@@ -48,13 +53,69 @@ def test_handle_limit_reached_success(mock_notify, mode, notification_type, serv
@pytest.mark.parametrize(
("mode", "notification_type"),
(
("screen_recording", "screenRecordingLimitReached"),
("transcript", "transcriptionLimitReached"),
("screen_recording", "screenRecordingFailed"),
("transcript", "transcriptionFailed"),
),
)
@mock.patch("core.utils.notify_participants")
def test_handle_limit_reached_error(mock_notify, mode, notification_type, service):
"""Test handle_limit_reached raises RecordingEventsError when notification fails."""
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.
"""
mock_notify.side_effect = NotificationError("Error notifying")
@@ -62,14 +123,15 @@ def test_handle_limit_reached_error(mock_notify, mode, notification_type, servic
with pytest.raises(
RecordingEventsError,
match=r"Failed to notify participants in room '.+' "
r"about recording limit reached \(recording_id=.+\)",
match=rf"Failed to notify participants in room '.+' "
rf"about recording {event} \(recording_id=.+\)",
):
service.handle_limit_reached(recording)
getattr(service, handler)(recording)
assert recording.status == "stopped"
assert recording.status == expected_status
mock_notify.assert_called_once_with(
room_name=str(recording.room.id), notification_data={"type": notification_type}
room_name=str(recording.room.id),
notification_data={"type": f"{notification_prefix}{notification_suffix}"},
)
@@ -82,19 +144,19 @@ def test_handle_limit_reached_error(mock_notify, mode, notification_type, servic
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
def test_handle_complete_saves_recording( # pylint: disable=too-many-arguments, too-many-positional-arguments
def test_handle_successful_saves_recording( # pylint: disable=too-many-arguments, too-many-positional-arguments
mock_notify_external_services,
notify_return_value,
expected_status,
status,
service,
):
"""Test handle_complete notifies external services and saves a savable recording."""
"""Test _handle_successful notifies external services and saves a savable recording."""
mock_notify_external_services.return_value = notify_return_value
recording = RecordingFactory(status=status)
service.handle_complete(recording)
service._handle_successful(recording)
mock_notify_external_services.assert_called_once_with(recording)
@@ -104,23 +166,293 @@ def test_handle_complete_saves_recording( # pylint: disable=too-many-arguments,
@pytest.mark.parametrize(
"status",
["initiated", "saved", "notification_succeeded", "aborted", "failed_to_start"],
[
"initiated",
"saved",
"notification_succeeded",
"aborted",
"failed",
"failed_to_start",
],
)
@mock.patch(
"core.recording.services.recording_events.notification_service."
"notify_external_services"
)
def test_handle_complete_non_savable_recording(
def test_handle_successful_non_savable_recording(
mock_notify_external_services, status, service
):
"""Test handle_complete refuses recordings that are already saved or in error."""
"""Test _handle_successful refuses recordings that are already saved or in error."""
recording = RecordingFactory(status=status)
with pytest.raises(RecordingNotSavableError):
service.handle_complete(recording)
service._handle_successful(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,6 +224,7 @@ 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,6 +4,7 @@ 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
@@ -154,9 +155,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):
@@ -175,6 +176,32 @@ 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,3 +312,34 @@ 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,6 +62,7 @@ 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,
}
@@ -75,9 +76,12 @@ def test_request_entry_anonymous(settings):
@freeze_time("2025-01-01 10:00:00")
def test_request_entry_authenticated_user(settings):
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
def test_request_entry_authenticated_user(settings, encryption_mode):
"""Authenticated users should be allowed to request entry."""
room = RoomFactory(access_level=RoomAccessLevel.RESTRICTED)
room = RoomFactory(
access_level=RoomAccessLevel.RESTRICTED, encryption_mode=encryption_mode
)
user = UserFactory()
client = APIClient()
client.force_login(user)
@@ -113,6 +117,7 @@ def test_request_entry_authenticated_user(settings):
"status": "waiting",
"color": "mocked-color",
"entered_at": "2025-01-01T10:00:00+00:00",
"is_authenticated": True,
"livekit": None,
}
@@ -189,6 +194,7 @@ 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,
}
@@ -243,6 +249,7 @@ 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"},
}
@@ -297,6 +304,7 @@ 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"},
}
@@ -354,6 +362,7 @@ 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"},
}
@@ -623,6 +632,7 @@ 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",
},
{
@@ -630,6 +640,7 @@ 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,6 +33,7 @@ def test_api_rooms_retrieve_anonymous_private_pk():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -52,6 +53,7 @@ def test_api_rooms_retrieve_anonymous_trusted_pk():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -70,6 +72,7 @@ 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,
}
@@ -86,6 +89,7 @@ def test_api_rooms_retrieve_anonymous_private_slug():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -102,6 +106,7 @@ 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,
}
@@ -217,6 +222,7 @@ 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()
@@ -263,6 +269,7 @@ 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(
@@ -273,6 +280,7 @@ def test_api_rooms_retrieve_authenticated_public(mock_token):
sources=["camera"],
role=None,
participant_id=None,
encryption_mode="none",
)
@@ -314,6 +322,7 @@ 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(
@@ -324,6 +333,7 @@ def test_api_rooms_retrieve_authenticated_trusted(mock_token):
sources=None,
role=None,
participant_id=None,
encryption_mode="none",
)
@@ -349,6 +359,7 @@ def test_api_rooms_retrieve_authenticated():
"id": str(room.id),
"name": room.name,
"slug": room.slug,
"encryption_mode": room.encryption_mode,
}
@@ -400,6 +411,7 @@ 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(
@@ -410,6 +422,7 @@ 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",
)
@@ -461,6 +474,7 @@ 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,
@@ -476,6 +490,7 @@ 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,
@@ -496,6 +511,7 @@ 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(
@@ -506,4 +522,25 @@ 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,3 +410,24 @@ 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,6 +3,7 @@ Test LiveKitEvents service.
"""
# pylint: disable=W0621,W0613, W0212, E0611
import logging
import uuid
from unittest import mock
@@ -10,13 +11,16 @@ 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
@@ -77,7 +81,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 notifies all participant."""
"""Should successfully stop recording and notify all participant."""
recording = RecordingFactory(worker_id="worker-1", mode=mode, status="active")
mock_data = mock.MagicMock()
@@ -104,7 +108,6 @@ 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")
@@ -125,29 +128,6 @@ 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"),
(
@@ -180,33 +160,38 @@ 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, service
mock_update_metadata, mock_notify, mock_notify_external_services, service
):
"""Should raise ActionFailedError when notification fails but still stop recording."""
"""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")
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
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"
service._handle_egress_ended(mock_data)
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")
@@ -232,17 +217,25 @@ 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_not_active(
mock_update_metadata, mock_notify, service
def test_handle_egress_ended_recording_should_not_be_saved(
mock_update_metadata, mock_notify, egress_status, service
):
"""Should ignore non-active recordings."""
"""Don't update status for recordings that must not be saved."""
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 = EgressStatus.EGRESS_LIMIT_REACHED
mock_data.egress_info.status = egress_status
service._handle_egress_ended(mock_data)
@@ -255,12 +248,24 @@ def test_handle_egress_ended_recording_not_active(
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_recording_not_limit_reached(
mock_update_metadata, mock_notify, service
def test_handle_egress_ended_complete_does_not_notify_participants(
mock_update_metadata, mock_notify, mock_notify_external_services, service
):
"""Should ignore egress non-limit-reached statuses."""
"""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
recording = RecordingFactory(worker_id="worker-1", status="stopped")
mock_data = mock.MagicMock()
@@ -273,7 +278,10 @@ def test_handle_egress_ended_recording_not_limit_reached(
mock_update_metadata.assert_called_once_with(
str(recording.room.id), remove_keys=["recording_mode", "recording_status"]
)
assert recording.status == "stopped"
mock_notify_external_services.assert_called_once_with(recording)
recording.refresh_from_db()
assert recording.status == "saved"
@mock.patch("core.services.livekit_events.MetadataCollectorService")
@@ -333,96 +341,163 @@ 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")
@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
def test_handle_egress_ended_does_not_log_error_on_successful_egress( # 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
"""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
"""Shouldn't log an egress error when LiveKit reports a successful egress."""
recording = RecordingFactory(worker_id="worker-1", status="active")
mock_data = mock.MagicMock()
mock_data.egress_info.egress_id = recording.worker_id
mock_data.egress_info.status = egress_status
service._handle_egress_ended(mock_data)
with caplog.at_level(logging.ERROR):
service._handle_egress_ended(mock_data)
mock_notify_external_services.assert_called_once_with(recording)
recording.refresh_from_db()
assert recording.status == recording_status
assert "Recording worker reported" not in caplog.text
@pytest.mark.parametrize(
"egress_status",
[
EgressStatus.EGRESS_STARTING,
EgressStatus.EGRESS_ACTIVE,
EgressStatus.EGRESS_ENDING,
EgressStatus.EGRESS_FAILED,
EgressStatus.EGRESS_ABORTED,
],
)
@mock.patch("core.utils.notify_participants")
@mock.patch("core.services.room_management.RoomManagement.update_metadata")
def test_handle_egress_ended_does_not_save_on_wrong_status(
mock_update_metadata, egress_status, service
def test_handle_egress_ended_logs_livekit_error_before_cleaning_up(
mock_update_metadata, mock_notify, service, caplog
):
"""Shouldn't save on invalid status."""
"""Should log the failure reason even when the cleanup fails afterwards."""
mock_update_metadata.side_effect = RuntimeError("LiveKit is unreachable")
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
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
service._handle_egress_ended(mock_data)
with caplog.at_level(logging.ERROR), pytest.raises(RuntimeError):
service._handle_egress_ended(mock_data)
recording.refresh_from_db()
assert recording.status == "active"
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
@pytest.mark.parametrize(
"status", ["failed_to_start", "aborted", "failed_to_stop", "saved", "initiated"]
("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),
),
)
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_ended_ignores_non_savable_recording(
mock_update_metadata, status, service
):
"""Should handle non-savable recordings idempotently without raising.
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."""
'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)
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_COMPLETE
mock_data.egress_info.egress_id = "worker-1"
mock_data.egress_info.status = 999
service._handle_egress_ended(mock_data)
service._handle_egress_updated(mock_data)
recording.refresh_from_db()
assert recording.status == status
mock_update_metadata.assert_not_called()
@mock.patch.object(LobbyService, "clear_room_cache")
+11 -1
View File
@@ -303,6 +303,7 @@ 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)
@@ -342,6 +343,7 @@ 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)
@@ -374,7 +376,12 @@ 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)
mock_enter.assert_called_once_with(
room.id,
participant_id,
username,
is_authenticated=request.user.is_authenticated,
)
lobby_service._get_participant.assert_called_once_with(room.id, participant_id)
@@ -442,6 +449,7 @@ 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)
@@ -483,6 +491,7 @@ 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)
@@ -886,6 +895,7 @@ 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,6 +127,7 @@ def test_api_users_retrieve_me_authenticated(settings):
"short_name": user.short_name,
"language": user.language,
"timezone": "UTC",
"default_encryption_mode": "none",
}
@@ -157,6 +157,7 @@ def test_models_recording_is_savable_normal():
@pytest.mark.parametrize(
"status",
[
RecordingStatusChoices.FAILED,
RecordingStatusChoices.FAILED_TO_STOP,
RecordingStatusChoices.FAILED_TO_START,
RecordingStatusChoices.ABORTED,
@@ -279,6 +280,7 @@ 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,3 +360,35 @@ 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
+18 -4
View File
@@ -57,21 +57,35 @@ def test_generate_token_authenticated_fallback_user_representation():
assert claims["name"] == str(user)
def test_generate_token_explicit_username_overrides_default():
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
def test_generate_token_explicit_username_overrides_default(encryption_mode):
"""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")
token = generate_token(
room="my-room",
user=user,
username="Custom Name",
encryption_mode=encryption_mode,
)
claims = decode_token(token)
assert claims["name"] == "Custom Name"
def test_authenticated_username_ignored_when_editing_disabled(settings):
@pytest.mark.parametrize("encryption_mode", ["none", "basic"])
def test_authenticated_username_ignored_when_editing_disabled(
settings, encryption_mode
):
"""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")
token = generate_token(
room="my-room",
user=user,
username="Custom Name",
encryption_mode=encryption_mode,
)
claims = decode_token(token)
assert claims["name"] == "Jane Doe"
+15 -3
View File
@@ -34,6 +34,9 @@ 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__)
@@ -69,6 +72,7 @@ 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.
@@ -91,10 +95,8 @@ def generate_token( # noqa: PLR0917
"""
is_admin_or_owner = role in ("owner", "administrator")
if is_admin_or_owner:
sources = settings.LIVEKIT_DEFAULT_SOURCES
if sources is None:
if is_admin_or_owner or sources is None:
sources = settings.LIVEKIT_DEFAULT_SOURCES
video_grants = VideoGrants(
@@ -141,6 +143,14 @@ 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()
@@ -152,6 +162,7 @@ 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.
@@ -184,6 +195,7 @@ def generate_livekit_config( # noqa: PLR0917
sources=sources,
role=role,
participant_id=participant_id,
encryption_mode=encryption_mode,
),
}
+4
View File
@@ -978,6 +978,10 @@ 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
+2 -1
View File
@@ -45,7 +45,8 @@ RUN npm run build
FROM nginxinc/nginx-unprivileged:1.30.4-alpine3.24 AS frontend-production
USER root
RUN apk del curl
RUN apk upgrade --no-cache libexpat && \
apk del curl
USER nginx
# Un-privileged user running the application
@@ -96,6 +96,10 @@ 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,11 +10,15 @@ 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,12 +22,20 @@ 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,6 +58,10 @@ 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:
+3 -2
View File
@@ -9,8 +9,9 @@ export enum RecordingStatus {
Stopped = 'stopped',
Saved = 'saved',
Aborted = 'aborted',
FailedToStart = 'failedToStart',
FailedToStop = 'failedToStop',
Failed = 'failed',
FailedToStart = 'failed_to_start',
FailedToStop = 'failed_to_stop',
NotificationSucceed = 'notification_succeeded',
ExternalProcessSuccessful = 'external_process_successful',
ExternalProcessFailed = 'external_process_failed',
@@ -37,6 +37,8 @@
"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": {
@@ -44,6 +46,8 @@
"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,6 +37,8 @@
"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": {
@@ -44,6 +46,8 @@
"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,6 +37,8 @@
"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": {
@@ -44,6 +46,8 @@
"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,6 +37,8 @@
"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": {
@@ -44,6 +46,8 @@
"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,6 +37,8 @@
"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": {
@@ -44,6 +46,8 @@
"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": {