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>
119 lines
4.0 KiB
Python
119 lines
4.0 KiB
Python
"""Play history repository — adapter over ``AsyncSession``."""
|
|
|
|
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
|
|
from app.infrastructure.db.models.play_history import PlayHistoryModel
|
|
|
|
|
|
def _to_entity(row: PlayHistoryModel) -> PlayHistoryEntry:
|
|
return PlayHistoryEntry(
|
|
id=row.id,
|
|
user_id=row.user_id,
|
|
track_id=row.track_id,
|
|
played_at=row.played_at,
|
|
play_duration_seconds=row.play_duration_seconds,
|
|
completed=row.completed,
|
|
)
|
|
|
|
|
|
class SqlAlchemyHistoryRepository:
|
|
def __init__(self, session: AsyncSession) -> None:
|
|
self._session = session
|
|
|
|
async def add(
|
|
self,
|
|
*,
|
|
user_id: uuid.UUID,
|
|
track_id: uuid.UUID,
|
|
played_at: dt.datetime,
|
|
play_duration_seconds: int | None,
|
|
completed: bool,
|
|
) -> PlayHistoryEntry:
|
|
row = PlayHistoryModel(
|
|
user_id=user_id,
|
|
track_id=track_id,
|
|
played_at=played_at,
|
|
play_duration_seconds=play_duration_seconds,
|
|
completed=completed,
|
|
)
|
|
self._session.add(row)
|
|
await self._session.flush()
|
|
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 = (
|
|
(
|
|
await self._session.execute(
|
|
select(PlayHistoryModel)
|
|
.where(PlayHistoryModel.user_id == user_id)
|
|
.order_by(PlayHistoryModel.played_at.desc())
|
|
.limit(limit)
|
|
.offset(offset)
|
|
)
|
|
)
|
|
.scalars()
|
|
.all()
|
|
)
|
|
return [_to_entity(r) for r in rows]
|
|
|
|
async def count(self, *, user_id: uuid.UUID) -> int:
|
|
return (
|
|
await self._session.execute(
|
|
select(func.count())
|
|
.select_from(PlayHistoryModel)
|
|
.where(PlayHistoryModel.user_id == user_id)
|
|
)
|
|
).scalar_one()
|