refactor(portals): extract the destructive Xtream refresh into one flow

`PlaylistRefreshActionService.refreshXtream()` and
`RecentPlaylistsComponent.refreshXtreamPlaylist()` were two independent ~60-line
implementations of the same sequence: confirm, reset the connectivity guard,
delete the cached catalog while collecting playback positions and stamping the
update date, park the restore state, dispatch the meta update, navigate to
re-import. That duplication is what let the guard reset land in one of them and
look landed in both, which the previous commit had to fix separately.

`XtreamRefreshFlowService` now owns the sequence once. The two entry points
differ only in how they report progress — the header action drives one global
preparation signal, the sources page drives per-row busy indicators shared with
deletion — so that part is injected as an `XtreamRefreshProgressReporter`. Its
optional `waitForVisibleProgress` hook keeps the action's paint delay without
imposing it on the sources page, which never had one. A reporter cannot reach
the guard reset, so a third entry point gets it for free.

The sources page now awaits `router.navigate` before clearing its busy row,
matching the action service; it previously cleared the row while navigation was
still in flight.

Both guard-reset regression tests still fail on their own when the reset is
removed from the shared flow. The new spec pins the reporter contract:
re-entry after confirmation, busy-before-reset ordering, progress forwarding
under one run, abort vs error, and a reporter without the optional paint hook.
The sources spec's abort test counted a fixed number of microtasks to reach the
delete; the shared flow adds one await to that path, so it now waits on a
deterministic signal instead.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
4grayandClaude Opus 5 committed 2026-08-13 08:25:21 +02:00
1 parent 0cba49f3e2
commit e3fc3c8802
8 files changed
+593 -257

No files matched your search

+1 -1
View File
@@ -698,7 +698,7 @@ This project uses modern Angular signal-based APIs and patterns. **ALWAYS** use
- `epg.events.ts` - EPG IPC registration; freshness/fetch orchestration lives in `epg-fetch.service.ts`, manual channel-mapping resolution and CRUD in `epg-mapping.service.ts`, worker lifecycle in `epg-worker.service.ts`, DB lookups in `epg-query.service.ts`
- `xtream.events.ts` - Xtream Codes API
- `stalker.events.ts` - Stalker portal API
- `connectivity-guard.events.ts` - `CONNECTIVITY_GUARD_RESET`: forgets the connection failures recorded for a portal host. Both portal handlers above run every request through the per-host circuit breaker in `util/host-connectivity-guard.ts` — after 2 consecutive connection-level failures (no HTTP response; `ETIMEDOUT`/`ENOTFOUND`/`ECONNREFUSED`/… but never `ECONNRESET`) requests to that endpoint fail immediately for 30 s. The key is `URL.origin`, not `URL.host`, which would give `http://panel` and `https://panel` one shared record and let a dead TLS listener fast-fail the working HTTP one instead of hanging the full 30 s/15 s axios timeout again, with one half-open trial request afterwards. Any HTTP response (4xx and 5xx included) clears the record. The refusal is a real `Error` whose wording is a renderer contract (`buildHostConnectivityFastFailMessage` in `libs/shared/interfaces`): it must carry no `HTTP Error <code>`, no timeout wording and none of the auth phrases, or Stalker endpoint discovery misclassifies it and lazy portal repair fires against a host just declared dead. Discovery probes are exempt via the `skipConnectionGuard` payload flag (bypass + no failure counting, but successes still clear the record). Every user-driven retry/refresh that issues portal requests must reset BEFORE its first request, or the affordance fast-fails and looks broken; automatic and first-load paths deliberately do not reset. Current senders: Xtream content-gate Retry, Stalker catalog append retry (`retryContentPage`), Stalker search-page retry, `StalkerItvCacheService.refresh()` (Live TV refresh), both account-info dialogs' Retry, both destructive Xtream refresh implementations (`PlaylistRefreshActionService.refreshXtream()` and `RecentPlaylistsComponent.refreshXtreamPlaylist()`, before they delete the cached catalog), `StalkerPortalDiscoveryService.discover()`, and `PortalStatusService` on `skipCache`. Kill switch: `IPTVNATOR_DISABLE_CONNECTIVITY_GUARD=1`. Contract: `docs/architecture/host-connectivity-guard.md`
- `connectivity-guard.events.ts` - `CONNECTIVITY_GUARD_RESET`: forgets the connection failures recorded for a portal host. Both portal handlers above run every request through the per-host circuit breaker in `util/host-connectivity-guard.ts` — after 2 consecutive connection-level failures (no HTTP response; `ETIMEDOUT`/`ENOTFOUND`/`ECONNREFUSED`/… but never `ECONNRESET`) requests to that endpoint fail immediately for 30 s. The key is `URL.origin`, not `URL.host`, which would give `http://panel` and `https://panel` one shared record and let a dead TLS listener fast-fail the working HTTP one instead of hanging the full 30 s/15 s axios timeout again, with one half-open trial request afterwards. Any HTTP response (4xx and 5xx included) clears the record. The refusal is a real `Error` whose wording is a renderer contract (`buildHostConnectivityFastFailMessage` in `libs/shared/interfaces`): it must carry no `HTTP Error <code>`, no timeout wording and none of the auth phrases, or Stalker endpoint discovery misclassifies it and lazy portal repair fires against a host just declared dead. Discovery probes are exempt via the `skipConnectionGuard` payload flag (bypass + no failure counting, but successes still clear the record). Every user-driven retry/refresh that issues portal requests must reset BEFORE its first request, or the affordance fast-fails and looks broken; automatic and first-load paths deliberately do not reset. Current senders: Xtream content-gate Retry, Stalker catalog append retry (`retryContentPage`), Stalker search-page retry, `StalkerItvCacheService.refresh()` (Live TV refresh), both account-info dialogs' Retry, the destructive Xtream refresh (`XtreamRefreshFlowService`, before it deletes the cached catalog — one flow shared by both entry points, `PlaylistRefreshActionService.refreshXtream()` and `RecentPlaylistsComponent.refreshXtreamPlaylist()`, which supply only a progress reporter), `StalkerPortalDiscoveryService.discover()`, and `PortalStatusService` on `skipCache`. Kill switch: `IPTVNATOR_DISABLE_CONNECTIVITY_GUARD=1`. Contract: `docs/architecture/host-connectivity-guard.md`
- `player.events.ts` - External player IPC registration; MPV/VLC lifecycle logic lives in `mpv-session.service.ts`, `vlc-session.service.ts`, and shared `external-player-*` helpers
- `settings.events.ts` - App settings
- `electron.events.ts` - App version, etc.
+8 -4
View File
@@ -200,11 +200,15 @@ Call sites:
- The destructive Xtream refresh — before anything is deleted. It removes the
cached catalog and then bootstraps a re-import whose status request an open
guard would fast-fail, leaving the user with no catalog at all until the
cooldown expires. There are **two independent implementations** of this flow
and both need the reset: `PlaylistRefreshActionService.refreshXtream()` and
cooldown expires. The reset lives in `XtreamRefreshFlowService.runRefresh()`
(`libs/playlist/shared/ui`), which owns the whole flow for both of its entry
points: `PlaylistRefreshActionService.refreshXtream()` and
`RecentPlaylistsComponent.refreshXtreamPlaylist()` (the Workspace sources
page). They are near-duplicates of each other, which is exactly why the second
one was missed first time round.
page). Those two used to be independent near-duplicates, which is exactly why
the second one was missed first time round; they now differ only in the
`XtreamRefreshProgressReporter` they hand over, and a reporter cannot reach
the reset. Keep it that way — a third entry point should pass a reporter, not
copy the sequence.
- `PortalStatusService.checkPortalStatusDetails` when `skipCache` is set — the
user-initiated "Test Connection".
+1
View File
@@ -5,3 +5,4 @@ export * from './lib/recent-playlists/playlist-info/stalker-playlist-connection-
export * from './lib/recent-playlists/recent-playlists.component';
export * from './lib/recent-playlists/empty-state/empty-state.component';
export * from './lib/playlist-refresh-action.service';
export * from './lib/xtream-refresh-flow.service';
@@ -1,5 +1,4 @@
import { inject, Injectable, signal } from '@angular/core';
import { Router } from '@angular/router';
import { Store } from '@ngrx/store';
import { TranslateService } from '@ngx-translate/core';
import { MatSnackBar } from '@angular/material/snack-bar';
@@ -9,12 +8,9 @@ import {
DatabaseService,
type DbOperationEvent,
isDbAbortError,
PlaybackPositionService,
PlaylistRefreshService,
resetHostConnectivityGuard,
RuntimeCapabilitiesService,
SettingsStore,
XtreamPendingRestoreService,
} from '@iptvnator/services';
import { ChannelActions, PlaylistActions } from '@iptvnator/m3u-state';
import {
@@ -24,11 +20,11 @@ import {
PLAYLIST_UPDATE,
PlaylistMeta,
} from '@iptvnator/shared/interfaces';
import {
measureRendererPerformancePhase,
RENDERER_PERFORMANCE_PHASE,
} from '@iptvnator/shared/logging';
import { PlaylistContextFacade } from '@iptvnator/playlist/shared/util';
import {
XtreamRefreshFlowService,
type XtreamRefreshProgressReporter,
} from './xtream-refresh-flow.service';
export interface XtreamRefreshPreparationState {
playlistId: string;
@@ -40,21 +36,17 @@ export interface XtreamRefreshPreparationState {
@Injectable({ providedIn: 'root' })
export class PlaylistRefreshActionService {
private readonly router = inject(Router);
private readonly store = inject(Store);
private readonly translate = inject(TranslateService);
private readonly snackBar = inject(MatSnackBar);
private readonly dialogService = inject(DialogService);
private readonly databaseService = inject(DatabaseService);
private readonly dataService = inject(DataService);
private readonly playbackPositionService = inject(PlaybackPositionService);
private readonly playlistRefreshService = inject(PlaylistRefreshService);
private readonly runtime = inject(RuntimeCapabilitiesService);
private readonly settingsStore = inject(SettingsStore);
private readonly playlistContext = inject(PlaylistContextFacade);
private readonly pendingRestoreService = inject(
XtreamPendingRestoreService
);
private readonly xtreamRefreshFlow = inject(XtreamRefreshFlowService);
private readonly refreshPreparationState =
signal<XtreamRefreshPreparationState | null>(null);
@@ -102,122 +94,41 @@ export class PlaylistRefreshActionService {
}
private refreshXtream(item: PlaylistMeta): void {
this.dialogService.openConfirmDialog({
title: this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.TITLE'
),
message: this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.MESSAGE'
),
width: '400px',
onConfirm: async () => {
if (this.isRefreshing()) {
return;
}
this.isRefreshing.set(true);
const operationId =
this.databaseService.createOperationId('xtream-refresh');
this.refreshPreparationState.set({
playlistId: item._id,
operationId,
phase: 'collecting-user-data',
});
try {
// Before anything destructive: this refresh deletes the
// cached catalog and then forces a route bootstrap whose
// status request would be fast-failed by an open
// connectivity guard, leaving the user with no catalog at
// all until the cooldown expires.
await resetHostConnectivityGuard(
this.dataService,
item.serverUrl
);
this.snackBar.open(
this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.STARTED'
),
undefined,
{ duration: 2000 }
);
await this.waitForRefreshPreparationPaint();
const updateDate = Date.now();
const [restoreState, playbackPositions] = await Promise.all(
[
this.databaseService.deleteXtreamPlaylistContent(
item._id,
{
operationId,
onEvent: (event) =>
this.updateRefreshPreparationFromEvent(
item._id,
operationId,
event
),
}
),
this.playbackPositionService.getAllPlaybackPositions(
item._id
),
this.databaseService.updateXtreamPlaylistDetails({
id: item._id,
updateDate,
}),
]
);
if (
!this.pendingRestoreService.set(item._id, {
...restoreState,
playbackPositions,
})
) {
throw new Error(
`Parking pending restore state for "${item._id}" failed.`
);
}
measureRendererPerformancePhase(
RENDERER_PERFORMANCE_PHASE.XTREAM_REFRESH_META,
() =>
this.store.dispatch(
PlaylistActions.updatePlaylistMeta({
playlist: { ...item, updateDate },
})
),
() => ({ items: 1 })
);
await this.router.navigate([
'/workspace',
'xtreams',
item._id,
]);
} catch (error) {
if (!isDbAbortError(error)) {
console.error(
'Error refreshing Xtream playlist:',
error
);
this.snackBar.open(
this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.ERROR'
),
undefined,
{ duration: 3000 }
);
}
} finally {
this.clearRefreshPreparation(operationId);
this.isRefreshing.set(false);
}
},
});
this.xtreamRefreshFlow.confirmAndRefresh(
item,
this.xtreamRefreshReporter
);
}
/**
* Reports the shared destructive refresh as one global preparation state:
* this action has no row to attach progress to, so the flow is either
* refreshing or it is not. Events are scoped to their operation, since a
* cancelled run's worker can still emit after the next one has started.
*/
private readonly xtreamRefreshReporter: XtreamRefreshProgressReporter = {
isBusy: () => this.isRefreshing(),
begin: ({ playlistId, operationId }) => {
this.isRefreshing.set(true);
this.refreshPreparationState.set({
playlistId,
operationId,
phase: 'collecting-user-data',
});
},
report: ({ playlistId, operationId }, event) =>
this.updateRefreshPreparationFromEvent(
playlistId,
operationId,
event
),
end: ({ operationId }) => {
this.clearRefreshPreparation(operationId);
this.isRefreshing.set(false);
},
waitForVisibleProgress: () => this.waitForRefreshPreparationPaint(),
};
private async refreshM3u(item: PlaylistMeta): Promise<void> {
const isActiveM3uRoute =
this.playlistContext.routeProvider() === 'playlists' &&
@@ -418,6 +418,12 @@ describe('RecentPlaylistsComponent busy state', () => {
}>;
}>();
let confirmPromise: Promise<void> | undefined;
// The shared flow awaits the connectivity-guard reset before it starts
// deleting, so the delete is several ticks away and counting them is
// fragile. The mock fires onEvent synchronously and the component
// applies it synchronously, so resolving this deferred inside the mock
// is a deterministic "the busy row has been updated" signal.
const workerEventDelivered = createDeferred<void>();
dialogService.openConfirmDialog.mockImplementation(
({ onConfirm }: { onConfirm?: () => Promise<void> }) => {
@@ -441,16 +447,14 @@ describe('RecentPlaylistsComponent busy state', () => {
current: 1,
total: 4,
});
workerEventDelivered.resolve();
return refresh.promise;
}
);
component.refreshXtreamPlaylist(item);
// The refresh clears the connectivity guard before it starts deleting,
// so drain the microtask queue rather than counting exact ticks.
for (let index = 0; index < 6; index += 1) {
await Promise.resolve();
}
await workerEventDelivered.promise;
expect(component.isRefreshPending(item._id)).toBe(true);
expect(component.getBusyMessage(item)).toBe(
@@ -33,14 +33,11 @@ import {
DataService,
DbOperationEvent,
isDbAbortError,
PlaybackPositionService,
PlaylistDeleteActionService,
PlaylistRefreshService,
resetHostConnectivityGuard,
RuntimeCapabilitiesService,
SortBy,
SortService,
XtreamPendingRestoreService,
} from '@iptvnator/services';
import {
PLAYLIST_UPDATE,
@@ -52,6 +49,10 @@ import {
RENDERER_PERFORMANCE_PHASE,
} from '@iptvnator/shared/logging';
import {
XtreamRefreshFlowService,
type XtreamRefreshProgressReporter,
} from '../xtream-refresh-flow.service';
import { EmptyStateComponent } from './empty-state/empty-state.component';
import type { PlaylistType } from '../add-playlist-menu/playlist-type';
import { PlaylistInfoComponent } from './playlist-info/playlist-info.component';
@@ -83,7 +84,6 @@ export class RecentPlaylistsComponent {
private readonly dialog = inject(MatDialog);
private readonly dialogService = inject(DialogService);
private readonly dataService = inject(DataService);
private readonly playbackPositionService = inject(PlaybackPositionService);
private readonly playlistRefreshService = inject(PlaylistRefreshService);
private readonly router = inject(Router);
private readonly snackBar = inject(MatSnackBar);
@@ -93,9 +93,7 @@ export class RecentPlaylistsComponent {
private readonly translate = inject(TranslateService);
private readonly playlistContext = inject(PlaylistContextFacade);
private readonly playlistDeleteAction = inject(PlaylistDeleteActionService);
private readonly pendingRestoreService = inject(
XtreamPendingRestoreService
);
private readonly xtreamRefreshFlow = inject(XtreamRefreshFlowService);
readonly sidebarMode = input(false);
readonly searchQueryInput = input<string>('');
@@ -341,125 +339,35 @@ export class RecentPlaylistsComponent {
* Refresh Xtream playlist by deleting all data and re-importing from remote
* @param item Xtream playlist to refresh
*/
async refreshXtreamPlaylist(item: PlaylistMeta) {
refreshXtreamPlaylist(item: PlaylistMeta): void {
if (this.isDeletePending(item._id) || this.isRefreshPending(item._id)) {
return;
}
this.dialogService.openConfirmDialog({
title: this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.TITLE'
),
message: this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.MESSAGE'
),
width: '400px',
onConfirm: async () => {
if (
this.isDeletePending(item._id) ||
this.isRefreshPending(item._id)
) {
return;
}
this.setPendingRefresh(item._id, true);
const operationId =
this.databaseService.createOperationId('xtream-refresh');
try {
// Before anything destructive, and for the same reason as
// the shared refresh action: this deletes the cached catalog
// and then navigates to re-import it, and an open
// connectivity guard would fast-fail that bootstrap — the
// user would be left with no catalog at all.
await resetHostConnectivityGuard(
this.dataService,
item.serverUrl
);
// Show immediate feedback — deletion can take several seconds
// for large playlists.
this.snackBar.open(
this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.STARTED'
),
undefined,
{ duration: 2000 }
);
// Delete content/categories and update the timestamp in
// parallel — both operations are fully independent.
const updateDate = Date.now();
const [restoreState, playbackPositions] = await Promise.all(
[
this.databaseService.deleteXtreamPlaylistContent(
item._id,
{
operationId,
onEvent: (workerEvent) =>
this.updateBusyOperation(
item._id,
workerEvent
),
}
),
this.playbackPositionService.getAllPlaybackPositions(
item._id
),
this.databaseService.updateXtreamPlaylistDetails({
id: item._id,
updateDate,
}),
]
);
if (
!this.pendingRestoreService.set(item._id, {
...restoreState,
playbackPositions,
})
) {
throw new Error(
`Parking pending restore state for "${item._id}" failed.`
);
}
// Update the timestamp in NgRx / IndexedDB
measureRendererPerformancePhase(
RENDERER_PERFORMANCE_PHASE.XTREAM_REFRESH_META,
() =>
this.store.dispatch(
PlaylistActions.updatePlaylistMeta({
playlist: { ...item, updateDate },
})
),
() => ({ items: 1 })
);
// Navigate to the playlist to trigger re-import
this.router.navigate(['/workspace', 'xtreams', item._id]);
} catch (error) {
if (!isDbAbortError(error)) {
console.error(
'Error refreshing Xtream playlist:',
error
);
this.snackBar.open(
this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.ERROR'
),
undefined,
{ duration: 3000 }
);
}
} finally {
this.clearBusyOperation(item._id);
this.setPendingRefresh(item._id, false);
}
},
});
this.xtreamRefreshFlow.confirmAndRefresh(
item,
this.xtreamRefreshReporter
);
}
/**
* Reports the shared destructive refresh into this page's per-row busy
* indicators, which are keyed by playlist and shared with deletion — hence
* the same `updateBusyOperation`/`clearBusyOperation` pair used there.
*/
private readonly xtreamRefreshReporter: XtreamRefreshProgressReporter = {
isBusy: (playlistId) =>
this.isDeletePending(playlistId) ||
this.isRefreshPending(playlistId),
begin: ({ playlistId }) => this.setPendingRefresh(playlistId, true),
report: ({ playlistId }, event) =>
this.updateBusyOperation(playlistId, event),
end: ({ playlistId }) => {
this.clearBusyOperation(playlistId);
this.setPendingRefresh(playlistId, false);
},
};
private async refreshM3uPlaylist(item: PlaylistMeta): Promise<void> {
if (this.isDeletePending(item._id) || this.isRefreshPending(item._id)) {
return;
@@ -0,0 +1,316 @@
import { TestBed } from '@angular/core/testing';
import { MatSnackBar } from '@angular/material/snack-bar';
import { Router } from '@angular/router';
import { Store } from '@ngrx/store';
import { TranslateService } from '@ngx-translate/core';
import { DialogService } from '@iptvnator/ui/components';
import {
DataService,
DatabaseService,
DbOperationEvent,
PlaybackPositionService,
} from '@iptvnator/services';
import {
CONNECTIVITY_GUARD_RESET,
PlaylistMeta,
} from '@iptvnator/shared/interfaces';
import {
XtreamRefreshFlowService,
type XtreamRefreshProgressReporter,
type XtreamRefreshRun,
} from './xtream-refresh-flow.service';
function createAbortError(): Error {
const error = new Error('Cancelled');
error.name = 'AbortError';
return error;
}
function createPlaylistMeta(
overrides: Partial<PlaylistMeta> = {}
): PlaylistMeta {
return {
_id: 'playlist-1',
title: 'Xtream Playlist',
serverUrl: 'http://panel.example:8080',
username: 'demo',
password: 'secret',
...overrides,
} as PlaylistMeta;
}
/**
* Records the reporter contract as the flow exercises it. `overrides` replaces
* individual hooks — notably `isBusy`, and `waitForVisibleProgress`, which the
* sources page deliberately does not implement.
*/
function createReporter(
overrides: Partial<XtreamRefreshProgressReporter> = {}
) {
const calls: string[] = [];
const runs: XtreamRefreshRun[] = [];
const events: DbOperationEvent[] = [];
const reporter: XtreamRefreshProgressReporter = {
isBusy: () => false,
begin: (run) => {
calls.push('begin');
runs.push(run);
},
report: (run, event) => {
calls.push('report');
runs.push(run);
events.push(event);
},
end: (run) => {
calls.push('end');
runs.push(run);
},
...overrides,
};
return { reporter, calls, runs, events };
}
describe('XtreamRefreshFlowService', () => {
let service: XtreamRefreshFlowService;
let confirmPromise: Promise<void> | undefined;
let order: string[];
let databaseService: {
createOperationId: jest.Mock;
deleteXtreamPlaylistContent: jest.Mock;
updateXtreamPlaylistDetails: jest.Mock;
};
let dataService: { sendIpcEvent: jest.Mock };
let dialogService: { openConfirmDialog: jest.Mock };
let playbackPositionService: { getAllPlaybackPositions: jest.Mock };
let router: { navigate: jest.Mock };
let snackBar: { open: jest.Mock };
let store: { dispatch: jest.Mock };
beforeEach(() => {
localStorage.clear();
confirmPromise = undefined;
order = [];
databaseService = {
createOperationId: jest.fn((prefix: string) => `${prefix}-op`),
deleteXtreamPlaylistContent: jest.fn(() => {
order.push('delete');
return Promise.resolve({
success: true,
favorites: [],
recentlyViewed: [],
hiddenCategories: [],
});
}),
updateXtreamPlaylistDetails: jest.fn().mockResolvedValue(true),
};
dataService = {
sendIpcEvent: jest.fn((event: string) => {
order.push(`ipc:${event}`);
return Promise.resolve({ success: true });
}),
};
dialogService = {
openConfirmDialog: jest.fn(
({ onConfirm }: { onConfirm?: () => Promise<void> }) => {
confirmPromise = onConfirm?.();
}
),
};
playbackPositionService = {
getAllPlaybackPositions: jest.fn().mockResolvedValue([]),
};
router = { navigate: jest.fn().mockResolvedValue(true) };
snackBar = { open: jest.fn() };
store = { dispatch: jest.fn() };
TestBed.configureTestingModule({
providers: [
XtreamRefreshFlowService,
{ provide: Router, useValue: router },
{ provide: Store, useValue: store },
{
provide: TranslateService,
useValue: { instant: jest.fn((key: string) => key) },
},
{ provide: MatSnackBar, useValue: snackBar },
{ provide: DialogService, useValue: dialogService },
{ provide: DatabaseService, useValue: databaseService },
{ provide: DataService, useValue: dataService },
{
provide: PlaybackPositionService,
useValue: playbackPositionService,
},
],
});
service = TestBed.inject(XtreamRefreshFlowService);
});
afterEach(() => {
jest.restoreAllMocks();
localStorage.clear();
});
it('ignores a confirmation that arrives while the playlist is already busy', async () => {
const { reporter, calls } = createReporter({ isBusy: () => true });
service.confirmAndRefresh(createPlaylistMeta(), reporter);
await confirmPromise;
expect(calls).toEqual([]);
expect(databaseService.createOperationId).not.toHaveBeenCalled();
expect(
databaseService.deleteXtreamPlaylistContent
).not.toHaveBeenCalled();
expect(dataService.sendIpcEvent).not.toHaveBeenCalled();
});
it('marks the run busy before the guard reset and any destructive work', async () => {
const { reporter, runs } = createReporter({
begin: (run) => {
order.push('begin');
runs.push(run);
},
});
service.confirmAndRefresh(createPlaylistMeta(), reporter);
// Everything up to the first await has run: the reporter is already
// busy (so an entry point rendering from `begin()` can paint) and the
// guard reset is already on the wire, while nothing destructive has
// started. `resetHostConnectivityGuard` dispatches its IPC
// synchronously and only awaits the reply, hence the reset appearing
// here rather than after the first tick.
expect(order).toEqual(['begin', `ipc:${CONNECTIVITY_GUARD_RESET}`]);
expect(runs[0]).toEqual({
playlistId: 'playlist-1',
operationId: 'xtream-refresh-op',
});
expect(
databaseService.deleteXtreamPlaylistContent
).not.toHaveBeenCalled();
await confirmPromise;
});
it('resets the connectivity guard, then waits for the reporter, then deletes', async () => {
const { reporter } = createReporter({
waitForVisibleProgress: () => {
order.push('wait');
return Promise.resolve();
},
});
const item = createPlaylistMeta();
service.confirmAndRefresh(item, reporter);
await confirmPromise;
expect(order).toEqual([
`ipc:${CONNECTIVITY_GUARD_RESET}`,
'wait',
'delete',
]);
expect(dataService.sendIpcEvent).toHaveBeenCalledWith(
CONNECTIVITY_GUARD_RESET,
{ url: item.serverUrl }
);
});
it('completes for a reporter that implements no paint hook', async () => {
const { reporter, calls } = createReporter();
const item = createPlaylistMeta();
service.confirmAndRefresh(item, reporter);
await confirmPromise;
expect(order).toEqual([`ipc:${CONNECTIVITY_GUARD_RESET}`, 'delete']);
expect(calls).toEqual(['begin', 'end']);
expect(router.navigate).toHaveBeenCalledWith([
'/workspace',
'xtreams',
item._id,
]);
});
it('forwards delete progress to the reporter under the same run', async () => {
const { reporter, events, runs } = createReporter();
const progressEvent: DbOperationEvent = {
operation: 'delete-xtream-content',
operationId: 'xtream-refresh-op',
status: 'progress',
phase: 'deleting-content',
current: 50,
total: 100,
};
databaseService.deleteXtreamPlaylistContent.mockImplementation(
(
_playlistId: string,
options?: {
operationId?: string;
onEvent?: (event: DbOperationEvent) => void;
}
) => {
options?.onEvent?.(progressEvent);
return Promise.resolve({
success: true,
favorites: [],
recentlyViewed: [],
hiddenCategories: [],
});
}
);
service.confirmAndRefresh(createPlaylistMeta(), reporter);
await confirmPromise;
expect(events).toEqual([progressEvent]);
// Every hook sees the same run, so a reporter can scope its state to
// this operation and ignore events from an older one.
expect(new Set(runs.map(({ operationId }) => operationId))).toEqual(
new Set(['xtream-refresh-op'])
);
});
it('ends the run without an error toast when the delete is aborted', async () => {
const { reporter, calls } = createReporter();
databaseService.deleteXtreamPlaylistContent.mockRejectedValue(
createAbortError()
);
service.confirmAndRefresh(createPlaylistMeta(), reporter);
await confirmPromise;
expect(calls).toEqual(['begin', 'end']);
expect(snackBar.open).not.toHaveBeenCalledWith(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.ERROR',
undefined,
{ duration: 3000 }
);
expect(router.navigate).not.toHaveBeenCalled();
});
it('reports a failed delete and still ends the run', async () => {
const consoleErrorSpy = jest
.spyOn(console, 'error')
.mockImplementation();
const { reporter, calls } = createReporter();
databaseService.deleteXtreamPlaylistContent.mockRejectedValue(
new Error('Refresh failed')
);
service.confirmAndRefresh(createPlaylistMeta(), reporter);
await confirmPromise;
expect(calls).toEqual(['begin', 'end']);
expect(snackBar.open).toHaveBeenCalledWith(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.ERROR',
undefined,
{ duration: 3000 }
);
expect(router.navigate).not.toHaveBeenCalled();
expect(consoleErrorSpy).toHaveBeenCalled();
});
});
@@ -0,0 +1,192 @@
import { inject, Injectable } from '@angular/core';
import { Router } from '@angular/router';
import { Store } from '@ngrx/store';
import { TranslateService } from '@ngx-translate/core';
import { MatSnackBar } from '@angular/material/snack-bar';
import { DialogService } from '@iptvnator/ui/components';
import {
DataService,
DatabaseService,
type DbOperationEvent,
isDbAbortError,
PlaybackPositionService,
resetHostConnectivityGuard,
XtreamPendingRestoreService,
} from '@iptvnator/services';
import { PlaylistActions } from '@iptvnator/m3u-state';
import { PlaylistMeta } from '@iptvnator/shared/interfaces';
import {
measureRendererPerformancePhase,
RENDERER_PERFORMANCE_PHASE,
} from '@iptvnator/shared/logging';
/** Identifies one confirmed run, so a reporter can scope its own state to it. */
export interface XtreamRefreshRun {
readonly playlistId: string;
readonly operationId: string;
}
/**
* How an entry point shows that the refresh is working. The flow itself owns no
* busy state: the header action renders one global preparation overlay while
* the sources page renders a per-row indicator, and those cannot be unified
* without changing what either page looks like.
*/
export interface XtreamRefreshProgressReporter {
/**
* Re-checked after the user confirms. The dialog is asynchronous, so the
* playlist may have picked up another operation (or this very refresh) in
* the meantime.
*/
isBusy(playlistId: string): boolean;
/** Marks the run busy. Runs synchronously, before anything is awaited. */
begin(run: XtreamRefreshRun): void;
/** Progress from the delete worker. */
report(run: XtreamRefreshRun, event: DbOperationEvent): void;
/** Clears the busy state. Always runs, including after an abort. */
end(run: XtreamRefreshRun): void;
/**
* Optional pause between `begin()` and the destructive work, for a reporter
* whose indicator needs a paint opportunity — a small cached playlist can
* otherwise finish deleting before its progress UI ever appears.
*/
waitForVisibleProgress?(): Promise<void>;
}
/**
* The destructive Xtream refresh: drop the cached catalog, park the user data
* that must survive it, and navigate so the portal route re-imports everything.
*
* This lives on its own because it has two entry points — the shared playlist
* refresh action and the Workspace sources page — that differ only in how they
* report progress. They used to be two near-identical implementations, and a
* connectivity-guard fix applied to one of them looked applied to both (#1421).
*/
@Injectable({ providedIn: 'root' })
export class XtreamRefreshFlowService {
private readonly router = inject(Router);
private readonly store = inject(Store);
private readonly translate = inject(TranslateService);
private readonly snackBar = inject(MatSnackBar);
private readonly dialogService = inject(DialogService);
private readonly databaseService = inject(DatabaseService);
private readonly dataService = inject(DataService);
private readonly playbackPositionService = inject(PlaybackPositionService);
private readonly pendingRestoreService = inject(
XtreamPendingRestoreService
);
/**
* Asks for confirmation, then deletes and re-imports the playlist. Callers
* keep their own entry-point guard (a disabled button, a pending row) —
* this only re-checks `reporter.isBusy()` once the dialog is confirmed.
*/
confirmAndRefresh(
item: PlaylistMeta,
reporter: XtreamRefreshProgressReporter
): void {
this.dialogService.openConfirmDialog({
title: this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.TITLE'
),
message: this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.MESSAGE'
),
width: '400px',
onConfirm: () => this.runRefresh(item, reporter),
});
}
private async runRefresh(
item: PlaylistMeta,
reporter: XtreamRefreshProgressReporter
): Promise<void> {
if (reporter.isBusy(item._id)) {
return;
}
const operationId =
this.databaseService.createOperationId('xtream-refresh');
const run: XtreamRefreshRun = { playlistId: item._id, operationId };
reporter.begin(run);
try {
// Before anything destructive: this deletes the cached catalog and
// then forces a route bootstrap whose status request an open
// connectivity guard would fast-fail, leaving the user with no
// catalog at all until the cooldown expires. Both entry points
// reach the reset through here, which is the point of the shared
// flow — see docs/architecture/host-connectivity-guard.md.
await resetHostConnectivityGuard(this.dataService, item.serverUrl);
// Show immediate feedback — deletion can take several seconds for
// large playlists.
this.snackBar.open(
this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.STARTED'
),
undefined,
{ duration: 2000 }
);
await reporter.waitForVisibleProgress?.();
// Delete content/categories and update the timestamp in parallel —
// all three operations are fully independent.
const updateDate = Date.now();
const [restoreState, playbackPositions] = await Promise.all([
this.databaseService.deleteXtreamPlaylistContent(item._id, {
operationId,
onEvent: (event) => reporter.report(run, event),
}),
this.playbackPositionService.getAllPlaybackPositions(item._id),
this.databaseService.updateXtreamPlaylistDetails({
id: item._id,
updateDate,
}),
]);
if (
!this.pendingRestoreService.set(item._id, {
...restoreState,
playbackPositions,
})
) {
throw new Error(
`Parking pending restore state for "${item._id}" failed.`
);
}
// Update the timestamp in NgRx / IndexedDB
measureRendererPerformancePhase(
RENDERER_PERFORMANCE_PHASE.XTREAM_REFRESH_META,
() =>
this.store.dispatch(
PlaylistActions.updatePlaylistMeta({
playlist: { ...item, updateDate },
})
),
() => ({ items: 1 })
);
// Navigate to the playlist to trigger re-import
await this.router.navigate(['/workspace', 'xtreams', item._id]);
} catch (error) {
if (!isDbAbortError(error)) {
console.error('Error refreshing Xtream playlist:', error);
this.snackBar.open(
this.translate.instant(
'HOME.PLAYLISTS.REFRESH_XTREAM_DIALOG.ERROR'
),
undefined,
{ duration: 3000 }
);
}
} finally {
reporter.end(run);
}
}
}