From d8ab229813c6535f1994aab392f52ff2718f4157 Mon Sep 17 00:00:00 2001 From: leo <260626284+cameledev@users.noreply.github.com> Date: Tue, 4 Aug 2026 11:36:22 +0200 Subject: [PATCH] wip --- CHANGELOG.md | 1 + Makefile | 6 + compose.yml | 30 ++- docker/auth/realm.json | 136 +++++++++++++ docs/features/recording.md | 55 ++++++ env.d/development/common.dist | 12 +- env.d/development/metadata_collector.dist | 2 +- env.d/development/summary.dist | 2 +- src/backend/core/api/viewsets.py | 32 ++++ src/backend/core/models.py | 44 +++++ .../core/recording/event/notification.py | 36 +++- .../recording/services/recording_events.py | 2 +- src/backend/core/services/drive.py | 181 ++++++++++++++++++ src/backend/core/tasks/__init__.py | 11 ++ src/backend/core/tasks/_task.py | 1 + src/backend/core/tasks/push_recording.py | 110 +++++++++++ .../core/tests/files/test_api_files_create.py | 2 +- src/backend/core/utils.py | 38 ++++ src/backend/meet/settings.py | 37 ++++ 19 files changed, 724 insertions(+), 14 deletions(-) create mode 100644 src/backend/core/services/drive.py create mode 100644 src/backend/core/tasks/push_recording.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 3cd8f588..f4e776fa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -17,6 +17,7 @@ and this project adheres to - ✨(backend) add roomkit viewset to start a room without WebRTC join - ✨(frontend) let users set default configuration for generated links - ✨(frontend) expose media state to external gateways +- ✨(backend) push recordings to the owner's Drive (POC) ### Changed diff --git a/Makefile b/Makefile index 2e5ae839..81c5b824 100644 --- a/Makefile +++ b/Makefile @@ -93,6 +93,7 @@ bootstrap: \ data/media \ data/static \ create-env-files \ + create-docker-network \ build \ migrate \ demo \ @@ -126,11 +127,16 @@ down: ## stop and remove containers, networks, images, and volumes @$(COMPOSE) down .PHONY: down +create-docker-network: ## create the shared lasuite-network if it doesn't exist + @docker network create lasuite-network || true +.PHONY: create-docker-network + logs: ## display app-dev logs (follow mode) @$(COMPOSE) logs -f app-dev .PHONY: logs run-backend: ## start only the backend application and all needed services + @$(MAKE) create-docker-network @$(COMPOSE) up --force-recreate -d celery-dev --remove-orphans @$(COMPOSE) up --force-recreate -d nginx @$(COMPOSE) up -d livekit diff --git a/compose.yml b/compose.yml index f60537e4..3499ee79 100644 --- a/compose.yml +++ b/compose.yml @@ -14,6 +14,9 @@ services: image: sj26/mailcatcher:latest ports: - "1081:1080" + networks: + - default + - lasuite minio: user: ${DOCKER_USER:-1000} @@ -33,6 +36,10 @@ services: command: minio server --console-address :9001 /data volumes: - ./data/media:/data + networks: + default: + aliases: + - meet-minio createbuckets: image: minio/mc @@ -92,6 +99,7 @@ services: networks: - resource-server - default + - lasuite celery-dev: user: ${DOCKER_USER:-1000} @@ -108,6 +116,9 @@ services: - /app/.venv depends_on: - app-dev + networks: + - default + - lasuite app: build: @@ -195,32 +206,32 @@ services: - env.d/development/kc_postgresql keycloak: - image: quay.io/keycloak/keycloak:20.0.1 + image: quay.io/keycloak/keycloak:26.3.2 volumes: - ./docker/auth/realm.json:/opt/keycloak/data/import/realm.json command: - start-dev - --features=preview - --import-realm - - --proxy=edge - - --hostname-url=http://localhost:8083 - - --hostname-admin-url=http://localhost:8083/ + - --proxy-headers=xforwarded + - --hostname=http://localhost:8083 - --hostname-strict=false - - --hostname-strict-https=false environment: - KEYCLOAK_ADMIN: admin - KEYCLOAK_ADMIN_PASSWORD: admin + KC_BOOTSTRAP_ADMIN_USERNAME: admin + KC_BOOTSTRAP_ADMIN_PASSWORD: admin KC_DB: postgres KC_DB_URL_HOST: kc_postgresql KC_DB_URL_DATABASE: keycloak KC_DB_PASSWORD: pass KC_DB_USERNAME: meet KC_DB_SCHEMA: public - PROXY_ADDRESS_FORWARDING: 'true' ports: - "8080:8080" depends_on: - kc_postgresql + networks: + - default + - lasuite livekit: image: livekit/livekit-server @@ -337,3 +348,6 @@ services: networks: default: resource-server: + lasuite: + name: lasuite-network + external: true diff --git a/docker/auth/realm.json b/docker/auth/realm.json index 2746c781..5a686b36 100644 --- a/docker/auth/realm.json +++ b/docker/auth/realm.json @@ -845,6 +845,142 @@ "offline_access", "microprofile-jwt" ] + }, + { + "clientId": "drive", + "name": "", + "description": "", + "rootUrl": "", + "adminUrl": "", + "baseUrl": "", + "surrogateAuthRequired": false, + "enabled": "true", + "alwaysDisplayInConsole": false, + "clientAuthenticatorType": "client-secret", + "secret": "ThisIsAnExampleKeyForDevPurposeOnly", + "redirectUris": [ + "http://localhost:3100/*", + "http://localhost:8171/*", + "http://localhost:8085/*" + ], + "webOrigins": [ + "http://localhost:3100", + "http://localhost:8171", + "http://localhost:8085" + ], + "notBefore": 0, + "bearerOnly": false, + "consentRequired": false, + "standardFlowEnabled": true, + "implicitFlowEnabled": false, + "directAccessGrantsEnabled": false, + "serviceAccountsEnabled": false, + "publicClient": false, + "frontchannelLogout": true, + "protocol": "openid-connect", + "attributes": { + "access.token.lifespan": "-1", + "client.secret.creation.time": "1707820779", + "user.info.response.signature.alg": "RS256", + "post.logout.redirect.uris": "http://localhost:3100/*##http://localhost:8171/*##http://localhost:8085/*", + "oauth2.device.authorization.grant.enabled": "false", + "use.jwks.url": "false", + "backchannel.logout.revoke.offline.tokens": "false", + "use.refresh.tokens": "true", + "tls-client-certificate-bound-access-tokens": "false", + "oidc.ciba.grant.enabled": "false", + "backchannel.logout.session.required": "true", + "client_credentials.use_refresh_token": "false", + "acr.loa.map": "{}", + "require.pushed.authorization.requests": "false", + "display.on.consent.screen": "false", + "client.session.idle.timeout": "-1", + "token.response.type.bearer.lower-case": "false" + }, + "authenticationFlowBindingOverrides": {}, + "fullScopeAllowed": true, + "nodeReRegistrationTimeout": -1, + "defaultClientScopes": [ + "web-origins", + "acr", + "roles", + "profile", + "email" + ], + "optionalClientScopes": [ + "address", + "phone", + "offline_access", + "microprofile-jwt" + ] + }, + { + "clientId": "deploycenter", + "name": "", + "description": "", + "rootUrl": "", + "adminUrl": "", + "baseUrl": "", + "surrogateAuthRequired": false, + "enabled": "true", + "alwaysDisplayInConsole": false, + "clientAuthenticatorType": "client-secret", + "secret": "ThisIsAnExampleKeyForDevPurposeOnly", + "redirectUris": [ + "http://localhost:3100/*", + "http://localhost:8171/*", + "http://localhost:8085/*" + ], + "webOrigins": [ + "http://localhost:3100", + "http://localhost:8171", + "http://localhost:8085" + ], + "notBefore": 0, + "bearerOnly": false, + "consentRequired": false, + "standardFlowEnabled": true, + "implicitFlowEnabled": false, + "directAccessGrantsEnabled": false, + "serviceAccountsEnabled": false, + "publicClient": false, + "frontchannelLogout": true, + "protocol": "openid-connect", + "attributes": { + "access.token.lifespan": "-1", + "client.secret.creation.time": "1707820779", + "user.info.response.signature.alg": "RS256", + "post.logout.redirect.uris": "http://localhost:3100/*##http://localhost:8171/*##http://localhost:8085/*", + "oauth2.device.authorization.grant.enabled": "false", + "use.jwks.url": "false", + "backchannel.logout.revoke.offline.tokens": "false", + "use.refresh.tokens": "true", + "tls-client-certificate-bound-access-tokens": "false", + "oidc.ciba.grant.enabled": "false", + "backchannel.logout.session.required": "true", + "client_credentials.use_refresh_token": "false", + "acr.loa.map": "{}", + "require.pushed.authorization.requests": "false", + "display.on.consent.screen": "false", + "client.session.idle.timeout": "-1", + "token.response.type.bearer.lower-case": "false" + }, + "authenticationFlowBindingOverrides": {}, + "fullScopeAllowed": true, + "nodeReRegistrationTimeout": -1, + "defaultClientScopes": [ + "web-origins", + "acr", + "roles", + "profile", + "email" + ], + "optionalClientScopes": [ + "address", + "phone", + "offline_access", + "microprofile-jwt" + ] } ], "clientScopes": [ diff --git a/docs/features/recording.md b/docs/features/recording.md index 13cc06cb..db2df3ca 100644 --- a/docs/features/recording.md +++ b/docs/features/recording.md @@ -126,6 +126,61 @@ RECORDING_STORAGE_EVENT_TOKEN = > Questions? Open an issue on [GitHub](https://github.com/suitenumerique/meet/issues/new?assignees=&labels=bug&template=Bug_report.md) or join our [Matrix community](https://matrix.to/#/#meet-official:matrix.org). +## Push recordings to Drive + +Once a recording is over, it can be pushed (streamed) to the user's [Drive](https://github.com/suitenumerique/drive), in addition to staying in the object storage. + +Drive is called as a resource server, following its +[resource server documentation](https://github.com/suitenumerique/drive/blob/main/docs/resource_server.md): + + +### Special requirements + +- Drive configured as an OIDC resource server, accepting Meet's audience + (`OIDC_RS_ALLOWED_AUDIENCES` must contain Meet's client id), with the `items` + endpoint allowing the `list`, `children` and `upload_ended` actions. +- `OIDC_STORE_ACCESS_TOKEN` enabled on Meet, along with + `OIDC_STORE_REFRESH_TOKEN_KEY`, the Fernet key encrypting the stored token. + +> [!CAUTION] +> This is a proof of concept: the access token is captured when the recording +> starts and assumed to still be valid when the recording ends. Long recordings +> may therefore fail to be pushed. Exchanging it for a long-lived, narrowly +> scoped token is the intended follow-up. + +### Configuration options + +| Option | Type | Default | Description | +| ----------------------------------------------------- | ----------- | ------- | -------------------------------------------------------------------------------------------------------------------------------------------------- | +| **RECORDING_PUSH_TO_DRIVE_ENABLED** | Boolean | `False` | Enable pushing recordings to the owner's Drive. | +| **DRIVE_API_BASE_URL** | String | `None` | Base URL of Drive's external API, e.g. `https://fichiers.numerique.gouv.fr/external_api/v1.0`. | +| **RECORDING_PUSH_TO_DRIVE_SIGNED_URL_EXPIRY_SECONDS** | Integer | `3600` | Lifetime of the signed URL the worker downloads the recording from. | +| **OIDC_STORE_ACCESS_TOKEN** | Boolean | `False` | Keep the user's access token in the session, required to call Drive on their behalf. | +| **OIDC_STORE_REFRESH_TOKEN_KEY** | Secret/File | `None` | Fernet key encrypting OIDC tokens at rest. Generate one with `Fernet.generate_key()`. | +| **DRIVE_UPLOAD_STORAGE_NETLOC** | String | `None` | Development only: `host:port` to reach Drive's object storage at, when the domain Drive signs its upload URLs with only resolves from a browser. | + +### Local development + +Meet and Drive run as two separate compose projects, joined by the external +`lasuite-network` (`make create-docker-network`). Meet's backend containers reach +Drive's nginx at `drive-nginx:8083` and its object storage at `drive-minio:9000`. + +On the Drive side: + +```bash +OIDC_RESOURCE_SERVER_ENABLED=True +OIDC_RS_CLIENT_ID=drive +OIDC_RS_CLIENT_SECRET=ThisIsAnExampleKeyForDevPurposeOnly +OIDC_RS_AUDIENCE_CLAIM=client_id +OIDC_RS_ALLOWED_AUDIENCES=meet +``` + +`DRIVE_UPLOAD_STORAGE_NETLOC` is needed because Drive signs its upload URLs +with `localhost:9100`, which does not resolve from Meet's containers. The +presigned signature covers the `Host` header, so the backend keeps announcing the +signed host and only swaps the address it connects to. + + ## LiveKit Egress La Suite Meet uses LiveKit Egress to record room sessions. For reference, see the [LiveKit Egress repository](https://github.com/livekit/egress) and the [official documentation](https://docs.livekit.io/home/egress/overview/). diff --git a/env.d/development/common.dist b/env.d/development/common.dist index a63f2ed4..23085e6d 100644 --- a/env.d/development/common.dist +++ b/env.d/development/common.dist @@ -24,7 +24,7 @@ MEET_BASE_URL="http://localhost:8072" # Media STORAGES_STATICFILES_BACKEND=django.contrib.staticfiles.storage.StaticFilesStorage AWS_S3_DOMAIN_REPLACE=http://localhost:9000 -AWS_S3_ENDPOINT_URL=http://minio:9000 +AWS_S3_ENDPOINT_URL=http://meet-minio:9000 AWS_S3_ACCESS_KEY_ID=meet AWS_S3_SECRET_ACCESS_KEY=password MEDIA_BASE_URL=http://localhost:3000 @@ -50,6 +50,10 @@ LOGOUT_REDIRECT_URL=http://localhost:3000 OIDC_REDIRECT_ALLOWED_HOSTS=localhost:8083,localhost:3000 OIDC_AUTH_REQUEST_EXTRA_PARAMS={"acr_values": "eidas1"} +# The key below encrypts OIDC tokens at rest. +OIDC_STORE_ACCESS_TOKEN=True +OIDC_STORE_REFRESH_TOKEN_KEY= + OIDC_RS_CLIENT_ID=meet OIDC_RS_CLIENT_SECRET=ThisIsAnExampleKeyForDevPurposeOnly @@ -70,6 +74,12 @@ SUMMARY_SERVICE_API_TOKEN=password SUMMARY_SERVICE_WEBHOOK_API_TOKEN=webhook-password RECORDING_DOWNLOAD_BASE_URL=http://localhost:3000/recording +# Push recordings to Drive +# DRIVE_UPLOAD_STORAGE_NETLOC is a development-only workaround. +RECORDING_PUSH_TO_DRIVE_ENABLED=True +DRIVE_API_BASE_URL=http://drive-app-dev:8000/external_api/v1.0 +DRIVE_UPLOAD_STORAGE_NETLOC=drive-minio:9000 + # Recording encoding (LiveKit Egress advanced options). # When RECORDING_ENCODING_ENABLED is False (default), LiveKit uses its built-in # H264_720P_30 preset (1280x720, 30fps, 3000 kbps). Enable and tune to reduce diff --git a/env.d/development/metadata_collector.dist b/env.d/development/metadata_collector.dist index ef6e2669..932cd1d4 100644 --- a/env.d/development/metadata_collector.dist +++ b/env.d/development/metadata_collector.dist @@ -2,7 +2,7 @@ LIVEKIT_URL=ws://livekit:7880 LIVEKIT_API_KEY=devkey LIVEKIT_API_SECRET=secret -AWS_S3_ENDPOINT_URL=minio:9000 +AWS_S3_ENDPOINT_URL=meet-minio:9000 AWS_S3_ACCESS_KEY_ID=meet AWS_S3_SECRET_ACCESS_KEY=password AWS_STORAGE_BUCKET_NAME=meet-media-storage diff --git a/env.d/development/summary.dist b/env.d/development/summary.dist index bfd56404..d1dc86f9 100644 --- a/env.d/development/summary.dist +++ b/env.d/development/summary.dist @@ -2,7 +2,7 @@ APP_NAME="meet-app-summary-dev" APP_API_TOKEN="password" AWS_STORAGE_BUCKET_NAME="meet-media-storage" -AWS_S3_ENDPOINT_URL="minio:9000" +AWS_S3_ENDPOINT_URL="meet-minio:9000" AWS_S3_SECURE_ACCESS=false AWS_S3_ACCESS_KEY_ID="meet" diff --git a/src/backend/core/api/viewsets.py b/src/backend/core/api/viewsets.py index 8a088ace..6eb044cd 100644 --- a/src/backend/core/api/viewsets.py +++ b/src/backend/core/api/viewsets.py @@ -383,6 +383,37 @@ class RoomViewSet( room.id, ) + @staticmethod + def _store_drive_credentials(request, recording): + """Keep the OIDC access token needed to push the recording to Drive later. + + Pushing happens long after this request, when the egress is over and the + user may be gone, so the token has to be stored. + + POC limitation: we assume the token is still valid by then. The target + design is a token exchange performed here, to get a long-lived + token narrowly scoped to that upload. + """ + + if not settings.RECORDING_PUSH_TO_DRIVE_ENABLED: + return + + if recording.mode != models.RecordingModeChoices.SCREEN_RECORDING: + # Only videos are pushed to Drive, no need for a token otherwise. + return + + access_token = request.session.get("oidc_access_token") + + if not access_token: + logger.warning( + "No OIDC access token in session, recording %s will not be pushed " + "to Drive. Is OIDC_STORE_ACCESS_TOKEN enabled?", + recording.id, + ) + return + + recording.set_owner_access_token(access_token) + @decorators.action( detail=True, methods=["post"], @@ -418,6 +449,7 @@ class RoomViewSet( role=models.RoleChoices.OWNER, recording=recording, ) + self._store_drive_credentials(request, recording) except (DjangoValidationError, IntegrityError): # DjangoValidationError covers the Python-level check (full_clean); diff --git a/src/backend/core/models.py b/src/backend/core/models.py index e8fbb316..d243ecf1 100644 --- a/src/backend/core/models.py +++ b/src/backend/core/models.py @@ -639,6 +639,17 @@ class Recording(BaseModel): verbose_name=_("External Process ID"), help_text=_("ID of the external process associated with the recording."), ) + owner_access_token = models.TextField( + null=True, + blank=True, + editable=False, + verbose_name=_("Owner access token"), + help_text=_( + "Encrypted OIDC access token of the user who started the recording, " + "used to push the recording to their Drive on their behalf. " + "Dropped as soon as the push has been attempted." + ), + ) class Meta: db_table = "meet_recording" @@ -741,6 +752,39 @@ class Recording(BaseModel): return self.expired_at < timezone.now() + def set_owner_access_token(self, access_token: str) -> None: + """Store the OIDC access token of the user who started the recording. + + It is stored encrypted, and only long enough for the worker to push the + recording to that user's Drive once the recording is over. + """ + + self.owner_access_token = utils.encrypt_secret(access_token) + self.save(update_fields=["owner_access_token", "updated_at"]) + + def get_owner_access_token(self) -> Optional[str]: + """Return the stored OIDC access token, or None if there is none left.""" + + if not self.owner_access_token: + return None + + try: + return utils.decrypt_secret(self.owner_access_token) + except utils.SecretDecryptionError: + logger.exception( + "Could not decrypt the access token of recording %s", self.id + ) + return None + + def clear_owner_access_token(self) -> None: + """Drop the stored access token, it is a user credential.""" + + if self.owner_access_token is None: + return + + self.owner_access_token = None + self.save(update_fields=["owner_access_token", "updated_at"]) + class RecordingAccess(BaseAccess): """Relation model to give access to a recording for a user or a team with a role.""" diff --git a/src/backend/core/recording/event/notification.py b/src/backend/core/recording/event/notification.py index 46bd1dc4..e5eea12e 100644 --- a/src/backend/core/recording/event/notification.py +++ b/src/backend/core/recording/event/notification.py @@ -19,6 +19,7 @@ from livekit import api as livekit_api from core import models, utils from core.analytics import UserFeatureFlag, is_user_feature_flag_enabled +from core.tasks.push_recording import push_recording from core.utils import generate_download_s3_url logger = logging.getLogger(__name__) @@ -45,7 +46,15 @@ class NotificationService: """Service for processing recordings and notifying external services.""" def notify_external_services(self, recording): - """Process a recording based on its mode.""" + """Process a recording, then push the video to the owner's Drive.""" + + try: + return self._notify_by_mode(recording) + finally: + self._push_recording_to_drive(recording) + + def _notify_by_mode(self, recording): + """Route a recording to the services its mode calls for.""" if recording.mode == models.RecordingModeChoices.TRANSCRIPT: return self._notify_summary_service(recording) @@ -222,6 +231,31 @@ class NotificationService: f"Unknown summary service version: {settings.SUMMARY_SERVICE_VERSION}" ) + @staticmethod + def _push_recording_to_drive(recording: models.Recording): + """Hand the recording over to the task pushing it to the owner's Drive.""" + + if not settings.RECORDING_PUSH_TO_DRIVE_ENABLED: + return + + if recording.mode != models.RecordingModeChoices.SCREEN_RECORDING: + recording.clear_owner_access_token() + return + + if not recording.owner_access_token: + logger.warning( + "No access token stored for recording %s, skipping the Drive push", + recording.id, + ) + return + + try: + push_recording.delay(str(recording.id)) + except Exception: # pylint: disable=broad-except + logger.exception( + "Could not schedule the Drive push of recording %s", recording.id + ) + @staticmethod def _notify_summary_service_v1(recording: models.Recording): """Notify summary service about a new recording.""" diff --git a/src/backend/core/recording/services/recording_events.py b/src/backend/core/recording/services/recording_events.py index 3b11315f..da60a126 100644 --- a/src/backend/core/recording/services/recording_events.py +++ b/src/backend/core/recording/services/recording_events.py @@ -107,4 +107,4 @@ class RecordingEventsService: if notification_succeeded else models.RecordingStatusChoices.SAVED ) - recording.save() + recording.save(update_fields=["status", "updated_at"]) diff --git a/src/backend/core/services/drive.py b/src/backend/core/services/drive.py new file mode 100644 index 00000000..adcc0d48 --- /dev/null +++ b/src/backend/core/services/drive.py @@ -0,0 +1,181 @@ +"""Client for La Suite Drive's external API (OIDC resource server). + +Drive exposes `/external_api/v1.0/*` to applications holding an user's OIDC +access token. +""" + +import logging +from urllib.parse import urlparse, urlunparse + +from django.conf import settings + +import requests + +logger = logging.getLogger(__name__) + +# (connect, read) timeouts, in seconds. The upload one covers a +# whole recording being relayed to Drive's object storage. +API_TIMEOUT = (10, 30) +UPLOAD_TIMEOUT = (10, 1800) + + +class DriveError(Exception): + """Raised when Drive's external API cannot fulfill a request.""" + + +class SizedStream: + """Read-only byte stream of a known size, suitable as a `requests` body. + + `requests` falls back to a chunked transfer encoding when it cannot guess the + body size upfront, which presigned S3 uploads reject. Advertising the size + through `__len__` makes it send a plain `Content-Length` instead, while the + underlying stream is still consumed chunk by chunk. + """ + + def __init__(self, stream, length: int): + """Wrap `stream`, whose full content is `length` bytes long.""" + self._stream = stream + self._length = length + + def __len__(self) -> int: + """Return the total size of the stream, in bytes.""" + return self._length + + def __iter__(self): + """Iterate over the stream, required for `requests` to stream the body.""" + return iter(self._stream) + + def read(self, amt=None) -> bytes: + """Read up to `amt` bytes from the stream.""" + return self._stream.read(amt) + + +class DriveClient: + """Access Drive's external API on behalf of a user. + + The client is bound to a single user access token: every call is performed + as that user, and Drive applies its own permissions accordingly. + """ + + def __init__(self, access_token: str, *, base_url: str | None = None): + """Prepare a session authenticated with the user's OIDC access token.""" + + self._base_url = (base_url or settings.DRIVE_API_BASE_URL or "").rstrip("/") + + if not self._base_url: + raise DriveError( + "Drive API is not configured, set DRIVE_API_BASE_URL to enable it." + ) + + if not access_token: + raise DriveError("An access token is required to call Drive.") + + self._session = requests.Session() + self._session.headers.update( + { + "Authorization": f"Bearer {access_token}", + "Content-Type": "application/json", + } + ) + + def __enter__(self): + """Allow use as a context manager, closing the session on exit.""" + return self + + def __exit__(self, *args): + """Close the underlying HTTP session.""" + self.close() + + def close(self): + """Release the underlying HTTP session.""" + self._session.close() + + def _request(self, method, path, **kwargs): + """Perform an authenticated call to the external API and return its body.""" + + url = f"{self._base_url}{path}" + kwargs.setdefault("timeout", API_TIMEOUT) + + try: + response = self._session.request(method, url, **kwargs) + response.raise_for_status() + except requests.RequestException as exc: + raise DriveError(f"Drive call failed: {method} {url}") from exc + + if not response.content: + return None + + try: + return response.json() + except ValueError as exc: + raise DriveError(f"Drive returned a non-JSON body for {url}") from exc + + def create_file(self, *, filename: str) -> dict: + """Create a file item at the root of the user's Drive and return it. + + The returned item carries a `policy`: the presigned URL the content has + to be uploaded to. + """ + + item = self._request( + "POST", + "/items/", + json={"type": "file", "filename": filename}, + ) + + if not item or not item.get("policy"): + raise DriveError( + f"Drive did not return an upload policy for file '{filename}'." + ) + + return item + + @staticmethod + def _resolve_upload_target(policy_url: str) -> tuple[str, str | None]: + """Return the address to connect to, and the `Host` header to send. + + Drive signs its upload URLs with the object storage domain meant for + browsers, which may not resolve from dev split docker compose setup. + The signature covers the `Host` header, so we swap the address we connect + to but keep announcing the original host. + """ + + override = settings.DRIVE_UPLOAD_STORAGE_NETLOC + + if not override: + return policy_url, None + + parsed = urlparse(policy_url) + return urlunparse(parsed._replace(netloc=override)), parsed.netloc + + def upload_content(self, *, policy_url: str, stream, content_length, content_type): + """Push `stream` to the presigned URL, without buffering it as a whole.""" + + url, host_header = self._resolve_upload_target(policy_url) + + headers = { + "Content-Type": content_type, + "Content-Length": str(content_length), + "x-amz-acl": "private", + } + + if host_header: + headers["Host"] = host_header + + try: + # A bare `requests.put`, not the authenticated session: the presigned + # URL carries its own credentials + response = requests.put( + url, + data=SizedStream(stream, content_length), + headers=headers, + timeout=UPLOAD_TIMEOUT, + ) + response.raise_for_status() + except requests.RequestException as exc: + raise DriveError("Upload to Drive's object storage failed.") from exc + + def complete_upload(self, item_id: str) -> None: + """Inform Drive that the upload is over, making the file available.""" + + self._request("POST", f"/items/{item_id}/upload-ended/", json={}) diff --git a/src/backend/core/tasks/__init__.py b/src/backend/core/tasks/__init__.py index e69de29b..1aeb2666 100644 --- a/src/backend/core/tasks/__init__.py +++ b/src/backend/core/tasks/__init__.py @@ -0,0 +1,11 @@ +"""Asynchronous tasks of the core application. + +Importing the task modules here is what makes Celery's `autodiscover_tasks` +register them: it only imports the `core.tasks` package itself, never its +submodules. +""" + +from core.tasks.file import process_file_deletion +from core.tasks.push_recording import push_recording + +__all__ = ["process_file_deletion", "push_recording"] diff --git a/src/backend/core/tasks/_task.py b/src/backend/core/tasks/_task.py index ede151be..4fe17b5a 100644 --- a/src/backend/core/tasks/_task.py +++ b/src/backend/core/tasks/_task.py @@ -4,6 +4,7 @@ Celery task decorator that degrades to a synchronous call when Celery is off. # The Celery app is imported lazily so that importing this module does not pull # in Celery when CELERY_ENABLED is false. + # ruff: noqa: PLC0415 # pylint: disable=import-outside-toplevel diff --git a/src/backend/core/tasks/push_recording.py b/src/backend/core/tasks/push_recording.py new file mode 100644 index 00000000..034bd406 --- /dev/null +++ b/src/backend/core/tasks/push_recording.py @@ -0,0 +1,110 @@ +"""Task pushing a finished recording to its owner's Drive.""" + +import logging + +from django.conf import settings + +import requests + +from core import models, utils +from core.services.drive import API_TIMEOUT, DriveClient, DriveError +from core.tasks._task import task + +logger = logging.getLogger(__name__) + +# (connect, read) timeouts for the download, in seconds. The read one has to +# accommodate a whole recording being relayed. +DOWNLOAD_TIMEOUT = (API_TIMEOUT[0], 1800) + + +def _build_filename(recording: models.Recording) -> str: + """Return a filename for the Drive item.""" + + return ( + f"{recording.room.slug}-" + f"{recording.created_at:%Y-%m-%d-%H-%M}." + f"{recording.extension}" + ) + + +@task +def push_recording(recording_id: str) -> bool: + """Push a recording to the Drive of the user who started it. + + The recording is streamed from object storage to Drive's presigned + URL. It is NOT fully downloaded to the worker's disk or memory. + + The access token stored when the recording started is consumed here and + dropped afterwards whatever the outcome. + + Mostly taken from: https://github.com/suitenumerique/drive/blob/main/docs/resource_server.md + """ + + try: + recording = models.Recording.objects.select_related("room").get(pk=recording_id) + except models.Recording.DoesNotExist: + logger.error( + "Recording %s does not exist, cannot push it to Drive", recording_id + ) + return False + + access_token = recording.get_owner_access_token() + + if not access_token: + logger.error( + "No access token stored for recording %s, cannot push it to Drive. " + "Was OIDC_STORE_ACCESS_TOKEN enabled when the recording started?", + recording_id, + ) + return False + + download_url = utils.generate_download_s3_url( + recording.key, + expires_in=settings.RECORDING_PUSH_TO_DRIVE_SIGNED_URL_EXPIRY_SECONDS, + override_domain=False, + ) + filename = _build_filename(recording) + + try: + with DriveClient(access_token) as drive: + item = drive.create_file(filename=filename) + + # The bytes are relayed chunk by chunk: the recording is never held + # in memory as a whole. + with requests.get( + download_url, stream=True, timeout=DOWNLOAD_TIMEOUT + ) as download: + download.raise_for_status() + + content_length = download.headers.get("Content-Length") + if content_length is None: + raise DriveError( + "Object storage did not return the recording size, " + "cannot stream it to Drive." + ) + + drive.upload_content( + policy_url=item["policy"], + stream=download.raw, + content_length=int(content_length), + content_type=download.headers.get( + "Content-Type", "application/octet-stream" + ), + ) + + drive.complete_upload(item["id"]) + + except (DriveError, requests.RequestException): + logger.exception("Failed to push recording %s to Drive", recording_id) + return False + + finally: + recording.clear_owner_access_token() + + logger.info( + "Recording %s pushed to Drive as '%s' (item %s)", + recording_id, + filename, + item["id"], + ) + return True diff --git a/src/backend/core/tests/files/test_api_files_create.py b/src/backend/core/tests/files/test_api_files_create.py index 3bfbae5a..6ec18d5e 100644 --- a/src/backend/core/tests/files/test_api_files_create.py +++ b/src/backend/core/tests/files/test_api_files_create.py @@ -117,7 +117,7 @@ def test_api_files_create_file_authenticated_success(): policy_parsed = urlparse(policy) assert policy_parsed.scheme == "http" - assert policy_parsed.netloc in ["minio:9000", "localhost:9000"] + assert policy_parsed.netloc in ["meet-minio:9000", "minio:9000", "localhost:9000"] assert policy_parsed.path == f"/meet-media-storage/tmp/files/{file.id!s}.png" query_params = parse_qs(policy_parsed.query) diff --git a/src/backend/core/utils.py b/src/backend/core/utils.py index b35b5bcb..0e76c3c2 100644 --- a/src/backend/core/utils.py +++ b/src/backend/core/utils.py @@ -17,6 +17,7 @@ from typing import List, Optional from uuid import uuid4 from django.conf import settings +from django.core.exceptions import ImproperlyConfigured from django.core.files.storage import default_storage import aiohttp @@ -25,6 +26,7 @@ import botocore import magic import phonenumbers from asgiref.sync import async_to_sync +from cryptography.fernet import Fernet, InvalidToken from livekit.api import ( # pylint: disable=E0611 AccessToken, ListRoomsRequest, @@ -417,6 +419,42 @@ def generate_upload_policy(file): return policy +class SecretDecryptionError(Exception): + """Raised when a stored secret cannot be decrypted.""" + + +@lru_cache(maxsize=1) +def get_cipher_suite(): + """Return the Fernet key used to encrypt secrets at rest. + + The same key as django-lasuite's OIDC token storage. + """ + + key = settings.OIDC_STORE_REFRESH_TOKEN_KEY + + if not key: + raise ImproperlyConfigured("OIDC_STORE_REFRESH_TOKEN_KEY setting is required.") + + return Fernet(key) + + +def encrypt_secret(value: str) -> str: + """Encrypt a secret meant to be stored at rest.""" + + return get_cipher_suite().encrypt(value.encode()).decode() + + +def decrypt_secret(value: str) -> str: + """Decrypt a secret stored by `encrypt_secret`.""" + + try: + return get_cipher_suite().decrypt(value.encode()).decode() + except InvalidToken as exc: + raise SecretDecryptionError( + "The stored secret could not be decrypted." + ) from exc + + def generate_download_s3_url( key: str, *, expires_in: int, override_domain: bool = True ): diff --git a/src/backend/meet/settings.py b/src/backend/meet/settings.py index 02ff0187..9b9b993b 100755 --- a/src/backend/meet/settings.py +++ b/src/backend/meet/settings.py @@ -573,6 +573,20 @@ class Base(Configuration): OIDC_STORE_ID_TOKEN = values.BooleanValue( default=True, environ_name="OIDC_STORE_ID_TOKEN", environ_prefix=None ) + # Required to call other La Suite applications on behalf of the user, e.g. + # to push a recording to their Drive. + OIDC_STORE_ACCESS_TOKEN = values.BooleanValue( + default=False, environ_name="OIDC_STORE_ACCESS_TOKEN", environ_prefix=None + ) + OIDC_STORE_REFRESH_TOKEN = values.BooleanValue( + default=False, environ_name="OIDC_STORE_REFRESH_TOKEN", environ_prefix=None + ) + # Fernet key used to encrypt OIDC tokens at rest, both the refresh token + # django-lasuite stores in the session and the access token stored on a + # recording. Generate one with `Fernet.generate_key()`. + OIDC_STORE_REFRESH_TOKEN_KEY = SecretFileValue( + None, environ_name="OIDC_STORE_REFRESH_TOKEN_KEY", environ_prefix=None + ) ALLOW_LOGOUT_GET_METHOD = values.BooleanValue( default=True, environ_name="ALLOW_LOGOUT_GET_METHOD", environ_prefix=None ) @@ -725,6 +739,29 @@ class Base(Configuration): None, environ_name="RECORDING_MAX_DURATION", environ_prefix=None ) + # Push recordings to Drive + # Once a recording is over, it is pushed to the user's + # started it, using their OIDC access token. It requires OIDC_STORE_ACCESS_TOKEN, + # and Drive to be configured as an OIDC resource server accepting Meet's audience. + RECORDING_PUSH_TO_DRIVE_ENABLED = values.BooleanValue( + False, environ_name="RECORDING_PUSH_TO_DRIVE_ENABLED", environ_prefix=None + ) + # Base URL of Drive's external API, e.g. https://fichiers.numerique.gouv.fr/external_api/v1.0 + DRIVE_API_BASE_URL = values.Value( + None, environ_name="DRIVE_API_BASE_URL", environ_prefix=None + ) + # Lifetime of the signed URL the worker downloads the recording from. + RECORDING_PUSH_TO_DRIVE_SIGNED_URL_EXPIRY_SECONDS = values.PositiveIntegerValue( + 60 * 60, + environ_name="RECORDING_PUSH_TO_DRIVE_SIGNED_URL_EXPIRY_SECONDS", + environ_prefix=None, + ) + # Development only: host:port to reach Drive's object storage at, when the + # domain Drive signs its upload URLs with is only resolvable from a browser. + DRIVE_UPLOAD_STORAGE_NETLOC = values.Value( + None, environ_name="DRIVE_UPLOAD_STORAGE_NETLOC", environ_prefix=None + ) + # Recording encoding options for LiveKit Egress (video composite egress only). # These settings affect screen recordings handled by VideoCompositeEgressService; # they are silently ignored by AudioCompositeEgressService (audio-only transcript