525 lines
22 KiB
Python
525 lines
22 KiB
Python
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)
|