From 313af3a0706a296c325e761aeff7616167ebd2db Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=A6=D0=B2=D1=8B=D0=BB=D0=B5=D0=B2=20=D0=90=D0=BB=D0=B5?= =?UTF-8?q?=D0=BA=D1=81=D0=B0=D0=BD=D0=B4=D1=80=20=D0=92=D0=B0=D0=B4=D0=B8?= =?UTF-8?q?=D0=BC=D0=BE=D0=B2=D0=B8=D1=87?= Date: Tue, 28 Jul 2026 21:19:36 +0300 Subject: [PATCH] =?UTF-8?q?feat(transcode):=20optimize=20+=20on-the-fly=20?= =?UTF-8?q?quality=20+=20HLS=20streaming=20(=C2=A76.6)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- app/api/deps.py | 11 +++ app/api/schemas/transcode.py | 11 +++ app/api/v1/streaming.py | 70 +++++++++++++++-- app/api/v1/tracks.py | 28 +++++-- app/application/transcode_service.py | 99 ++++++++++++++++++++++++ app/core/config.py | 3 + app/domain/errors.py | 7 ++ app/domain/ports.py | 10 +++ app/infrastructure/transcode/__init__.py | 1 + app/infrastructure/transcode/ffmpeg.py | 67 ++++++++++++++++ app/workers/arq_worker.py | 2 + app/workers/queue.py | 22 ++++++ app/workers/tasks/transcode_task.py | 64 +++++++++++++++ 13 files changed, 384 insertions(+), 11 deletions(-) create mode 100644 app/api/schemas/transcode.py create mode 100644 app/application/transcode_service.py create mode 100644 app/infrastructure/transcode/__init__.py create mode 100644 app/infrastructure/transcode/ffmpeg.py create mode 100644 app/workers/tasks/transcode_task.py diff --git a/app/api/deps.py b/app/api/deps.py index 51ae72d..872bf83 100644 --- a/app/api/deps.py +++ b/app/api/deps.py @@ -24,6 +24,7 @@ from app.application.remote_library_service import RemoteLibraryService from app.application.streaming_service import StreamingService from app.application.subsonic_auth_service import SubsonicAuthService from app.application.sync_service import SyncService +from app.application.transcode_service import TranscodeService from app.application.upload_service import UploadService from app.application.user_service import UserService 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: """Wires the LRCLIB lyrics provider + cache repo (plan §6.7). LRCLIB is 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)] MetadataServiceDep = Annotated[MetadataEnrichmentService, Depends(get_metadata_service)] LyricsServiceDep = Annotated[LyricsService, Depends(get_lyrics_service)] +TranscodeServiceDep = Annotated[TranscodeService, Depends(get_transcode_service)] DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)] RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)] diff --git a/app/api/schemas/transcode.py b/app/api/schemas/transcode.py new file mode 100644 index 0000000..83ef09b --- /dev/null +++ b/app/api/schemas/transcode.py @@ -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 diff --git a/app/api/v1/streaming.py b/app/api/v1/streaming.py index 9584e0e..f9600c9 100644 --- a/app/api/v1/streaming.py +++ b/app/api/v1/streaming.py @@ -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 from typing import Annotated -from fastapi import APIRouter, Header -from fastapi.responses import StreamingResponse +from fastapi import APIRouter, Header, Query, Response +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"]) +_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}") async def stream_track( track_id: uuid.UUID, service: StreamingServiceDep, + transcode: TranscodeServiceDep, _user: StreamUser, 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) headers = { "Accept-Ranges": "bytes", "Content-Length": str(result.content_length), } - if result.is_partial: headers["Content-Range"] = f"bytes {result.start}-{result.end}/{result.total_size}" status_code = 206 @@ -37,3 +60,38 @@ async def stream_track( headers=headers, 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) diff --git a/app/api/v1/tracks.py b/app/api/v1/tracks.py index 1db3b2c..c57ae93 100644 --- a/app/api/v1/tracks.py +++ b/app/api/v1/tracks.py @@ -1,7 +1,7 @@ """Track endpoints.""" import uuid -from typing import Any +from typing import Annotated, Any from fastapi import APIRouter, Query, Response from fastapi.responses import StreamingResponse @@ -30,10 +30,12 @@ from app.api.schemas.track import ( TrackOut, 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.track import Artist, Track -from app.domain.errors import NotFoundError -from app.workers.queue import enqueue +from app.domain.errors import NotFoundError, ValidationError +from app.workers.queue import enqueue, enqueue_transcode 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: ... -@router.post("/{track_id}/optimize") -async def optimize_track(track_id: uuid.UUID, _: CurrentUser) -> Any: ... +@router.post("/{track_id}/optimize", status_code=202) +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") diff --git a/app/application/transcode_service.py b/app/application/transcode_service.py new file mode 100644 index 0000000..dde82b6 --- /dev/null +++ b/app/application/transcode_service.py @@ -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 diff --git a/app/core/config.py b/app/core/config.py index 2c7a673..95bcbcf 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -81,6 +81,9 @@ class Settings(BaseSettings): # -- media / storage -------------------------------------------------- media_path: Path = Path("/data/media") 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 # How many times the download worker retries a failed fetch (yt-dlp fails # often) before marking the job ``failed`` — exponential backoff between tries. diff --git a/app/domain/errors.py b/app/domain/errors.py index e2b1b69..27d1ad8 100644 --- a/app/domain/errors.py +++ b/app/domain/errors.py @@ -76,6 +76,13 @@ class StorageError(DomainError): 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): """Requested byte range cannot be satisfied.""" diff --git a/app/domain/ports.py b/app/domain/ports.py index ff71630..956c266 100644 --- a/app/domain/ports.py +++ b/app/domain/ports.py @@ -497,6 +497,16 @@ class CoverArtProvider(Protocol): 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): """Fetches lyrics from an external database (LRCLIB) by artist/title/album/ duration. Returns a hit or ``None`` (no match / service down), never raising.""" diff --git a/app/infrastructure/transcode/__init__.py b/app/infrastructure/transcode/__init__.py new file mode 100644 index 0000000..e9bff59 --- /dev/null +++ b/app/infrastructure/transcode/__init__.py @@ -0,0 +1 @@ +"""ffmpeg-based transcoding adapters (Group B / plan §6.6).""" diff --git a/app/infrastructure/transcode/ffmpeg.py b/app/infrastructure/transcode/ffmpeg.py new file mode 100644 index 0000000..404dc43 --- /dev/null +++ b/app/infrastructure/transcode/ffmpeg.py @@ -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}") diff --git a/app/workers/arq_worker.py b/app/workers/arq_worker.py index 94bdb41..ea51de9 100644 --- a/app/workers/arq_worker.py +++ b/app/workers/arq_worker.py @@ -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.import_task import scan_local_folder from app.workers.tasks.materialize_task import materialize_track +from app.workers.tasks.transcode_task import transcode_track log = get_logger("worker") @@ -35,6 +36,7 @@ class WorkerSettings: download_track, materialize_track, cleanup_storage, + transcode_track, ] on_startup = startup on_shutdown = shutdown diff --git a/app/workers/queue.py b/app/workers/queue.py index 75682e8..ca54c7e 100644 --- a/app/workers/queue.py +++ b/app/workers/queue.py @@ -59,6 +59,28 @@ async def enqueue_materialize(job_id: uuid.UUID) -> None: 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: """Best-effort enqueue of metadata enrichment for a freshly stored track. diff --git a/app/workers/tasks/transcode_task.py b/app/workers/tasks/transcode_task.py new file mode 100644 index 0000000..22b4dbd --- /dev/null +++ b/app/workers/tasks/transcode_task.py @@ -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}