Implement incremental sync and offline recovery
This commit is contained in:
@@ -395,18 +395,7 @@
|
||||
"/api/v1/sync/bootstrap": {
|
||||
"get": {
|
||||
"operationId": "SyncController_bootstrap_v1",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "deviceId",
|
||||
"required": false,
|
||||
"in": "query",
|
||||
"description": "Optional client metadata. Authorization: Bearer <deviceAccessToken> is required and determines access.",
|
||||
"schema": {
|
||||
"format": "uuid",
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
],
|
||||
"parameters": [],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "",
|
||||
@@ -434,23 +423,24 @@
|
||||
"operationId": "SyncController_changes_v1",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "deviceId",
|
||||
"required": false,
|
||||
"in": "query",
|
||||
"description": "Optional client metadata. Authorization: Bearer <deviceAccessToken> is required and determines access.",
|
||||
"schema": {
|
||||
"format": "uuid",
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "after",
|
||||
"name": "cursor",
|
||||
"required": false,
|
||||
"in": "query",
|
||||
"schema": {
|
||||
"example": "0",
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "limit",
|
||||
"required": false,
|
||||
"in": "query",
|
||||
"schema": {
|
||||
"minimum": 1,
|
||||
"maximum": 500,
|
||||
"example": 100,
|
||||
"type": "number"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
@@ -918,26 +908,37 @@
|
||||
"tracks"
|
||||
]
|
||||
},
|
||||
"LibraryTrackDto": {
|
||||
"SyncBootstrapResponseDto": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"id": {
|
||||
"nextCursor": {
|
||||
"type": "string",
|
||||
"format": "uuid"
|
||||
"example": "7"
|
||||
},
|
||||
"title": {
|
||||
"type": "string",
|
||||
"example": "Placeholder Track"
|
||||
"tracks": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/RemoteLibraryTrackDto"
|
||||
}
|
||||
},
|
||||
"artist": {
|
||||
"serverTime": {
|
||||
"type": "string",
|
||||
"example": "Velody"
|
||||
"example": "2026-06-15T12:00:00.000Z"
|
||||
}
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
"nextCursor",
|
||||
"tracks",
|
||||
"serverTime"
|
||||
]
|
||||
},
|
||||
"SyncEventDto": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"cursor": {
|
||||
"type": "string",
|
||||
"example": "3"
|
||||
},
|
||||
"entityType": {
|
||||
"type": "string",
|
||||
"example": "TRACK"
|
||||
@@ -948,56 +949,33 @@
|
||||
},
|
||||
"action": {
|
||||
"type": "string",
|
||||
"example": "CREATED"
|
||||
"example": "UPDATED"
|
||||
},
|
||||
"eventId": {
|
||||
"track": {
|
||||
"nullable": true,
|
||||
"type": "object",
|
||||
"allOf": [
|
||||
{
|
||||
"$ref": "#/components/schemas/RemoteLibraryTrackDto"
|
||||
}
|
||||
]
|
||||
},
|
||||
"deletedTrackId": {
|
||||
"type": "object",
|
||||
"format": "uuid",
|
||||
"nullable": true
|
||||
},
|
||||
"createdAt": {
|
||||
"type": "string",
|
||||
"example": "0"
|
||||
"example": "2026-06-15T12:00:00.000Z"
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
"cursor",
|
||||
"entityType",
|
||||
"entityId",
|
||||
"action",
|
||||
"eventId"
|
||||
]
|
||||
},
|
||||
"SyncBootstrapResponseDto": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"nextCursor": {
|
||||
"type": "string",
|
||||
"example": "0"
|
||||
},
|
||||
"tracks": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/LibraryTrackDto"
|
||||
}
|
||||
},
|
||||
"events": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/SyncEventDto"
|
||||
}
|
||||
},
|
||||
"deletedTrackIds": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
"serverTime": {
|
||||
"type": "string",
|
||||
"example": "2026-05-24T20:00:00.000Z"
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
"nextCursor",
|
||||
"tracks",
|
||||
"events",
|
||||
"deletedTrackIds",
|
||||
"serverTime"
|
||||
"createdAt"
|
||||
]
|
||||
},
|
||||
"SyncChangesResponseDto": {
|
||||
@@ -1005,13 +983,20 @@
|
||||
"properties": {
|
||||
"nextCursor": {
|
||||
"type": "string",
|
||||
"example": "0"
|
||||
"example": "7"
|
||||
},
|
||||
"tracks": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/LibraryTrackDto"
|
||||
}
|
||||
"hasMore": {
|
||||
"type": "boolean",
|
||||
"example": false
|
||||
},
|
||||
"requiresBootstrap": {
|
||||
"type": "boolean",
|
||||
"example": false
|
||||
},
|
||||
"reason": {
|
||||
"type": "object",
|
||||
"nullable": true,
|
||||
"example": "cursor_too_old"
|
||||
},
|
||||
"events": {
|
||||
"type": "array",
|
||||
@@ -1019,22 +1004,16 @@
|
||||
"$ref": "#/components/schemas/SyncEventDto"
|
||||
}
|
||||
},
|
||||
"deletedTrackIds": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
"serverTime": {
|
||||
"type": "string",
|
||||
"example": "2026-05-24T20:00:00.000Z"
|
||||
"example": "2026-06-15T12:00:00.000Z"
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
"nextCursor",
|
||||
"tracks",
|
||||
"hasMore",
|
||||
"requiresBootstrap",
|
||||
"events",
|
||||
"deletedTrackIds",
|
||||
"serverTime"
|
||||
]
|
||||
}
|
||||
|
||||
+128
@@ -0,0 +1,128 @@
|
||||
ALTER TABLE "users"
|
||||
ADD COLUMN "library_cursor" BIGINT NOT NULL DEFAULT 0;
|
||||
|
||||
ALTER TABLE "library_events"
|
||||
ADD COLUMN "cursor" BIGINT,
|
||||
ADD COLUMN "payload" JSONB;
|
||||
|
||||
WITH ranked_events AS (
|
||||
SELECT
|
||||
"id",
|
||||
ROW_NUMBER() OVER (
|
||||
PARTITION BY "user_id"
|
||||
ORDER BY "created_at" ASC, "id" ASC
|
||||
)::BIGINT AS next_cursor
|
||||
FROM "library_events"
|
||||
)
|
||||
UPDATE "library_events" AS "event"
|
||||
SET "cursor" = ranked_events.next_cursor
|
||||
FROM ranked_events
|
||||
WHERE ranked_events."id" = "event"."id";
|
||||
|
||||
UPDATE "library_events" AS "event"
|
||||
SET "payload" = CASE
|
||||
WHEN track."status" = 'DELETED'::"TrackStatus" THEN jsonb_build_object(
|
||||
'deletedTrackId',
|
||||
track."id"
|
||||
)
|
||||
WHEN audio_asset."id" IS NOT NULL THEN jsonb_build_object(
|
||||
'track',
|
||||
jsonb_build_object(
|
||||
'trackId',
|
||||
track."id",
|
||||
'title',
|
||||
track."title",
|
||||
'artist',
|
||||
track."artist",
|
||||
'durationSeconds',
|
||||
GREATEST(
|
||||
0,
|
||||
ROUND(COALESCE(track."duration_ms", audio_asset."duration_ms", 0)::numeric / 1000.0)
|
||||
)::INTEGER,
|
||||
'sha256',
|
||||
audio_asset."sha256",
|
||||
'assetId',
|
||||
audio_asset."id",
|
||||
'createdAt',
|
||||
to_char(track."created_at" AT TIME ZONE 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.MS"Z"'),
|
||||
'updatedAt',
|
||||
to_char(track."updated_at" AT TIME ZONE 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.MS"Z"'),
|
||||
'artwork',
|
||||
CASE
|
||||
WHEN artwork_asset."id" IS NULL THEN NULL
|
||||
ELSE jsonb_build_object(
|
||||
'artworkId',
|
||||
artwork_asset."id",
|
||||
'sha256',
|
||||
artwork_asset."sha256",
|
||||
'mimeType',
|
||||
artwork_asset."mime_type",
|
||||
'width',
|
||||
artwork_asset."width",
|
||||
'height',
|
||||
artwork_asset."height"
|
||||
)
|
||||
END
|
||||
)
|
||||
)
|
||||
ELSE '{}'::jsonb
|
||||
END
|
||||
FROM "tracks" AS track
|
||||
LEFT JOIN "audio_assets" AS audio_asset
|
||||
ON audio_asset."id" = track."primary_audio_asset_id"
|
||||
LEFT JOIN "artwork_assets" AS artwork_asset
|
||||
ON artwork_asset."id" = track."artwork_asset_id"
|
||||
WHERE "event"."entity_type" = 'TRACK'::"EntityType"
|
||||
AND "event"."entity_id" = track."id";
|
||||
|
||||
UPDATE "library_events"
|
||||
SET "payload" = '{}'::jsonb
|
||||
WHERE "payload" IS NULL;
|
||||
|
||||
ALTER TABLE "library_events"
|
||||
ALTER COLUMN "cursor" SET NOT NULL,
|
||||
ALTER COLUMN "payload" SET NOT NULL;
|
||||
|
||||
CREATE UNIQUE INDEX "library_events_user_id_cursor_key"
|
||||
ON "library_events"("user_id", "cursor");
|
||||
|
||||
CREATE INDEX "library_events_user_id_cursor_idx"
|
||||
ON "library_events"("user_id", "cursor");
|
||||
|
||||
UPDATE "users" AS "user"
|
||||
SET "library_cursor" = COALESCE(cursor_summary."max_cursor", 0)
|
||||
FROM (
|
||||
SELECT
|
||||
"user_id",
|
||||
MAX("cursor") AS "max_cursor"
|
||||
FROM "library_events"
|
||||
GROUP BY "user_id"
|
||||
) AS cursor_summary
|
||||
WHERE cursor_summary."user_id" = "user"."id";
|
||||
|
||||
ALTER TABLE "device_sync_cursors"
|
||||
ADD COLUMN "user_id" UUID,
|
||||
ADD COLUMN "cursor" BIGINT NOT NULL DEFAULT 0;
|
||||
|
||||
UPDATE "device_sync_cursors" AS "sync_cursor"
|
||||
SET
|
||||
"user_id" = device."user_id",
|
||||
"cursor" = "sync_cursor"."last_event_id"
|
||||
FROM "devices" AS device
|
||||
WHERE device."id" = "sync_cursor"."device_id";
|
||||
|
||||
ALTER TABLE "device_sync_cursors"
|
||||
ALTER COLUMN "user_id" SET NOT NULL;
|
||||
|
||||
ALTER TABLE "device_sync_cursors"
|
||||
DROP COLUMN "last_event_id",
|
||||
DROP COLUMN "last_full_sync_at";
|
||||
|
||||
CREATE INDEX "device_sync_cursors_user_id_idx"
|
||||
ON "device_sync_cursors"("user_id");
|
||||
|
||||
ALTER TABLE "device_sync_cursors"
|
||||
ADD CONSTRAINT "device_sync_cursors_user_id_fkey"
|
||||
FOREIGN KEY ("user_id") REFERENCES "users"("id")
|
||||
ON DELETE CASCADE
|
||||
ON UPDATE CASCADE;
|
||||
@@ -12,6 +12,7 @@ model User {
|
||||
slug String @unique
|
||||
displayName String @map("display_name")
|
||||
isDefault Boolean @default(false) @map("is_default")
|
||||
libraryCursor BigInt @default(0) @map("library_cursor")
|
||||
createdAt DateTime @default(now()) @map("created_at")
|
||||
updatedAt DateTime @updatedAt @map("updated_at")
|
||||
devices Device[]
|
||||
@@ -20,6 +21,7 @@ model User {
|
||||
artworkAssets ArtworkAsset[]
|
||||
uploadSessions UploadSession[]
|
||||
libraryEvents LibraryEvent[]
|
||||
syncCursors DeviceSyncCursor[]
|
||||
|
||||
@@map("users")
|
||||
}
|
||||
@@ -145,24 +147,30 @@ model UploadSession {
|
||||
model LibraryEvent {
|
||||
id BigInt @id @default(autoincrement())
|
||||
userId String @db.Uuid @map("user_id")
|
||||
cursor BigInt
|
||||
entityType EntityType @map("entity_type")
|
||||
entityId String @db.Uuid @map("entity_id")
|
||||
action EventAction
|
||||
payload Json
|
||||
payloadVersion Int @default(1) @map("payload_version")
|
||||
createdAt DateTime @default(now()) @map("created_at")
|
||||
user User @relation(fields: [userId], references: [id], onDelete: Cascade, onUpdate: Cascade)
|
||||
|
||||
@@index([userId])
|
||||
@@unique([userId, cursor])
|
||||
@@index([userId, cursor])
|
||||
@@map("library_events")
|
||||
}
|
||||
|
||||
model DeviceSyncCursor {
|
||||
deviceId String @id @db.Uuid @map("device_id")
|
||||
lastEventId BigInt @default(0) @map("last_event_id")
|
||||
lastFullSyncAt DateTime? @map("last_full_sync_at")
|
||||
updatedAt DateTime @updatedAt @map("updated_at")
|
||||
device Device @relation(fields: [deviceId], references: [id], onDelete: Cascade, onUpdate: Cascade)
|
||||
deviceId String @id @db.Uuid @map("device_id")
|
||||
userId String @db.Uuid @map("user_id")
|
||||
cursor BigInt @default(0)
|
||||
updatedAt DateTime @updatedAt @map("updated_at")
|
||||
device Device @relation(fields: [deviceId], references: [id], onDelete: Cascade, onUpdate: Cascade)
|
||||
user User @relation(fields: [userId], references: [id], onDelete: Cascade, onUpdate: Cascade)
|
||||
|
||||
@@index([userId])
|
||||
@@map("device_sync_cursors")
|
||||
}
|
||||
|
||||
|
||||
@@ -41,6 +41,12 @@ export class LibraryService {
|
||||
const { userId: ownerUserId } =
|
||||
this.deviceAuthService.getAuthenticatedDeviceOrThrow();
|
||||
|
||||
return this.getRemoteLibraryTracksForUser(ownerUserId);
|
||||
}
|
||||
|
||||
async getRemoteLibraryTracksForUser(
|
||||
ownerUserId: string,
|
||||
): Promise<RemoteLibraryTrackDto[]> {
|
||||
const tracks = await this.prismaService.track.findMany({
|
||||
where: {
|
||||
userId: ownerUserId,
|
||||
|
||||
@@ -2,7 +2,6 @@ import { Controller, Get, Query, UseGuards } from '@nestjs/common';
|
||||
import { ApiBearerAuth, ApiOkResponse, ApiTags } from '@nestjs/swagger';
|
||||
import { DeviceAuthGuard } from '../auth/device-auth.guard';
|
||||
import {
|
||||
SyncBootstrapQueryDto,
|
||||
SyncBootstrapResponseDto,
|
||||
SyncChangesQueryDto,
|
||||
SyncChangesResponseDto,
|
||||
@@ -21,9 +20,7 @@ export class SyncController {
|
||||
|
||||
@Get('bootstrap')
|
||||
@ApiOkResponse({ type: SyncBootstrapResponseDto })
|
||||
async bootstrap(
|
||||
@Query() _query?: SyncBootstrapQueryDto,
|
||||
): Promise<SyncBootstrapResponseDto> {
|
||||
async bootstrap(): Promise<SyncBootstrapResponseDto> {
|
||||
return this.syncService.bootstrap();
|
||||
}
|
||||
|
||||
@@ -32,6 +29,6 @@ export class SyncController {
|
||||
async changes(
|
||||
@Query() query: SyncChangesQueryDto,
|
||||
): Promise<SyncChangesResponseDto> {
|
||||
return this.syncService.changes(query.after ?? '0');
|
||||
return this.syncService.changes(query.cursor ?? '0', query.limit);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
import { ApiProperty } from '@nestjs/swagger';
|
||||
import { IsOptional, IsString, IsUUID, Matches } from 'class-validator';
|
||||
import { Type } from 'class-transformer';
|
||||
import { IsInt, IsOptional, IsString, Matches, Max, Min } from 'class-validator';
|
||||
import { RemoteLibraryTrackDto } from '../library/library.dto';
|
||||
|
||||
export class LibraryTrackDto {
|
||||
@ApiProperty({ format: 'uuid', required: false })
|
||||
@@ -13,54 +15,83 @@ export class LibraryTrackDto {
|
||||
}
|
||||
|
||||
export class SyncEventDto {
|
||||
@ApiProperty({ example: '3' })
|
||||
cursor!: string;
|
||||
|
||||
@ApiProperty({ example: 'TRACK' })
|
||||
entityType!: string;
|
||||
|
||||
@ApiProperty({ format: 'uuid' })
|
||||
entityId!: string;
|
||||
|
||||
@ApiProperty({ example: 'CREATED' })
|
||||
@ApiProperty({ example: 'UPDATED' })
|
||||
action!: string;
|
||||
|
||||
@ApiProperty({ example: '0' })
|
||||
eventId!: string;
|
||||
}
|
||||
@ApiProperty({
|
||||
type: RemoteLibraryTrackDto,
|
||||
required: false,
|
||||
nullable: true,
|
||||
})
|
||||
track!: RemoteLibraryTrackDto | null;
|
||||
|
||||
export class SyncBootstrapResponseDto {
|
||||
@ApiProperty({ example: '0' })
|
||||
nextCursor!: string;
|
||||
|
||||
@ApiProperty({ type: [LibraryTrackDto] })
|
||||
tracks!: LibraryTrackDto[];
|
||||
|
||||
@ApiProperty({ type: [SyncEventDto] })
|
||||
events!: SyncEventDto[];
|
||||
|
||||
@ApiProperty({ type: [String] })
|
||||
deletedTrackIds!: string[];
|
||||
|
||||
@ApiProperty({ example: '2026-05-24T20:00:00.000Z' })
|
||||
serverTime!: string;
|
||||
}
|
||||
|
||||
export class SyncBootstrapQueryDto {
|
||||
@ApiProperty({
|
||||
format: 'uuid',
|
||||
required: false,
|
||||
description:
|
||||
'Optional client metadata. Authorization: Bearer <deviceAccessToken> is required and determines access.',
|
||||
nullable: true,
|
||||
})
|
||||
@IsOptional()
|
||||
@IsUUID()
|
||||
deviceId?: string;
|
||||
deletedTrackId!: string | null;
|
||||
|
||||
@ApiProperty({ example: '2026-06-15T12:00:00.000Z' })
|
||||
createdAt!: string;
|
||||
}
|
||||
|
||||
export class SyncChangesQueryDto extends SyncBootstrapQueryDto {
|
||||
export class SyncBootstrapResponseDto {
|
||||
@ApiProperty({ example: '7' })
|
||||
nextCursor!: string;
|
||||
|
||||
@ApiProperty({ type: [RemoteLibraryTrackDto] })
|
||||
tracks!: RemoteLibraryTrackDto[];
|
||||
|
||||
@ApiProperty({ example: '2026-06-15T12:00:00.000Z' })
|
||||
serverTime!: string;
|
||||
}
|
||||
|
||||
export class SyncChangesQueryDto {
|
||||
@ApiProperty({ required: false, example: '0' })
|
||||
@IsOptional()
|
||||
@IsString()
|
||||
@Matches(/^\d+$/)
|
||||
after?: string;
|
||||
cursor?: string;
|
||||
|
||||
@ApiProperty({ required: false, example: 100, minimum: 1, maximum: 500 })
|
||||
@IsOptional()
|
||||
@Type(() => Number)
|
||||
@IsInt()
|
||||
@Min(1)
|
||||
@Max(500)
|
||||
limit?: number;
|
||||
}
|
||||
|
||||
export class SyncChangesResponseDto extends SyncBootstrapResponseDto {}
|
||||
export class SyncChangesResponseDto {
|
||||
@ApiProperty({ example: '7' })
|
||||
nextCursor!: string;
|
||||
|
||||
@ApiProperty({ example: false })
|
||||
hasMore!: boolean;
|
||||
|
||||
@ApiProperty({ example: false })
|
||||
requiresBootstrap!: boolean;
|
||||
|
||||
@ApiProperty({
|
||||
required: false,
|
||||
nullable: true,
|
||||
example: 'cursor_too_old',
|
||||
})
|
||||
reason!: string | null;
|
||||
|
||||
@ApiProperty({ type: [SyncEventDto] })
|
||||
events!: SyncEventDto[];
|
||||
|
||||
@ApiProperty({ example: '2026-06-15T12:00:00.000Z' })
|
||||
serverTime!: string;
|
||||
}
|
||||
|
||||
@@ -1,58 +1,300 @@
|
||||
import { Test } from '@nestjs/testing';
|
||||
import { PrismaService } from '../../infrastructure/database/prisma.service';
|
||||
import { LibraryService } from '../library/library.service';
|
||||
import { OwnerContext } from '../users/owner-context.service';
|
||||
import { SyncService } from './sync.service';
|
||||
|
||||
function makeEvent(params: {
|
||||
userId: string;
|
||||
cursor: bigint;
|
||||
entityType?: string;
|
||||
entityId?: string;
|
||||
action?: string;
|
||||
payload?: Record<string, unknown>;
|
||||
createdAt?: Date;
|
||||
}) {
|
||||
return {
|
||||
id: params.cursor,
|
||||
userId: params.userId,
|
||||
cursor: params.cursor,
|
||||
entityType: params.entityType ?? 'TRACK',
|
||||
entityId: params.entityId ?? `track-${params.cursor.toString()}`,
|
||||
action: params.action ?? 'UPDATED',
|
||||
payloadVersion: 1,
|
||||
payload: params.payload ?? {},
|
||||
createdAt:
|
||||
params.createdAt ??
|
||||
new Date(`2026-06-15T12:00:0${params.cursor.toString()}.000Z`),
|
||||
};
|
||||
}
|
||||
|
||||
describe('SyncService', () => {
|
||||
it('uses OwnerContext to scope the bootstrap cursor lookup', async () => {
|
||||
const ownerContextMock = {
|
||||
resolve: jest.fn().mockResolvedValue({
|
||||
userId: 'bootstrap-owner-id',
|
||||
}),
|
||||
};
|
||||
const prismaMock = {
|
||||
libraryEvent: {
|
||||
findFirst: jest.fn().mockResolvedValue({
|
||||
id: 7n,
|
||||
it('returns the bootstrap snapshot and persists the device cursor', async () => {
|
||||
const upsert = jest.fn();
|
||||
const service = new SyncService(
|
||||
{
|
||||
libraryEvent: {
|
||||
findFirst: jest.fn().mockResolvedValue({ cursor: 7n }),
|
||||
},
|
||||
deviceSyncCursor: {
|
||||
upsert,
|
||||
},
|
||||
} as any,
|
||||
{
|
||||
getRemoteLibraryTracksForUser: jest.fn().mockResolvedValue([
|
||||
{
|
||||
trackId: 'track-123',
|
||||
title: 'Remote Title',
|
||||
artist: 'Remote Artist',
|
||||
durationSeconds: 245,
|
||||
sha256: 'a'.repeat(64),
|
||||
assetId: 'asset-123',
|
||||
createdAt: '2026-06-15T10:00:00.000Z',
|
||||
updatedAt: '2026-06-15T10:05:00.000Z',
|
||||
artwork: null,
|
||||
},
|
||||
]),
|
||||
} as any,
|
||||
{
|
||||
getAuthenticatedDeviceOrThrow: jest.fn().mockReturnValue({
|
||||
deviceId: 'device-123',
|
||||
userId: 'owner-123',
|
||||
}),
|
||||
},
|
||||
};
|
||||
const libraryServiceMock = {
|
||||
getBootstrapTracks: jest.fn().mockResolvedValue([]),
|
||||
};
|
||||
} as any,
|
||||
);
|
||||
|
||||
const moduleRef = await Test.createTestingModule({
|
||||
providers: [
|
||||
SyncService,
|
||||
{
|
||||
provide: PrismaService,
|
||||
useValue: prismaMock,
|
||||
},
|
||||
{
|
||||
provide: LibraryService,
|
||||
useValue: libraryServiceMock,
|
||||
},
|
||||
{
|
||||
provide: OwnerContext,
|
||||
useValue: ownerContextMock,
|
||||
},
|
||||
],
|
||||
}).compile();
|
||||
|
||||
const service = moduleRef.get(SyncService);
|
||||
|
||||
await expect(service.changes('0')).resolves.toMatchObject({
|
||||
await expect(service.bootstrap()).resolves.toMatchObject({
|
||||
nextCursor: '7',
|
||||
tracks: [
|
||||
expect.objectContaining({
|
||||
trackId: 'track-123',
|
||||
}),
|
||||
],
|
||||
});
|
||||
expect(ownerContextMock.resolve).toHaveBeenCalledTimes(1);
|
||||
expect(prismaMock.libraryEvent.findFirst).toHaveBeenCalledWith({
|
||||
expect(upsert).toHaveBeenCalledWith({
|
||||
where: {
|
||||
userId: 'bootstrap-owner-id',
|
||||
deviceId: 'device-123',
|
||||
},
|
||||
orderBy: {
|
||||
id: 'desc',
|
||||
update: {
|
||||
userId: 'owner-123',
|
||||
cursor: 7n,
|
||||
},
|
||||
create: {
|
||||
deviceId: 'device-123',
|
||||
userId: 'owner-123',
|
||||
cursor: 7n,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
it('returns ordered changes after the requested cursor', async () => {
|
||||
const ownerId = 'owner-123';
|
||||
const foreignId = 'owner-999';
|
||||
const events = [
|
||||
makeEvent({
|
||||
userId: foreignId,
|
||||
cursor: 1n,
|
||||
}),
|
||||
makeEvent({
|
||||
userId: ownerId,
|
||||
cursor: 2n,
|
||||
payload: {
|
||||
track: {
|
||||
trackId: 'track-2',
|
||||
title: 'Two',
|
||||
artist: 'Owner',
|
||||
durationSeconds: 200,
|
||||
sha256: 'b'.repeat(64),
|
||||
assetId: 'asset-2',
|
||||
createdAt: '2026-06-15T10:00:02.000Z',
|
||||
updatedAt: '2026-06-15T10:00:02.000Z',
|
||||
artwork: null,
|
||||
},
|
||||
},
|
||||
}),
|
||||
makeEvent({
|
||||
userId: ownerId,
|
||||
cursor: 3n,
|
||||
action: 'DELETED',
|
||||
payload: {
|
||||
deletedTrackId: 'track-3',
|
||||
},
|
||||
}),
|
||||
];
|
||||
|
||||
const upsert = jest.fn();
|
||||
const findFirst = jest.fn().mockImplementation(async ({ where, orderBy }) => {
|
||||
const filteredEvents = events.filter((event) => event.userId === where.userId);
|
||||
const direction = orderBy.cursor;
|
||||
const sortedEvents = [...filteredEvents].sort((lhs, rhs) =>
|
||||
direction === 'asc'
|
||||
? Number(lhs.cursor - rhs.cursor)
|
||||
: Number(rhs.cursor - lhs.cursor),
|
||||
);
|
||||
return sortedEvents[0] ?? null;
|
||||
});
|
||||
const findMany = jest.fn().mockImplementation(async ({ where, take }) => {
|
||||
return events
|
||||
.filter(
|
||||
(event) =>
|
||||
event.userId === where.userId && event.cursor > where.cursor.gt,
|
||||
)
|
||||
.sort((lhs, rhs) => Number(lhs.cursor - rhs.cursor))
|
||||
.slice(0, take);
|
||||
});
|
||||
const service = new SyncService(
|
||||
{
|
||||
libraryEvent: {
|
||||
findFirst,
|
||||
findMany,
|
||||
},
|
||||
deviceSyncCursor: {
|
||||
upsert,
|
||||
},
|
||||
} as any,
|
||||
{} as any,
|
||||
{
|
||||
getAuthenticatedDeviceOrThrow: jest.fn().mockReturnValue({
|
||||
deviceId: 'device-123',
|
||||
userId: ownerId,
|
||||
}),
|
||||
} as any,
|
||||
);
|
||||
|
||||
const response = await service.changes('1');
|
||||
|
||||
expect(response.requiresBootstrap).toBe(false);
|
||||
expect(response.hasMore).toBe(false);
|
||||
expect(response.nextCursor).toBe('3');
|
||||
expect(response.events.map((event) => event.cursor)).toEqual(['2', '3']);
|
||||
expect(response.events[0]?.track?.trackId).toBe('track-2');
|
||||
expect(response.events[1]?.deletedTrackId).toBe('track-3');
|
||||
expect(upsert).toHaveBeenCalledWith({
|
||||
where: {
|
||||
deviceId: 'device-123',
|
||||
},
|
||||
update: {
|
||||
userId: ownerId,
|
||||
cursor: 3n,
|
||||
},
|
||||
create: {
|
||||
deviceId: 'device-123',
|
||||
userId: ownerId,
|
||||
cursor: 3n,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
it('paginates deterministically and reports hasMore', async () => {
|
||||
jest.useFakeTimers().setSystemTime(new Date('2026-06-15T08:24:36.011Z'));
|
||||
|
||||
try {
|
||||
const ownerId = 'owner-123';
|
||||
const events = [1n, 2n, 3n].map((cursor) =>
|
||||
makeEvent({
|
||||
userId: ownerId,
|
||||
cursor,
|
||||
payload: {
|
||||
track: {
|
||||
trackId: `track-${cursor.toString()}`,
|
||||
title: `Track ${cursor.toString()}`,
|
||||
artist: 'Owner',
|
||||
durationSeconds: 180,
|
||||
sha256: 'c'.repeat(64),
|
||||
assetId: `asset-${cursor.toString()}`,
|
||||
createdAt: '2026-06-15T10:00:00.000Z',
|
||||
updatedAt: '2026-06-15T10:00:00.000Z',
|
||||
artwork: null,
|
||||
},
|
||||
},
|
||||
}),
|
||||
);
|
||||
|
||||
const findFirst = jest.fn().mockImplementation(async ({ where, orderBy }) => {
|
||||
const filteredEvents = events.filter((event) => event.userId === where.userId);
|
||||
const direction = orderBy.cursor;
|
||||
const sortedEvents = [...filteredEvents].sort((lhs, rhs) =>
|
||||
direction === 'asc'
|
||||
? Number(lhs.cursor - rhs.cursor)
|
||||
: Number(rhs.cursor - lhs.cursor),
|
||||
);
|
||||
return sortedEvents[0] ?? null;
|
||||
});
|
||||
const findMany = jest.fn().mockImplementation(async ({ where, take }) => {
|
||||
return events
|
||||
.filter(
|
||||
(event) =>
|
||||
event.userId === where.userId && event.cursor > where.cursor.gt,
|
||||
)
|
||||
.sort((lhs, rhs) => Number(lhs.cursor - rhs.cursor))
|
||||
.slice(0, take);
|
||||
});
|
||||
const upsert = jest.fn();
|
||||
const service = new SyncService(
|
||||
{
|
||||
libraryEvent: {
|
||||
findFirst,
|
||||
findMany,
|
||||
},
|
||||
deviceSyncCursor: {
|
||||
upsert,
|
||||
},
|
||||
} as any,
|
||||
{} as any,
|
||||
{
|
||||
getAuthenticatedDeviceOrThrow: jest.fn().mockReturnValue({
|
||||
deviceId: 'device-123',
|
||||
userId: ownerId,
|
||||
}),
|
||||
} as any,
|
||||
);
|
||||
|
||||
const firstResponse = await service.changes('0', 2);
|
||||
const replayResponse = await service.changes('0', 2);
|
||||
|
||||
expect(firstResponse.hasMore).toBe(true);
|
||||
expect(firstResponse.nextCursor).toBe('2');
|
||||
expect(firstResponse.events.map((event) => event.cursor)).toEqual(['1', '2']);
|
||||
expect(replayResponse).toEqual(firstResponse);
|
||||
} finally {
|
||||
jest.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it('requires bootstrap when the requested cursor is older than retained history', async () => {
|
||||
const service = new SyncService(
|
||||
{
|
||||
libraryEvent: {
|
||||
findFirst: jest.fn().mockImplementation(async ({ where, orderBy }) => {
|
||||
if (orderBy.cursor === 'asc') {
|
||||
return makeEvent({
|
||||
userId: where.userId,
|
||||
cursor: 5n,
|
||||
});
|
||||
}
|
||||
|
||||
return makeEvent({
|
||||
userId: where.userId,
|
||||
cursor: 9n,
|
||||
});
|
||||
}),
|
||||
findMany: jest.fn(),
|
||||
},
|
||||
deviceSyncCursor: {
|
||||
upsert: jest.fn(),
|
||||
},
|
||||
} as any,
|
||||
{} as any,
|
||||
{
|
||||
getAuthenticatedDeviceOrThrow: jest.fn().mockReturnValue({
|
||||
deviceId: 'device-123',
|
||||
userId: 'owner-123',
|
||||
}),
|
||||
} as any,
|
||||
);
|
||||
|
||||
await expect(service.changes('3')).resolves.toMatchObject({
|
||||
requiresBootstrap: true,
|
||||
reason: 'cursor_too_old',
|
||||
events: [],
|
||||
hasMore: false,
|
||||
nextCursor: '3',
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,57 +1,209 @@
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { PrismaService } from '../../infrastructure/database/prisma.service';
|
||||
import { DeviceAuthService } from '../auth/device-auth.service';
|
||||
import { RemoteLibraryTrackDto } from '../library/library.dto';
|
||||
import { LibraryService } from '../library/library.service';
|
||||
import { OwnerContext } from '../users/owner-context.service';
|
||||
import { SyncBootstrapResponseDto, SyncChangesResponseDto } from './sync.dto';
|
||||
import {
|
||||
SyncBootstrapResponseDto,
|
||||
SyncChangesResponseDto,
|
||||
SyncEventDto,
|
||||
} from './sync.dto';
|
||||
|
||||
interface LibraryEventPayload {
|
||||
track?: RemoteLibraryTrackDto | null;
|
||||
deletedTrackId?: string | null;
|
||||
}
|
||||
|
||||
const DEFAULT_SYNC_PAGE_SIZE = 100;
|
||||
|
||||
@Injectable()
|
||||
export class SyncService {
|
||||
constructor(
|
||||
private readonly prismaService: PrismaService,
|
||||
private readonly libraryService: LibraryService,
|
||||
private readonly ownerContext: OwnerContext,
|
||||
private readonly deviceAuthService: DeviceAuthService,
|
||||
) {}
|
||||
|
||||
async bootstrap(): Promise<SyncBootstrapResponseDto> {
|
||||
const latestCursor = await this.getLatestCursor();
|
||||
const device = this.deviceAuthService.getAuthenticatedDeviceOrThrow();
|
||||
const [tracks, latestCursor] = await Promise.all([
|
||||
this.libraryService.getRemoteLibraryTracksForUser(device.userId),
|
||||
this.getLatestCursor(device.userId),
|
||||
]);
|
||||
|
||||
return {
|
||||
const response: SyncBootstrapResponseDto = {
|
||||
nextCursor: latestCursor,
|
||||
tracks: await this.libraryService.getBootstrapTracks(),
|
||||
events: [],
|
||||
deletedTrackIds: [],
|
||||
tracks,
|
||||
serverTime: new Date().toISOString(),
|
||||
};
|
||||
|
||||
await this.updateDeviceSyncCursor(
|
||||
device.deviceId,
|
||||
device.userId,
|
||||
response.nextCursor,
|
||||
);
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
async changes(after: string): Promise<SyncChangesResponseDto> {
|
||||
const latestCursor = await this.getLatestCursor();
|
||||
const normalizedCursor =
|
||||
BigInt(latestCursor) > BigInt(after) ? latestCursor : after;
|
||||
async changes(
|
||||
cursor: string,
|
||||
limit = DEFAULT_SYNC_PAGE_SIZE,
|
||||
): Promise<SyncChangesResponseDto> {
|
||||
const device = this.deviceAuthService.getAuthenticatedDeviceOrThrow();
|
||||
const requestedCursor = BigInt(cursor);
|
||||
const earliestRetainedCursor = await this.getEarliestRetainedCursor(
|
||||
device.userId,
|
||||
);
|
||||
|
||||
return {
|
||||
nextCursor: normalizedCursor,
|
||||
tracks: [],
|
||||
events: [],
|
||||
deletedTrackIds: [],
|
||||
serverTime: new Date().toISOString(),
|
||||
};
|
||||
}
|
||||
if (
|
||||
requestedCursor > 0n &&
|
||||
earliestRetainedCursor !== null &&
|
||||
requestedCursor < earliestRetainedCursor - 1n
|
||||
) {
|
||||
return {
|
||||
nextCursor: cursor,
|
||||
hasMore: false,
|
||||
requiresBootstrap: true,
|
||||
reason: 'cursor_too_old',
|
||||
events: [],
|
||||
serverTime: new Date().toISOString(),
|
||||
};
|
||||
}
|
||||
|
||||
private async getLatestCursor(): Promise<string> {
|
||||
const owner = await this.ownerContext.resolve({
|
||||
allowLegacyDeviceFallback: false,
|
||||
allowBootstrapFallback: false,
|
||||
});
|
||||
const latest = await this.prismaService.libraryEvent.findFirst({
|
||||
const events = await this.prismaService.libraryEvent.findMany({
|
||||
where: {
|
||||
userId: owner.userId,
|
||||
userId: device.userId,
|
||||
cursor: {
|
||||
gt: requestedCursor,
|
||||
},
|
||||
},
|
||||
orderBy: {
|
||||
id: 'desc',
|
||||
cursor: 'asc',
|
||||
},
|
||||
take: limit + 1,
|
||||
});
|
||||
|
||||
const hasMore = events.length > limit;
|
||||
const visibleEvents = hasMore ? events.slice(0, limit) : events;
|
||||
const nextCursor =
|
||||
visibleEvents.at(-1)?.cursor.toString() ?? requestedCursor.toString();
|
||||
|
||||
const response: SyncChangesResponseDto = {
|
||||
nextCursor,
|
||||
hasMore,
|
||||
requiresBootstrap: false,
|
||||
reason: null,
|
||||
events: visibleEvents.map((event) => this.toSyncEventDto(event)),
|
||||
serverTime: new Date().toISOString(),
|
||||
};
|
||||
|
||||
await this.updateDeviceSyncCursor(
|
||||
device.deviceId,
|
||||
device.userId,
|
||||
response.nextCursor,
|
||||
);
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
private async getLatestCursor(userId: string): Promise<string> {
|
||||
const latest = await this.prismaService.libraryEvent.findFirst({
|
||||
where: {
|
||||
userId,
|
||||
},
|
||||
orderBy: {
|
||||
cursor: 'desc',
|
||||
},
|
||||
select: {
|
||||
cursor: true,
|
||||
},
|
||||
});
|
||||
|
||||
return latest?.id.toString() ?? '0';
|
||||
return latest?.cursor.toString() ?? '0';
|
||||
}
|
||||
|
||||
private async getEarliestRetainedCursor(
|
||||
userId: string,
|
||||
): Promise<bigint | null> {
|
||||
const earliest = await this.prismaService.libraryEvent.findFirst({
|
||||
where: {
|
||||
userId,
|
||||
},
|
||||
orderBy: {
|
||||
cursor: 'asc',
|
||||
},
|
||||
select: {
|
||||
cursor: true,
|
||||
},
|
||||
});
|
||||
|
||||
return earliest?.cursor ?? null;
|
||||
}
|
||||
|
||||
private async updateDeviceSyncCursor(
|
||||
deviceId: string,
|
||||
userId: string,
|
||||
cursor: string,
|
||||
): Promise<void> {
|
||||
await this.prismaService.deviceSyncCursor.upsert({
|
||||
where: {
|
||||
deviceId,
|
||||
},
|
||||
update: {
|
||||
userId,
|
||||
cursor: BigInt(cursor),
|
||||
},
|
||||
create: {
|
||||
deviceId,
|
||||
userId,
|
||||
cursor: BigInt(cursor),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
private toSyncEventDto(event: {
|
||||
cursor: bigint;
|
||||
entityType: string;
|
||||
entityId: string;
|
||||
action: string;
|
||||
payload: unknown;
|
||||
createdAt: Date;
|
||||
}): SyncEventDto {
|
||||
const payload = this.parsePayload(event.payload);
|
||||
const deletedTrackId =
|
||||
payload.deletedTrackId ??
|
||||
(event.action === 'DELETED' && event.entityType === 'TRACK'
|
||||
? event.entityId
|
||||
: null);
|
||||
|
||||
return {
|
||||
cursor: event.cursor.toString(),
|
||||
entityType: event.entityType,
|
||||
entityId: event.entityId,
|
||||
action: event.action,
|
||||
track: payload.track ?? null,
|
||||
deletedTrackId,
|
||||
createdAt: event.createdAt.toISOString(),
|
||||
};
|
||||
}
|
||||
|
||||
private parsePayload(payload: unknown): LibraryEventPayload {
|
||||
if (!payload || typeof payload !== 'object' || Array.isArray(payload)) {
|
||||
return {};
|
||||
}
|
||||
|
||||
const record = payload as Record<string, unknown>;
|
||||
const track =
|
||||
record.track && typeof record.track === 'object' && !Array.isArray(record.track)
|
||||
? (record.track as RemoteLibraryTrackDto)
|
||||
: null;
|
||||
const deletedTrackId =
|
||||
typeof record.deletedTrackId === 'string' ? record.deletedTrackId : null;
|
||||
|
||||
return {
|
||||
track,
|
||||
deletedTrackId,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,6 +28,7 @@ function createPrismaMock() {
|
||||
slug: 'default-owner',
|
||||
displayName: 'Default Owner',
|
||||
isDefault: true,
|
||||
libraryCursor: 0n,
|
||||
createdAt: new Date(),
|
||||
updatedAt: new Date(),
|
||||
};
|
||||
@@ -38,6 +39,24 @@ function createPrismaMock() {
|
||||
$transaction: jest.fn().mockImplementation(async (callback: any) => callback(prismaMock)),
|
||||
user: {
|
||||
upsert: jest.fn().mockResolvedValue(defaultUser),
|
||||
update: jest.fn().mockImplementation(async ({ where, data, select }) => {
|
||||
const current = users.get(where.id);
|
||||
const incrementBy = BigInt(data.libraryCursor?.increment ?? 0);
|
||||
const updated = {
|
||||
...current,
|
||||
libraryCursor: BigInt(current.libraryCursor ?? 0) + incrementBy,
|
||||
updatedAt: new Date(),
|
||||
};
|
||||
users.set(where.id, updated);
|
||||
|
||||
if (select?.libraryCursor) {
|
||||
return {
|
||||
libraryCursor: updated.libraryCursor,
|
||||
};
|
||||
}
|
||||
|
||||
return updated;
|
||||
}),
|
||||
},
|
||||
device: {
|
||||
findUnique: jest.fn().mockImplementation(async ({ where }) => {
|
||||
@@ -218,11 +237,17 @@ function createPrismaMock() {
|
||||
nextLibraryEventId += 1n;
|
||||
return record;
|
||||
}),
|
||||
findFirst: jest.fn().mockImplementation(async ({ where }) => {
|
||||
findFirst: jest.fn().mockImplementation(async ({ where, orderBy }) => {
|
||||
const filteredEvents = [...libraryEvents.values()].filter((event) =>
|
||||
where?.userId ? event.userId === where.userId : true,
|
||||
);
|
||||
return filteredEvents.sort((lhs, rhs) => Number(rhs.id - lhs.id))[0] ?? null;
|
||||
const direction = orderBy?.cursor ?? 'desc';
|
||||
return filteredEvents
|
||||
.sort((lhs, rhs) =>
|
||||
direction === 'asc'
|
||||
? Number(lhs.cursor - rhs.cursor)
|
||||
: Number(rhs.cursor - lhs.cursor),
|
||||
)[0] ?? null;
|
||||
}),
|
||||
},
|
||||
state: {
|
||||
@@ -463,15 +488,27 @@ describe('UploadsService', () => {
|
||||
expect(finalizeResponse.assetId).toBeDefined();
|
||||
expect(state.tracks.size).toBe(1);
|
||||
expect(state.audioAssets.size).toBe(1);
|
||||
expect(state.libraryEvents.size).toBe(1);
|
||||
expect(state.libraryEvents.size).toBe(2);
|
||||
|
||||
const track = [...state.tracks.values()][0];
|
||||
const audioAsset = [...state.audioAssets.values()][0];
|
||||
const libraryEvent = [...state.libraryEvents.values()][0];
|
||||
const libraryEvents = [...state.libraryEvents.values()].sort((lhs, rhs) =>
|
||||
Number(lhs.cursor - rhs.cursor),
|
||||
);
|
||||
|
||||
expect(track.userId).toBe(state.defaultUser.id);
|
||||
expect(audioAsset.userId).toBe(state.defaultUser.id);
|
||||
expect(libraryEvent.userId).toBe(state.defaultUser.id);
|
||||
expect(libraryEvents.map((event) => event.entityType)).toEqual([
|
||||
'TRACK',
|
||||
'AUDIO_ASSET',
|
||||
]);
|
||||
expect(libraryEvents.map((event) => event.action)).toEqual([
|
||||
'CREATED',
|
||||
'CREATED',
|
||||
]);
|
||||
expect(libraryEvents[0]?.payload.track.trackId).toBe(track.id);
|
||||
expect(libraryEvents[1]?.payload.track.assetId).toBe(audioAsset.id);
|
||||
expect(state.users.get(state.defaultUser.id)?.libraryCursor).toBe(2n);
|
||||
|
||||
const session = state.uploadSessions.get(response.uploadId!);
|
||||
expect(session.finalizedAt).toBeInstanceOf(Date);
|
||||
|
||||
@@ -7,6 +7,7 @@ import {
|
||||
import {
|
||||
EntityType,
|
||||
EventAction,
|
||||
Prisma,
|
||||
type UploadSession,
|
||||
UploadSessionStatus,
|
||||
} from '@prisma/client';
|
||||
@@ -18,6 +19,7 @@ import { extname } from 'node:path';
|
||||
import { PrismaService } from '../../infrastructure/database/prisma.service';
|
||||
import { AppConfigService } from '../config/config.service';
|
||||
import { DeviceAuthService } from '../auth/device-auth.service';
|
||||
import { RemoteLibraryTrackDto } from '../library/library.dto';
|
||||
import { LocalFilesystemStorageService } from '../storage/storage.service';
|
||||
import { OwnerContext } from '../users/owner-context.service';
|
||||
import {
|
||||
@@ -38,6 +40,11 @@ interface PreparedArtworkAssetInput {
|
||||
fileSizeBytes: bigint;
|
||||
}
|
||||
|
||||
interface LibraryEventPayload {
|
||||
track?: RemoteLibraryTrackDto;
|
||||
deletedTrackId?: string;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class UploadsService {
|
||||
constructor(
|
||||
@@ -367,6 +374,8 @@ export class UploadsService {
|
||||
}
|
||||
|
||||
const createdTrack = !track;
|
||||
let trackMetadataChanged = false;
|
||||
|
||||
if (!track) {
|
||||
track = await tx.track.create({
|
||||
data: {
|
||||
@@ -378,8 +387,34 @@ export class UploadsService {
|
||||
status: 'ACTIVE',
|
||||
},
|
||||
});
|
||||
} else {
|
||||
const nextTrackDurationMs = body.durationMs ?? track.durationMs;
|
||||
const shouldUpdateTrack =
|
||||
track.title !== title ||
|
||||
track.artist !== artist ||
|
||||
(track.album ?? null) !== album ||
|
||||
(track.durationMs ?? null) !== (nextTrackDurationMs ?? null) ||
|
||||
track.status !== 'ACTIVE' ||
|
||||
track.deletedAt !== null;
|
||||
|
||||
if (shouldUpdateTrack) {
|
||||
track = await tx.track.update({
|
||||
where: { id: track.id },
|
||||
data: {
|
||||
title,
|
||||
artist,
|
||||
album,
|
||||
durationMs: nextTrackDurationMs,
|
||||
status: 'ACTIVE',
|
||||
deletedAt: null,
|
||||
},
|
||||
});
|
||||
trackMetadataChanged = true;
|
||||
}
|
||||
}
|
||||
|
||||
const createdAudioAsset = !audioAsset;
|
||||
let audioAssetChanged = createdAudioAsset;
|
||||
if (audioAsset) {
|
||||
const nextDurationMs = body.durationMs ?? audioAsset.durationMs;
|
||||
const shouldUpdateAsset =
|
||||
@@ -404,6 +439,7 @@ export class UploadsService {
|
||||
durationMs: nextDurationMs,
|
||||
},
|
||||
});
|
||||
audioAssetChanged = true;
|
||||
}
|
||||
} else {
|
||||
audioAsset = await tx.audioAsset.create({
|
||||
@@ -422,6 +458,7 @@ export class UploadsService {
|
||||
});
|
||||
}
|
||||
|
||||
let primaryAudioAssetChanged = false;
|
||||
if (track.primaryAudioAssetId !== audioAsset.id) {
|
||||
track = await tx.track.update({
|
||||
where: { id: track.id },
|
||||
@@ -429,18 +466,21 @@ export class UploadsService {
|
||||
primaryAudioAssetId: audioAsset.id,
|
||||
},
|
||||
});
|
||||
primaryAudioAssetChanged = true;
|
||||
}
|
||||
|
||||
const artworkAssetId = preparedArtwork
|
||||
? (
|
||||
await this.findOrCreateArtworkAsset(
|
||||
tx,
|
||||
ownerUserId,
|
||||
preparedArtwork,
|
||||
)
|
||||
).id
|
||||
const priorArtworkAssetId = track.artworkAssetId ?? null;
|
||||
const artworkResult = preparedArtwork
|
||||
? await this.findOrCreateArtworkAsset(
|
||||
tx,
|
||||
ownerUserId,
|
||||
preparedArtwork,
|
||||
)
|
||||
: null;
|
||||
const artworkAsset = artworkResult?.artworkAsset ?? null;
|
||||
const artworkAssetId = artworkAsset?.id ?? null;
|
||||
|
||||
let artworkLinkChanged = false;
|
||||
if ((track.artworkAssetId ?? null) !== artworkAssetId) {
|
||||
track = await tx.track.update({
|
||||
where: { id: track.id },
|
||||
@@ -448,16 +488,57 @@ export class UploadsService {
|
||||
artworkAssetId,
|
||||
},
|
||||
});
|
||||
artworkLinkChanged = true;
|
||||
}
|
||||
|
||||
await tx.libraryEvent.create({
|
||||
data: {
|
||||
const finalTrackSnapshot = this.buildRemoteLibraryTrackDto(
|
||||
track,
|
||||
audioAsset,
|
||||
artworkAssetId ? artworkAsset : null,
|
||||
);
|
||||
const eventPayload: LibraryEventPayload = {
|
||||
track: finalTrackSnapshot,
|
||||
};
|
||||
|
||||
if (createdTrack || trackMetadataChanged) {
|
||||
await this.appendLibraryEvent(tx, {
|
||||
userId: ownerUserId,
|
||||
entityType: EntityType.TRACK,
|
||||
entityId: track.id,
|
||||
action: createdTrack ? EventAction.CREATED : EventAction.UPDATED,
|
||||
},
|
||||
});
|
||||
payload: eventPayload,
|
||||
});
|
||||
}
|
||||
|
||||
if (audioAssetChanged || primaryAudioAssetChanged) {
|
||||
await this.appendLibraryEvent(tx, {
|
||||
userId: ownerUserId,
|
||||
entityType: EntityType.AUDIO_ASSET,
|
||||
entityId: audioAsset.id,
|
||||
action: createdAudioAsset ? EventAction.CREATED : EventAction.UPDATED,
|
||||
payload: eventPayload,
|
||||
});
|
||||
}
|
||||
|
||||
if (
|
||||
artworkResult?.wasCreated ||
|
||||
artworkResult?.wasUpdated ||
|
||||
artworkLinkChanged ||
|
||||
(priorArtworkAssetId !== null && artworkAssetId === null)
|
||||
) {
|
||||
await this.appendLibraryEvent(tx, {
|
||||
userId: ownerUserId,
|
||||
entityType: EntityType.ARTWORK_ASSET,
|
||||
entityId: artworkAssetId ?? priorArtworkAssetId!,
|
||||
action:
|
||||
artworkAssetId == null
|
||||
? EventAction.DELETED
|
||||
: artworkResult?.wasCreated
|
||||
? EventAction.CREATED
|
||||
: EventAction.UPDATED,
|
||||
payload: eventPayload,
|
||||
});
|
||||
}
|
||||
|
||||
await tx.uploadSession.update({
|
||||
where: { id: currentSession.id },
|
||||
@@ -542,6 +623,87 @@ export class UploadsService {
|
||||
}
|
||||
}
|
||||
|
||||
private buildRemoteLibraryTrackDto(
|
||||
track: {
|
||||
id: string;
|
||||
title: string;
|
||||
artist: string;
|
||||
durationMs: number | null;
|
||||
createdAt: Date;
|
||||
updatedAt: Date;
|
||||
},
|
||||
audioAsset: {
|
||||
id: string;
|
||||
sha256: string;
|
||||
durationMs: number | null;
|
||||
},
|
||||
artworkAsset: {
|
||||
id: string;
|
||||
sha256: string;
|
||||
mimeType: string;
|
||||
width: number | null;
|
||||
height: number | null;
|
||||
} | null,
|
||||
): RemoteLibraryTrackDto {
|
||||
const durationMs = track.durationMs ?? audioAsset.durationMs ?? 0;
|
||||
|
||||
return {
|
||||
trackId: track.id,
|
||||
title: track.title,
|
||||
artist: track.artist,
|
||||
durationSeconds: Math.max(0, Math.round(durationMs / 1000)),
|
||||
sha256: audioAsset.sha256,
|
||||
assetId: audioAsset.id,
|
||||
createdAt: track.createdAt.toISOString(),
|
||||
updatedAt: track.updatedAt.toISOString(),
|
||||
artwork: artworkAsset
|
||||
? {
|
||||
artworkId: artworkAsset.id,
|
||||
sha256: artworkAsset.sha256,
|
||||
mimeType: artworkAsset.mimeType,
|
||||
width: artworkAsset.width,
|
||||
height: artworkAsset.height,
|
||||
}
|
||||
: null,
|
||||
};
|
||||
}
|
||||
|
||||
private async appendLibraryEvent(
|
||||
tx: Pick<PrismaService, 'user' | 'libraryEvent'>,
|
||||
params: {
|
||||
userId: string;
|
||||
entityType: EntityType;
|
||||
entityId: string;
|
||||
action: EventAction;
|
||||
payload: LibraryEventPayload;
|
||||
},
|
||||
): Promise<void> {
|
||||
const owner = await tx.user.update({
|
||||
where: {
|
||||
id: params.userId,
|
||||
},
|
||||
data: {
|
||||
libraryCursor: {
|
||||
increment: 1,
|
||||
},
|
||||
},
|
||||
select: {
|
||||
libraryCursor: true,
|
||||
},
|
||||
});
|
||||
|
||||
await tx.libraryEvent.create({
|
||||
data: {
|
||||
userId: params.userId,
|
||||
cursor: owner.libraryCursor,
|
||||
entityType: params.entityType,
|
||||
entityId: params.entityId,
|
||||
action: params.action,
|
||||
payload: params.payload as Prisma.InputJsonValue,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
private toStatusResponse(
|
||||
uploadSession: Pick<
|
||||
UploadSession,
|
||||
@@ -636,12 +798,22 @@ export class UploadsService {
|
||||
fileSizeBytes: artwork.fileSizeBytes,
|
||||
},
|
||||
});
|
||||
|
||||
return {
|
||||
artworkAsset,
|
||||
wasCreated: false,
|
||||
wasUpdated: true,
|
||||
};
|
||||
}
|
||||
|
||||
return artworkAsset;
|
||||
return {
|
||||
artworkAsset,
|
||||
wasCreated: false,
|
||||
wasUpdated: false,
|
||||
};
|
||||
}
|
||||
|
||||
return tx.artworkAsset.create({
|
||||
const createdArtworkAsset = await tx.artworkAsset.create({
|
||||
data: {
|
||||
userId,
|
||||
sha256: artwork.sha256,
|
||||
@@ -652,6 +824,12 @@ export class UploadsService {
|
||||
fileSizeBytes: artwork.fileSizeBytes,
|
||||
},
|
||||
});
|
||||
|
||||
return {
|
||||
artworkAsset: createdArtworkAsset,
|
||||
wasCreated: true,
|
||||
wasUpdated: false,
|
||||
};
|
||||
}
|
||||
|
||||
private assertMp3Filename(filename: string): void {
|
||||
|
||||
@@ -8,6 +8,7 @@ describe('DefaultUserService', () => {
|
||||
slug: DefaultUserService.defaultOwnerSlug,
|
||||
displayName: DefaultUserService.defaultOwnerDisplayName,
|
||||
isDefault: true,
|
||||
libraryCursor: 0n,
|
||||
createdAt: new Date(),
|
||||
updatedAt: new Date(),
|
||||
};
|
||||
@@ -48,6 +49,7 @@ describe('DefaultUserService', () => {
|
||||
slug: DefaultUserService.defaultOwnerSlug,
|
||||
displayName: DefaultUserService.defaultOwnerDisplayName,
|
||||
isDefault: true,
|
||||
libraryCursor: 0n,
|
||||
createdAt: new Date(),
|
||||
updatedAt: new Date(),
|
||||
});
|
||||
|
||||
@@ -84,6 +84,7 @@ async function streamToBuffer(stream: NodeJS.ReadableStream): Promise<Buffer> {
|
||||
function createPrismaMock() {
|
||||
const users = new Map<string, any>();
|
||||
const devices = new Map<string, any>();
|
||||
const deviceSyncCursors = new Map<string, any>();
|
||||
const tracks = new Map<string, any>();
|
||||
const audioAssets = new Map<string, any>();
|
||||
const artworkAssets = new Map<string, any>();
|
||||
@@ -91,21 +92,81 @@ function createPrismaMock() {
|
||||
const libraryEvents = new Map<bigint, any>();
|
||||
let nextLibraryEventId = 1n;
|
||||
|
||||
const defaultUser = {
|
||||
const createUserRecord = (data: Record<string, any>) => {
|
||||
const now = new Date();
|
||||
return {
|
||||
id: data.id ?? randomUUID(),
|
||||
createdAt: data.createdAt ?? now,
|
||||
updatedAt: data.updatedAt ?? now,
|
||||
...data,
|
||||
libraryCursor:
|
||||
data.libraryCursor == null ? 0n : BigInt(data.libraryCursor),
|
||||
};
|
||||
};
|
||||
|
||||
const defaultUser = createUserRecord({
|
||||
id: randomUUID(),
|
||||
slug: 'default-owner',
|
||||
displayName: 'Default Owner',
|
||||
isDefault: true,
|
||||
createdAt: new Date(),
|
||||
updatedAt: new Date(),
|
||||
};
|
||||
libraryCursor: 0n,
|
||||
});
|
||||
users.set(defaultUser.id, defaultUser);
|
||||
|
||||
const prismaMock: any = {
|
||||
$queryRawUnsafe: jest.fn().mockResolvedValue([{ '?column?': 1 }]),
|
||||
$transaction: jest.fn().mockImplementation(async (callback: any) => callback(prismaMock)),
|
||||
user: {
|
||||
upsert: jest.fn().mockResolvedValue(defaultUser),
|
||||
upsert: jest.fn().mockImplementation(async ({ where, update, create }) => {
|
||||
const current =
|
||||
[...users.values()].find((user) => user.slug === where.slug) ?? null;
|
||||
|
||||
if (current) {
|
||||
const updated = createUserRecord({
|
||||
...current,
|
||||
...update,
|
||||
id: current.id,
|
||||
createdAt: current.createdAt,
|
||||
updatedAt: new Date(),
|
||||
libraryCursor: current.libraryCursor,
|
||||
});
|
||||
users.set(updated.id, updated);
|
||||
return updated;
|
||||
}
|
||||
|
||||
const created = createUserRecord(create);
|
||||
users.set(created.id, created);
|
||||
return created;
|
||||
}),
|
||||
create: jest.fn().mockImplementation(async ({ data }) => {
|
||||
const created = createUserRecord(data);
|
||||
users.set(created.id, created);
|
||||
return created;
|
||||
}),
|
||||
update: jest.fn().mockImplementation(async ({ where, data, select }) => {
|
||||
const current = users.get(where.id);
|
||||
if (!current) {
|
||||
throw new Error(
|
||||
`Test Prisma mock invariant failed: user ${where.id} not found for update`,
|
||||
);
|
||||
}
|
||||
|
||||
const incrementBy = BigInt(data.libraryCursor?.increment ?? 0);
|
||||
const updated = createUserRecord({
|
||||
...current,
|
||||
updatedAt: new Date(),
|
||||
libraryCursor: current.libraryCursor + incrementBy,
|
||||
});
|
||||
users.set(where.id, updated);
|
||||
|
||||
if (select?.libraryCursor) {
|
||||
return {
|
||||
libraryCursor: updated.libraryCursor,
|
||||
};
|
||||
}
|
||||
|
||||
return updated;
|
||||
}),
|
||||
},
|
||||
device: {
|
||||
create: jest.fn().mockImplementation(async ({ data }) => {
|
||||
@@ -314,11 +375,44 @@ function createPrismaMock() {
|
||||
nextLibraryEventId += 1n;
|
||||
return record;
|
||||
}),
|
||||
findFirst: jest.fn().mockImplementation(async ({ where }) => {
|
||||
findFirst: jest.fn().mockImplementation(async ({ where, orderBy }) => {
|
||||
const filteredEvents = [...libraryEvents.values()].filter((event) =>
|
||||
where?.userId ? event.userId === where.userId : true,
|
||||
);
|
||||
return filteredEvents.sort((lhs, rhs) => Number(rhs.id - lhs.id))[0] ?? null;
|
||||
const direction = orderBy?.cursor ?? 'desc';
|
||||
return filteredEvents
|
||||
.sort((lhs, rhs) =>
|
||||
direction === 'asc'
|
||||
? Number(lhs.cursor - rhs.cursor)
|
||||
: Number(rhs.cursor - lhs.cursor),
|
||||
)[0] ?? null;
|
||||
}),
|
||||
findMany: jest.fn().mockImplementation(async ({ where, take }) => {
|
||||
return [...libraryEvents.values()]
|
||||
.filter(
|
||||
(event) =>
|
||||
(where?.userId ? event.userId === where.userId : true) &&
|
||||
(where?.cursor?.gt != null ? event.cursor > where.cursor.gt : true),
|
||||
)
|
||||
.sort((lhs, rhs) => Number(lhs.cursor - rhs.cursor))
|
||||
.slice(0, take);
|
||||
}),
|
||||
},
|
||||
deviceSyncCursor: {
|
||||
upsert: jest.fn().mockImplementation(async ({ where, update, create }) => {
|
||||
const current = deviceSyncCursors.get(where.deviceId);
|
||||
const nextRecord = current
|
||||
? {
|
||||
...current,
|
||||
...update,
|
||||
updatedAt: new Date(),
|
||||
}
|
||||
: {
|
||||
...create,
|
||||
updatedAt: new Date(),
|
||||
};
|
||||
deviceSyncCursors.set(where.deviceId, nextRecord);
|
||||
return nextRecord;
|
||||
}),
|
||||
},
|
||||
};
|
||||
@@ -327,7 +421,9 @@ function createPrismaMock() {
|
||||
prismaMock,
|
||||
state: {
|
||||
defaultUser,
|
||||
users,
|
||||
devices,
|
||||
deviceSyncCursors,
|
||||
tracks,
|
||||
audioAssets,
|
||||
artworkAssets,
|
||||
@@ -696,12 +792,20 @@ describe('Velody API wiring (e2e)', () => {
|
||||
syncController.bootstrap(),
|
||||
);
|
||||
const changesResponse = await runAsDevice(device.deviceAccessToken, () =>
|
||||
syncController.changes({ after: '0' }),
|
||||
syncController.changes({ cursor: '0' }),
|
||||
);
|
||||
|
||||
expect(bootstrapResponse.tracks).toEqual([]);
|
||||
expect(bootstrapResponse.nextCursor).toBe('0');
|
||||
expect(changesResponse.events).toEqual([]);
|
||||
expect(changesResponse.nextCursor).toBe('0');
|
||||
expect(changesResponse.hasMore).toBe(false);
|
||||
expect(changesResponse.requiresBootstrap).toBe(false);
|
||||
expect(prismaState.deviceSyncCursors.get(device.deviceId)).toMatchObject({
|
||||
deviceId: device.deviceId,
|
||||
userId: prismaState.defaultUser.id,
|
||||
cursor: 0n,
|
||||
});
|
||||
});
|
||||
|
||||
it('sync bootstrap and changes do not expose foreign-owner data', async () => {
|
||||
@@ -734,10 +838,24 @@ describe('Velody API wiring (e2e)', () => {
|
||||
});
|
||||
prismaState.libraryEvents.set(1n, {
|
||||
id: 1n,
|
||||
cursor: 1n,
|
||||
userId: foreignUserId,
|
||||
entityType: 'TRACK',
|
||||
entityId: foreignTrackId,
|
||||
action: 'CREATED',
|
||||
payload: {
|
||||
track: {
|
||||
trackId: foreignTrackId,
|
||||
title: 'Foreign Bootstrap Track',
|
||||
artist: 'Elsewhere',
|
||||
durationSeconds: 180,
|
||||
sha256: 'f'.repeat(64),
|
||||
assetId: randomUUID(),
|
||||
createdAt: '2026-05-29T08:00:00.000Z',
|
||||
updatedAt: '2026-05-29T08:01:00.000Z',
|
||||
artwork: null,
|
||||
},
|
||||
},
|
||||
payloadVersion: 1,
|
||||
createdAt: new Date('2026-05-29T08:02:00.000Z'),
|
||||
});
|
||||
@@ -746,7 +864,7 @@ describe('Velody API wiring (e2e)', () => {
|
||||
syncController.bootstrap(),
|
||||
);
|
||||
const changesResponse = await runAsDevice(device.deviceAccessToken, () =>
|
||||
syncController.changes({ after: '0' }),
|
||||
syncController.changes({ cursor: '0' }),
|
||||
);
|
||||
|
||||
expect(bootstrapResponse.tracks).toEqual([]);
|
||||
@@ -1590,6 +1708,18 @@ describe('Velody API wiring (e2e)', () => {
|
||||
|
||||
it('makes an upload from one device visible to another linked device under the same owner', async () => {
|
||||
const identityUserId = randomUUID();
|
||||
prismaState.users.set(
|
||||
identityUserId,
|
||||
{
|
||||
id: identityUserId,
|
||||
slug: `identity-${identityUserId}`,
|
||||
displayName: 'Identity Owner',
|
||||
isDefault: false,
|
||||
libraryCursor: 0n,
|
||||
createdAt: new Date(),
|
||||
updatedAt: new Date(),
|
||||
},
|
||||
);
|
||||
const primaryDevice = seedDevice({
|
||||
userId: identityUserId,
|
||||
deviceAccessToken: 'linked-upload-primary-token',
|
||||
@@ -1713,7 +1843,7 @@ describe('Velody API wiring (e2e)', () => {
|
||||
expect(duplicatePrepare.status).toBe('exists');
|
||||
expect(duplicatePrepare.uploadId).toBeDefined();
|
||||
expect(prismaState.audioAssets.size).toBe(1);
|
||||
expect(prismaState.libraryEvents.size).toBe(1);
|
||||
expect(prismaState.libraryEvents.size).toBe(2);
|
||||
});
|
||||
|
||||
it('supports upload finalize with embedded artwork and exposes remote artwork metadata', async () => {
|
||||
|
||||
Reference in New Issue
Block a user