Compare commits

..

5 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
14 changed files with 861 additions and 71 deletions
+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 ``stream`` and ``download`` reuse :class:`StreamingService` (honouring HTTP
Range) — they return raw bytes, not the Subsonic envelope. Transcoding params Range) — they return raw bytes, not the Subsonic envelope. Transcoding params
(``maxBitRate``/``format``) are accepted but ignored; the original file is served (``maxBitRate``/``format``) are accepted but ignored; the original file is served
(no in-request ffmpeg — CLAUDE.md). ``getCoverArt`` returns a placeholder until (no in-request ffmpeg — CLAUDE.md). ``getCoverArt`` serves the album cover (a
the cover pipeline lands (the ``/api/v1`` cover endpoints are still stubs). placeholder when there's none). ``getLyricsBySongId`` adapts the native
``LyricsService`` into the OpenSubsonic structured-lyrics shape.
""" """
import base64 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.covers import resolve_album_for_track, stream_cover
from app.api.deps import ( from app.api.deps import (
AlbumRepoDep, AlbumRepoDep,
ArtistRepoDep,
FileStorageDep, FileStorageDep,
LyricsServiceDep,
StreamingServiceDep, StreamingServiceDep,
SubsonicFormat,
SubsonicUser, SubsonicUser,
TrackRepoDep, TrackRepoDep,
) )
from app.api.rest.envelope import subsonic_response
from app.api.rest.ids import IdKind, decode_track, parse 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.entities.album import Album
from app.domain.errors import NotFoundError, StorageError 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: if album is not None and album.cover_path:
try: try:
return await stream_cover(storage, album.cover_path) return await stream_cover(storage, album.cover_path)
except NotFoundError, StorageError: except (NotFoundError, StorageError):
pass pass
return Response(content=_PLACEHOLDER_PNG, media_type="image/png") 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 datetime as dt
import re
from typing import Any from typing import Any
from app.api.rest.ids import encode_album, encode_artist, encode_track from app.api.rest.ids import encode_album, encode_artist, encode_track
from app.domain.entities import Album, Artist, 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 # Suffix → MIME, for the ``contentType``/``suffix`` song attributes. A
# presentation detail (mirrors StreamingService's content-type negotiation). # presentation detail (mirrors StreamingService's content-type negotiation).
@@ -90,3 +97,52 @@ def song_dict(
"type": "music", "type": "music",
"isVideo": False, "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,
}
]
}
}
+2 -1
View File
@@ -10,6 +10,7 @@ import re
import uuid import uuid
from typing import Annotated from typing import Annotated
import anyio
from fastapi import APIRouter, Header, Query, Response from fastapi import APIRouter, Header, Query, Response
from fastapi.responses import FileResponse, StreamingResponse from fastapi.responses import FileResponse, StreamingResponse
@@ -78,7 +79,7 @@ async def stream_hls_playlist(
await enqueue_transcode_quiet(track_id, hls=True) await enqueue_transcode_quiet(track_id, hls=True)
raise NotFoundError("HLS rendition is being prepared; retry shortly.") raise NotFoundError("HLS rendition is being prepared; retry shortly.")
body = path.read_text() body = await anyio.to_thread.run_sync(path.read_text)
if token: if token:
body = _SEGMENT_LINE_RE.sub(rf"\1?token={token}", body) body = _SEGMENT_LINE_RE.sub(rf"\1?token={token}", body)
return Response(body, media_type=_HLS_PLAYLIST_TYPE) return Response(body, media_type=_HLS_PLAYLIST_TYPE)
+7 -1
View File
@@ -3,6 +3,7 @@
import uuid import uuid
from typing import Annotated from typing import Annotated
import anyio
from fastapi import APIRouter, Query, Response from fastapi import APIRouter, Query, Response
from fastapi.responses import StreamingResponse from fastapi.responses import StreamingResponse
@@ -33,7 +34,8 @@ from app.api.schemas.track import (
TrackUpdate, TrackUpdate,
) )
from app.api.schemas.transcode import OptimizeEnqueuedOut from app.api.schemas.transcode import OptimizeEnqueuedOut
from app.application.transcode_service import bitrate_for_quality 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.album import Album
from app.domain.entities.track import Artist, Track from app.domain.entities.track import Artist, Track
from app.domain.errors import NotFoundError, ValidationError from app.domain.errors import NotFoundError, ValidationError
@@ -219,6 +221,10 @@ async def delete_track(
await track_repo.delete(track_id) await track_repo.delete(track_id)
if track.storage_uri is not None: if track.storage_uri is not None:
await storage.delete(track.storage_uri) 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) return Response(status_code=204)
+6 -4
View File
@@ -82,8 +82,8 @@ class RecommendationService:
if self._recommender.is_available(): if self._recommender.is_available():
ids = await self._recommender.similar_artist_ids(artist_id, limit=limit) ids = await self._recommender.similar_artist_ids(artist_id, limit=limit)
if ids is not None: if ids is not None:
found = [a for i in ids if (a := await self._artists.get_by_id(i))] by_id = {a.id: a for a in await self._artists.get_many(ids)}
return REASON_ML, found 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) found = await self._artists.list_similar(artist_id=artist_id, limit=limit)
return REASON_SIMILAR, found return REASON_SIMILAR, found
@@ -183,5 +183,7 @@ class RecommendationService:
return None, REASON_DISCOVER return None, REASON_DISCOVER
async def _hydrate_tracks(self, ids: list[uuid.UUID]) -> list[Track]: async def _hydrate_tracks(self, ids: list[uuid.UUID]) -> list[Track]:
"""Resolve ids → tracks preserving order, skipping any that vanished.""" """Resolve ids → tracks preserving order, skipping any that vanished.
return [t for i in ids if (t := await self._tracks.get_by_id(i))] 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]
+13 -2
View File
@@ -13,9 +13,12 @@ importing the other.
""" """
import re import re
import shutil
import uuid import uuid
from pathlib import Path from pathlib import Path
import anyio
from app.domain.errors import NotFoundError from app.domain.errors import NotFoundError
from app.domain.ports import TrackRepository from app.domain.ports import TrackRepository
@@ -60,6 +63,12 @@ def hls_segment_path(root: Path, track_id: uuid.UUID, name: str) -> Path | None:
return hls_dir(root, track_id) / name 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: class TranscodeService:
"""Request-side cache lookups for transcoded renditions.""" """Request-side cache lookups for transcoded renditions."""
@@ -83,14 +92,16 @@ class TranscodeService:
if bitrate is None: if bitrate is None:
return None return None
path = opus_path(self._root, track_id, bitrate) path = opus_path(self._root, track_id, bitrate)
return path if path.exists() else None 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: async def hls_playlist(self, track_id: uuid.UUID) -> Path | None:
"""Cached HLS playlist if generated, else ``None`` (validates the track """Cached HLS playlist if generated, else ``None`` (validates the track
exists so an unknown id 404s rather than silently missing).""" exists so an unknown id 404s rather than silently missing)."""
await self._require_streamable(track_id) await self._require_streamable(track_id)
path = hls_playlist_path(self._root, track_id) path = hls_playlist_path(self._root, track_id)
return path if path.exists() else None 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: def hls_segment(self, track_id: uuid.UUID, name: str) -> Path | None:
path = hls_segment_path(self._root, track_id, name) path = hls_segment_path(self._root, track_id, name)
+5
View File
@@ -142,6 +142,11 @@ class ArtistRepository(Protocol):
class TrackRepository(Protocol): class TrackRepository(Protocol):
async def get_by_id(self, track_id: uuid.UUID) -> Track | None: ... 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 get_by_source(self, source: str, source_id: str) -> Track | None: ...
async def add( async def add(
self, self,
@@ -6,7 +6,7 @@ Likes are an append-only event log. Current state = latest event per (user, trac
import datetime as dt import datetime as dt
import uuid import uuid
from sqlalchemy import func, select from sqlalchemy import Subquery, func, select
from sqlalchemy.dialects.postgresql import insert as pg_insert from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
@@ -105,31 +105,52 @@ class SqlAlchemyLikeRepository:
rows = (await self._session.execute(stmt)).scalars().all() rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows] 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( async def get_latest_state(
self, *, user_id: uuid.UUID, track_ids: list[uuid.UUID] self, *, user_id: uuid.UUID, track_ids: list[uuid.UUID]
) -> list[Like]: ) -> list[Like]:
if not track_ids: if not track_ids:
return [] 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 = ( rows = (
( (
await self._session.execute( await self._session.execute(
select(LikeModel) select(LikeModel)
.join( .where(
max_sq, LikeModel.user_id == user_id,
(LikeModel.track_id == max_sq.c.track_id) LikeModel.track_id.in_(track_ids),
& (LikeModel.created_at == max_sq.c.latest), )
.distinct(LikeModel.track_id)
.order_by(
LikeModel.track_id,
LikeModel.created_at.desc(),
LikeModel.id.desc(),
) )
.where(LikeModel.user_id == user_id)
) )
) )
.scalars() .scalars()
@@ -141,31 +162,14 @@ class SqlAlchemyLikeRepository:
self, *, user_id: uuid.UUID, limit: int, offset: int self, *, user_id: uuid.UUID, limit: int, offset: int
) -> list[Track]: ) -> list[Track]:
# Tracks where the latest like event has value='like', ordered by like time desc # Tracks where the latest like event has value='like', ordered by like time desc
max_sq = ( latest_sq = self._latest_events_sq(user_id, None)
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()
)
rows = ( rows = (
( (
await self._session.execute( await self._session.execute(
select(TrackModel) select(TrackModel)
.join(liked_sq, TrackModel.id == liked_sq.c.track_id) .join(latest_sq, TrackModel.id == latest_sq.c.track_id)
.order_by(liked_sq.c.created_at.desc()) .where(latest_sq.c.value == "like")
.order_by(latest_sq.c.created_at.desc())
.limit(limit) .limit(limit)
.offset(offset) .offset(offset)
) )
@@ -176,25 +180,11 @@ class SqlAlchemyLikeRepository:
return [_track_to_entity(r) for r in rows] return [_track_to_entity(r) for r in rows]
async def count_liked_tracks(self, *, user_id: uuid.UUID) -> int: async def count_liked_tracks(self, *, user_id: uuid.UUID) -> int:
max_sq = ( latest_sq = self._latest_events_sq(user_id, None)
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()
)
return ( 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() ).scalar_one()
@@ -46,6 +46,16 @@ class SqlAlchemyTrackRepository:
row = await self._session.get(TrackModel, track_id) row = await self._session.get(TrackModel, track_id)
return _to_entity(row) if row is not None else None 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: async def get_by_source(self, source: str, source_id: str) -> Track | None:
row = ( row = (
await self._session.execute( await self._session.execute(
+5 -1
View File
@@ -7,6 +7,7 @@ quality renditions; HLS segments are AAC-in-MPEG-TS (broad player support).
""" """
import asyncio import asyncio
import uuid
from pathlib import Path from pathlib import Path
import anyio import anyio
@@ -30,7 +31,10 @@ class FfmpegTranscoder:
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None: async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None:
await anyio.to_thread.run_sync(_mkdir, dest.parent) await anyio.to_thread.run_sync(_mkdir, dest.parent)
tmp = dest.with_suffix(dest.suffix + ".part") # 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( await self._run(
"-i", str(src), "-i", str(src),
"-vn", "-c:a", "libopus", "-b:a", f"{bitrate_kbps}k", "-vn", "-c:a", "libopus", "-b:a", f"{bitrate_kbps}k",
+19 -1
View File
@@ -7,9 +7,13 @@ enqueues (e.g. several plays before the first finishes) are cheap. The DB sessio
is released before ffmpeg runs; the track entity is a detached value object. is released before ffmpeg runs; the track entity is a detached value object.
""" """
import shutil
import uuid import uuid
from pathlib import Path
from typing import Any from typing import Any
import anyio
from app.application.transcode_service import ( from app.application.transcode_service import (
HLS_BITRATE, HLS_BITRATE,
bitrate_for_quality, bitrate_for_quality,
@@ -27,6 +31,15 @@ from app.infrastructure.transcode.ffmpeg import FfmpegTranscoder
log = get_logger("worker.transcode") 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( async def transcode_track(
_ctx: dict[str, Any], _ctx: dict[str, Any],
*, *,
@@ -57,7 +70,12 @@ async def transcode_track(
await transcoder.to_opus(src, dest, bitrate_kbps=bitrate) await transcoder.to_opus(src, dest, bitrate_kbps=bitrate)
did_opus = True did_opus = True
if hls and not hls_playlist_path(root, tid).exists(): if hls and not hls_playlist_path(root, tid).exists():
await transcoder.to_hls(src, hls_dir(root, tid), bitrate_kbps=HLS_BITRATE) # 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 did_hls = True
log.info("transcode_done", track_id=track_id, opus=did_opus, hls=did_hls) log.info("transcode_done", 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