| from __future__ import annotations |
|
|
| import concurrent.futures |
| import http.client |
| import json |
| import shutil |
| import threading |
| 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, |
| ) |
|
|
| 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)) |
|
|
| def submit_task( |
| self, |
| handler: Callable[[str, Callable[[str], None]], object], |
| *, |
| callback_url: str | None = None, |
| export_target: str | None = None, |
| ) -> str: |
| return self._submit(handler, callback_url=callback_url, export_target=export_target) |
|
|
| 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 _submit( |
| self, |
| handler: Callable[[str, Callable[[str], None]], object], |
| callback_url: str | None = None, |
| export_target: str | None = None, |
| ) -> str: |
| job_id = new_id() |
| record = JobRecord( |
| job_id=job_id, |
| state="PENDING", |
| created_at=now(), |
| updated_at=now(), |
| download_token=create_download_token(self.settings.signing_secret, job_id), |
| callback_url=callback_url, |
| export_target=export_target, |
| ) |
| 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 |
| while attempts < self.settings.max_retries: |
| 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 |
| 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=getattr(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 >= self.settings.max_retries: |
| 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 |
| payload = json.dumps(asdict(record), default=str).encode("utf-8") |
| request = urllib.request.Request( |
| record.callback_url, |
| data=payload, |
| headers={"Content-Type": "application/json", "User-Agent": "ava2lon-studio-callback/2.0"}, |
| method="POST", |
| ) |
| try: |
| urllib.request.urlopen(request, timeout=10).read() |
| except Exception as exc: |
| self.append_log(job_id, f"Callback delivery failed: {exc}") |
|
|
|
|
| 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}") |
|
|