Spaces:
Sleeping
Sleeping
| from __future__ import annotations | |
| import hashlib | |
| import json | |
| from pathlib import Path | |
| from app.core.config import Settings | |
| from app.core.logger import get_logger | |
| from app.projects.editor_schemas import ( | |
| AudioClip, | |
| EditorDocument, | |
| MediaClip, | |
| ProjectRenderCreate, | |
| ProjectRenderResponse, | |
| ) | |
| from app.projects.errors import ( | |
| ProjectRenderConflictError, | |
| ProjectRenderInvalidError, | |
| ProjectRenderLimitError, | |
| ) | |
| from app.projects.models import ProjectRenderJob | |
| from app.projects.repositories.editor_repository import ProjectEditorRepository | |
| from app.projects.repositories.render_repository import ProjectRenderRepository | |
| from app.projects.services.render_compiler import validate_renderable_document | |
| from app.security.assets import CanonicalAssetNotFoundError, CanonicalAssetService | |
| from app.security.audit import AuditService | |
| from app.services.cleanup import CleanupService | |
| logger = get_logger(__name__) | |
| class ProjectRenderService: | |
| def __init__( | |
| self, | |
| settings: Settings, | |
| repository: ProjectRenderRepository, | |
| editor_repository: ProjectEditorRepository, | |
| assets: CanonicalAssetService, | |
| cleanup: CleanupService, | |
| audit: AuditService, | |
| ) -> None: | |
| self.settings = settings | |
| self.repository = repository | |
| self.editor_repository = editor_repository | |
| self.assets = assets | |
| self.cleanup = cleanup | |
| self.audit = audit | |
| def _document_from_job(job) -> EditorDocument: | |
| return EditorDocument.model_validate(job.editor_state_json) | |
| async def create( | |
| self, | |
| *, | |
| workspace_id: str, | |
| user_id: str, | |
| api_key_id: str, | |
| request_id: str, | |
| project_id: str, | |
| payload: ProjectRenderCreate, | |
| idempotency_key: str, | |
| ) -> ProjectRenderResponse: | |
| if len(idempotency_key.strip()) > 255 or not idempotency_key.strip(): | |
| raise ProjectRenderInvalidError("An Idempotency-Key header is required for rendering.") | |
| editor = await self.editor_repository.get(workspace_id, project_id, user_id=user_id) | |
| if editor.revision != payload.editor_revision: | |
| raise ProjectRenderInvalidError( | |
| "The requested editor revision is not current or available." | |
| ) | |
| document = EditorDocument.model_validate(editor.state_json) | |
| if document.project_id != project_id: | |
| raise ProjectRenderInvalidError("Saved editor state belongs to another project.") | |
| validate_renderable_document(document) | |
| if document.duration_ms() > self.settings.render_max_duration_seconds * 1000: | |
| raise ProjectRenderLimitError("The editor timeline exceeds the render duration limit.") | |
| if len(document.timeline.tracks) > self.settings.render_max_tracks: | |
| raise ProjectRenderLimitError("The editor timeline exceeds the render track limit.") | |
| if ( | |
| sum(len(track.clips) for track in document.timeline.tracks) | |
| > self.settings.render_max_clips | |
| ): | |
| raise ProjectRenderLimitError("The editor timeline exceeds the render clip limit.") | |
| if payload.width * payload.height > self.settings.max_resolution_pixels: | |
| raise ProjectRenderLimitError( | |
| "The requested render resolution exceeds the configured limit." | |
| ) | |
| await self._validate_assets(workspace_id, user_id, document) | |
| render_settings = { | |
| "format": payload.output_format, | |
| "width": payload.width, | |
| "height": payload.height, | |
| "frameRate": payload.frame_rate, | |
| "quality": payload.quality, | |
| "preset": payload.preset, | |
| } | |
| fingerprint = hashlib.sha256( | |
| json.dumps( | |
| { | |
| "revision": editor.revision, | |
| "state": document.model_dump(by_alias=True), | |
| "settings": render_settings, | |
| }, | |
| sort_keys=True, | |
| separators=(",", ":"), | |
| ).encode() | |
| ).hexdigest() | |
| existing = await self.repository.get_by_idempotency( | |
| workspace_id, | |
| project_id, | |
| editor.revision, | |
| idempotency_key.strip(), | |
| user_id=user_id, | |
| ) | |
| if existing is not None: | |
| if existing.request_fingerprint != fingerprint: | |
| raise ProjectRenderConflictError( | |
| "The idempotency key was already used for different render settings." | |
| ) | |
| return self._response(existing) | |
| job, created = await self.repository.create( | |
| workspace_id=workspace_id, | |
| project_id=project_id, | |
| user_id=user_id, | |
| editor_revision=editor.revision, | |
| document=document, | |
| render_settings=render_settings, | |
| request_fingerprint=fingerprint, | |
| idempotency_key=idempotency_key.strip(), | |
| max_attempts=max(1, self.settings.render_job_retry_limit + 1), | |
| max_active_jobs=self.settings.render_max_active_jobs_per_project, | |
| ) | |
| if created: | |
| await self.audit.record_event( | |
| workspace_id=workspace_id, | |
| user_id=user_id, | |
| api_key_id=api_key_id, | |
| request_id=request_id, | |
| event_type="project.render_requested", | |
| entity_type="project_render", | |
| entity_id=job.id, | |
| metadata={"project_id": project_id, "editor_revision": editor.revision}, | |
| ) | |
| logger.info( | |
| "project render requested", | |
| extra={ | |
| "operation": "render.create", | |
| "project_id": project_id, | |
| "render_id": job.id, | |
| "revision": editor.revision, | |
| }, | |
| ) | |
| return self._response(job) | |
| async def get( | |
| self, *, workspace_id: str, user_id: str, project_id: str, render_id: str | |
| ) -> ProjectRenderResponse: | |
| return self._response( | |
| await self.repository.get(workspace_id, project_id, render_id, user_id=user_id) | |
| ) | |
| async def list( | |
| self, *, workspace_id: str, user_id: str, project_id: str | |
| ) -> list[ProjectRenderResponse]: | |
| return [ | |
| self._response(item) | |
| for item in await self.repository.list(workspace_id, project_id, user_id=user_id) | |
| ] | |
| async def cancel( | |
| self, | |
| *, | |
| workspace_id: str, | |
| user_id: str, | |
| api_key_id: str, | |
| request_id: str, | |
| project_id: str, | |
| render_id: str, | |
| ) -> ProjectRenderResponse: | |
| job, changed = await self.repository.request_cancel( | |
| workspace_id, project_id, render_id, user_id=user_id | |
| ) | |
| # Queued jobs are cancelled immediately. Processing jobs are audited | |
| # only after the worker has actually stopped FFmpeg. | |
| if changed and job.status == "cancelled": | |
| await self.audit.record_event( | |
| workspace_id=workspace_id, | |
| user_id=user_id, | |
| api_key_id=api_key_id, | |
| request_id=request_id, | |
| event_type="project.render_cancelled", | |
| entity_type="project_render", | |
| entity_id=render_id, | |
| metadata={"project_id": project_id}, | |
| ) | |
| logger.info( | |
| "project render cancellation handled", | |
| extra={ | |
| "operation": "render.cancel", | |
| "project_id": project_id, | |
| "render_id": render_id, | |
| "status": job.status, | |
| "changed": changed, | |
| }, | |
| ) | |
| return self._response(job) | |
| async def resolve_asset_paths(self, job: ProjectRenderJob) -> dict[str, tuple[Path, str]]: | |
| document = EditorDocument.model_validate(job.editor_state_json) | |
| paths: dict[str, tuple[Path, str]] = {} | |
| for asset_id in document.asset_ids(): | |
| asset = await self.assets.get_owned_by_id( | |
| workspace_id=job.workspace_id, user_id=job.requested_by, asset_id=asset_id | |
| ) | |
| if asset.project_id != job.project_id: | |
| raise CanonicalAssetNotFoundError("Asset is no longer attached to this project.") | |
| path = self.cleanup.resolve_download(asset.request_id, asset.filename) | |
| await self.assets.verify_file(asset, path) | |
| paths[asset_id] = (path, asset.mime_type) | |
| return paths | |
| async def _validate_assets( | |
| self, workspace_id: str, user_id: str, document: EditorDocument | |
| ) -> None: | |
| total_bytes = 0 | |
| owned = {} | |
| for asset_id in document.asset_ids(): | |
| try: | |
| asset = await self.assets.get_owned_by_id( | |
| workspace_id=workspace_id, user_id=user_id, asset_id=asset_id | |
| ) | |
| if asset.project_id != document.project_id: | |
| raise CanonicalAssetNotFoundError("Asset is not attached to this project.") | |
| total_bytes += asset.file_size | |
| owned[asset_id] = asset | |
| except CanonicalAssetNotFoundError as exc: | |
| raise ProjectRenderInvalidError( | |
| "A referenced asset is not accessible in this workspace." | |
| ) from exc | |
| if total_bytes > self.settings.render_max_input_bytes: | |
| raise ProjectRenderLimitError( | |
| "The referenced render inputs exceed the configured byte limit." | |
| ) | |
| for track in document.timeline.tracks: | |
| for clip in track.clips: | |
| if not isinstance(clip, (MediaClip, AudioClip)): | |
| continue | |
| mime = owned[clip.asset_id].mime_type | |
| if isinstance(clip, AudioClip) and not mime.startswith("audio/"): | |
| raise ProjectRenderInvalidError("An audio clip must reference an audio asset.") | |
| if isinstance(clip, MediaClip) and not mime.startswith(f"{clip.media_type}/"): | |
| raise ProjectRenderInvalidError( | |
| f"A {clip.media_type} clip must reference a {clip.media_type} asset." | |
| ) | |
| def _response(job: ProjectRenderJob) -> ProjectRenderResponse: | |
| return ProjectRenderResponse( | |
| id=job.id, | |
| project_id=job.project_id, | |
| editor_revision=job.editor_revision, | |
| status=job.status, | |
| render_settings=dict(job.render_settings_json or {}), | |
| output_asset_id=job.output_asset_id, | |
| error_code=job.error_code, | |
| error_message=job.error_message, | |
| attempt_count=job.attempt_count, | |
| created_at=job.created_at, | |
| started_at=job.started_at, | |
| completed_at=job.completed_at, | |
| cancelled_at=job.cancelled_at, | |
| updated_at=job.updated_at, | |
| ) | |