Compare commits
5 Commits
591a938e71
...
sloppyslop
| Author | SHA1 | Date | |
|---|---|---|---|
| fb7827d09c | |||
| c5a473fddf | |||
| 9a78cf5261 | |||
| d16c6085c9 | |||
| ed77acf0fe |
+39
-4
@@ -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,
|
||||||
|
)
|
||||||
|
|||||||
@@ -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,
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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,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)
|
||||||
|
|||||||
@@ -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(
|
||||||
|
|||||||
@@ -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",
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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"}]
|
||||||
@@ -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)
|
||||||
@@ -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
|
||||||
Reference in New Issue
Block a user