feat(transcode): optimize + on-the-fly quality + HLS streaming (§6.6)

POST /tracks/{id}/optimize enqueues a transcode; GET /stream/{id}?quality=
serves a cached Opus rendition (miss → master + background warm); GET
/stream/{id}/hls/{playlist.m3u8,segment} serves the worker-generated HLS
rendition (AAC-in-TS), with ?token= propagated onto segment URLs for players
that can't set headers. Hexagonal: Transcoder port, FfmpegTranscoder adapter,
TranscodeService + cache-path helpers, transcode_track worker (idempotent),
schemas + deps wiring. ffmpeg already in the image.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Цвылев Александр Вадимович
2026-07-28 21:19:36 +03:00
parent 8271de34eb
commit 313af3a070
13 changed files with 384 additions and 11 deletions
+11
View File
@@ -24,6 +24,7 @@ from app.application.remote_library_service import RemoteLibraryService
from app.application.streaming_service import StreamingService from app.application.streaming_service import StreamingService
from app.application.subsonic_auth_service import SubsonicAuthService from app.application.subsonic_auth_service import SubsonicAuthService
from app.application.sync_service import SyncService from app.application.sync_service import SyncService
from app.application.transcode_service import TranscodeService
from app.application.upload_service import UploadService from app.application.upload_service import UploadService
from app.application.user_service import UserService from app.application.user_service import UserService
from app.application.user_settings_service import UserSettingsService from app.application.user_settings_service import UserSettingsService
@@ -178,6 +179,15 @@ def get_metadata_service(session: SessionDep, storage: FileStorageDep) -> Metada
) )
def get_transcode_service(session: SessionDep) -> TranscodeService:
"""Request-side cache lookups for transcoded renditions (§6.6). Generation
itself runs in the ``transcode_track`` worker, never here."""
return TranscodeService(
tracks=SqlAlchemyTrackRepository(session),
cache_root=get_settings().transcode_cache_path,
)
def get_lyrics_service(session: SessionDep) -> LyricsService: def get_lyrics_service(session: SessionDep) -> LyricsService:
"""Wires the LRCLIB lyrics provider + cache repo (plan §6.7). LRCLIB is """Wires the LRCLIB lyrics provider + cache repo (plan §6.7). LRCLIB is
keyless, so this is always available; failures degrade to ``not_found``.""" keyless, so this is always available; failures degrade to ``not_found``."""
@@ -215,6 +225,7 @@ UploadServiceDep = Annotated[UploadService, Depends(get_upload_service)]
StreamingServiceDep = Annotated[StreamingService, Depends(get_streaming_service)] StreamingServiceDep = Annotated[StreamingService, Depends(get_streaming_service)]
MetadataServiceDep = Annotated[MetadataEnrichmentService, Depends(get_metadata_service)] MetadataServiceDep = Annotated[MetadataEnrichmentService, Depends(get_metadata_service)]
LyricsServiceDep = Annotated[LyricsService, Depends(get_lyrics_service)] LyricsServiceDep = Annotated[LyricsService, Depends(get_lyrics_service)]
TranscodeServiceDep = Annotated[TranscodeService, Depends(get_transcode_service)]
DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)] DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)]
RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)] RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)]
+11
View File
@@ -0,0 +1,11 @@
"""Transcode/optimize response schemas (§6.6)."""
from pydantic import BaseModel
class OptimizeEnqueuedOut(BaseModel):
"""Acknowledgement that a transcode job was queued (it runs in the worker)."""
status: str
job_id: str
quality: str
+64 -6
View File
@@ -1,30 +1,53 @@
"""Audio streaming endpoint — direct stream with Range support.""" """Audio streaming — direct byte-range stream, transcoded quality, and HLS.
``GET /stream/{id}`` streams the master with Range support, or a cached Opus
rendition when ``?quality=`` is set (a cache miss falls back to the master and
warms the cache in the background — playback never waits on ffmpeg). ``/hls/*``
serves the cached HLS rendition (generated by the ``transcode_track`` worker).
"""
import re
import uuid import uuid
from typing import Annotated from typing import Annotated
from fastapi import APIRouter, Header from fastapi import APIRouter, Header, Query, Response
from fastapi.responses import StreamingResponse from fastapi.responses import FileResponse, StreamingResponse
from app.api.deps import StreamingServiceDep, StreamUser from app.api.deps import StreamingServiceDep, StreamUser, TranscodeServiceDep
from app.domain.errors import NotFoundError
from app.workers.queue import enqueue_transcode_quiet
router = APIRouter(prefix="/stream", tags=["streaming"]) router = APIRouter(prefix="/stream", tags=["streaming"])
_HLS_PLAYLIST_TYPE = "application/vnd.apple.mpegurl"
_HLS_SEGMENT_TYPE = "video/mp2t"
_OPUS_TYPE = "audio/ogg"
_SEGMENT_LINE_RE = re.compile(r"^(seg_\d+\.ts)$", re.MULTILINE)
@router.get("/{track_id}") @router.get("/{track_id}")
async def stream_track( async def stream_track(
track_id: uuid.UUID, track_id: uuid.UUID,
service: StreamingServiceDep, service: StreamingServiceDep,
transcode: TranscodeServiceDep,
_user: StreamUser, _user: StreamUser,
range_header: Annotated[str | None, Header(alias="Range")] = None, range_header: Annotated[str | None, Header(alias="Range")] = None,
) -> StreamingResponse: quality: Annotated[str | None, Query()] = None,
) -> Response:
# A quality rendition, if one is cached; otherwise fall back to the master
# and enqueue generation so the next play gets it (graceful degradation).
if quality and quality != "original":
cached = await transcode.resolve_quality_file(track_id, quality)
if cached is not None:
return FileResponse(cached, media_type=_OPUS_TYPE)
await enqueue_transcode_quiet(track_id, quality=quality, hls=False)
result = await service.open_stream(track_id, range_header) result = await service.open_stream(track_id, range_header)
headers = { headers = {
"Accept-Ranges": "bytes", "Accept-Ranges": "bytes",
"Content-Length": str(result.content_length), "Content-Length": str(result.content_length),
} }
if result.is_partial: if result.is_partial:
headers["Content-Range"] = f"bytes {result.start}-{result.end}/{result.total_size}" headers["Content-Range"] = f"bytes {result.start}-{result.end}/{result.total_size}"
status_code = 206 status_code = 206
@@ -37,3 +60,38 @@ async def stream_track(
headers=headers, headers=headers,
media_type=result.content_type, media_type=result.content_type,
) )
@router.get("/{track_id}/hls/playlist.m3u8")
async def stream_hls_playlist(
track_id: uuid.UUID,
transcode: TranscodeServiceDep,
_user: StreamUser,
token: Annotated[str | None, Query()] = None,
) -> Response:
"""Serve the cached HLS playlist. On a miss, kick off generation and 404 so
the client retries. Segment URLs are relative; when the request carried a
``?token=`` (players can't set an Authorization header), it's appended to
each segment line so the segment requests authenticate the same way."""
path = await transcode.hls_playlist(track_id)
if path is None:
await enqueue_transcode_quiet(track_id, hls=True)
raise NotFoundError("HLS rendition is being prepared; retry shortly.")
body = path.read_text()
if token:
body = _SEGMENT_LINE_RE.sub(rf"\1?token={token}", body)
return Response(body, media_type=_HLS_PLAYLIST_TYPE)
@router.get("/{track_id}/hls/{segment}")
async def stream_hls_segment(
track_id: uuid.UUID,
segment: str,
transcode: TranscodeServiceDep,
_user: StreamUser,
) -> FileResponse:
path = transcode.hls_segment(track_id, segment)
if path is None:
raise NotFoundError("Segment not found.")
return FileResponse(path, media_type=_HLS_SEGMENT_TYPE)
+23 -5
View File
@@ -1,7 +1,7 @@
"""Track endpoints.""" """Track endpoints."""
import uuid import uuid
from typing import Any from typing import Annotated, Any
from fastapi import APIRouter, Query, Response from fastapi import APIRouter, Query, Response
from fastapi.responses import StreamingResponse from fastapi.responses import StreamingResponse
@@ -30,10 +30,12 @@ from app.api.schemas.track import (
TrackOut, TrackOut,
TrackUpdate, TrackUpdate,
) )
from app.api.schemas.transcode import OptimizeEnqueuedOut
from app.application.transcode_service import bitrate_for_quality
from app.domain.entities.album import Album from app.domain.entities.album import Album
from app.domain.entities.track import Artist, Track from app.domain.entities.track import Artist, Track
from app.domain.errors import NotFoundError from app.domain.errors import NotFoundError, ValidationError
from app.workers.queue import enqueue from app.workers.queue import enqueue, enqueue_transcode
router = APIRouter(prefix="/tracks", tags=["tracks"]) router = APIRouter(prefix="/tracks", tags=["tracks"])
@@ -222,8 +224,24 @@ async def delete_track(
async def get_similar_tracks(track_id: uuid.UUID, _: CurrentUser) -> Any: ... async def get_similar_tracks(track_id: uuid.UUID, _: CurrentUser) -> Any: ...
@router.post("/{track_id}/optimize") @router.post("/{track_id}/optimize", status_code=202)
async def optimize_track(track_id: uuid.UUID, _: CurrentUser) -> Any: ... async def optimize_track(
track_id: uuid.UUID,
track_repo: TrackRepoDep,
_: CurrentUser,
quality: Annotated[str, Query()] = "high",
) -> OptimizeEnqueuedOut:
"""Enqueue transcoding of a track into a cached Opus rendition + HLS (§6.6).
Heavy ffmpeg work runs in the worker; this only queues it."""
if bitrate_for_quality(quality) is None:
raise ValidationError(
f"Unknown quality '{quality}'; expected one of high, medium, low."
)
track = await track_repo.get_by_id(track_id)
if track is None:
raise NotFoundError(f"Track {track_id} not found.")
job_id = await enqueue_transcode(track_id, quality=quality, hls=True)
return OptimizeEnqueuedOut(status="enqueued", job_id=job_id, quality=quality)
@router.get("/{track_id}/lyrics") @router.get("/{track_id}/lyrics")
+99
View File
@@ -0,0 +1,99 @@
"""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 uuid
from pathlib import Path
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
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)
return path if path.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)
return path if path.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
+3
View File
@@ -81,6 +81,9 @@ class Settings(BaseSettings):
# -- media / storage -------------------------------------------------- # -- media / storage --------------------------------------------------
media_path: Path = Path("/data/media") media_path: Path = Path("/data/media")
transcode_cache_path: Path = Path("/data/transcode-cache") transcode_cache_path: Path = Path("/data/transcode-cache")
# ffmpeg binary for transcoding/HLS (on PATH in the image); override for a
# non-standard location.
ffmpeg_path: str = "ffmpeg"
max_parallel_downloads: int = 2 max_parallel_downloads: int = 2
# How many times the download worker retries a failed fetch (yt-dlp fails # How many times the download worker retries a failed fetch (yt-dlp fails
# often) before marking the job ``failed`` — exponential backoff between tries. # often) before marking the job ``failed`` — exponential backoff between tries.
+7
View File
@@ -76,6 +76,13 @@ class StorageError(DomainError):
code = "storage_error" code = "storage_error"
class TranscodeError(DomainError):
"""Transcoding (ffmpeg) failed. Raised in the worker; a play falls back to
the original file rather than surfacing this."""
code = "transcode_error"
class RangeNotSatisfiableError(DomainError): class RangeNotSatisfiableError(DomainError):
"""Requested byte range cannot be satisfied.""" """Requested byte range cannot be satisfied."""
+10
View File
@@ -497,6 +497,16 @@ class CoverArtProvider(Protocol):
async def fetch_release_group(self, release_group_mbid: str) -> CoverArt | None: ... async def fetch_release_group(self, release_group_mbid: str) -> CoverArt | None: ...
class Transcoder(Protocol):
"""Transcodes an audio file with ffmpeg (plan §6.6 / Group B). ``to_opus``
writes a single Opus rendition; ``to_hls`` writes an HLS playlist + segments
(AAC-in-TS) into ``out_dir``. Both raise ``TranscodeError`` on failure; heavy
work always runs in a worker, never the request cycle."""
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None: ...
async def to_hls(self, src: Path, out_dir: Path, *, bitrate_kbps: int) -> None: ...
class LyricsProvider(Protocol): class LyricsProvider(Protocol):
"""Fetches lyrics from an external database (LRCLIB) by artist/title/album/ """Fetches lyrics from an external database (LRCLIB) by artist/title/album/
duration. Returns a hit or ``None`` (no match / service down), never raising.""" duration. Returns a hit or ``None`` (no match / service down), never raising."""
+1
View File
@@ -0,0 +1 @@
"""ffmpeg-based transcoding adapters (Group B / plan §6.6)."""
+67
View File
@@ -0,0 +1,67 @@
"""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
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)
tmp = dest.with_suffix(dest.suffix + ".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}")
+2
View File
@@ -14,6 +14,7 @@ from app.workers.tasks.download_task import download_track
from app.workers.tasks.enrich_task import enrich_track from app.workers.tasks.enrich_task import enrich_track
from app.workers.tasks.import_task import scan_local_folder from app.workers.tasks.import_task import scan_local_folder
from app.workers.tasks.materialize_task import materialize_track from app.workers.tasks.materialize_task import materialize_track
from app.workers.tasks.transcode_task import transcode_track
log = get_logger("worker") log = get_logger("worker")
@@ -35,6 +36,7 @@ class WorkerSettings:
download_track, download_track,
materialize_track, materialize_track,
cleanup_storage, cleanup_storage,
transcode_track,
] ]
on_startup = startup on_startup = startup
on_shutdown = shutdown on_shutdown = shutdown
+22
View File
@@ -59,6 +59,28 @@ async def enqueue_materialize(job_id: uuid.UUID) -> None:
log.warning("materialize_enqueue_failed", job_id=str(job_id)) log.warning("materialize_enqueue_failed", job_id=str(job_id))
async def enqueue_transcode(
track_id: uuid.UUID, *, quality: str = "high", hls: bool = True
) -> str:
"""Enqueue a transcode job (§6.6). Unlike the best-effort follow-ups below,
this is user/stream-driven, so it surfaces a job id (and a 503 via the caller
if the queue is down) rather than swallowing failures."""
return await enqueue(
"transcode_track", track_id=str(track_id), quality=quality, hls=hls
)
async def enqueue_transcode_quiet(
track_id: uuid.UUID, *, quality: str = "high", hls: bool = False
) -> None:
"""Best-effort transcode enqueue for a streaming cache miss — never blocks or
fails the play (the master is served meanwhile)."""
try:
await enqueue_transcode(track_id, quality=quality, hls=hls)
except DependencyUnavailableError:
log.warning("transcode_enqueue_failed", track_id=str(track_id))
async def enqueue_enrich(track_id: uuid.UUID) -> None: async def enqueue_enrich(track_id: uuid.UUID) -> None:
"""Best-effort enqueue of metadata enrichment for a freshly stored track. """Best-effort enqueue of metadata enrichment for a freshly stored track.
+64
View File
@@ -0,0 +1,64 @@
"""arq task: transcode a track to a cached Opus rendition and/or HLS (§6.6).
Triggered by ``POST /tracks/{id}/optimize`` and by a streaming cache miss. Reads
the master once (via ``storage.as_local_path``) and writes into
``transcode_cache_path``. Idempotent — existing outputs are skipped, so repeated
enqueues (e.g. several plays before the first finishes) are cheap. The DB session
is released before ffmpeg runs; the track entity is a detached value object.
"""
import uuid
from typing import Any
from app.application.transcode_service import (
HLS_BITRATE,
bitrate_for_quality,
hls_dir,
hls_playlist_path,
opus_path,
)
from app.core.config import get_settings
from app.core.logging import get_logger
from app.infrastructure.db import session_scope
from app.infrastructure.db.repositories import SqlAlchemyTrackRepository
from app.infrastructure.storage.provider import get_file_storage
from app.infrastructure.transcode.ffmpeg import FfmpegTranscoder
log = get_logger("worker.transcode")
async def transcode_track(
_ctx: dict[str, Any],
*,
track_id: str,
quality: str = "high",
hls: bool = True,
) -> dict[str, Any]:
settings = get_settings()
tid = uuid.UUID(track_id)
async with session_scope() as session:
track = await SqlAlchemyTrackRepository(session).get_by_id(tid)
if track is None or track.storage_uri is None:
log.warning("transcode_skip_no_file", track_id=track_id)
return {"track_id": track_id, "opus": False, "hls": False}
storage = get_file_storage()
transcoder = FfmpegTranscoder(settings.ffmpeg_path)
root = settings.transcode_cache_path
did_opus = False
did_hls = False
async with storage.as_local_path(track.storage_uri) as src:
bitrate = bitrate_for_quality(quality)
if bitrate is not None:
dest = opus_path(root, tid, bitrate)
if not dest.exists():
await transcoder.to_opus(src, dest, bitrate_kbps=bitrate)
did_opus = True
if hls and not hls_playlist_path(root, tid).exists():
await transcoder.to_hls(src, hls_dir(root, tid), bitrate_kbps=HLS_BITRATE)
did_hls = True
log.info("transcode_done", track_id=track_id, opus=did_opus, hls=did_hls)
return {"track_id": track_id, "opus": did_opus, "hls": did_hls}