MediaRouter / app /projects /workers /render_worker.py
basyx's picture
Upload 340 files
3493993 verified
Raw
History Blame Contribute Delete
10.5 kB
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