048854f92a
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>
141 lines
4.3 KiB
Python
141 lines
4.3 KiB
Python
"""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
|
|
)
|