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 @staticmethod 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." ) @staticmethod 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, )