from __future__ import annotations import json import os import queue import shutil import subprocess import threading from datetime import datetime, timedelta, timezone from pathlib import Path from typing import Any from uuid import UUID, uuid4 from app.animations.store import AnimationStore from app.persistence import atomic_write_bytes from app.templates.store import TemplateConflictError, TemplateStore from .converter import ( MAX_MERGED_FRAMES, MediaConversionError, analyze, converted_frames, rgb_scene, validate_settings_for_metadata, ) JOB_SCHEMA_VERSION = 1 MAX_SOURCE_BYTES = 4 * 1024 * 1024 * 1024 MIN_FREE_BYTES = 2 * 1024 * 1024 * 1024 RETENTION_HOURS = 24 STATES = { "uploading", "analyzing", "awaiting_settings", "queued", "converting", "failed", } JOB_FIELDS = { "id", "filename", "name", "source_size", "state", "action", "created_at", "updated_at", "expires_at", "progress", "metadata", "settings", "preview_count", "error", "queue_sequence", } class MediaImportError(ValueError): def __init__(self, message: str, *, status_code: int = 400) -> None: super().__init__(message) self.status_code = status_code def _now() -> datetime: return datetime.now(timezone.utc) def _iso(value: datetime) -> str: return value.isoformat() def _job_id(value: Any) -> str: try: return str(UUID(str(value))) except (ValueError, TypeError, AttributeError) as exc: raise MediaImportError("media import job not found", status_code=404) from exc def _output_name(value: Any) -> str: if not isinstance(value, str): raise MediaImportError("output name must be a string") checked = value.strip() if not checked or len(checked) > 80: raise MediaImportError("output name must contain 1..80 characters") return checked def _default_output_name(filename: str) -> str: leaf = filename.replace("\\", "/").rsplit("/", 1)[-1].strip() stem = leaf.rsplit(".", 1)[0].strip() if "." in leaf else leaf return (stem[:80].strip() or "未命名媒体") class MediaImportManager: """Persistent single-worker media queue. The record is authoritative, so analyzing/queued/converting work can be restored after either service or device restart. """ def __init__(self, data_dir: Path, templates: TemplateStore, animations: AnimationStore) -> None: self.root = Path(data_dir) / "media-import" / "jobs" self.root.mkdir(parents=True, exist_ok=True) self.templates = templates self.animations = animations self._lock = threading.RLock() self._queue: queue.PriorityQueue[tuple[int, str, str]] = queue.PriorityQueue() self._thread: threading.Thread | None = None self._stop = threading.Event() self._canceled: set[str] = set() self._external_worker = os.environ.get("MATRIX_MEDIA_WORKER_MODE") == "systemd" self._sequence = 0 with self._lock: for child in self.root.iterdir(): if not child.is_dir(): continue record = self._read(child.name) self._sequence = max(self._sequence, record["queue_sequence"]) if record["state"] == "uploading": record.update({ "state": "failed", "action": "analyze", "error": "upload was interrupted", "progress": 0, "updated_at": _iso(_now()), "expires_at": _iso(_now() + timedelta(hours=RETENTION_HOURS)), }) self._write(record) elif record["state"] in {"analyzing", "queued", "converting"}: action = "analyze" if record["action"] == "analyze" else "convert" record["state"] = "analyzing" if action == "analyze" else "queued" record["progress"] = 0 self._write(record) self._queue.put((record["queue_sequence"], record["id"], action)) def _dir(self, job_id: str) -> Path: return self.root / job_id def _path(self, job_id: str) -> Path: return self._dir(job_id) / "job.json" def source_path(self, job_id: Any) -> Path: return self._dir(_job_id(job_id)) / "source.bin" def preview_path(self, job_id: Any, index: int) -> Path: normalized = _job_id(job_id) record = self._read(normalized) if index < 0 or index >= record["preview_count"]: raise MediaImportError("media preview not found", status_code=404) path = self._dir(normalized) / "previews" / f"{index}.png" if not path.is_file(): raise MediaImportError("media preview not found", status_code=404) return path def _write(self, record: dict[str, Any]) -> None: payload = json.dumps( {"schema_version": JOB_SCHEMA_VERSION, **record}, ensure_ascii=False, sort_keys=True, indent=2, ).encode("utf-8") + b"\n" atomic_write_bytes(self._path(record["id"]), payload) def _validate(self, value: Any, expected_id: str) -> dict[str, Any]: if not isinstance(value, dict): raise MediaImportError("media job record is invalid", status_code=500) document = dict(value) if document.pop("schema_version", None) != 1 or set(document) != JOB_FIELDS: raise MediaImportError("media job schema is unsupported", status_code=500) if document.get("id") != expected_id or document.get("state") not in STATES: raise MediaImportError("media job record is invalid", status_code=500) if document.get("action") not in {"analyze", "convert"}: raise MediaImportError("media job action is invalid", status_code=500) if not isinstance(document.get("settings"), dict): raise MediaImportError("media job settings are invalid", status_code=500) return document def _read(self, job_id: Any) -> dict[str, Any]: normalized = _job_id(job_id) path = self._path(normalized) if not path.is_file(): raise MediaImportError("media import job not found", status_code=404) try: value = json.loads(path.read_text(encoding="utf-8")) except (OSError, UnicodeError, json.JSONDecodeError) as exc: raise MediaImportError("media job record is unreadable", status_code=500) from exc return self._validate(value, normalized) def _public(self, record: dict[str, Any]) -> dict[str, Any]: queued = [] for _, queued_id, action in list(self._queue.queue): if action == "convert": queued.append(queued_id) queue_position = queued.index(record["id"]) + 1 if record["id"] in queued else None return { key: record[key] for key in ( "id", "filename", "name", "source_size", "state", "created_at", "updated_at", "expires_at", "progress", "metadata", "settings", "error", ) } | { "queue_position": queue_position, "previews": [ f"/api/media-imports/{record['id']}/previews/{index}" for index in range(record["preview_count"]) ], } def _reconcile_external(self, record: dict[str, Any]) -> dict[str, Any]: if not self._external_worker or record["state"] not in {"analyzing", "queued", "converting"}: return record result = subprocess.run( ["systemctl", "show", f"matrix-screen-converter@{record['id']}.service", "-p", "ActiveState", "-p", "Result", "--value"], stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, check=False, timeout=5, ) values = set(result.stdout.decode("utf-8", "replace").split()) if result.returncode == 0 and "failed" in values: now = _now() record.update({ "state": "failed", "progress": 0, "error": "isolated converter process exited unexpectedly", "updated_at": _iso(now), "expires_at": _iso(now + timedelta(hours=RETENTION_HOURS)), }) self._write(record) return record def _ensure_space(self, incoming: int = 0) -> None: free = shutil.disk_usage(self.root).free if free - incoming < MIN_FREE_BYTES: raise MediaImportError("device must retain at least 2 GiB free space", status_code=507) def begin_upload(self, filename: Any, name: Any | None, size: int) -> tuple[dict[str, Any], Path]: if not isinstance(filename, str) or not filename.strip() or len(filename) > 255: raise MediaImportError("filename must contain 1..255 characters") checked_name = _default_output_name(filename) if name is None else _output_name(name) if isinstance(size, bool) or not isinstance(size, int) or size <= 0: raise MediaImportError("Content-Length must be a positive integer", status_code=411) if size > MAX_SOURCE_BYTES: raise MediaImportError("source file exceeds 4 GiB", status_code=413) with self._lock: self._ensure_space(size) job_id, now = str(uuid4()), _now() self._sequence += 1 directory = self._dir(job_id) directory.mkdir(parents=True, exist_ok=False) record = { "id": job_id, "filename": filename.strip(), "name": checked_name, "source_size": size, "state": "uploading", "action": "analyze", "created_at": _iso(now), "updated_at": _iso(now), "expires_at": _iso(now + timedelta(hours=RETENTION_HOURS)), "progress": 0, "metadata": None, "settings": { "fit_mode": "crop", "center_x": 0.5, "center_y": 0.5, "zoom": 1.0, "transparency_color": "#000000", "padding_color": "#000000", }, "preview_count": 0, "error": None, "queue_sequence": self._sequence, } self._write(record) return self._public(record), directory / "source.bin" def finish_upload(self, job_id: Any, received: int) -> dict[str, Any]: with self._lock: record = self._read(job_id) if record["state"] != "uploading": raise MediaImportError("media upload is not active", status_code=409) if received != record["source_size"]: shutil.rmtree(self._dir(record["id"]), ignore_errors=True) raise MediaImportError("uploaded byte count does not match Content-Length") record.update({ "state": "analyzing", "action": "analyze", "progress": 0, "updated_at": _iso(_now()), "error": None, }) self._write(record) self._schedule(record["id"], "analyze", record["queue_sequence"]) return self._public(record) def fail_upload(self, job_id: Any) -> None: with self._lock: shutil.rmtree(self._dir(_job_id(job_id)), ignore_errors=True) def start(self) -> None: with self._lock: if self._external_worker: pending = [] while not self._queue.empty(): pending.append(self._queue.get_nowait()) for sequence, job_id, action in pending: self._schedule(job_id, action, sequence) return if self._thread and self._thread.is_alive(): return self._stop.clear() self._thread = threading.Thread(target=self._run, name="media-import-queue", daemon=True) self._thread.start() def stop(self) -> None: if self._external_worker: return self._stop.set() self._queue.put((2**63 - 1, "", "stop")) if self._thread: self._thread.join(timeout=10) def _schedule(self, job_id: str, action: str, sequence: int) -> None: if not self._external_worker: self._queue.put((sequence, job_id, action)) return result = subprocess.run( ["systemctl", "start", "--no-block", f"matrix-screen-converter@{job_id}.service"], stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.PIPE, check=False, timeout=15, ) if result.returncode: raise MediaImportError( "cannot start isolated converter: " + result.stderr.decode("utf-8", "replace")[-800:].strip(), status_code=503, ) def _run(self) -> None: while not self._stop.is_set(): _, job_id, action = self._queue.get() if not job_id or self._stop.is_set(): return with self._lock: if job_id in self._canceled or not self._path(job_id).is_file(): continue try: if action == "analyze": self._analyze(job_id) else: self._convert(job_id) except Exception as exc: with self._lock: if job_id in self._canceled or not self._path(job_id).is_file(): continue record = self._read(job_id) now = _now() record.update({ "state": "failed", "progress": 0, "error": str(exc)[:2000], "updated_at": _iso(now), "expires_at": _iso(now + timedelta(hours=RETENTION_HOURS)), }) self._write(record) def run_job(self, job_id: Any) -> None: """Run one persisted action; used by the isolated systemd instance.""" normalized = _job_id(job_id) record = self._read(normalized) try: if record["action"] == "analyze": self._analyze(normalized) else: self._convert(normalized) except Exception as exc: if not self._path(normalized).is_file(): return record = self._read(normalized) now = _now() record.update({ "state": "failed", "progress": 0, "error": str(exc)[:2000], "updated_at": _iso(now), "expires_at": _iso(now + timedelta(hours=RETENTION_HOURS)), }) self._write(record) raise def _analyze(self, job_id: str) -> None: metadata = analyze(self.source_path(job_id), self._dir(job_id) / "previews") with self._lock: if job_id in self._canceled: return record = self._read(job_id) record.update({ "state": "awaiting_settings", "metadata": metadata, "preview_count": metadata["preview_count"], "progress": 100, "updated_at": _iso(_now()), "error": None, }) self._write(record) def _library_name_conflict(self, name: str) -> bool: folded = name.casefold() if any(item["name"].casefold() == folded for item in self.templates.list()["templates"]): return True return not self.animations.name_available(name) def _reserved_name_conflict(self, name: str, *, exclude_job_id: str) -> bool: folded = name.casefold() for child in self.root.iterdir(): if not child.is_dir() or child.name == exclude_job_id: continue other = self._read(child.name) if other["state"] in {"queued", "converting"} and other["name"].casefold() == folded: return True return False def _name_conflict(self, record: dict[str, Any], *, include_reservations: bool) -> bool: if self._library_name_conflict(record["name"]): return True return include_reservations and self._reserved_name_conflict( record["name"], exclude_job_id=record["id"], ) def _convert(self, job_id: str) -> None: with self._lock: self._ensure_space() record = self._read(job_id) if self._name_conflict(record, include_reservations=False): raise TemplateConflictError("output name already exists") record.update({ "state": "converting", "action": "convert", "progress": 1, "updated_at": _iso(_now()), "error": None, }) self._write(record) frames = converted_frames( self.source_path(job_id), record["metadata"], record["settings"], ) if record["metadata"]["dynamic"]: def scenes(): for index, (rgb, duration) in enumerate(frames): with self._lock: if job_id in self._canceled: raise MediaConversionError("media conversion was canceled") current = self._read(job_id) current["progress"] = min(99, 1 + index * 98 // MAX_MERGED_FRAMES) self._write(current) yield { "scene": rgb_scene(rgb), "duration_ms": duration, "name": None, } self.animations.create_from_frames(record["name"], scenes()) else: rgb, _ = next(iter(frames)) with self._lock: if job_id in self._canceled: raise MediaConversionError("media conversion was canceled") self.templates.create(record["name"], rgb_scene(rgb)) with self._lock: if job_id not in self._canceled: shutil.rmtree(self._dir(job_id), ignore_errors=True) def list(self) -> dict[str, Any]: self.cleanup_expired() with self._lock: records = sorted( (self._reconcile_external(self._read(child.name)) for child in self.root.iterdir() if child.is_dir()), key=lambda record: record["created_at"], ) return {"jobs": [self._public(record) for record in records]} def get(self, job_id: Any) -> dict[str, Any]: self.cleanup_expired() with self._lock: return self._public(self._reconcile_external(self._read(job_id))) def update_settings(self, job_id: Any, value: Any) -> dict[str, Any]: if not isinstance(value, dict): raise MediaImportError("settings must be an object") allowed = { "name", "fit_mode", "center_x", "center_y", "zoom", "transparency_color", "padding_color", } if set(value) - allowed: raise MediaImportError("unknown media settings field") with self._lock: record = self._read(job_id) if record["state"] not in {"awaiting_settings", "failed"}: raise MediaImportError("media settings cannot be changed in the current state", status_code=409) settings = {**record["settings"], **{k: v for k, v in value.items() if k != "name"}} try: record["settings"] = validate_settings_for_metadata(settings, record["metadata"]) except MediaConversionError as exc: raise MediaImportError(str(exc)) from exc if "name" in value: proposed_name = _output_name(value["name"]) proposed = {**record, "name": proposed_name} if self._name_conflict(proposed, include_reservations=True): raise MediaImportError("output name already exists", status_code=409) record["name"] = proposed_name record.update({"updated_at": _iso(_now()), "error": None}) if record["state"] == "failed" and record["metadata"] is not None: record["state"] = "awaiting_settings" self._write(record) return self._public(record) def enqueue_conversion(self, job_id: Any) -> dict[str, Any]: with self._lock: record = self._read(job_id) if record["state"] != "awaiting_settings": raise MediaImportError("media job is not ready for conversion", status_code=409) if self._name_conflict(record, include_reservations=True): raise MediaImportError("output name already exists", status_code=409) self._sequence += 1 record.update({ "state": "queued", "action": "convert", "progress": 0, "queue_sequence": self._sequence, "updated_at": _iso(_now()), "error": None, }) self._write(record) self._schedule(record["id"], "convert", record["queue_sequence"]) return self._public(record) def retry(self, job_id: Any) -> dict[str, Any]: with self._lock: record = self._read(job_id) if record["state"] != "failed": raise MediaImportError("only failed jobs can be retried", status_code=409) action = record["action"] if action == "convert" and record["metadata"] is not None and self._name_conflict( record, include_reservations=True, ): raise MediaImportError("output name already exists", status_code=409) self._sequence += 1 record.update({ "state": "analyzing" if action == "analyze" else "queued", "queue_sequence": self._sequence, "progress": 0, "updated_at": _iso(_now()), "error": None, }) self._write(record) self._schedule(record["id"], action, record["queue_sequence"]) return self._public(record) def cancel(self, job_id: Any) -> None: normalized = _job_id(job_id) with self._lock: self._read(normalized) self._canceled.add(normalized) shutil.rmtree(self._dir(normalized), ignore_errors=True) def cleanup_expired(self) -> None: with self._lock: now = _now() for child in list(self.root.iterdir()): if not child.is_dir(): continue record = self._read(child.name) if record["state"] in {"failed", "awaiting_settings"}: try: expires = datetime.fromisoformat(record["expires_at"]) except (TypeError, ValueError): continue if expires <= now: self._canceled.add(record["id"]) shutil.rmtree(child, ignore_errors=True)