From 9be009d809666e37c0649bdf14c5ba4e32bcfa92 Mon Sep 17 00:00:00 2001 From: 4gray Date: Thu, 30 Jul 2026 01:29:52 +0200 Subject: [PATCH] fix(portals): serialize Xtream restore revisions --- docs/architecture/vod-multi-source.md | 9 +- .../lib/playlist-refresh-action.service.ts | 14 +- .../recent-playlists.component.ts | 14 +- .../with-content.feature.spec-helpers.ts | 70 ++++++ .../features/with-content.feature.spec.ts | 31 ++- .../stores/features/with-content.feature.ts | 74 +++--- .../lib/playlist-backup.service.pins.spec.ts | 27 ++ .../playlist-backup.service.test-helpers.ts | 63 ++++- .../src/lib/playlist-backup.service.ts | 20 +- ...list-backup.service.xtream-restore.spec.ts | 13 + .../lib/playlist-backup.xtream-fixtures.ts | 60 ++++- .../xtream-pending-restore.service.spec.ts | 235 +++++++++++++++++- .../src/lib/xtream-pending-restore.service.ts | 181 ++++++++++++-- 13 files changed, 732 insertions(+), 79 deletions(-) diff --git a/docs/architecture/vod-multi-source.md b/docs/architecture/vod-multi-source.md index a95ca5974..16c205a26 100644 --- a/docs/architecture/vod-multi-source.md +++ b/docs/architecture/vod-multi-source.md @@ -684,8 +684,13 @@ prefix from being visible before the parked state is retried. The parked state is consumed only after that atomic replacement succeeds, and its removal is verified (with an empty tombstone as the safe fallback when storage removal fails). Consumption compares the current parked snapshot with -the one that was applied, and backup imports are serialized so an older restore -cannot finish after and overwrite a newer one. Either restore path reports a +the one that was applied. Store-owned replay and direct backup restore share a +playlist-scoped FIFO coordinator and consume a revision captured for the +content generation they restored. Duplicate consumers of that revision +coalesce, while a newer revision remains parked for its own post-import +consumer; an older asynchronous replacement therefore cannot clear or finish +after a newer one. Parking is verified before an import can report success, and +whole backup imports are serialized as well. Either restore path reports a failed import instead of claiming success when consumption cannot be confirmed. Fresh-import initialization stays blocked and retryable on either failure. Cached content is not exposed while parked state exists: a complete diff --git a/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts b/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts index 214cc380d..dd911f832 100644 --- a/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts +++ b/libs/playlist/shared/ui/src/lib/playlist-refresh-action.service.ts @@ -158,10 +158,16 @@ export class PlaylistRefreshActionService { ] ); - this.pendingRestoreService.set(item._id, { - ...restoreState, - playbackPositions, - }); + if ( + !this.pendingRestoreService.set(item._id, { + ...restoreState, + playbackPositions, + }) + ) { + throw new Error( + `Parking pending restore state for "${item._id}" failed.` + ); + } measureRendererPerformancePhase( RENDERER_PERFORMANCE_PHASE.XTREAM_REFRESH_META, diff --git a/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts b/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts index 45d88f585..5cb08a43b 100644 --- a/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts +++ b/libs/playlist/shared/ui/src/lib/recent-playlists/recent-playlists.component.ts @@ -402,10 +402,16 @@ export class RecentPlaylistsComponent { ] ); - this.pendingRestoreService.set(item._id, { - ...restoreState, - playbackPositions, - }); + if ( + !this.pendingRestoreService.set(item._id, { + ...restoreState, + playbackPositions, + }) + ) { + throw new Error( + `Parking pending restore state for "${item._id}" failed.` + ); + } // Update the timestamp in NgRx / IndexedDB measureRendererPerformancePhase( diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts index 7a43f8677..8f4d2c204 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts @@ -4,6 +4,7 @@ import { RENDERER_PERFORMANCE_PHASE_HOOK_KEY, type RendererPerformancePhaseEvent, } from '@iptvnator/shared/logging'; +import { XtreamPendingRestoreSnapshot } from '@iptvnator/services'; import { XtreamPlaylistData } from '../../data-sources/xtream-data-source.interface'; import { PortalStatusType } from '../../xtream-state'; import { withContent } from './with-content.feature'; @@ -25,6 +26,75 @@ export const PENDING_RESTORE_STATE: XtreamPendingRestoreState = { sourcePins: [{ matchKey: 'tmdb:603', contentId: 501 }], }; +export function createPendingRestoreServiceMock() { + let nextRevision = 0; + let activeSnapshot: XtreamPendingRestoreSnapshot | null = null; + let activeSerializedState: string | null = null; + let lastConsumedRevision: number | null = null; + const mock = { + getOrThrow: jest.fn().mockReturnValue(null), + getSnapshotOrThrow: jest.fn(), + clear: jest.fn().mockReturnValue(true), + applyAndConsume: jest.fn(), + }; + mock.getSnapshotOrThrow.mockImplementation((playlistId: string) => { + const state = mock.getOrThrow( + playlistId + ) as XtreamPendingRestoreState | null; + if (!state) { + activeSnapshot = null; + activeSerializedState = null; + return null; + } + + const serializedState = JSON.stringify(state); + if ( + !activeSnapshot || + activeSerializedState !== serializedState + ) { + activeSnapshot = { + playlistId, + revision: ++nextRevision, + state, + }; + activeSerializedState = serializedState; + lastConsumedRevision = null; + } + return activeSnapshot; + }); + mock.applyAndConsume.mockImplementation( + async ( + playlistId: string, + expectedSnapshot: XtreamPendingRestoreSnapshot, + apply: (state: XtreamPendingRestoreState) => Promise + ) => { + const currentSnapshot = + mock.getSnapshotOrThrow(playlistId); + if (!currentSnapshot) { + return lastConsumedRevision === expectedSnapshot.revision + ? 'consumed' + : 'superseded'; + } + if ( + currentSnapshot.revision !== expectedSnapshot.revision + ) { + return 'superseded'; + } + + await apply(currentSnapshot.state); + if (!mock.clear(playlistId, currentSnapshot.state)) { + return 'consume-failed'; + } + lastConsumedRevision = currentSnapshot.revision; + activeSnapshot = null; + activeSerializedState = null; + mock.getOrThrow.mockReturnValue(null); + return 'consumed'; + } + ); + return mock; +} + export function createContentTestStore( checkPortalStatus: () => Promise ) { diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts index 0e69f4531..8acb45a06 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec.ts @@ -10,6 +10,7 @@ import { createAbortError, createContentTestStore, createDeferred, + createPendingRestoreServiceMock, PENDING_RESTORE_STATE, readPendingRestoreBlocked, setPerformanceHook, @@ -56,10 +57,9 @@ describe('withContent import state', () => { let xtreamApiService: { cancelSession: jest.Mock; }; - let pendingRestoreService: { - getOrThrow: jest.Mock; - clear: jest.Mock; - }; + let pendingRestoreService: ReturnType< + typeof createPendingRestoreServiceMock + >; beforeEach(() => { localStorage.clear(); @@ -94,10 +94,7 @@ describe('withContent import state', () => { xtreamApiService = { cancelSession: jest.fn().mockResolvedValue(true), }; - pendingRestoreService = { - getOrThrow: jest.fn().mockReturnValue(null), - clear: jest.fn().mockReturnValue(true), - }; + pendingRestoreService = createPendingRestoreServiceMock(); checkPortalStatusMock = jest.fn().mockResolvedValue('active'); TestBed.configureTestingModule({ @@ -1069,7 +1066,7 @@ describe('withContent import state', () => { expect(store.contentInitBlockReason()).toBeNull(); }); - it('consumes pending state parked during an active import', async () => { + it('leaves state parked during an active import for the next content generation', async () => { const pendingSeries = createDeferred(); dataSource.getContent.mockImplementation( (_playlistId: string, _credentials: unknown, type: ContentType) => @@ -1092,18 +1089,30 @@ describe('withContent import state', () => { pendingSeries.resolve([]); await initialization; + expect(dataSource.restoreUserData).not.toHaveBeenCalled(); + expect(readPendingRestoreBlocked(store)).toBe(true); + expect(store.isContentInitialized()).toBe(false); + + dataSource.getContent.mockResolvedValue([]); + await store.initializeContent(); + expect(dataSource.restoreUserData).toHaveBeenCalledWith( PLAYLIST.id, PENDING_RESTORE_STATE, expect.any(Object) ); + expect(pendingRestoreService.applyAndConsume).toHaveBeenCalledWith( + PLAYLIST.id, + expect.objectContaining({ + state: PENDING_RESTORE_STATE, + }), + expect.any(Function) + ); expect(pendingRestoreService.clear).toHaveBeenCalledWith( PLAYLIST.id, PENDING_RESTORE_STATE ); expect(readPendingRestoreBlocked(store)).toBe(false); - expect(store.isImporting()).toBe(false); - expect(store.activeImportSessionId()).toBeNull(); expect(store.isContentInitialized()).toBe(true); }); diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts index dfb5bb54b..4dbe1492d 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.ts @@ -855,11 +855,13 @@ export function withContent() { // content. A retry may load each type from the DB without // emitting an import phase, so the store-owned gate is what // prevents source-pin edits until replay is consumed. - const initialRestoreData = pendingRestoreService.getOrThrow( - ctx.playlistId - ); + const initialRestoreSnapshot = + pendingRestoreService.getSnapshotOrThrow( + ctx.playlistId + ); patchState(store, { - isPendingRestoreBlocked: initialRestoreData !== null, + isPendingRestoreBlocked: + initialRestoreSnapshot !== null, }); // Electron content persistence maps remote category IDs @@ -876,19 +878,21 @@ export function withContent() { }); throwIfImportCancelled(importSessionId); - // A backup may be imported while this initialization is - // awaiting network or DB work. Use the latest snapshot; - // restoreUserData applies category visibility too, so a - // late arrival is complete rather than pin-only. - const restoreData = pendingRestoreService.getOrThrow( - ctx.playlistId - ); + // Only the revision captured before import belongs to this + // content generation. A newer revision may come from a + // refresh that has deleted the catalog and must remain + // parked until its own replacement import completes. + const currentRestoreSnapshot = + pendingRestoreService.getSnapshotOrThrow( + ctx.playlistId + ); patchState(store, { - isPendingRestoreBlocked: restoreData !== null, + isPendingRestoreBlocked: + currentRestoreSnapshot !== null, }); // Restore user data if needed - if (restoreData) { + if (initialRestoreSnapshot) { try { throwIfImportCancelled(importSessionId); const restoreOperationId = @@ -899,28 +903,29 @@ export function withContent() { patchState(store, { importPhase: 'restoring-favorites', }); - await dataSource.restoreUserData( - ctx.playlistId, - restoreData, - { - onEvent: trackImportEvent, - operationId: restoreOperationId, - } - ); - throwIfImportCancelled(importSessionId); - if ( - !pendingRestoreService.clear( + const restoreResult = + await pendingRestoreService.applyAndConsume( ctx.playlistId, - restoreData - ) - ) { + initialRestoreSnapshot, + async (pendingState) => { + throwIfImportCancelled(importSessionId); + await dataSource.restoreUserData( + ctx.playlistId, + pendingState, + { + onEvent: trackImportEvent, + operationId: + restoreOperationId, + } + ); + throwIfImportCancelled(importSessionId); + } + ); + if (restoreResult === 'consume-failed') { throw new Error( `Clearing pending restore state for "${ctx.playlistId}" failed.` ); } - patchState(store, { - isPendingRestoreBlocked: false, - }); } catch (err) { throwIfImportCancelled(importSessionId); @@ -932,6 +937,15 @@ export function withContent() { } } + const isRestoreStillPending = + hasPendingRestoreOrReadFailure(ctx.playlistId); + patchState(store, { + isPendingRestoreBlocked: isRestoreStillPending, + }); + if (isRestoreStillPending) { + return; + } + throwIfImportCancelled(importSessionId); // Mark as initialized so next routings won't re-trigger it diff --git a/libs/services/src/lib/playlist-backup.service.pins.spec.ts b/libs/services/src/lib/playlist-backup.service.pins.spec.ts index d97cd630c..ea7c7efca 100644 --- a/libs/services/src/lib/playlist-backup.service.pins.spec.ts +++ b/libs/services/src/lib/playlist-backup.service.pins.spec.ts @@ -217,6 +217,33 @@ describe('PlaylistBackupService Xtream source pins', () => { expect(summary.errors[0]).toMatch(/pending restore state/i); }); + it('reports a restore whose pending state could not be parked', async () => { + const collaborators = createRestoreCollaborators(); + collaborators.pendingRestoreService.set.mockReturnValueOnce(null); + const replacePins = jest.fn().mockResolvedValue(true); + const service = createPlaylistBackupService({ + ...collaborators, + vodSourcePinService: { + isAvailable: true, + listForPlaylistOrThrow: jest.fn().mockResolvedValue([]), + replaceForPlaylist: replacePins, + }, + }); + + const manifest = createXtreamManifest( + [], + [{ matchKey: 'tmdb:603', contentId: 501 }] + ); + + const summary = await service.importBackup(JSON.stringify(manifest)); + + expect(summary).toEqual( + expect.objectContaining({ merged: 0, failed: 1 }) + ); + expect(summary.errors[0]).toMatch(/parking pending restore state/i); + expect(replacePins).not.toHaveBeenCalled(); + }); + it('drops pins the backup does not contain', async () => { const collaborators = createRestoreCollaborators(); const clearPins = jest.fn().mockResolvedValue(true); diff --git a/libs/services/src/lib/playlist-backup.service.test-helpers.ts b/libs/services/src/lib/playlist-backup.service.test-helpers.ts index c6fb18849..684380ab6 100644 --- a/libs/services/src/lib/playlist-backup.service.test-helpers.ts +++ b/libs/services/src/lib/playlist-backup.service.test-helpers.ts @@ -5,9 +5,13 @@ import { VodSourcePin, XtreamBackupFavoriteItem, XtreamBackupRecentlyViewedItem, + XtreamPendingRestoreState, } from '@iptvnator/shared/interfaces'; import { PlaylistBackupService } from './playlist-backup.service'; -import { XtreamPendingRestoreService } from './xtream-pending-restore.service'; +import { + XtreamPendingRestoreService, + XtreamPendingRestoreSnapshot, +} from './xtream-pending-restore.service'; /** * Shared factory for PlaylistBackupService specs. Instantiates the service @@ -20,6 +24,58 @@ export function createPlaylistBackupService( const service = Object.create( PlaylistBackupService.prototype ) as PlaylistBackupService; + let nextRevision = 0; + let pendingSnapshot: XtreamPendingRestoreSnapshot | null = null; + let lastConsumedRevision: number | null = null; + const pendingRestoreService = { + set: jest.fn( + ( + playlistId: string, + state: XtreamPendingRestoreState + ): XtreamPendingRestoreSnapshot | null => { + pendingSnapshot = { + playlistId, + revision: ++nextRevision, + state, + }; + lastConsumedRevision = null; + return pendingSnapshot; + } + ), + clear: jest.fn().mockReturnValue(true), + applyAndConsume: jest.fn(), + }; + pendingRestoreService.applyAndConsume.mockImplementation( + async ( + playlistId: string, + expectedSnapshot: XtreamPendingRestoreSnapshot, + apply: (state: XtreamPendingRestoreState) => Promise + ) => { + if (!pendingSnapshot) { + return lastConsumedRevision === expectedSnapshot.revision + ? 'consumed' + : 'superseded'; + } + if ( + pendingSnapshot.revision !== expectedSnapshot.revision + ) { + return 'superseded'; + } + + await apply(pendingSnapshot.state); + if ( + !pendingRestoreService.clear( + playlistId, + pendingSnapshot.state + ) + ) { + return 'consume-failed'; + } + lastConsumedRevision = pendingSnapshot.revision; + pendingSnapshot = null; + return 'consumed'; + } + ); Object.assign(service as object, { playlistsService: { @@ -68,10 +124,7 @@ export function createPlaylistBackupService( clear: jest.fn().mockResolvedValue(true), clearForPlaylist: jest.fn().mockResolvedValue(true), }, - pendingRestoreService: { - set: jest.fn(), - clear: jest.fn().mockReturnValue(true), - }, + pendingRestoreService, ...overrides, }); diff --git a/libs/services/src/lib/playlist-backup.service.ts b/libs/services/src/lib/playlist-backup.service.ts index c62b8f3c9..479490003 100644 --- a/libs/services/src/lib/playlist-backup.service.ts +++ b/libs/services/src/lib/playlist-backup.service.ts @@ -841,7 +841,15 @@ export class PlaylistBackupService { const restoreState: XtreamPendingRestoreState = normalizeXtreamPendingRestoreState(entry.userState); - this.pendingRestoreService.set(playlistId, restoreState); + const restoreSnapshot = this.pendingRestoreService.set( + playlistId, + restoreState + ); + if (!restoreSnapshot) { + throw new PlaylistBackupError( + `Parking pending restore state for "${playlistId}" failed.` + ); + } if (!this.hasElectronApi()) { return; @@ -851,8 +859,14 @@ export class PlaylistBackupService { return; } - await this.applyXtreamRestoreState(playlistId, restoreState); - if (!this.pendingRestoreService.clear(playlistId, restoreState)) { + const restoreResult = + await this.pendingRestoreService.applyAndConsume( + playlistId, + restoreSnapshot, + (pendingState) => + this.applyXtreamRestoreState(playlistId, pendingState) + ); + if (restoreResult !== 'consumed') { throw new PlaylistBackupError( `Clearing pending restore state for "${playlistId}" failed.` ); diff --git a/libs/services/src/lib/playlist-backup.service.xtream-restore.spec.ts b/libs/services/src/lib/playlist-backup.service.xtream-restore.spec.ts index ca8ed5383..a708a7566 100644 --- a/libs/services/src/lib/playlist-backup.service.xtream-restore.spec.ts +++ b/libs/services/src/lib/playlist-backup.service.xtream-restore.spec.ts @@ -86,6 +86,19 @@ describe('PlaylistBackupService Xtream hidden categories (issue #1017)', () => { expect( collaborators.databaseService.updateCategoryVisibility ).toHaveBeenCalledTimes(3); + expect( + collaborators.pendingRestoreService.applyAndConsume + ).toHaveBeenCalledWith( + 'xtream-1', + expect.objectContaining({ + state: expect.objectContaining({ + hiddenCategories: [ + { categoryType: 'live', xtreamId: 101 }, + ], + }), + }), + expect.any(Function) + ); expect(collaborators.pendingRestoreService.clear).toHaveBeenCalledWith( 'xtream-1', expect.objectContaining({ diff --git a/libs/services/src/lib/playlist-backup.xtream-fixtures.ts b/libs/services/src/lib/playlist-backup.xtream-fixtures.ts index c374a5ec1..17de2a304 100644 --- a/libs/services/src/lib/playlist-backup.xtream-fixtures.ts +++ b/libs/services/src/lib/playlist-backup.xtream-fixtures.ts @@ -4,8 +4,10 @@ import { PlaylistBackupManifestV1, PLAYLIST_BACKUP_KIND, PLAYLIST_BACKUP_VERSION, + XtreamPendingRestoreState, XtreamPlaylistBackupEntry, } from '@iptvnator/shared/interfaces'; +import { XtreamPendingRestoreSnapshot } from './xtream-pending-restore.service'; /** * The Xtream restore scaffolding both backup suites need — the existing @@ -91,6 +93,59 @@ export function createXtreamManifest( } export function createRestoreCollaborators() { + let nextRevision = 0; + let pendingSnapshot: XtreamPendingRestoreSnapshot | null = null; + let lastConsumedRevision: number | null = null; + const pendingRestoreService = { + set: jest.fn( + ( + playlistId: string, + state: XtreamPendingRestoreState + ): XtreamPendingRestoreSnapshot | null => { + pendingSnapshot = { + playlistId, + revision: ++nextRevision, + state, + }; + lastConsumedRevision = null; + return pendingSnapshot; + } + ), + clear: jest.fn().mockReturnValue(true), + applyAndConsume: jest.fn(), + }; + pendingRestoreService.applyAndConsume.mockImplementation( + async ( + playlistId: string, + expectedSnapshot: XtreamPendingRestoreSnapshot, + apply: (state: XtreamPendingRestoreState) => Promise + ) => { + if (!pendingSnapshot) { + return lastConsumedRevision === expectedSnapshot.revision + ? 'consumed' + : 'superseded'; + } + if ( + pendingSnapshot.revision !== expectedSnapshot.revision + ) { + return 'superseded'; + } + + await apply(pendingSnapshot.state); + if ( + !pendingRestoreService.clear( + playlistId, + pendingSnapshot.state + ) + ) { + return 'consume-failed'; + } + lastConsumedRevision = pendingSnapshot.revision; + pendingSnapshot = null; + return 'consumed'; + } + ); + return { playlistsService: { addPlaylist: jest.fn((playlist: Playlist) => of(playlist)), @@ -111,9 +166,6 @@ export function createRestoreCollaborators() { restoreXtreamUserData: jest.fn().mockResolvedValue(undefined), updateCategoryVisibility: jest.fn().mockResolvedValue(true), }, - pendingRestoreService: { - set: jest.fn(), - clear: jest.fn().mockReturnValue(true), - }, + pendingRestoreService, }; } diff --git a/libs/services/src/lib/xtream-pending-restore.service.spec.ts b/libs/services/src/lib/xtream-pending-restore.service.spec.ts index c0e4738aa..43eefb9e7 100644 --- a/libs/services/src/lib/xtream-pending-restore.service.spec.ts +++ b/libs/services/src/lib/xtream-pending-restore.service.spec.ts @@ -1,5 +1,17 @@ import { getXtreamPendingRestoreStorageKey } from '@iptvnator/shared/interfaces'; -import { XtreamPendingRestoreService } from './xtream-pending-restore.service'; +import { + XtreamPendingRestoreService, + XtreamPendingRestoreSnapshot, +} from './xtream-pending-restore.service'; + +function requireSnapshot( + snapshot: XtreamPendingRestoreSnapshot | null +): XtreamPendingRestoreSnapshot { + if (!snapshot) { + throw new Error('Expected pending restore snapshot'); + } + return snapshot; +} describe('XtreamPendingRestoreService', () => { const playlistId = 'playlist-1'; @@ -39,7 +51,7 @@ describe('XtreamPendingRestoreService', () => { }); it('normalizes state on write', () => { - service.set(playlistId, { + const snapshot = service.set(playlistId, { hiddenCategories: [ { categoryType: 'live', xtreamId: 101 }, { categoryType: 'live' } as never, @@ -49,6 +61,7 @@ describe('XtreamPendingRestoreService', () => { playbackPositions: [], }); + expect(snapshot).not.toBeNull(); const persisted = JSON.parse( localStorage.getItem(storageKey) ?? 'null' ); @@ -152,4 +165,222 @@ describe('XtreamPendingRestoreService', () => { expect(service.clear(playlistId, restoredState)).toBe(false); expect(service.getOrThrow(playlistId)).toEqual(newerState); }); + + it('leaves a newer snapshot for its own queued consumer', async () => { + const olderState = { + hiddenCategories: [], + favorites: [{ xtreamId: 101, contentType: 'live' as const }], + recentlyViewed: [], + playbackPositions: [], + }; + const newerState = { + ...olderState, + favorites: [{ xtreamId: 202, contentType: 'movie' as const }], + }; + let releaseOlder!: () => void; + const olderCanFinish = new Promise((resolve) => { + releaseOlder = resolve; + }); + let signalOlderStarted!: () => void; + const olderStarted = new Promise((resolve) => { + signalOlderStarted = resolve; + }); + const appliedSnapshots: number[] = []; + + const olderSnapshot = requireSnapshot( + service.set(playlistId, olderState) + ); + const olderApplication = service.applyAndConsume( + playlistId, + olderSnapshot, + async (state) => { + const xtreamId = state.favorites[0]?.xtreamId; + signalOlderStarted(); + await olderCanFinish; + appliedSnapshots.push(xtreamId); + } + ); + await olderStarted; + + const newerSnapshot = requireSnapshot( + service.set(playlistId, newerState) + ); + const newerApply = jest.fn(async (state) => { + appliedSnapshots.push(state.favorites[0]?.xtreamId); + }); + const newerApplication = service.applyAndConsume( + playlistId, + newerSnapshot, + newerApply + ); + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(newerApply).not.toHaveBeenCalled(); + + releaseOlder(); + await expect(olderApplication).resolves.toBe('superseded'); + await expect(newerApplication).resolves.toBe('consumed'); + expect(appliedSnapshots).toEqual([101, 202]); + expect(newerApply).toHaveBeenCalledTimes(1); + expect(service.getOrThrow(playlistId)).toBeNull(); + }); + + it('coalesces concurrent consumers of the same snapshot', async () => { + const state = { + hiddenCategories: [], + favorites: [{ xtreamId: 101, contentType: 'live' as const }], + recentlyViewed: [], + playbackPositions: [], + }; + let releaseFirst!: () => void; + const firstCanFinish = new Promise((resolve) => { + releaseFirst = resolve; + }); + let signalFirstStarted!: () => void; + const firstStarted = new Promise((resolve) => { + signalFirstStarted = resolve; + }); + const firstApply = jest.fn(async () => { + signalFirstStarted(); + await firstCanFinish; + }); + const secondApply = jest.fn().mockResolvedValue(undefined); + + const snapshot = requireSnapshot(service.set(playlistId, state)); + const firstApplication = service.applyAndConsume( + playlistId, + snapshot, + firstApply + ); + await firstStarted; + const secondApplication = service.applyAndConsume( + playlistId, + snapshot, + secondApply + ); + + releaseFirst(); + await expect(firstApplication).resolves.toBe('consumed'); + await expect(secondApplication).resolves.toBe('consumed'); + expect(firstApply).toHaveBeenCalledTimes(1); + expect(secondApply).not.toHaveBeenCalled(); + }); + + it('does not consume an identical snapshot parked by a newer producer', async () => { + const state = { + hiddenCategories: [], + favorites: [{ xtreamId: 101, contentType: 'live' as const }], + recentlyViewed: [], + playbackPositions: [], + }; + let releaseOlder!: () => void; + const olderCanFinish = new Promise((resolve) => { + releaseOlder = resolve; + }); + let signalOlderStarted!: () => void; + const olderStarted = new Promise((resolve) => { + signalOlderStarted = resolve; + }); + const olderSnapshot = requireSnapshot( + service.set(playlistId, state) + ); + const olderApplication = service.applyAndConsume( + playlistId, + olderSnapshot, + async () => { + signalOlderStarted(); + await olderCanFinish; + } + ); + await olderStarted; + + const newerSnapshot = requireSnapshot( + service.set(playlistId, state) + ); + expect(newerSnapshot.revision).not.toBe(olderSnapshot.revision); + releaseOlder(); + + await expect( + olderApplication + ).resolves.toBe('superseded'); + expect(service.getSnapshotOrThrow(playlistId)).toEqual(newerSnapshot); + }); + + it('continues the playlist queue after an application rejects', async () => { + const olderState = { + hiddenCategories: [], + favorites: [{ xtreamId: 101, contentType: 'live' as const }], + recentlyViewed: [], + playbackPositions: [], + }; + const newerState = { + ...olderState, + favorites: [{ xtreamId: 202, contentType: 'movie' as const }], + }; + const olderSnapshot = requireSnapshot( + service.set(playlistId, olderState) + ); + + await expect( + service.applyAndConsume( + playlistId, + olderSnapshot, + async () => { + throw new Error('database is locked'); + } + ) + ).rejects.toThrow('database is locked'); + const newerSnapshot = requireSnapshot( + service.set(playlistId, newerState) + ); + const applyNewer = jest.fn().mockResolvedValue(undefined); + + await expect( + service.applyAndConsume( + playlistId, + newerSnapshot, + applyNewer + ) + ).resolves.toBe('consumed'); + expect(applyNewer).toHaveBeenCalledTimes(1); + }); + + it('reports failure when an applied snapshot cannot be consumed', async () => { + const state = { + hiddenCategories: [], + favorites: [{ xtreamId: 101, contentType: 'live' as const }], + recentlyViewed: [], + playbackPositions: [], + }; + const snapshot = requireSnapshot(service.set(playlistId, state)); + jest.spyOn(Storage.prototype, 'setItem').mockImplementation(() => { + throw new Error('storage is locked'); + }); + jest.spyOn(Storage.prototype, 'removeItem').mockImplementation(() => { + throw new Error('storage is locked'); + }); + const apply = jest.fn().mockResolvedValue(undefined); + + await expect( + service.applyAndConsume(playlistId, snapshot, apply) + ).resolves.toBe('consume-failed'); + expect(apply).toHaveBeenCalledWith(state); + expect(service.getOrThrow(playlistId)).toEqual(state); + }); + + it('reports failure when pending state cannot be parked', () => { + jest.spyOn(Storage.prototype, 'setItem').mockImplementation(() => { + throw new Error('storage is locked'); + }); + + expect( + service.set(playlistId, { + hiddenCategories: [], + favorites: [], + recentlyViewed: [], + playbackPositions: [], + }) + ).toBeNull(); + expect(service.getOrThrow(playlistId)).toBeNull(); + }); }); diff --git a/libs/services/src/lib/xtream-pending-restore.service.ts b/libs/services/src/lib/xtream-pending-restore.service.ts index 277e4d1b0..ca568c9df 100644 --- a/libs/services/src/lib/xtream-pending-restore.service.ts +++ b/libs/services/src/lib/xtream-pending-restore.service.ts @@ -5,10 +5,32 @@ import { XtreamPendingRestoreState, } from '@iptvnator/shared/interfaces'; +export interface XtreamPendingRestoreSnapshot { + readonly playlistId: string; + readonly revision: number; + readonly state: XtreamPendingRestoreState; +} + +export type XtreamPendingRestoreApplicationResult = + | 'consumed' + | 'superseded' + | 'consume-failed'; + @Injectable({ providedIn: 'root', }) export class XtreamPendingRestoreService { + private readonly restoreApplicationTails = new Map< + string, + Promise + >(); + private readonly restoreRevisions = new Map< + string, + { revision: number; serializedState: string } + >(); + private readonly lastConsumedRevisions = new Map(); + private nextRestoreRevision = 0; + get(playlistId: string): XtreamPendingRestoreState | null { try { return this.getOrThrow(playlistId); @@ -44,19 +66,119 @@ export class XtreamPendingRestoreService { } } - set(playlistId: string, state: XtreamPendingRestoreState): void { - if (!playlistId) { - return; + getSnapshotOrThrow( + playlistId: string + ): XtreamPendingRestoreSnapshot | null { + const state = this.getOrThrow(playlistId); + if (!state) { + this.restoreRevisions.delete(playlistId); + return null; } - try { - localStorage.setItem( - getXtreamPendingRestoreStorageKey(playlistId), - JSON.stringify(normalizeXtreamPendingRestoreState(state)) - ); - } catch { - // Ignore local storage write failures. + const serializedState = JSON.stringify(state); + let revision = this.restoreRevisions.get(playlistId); + if (!revision || revision.serializedState !== serializedState) { + revision = { + revision: ++this.nextRestoreRevision, + serializedState, + }; + this.restoreRevisions.set(playlistId, revision); + this.lastConsumedRevisions.delete(playlistId); } + + return { + playlistId, + revision: revision.revision, + state, + }; + } + + set( + playlistId: string, + state: XtreamPendingRestoreState + ): XtreamPendingRestoreSnapshot | null { + if (!playlistId) { + return null; + } + + const normalizedState = normalizeXtreamPendingRestoreState(state); + const serializedState = JSON.stringify(normalizedState); + + try { + const storageKey = getXtreamPendingRestoreStorageKey(playlistId); + localStorage.setItem(storageKey, serializedState); + if (localStorage.getItem(storageKey) !== serializedState) { + return null; + } + } catch { + return null; + } + + const revision = ++this.nextRestoreRevision; + this.restoreRevisions.set(playlistId, { + revision, + serializedState, + }); + this.lastConsumedRevisions.delete(playlistId); + return { + playlistId, + revision, + state: normalizedState, + }; + } + + /** + * Serializes state-bound restore consumers for one playlist. Consumers of + * the same revision coalesce after the first succeeds, while a newer + * revision remains parked for the importer whose content generation is + * ready to accept it. + */ + async applyAndConsume( + playlistId: string, + expectedSnapshot: XtreamPendingRestoreSnapshot, + apply: (state: XtreamPendingRestoreState) => Promise + ): Promise { + if ( + !playlistId || + expectedSnapshot.playlistId !== playlistId + ) { + return 'consume-failed'; + } + + return this.enqueueRestoreApplication(playlistId, async () => { + const currentSnapshot = this.getSnapshotOrThrow(playlistId); + if (!currentSnapshot) { + return this.lastConsumedRevisions.get(playlistId) === + expectedSnapshot.revision + ? 'consumed' + : 'superseded'; + } + if ( + currentSnapshot.revision !== expectedSnapshot.revision + ) { + return 'superseded'; + } + + await apply(currentSnapshot.state); + + const snapshotAfterApply = + this.getSnapshotOrThrow(playlistId); + if ( + !snapshotAfterApply || + snapshotAfterApply.revision !== expectedSnapshot.revision + ) { + return 'superseded'; + } + if (!this.clear(playlistId, snapshotAfterApply.state)) { + return 'consume-failed'; + } + + this.lastConsumedRevisions.set( + playlistId, + expectedSnapshot.revision + ); + return 'consumed'; + }); } clear( @@ -77,10 +199,7 @@ export class XtreamPendingRestoreService { if ( !currentState || - JSON.stringify(currentState) !== - JSON.stringify( - normalizeXtreamPendingRestoreState(expectedState) - ) + !this.restoreStatesEqual(currentState, expectedState) ) { return false; } @@ -102,12 +221,46 @@ export class XtreamPendingRestoreService { try { localStorage.removeItem(storageKey); if (localStorage.getItem(storageKey) === null) { + this.restoreRevisions.delete(playlistId); return true; } } catch { // A verified tombstone is sufficient even if removal fails. } + if (tombstoneStored) { + this.restoreRevisions.delete(playlistId); + } return tombstoneStored; } + + private restoreStatesEqual( + left: XtreamPendingRestoreState, + right: XtreamPendingRestoreState + ): boolean { + return ( + JSON.stringify(normalizeXtreamPendingRestoreState(left)) === + JSON.stringify(normalizeXtreamPendingRestoreState(right)) + ); + } + + private enqueueRestoreApplication( + playlistId: string, + apply: () => Promise + ): Promise { + const previous = + this.restoreApplicationTails.get(playlistId) ?? Promise.resolve(); + const result = previous.then(apply); + const tail = result.then( + () => undefined, + () => undefined + ); + this.restoreApplicationTails.set(playlistId, tail); + void tail.then(() => { + if (this.restoreApplicationTails.get(playlistId) === tail) { + this.restoreApplicationTails.delete(playlistId); + } + }); + return result; + } }