| from __future__ import annotations |
|
|
| import concurrent.futures |
| import hashlib |
| import hmac |
| import http.client |
| import json |
| import mimetypes |
| import shutil |
| import threading |
| import time |
| import urllib.request |
| from urllib.parse import urlparse |
| from dataclasses import asdict |
| from pathlib import Path |
| from typing import Callable |
|
|
| from renderer.core.config import Settings |
| from renderer.core.models import AIReelsRequest, JobRecord, RenderRequest |
| from renderer.core.render_engine import RenderEngine |
| from renderer.core.security import create_download_token |
| from renderer.core.utils import new_id, now, read_json, write_json |
|
|
|
|
| class JobManager: |
| def __init__(self, settings: Settings | None = None) -> None: |
| self.settings = settings or Settings() |
| self.settings.ensure_dirs() |
| self.executor = concurrent.futures.ThreadPoolExecutor(max_workers=self.settings.max_workers) |
| self.lock = threading.Lock() |
|
|
| def submit_render(self, request: RenderRequest) -> str: |
| return self._submit( |
| lambda job_id, log: RenderEngine(self.settings, log=log).render(request, job_id), |
| callback_url=request.callback_url, |
| export_target=request.export_target, |
| metadata=request.metadata, |
| job_type="render", |
| ) |
|
|
| def submit_ai_reels(self, request: AIReelsRequest) -> str: |
| return self._submit( |
| lambda job_id, log: RenderEngine(self.settings, log=log).ai_reels(request, job_id), |
| callback_url=request.callback_url, |
| export_target=request.export_target, |
| job_type="render", |
| ) |
|
|
| def submit_task( |
| self, |
| handler: Callable[[str, Callable[[str], None]], object], |
| *, |
| callback_url: str | None = None, |
| export_target: str | None = None, |
| metadata: dict | None = None, |
| job_type: str = "task", |
| ) -> str: |
| return self._submit(handler, callback_url=callback_url, export_target=export_target, metadata=metadata, job_type=job_type) |
|
|
| def submit_batch(self, requests: list[RenderRequest]) -> list[str]: |
| job_ids: list[str] = [] |
| for request in requests: |
| job_ids.append(self.submit_render(request)) |
| return job_ids |
|
|
| def get(self, job_id: str) -> JobRecord: |
| data = read_json(self._record_path(job_id), None) |
| if data is None: |
| raise KeyError(job_id) |
| return JobRecord(**data) |
|
|
| def cancel(self, job_id: str) -> JobRecord: |
| record = self.get(job_id) |
| if record.state == "PENDING": |
| self._update(job_id, state="CANCELLED", failure_reason="Cancelled before execution") |
| elif record.state == "RUNNING": |
| self._update(job_id, state="CANCEL_REQUESTED", failure_reason="Cancellation requested") |
| return self.get(job_id) |
|
|
| def cleanup(self, older_than_seconds: int | None = None) -> dict[str, int]: |
| cutoff = now() - (older_than_seconds or self.settings.job_retention_seconds) |
| removed_jobs = 0 |
| removed_exports = 0 |
| for path in self.settings.jobs_dir.glob("*.json"): |
| record = JobRecord(**read_json(path, {})) |
| if record.updated_at >= cutoff or record.state in {"PENDING", "RUNNING", "CANCEL_REQUESTED"}: |
| continue |
| if record.output_path: |
| output = Path(record.output_path) |
| if output.exists(): |
| output.unlink() |
| removed_exports += 1 |
| path.unlink(missing_ok=True) |
| removed_jobs += 1 |
| uploads = self.settings.temp_dir / "uploads" |
| if uploads.exists(): |
| shutil.rmtree(uploads, ignore_errors=True) |
| return {"removed_jobs": removed_jobs, "removed_exports": removed_exports} |
|
|
| def summary(self) -> dict: |
| records: list[JobRecord] = [] |
| for path in self.settings.jobs_dir.glob("*.json"): |
| try: |
| records.append(JobRecord(**read_json(path, {}))) |
| except Exception: |
| continue |
| state_counts: dict[str, int] = {} |
| for record in records: |
| state_counts[record.state] = state_counts.get(record.state, 0) + 1 |
| active = [ |
| { |
| "job_id": record.job_id, |
| "state": record.state, |
| "created_at": record.created_at, |
| "updated_at": record.updated_at, |
| "metrics": record.metrics, |
| } |
| for record in sorted(records, key=lambda item: item.updated_at, reverse=True) |
| if record.state in {"PENDING", "RUNNING", "CANCEL_REQUESTED"} |
| ] |
| return { |
| "total_jobs": len(records), |
| "state_counts": state_counts, |
| "active_jobs": active, |
| "max_workers": self.settings.max_workers, |
| } |
|
|
| def group(self, group_id: str) -> list[JobRecord]: |
| records: list[JobRecord] = [] |
| for path in self.settings.jobs_dir.glob("*.json"): |
| try: |
| record = JobRecord(**read_json(path, {})) |
| except Exception: |
| continue |
| if record.metadata.get("group_id") == group_id: |
| records.append(record) |
| return sorted(records, key=lambda item: item.created_at) |
|
|
| def _submit( |
| self, |
| handler: Callable[[str, Callable[[str], None]], object], |
| callback_url: str | None = None, |
| export_target: str | None = None, |
| metadata: dict | None = None, |
| job_type: str = "task", |
| ) -> str: |
| job_id = new_id() |
| record = JobRecord( |
| job_id=job_id, |
| state="PENDING", |
| created_at=now(), |
| updated_at=now(), |
| job_type=job_type, |
| download_token=create_download_token(self.settings.signing_secret, job_id), |
| callback_url=callback_url, |
| export_target=export_target, |
| metadata=metadata or {}, |
| ) |
| self._save(record) |
| self.executor.submit(self._run_with_retries, job_id, handler) |
| return job_id |
|
|
| def _run_with_retries(self, job_id: str, handler: Callable[[str, Callable[[str], None]], object]) -> None: |
| attempts = 0 |
| max_attempts = max(1, self.settings.max_retries) |
| while attempts < max_attempts: |
| attempts += 1 |
| if self.get(job_id).state == "CANCELLED": |
| self._send_callback(job_id) |
| return |
| self._update(job_id, state="RUNNING", metrics={"attempt": attempts}) |
| try: |
| if self.get(job_id).state == "CANCEL_REQUESTED": |
| self._update(job_id, state="CANCELLED", failure_reason="Cancelled before render started") |
| self._send_callback(job_id) |
| return |
| result = handler(job_id, lambda message: self.append_log(job_id, message)) |
| record = self.get(job_id) |
| output_path = getattr(result, "output_path", None) |
| export_path = self._export_copy(output_path, record.export_target, job_id) if output_path else None |
| result_metrics = getattr(result, "metrics", {}) or {} |
| if output_path: |
| result_metrics = result_metrics | {"artifact": _artifact_metadata(Path(output_path))} |
| self._update( |
| job_id, |
| state="COMPLETED", |
| output_path=str(output_path) if output_path else None, |
| export_path=export_path, |
| commands=getattr(result, "commands", []), |
| logs=getattr(result, "logs", []), |
| metrics=result_metrics | {"attempt": attempts}, |
| ) |
| self._send_callback(job_id) |
| return |
| except Exception as exc: |
| self.append_log(job_id, f"Attempt {attempts} failed: {exc}") |
| if attempts >= max_attempts: |
| self._update(job_id, state="FAILED", failure_reason=str(exc), metrics={"attempt": attempts}) |
| self._send_callback(job_id) |
|
|
| def append_log(self, job_id: str, message: str) -> None: |
| with self.lock: |
| record = self.get(job_id) |
| record.logs.append(message) |
| record.updated_at = now() |
| self._save(record) |
|
|
| def _update(self, job_id: str, **changes) -> None: |
| with self.lock: |
| record = self.get(job_id) |
| for key, value in changes.items(): |
| if key == "metrics" and record.metrics and isinstance(value, dict): |
| record.metrics.update(value) |
| else: |
| setattr(record, key, value) |
| record.updated_at = now() |
| self._save(record) |
|
|
| def _record_path(self, job_id: str) -> Path: |
| return self.settings.jobs_dir / f"{job_id}.json" |
|
|
| def _save(self, record: JobRecord) -> None: |
| write_json(self._record_path(record.job_id), asdict(record)) |
|
|
| def _export_copy(self, output_path: Path, export_target: str | None, job_id: str) -> str | None: |
| if not export_target: |
| return None |
| if export_target != "local": |
| if export_target.startswith(("http://", "https://")): |
| _put_file(export_target, output_path) |
| return export_target |
| return None |
| target_dir = self.settings.storage_dir / job_id |
| target_dir.mkdir(parents=True, exist_ok=True) |
| target = target_dir / output_path.name |
| shutil.copy2(output_path, target) |
| return str(target) |
|
|
| def _send_callback(self, job_id: str) -> None: |
| try: |
| record = self.get(job_id) |
| except KeyError: |
| return |
| if not record.callback_url: |
| return |
| event = _event_name(record.job_type, record.state) |
| data = asdict(record) |
| data.update( |
| { |
| "event": event, |
| "event_id": f"{job_id}:{record.state.lower()}:{int(record.updated_at)}", |
| "occurred_at": record.updated_at, |
| "status_url": self._public_url(f"/status/{job_id}"), |
| "download_url": self._public_url(f"/download/{job_id}?token={record.download_token}"), |
| } |
| ) |
| payload = json.dumps(data, default=str).encode("utf-8") |
| signature = hmac.new(self.settings.signing_secret.encode("utf-8"), payload, hashlib.sha256).hexdigest() |
| attempts = max(1, self.settings.callback_max_retries) |
| last_error: str | None = None |
| for attempt in range(1, attempts + 1): |
| request = urllib.request.Request( |
| record.callback_url, |
| data=payload, |
| headers={ |
| "Content-Type": "application/json", |
| "User-Agent": "ava2lon-studio-callback/2.0", |
| "X-Ava2lon-Event": event, |
| "X-Ava2lon-Signature": f"sha256={signature}", |
| }, |
| method="POST", |
| ) |
| try: |
| with urllib.request.urlopen(request, timeout=10) as response: |
| response.read() |
| if response.status >= 400: |
| raise RuntimeError(f"HTTP {response.status}") |
| self._update( |
| job_id, |
| callback_attempts=attempt, |
| callback_delivered_at=now(), |
| callback_last_error=None, |
| ) |
| return |
| except Exception as exc: |
| last_error = str(exc) |
| self._update(job_id, callback_attempts=attempt, callback_last_error=last_error) |
| if attempt < attempts: |
| time.sleep(max(0.0, self.settings.callback_retry_seconds) * (2 ** (attempt - 1))) |
| self.append_log(job_id, f"Callback delivery failed after {attempts} attempts: {last_error}") |
|
|
| def _public_url(self, path: str) -> str: |
| return f"{self.settings.public_base_url}{path}" if self.settings.public_base_url else path |
|
|
|
|
| def _put_file(url: str, path: Path) -> None: |
| parsed = urlparse(url) |
| connection_cls = http.client.HTTPSConnection if parsed.scheme == "https" else http.client.HTTPConnection |
| connection = connection_cls(parsed.netloc, timeout=60) |
| target = parsed.path or "/" |
| if parsed.query: |
| target += f"?{parsed.query}" |
| headers = { |
| "Content-Type": "video/mp4", |
| "Content-Length": str(path.stat().st_size), |
| "User-Agent": "ava2lon-studio-export/2.0", |
| } |
| connection.putrequest("PUT", target) |
| for key, value in headers.items(): |
| connection.putheader(key, value) |
| connection.endheaders() |
| with path.open("rb") as source: |
| while True: |
| chunk = source.read(1024 * 1024) |
| if not chunk: |
| break |
| connection.send(chunk) |
| response = connection.getresponse() |
| body = response.read() |
| connection.close() |
| if response.status >= 400: |
| raise RuntimeError(f"Export upload failed with HTTP {response.status}: {body[:500]!r}") |
|
|
|
|
| def _event_name(job_type: str, state: str) -> str: |
| return f"{job_type}.{state.lower()}" |
|
|
|
|
| def _artifact_metadata(path: Path) -> dict: |
| digest = hashlib.sha256() |
| with path.open("rb") as source: |
| for chunk in iter(lambda: source.read(1024 * 1024), b""): |
| digest.update(chunk) |
| return { |
| "filename": path.name, |
| "size_bytes": path.stat().st_size, |
| "mime_type": mimetypes.guess_type(path.name)[0] or "application/octet-stream", |
| "sha256": digest.hexdigest(), |
| } |
|
|