This commit is contained in:
root
2026-09-27 15:55:38 +08:00
commit 1b6f2e6831
14 changed files with 1379 additions and 0 deletions
+475
View File
@@ -0,0 +1,475 @@
#!/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()