from __future__ import annotations import asyncio from pathlib import Path from app.core.config import Settings from app.core.exceptions import MediaAPIError, ProcessingError from app.core.logger import get_logger from app.projects.repositories.render_repository import ProjectRenderRepository from app.projects.services.render_compiler import compile_render from app.projects.services.render_service import ProjectRenderService from app.security.assets import CanonicalAssetNotFoundError from app.security.audit import AuditService from app.services.cleanup import CleanupService from app.services.ffmpeg_service import FFmpegService logger = get_logger(__name__) class ProjectRenderWorker: """Small cooperative worker; durable state remains the source of truth.""" def __init__( self, settings: Settings, repository: ProjectRenderRepository, renders: ProjectRenderService, ffmpeg: FFmpegService, cleanup: CleanupService, audit: AuditService, ) -> None: self.settings = settings self.repository = repository self.renders = renders self.ffmpeg = ffmpeg self.cleanup = cleanup self.audit = audit self._task: asyncio.Task[None] | None = None self._stopping = asyncio.Event() async def start(self) -> None: if not self.settings.render_worker_enabled or self._task is not None: return self._stopping.clear() self._task = asyncio.create_task(self._run(), name="project-render-worker") async def stop(self) -> None: self._stopping.set() if self._task is not None: self._task.cancel() await asyncio.gather(self._task, return_exceptions=True) self._task = None async def _run(self) -> None: while not self._stopping.is_set(): job = await self.repository.claim_next( stale_after_seconds=self.settings.render_job_stale_after_seconds ) if job is None: try: await asyncio.wait_for( self._stopping.wait(), self.settings.render_worker_interval_seconds ) except asyncio.TimeoutError: pass continue try: await self._execute(job) except Exception: logger.exception( "project render worker execution failed", extra={"render_id": job.id, "project_id": job.project_id}, ) async def _execute(self, job) -> None: workspace = await self.cleanup.create_workspace(job.id) staging = ( workspace.outputs / f"render-{job.id}.staging.{job.render_settings_json.get('format', 'mp4')}" ) cancel_event = asyncio.Event() render_task: asyncio.Task[None] | None = None try: await self.repository.ensure_project_active(job.id) paths = await self.renders.resolve_asset_paths(job) document = self.renders._document_from_job(job) settings = job.render_settings_json plan = compile_render( document, asset_paths=paths, width=int(settings["width"]), height=int(settings["height"]), frame_rate=float(settings["frameRate"]), output_format=str(settings["format"]), quality=str(settings["quality"]), preset=str(settings["preset"]), ) logger.info( "project render started", extra={ "operation": "render.started", "render_id": job.id, "project_id": job.project_id, "workspace_id": job.workspace_id, "attempt": job.attempt_count, }, ) render_task = asyncio.create_task( self.ffmpeg.run( [*plan.args, staging], operation="project.render", timeout=self.settings.render_job_timeout_seconds, cancel_event=cancel_event, ) ) while not render_task.done(): latest = await self.repository.heartbeat(job.id) if latest is not None and latest.status == "cancelling": cancel_event.set() break await asyncio.sleep(1) try: await render_task except ProcessingError as exc: if cancel_event.is_set(): await self.repository.cancel_worker(job.id) await self._audit(job, "project.render_cancelled", {}) return failed = await self.repository.fail( job.id, code="RENDER_PROCESSING_FAILED", message=str(exc), retryable=True, ) if failed is not None and failed.status == "cancelled": await self._audit(job, "project.render_cancelled", {}) elif failed is not None and failed.status == "failed": await self._audit( job, "project.render_failed", {"error_code": "RENDER_PROCESSING_FAILED"}, ) return latest = await self.repository.get_worker(job.id) if latest is not None and latest.status == "cancelling": await self.repository.cancel_worker(job.id) await self._audit(job, "project.render_cancelled", {}) return filename = f"render-{job.id}.{plan.output_extension}" published = await self.cleanup.publish_new(job.id, staging, filename) if published is None: # A prior attempt may have published before losing its worker # lease or database connection. Reconcile that deterministic # path instead of overwriting it or failing a safe retry. published = self.cleanup.resolve_download(job.id, filename) asset = await self.renders.assets.register_output( workspace_id=job.workspace_id, user_id=job.requested_by, request_id=job.id, path=published, mime_type="video/webm" if plan.output_extension == "webm" else "video/mp4", metadata={"project_render_id": job.id, "editor_revision": job.editor_revision}, project_id=job.project_id, ) completed = await self.repository.complete(job.id, asset.id) if completed.status == "cancelled": await self.renders.assets.discard_output( workspace_id=job.workspace_id, asset_id=asset.id, request_id=job.id, filename=published.name, ) await self.cleanup.remove_request(job.id) await self._audit(job, "project.render_cancelled", {}) return await self._audit(job, "project.render_completed", {"output_asset_id": asset.id}) except asyncio.CancelledError: cancel_event.set() if render_task is not None: await asyncio.gather(render_task, return_exceptions=True) stopped = await self.repository.fail( job.id, code="RENDER_WORKER_STOPPED", message="Render worker stopped before completion.", retryable=True, ) if stopped is not None and stopped.status == "cancelled": await self._audit(job, "project.render_cancelled", {}) raise except (MediaAPIError, CanonicalAssetNotFoundError, ValueError) as exc: code = getattr(exc, "code", "RENDER_INVALID") failed = await self.repository.fail( job.id, code=code, message=str(exc), retryable=False ) if failed is not None and failed.status == "cancelled": await self._audit(job, "project.render_cancelled", {}) else: await self._audit(job, "project.render_failed", {"error_code": code}) except Exception: logger.exception( "project render failed unexpectedly", extra={"render_id": job.id, "project_id": job.project_id}, ) failed = await self.repository.fail( job.id, code="RENDER_INTERNAL_ERROR", message="Render processing failed.", retryable=True, ) if failed is not None and failed.status == "cancelled": await self._audit(job, "project.render_cancelled", {}) elif failed is not None and failed.status == "failed": await self._audit( job, "project.render_failed", {"error_code": "RENDER_INTERNAL_ERROR"} ) finally: await self._remove(staging) await self.cleanup.remove_temporary_request(job.id) await self.cleanup.complete(job.id) async def _audit(self, job, event_type: str, metadata: dict[str, object]) -> None: await self.audit.record_event( workspace_id=job.workspace_id, user_id=job.requested_by, api_key_id=None, request_id=job.id, event_type=event_type, entity_type="project_render", entity_id=job.id, metadata={"project_id": job.project_id, **metadata}, ) logger.info( "project render lifecycle event", extra={ "operation": event_type, "render_id": job.id, "project_id": job.project_id, "workspace_id": job.workspace_id, "error_code": metadata.get("error_code"), "result": event_type.rsplit(".", 1)[-1], }, ) @staticmethod async def _remove(path: Path) -> None: try: await asyncio.to_thread(path.unlink, missing_ok=True) except OSError: pass