Found by running the round-8 fix in the real container instead of trusting the unit test: a directory created 2ms AFTER `_PROCESS_START` was stamped with an mtime 2ms BEFORE it, so the sweep deleted it as a leftover. The two sides come from different clocks — a file's mtime from the kernel's coarse clock (one timer tick of lag), `time.time()` from the fine one — and the unit test could not see it because it supplies an artificial cutoff. The window is small and lands exactly where it hurts: the sweep runs at boot, so the session at risk is one started in the first instants after startup — the case the guard exists for. Keep-side slack of 2s; the cost is that a leftover from the last seconds before a restart survives until the next one. Verified in the container: a backdated leftover is swept, a directory created by this run is kept, and a registered live session with an old mtime is kept — 1 of 3. Test pins the skew (mutation-checked: with the grace at 0 it goes red).
341 lines
16 KiB
Python
341 lines
16 KiB
Python
"""On-the-fly HLS remux for Plex playback (seek-restart session model).
|
|
|
|
Playback is from the LOCAL physical file. Browser-compatible files (playable="direct") are served
|
|
raw with HTTP range requests. Everything else that is h264 video (playable="remux") is remuxed on
|
|
the fly to HLS with **video stream-copy** (cheap, I/O-bound — no video re-encode) and audio to AAC
|
|
when needed. HEVC/VP9 (playable="transcode") needs a full re-encode and is deferred to P3.
|
|
|
|
Seek model (Jellyfin-style): one ffmpeg session per item, started at a given offset via `-ss`
|
|
(fast keyframe seek). A seek beyond the generated region restarts the session at the new offset,
|
|
so seeking is responsive even on long movies without pre-generating the whole file. ffmpeg's own
|
|
HLS muxer does the segmentation (the only reliable way to cut a stream-copy at keyframes).
|
|
|
|
Sessions live in this (single) API process; a periodic reaper kills idle ones and frees the temp
|
|
segments. The remux is CPU-light (video copy + a tiny audio transcode), so it runs fine on the
|
|
CPU-only prod host; only P3's full transcode is CPU-heavy.
|
|
"""
|
|
import logging
|
|
import re
|
|
import shutil
|
|
import subprocess
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
|
|
from sqlalchemy.orm import Session as DbSession
|
|
|
|
from app import sysconfig
|
|
from app.config import settings
|
|
from app.models import PlexItem
|
|
from app.plex import paths
|
|
|
|
log = logging.getLogger("siftlode.plex")
|
|
|
|
_HLS_ROOT = Path(settings.plex_hls_dir) if settings.plex_hls_dir else Path(settings.download_root) / ".plex-hls"
|
|
_SEG_SECONDS = 6
|
|
_SESSION_IDLE_S = 600 # reap a session with no access for this long
|
|
# A session directory is named `{rating_key}_{int(start_s)}_{tag}_{aoff:+.2f}` (see start_session);
|
|
# before multi-audio it was just `{rating_key}_{int(start_s)}`. `sweep_orphans` deletes ONLY names
|
|
# matching one of those two shapes: `_HLS_ROOT` is admin-configurable (PLEX_HLS_DIR), so it may
|
|
# legitimately point at a shared scratch path we do not own, and a blanket "delete every directory
|
|
# under the root" would take the neighbours with it.
|
|
#
|
|
# The KEY is anchored to digits (a Plex rating_key is a number) — that is what keeps foreign names
|
|
# out. A `.+` there let `snapshot_2024_backup_+1.00` through, and pairing `.+` with an optional tail
|
|
# degenerated all the way to "anything ending in _<digits>" (`backup_2024`, `release_10`). A
|
|
# non-numeric rating_key would merely leave its directory unswept, which is the safe direction to
|
|
# fail. Covered by test_plex_stream.py — this regex is the only thing standing between the sweep
|
|
# and someone else's data.
|
|
_SESSION_DIR_RE = re.compile(r"^\d+_\d+(_[^_]+_[+-]\d+\.\d{2})?$")
|
|
# When this process started. `sweep_orphans` only deletes directories older than this: everything
|
|
# newer belongs to a session THIS run created, and since `start_session` reuses a directory path
|
|
# for the same (rating_key, offset, tag), "stale name" is not the same as "stale directory".
|
|
_PROCESS_START = time.time()
|
|
# Slack on that comparison, because the two sides come from DIFFERENT clocks. A file's mtime is
|
|
# stamped from the kernel's COARSE clock (updated once per timer tick), while `time.time()` reads
|
|
# the fine-grained one — measured in this very container, a directory created 2ms AFTER
|
|
# `_PROCESS_START` got an mtime 2ms BEFORE it. Without slack the sweep would classify a session
|
|
# started in the first instants after boot as a leftover and delete it while ffmpeg writes into it
|
|
# — and "the first instants after boot" is exactly when the sweep runs. Erring on the keep side
|
|
# costs nothing: a genuine leftover from the last seconds before the restart simply survives until
|
|
# the next one.
|
|
_MTIME_GRACE_S = 2.0
|
|
|
|
_lock = threading.Lock()
|
|
_sessions: dict[str, "HlsSession"] = {}
|
|
|
|
|
|
class HlsSession:
|
|
def __init__(self, key: str, directory: Path, proc: subprocess.Popen, start_s: float, entry: str):
|
|
self.key = key
|
|
self.dir = directory
|
|
self.proc = proc
|
|
self.start_s = start_s # the REQUESTED seek offset
|
|
# The REAL media start: with `-ss X -c:v copy` ffmpeg can only cut at the nearest earlier
|
|
# keyframe K<=X, and hls.js zero-bases the timeline to K (so video.currentTime=0 is K, not X).
|
|
# We measure K from the first segment's absolute PTS (kept absolute by `-copyts`) and report
|
|
# THAT as the session start, so the client's absolute clock (start + currentTime) and the
|
|
# subtitle cue-shift both key off the true content position — no drift. Provisional = X until
|
|
# the first segment is probed.
|
|
self.media_start_s = start_s
|
|
self.entry = entry # the playlist filename hls.js should load (master.m3u8 with subs, else index.m3u8)
|
|
self.last_access = time.time()
|
|
|
|
|
|
def _seek_opts(start_s: float) -> list[str]:
|
|
# -noaccurate_seek: with `-ss` before -i, video (stream-copy) backs up to the keyframe K<=X,
|
|
# but the re-encoded audio is normally trimmed to the exact X — so audio starts (X-K) AFTER
|
|
# video → seconds of silence at the start of every seek/audio-switch. noaccurate_seek makes
|
|
# audio ALSO start at K, aligned with the video (K is our real clock anchor anyway).
|
|
return ["-noaccurate_seek", "-ss", f"{start_s:.3f}"] if start_s > 0 else []
|
|
|
|
|
|
_HLS_TAIL = [
|
|
"-f", "hls",
|
|
"-hls_time", str(_SEG_SECONDS),
|
|
"-hls_list_size", "0",
|
|
# EVENT (not VOD): ffmpeg writes/appends the playlist AS segments complete, so playback can
|
|
# start immediately from this session's offset. VOD only writes the playlist at the end. The
|
|
# full seekbar comes from our known duration + the seek-restart model, not from the playlist.
|
|
"-hls_playlist_type", "event",
|
|
"-hls_segment_type", "mpegts",
|
|
"-hls_flags", "independent_segments+temp_file",
|
|
]
|
|
|
|
|
|
def _probe_audio_langs(src: Path) -> list[str]:
|
|
"""One entry per audio stream (its language tag, or "" if untagged). len() = audio-track count.
|
|
Used to build the multi-rendition var_stream_map. Empty on probe failure → single-audio path."""
|
|
try:
|
|
out = subprocess.run(
|
|
["ffprobe", "-v", "error", "-select_streams", "a", "-show_entries", "stream_tags=language",
|
|
"-of", "csv=p=0", str(src)],
|
|
capture_output=True, text=True, timeout=10,
|
|
)
|
|
return [ln.strip() for ln in out.stdout.splitlines()]
|
|
except (subprocess.SubprocessError, OSError):
|
|
return []
|
|
|
|
|
|
def _ffmpeg_cmd(
|
|
src: Path, start_s: float, out_dir: Path, audio_ord: int | None,
|
|
aoff: float = 0.0, audio_langs: list[str] | None = None,
|
|
) -> list[str]:
|
|
args = ["ffmpeg", "-nostdin", "-loglevel", "error"]
|
|
# Keep ORIGINAL timestamps on the output so the first segment carries the true absolute PTS of the
|
|
# keyframe ffmpeg actually seeked to. hls.js still zero-bases playback to that PTS; we read it back
|
|
# (see _probe_first_pts) to learn the real content offset. Without -copyts the segment PTS is
|
|
# rewritten toward 0 and the true offset is unknowable, which caused the ~seconds subtitle lead.
|
|
args += _seek_opts(start_s) + ["-copyts", "-i", str(src)]
|
|
ai = 0 # input index the audio maps come from
|
|
if aoff:
|
|
# A/V-sync offset: read the file a SECOND time for audio only, with `-itsoffset` shifting the
|
|
# audio PTS by ±aoff relative to the (input-0) video. Maps audio from input 1 then.
|
|
args += ["-itsoffset", f"{aoff:.3f}"] + _seek_opts(start_s) + ["-copyts", "-i", str(src)]
|
|
ai = 1
|
|
# Video always stream-copied; audio → AAC. Subtitles are NOT muxed here — they're separate WebVTT
|
|
# tracks (GET /subtitle) the browser overlays, so choosing a subtitle doesn't restart this session.
|
|
if audio_langs:
|
|
# MULTI-RENDITION: map EVERY audio track as an alternate HLS rendition in one group, so hls.js
|
|
# switches audio CLIENT-SIDE on the same timeline (no session restart, no drift) — the premium
|
|
# audio-switch path. Produces master.m3u8 + stream_%v.m3u8 (v0=video, v1..=audio) + seg_%v_%d.ts.
|
|
args += ["-map", "0:v:0"]
|
|
for i in range(len(audio_langs)):
|
|
args += ["-map", f"{ai}:a:{i}?"]
|
|
args += ["-sn", "-c:v", "copy", "-c:a", "aac", "-ac", "2", "-b:a", "192k"]
|
|
var = ["v:0,agroup:aud"]
|
|
for i, lang in enumerate(audio_langs):
|
|
entry = f"a:{i},agroup:aud,name:a{i}"
|
|
if lang:
|
|
entry += f",language:{lang}"
|
|
if i == 0:
|
|
entry += ",default:yes"
|
|
var.append(entry)
|
|
args += _HLS_TAIL + [
|
|
"-master_pl_name", "master.m3u8",
|
|
"-var_stream_map", " ".join(var),
|
|
"-hls_segment_filename", str(out_dir / "seg_%v_%d.ts"), str(out_dir / "stream_%v.m3u8"),
|
|
]
|
|
else:
|
|
ao = audio_ord if audio_ord is not None else 0
|
|
args += ["-map", "0:v:0", "-map", f"{ai}:a:{ao}?", "-sn"]
|
|
args += ["-c:v", "copy", "-c:a", "aac", "-ac", "2", "-b:a", "192k"]
|
|
args += _HLS_TAIL + ["-hls_segment_filename", str(out_dir / "seg_%d.ts"), str(out_dir / "index.m3u8")]
|
|
return args
|
|
|
|
|
|
def _probe_first_pts(seg: Path) -> float | None:
|
|
"""Read the first video packet's absolute PTS from a finished HLS segment (`-copyts` kept it
|
|
absolute). This is the true content position of the segment's first frame — the value hls.js
|
|
zero-bases the timeline to. Returns None on any failure (caller falls back to the requested offset)."""
|
|
try:
|
|
out = subprocess.run(
|
|
["ffprobe", "-v", "error", "-select_streams", "v:0", "-read_intervals", "%+#1",
|
|
"-show_entries", "packet=pts_time", "-of", "csv=p=0", str(seg)],
|
|
capture_output=True, text=True, timeout=10,
|
|
)
|
|
for line in out.stdout.splitlines():
|
|
tok = line.strip().strip(",")
|
|
if tok:
|
|
return float(tok)
|
|
except (subprocess.SubprocessError, ValueError, OSError):
|
|
pass
|
|
return None
|
|
|
|
|
|
def _kill(s: HlsSession) -> None:
|
|
try:
|
|
s.proc.terminate()
|
|
try:
|
|
s.proc.wait(timeout=3)
|
|
except subprocess.TimeoutExpired:
|
|
s.proc.kill()
|
|
except Exception:
|
|
pass
|
|
shutil.rmtree(s.dir, ignore_errors=True)
|
|
|
|
|
|
def _enforce_cap(cap: int) -> None:
|
|
# Caller holds _lock. Drop the least-recently-accessed sessions over the cap.
|
|
if len(_sessions) <= cap:
|
|
return
|
|
for key, s in sorted(_sessions.items(), key=lambda kv: kv[1].last_access)[: len(_sessions) - cap]:
|
|
_kill(s)
|
|
_sessions.pop(key, None)
|
|
|
|
|
|
def start_session(
|
|
db: DbSession,
|
|
item: PlexItem,
|
|
start_s: float,
|
|
audio_ord: int | None = None,
|
|
aoff: float = 0.0,
|
|
multi: bool = False,
|
|
) -> HlsSession | None:
|
|
"""(Re)start the HLS remux for an item at the given offset. `multi` (item has >1 audio track) maps
|
|
every audio track as an HLS rendition so the client switches audio without a restart; otherwise a
|
|
single audio track (`audio_ord`) is muxed. Subtitles are separate WebVTT tracks (not muxed here).
|
|
Returns None if the local file can't be read."""
|
|
src = paths.local_media_path(db, item.file_path)
|
|
if src is None:
|
|
return None
|
|
key = item.rating_key
|
|
start_s = max(0.0, float(start_s))
|
|
audio_langs = _probe_audio_langs(src) if multi else []
|
|
is_multi = len(audio_langs) > 1 # only worth a master playlist when there really are ≥2 tracks
|
|
with _lock:
|
|
old = _sessions.pop(key, None)
|
|
if old is not None:
|
|
_kill(old)
|
|
tag = "multi" if is_multi else str(audio_ord)
|
|
directory = _HLS_ROOT / f"{key}_{int(start_s)}_{tag}_{aoff:+.2f}"
|
|
shutil.rmtree(directory, ignore_errors=True)
|
|
directory.mkdir(parents=True, exist_ok=True)
|
|
proc = subprocess.Popen(
|
|
_ffmpeg_cmd(src, start_s, directory, audio_ord, aoff, audio_langs if is_multi else None),
|
|
stdout=subprocess.DEVNULL,
|
|
stderr=subprocess.DEVNULL,
|
|
)
|
|
s = HlsSession(key, directory, proc, start_s, "master.m3u8" if is_multi else "index.m3u8")
|
|
_sessions[key] = s
|
|
# Admin-configurable concurrency cap (Configuration → Plex → Max concurrent transcodes);
|
|
# at least 1 so playback can't be capped to zero.
|
|
_enforce_cap(max(1, sysconfig.get_int(db, "plex_max_transcodes")))
|
|
# Learn the REAL start offset K from the first VIDEO segment (outside the lock — this blocks on
|
|
# ffmpeg producing it, and we must not hold up other sessions). The client reads `media_start_s`, so
|
|
# its absolute clock + subtitle shift key off the true keyframe, not the requested offset. On
|
|
# timeout/probe failure we keep the provisional requested offset (old behaviour, small drift).
|
|
if start_s > 0:
|
|
seg0 = directory / ("seg_0_0.ts" if is_multi else "seg_0.ts") # multi: video variant = v0
|
|
if wait_for(seg0, timeout=25.0):
|
|
k = _probe_first_pts(seg0)
|
|
if k is not None and k >= 0:
|
|
s.media_start_s = k
|
|
log.info(
|
|
"plex hls session start key=%s req=%.1f real=%.3f multi=%s audio=%s",
|
|
key, start_s, s.media_start_s, is_multi, audio_ord,
|
|
)
|
|
return s
|
|
|
|
|
|
def current_session(key: str) -> HlsSession | None:
|
|
with _lock:
|
|
s = _sessions.get(key)
|
|
if s is not None:
|
|
s.last_access = time.time()
|
|
return s
|
|
|
|
|
|
def wait_for(path: Path, timeout: float = 20.0) -> bool:
|
|
"""Wait until a session file (playlist / segment) exists and is non-empty. Segments are
|
|
produced ~faster than realtime, so a segment just ahead of playback appears quickly; a segment
|
|
far beyond the generated region won't (the frontend restarts the session at a seek instead)."""
|
|
end = time.time() + timeout
|
|
while time.time() < end:
|
|
try:
|
|
if path.exists() and path.stat().st_size > 0:
|
|
return True
|
|
except OSError:
|
|
pass
|
|
time.sleep(0.15)
|
|
return path.exists()
|
|
|
|
|
|
def sweep_orphans(older_than: float | None = None) -> int:
|
|
"""Delete OUR segment directories that no live session owns.
|
|
|
|
Sessions live only in THIS process, so at startup every session directory under `_HLS_ROOT` is
|
|
a leftover from a crashed or killed run — and nothing else ever reaps them (the download GC
|
|
deliberately skips `.plex-hls` as a system tree), so a few restarts mid-playback leave
|
|
gigabytes of `.ts` segments behind forever.
|
|
|
|
Only entries whose name matches `_SESSION_DIR_RE` are touched: the root is admin-configurable,
|
|
so anything else under it belongs to someone else and is left alone. Blocking (rmtree) — call
|
|
it off the event loop.
|
|
|
|
TWO guards against deleting a LIVE session, because this now runs CONCURRENTLY with request
|
|
handling (the caller stopped awaiting it so startup isn't held up):
|
|
* `older_than` (default: this process's start, minus `_MTIME_GRACE_S` of cross-clock slack) —
|
|
a directory touched since we booted was made by THIS run, so it is by definition not a
|
|
leftover. This is the guard that matters, because `start_session` REUSES the same path for
|
|
the same (key, offset, tag): a stale directory can become a live one at any moment.
|
|
* the live-session check and the rmtree happen together under `_lock`, so a session can't be
|
|
registered in the gap between them. The lock is held for one directory at a time — long
|
|
enough to be atomic, short enough that a boot-time sweep of a big backlog doesn't stall
|
|
playback."""
|
|
cutoff = (_PROCESS_START - _MTIME_GRACE_S) if older_than is None else older_than
|
|
dropped = 0
|
|
try:
|
|
entries = list(_HLS_ROOT.iterdir())
|
|
except OSError: # the root doesn't exist yet — nothing to sweep
|
|
return 0
|
|
for p in entries:
|
|
if not _SESSION_DIR_RE.match(p.name):
|
|
continue
|
|
if not p.is_dir():
|
|
continue
|
|
try:
|
|
if p.stat().st_mtime >= cutoff:
|
|
continue # this run made (or reused) it — never ours to delete
|
|
except OSError: # vanished under us (a concurrent restart of that session)
|
|
continue
|
|
with _lock:
|
|
if any(s.dir == p for s in _sessions.values()):
|
|
continue
|
|
shutil.rmtree(p, ignore_errors=True)
|
|
dropped += 1
|
|
return dropped
|
|
|
|
|
|
def reap_idle() -> int:
|
|
now = time.time()
|
|
dropped = 0
|
|
with _lock:
|
|
for key, s in list(_sessions.items()):
|
|
done = s.proc.poll() is not None
|
|
if now - s.last_access > _SESSION_IDLE_S or (done and now - s.last_access > 30):
|
|
_kill(s)
|
|
_sessions.pop(key, None)
|
|
dropped += 1
|
|
return dropped
|