fix(portals): preserve Xtream restore retry state

This commit is contained in:
4gray committed 2026-07-30 00:58:38 +02:00
1 parent b0598f8f70
commit 83fb06e79a
24 files changed
+913 -115

No files matched your search

+4 -2
View File
@@ -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.
@@ -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<void>;
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();
});
});
@@ -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());
+16
View File
@@ -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
@@ -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),
@@ -566,6 +566,48 @@ export class ElectronXtreamDataSource implements IXtreamDataSource {
restoreState: XtreamPendingRestoreState,
options?: XtreamOperationOptions
): Promise<void> {
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,
@@ -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);
@@ -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<PortalStatusType>
) {
return signalStore(
withState({
playlistId: TEST_PLAYLIST.id,
currentPlaylist: TEST_PLAYLIST,
portalStatus: 'active' as PortalStatusType,
selectedContentType: 'vod' as const,
}),
withMethods((store) => ({
async checkPortalStatus(): Promise<PortalStatusType> {
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<symbol, unknown>;
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<T>() {
let resolve!: (value: T) => void;
let reject!: (reason?: unknown) => void;
const promise = new Promise<T>((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<void> {
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');
}
@@ -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<symbol, unknown>;
if (hook === null) {
delete target[performanceHookSymbol];
} else {
target[performanceHookSymbol] = hook;
}
}
const PLAYLIST = TEST_PLAYLIST;
let checkPortalStatusMock: jest.Mock<Promise<PortalStatusType>, []>;
const TestContentStore = signalStore(
withState({
playlistId: PLAYLIST.id,
currentPlaylist: PLAYLIST,
portalStatus: 'active' as PortalStatusType,
selectedContentType: 'vod' as const,
}),
withMethods((store) => ({
async checkPortalStatus(): Promise<PortalStatusType> {
const status = await checkPortalStatusMock();
patchState(store, { portalStatus: status });
return status;
},
})),
withContent()
);
function createDeferred<T>() {
let resolve!: (value: T) => void;
let reject!: (reason?: unknown) => void;
const promise = new Promise<T>((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<void> {
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<typeof TestContentStore>;
@@ -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<any[]>();
const cachedContent = createDeferred<any[]>();
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<any[]>();
const cachedContent = createDeferred<any[]>();
@@ -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<any[]>();
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<void>();
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<any[]>(),
@@ -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 = <T>(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<boolean> => {
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<ContentType>();
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<boolean> {
@@ -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(
@@ -65,6 +65,7 @@ export interface XtreamState {
epgItems: EpgItem[];
hideExternalInfoDialog: boolean;
portalStatus: PortalStatusType;
isPendingRestoreBlocked: boolean;
contentInitBlockReason: XtreamContentInitBlockReason | null;
globalSearchResults: GlobalSearchResult[];
streamUrl: string;
@@ -32,6 +32,7 @@ describe('XtreamContentGateComponent', () => {
const contentInitBlockReason =
signal<XtreamContentInitBlockReason | null>(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', () => {
@@ -24,7 +24,7 @@ import { XtreamCachedOfflineNoticeComponent } from './xtream-cached-offline-noti
XtreamCachedOfflineNoticeComponent,
],
template: `
@if (contentInitBlockReason(); as blockReason) {
@if (effectiveBlockReason(); as blockReason) {
<div class="xtream-content-gate">
<app-playlist-error-view
[title]="titleKey() | translate"
@@ -69,8 +69,14 @@ export class XtreamContentGateComponent {
private readonly xtreamStore = inject(XtreamStore);
readonly contentInitBlockReason = this.xtreamStore.contentInitBlockReason;
readonly isPendingRestoreBlocked = this.xtreamStore.isPendingRestoreBlocked;
readonly effectiveBlockReason = computed(
() =>
this.contentInitBlockReason() ??
(this.isPendingRestoreBlocked() ? 'error' : null)
);
private readonly errorViewKey = computed(() => {
switch (this.contentInitBlockReason()) {
switch (this.effectiveBlockReason()) {
case 'cancelled':
return 'IMPORT_CANCELLED';
case 'expired':
@@ -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');
@@ -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);
@@ -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<boolean>((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);
@@ -70,7 +70,7 @@ export function createPlaylistBackupService(
},
pendingRestoreService: {
set: jest.fn(),
clear: jest.fn(),
clear: jest.fn().mockReturnValue(true),
},
...overrides,
});
@@ -65,6 +65,7 @@ export class PlaylistBackupService {
private readonly pendingRestoreService = inject(
XtreamPendingRestoreService
);
private backupImportTail?: Promise<void>;
async exportBackup(): Promise<PlaylistBackupExportPayload> {
const playlists = await firstValueFrom(
@@ -94,7 +95,20 @@ export class PlaylistBackupService {
};
}
async importBackup(json: string): Promise<PlaylistBackupImportSummary> {
importBackup(json: string): Promise<PlaylistBackupImportSummary> {
const importResult = (
this.backupImportTail ?? Promise.resolve()
).then(() => this.executeImportBackup(json));
this.backupImportTail = importResult.then(
() => undefined,
() => undefined
);
return importResult;
}
private async executeImportBackup(
json: string
): Promise<PlaylistBackupImportSummary> {
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(
@@ -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 }],
})
);
});
@@ -113,7 +113,7 @@ export function createRestoreCollaborators() {
},
pendingRestoreService: {
set: jest.fn(),
clear: jest.fn(),
clear: jest.fn().mockReturnValue(true),
},
};
}
@@ -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);
});
});
@@ -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;
}
}
@@ -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())
@@ -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({