"""Video streaming: preferred-audio remux and live channel concat.""" import json import logging import subprocess import tempfile from pathlib import Path from urllib.parse import unquote from .config import ( VIDEO_EXTENSIONS, MOVIES_ROOT, SERIES_ROOT, NORMALIZE_LIVE_STREAM, NORMALIZE_TARGET_RES, ) logger = logging.getLogger(__name__) def path_to_absolute(path_param: str) -> Path | None: """ Convert a path query parameter to an absolute Path under /movies or /series. Path must start with 'movies/' or 'series/'. Prevents path traversal. """ path_param = unquote(path_param).strip("/") if ".." in path_param or path_param.startswith("/"): return None if path_param.startswith("movies/"): sub = path_param[7:] root = MOVIES_ROOT elif path_param.startswith("series/"): sub = path_param[7:] root = SERIES_ROOT else: return None candidate = (root / sub).resolve() try: candidate.relative_to(root) except ValueError: return None if ( candidate.exists() and candidate.is_file() and candidate.suffix.lower() in VIDEO_EXTENSIONS ): return candidate return None def get_preferred_audio_stream_index(path: Path) -> int | None: """ Probe the file with ffprobe and return the 0-based audio stream index to use: Dutch (nld/dut/nl) if available, else Flemish (vls), else default (first audio). Returns None on probe failure or no audio streams. """ try: out = subprocess.run( [ "ffprobe", "-v", "quiet", "-print_format", "json", "-show_entries", "stream=index,codec_type", "-show_entries", "stream_tags=language", str(path), ], capture_output=True, text=True, timeout=10, ) if out.returncode != 0 or not out.stdout: return None data = json.loads(out.stdout) streams = data.get("streams") or [] audio_streams = [] for s in streams: if s.get("codec_type") != "audio": continue idx = s.get("index") if idx is None: continue tags = s.get("tags") or {} raw = (tags.get("language") or "").strip().lower() lang = raw[:3] if len(raw) >= 3 else raw audio_streams.append((idx, lang, raw)) if not audio_streams: return None dutch = ("nld", "dut", "nl") flemish = ("vls", "nl-be", "nlb") for idx, lang, raw in audio_streams: if lang in dutch or raw in dutch: return idx for idx, lang, raw in audio_streams: if lang in flemish or raw in flemish: return idx return audio_streams[0][0] except ( subprocess.TimeoutExpired, ValueError, FileNotFoundError, json.JSONDecodeError, ): return None def stream_file_with_preferred_audio(path: Path, wfile) -> bool: """ Stream the file via FFmpeg with only video and the preferred audio track. Returns True if streaming was started, False if caller should fall back to raw file. """ audio_idx = get_preferred_audio_stream_index(path) if audio_idx is None: return False try: proc = subprocess.Popen( [ "ffmpeg", "-i", str(path), "-map", "0:v", "-map", f"0:{audio_idx}", "-c", "copy", "-f", "mp4", "-movflags", "frag_keyframe+empty_moov+default_base_moof", "-", ], stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, start_new_session=True, ) while True: chunk = proc.stdout.read(65536) if not chunk: break wfile.write(chunk) wfile.flush() proc.wait(timeout=1) except (OSError, subprocess.TimeoutExpired): try: proc.terminate() proc.wait(timeout=2) except Exception: pass return False return True def _escape_path(p: Path) -> str: return str(p.resolve()).replace("\\", "\\\\").replace("'", "'\\''") def write_concat_list(paths: list[Path], fd) -> None: """Write FFmpeg concat demuxer list to file.""" for p in paths: fd.write(f"file '{_escape_path(p)}'\n") fd.flush() def _schedule_epoch_utc(): from datetime import datetime, timezone return datetime.now(timezone.utc).replace(hour=0, minute=0, second=0, microsecond=0) def _current_position_in_cycle(durations: list[float]) -> tuple[int, float]: """ Return (file_index, offset_sec_in_file) for 'now' in the schedule cycle. Schedule epoch = midnight UTC; cycle = one full playlist. """ from datetime import datetime, timezone epoch = _schedule_epoch_utc() now = datetime.now(timezone.utc) cycle_duration = sum(durations) if cycle_duration <= 0: return 0, 0.0 offset_sec = (now - epoch).total_seconds() % cycle_duration cumul = 0.0 for i, dur in enumerate(durations): if offset_sec < cumul + dur: return i, offset_sec - cumul cumul += dur return len(durations) - 1, durations[-1] def write_concat_list_from_current(paths: list[Path], durations: list[float], fd) -> None: """ Write a concat list that starts at the current position in the schedule and loops seamlessly. Uses inpoint/outpoint so we join mid-file, then repeat. """ n = len(paths) if n == 0: return idx, offset_in_file = _current_position_in_cycle(durations) for i in range(n): j = (idx + i) % n fd.write(f"file '{_escape_path(paths[j])}'\n") if i == 0 and offset_in_file >= 0.5: fd.write(f"inpoint {offset_in_file:.2f}\n") if offset_in_file >= 0.5: fd.write(f"file '{_escape_path(paths[idx])}'\n") fd.write(f"outpoint {offset_in_file:.2f}\n") fd.flush() def _parse_resolution(res: str) -> tuple[int, int]: """Parse 'WxH' into (width, height). Default 1920x1080 on parse error.""" try: w, h = res.strip().lower().split("x") return int(w), int(h) except (ValueError, AttributeError): return 1920, 1080 def stream_live_channel(paths: list[Path], wfile, durations: list[float]) -> None: """ Run FFmpeg to output a continuous MPEG-TS stream. Concat list starts at current position in schedule (midnight UTC cycle) so each new client joins where the 'channel' is now; then the list loops. EPG and stream stay in sync. When NORMALIZE_LIVE_STREAM is enabled, videos are re-encoded to a uniform format (resolution, codec). This prevents crashes when switching between episodes with different resolutions or codec parameters. """ if not paths: return with tempfile.NamedTemporaryFile(mode="w", suffix=".txt", delete=False) as f: write_concat_list_from_current(paths, durations, f) list_path = f.name try: base_args = [ "ffmpeg", "-fflags", "+nobuffer+flush_packets", "-stream_loop", "-1", "-f", "concat", "-safe", "0", "-i", list_path, ] if NORMALIZE_LIVE_STREAM: w, h = _parse_resolution(NORMALIZE_TARGET_RES) # Scale to fit target res, pad to exact size. Re-encode for uniform output. vf = f"scale={w}:{h}:force_original_aspect_ratio=decrease,pad={w}:{h}:(ow-iw)/2:(oh-ih)/2" output_args = [ "-vf", vf, "-c:v", "libx264", "-preset", "ultrafast", "-tune", "zerolatency", "-c:a", "aac", "-ac", "2", "-b:a", "128k", "-avoid_negative_ts", "make_zero", "-muxdelay", "0", "-muxpreload", "0", "-max_muxing_queue_size", "1024", "-f", "mpegts", "-", ] else: output_args = [ "-c", "copy", "-avoid_negative_ts", "make_zero", "-muxdelay", "0", "-muxpreload", "0", "-max_muxing_queue_size", "1024", "-f", "mpegts", "-", ] proc = subprocess.Popen( base_args + output_args, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, start_new_session=True, ) try: while True: chunk = proc.stdout.read(65536) if not chunk: break wfile.write(chunk) wfile.flush() except (BrokenPipeError, ConnectionResetError, OSError): pass finally: proc.terminate() try: proc.wait(timeout=5) except subprocess.TimeoutExpired: proc.kill() finally: Path(list_path).unlink(missing_ok=True)