(visio) use compatible with summary v2

This commit introduces the compatibility with summary v2.
It tecnically doesn't break the compatibility with v1 as
v1 params are still sent. But we advise people using the
transcribe feature in their own deployments to adapt to the
new v2 API, as this compatibility will be removed in a
future major version.

* RecordingStatusChoices now has
EXTERNAL_PROCESS_SUCCESSFUL
& EXTERNAL_PROCESS_FAILED
Values, which are changed by a new webhook that
can be called by the transcribe service.
This webhook is protected by its own bearer token.
* Title for the document is computed in visio,
* Tests are added / updated accordingly
This commit is contained in:
Florent Chehab
2026-05-28 21:45:37 +02:00
committed by aleb_the_flash
parent ecf8f0fe3f
commit 7fa172d100
21 changed files with 1282 additions and 416 deletions
+2 -2
View File
@@ -11,7 +11,7 @@ from core.recording.event import notification
from . import models
from .tasks.file import process_file_deletion
from .utils import generate_download_file_url
from .utils import generate_download_s3_url
def hard_delete_file(file):
@@ -242,7 +242,7 @@ class FileAdmin(admin.ModelAdmin):
"""Return a clickable preview URL for the file."""
if not obj.is_ready:
return "-"
url = generate_download_file_url(obj, expires_in=60 * 60)
url = generate_download_s3_url(obj.key, expires_in=60 * 60)
return format_html(
'<a href="{}" target="_blank" rel="noopener noreferrer">Open File</a>', url
+10
View File
@@ -563,3 +563,13 @@ class RenameParticipantSerializer(BaseValidationOnlySerializer):
"""Serializer for renaming a participant in a room."""
name = serializers.CharField(min_length=1, max_length=255, allow_blank=False)
class ExternalProcessEventSerializer(BaseValidationOnlySerializer):
"""Validate external process event data."""
job_id = serializers.CharField(required=True)
# We are not strict on purpose on those fields to avoid
# useless bad requests
type = serializers.CharField(required=False, allow_null=True, allow_blank=True)
status = serializers.CharField(required=False, allow_null=True, allow_blank=True)
+58 -1
View File
@@ -39,7 +39,10 @@ from core import analytics, enums, models, utils
from core.api.filters import ListFileFilter
from core.enums import MEDIA_STORAGE_URL_PATTERN
from core.recording.enums import FileExtension
from core.recording.event.authentication import StorageEventAuthentication
from core.recording.event.authentication import (
RecordingProcessWebhookAuthentication,
StorageEventAuthentication,
)
from core.recording.event.exceptions import (
InvalidBucketError,
InvalidFilepathError,
@@ -999,6 +1002,60 @@ class RecordingViewSet(
{"message": "Event processed."},
)
@decorators.action(
detail=False,
methods=["post"],
url_path="external-process-hook",
authentication_classes=[RecordingProcessWebhookAuthentication],
serializer_class=serializers.ExternalProcessEventSerializer,
)
def on_external_process_event_received(self, request, pk=None): # pylint: disable=unused-argument
"""Handle incoming external process events for recordings."""
logger.debug("Processing external process event %s", request.data)
serializer = self.get_serializer(data=request.data)
serializer.is_valid(raise_exception=True)
ok_response = drf_response.Response(
{"message": "Event processed."},
)
validated_data = serializer.validated_data
job_id = validated_data["job_id"]
try:
recording = models.Recording.objects.get(external_process_id=job_id)
except models.Recording.DoesNotExist as e:
logger.warning("No recording found for job_id %s: %s", job_id, e)
return ok_response
if validated_data.get("type") == "transcript":
if validated_data.get("status") == "success":
logger.info(
"External process transcript success received for recording %s",
job_id,
)
recording.status = (
models.RecordingStatusChoices.EXTERNAL_PROCESS_SUCCESSFUL
)
recording.save()
return ok_response
if validated_data.get("status") == "failure":
logger.info(
"External process transcript failure received for recording %s",
job_id,
)
recording.status = models.RecordingStatusChoices.EXTERNAL_PROCESS_FAILED
recording.save()
return ok_response
logger.info(
"No changes to save for external process id %s and payload %s",
job_id,
validated_data,
)
return ok_response
def _auth_get_original_url(self, request):
"""
Extracts and parses the original URL from the "HTTP_X_ORIGINAL_URL" header.
@@ -0,0 +1,23 @@
# Generated by Django 5.2.14 on 2026-06-22 08:26
from django.db import migrations, models
class Migration(migrations.Migration):
dependencies = [
('core', '0020_alter_file_upload_state'),
]
operations = [
migrations.AddField(
model_name='recording',
name='external_process_id',
field=models.CharField(blank=True, help_text='ID of the external process associated with the recording.', max_length=255, null=True, unique=True, verbose_name='External Process ID'),
),
migrations.AlterField(
model_name='recording',
name='status',
field=models.CharField(choices=[('initiated', 'Initiated'), ('active', 'Active'), ('stopped', 'Stopped'), ('saved', 'Saved'), ('aborted', 'Aborted'), ('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),
),
]
+17
View File
@@ -60,6 +60,11 @@ class RecordingStatusChoices(models.TextChoices):
FAILED_TO_START = "failed_to_start", _("Failed to Start")
FAILED_TO_STOP = "failed_to_stop", _("Failed to Stop")
NOTIFICATION_SUCCEEDED = "notification_succeeded", _("Notification succeeded")
EXTERNAL_PROCESS_SUCCESSFUL = (
"external_process_successful",
_("External process successful"),
)
EXTERNAL_PROCESS_FAILED = "external_process_failed", _("External process failed")
@classmethod
def is_final(cls, status):
@@ -73,6 +78,8 @@ class RecordingStatusChoices(models.TextChoices):
cls.STOPPED,
cls.SAVED,
cls.ABORTED,
cls.EXTERNAL_PROCESS_SUCCESSFUL,
cls.EXTERNAL_PROCESS_FAILED,
cls.FAILED_TO_START,
cls.FAILED_TO_STOP,
}
@@ -598,6 +605,14 @@ class Recording(BaseModel):
verbose_name=_("Recording options"),
help_text=_("Recording options"),
)
external_process_id = models.CharField(
max_length=255,
null=True,
blank=True,
unique=True,
verbose_name=_("External Process ID"),
help_text=_("ID of the external process associated with the recording."),
)
class Meta:
db_table = "meet_recording"
@@ -653,6 +668,8 @@ class Recording(BaseModel):
return self.status in {
RecordingStatusChoices.NOTIFICATION_SUCCEEDED,
RecordingStatusChoices.SAVED,
RecordingStatusChoices.EXTERNAL_PROCESS_SUCCESSFUL,
RecordingStatusChoices.EXTERNAL_PROCESS_FAILED,
}
@property
@@ -14,9 +14,9 @@ logger = logging.getLogger(__name__)
class MachineUser:
"""Represent a non-interactive system user for automated storage operations."""
def __init__(self) -> None:
def __init__(self, username: str = "storage_event_user") -> None:
self.pk = None
self.username = "storage_event_user"
self.username = username
self.is_active = True
@property
@@ -91,3 +91,46 @@ class StorageEventAuthentication(BaseAuthentication):
def authenticate_header(self, request):
"""Return the WWW-Authenticate header value."""
return f"{self.TOKEN_TYPE} realm='Storage event API'"
class RecordingProcessWebhookAuthentication(BaseAuthentication):
"""
Custom authentication class for recording process webhook requests.
Validates the API key in the Authorization header.
"""
AUTH_HEADER = "Authorization"
TOKEN_TYPE = "Bearer" # noqa S105
def authenticate(self, request):
"""
Authenticate the request and return a two-tuple of (user, token).
"""
required_token = settings.SUMMARY_SERVICE_WEBHOOK_API_TOKEN
if not required_token:
raise AuthenticationFailed("Webhook authentication is not configured.")
auth_header: str = request.headers.get("Authorization") or ""
if not auth_header.startswith("Bearer "):
logger.warning(
"Authentication failed: Invalid authorization header format (ip: %s)",
request.META.get("REMOTE_ADDR"),
)
raise AuthenticationFailed("Invalid authorization header format.")
token = auth_header[7:] # len("Bearer ") == 7
if not secrets.compare_digest(
token,
required_token,
):
logger.warning(
"Authentication failed: Bad Authorization header (ip: %s)",
request.META.get("REMOTE_ADDR"),
)
raise AuthenticationFailed()
return MachineUser("external_process_user"), None
def authenticate_header(self, request):
"""Return the WWW-Authenticate header value."""
return f"{self.TOKEN_TYPE} realm='External process webhook API'"
@@ -4,6 +4,7 @@ import asyncio
import logging
import smtplib
from datetime import datetime, timezone
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
from django.conf import settings
from django.core.mail import send_mail
@@ -17,6 +18,7 @@ from asgiref.sync import async_to_sync
from livekit import api as livekit_api
from core import models, utils
from core.utils import generate_download_s3_url
logger = logging.getLogger(__name__)
@@ -178,8 +180,49 @@ class NotificationService:
return _ns_to_utc(file_result.started_at), _ns_to_utc(file_result.ended_at)
@staticmethod
def _generate_title(
*,
locale: str,
room: str,
recording_datetime: datetime | None,
owner_timezone: str | None,
) -> str:
"""Generate title from context or return default."""
if recording_datetime is None:
with override(locale):
return _("Transcription")
dt = recording_datetime
if owner_timezone:
try:
dt = recording_datetime.astimezone(ZoneInfo(owner_timezone))
except (KeyError, ZoneInfoNotFoundError):
pass # Keep the original UTC datetime
with override(locale):
translated_template = _(
'Meeting "{room}" on {room_recording_date} at {room_recording_time}'
)
return translated_template.format(
room=room,
room_recording_date=dt.strftime("%Y-%m-%d"),
room_recording_time=dt.strftime("%H:%M"),
)
@staticmethod
def _notify_summary_service(recording: models.Recording):
if settings.SUMMARY_SERVICE_VERSION == 1:
return NotificationService._notify_summary_service_v1(recording)
if settings.SUMMARY_SERVICE_VERSION == 2:
return NotificationService._notify_summary_service_v2(recording)
raise NotImplementedError(
f"Unknown summary service version: {settings.SUMMARY_SERVICE_VERSION}"
)
@staticmethod
def _notify_summary_service_v1(recording: models.Recording):
"""Notify summary service about a new recording."""
if (
@@ -253,5 +296,119 @@ class NotificationService:
return True
@staticmethod
def _notify_summary_service_v2(recording: models.Recording):
"""Notify summary service about a new recording."""
if (
not settings.SUMMARY_SERVICE_ENDPOINT
or not settings.SUMMARY_SERVICE_API_TOKEN
):
logger.error("Summary service not configured")
return False
owner_access = (
models.RecordingAccess.objects.select_related("user")
.filter(
role=models.RoleChoices.OWNER,
recording_id=recording.id,
)
.first()
)
metadata_filename: None | str = None
if settings.METADATA_COLLECTOR_ENABLED and recording.options.get(
"collect_metadata", False
):
output_folder = settings.METADATA_COLLECTOR_OUTPUT_FOLDER
metadata_filename = f"{output_folder}/{recording.id}-metadata.json"
if not owner_access:
logger.error("No owner found for recording %s", recording.id)
return False
started_at, ended_at = async_to_sync(
NotificationService._get_recording_timestamps
)(recording.worker_id)
form_base_url = settings.TRANSCRIPTION_SATISFACTION_FORM_BASE_URL
form_link = (
f"{form_base_url}?room_id={recording.room.id}"
if (form_base_url and metadata_filename is not None)
else None
)
metadata_payload = None
if started_at and ended_at and metadata_filename:
metadata_payload = {
"cloud_storage_url": generate_download_s3_url(
metadata_filename,
expires_in=settings.SUMMARY_SERVICE_CLOUD_STORAGE_SIGNED_URL_EXPIRY_SECONDS,
override_domain=False,
),
"started_at": started_at.isoformat(),
"ended_at": ended_at.isoformat(),
}
payload = {
"user_sub": owner_access.user.sub,
"user_email": owner_access.user.email,
"cloud_storage_url": generate_download_s3_url(
recording.key,
expires_in=settings.SUMMARY_SERVICE_CLOUD_STORAGE_SIGNED_URL_EXPIRY_SECONDS,
override_domain=False,
),
"language": recording.options.get(
"language", get_language().split("-")[0].lower()
),
"context_language": owner_access.user.language,
"push_to_docs_config": {
"user_email": owner_access.user.email,
"title": NotificationService._generate_title(
locale=owner_access.user.language
or recording.options.get("language", get_language()),
room=recording.room.name,
recording_datetime=started_at,
owner_timezone=str(owner_access.user.timezone),
),
"download_link": f"{get_recording_download_base_url()}/{recording.id}",
"form_link": form_link,
# For now the feature flag logic is handled on summary side
"auto_create_summary": True,
},
"metadata": metadata_payload,
}
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {settings.SUMMARY_SERVICE_API_TOKEN}",
}
try:
response = requests.post(
settings.SUMMARY_SERVICE_ENDPOINT,
json=payload,
headers=headers,
timeout=30,
)
response.raise_for_status()
response_json = response.json()
# We do not require a job_id to avoid a breaking change
job_id = response_json.get("job_id")
if not isinstance(job_id, str):
raise ValueError("job_id is not a string")
recording.external_process_id = job_id
recording.save()
except requests.RequestException as exc:
logger.exception(
"Summary service error for recording %s. URL: %s. Exception: %s",
recording.id,
settings.SUMMARY_SERVICE_ENDPOINT,
exc,
)
return False
return True
notification_service = NotificationService()
@@ -243,3 +243,168 @@ def test_notify_user_by_email_smtp_exception(mocked_current_site, caplog):
assert result is False
assert mock_send_mail.call_count == 2
assert "notification could not be sent:" in caplog.text
@mock.patch("core.recording.event.notification.requests.post")
@mock.patch("core.recording.event.notification.generate_download_s3_url")
@mock.patch.object(
NotificationService, "_get_recording_timestamps", new_callable=mock.AsyncMock
)
def test_notify_summary_service_post_args_with_metadata(
mock_get_recording_timestamps,
mock_generate_download_s3_url,
mock_post,
settings,
):
"""Test summary notification computed request args when metadata is enabled."""
settings.SUMMARY_SERVICE_VERSION = 2
settings.SUMMARY_SERVICE_ENDPOINT = "https://summary.test/api/v2/tasks"
settings.SUMMARY_SERVICE_API_TOKEN = "summary-token"
settings.RECORDING_DOWNLOAD_BASE_URL = "https://app.test/recordings"
settings.SCREEN_RECORDING_BASE_URL = None
settings.METADATA_COLLECTOR_ENABLED = True
settings.METADATA_COLLECTOR_OUTPUT_FOLDER = "recordings-metadata"
recording = factories.RecordingFactory(
room__name="Engineering Sync",
worker_id="egress-1",
options={"collect_metadata": True, "language": "en-us"},
)
owner = factories.UserFactory(
email="owner@test.com",
sub="owner-sub",
language="fr-fr",
timezone="Europe/Paris",
)
factories.UserRecordingAccessFactory(
recording=recording, role=models.RoleChoices.OWNER, user=owner
)
started_at = datetime.datetime(2026, 1, 2, 10, 30, tzinfo=datetime.timezone.utc)
ended_at = datetime.datetime(2026, 1, 2, 11, 45, tzinfo=datetime.timezone.utc)
mock_get_recording_timestamps.return_value = (started_at, ended_at)
mock_generate_download_s3_url.side_effect = [
"https://storage.test/metadata.json",
"https://storage.test/recording.ogg",
]
mock_response = mock.Mock()
mock_response.raise_for_status.return_value = None
mock_response.json.return_value = {"job_id": "job-42"}
mock_post.return_value = mock_response
result = NotificationService._notify_summary_service(recording)
recording.refresh_from_db()
assert result is True
assert recording.external_process_id == "job-42"
metadata_filename = (
f"{settings.METADATA_COLLECTOR_OUTPUT_FOLDER}/{recording.id}-metadata.json"
)
expected_payload = {
"user_sub": owner.sub,
"user_email": owner.email,
"cloud_storage_url": "https://storage.test/recording.ogg",
"language": "en-us",
"context_language": owner.language,
"push_to_docs_config": {
"user_email": owner.email,
"title": 'Réunion "Engineering Sync" du 2026-01-02 à 11:30',
"download_link": f"{settings.RECORDING_DOWNLOAD_BASE_URL}/{recording.id}",
"auto_create_summary": True,
"form_link": None,
},
"metadata": {
"cloud_storage_url": "https://storage.test/metadata.json",
"started_at": started_at.isoformat(),
"ended_at": ended_at.isoformat(),
},
}
expected_headers = {
"Content-Type": "application/json",
"Authorization": "Bearer summary-token",
}
mock_post.assert_called_once_with(
"https://summary.test/api/v2/tasks",
json=expected_payload,
headers=expected_headers,
timeout=30,
)
assert mock_generate_download_s3_url.call_args_list == [
mock.call(metadata_filename, expires_in=60 * 60 * 24, override_domain=False),
mock.call(recording.key, expires_in=60 * 60 * 24, override_domain=False),
]
mock_get_recording_timestamps.assert_awaited_once_with("egress-1")
@mock.patch("core.recording.event.notification.requests.post")
@mock.patch("core.recording.event.notification.generate_download_s3_url")
@mock.patch.object(
NotificationService, "_get_recording_timestamps", new_callable=mock.AsyncMock
)
def test_notify_summary_service_post_args_without_metadata(
mock_get_recording_timestamps,
mock_generate_download_s3_url,
mock_post,
settings,
):
"""Test summary notification computed request args when metadata is not available."""
settings.SUMMARY_SERVICE_VERSION = 2
settings.SUMMARY_SERVICE_ENDPOINT = "https://summary.test/api/v2/tasks"
settings.SUMMARY_SERVICE_API_TOKEN = "summary-token"
settings.RECORDING_DOWNLOAD_BASE_URL = "https://app.test/recordings"
settings.SCREEN_RECORDING_BASE_URL = None
settings.METADATA_COLLECTOR_ENABLED = False
recording = factories.RecordingFactory(room__name="Daily")
owner = factories.UserFactory(
email="owner@test.com",
sub="owner-sub",
language="en-us",
timezone="UTC",
)
factories.UserRecordingAccessFactory(
recording=recording, role=models.RoleChoices.OWNER, user=owner
)
mock_get_recording_timestamps.return_value = (None, None)
mock_generate_download_s3_url.return_value = "https://storage.test/recording.mp4"
mock_response = mock.Mock()
mock_response.raise_for_status.return_value = None
mock_response.json.return_value = {"job_id": "job-51"}
mock_post.return_value = mock_response
result = NotificationService._notify_summary_service(recording)
assert result is True
expected_payload = {
"user_sub": owner.sub,
"user_email": owner.email,
"cloud_storage_url": "https://storage.test/recording.mp4",
"language": "en",
"context_language": owner.language,
"push_to_docs_config": {
"user_email": owner.email,
"title": "Transcription",
"download_link": f"{settings.RECORDING_DOWNLOAD_BASE_URL}/{recording.id}",
"auto_create_summary": True,
"form_link": None,
},
"metadata": None,
}
expected_headers = {
"Content-Type": "application/json",
"Authorization": "Bearer summary-token",
}
mock_post.assert_called_once_with(
"https://summary.test/api/v2/tasks",
json=expected_payload,
headers=expected_headers,
timeout=30,
)
mock_generate_download_s3_url.assert_called_once_with(
recording.key, expires_in=60 * 60 * 24, override_domain=False
)
mock_get_recording_timestamps.assert_awaited_once_with(recording.worker_id)
@@ -0,0 +1,134 @@
"""
Test recordings API endpoints: external process hook.
"""
# pylint: disable=redefined-outer-name,unused-argument
import pytest
from ...factories import RecordingFactory
from ...models import RecordingStatusChoices
pytestmark = pytest.mark.django_db
@pytest.fixture
def external_process_settings(settings):
"""Configure authentication token for the external process webhook."""
settings.SUMMARY_SERVICE_WEBHOOK_API_TOKEN = "testWebhookToken"
return settings
def test_external_process_event_missing_authorization_header(
external_process_settings, client
):
"""Requests without authorization must be rejected."""
response = client.post(
"/api/v1.0/recordings/external-process-hook/",
{"job_id": "job-1", "type": "transcript", "status": "success"},
)
assert response.status_code == 401
def test_external_process_event_wrong_bearer_token(external_process_settings, client):
"""Requests with invalid bearer token must be rejected."""
response = client.post(
"/api/v1.0/recordings/external-process-hook/",
{"job_id": "job-1", "type": "transcript", "status": "success"},
HTTP_AUTHORIZATION="Bearer wrongToken",
)
assert response.status_code == 401
def test_external_process_event_missing_job_id(external_process_settings, client):
"""Payload without job_id must fail validation."""
response = client.post(
"/api/v1.0/recordings/external-process-hook/",
{"type": "transcript", "status": "success"},
HTTP_AUTHORIZATION="Bearer testWebhookToken",
)
assert response.status_code == 400
assert response.json() == {"job_id": ["This field is required."]}
def test_external_process_event_success_updates_recording_status(
external_process_settings, client
):
"""A successful transcript process should update recording status."""
recording = RecordingFactory(
status=RecordingStatusChoices.SAVED,
external_process_id="job-123",
)
response = client.post(
"/api/v1.0/recordings/external-process-hook/",
{"job_id": "job-123", "type": "transcript", "status": "success"},
HTTP_AUTHORIZATION="Bearer testWebhookToken",
)
assert response.status_code == 200
assert response.json() == {"message": "Event processed."}
recording.refresh_from_db()
assert recording.status == RecordingStatusChoices.EXTERNAL_PROCESS_SUCCESSFUL
def test_external_process_event_failure_updates_recording_status(
external_process_settings, client
):
"""A failing transcript process should update recording status."""
recording = RecordingFactory(
status=RecordingStatusChoices.SAVED,
external_process_id="job-456",
)
response = client.post(
"/api/v1.0/recordings/external-process-hook/",
{"job_id": "job-456", "type": "transcript", "status": "failure"},
HTTP_AUTHORIZATION="Bearer testWebhookToken",
)
assert response.status_code == 200
assert response.json() == {"message": "Event processed."}
recording.refresh_from_db()
assert recording.status == RecordingStatusChoices.EXTERNAL_PROCESS_FAILED
def test_external_process_event_unknown_recording_is_ignored(
external_process_settings, client
):
"""Unknown job_id should not fail the webhook processing."""
response = client.post(
"/api/v1.0/recordings/external-process-hook/",
{"job_id": "missing-job", "type": "transcript", "status": "success"},
HTTP_AUTHORIZATION="Bearer testWebhookToken",
)
assert response.status_code == 200
assert response.json() == {"message": "Event processed."}
def test_external_process_event_non_transcript_event_does_not_change_status(
external_process_settings, client
):
"""Only transcript events should update recording status."""
recording = RecordingFactory(
status=RecordingStatusChoices.SAVED,
external_process_id="job-789",
)
response = client.post(
"/api/v1.0/recordings/external-process-hook/",
{"job_id": "job-789", "type": "thumbnail", "status": "success"},
HTTP_AUTHORIZATION="Bearer testWebhookToken",
)
assert response.status_code == 200
assert response.json() == {"message": "Event processed."}
recording.refresh_from_db()
assert recording.status == RecordingStatusChoices.SAVED
+6 -4
View File
@@ -459,12 +459,14 @@ def generate_upload_policy(file):
return policy
def generate_download_file_url(file, *, expires_in: int, override_domain: bool = True):
def generate_download_s3_url(
key: str, *, expires_in: int, override_domain: bool = True
):
"""
Generate a S3 signed download url for a given file.
Generate a S3 signed download url for a given key.
"""
key = file.file_key
if not key:
raise ValueError("key cannot be empty")
# This setting should be used if the backend application and the frontend application
# can't connect to the object storage with the same domain. This is the case in the