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>
This commit is contained in:
Цвылев Александр Вадимович
2026-07-28 14:33:12 +03:00
parent df580578f6
commit a263272935
27 changed files with 1362 additions and 27 deletions
+10
View File
@@ -22,6 +22,7 @@ from app.application.streaming_service import StreamingService
from app.application.subsonic_auth_service import SubsonicAuthService
from app.application.upload_service import UploadService
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.security import Argon2PasswordHasher, JwtTokenService, SubsonicPasswordCipher
from app.domain.entities import User
@@ -38,6 +39,7 @@ from app.infrastructure.db.repositories import (
SqlAlchemyRefreshTokenRepository,
SqlAlchemyTrackRepository,
SqlAlchemyUserRepository,
SqlAlchemyUserSettingsRepository,
)
from app.infrastructure.metadata.acoustid import AcoustIdHttpClient
from app.infrastructure.metadata.fingerprint import FpcalcFingerprinter
@@ -112,9 +114,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)]
UserServiceDep = Annotated[UserService, Depends(get_user_service)]
SubsonicAuthServiceDep = Annotated[SubsonicAuthService, Depends(get_subsonic_auth_service)]
UserSettingsServiceDep = Annotated[UserSettingsService, Depends(get_user_settings_service)]
# -- file storage (process-cached) ---------------------------------------------
+2
View File
@@ -18,6 +18,7 @@ from app.domain.errors import (
DependencyUnavailableError,
DomainError,
NotFoundError,
NotSupportedError,
PermissionDeniedError,
RangeNotSatisfiableError,
StorageError,
@@ -33,6 +34,7 @@ _STATUS_BY_ERROR: dict[type[DomainError], int] = {
ValidationError: status.HTTP_422_UNPROCESSABLE_CONTENT,
AuthenticationError: status.HTTP_401_UNAUTHORIZED,
PermissionDeniedError: status.HTTP_403_FORBIDDEN,
NotSupportedError: status.HTTP_501_NOT_IMPLEMENTED,
DependencyUnavailableError: status.HTTP_503_SERVICE_UNAVAILABLE,
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
+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 app.api.schemas.track import TrackOut
class DiskUsageOut(BaseModel):
total: int
@@ -43,3 +45,17 @@ class StorageStatsOut(BaseModel):
# backing volume (``None`` for object-store backends)
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
+71 -11
View File
@@ -5,11 +5,18 @@ sign-up (plan §6.4).
"""
import uuid
from typing import Any
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.user import (
CreateUserRequest,
@@ -17,6 +24,9 @@ from app.api.schemas.user import (
UpdateUserRequest,
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"])
@@ -91,24 +101,74 @@ async def rotate_user_subsonic_password(
@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")
async def list_admin_sources(_admin: SuperUser) -> Any: ...
@router.patch("/sources/{source}")
async def update_admin_source(source: str, _admin: SuperUser) -> Any: ...
async def list_admin_sources(
_admin: SuperUser, registry: SourceRegistryDep
) -> list[SourceInfoOut]:
"""Configured sources and their live availability (same view as
``/sources``, admin-scoped)."""
return [SourceInfoOut.from_entity(info) for info in registry.infos()]
@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")
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")
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."
)
+15 -2
View File
@@ -1,15 +1,18 @@
"""Playlist endpoints."""
import uuid
from typing import Any
from fastapi import APIRouter, Query, Response
from fastapi.responses import StreamingResponse
from app.api.covers import stream_cover
from app.api.deps import (
AlbumRepoDep,
ArtistRepoDep,
CurrentUser,
FileStorageDep,
PlaylistRepoDep,
StreamUser,
TrackRepoDep,
)
from app.api.schemas.pagination import PagedResponse
@@ -217,4 +220,14 @@ async def reorder_playlist_tracks(
@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)
+73 -8
View File
@@ -1,22 +1,28 @@
"""Storage analysis and cleanup endpoints."""
from typing import Any
from fastapi import APIRouter
from fastapi import APIRouter, Query
from app.api.deps import (
AlbumRepoDep,
ArtistRepoDep,
CurrentUser,
FileStorageDep,
SuperUser,
TrackRepoDep,
)
from app.api.schemas.pagination import PagedResponse
from app.api.schemas.storage import (
CleanupEnqueuedOut,
DiskUsageOut,
DuplicateGroupOut,
FormatBreakdownOut,
GenreCountOut,
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"])
@@ -24,6 +30,18 @@ router = APIRouter(prefix="/storage", tags=["storage"])
_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("")
async def get_storage_stats(
track_repo: TrackRepoDep,
@@ -70,16 +88,63 @@ async def get_storage_stats(
@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")
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")
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")
async def run_cleanup() -> Any: ...
@router.post("/cleanup", status_code=202)
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)
+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 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"])
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("")
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("")
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")
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")
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)
+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)
+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,
)
+7
View File
@@ -54,6 +54,13 @@ class PermissionDeniedError(DomainError):
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):
"""An external dependency (source, ML, MusicBrainz) is unavailable.
+13
View File
@@ -29,6 +29,7 @@ from app.domain.entities import (
SubsonicCredentials,
User,
)
from app.domain.entities.settings import UserSettings
from app.domain.entities.track import Artist, Track
from app.domain.sources import DownloadResult, RawMetadata, SearchResult, SourceFile, SourceInfo
from app.domain.tokens import IssuedToken, TokenClaims, TokenType
@@ -56,6 +57,11 @@ class UserRepository(Protocol):
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):
"""Symmetric encrypt/decrypt for the recoverable Subsonic app-password."""
@@ -164,6 +170,12 @@ class TrackRepository(Protocol):
# AlbumRepository below).
async def genres(self) -> list[tuple[str, int]]: ...
async def library_stats(self) -> LibraryStats: ...
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(
self,
*,
@@ -287,6 +299,7 @@ class PlaylistRepository(Protocol):
async def reorder_tracks(
self, playlist_id: uuid.UUID, ordered_track_ids: list[uuid.UUID]
) -> None: ...
async def get_cover_path(self, playlist_id: uuid.UUID) -> str | None: ...
# 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]: ...
+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.track import TrackModel
from app.infrastructure.db.models.user import RefreshTokenModel, UserModel
from app.infrastructure.db.models.user_settings import UserSettingsModel
__all__ = [
"AlbumModel",
@@ -27,4 +28,5 @@ __all__ = [
"RefreshTokenModel",
"TrackModel",
"UserModel",
"UserSettingsModel",
]
@@ -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)
@@ -13,6 +13,9 @@ from app.infrastructure.db.repositories.refresh_token_repository import (
)
from app.infrastructure.db.repositories.track_repository import SqlAlchemyTrackRepository
from app.infrastructure.db.repositories.user_repository import SqlAlchemyUserRepository
from app.infrastructure.db.repositories.user_settings_repository import (
SqlAlchemyUserSettingsRepository,
)
__all__ = [
"SqlAlchemyAlbumRepository",
@@ -24,4 +27,5 @@ __all__ = [
"SqlAlchemyRefreshTokenRepository",
"SqlAlchemyTrackRepository",
"SqlAlchemyUserRepository",
"SqlAlchemyUserSettingsRepository",
]
@@ -207,6 +207,15 @@ class SqlAlchemyPlaylistRepository:
).scalar_one_or_none()
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()
# 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]:
rows = (
@@ -194,6 +194,78 @@ class SqlAlchemyTrackRepository:
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(
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)
+4
View File
@@ -48,6 +48,10 @@ class SourceRegistry:
"""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")]
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]:
return [backend.info() for backend in self._by_name.values()]
+2
View File
@@ -9,6 +9,7 @@ from arq.connections import RedisSettings
from app.core.config import get_settings
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.enrich_task import enrich_track
from app.workers.tasks.import_task import scan_local_folder
@@ -33,6 +34,7 @@ class WorkerSettings:
enrich_track,
download_track,
materialize_track,
cleanup_storage,
]
on_startup = startup
on_shutdown = shutdown
+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}