Compare commits

..

8 Commits

Author SHA1 Message Date
Цвылев Александр Вадимович fb7827d09c test: cover lyrics, transcode, and recommendation
DB-free unit tests for the three previously-untested features:
- lyrics: get-or-fetch caching, not_found TTL, force refetch, graceful miss;
  plus the LRC -> structured-lyrics serializer (timing, plain fallback, empty)
- transcode: path helpers, segment-name traversal guard, cache hit/miss, the
  unknown/not-downloaded 404 paths
- reco: ML path (order-preserving batched hydration) + metadata fallback for
  similar/radio, exclude handling, reason codes

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-29 11:07:22 +03:00
Цвылев Александр Вадимович c5a473fddf feat(subsonic): getLyricsBySongId over the native LyricsService
Adapts the cached LyricsService into the OpenSubsonic structured-lyrics shape:
synced LRC parsed into timed lines (ms), plain text as untimed fallback, empty
lyricsList when a track has none. Thin adapter — no logic duplicated. Also
parenthesize the getCoverArt except-tuple for clarity.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-29 11:07:21 +03:00
Цвылев Александр Вадимович 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
Цвылев Александр Вадимович d16c6085c9 perf(reco): batch id->entity hydration
radio/similar resolved ids one get_by_id at a time (N round-trips). Add
TrackRepository.get_many and use it (plus ArtistRepository.get_many) to
hydrate in a single query, preserving requested order.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-29 10:53:35 +03:00
Цвылев Александр Вадимович ed77acf0fe fix(likes): dedupe latest state with DISTINCT ON
max(created_at)+equality-join returned both rows when two like events shared
an identical created_at (realistic: likes carry a client-supplied timestamp
from offline sync), double-counting the track. Replace with DISTINCT ON
(track_id) + deterministic tiebreaker (created_at desc, id desc) across
get_latest_state / liked-tracks / count.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-29 10:53:34 +03:00
Цвылев Александр Вадимович 591a938e71 feat(reco): radio + similar with metadata fallback (§6.5)
POST /radio + /radio/next (stateless infinite feed: seed track / from-likes,
exploration mix, client-passed exclude_ids) and GET /tracks|artists/{id}/similar,
replacing the stubs. Recommender port abstracts the (future) ML service —
NullRecommender is wired now so RecommendationService always uses its metadata
heuristics (genre/artist similarity, random exploration filler), never a hard ML
dependency. Adds TrackRepository.list_similar/sample_playable + Artist.list_similar,
reason codes for the client, RemoteRecommender skeleton (TODO: ML contract).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 21:59:15 +03:00
Цвылев Александр Вадимович 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
33 changed files with 2233 additions and 81 deletions
+45
View File
@@ -18,11 +18,14 @@ 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.recommendation_service import RecommendationService
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 +41,7 @@ from app.infrastructure.db.repositories import (
SqlAlchemyDownloadJobRepository,
SqlAlchemyHistoryRepository,
SqlAlchemyLikeRepository,
SqlAlchemyLyricsRepository,
SqlAlchemyPlaylistRepository,
SqlAlchemyRefreshTokenRepository,
SqlAlchemyTrackRepository,
@@ -46,7 +50,9 @@ 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.ml.recommender import NullRecommender
from app.infrastructure.sources.registry import SourceRegistry, build_source_registry
from app.infrastructure.storage.provider import get_file_storage
from app.workers.queue import enqueue_download, enqueue_enrich, enqueue_materialize
@@ -175,6 +181,40 @@ 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_recommendation_service(session: SessionDep) -> RecommendationService:
"""Radio + similarity (§6.5). ML is optional and no service/contract exists
yet, so we wire ``NullRecommender`` — the service then uses its metadata
fallback. Swap in ``RemoteRecommender(ml_service_url)`` once ML lands."""
return RecommendationService(
recommender=NullRecommender(),
tracks=SqlAlchemyTrackRepository(session),
artists=SqlAlchemyArtistRepository(session),
likes=SqlAlchemyLikeRepository(session),
)
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 +238,11 @@ 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)]
RecommendationServiceDep = Annotated[
RecommendationService, Depends(get_recommendation_service)
]
DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)]
RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)]
+39 -4
View File
@@ -1,10 +1,11 @@
"""Subsonic media endpoints: stream, download, cover art.
"""Subsonic media endpoints: stream, download, cover art, lyrics.
``stream`` and ``download`` reuse :class:`StreamingService` (honouring HTTP
Range) — they return raw bytes, not the Subsonic envelope. Transcoding params
(``maxBitRate``/``format``) are accepted but ignored; the original file is served
(no in-request ffmpeg — CLAUDE.md). ``getCoverArt`` returns a placeholder until
the cover pipeline lands (the ``/api/v1`` cover endpoints are still stubs).
(no in-request ffmpeg — CLAUDE.md). ``getCoverArt`` serves the album cover (a
placeholder when there's none). ``getLyricsBySongId`` adapts the native
``LyricsService`` into the OpenSubsonic structured-lyrics shape.
"""
import base64
@@ -16,12 +17,17 @@ from fastapi.responses import Response, StreamingResponse
from app.api.covers import resolve_album_for_track, stream_cover
from app.api.deps import (
AlbumRepoDep,
ArtistRepoDep,
FileStorageDep,
LyricsServiceDep,
StreamingServiceDep,
SubsonicFormat,
SubsonicUser,
TrackRepoDep,
)
from app.api.rest.envelope import subsonic_response
from app.api.rest.ids import IdKind, decode_track, parse
from app.api.rest.serializers import structured_lyrics
from app.domain.entities.album import Album
from app.domain.errors import NotFoundError, StorageError
@@ -99,6 +105,35 @@ async def get_cover_art(
if album is not None and album.cover_path:
try:
return await stream_cover(storage, album.cover_path)
except NotFoundError, StorageError:
except (NotFoundError, StorageError):
pass
return Response(content=_PLACEHOLDER_PNG, media_type="image/png")
@router.api_route("/getLyricsBySongId", methods=["GET", "POST"])
@router.api_route("/getLyricsBySongId.view", methods=["GET", "POST"])
async def get_lyrics_by_song_id(
_user: SubsonicUser,
fmt: SubsonicFormat,
lyrics_service: LyricsServiceDep,
track_repo: TrackRepoDep,
artist_repo: ArtistRepoDep,
id: Annotated[str, Query()],
) -> Response:
# OpenSubsonic structured lyrics over the native LyricsService (§6.7). A miss
# is a normal empty ``lyricsList`` — the service degrades to not_found rather
# than raising, so clients get 200 either way.
track_id = decode_track(id)
track = await track_repo.get_by_id(track_id)
if track is None:
raise NotFoundError("Song not found.")
lyrics = await lyrics_service.get_lyrics(track_id)
artist = await artist_repo.get_by_id(track.artist_id)
return subsonic_response(
structured_lyrics(
lyrics,
display_artist=artist.name if artist is not None else "",
display_title=track.title,
),
fmt=fmt,
)
+56
View File
@@ -6,10 +6,17 @@ JSON equivalents). No business logic — they only reshape and rename.
"""
import datetime as dt
import re
from typing import Any
from app.api.rest.ids import encode_album, encode_artist, encode_track
from app.domain.entities import Album, Artist, Track
from app.domain.entities.lyrics import Lyrics
# One LRC timecode: ``[mm:ss.xx]`` / ``[mm:ss.xxx]`` (fraction optional). A line
# may carry several (the same words repeat at multiple times); metadata tags like
# ``[ar:..]`` don't match, so they're ignored.
_LRC_TAG_RE = re.compile(r"\[(\d+):(\d{1,2})(?:[.:](\d{1,3}))?\]")
# Suffix → MIME, for the ``contentType``/``suffix`` song attributes. A
# presentation detail (mirrors StreamingService's content-type negotiation).
@@ -90,3 +97,52 @@ def song_dict(
"type": "music",
"isVideo": False,
}
def _parse_lrc(synced: str) -> list[dict[str, Any]]:
"""LRC text → OpenSubsonic ``line`` dicts (``start`` in ms, ``value`` text),
ordered by time. Lines with no timecode (blank lines, metadata tags) drop out;
a timecode carrying several stamps yields one line per stamp."""
lines: list[tuple[int, str]] = []
for raw in synced.splitlines():
stamps = list(_LRC_TAG_RE.finditer(raw))
if not stamps:
continue
text = _LRC_TAG_RE.sub("", raw).strip()
for m in stamps:
minutes, seconds = int(m.group(1)), int(m.group(2))
# LRC fractions are centiseconds (2 digits) or ms (3); pad to ms.
ms = int((m.group(3) or "0").ljust(3, "0")[:3])
lines.append(((minutes * 60 + seconds) * 1000 + ms, text))
lines.sort(key=lambda pair: pair[0])
return [{"start": start, "value": text} for start, text in lines]
def structured_lyrics(
lyrics: Lyrics, *, display_artist: str, display_title: str
) -> dict[str, Any]:
"""OpenSubsonic ``getLyricsBySongId`` payload. Prefers synced (LRC) lines and
falls back to plain text; an empty ``lyricsList`` when the track has none."""
lines: list[dict[str, Any]] = []
synced = False
if lyrics.synced:
lines = _parse_lrc(lyrics.synced)
synced = bool(lines)
if not lines and lyrics.plain:
lines = [{"value": line} for line in lyrics.plain.splitlines()]
if not lines:
return {"lyricsList": {}}
return {
"lyricsList": {
"structuredLyrics": [
{
"displayArtist": display_artist,
"displayTitle": display_title,
"lang": "xxx",
"offset": 0,
"synced": synced,
"line": lines,
}
]
}
}
+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,
)
+43
View File
@@ -0,0 +1,43 @@
"""Radio + similarity response schemas (§6.5)."""
import uuid
from pydantic import BaseModel, Field
from app.api.schemas.artist import ArtistOut
from app.api.schemas.track import TrackOut
class RadioRequest(BaseModel):
"""Start or continue a radio. ``seed_track_id`` seeds from a track;
``from_likes`` seeds from the caller's likes. ``exclude_ids`` are already-
queued tracks to skip (the client drives the infinite feed). ``exploration``
biases familiar↔new."""
seed_track_id: uuid.UUID | None = None
from_likes: bool = False
exploration: float = Field(default=0.25, ge=0.0, le=1.0)
count: int = Field(default=20, ge=1, le=50)
exclude_ids: list[uuid.UUID] = Field(default_factory=list)
class RadioTrackOut(BaseModel):
track: TrackOut
# Short code the client localizes: ml | similar | from_likes | discover.
reason: str
class RadioResponse(BaseModel):
# Where the picks came from: "ml" or "metadata" (fallback).
source: str
tracks: list[RadioTrackOut]
class SimilarTracksOut(BaseModel):
source: str
tracks: list[TrackOut]
class SimilarArtistsOut(BaseModel):
source: str
artists: list[ArtistOut]
+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
+29 -3
View File
@@ -1,14 +1,20 @@
"""Artist endpoints."""
import uuid
from typing import Any
from fastapi import APIRouter, Query
from app.api.deps import AlbumRepoDep, ArtistRepoDep, CurrentUser, TrackRepoDep
from app.api.deps import (
AlbumRepoDep,
ArtistRepoDep,
CurrentUser,
RecommendationServiceDep,
TrackRepoDep,
)
from app.api.schemas.album import AlbumOut
from app.api.schemas.artist import ArtistOut
from app.api.schemas.pagination import PagedResponse
from app.api.schemas.radio import SimilarArtistsOut
from app.api.schemas.track import TrackOut
from app.api.v1.albums import _build_album_out
from app.api.v1.tracks import _build_track_out
@@ -124,4 +130,24 @@ async def get_artist_tracks(
@router.get("/{artist_id}/similar")
async def get_similar_artists(artist_id: uuid.UUID, _: CurrentUser) -> Any: ...
async def get_similar_artists(
artist_id: uuid.UUID,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
_: CurrentUser,
limit: int = Query(20, ge=1, le=100),
) -> SimilarArtistsOut:
"""Artists similar to this one (§6.5). ML when configured, else a shared-
genre metadata heuristic."""
source, artists = await service.similar_artists(artist_id, limit=limit)
items = [
ArtistOut(
id=a.id,
name=a.name,
album_count=await artist_repo.album_count(a.id),
track_count=await artist_repo.track_count(a.id),
created_at=a.created_at,
)
for a in artists
]
return SimilarArtistsOut(source=source, artists=items)
+72 -4
View File
@@ -1,15 +1,83 @@
"""Radio / continuous-mix endpoints. Degrades gracefully when ML service is down."""
"""Radio / continuous-mix endpoints (§6.5).
from typing import Any
Stateless: the client passes the seed + already-queued ids and pulls more as the
queue drains (offline-first infinite feed). Degrades gracefully when no ML
service is configured — the recommendation service falls back to metadata.
"""
from fastapi import APIRouter
from app.api.deps import (
AlbumRepoDep,
ArtistRepoDep,
CurrentUser,
RecommendationServiceDep,
)
from app.api.schemas.radio import RadioRequest, RadioResponse, RadioTrackOut
from app.api.v1.tracks import _build_track_out
from app.application.recommendation_service import RadioPick
router = APIRouter(prefix="/radio", tags=["radio"])
async def _to_response(
source: str,
picks: list[RadioPick],
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
tracks = [p.track for p in picks]
artist_ids = list({t.artist_id for t in tracks})
album_ids = list({t.album_id for t in tracks if t.album_id is not None})
artists = {a.id: a for a in await artist_repo.get_many(artist_ids)}
albums = {a.id: a for a in await album_repo.get_many(album_ids)}
outs = await _build_track_out(tracks, artists, albums)
return RadioResponse(
source=source,
tracks=[
RadioTrackOut(track=out, reason=pick.reason)
for out, pick in zip(outs, picks, strict=True)
],
)
async def _run_radio(
body: RadioRequest,
user: CurrentUser,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
source, picks = await service.radio(
user_id=user.id,
seed_track_id=body.seed_track_id,
from_likes=body.from_likes,
exploration=body.exploration,
limit=body.count,
exclude_ids=body.exclude_ids,
)
return await _to_response(source, picks, artist_repo, album_repo)
@router.post("")
async def start_radio() -> Any: ...
async def start_radio(
body: RadioRequest,
user: CurrentUser,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
"""Start a radio from a seed track or the caller's likes."""
return await _run_radio(body, user, service, artist_repo, album_repo)
@router.post("/next")
async def next_radio_track() -> Any: ...
async def next_radio_track(
body: RadioRequest,
user: CurrentUser,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
"""Fetch more tracks as the radio queue drains (pass ``exclude_ids``)."""
return await _run_radio(body, user, service, artist_repo, album_repo)
+65 -6
View File
@@ -1,30 +1,54 @@
"""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
import anyio
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 +61,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 = await anyio.to_thread.run_sync(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)
+68 -6
View File
@@ -1,8 +1,9 @@
"""Track endpoints."""
import uuid
from typing import Any
from typing import Annotated
import anyio
from fastapi import APIRouter, Query, Response
from fastapi.responses import StreamingResponse
@@ -12,13 +13,17 @@ from app.api.deps import (
ArtistRepoDep,
CurrentUser,
FileStorageDep,
LyricsServiceDep,
MetadataServiceDep,
RecommendationServiceDep,
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.radio import SimilarTracksOut
from app.api.schemas.track import (
MaterializeResponse,
MetadataApply,
@@ -28,10 +33,13 @@ from app.api.schemas.track import (
TrackOut,
TrackUpdate,
)
from app.api.schemas.transcode import OptimizeEnqueuedOut
from app.application.transcode_service import bitrate_for_quality, remove_track_cache
from app.core.config import get_settings
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"])
@@ -213,15 +221,69 @@ async def delete_track(
await track_repo.delete(track_id)
if track.storage_uri is not None:
await storage.delete(track.storage_uri)
# Drop any cached transcode renditions (Opus + HLS) so they don't dangle.
await anyio.to_thread.run_sync(
remove_track_cache, get_settings().transcode_cache_path, track_id
)
return Response(status_code=204)
@router.get("/{track_id}/similar")
async def get_similar_tracks(track_id: uuid.UUID, _: CurrentUser) -> Any: ...
async def get_similar_tracks(
track_id: uuid.UUID,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
_: CurrentUser,
limit: Annotated[int, Query(ge=1, le=100)] = 20,
) -> SimilarTracksOut:
"""Tracks similar to this one (§6.5). Uses ML when configured, else a
genre/artist metadata heuristic."""
source, tracks = await service.similar_tracks(track_id, limit=limit)
artist_ids = list({t.artist_id for t in tracks})
album_ids = list({t.album_id for t in tracks if t.album_id is not None})
artists = {a.id: a for a in await artist_repo.get_many(artist_ids)}
albums = {a.id: a for a in await album_repo.get_many(album_ids)}
outs = await _build_track_out(tracks, artists, albums)
return SimilarTracksOut(source=source, tracks=outs)
@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,
)
+189
View File
@@ -0,0 +1,189 @@
"""Recommendation / radio service (plan §6.5).
Tries the external ML recommender first; when it's unavailable or declines
(returns ``None``), falls back to metadata heuristics over the catalogue — so
similar/radio always work, worse, without ML (graceful-degradation invariant).
``reason`` values are short codes (``ml`` / ``similar`` / ``from_likes`` /
``discover``) the client localizes for the "why is this playing?" affordance.
"""
import random
import uuid
from dataclasses import dataclass
from app.domain.entities.track import Artist, Track
from app.domain.errors import NotFoundError
from app.domain.ports import (
ArtistRepository,
LikeRepository,
Recommender,
TrackRepository,
)
REASON_ML = "ml"
REASON_SIMILAR = "similar"
REASON_FROM_LIKES = "from_likes"
REASON_DISCOVER = "discover"
_LIKED_SEED_POOL = 50
@dataclass(frozen=True, slots=True)
class RadioPick:
track: Track
reason: str
class RecommendationService:
def __init__(
self,
*,
recommender: Recommender,
tracks: TrackRepository,
artists: ArtistRepository,
likes: LikeRepository,
) -> None:
self._recommender = recommender
self._tracks = tracks
self._artists = artists
self._likes = likes
# -- similar ---------------------------------------------------------------
async def similar_tracks(
self, track_id: uuid.UUID, *, limit: int
) -> tuple[str, list[Track]]:
seed = await self._tracks.get_by_id(track_id)
if seed is None:
raise NotFoundError(f"Track {track_id} not found.")
if self._recommender.is_available():
ids = await self._recommender.similar_track_ids(
track_id, limit=limit, exclude_ids=[track_id]
)
if ids is not None:
return REASON_ML, await self._hydrate_tracks(ids)
found = await self._tracks.list_similar(
genre=seed.genre,
artist_id=seed.artist_id,
exclude_ids=[track_id],
limit=limit,
)
return REASON_SIMILAR, found
async def similar_artists(
self, artist_id: uuid.UUID, *, limit: int
) -> tuple[str, list[Artist]]:
if await self._artists.get_by_id(artist_id) is None:
raise NotFoundError(f"Artist {artist_id} not found.")
if self._recommender.is_available():
ids = await self._recommender.similar_artist_ids(artist_id, limit=limit)
if ids is not None:
by_id = {a.id: a for a in await self._artists.get_many(ids)}
return REASON_ML, [by_id[i] for i in ids if i in by_id]
found = await self._artists.list_similar(artist_id=artist_id, limit=limit)
return REASON_SIMILAR, found
# -- radio -----------------------------------------------------------------
async def radio(
self,
*,
user_id: uuid.UUID,
seed_track_id: uuid.UUID | None,
from_likes: bool,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> tuple[str, list[RadioPick]]:
exploration = min(1.0, max(0.0, exploration))
if self._recommender.is_available():
ids = await self._recommender.radio_track_ids(
seed_track_id=seed_track_id,
exploration=exploration,
limit=limit,
exclude_ids=exclude_ids,
)
if ids is not None:
picks = [
RadioPick(track=t, reason=REASON_ML)
for t in await self._hydrate_tracks(ids)
]
return REASON_ML, picks
return "metadata", await self._radio_fallback(
user_id=user_id,
seed_track_id=seed_track_id,
from_likes=from_likes,
exploration=exploration,
limit=limit,
exclude_ids=exclude_ids,
)
async def _radio_fallback(
self,
*,
user_id: uuid.UUID,
seed_track_id: uuid.UUID | None,
from_likes: bool,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[RadioPick]:
exclude = list(dict.fromkeys(exclude_ids)) # de-dupe, keep order
explore_n = round(limit * exploration)
similar_n = limit - explore_n
picks: list[RadioPick] = []
seed, seed_reason = await self._resolve_seed(
user_id, seed_track_id, from_likes
)
if seed is not None and similar_n > 0:
for track in await self._tracks.list_similar(
genre=seed.genre,
artist_id=seed.artist_id,
exclude_ids=exclude,
limit=similar_n,
):
picks.append(RadioPick(track=track, reason=seed_reason))
exclude.append(track.id)
# Fill the remainder (exploration + any similarity shortfall) with random
# playable tracks — this is also the total fallback when there's no seed.
remaining = limit - len(picks)
if remaining > 0:
for track in await self._tracks.sample_playable(
exclude_ids=exclude, limit=remaining
):
picks.append(RadioPick(track=track, reason=REASON_DISCOVER))
exclude.append(track.id)
random.shuffle(picks)
return picks
async def _resolve_seed(
self,
user_id: uuid.UUID,
seed_track_id: uuid.UUID | None,
from_likes: bool,
) -> tuple[Track | None, str]:
if seed_track_id is not None:
return await self._tracks.get_by_id(seed_track_id), REASON_SIMILAR
if from_likes:
liked = await self._likes.list_liked_tracks(
user_id=user_id, limit=_LIKED_SEED_POOL, offset=0
)
if liked:
return random.choice(liked), REASON_FROM_LIKES
return None, REASON_DISCOVER
async def _hydrate_tracks(self, ids: list[uuid.UUID]) -> list[Track]:
"""Resolve ids → tracks preserving order, skipping any that vanished.
One batched query rather than N per-id round-trips."""
by_id = {t.id: t for t in await self._tracks.get_many(ids)}
return [by_id[i] for i in ids if i in by_id]
+110
View File
@@ -0,0 +1,110 @@
"""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
+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."""
+95
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
@@ -127,6 +128,12 @@ class ArtistRepository(Protocol):
async def get_by_id(self, artist_id: uuid.UUID) -> Artist | None: ...
async def get_many(self, ids: list[uuid.UUID]) -> list[Artist]: ...
async def list_similar(self, *, artist_id: uuid.UUID, limit: int) -> list[Artist]:
"""Artists sharing the seed artist's genres, ranked by overlap. Metadata
fallback for ``GET /artists/{id}/similar``. Defined before ``list`` so the
``list[Artist]`` annotation isn't shadowed by the method named ``list``."""
...
async def list(self, *, q: str | None, limit: int, offset: int) -> list[Artist]: ...
async def count(self, *, q: str | None) -> int: ...
async def album_count(self, artist_id: uuid.UUID) -> int: ...
@@ -135,6 +142,11 @@ class ArtistRepository(Protocol):
class TrackRepository(Protocol):
async def get_by_id(self, track_id: uuid.UUID) -> Track | None: ...
async def get_many(self, ids: list[uuid.UUID]) -> list[Track]:
"""Resolve multiple ids in one query (unordered) — batches the per-id
lookups radio/similar would otherwise fan out into N round-trips."""
...
async def get_by_source(self, source: str, source_id: str) -> Track | None: ...
async def add(
self,
@@ -170,6 +182,24 @@ class TrackRepository(Protocol):
# AlbumRepository below).
async def genres(self) -> list[tuple[str, int]]: ...
async def library_stats(self) -> LibraryStats: ...
async def list_similar(
self,
*,
genre: str | None,
artist_id: uuid.UUID,
exclude_ids: list[uuid.UUID],
limit: int,
) -> list[Track]:
"""Playable tracks resembling a seed (same genre and/or artist), ranked
by match strength then shuffled. The metadata fallback for §6.5 radio /
similar when no ML service is configured."""
...
async def sample_playable(
self, *, exclude_ids: list[uuid.UUID], limit: int
) -> list[Track]:
"""Random playable tracks — the exploration filler for radio."""
...
async def find_duplicate_groups(self) -> list[tuple[str, list[Track]]]: ...
async def list_by_metadata_status(
self, status: str, *, limit: int, offset: int
@@ -494,3 +524,68 @@ class CoverArtProvider(Protocol):
def is_available(self) -> bool: ...
async def fetch_release_group(self, release_group_mbid: str) -> CoverArt | None: ...
class Recommender(Protocol):
"""External ML recommender (plan §6.5, ``ML_SERVICE_URL``). Returns ordered
track/artist ids, or ``None`` when unavailable/erroring so the service falls
back to metadata heuristics — ML is never a hard dependency (invariant)."""
def is_available(self) -> bool: ...
async def similar_track_ids(
self, track_id: uuid.UUID, *, limit: int, exclude_ids: list[uuid.UUID]
) -> list[uuid.UUID] | None: ...
async def similar_artist_ids(
self, artist_id: uuid.UUID, *, limit: int
) -> list[uuid.UUID] | None: ...
async def radio_track_ids(
self,
*,
seed_track_id: uuid.UUID | None,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[uuid.UUID] | 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",
@@ -77,6 +77,26 @@ class SqlAlchemyArtistRepository:
)
return [_to_entity(r) for r in rows]
async def list_similar(self, *, artist_id: uuid.UUID, limit: int) -> list[Artist]:
# Artists whose tracks fall in the seed artist's genres, ranked by how
# many such tracks they have. Defined before ``list`` so the ``list[Artist]``
# return annotation isn't shadowed by the method named ``list``.
seed_genres = (
select(TrackModel.genre)
.where(TrackModel.artist_id == artist_id, TrackModel.genre.is_not(None))
.distinct()
)
stmt = (
select(ArtistModel)
.join(TrackModel, TrackModel.artist_id == ArtistModel.id)
.where(TrackModel.genre.in_(seed_genres), ArtistModel.id != artist_id)
.group_by(ArtistModel.id)
.order_by(func.count(TrackModel.id).desc())
.limit(limit)
)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def list(self, *, q: str | None, limit: int, offset: int) -> list[Artist]:
stmt = select(ArtistModel)
if q:
@@ -108,3 +128,4 @@ class SqlAlchemyArtistRepository:
.where(TrackModel.artist_id == artist_id)
)
).scalar_one()
@@ -6,7 +6,7 @@ Likes are an append-only event log. Current state = latest event per (user, trac
import datetime as dt
import uuid
from sqlalchemy import func, select
from sqlalchemy import Subquery, func, select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession
@@ -105,31 +105,52 @@ class SqlAlchemyLikeRepository:
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
def _latest_events_sq(
self, user_id: uuid.UUID, track_ids: list[uuid.UUID] | None
) -> Subquery:
"""The latest like event per ``track_id`` for a user, as a subquery.
``DISTINCT ON (track_id)`` with a deterministic tiebreaker (``created_at``
then ``id``) picks exactly one row per track even when two events share an
identical ``created_at`` — likes carry a client-supplied timestamp from
offline sync, so ties are realistic and a plain ``max()``+equality-join
would return both rows (double-counting the track)."""
stmt = select(
LikeModel.track_id,
LikeModel.value.label("value"),
LikeModel.created_at.label("created_at"),
).where(LikeModel.user_id == user_id)
if track_ids is not None:
stmt = stmt.where(LikeModel.track_id.in_(track_ids))
return (
stmt.distinct(LikeModel.track_id)
.order_by(
LikeModel.track_id,
LikeModel.created_at.desc(),
LikeModel.id.desc(),
)
.subquery()
)
async def get_latest_state(
self, *, user_id: uuid.UUID, track_ids: list[uuid.UUID]
) -> list[Like]:
if not track_ids:
return []
# Subquery: max(created_at) per track for this user
max_sq = (
select(
LikeModel.track_id,
func.max(LikeModel.created_at).label("latest"),
)
.where(LikeModel.user_id == user_id, LikeModel.track_id.in_(track_ids))
.group_by(LikeModel.track_id)
.subquery()
)
rows = (
(
await self._session.execute(
select(LikeModel)
.join(
max_sq,
(LikeModel.track_id == max_sq.c.track_id)
& (LikeModel.created_at == max_sq.c.latest),
.where(
LikeModel.user_id == user_id,
LikeModel.track_id.in_(track_ids),
)
.distinct(LikeModel.track_id)
.order_by(
LikeModel.track_id,
LikeModel.created_at.desc(),
LikeModel.id.desc(),
)
.where(LikeModel.user_id == user_id)
)
)
.scalars()
@@ -141,31 +162,14 @@ class SqlAlchemyLikeRepository:
self, *, user_id: uuid.UUID, limit: int, offset: int
) -> list[Track]:
# Tracks where the latest like event has value='like', ordered by like time desc
max_sq = (
select(
LikeModel.track_id,
func.max(LikeModel.created_at).label("latest"),
)
.where(LikeModel.user_id == user_id)
.group_by(LikeModel.track_id)
.subquery()
)
liked_sq = (
select(LikeModel.track_id, LikeModel.created_at)
.join(
max_sq,
(LikeModel.track_id == max_sq.c.track_id)
& (LikeModel.created_at == max_sq.c.latest),
)
.where(LikeModel.user_id == user_id, LikeModel.value == "like")
.subquery()
)
latest_sq = self._latest_events_sq(user_id, None)
rows = (
(
await self._session.execute(
select(TrackModel)
.join(liked_sq, TrackModel.id == liked_sq.c.track_id)
.order_by(liked_sq.c.created_at.desc())
.join(latest_sq, TrackModel.id == latest_sq.c.track_id)
.where(latest_sq.c.value == "like")
.order_by(latest_sq.c.created_at.desc())
.limit(limit)
.offset(offset)
)
@@ -176,25 +180,11 @@ class SqlAlchemyLikeRepository:
return [_track_to_entity(r) for r in rows]
async def count_liked_tracks(self, *, user_id: uuid.UUID) -> int:
max_sq = (
select(
LikeModel.track_id,
func.max(LikeModel.created_at).label("latest"),
)
.where(LikeModel.user_id == user_id)
.group_by(LikeModel.track_id)
.subquery()
)
liked_sq = (
select(LikeModel.track_id)
.join(
max_sq,
(LikeModel.track_id == max_sq.c.track_id)
& (LikeModel.created_at == max_sq.c.latest),
)
.where(LikeModel.user_id == user_id, LikeModel.value == "like")
.subquery()
)
latest_sq = self._latest_events_sq(user_id, None)
return (
await self._session.execute(select(func.count()).select_from(liked_sq))
await self._session.execute(
select(func.count())
.select_from(latest_sq)
.where(latest_sq.c.value == "like")
)
).scalar_one()
@@ -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)
@@ -3,7 +3,7 @@
import datetime as dt
import uuid
from sqlalchemy import func, select
from sqlalchemy import case, func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.domain.entities.storage import FormatBreakdown, LibraryStats
@@ -46,6 +46,16 @@ class SqlAlchemyTrackRepository:
row = await self._session.get(TrackModel, track_id)
return _to_entity(row) if row is not None else None
async def get_many(self, ids: list[uuid.UUID]) -> list[Track]:
if not ids:
return []
rows = (
(await self._session.execute(select(TrackModel).where(TrackModel.id.in_(ids))))
.scalars()
.all()
)
return [_to_entity(r) for r in rows]
async def get_by_source(self, source: str, source_id: str) -> Track | None:
row = (
await self._session.execute(
@@ -136,6 +146,42 @@ class SqlAlchemyTrackRepository:
).all()
return [(row.genre, row.cnt) for row in rows]
async def list_similar(
self,
*,
genre: str | None,
artist_id: uuid.UUID,
exclude_ids: list[uuid.UUID],
limit: int,
) -> list[Track]:
# Rank a same-genre hit above a same-artist hit; shuffle within a tier so
# the mix varies. Only playable (locally-stored) tracks are candidates.
if genre is not None:
match = or_(TrackModel.genre == genre, TrackModel.artist_id == artist_id)
score = case((TrackModel.genre == genre, 2), else_=0) + case(
(TrackModel.artist_id == artist_id, 1), else_=0
)
else:
match = TrackModel.artist_id == artist_id
score = case((TrackModel.artist_id == artist_id, 1), else_=0)
stmt = select(TrackModel).where(TrackModel.storage_uri.is_not(None), match)
if exclude_ids:
stmt = stmt.where(TrackModel.id.not_in(exclude_ids))
stmt = stmt.order_by(score.desc(), func.random()).limit(limit)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def sample_playable(
self, *, exclude_ids: list[uuid.UUID], limit: int
) -> list[Track]:
stmt = select(TrackModel).where(TrackModel.storage_uri.is_not(None))
if exclude_ids:
stmt = stmt.where(TrackModel.id.not_in(exclude_ids))
stmt = stmt.order_by(func.random()).limit(limit)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def library_stats(self) -> LibraryStats:
"""One-shot aggregate over the whole catalogue (no pagination). Defined
before ``list`` for the same shadowing reason as ``genres``."""
+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 @@
"""ML/recommender adapters (plan §6.5). ML is optional — see Recommender port."""
+110
View File
@@ -0,0 +1,110 @@
"""Recommender adapters (plan §6.5).
``NullRecommender`` is the default: no ML service, so every method reports
unavailable and the ``RecommendationService`` uses its metadata fallback. When
an embedding service exists, wire ``RemoteRecommender`` (skeleton below) to
``ML_SERVICE_URL`` — its exact request/response contract is TODO pending that
service. Both keep the invariant: ML is optional, never a hard dependency.
"""
import uuid
import httpx
from app.core.logging import get_logger
log = get_logger(__name__)
class NullRecommender:
"""No ML configured — always unavailable, always ``None`` (→ fallback)."""
def is_available(self) -> bool:
return False
async def similar_track_ids(
self, track_id: uuid.UUID, *, limit: int, exclude_ids: list[uuid.UUID]
) -> list[uuid.UUID] | None:
return None
async def similar_artist_ids(
self, artist_id: uuid.UUID, *, limit: int
) -> list[uuid.UUID] | None:
return None
async def radio_track_ids(
self,
*,
seed_track_id: uuid.UUID | None,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[uuid.UUID] | None:
return None
_TIMEOUT_SECONDS = 5.0
class RemoteRecommender:
"""HTTP client for an external embedding/recommender service.
TODO: the request/response schema below is a placeholder — align it with the
real ML service once its contract is known. Until then this stays unused
(``deps`` wires ``NullRecommender``). Every call is defensive: any error or a
malformed body returns ``None`` so the service degrades to metadata, matching
the graceful-degradation invariant.
"""
def __init__(self, base_url: str) -> None:
self._base_url = base_url.rstrip("/")
def is_available(self) -> bool:
return True
async def _post_ids(self, path: str, payload: dict[str, object]) -> list[uuid.UUID] | None:
try:
async with httpx.AsyncClient(timeout=_TIMEOUT_SECONDS) as client:
resp = await client.post(f"{self._base_url}{path}", json=payload)
resp.raise_for_status()
data = resp.json()
ids = data.get("track_ids") if isinstance(data, dict) else None
if not isinstance(ids, list):
return None
return [uuid.UUID(str(i)) for i in ids]
except (httpx.HTTPError, ValueError, KeyError) as exc:
log.warning("recommender.remote_failed", path=path, error=str(exc))
return None
async def similar_track_ids(
self, track_id: uuid.UUID, *, limit: int, exclude_ids: list[uuid.UUID]
) -> list[uuid.UUID] | None:
return await self._post_ids(
"/similar/tracks",
{"track_id": str(track_id), "limit": limit,
"exclude": [str(i) for i in exclude_ids]},
)
async def similar_artist_ids(
self, artist_id: uuid.UUID, *, limit: int
) -> list[uuid.UUID] | None:
# Artist recommendations aren't part of the placeholder track contract.
return None
async def radio_track_ids(
self,
*,
seed_track_id: uuid.UUID | None,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[uuid.UUID] | None:
return await self._post_ids(
"/radio",
{
"seed_track_id": str(seed_track_id) if seed_track_id else None,
"exploration": exploration,
"limit": limit,
"exclude": [str(i) for i in exclude_ids],
},
)
+1
View File
@@ -0,0 +1 @@
"""ffmpeg-based transcoding adapters (Group B / plan §6.6)."""
+71
View File
@@ -0,0 +1,71 @@
"""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}")
+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.
+82
View File
@@ -0,0 +1,82 @@
"""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 shutil
import uuid
from pathlib import Path
from typing import Any
import anyio
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")
def _swap_dir(tmp_dir: Path, final_dir: Path) -> None:
"""Publish a freshly-built HLS rendition atomically: replace the final dir in
one rename so a reader never sees a playlist referencing half-written
segments. Any stale partial at the destination is cleared first."""
if final_dir.exists():
shutil.rmtree(final_dir)
tmp_dir.replace(final_dir)
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():
# Build into a private temp dir, then swap it in atomically — the
# cached playlist only becomes visible once every segment is written.
final_dir = hls_dir(root, tid)
tmp_dir = final_dir.with_name(f"hls.{uuid.uuid4().hex}.tmp")
await transcoder.to_hls(src, tmp_dir, bitrate_kbps=HLS_BITRATE)
await anyio.to_thread.run_sync(_swap_dir, tmp_dir, final_dir)
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}
+247
View File
@@ -0,0 +1,247 @@
"""Lyrics service (§6.7) + OpenSubsonic structured-lyrics serializer — DB-free.
Service: get-or-fetch caching, not_found TTL, force refetch, graceful miss.
Serializer: LRC → timed lines, plain fallback, empty payload when absent.
"""
import datetime as dt
import uuid
import pytest
from app.api.rest.serializers import _parse_lrc, structured_lyrics
from app.application.lyrics_service import LyricsService
from app.domain.entities import Artist, Track
from app.domain.entities.album import Album
from app.domain.entities.lyrics import Lyrics, LyricsResult
from app.domain.errors import NotFoundError
def _now() -> dt.datetime:
return dt.datetime.now(dt.UTC)
def _lyrics(*, synced=None, plain=None, source=None, status="found", age_days=0) -> Lyrics:
return Lyrics(
track_id=uuid.uuid4(),
synced=synced,
plain=plain,
source=source,
status=status,
fetched_at=_now() - dt.timedelta(days=age_days),
)
def _track() -> Track:
now = _now()
return Track(
id=uuid.uuid4(),
title="Song",
artist_id=uuid.uuid4(),
album_id=None,
storage_uri="tracks/x.mp3",
file_format="mp3",
file_size=1,
source="upload",
source_id="x",
duration_seconds=200,
genre=None,
year=None,
track_number=None,
metadata_status="pending",
metadata_error=None,
enriched_at=None,
availability="local",
created_at=now,
updated_at=now,
)
class FakeLyricsRepo:
def __init__(self, cached: Lyrics | None) -> None:
self._cached = cached
self.upserts: list[dict[str, object]] = []
async def get(self, track_id: uuid.UUID) -> Lyrics | None:
return self._cached
async def upsert(self, *, track_id, synced, plain, source, status) -> Lyrics:
self.upserts.append({"status": status, "source": source})
return Lyrics(
track_id=track_id,
synced=synced,
plain=plain,
source=source,
status=status,
fetched_at=_now(),
)
class FakeTrackRepo:
def __init__(self, track: Track | None) -> None:
self._track = track
async def get_by_id(self, track_id: uuid.UUID) -> Track | None:
return self._track
class FakeArtistRepo:
def __init__(self, artist: Artist | None = None) -> None:
self._artist = artist
async def get_by_id(self, artist_id: uuid.UUID) -> Artist | None:
return self._artist
class FakeAlbumRepo:
async def get_by_id(self, album_id: uuid.UUID) -> Album | None:
return None
class FakeProvider:
def __init__(self, result: LyricsResult | None) -> None:
self._result = result
self.calls = 0
async def fetch(self, *, artist, title, album, duration_seconds) -> LyricsResult | None:
self.calls += 1
return self._result
def _service(*, cached, track, provider) -> tuple[LyricsService, FakeLyricsRepo, FakeProvider]:
repo = FakeLyricsRepo(cached)
prov = FakeProvider(provider)
svc = LyricsService(
lyrics=repo,
tracks=FakeTrackRepo(track),
artists=FakeArtistRepo(),
albums=FakeAlbumRepo(),
provider=prov,
)
return svc, repo, prov
# -- service ------------------------------------------------------------------
async def test_fresh_found_is_served_from_cache_without_provider() -> None:
cached = _lyrics(plain="hi", status="found")
svc, _repo, prov = _service(cached=cached, track=_track(), provider=None)
out = await svc.get_lyrics(cached.track_id)
assert out is cached
assert prov.calls == 0
async def test_miss_fetches_and_caches_found() -> None:
track = _track()
result = LyricsResult(synced="[00:01.00]hi", plain="hi", source="lrclib")
svc, repo, prov = _service(cached=None, track=track, provider=result)
out = await svc.get_lyrics(track.id)
assert prov.calls == 1
assert out.status == "found"
assert repo.upserts[-1]["status"] == "found"
async def test_provider_miss_caches_not_found() -> None:
track = _track()
svc, repo, prov = _service(cached=None, track=track, provider=None)
out = await svc.get_lyrics(track.id)
assert prov.calls == 1
assert out.status == "not_found"
assert repo.upserts[-1]["status"] == "not_found"
async def test_stale_not_found_is_refetched() -> None:
# A not_found older than the 7-day TTL is retried against the provider.
stale = _lyrics(status="not_found", age_days=8)
result = LyricsResult(synced=None, plain="found now", source="lrclib")
svc, _repo, prov = _service(cached=stale, track=_track(), provider=result)
out = await svc.get_lyrics(stale.track_id)
assert prov.calls == 1
assert out.status == "found"
async def test_fresh_not_found_is_not_refetched() -> None:
fresh = _lyrics(status="not_found", age_days=1)
svc, _repo, prov = _service(cached=fresh, track=_track(), provider=None)
out = await svc.get_lyrics(fresh.track_id)
assert out is fresh
assert prov.calls == 0
async def test_force_bypasses_a_fresh_cache() -> None:
cached = _lyrics(plain="old", status="found")
result = LyricsResult(synced=None, plain="new", source="lrclib")
svc, _repo, prov = _service(cached=cached, track=_track(), provider=result)
out = await svc.get_lyrics(cached.track_id, force=True)
assert prov.calls == 1
assert out.plain == "new"
async def test_unknown_track_raises() -> None:
svc, _repo, _prov = _service(cached=None, track=None, provider=None)
with pytest.raises(NotFoundError):
await svc.get_lyrics(uuid.uuid4())
# -- serializer ---------------------------------------------------------------
def test_parse_lrc_times_and_order() -> None:
lines = _parse_lrc("[00:12.34]second\n[00:01.00]first\n[bad]meta\n\n")
# Metadata/blank lines dropped; output ordered by time; centiseconds → ms.
assert lines == [
{"start": 1000, "value": "first"},
{"start": 12340, "value": "second"},
]
def test_parse_lrc_repeated_stamps_on_one_line() -> None:
lines = _parse_lrc("[00:01.00][00:05.00]chorus")
assert lines == [
{"start": 1000, "value": "chorus"},
{"start": 5000, "value": "chorus"},
]
def test_structured_lyrics_prefers_synced() -> None:
out = structured_lyrics(
_lyrics(synced="[00:01.00]hi", plain="hi", status="found"),
display_artist="A",
display_title="T",
)
entry = out["lyricsList"]["structuredLyrics"][0]
assert entry["synced"] is True
assert entry["displayArtist"] == "A"
assert entry["line"] == [{"start": 1000, "value": "hi"}]
def test_structured_lyrics_falls_back_to_plain() -> None:
out = structured_lyrics(
_lyrics(synced=None, plain="line one\nline two", status="found"),
display_artist="A",
display_title="T",
)
entry = out["lyricsList"]["structuredLyrics"][0]
assert entry["synced"] is False
assert entry["line"] == [{"value": "line one"}, {"value": "line two"}]
def test_structured_lyrics_empty_when_absent() -> None:
out = structured_lyrics(
_lyrics(synced=None, plain=None, status="not_found"),
display_artist="A",
display_title="T",
)
assert out == {"lyricsList": {}}
def test_structured_lyrics_unparseable_synced_falls_back_to_plain() -> None:
# synced present but no valid timecodes → treat as plain, synced=False.
out = structured_lyrics(
_lyrics(synced="no timecodes here", plain="plain text", status="found"),
display_artist="A",
display_title="T",
)
entry = out["lyricsList"]["structuredLyrics"][0]
assert entry["synced"] is False
assert entry["line"] == [{"value": "plain text"}]
+248
View File
@@ -0,0 +1,248 @@
"""RecommendationService (§6.5) — DB-free, in-memory fakes.
Covers the two paths that matter: the ML recommender when available (reason
``ml``, order-preserving batched hydration) and the metadata fallback when it
declines (similar/radio still work — the graceful-degradation invariant).
"""
import datetime as dt
import uuid
import pytest
from app.application.recommendation_service import (
REASON_FROM_LIKES,
REASON_ML,
REASON_SIMILAR,
RecommendationService,
)
from app.domain.entities import Artist, Track
from app.domain.errors import NotFoundError
from app.infrastructure.ml.recommender import NullRecommender
def _now() -> dt.datetime:
return dt.datetime.now(dt.UTC)
def _track(*, genre: str | None = "rock", artist_id: uuid.UUID | None = None) -> Track:
now = _now()
return Track(
id=uuid.uuid4(),
title="T",
artist_id=artist_id or uuid.uuid4(),
album_id=None,
storage_uri="tracks/x.mp3",
file_format="mp3",
file_size=1,
source="upload",
source_id="x",
duration_seconds=None,
genre=genre,
year=None,
track_number=None,
metadata_status="pending",
metadata_error=None,
enriched_at=None,
availability="local",
created_at=now,
updated_at=now,
)
def _artist() -> Artist:
now = _now()
return Artist(
id=uuid.uuid4(), name="A", source=None, source_id=None, created_at=now, updated_at=now
)
class FakeTrackRepo:
def __init__(self, tracks: list[Track]) -> None:
self._by_id = {t.id: t for t in tracks}
self.similar: list[Track] = []
self.sample: list[Track] = []
self.similar_calls: list[dict[str, object]] = []
async def get_by_id(self, track_id: uuid.UUID) -> Track | None:
return self._by_id.get(track_id)
async def get_many(self, ids: list[uuid.UUID]) -> list[Track]:
# Deliberately unordered (mirrors a real ``WHERE id IN`` query) so the
# service is responsible for restoring request order.
return [self._by_id[i] for i in reversed(ids) if i in self._by_id]
async def list_similar(self, *, genre, artist_id, exclude_ids, limit) -> list[Track]:
self.similar_calls.append({"exclude_ids": list(exclude_ids), "limit": limit})
return [t for t in self.similar if t.id not in exclude_ids][:limit]
async def sample_playable(self, *, exclude_ids, limit) -> list[Track]:
return [t for t in self.sample if t.id not in exclude_ids][:limit]
class FakeArtistRepo:
def __init__(self, artists: list[Artist]) -> None:
self._by_id = {a.id: a for a in artists}
self.similar: list[Artist] = []
async def get_by_id(self, artist_id: uuid.UUID) -> Artist | None:
return self._by_id.get(artist_id)
async def get_many(self, ids: list[uuid.UUID]) -> list[Artist]:
return [self._by_id[i] for i in reversed(ids) if i in self._by_id]
async def list_similar(self, *, artist_id, limit) -> list[Artist]:
return self.similar[:limit]
class FakeLikeRepo:
def __init__(self, liked: list[Track]) -> None:
self.liked = liked
async def list_liked_tracks(self, *, user_id, limit, offset) -> list[Track]:
return self.liked[offset : offset + limit]
class StubRecommender:
"""An 'available' ML recommender returning fixed ids for the ML path."""
def __init__(self, ids: list[uuid.UUID] | None) -> None:
self._ids = ids
def is_available(self) -> bool:
return True
async def similar_track_ids(self, track_id, *, limit, exclude_ids):
return self._ids
async def similar_artist_ids(self, artist_id, *, limit):
return self._ids
async def radio_track_ids(self, *, seed_track_id, exploration, limit, exclude_ids):
return self._ids
def _service(tracks, artists, likes, recommender) -> RecommendationService:
return RecommendationService(
recommender=recommender, tracks=tracks, artists=artists, likes=likes
)
# -- similar ------------------------------------------------------------------
async def test_similar_tracks_unknown_seed_raises() -> None:
svc = _service(FakeTrackRepo([]), FakeArtistRepo([]), FakeLikeRepo([]), NullRecommender())
with pytest.raises(NotFoundError):
await svc.similar_tracks(uuid.uuid4(), limit=5)
async def test_similar_tracks_falls_back_to_metadata() -> None:
seed = _track()
neighbours = [_track(), _track()]
repo = FakeTrackRepo([seed, *neighbours])
repo.similar = neighbours
svc = _service(repo, FakeArtistRepo([]), FakeLikeRepo([]), NullRecommender())
reason, found = await svc.similar_tracks(seed.id, limit=5)
assert reason == REASON_SIMILAR
assert [t.id for t in found] == [n.id for n in neighbours]
# The seed itself is always excluded from its own neighbours.
assert repo.similar_calls[-1]["exclude_ids"] == [seed.id]
async def test_similar_tracks_uses_ml_and_preserves_order() -> None:
seed = _track()
a, b, c = _track(), _track(), _track()
repo = FakeTrackRepo([seed, a, b, c])
# ML returns a specific order; hydration must preserve it despite get_many
# returning rows unordered, and skip ids that no longer exist.
ml_ids = [c.id, uuid.uuid4(), a.id, b.id]
svc = _service(repo, FakeArtistRepo([]), FakeLikeRepo([]), StubRecommender(ml_ids))
reason, found = await svc.similar_tracks(seed.id, limit=10)
assert reason == REASON_ML
assert [t.id for t in found] == [c.id, a.id, b.id]
async def test_similar_artists_falls_back_to_metadata() -> None:
seed = _artist()
neighbours = [_artist(), _artist()]
repo = FakeArtistRepo([seed])
repo.similar = neighbours
svc = _service(FakeTrackRepo([]), repo, FakeLikeRepo([]), NullRecommender())
reason, found = await svc.similar_artists(seed.id, limit=5)
assert reason == REASON_SIMILAR
assert [a.id for a in found] == [n.id for n in neighbours]
# -- radio --------------------------------------------------------------------
async def test_radio_from_likes_seeds_and_fills(monkeypatch) -> None:
liked = _track()
similar = [_track(), _track()]
explore = [_track(), _track(), _track()]
repo = FakeTrackRepo([liked, *similar, *explore])
repo.similar = similar
repo.sample = explore
svc = _service(repo, FakeArtistRepo([]), FakeLikeRepo([liked]), NullRecommender())
# Deterministic seed choice + no shuffling for a stable assertion.
monkeypatch.setattr("app.application.recommendation_service.random.choice", lambda s: s[0])
monkeypatch.setattr("app.application.recommendation_service.random.shuffle", lambda s: None)
reason, picks = await svc.radio(
user_id=uuid.uuid4(),
seed_track_id=None,
from_likes=True,
exploration=0.5,
limit=4,
exclude_ids=[],
)
assert reason == "metadata"
ids = [p.track.id for p in picks]
# 4 picks, all distinct, no seed repeated, mix of similar + discover.
assert len(ids) == 4
assert len(set(ids)) == 4
reasons = {p.reason for p in picks}
assert REASON_FROM_LIKES in reasons # similar picks carry the seed's reason
async def test_radio_respects_exclude_ids(monkeypatch) -> None:
similar = [_track(), _track()]
explore = [_track(), _track()]
repo = FakeTrackRepo([*similar, *explore])
repo.similar = similar
repo.sample = explore
svc = _service(repo, FakeArtistRepo([]), FakeLikeRepo([]), NullRecommender())
monkeypatch.setattr("app.application.recommendation_service.random.shuffle", lambda s: None)
excluded = similar[0].id
_, picks = await svc.radio(
user_id=uuid.uuid4(),
seed_track_id=None,
from_likes=False,
exploration=1.0,
limit=4,
exclude_ids=[excluded],
)
assert excluded not in {p.track.id for p in picks}
async def test_radio_uses_ml_when_available(monkeypatch) -> None:
a, b = _track(), _track()
repo = FakeTrackRepo([a, b])
svc = _service(repo, FakeArtistRepo([]), FakeLikeRepo([]), StubRecommender([b.id, a.id]))
monkeypatch.setattr("app.application.recommendation_service.random.shuffle", lambda s: None)
reason, picks = await svc.radio(
user_id=uuid.uuid4(),
seed_track_id=a.id,
from_likes=False,
exploration=0.3,
limit=5,
exclude_ids=[],
)
assert reason == REASON_ML
assert [p.track.id for p in picks] == [b.id, a.id]
assert all(p.reason == REASON_ML for p in picks)
+157
View File
@@ -0,0 +1,157 @@
"""Transcode cache-path helpers + request-side service (§6.6) — DB-free.
Covers the pure path math, the segment-name traversal guard, cache
hit/miss lookups, and the "not yet downloaded" / unknown-track 404 paths.
"""
import datetime as dt
import uuid
from pathlib import Path
import pytest
from app.application.transcode_service import (
QUALITY_BITRATE,
TranscodeService,
bitrate_for_quality,
hls_playlist_path,
hls_segment_path,
opus_path,
remove_track_cache,
track_cache_dir,
)
from app.domain.entities import Track
from app.domain.errors import NotFoundError
def _track(*, storage_uri: str | None = "tracks/aa/song.mp3") -> Track:
now = dt.datetime.now(dt.UTC)
return Track(
id=uuid.uuid4(),
title="Song",
artist_id=uuid.uuid4(),
album_id=None,
storage_uri=storage_uri,
file_format="mp3",
file_size=1,
source="upload",
source_id="x",
duration_seconds=None,
genre=None,
year=None,
track_number=None,
metadata_status="pending",
metadata_error=None,
enriched_at=None,
availability="local",
created_at=now,
updated_at=now,
)
class FakeTrackRepo:
def __init__(self, track: Track | None) -> None:
self._track = track
async def get_by_id(self, track_id: uuid.UUID) -> Track | None:
return self._track
# -- pure helpers -------------------------------------------------------------
def test_bitrate_for_quality_known_and_unknown() -> None:
assert bitrate_for_quality("high") == QUALITY_BITRATE["high"]
assert bitrate_for_quality("medium") == 96
assert bitrate_for_quality("low") == 64
# "original" and anything unrecognised mean "serve the master, no transcode".
assert bitrate_for_quality("original") is None
assert bitrate_for_quality("nonsense") is None
def test_path_helpers_are_under_the_track_dir(tmp_path: Path) -> None:
tid = uuid.uuid4()
cache = track_cache_dir(tmp_path, tid)
assert cache == tmp_path / str(tid)
assert opus_path(tmp_path, tid, 128) == cache / "opus_128.opus"
assert hls_playlist_path(tmp_path, tid) == cache / "hls" / "playlist.m3u8"
@pytest.mark.parametrize(
"name",
["seg_000.ts", "seg_1234.ts"],
)
def test_hls_segment_path_accepts_valid_names(tmp_path: Path, name: str) -> None:
tid = uuid.uuid4()
resolved = hls_segment_path(tmp_path, tid, name)
assert resolved == track_cache_dir(tmp_path, tid) / "hls" / name
@pytest.mark.parametrize(
"name",
["../../etc/passwd", "seg_.ts", "seg_00.tsx", "playlist.m3u8", "seg_00.ts", "..", ""],
)
def test_hls_segment_path_rejects_traversal_and_junk(tmp_path: Path, name: str) -> None:
assert hls_segment_path(tmp_path, uuid.uuid4(), name) is None
def test_remove_track_cache_is_best_effort(tmp_path: Path) -> None:
tid = uuid.uuid4()
d = track_cache_dir(tmp_path, tid)
(d / "hls").mkdir(parents=True)
(d / "opus_128.opus").write_bytes(b"x")
remove_track_cache(tmp_path, tid)
assert not d.exists()
# Idempotent: deleting an already-absent cache doesn't raise.
remove_track_cache(tmp_path, tid)
# -- request-side service -----------------------------------------------------
async def test_resolve_quality_file_hit_and_miss(tmp_path: Path) -> None:
track = _track()
svc = TranscodeService(tracks=FakeTrackRepo(track), cache_root=tmp_path)
# Miss: nothing on disk yet.
assert await svc.resolve_quality_file(track.id, "high") is None
# "original" never has a rendition.
assert await svc.resolve_quality_file(track.id, "original") is None
# Hit: once the rendition exists it's returned.
path = opus_path(tmp_path, track.id, 128)
path.parent.mkdir(parents=True)
path.write_bytes(b"opus")
assert await svc.resolve_quality_file(track.id, "high") == path
async def test_hls_playlist_requires_downloaded_track(tmp_path: Path) -> None:
tid = uuid.uuid4()
# Unknown track → 404.
with pytest.raises(NotFoundError):
await TranscodeService(tracks=FakeTrackRepo(None), cache_root=tmp_path).hls_playlist(tid)
# Known but not downloaded (no storage_uri) → 404.
pending = _track(storage_uri=None)
with pytest.raises(NotFoundError):
await TranscodeService(tracks=FakeTrackRepo(pending), cache_root=tmp_path).hls_playlist(
pending.id
)
async def test_hls_playlist_miss_then_hit(tmp_path: Path) -> None:
track = _track()
svc = TranscodeService(tracks=FakeTrackRepo(track), cache_root=tmp_path)
assert await svc.hls_playlist(track.id) is None
playlist = hls_playlist_path(tmp_path, track.id)
playlist.parent.mkdir(parents=True)
playlist.write_text("#EXTM3U")
assert await svc.hls_playlist(track.id) == playlist
def test_hls_segment_returns_none_for_missing_file(tmp_path: Path) -> None:
svc = TranscodeService(tracks=FakeTrackRepo(_track()), cache_root=tmp_path)
# Valid name but no file on disk.
assert svc.hls_segment(uuid.uuid4(), "seg_000.ts") is None
# Junk name is refused outright.
assert svc.hls_segment(uuid.uuid4(), "../secret") is None