476 lines
18 KiB
Python
476 lines
18 KiB
Python
#!/usr/bin/env python3
|
|
"""Prebuild and serve a read-only MP4 catalog snapshot."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import re
|
|
import shutil
|
|
import subprocess
|
|
import threading
|
|
import uuid
|
|
from dataclasses import dataclass
|
|
from datetime import UTC, datetime
|
|
from http import HTTPStatus
|
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
|
from pathlib import Path
|
|
from urllib.parse import quote, urlsplit
|
|
|
|
|
|
LOG = logging.getLogger(__name__)
|
|
COVER_NAME = re.compile(r"^[0-9a-f]{64}\.jpg$")
|
|
DEFAULT_REFRESH_INTERVAL_SECONDS = 6 * 60 * 60
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ProbeResult:
|
|
duration: float
|
|
attached_picture_stream: int | None
|
|
|
|
|
|
def inspect_media(source: Path) -> ProbeResult | None:
|
|
"""Validate a source file against the HLS service's supported streams."""
|
|
command = [
|
|
"ffprobe",
|
|
"-v",
|
|
"error",
|
|
"-show_entries",
|
|
"format=duration:stream=index,codec_name,codec_type:stream_disposition=attached_pic",
|
|
"-of",
|
|
"json",
|
|
str(source),
|
|
]
|
|
try:
|
|
result = subprocess.run(command, capture_output=True, check=False, text=True, timeout=30)
|
|
except (OSError, subprocess.TimeoutExpired) as error:
|
|
LOG.warning("Unable to inspect %s: %s", source, error)
|
|
return None
|
|
|
|
if result.returncode != 0:
|
|
LOG.warning("Ignoring unreadable media %s: %s", source, result.stderr.strip())
|
|
return None
|
|
|
|
try:
|
|
payload = json.loads(result.stdout)
|
|
duration = float(payload["format"]["duration"])
|
|
except (KeyError, TypeError, ValueError, json.JSONDecodeError):
|
|
LOG.warning("Ignoring media with incomplete metadata: %s", source)
|
|
return None
|
|
|
|
video_count = 0
|
|
audio_count = 0
|
|
attached_picture_stream: int | None = None
|
|
for stream in payload.get("streams", []):
|
|
codec_type = stream.get("codec_type")
|
|
codec_name = stream.get("codec_name")
|
|
attached_picture = stream.get("disposition", {}).get("attached_pic") == 1
|
|
if codec_type == "video" and attached_picture:
|
|
attached_picture_stream = stream.get("index")
|
|
elif codec_type == "video" and codec_name == "h264":
|
|
video_count += 1
|
|
elif codec_type == "audio" and codec_name == "aac":
|
|
audio_count += 1
|
|
else:
|
|
LOG.warning("Ignoring unsupported media %s", source)
|
|
return None
|
|
|
|
if video_count != 1 or audio_count != 1 or duration < 0:
|
|
LOG.warning("Ignoring unsupported media %s", source)
|
|
return None
|
|
return ProbeResult(duration=duration, attached_picture_stream=attached_picture_stream)
|
|
|
|
|
|
def playlist_url(relative_path: Path) -> str:
|
|
encoded_path = "/".join(quote(part, safe="") for part in relative_path.parts)
|
|
return f"/hls/{encoded_path}/master.m3u8"
|
|
|
|
|
|
def cache_key(relative_path: Path, stat_result: os.stat_result) -> str:
|
|
source_fingerprint = f"{relative_path.as_posix()}:{stat_result.st_size}:{stat_result.st_mtime_ns}"
|
|
return hashlib.sha256(source_fingerprint.encode("utf-8")).hexdigest()
|
|
|
|
|
|
class Catalog:
|
|
def __init__(
|
|
self,
|
|
media_dir: Path,
|
|
cache_dir: Path,
|
|
refresh_interval: float = DEFAULT_REFRESH_INTERVAL_SECONDS,
|
|
auto_refresh: bool = True,
|
|
) -> None:
|
|
self.media_dir = media_dir.resolve()
|
|
self.cache_dir = cache_dir.resolve()
|
|
self.releases_dir = self.cache_dir / "releases"
|
|
self.current_link = self.cache_dir / "current"
|
|
self.releases_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
self._state_lock = threading.Lock()
|
|
self._refresh_event = threading.Event()
|
|
self._stop_event = threading.Event()
|
|
self._build_active = False
|
|
self._refresh_pending = False
|
|
self._current_release: Path | None = None
|
|
self._media_payload: bytes | None = None
|
|
self._status: dict[str, object] = {
|
|
"state": "building",
|
|
"lastSuccessAt": None,
|
|
"validFiles": 0,
|
|
"skippedFiles": 0,
|
|
"lastError": None,
|
|
}
|
|
self._load_current_snapshot()
|
|
self._refresh_interval = refresh_interval
|
|
self._worker: threading.Thread | None = None
|
|
if auto_refresh:
|
|
self._worker = threading.Thread(target=self._worker_loop, name="catalog-refresh", daemon=True)
|
|
self._worker.start()
|
|
self.request_rescan()
|
|
|
|
def ready(self) -> bool:
|
|
with self._state_lock:
|
|
return self._media_payload is not None and self._current_release is not None
|
|
|
|
def media_payload(self) -> bytes | None:
|
|
with self._state_lock:
|
|
return self._media_payload
|
|
|
|
def cover_payload(self, cover_name: str) -> bytes | None:
|
|
if not COVER_NAME.fullmatch(cover_name):
|
|
return None
|
|
with self._state_lock:
|
|
release = self._current_release
|
|
if release is None:
|
|
return None
|
|
try:
|
|
return (release / "covers" / cover_name).read_bytes()
|
|
except OSError:
|
|
return None
|
|
|
|
def status(self) -> dict[str, object]:
|
|
with self._state_lock:
|
|
return {"ready": self._media_payload is not None, **self._status}
|
|
|
|
def request_rescan(self) -> dict[str, object]:
|
|
with self._state_lock:
|
|
if not self._build_active and not self._refresh_pending:
|
|
self._refresh_pending = True
|
|
self._status["state"] = "building"
|
|
self._status["lastError"] = None
|
|
self._refresh_event.set()
|
|
return {"ready": self._media_payload is not None, **self._status}
|
|
|
|
def refresh_now(self) -> None:
|
|
"""Synchronously build a snapshot. Intended for tests and worker use."""
|
|
with self._state_lock:
|
|
self._build_active = True
|
|
self._refresh_pending = False
|
|
self._status["state"] = "building"
|
|
self._status["lastError"] = None
|
|
try:
|
|
release, payload, metadata = self._build_snapshot()
|
|
self._publish_snapshot(release, payload, metadata)
|
|
except Exception as error: # Keep the previous release available on a failed build.
|
|
LOG.exception("Catalog refresh failed")
|
|
with self._state_lock:
|
|
self._status["state"] = "ready" if self._media_payload is not None else "failed"
|
|
self._status["lastError"] = str(error)
|
|
finally:
|
|
with self._state_lock:
|
|
self._build_active = False
|
|
|
|
def _worker_loop(self) -> None:
|
|
while not self._stop_event.is_set():
|
|
requested = self._refresh_event.wait(self._refresh_interval)
|
|
if self._stop_event.is_set():
|
|
return
|
|
self._refresh_event.clear()
|
|
with self._state_lock:
|
|
if not requested and not self._refresh_pending:
|
|
self._refresh_pending = True
|
|
self._status["state"] = "building"
|
|
if self._build_active or not self._refresh_pending:
|
|
continue
|
|
self.refresh_now()
|
|
|
|
def _build_snapshot(self) -> tuple[Path, bytes, dict[str, object]]:
|
|
release_id = f"{datetime.now(UTC).strftime('%Y%m%dT%H%M%SZ')}-{uuid.uuid4().hex}"
|
|
staging = self.releases_dir / f".staging-{release_id}"
|
|
release = self.releases_dir / release_id
|
|
covers_dir = staging / "covers"
|
|
covers_dir.mkdir(parents=True)
|
|
valid_files = 0
|
|
skipped_files = 0
|
|
items: list[dict[str, str]] = []
|
|
|
|
try:
|
|
for source in self._media_files():
|
|
try:
|
|
relative_path = source.relative_to(self.media_dir)
|
|
before = source.stat()
|
|
except OSError:
|
|
skipped_files += 1
|
|
continue
|
|
|
|
probe = inspect_media(source)
|
|
if probe is None:
|
|
skipped_files += 1
|
|
continue
|
|
|
|
cover_name = f"{cache_key(relative_path, before)}.jpg"
|
|
if not self._generate_cover(source, covers_dir / cover_name, probe):
|
|
skipped_files += 1
|
|
continue
|
|
|
|
try:
|
|
after = source.stat()
|
|
except OSError:
|
|
skipped_files += 1
|
|
continue
|
|
if (before.st_size, before.st_mtime_ns) != (after.st_size, after.st_mtime_ns):
|
|
LOG.warning("Ignoring media changed while processing: %s", source)
|
|
(covers_dir / cover_name).unlink(missing_ok=True)
|
|
skipped_files += 1
|
|
continue
|
|
|
|
valid_files += 1
|
|
items.append(
|
|
{
|
|
"path": relative_path.as_posix(),
|
|
"name": source.stem,
|
|
"playlistUrl": playlist_url(relative_path),
|
|
"coverUrl": f"/api/media/covers/{cover_name}",
|
|
}
|
|
)
|
|
|
|
items.sort(key=lambda item: (item["name"].casefold(), item["path"].casefold()))
|
|
payload = json.dumps(items, ensure_ascii=False, separators=(",", ":")).encode("utf-8")
|
|
generated_at = datetime.now(UTC).isoformat().replace("+00:00", "Z")
|
|
metadata: dict[str, object] = {
|
|
"generatedAt": generated_at,
|
|
"validFiles": valid_files,
|
|
"skippedFiles": skipped_files,
|
|
}
|
|
(staging / "media.json").write_bytes(payload)
|
|
(staging / "metadata.json").write_text(json.dumps(metadata, separators=(",", ":")), encoding="utf-8")
|
|
staging.replace(release)
|
|
return release, payload, metadata
|
|
except Exception:
|
|
shutil.rmtree(staging, ignore_errors=True)
|
|
raise
|
|
|
|
def _publish_snapshot(self, release: Path, payload: bytes, metadata: dict[str, object]) -> None:
|
|
pending_link = self.cache_dir / f".current-{uuid.uuid4().hex}"
|
|
os.symlink(str(release.relative_to(self.cache_dir)), pending_link)
|
|
try:
|
|
os.replace(pending_link, self.current_link)
|
|
finally:
|
|
pending_link.unlink(missing_ok=True)
|
|
|
|
with self._state_lock:
|
|
self._current_release = release
|
|
self._media_payload = payload
|
|
self._status.update(
|
|
{
|
|
"state": "ready",
|
|
"lastSuccessAt": metadata["generatedAt"],
|
|
"validFiles": metadata["validFiles"],
|
|
"skippedFiles": metadata["skippedFiles"],
|
|
"lastError": None,
|
|
}
|
|
)
|
|
self._cleanup_old_releases(release)
|
|
|
|
def _load_current_snapshot(self) -> None:
|
|
try:
|
|
release = self.current_link.resolve(strict=True)
|
|
payload = (release / "media.json").read_bytes()
|
|
items = json.loads(payload)
|
|
metadata = json.loads((release / "metadata.json").read_text(encoding="utf-8"))
|
|
if not isinstance(items, list) or not isinstance(metadata, dict):
|
|
raise ValueError("invalid snapshot")
|
|
except (OSError, ValueError, json.JSONDecodeError) as error:
|
|
if self.current_link.exists() or self.current_link.is_symlink():
|
|
LOG.warning("Ignoring invalid catalog snapshot: %s", error)
|
|
return
|
|
|
|
self._current_release = release
|
|
self._media_payload = payload
|
|
self._status.update(
|
|
{
|
|
"state": "ready",
|
|
"lastSuccessAt": metadata.get("generatedAt"),
|
|
"validFiles": metadata.get("validFiles", 0),
|
|
"skippedFiles": metadata.get("skippedFiles", 0),
|
|
}
|
|
)
|
|
|
|
def _cleanup_old_releases(self, active_release: Path) -> None:
|
|
releases = sorted(
|
|
(path for path in self.releases_dir.iterdir() if path.is_dir() and not path.name.startswith(".staging-")),
|
|
key=lambda path: path.stat().st_mtime_ns,
|
|
reverse=True,
|
|
)
|
|
keep = {active_release, *releases[:2]}
|
|
for release in releases:
|
|
if release not in keep:
|
|
shutil.rmtree(release, ignore_errors=True)
|
|
|
|
def _generate_cover(self, source: Path, cover_path: Path, probe: ProbeResult) -> bool:
|
|
temporary_path = cover_path.with_name(f".{cover_path.stem}.{os.getpid()}.jpg")
|
|
if probe.attached_picture_stream is not None:
|
|
command = [
|
|
"ffmpeg",
|
|
"-nostdin",
|
|
"-v",
|
|
"error",
|
|
"-i",
|
|
str(source),
|
|
"-map",
|
|
f"0:{probe.attached_picture_stream}",
|
|
"-frames:v",
|
|
"1",
|
|
"-vf",
|
|
"scale=480:-2",
|
|
"-q:v",
|
|
"2",
|
|
"-y",
|
|
str(temporary_path),
|
|
]
|
|
else:
|
|
timestamp = min(10.0, probe.duration * 0.1)
|
|
command = [
|
|
"ffmpeg",
|
|
"-nostdin",
|
|
"-v",
|
|
"error",
|
|
"-ss",
|
|
f"{timestamp:.3f}",
|
|
"-i",
|
|
str(source),
|
|
"-frames:v",
|
|
"1",
|
|
"-vf",
|
|
"scale=480:-2",
|
|
"-q:v",
|
|
"2",
|
|
"-y",
|
|
str(temporary_path),
|
|
]
|
|
try:
|
|
result = subprocess.run(command, capture_output=True, check=False, text=True, timeout=60)
|
|
if result.returncode != 0 or not temporary_path.is_file() or temporary_path.stat().st_size == 0:
|
|
LOG.warning("Unable to generate cover for %s: %s", source, result.stderr.strip())
|
|
return False
|
|
temporary_path.replace(cover_path)
|
|
return True
|
|
except (OSError, subprocess.TimeoutExpired) as error:
|
|
LOG.warning("Unable to generate cover for %s: %s", source, error)
|
|
return False
|
|
finally:
|
|
temporary_path.unlink(missing_ok=True)
|
|
|
|
def _media_files(self) -> list[Path]:
|
|
files: list[Path] = []
|
|
for root, directories, filenames in os.walk(self.media_dir, followlinks=False):
|
|
directories[:] = [name for name in directories if not (Path(root) / name).is_symlink()]
|
|
for filename in filenames:
|
|
source = Path(root) / filename
|
|
if source.suffix.lower() == ".mp4" and source.is_file() and not source.is_symlink():
|
|
files.append(source)
|
|
return sorted(files, key=lambda path: path.relative_to(self.media_dir).as_posix().casefold())
|
|
|
|
|
|
class CatalogRequestHandler(BaseHTTPRequestHandler):
|
|
server: "CatalogHTTPServer"
|
|
|
|
def do_GET(self) -> None: # noqa: N802
|
|
request = urlsplit(self.path)
|
|
if request.path == "/healthz":
|
|
if self.server.catalog.ready():
|
|
self._send_bytes(HTTPStatus.OK, b"ok\n", "text/plain; charset=utf-8")
|
|
else:
|
|
self._send_json(HTTPStatus.SERVICE_UNAVAILABLE, {"error": "catalog is building"})
|
|
elif request.path == "/api/media/status":
|
|
self._send_json(HTTPStatus.OK, self.server.catalog.status())
|
|
elif request.path == "/api/media":
|
|
payload = self.server.catalog.media_payload()
|
|
if payload is None:
|
|
self._send_json(HTTPStatus.SERVICE_UNAVAILABLE, {"error": "catalog is building"})
|
|
else:
|
|
self._send_bytes(HTTPStatus.OK, payload, "application/json; charset=utf-8")
|
|
elif request.path.startswith("/api/media/covers/"):
|
|
if not self.server.catalog.ready():
|
|
self._send_json(HTTPStatus.SERVICE_UNAVAILABLE, {"error": "catalog is building"})
|
|
else:
|
|
self._send_cover(request.path.rsplit("/", 1)[-1])
|
|
else:
|
|
self.send_error(HTTPStatus.NOT_FOUND)
|
|
|
|
def do_HEAD(self) -> None: # noqa: N802
|
|
self.do_GET()
|
|
|
|
def do_POST(self) -> None: # noqa: N802
|
|
if urlsplit(self.path).path != "/api/media/rescan":
|
|
self.send_error(HTTPStatus.NOT_FOUND)
|
|
return
|
|
self._send_json(HTTPStatus.ACCEPTED, self.server.catalog.request_rescan())
|
|
|
|
def do_OPTIONS(self) -> None: # noqa: N802
|
|
self.send_response(HTTPStatus.NO_CONTENT)
|
|
self.send_header("Allow", "GET, HEAD, POST, OPTIONS")
|
|
self.end_headers()
|
|
|
|
def _send_cover(self, cover_name: str) -> None:
|
|
payload = self.server.catalog.cover_payload(cover_name)
|
|
if payload is None:
|
|
self.send_error(HTTPStatus.NOT_FOUND)
|
|
return
|
|
self._send_bytes(HTTPStatus.OK, payload, "image/jpeg")
|
|
|
|
def _send_json(self, status: HTTPStatus, payload: object) -> None:
|
|
self._send_bytes(
|
|
status,
|
|
json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8"),
|
|
"application/json; charset=utf-8",
|
|
)
|
|
|
|
def _send_bytes(self, status: HTTPStatus, body: bytes, content_type: str) -> None:
|
|
self.send_response(status)
|
|
self.send_header("Content-Type", content_type)
|
|
self.send_header("Content-Length", str(len(body)))
|
|
self.send_header("Cache-Control", "no-store")
|
|
self.end_headers()
|
|
if self.command != "HEAD":
|
|
self.wfile.write(body)
|
|
|
|
def log_message(self, format_string: str, *args: object) -> None:
|
|
LOG.info("%s - %s", self.address_string(), format_string % args)
|
|
|
|
|
|
class CatalogHTTPServer(ThreadingHTTPServer):
|
|
daemon_threads = True
|
|
|
|
def __init__(self, address: tuple[str, int], catalog: Catalog) -> None:
|
|
super().__init__(address, CatalogRequestHandler)
|
|
self.catalog = catalog
|
|
|
|
|
|
def main() -> None:
|
|
logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO"), format="%(asctime)s %(levelname)s %(message)s")
|
|
media_dir = Path(os.environ.get("MEDIA_DIR", "/media"))
|
|
cache_dir = Path(os.environ.get("CATALOG_CACHE_DIR", "/cache"))
|
|
refresh_interval = float(os.environ.get("CATALOG_REFRESH_INTERVAL_SECONDS", str(DEFAULT_REFRESH_INTERVAL_SECONDS)))
|
|
port = int(os.environ.get("CATALOG_PORT", "8081"))
|
|
if not media_dir.is_dir():
|
|
raise SystemExit(f"Media directory does not exist: {media_dir}")
|
|
if refresh_interval <= 0:
|
|
raise SystemExit("CATALOG_REFRESH_INTERVAL_SECONDS must be positive")
|
|
CatalogHTTPServer(("0.0.0.0", port), Catalog(media_dir, cache_dir, refresh_interval)).serve_forever()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|