Spaces:
Sleeping
Sleeping
| 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], | |
| }, | |
| ) | |
| async def _remove(path: Path) -> None: | |
| try: | |
| await asyncio.to_thread(path.unlink, missing_ok=True) | |
| except OSError: | |
| pass | |