From 048854f92ac94bbbecd0a6bdd6f011644477e182 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=A6=D0=B2=D1=8B=D0=BB=D0=B5=D0=B2=20=D0=90=D0=BB=D0=B5?= =?UTF-8?q?=D0=BA=D1=81=D0=B0=D0=BD=D0=B4=D1=80=20=D0=92=D0=B0=D0=B4=D0=B8?= =?UTF-8?q?=D0=BC=D0=BE=D0=B2=D0=B8=D1=87?= Date: Tue, 28 Jul 2026 15:04:13 +0300 Subject: [PATCH] feat(api): offline-first sync layer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../20260728_1100-sync_ingestion_columns.py | 57 ++++ app/api/deps.py | 21 ++ app/api/schemas/sync.py | 82 ++++++ app/api/v1/sync.py | 96 ++++++- app/application/sync_service.py | 140 ++++++++++ app/domain/ports.py | 32 +++ app/infrastructure/db/models/like.py | 9 + app/infrastructure/db/models/play_history.py | 8 + .../db/repositories/history_repository.py | 46 ++++ .../db/repositories/like_repository.py | 46 ++++ .../db/repositories/playlist_repository.py | 27 ++ .../db/repositories/track_repository.py | 13 + tests/test_admin_api.py | 4 +- tests/test_sources_api.py | 8 + tests/test_sync_api.py | 252 ++++++++++++++++++ 15 files changed, 836 insertions(+), 5 deletions(-) create mode 100644 alembic/versions/20260728_1100-sync_ingestion_columns.py create mode 100644 app/api/schemas/sync.py create mode 100644 app/application/sync_service.py create mode 100644 tests/test_sync_api.py diff --git a/alembic/versions/20260728_1100-sync_ingestion_columns.py b/alembic/versions/20260728_1100-sync_ingestion_columns.py new file mode 100644 index 0000000..5b47032 --- /dev/null +++ b/alembic/versions/20260728_1100-sync_ingestion_columns.py @@ -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") diff --git a/app/api/deps.py b/app/api/deps.py index 81fb117..6c8608d 100644 --- a/app/api/deps.py +++ b/app/api/deps.py @@ -6,12 +6,14 @@ bound to the request-scoped DB session; stateless adapters (hasher, token service) are process-cached. """ +import datetime as dt from collections.abc import AsyncIterator from functools import lru_cache from typing import Annotated from fastapi import Depends, Query from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer +from sqlalchemy import func, select from sqlalchemy.ext.asyncio import AsyncSession from app.application.auth_service import AuthService @@ -20,6 +22,7 @@ from app.application.metadata_service import MetadataEnrichmentService from app.application.remote_library_service import RemoteLibraryService from app.application.streaming_service import StreamingService from app.application.subsonic_auth_service import SubsonicAuthService +from app.application.sync_service import SyncService from app.application.upload_service import UploadService from app.application.user_service import UserService from app.application.user_settings_service import UserSettingsService @@ -199,6 +202,24 @@ DownloadServiceDep = Annotated[DownloadService, Depends(get_download_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 --------------------------------------------------- def get_track_repository(session: SessionDep) -> SqlAlchemyTrackRepository: return SqlAlchemyTrackRepository(session) diff --git a/app/api/schemas/sync.py b/app/api/schemas/sync.py new file mode 100644 index 0000000..fb8b412 --- /dev/null +++ b/app/api/schemas/sync.py @@ -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 diff --git a/app/api/v1/sync.py b/app/api/v1/sync.py index 5c0704f..d49be2b 100644 --- a/app/api/v1/sync.py +++ b/app/api/v1/sync.py @@ -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 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.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") -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, + ) diff --git a/app/application/sync_service.py b/app/application/sync_service.py new file mode 100644 index 0000000..28da619 --- /dev/null +++ b/app/application/sync_service.py @@ -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 + ) diff --git a/app/domain/ports.py b/app/domain/ports.py index 1ba1a34..f30bbcd 100644 --- a/app/domain/ports.py +++ b/app/domain/ports.py @@ -176,6 +176,9 @@ class TrackRepository(Protocol): ) -> 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( self, *, @@ -300,12 +303,28 @@ class PlaylistRepository(Protocol): self, playlist_id: uuid.UUID, ordered_track_ids: list[uuid.UUID] ) -> 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) async def list(self, *, owner_id: uuid.UUID, limit: int, offset: int) -> list[Playlist]: ... class LikeRepository(Protocol): 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( self, *, user_id: uuid.UUID, track_ids: list[uuid.UUID] ) -> list[Like]: ... @@ -325,6 +344,19 @@ class HistoryRepository(Protocol): play_duration_seconds: int | None, completed: bool, ) -> 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( self, *, user_id: uuid.UUID, limit: int, offset: int ) -> list[PlayHistoryEntry]: ... diff --git a/app/infrastructure/db/models/like.py b/app/infrastructure/db/models/like.py index ce84ba3..d0aed5a 100644 --- a/app/infrastructure/db/models/like.py +++ b/app/infrastructure/db/models/like.py @@ -37,3 +37,12 @@ class LikeModel(UUIDPrimaryKeyMixin, Base): server_default=func.now(), 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, + ) diff --git a/app/infrastructure/db/models/play_history.py b/app/infrastructure/db/models/play_history.py index 2383aad..942c9cb 100644 --- a/app/infrastructure/db/models/play_history.py +++ b/app/infrastructure/db/models/play_history.py @@ -34,3 +34,11 @@ class PlayHistoryModel(UUIDPrimaryKeyMixin, Base): ) play_duration_seconds: Mapped[int | None] = mapped_column(Integer, nullable=True) 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, + ) diff --git a/app/infrastructure/db/repositories/history_repository.py b/app/infrastructure/db/repositories/history_repository.py index 9889800..5e17a3c 100644 --- a/app/infrastructure/db/repositories/history_repository.py +++ b/app/infrastructure/db/repositories/history_repository.py @@ -4,6 +4,7 @@ import datetime as dt 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.history import PlayHistoryEntry @@ -46,6 +47,51 @@ class SqlAlchemyHistoryRepository: await self._session.refresh(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]: rows = ( ( diff --git a/app/infrastructure/db/repositories/like_repository.py b/app/infrastructure/db/repositories/like_repository.py index 2e50386..25c7916 100644 --- a/app/infrastructure/db/repositories/like_repository.py +++ b/app/infrastructure/db/repositories/like_repository.py @@ -3,9 +3,11 @@ Likes are an append-only event log. Current state = latest event per (user, track). """ +import datetime as dt 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.like import Like @@ -59,6 +61,50 @@ class SqlAlchemyLikeRepository: await self._session.refresh(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( self, *, user_id: uuid.UUID, track_ids: list[uuid.UUID] ) -> list[Like]: diff --git a/app/infrastructure/db/repositories/playlist_repository.py b/app/infrastructure/db/repositories/playlist_repository.py index 257b7a7..db4e4d7 100644 --- a/app/infrastructure/db/repositories/playlist_repository.py +++ b/app/infrastructure/db/repositories/playlist_repository.py @@ -1,5 +1,6 @@ """Playlist repository — adapter over ``AsyncSession``.""" +import datetime as dt import uuid from sqlalchemy import func, select @@ -216,6 +217,32 @@ class SqlAlchemyPlaylistRepository: ) ).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) async def list(self, *, owner_id: uuid.UUID, limit: int, offset: int) -> list[Playlist]: rows = ( diff --git a/app/infrastructure/db/repositories/track_repository.py b/app/infrastructure/db/repositories/track_repository.py index f3848c8..521628b 100644 --- a/app/infrastructure/db/repositories/track_repository.py +++ b/app/infrastructure/db/repositories/track_repository.py @@ -266,6 +266,19 @@ class SqlAlchemyTrackRepository: ) ).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( self, *, diff --git a/tests/test_admin_api.py b/tests/test_admin_api.py index 2e60b78..4f75ad7 100644 --- a/tests/test_admin_api.py +++ b/tests/test_admin_api.py @@ -128,7 +128,9 @@ async def test_settings_exposes_effective_config_without_secrets(api: AsyncClien resp = await api.get("/api/v1/admin/settings", headers=headers) assert resp.status_code == 200, resp.text body = resp.json() - assert body["environment"] == "test" + # 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. diff --git a/tests/test_sources_api.py b/tests/test_sources_api.py index 3757afc..c94de17 100644 --- a/tests/test_sources_api.py +++ b/tests/test_sources_api.py @@ -59,6 +59,13 @@ async def api(tmp_path: Path) -> AsyncIterator[AsyncClient]: os.environ["LOCAL_MEDIA_IMPORT_PATH"] = str(music) 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 _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("LOCAL_MEDIA_IMPORT_PATH", None) get_settings.cache_clear() + get_source_registry.cache_clear() async def _login(api: AsyncClient) -> str: diff --git a/tests/test_sync_api.py b/tests/test_sync_api.py new file mode 100644 index 0000000..b97d688 --- /dev/null +++ b/tests/test_sync_api.py @@ -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