mirror of
https://github.com/myronblair/kino-app
synced 2026-07-27 21:18:43 -05:00
8fa1291e0f
-progress pipe:1 -nostats belong to ffmpeg, not nice — they were landing right after cmd[0] (nice) instead of after the ffmpeg token, which made nice choke on an unrecognized option and fail before ffmpeg ever started.
253 lines
9.8 KiB
Python
253 lines
9.8 KiB
Python
"""HLS transcoding via ffmpeg.
|
|
|
|
Two modes:
|
|
- quick (stream-copy): instant, single bitrate, no quality loss. Source must be H.264/AAC.
|
|
- abr (adaptive): re-encode to multiple bitrates with master playlist for ABR streaming.
|
|
"""
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import shutil
|
|
from pathlib import Path
|
|
from typing import Optional, List, Tuple, Callable, Awaitable
|
|
|
|
logger = logging.getLogger("kino.transcode")
|
|
|
|
# Moderate priority (was 19, the absolute lowest) and a thread cap so an ABR encode can't
|
|
# starve the rest of the stack (Mongo, *arr apps, the API itself) on this 8-core host.
|
|
TRANSCODE_NICE = 10
|
|
TRANSCODE_THREADS = 6
|
|
|
|
|
|
async def probe_video(source: Path) -> Tuple[int, int]:
|
|
"""Return (width, height) using ffprobe; (0, 0) on failure."""
|
|
try:
|
|
proc = await asyncio.create_subprocess_exec(
|
|
"ffprobe", "-v", "error", "-select_streams", "v:0",
|
|
"-show_entries", "stream=width,height", "-of", "json", str(source),
|
|
stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE,
|
|
)
|
|
stdout, _ = await proc.communicate()
|
|
if proc.returncode != 0:
|
|
return (0, 0)
|
|
data = json.loads(stdout.decode("utf-8", errors="ignore") or "{}")
|
|
s = (data.get("streams") or [{}])[0]
|
|
return (int(s.get("width") or 0), int(s.get("height") or 0))
|
|
except Exception:
|
|
return (0, 0)
|
|
|
|
|
|
async def probe_duration(source: Path) -> float:
|
|
"""Return source duration in seconds using ffprobe; 0 on failure."""
|
|
try:
|
|
proc = await asyncio.create_subprocess_exec(
|
|
"ffprobe", "-v", "error", "-show_entries", "format=duration",
|
|
"-of", "json", str(source),
|
|
stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE,
|
|
)
|
|
stdout, _ = await proc.communicate()
|
|
if proc.returncode != 0:
|
|
return 0.0
|
|
data = json.loads(stdout.decode("utf-8", errors="ignore") or "{}")
|
|
return float((data.get("format") or {}).get("duration") or 0.0)
|
|
except Exception:
|
|
return 0.0
|
|
|
|
|
|
# (height, video_bitrate, max_bitrate, buffer, audio_bitrate)
|
|
ALL_VARIANTS = [
|
|
(1080, "5000k", "5500k", "7500k", "192k"),
|
|
(720, "2800k", "3000k", "4200k", "128k"),
|
|
(480, "1400k", "1500k", "2100k", "96k"),
|
|
(360, "800k", "856k", "1200k", "64k"),
|
|
]
|
|
|
|
|
|
def variants_for_source(src_height: int) -> List[Tuple[int, str, str, str, str]]:
|
|
"""Pick variants ≤ source height. Always include at least one (smallest)."""
|
|
if src_height <= 0:
|
|
return [ALL_VARIANTS[2]] # safe default 480p
|
|
chosen = [v for v in ALL_VARIANTS if v[0] <= src_height]
|
|
if not chosen:
|
|
chosen = [ALL_VARIANTS[-1]]
|
|
return chosen
|
|
|
|
|
|
async def transcode_audio_fix(source: Path, out_dir: Path, on_status: Callable[..., Awaitable]) -> None:
|
|
"""Video stream-copied untouched, audio re-encoded to AAC — fixes browsers' inability to
|
|
play AC3/DTS/EAC3 (common on BDRips) without paying for a full video re-encode. Same
|
|
H.264-source assumption as transcode_quick; HEVC/x265 sources still need the full ABR path
|
|
since the video itself isn't browser-playable either in that case."""
|
|
if not source.is_file():
|
|
await on_status("failed", error=f"Source missing: {source}")
|
|
return
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
cmd = [
|
|
"nice", "-n", str(TRANSCODE_NICE),
|
|
"ffmpeg", "-y", "-i", str(source),
|
|
"-c:v", "copy", "-c:a", "aac", "-b:a", "192k", "-ac", "2",
|
|
"-bsf:v", "h264_mp4toannexb",
|
|
"-f", "hls", "-hls_time", "6",
|
|
"-hls_list_size", "0", "-hls_playlist_type", "vod",
|
|
"-hls_segment_filename", str(out_dir / "seg_%04d.ts"),
|
|
str(out_dir / "playlist.m3u8"),
|
|
]
|
|
await _run(cmd, on_status, out_dir, "playlist.m3u8", source)
|
|
|
|
|
|
async def transcode_quick(source: Path, out_dir: Path, on_status: Callable[..., Awaitable]) -> None:
|
|
"""Stream-copy to single-rate HLS. Output filename: playlist.m3u8."""
|
|
if not source.is_file():
|
|
await on_status("failed", error=f"Source missing: {source}")
|
|
return
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
cmd = [
|
|
"nice", "-n", str(TRANSCODE_NICE),
|
|
"ffmpeg", "-y", "-i", str(source),
|
|
"-c:v", "copy", "-c:a", "copy",
|
|
"-bsf:v", "h264_mp4toannexb",
|
|
"-f", "hls", "-hls_time", "6",
|
|
"-hls_list_size", "0", "-hls_playlist_type", "vod",
|
|
"-hls_segment_filename", str(out_dir / "seg_%04d.ts"),
|
|
str(out_dir / "playlist.m3u8"),
|
|
]
|
|
await _run(cmd, on_status, out_dir, "playlist.m3u8", source)
|
|
|
|
|
|
async def transcode_abr(source: Path, out_dir: Path, on_status: Callable[..., Awaitable]) -> None:
|
|
"""Multi-bitrate ABR HLS. Output entry filename: master.m3u8."""
|
|
if not source.is_file():
|
|
await on_status("failed", error=f"Source missing: {source}")
|
|
return
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
_, height = await probe_video(source)
|
|
variants = variants_for_source(height)
|
|
n = len(variants)
|
|
|
|
# Build filter graph
|
|
splits = "".join(f"[v{i}]" for i in range(n))
|
|
fc_parts = [f"[0:v]split={n}{splits}"]
|
|
for i, (h, *_rest) in enumerate(variants):
|
|
fc_parts.append(f"[v{i}]scale=w=-2:h={h}[v{i}out]")
|
|
filter_complex = ";".join(fc_parts)
|
|
|
|
# Per-variant thread count, not the full cap each — x264 threads are set per encoder
|
|
# instance, so giving every one of the (up to 4) variants the full TRANSCODE_THREADS
|
|
# budget lets them all run near-full-tilt simultaneously and still saturate the host.
|
|
per_variant_threads = max(1, TRANSCODE_THREADS // n)
|
|
|
|
cmd: List[str] = ["nice", "-n", str(TRANSCODE_NICE), "ffmpeg", "-y", "-threads", str(per_variant_threads), "-i", str(source), "-filter_complex", filter_complex]
|
|
|
|
for i, (_h, vb, maxr, buf, _ab) in enumerate(variants):
|
|
cmd += [
|
|
"-map", f"[v{i}out]",
|
|
f"-c:v:{i}", "libx264",
|
|
f"-preset:v:{i}", "veryfast",
|
|
f"-threads:v:{i}", str(per_variant_threads),
|
|
f"-profile:v:{i}", "main",
|
|
f"-pix_fmt:v:{i}", "yuv420p",
|
|
f"-b:v:{i}", vb,
|
|
f"-maxrate:v:{i}", maxr,
|
|
f"-bufsize:v:{i}", buf,
|
|
f"-g", "48", f"-keyint_min", "48", f"-sc_threshold", "0",
|
|
]
|
|
|
|
# Audio: same source mapped N times, one per variant
|
|
for i in range(n):
|
|
cmd += ["-map", "a:0?"]
|
|
for i, (*_v, ab) in enumerate(variants):
|
|
cmd += [f"-c:a:{i}", "aac", f"-b:a:{i}", ab, f"-ac:a:{i}", "2"]
|
|
|
|
var_stream_map = " ".join(f"v:{i},a:{i}" for i in range(n))
|
|
cmd += [
|
|
"-f", "hls",
|
|
"-hls_time", "6",
|
|
"-hls_playlist_type", "vod",
|
|
"-hls_flags", "independent_segments",
|
|
"-hls_segment_filename", str(out_dir / "v%v" / "seg_%04d.ts"),
|
|
"-master_pl_name", "master.m3u8",
|
|
"-var_stream_map", var_stream_map,
|
|
str(out_dir / "v%v" / "playlist.m3u8"),
|
|
]
|
|
|
|
# Pre-create variant subdirs (some ffmpeg builds need them)
|
|
for i in range(n):
|
|
(out_dir / f"v{i}").mkdir(exist_ok=True)
|
|
|
|
await _run(cmd, on_status, out_dir, "master.m3u8", source)
|
|
|
|
|
|
def _parse_out_time_seconds(line: str) -> Optional[float]:
|
|
"""Parse an `out_time=HH:MM:SS.microseconds` line from ffmpeg's -progress output."""
|
|
_, _, value = line.partition("=")
|
|
value = value.strip()
|
|
if not value or value == "N/A":
|
|
return None
|
|
try:
|
|
h, m, s = value.split(":")
|
|
return int(h) * 3600 + int(m) * 60 + float(s)
|
|
except ValueError:
|
|
return None
|
|
|
|
|
|
async def _run(cmd: List[str], on_status, out_dir: Path, entry_filename: str, source: Optional[Path] = None) -> None:
|
|
duration = await probe_duration(source) if source else 0.0
|
|
await on_status("running", progress=0.0)
|
|
# Machine-readable progress on stdout, separate from ffmpeg's normal stderr logging.
|
|
# Insert right after the "ffmpeg" token, not cmd[0] — cmd is prefixed with `nice -n N`,
|
|
# and these flags belong to ffmpeg, not to nice.
|
|
ffmpeg_idx = cmd.index("ffmpeg")
|
|
cmd = cmd[:ffmpeg_idx + 1] + ["-progress", "pipe:1", "-nostats"] + cmd[ffmpeg_idx + 1:]
|
|
try:
|
|
proc = await asyncio.create_subprocess_exec(
|
|
*cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE,
|
|
)
|
|
|
|
stderr_chunks: List[bytes] = []
|
|
|
|
async def _drain_stderr():
|
|
async for line in proc.stderr:
|
|
stderr_chunks.append(line)
|
|
|
|
stderr_task = asyncio.create_task(_drain_stderr())
|
|
|
|
last_reported = -1.0
|
|
async for raw in proc.stdout:
|
|
line = raw.decode("utf-8", errors="ignore").strip()
|
|
if line.startswith("out_time=") and duration > 0:
|
|
seconds = _parse_out_time_seconds(line)
|
|
if seconds is not None:
|
|
pct = round(min(99.0, seconds / duration * 100), 1)
|
|
if pct != last_reported:
|
|
last_reported = pct
|
|
await on_status("running", progress=pct)
|
|
|
|
await proc.wait()
|
|
await stderr_task
|
|
stderr = b"".join(stderr_chunks)
|
|
if proc.returncode != 0:
|
|
err = stderr.decode("utf-8", errors="ignore")[-500:]
|
|
logger.error(f"ffmpeg failed: {err}")
|
|
await on_status("failed", error=err[:200])
|
|
shutil.rmtree(out_dir, ignore_errors=True)
|
|
return
|
|
await on_status("done", entry=entry_filename, progress=100.0)
|
|
except FileNotFoundError:
|
|
await on_status("failed", error="ffmpeg not installed")
|
|
except Exception as e:
|
|
logger.exception("transcode crashed")
|
|
await on_status("failed", error=str(e))
|
|
shutil.rmtree(out_dir, ignore_errors=True)
|
|
|
|
|
|
def srt_to_vtt(srt_text: str) -> str:
|
|
lines = srt_text.replace("\r\n", "\n").split("\n")
|
|
out = ["WEBVTT", ""]
|
|
for line in lines:
|
|
if "-->" in line:
|
|
out.append(line.replace(",", "."))
|
|
else:
|
|
out.append(line)
|
|
return "\n".join(out)
|