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>
This commit is contained in:
Цвылев Александр Вадимович
2026-07-28 15:04:13 +03:00
parent c47242aa3a
commit 048854f92a
15 changed files with 836 additions and 5 deletions
+21
View File
@@ -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)
+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
+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 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,
)
+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
)
+32
View File
@@ -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]: ...
+9
View File
@@ -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,
)
@@ -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,
)
@@ -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 = (
(
@@ -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]:
@@ -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 = (
@@ -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,
*,