Compare commits

..

2 Commits

Author SHA1 Message Date
Цвылев Александр Вадимович 313af3a070 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>
2026-07-28 21:19:36 +03:00
Цвылев Александр Вадимович 8271de34eb feat(lyrics): LRCLIB provider + cached lyrics endpoints (§6.7)
GET /tracks/{id}/lyrics (get-or-fetch, caches found/not_found with a 7-day
miss TTL) and POST /tracks/{id}/lyrics/refetch (force). Hexagonal wiring:
LyricsProvider/LyricsRepository ports, LrclibHttpClient adapter (keyless,
degrades to not_found on error), SqlAlchemyLyricsRepository (upsert on the
existing lyrics table), LyricsService, LyricsOut schema, deps.py factory.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 20:58:28 +03:00
19 changed files with 797 additions and 11 deletions
+28
View File
@@ -18,11 +18,13 @@ from sqlalchemy.ext.asyncio import AsyncSession
from app.application.auth_service import AuthService
from app.application.download_service import DownloadService
from app.application.lyrics_service import LyricsService
from app.application.metadata_service import MetadataEnrichmentService
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
@@ -38,6 +40,7 @@ from app.infrastructure.db.repositories import (
SqlAlchemyDownloadJobRepository,
SqlAlchemyHistoryRepository,
SqlAlchemyLikeRepository,
SqlAlchemyLyricsRepository,
SqlAlchemyPlaylistRepository,
SqlAlchemyRefreshTokenRepository,
SqlAlchemyTrackRepository,
@@ -46,6 +49,7 @@ from app.infrastructure.db.repositories import (
)
from app.infrastructure.metadata.acoustid import AcoustIdHttpClient
from app.infrastructure.metadata.fingerprint import FpcalcFingerprinter
from app.infrastructure.metadata.lrclib import LrclibHttpClient
from app.infrastructure.metadata.tags import MutagenTagReader
from app.infrastructure.sources.registry import SourceRegistry, build_source_registry
from app.infrastructure.storage.provider import get_file_storage
@@ -175,6 +179,28 @@ 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``."""
settings = get_settings()
return LyricsService(
lyrics=SqlAlchemyLyricsRepository(session),
tracks=SqlAlchemyTrackRepository(session),
artists=SqlAlchemyArtistRepository(session),
albums=SqlAlchemyAlbumRepository(session),
provider=LrclibHttpClient(user_agent=settings.musicbrainz_user_agent),
)
def get_download_service(session: SessionDep, storage: FileStorageDep) -> DownloadService:
return DownloadService(
jobs=SqlAlchemyDownloadJobRepository(session),
@@ -198,6 +224,8 @@ def get_remote_library_service(session: SessionDep) -> RemoteLibraryService:
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)]
+33
View File
@@ -0,0 +1,33 @@
"""Lyrics response schema (§6.7 / Now Playing lyrics panel).
Returns the raw LRC (``synced``) and/or ``plain`` text; the client parses LRC
timestamps for synced highlighting. A miss is a normal 200 with
``status="not_found"`` and null text — not an error — so the panel can render a
"no lyrics" state.
"""
import uuid
from pydantic import BaseModel
from app.domain.entities.lyrics import Lyrics
class LyricsOut(BaseModel):
track_id: uuid.UUID
status: str
source: str | None
synced: str | None
plain: str | None
synced_available: bool
@classmethod
def from_entity(cls, lyrics: Lyrics) -> LyricsOut:
return cls(
track_id=lyrics.track_id,
status=lyrics.status,
source=lyrics.source,
synced=lyrics.synced,
plain=lyrics.plain,
synced_available=lyrics.synced is not None,
)
+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
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)
+43 -5
View File
@@ -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
@@ -12,12 +12,14 @@ from app.api.deps import (
ArtistRepoDep,
CurrentUser,
FileStorageDep,
LyricsServiceDep,
MetadataServiceDep,
RemoteLibraryServiceDep,
StreamUser,
TrackRepoDep,
)
from app.api.schemas.download import DownloadJobOut
from app.api.schemas.lyrics import LyricsOut
from app.api.schemas.pagination import PagedResponse
from app.api.schemas.track import (
MaterializeResponse,
@@ -28,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"])
@@ -220,8 +224,42 @@ 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")
async def get_track_lyrics(
track_id: uuid.UUID, lyrics: LyricsServiceDep, _: CurrentUser
) -> LyricsOut:
"""Cached lyrics for the Now Playing panel (§6.7). A miss is a normal 200
with ``status="not_found"`` — the provider (LRCLIB) is queried at most once,
then the outcome is cached."""
return LyricsOut.from_entity(await lyrics.get_lyrics(track_id))
@router.post("/{track_id}/lyrics/refetch")
async def refetch_track_lyrics(
track_id: uuid.UUID, lyrics: LyricsServiceDep, _: CurrentUser
) -> LyricsOut:
"""Force a fresh provider lookup, bypassing the cache (user-triggered)."""
return LyricsOut.from_entity(await lyrics.get_lyrics(track_id, force=True))
@router.get("/{track_id}/cover")
+97
View File
@@ -0,0 +1,97 @@
"""Lyrics service (plan §6.7).
Get-or-fetch with caching: a track's lyrics are served from the DB when present;
on a miss (or an expired ``not_found``) we ask the provider (LRCLIB) once, then
cache the outcome. ``not_found`` is cached with a TTL so tracks that genuinely
have no lyrics aren't looked up on every play, but can eventually be retried.
Degrades gracefully: if the provider is unreachable the lookup just yields a
``not_found`` — the endpoint still returns 200 with empty lyrics, never an error.
"""
import datetime as dt
import uuid
from app.domain.entities.lyrics import Lyrics
from app.domain.errors import NotFoundError
from app.domain.ports import (
AlbumRepository,
ArtistRepository,
LyricsProvider,
LyricsRepository,
TrackRepository,
)
# Re-lookup a cached "not_found" only after this long — long enough not to spam
# the provider, short enough that lyrics added upstream eventually surface.
_NOT_FOUND_TTL = dt.timedelta(days=7)
_STATUS_FOUND = "found"
_STATUS_NOT_FOUND = "not_found"
class LyricsService:
def __init__(
self,
*,
lyrics: LyricsRepository,
tracks: TrackRepository,
artists: ArtistRepository,
albums: AlbumRepository,
provider: LyricsProvider,
) -> None:
self._lyrics = lyrics
self._tracks = tracks
self._artists = artists
self._albums = albums
self._provider = provider
async def get_lyrics(self, track_id: uuid.UUID, *, force: bool = False) -> Lyrics:
"""Return cached lyrics, fetching from the provider on a miss/expiry.
``force`` (the refetch endpoint) bypasses the cache entirely."""
cached = await self._lyrics.get(track_id)
if not force and cached is not None and self._is_fresh(cached):
return cached
return await self._fetch_and_cache(track_id)
def _is_fresh(self, cached: Lyrics) -> bool:
if cached.status == _STATUS_FOUND:
return True
if cached.status == _STATUS_NOT_FOUND:
return dt.datetime.now(dt.UTC) - cached.fetched_at < _NOT_FOUND_TTL
# "pending" (never fetched) → not fresh, go fetch.
return False
async def _fetch_and_cache(self, track_id: uuid.UUID) -> Lyrics:
track = await self._tracks.get_by_id(track_id)
if track is None:
raise NotFoundError(f"Track {track_id} not found.")
artist = await self._artists.get_by_id(track.artist_id)
album = (
await self._albums.get_by_id(track.album_id)
if track.album_id is not None
else None
)
result = await self._provider.fetch(
artist=artist.name if artist else "",
title=track.title,
album=album.title if album else None,
duration_seconds=track.duration_seconds,
)
if result is None:
return await self._lyrics.upsert(
track_id=track_id,
synced=None,
plain=None,
source=None,
status=_STATUS_NOT_FOUND,
)
return await self._lyrics.upsert(
track_id=track_id,
synced=result.synced,
plain=result.plain,
source=result.source,
status=_STATUS_FOUND,
)
+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_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.
+41
View File
@@ -0,0 +1,41 @@
"""Lyrics value objects (plan §6.7).
``Lyrics`` is the cached row for a track; ``LyricsResult`` is what a provider
(LRCLIB) returns for a lookup. Both cross the domain boundary — no framework
imports. Status values mirror ``LyricsStatus`` in the ORM enum ("found" /
"not_found" / "pending") but are kept as plain strings here so the domain stays
independent of the persistence layer.
"""
import datetime as dt
import uuid
from dataclasses import dataclass
@dataclass(frozen=True, slots=True)
class LyricsResult:
"""A provider hit: synced (timestamped LRC) and/or plain text."""
synced: str | None
plain: str | None
source: str
@dataclass(frozen=True, slots=True)
class Lyrics:
"""Cached lyrics for one track. ``status`` is ``found`` / ``not_found`` /
``pending``; ``not_found`` is cached too (with a TTL in the service) so a
track with no lyrics doesn't hammer the provider on every play."""
track_id: uuid.UUID
synced: str | None
plain: str | None
source: str | None
status: str
fetched_at: dt.datetime
@property
def has_lyrics(self) -> bool:
return self.status == "found" and (
self.synced is not None or self.plain is not None
)
+7
View File
@@ -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."""
+41
View File
@@ -29,6 +29,7 @@ from app.domain.entities import (
SubsonicCredentials,
User,
)
from app.domain.entities.lyrics import Lyrics, LyricsResult
from app.domain.entities.settings import UserSettings
from app.domain.entities.track import Artist, Track
from app.domain.sources import DownloadResult, RawMetadata, SearchResult, SourceFile, SourceInfo
@@ -494,3 +495,43 @@ class CoverArtProvider(Protocol):
def is_available(self) -> bool: ...
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."""
async def fetch(
self,
*,
artist: str,
title: str,
album: str | None,
duration_seconds: int | None,
) -> LyricsResult | None: ...
class LyricsRepository(Protocol):
"""Cached lyrics, one row per track. ``upsert`` also caches a ``not_found``
(empty text) so misses aren't re-fetched until the service's TTL lapses."""
async def get(self, track_id: uuid.UUID) -> Lyrics | None: ...
async def upsert(
self,
*,
track_id: uuid.UUID,
synced: str | None,
plain: str | None,
source: str | None,
status: str,
) -> Lyrics: ...
@@ -7,6 +7,7 @@ from app.infrastructure.db.repositories.download_job_repository import (
)
from app.infrastructure.db.repositories.history_repository import SqlAlchemyHistoryRepository
from app.infrastructure.db.repositories.like_repository import SqlAlchemyLikeRepository
from app.infrastructure.db.repositories.lyrics_repository import SqlAlchemyLyricsRepository
from app.infrastructure.db.repositories.playlist_repository import SqlAlchemyPlaylistRepository
from app.infrastructure.db.repositories.refresh_token_repository import (
SqlAlchemyRefreshTokenRepository,
@@ -23,6 +24,7 @@ __all__ = [
"SqlAlchemyDownloadJobRepository",
"SqlAlchemyHistoryRepository",
"SqlAlchemyLikeRepository",
"SqlAlchemyLyricsRepository",
"SqlAlchemyPlaylistRepository",
"SqlAlchemyRefreshTokenRepository",
"SqlAlchemyTrackRepository",
@@ -0,0 +1,72 @@
"""Lyrics repository — adapter over ``AsyncSession``.
One cached row per track (``track_id`` unique). ``upsert`` refreshes the row and
bumps ``fetched_at`` so the service's TTL is measured from the last fetch.
"""
import uuid
from sqlalchemy import func, select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession
from app.domain.entities.lyrics import Lyrics
from app.infrastructure.db.models.lyrics import LyricsModel
def _to_entity(row: LyricsModel) -> Lyrics:
return Lyrics(
track_id=row.track_id,
synced=row.synced,
plain=row.plain,
source=row.source,
status=row.status,
fetched_at=row.fetched_at,
)
class SqlAlchemyLyricsRepository:
def __init__(self, session: AsyncSession) -> None:
self._session = session
async def get(self, track_id: uuid.UUID) -> Lyrics | None:
row = await self._session.scalar(
select(LyricsModel).where(LyricsModel.track_id == track_id)
)
return _to_entity(row) if row is not None else None
async def upsert(
self,
*,
track_id: uuid.UUID,
synced: str | None,
plain: str | None,
source: str | None,
status: str,
) -> Lyrics:
values = {
"track_id": track_id,
"synced": synced,
"plain": plain,
"source": source,
"status": status,
"fetched_at": func.now(),
}
stmt = (
pg_insert(LyricsModel)
.values(**values)
.on_conflict_do_update(
index_elements=[LyricsModel.track_id],
set_={
"synced": synced,
"plain": plain,
"source": source,
"status": status,
"fetched_at": func.now(),
},
)
.returning(LyricsModel)
)
row = (await self._session.scalars(stmt)).one()
await self._session.flush()
return _to_entity(row)
+100
View File
@@ -0,0 +1,100 @@
"""LrclibHttpClient — fetches lyrics from LRCLIB (plan §6.7).
LRCLIB is a free, keyless lyrics database. ``/api/get`` does an exact match on
artist+track+album+duration; if that misses we fall back to ``/api/search`` and
take the best-scoring hit. Graceful degradation: any network/parse error →
``fetch`` returns ``None`` (the service then caches a ``not_found``), never
raising. No API key is needed, so this provider is always "available".
"""
import httpx
from app.core.logging import get_logger
from app.domain.entities.lyrics import LyricsResult
log = get_logger(__name__)
_BASE_URL = "https://lrclib.net"
_TIMEOUT_SECONDS = 10.0
_SOURCE = "lrclib"
class LrclibHttpClient:
"""Implements :class:`app.domain.ports.LyricsProvider`."""
def __init__(self, *, user_agent: str, base_url: str = _BASE_URL) -> None:
self._user_agent = user_agent
self._base_url = base_url.rstrip("/")
async def fetch(
self,
*,
artist: str,
title: str,
album: str | None,
duration_seconds: int | None,
) -> LyricsResult | None:
try:
async with httpx.AsyncClient(
timeout=_TIMEOUT_SECONDS,
headers={"User-Agent": self._user_agent},
base_url=self._base_url,
) as client:
hit = await self._get(client, artist, title, album, duration_seconds)
if hit is None:
hit = await self._search(client, artist, title)
except (httpx.HTTPError, ValueError) as exc:
log.warning("lrclib.fetch_failed", error=str(exc))
return None
return hit
async def _get(
self,
client: httpx.AsyncClient,
artist: str,
title: str,
album: str | None,
duration_seconds: int | None,
) -> LyricsResult | None:
"""Exact match via ``/api/get`` (404 when nothing matches exactly)."""
params = {"artist_name": artist, "track_name": title}
if album:
params["album_name"] = album
if duration_seconds is not None:
params["duration"] = str(duration_seconds)
resp = await client.get("/api/get", params=params)
if resp.status_code == httpx.codes.NOT_FOUND:
return None
resp.raise_for_status()
return _to_result(resp.json())
async def _search(
self, client: httpx.AsyncClient, artist: str, title: str
) -> LyricsResult | None:
"""Fuzzy fallback via ``/api/search`` — take the first usable hit."""
resp = await client.get(
"/api/search", params={"artist_name": artist, "track_name": title}
)
resp.raise_for_status()
results = resp.json()
if not isinstance(results, list):
return None
for item in results:
result = _to_result(item)
if result is not None:
return result
return None
def _to_result(payload: object) -> LyricsResult | None:
"""Map an LRCLIB record to a ``LyricsResult``. Instrumental tracks and empty
records yield ``None`` (nothing worth caching as "found")."""
if not isinstance(payload, dict):
return None
if payload.get("instrumental"):
return None
synced = payload.get("syncedLyrics") or None
plain = payload.get("plainLyrics") or None
if synced is None and plain is None:
return None
return LyricsResult(synced=synced, plain=plain, source=_SOURCE)
+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.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
+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))
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.
+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}