diff --git a/.changes/portals-atomic-fresh-import-pins.md b/.changes/portals-atomic-fresh-import-pins.md index 2c9fb55e0..9ba031ca1 100644 --- a/.changes/portals-atomic-fresh-import-pins.md +++ b/.changes/portals-atomic-fresh-import-pins.md @@ -3,5 +3,7 @@ type: fix area: portals --- -Importing an Xtream backup now restores all preferred VOD sources together, so -a database failure cannot leave only some source preferences applied. +Importing an Xtream backup now restores all preferred VOD sources together. If +restoration or cleanup fails, the catalog stays retryable and remains blocked +until the parked state is safely consumed, preventing stale replay from +overwriting newer choices. diff --git a/apps/web/src/app/settings/settings-backup.facade.spec.ts b/apps/web/src/app/settings/settings-backup.facade.spec.ts index e4e70ea9d..a14201475 100644 --- a/apps/web/src/app/settings/settings-backup.facade.spec.ts +++ b/apps/web/src/app/settings/settings-backup.facade.spec.ts @@ -4,6 +4,7 @@ import { selectAllPlaylistsMeta, selectIsEpgAvailable, } from '@iptvnator/m3u-state'; +import { XtreamStore } from '@iptvnator/portal/xtream/data-access'; import { PlaylistBackupService } from '@iptvnator/services'; import { provideMockStore } from '@ngrx/store/testing'; import { TranslateModule } from '@ngx-translate/core'; @@ -36,6 +37,9 @@ describe('SettingsBackupFacade', () => { .mockResolvedValue(BACKUP_EXPORT_RESULT), importBackup: jest.fn(), }), + MockProvider(XtreamStore, { + reconcilePendingRestoreBlock: jest.fn(), + }), provideMockStore({ selectors: [ { selector: selectAllPlaylistsMeta, value: [] }, @@ -154,4 +158,46 @@ describe('SettingsBackupFacade', () => { expect.objectContaining({ panelClass: ['settings-snackbar'] }) ); }); + + it('reconciles a pending Xtream restore before completing an import', async () => { + configure(); + const input = document.createElement('input'); + const file = { + text: jest.fn().mockResolvedValue('{}'), + } as unknown as File; + Object.defineProperty(input, 'files', { value: [file] }); + const addEventListener = jest.spyOn(input, 'addEventListener'); + const createElement = jest + .spyOn(document, 'createElement') + .mockReturnValue(input); + jest.spyOn(input, 'click').mockImplementation(); + const summary = { + imported: 0, + merged: 0, + skipped: 0, + failed: 1, + errors: ['pending restore state could not be consumed'], + }; + (playlistBackupService.importBackup as jest.Mock).mockResolvedValue( + summary + ); + const onImported = jest.fn(); + const xtreamStore = TestBed.inject(XtreamStore); + + facade.importData(onImported); + const changeListener = addEventListener.mock.calls.find( + ([type]) => type === 'change' + )?.[1] as (event: Event) => Promise; + await changeListener({ target: input } as unknown as Event); + + expect(xtreamStore.reconcilePendingRestoreBlock).toHaveBeenCalledTimes( + 1 + ); + expect(onImported).toHaveBeenCalledTimes(1); + expect( + (xtreamStore.reconcilePendingRestoreBlock as jest.Mock).mock + .invocationCallOrder[0] + ).toBeLessThan(onImported.mock.invocationCallOrder[0]); + createElement.mockRestore(); + }); }); diff --git a/apps/web/src/app/settings/settings-backup.facade.ts b/apps/web/src/app/settings/settings-backup.facade.ts index 6aa377b4c..bebfc6004 100644 --- a/apps/web/src/app/settings/settings-backup.facade.ts +++ b/apps/web/src/app/settings/settings-backup.facade.ts @@ -1,7 +1,8 @@ -import { inject, Injectable, signal } from '@angular/core'; +import { inject, Injectable, Injector, signal } from '@angular/core'; import { Store } from '@ngrx/store'; import { TranslateService } from '@ngx-translate/core'; import { PlaylistActions } from '@iptvnator/m3u-state'; +import { XtreamStore } from '@iptvnator/portal/xtream/data-access'; import { PlaylistBackupImportSummary, PlaylistBackupService, @@ -16,6 +17,7 @@ export class SettingsBackupFacade { private readonly settingsSnackbar = inject(SettingsSnackbarService); private readonly store = inject(Store); private readonly translate = inject(TranslateService); + private readonly injector = inject(Injector); readonly isExportingData = signal(false); @@ -79,6 +81,9 @@ export class SettingsBackupFacade { const summary = await this.playlistBackupService.importBackup( await file.text() ); + this.injector + .get(XtreamStore, null) + ?.reconcilePendingRestoreBlock(); if (summary.imported > 0 || summary.merged > 0) { this.store.dispatch(PlaylistActions.removeAllPlaylists()); diff --git a/docs/architecture/vod-multi-source.md b/docs/architecture/vod-multi-source.md index 63e0f12e7..a95ca5974 100644 --- a/docs/architecture/vod-multi-source.md +++ b/docs/architecture/vod-multi-source.md @@ -681,6 +681,22 @@ direct restore, that preserves the existing pins if an insert fails. A fresh import has no pre-existing pins to lose; there it prevents a partially applied 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 +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 +active initialization must restore and retire the snapshot before the catalog +opens, because a scope-limited offline cache may not contain every identity +referenced by the backup. Route bootstrap and Settings import completion +reconcile that gate before content can be edited. The active content route also +stays unmounted while replay is pending, and the import overlay follows the +full initialization session (including cache-only retries), not only +remote-download events. + ## Which engines can fail over Only the built-in web players (HTML5, Video.js, ArtPlayer) raise the playback diff --git a/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.test-helpers.ts b/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.test-helpers.ts index a39f632d4..5a1005632 100644 --- a/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.test-helpers.ts +++ b/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.test-helpers.ts @@ -33,7 +33,7 @@ export function createDbServiceMock() { getXtreamCategories: jest.fn().mockResolvedValue([]), saveXtreamCategories: jest.fn().mockResolvedValue(undefined), getAllXtreamCategories: jest.fn().mockResolvedValue([]), - updateCategoryVisibility: jest.fn().mockResolvedValue(undefined), + updateCategoryVisibility: jest.fn().mockResolvedValue(true), hasXtreamContent: jest.fn().mockResolvedValue(false), getXtreamContent: jest.fn().mockResolvedValue([]), saveXtreamContent: jest.fn().mockResolvedValue(0), diff --git a/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.ts b/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.ts index 9ee72e4bf..7dac4bdcc 100644 --- a/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.ts +++ b/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.ts @@ -566,6 +566,48 @@ export class ElectronXtreamDataSource implements IXtreamDataSource { restoreState: XtreamPendingRestoreState, options?: XtreamOperationOptions ): Promise { + const categoriesByType = await Promise.all([ + this.dbService.getAllXtreamCategories(playlistId, 'live'), + this.dbService.getAllXtreamCategories(playlistId, 'movies'), + this.dbService.getAllXtreamCategories(playlistId, 'series'), + ]); + for (const categories of categoriesByType) { + if (categories.length === 0) { + continue; + } + + const reset = await this.dbService.updateCategoryVisibility( + categories.map((category) => category.id), + false + ); + if (!reset) { + throw new Error( + `Resetting category visibility for "${playlistId}" failed.` + ); + } + + const hiddenCategoryIds = categories + .filter((category) => + restoreState.hiddenCategories.some( + (hiddenCategory) => + hiddenCategory.categoryType === category.type && + hiddenCategory.xtreamId === category.xtream_id + ) + ) + .map((category) => category.id); + if ( + hiddenCategoryIds.length > 0 && + !(await this.dbService.updateCategoryVisibility( + hiddenCategoryIds, + true + )) + ) { + throw new Error( + `Restoring category visibility for "${playlistId}" failed.` + ); + } + } + await this.dbService.restoreXtreamUserData( playlistId, restoreState.favorites, diff --git a/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.user-data.spec.ts b/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.user-data.spec.ts index efd317d2d..bddcda14d 100644 --- a/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.user-data.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/data-sources/electron-xtream-data-source.user-data.spec.ts @@ -258,12 +258,31 @@ describe('ElectronXtreamDataSource (user data delegation)', () => { const positionA = { contentXtreamId: 1 } as never; const positionB = { contentXtreamId: 2 } as never; const restoreState = { - hiddenCategories: [], + hiddenCategories: [{ categoryType: 'live', xtreamId: 101 }], favorites: [{ xtreamId: 202, type: 'movie' }], recentlyViewed: [{ xtreamId: 101, type: 'live' }], playbackPositions: [positionA, positionB], } as never; const options = { operationId: 'op-1' }; + harness.dbService.getAllXtreamCategories.mockImplementation( + (_playlistId: string, type: string) => + Promise.resolve( + type === 'live' + ? [ + { + id: 11, + type: 'live', + xtream_id: 101, + }, + { + id: 12, + type: 'live', + xtream_id: 102, + }, + ] + : [] + ) + ); await harness.dataSource.restoreUserData( playlistId, @@ -279,6 +298,12 @@ describe('ElectronXtreamDataSource (user data delegation)', () => { [{ xtreamId: 101, type: 'live' }], options ); + expect( + harness.dbService.updateCategoryVisibility.mock.calls + ).toEqual([ + [[11, 12], false], + [[11], true], + ]); expect( harness.playbackService.clearAllPlaybackPositions ).toHaveBeenCalledWith(playlistId); 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 new file mode 100644 index 000000000..7a43f8677 --- /dev/null +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-content.feature.spec-helpers.ts @@ -0,0 +1,101 @@ +import { patchState, signalStore, withMethods, withState } from '@ngrx/signals'; +import { XtreamPendingRestoreState } from '@iptvnator/shared/interfaces'; +import { + RENDERER_PERFORMANCE_PHASE_HOOK_KEY, + type RendererPerformancePhaseEvent, +} from '@iptvnator/shared/logging'; +import { XtreamPlaylistData } from '../../data-sources/xtream-data-source.interface'; +import { PortalStatusType } from '../../xtream-state'; +import { withContent } from './with-content.feature'; + +export const TEST_PLAYLIST: XtreamPlaylistData = { + id: 'playlist-1', + name: 'Test Xtream', + serverUrl: 'http://localhost:8080', + username: 'demo', + password: 'secret', + type: 'xtream', +}; + +export const PENDING_RESTORE_STATE: XtreamPendingRestoreState = { + hiddenCategories: [], + favorites: [{ contentType: 'live', xtreamId: 77 }], + recentlyViewed: [], + playbackPositions: [], + sourcePins: [{ matchKey: 'tmdb:603', contentId: 501 }], +}; + +export function createContentTestStore( + checkPortalStatus: () => Promise +) { + return signalStore( + withState({ + playlistId: TEST_PLAYLIST.id, + currentPlaylist: TEST_PLAYLIST, + portalStatus: 'active' as PortalStatusType, + selectedContentType: 'vod' as const, + }), + withMethods((store) => ({ + async checkPortalStatus(): Promise { + const status = await checkPortalStatus(); + patchState(store, { portalStatus: status }); + return status; + }, + })), + withContent() + ); +} + +const performanceHookSymbol = Symbol.for(RENDERER_PERFORMANCE_PHASE_HOOK_KEY); + +export function setPerformanceHook( + hook: ((event: RendererPerformancePhaseEvent) => void) | null +): void { + const target = globalThis as unknown as Record; + if (hook === null) { + delete target[performanceHookSymbol]; + } else { + target[performanceHookSymbol] = hook; + } +} + +export function readPendingRestoreBlocked(store: unknown): boolean { + return ( + store as { + isPendingRestoreBlocked: () => boolean; + } + ).isPendingRestoreBlocked(); +} + +export function createDeferred() { + let resolve!: (value: T) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + + return { promise, resolve, reject }; +} + +export function createAbortError(): Error { + const error = new Error('Import cancelled'); + error.name = 'AbortError'; + return error; +} + +export async function waitForCondition( + predicate: () => boolean, + attempts = 20 +): Promise { + for (let index = 0; index < attempts; index += 1) { + if (predicate()) { + return; + } + + await Promise.resolve(); + await new Promise((resolve) => setTimeout(resolve, 0)); + } + + throw new Error('Timed out waiting for test condition'); +} 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 d210d259e..0e69f4531 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 @@ -1,17 +1,21 @@ import { TestBed } from '@angular/core/testing'; -import { patchState, signalStore, withMethods, withState } from '@ngrx/signals'; -import { DatabaseService } from '@iptvnator/services'; import { - RENDERER_PERFORMANCE_PHASE_HOOK_KEY, - type RendererPerformancePhaseEvent, -} from '@iptvnator/shared/logging'; -import { - XTREAM_DATA_SOURCE, - XtreamPlaylistData, -} from '../../data-sources/xtream-data-source.interface'; + DatabaseService, + XtreamPendingRestoreService, +} from '@iptvnator/services'; +import { XTREAM_DATA_SOURCE } from '../../data-sources/xtream-data-source.interface'; import { XtreamApiService } from '../../services/xtream-api.service'; import { PortalStatusType } from '../../xtream-state'; -import { withContent } from './with-content.feature'; +import { + createAbortError, + createContentTestStore, + createDeferred, + PENDING_RESTORE_STATE, + readPendingRestoreBlocked, + setPerformanceHook, + TEST_PLAYLIST, + waitForCondition, +} from './with-content.feature.spec-helpers'; jest.mock('@iptvnator/portal/shared/util', () => ({ createLogger: () => ({ @@ -24,78 +28,11 @@ jest.mock('@iptvnator/portal/shared/util', () => ({ type ContentType = 'live' | 'movie' | 'series'; -const PLAYLIST: XtreamPlaylistData = { - id: 'playlist-1', - name: 'Test Xtream', - serverUrl: 'http://localhost:8080', - username: 'demo', - password: 'secret', - type: 'xtream', -}; -const performanceHookSymbol = Symbol.for(RENDERER_PERFORMANCE_PHASE_HOOK_KEY); - -function setPerformanceHook( - hook: ((event: RendererPerformancePhaseEvent) => void) | null -): void { - const target = globalThis as unknown as Record; - if (hook === null) { - delete target[performanceHookSymbol]; - } else { - target[performanceHookSymbol] = hook; - } -} +const PLAYLIST = TEST_PLAYLIST; let checkPortalStatusMock: jest.Mock, []>; -const TestContentStore = signalStore( - withState({ - playlistId: PLAYLIST.id, - currentPlaylist: PLAYLIST, - portalStatus: 'active' as PortalStatusType, - selectedContentType: 'vod' as const, - }), - withMethods((store) => ({ - async checkPortalStatus(): Promise { - const status = await checkPortalStatusMock(); - patchState(store, { portalStatus: status }); - return status; - }, - })), - withContent() -); - -function createDeferred() { - let resolve!: (value: T) => void; - let reject!: (reason?: unknown) => void; - const promise = new Promise((res, rej) => { - resolve = res; - reject = rej; - }); - - return { promise, resolve, reject }; -} - -function createAbortError(): Error { - const error = new Error('Import cancelled'); - error.name = 'AbortError'; - return error; -} - -async function waitForCondition( - predicate: () => boolean, - attempts = 20 -): Promise { - for (let index = 0; index < attempts; index += 1) { - if (predicate()) { - return; - } - - await Promise.resolve(); - await new Promise((resolve) => setTimeout(resolve, 0)); - } - - throw new Error('Timed out waiting for test condition'); -} +const TestContentStore = createContentTestStore(() => checkPortalStatusMock()); describe('withContent import state', () => { let store: InstanceType; @@ -119,6 +56,10 @@ describe('withContent import state', () => { let xtreamApiService: { cancelSession: jest.Mock; }; + let pendingRestoreService: { + getOrThrow: jest.Mock; + clear: jest.Mock; + }; beforeEach(() => { localStorage.clear(); @@ -153,6 +94,10 @@ describe('withContent import state', () => { xtreamApiService = { cancelSession: jest.fn().mockResolvedValue(true), }; + pendingRestoreService = { + getOrThrow: jest.fn().mockReturnValue(null), + clear: jest.fn().mockReturnValue(true), + }; checkPortalStatusMock = jest.fn().mockResolvedValue('active'); TestBed.configureTestingModule({ @@ -170,6 +115,10 @@ describe('withContent import state', () => { provide: XtreamApiService, useValue: xtreamApiService, }, + { + provide: XtreamPendingRestoreService, + useValue: pendingRestoreService, + }, ], }); @@ -614,6 +563,63 @@ describe('withContent import state', () => { expect(databaseService.getXtreamImportStatus).not.toHaveBeenCalled(); }); + it('does not expose offline cache while a restore is pending', async () => { + pendingRestoreService.getOrThrow.mockReturnValue(PENDING_RESTORE_STATE); + dataSource.hasCategories.mockResolvedValue(true); + dataSource.hasContent.mockResolvedValue(true); + + await expect(store.hasUsableOfflineCache('vod')).resolves.toBe(false); + await store.hydrateCachedContent('vod'); + + expect(dataSource.hasCategories).not.toHaveBeenCalled(); + expect(dataSource.hasContent).not.toHaveBeenCalled(); + expect(dataSource.getCachedCategories).not.toHaveBeenCalled(); + expect(dataSource.getCachedContent).not.toHaveBeenCalled(); + expect(dataSource.restoreUserData).not.toHaveBeenCalled(); + expect(store.isContentInitialized()).toBe(false); + expect(store.contentInitBlockReason()).toBe('error'); + }); + + it('blocks active initialization when pending restore storage is unreadable', async () => { + dataSource.getContent.mockResolvedValue([]); + pendingRestoreService.getOrThrow.mockImplementation(() => { + throw new Error('storage is locked'); + }); + + await store.initializeContent(); + + expect(dataSource.restoreUserData).not.toHaveBeenCalled(); + expect(store.isContentInitialized()).toBe(false); + expect(store.contentInitBlockReason()).toBe('error'); + }); + + it('fails closed when pending restore storage is unreadable', async () => { + pendingRestoreService.getOrThrow.mockImplementation(() => { + throw new Error('storage is locked'); + }); + dataSource.hasCategories.mockResolvedValue(true); + dataSource.hasContent.mockResolvedValue(true); + + await expect(store.hasUsableOfflineCache('vod')).resolves.toBe(false); + await store.hydrateCachedContent('vod'); + + expect(dataSource.hasCategories).not.toHaveBeenCalled(); + expect(dataSource.hasContent).not.toHaveBeenCalled(); + expect(dataSource.getCachedContent).not.toHaveBeenCalled(); + expect(store.isContentInitialized()).toBe(false); + expect(store.contentInitBlockReason()).toBe('error'); + + pendingRestoreService.getOrThrow.mockReturnValue(null); + checkPortalStatusMock.mockResolvedValue('unavailable'); + dataSource.hasContent.mockResolvedValue(true); + await store.retryContentInitialization(); + + expect(readPendingRestoreBlocked(store)).toBe(false); + expect(dataSource.getCachedContent).toHaveBeenCalled(); + expect(store.isContentInitialized()).toBe(true); + expect(store.contentInitBlockReason()).toBeNull(); + }); + it('does not require every content type for aggregate cached sections', async () => { dataSource.hasContent .mockResolvedValueOnce(false) @@ -738,6 +744,33 @@ describe('withContent import state', () => { expect(store.vodStreams()).toHaveLength(1); }); + it('blocks cache publishing when pending state appears during hydration', async () => { + const cachedCategories = createDeferred(); + const cachedContent = createDeferred(); + dataSource.getCachedCategories.mockReturnValueOnce( + cachedCategories.promise + ); + dataSource.getCachedContent.mockReturnValueOnce(cachedContent.promise); + + const hydration = store.hydrateCachedContent('vod'); + await waitForCondition( + () => dataSource.getCachedContent.mock.calls.length === 1 + ); + + pendingRestoreService.getOrThrow.mockReturnValue(PENDING_RESTORE_STATE); + cachedCategories.resolve([]); + cachedContent.resolve([]); + + await hydration; + + expect(dataSource.restoreUserData).not.toHaveBeenCalled(); + expect(pendingRestoreService.clear).not.toHaveBeenCalled(); + expect(store.vodStreams()).toEqual([]); + expect(store.isContentInitialized()).toBe(false); + expect(store.isCachedContentScopeReady('vod')).toBe(false); + expect(store.contentInitBlockReason()).toBe('error'); + }); + it('coalesces concurrent cached hydration calls for the same scope', async () => { const cachedCategories = createDeferred(); const cachedContent = createDeferred(); @@ -991,6 +1024,133 @@ describe('withContent import state', () => { expect(dataSource.getContent).toHaveBeenCalledTimes(3); }); + it('keeps a failed pending restore blocked until an explicit retry succeeds', async () => { + dataSource.getContent.mockResolvedValue([]); + dataSource.restoreUserData + .mockRejectedValueOnce(new Error('database is locked')) + .mockResolvedValueOnce(undefined); + pendingRestoreService.getOrThrow.mockReturnValue(PENDING_RESTORE_STATE); + pendingRestoreService.clear.mockImplementation(() => { + pendingRestoreService.getOrThrow.mockReturnValue(null); + return true; + }); + + await store.initializeContent(); + + expect(dataSource.restoreUserData).toHaveBeenCalledTimes(1); + expect(pendingRestoreService.clear).not.toHaveBeenCalled(); + expect(store.isContentInitialized()).toBe(false); + expect(store.contentInitBlockReason()).toBe('error'); + + await store.initializeContent(); + expect(dataSource.restoreUserData).toHaveBeenCalledTimes(1); + + // Ready cache rows are not safe to expose while replay is pending: + // a later full replacement could otherwise erase pins changed in the + // meantime. + expect(store.isCachedContentScopeReady('vod')).toBe(false); + checkPortalStatusMock.mockResolvedValue('unavailable'); + dataSource.hasContent.mockResolvedValue(true); + await store.retryContentInitialization(); + + expect(dataSource.restoreUserData).toHaveBeenCalledTimes(1); + expect(store.isContentInitialized()).toBe(false); + expect(store.contentInitBlockReason()).toBe('unavailable'); + + checkPortalStatusMock.mockResolvedValue('active'); + await store.retryContentInitialization(); + + expect(dataSource.restoreUserData).toHaveBeenCalledTimes(2); + expect(pendingRestoreService.clear).toHaveBeenCalledWith( + PLAYLIST.id, + PENDING_RESTORE_STATE + ); + expect(store.isContentInitialized()).toBe(true); + expect(store.contentInitBlockReason()).toBeNull(); + }); + + it('consumes pending state parked during an active import', async () => { + const pendingSeries = createDeferred(); + dataSource.getContent.mockImplementation( + (_playlistId: string, _credentials: unknown, type: ContentType) => + type === 'series' ? pendingSeries.promise : Promise.resolve([]) + ); + + const initialization = store.initializeContent(); + await waitForCondition( + () => store.contentLoadStateByType().vod === 'ready' + ); + + pendingRestoreService.getOrThrow.mockReturnValue(PENDING_RESTORE_STATE); + expect(store.reconcilePendingRestoreBlock()).toBe(true); + expect(readPendingRestoreBlocked(store)).toBe(true); + expect(store.isImporting()).toBe(false); + expect(store.activeImportSessionId()).not.toBeNull(); + expect(store.isContentInitialized()).toBe(false); + expect(dataSource.restoreUserData).not.toHaveBeenCalled(); + + pendingSeries.resolve([]); + await initialization; + + expect(dataSource.restoreUserData).toHaveBeenCalledWith( + PLAYLIST.id, + PENDING_RESTORE_STATE, + expect.any(Object) + ); + 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); + }); + + it('retries when pending state could not be consumed', async () => { + dataSource.getContent.mockResolvedValue([]); + pendingRestoreService.getOrThrow.mockReturnValue(PENDING_RESTORE_STATE); + pendingRestoreService.clear + .mockReturnValueOnce(false) + .mockImplementationOnce(() => { + pendingRestoreService.getOrThrow.mockReturnValue(null); + return true; + }); + + await store.initializeContent(); + + expect(dataSource.restoreUserData).toHaveBeenCalledTimes(1); + expect(store.isContentInitialized()).toBe(false); + expect(store.contentInitBlockReason()).toBe('error'); + + await store.retryContentInitialization(); + + expect(dataSource.restoreUserData).toHaveBeenCalledTimes(2); + expect(pendingRestoreService.clear).toHaveBeenCalledTimes(2); + expect(store.isContentInitialized()).toBe(true); + expect(store.contentInitBlockReason()).toBeNull(); + }); + + it('keeps cancellation authoritative while a pending restore finishes', async () => { + const pendingRestore = createDeferred(); + dataSource.getContent.mockResolvedValue([]); + dataSource.restoreUserData.mockReturnValueOnce(pendingRestore.promise); + pendingRestoreService.getOrThrow.mockReturnValue(PENDING_RESTORE_STATE); + + const initialization = store.initializeContent(); + await waitForCondition( + () => dataSource.restoreUserData.mock.calls.length === 1 + ); + + await store.cancelImport(); + pendingRestore.reject(new Error('pin replacement failed')); + await expect(initialization).resolves.toBeUndefined(); + + expect(pendingRestoreService.clear).not.toHaveBeenCalled(); + expect(store.isContentInitialized()).toBe(false); + expect(store.contentInitBlockReason()).toBe('cancelled'); + }); + it('stops before content fetch if cancel lands between categories and content phases', async () => { const pendingCategories = { live: createDeferred(), 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 15ac0f490..dfb5bb54b 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 @@ -106,6 +106,7 @@ export interface ContentState { activeImportSessionId: string | null; activeImportOperationIds: string[]; isContentInitialized: boolean; + isPendingRestoreBlocked: boolean; contentInitBlockReason: XtreamContentInitBlockReason | null; } @@ -141,6 +142,7 @@ const initialContentState: ContentState = { activeImportSessionId: null, activeImportOperationIds: [], isContentInitialized: false, + isPendingRestoreBlocked: false, contentInitBlockReason: null, }; @@ -270,6 +272,18 @@ export function withContent() { const asCachedContent = (content: unknown): T[] => content as T[]; + const hasPendingRestoreOrReadFailure = ( + playlistId: string + ): boolean => { + try { + return ( + pendingRestoreService.getOrThrow(playlistId) !== null + ); + } catch { + return true; + } + }; + const markContentScopeLoading = ( scope?: XtreamCachedContentScope | null, options?: { preserveInitialized?: boolean } @@ -404,6 +418,10 @@ export function withContent() { playlistId: string, scope?: XtreamCachedContentScope | null ): Promise => { + if (hasPendingRestoreOrReadFailure(playlistId)) { + return false; + } + const types = getTypesForCacheScope(scope); if ( @@ -419,10 +437,16 @@ export function withContent() { ) ) ); - return checks.some(Boolean); + return ( + checks.some(Boolean) && + !hasPendingRestoreOrReadFailure(playlistId) + ); } - return hasCachedContentForType(playlistId, scope); + return ( + (await hasCachedContentForType(playlistId, scope)) && + !hasPendingRestoreOrReadFailure(playlistId) + ); }; const isCurrentCachedHydrationContext = ( @@ -451,6 +475,35 @@ export function withContent() { return types.every((type) => loadStates[type] === 'ready'); }; + const blockCacheForPendingRestore = ( + playlistId: string, + scope?: XtreamCachedContentScope | null + ): boolean => { + if (!hasPendingRestoreOrReadFailure(playlistId)) { + return false; + } + + patchState(store, (state) => { + const nextLoadStates = { + ...state.contentLoadStateByType, + }; + for (const type of getTypesForCacheScope(scope)) { + nextLoadStates[type] = 'error'; + } + + return { + isLoadingCategories: false, + isLoadingContent: false, + isContentInitialized: false, + isPendingRestoreBlocked: true, + contentInitBlockReason: + state.contentInitBlockReason ?? 'error', + contentLoadStateByType: nextLoadStates, + }; + }); + return true; + }; + const executeCachedContentHydration = async ( playlistId: string, scope: XtreamCachedContentScope | null | undefined, @@ -520,6 +573,10 @@ export function withContent() { return; } + if (blockCacheForPendingRestore(playlistId, scope)) { + return; + } + patchState(store, (state) => { const nextLoadStates = { ...state.contentLoadStateByType, @@ -529,6 +586,7 @@ export function withContent() { isLoadingContent: false, isImporting: false, isContentInitialized: true, + isPendingRestoreBlocked: false, contentInitBlockReason: null, }; @@ -573,11 +631,16 @@ export function withContent() { const ctx = getCredentialsFromStore(); if (!ctx) return; + if (blockCacheForPendingRestore(ctx.playlistId, scope)) { + return; + } + if (isCachedContentScopeReady(scope)) { patchState(store, { isLoadingCategories: false, isLoadingContent: false, isContentInitialized: true, + isPendingRestoreBlocked: false, contentInitBlockReason: null, }); return; @@ -788,6 +851,17 @@ export function withContent() { const completedTypes = new Set(); try { + // Capture parked state before publishing any imported + // 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 + ); + patchState(store, { + isPendingRestoreBlocked: initialRestoreData !== null, + }); + // Electron content persistence maps remote category IDs // to internal DB category rows, so categories must exist // before content import starts. @@ -802,10 +876,18 @@ export function withContent() { }); throwIfImportCancelled(importSessionId); - // Restore user data if needed - const restoreData = pendingRestoreService.get( + // 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 ); + patchState(store, { + isPendingRestoreBlocked: restoreData !== null, + }); + + // Restore user data if needed if (restoreData) { try { throwIfImportCancelled(importSessionId); @@ -826,11 +908,27 @@ export function withContent() { } ); throwIfImportCancelled(importSessionId); - pendingRestoreService.clear(ctx.playlistId); + if ( + !pendingRestoreService.clear( + ctx.playlistId, + restoreData + ) + ) { + throw new Error( + `Clearing pending restore state for "${ctx.playlistId}" failed.` + ); + } + patchState(store, { + isPendingRestoreBlocked: false, + }); } catch (err) { + throwIfImportCancelled(importSessionId); + if (!isDbAbortError(err)) { logger.error('Error restoring user data', err); } + + throw err; } } @@ -1183,6 +1281,25 @@ export function withContent() { await runContentInitialization(); }, + reconcilePendingRestoreBlock(): boolean { + const ctx = getCredentialsFromStore(); + if (!ctx) { + return false; + } + + const isBlocked = hasPendingRestoreOrReadFailure( + ctx.playlistId + ); + patchState(store, (state) => ({ + isPendingRestoreBlocked: isBlocked, + contentInitBlockReason: + isBlocked && !state.activeImportSessionId + ? (state.contentInitBlockReason ?? 'error') + : state.contentInitBlockReason, + })); + return isBlocked; + }, + async hasUsableOfflineCache( scope?: XtreamCachedContentScope | null ): Promise { @@ -1203,7 +1320,12 @@ export function withContent() { isCachedContentScopeReady( scope?: XtreamCachedContentScope | null ): boolean { - return isCachedContentScopeReady(scope); + const ctx = getCredentialsFromStore(); + return ( + (!ctx || + !hasPendingRestoreOrReadFailure(ctx.playlistId)) && + isCachedContentScopeReady(scope) + ); }, async hydrateCachedContent( diff --git a/libs/portal/xtream/data-access/src/lib/xtream-state.ts b/libs/portal/xtream/data-access/src/lib/xtream-state.ts index 700a8dddf..4c3becfd2 100644 --- a/libs/portal/xtream/data-access/src/lib/xtream-state.ts +++ b/libs/portal/xtream/data-access/src/lib/xtream-state.ts @@ -65,6 +65,7 @@ export interface XtreamState { epgItems: EpgItem[]; hideExternalInfoDialog: boolean; portalStatus: PortalStatusType; + isPendingRestoreBlocked: boolean; contentInitBlockReason: XtreamContentInitBlockReason | null; globalSearchResults: GlobalSearchResult[]; streamUrl: string; diff --git a/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.spec.ts b/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.spec.ts index 30ae79211..255833334 100644 --- a/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.spec.ts +++ b/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.spec.ts @@ -32,6 +32,7 @@ describe('XtreamContentGateComponent', () => { const contentInitBlockReason = signal(null); const isContentInitialized = signal(false); + const isPendingRestoreBlocked = signal(false); const portalStatus = signal<'active' | 'inactive' | 'expired' | 'unavailable'>( 'active' ); @@ -40,6 +41,7 @@ describe('XtreamContentGateComponent', () => { beforeEach(async () => { contentInitBlockReason.set(null); isContentInitialized.set(false); + isPendingRestoreBlocked.set(false); portalStatus.set('active'); retryContentInitialization.mockClear(); @@ -65,6 +67,7 @@ describe('XtreamContentGateComponent', () => { useValue: { contentInitBlockReason, isContentInitialized, + isPendingRestoreBlocked, portalStatus, retryContentInitialization, }, @@ -110,10 +113,19 @@ describe('XtreamContentGateComponent', () => { it('keeps the child outlet available when there is no block reason', () => { fixture.detectChanges(); + expect(fixture.nativeElement.querySelector('.mock-error')).toBeNull(); + expect(fixture.nativeElement.querySelector('router-outlet')).not.toBeNull(); + }); + + it('keeps the child outlet unavailable while parked state is pending', () => { + isPendingRestoreBlocked.set(true); + fixture.detectChanges(); + + expect(fixture.nativeElement.querySelector('router-outlet')).toBeNull(); expect( fixture.nativeElement.querySelector('.mock-error') - ).toBeNull(); - expect(fixture.nativeElement.querySelector('router-outlet')).not.toBeNull(); + ).not.toBeNull(); + expect(fixture.nativeElement.querySelector('button')).not.toBeNull(); }); it('shows an inline warning when cached content remains available offline', () => { diff --git a/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.ts b/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.ts index e53f35e12..421f08fe2 100644 --- a/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.ts +++ b/libs/portal/xtream/feature/src/lib/xtream-content-gate.component.ts @@ -24,7 +24,7 @@ import { XtreamCachedOfflineNoticeComponent } from './xtream-cached-offline-noti XtreamCachedOfflineNoticeComponent, ], template: ` - @if (contentInitBlockReason(); as blockReason) { + @if (effectiveBlockReason(); as blockReason) {
+ this.contentInitBlockReason() ?? + (this.isPendingRestoreBlocked() ? 'error' : null) + ); private readonly errorViewKey = computed(() => { - switch (this.contentInitBlockReason()) { + switch (this.effectiveBlockReason()) { case 'cancelled': return 'IMPORT_CANCELLED'; case 'expired': diff --git a/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.spec.ts b/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.spec.ts index a8aca7734..fb4525532 100644 --- a/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.spec.ts +++ b/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.spec.ts @@ -116,6 +116,7 @@ describe('XtreamWorkspaceRouteSession', () => { setCurrentPlaylist: jest.fn((playlist: XtreamPlaylistData | null) => { currentPlaylist.set(playlist); }), + reconcilePendingRestoreBlock: jest.fn().mockReturnValue(false), fetchXtreamPlaylist: jest.fn().mockResolvedValue(undefined), checkPortalStatus: jest.fn(), hasUsableOfflineCache: jest.fn().mockImplementation(async () => { @@ -236,6 +237,7 @@ describe('XtreamWorkspaceRouteSession', () => { xtreamStore.resetStore.mockClear(); xtreamStore.setCurrentPlaylist.mockClear(); + xtreamStore.reconcilePendingRestoreBlock.mockClear(); xtreamStore.fetchXtreamPlaylist.mockClear(); xtreamStore.checkPortalStatus.mockReset(); xtreamStore.hasUsableOfflineCache.mockClear(); @@ -361,6 +363,16 @@ describe('XtreamWorkspaceRouteSession', () => { await flushEffects(); expect(xtreamStore.resetStore).toHaveBeenCalledWith(PLAYLIST_ID); + expect( + xtreamStore.reconcilePendingRestoreBlock.mock.invocationCallOrder[0] + ).toBeGreaterThan( + xtreamStore.setCurrentPlaylist.mock.invocationCallOrder[0] + ); + expect( + xtreamStore.reconcilePendingRestoreBlock.mock.invocationCallOrder[0] + ).toBeLessThan( + xtreamStore.fetchXtreamPlaylist.mock.invocationCallOrder[0] + ); expect(xtreamStore.setSelectedContentType).toHaveBeenCalledWith('live'); expect(xtreamStore.prepareContentLoading).toHaveBeenCalledWith('live'); expect(selectedContentType()).toBe('live'); diff --git a/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.ts b/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.ts index 90aba19d2..0b48352a7 100644 --- a/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.ts +++ b/libs/portal/xtream/feature/src/lib/xtream-workspace-route-session.service.ts @@ -298,6 +298,7 @@ export class XtreamWorkspaceRouteSession { didBootstrapPlaylist = true; this.xtreamStore.setCurrentPlaylist(routePlaylist); + this.xtreamStore.reconcilePendingRestoreBlock(); section = this.syncRouteState(routeSection); if (isImportDrivenSection(section)) { this.xtreamStore.prepareContentLoading(cacheScope); 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 98eea4d2c..d97cd630c 100644 --- a/libs/services/src/lib/playlist-backup.service.pins.spec.ts +++ b/libs/services/src/lib/playlist-backup.service.pins.spec.ts @@ -116,6 +116,56 @@ describe('PlaylistBackupService Xtream source pins', () => { ]); }); + it('serializes overlapping restores so an older snapshot cannot finish last', async () => { + const collaborators = createRestoreCollaborators(); + let finishFirstReplacement!: (replaced: boolean) => void; + const firstReplacement = new Promise((resolve) => { + finishFirstReplacement = resolve; + }); + const replacePins = jest + .fn() + .mockReturnValueOnce(firstReplacement) + .mockResolvedValueOnce(true); + const service = createPlaylistBackupService({ + ...collaborators, + vodSourcePinService: { + isAvailable: true, + listForPlaylistOrThrow: jest.fn().mockResolvedValue([]), + replaceForPlaylist: replacePins, + }, + }); + const firstImport = service.importBackup( + JSON.stringify( + createXtreamManifest( + [], + [{ matchKey: 'tmdb:603', contentId: 501 }] + ) + ) + ); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(replacePins).toHaveBeenCalledTimes(1); + + const secondImport = service.importBackup( + JSON.stringify( + createXtreamManifest( + [], + [{ matchKey: 'tmdb:603', contentId: 502 }] + ) + ) + ); + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(replacePins).toHaveBeenCalledTimes(1); + + finishFirstReplacement(true); + await Promise.all([firstImport, secondImport]); + expect( + replacePins.mock.calls.map( + ([, pins]) => (pins as Array<{ contentId: number }>)[0].contentId + ) + ).toEqual([501, 502]); + }); + it('reports pins that could not be written', async () => { const collaborators = createRestoreCollaborators(); const service = createPlaylistBackupService({ @@ -142,6 +192,31 @@ describe('PlaylistBackupService Xtream source pins', () => { ); }); + it('reports a restore whose parked state could not be consumed', async () => { + const collaborators = createRestoreCollaborators(); + collaborators.pendingRestoreService.clear.mockReturnValue(false); + const service = createPlaylistBackupService({ + ...collaborators, + vodSourcePinService: { + isAvailable: true, + listForPlaylistOrThrow: jest.fn().mockResolvedValue([]), + replaceForPlaylist: jest.fn().mockResolvedValue(true), + }, + }); + + 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(/pending restore state/i); + }); + 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 a14f9f3af..c6fb18849 100644 --- a/libs/services/src/lib/playlist-backup.service.test-helpers.ts +++ b/libs/services/src/lib/playlist-backup.service.test-helpers.ts @@ -70,7 +70,7 @@ export function createPlaylistBackupService( }, pendingRestoreService: { set: jest.fn(), - clear: jest.fn(), + clear: jest.fn().mockReturnValue(true), }, ...overrides, }); diff --git a/libs/services/src/lib/playlist-backup.service.ts b/libs/services/src/lib/playlist-backup.service.ts index 881041f5f..c62b8f3c9 100644 --- a/libs/services/src/lib/playlist-backup.service.ts +++ b/libs/services/src/lib/playlist-backup.service.ts @@ -65,6 +65,7 @@ export class PlaylistBackupService { private readonly pendingRestoreService = inject( XtreamPendingRestoreService ); + private backupImportTail?: Promise; async exportBackup(): Promise { const playlists = await firstValueFrom( @@ -94,7 +95,20 @@ export class PlaylistBackupService { }; } - async importBackup(json: string): Promise { + importBackup(json: string): Promise { + const importResult = ( + this.backupImportTail ?? Promise.resolve() + ).then(() => this.executeImportBackup(json)); + this.backupImportTail = importResult.then( + () => undefined, + () => undefined + ); + return importResult; + } + + private async executeImportBackup( + json: string + ): Promise { const manifest = this.parseManifest(json); const existingPlaylists = await firstValueFrom( this.playlistsService.getAllData() @@ -838,7 +852,11 @@ export class PlaylistBackupService { } await this.applyXtreamRestoreState(playlistId, restoreState); - this.pendingRestoreService.clear(playlistId); + if (!this.pendingRestoreService.clear(playlistId, restoreState)) { + throw new PlaylistBackupError( + `Clearing pending restore state for "${playlistId}" failed.` + ); + } } private async hasCompletedOfflineCache( 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 61f8fe308..ca8ed5383 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 @@ -87,7 +87,10 @@ describe('PlaylistBackupService Xtream hidden categories (issue #1017)', () => { collaborators.databaseService.updateCategoryVisibility ).toHaveBeenCalledTimes(3); expect(collaborators.pendingRestoreService.clear).toHaveBeenCalledWith( - 'xtream-1' + 'xtream-1', + expect.objectContaining({ + hiddenCategories: [{ categoryType: 'live', xtreamId: 101 }], + }) ); }); diff --git a/libs/services/src/lib/playlist-backup.xtream-fixtures.ts b/libs/services/src/lib/playlist-backup.xtream-fixtures.ts index 7e7e95004..c374a5ec1 100644 --- a/libs/services/src/lib/playlist-backup.xtream-fixtures.ts +++ b/libs/services/src/lib/playlist-backup.xtream-fixtures.ts @@ -113,7 +113,7 @@ export function createRestoreCollaborators() { }, pendingRestoreService: { set: jest.fn(), - clear: jest.fn(), + clear: jest.fn().mockReturnValue(true), }, }; } 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 772125fab..c0e4738aa 100644 --- a/libs/services/src/lib/xtream-pending-restore.service.spec.ts +++ b/libs/services/src/lib/xtream-pending-restore.service.spec.ts @@ -12,6 +12,7 @@ describe('XtreamPendingRestoreService', () => { }); afterEach(() => { + jest.restoreAllMocks(); localStorage.clear(); }); @@ -62,4 +63,93 @@ describe('XtreamPendingRestoreService', () => { localStorage.setItem(storageKey, '{not json'); expect(service.get(playlistId)).toBeNull(); }); + + it('offers a strict read that reports storage access failures', () => { + jest.spyOn(Storage.prototype, 'getItem').mockImplementation(() => { + throw new Error('storage is locked'); + }); + + expect(service.get(playlistId)).toBeNull(); + expect(() => service.getOrThrow(playlistId)).toThrow( + 'storage is locked' + ); + }); + + it('clears persisted state and reports that it was consumed', () => { + service.set(playlistId, { + hiddenCategories: [], + favorites: [], + recentlyViewed: [], + playbackPositions: [], + }); + + expect(service.clear(playlistId)).toBe(true); + expect(service.get(playlistId)).toBeNull(); + expect(localStorage.getItem(storageKey)).toBeNull(); + }); + + it.each([ + { + failure: 'throws', + remove: () => { + throw new Error('storage is locked'); + }, + }, + { + failure: 'does not remove the key', + remove: () => undefined, + }, + ])( + 'consumes state with a tombstone when removeItem $failure', + ({ remove }) => { + service.set(playlistId, { + hiddenCategories: [], + favorites: [], + recentlyViewed: [], + playbackPositions: [], + }); + jest.spyOn(Storage.prototype, 'removeItem').mockImplementation( + remove + ); + + expect(service.clear(playlistId)).toBe(true); + expect(service.get(playlistId)).toBeNull(); + expect(localStorage.getItem(storageKey)).toBe(''); + } + ); + + it('reports failure when state can neither be tombstoned nor removed', () => { + service.set(playlistId, { + hiddenCategories: [], + favorites: [], + recentlyViewed: [], + playbackPositions: [], + }); + jest.spyOn(Storage.prototype, 'setItem').mockImplementation(() => { + throw new Error('storage is locked'); + }); + jest.spyOn(Storage.prototype, 'removeItem').mockImplementation(() => { + throw new Error('storage is locked'); + }); + + expect(service.clear(playlistId)).toBe(false); + expect(service.get(playlistId)).not.toBeNull(); + }); + + it('does not clear a newer snapshot than the one already restored', () => { + const restoredState = { + hiddenCategories: [], + favorites: [{ xtreamId: 101, contentType: 'live' as const }], + recentlyViewed: [], + playbackPositions: [], + }; + const newerState = { + ...restoredState, + favorites: [{ xtreamId: 202, contentType: 'movie' as const }], + }; + service.set(playlistId, newerState); + + expect(service.clear(playlistId, restoredState)).toBe(false); + expect(service.getOrThrow(playlistId)).toEqual(newerState); + }); }); diff --git a/libs/services/src/lib/xtream-pending-restore.service.ts b/libs/services/src/lib/xtream-pending-restore.service.ts index 59b6178f6..277e4d1b0 100644 --- a/libs/services/src/lib/xtream-pending-restore.service.ts +++ b/libs/services/src/lib/xtream-pending-restore.service.ts @@ -10,19 +10,31 @@ import { }) export class XtreamPendingRestoreService { get(playlistId: string): XtreamPendingRestoreState | null { + try { + return this.getOrThrow(playlistId); + } catch { + return null; + } + } + + /** + * Strict restore-path read. Storage access failures are not equivalent to + * an absent snapshot: callers that could expose mutable content must fail + * closed until they can prove the state is missing or consumed. + */ + getOrThrow(playlistId: string): XtreamPendingRestoreState | null { if (!playlistId) { return null; } + const rawState = localStorage.getItem( + getXtreamPendingRestoreStorageKey(playlistId) + ); + if (!rawState) { + return null; + } + try { - const rawState = localStorage.getItem( - getXtreamPendingRestoreStorageKey(playlistId) - ); - - if (!rawState) { - return null; - } - // Persisted state may predate the current build (e.g. entries // written by versions affected by issue #1017), so it is // re-normalized on every read, not only on write. @@ -47,17 +59,55 @@ export class XtreamPendingRestoreService { } } - clear(playlistId: string): void { + clear( + playlistId: string, + expectedState?: XtreamPendingRestoreState + ): boolean { if (!playlistId) { - return; + return false; + } + + if (expectedState) { + let currentState: XtreamPendingRestoreState | null; + try { + currentState = this.getOrThrow(playlistId); + } catch { + return false; + } + + if ( + !currentState || + JSON.stringify(currentState) !== + JSON.stringify( + normalizeXtreamPendingRestoreState(expectedState) + ) + ) { + return false; + } + } + + const storageKey = getXtreamPendingRestoreStorageKey(playlistId); + let tombstoneStored = false; + + // An empty value is already treated as "no pending restore" by get(). + // Persist it first so a throwing or no-op remove cannot leave the old + // snapshot eligible for a later destructive replay. + try { + localStorage.setItem(storageKey, ''); + tombstoneStored = localStorage.getItem(storageKey) === ''; + } catch { + // Removal may still work when writing is unavailable. } try { - localStorage.removeItem( - getXtreamPendingRestoreStorageKey(playlistId) - ); + localStorage.removeItem(storageKey); + if (localStorage.getItem(storageKey) === null) { + return true; + } } catch { - // Ignore local storage remove failures. + // A verified tombstone is sufficient even if removal fails. } + + return tombstoneStored; } } diff --git a/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell-xtream-import.service.ts b/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell-xtream-import.service.ts index 7ffc7fe3b..35c64202a 100644 --- a/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell-xtream-import.service.ts +++ b/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell-xtream-import.service.ts @@ -67,7 +67,6 @@ export class WorkspaceShellXtreamImportService { readonly canCancelXtreamImport = computed( () => this.isElectron && - this.xtreamStore.isImporting() && Boolean(this.xtreamStore.activeImportSessionId()) && !this.xtreamStore.isCancellingImport() ); @@ -193,7 +192,7 @@ export class WorkspaceShellXtreamImportService { readonly isImportRunning = computed( () => !this.xtreamStore.contentInitBlockReason() && - this.xtreamStore.isImporting() + Boolean(this.xtreamStore.activeImportSessionId()) ); readonly isRefreshPreparationRunning = computed(() => Boolean(this.activeRefreshPreparation()) diff --git a/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell.facade.spec.ts b/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell.facade.spec.ts index 06585d642..eaa580238 100644 --- a/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell.facade.spec.ts +++ b/libs/workspace/shell/feature/src/lib/workspace-shell/services/workspace-shell.facade.spec.ts @@ -444,6 +444,18 @@ describe('WorkspaceShellFacade', () => { expect(facade.showXtreamImportOverlay()).toBe(true); }); + it('keeps a cache-only Xtream import session blocking and cancellable', () => { + const xtreamStore = TestBed.inject( + XtreamStore + ) as unknown as MockXtreamStore; + const xtreamImport = TestBed.inject(WorkspaceShellXtreamImportService); + xtreamStore.activeImportSessionId.set('xtream-import-session'); + + expect(xtreamStore.isImporting()).toBe(false); + expect(facade.showXtreamImportOverlay()).toBe(true); + expect(xtreamImport.canCancelXtreamImport()).toBe(true); + }); + it('shows the Xtream overlay during refresh preparation on the dashboard', () => { facade.currentUrl.set('/workspace/dashboard'); refreshPreparationSignal.set({