Compare commits

...

6 Commits

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

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 21:59:15 +03:00
Цвылев Александр Вадимович 313af3a070 feat(transcode): optimize + on-the-fly quality + HLS streaming (§6.6)
POST /tracks/{id}/optimize enqueues a transcode; GET /stream/{id}?quality=
serves a cached Opus rendition (miss → master + background warm); GET
/stream/{id}/hls/{playlist.m3u8,segment} serves the worker-generated HLS
rendition (AAC-in-TS), with ?token= propagated onto segment URLs for players
that can't set headers. Hexagonal: Transcoder port, FfmpegTranscoder adapter,
TranscodeService + cache-path helpers, transcode_track worker (idempotent),
schemas + deps wiring. ffmpeg already in the image.

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

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 20:58:28 +03:00
Цвылев Александр Вадимович 048854f92a feat(api): offline-first sync layer
Implements the stubbed /sync endpoints:
- GET /sync/changes — delta pull (likes, plays, changed playlists with
  their track ids, changed catalogue tracks) over the half-open window
  (since, cursor]; the cursor is the DB clock, so it's immune to app/DB
  skew.
- POST /sync/push — idempotent append of client like/play events
  (ON CONFLICT DO NOTHING by client-supplied id); events for tracks the
  server doesn't have are skipped (graceful degradation).

Adds a server-ingestion column `synced_at` to the likes + play_history
event logs (migration) as the delta ordering key, so an event pushed
with an older event time still surfaces for other devices. Repos gain
list_since/add_event (likes, history), list_changed_since (playlists,
tracks) and playlist.track_ids; wired via SyncService in deps.

Also fixes test isolation exposed by the registry-backed /admin/sources
endpoint: test_sources_api clears the process-cached source registry,
and test_admin_api no longer hardcodes the environment name.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 15:04:13 +03:00
Цвылев Александр Вадимович c47242aa3a feat(api): add configurable CORS middleware
The web UI is multi-instance and can connect to the backend at a
different origin (the direct :8000 port, a LAN IP, 127.0.0.1 vs
localhost), which the browser blocks without CORS headers. Adds
CORSMiddleware driven by a new cors_allow_origins setting (default "*",
safe here: bearer-token auth with allow_credentials=False). Accepts a
comma-separated string in .env.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 14:33:24 +03:00
Цвылев Александр Вадимович a263272935 feat(api): finish Group A stubbed endpoints
Implements previously-stubbed /api/v1 endpoints (hexagonal: ports -> repos
-> services -> routers wired in deps):

- playlists: GET /{id}/cover (serves stored cover, 404 when absent)
- settings: GET/PATCH /settings + GET/PUT /settings/scrobbling — lazy
  per-user row, write-only Fernet-encrypted scrobble session key; adds
  user_settings table + migration (chains off dc126696f5a6)
- storage: GET /duplicates, /broken, /missing-metadata + admin POST
  /cleanup (arq cleanup_storage worker; reconciles local refs only,
  guarded against a storage-outage mass delete)
- admin: GET /services, /sources, /settings + POST /reindex; PATCH
  /settings and /sources/{source} return 501 (config is env-managed)

Adds NotSupportedError (-> HTTP 501). Integration tests for each surface.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-28 14:33:12 +03:00
59 changed files with 3613 additions and 51 deletions
@@ -0,0 +1,58 @@
"""user_settings: per-user preferences + scrobbling config
Revision ID: 20260728_user_settings
Revises: dc126696f5a6
Create Date: 2026-07-28 10:00:00.000000
Adds the ``user_settings`` table (1:1 with ``users``, PK = user_id): general
preferences (theme, stream quality) plus scrobbling config. The scrobbler
session key is stored Fernet-encrypted, never in plaintext.
"""
from __future__ import annotations
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
revision: str = "20260728_user_settings"
down_revision: str | None = "dc126696f5a6"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.create_table(
"user_settings",
sa.Column("user_id", sa.Uuid(), nullable=False),
sa.Column("theme", sa.String(length=16), nullable=False),
sa.Column("stream_quality", sa.String(length=16), nullable=False),
sa.Column("scrobble_enabled", sa.Boolean(), nullable=False),
sa.Column("scrobble_provider", sa.String(length=16), nullable=True),
sa.Column("scrobble_username", sa.String(length=255), nullable=True),
sa.Column("scrobble_session_key_enc", sa.String(length=512), nullable=True),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
sa.Column(
"updated_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
sa.ForeignKeyConstraint(
["user_id"],
["users.id"],
name=op.f("fk_user_settings_user_id_users"),
ondelete="CASCADE",
),
sa.PrimaryKeyConstraint("user_id", name=op.f("pk_user_settings")),
)
def downgrade() -> None:
op.drop_table("user_settings")
@@ -0,0 +1,57 @@
"""sync: server-ingestion columns on the event logs
Revision ID: 20260728_sync_synced_at
Revises: 20260728_user_settings
Create Date: 2026-07-28 11:00:00.000000
Adds ``synced_at`` to ``likes`` and ``play_history`` — the server-side ingestion
time used as the delta-sync ordering key (distinct from the event time
``created_at``/``played_at``, which a sync push preserves from the client even
when the event happened offline earlier). Existing rows are backfilled from
their event time so a first sync after upgrade behaves sensibly.
"""
from __future__ import annotations
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
revision: str = "20260728_sync_synced_at"
down_revision: str | None = "20260728_user_settings"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.add_column(
"likes",
sa.Column(
"synced_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
)
op.execute("UPDATE likes SET synced_at = created_at")
op.create_index(op.f("ix_likes_synced_at"), "likes", ["synced_at"])
op.add_column(
"play_history",
sa.Column(
"synced_at",
sa.DateTime(timezone=True),
server_default=sa.text("now()"),
nullable=False,
),
)
op.execute("UPDATE play_history SET synced_at = played_at")
op.create_index(op.f("ix_play_history_synced_at"), "play_history", ["synced_at"])
def downgrade() -> None:
op.drop_index(op.f("ix_play_history_synced_at"), table_name="play_history")
op.drop_column("play_history", "synced_at")
op.drop_index(op.f("ix_likes_synced_at"), table_name="likes")
op.drop_column("likes", "synced_at")
+76
View File
@@ -6,22 +6,29 @@ bound to the request-scoped DB session; stateless adapters (hasher, token
service) are process-cached. service) are process-cached.
""" """
import datetime as dt
from collections.abc import AsyncIterator from collections.abc import AsyncIterator
from functools import lru_cache from functools import lru_cache
from typing import Annotated from typing import Annotated
from fastapi import Depends, Query from fastapi import Depends, Query
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.application.auth_service import AuthService from app.application.auth_service import AuthService
from app.application.download_service import DownloadService from app.application.download_service import DownloadService
from app.application.lyrics_service import LyricsService
from app.application.metadata_service import MetadataEnrichmentService from app.application.metadata_service import MetadataEnrichmentService
from app.application.recommendation_service import RecommendationService
from app.application.remote_library_service import RemoteLibraryService from app.application.remote_library_service import RemoteLibraryService
from app.application.streaming_service import StreamingService from app.application.streaming_service import StreamingService
from app.application.subsonic_auth_service import SubsonicAuthService from app.application.subsonic_auth_service import SubsonicAuthService
from app.application.sync_service import SyncService
from app.application.transcode_service import TranscodeService
from app.application.upload_service import UploadService from app.application.upload_service import UploadService
from app.application.user_service import UserService from app.application.user_service import UserService
from app.application.user_settings_service import UserSettingsService
from app.core.config import get_settings from app.core.config import get_settings
from app.core.security import Argon2PasswordHasher, JwtTokenService, SubsonicPasswordCipher from app.core.security import Argon2PasswordHasher, JwtTokenService, SubsonicPasswordCipher
from app.domain.entities import User from app.domain.entities import User
@@ -34,14 +41,18 @@ from app.infrastructure.db.repositories import (
SqlAlchemyDownloadJobRepository, SqlAlchemyDownloadJobRepository,
SqlAlchemyHistoryRepository, SqlAlchemyHistoryRepository,
SqlAlchemyLikeRepository, SqlAlchemyLikeRepository,
SqlAlchemyLyricsRepository,
SqlAlchemyPlaylistRepository, SqlAlchemyPlaylistRepository,
SqlAlchemyRefreshTokenRepository, SqlAlchemyRefreshTokenRepository,
SqlAlchemyTrackRepository, SqlAlchemyTrackRepository,
SqlAlchemyUserRepository, SqlAlchemyUserRepository,
SqlAlchemyUserSettingsRepository,
) )
from app.infrastructure.metadata.acoustid import AcoustIdHttpClient from app.infrastructure.metadata.acoustid import AcoustIdHttpClient
from app.infrastructure.metadata.fingerprint import FpcalcFingerprinter from app.infrastructure.metadata.fingerprint import FpcalcFingerprinter
from app.infrastructure.metadata.lrclib import LrclibHttpClient
from app.infrastructure.metadata.tags import MutagenTagReader from app.infrastructure.metadata.tags import MutagenTagReader
from app.infrastructure.ml.recommender import NullRecommender
from app.infrastructure.sources.registry import SourceRegistry, build_source_registry from app.infrastructure.sources.registry import SourceRegistry, build_source_registry
from app.infrastructure.storage.provider import get_file_storage from app.infrastructure.storage.provider import get_file_storage
from app.workers.queue import enqueue_download, enqueue_enrich, enqueue_materialize from app.workers.queue import enqueue_download, enqueue_enrich, enqueue_materialize
@@ -112,9 +123,17 @@ def get_subsonic_auth_service(session: SessionDep) -> SubsonicAuthService:
) )
def get_user_settings_service(session: SessionDep) -> UserSettingsService:
return UserSettingsService(
settings=SqlAlchemyUserSettingsRepository(session),
cipher=get_subsonic_cipher(),
)
AuthServiceDep = Annotated[AuthService, Depends(get_auth_service)] AuthServiceDep = Annotated[AuthService, Depends(get_auth_service)]
UserServiceDep = Annotated[UserService, Depends(get_user_service)] UserServiceDep = Annotated[UserService, Depends(get_user_service)]
SubsonicAuthServiceDep = Annotated[SubsonicAuthService, Depends(get_subsonic_auth_service)] SubsonicAuthServiceDep = Annotated[SubsonicAuthService, Depends(get_subsonic_auth_service)]
UserSettingsServiceDep = Annotated[UserSettingsService, Depends(get_user_settings_service)]
# -- file storage (process-cached) --------------------------------------------- # -- file storage (process-cached) ---------------------------------------------
@@ -162,6 +181,40 @@ def get_metadata_service(session: SessionDep, storage: FileStorageDep) -> Metada
) )
def get_transcode_service(session: SessionDep) -> TranscodeService:
"""Request-side cache lookups for transcoded renditions (§6.6). Generation
itself runs in the ``transcode_track`` worker, never here."""
return TranscodeService(
tracks=SqlAlchemyTrackRepository(session),
cache_root=get_settings().transcode_cache_path,
)
def get_recommendation_service(session: SessionDep) -> RecommendationService:
"""Radio + similarity (§6.5). ML is optional and no service/contract exists
yet, so we wire ``NullRecommender`` — the service then uses its metadata
fallback. Swap in ``RemoteRecommender(ml_service_url)`` once ML lands."""
return RecommendationService(
recommender=NullRecommender(),
tracks=SqlAlchemyTrackRepository(session),
artists=SqlAlchemyArtistRepository(session),
likes=SqlAlchemyLikeRepository(session),
)
def get_lyrics_service(session: SessionDep) -> LyricsService:
"""Wires the LRCLIB lyrics provider + cache repo (plan §6.7). LRCLIB is
keyless, so this is always available; failures degrade to ``not_found``."""
settings = get_settings()
return LyricsService(
lyrics=SqlAlchemyLyricsRepository(session),
tracks=SqlAlchemyTrackRepository(session),
artists=SqlAlchemyArtistRepository(session),
albums=SqlAlchemyAlbumRepository(session),
provider=LrclibHttpClient(user_agent=settings.musicbrainz_user_agent),
)
def get_download_service(session: SessionDep, storage: FileStorageDep) -> DownloadService: def get_download_service(session: SessionDep, storage: FileStorageDep) -> DownloadService:
return DownloadService( return DownloadService(
jobs=SqlAlchemyDownloadJobRepository(session), jobs=SqlAlchemyDownloadJobRepository(session),
@@ -185,10 +238,33 @@ def get_remote_library_service(session: SessionDep) -> RemoteLibraryService:
UploadServiceDep = Annotated[UploadService, Depends(get_upload_service)] UploadServiceDep = Annotated[UploadService, Depends(get_upload_service)]
StreamingServiceDep = Annotated[StreamingService, Depends(get_streaming_service)] StreamingServiceDep = Annotated[StreamingService, Depends(get_streaming_service)]
MetadataServiceDep = Annotated[MetadataEnrichmentService, Depends(get_metadata_service)] MetadataServiceDep = Annotated[MetadataEnrichmentService, Depends(get_metadata_service)]
LyricsServiceDep = Annotated[LyricsService, Depends(get_lyrics_service)]
TranscodeServiceDep = Annotated[TranscodeService, Depends(get_transcode_service)]
RecommendationServiceDep = Annotated[
RecommendationService, Depends(get_recommendation_service)
]
DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)] DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)]
RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)] RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)]
async def _db_now(session: AsyncSession) -> dt.datetime:
"""The database clock — the sync cursor watermark (avoids app/DB skew)."""
return (await session.execute(select(func.now()))).scalar_one()
def get_sync_service(session: SessionDep) -> SyncService:
return SyncService(
likes=SqlAlchemyLikeRepository(session),
history=SqlAlchemyHistoryRepository(session),
playlists=SqlAlchemyPlaylistRepository(session),
tracks=SqlAlchemyTrackRepository(session),
now=lambda: _db_now(session),
)
SyncServiceDep = Annotated[SyncService, Depends(get_sync_service)]
# -- library repository deps --------------------------------------------------- # -- library repository deps ---------------------------------------------------
def get_track_repository(session: SessionDep) -> SqlAlchemyTrackRepository: def get_track_repository(session: SessionDep) -> SqlAlchemyTrackRepository:
return SqlAlchemyTrackRepository(session) return SqlAlchemyTrackRepository(session)
+2
View File
@@ -18,6 +18,7 @@ from app.domain.errors import (
DependencyUnavailableError, DependencyUnavailableError,
DomainError, DomainError,
NotFoundError, NotFoundError,
NotSupportedError,
PermissionDeniedError, PermissionDeniedError,
RangeNotSatisfiableError, RangeNotSatisfiableError,
StorageError, StorageError,
@@ -33,6 +34,7 @@ _STATUS_BY_ERROR: dict[type[DomainError], int] = {
ValidationError: status.HTTP_422_UNPROCESSABLE_CONTENT, ValidationError: status.HTTP_422_UNPROCESSABLE_CONTENT,
AuthenticationError: status.HTTP_401_UNAUTHORIZED, AuthenticationError: status.HTTP_401_UNAUTHORIZED,
PermissionDeniedError: status.HTTP_403_FORBIDDEN, PermissionDeniedError: status.HTTP_403_FORBIDDEN,
NotSupportedError: status.HTTP_501_NOT_IMPLEMENTED,
DependencyUnavailableError: status.HTTP_503_SERVICE_UNAVAILABLE, DependencyUnavailableError: status.HTTP_503_SERVICE_UNAVAILABLE,
StorageError: status.HTTP_500_INTERNAL_SERVER_ERROR, StorageError: status.HTTP_500_INTERNAL_SERVER_ERROR,
} }
+39
View File
@@ -0,0 +1,39 @@
"""Admin (instance-management) response schemas."""
from pydantic import BaseModel
from app.api.health import CheckStatus
class ServicesStatusOut(BaseModel):
"""Backing-dependency health for the admin dashboard (mirrors readiness)."""
database: CheckStatus
redis: CheckStatus
ml: CheckStatus
class ReindexJob(BaseModel):
source: str
job_id: str
class ReindexResponse(BaseModel):
"""The scan jobs enqueued by a re-index, one per indexable source."""
jobs: list[ReindexJob]
class AdminSettingsOut(BaseModel):
"""Effective, non-secret instance configuration. Secrets and connection
strings are never exposed — only whether an optional integration is set up."""
environment: str
allow_registration: bool
storage_backend: str
media_path: str
youtube_enabled: bool
coverart_enabled: bool
ml_configured: bool
acoustid_configured: bool
local_import_configured: bool
+33
View File
@@ -0,0 +1,33 @@
"""Lyrics response schema (§6.7 / Now Playing lyrics panel).
Returns the raw LRC (``synced``) and/or ``plain`` text; the client parses LRC
timestamps for synced highlighting. A miss is a normal 200 with
``status="not_found"`` and null text — not an error — so the panel can render a
"no lyrics" state.
"""
import uuid
from pydantic import BaseModel
from app.domain.entities.lyrics import Lyrics
class LyricsOut(BaseModel):
track_id: uuid.UUID
status: str
source: str | None
synced: str | None
plain: str | None
synced_available: bool
@classmethod
def from_entity(cls, lyrics: Lyrics) -> LyricsOut:
return cls(
track_id=lyrics.track_id,
status=lyrics.status,
source=lyrics.source,
synced=lyrics.synced,
plain=lyrics.plain,
synced_available=lyrics.synced is not None,
)
+43
View File
@@ -0,0 +1,43 @@
"""Radio + similarity response schemas (§6.5)."""
import uuid
from pydantic import BaseModel, Field
from app.api.schemas.artist import ArtistOut
from app.api.schemas.track import TrackOut
class RadioRequest(BaseModel):
"""Start or continue a radio. ``seed_track_id`` seeds from a track;
``from_likes`` seeds from the caller's likes. ``exclude_ids`` are already-
queued tracks to skip (the client drives the infinite feed). ``exploration``
biases familiar↔new."""
seed_track_id: uuid.UUID | None = None
from_likes: bool = False
exploration: float = Field(default=0.25, ge=0.0, le=1.0)
count: int = Field(default=20, ge=1, le=50)
exclude_ids: list[uuid.UUID] = Field(default_factory=list)
class RadioTrackOut(BaseModel):
track: TrackOut
# Short code the client localizes: ml | similar | from_likes | discover.
reason: str
class RadioResponse(BaseModel):
# Where the picks came from: "ml" or "metadata" (fallback).
source: str
tracks: list[RadioTrackOut]
class SimilarTracksOut(BaseModel):
source: str
tracks: list[TrackOut]
class SimilarArtistsOut(BaseModel):
source: str
artists: list[ArtistOut]
+49
View File
@@ -0,0 +1,49 @@
"""User-settings request/response schemas.
Enums are enforced at the API boundary (Pydantic ``Literal`` → 422 on bad
input), so the service can trust the values it receives.
"""
from typing import Literal
from pydantic import BaseModel, model_validator
Theme = Literal["system", "light", "dark"]
# Playback quality preference. ``original`` = no transcode; the lower tiers are
# consumed by the (upcoming) transcoding pipeline.
StreamQuality = Literal["original", "high", "medium", "low"]
ScrobbleProvider = Literal["lastfm", "listenbrainz"]
class SettingsOut(BaseModel):
theme: Theme
stream_quality: StreamQuality
class SettingsUpdate(BaseModel):
"""Partial update — omitted fields keep their current value."""
theme: Theme | None = None
stream_quality: StreamQuality | None = None
class ScrobblingOut(BaseModel):
enabled: bool
provider: ScrobbleProvider | None
username: str | None
# Whether a session key is stored. The key itself is never returned.
configured: bool
class ScrobblingUpdate(BaseModel):
enabled: bool = False
provider: ScrobbleProvider | None = None
username: str | None = None
# Write-only scrobbler session key / user token. Omit to keep the stored one.
session_key: str | None = None
@model_validator(mode="after")
def _provider_required_when_enabled(self) -> ScrobblingUpdate:
if self.enabled and self.provider is None:
raise ValueError("provider is required when scrobbling is enabled")
return self
+16
View File
@@ -4,6 +4,8 @@ import datetime as dt
from pydantic import BaseModel from pydantic import BaseModel
from app.api.schemas.track import TrackOut
class DiskUsageOut(BaseModel): class DiskUsageOut(BaseModel):
total: int total: int
@@ -43,3 +45,17 @@ class StorageStatsOut(BaseModel):
# backing volume (``None`` for object-store backends) # backing volume (``None`` for object-store backends)
disk: DiskUsageOut | None disk: DiskUsageOut | None
class DuplicateGroupOut(BaseModel):
"""Tracks sharing one acoustic fingerprint — candidates for de-duplication."""
fingerprint: str
tracks: list[TrackOut]
class CleanupEnqueuedOut(BaseModel):
"""Acknowledgement that a cleanup job was queued (it runs in the worker)."""
status: str
job_id: str
+82
View File
@@ -0,0 +1,82 @@
"""Offline-first sync schemas (delta pull + idempotent push).
The client keeps an opaque ``cursor`` (a server-clock timestamp). It pulls
everything changed in the half-open window ``(since, cursor]`` and pushes the
append-only events it accumulated offline. Events carry a client-generated
``id`` so a replay is idempotent.
"""
import datetime as dt
import uuid
from typing import Literal
from pydantic import BaseModel, Field
from app.api.schemas.track import TrackOut
LikeValue = Literal["like", "dislike", "neutral"]
# -- pull (server -> client) --------------------------------------------------
class LikeEventOut(BaseModel):
id: uuid.UUID
track_id: uuid.UUID
value: str
created_at: dt.datetime
class PlayEventOut(BaseModel):
id: uuid.UUID
track_id: uuid.UUID
played_at: dt.datetime
play_duration_seconds: int | None
completed: bool
class PlaylistSyncOut(BaseModel):
id: uuid.UUID
name: str
description: str | None
version: int
updated_at: dt.datetime
track_ids: list[uuid.UUID]
class SyncChangesOut(BaseModel):
"""Everything that changed for the caller since their last cursor. Feed
``cursor`` back as ``?since=`` on the next pull."""
cursor: dt.datetime
likes: list[LikeEventOut]
plays: list[PlayEventOut]
playlists: list[PlaylistSyncOut]
tracks: list[TrackOut]
# -- push (client -> server) --------------------------------------------------
class LikeEventIn(BaseModel):
id: uuid.UUID
track_id: uuid.UUID
value: LikeValue
created_at: dt.datetime
class PlayEventIn(BaseModel):
id: uuid.UUID
track_id: uuid.UUID
played_at: dt.datetime
play_duration_seconds: int | None = None
completed: bool = False
class SyncPushIn(BaseModel):
likes: list[LikeEventIn] = Field(default_factory=list)
plays: list[PlayEventIn] = Field(default_factory=list)
class SyncPushOut(BaseModel):
"""How many events were newly stored (a replay reports 0) + a fresh cursor."""
cursor: dt.datetime
accepted_likes: int
accepted_plays: int
+11
View File
@@ -0,0 +1,11 @@
"""Transcode/optimize response schemas (§6.6)."""
from pydantic import BaseModel
class OptimizeEnqueuedOut(BaseModel):
"""Acknowledgement that a transcode job was queued (it runs in the worker)."""
status: str
job_id: str
quality: str
+71 -11
View File
@@ -5,11 +5,18 @@ sign-up (plan §6.4).
""" """
import uuid import uuid
from typing import Any
from fastapi import APIRouter, Query, status from fastapi import APIRouter, Query, status
from app.api.deps import SubsonicAuthServiceDep, SuperUser, UserServiceDep from app.api.deps import SourceRegistryDep, SubsonicAuthServiceDep, SuperUser, UserServiceDep
from app.api.health import _check_db, _check_ml, _check_redis
from app.api.schemas.admin import (
AdminSettingsOut,
ReindexJob,
ReindexResponse,
ServicesStatusOut,
)
from app.api.schemas.source import SourceInfoOut
from app.api.schemas.subsonic import SubsonicPasswordResponse from app.api.schemas.subsonic import SubsonicPasswordResponse
from app.api.schemas.user import ( from app.api.schemas.user import (
CreateUserRequest, CreateUserRequest,
@@ -17,6 +24,9 @@ from app.api.schemas.user import (
UpdateUserRequest, UpdateUserRequest,
UserResponse, UserResponse,
) )
from app.core.config import get_settings
from app.domain.errors import DependencyUnavailableError, NotSupportedError
from app.workers.queue import enqueue
router = APIRouter(prefix="/admin", tags=["admin"]) router = APIRouter(prefix="/admin", tags=["admin"])
@@ -91,24 +101,74 @@ async def rotate_user_subsonic_password(
@router.get("/services") @router.get("/services")
async def list_services(_admin: SuperUser) -> Any: ... async def list_services(_admin: SuperUser) -> ServicesStatusOut:
"""Backing-dependency health for the admin dashboard — same probes as the
readiness endpoint (DB + Redis required, ML optional)."""
database = await _check_db()
redis = await _check_redis()
ml = await _check_ml()
return ServicesStatusOut(database=database, redis=redis, ml=ml)
@router.get("/sources") @router.get("/sources")
async def list_admin_sources(_admin: SuperUser) -> Any: ... async def list_admin_sources(
_admin: SuperUser, registry: SourceRegistryDep
) -> list[SourceInfoOut]:
@router.patch("/sources/{source}") """Configured sources and their live availability (same view as
async def update_admin_source(source: str, _admin: SuperUser) -> Any: ... ``/sources``, admin-scoped)."""
return [SourceInfoOut.from_entity(info) for info in registry.infos()]
@router.post("/reindex") @router.post("/reindex")
async def trigger_reindex(_admin: SuperUser) -> Any: ... async def trigger_reindex(admin: SuperUser, registry: SourceRegistryDep) -> ReindexResponse:
"""Enqueue a full re-scan of every indexable source. The walk + file copies
run in the worker (never the request cycle); re-scans are idempotent."""
indexables = registry.indexables()
if not indexables:
raise DependencyUnavailableError("No indexable source is configured.")
jobs: list[ReindexJob] = []
for backend in indexables:
job_id = await enqueue("scan_local_folder", source=backend.name, added_by=str(admin.id))
jobs.append(ReindexJob(source=backend.name, job_id=job_id))
return ReindexResponse(jobs=jobs)
@router.get("/settings") @router.get("/settings")
async def get_admin_settings(_admin: SuperUser) -> Any: ... async def get_admin_settings(_admin: SuperUser) -> AdminSettingsOut:
"""Effective, non-secret instance configuration. Reflects the environment the
process booted with; secrets/connection strings are never returned — only
whether each optional integration is configured."""
settings = get_settings()
return AdminSettingsOut(
environment=settings.environment,
allow_registration=settings.allow_registration,
storage_backend=settings.storage_backend,
media_path=str(settings.media_path),
youtube_enabled=settings.youtube_enabled,
coverart_enabled=settings.coverart_enabled,
ml_configured=settings.ml_service_url is not None,
acoustid_configured=settings.acoustid_api_key is not None,
local_import_configured=settings.local_media_import_path is not None,
)
# -- runtime config mutation (intentionally unsupported) ----------------------
# The instance is env-configured (CLAUDE.md: nothing hardcoded, all from env) and
# get_settings() is a cached singleton, so config is not mutable at runtime.
# These endpoints answer 501 with a clear reason rather than silently no-op'ing;
# a persistent override layer that shadows env would be a deliberate future
# departure. Read the effective config via GET /admin/settings.
@router.patch("/sources/{source}")
async def update_admin_source(source: str, _admin: SuperUser) -> None:
raise NotSupportedError(
"Sources are configured via environment variables (e.g. YOUTUBE_ENABLED, "
"LOCAL_MEDIA_IMPORT_PATH); runtime changes are not supported."
)
@router.patch("/settings") @router.patch("/settings")
async def update_admin_settings(_admin: SuperUser) -> Any: ... async def update_admin_settings(_admin: SuperUser) -> None:
raise NotSupportedError(
"Instance settings are managed via environment configuration; "
"runtime changes are not supported."
)
+29 -3
View File
@@ -1,14 +1,20 @@
"""Artist endpoints.""" """Artist endpoints."""
import uuid import uuid
from typing import Any
from fastapi import APIRouter, Query from fastapi import APIRouter, Query
from app.api.deps import AlbumRepoDep, ArtistRepoDep, CurrentUser, TrackRepoDep from app.api.deps import (
AlbumRepoDep,
ArtistRepoDep,
CurrentUser,
RecommendationServiceDep,
TrackRepoDep,
)
from app.api.schemas.album import AlbumOut from app.api.schemas.album import AlbumOut
from app.api.schemas.artist import ArtistOut from app.api.schemas.artist import ArtistOut
from app.api.schemas.pagination import PagedResponse from app.api.schemas.pagination import PagedResponse
from app.api.schemas.radio import SimilarArtistsOut
from app.api.schemas.track import TrackOut from app.api.schemas.track import TrackOut
from app.api.v1.albums import _build_album_out from app.api.v1.albums import _build_album_out
from app.api.v1.tracks import _build_track_out from app.api.v1.tracks import _build_track_out
@@ -124,4 +130,24 @@ async def get_artist_tracks(
@router.get("/{artist_id}/similar") @router.get("/{artist_id}/similar")
async def get_similar_artists(artist_id: uuid.UUID, _: CurrentUser) -> Any: ... async def get_similar_artists(
artist_id: uuid.UUID,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
_: CurrentUser,
limit: int = Query(20, ge=1, le=100),
) -> SimilarArtistsOut:
"""Artists similar to this one (§6.5). ML when configured, else a shared-
genre metadata heuristic."""
source, artists = await service.similar_artists(artist_id, limit=limit)
items = [
ArtistOut(
id=a.id,
name=a.name,
album_count=await artist_repo.album_count(a.id),
track_count=await artist_repo.track_count(a.id),
created_at=a.created_at,
)
for a in artists
]
return SimilarArtistsOut(source=source, artists=items)
+15 -2
View File
@@ -1,15 +1,18 @@
"""Playlist endpoints.""" """Playlist endpoints."""
import uuid import uuid
from typing import Any
from fastapi import APIRouter, Query, Response from fastapi import APIRouter, Query, Response
from fastapi.responses import StreamingResponse
from app.api.covers import stream_cover
from app.api.deps import ( from app.api.deps import (
AlbumRepoDep, AlbumRepoDep,
ArtistRepoDep, ArtistRepoDep,
CurrentUser, CurrentUser,
FileStorageDep,
PlaylistRepoDep, PlaylistRepoDep,
StreamUser,
TrackRepoDep, TrackRepoDep,
) )
from app.api.schemas.pagination import PagedResponse from app.api.schemas.pagination import PagedResponse
@@ -217,4 +220,14 @@ async def reorder_playlist_tracks(
@router.get("/{playlist_id}/cover") @router.get("/{playlist_id}/cover")
async def get_playlist_cover(playlist_id: uuid.UUID, _: CurrentUser) -> Any: ... async def get_playlist_cover(
playlist_id: uuid.UUID,
playlist_repo: PlaylistRepoDep,
storage: FileStorageDep,
_: StreamUser,
) -> StreamingResponse:
# ``<img>`` can't send a bearer header → StreamUser accepts ``?token=``.
cover_path = await playlist_repo.get_cover_path(playlist_id)
if not cover_path:
raise NotFoundError("Cover not found.")
return await stream_cover(storage, cover_path)
+72 -4
View File
@@ -1,15 +1,83 @@
"""Radio / continuous-mix endpoints. Degrades gracefully when ML service is down.""" """Radio / continuous-mix endpoints (§6.5).
from typing import Any Stateless: the client passes the seed + already-queued ids and pulls more as the
queue drains (offline-first infinite feed). Degrades gracefully when no ML
service is configured — the recommendation service falls back to metadata.
"""
from fastapi import APIRouter from fastapi import APIRouter
from app.api.deps import (
AlbumRepoDep,
ArtistRepoDep,
CurrentUser,
RecommendationServiceDep,
)
from app.api.schemas.radio import RadioRequest, RadioResponse, RadioTrackOut
from app.api.v1.tracks import _build_track_out
from app.application.recommendation_service import RadioPick
router = APIRouter(prefix="/radio", tags=["radio"]) router = APIRouter(prefix="/radio", tags=["radio"])
async def _to_response(
source: str,
picks: list[RadioPick],
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
tracks = [p.track for p in picks]
artist_ids = list({t.artist_id for t in tracks})
album_ids = list({t.album_id for t in tracks if t.album_id is not None})
artists = {a.id: a for a in await artist_repo.get_many(artist_ids)}
albums = {a.id: a for a in await album_repo.get_many(album_ids)}
outs = await _build_track_out(tracks, artists, albums)
return RadioResponse(
source=source,
tracks=[
RadioTrackOut(track=out, reason=pick.reason)
for out, pick in zip(outs, picks, strict=True)
],
)
async def _run_radio(
body: RadioRequest,
user: CurrentUser,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
source, picks = await service.radio(
user_id=user.id,
seed_track_id=body.seed_track_id,
from_likes=body.from_likes,
exploration=body.exploration,
limit=body.count,
exclude_ids=body.exclude_ids,
)
return await _to_response(source, picks, artist_repo, album_repo)
@router.post("") @router.post("")
async def start_radio() -> Any: ... async def start_radio(
body: RadioRequest,
user: CurrentUser,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
"""Start a radio from a seed track or the caller's likes."""
return await _run_radio(body, user, service, artist_repo, album_repo)
@router.post("/next") @router.post("/next")
async def next_radio_track() -> Any: ... async def next_radio_track(
body: RadioRequest,
user: CurrentUser,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
) -> RadioResponse:
"""Fetch more tracks as the radio queue drains (pass ``exclude_ids``)."""
return await _run_radio(body, user, service, artist_repo, album_repo)
+73 -8
View File
@@ -1,22 +1,28 @@
"""Storage analysis and cleanup endpoints.""" """Storage analysis and cleanup endpoints."""
from typing import Any from fastapi import APIRouter, Query
from fastapi import APIRouter
from app.api.deps import ( from app.api.deps import (
AlbumRepoDep, AlbumRepoDep,
ArtistRepoDep, ArtistRepoDep,
CurrentUser, CurrentUser,
FileStorageDep, FileStorageDep,
SuperUser,
TrackRepoDep, TrackRepoDep,
) )
from app.api.schemas.pagination import PagedResponse
from app.api.schemas.storage import ( from app.api.schemas.storage import (
CleanupEnqueuedOut,
DiskUsageOut, DiskUsageOut,
DuplicateGroupOut,
FormatBreakdownOut, FormatBreakdownOut,
GenreCountOut, GenreCountOut,
StorageStatsOut, StorageStatsOut,
) )
from app.api.schemas.track import TrackOut
from app.api.v1.tracks import _build_track_out
from app.domain.entities.track import Track
from app.workers.queue import enqueue
router = APIRouter(prefix="/storage", tags=["storage"]) router = APIRouter(prefix="/storage", tags=["storage"])
@@ -24,6 +30,18 @@ router = APIRouter(prefix="/storage", tags=["storage"])
_TOP_GENRES = 8 _TOP_GENRES = 8
async def _tracks_to_out(
tracks: list[Track], artist_repo: ArtistRepoDep, album_repo: AlbumRepoDep
) -> list[TrackOut]:
"""Hydrate a batch of tracks into ``TrackOut`` (artist/album names + cover
flag), resolving each referenced artist/album in a single query."""
artist_ids = list({t.artist_id for t in tracks})
album_ids = list({t.album_id for t in tracks if t.album_id is not None})
artists = {a.id: a for a in await artist_repo.get_many(artist_ids)}
albums = {a.id: a for a in await album_repo.get_many(album_ids)}
return await _build_track_out(tracks, artists, albums)
@router.get("") @router.get("")
async def get_storage_stats( async def get_storage_stats(
track_repo: TrackRepoDep, track_repo: TrackRepoDep,
@@ -70,16 +88,63 @@ async def get_storage_stats(
@router.get("/duplicates") @router.get("/duplicates")
async def get_duplicates() -> Any: ... async def get_duplicates(
track_repo: TrackRepoDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
_: CurrentUser,
) -> list[DuplicateGroupOut]:
"""Tracks sharing an acoustic fingerprint, grouped — the library's real
duplicates (``(source, source_id)`` is already unique). Cheap DB GROUP BY."""
groups = await track_repo.find_duplicate_groups()
all_tracks = [track for _, tracks in groups for track in tracks]
out = await _tracks_to_out(all_tracks, artist_repo, album_repo)
by_id = {item.id: item for item in out}
return [
DuplicateGroupOut(fingerprint=fingerprint, tracks=[by_id[t.id] for t in tracks])
for fingerprint, tracks in groups
]
@router.get("/broken") @router.get("/broken")
async def get_broken_files() -> Any: ... async def get_broken_files(
track_repo: TrackRepoDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
_: CurrentUser,
limit: int = Query(50, ge=1, le=200),
offset: int = Query(0, ge=0),
) -> PagedResponse[TrackOut]:
"""Tracks whose last enrichment run failed (``metadata_status=failed``) —
each carries its ``metadata_error``. A file gone missing on disk is instead
reconciled by ``POST /storage/cleanup`` (that needs a filesystem scan)."""
tracks = await track_repo.list_by_metadata_status("failed", limit=limit, offset=offset)
total = await track_repo.count_by_metadata_status("failed")
items = await _tracks_to_out(tracks, artist_repo, album_repo)
return PagedResponse(items=items, total=total, limit=limit, offset=offset)
@router.get("/missing-metadata") @router.get("/missing-metadata")
async def get_missing_metadata() -> Any: ... async def get_missing_metadata(
track_repo: TrackRepoDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
_: CurrentUser,
limit: int = Query(50, ge=1, le=200),
offset: int = Query(0, ge=0),
) -> PagedResponse[TrackOut]:
"""Tracks still awaiting enrichment (``metadata_status=pending``) — imported
but never identified."""
tracks = await track_repo.list_by_metadata_status("pending", limit=limit, offset=offset)
total = await track_repo.count_by_metadata_status("pending")
items = await _tracks_to_out(tracks, artist_repo, album_repo)
return PagedResponse(items=items, total=total, limit=limit, offset=offset)
@router.post("/cleanup") @router.post("/cleanup", status_code=202)
async def run_cleanup() -> Any: ... async def run_cleanup(_: SuperUser) -> CleanupEnqueuedOut:
"""Admin: enqueue the storage reconciliation job. It scans the catalogue and
removes rows whose backing file has vanished (dangling references). Runs in
the worker — the filesystem scan must not block the request cycle."""
job_id = await enqueue("cleanup_storage")
return CleanupEnqueuedOut(status="enqueued", job_id=job_id)
+64 -6
View File
@@ -1,30 +1,53 @@
"""Audio streaming endpoint — direct stream with Range support.""" """Audio streaming — direct byte-range stream, transcoded quality, and HLS.
``GET /stream/{id}`` streams the master with Range support, or a cached Opus
rendition when ``?quality=`` is set (a cache miss falls back to the master and
warms the cache in the background — playback never waits on ffmpeg). ``/hls/*``
serves the cached HLS rendition (generated by the ``transcode_track`` worker).
"""
import re
import uuid import uuid
from typing import Annotated from typing import Annotated
from fastapi import APIRouter, Header from fastapi import APIRouter, Header, Query, Response
from fastapi.responses import StreamingResponse from fastapi.responses import FileResponse, StreamingResponse
from app.api.deps import StreamingServiceDep, StreamUser from app.api.deps import StreamingServiceDep, StreamUser, TranscodeServiceDep
from app.domain.errors import NotFoundError
from app.workers.queue import enqueue_transcode_quiet
router = APIRouter(prefix="/stream", tags=["streaming"]) router = APIRouter(prefix="/stream", tags=["streaming"])
_HLS_PLAYLIST_TYPE = "application/vnd.apple.mpegurl"
_HLS_SEGMENT_TYPE = "video/mp2t"
_OPUS_TYPE = "audio/ogg"
_SEGMENT_LINE_RE = re.compile(r"^(seg_\d+\.ts)$", re.MULTILINE)
@router.get("/{track_id}") @router.get("/{track_id}")
async def stream_track( async def stream_track(
track_id: uuid.UUID, track_id: uuid.UUID,
service: StreamingServiceDep, service: StreamingServiceDep,
transcode: TranscodeServiceDep,
_user: StreamUser, _user: StreamUser,
range_header: Annotated[str | None, Header(alias="Range")] = None, range_header: Annotated[str | None, Header(alias="Range")] = None,
) -> StreamingResponse: quality: Annotated[str | None, Query()] = None,
) -> Response:
# A quality rendition, if one is cached; otherwise fall back to the master
# and enqueue generation so the next play gets it (graceful degradation).
if quality and quality != "original":
cached = await transcode.resolve_quality_file(track_id, quality)
if cached is not None:
return FileResponse(cached, media_type=_OPUS_TYPE)
await enqueue_transcode_quiet(track_id, quality=quality, hls=False)
result = await service.open_stream(track_id, range_header) result = await service.open_stream(track_id, range_header)
headers = { headers = {
"Accept-Ranges": "bytes", "Accept-Ranges": "bytes",
"Content-Length": str(result.content_length), "Content-Length": str(result.content_length),
} }
if result.is_partial: if result.is_partial:
headers["Content-Range"] = f"bytes {result.start}-{result.end}/{result.total_size}" headers["Content-Range"] = f"bytes {result.start}-{result.end}/{result.total_size}"
status_code = 206 status_code = 206
@@ -37,3 +60,38 @@ async def stream_track(
headers=headers, headers=headers,
media_type=result.content_type, media_type=result.content_type,
) )
@router.get("/{track_id}/hls/playlist.m3u8")
async def stream_hls_playlist(
track_id: uuid.UUID,
transcode: TranscodeServiceDep,
_user: StreamUser,
token: Annotated[str | None, Query()] = None,
) -> Response:
"""Serve the cached HLS playlist. On a miss, kick off generation and 404 so
the client retries. Segment URLs are relative; when the request carried a
``?token=`` (players can't set an Authorization header), it's appended to
each segment line so the segment requests authenticate the same way."""
path = await transcode.hls_playlist(track_id)
if path is None:
await enqueue_transcode_quiet(track_id, hls=True)
raise NotFoundError("HLS rendition is being prepared; retry shortly.")
body = path.read_text()
if token:
body = _SEGMENT_LINE_RE.sub(rf"\1?token={token}", body)
return Response(body, media_type=_HLS_PLAYLIST_TYPE)
@router.get("/{track_id}/hls/{segment}")
async def stream_hls_segment(
track_id: uuid.UUID,
segment: str,
transcode: TranscodeServiceDep,
_user: StreamUser,
) -> FileResponse:
path = transcode.hls_segment(track_id, segment)
if path is None:
raise NotFoundError("Segment not found.")
return FileResponse(path, media_type=_HLS_SEGMENT_TYPE)
+92 -4
View File
@@ -1,15 +1,103 @@
"""Client sync endpoints (offline-first event log).""" """Client sync endpoints (offline-first event log).
from typing import Any ``GET /sync/changes`` pulls everything the caller changed since their cursor;
``POST /sync/push`` uploads the like/play events a client accumulated offline
(idempotent — replays are no-ops). See :mod:`app.application.sync_service`.
"""
import datetime as dt
from fastapi import APIRouter from fastapi import APIRouter
from app.api.deps import AlbumRepoDep, ArtistRepoDep, CurrentUser, SyncServiceDep
from app.api.schemas.sync import (
LikeEventOut,
PlayEventOut,
PlaylistSyncOut,
SyncChangesOut,
SyncPushIn,
SyncPushOut,
)
from app.api.v1.tracks import _build_track_out
from app.application.sync_service import LikeEvent, PlayEvent
router = APIRouter(prefix="/sync", tags=["sync"]) router = APIRouter(prefix="/sync", tags=["sync"])
@router.get("/changes") @router.get("/changes")
async def get_changes() -> Any: ... async def get_changes(
service: SyncServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
user: CurrentUser,
since: dt.datetime | None = None,
) -> SyncChangesOut:
"""Delta since ``since`` (omit for a full snapshot). Persist ``cursor`` from
the response and pass it back as ``?since=`` next time."""
changes = await service.get_changes(user.id, since=since)
artist_ids = list({t.artist_id for t in changes.tracks})
album_ids = list({t.album_id for t in changes.tracks if t.album_id is not None})
artists = {a.id: a for a in await artist_repo.get_many(artist_ids)}
albums = {a.id: a for a in await album_repo.get_many(album_ids)}
tracks_out = await _build_track_out(changes.tracks, artists, albums)
return SyncChangesOut(
cursor=changes.cursor,
likes=[
LikeEventOut(
id=lk.id, track_id=lk.track_id, value=lk.value, created_at=lk.created_at
)
for lk in changes.likes
],
plays=[
PlayEventOut(
id=p.id,
track_id=p.track_id,
played_at=p.played_at,
play_duration_seconds=p.play_duration_seconds,
completed=p.completed,
)
for p in changes.plays
],
playlists=[
PlaylistSyncOut(
id=d.playlist.id,
name=d.playlist.name,
description=d.playlist.description,
version=d.playlist.version,
updated_at=d.playlist.updated_at,
track_ids=d.track_ids,
)
for d in changes.playlists
],
tracks=tracks_out,
)
@router.post("/push") @router.post("/push")
async def push_changes() -> Any: ... async def push_changes(
body: SyncPushIn, service: SyncServiceDep, user: CurrentUser
) -> SyncPushOut:
result = await service.push(
user.id,
likes=[
LikeEvent(id=e.id, track_id=e.track_id, value=e.value, created_at=e.created_at)
for e in body.likes
],
plays=[
PlayEvent(
id=e.id,
track_id=e.track_id,
played_at=e.played_at,
play_duration_seconds=e.play_duration_seconds,
completed=e.completed,
)
for e in body.plays
],
)
return SyncPushOut(
cursor=result.cursor,
accepted_likes=result.accepted_likes,
accepted_plays=result.accepted_plays,
)
+62 -6
View File
@@ -1,7 +1,7 @@
"""Track endpoints.""" """Track endpoints."""
import uuid import uuid
from typing import Any from typing import Annotated
from fastapi import APIRouter, Query, Response from fastapi import APIRouter, Query, Response
from fastapi.responses import StreamingResponse from fastapi.responses import StreamingResponse
@@ -12,13 +12,17 @@ from app.api.deps import (
ArtistRepoDep, ArtistRepoDep,
CurrentUser, CurrentUser,
FileStorageDep, FileStorageDep,
LyricsServiceDep,
MetadataServiceDep, MetadataServiceDep,
RecommendationServiceDep,
RemoteLibraryServiceDep, RemoteLibraryServiceDep,
StreamUser, StreamUser,
TrackRepoDep, TrackRepoDep,
) )
from app.api.schemas.download import DownloadJobOut from app.api.schemas.download import DownloadJobOut
from app.api.schemas.lyrics import LyricsOut
from app.api.schemas.pagination import PagedResponse from app.api.schemas.pagination import PagedResponse
from app.api.schemas.radio import SimilarTracksOut
from app.api.schemas.track import ( from app.api.schemas.track import (
MaterializeResponse, MaterializeResponse,
MetadataApply, MetadataApply,
@@ -28,10 +32,12 @@ from app.api.schemas.track import (
TrackOut, TrackOut,
TrackUpdate, TrackUpdate,
) )
from app.api.schemas.transcode import OptimizeEnqueuedOut
from app.application.transcode_service import bitrate_for_quality
from app.domain.entities.album import Album from app.domain.entities.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 from app.domain.errors import NotFoundError, ValidationError
from app.workers.queue import enqueue from app.workers.queue import enqueue, enqueue_transcode
router = APIRouter(prefix="/tracks", tags=["tracks"]) router = APIRouter(prefix="/tracks", tags=["tracks"])
@@ -217,11 +223,61 @@ async def delete_track(
@router.get("/{track_id}/similar") @router.get("/{track_id}/similar")
async def get_similar_tracks(track_id: uuid.UUID, _: CurrentUser) -> Any: ... async def get_similar_tracks(
track_id: uuid.UUID,
service: RecommendationServiceDep,
artist_repo: ArtistRepoDep,
album_repo: AlbumRepoDep,
_: CurrentUser,
limit: Annotated[int, Query(ge=1, le=100)] = 20,
) -> SimilarTracksOut:
"""Tracks similar to this one (§6.5). Uses ML when configured, else a
genre/artist metadata heuristic."""
source, tracks = await service.similar_tracks(track_id, limit=limit)
artist_ids = list({t.artist_id for t in tracks})
album_ids = list({t.album_id for t in tracks if t.album_id is not None})
artists = {a.id: a for a in await artist_repo.get_many(artist_ids)}
albums = {a.id: a for a in await album_repo.get_many(album_ids)}
outs = await _build_track_out(tracks, artists, albums)
return SimilarTracksOut(source=source, tracks=outs)
@router.post("/{track_id}/optimize") @router.post("/{track_id}/optimize", status_code=202)
async def optimize_track(track_id: uuid.UUID, _: CurrentUser) -> Any: ... async def optimize_track(
track_id: uuid.UUID,
track_repo: TrackRepoDep,
_: CurrentUser,
quality: Annotated[str, Query()] = "high",
) -> OptimizeEnqueuedOut:
"""Enqueue transcoding of a track into a cached Opus rendition + HLS (§6.6).
Heavy ffmpeg work runs in the worker; this only queues it."""
if bitrate_for_quality(quality) is None:
raise ValidationError(
f"Unknown quality '{quality}'; expected one of high, medium, low."
)
track = await track_repo.get_by_id(track_id)
if track is None:
raise NotFoundError(f"Track {track_id} not found.")
job_id = await enqueue_transcode(track_id, quality=quality, hls=True)
return OptimizeEnqueuedOut(status="enqueued", job_id=job_id, quality=quality)
@router.get("/{track_id}/lyrics")
async def get_track_lyrics(
track_id: uuid.UUID, lyrics: LyricsServiceDep, _: CurrentUser
) -> LyricsOut:
"""Cached lyrics for the Now Playing panel (§6.7). A miss is a normal 200
with ``status="not_found"`` — the provider (LRCLIB) is queried at most once,
then the outcome is cached."""
return LyricsOut.from_entity(await lyrics.get_lyrics(track_id))
@router.post("/{track_id}/lyrics/refetch")
async def refetch_track_lyrics(
track_id: uuid.UUID, lyrics: LyricsServiceDep, _: CurrentUser
) -> LyricsOut:
"""Force a fresh provider lookup, bypassing the cache (user-triggered)."""
return LyricsOut.from_entity(await lyrics.get_lyrics(track_id, force=True))
@router.get("/{track_id}/cover") @router.get("/{track_id}/cover")
+50 -6
View File
@@ -1,23 +1,67 @@
"""User settings endpoints, including scrobbling configuration.""" """User settings endpoints, including scrobbling configuration.
from typing import Any Settings are per-caller and created lazily, so a first read returns defaults.
The scrobbler session key is write-only — accepted on ``PUT`` but never returned.
"""
from fastapi import APIRouter from fastapi import APIRouter
from app.api.deps import CurrentUser, UserSettingsServiceDep
from app.api.schemas.settings import (
ScrobblingOut,
ScrobblingUpdate,
SettingsOut,
SettingsUpdate,
)
from app.domain.entities.settings import UserSettings
router = APIRouter(prefix="/settings", tags=["settings"]) router = APIRouter(prefix="/settings", tags=["settings"])
def _to_settings_out(settings: UserSettings) -> SettingsOut:
return SettingsOut(theme=settings.theme, stream_quality=settings.stream_quality)
def _to_scrobbling_out(settings: UserSettings) -> ScrobblingOut:
return ScrobblingOut(
enabled=settings.scrobble_enabled,
provider=settings.scrobble_provider,
username=settings.scrobble_username,
configured=settings.scrobble_session_key_enc is not None,
)
@router.get("") @router.get("")
async def get_settings() -> Any: ... async def get_settings(user: CurrentUser, service: UserSettingsServiceDep) -> SettingsOut:
return _to_settings_out(await service.get(user.id))
@router.patch("") @router.patch("")
async def update_settings() -> Any: ... async def update_settings(
body: SettingsUpdate, user: CurrentUser, service: UserSettingsServiceDep
) -> SettingsOut:
settings = await service.update_general(
user.id, theme=body.theme, stream_quality=body.stream_quality
)
return _to_settings_out(settings)
@router.get("/scrobbling") @router.get("/scrobbling")
async def get_scrobbling_settings() -> Any: ... async def get_scrobbling_settings(
user: CurrentUser, service: UserSettingsServiceDep
) -> ScrobblingOut:
return _to_scrobbling_out(await service.get(user.id))
@router.put("/scrobbling") @router.put("/scrobbling")
async def set_scrobbling_settings() -> Any: ... async def set_scrobbling_settings(
body: ScrobblingUpdate, user: CurrentUser, service: UserSettingsServiceDep
) -> ScrobblingOut:
settings = await service.set_scrobbling(
user.id,
enabled=body.enabled,
provider=body.provider,
username=body.username,
session_key=body.session_key,
)
return _to_scrobbling_out(settings)
+97
View File
@@ -0,0 +1,97 @@
"""Lyrics service (plan §6.7).
Get-or-fetch with caching: a track's lyrics are served from the DB when present;
on a miss (or an expired ``not_found``) we ask the provider (LRCLIB) once, then
cache the outcome. ``not_found`` is cached with a TTL so tracks that genuinely
have no lyrics aren't looked up on every play, but can eventually be retried.
Degrades gracefully: if the provider is unreachable the lookup just yields a
``not_found`` — the endpoint still returns 200 with empty lyrics, never an error.
"""
import datetime as dt
import uuid
from app.domain.entities.lyrics import Lyrics
from app.domain.errors import NotFoundError
from app.domain.ports import (
AlbumRepository,
ArtistRepository,
LyricsProvider,
LyricsRepository,
TrackRepository,
)
# Re-lookup a cached "not_found" only after this long — long enough not to spam
# the provider, short enough that lyrics added upstream eventually surface.
_NOT_FOUND_TTL = dt.timedelta(days=7)
_STATUS_FOUND = "found"
_STATUS_NOT_FOUND = "not_found"
class LyricsService:
def __init__(
self,
*,
lyrics: LyricsRepository,
tracks: TrackRepository,
artists: ArtistRepository,
albums: AlbumRepository,
provider: LyricsProvider,
) -> None:
self._lyrics = lyrics
self._tracks = tracks
self._artists = artists
self._albums = albums
self._provider = provider
async def get_lyrics(self, track_id: uuid.UUID, *, force: bool = False) -> Lyrics:
"""Return cached lyrics, fetching from the provider on a miss/expiry.
``force`` (the refetch endpoint) bypasses the cache entirely."""
cached = await self._lyrics.get(track_id)
if not force and cached is not None and self._is_fresh(cached):
return cached
return await self._fetch_and_cache(track_id)
def _is_fresh(self, cached: Lyrics) -> bool:
if cached.status == _STATUS_FOUND:
return True
if cached.status == _STATUS_NOT_FOUND:
return dt.datetime.now(dt.UTC) - cached.fetched_at < _NOT_FOUND_TTL
# "pending" (never fetched) → not fresh, go fetch.
return False
async def _fetch_and_cache(self, track_id: uuid.UUID) -> Lyrics:
track = await self._tracks.get_by_id(track_id)
if track is None:
raise NotFoundError(f"Track {track_id} not found.")
artist = await self._artists.get_by_id(track.artist_id)
album = (
await self._albums.get_by_id(track.album_id)
if track.album_id is not None
else None
)
result = await self._provider.fetch(
artist=artist.name if artist else "",
title=track.title,
album=album.title if album else None,
duration_seconds=track.duration_seconds,
)
if result is None:
return await self._lyrics.upsert(
track_id=track_id,
synced=None,
plain=None,
source=None,
status=_STATUS_NOT_FOUND,
)
return await self._lyrics.upsert(
track_id=track_id,
synced=result.synced,
plain=result.plain,
source=result.source,
status=_STATUS_FOUND,
)
+187
View File
@@ -0,0 +1,187 @@
"""Recommendation / radio service (plan §6.5).
Tries the external ML recommender first; when it's unavailable or declines
(returns ``None``), falls back to metadata heuristics over the catalogue — so
similar/radio always work, worse, without ML (graceful-degradation invariant).
``reason`` values are short codes (``ml`` / ``similar`` / ``from_likes`` /
``discover``) the client localizes for the "why is this playing?" affordance.
"""
import random
import uuid
from dataclasses import dataclass
from app.domain.entities.track import Artist, Track
from app.domain.errors import NotFoundError
from app.domain.ports import (
ArtistRepository,
LikeRepository,
Recommender,
TrackRepository,
)
REASON_ML = "ml"
REASON_SIMILAR = "similar"
REASON_FROM_LIKES = "from_likes"
REASON_DISCOVER = "discover"
_LIKED_SEED_POOL = 50
@dataclass(frozen=True, slots=True)
class RadioPick:
track: Track
reason: str
class RecommendationService:
def __init__(
self,
*,
recommender: Recommender,
tracks: TrackRepository,
artists: ArtistRepository,
likes: LikeRepository,
) -> None:
self._recommender = recommender
self._tracks = tracks
self._artists = artists
self._likes = likes
# -- similar ---------------------------------------------------------------
async def similar_tracks(
self, track_id: uuid.UUID, *, limit: int
) -> tuple[str, list[Track]]:
seed = await self._tracks.get_by_id(track_id)
if seed is None:
raise NotFoundError(f"Track {track_id} not found.")
if self._recommender.is_available():
ids = await self._recommender.similar_track_ids(
track_id, limit=limit, exclude_ids=[track_id]
)
if ids is not None:
return REASON_ML, await self._hydrate_tracks(ids)
found = await self._tracks.list_similar(
genre=seed.genre,
artist_id=seed.artist_id,
exclude_ids=[track_id],
limit=limit,
)
return REASON_SIMILAR, found
async def similar_artists(
self, artist_id: uuid.UUID, *, limit: int
) -> tuple[str, list[Artist]]:
if await self._artists.get_by_id(artist_id) is None:
raise NotFoundError(f"Artist {artist_id} not found.")
if self._recommender.is_available():
ids = await self._recommender.similar_artist_ids(artist_id, limit=limit)
if ids is not None:
found = [a for i in ids if (a := await self._artists.get_by_id(i))]
return REASON_ML, found
found = await self._artists.list_similar(artist_id=artist_id, limit=limit)
return REASON_SIMILAR, found
# -- radio -----------------------------------------------------------------
async def radio(
self,
*,
user_id: uuid.UUID,
seed_track_id: uuid.UUID | None,
from_likes: bool,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> tuple[str, list[RadioPick]]:
exploration = min(1.0, max(0.0, exploration))
if self._recommender.is_available():
ids = await self._recommender.radio_track_ids(
seed_track_id=seed_track_id,
exploration=exploration,
limit=limit,
exclude_ids=exclude_ids,
)
if ids is not None:
picks = [
RadioPick(track=t, reason=REASON_ML)
for t in await self._hydrate_tracks(ids)
]
return REASON_ML, picks
return "metadata", await self._radio_fallback(
user_id=user_id,
seed_track_id=seed_track_id,
from_likes=from_likes,
exploration=exploration,
limit=limit,
exclude_ids=exclude_ids,
)
async def _radio_fallback(
self,
*,
user_id: uuid.UUID,
seed_track_id: uuid.UUID | None,
from_likes: bool,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[RadioPick]:
exclude = list(dict.fromkeys(exclude_ids)) # de-dupe, keep order
explore_n = round(limit * exploration)
similar_n = limit - explore_n
picks: list[RadioPick] = []
seed, seed_reason = await self._resolve_seed(
user_id, seed_track_id, from_likes
)
if seed is not None and similar_n > 0:
for track in await self._tracks.list_similar(
genre=seed.genre,
artist_id=seed.artist_id,
exclude_ids=exclude,
limit=similar_n,
):
picks.append(RadioPick(track=track, reason=seed_reason))
exclude.append(track.id)
# Fill the remainder (exploration + any similarity shortfall) with random
# playable tracks — this is also the total fallback when there's no seed.
remaining = limit - len(picks)
if remaining > 0:
for track in await self._tracks.sample_playable(
exclude_ids=exclude, limit=remaining
):
picks.append(RadioPick(track=track, reason=REASON_DISCOVER))
exclude.append(track.id)
random.shuffle(picks)
return picks
async def _resolve_seed(
self,
user_id: uuid.UUID,
seed_track_id: uuid.UUID | None,
from_likes: bool,
) -> tuple[Track | None, str]:
if seed_track_id is not None:
return await self._tracks.get_by_id(seed_track_id), REASON_SIMILAR
if from_likes:
liked = await self._likes.list_liked_tracks(
user_id=user_id, limit=_LIKED_SEED_POOL, offset=0
)
if liked:
return random.choice(liked), REASON_FROM_LIKES
return None, REASON_DISCOVER
async def _hydrate_tracks(self, ids: list[uuid.UUID]) -> list[Track]:
"""Resolve ids → tracks preserving order, skipping any that vanished."""
return [t for i in ids if (t := await self._tracks.get_by_id(i))]
+140
View File
@@ -0,0 +1,140 @@
"""Offline-first sync use cases: delta pull + idempotent push.
The cursor is a server-clock timestamp obtained from the DB (``now``), so it is
immune to app/DB clock skew. A pull returns everything a user changed in the
half-open window ``(since, cursor]``; a push appends the event-log entries a
client accumulated offline. Events carry a client-generated id, so a replay is a
no-op (append is ``ON CONFLICT DO NOTHING``). Events for tracks this server does
not have are skipped rather than rejected (graceful degradation).
Known limitation (v1): the delta is unpaginated and, being wall-clock based, a
write committing right on the cursor boundary under concurrency can slip a cycle
— a periodic full resync (``since=None``) heals it.
"""
import datetime as dt
import uuid
from collections.abc import Awaitable, Callable
from dataclasses import dataclass
from app.domain.entities.history import PlayHistoryEntry
from app.domain.entities.like import Like
from app.domain.entities.playlist import Playlist
from app.domain.entities.track import Track
from app.domain.ports import (
HistoryRepository,
LikeRepository,
PlaylistRepository,
TrackRepository,
)
@dataclass(frozen=True, slots=True)
class LikeEvent:
id: uuid.UUID
track_id: uuid.UUID
value: str
created_at: dt.datetime
@dataclass(frozen=True, slots=True)
class PlayEvent:
id: uuid.UUID
track_id: uuid.UUID
played_at: dt.datetime
play_duration_seconds: int | None
completed: bool
@dataclass(frozen=True, slots=True)
class PlaylistDelta:
playlist: Playlist
track_ids: list[uuid.UUID]
@dataclass(frozen=True, slots=True)
class SyncChanges:
cursor: dt.datetime
likes: list[Like]
plays: list[PlayHistoryEntry]
playlists: list[PlaylistDelta]
tracks: list[Track]
@dataclass(frozen=True, slots=True)
class SyncPushResult:
cursor: dt.datetime
accepted_likes: int
accepted_plays: int
class SyncService:
def __init__(
self,
*,
likes: LikeRepository,
history: HistoryRepository,
playlists: PlaylistRepository,
tracks: TrackRepository,
now: Callable[[], Awaitable[dt.datetime]],
) -> None:
self._likes = likes
self._history = history
self._playlists = playlists
self._tracks = tracks
self._now = now
async def get_changes(self, user_id: uuid.UUID, *, since: dt.datetime | None) -> SyncChanges:
until = await self._now()
likes = await self._likes.list_since(user_id, since=since, until=until)
plays = await self._history.list_since(user_id, since=since, until=until)
changed = await self._playlists.list_changed_since(
owner_id=user_id, since=since, until=until
)
playlists = [
PlaylistDelta(playlist=p, track_ids=await self._playlists.track_ids(p.id))
for p in changed
]
tracks = await self._tracks.list_changed_since(since=since, until=until)
return SyncChanges(
cursor=until, likes=likes, plays=plays, playlists=playlists, tracks=tracks
)
async def push(
self,
user_id: uuid.UUID,
*,
likes: list[LikeEvent],
plays: list[PlayEvent],
) -> SyncPushResult:
accepted_likes = 0
for event in likes:
if await self._tracks.get_by_id(event.track_id) is None:
continue # skip events for tracks this server doesn't have
if await self._likes.add_event(
id=event.id,
user_id=user_id,
track_id=event.track_id,
value=event.value,
created_at=event.created_at,
):
accepted_likes += 1
accepted_plays = 0
for play in plays:
if await self._tracks.get_by_id(play.track_id) is None:
continue
if await self._history.add_event(
id=play.id,
user_id=user_id,
track_id=play.track_id,
played_at=play.played_at,
play_duration_seconds=play.play_duration_seconds,
completed=play.completed,
):
accepted_plays += 1
cursor = await self._now()
return SyncPushResult(
cursor=cursor, accepted_likes=accepted_likes, accepted_plays=accepted_plays
)
+99
View File
@@ -0,0 +1,99 @@
"""Transcode service + cache-path helpers (Group B / plan §6.6).
Cache layout under ``transcode_cache_path``::
{track_id}/opus_{kbps}.opus # direct quality renditions
{track_id}/hls/playlist.m3u8 # HLS rendition (AAC-in-TS)
{track_id}/hls/seg_000.ts …
The request side (streaming router) only *reads* the cache — misses fall back to
the original file and enqueue generation. The worker (``transcode_task``) writes
it. Path helpers are module-level so both sides agree on locations without one
importing the other.
"""
import re
import uuid
from pathlib import Path
from app.domain.errors import NotFoundError
from app.domain.ports import TrackRepository
# Stream-quality name (matches the user-settings ``StreamQuality``) → Opus
# bitrate. ``original`` is absent: it means "serve the master, no transcode".
QUALITY_BITRATE: dict[str, int] = {"high": 128, "medium": 96, "low": 64}
# Single HLS rendition bitrate (AAC). One rendition keeps the MVP simple; a
# multi-bitrate ladder can come later.
HLS_BITRATE = 128
# Only these segment names may be served, guarding the segment route against
# path traversal.
_SEGMENT_RE = re.compile(r"^seg_\d{3,}\.ts$")
def bitrate_for_quality(quality: str) -> int | None:
"""Opus bitrate for a quality name, or ``None`` for ``original``/unknown."""
return QUALITY_BITRATE.get(quality)
def track_cache_dir(root: Path, track_id: uuid.UUID) -> Path:
return root / str(track_id)
def opus_path(root: Path, track_id: uuid.UUID, bitrate_kbps: int) -> Path:
return track_cache_dir(root, track_id) / f"opus_{bitrate_kbps}.opus"
def hls_dir(root: Path, track_id: uuid.UUID) -> Path:
return track_cache_dir(root, track_id) / "hls"
def hls_playlist_path(root: Path, track_id: uuid.UUID) -> Path:
return hls_dir(root, track_id) / "playlist.m3u8"
def hls_segment_path(root: Path, track_id: uuid.UUID, name: str) -> Path | None:
"""Resolve a segment file, or ``None`` if the name is not a valid segment."""
if not _SEGMENT_RE.fullmatch(name):
return None
return hls_dir(root, track_id) / name
class TranscodeService:
"""Request-side cache lookups for transcoded renditions."""
def __init__(self, *, tracks: TrackRepository, cache_root: Path) -> None:
self._tracks = tracks
self._root = cache_root
async def _require_streamable(self, track_id: uuid.UUID) -> None:
track = await self._tracks.get_by_id(track_id)
if track is None:
raise NotFoundError("Track not found.")
if track.storage_uri is None:
raise NotFoundError("Track is not yet downloaded.")
async def resolve_quality_file(
self, track_id: uuid.UUID, quality: str
) -> Path | None:
"""Cached Opus file for ``quality`` if present, else ``None`` (caller
falls back to the master and enqueues generation). ``original`` → None."""
bitrate = bitrate_for_quality(quality)
if bitrate is None:
return None
path = opus_path(self._root, track_id, bitrate)
return path if path.exists() else None
async def hls_playlist(self, track_id: uuid.UUID) -> Path | None:
"""Cached HLS playlist if generated, else ``None`` (validates the track
exists so an unknown id 404s rather than silently missing)."""
await self._require_streamable(track_id)
path = hls_playlist_path(self._root, track_id)
return path if path.exists() else None
def hls_segment(self, track_id: uuid.UUID, name: str) -> Path | None:
path = hls_segment_path(self._root, track_id, name)
if path is None or not path.exists():
return None
return path
+62
View File
@@ -0,0 +1,62 @@
"""User-settings use cases: general preferences + scrobbling configuration.
Settings rows are created lazily — a user who never saved anything reads clean
defaults. Partial updates are merged against current values here, so the
repository always persists the complete desired state. The scrobbler session
key is encrypted before it touches the DB (never stored or returned in plain).
"""
import uuid
from dataclasses import replace
from app.domain.entities.settings import UserSettings
from app.domain.ports import SubsonicCipher, UserSettingsRepository
class UserSettingsService:
def __init__(self, *, settings: UserSettingsRepository, cipher: SubsonicCipher) -> None:
self._settings = settings
# Same Fernet cipher used for the Subsonic app-password — reused here to
# encrypt the scrobbler session key at rest (symmetric, recoverable).
self._cipher = cipher
async def get(self, user_id: uuid.UUID) -> UserSettings:
return await self._settings.get(user_id) or UserSettings.defaults(user_id)
async def update_general(
self, user_id: uuid.UUID, *, theme: str | None, stream_quality: str | None
) -> UserSettings:
current = await self.get(user_id)
merged = replace(
current,
theme=theme if theme is not None else current.theme,
stream_quality=stream_quality if stream_quality is not None else current.stream_quality,
)
return await self._settings.upsert(merged)
async def set_scrobbling(
self,
user_id: uuid.UUID,
*,
enabled: bool,
provider: str | None,
username: str | None,
session_key: str | None,
) -> UserSettings:
"""Replace the scrobbling config. ``session_key`` is write-only: a new
value is encrypted and stored; omitting it keeps the existing key (the
client can't read it back to re-send it)."""
current = await self.get(user_id)
session_key_enc: str | None
if session_key is not None:
session_key_enc = self._cipher.encrypt(session_key)
else:
session_key_enc = current.scrobble_session_key_enc
merged = replace(
current,
scrobble_enabled=enabled,
scrobble_provider=provider,
scrobble_username=username,
scrobble_session_key_enc=session_key_enc,
)
return await self._settings.upsert(merged)
+23
View File
@@ -61,6 +61,17 @@ class Settings(BaseSettings):
# admin-only (POST /admin/users). Registered users are never superusers. # admin-only (POST /admin/users). Registered users are never superusers.
allow_registration: bool = True allow_registration: bool = True
# -- CORS -------------------------------------------------------------
# Origins allowed to call the API from a browser. The web UI is multi-
# instance — it connects to whatever origin the operator types on the
# connect screen — so a page served from origin A may call this backend at
# origin B (e.g. the direct :8000 port, a LAN IP, or 127.0.0.1 vs localhost).
# Auth rides in the ``Authorization`` bearer header (not cookies), so the
# wildcard default is safe here — it is paired with ``allow_credentials=False``.
# Set explicit origins in hardened deployments. Accepts a comma-separated
# string in ``.env`` (``CORS_ALLOW_ORIGINS=https://a,https://b``) or ``*``.
cors_allow_origins: list[str] = Field(default_factory=lambda: ["*"])
# -- subsonic --------------------------------------------------------- # -- subsonic ---------------------------------------------------------
# Symmetric key (any string) used to encrypt each user's recoverable # Symmetric key (any string) used to encrypt each user's recoverable
# Subsonic app-password at rest. A Fernet key is derived from it; rotating # Subsonic app-password at rest. A Fernet key is derived from it; rotating
@@ -70,6 +81,9 @@ class Settings(BaseSettings):
# -- media / storage -------------------------------------------------- # -- media / storage --------------------------------------------------
media_path: Path = Path("/data/media") media_path: Path = Path("/data/media")
transcode_cache_path: Path = Path("/data/transcode-cache") transcode_cache_path: Path = Path("/data/transcode-cache")
# ffmpeg binary for transcoding/HLS (on PATH in the image); override for a
# non-standard location.
ffmpeg_path: str = "ffmpeg"
max_parallel_downloads: int = 2 max_parallel_downloads: int = 2
# How many times the download worker retries a failed fetch (yt-dlp fails # How many times the download worker retries a failed fetch (yt-dlp fails
# often) before marking the job ``failed`` — exponential backoff between tries. # often) before marking the job ``failed`` — exponential backoff between tries.
@@ -127,6 +141,15 @@ class Settings(BaseSettings):
raise ValueError("database_url must use the asyncpg driver: postgresql+asyncpg://") raise ValueError("database_url must use the asyncpg driver: postgresql+asyncpg://")
return v return v
@field_validator("cors_allow_origins", mode="before")
@classmethod
def _split_cors_origins(cls, v: object) -> object:
# Allow a plain comma-separated string in .env (pydantic would otherwise
# try to JSON-decode a list field): "a, b" -> ["a", "b"]; "*" -> ["*"].
if isinstance(v, str):
return [origin.strip() for origin in v.split(",") if origin.strip()]
return v
@property @property
def is_prod(self) -> bool: def is_prod(self) -> bool:
return self.environment == "prod" return self.environment == "prod"
+41
View File
@@ -0,0 +1,41 @@
"""Lyrics value objects (plan §6.7).
``Lyrics`` is the cached row for a track; ``LyricsResult`` is what a provider
(LRCLIB) returns for a lookup. Both cross the domain boundary — no framework
imports. Status values mirror ``LyricsStatus`` in the ORM enum ("found" /
"not_found" / "pending") but are kept as plain strings here so the domain stays
independent of the persistence layer.
"""
import datetime as dt
import uuid
from dataclasses import dataclass
@dataclass(frozen=True, slots=True)
class LyricsResult:
"""A provider hit: synced (timestamped LRC) and/or plain text."""
synced: str | None
plain: str | None
source: str
@dataclass(frozen=True, slots=True)
class Lyrics:
"""Cached lyrics for one track. ``status`` is ``found`` / ``not_found`` /
``pending``; ``not_found`` is cached too (with a TTL in the service) so a
track with no lyrics doesn't hammer the provider on every play."""
track_id: uuid.UUID
synced: str | None
plain: str | None
source: str | None
status: str
fetched_at: dt.datetime
@property
def has_lyrics(self) -> bool:
return self.status == "found" and (
self.synced is not None or self.plain is not None
)
+34
View File
@@ -0,0 +1,34 @@
"""User settings domain entity (general preferences + scrobbling config)."""
import uuid
from dataclasses import dataclass
# Defaults for a user who has never saved settings — the row is created lazily,
# so reads return these before the first write.
DEFAULT_THEME = "system"
DEFAULT_STREAM_QUALITY = "original"
@dataclass(frozen=True, slots=True)
class UserSettings:
user_id: uuid.UUID
theme: str
stream_quality: str
scrobble_enabled: bool
scrobble_provider: str | None
scrobble_username: str | None
# Scrobbler session key / user token, encrypted at rest (never leaves the
# server in plaintext). ``None`` until the user configures scrobbling.
scrobble_session_key_enc: str | None
@classmethod
def defaults(cls, user_id: uuid.UUID) -> UserSettings:
return cls(
user_id=user_id,
theme=DEFAULT_THEME,
stream_quality=DEFAULT_STREAM_QUALITY,
scrobble_enabled=False,
scrobble_provider=None,
scrobble_username=None,
scrobble_session_key_enc=None,
)
+14
View File
@@ -54,6 +54,13 @@ class PermissionDeniedError(DomainError):
code = "permission_denied" code = "permission_denied"
class NotSupportedError(DomainError):
"""Operation is intentionally unsupported (e.g. a config knob that is managed
via environment, not mutable at runtime)."""
code = "not_supported"
class DependencyUnavailableError(DomainError): class DependencyUnavailableError(DomainError):
"""An external dependency (source, ML, MusicBrainz) is unavailable. """An external dependency (source, ML, MusicBrainz) is unavailable.
@@ -69,6 +76,13 @@ class StorageError(DomainError):
code = "storage_error" code = "storage_error"
class TranscodeError(DomainError):
"""Transcoding (ffmpeg) failed. Raised in the worker; a play falls back to
the original file rather than surfacing this."""
code = "transcode_error"
class RangeNotSatisfiableError(DomainError): class RangeNotSatisfiableError(DomainError):
"""Requested byte range cannot be satisfied.""" """Requested byte range cannot be satisfied."""
+135
View File
@@ -29,6 +29,8 @@ from app.domain.entities import (
SubsonicCredentials, SubsonicCredentials,
User, User,
) )
from app.domain.entities.lyrics import Lyrics, LyricsResult
from app.domain.entities.settings import UserSettings
from app.domain.entities.track import Artist, Track from app.domain.entities.track import Artist, Track
from app.domain.sources import DownloadResult, RawMetadata, SearchResult, SourceFile, SourceInfo from app.domain.sources import DownloadResult, RawMetadata, SearchResult, SourceFile, SourceInfo
from app.domain.tokens import IssuedToken, TokenClaims, TokenType from app.domain.tokens import IssuedToken, TokenClaims, TokenType
@@ -56,6 +58,11 @@ class UserRepository(Protocol):
async def set_subsonic_password_enc(self, user_id: uuid.UUID, password_enc: str) -> None: ... async def set_subsonic_password_enc(self, user_id: uuid.UUID, password_enc: str) -> None: ...
class UserSettingsRepository(Protocol):
async def get(self, user_id: uuid.UUID) -> UserSettings | None: ...
async def upsert(self, settings: UserSettings) -> UserSettings: ...
class SubsonicCipher(Protocol): class SubsonicCipher(Protocol):
"""Symmetric encrypt/decrypt for the recoverable Subsonic app-password.""" """Symmetric encrypt/decrypt for the recoverable Subsonic app-password."""
@@ -121,6 +128,12 @@ class ArtistRepository(Protocol):
async def get_by_id(self, artist_id: uuid.UUID) -> Artist | None: ... async def get_by_id(self, artist_id: uuid.UUID) -> Artist | None: ...
async def get_many(self, ids: list[uuid.UUID]) -> list[Artist]: ... async def get_many(self, ids: list[uuid.UUID]) -> list[Artist]: ...
async def list_similar(self, *, artist_id: uuid.UUID, limit: int) -> list[Artist]:
"""Artists sharing the seed artist's genres, ranked by overlap. Metadata
fallback for ``GET /artists/{id}/similar``. Defined before ``list`` so the
``list[Artist]`` annotation isn't shadowed by the method named ``list``."""
...
async def list(self, *, q: str | None, limit: int, offset: int) -> list[Artist]: ... async def list(self, *, q: str | None, limit: int, offset: int) -> list[Artist]: ...
async def count(self, *, q: str | None) -> int: ... async def count(self, *, q: str | None) -> int: ...
async def album_count(self, artist_id: uuid.UUID) -> int: ... async def album_count(self, artist_id: uuid.UUID) -> int: ...
@@ -164,6 +177,33 @@ class TrackRepository(Protocol):
# AlbumRepository below). # AlbumRepository below).
async def genres(self) -> list[tuple[str, int]]: ... async def genres(self) -> list[tuple[str, int]]: ...
async def library_stats(self) -> LibraryStats: ... async def library_stats(self) -> LibraryStats: ...
async def list_similar(
self,
*,
genre: str | None,
artist_id: uuid.UUID,
exclude_ids: list[uuid.UUID],
limit: int,
) -> list[Track]:
"""Playable tracks resembling a seed (same genre and/or artist), ranked
by match strength then shuffled. The metadata fallback for §6.5 radio /
similar when no ML service is configured."""
...
async def sample_playable(
self, *, exclude_ids: list[uuid.UUID], limit: int
) -> list[Track]:
"""Random playable tracks — the exploration filler for radio."""
...
async def find_duplicate_groups(self) -> list[tuple[str, list[Track]]]: ...
async def list_by_metadata_status(
self, status: str, *, limit: int, offset: int
) -> list[Track]: ...
async def all_storage_refs(self) -> list[tuple[uuid.UUID, str]]: ...
async def count_by_metadata_status(self, status: str) -> int: ...
async def list_changed_since(
self, *, since: dt.datetime | None, until: dt.datetime
) -> list[Track]: ...
async def list( async def list(
self, self,
*, *,
@@ -287,12 +327,29 @@ class PlaylistRepository(Protocol):
async def reorder_tracks( async def reorder_tracks(
self, playlist_id: uuid.UUID, ordered_track_ids: list[uuid.UUID] self, playlist_id: uuid.UUID, ordered_track_ids: list[uuid.UUID]
) -> None: ... ) -> None: ...
async def get_cover_path(self, playlist_id: uuid.UUID) -> str | None: ...
async def list_changed_since(
self, *, owner_id: uuid.UUID, since: dt.datetime | None, until: dt.datetime
) -> list[Playlist]: ...
async def track_ids(self, playlist_id: uuid.UUID) -> list[uuid.UUID]: ...
# list must come after any method using list[...] in its signature (name shadowing) # list must come after any method using list[...] in its signature (name shadowing)
async def list(self, *, owner_id: uuid.UUID, limit: int, offset: int) -> list[Playlist]: ... async def list(self, *, owner_id: uuid.UUID, limit: int, offset: int) -> list[Playlist]: ...
class LikeRepository(Protocol): class LikeRepository(Protocol):
async def add(self, *, user_id: uuid.UUID, track_id: uuid.UUID, value: str) -> Like: ... async def add(self, *, user_id: uuid.UUID, track_id: uuid.UUID, value: str) -> Like: ...
async def add_event(
self,
*,
id: uuid.UUID,
user_id: uuid.UUID,
track_id: uuid.UUID,
value: str,
created_at: dt.datetime,
) -> bool: ...
async def list_since(
self, user_id: uuid.UUID, *, since: dt.datetime | None, until: dt.datetime
) -> list[Like]: ...
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]: ...
@@ -312,6 +369,19 @@ class HistoryRepository(Protocol):
play_duration_seconds: int | None, play_duration_seconds: int | None,
completed: bool, completed: bool,
) -> PlayHistoryEntry: ... ) -> PlayHistoryEntry: ...
async def add_event(
self,
*,
id: uuid.UUID,
user_id: uuid.UUID,
track_id: uuid.UUID,
played_at: dt.datetime,
play_duration_seconds: int | None,
completed: bool,
) -> bool: ...
async def list_since(
self, user_id: uuid.UUID, *, since: dt.datetime | None, until: dt.datetime
) -> list[PlayHistoryEntry]: ...
async def list( async def list(
self, *, user_id: uuid.UUID, limit: int, offset: int self, *, user_id: uuid.UUID, limit: int, offset: int
) -> list[PlayHistoryEntry]: ... ) -> list[PlayHistoryEntry]: ...
@@ -449,3 +519,68 @@ class CoverArtProvider(Protocol):
def is_available(self) -> bool: ... def is_available(self) -> bool: ...
async def fetch_release_group(self, release_group_mbid: str) -> CoverArt | None: ... async def fetch_release_group(self, release_group_mbid: str) -> CoverArt | None: ...
class Recommender(Protocol):
"""External ML recommender (plan §6.5, ``ML_SERVICE_URL``). Returns ordered
track/artist ids, or ``None`` when unavailable/erroring so the service falls
back to metadata heuristics — ML is never a hard dependency (invariant)."""
def is_available(self) -> bool: ...
async def similar_track_ids(
self, track_id: uuid.UUID, *, limit: int, exclude_ids: list[uuid.UUID]
) -> list[uuid.UUID] | None: ...
async def similar_artist_ids(
self, artist_id: uuid.UUID, *, limit: int
) -> list[uuid.UUID] | None: ...
async def radio_track_ids(
self,
*,
seed_track_id: uuid.UUID | None,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[uuid.UUID] | None: ...
class Transcoder(Protocol):
"""Transcodes an audio file with ffmpeg (plan §6.6 / Group B). ``to_opus``
writes a single Opus rendition; ``to_hls`` writes an HLS playlist + segments
(AAC-in-TS) into ``out_dir``. Both raise ``TranscodeError`` on failure; heavy
work always runs in a worker, never the request cycle."""
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None: ...
async def to_hls(self, src: Path, out_dir: Path, *, bitrate_kbps: int) -> None: ...
class LyricsProvider(Protocol):
"""Fetches lyrics from an external database (LRCLIB) by artist/title/album/
duration. Returns a hit or ``None`` (no match / service down), never raising."""
async def fetch(
self,
*,
artist: str,
title: str,
album: str | None,
duration_seconds: int | None,
) -> LyricsResult | None: ...
class LyricsRepository(Protocol):
"""Cached lyrics, one row per track. ``upsert`` also caches a ``not_found``
(empty text) so misses aren't re-fetched until the service's TTL lapses."""
async def get(self, track_id: uuid.UUID) -> Lyrics | None: ...
async def upsert(
self,
*,
track_id: uuid.UUID,
synced: str | None,
plain: str | None,
source: str | None,
status: str,
) -> Lyrics: ...
+2
View File
@@ -14,6 +14,7 @@ from app.infrastructure.db.models.play_history import PlayHistoryModel
from app.infrastructure.db.models.playlist import PlaylistModel, PlaylistTrackModel from app.infrastructure.db.models.playlist import PlaylistModel, PlaylistTrackModel
from app.infrastructure.db.models.track import TrackModel from app.infrastructure.db.models.track import TrackModel
from app.infrastructure.db.models.user import RefreshTokenModel, UserModel from app.infrastructure.db.models.user import RefreshTokenModel, UserModel
from app.infrastructure.db.models.user_settings import UserSettingsModel
__all__ = [ __all__ = [
"AlbumModel", "AlbumModel",
@@ -27,4 +28,5 @@ __all__ = [
"RefreshTokenModel", "RefreshTokenModel",
"TrackModel", "TrackModel",
"UserModel", "UserModel",
"UserSettingsModel",
] ]
+9
View File
@@ -37,3 +37,12 @@ class LikeModel(UUIDPrimaryKeyMixin, Base):
server_default=func.now(), server_default=func.now(),
nullable=False, nullable=False,
) )
# Server ingestion time — the delta-sync ordering key. Set once at insert and
# never changed; distinct from ``created_at`` (the real event time, which a
# sync push preserves from the client even when it happened offline earlier).
synced_at: Mapped[dt.datetime] = mapped_column(
DateTime(timezone=True),
server_default=func.now(),
nullable=False,
index=True,
)
@@ -34,3 +34,11 @@ class PlayHistoryModel(UUIDPrimaryKeyMixin, Base):
) )
play_duration_seconds: Mapped[int | None] = mapped_column(Integer, nullable=True) play_duration_seconds: Mapped[int | None] = mapped_column(Integer, nullable=True)
completed: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False) completed: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
# Server ingestion time — the delta-sync ordering key (see LikeModel). Distinct
# from ``played_at`` (the real play time), which a sync push preserves.
synced_at: Mapped[dt.datetime] = mapped_column(
DateTime(timezone=True),
server_default=func.now(),
nullable=False,
index=True,
)
@@ -0,0 +1,29 @@
"""ORM model for per-user settings (general preferences + scrobbling)."""
import uuid
from sqlalchemy import Boolean, ForeignKey, String
from sqlalchemy.orm import Mapped, mapped_column
from app.infrastructure.db.base import Base
from app.infrastructure.db.models.mixins import TimestampMixin
class UserSettingsModel(TimestampMixin, Base):
"""One row per user, created lazily on first save. The primary key *is* the
user id (a 1:1 extension of ``users``), so there's no separate surrogate id."""
__tablename__ = "user_settings"
user_id: Mapped[uuid.UUID] = mapped_column(
ForeignKey("users.id", ondelete="CASCADE"),
primary_key=True,
)
theme: Mapped[str] = mapped_column(String(16), default="system", nullable=False)
stream_quality: Mapped[str] = mapped_column(String(16), default="original", nullable=False)
scrobble_enabled: Mapped[bool] = mapped_column(Boolean, default=False, nullable=False)
scrobble_provider: Mapped[str | None] = mapped_column(String(16), nullable=True)
scrobble_username: Mapped[str | None] = mapped_column(String(255), nullable=True)
# Fernet-encrypted scrobbler session key / token (see core.security). Never
# the plaintext — mirrors how the Subsonic app-password is stored.
scrobble_session_key_enc: Mapped[str | None] = mapped_column(String(512), nullable=True)
@@ -7,12 +7,16 @@ from app.infrastructure.db.repositories.download_job_repository import (
) )
from app.infrastructure.db.repositories.history_repository import SqlAlchemyHistoryRepository from app.infrastructure.db.repositories.history_repository import SqlAlchemyHistoryRepository
from app.infrastructure.db.repositories.like_repository import SqlAlchemyLikeRepository from app.infrastructure.db.repositories.like_repository import SqlAlchemyLikeRepository
from app.infrastructure.db.repositories.lyrics_repository import SqlAlchemyLyricsRepository
from app.infrastructure.db.repositories.playlist_repository import SqlAlchemyPlaylistRepository from app.infrastructure.db.repositories.playlist_repository import SqlAlchemyPlaylistRepository
from app.infrastructure.db.repositories.refresh_token_repository import ( from app.infrastructure.db.repositories.refresh_token_repository import (
SqlAlchemyRefreshTokenRepository, SqlAlchemyRefreshTokenRepository,
) )
from app.infrastructure.db.repositories.track_repository import SqlAlchemyTrackRepository from app.infrastructure.db.repositories.track_repository import SqlAlchemyTrackRepository
from app.infrastructure.db.repositories.user_repository import SqlAlchemyUserRepository from app.infrastructure.db.repositories.user_repository import SqlAlchemyUserRepository
from app.infrastructure.db.repositories.user_settings_repository import (
SqlAlchemyUserSettingsRepository,
)
__all__ = [ __all__ = [
"SqlAlchemyAlbumRepository", "SqlAlchemyAlbumRepository",
@@ -20,8 +24,10 @@ __all__ = [
"SqlAlchemyDownloadJobRepository", "SqlAlchemyDownloadJobRepository",
"SqlAlchemyHistoryRepository", "SqlAlchemyHistoryRepository",
"SqlAlchemyLikeRepository", "SqlAlchemyLikeRepository",
"SqlAlchemyLyricsRepository",
"SqlAlchemyPlaylistRepository", "SqlAlchemyPlaylistRepository",
"SqlAlchemyRefreshTokenRepository", "SqlAlchemyRefreshTokenRepository",
"SqlAlchemyTrackRepository", "SqlAlchemyTrackRepository",
"SqlAlchemyUserRepository", "SqlAlchemyUserRepository",
"SqlAlchemyUserSettingsRepository",
] ]
@@ -77,6 +77,26 @@ class SqlAlchemyArtistRepository:
) )
return [_to_entity(r) for r in rows] return [_to_entity(r) for r in rows]
async def list_similar(self, *, artist_id: uuid.UUID, limit: int) -> list[Artist]:
# Artists whose tracks fall in the seed artist's genres, ranked by how
# many such tracks they have. Defined before ``list`` so the ``list[Artist]``
# return annotation isn't shadowed by the method named ``list``.
seed_genres = (
select(TrackModel.genre)
.where(TrackModel.artist_id == artist_id, TrackModel.genre.is_not(None))
.distinct()
)
stmt = (
select(ArtistModel)
.join(TrackModel, TrackModel.artist_id == ArtistModel.id)
.where(TrackModel.genre.in_(seed_genres), ArtistModel.id != artist_id)
.group_by(ArtistModel.id)
.order_by(func.count(TrackModel.id).desc())
.limit(limit)
)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def list(self, *, q: str | None, limit: int, offset: int) -> list[Artist]: async def list(self, *, q: str | None, limit: int, offset: int) -> list[Artist]:
stmt = select(ArtistModel) stmt = select(ArtistModel)
if q: if q:
@@ -108,3 +128,4 @@ class SqlAlchemyArtistRepository:
.where(TrackModel.artist_id == artist_id) .where(TrackModel.artist_id == artist_id)
) )
).scalar_one() ).scalar_one()
@@ -4,6 +4,7 @@ import datetime as dt
import uuid import uuid
from sqlalchemy import func, select from sqlalchemy import func, select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.domain.entities.history import PlayHistoryEntry from app.domain.entities.history import PlayHistoryEntry
@@ -46,6 +47,51 @@ class SqlAlchemyHistoryRepository:
await self._session.refresh(row) await self._session.refresh(row)
return _to_entity(row) return _to_entity(row)
async def add_event(
self,
*,
id: uuid.UUID,
user_id: uuid.UUID,
track_id: uuid.UUID,
played_at: dt.datetime,
play_duration_seconds: int | None,
completed: bool,
) -> bool:
"""Idempotent append for sync push: insert a client-generated play event,
skipping it if the ``id`` already exists (a replay). Returns whether a new
row was stored. Defined before ``list`` (name-shadowing)."""
stmt = (
pg_insert(PlayHistoryModel)
.values(
id=id,
user_id=user_id,
track_id=track_id,
played_at=played_at,
play_duration_seconds=play_duration_seconds,
completed=completed,
)
.on_conflict_do_nothing(index_elements=["id"])
.returning(PlayHistoryModel.id)
)
inserted = (await self._session.execute(stmt)).scalar_one_or_none()
return inserted is not None
async def list_since(
self, user_id: uuid.UUID, *, since: dt.datetime | None, until: dt.datetime
) -> list[PlayHistoryEntry]:
"""Play events for a user ingested in the half-open window ``(since,
until]`` (by ``synced_at``, the server-side sync key — so events pushed
with an older ``played_at`` still surface), oldest first. ``since=None``
returns everything up to ``until``."""
stmt = select(PlayHistoryModel).where(
PlayHistoryModel.user_id == user_id, PlayHistoryModel.synced_at <= until
)
if since is not None:
stmt = stmt.where(PlayHistoryModel.synced_at > since)
stmt = stmt.order_by(PlayHistoryModel.synced_at)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def list(self, *, user_id: uuid.UUID, limit: int, offset: int) -> list[PlayHistoryEntry]: async def list(self, *, user_id: uuid.UUID, limit: int, offset: int) -> list[PlayHistoryEntry]:
rows = ( rows = (
( (
@@ -3,9 +3,11 @@
Likes are an append-only event log. Current state = latest event per (user, track). Likes are an append-only event log. Current state = latest event per (user, track).
""" """
import datetime as dt
import uuid import uuid
from sqlalchemy import func, select from sqlalchemy import func, select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.domain.entities.like import Like from app.domain.entities.like import Like
@@ -59,6 +61,50 @@ class SqlAlchemyLikeRepository:
await self._session.refresh(row) await self._session.refresh(row)
return _to_entity(row) return _to_entity(row)
async def add_event(
self,
*,
id: uuid.UUID,
user_id: uuid.UUID,
track_id: uuid.UUID,
value: str,
created_at: dt.datetime,
) -> bool:
"""Idempotent append for sync push: insert a client-generated like event,
skipping it if the ``id`` already exists (a replay). Preserves the
client's ``created_at`` (the event happened offline earlier). Returns
whether a new row was stored."""
stmt = (
pg_insert(LikeModel)
.values(
id=id,
user_id=user_id,
track_id=track_id,
value=value,
created_at=created_at,
)
.on_conflict_do_nothing(index_elements=["id"])
.returning(LikeModel.id)
)
inserted = (await self._session.execute(stmt)).scalar_one_or_none()
return inserted is not None
async def list_since(
self, user_id: uuid.UUID, *, since: dt.datetime | None, until: dt.datetime
) -> list[Like]:
"""Like events for a user ingested in the half-open window ``(since,
until]`` (by ``synced_at``, the server-side sync key — so events pushed
with an older ``created_at`` still surface), oldest first. ``since=None``
returns everything up to ``until``."""
stmt = select(LikeModel).where(
LikeModel.user_id == user_id, LikeModel.synced_at <= until
)
if since is not None:
stmt = stmt.where(LikeModel.synced_at > since)
stmt = stmt.order_by(LikeModel.synced_at)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
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]:
@@ -0,0 +1,72 @@
"""Lyrics repository — adapter over ``AsyncSession``.
One cached row per track (``track_id`` unique). ``upsert`` refreshes the row and
bumps ``fetched_at`` so the service's TTL is measured from the last fetch.
"""
import uuid
from sqlalchemy import func, select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession
from app.domain.entities.lyrics import Lyrics
from app.infrastructure.db.models.lyrics import LyricsModel
def _to_entity(row: LyricsModel) -> Lyrics:
return Lyrics(
track_id=row.track_id,
synced=row.synced,
plain=row.plain,
source=row.source,
status=row.status,
fetched_at=row.fetched_at,
)
class SqlAlchemyLyricsRepository:
def __init__(self, session: AsyncSession) -> None:
self._session = session
async def get(self, track_id: uuid.UUID) -> Lyrics | None:
row = await self._session.scalar(
select(LyricsModel).where(LyricsModel.track_id == track_id)
)
return _to_entity(row) if row is not None else None
async def upsert(
self,
*,
track_id: uuid.UUID,
synced: str | None,
plain: str | None,
source: str | None,
status: str,
) -> Lyrics:
values = {
"track_id": track_id,
"synced": synced,
"plain": plain,
"source": source,
"status": status,
"fetched_at": func.now(),
}
stmt = (
pg_insert(LyricsModel)
.values(**values)
.on_conflict_do_update(
index_elements=[LyricsModel.track_id],
set_={
"synced": synced,
"plain": plain,
"source": source,
"status": status,
"fetched_at": func.now(),
},
)
.returning(LyricsModel)
)
row = (await self._session.scalars(stmt)).one()
await self._session.flush()
return _to_entity(row)
@@ -1,5 +1,6 @@
"""Playlist repository — adapter over ``AsyncSession``.""" """Playlist repository — adapter over ``AsyncSession``."""
import datetime as dt
import uuid import uuid
from sqlalchemy import func, select from sqlalchemy import func, select
@@ -207,6 +208,41 @@ class SqlAlchemyPlaylistRepository:
).scalar_one_or_none() ).scalar_one_or_none()
return float(result) if result is not None else 0.0 return float(result) if result is not None else 0.0
async def get_cover_path(self, playlist_id: uuid.UUID) -> str | None:
"""The playlist's stored cover key, or ``None`` (missing playlist or no
cover). Read directly — the entity doesn't carry the storage key."""
return (
await self._session.execute(
select(PlaylistModel.cover_path).where(PlaylistModel.id == playlist_id)
)
).scalar_one_or_none()
async def list_changed_since(
self, *, owner_id: uuid.UUID, since: dt.datetime | None, until: dt.datetime
) -> list[Playlist]:
"""A user's playlists changed in the window ``(since, until]`` (metadata
or membership — every mutation bumps ``updated_at``/``version``). Defined
before ``list`` (name-shadowing)."""
stmt = select(PlaylistModel).where(
PlaylistModel.owner_id == owner_id, PlaylistModel.updated_at <= until
)
if since is not None:
stmt = stmt.where(PlaylistModel.updated_at > since)
stmt = stmt.order_by(PlaylistModel.updated_at)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def track_ids(self, playlist_id: uuid.UUID) -> list[uuid.UUID]:
"""Ordered track ids of a playlist (by position) — for sync payloads."""
rows = (
await self._session.execute(
select(PlaylistTrackModel.track_id)
.where(PlaylistTrackModel.playlist_id == playlist_id)
.order_by(PlaylistTrackModel.position)
)
).scalars().all()
return list(rows)
# list must come after methods using list[...] in signatures (builtin name shadowing) # list must come after methods using list[...] in signatures (builtin name shadowing)
async def list(self, *, owner_id: uuid.UUID, limit: int, offset: int) -> list[Playlist]: async def list(self, *, owner_id: uuid.UUID, limit: int, offset: int) -> list[Playlist]:
rows = ( rows = (
@@ -3,7 +3,7 @@
import datetime as dt import datetime as dt
import uuid import uuid
from sqlalchemy import func, select from sqlalchemy import case, func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.domain.entities.storage import FormatBreakdown, LibraryStats from app.domain.entities.storage import FormatBreakdown, LibraryStats
@@ -136,6 +136,42 @@ class SqlAlchemyTrackRepository:
).all() ).all()
return [(row.genre, row.cnt) for row in rows] return [(row.genre, row.cnt) for row in rows]
async def list_similar(
self,
*,
genre: str | None,
artist_id: uuid.UUID,
exclude_ids: list[uuid.UUID],
limit: int,
) -> list[Track]:
# Rank a same-genre hit above a same-artist hit; shuffle within a tier so
# the mix varies. Only playable (locally-stored) tracks are candidates.
if genre is not None:
match = or_(TrackModel.genre == genre, TrackModel.artist_id == artist_id)
score = case((TrackModel.genre == genre, 2), else_=0) + case(
(TrackModel.artist_id == artist_id, 1), else_=0
)
else:
match = TrackModel.artist_id == artist_id
score = case((TrackModel.artist_id == artist_id, 1), else_=0)
stmt = select(TrackModel).where(TrackModel.storage_uri.is_not(None), match)
if exclude_ids:
stmt = stmt.where(TrackModel.id.not_in(exclude_ids))
stmt = stmt.order_by(score.desc(), func.random()).limit(limit)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def sample_playable(
self, *, exclude_ids: list[uuid.UUID], limit: int
) -> list[Track]:
stmt = select(TrackModel).where(TrackModel.storage_uri.is_not(None))
if exclude_ids:
stmt = stmt.where(TrackModel.id.not_in(exclude_ids))
stmt = stmt.order_by(func.random()).limit(limit)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def library_stats(self) -> LibraryStats: async def library_stats(self) -> LibraryStats:
"""One-shot aggregate over the whole catalogue (no pagination). Defined """One-shot aggregate over the whole catalogue (no pagination). Defined
before ``list`` for the same shadowing reason as ``genres``.""" before ``list`` for the same shadowing reason as ``genres``."""
@@ -194,6 +230,91 @@ class SqlAlchemyTrackRepository:
by_source={source: cnt for source, cnt in source_rows}, by_source={source: cnt for source, cnt in source_rows},
) )
async def find_duplicate_groups(self) -> list[tuple[str, list[Track]]]:
"""Tracks that share an ``acoustid_fingerprint`` (the dedup key), grouped
by it — only fingerprints with more than one track. Empty when clean.
Defined before ``list`` for the same name-shadowing reason as ``genres``."""
dup_fps = (
select(TrackModel.acoustid_fingerprint)
.where(TrackModel.acoustid_fingerprint.is_not(None))
.group_by(TrackModel.acoustid_fingerprint)
.having(func.count(TrackModel.id) > 1)
.scalar_subquery()
)
rows = (
(
await self._session.execute(
select(TrackModel)
.where(TrackModel.acoustid_fingerprint.in_(dup_fps))
.order_by(TrackModel.acoustid_fingerprint, TrackModel.created_at)
)
)
.scalars()
.all()
)
groups: dict[str, list[Track]] = {}
for row in rows:
fingerprint = row.acoustid_fingerprint
assert fingerprint is not None # filtered to non-null above
groups.setdefault(fingerprint, []).append(_to_entity(row))
return list(groups.items())
async def list_by_metadata_status(
self, status: str, *, limit: int, offset: int
) -> list[Track]:
"""Tracks in a given ``metadata_status`` (e.g. ``pending``/``failed``),
newest first. Defined before ``list`` (name-shadowing)."""
rows = (
(
await self._session.execute(
select(TrackModel)
.where(TrackModel.metadata_status == status)
.order_by(TrackModel.created_at.desc())
.limit(limit)
.offset(offset)
)
)
.scalars()
.all()
)
return [_to_entity(r) for r in rows]
async def all_storage_refs(self) -> list[tuple[uuid.UUID, str]]:
"""``(id, storage_uri)`` for every *local* track — for the cleanup
worker's filesystem reconciliation. Remote placeholders have no local
file (``availability != local``) and are skipped. No entity hydration."""
rows = (
await self._session.execute(
select(TrackModel.id, TrackModel.storage_uri).where(
TrackModel.availability == TrackAvailability.LOCAL.value,
TrackModel.storage_uri.is_not(None),
)
)
).all()
return [(row.id, row.storage_uri) for row in rows]
async def count_by_metadata_status(self, status: str) -> int:
return (
await self._session.execute(
select(func.count())
.select_from(TrackModel)
.where(TrackModel.metadata_status == status)
)
).scalar_one()
async def list_changed_since(
self, *, since: dt.datetime | None, until: dt.datetime
) -> list[Track]:
"""Catalogue tracks changed in the window ``(since, until]`` (by
``updated_at``), oldest first — the delta a client caches for offline use.
Defined before ``list`` (name-shadowing)."""
stmt = select(TrackModel).where(TrackModel.updated_at <= until)
if since is not None:
stmt = stmt.where(TrackModel.updated_at > since)
stmt = stmt.order_by(TrackModel.updated_at)
rows = (await self._session.execute(stmt)).scalars().all()
return [_to_entity(r) for r in rows]
async def list( async def list(
self, self,
*, *,
@@ -0,0 +1,47 @@
"""User-settings repository — adapter over ``AsyncSession``."""
import uuid
from sqlalchemy.ext.asyncio import AsyncSession
from app.domain.entities.settings import UserSettings
from app.infrastructure.db.models.user_settings import UserSettingsModel
def _to_entity(row: UserSettingsModel) -> UserSettings:
return UserSettings(
user_id=row.user_id,
theme=row.theme,
stream_quality=row.stream_quality,
scrobble_enabled=row.scrobble_enabled,
scrobble_provider=row.scrobble_provider,
scrobble_username=row.scrobble_username,
scrobble_session_key_enc=row.scrobble_session_key_enc,
)
class SqlAlchemyUserSettingsRepository:
def __init__(self, session: AsyncSession) -> None:
self._session = session
async def get(self, user_id: uuid.UUID) -> UserSettings | None:
row = await self._session.get(UserSettingsModel, user_id)
return _to_entity(row) if row is not None else None
async def upsert(self, settings: UserSettings) -> UserSettings:
"""Create or replace the caller's settings row with ``settings`` in full.
The service merges partial updates against current values before calling
this, so the write always carries the complete desired state."""
row = await self._session.get(UserSettingsModel, settings.user_id)
if row is None:
row = UserSettingsModel(user_id=settings.user_id)
self._session.add(row)
row.theme = settings.theme
row.stream_quality = settings.stream_quality
row.scrobble_enabled = settings.scrobble_enabled
row.scrobble_provider = settings.scrobble_provider
row.scrobble_username = settings.scrobble_username
row.scrobble_session_key_enc = settings.scrobble_session_key_enc
await self._session.flush()
await self._session.refresh(row)
return _to_entity(row)
+100
View File
@@ -0,0 +1,100 @@
"""LrclibHttpClient — fetches lyrics from LRCLIB (plan §6.7).
LRCLIB is a free, keyless lyrics database. ``/api/get`` does an exact match on
artist+track+album+duration; if that misses we fall back to ``/api/search`` and
take the best-scoring hit. Graceful degradation: any network/parse error →
``fetch`` returns ``None`` (the service then caches a ``not_found``), never
raising. No API key is needed, so this provider is always "available".
"""
import httpx
from app.core.logging import get_logger
from app.domain.entities.lyrics import LyricsResult
log = get_logger(__name__)
_BASE_URL = "https://lrclib.net"
_TIMEOUT_SECONDS = 10.0
_SOURCE = "lrclib"
class LrclibHttpClient:
"""Implements :class:`app.domain.ports.LyricsProvider`."""
def __init__(self, *, user_agent: str, base_url: str = _BASE_URL) -> None:
self._user_agent = user_agent
self._base_url = base_url.rstrip("/")
async def fetch(
self,
*,
artist: str,
title: str,
album: str | None,
duration_seconds: int | None,
) -> LyricsResult | None:
try:
async with httpx.AsyncClient(
timeout=_TIMEOUT_SECONDS,
headers={"User-Agent": self._user_agent},
base_url=self._base_url,
) as client:
hit = await self._get(client, artist, title, album, duration_seconds)
if hit is None:
hit = await self._search(client, artist, title)
except (httpx.HTTPError, ValueError) as exc:
log.warning("lrclib.fetch_failed", error=str(exc))
return None
return hit
async def _get(
self,
client: httpx.AsyncClient,
artist: str,
title: str,
album: str | None,
duration_seconds: int | None,
) -> LyricsResult | None:
"""Exact match via ``/api/get`` (404 when nothing matches exactly)."""
params = {"artist_name": artist, "track_name": title}
if album:
params["album_name"] = album
if duration_seconds is not None:
params["duration"] = str(duration_seconds)
resp = await client.get("/api/get", params=params)
if resp.status_code == httpx.codes.NOT_FOUND:
return None
resp.raise_for_status()
return _to_result(resp.json())
async def _search(
self, client: httpx.AsyncClient, artist: str, title: str
) -> LyricsResult | None:
"""Fuzzy fallback via ``/api/search`` — take the first usable hit."""
resp = await client.get(
"/api/search", params={"artist_name": artist, "track_name": title}
)
resp.raise_for_status()
results = resp.json()
if not isinstance(results, list):
return None
for item in results:
result = _to_result(item)
if result is not None:
return result
return None
def _to_result(payload: object) -> LyricsResult | None:
"""Map an LRCLIB record to a ``LyricsResult``. Instrumental tracks and empty
records yield ``None`` (nothing worth caching as "found")."""
if not isinstance(payload, dict):
return None
if payload.get("instrumental"):
return None
synced = payload.get("syncedLyrics") or None
plain = payload.get("plainLyrics") or None
if synced is None and plain is None:
return None
return LyricsResult(synced=synced, plain=plain, source=_SOURCE)
+1
View File
@@ -0,0 +1 @@
"""ML/recommender adapters (plan §6.5). ML is optional — see Recommender port."""
+110
View File
@@ -0,0 +1,110 @@
"""Recommender adapters (plan §6.5).
``NullRecommender`` is the default: no ML service, so every method reports
unavailable and the ``RecommendationService`` uses its metadata fallback. When
an embedding service exists, wire ``RemoteRecommender`` (skeleton below) to
``ML_SERVICE_URL`` — its exact request/response contract is TODO pending that
service. Both keep the invariant: ML is optional, never a hard dependency.
"""
import uuid
import httpx
from app.core.logging import get_logger
log = get_logger(__name__)
class NullRecommender:
"""No ML configured — always unavailable, always ``None`` (→ fallback)."""
def is_available(self) -> bool:
return False
async def similar_track_ids(
self, track_id: uuid.UUID, *, limit: int, exclude_ids: list[uuid.UUID]
) -> list[uuid.UUID] | None:
return None
async def similar_artist_ids(
self, artist_id: uuid.UUID, *, limit: int
) -> list[uuid.UUID] | None:
return None
async def radio_track_ids(
self,
*,
seed_track_id: uuid.UUID | None,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[uuid.UUID] | None:
return None
_TIMEOUT_SECONDS = 5.0
class RemoteRecommender:
"""HTTP client for an external embedding/recommender service.
TODO: the request/response schema below is a placeholder — align it with the
real ML service once its contract is known. Until then this stays unused
(``deps`` wires ``NullRecommender``). Every call is defensive: any error or a
malformed body returns ``None`` so the service degrades to metadata, matching
the graceful-degradation invariant.
"""
def __init__(self, base_url: str) -> None:
self._base_url = base_url.rstrip("/")
def is_available(self) -> bool:
return True
async def _post_ids(self, path: str, payload: dict[str, object]) -> list[uuid.UUID] | None:
try:
async with httpx.AsyncClient(timeout=_TIMEOUT_SECONDS) as client:
resp = await client.post(f"{self._base_url}{path}", json=payload)
resp.raise_for_status()
data = resp.json()
ids = data.get("track_ids") if isinstance(data, dict) else None
if not isinstance(ids, list):
return None
return [uuid.UUID(str(i)) for i in ids]
except (httpx.HTTPError, ValueError, KeyError) as exc:
log.warning("recommender.remote_failed", path=path, error=str(exc))
return None
async def similar_track_ids(
self, track_id: uuid.UUID, *, limit: int, exclude_ids: list[uuid.UUID]
) -> list[uuid.UUID] | None:
return await self._post_ids(
"/similar/tracks",
{"track_id": str(track_id), "limit": limit,
"exclude": [str(i) for i in exclude_ids]},
)
async def similar_artist_ids(
self, artist_id: uuid.UUID, *, limit: int
) -> list[uuid.UUID] | None:
# Artist recommendations aren't part of the placeholder track contract.
return None
async def radio_track_ids(
self,
*,
seed_track_id: uuid.UUID | None,
exploration: float,
limit: int,
exclude_ids: list[uuid.UUID],
) -> list[uuid.UUID] | None:
return await self._post_ids(
"/radio",
{
"seed_track_id": str(seed_track_id) if seed_track_id else None,
"exploration": exploration,
"limit": limit,
"exclude": [str(i) for i in exclude_ids],
},
)
+4
View File
@@ -48,6 +48,10 @@ class SourceRegistry:
"""Every registered source that supports search (for cross-source search).""" """Every registered source that supports search (for cross-source search)."""
return [cast(SearchableSource, b) for b in self._by_name.values() if hasattr(b, "search")] return [cast(SearchableSource, b) for b in self._by_name.values() if hasattr(b, "search")]
def indexables(self) -> list[IndexableSource]:
"""Every registered source that can be indexed (for a full re-scan)."""
return [cast(IndexableSource, b) for b in self._by_name.values() if hasattr(b, "scan")]
def infos(self) -> list[SourceInfo]: def infos(self) -> list[SourceInfo]:
return [backend.info() for backend in self._by_name.values()] return [backend.info() for backend in self._by_name.values()]
+1
View File
@@ -0,0 +1 @@
"""ffmpeg-based transcoding adapters (Group B / plan §6.6)."""
+67
View File
@@ -0,0 +1,67 @@
"""FfmpegTranscoder — Opus + HLS renditions via the ffmpeg CLI.
Implements :class:`app.domain.ports.Transcoder`. Runs ffmpeg as a subprocess
(only ever from a worker — CLAUDE.md: no heavy work in the request cycle) and
raises :class:`TranscodeError` on a non-zero exit. Opus is used for the direct
quality renditions; HLS segments are AAC-in-MPEG-TS (broad player support).
"""
import asyncio
from pathlib import Path
import anyio
from app.core.logging import get_logger
from app.domain.errors import TranscodeError
log = get_logger(__name__)
# HLS segment length. 10s is a common VOD default — few requests, quick seeks.
_HLS_SEGMENT_SECONDS = "10"
def _mkdir(path: Path) -> None:
path.mkdir(parents=True, exist_ok=True)
class FfmpegTranscoder:
def __init__(self, ffmpeg_path: str = "ffmpeg") -> None:
self._ffmpeg = ffmpeg_path
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None:
await anyio.to_thread.run_sync(_mkdir, dest.parent)
tmp = dest.with_suffix(dest.suffix + ".part")
await self._run(
"-i", str(src),
"-vn", "-c:a", "libopus", "-b:a", f"{bitrate_kbps}k",
"-f", "opus", str(tmp),
)
# Publish atomically so a concurrent reader never sees a half-written file.
await anyio.to_thread.run_sync(tmp.replace, dest)
async def to_hls(self, src: Path, out_dir: Path, *, bitrate_kbps: int) -> None:
await anyio.to_thread.run_sync(_mkdir, out_dir)
playlist = out_dir / "playlist.m3u8"
segments = out_dir / "seg_%03d.ts"
await self._run(
"-i", str(src),
"-vn", "-c:a", "aac", "-b:a", f"{bitrate_kbps}k",
"-f", "hls",
"-hls_time", _HLS_SEGMENT_SECONDS,
"-hls_playlist_type", "vod",
"-hls_segment_filename", str(segments),
str(playlist),
)
async def _run(self, *args: str) -> None:
cmd = [self._ffmpeg, "-y", "-nostdin", "-loglevel", "error", *args]
proc = await asyncio.create_subprocess_exec(
*cmd,
stdout=asyncio.subprocess.DEVNULL,
stderr=asyncio.subprocess.PIPE,
)
_, stderr = await proc.communicate()
if proc.returncode != 0:
detail = stderr.decode(errors="replace").strip()[-500:]
log.error("ffmpeg_failed", returncode=proc.returncode, error=detail)
raise TranscodeError(f"ffmpeg exited {proc.returncode}: {detail}")
+14
View File
@@ -4,6 +4,7 @@ from collections.abc import AsyncIterator
from contextlib import asynccontextmanager from contextlib import asynccontextmanager
from fastapi import FastAPI, WebSocket from fastapi import FastAPI, WebSocket
from fastapi.middleware.cors import CORSMiddleware
from app.api.errors import register_exception_handlers from app.api.errors import register_exception_handlers
from app.api.health import router as health_router from app.api.health import router as health_router
@@ -40,6 +41,19 @@ def create_app() -> FastAPI:
) )
app.add_middleware(CorrelationIdMiddleware) app.add_middleware(CorrelationIdMiddleware)
# CORS added last → outermost, so browser preflight (OPTIONS) is answered
# before anything else. The web UI can connect cross-origin (direct :8000,
# a LAN IP, 127.0.0.1 vs localhost), which needs these headers. Bearer-token
# auth (no cookies) → wildcard origins are safe with allow_credentials=False.
if settings.cors_allow_origins:
app.add_middleware(
CORSMiddleware,
allow_origins=settings.cors_allow_origins,
allow_credentials=False,
allow_methods=["*"],
allow_headers=["*"],
expose_headers=["Content-Range", "Accept-Ranges", "Content-Length"],
)
register_exception_handlers(app) register_exception_handlers(app)
app.include_router(health_router) app.include_router(health_router)
+4
View File
@@ -9,10 +9,12 @@ from arq.connections import RedisSettings
from app.core.config import get_settings from app.core.config import get_settings
from app.core.logging import configure_logging, get_logger from app.core.logging import configure_logging, get_logger
from app.workers.tasks.cleanup_task import cleanup_storage
from app.workers.tasks.download_task import download_track from app.workers.tasks.download_task import download_track
from app.workers.tasks.enrich_task import enrich_track from app.workers.tasks.enrich_task import enrich_track
from app.workers.tasks.import_task import scan_local_folder from app.workers.tasks.import_task import scan_local_folder
from app.workers.tasks.materialize_task import materialize_track from app.workers.tasks.materialize_task import materialize_track
from app.workers.tasks.transcode_task import transcode_track
log = get_logger("worker") log = get_logger("worker")
@@ -33,6 +35,8 @@ class WorkerSettings:
enrich_track, enrich_track,
download_track, download_track,
materialize_track, materialize_track,
cleanup_storage,
transcode_track,
] ]
on_startup = startup on_startup = startup
on_shutdown = shutdown on_shutdown = shutdown
+22
View File
@@ -59,6 +59,28 @@ async def enqueue_materialize(job_id: uuid.UUID) -> None:
log.warning("materialize_enqueue_failed", job_id=str(job_id)) log.warning("materialize_enqueue_failed", job_id=str(job_id))
async def enqueue_transcode(
track_id: uuid.UUID, *, quality: str = "high", hls: bool = True
) -> str:
"""Enqueue a transcode job (§6.6). Unlike the best-effort follow-ups below,
this is user/stream-driven, so it surfaces a job id (and a 503 via the caller
if the queue is down) rather than swallowing failures."""
return await enqueue(
"transcode_track", track_id=str(track_id), quality=quality, hls=hls
)
async def enqueue_transcode_quiet(
track_id: uuid.UUID, *, quality: str = "high", hls: bool = False
) -> None:
"""Best-effort transcode enqueue for a streaming cache miss — never blocks or
fails the play (the master is served meanwhile)."""
try:
await enqueue_transcode(track_id, quality=quality, hls=hls)
except DependencyUnavailableError:
log.warning("transcode_enqueue_failed", track_id=str(track_id))
async def enqueue_enrich(track_id: uuid.UUID) -> None: async def enqueue_enrich(track_id: uuid.UUID) -> None:
"""Best-effort enqueue of metadata enrichment for a freshly stored track. """Best-effort enqueue of metadata enrichment for a freshly stored track.
+40
View File
@@ -0,0 +1,40 @@
"""arq task: reconcile the catalogue against storage.
Scans every *local* track and drops rows whose backing file has vanished from
storage (a "dangling" reference). Remote placeholders have no local file yet and
are skipped by the repository query. The filesystem checks (one ``exists()`` per
track) are heavy, so this runs off the request cycle (CLAUDE.md). Guarded against
a storage outage: if an implausibly large share of files look missing, it assumes
the backend is down and aborts without deleting anything.
"""
from typing import Any
from app.core.logging import get_logger
from app.infrastructure.db import session_scope
from app.infrastructure.db.repositories import SqlAlchemyTrackRepository
from app.infrastructure.storage.provider import get_file_storage
log = get_logger("worker.cleanup")
# If more than this fraction of tracks look missing, assume the storage backend
# is unavailable (not that the library really evaporated) and refuse to delete.
_OUTAGE_GUARD_FRACTION = 0.5
async def cleanup_storage(_ctx: dict[str, Any]) -> dict[str, Any]:
storage = get_file_storage()
async with session_scope() as session:
tracks = SqlAlchemyTrackRepository(session)
refs = await tracks.all_storage_refs()
missing = [track_id for track_id, uri in refs if not await storage.exists(uri)]
if refs and len(missing) > len(refs) * _OUTAGE_GUARD_FRACTION:
log.warning("cleanup_aborted_outage_guard", scanned=len(refs), missing=len(missing))
return {"scanned": len(refs), "removed": 0, "aborted": True}
for track_id in missing:
await tracks.delete(track_id)
log.info("cleanup_done", scanned=len(refs), removed=len(missing))
return {"scanned": len(refs), "removed": len(missing), "aborted": False}
+64
View File
@@ -0,0 +1,64 @@
"""arq task: transcode a track to a cached Opus rendition and/or HLS (§6.6).
Triggered by ``POST /tracks/{id}/optimize`` and by a streaming cache miss. Reads
the master once (via ``storage.as_local_path``) and writes into
``transcode_cache_path``. Idempotent — existing outputs are skipped, so repeated
enqueues (e.g. several plays before the first finishes) are cheap. The DB session
is released before ffmpeg runs; the track entity is a detached value object.
"""
import uuid
from typing import Any
from app.application.transcode_service import (
HLS_BITRATE,
bitrate_for_quality,
hls_dir,
hls_playlist_path,
opus_path,
)
from app.core.config import get_settings
from app.core.logging import get_logger
from app.infrastructure.db import session_scope
from app.infrastructure.db.repositories import SqlAlchemyTrackRepository
from app.infrastructure.storage.provider import get_file_storage
from app.infrastructure.transcode.ffmpeg import FfmpegTranscoder
log = get_logger("worker.transcode")
async def transcode_track(
_ctx: dict[str, Any],
*,
track_id: str,
quality: str = "high",
hls: bool = True,
) -> dict[str, Any]:
settings = get_settings()
tid = uuid.UUID(track_id)
async with session_scope() as session:
track = await SqlAlchemyTrackRepository(session).get_by_id(tid)
if track is None or track.storage_uri is None:
log.warning("transcode_skip_no_file", track_id=track_id)
return {"track_id": track_id, "opus": False, "hls": False}
storage = get_file_storage()
transcoder = FfmpegTranscoder(settings.ffmpeg_path)
root = settings.transcode_cache_path
did_opus = False
did_hls = False
async with storage.as_local_path(track.storage_uri) as src:
bitrate = bitrate_for_quality(quality)
if bitrate is not None:
dest = opus_path(root, tid, bitrate)
if not dest.exists():
await transcoder.to_opus(src, dest, bitrate_kbps=bitrate)
did_opus = True
if hls and not hls_playlist_path(root, tid).exists():
await transcoder.to_hls(src, hls_dir(root, tid), bitrate_kbps=HLS_BITRATE)
did_hls = True
log.info("transcode_done", track_id=track_id, opus=did_opus, hls=did_hls)
return {"track_id": track_id, "opus": did_opus, "hls": did_hls}
+144
View File
@@ -0,0 +1,144 @@
"""Integration tests for the admin instance-management endpoints.
Covers ``/admin/services`` (dependency health), ``/admin/sources``,
``/admin/reindex`` (enqueue), and ``/admin/settings`` (effective config) plus
their admin gating. Requires a reachable Postgres; skips otherwise.
"""
import asyncio
import os
from collections.abc import AsyncIterator
from pathlib import Path
import pytest
from app.core.config import get_settings
from app.infrastructure.db import Base, dispose_engine, get_engine, session_scope
from app.infrastructure.db.repositories import (
SqlAlchemyRefreshTokenRepository,
SqlAlchemyUserRepository,
)
from asgi_lifespan import LifespanManager
from httpx import ASGITransport, AsyncClient
pytestmark = pytest.mark.asyncio
_db_reachable_cache: bool | None = None
async def _db_reachable() -> bool:
global _db_reachable_cache
if _db_reachable_cache is not None:
return _db_reachable_cache
from sqlalchemy import text
try:
async with asyncio.timeout(3):
async with get_engine().connect() as conn:
await conn.execute(text("SELECT 1"))
_db_reachable_cache = True
except Exception:
_db_reachable_cache = False
return _db_reachable_cache
@pytest.fixture
async def api(tmp_path: Path) -> AsyncIterator[AsyncClient]:
if not await _db_reachable():
pytest.skip("Postgres not reachable — integration test skipped.")
os.environ["MEDIA_PATH"] = str(tmp_path)
get_settings.cache_clear()
try:
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
from app.application.user_service import UserService
from app.core.security import Argon2PasswordHasher
async with session_scope() as session:
svc = UserService(
users=SqlAlchemyUserRepository(session),
refresh_tokens=SqlAlchemyRefreshTokenRepository(session),
hasher=Argon2PasswordHasher(),
)
await svc.create_user(username="user", password="testpass1", is_superuser=False)
await svc.create_user(username="admin", password="testpass1", is_superuser=True)
from app.main import create_app
app = create_app()
async with LifespanManager(app):
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
yield client
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await dispose_engine()
finally:
os.environ.pop("MEDIA_PATH", None)
get_settings.cache_clear()
async def _auth(api: AsyncClient, username: str) -> dict[str, str]:
resp = await api.post(
"/api/v1/auth/login", json={"username": username, "password": "testpass1"}
)
assert resp.status_code == 200, resp.text
return {"Authorization": f"Bearer {resp.json()['access_token']}"}
async def test_services_reports_dependency_health(api: AsyncClient) -> None:
headers = await _auth(api, "admin")
resp = await api.get("/api/v1/admin/services", headers=headers)
assert resp.status_code == 200, resp.text
body = resp.json()
assert body["database"] == "ok"
assert body["redis"] in ("ok", "down") # redis is up in CI/dev, but don't hard-require it
assert body["ml"] == "skipped" # no ML_SERVICE_URL configured
async def test_services_requires_admin(api: AsyncClient) -> None:
headers = await _auth(api, "user")
resp = await api.get("/api/v1/admin/services", headers=headers)
assert resp.status_code == 403
async def test_sources_list_is_admin_only(api: AsyncClient) -> None:
user = await _auth(api, "user")
assert (await api.get("/api/v1/admin/sources", headers=user)).status_code == 403
admin = await _auth(api, "admin")
resp = await api.get("/api/v1/admin/sources", headers=admin)
assert resp.status_code == 200, resp.text
assert isinstance(resp.json(), list)
async def test_reindex_without_indexable_source_is_503(api: AsyncClient) -> None:
# No LOCAL_MEDIA_IMPORT_PATH configured → nothing to index.
headers = await _auth(api, "admin")
resp = await api.post("/api/v1/admin/reindex", headers=headers)
assert resp.status_code == 503, resp.text
async def test_settings_exposes_effective_config_without_secrets(api: AsyncClient) -> None:
headers = await _auth(api, "admin")
resp = await api.get("/api/v1/admin/settings", headers=headers)
assert resp.status_code == 200, resp.text
body = resp.json()
# environment reflects however the process booted (test on host, dev in the
# container) — assert it's a valid value, not a specific one.
assert body["environment"] in ("dev", "test", "prod")
assert body["storage_backend"] == "local"
assert body["allow_registration"] is True
# No secret material should ever appear in the payload.
assert "jwt_secret" not in body
assert "subsonic_secret_key" not in body
async def test_settings_requires_admin(api: AsyncClient) -> None:
headers = await _auth(api, "user")
resp = await api.get("/api/v1/admin/settings", headers=headers)
assert resp.status_code == 403
+153
View File
@@ -0,0 +1,153 @@
"""Integration test for the playlist cover endpoint.
Master already covers playlist reorder in ``test_playlists_api.py``; this file
only exercises ``GET /playlists/{id}/cover`` (served when set, 404 otherwise).
Requires a reachable Postgres; skips otherwise.
"""
import asyncio
import os
import uuid
from collections.abc import AsyncIterator
from pathlib import Path
import pytest
from app.core.config import get_settings
from app.infrastructure.db import Base, dispose_engine, get_engine, session_scope
from app.infrastructure.db.repositories import (
SqlAlchemyPlaylistRepository,
SqlAlchemyRefreshTokenRepository,
SqlAlchemyUserRepository,
)
from app.infrastructure.storage.provider import get_file_storage
from asgi_lifespan import LifespanManager
from httpx import ASGITransport, AsyncClient
pytestmark = pytest.mark.asyncio
# A minimal valid 1x1 PNG.
_PNG_BYTES = bytes.fromhex(
"89504e470d0a1a0a0000000d4948445200000001000000010802000000907753"
"de0000000c4944415408d763f8cfc0f01f0005000155a2b4f60000000049454e44ae426082"
)
_db_reachable_cache: bool | None = None
async def _db_reachable() -> bool:
global _db_reachable_cache
if _db_reachable_cache is not None:
return _db_reachable_cache
from sqlalchemy import text
try:
async with asyncio.timeout(3):
async with get_engine().connect() as conn:
await conn.execute(text("SELECT 1"))
_db_reachable_cache = True
except Exception:
_db_reachable_cache = False
return _db_reachable_cache
async def _seed_playlist(*, owner_id: uuid.UUID, cover: Path | None) -> uuid.UUID:
"""Create a playlist owned by ``owner_id``; if ``cover`` is given, store it
and point the playlist's ``cover_path`` at it."""
from app.infrastructure.db.models.playlist import PlaylistModel
async with session_scope() as session:
playlist = await SqlAlchemyPlaylistRepository(session).add(
name="Mix", description=None, owner_id=owner_id
)
if cover is not None:
key = f"covers/playlists/{playlist.id}.png"
await get_file_storage().save_file(key, cover)
row = await session.get(PlaylistModel, playlist.id)
assert row is not None
row.cover_path = key
return playlist.id
@pytest.fixture
async def ctx(tmp_path: Path) -> AsyncIterator[tuple[AsyncClient, uuid.UUID, Path]]:
if not await _db_reachable():
pytest.skip("Postgres not reachable — integration test skipped.")
os.environ["MEDIA_PATH"] = str(tmp_path)
get_settings.cache_clear()
import app.infrastructure.storage.provider as _storage_provider
_storage_provider._storage = None
try:
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
from app.application.user_service import UserService
from app.core.security import Argon2PasswordHasher
async with session_scope() as session:
user = await UserService(
users=SqlAlchemyUserRepository(session),
refresh_tokens=SqlAlchemyRefreshTokenRepository(session),
hasher=Argon2PasswordHasher(),
).create_user(username="pluser", password="testpass1", is_superuser=False)
user_id = user.id
# A real source file for save_file (avoids the Windows NamedTemporaryFile
# re-open quirk).
src = tmp_path / "src_cover.png"
src.write_bytes(_PNG_BYTES)
from app.main import create_app
app = create_app()
async with LifespanManager(app):
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
yield client, user_id, src
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await dispose_engine()
finally:
_storage_provider._storage = None
os.environ.pop("MEDIA_PATH", None)
get_settings.cache_clear()
async def _token(api: AsyncClient) -> str:
resp = await api.post(
"/api/v1/auth/login", json={"username": "pluser", "password": "testpass1"}
)
assert resp.status_code == 200, resp.text
return str(resp.json()["access_token"])
async def test_playlist_cover_served(ctx: tuple[AsyncClient, uuid.UUID, Path]) -> None:
api, user_id, src = ctx
token = await _token(api)
playlist_id = await _seed_playlist(owner_id=user_id, cover=src)
resp = await api.get(f"/api/v1/playlists/{playlist_id}/cover?token={token}")
assert resp.status_code == 200, resp.text
assert resp.headers["content-type"] == "image/png"
assert resp.content == _PNG_BYTES
async def test_playlist_without_cover_is_404(ctx: tuple[AsyncClient, uuid.UUID, Path]) -> None:
api, user_id, _ = ctx
token = await _token(api)
playlist_id = await _seed_playlist(owner_id=user_id, cover=None)
resp = await api.get(f"/api/v1/playlists/{playlist_id}/cover?token={token}")
assert resp.status_code == 404
async def test_playlist_cover_requires_auth(ctx: tuple[AsyncClient, uuid.UUID, Path]) -> None:
api, user_id, src = ctx
playlist_id = await _seed_playlist(owner_id=user_id, cover=src)
resp = await api.get(f"/api/v1/playlists/{playlist_id}/cover")
assert resp.status_code == 401
+8
View File
@@ -59,6 +59,13 @@ async def api(tmp_path: Path) -> AsyncIterator[AsyncClient]:
os.environ["LOCAL_MEDIA_IMPORT_PATH"] = str(music) os.environ["LOCAL_MEDIA_IMPORT_PATH"] = str(music)
get_settings.cache_clear() get_settings.cache_clear()
# The source registry is process-cached (lru_cache); another test hitting a
# registry-backed endpoint (e.g. /admin/sources) may have built it before
# LOCAL_MEDIA_IMPORT_PATH was set. Clear it so this test sees a fresh one.
from app.api.deps import get_source_registry
get_source_registry.cache_clear()
import app.infrastructure.storage.provider as _storage_provider import app.infrastructure.storage.provider as _storage_provider
_storage_provider._storage = None _storage_provider._storage = None
@@ -94,6 +101,7 @@ async def api(tmp_path: Path) -> AsyncIterator[AsyncClient]:
os.environ.pop("MEDIA_PATH", None) os.environ.pop("MEDIA_PATH", None)
os.environ.pop("LOCAL_MEDIA_IMPORT_PATH", None) os.environ.pop("LOCAL_MEDIA_IMPORT_PATH", None)
get_settings.cache_clear() get_settings.cache_clear()
get_source_registry.cache_clear()
async def _login(api: AsyncClient) -> str: async def _login(api: AsyncClient) -> str:
+185
View File
@@ -0,0 +1,185 @@
"""Integration tests for the storage maintenance endpoints.
Covers ``/storage/duplicates`` (shared fingerprint), ``/broken`` (failed
enrichment), ``/missing-metadata`` (pending), and the admin-gated ``/cleanup``
enqueue. Requires a reachable Postgres; skips otherwise.
"""
import asyncio
import os
import uuid
from collections.abc import AsyncIterator
from pathlib import Path
import pytest
from app.core.config import get_settings
from app.infrastructure.db import Base, dispose_engine, get_engine, session_scope
from app.infrastructure.db.repositories import (
SqlAlchemyArtistRepository,
SqlAlchemyRefreshTokenRepository,
SqlAlchemyTrackRepository,
SqlAlchemyUserRepository,
)
from asgi_lifespan import LifespanManager
from httpx import ASGITransport, AsyncClient
pytestmark = pytest.mark.asyncio
_db_reachable_cache: bool | None = None
async def _db_reachable() -> bool:
global _db_reachable_cache
if _db_reachable_cache is not None:
return _db_reachable_cache
from sqlalchemy import text
try:
async with asyncio.timeout(3):
async with get_engine().connect() as conn:
await conn.execute(text("SELECT 1"))
_db_reachable_cache = True
except Exception:
_db_reachable_cache = False
return _db_reachable_cache
async def _seed_track(
*, title: str, source_id: str, metadata_status: str, fingerprint: str | None
) -> uuid.UUID:
async with session_scope() as session:
artist = await SqlAlchemyArtistRepository(session).get_or_create("Storage Artist")
tracks = SqlAlchemyTrackRepository(session)
tid = uuid.uuid4()
await tracks.add(
id=tid,
title=title,
artist_id=artist.id,
storage_uri=f"tracks/zz/{source_id}.mp3",
file_format="mp3",
file_size=10,
source="upload",
source_id=source_id,
metadata_status="pending",
added_by=None,
)
# apply_enrichment lets us set fingerprint + final status precisely.
await tracks.apply_enrichment(
tid,
title=title,
artist_id=artist.id,
album_id=None,
genre=None,
year=None,
track_number=None,
duration_seconds=1,
bitrate=None,
acoustid_fingerprint=fingerprint,
musicbrainz_id=None,
metadata_status=metadata_status,
)
return tid
@pytest.fixture
async def api(tmp_path: Path) -> AsyncIterator[AsyncClient]:
if not await _db_reachable():
pytest.skip("Postgres not reachable — integration test skipped.")
os.environ["MEDIA_PATH"] = str(tmp_path)
get_settings.cache_clear()
try:
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
from app.application.user_service import UserService
from app.core.security import Argon2PasswordHasher
async with session_scope() as session:
svc = UserService(
users=SqlAlchemyUserRepository(session),
refresh_tokens=SqlAlchemyRefreshTokenRepository(session),
hasher=Argon2PasswordHasher(),
)
await svc.create_user(username="user", password="testpass1", is_superuser=False)
await svc.create_user(username="admin", password="testpass1", is_superuser=True)
from app.main import create_app
app = create_app()
async with LifespanManager(app):
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
yield client
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await dispose_engine()
finally:
os.environ.pop("MEDIA_PATH", None)
get_settings.cache_clear()
async def _auth(api: AsyncClient, username: str) -> dict[str, str]:
resp = await api.post(
"/api/v1/auth/login", json={"username": username, "password": "testpass1"}
)
assert resp.status_code == 200, resp.text
return {"Authorization": f"Bearer {resp.json()['access_token']}"}
async def test_duplicates_grouped_by_fingerprint(api: AsyncClient) -> None:
headers = await _auth(api, "user")
await _seed_track(title="Dup A", source_id="a", metadata_status="enriched", fingerprint="FP1")
await _seed_track(title="Dup B", source_id="b", metadata_status="enriched", fingerprint="FP1")
# A lone fingerprint must NOT show up as a duplicate.
await _seed_track(title="Solo", source_id="c", metadata_status="enriched", fingerprint="FP2")
resp = await api.get("/api/v1/storage/duplicates", headers=headers)
assert resp.status_code == 200, resp.text
groups = resp.json()
assert len(groups) == 1
assert groups[0]["fingerprint"] == "FP1"
assert len(groups[0]["tracks"]) == 2
async def test_broken_lists_failed_tracks(api: AsyncClient) -> None:
headers = await _auth(api, "user")
await _seed_track(title="Bad", source_id="x", metadata_status="failed", fingerprint=None)
await _seed_track(title="Good", source_id="y", metadata_status="enriched", fingerprint=None)
resp = await api.get("/api/v1/storage/broken", headers=headers)
assert resp.status_code == 200, resp.text
body = resp.json()
assert body["total"] == 1
assert body["items"][0]["title"] == "Bad"
async def test_missing_metadata_lists_pending(api: AsyncClient) -> None:
headers = await _auth(api, "user")
await _seed_track(title="Pending", source_id="p", metadata_status="pending", fingerprint=None)
await _seed_track(title="Done", source_id="d", metadata_status="enriched", fingerprint=None)
resp = await api.get("/api/v1/storage/missing-metadata", headers=headers)
assert resp.status_code == 200, resp.text
body = resp.json()
assert body["total"] == 1
assert body["items"][0]["title"] == "Pending"
async def test_cleanup_requires_admin(api: AsyncClient) -> None:
user = await _auth(api, "user")
resp = await api.post("/api/v1/storage/cleanup", headers=user)
assert resp.status_code == 403
async def test_cleanup_enqueues_for_admin(api: AsyncClient) -> None:
admin = await _auth(api, "admin")
resp = await api.post("/api/v1/storage/cleanup", headers=admin)
# 202 with a job id when the queue is reachable; 503 if Redis is down.
assert resp.status_code in (202, 503), resp.text
if resp.status_code == 202:
assert resp.json()["status"] == "enqueued"
assert resp.json()["job_id"]
+252
View File
@@ -0,0 +1,252 @@
"""Integration tests for the offline-first sync endpoints.
Drives POST /sync/push (idempotent append of like/play events) and GET
/sync/changes (delta pull with a server-clock cursor). Requires a reachable
Postgres; skips otherwise.
"""
import asyncio
import os
import uuid
from collections.abc import AsyncIterator
from pathlib import Path
import pytest
from app.core.config import get_settings
from app.infrastructure.db import Base, dispose_engine, get_engine, session_scope
from app.infrastructure.db.repositories import (
SqlAlchemyArtistRepository,
SqlAlchemyRefreshTokenRepository,
SqlAlchemyTrackRepository,
SqlAlchemyUserRepository,
)
from asgi_lifespan import LifespanManager
from httpx import ASGITransport, AsyncClient
pytestmark = pytest.mark.asyncio
_db_reachable_cache: bool | None = None
async def _db_reachable() -> bool:
global _db_reachable_cache
if _db_reachable_cache is not None:
return _db_reachable_cache
from sqlalchemy import text
try:
async with asyncio.timeout(3):
async with get_engine().connect() as conn:
await conn.execute(text("SELECT 1"))
_db_reachable_cache = True
except Exception:
_db_reachable_cache = False
return _db_reachable_cache
async def _seed_track(source_id: str = "sync-1") -> uuid.UUID:
async with session_scope() as session:
artist = await SqlAlchemyArtistRepository(session).get_or_create("Sync Artist")
tid = uuid.uuid4()
await SqlAlchemyTrackRepository(session).add(
id=tid,
title="Sync Track",
artist_id=artist.id,
storage_uri=f"tracks/zz/{source_id}.mp3",
file_format="mp3",
file_size=10,
source="upload",
source_id=source_id,
metadata_status="enriched",
added_by=None,
)
return tid
@pytest.fixture
async def api(tmp_path: Path) -> AsyncIterator[AsyncClient]:
if not await _db_reachable():
pytest.skip("Postgres not reachable — integration test skipped.")
os.environ["MEDIA_PATH"] = str(tmp_path)
get_settings.cache_clear()
try:
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
from app.application.user_service import UserService
from app.core.security import Argon2PasswordHasher
async with session_scope() as session:
await UserService(
users=SqlAlchemyUserRepository(session),
refresh_tokens=SqlAlchemyRefreshTokenRepository(session),
hasher=Argon2PasswordHasher(),
).create_user(username="syncer", password="testpass1", is_superuser=False)
from app.main import create_app
app = create_app()
async with LifespanManager(app):
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
yield client
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await dispose_engine()
finally:
os.environ.pop("MEDIA_PATH", None)
get_settings.cache_clear()
async def _auth(api: AsyncClient) -> dict[str, str]:
resp = await api.post(
"/api/v1/auth/login", json={"username": "syncer", "password": "testpass1"}
)
assert resp.status_code == 200, resp.text
return {"Authorization": f"Bearer {resp.json()['access_token']}"}
async def test_push_is_idempotent_and_pull_returns_events(api: AsyncClient) -> None:
headers = await _auth(api)
track_id = await _seed_track()
like_id, play_id = str(uuid.uuid4()), str(uuid.uuid4())
push_body = {
"likes": [
{
"id": like_id,
"track_id": str(track_id),
"value": "like",
"created_at": "2020-01-01T00:00:00+00:00",
}
],
"plays": [
{
"id": play_id,
"track_id": str(track_id),
"played_at": "2020-01-01T00:05:00+00:00",
"play_duration_seconds": 180,
"completed": True,
}
],
}
first = await api.post("/api/v1/sync/push", json=push_body, headers=headers)
assert first.status_code == 200, first.text
assert first.json()["accepted_likes"] == 1
assert first.json()["accepted_plays"] == 1
# Replaying the exact same events must store nothing (idempotent by id).
replay = await api.post("/api/v1/sync/push", json=push_body, headers=headers)
assert replay.json()["accepted_likes"] == 0
assert replay.json()["accepted_plays"] == 0
# Full pull returns the events + the catalogue track, with a cursor.
changes = await api.get("/api/v1/sync/changes", headers=headers)
assert changes.status_code == 200, changes.text
body = changes.json()
assert [lk["id"] for lk in body["likes"]] == [like_id]
assert body["likes"][0]["created_at"].startswith("2020-01-01") # event time preserved
assert [p["id"] for p in body["plays"]] == [play_id]
assert str(track_id) in [t["id"] for t in body["tracks"]]
assert body["cursor"]
async def test_cursor_advances_and_empty_delta(api: AsyncClient) -> None:
headers = await _auth(api)
track_id = await _seed_track()
await api.post(
"/api/v1/sync/push",
json={
"likes": [
{
"id": str(uuid.uuid4()),
"track_id": str(track_id),
"value": "like",
"created_at": "2024-01-01T00:00:00+00:00",
}
]
},
headers=headers,
)
cursor = (await api.get("/api/v1/sync/changes", headers=headers)).json()["cursor"]
# Nothing changed since the cursor → empty delta (but a fresh cursor).
delta = await api.get("/api/v1/sync/changes", params={"since": cursor}, headers=headers)
assert delta.status_code == 200, delta.text
assert delta.json()["likes"] == []
assert delta.json()["plays"] == []
async def test_offline_event_with_old_time_still_syncs(api: AsyncClient) -> None:
# The delta is keyed by server-ingestion time, not the client event time, so
# an event pushed now with an old created_at still surfaces for a cursor that
# sits *after* that old event time.
headers = await _auth(api)
track_id = await _seed_track()
baseline = (await api.get("/api/v1/sync/changes", headers=headers)).json()["cursor"]
await api.post(
"/api/v1/sync/push",
json={
"likes": [
{
"id": str(uuid.uuid4()),
"track_id": str(track_id),
"value": "like",
"created_at": "2019-06-01T00:00:00+00:00", # long before `baseline`
}
]
},
headers=headers,
)
delta = await api.get("/api/v1/sync/changes", params={"since": baseline}, headers=headers)
assert len(delta.json()["likes"]) == 1 # surfaced despite the old event time
async def test_push_skips_unknown_track(api: AsyncClient) -> None:
headers = await _auth(api)
resp = await api.post(
"/api/v1/sync/push",
json={
"likes": [
{
"id": str(uuid.uuid4()),
"track_id": str(uuid.uuid4()), # not in the catalogue
"value": "like",
"created_at": "2024-01-01T00:00:00+00:00",
}
]
},
headers=headers,
)
assert resp.status_code == 200, resp.text
assert resp.json()["accepted_likes"] == 0
async def test_pull_includes_changed_playlists_with_track_ids(api: AsyncClient) -> None:
headers = await _auth(api)
track_id = await _seed_track()
created = await api.post("/api/v1/playlists", json={"name": "Sync Mix"}, headers=headers)
playlist_id = created.json()["id"]
await api.post(
f"/api/v1/playlists/{playlist_id}/tracks",
json={"track_id": str(track_id)},
headers=headers,
)
body = (await api.get("/api/v1/sync/changes", headers=headers)).json()
playlists = {p["id"]: p for p in body["playlists"]}
assert playlist_id in playlists
assert playlists[playlist_id]["track_ids"] == [str(track_id)]
async def test_sync_requires_auth(api: AsyncClient) -> None:
assert (await api.get("/api/v1/sync/changes")).status_code == 401
assert (await api.post("/api/v1/sync/push", json={})).status_code == 401
+174
View File
@@ -0,0 +1,174 @@
"""Integration tests for the user-settings + scrobbling endpoints.
Drives ``/api/v1/settings`` end to end: lazy defaults, partial update,
enum validation, and scrobbling config (write-only session key). Requires a
reachable Postgres; skips otherwise.
"""
import asyncio
import os
from collections.abc import AsyncIterator
from pathlib import Path
import pytest
from app.core.config import get_settings
from app.infrastructure.db import Base, dispose_engine, get_engine, session_scope
from app.infrastructure.db.repositories import (
SqlAlchemyRefreshTokenRepository,
SqlAlchemyUserRepository,
)
from asgi_lifespan import LifespanManager
from httpx import ASGITransport, AsyncClient
pytestmark = pytest.mark.asyncio
_db_reachable_cache: bool | None = None
async def _db_reachable() -> bool:
global _db_reachable_cache
if _db_reachable_cache is not None:
return _db_reachable_cache
from sqlalchemy import text
try:
async with asyncio.timeout(3):
async with get_engine().connect() as conn:
await conn.execute(text("SELECT 1"))
_db_reachable_cache = True
except Exception:
_db_reachable_cache = False
return _db_reachable_cache
@pytest.fixture
async def api(tmp_path: Path) -> AsyncIterator[AsyncClient]:
if not await _db_reachable():
pytest.skip("Postgres not reachable — integration test skipped.")
os.environ["MEDIA_PATH"] = str(tmp_path)
get_settings.cache_clear()
try:
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
from app.application.user_service import UserService
from app.core.security import Argon2PasswordHasher
async with session_scope() as session:
await UserService(
users=SqlAlchemyUserRepository(session),
refresh_tokens=SqlAlchemyRefreshTokenRepository(session),
hasher=Argon2PasswordHasher(),
).create_user(username="setuser", password="testpass1", is_superuser=False)
from app.main import create_app
app = create_app()
async with LifespanManager(app):
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
yield client
async with get_engine().begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await dispose_engine()
finally:
os.environ.pop("MEDIA_PATH", None)
get_settings.cache_clear()
async def _auth(api: AsyncClient) -> dict[str, str]:
resp = await api.post(
"/api/v1/auth/login", json={"username": "setuser", "password": "testpass1"}
)
assert resp.status_code == 200, resp.text
return {"Authorization": f"Bearer {resp.json()['access_token']}"}
async def test_get_settings_returns_defaults(api: AsyncClient) -> None:
headers = await _auth(api)
resp = await api.get("/api/v1/settings", headers=headers)
assert resp.status_code == 200, resp.text
assert resp.json() == {"theme": "system", "stream_quality": "original"}
async def test_patch_settings_persists(api: AsyncClient) -> None:
headers = await _auth(api)
resp = await api.patch("/api/v1/settings", json={"theme": "dark"}, headers=headers)
assert resp.status_code == 200, resp.text
assert resp.json() == {"theme": "dark", "stream_quality": "original"}
# Persisted + partial update leaves the untouched field alone.
again = await api.patch("/api/v1/settings", json={"stream_quality": "low"}, headers=headers)
assert again.json() == {"theme": "dark", "stream_quality": "low"}
async def test_invalid_theme_is_422(api: AsyncClient) -> None:
headers = await _auth(api)
resp = await api.patch("/api/v1/settings", json={"theme": "neon"}, headers=headers)
assert resp.status_code == 422
async def test_scrobbling_defaults_and_enable(api: AsyncClient) -> None:
headers = await _auth(api)
resp = await api.get("/api/v1/settings/scrobbling", headers=headers)
assert resp.status_code == 200, resp.text
assert resp.json() == {
"enabled": False,
"provider": None,
"username": None,
"configured": False,
}
# Enabling without a provider is rejected.
bad = await api.put("/api/v1/settings/scrobbling", json={"enabled": True}, headers=headers)
assert bad.status_code == 422
# Enable with a provider + secret token; the token is never echoed back.
ok = await api.put(
"/api/v1/settings/scrobbling",
json={
"enabled": True,
"provider": "listenbrainz",
"username": "me",
"session_key": "super-secret-token",
},
headers=headers,
)
assert ok.status_code == 200, ok.text
body = ok.json()
assert body == {
"enabled": True,
"provider": "listenbrainz",
"username": "me",
"configured": True,
}
assert "session_key" not in body
assert "super-secret-token" not in ok.text
async def test_scrobbling_keeps_stored_key_when_omitted(api: AsyncClient) -> None:
headers = await _auth(api)
await api.put(
"/api/v1/settings/scrobbling",
json={"enabled": True, "provider": "lastfm", "session_key": "k"},
headers=headers,
)
# A later update without a session_key must keep the stored one.
resp = await api.put(
"/api/v1/settings/scrobbling",
json={"enabled": True, "provider": "lastfm", "username": "changed"},
headers=headers,
)
assert resp.status_code == 200, resp.text
assert resp.json()["configured"] is True
assert resp.json()["username"] == "changed"
async def test_settings_require_auth(api: AsyncClient) -> None:
resp = await api.get("/api/v1/settings")
assert resp.status_code == 401