Files
Цвылев Александр Вадимович 9a78cf5261 refactor(transcode): thread fs ops + atomic cache publish
- run blocking Path.exists/read_text off the event loop via anyio.to_thread
- publish HLS renditions by building into a temp dir then atomic rename, so a
  reader never sees a playlist referencing half-written segments
- per-writer temp name for opus so two concurrent jobs can't interleave
- drop a track's cached renditions on delete so they don't dangle

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-29 10:53:35 +03:00

72 lines
2.7 KiB
Python

"""FfmpegTranscoder — Opus + HLS renditions via the ffmpeg CLI.
Implements :class:`app.domain.ports.Transcoder`. Runs ffmpeg as a subprocess
(only ever from a worker — CLAUDE.md: no heavy work in the request cycle) and
raises :class:`TranscodeError` on a non-zero exit. Opus is used for the direct
quality renditions; HLS segments are AAC-in-MPEG-TS (broad player support).
"""
import asyncio
import uuid
from pathlib import Path
import anyio
from app.core.logging import get_logger
from app.domain.errors import TranscodeError
log = get_logger(__name__)
# HLS segment length. 10s is a common VOD default — few requests, quick seeks.
_HLS_SEGMENT_SECONDS = "10"
def _mkdir(path: Path) -> None:
path.mkdir(parents=True, exist_ok=True)
class FfmpegTranscoder:
def __init__(self, ffmpeg_path: str = "ffmpeg") -> None:
self._ffmpeg = ffmpeg_path
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None:
await anyio.to_thread.run_sync(_mkdir, dest.parent)
# Per-writer temp name (not a shared ``.part``) so two concurrent jobs for
# the same rendition can't interleave into one file — each writes its own
# temp and the last atomic replace wins, both leaving a valid output.
tmp = dest.with_name(f"{dest.name}.{uuid.uuid4().hex}.part")
await self._run(
"-i", str(src),
"-vn", "-c:a", "libopus", "-b:a", f"{bitrate_kbps}k",
"-f", "opus", str(tmp),
)
# Publish atomically so a concurrent reader never sees a half-written file.
await anyio.to_thread.run_sync(tmp.replace, dest)
async def to_hls(self, src: Path, out_dir: Path, *, bitrate_kbps: int) -> None:
await anyio.to_thread.run_sync(_mkdir, out_dir)
playlist = out_dir / "playlist.m3u8"
segments = out_dir / "seg_%03d.ts"
await self._run(
"-i", str(src),
"-vn", "-c:a", "aac", "-b:a", f"{bitrate_kbps}k",
"-f", "hls",
"-hls_time", _HLS_SEGMENT_SECONDS,
"-hls_playlist_type", "vod",
"-hls_segment_filename", str(segments),
str(playlist),
)
async def _run(self, *args: str) -> None:
cmd = [self._ffmpeg, "-y", "-nostdin", "-loglevel", "error", *args]
proc = await asyncio.create_subprocess_exec(
*cmd,
stdout=asyncio.subprocess.DEVNULL,
stderr=asyncio.subprocess.PIPE,
)
_, stderr = await proc.communicate()
if proc.returncode != 0:
detail = stderr.decode(errors="replace").strip()[-500:]
log.error("ffmpeg_failed", returncode=proc.returncode, error=detail)
raise TranscodeError(f"ffmpeg exited {proc.returncode}: {detail}")