9a78cf5261
- 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>
111 lines
4.0 KiB
Python
111 lines
4.0 KiB
Python
"""Transcode service + cache-path helpers (Group B / plan §6.6).
|
|
|
|
Cache layout under ``transcode_cache_path``::
|
|
|
|
{track_id}/opus_{kbps}.opus # direct quality renditions
|
|
{track_id}/hls/playlist.m3u8 # HLS rendition (AAC-in-TS)
|
|
{track_id}/hls/seg_000.ts …
|
|
|
|
The request side (streaming router) only *reads* the cache — misses fall back to
|
|
the original file and enqueue generation. The worker (``transcode_task``) writes
|
|
it. Path helpers are module-level so both sides agree on locations without one
|
|
importing the other.
|
|
"""
|
|
|
|
import re
|
|
import shutil
|
|
import uuid
|
|
from pathlib import Path
|
|
|
|
import anyio
|
|
|
|
from app.domain.errors import NotFoundError
|
|
from app.domain.ports import TrackRepository
|
|
|
|
# Stream-quality name (matches the user-settings ``StreamQuality``) → Opus
|
|
# bitrate. ``original`` is absent: it means "serve the master, no transcode".
|
|
QUALITY_BITRATE: dict[str, int] = {"high": 128, "medium": 96, "low": 64}
|
|
|
|
# Single HLS rendition bitrate (AAC). One rendition keeps the MVP simple; a
|
|
# multi-bitrate ladder can come later.
|
|
HLS_BITRATE = 128
|
|
|
|
# Only these segment names may be served, guarding the segment route against
|
|
# path traversal.
|
|
_SEGMENT_RE = re.compile(r"^seg_\d{3,}\.ts$")
|
|
|
|
|
|
def bitrate_for_quality(quality: str) -> int | None:
|
|
"""Opus bitrate for a quality name, or ``None`` for ``original``/unknown."""
|
|
return QUALITY_BITRATE.get(quality)
|
|
|
|
|
|
def track_cache_dir(root: Path, track_id: uuid.UUID) -> Path:
|
|
return root / str(track_id)
|
|
|
|
|
|
def opus_path(root: Path, track_id: uuid.UUID, bitrate_kbps: int) -> Path:
|
|
return track_cache_dir(root, track_id) / f"opus_{bitrate_kbps}.opus"
|
|
|
|
|
|
def hls_dir(root: Path, track_id: uuid.UUID) -> Path:
|
|
return track_cache_dir(root, track_id) / "hls"
|
|
|
|
|
|
def hls_playlist_path(root: Path, track_id: uuid.UUID) -> Path:
|
|
return hls_dir(root, track_id) / "playlist.m3u8"
|
|
|
|
|
|
def hls_segment_path(root: Path, track_id: uuid.UUID, name: str) -> Path | None:
|
|
"""Resolve a segment file, or ``None`` if the name is not a valid segment."""
|
|
if not _SEGMENT_RE.fullmatch(name):
|
|
return None
|
|
return hls_dir(root, track_id) / name
|
|
|
|
|
|
def remove_track_cache(root: Path, track_id: uuid.UUID) -> None:
|
|
"""Delete every cached rendition for a track (Opus + HLS). Best-effort — used
|
|
when a track is deleted so its transcode cache doesn't dangle forever."""
|
|
shutil.rmtree(track_cache_dir(root, track_id), ignore_errors=True)
|
|
|
|
|
|
class TranscodeService:
|
|
"""Request-side cache lookups for transcoded renditions."""
|
|
|
|
def __init__(self, *, tracks: TrackRepository, cache_root: Path) -> None:
|
|
self._tracks = tracks
|
|
self._root = cache_root
|
|
|
|
async def _require_streamable(self, track_id: uuid.UUID) -> None:
|
|
track = await self._tracks.get_by_id(track_id)
|
|
if track is None:
|
|
raise NotFoundError("Track not found.")
|
|
if track.storage_uri is None:
|
|
raise NotFoundError("Track is not yet downloaded.")
|
|
|
|
async def resolve_quality_file(
|
|
self, track_id: uuid.UUID, quality: str
|
|
) -> Path | None:
|
|
"""Cached Opus file for ``quality`` if present, else ``None`` (caller
|
|
falls back to the master and enqueues generation). ``original`` → None."""
|
|
bitrate = bitrate_for_quality(quality)
|
|
if bitrate is None:
|
|
return None
|
|
path = opus_path(self._root, track_id, bitrate)
|
|
exists = await anyio.to_thread.run_sync(path.exists)
|
|
return path if exists else None
|
|
|
|
async def hls_playlist(self, track_id: uuid.UUID) -> Path | None:
|
|
"""Cached HLS playlist if generated, else ``None`` (validates the track
|
|
exists so an unknown id 404s rather than silently missing)."""
|
|
await self._require_streamable(track_id)
|
|
path = hls_playlist_path(self._root, track_id)
|
|
exists = await anyio.to_thread.run_sync(path.exists)
|
|
return path if exists else None
|
|
|
|
def hls_segment(self, track_id: uuid.UUID, name: str) -> Path | None:
|
|
path = hls_segment_path(self._root, track_id, name)
|
|
if path is None or not path.exists():
|
|
return None
|
|
return path
|