Files
peter 5fe6d32d38 fix(plex): give the sweep's boot-time cutoff cross-clock slack
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).
2026-07-24 04:52:41 +02:00

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