初始化奇妙小屏幕控制器项目
This commit is contained in:
@@ -0,0 +1,524 @@
|
||||
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)
|
||||
Reference in New Issue
Block a user