MediaRouter / app /projects /services /render_service.py
basyx's picture
Upload 340 files
3493993 verified
Raw
History Blame Contribute Delete
11 kB
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,
)