From e799fcec32c779cc10daefac39aacd56eab7738f Mon Sep 17 00:00:00 2001 From: 4gray Date: Sat, 5 Sep 2026 22:37:11 +0200 Subject: [PATCH] fix(epg): close source reconciliation review races --- AGENTS.md | 2 +- CLAUDE.md | 2 +- .../src/xtream-epg.e2e.ts | 96 +++++++++++++++++ .../src/app/events/epg-fetch.service.ts | 18 +++- .../src/app/events/epg-source-generation.ts | 20 +++- .../epg-source-settings.service.spec.ts | 41 ++++++- .../app/events/epg-source-settings.service.ts | 13 ++- .../src/app/events/epg-worker.service.ts | 15 ++- .../src/app/events/epg.events.spec.ts | 12 +++ docs/architecture/m3u-playlist-module.md | 7 +- docs/architecture/stalker-epg.md | 30 +++--- .../src/lib/epg-progress.service.spec.ts | 23 ++++ .../src/lib/epg-progress.service.ts | 6 ++ .../features/with-stalker-epg.feature.spec.ts | 101 +++++++++++++++++- .../features/with-stalker-epg.feature.ts | 47 ++++++-- .../src/lib/electron-api.interface.ts | 2 + .../epg-progress-panel.component.ts | 2 + 17 files changed, 396 insertions(+), 41 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index b5c7c0d19..c928a5d5b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -182,7 +182,7 @@ settings load and playlist migration; failed settings reads and incomplete playlist migration never authorize pruning. Removed sources retire queued and running imports before worker-owned deletion. Shared channel IDs survive while another source has programmes; manual mappings remain user preferences, but no -longer resolve deleted data. Renderer lookup generations and Xtream preview +longer resolve deleted data. Renderer lookup generations, Xtream previews and Stalker mapping-cache invalidation prevent late results from restoring removed programmes. Provider EPG is independent. See `docs/architecture/m3u-playlist-module.md` ("XMLTV source lifecycle"). diff --git a/CLAUDE.md b/CLAUDE.md index d95cf73c2..3c759fe1a 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1667,7 +1667,7 @@ settings load and playlist migration; failed settings reads and incomplete playlist migration never authorize pruning. Removed sources retire queued and running imports before worker-owned deletion. Shared channel IDs survive while another source has programmes; manual mappings remain user preferences, but no -longer resolve deleted data. Renderer lookup generations and Xtream preview +longer resolve deleted data. Renderer lookup generations, Xtream previews and Stalker mapping-cache invalidation prevent late results from restoring removed programmes. Provider EPG is independent. See `docs/architecture/m3u-playlist-module.md` ("XMLTV source lifecycle"). diff --git a/apps/electron-backend-e2e/src/xtream-epg.e2e.ts b/apps/electron-backend-e2e/src/xtream-epg.e2e.ts index d2a9ae7c4..02eeb28e3 100644 --- a/apps/electron-backend-e2e/src/xtream-epg.e2e.ts +++ b/apps/electron-backend-e2e/src/xtream-epg.e2e.ts @@ -116,6 +116,102 @@ test('@epg @xtream @electron removes uploaded guide data and restores provider E } }); +test('@epg @stalker @electron invalidates a loaded manual XMLTV mapping after source removal', async ({ + dataDir, + request, +}) => { + test.setTimeout(120000); + await resetMockServers(request, ['stalker']); + const fixture = await fetchStalkerCategoryFixture(request, 'itv'); + const item = fixture.items[0]; + const stamp = (date: Date) => + date.toISOString().replace(/[-:T]/g, '').slice(0, 14); + const source = await createMutableTextServer( + `Mapped Guide + Retired Stalker Bulletin`, + { + contentType: 'application/xml', + resourcePath: '/stalker.xml', + } + ); + const app = await launchElectronApp(dataDir); + try { + await app.mainWindow.route('https://test-streams.mux.dev/**', () => { + // Keep external media pending while the local guide is exercised. + }); + await addStalkerPortal(app.mainWindow, { + name: 'Stalker XMLTV Removal', + }); + await waitForStalkerCatalog(app.mainWindow); + const playlistId = new URL(app.mainWindow.url()).pathname.match( + /\/stalker\/([^/]+)/ + )?.[1]; + expect(playlistId).toBeTruthy(); + const key = `stalker:${decodeURIComponent(playlistId!)}:${String(item.id).trim()}`; + await app.mainWindow.evaluate( + (mappingKey) => + window.electron.setEpgMapping(mappingKey, 'stalker-mapped'), + key + ); + await openSettings(app.mainWindow); + await openSettingsSection(app.mainWindow, 'epg'); + await app.mainWindow + .getByRole('button', { name: 'Add EPG source' }) + .click(); + await app.mainWindow + .locator('.epg-source-row input') + .fill(source.resourceUrl); + await saveSettings(app.mainWindow); + await expect + .poll( + () => + app.mainWindow.evaluate(async () => + ( + await window.electron.getChannelPrograms( + 'stalker-mapped' + ) + ).map((p) => p.title) + ), + { timeout: 30000 } + ) + .toContain('Retired Stalker Bulletin'); + await openWorkspaceSection(app.mainWindow, 'Live TV'); + await clickCategoryByNameExact(app.mainWindow, fixture.categoryName); + const row = channelItemByTitle( + app.mainWindow, + item.o_name || item.name || '' + ).first(); + await expect(row.locator('.epg-title')).toHaveText( + 'Retired Stalker Bulletin' + ); + await row.click(); + await expect + .poll(() => timelineBlockTitles(app.mainWindow)) + .toContain('Retired Stalker Bulletin'); + await openSettings(app.mainWindow); + await openSettingsSection(app.mainWindow, 'epg'); + await app.mainWindow.locator('.epg-source-row button').nth(1).click(); + await saveSettings(app.mainWindow); + await openWorkspaceSection(app.mainWindow, 'Live TV'); + await clickCategoryByNameExact(app.mainWindow, fixture.categoryName); + await expect(row).toBeVisible(); + await row.click(); + await expect + .poll(() => timelineBlockTitles(app.mainWindow)) + .not.toContain('Retired Stalker Bulletin'); + await expect(row.locator('.epg-title')).toHaveCount(0); + expect( + await app.mainWindow.evaluate( + (mappingKey) => window.electron.getEpgMapping(mappingKey), + key + ) + ).toMatchObject({ epgChannelId: 'stalker-mapped' }); + } finally { + await closeElectronApp(app); + await source.close(); + } +}); + for (const timeZone of ['UTC', 'Europe/Berlin'] as const) { test(`@epg @xtream @electron renders Xtream EPG previews and the timeline schedule in ${timeZone}`, async ({ dataDir, diff --git a/apps/electron-backend/src/app/events/epg-fetch.service.ts b/apps/electron-backend/src/app/events/epg-fetch.service.ts index 1c54c7bdf..4674d772f 100644 --- a/apps/electron-backend/src/app/events/epg-fetch.service.ts +++ b/apps/electron-backend/src/app/events/epg-fetch.service.ts @@ -1,4 +1,4 @@ -import { epgSourceGeneration } from './epg-source-generation'; +import { epgSourceGeneration, requestEpgSource } from './epg-source-generation'; import { eq } from 'drizzle-orm'; import { ElectronBridgeTrustOptions } from '@iptvnator/shared/interfaces'; import { getDatabase } from '../database/connection'; @@ -102,7 +102,7 @@ export async function handleFetchEpg( .filter((url) => url?.trim()) .map((url) => url.trim()); const generations = new Map( - validUrls.map((url) => [url, epgSourceGeneration(url)]) + validUrls.map((url) => [url, requestEpgSource(url)]) ); if (validUrls.length === 0) { @@ -152,7 +152,19 @@ export async function handleFetchEpg( const errors: string[] = []; for (const url of urlsToFetch) { try { - if (generations.get(url) !== epgSourceGeneration(url)) continue; + if (generations.get(url) !== epgSourceGeneration(url)) { + epgWorkerService.sendProgressToRenderer( + url, + 'cancelled', + undefined, + undefined, + undefined, + undefined, + undefined, + generations.get(url) + ); + continue; + } await epgWorkerService.fetchEpgFromUrl(url, options); } catch (error) { console.error( diff --git a/apps/electron-backend/src/app/events/epg-source-generation.ts b/apps/electron-backend/src/app/events/epg-source-generation.ts index 5dbd9b721..be9aa620f 100644 --- a/apps/electron-backend/src/app/events/epg-source-generation.ts +++ b/apps/electron-backend/src/app/events/epg-source-generation.ts @@ -1,13 +1,23 @@ /** Retires queued imports as well as workers already parsing a removed URL. */ const generations = new Map(); +const requests = new Map(); export function epgSourceGeneration(url: string): number { - const key = url.trim(); - if (!generations.has(key)) generations.set(key, 0); - return generations.get(key)!; + return generations.get(url.trim()) ?? 0; +} +export function requestEpgSource(url: string): number { + requests.set(url.trim(), Symbol()); + return epgSourceGeneration(url); } export function retireEpgSource(url: string): void { generations.set(url.trim(), epgSourceGeneration(url) + 1); } -export function requestedEpgSources(): string[] { - return [...generations.keys()]; +export function requestedEpgSources(): Map { + return new Map(requests); +} +export function forgetEpgSourceRequest( + url: string, + request: symbol | undefined +): void { + // A request received during cleanup still needs reconciliation next time. + if (requests.get(url.trim()) === request) requests.delete(url.trim()); } diff --git a/apps/electron-backend/src/app/events/epg-source-settings.service.spec.ts b/apps/electron-backend/src/app/events/epg-source-settings.service.spec.ts index 2ffbdae9e..0be6d2991 100644 --- a/apps/electron-backend/src/app/events/epg-source-settings.service.spec.ts +++ b/apps/electron-backend/src/app/events/epg-source-settings.service.spec.ts @@ -6,7 +6,7 @@ import { } from '../database/schema'; import { getDatabase } from '../database/connection'; import { epgWorkerService } from './epg-worker.service'; -import { epgSourceGeneration } from './epg-source-generation'; +import { epgSourceGeneration, requestEpgSource } from './epg-source-generation'; import { reconcileEpgSources } from './epg-source-settings.service'; jest.mock('../database/connection', () => ({ getDatabase: jest.fn() })); @@ -65,6 +65,45 @@ describe('committed EPG source reconciliation', () => { } }); + it('does not clear historical request keys again after successful cleanup', async () => { + requestEpgSource('removed'); + await reconcileEpgSources([]); + rows.set(epgChannels, []); + rows.set(epgPrograms, []); + jest.clearAllMocks(); + await reconcileEpgSources([]); + expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalled(); + }); + + it('retries a failed cleanup even when the queued source had no database rows', async () => { + rows.set(epgChannels, []); + rows.set(epgPrograms, []); + requestEpgSource('queued-only'); + ( + epgWorkerService.clearEpgDataForSource as jest.Mock + ).mockRejectedValueOnce(new Error('worker failure')); + await expect(reconcileEpgSources([])).rejects.toThrow('worker failure'); + await reconcileEpgSources([]); + expect(epgWorkerService.clearEpgDataForSource).toHaveBeenCalledTimes(2); + }); + + it('preserves a new request received while the previous request is being cleared', async () => { + rows.set(epgChannels, []); + rows.set(epgPrograms, []); + requestEpgSource('requested-again'); + ( + epgWorkerService.clearEpgDataForSource as jest.Mock + ).mockImplementationOnce(async () => { + requestEpgSource('requested-again'); + }); + await reconcileEpgSources([]); + await reconcileEpgSources([]); + expect(epgWorkerService.clearEpgDataForSource).toHaveBeenCalledTimes(2); + jest.clearAllMocks(); + await reconcileEpgSources([]); + expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalled(); + }); + it('does not prune sources before playlist migration succeeds', async () => { rows.set(appState, []); await expect(reconcileEpgSources([])).rejects.toThrow( diff --git a/apps/electron-backend/src/app/events/epg-source-settings.service.ts b/apps/electron-backend/src/app/events/epg-source-settings.service.ts index b439bc228..f51bc77fb 100644 --- a/apps/electron-backend/src/app/events/epg-source-settings.service.ts +++ b/apps/electron-backend/src/app/events/epg-source-settings.service.ts @@ -7,7 +7,11 @@ import { playlists, } from '../database/schema'; import { epgWorkerService } from './epg-worker.service'; -import { requestedEpgSources, retireEpgSource } from './epg-source-generation'; +import { + forgetEpgSourceRequest, + requestedEpgSources, + retireEpgSource, +} from './epg-source-generation'; let reconciliation = Promise.resolve(); @@ -53,12 +57,13 @@ export function reconcileEpgSources(globalUrls: string[]): Promise { const programs = await db .selectDistinct({ url: epgPrograms.sourceUrl }) .from(epgPrograms); + const requested = requestedEpgSources(); const removed = [ ...new Set( [ ...channels.map((row) => row.url), ...programs.map((row) => row.url), - ...requestedEpgSources(), + ...requested.keys(), ].filter( (url): url is string => !!url?.trim() && !active.has(url.trim()) @@ -67,8 +72,10 @@ export function reconcileEpgSources(globalUrls: string[]): Promise { ]; // Fence the whole obsolete set before awaiting the first worker exit. removed.forEach(retireEpgSource); - for (const url of removed) + for (const url of removed) { await epgWorkerService.clearEpgDataForSource(url); + forgetEpgSourceRequest(url, requested.get(url.trim())); + } }); reconciliation = next; return next; diff --git a/apps/electron-backend/src/app/events/epg-worker.service.ts b/apps/electron-backend/src/app/events/epg-worker.service.ts index f52c62fb6..f5e7e143d 100644 --- a/apps/electron-backend/src/app/events/epg-worker.service.ts +++ b/apps/electron-backend/src/app/events/epg-worker.service.ts @@ -1,4 +1,8 @@ -import { epgSourceGeneration, retireEpgSource } from './epg-source-generation'; +import { + epgSourceGeneration, + requestEpgSource, + retireEpgSource, +} from './epg-source-generation'; import { app, BrowserWindow } from 'electron'; import * as path from 'path'; import { pathToFileURL } from 'url'; @@ -9,7 +13,8 @@ import { } from '@iptvnator/shared/interfaces'; import { resolveWorkerRuntimeBootstrap } from '../workers/worker-runtime-paths'; -export type EpgProgressStatus = 'queued' | 'loading' | 'complete' | 'error'; +export type EpgProgressStatus = + 'queued' | 'loading' | 'complete' | 'error' | 'cancelled'; export interface EpgProgressStats { totalChannels: number; @@ -59,12 +64,14 @@ export class EpgWorkerService { error?: string, queuePosition?: number, errorCode?: ElectronBridgeSecurityErrorCode, - errorHost?: string + errorHost?: string, + generation = epgSourceGeneration(url) ): void { const windows = BrowserWindow.getAllWindows(); windows.forEach((win) => { win.webContents.send('EPG_PROGRESS_UPDATE', { url, + generation, status, stats, error, @@ -114,7 +121,7 @@ export class EpgWorkerService { url: string, options: ElectronBridgeTrustOptions ): Promise { - const generation = epgSourceGeneration(url); + const generation = requestEpgSource(url); return new Promise((resolve, reject) => { let worker: Worker; try { diff --git a/apps/electron-backend/src/app/events/epg.events.spec.ts b/apps/electron-backend/src/app/events/epg.events.spec.ts index 81d847d87..a661a9d23 100644 --- a/apps/electron-backend/src/app/events/epg.events.spec.ts +++ b/apps/electron-backend/src/app/events/epg.events.spec.ts @@ -319,6 +319,7 @@ describe('EpgEvents', () => { const { handleFetchEpg } = await import('./epg-fetch.service'); const { retireEpgSource } = await import('./epg-source-generation'); const { epgWorkerService } = await import('./epg-worker.service'); + const progress = jest.spyOn(epgWorkerService, 'sendProgressToRenderer'); let finishFirst!: () => void; const fetch = jest .spyOn(epgWorkerService, 'fetchEpgFromUrl') @@ -338,6 +339,17 @@ describe('EpgEvents', () => { finishFirst(); await request; expect(fetch).toHaveBeenCalledTimes(1); + expect(progress).toHaveBeenCalledWith( + 'https://queued.example/guide.xml', + 'cancelled', + undefined, + undefined, + undefined, + undefined, + undefined, + expect.any(Number) + ); + progress.mockRestore(); fetch.mockRestore(); }); diff --git a/docs/architecture/m3u-playlist-module.md b/docs/architecture/m3u-playlist-module.md index 4d54d09ef..a0ab67b29 100644 --- a/docs/architecture/m3u-playlist-module.md +++ b/docs/architecture/m3u-playlist-module.md @@ -1331,7 +1331,9 @@ not global XMLTV owners. No source-discovery or provider matching policy changes Reconciliation finds old sources in both XMLTV tables and queued imports. It retires their generations before waiting for workers to exit, then uses the -existing source-clear worker. Programmes are deleted by source; a globally keyed +existing source-clear worker. Successfully cleared request candidates are forgotten +without resetting their generation fences; failed cleanups remain retryable. +Retired queued imports emit cancellation so progress rows disappear. Programmes are deleted by source; a globally keyed channel is retained while another source still has programmes, transferring its legacy owner to that remaining source. Manual mappings are preserved and can resolve another retained source sharing that channel ID. Legacy programmes with @@ -1341,7 +1343,8 @@ There is no new schema migration. Renderer reconciliation increments a data revision and cancels earlier lookup subscriptions, clears program caches and the selected M3U guide, and refreshes -Xtream selection and visible channel previews. A delayed startup import is +Xtream selection and visible channel previews, plus Stalker manual mapping +overrides and bulk guides. A delayed startup import is started only if its source still belongs to the reconciled configuration; its completion observer is installed after settings initialization. Provider EPG continues through its existing APIs. Playlist refresh is not EPG cache cleanup. diff --git a/docs/architecture/stalker-epg.md b/docs/architecture/stalker-epg.md index 32f321bfc..73f069877 100644 --- a/docs/architecture/stalker-epg.md +++ b/docs/architecture/stalker-epg.md @@ -24,13 +24,13 @@ Stalker now uses two EPG paths with different purposes: settled bulk guide cannot answer fall back to throttled per-channel `get_short_epg` through `StalkerEpgPreviewQueue` (see "Channel row preview flow"). - - Effect ordering matters: the eager-EPG effect is registered **after** the - playlist-change effect that calls `clearBulkItvEpgCache()`. On a portal - switch the cache is cleared first and then refilled; if the order is - reversed the clear clobbers the just-loaded bulk EPG on initial render. - - `ensureBulkItvEpg` de-duplicates (via `isLoadingBulkItvEpg` / - `bulkItvEpgLoaded` + matching playlist/period), so the eager trigger and the - play-time `loadEpgForChannel` path never double-fetch. + - Effect ordering matters: the eager-EPG effect is registered **after** the + playlist-change effect that calls `clearBulkItvEpgCache()`. On a portal + switch the cache is cleared first and then refilled; if the order is + reversed the clear clobbers the just-loaded bulk EPG on initial render. + - `ensureBulkItvEpg` de-duplicates (via `isLoadingBulkItvEpg` / + `bulkItvEpgLoaded` + matching playlist/period), so the eager trigger and the + play-time `loadEpgForChannel` path never double-fetch. - If a portal does not return usable bulk data for the selected channel, the active panel falls back to `get_short_epg`. @@ -196,12 +196,12 @@ Normalization rules: ### Key files -| File | Purpose | -| -------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------- | -| `libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts` | bulk cache and fallback handling | -| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.ts` | active-channel EPG loading and controlled `app-epg-timeline` wiring | -| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.html` | active panel template | -| `libs/ui/epg/src/lib/epg-timeline/epg-timeline.component.ts` | shared controlled EPG timeline with date navigator | +| File | Purpose | +| ---------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------- | +| `libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts` | bulk cache and fallback handling | +| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.ts` | active-channel EPG loading and controlled `app-epg-timeline` wiring | +| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.html` | active panel template | +| `libs/ui/epg/src/lib/epg-timeline/epg-timeline.component.ts` | shared controlled EPG timeline with date navigator | ### Store API @@ -283,6 +283,10 @@ rows stop consuming portal request capacity. - Bulk EPG is fetched once per playlist session - Channel switches only read from `bulkItvEpgByChannel` - The cache is cleared when the Stalker playlist changes +- Committed XMLTV source reconciliation also clears mapping overrides, checked + IDs and bulk data, reloads the guide and selected mapping, and fences pending + mapping/bulk replies. Saved mappings remain authoritative even when removal + leaves their guide empty; portal EPG is not mixed into that empty override. - This implementation does not add TTL-based refresh or background polling ## Authentication diff --git a/libs/epg/data-access/src/lib/epg-progress.service.spec.ts b/libs/epg/data-access/src/lib/epg-progress.service.spec.ts index 434194761..86915668c 100644 --- a/libs/epg/data-access/src/lib/epg-progress.service.spec.ts +++ b/libs/epg/data-access/src/lib/epg-progress.service.spec.ts @@ -63,6 +63,29 @@ describe('EpgProgressService', () => { expect(epgBridge.onProgress).not.toHaveBeenCalled(); }); + it('removes a retired queued import immediately on cancellation', () => { + epgBridge.supportsProgress = true; + const service = configureService(); + const listener = (epgBridge.onProgress as jest.Mock).mock.calls[0][0]; + const url = 'https://removed.example/guide.xml'; + listener({ url, status: 'queued' }); + expect(service.queuedCount()).toBe(1); + listener({ url, status: 'cancelled' }); + expect(service.queuedCount()).toBe(0); + expect(service.isVisible()).toBe(false); + }); + + it('keeps a replacement import when an older queued batch finally cancels', () => { + epgBridge.supportsProgress = true; + const service = configureService(); + const listener = (epgBridge.onProgress as jest.Mock).mock.calls[0][0]; + const url = 'https://readded.example/guide.xml'; + listener({ url, status: 'queued', generation: 0 }); + listener({ url, status: 'loading', generation: 2 }); + listener({ url, status: 'cancelled', generation: 0 }); + expect(service.activeCount()).toBe(1); + }); + it('does not force retry when EPG data management is disabled', () => { const service = configureService(); diff --git a/libs/epg/data-access/src/lib/epg-progress.service.ts b/libs/epg/data-access/src/lib/epg-progress.service.ts index 77292f314..eb71abe5c 100644 --- a/libs/epg/data-access/src/lib/epg-progress.service.ts +++ b/libs/epg/data-access/src/lib/epg-progress.service.ts @@ -108,6 +108,12 @@ export class EpgProgressService { } private updateProgress(progress: EpgImportProgress): void { + const current = this.importsMap().get(progress.url); + if ((progress.generation ?? 0) < (current?.generation ?? 0)) return; + if (progress.status === 'cancelled') { + this.removeImport(progress.url); + return; + } this.importsMap.update((current) => { const updated = new Map(current); updated.set(progress.url, progress); diff --git a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.spec.ts b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.spec.ts index 4573c561c..7b84bbd54 100644 --- a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.spec.ts +++ b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.spec.ts @@ -1,6 +1,10 @@ import { TestBed } from '@angular/core/testing'; import { signalStore, withState } from '@ngrx/signals'; -import { DataService, RuntimeCapabilitiesService } from '@iptvnator/services'; +import { + DataService, + EpgSourceSettingsService, + RuntimeCapabilitiesService, +} from '@iptvnator/services'; import { EpgRuntimeBridgeService } from '@iptvnator/epg/data-access'; import { EpgItem, Playlist } from '@iptvnator/shared/interfaces'; import { StalkerSessionService } from '../../stalker-session.service'; @@ -232,6 +236,44 @@ describe('withStalkerEpg', () => { expect(epgBridge.getEpgMappingsBatch).toHaveBeenCalledTimes(1); }); + it('drops removed-source overrides and reloads the selected mapping without a portal reset', async () => { + epgBridge.getEpgMappingsBatch.mockResolvedValue({ + 'stalker:playlist-1:10001': 'mapped.channel.id', + }); + epgBridge.getChannelPrograms.mockResolvedValue([MAPPED_PROGRAM]); + await store.applyMappedItvEpg(['10001']); + epgBridge.getChannelPrograms.mockResolvedValue([]); + const sources = TestBed.inject(EpgSourceSettingsService); + sources.revision.update((value) => value + 1); + sources.changed$.next(); + expect(store.selectedItvEpgPrograms()).toEqual([]); + await store.applyMappedItvEpg(['10001']); + expect(store.selectedItvEpgPrograms()).toEqual([]); + // Saved mappings remain authoritative even when their source is gone. + expect(store.hasItvEpgMappingOverride('10001')).toBe(true); + }); + + it('ignores a mapped-program response that completes after source invalidation', async () => { + epgBridge.getEpgMappingsBatch.mockResolvedValue({ + 'stalker:playlist-1:10001': 'mapped.channel.id', + }); + let resolvePrograms!: (value: unknown) => void; + epgBridge.getChannelPrograms.mockImplementationOnce( + () => + new Promise((resolve) => { + resolvePrograms = resolve; + }) + ); + const pending = store.applyMappedItvEpg(['10001']); + await Promise.resolve(); + const sources = TestBed.inject(EpgSourceSettingsService); + sources.revision.update((value) => value + 1); + sources.changed$.next(); + resolvePrograms([MAPPED_PROGRAM]); + await pending; + expect(store.selectedItvEpgPrograms()).toEqual([]); + }); + it('keeps overrides when ensureBulkItvEpg replaces the bulk record', async () => { epgBridge.getEpgMappingsBatch.mockResolvedValue({ 'stalker:playlist-1:10001': 'mapped.channel.id', @@ -259,6 +301,63 @@ describe('withStalkerEpg', () => { expect(record['10002']?.length).toBeGreaterThan(0); }); + it('does not mark later IDs checked after a stale mapped lookup rejects', async () => { + epgBridge.getEpgMappingsBatch.mockResolvedValue({ + 'stalker:playlist-1:10001': 'mapped.channel.id', + }); + let rejectPrograms!: (error: Error) => void; + epgBridge.getChannelPrograms.mockImplementationOnce( + () => + new Promise((_, reject) => { + rejectPrograms = reject; + }) + ); + const pending = store.applyMappedItvEpg(['10001', '10002']); + await Promise.resolve(); + store.clearBulkItvEpgCache(); + await store.applyMappedItvEpg(['10003']); + rejectPrograms(new Error('old request failed')); + await pending; + epgBridge.getEpgMappingsBatch.mockResolvedValue({ + 'stalker:playlist-1:10002': 'new.channel.id', + }); + epgBridge.getChannelPrograms.mockResolvedValue([MAPPED_PROGRAM]); + await store.applyMappedItvEpg(['10002']); + expect(store.bulkItvEpgByChannel()['10002']).toEqual([ + { ...MAPPED_PROGRAM, channel: '10002' }, + ]); + }); + + it('allows initial bulk loading and mapping lookup to finish concurrently', async () => { + let finishBulk!: (value: unknown) => void; + const response = new Promise((resolve) => { + finishBulk = resolve; + }); + dataService.sendIpcEvent.mockReturnValueOnce(response); + const bulk = store.ensureBulkItvEpg(); + epgBridge.getEpgMappingsBatch.mockResolvedValue({ + 'stalker:playlist-1:10001': 'mapped.channel.id', + }); + epgBridge.getChannelPrograms.mockResolvedValue([]); + await store.applyMappedItvEpg(['10001']); + finishBulk({ + js: { + '10001': [ + buildEntry( + '10001', + 'Portal Show', + 1744365600, + 1744367400 + ), + ], + }, + }); + await bulk; + expect(store.isLoadingBulkItvEpg()).toBe(false); + expect(store.bulkItvEpgLoaded()).toBe(true); + expect(store.selectedItvEpgPrograms()).toEqual([]); + }); + it('does nothing when the mapping bridge is unsupported', async () => { epgBridge.supportsEpgMapping = false; diff --git a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts index 1fb216a31..6ed086652 100644 --- a/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts +++ b/libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts @@ -3,12 +3,17 @@ import { patchState, signalStoreFeature, withComputed, + withHooks, withMethods, withState, } from '@ngrx/signals'; import { EpgRuntimeBridgeService } from '@iptvnator/epg/data-access'; import { createLogger } from '@iptvnator/portal/shared/util'; -import { DataService, RuntimeCapabilitiesService } from '@iptvnator/services'; +import { + DataService, + EpgSourceSettingsService, + RuntimeCapabilitiesService, +} from '@iptvnator/services'; import { buildStalkerEpgMappingKey, EpgItem, @@ -105,7 +110,8 @@ export function withStalkerEpg() { stalkerSession = inject(StalkerSessionService), portalRepair = inject(StalkerPortalRepairService), runtime = inject(RuntimeCapabilitiesService), - epgBridge = inject(EpgRuntimeBridgeService) + epgBridge = inject(EpgRuntimeBridgeService), + sources = inject(EpgSourceSettingsService) ) => { const storeContext = store as typeof store & StalkerEpgFeatureStoreContract; @@ -127,6 +133,7 @@ export function withStalkerEpg() { // an empty mapped guide must still keep the portal EPG out. const mappingOwnedIds = new Set(); let mappingPlaylistId: string | null = null; + let cacheGeneration = 0; const resetMappingOverrides = (): void => { mappingOverridesById.clear(); @@ -140,8 +147,8 @@ export function withStalkerEpg() { EpgProgram[] > => { const record: Record = {}; - for (const [id, programs] of mappingOverridesById) { - record[id] = programs; + for (const id of mappingOwnedIds) { + record[id] = mappingOverridesById.get(id) ?? []; } return record; }; @@ -222,6 +229,7 @@ export function withStalkerEpg() { } const playlistId = String(playlist._id); + const generation = cacheGeneration; if (!supportsEpg()) { patchState(store, { bulkItvEpgByChannel: {}, @@ -258,6 +266,7 @@ export function withStalkerEpg() { type: 'itv', period: String(periodHours), }); + if (generation !== cacheGeneration) return; const selectedChannelId = storeContext.selectedItvId() ?? null; const bulkPrograms = extractBulkEpgByChannel( @@ -274,6 +283,7 @@ export function withStalkerEpg() { isLoadingBulkItvEpg: false, }); } catch (error) { + if (generation !== cacheGeneration) return; logger.warn('Bulk Stalker EPG unavailable', error); patchState(store, { bulkItvEpgByChannel: mappingOverridesRecord(), @@ -312,8 +322,7 @@ export function withStalkerEpg() { channelIds .map((id) => normalizeStalkerEntityId(id)) .filter( - (id) => - id && !mappingCheckedIds.has(id) + (id) => id && !mappingCheckedIds.has(id) ) ), ]; @@ -325,7 +334,11 @@ export function withStalkerEpg() { // call can outlive a portal switch — bail out after // every await instead of writing portal A's data // into portal B's state. + const revision = sources.revision(); + const generation = cacheGeneration; const isStale = (): boolean => + revision !== sources.revision() || + generation !== cacheGeneration || mappingPlaylistId !== playlistId || String( storeContext.currentPlaylist()?._id ?? '' @@ -365,6 +378,7 @@ export function withStalkerEpg() { let changed = false; let ownershipChanged = false; for (const [channelId, key] of keyById) { + if (isStale()) return; const mappedEpgId = mappings[key]?.trim(); if (!mappedEpgId) { // No mapping for this channel — a stable @@ -443,12 +457,31 @@ export function withStalkerEpg() { }, clearBulkItvEpgCache(): void { + cacheGeneration++; resetMappingOverrides(); patchState(store, initialEpgState); }, }; } - ) + ), + withHooks((store) => { + const sources = inject(EpgSourceSettingsService); + let subscription: { unsubscribe(): void } | undefined; + return { + onInit: () => { + subscription = sources.changed$.subscribe(() => { + const context = store as typeof store & + StalkerEpgFeatureStoreContract; + const selectedId = context.selectedItvId(); + store.clearBulkItvEpgCache(); + void store.ensureBulkItvEpg(); + if (selectedId) + void store.applyMappedItvEpg([selectedId]); + }); + }, + onDestroy: () => subscription?.unsubscribe(), + }; + }) ); } diff --git a/libs/shared/interfaces/src/lib/electron-api.interface.ts b/libs/shared/interfaces/src/lib/electron-api.interface.ts index 54efe9331..4e8394a7c 100644 --- a/libs/shared/interfaces/src/lib/electron-api.interface.ts +++ b/libs/shared/interfaces/src/lib/electron-api.interface.ts @@ -85,6 +85,7 @@ export type ElectronBridgePlaylistType = (typeof ELECTRON_BRIDGE_PLAYLIST_TYPES)[keyof typeof ELECTRON_BRIDGE_PLAYLIST_TYPES]; export const ELECTRON_BRIDGE_EPG_PROGRESS_STATUSES = { + Cancelled: 'cancelled', Complete: 'complete', Error: 'error', Loading: 'loading', @@ -385,6 +386,7 @@ export interface ElectronBridgeEpgProgressStats { export interface ElectronBridgeEpgProgress { url: string; + generation?: number; status: ElectronBridgeEpgProgressStatus; stats?: ElectronBridgeEpgProgressStats; error?: string; diff --git a/libs/ui/epg/src/lib/epg-progress-panel/epg-progress-panel.component.ts b/libs/ui/epg/src/lib/epg-progress-panel/epg-progress-panel.component.ts index ce8cb2a8d..b83619288 100644 --- a/libs/ui/epg/src/lib/epg-progress-panel/epg-progress-panel.component.ts +++ b/libs/ui/epg/src/lib/epg-progress-panel/epg-progress-panel.component.ts @@ -96,6 +96,8 @@ export class EpgProgressPanelComponent { return 'check_circle'; case 'error': return 'error'; + case 'cancelled': + return 'cancel'; } }