diff --git a/apps/electron-backend-e2e/src/epg.e2e.ts b/apps/electron-backend-e2e/src/epg.e2e.ts index b3ce5bf6c..1a0a0d612 100644 --- a/apps/electron-backend-e2e/src/epg.e2e.ts +++ b/apps/electron-backend-e2e/src/epg.e2e.ts @@ -153,14 +153,16 @@ test.describe('Electron EPG', () => { try { await openSettings(app.mainWindow); await openSettingsSection(app.mainWindow, 'epg'); - for (const url of [first.resourceUrl, second.resourceUrl]) { + for (const [index, url] of [ + first.resourceUrl, + second.resourceUrl, + ].entries()) { await app.mainWindow .getByRole('button', { name: 'Add EPG source' }) .click(); - await app.mainWindow - .locator('.epg-source-row input') - .last() - .fill(url); + const inputs = app.mainWindow.locator('.epg-source-row input'); + await expect(inputs).toHaveCount(index + 1); + await inputs.nth(index).fill(url); } await saveSettings(app.mainWindow); const programs = () => 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 f5e7e143d..49f06c494 100644 --- a/apps/electron-backend/src/app/events/epg-worker.service.ts +++ b/apps/electron-backend/src/app/events/epg-worker.service.ts @@ -39,6 +39,7 @@ export class EpgWorkerService { private readonly fetchedUrls = new Set(); private readonly workers = new Map(); private readonly inFlightFetches = new Map>(); + private readonly inFlightSourceClears = new Map>(); constructor( private readonly loggerLabel = '[EPG Events]', @@ -86,6 +87,26 @@ export class EpgWorkerService { url: string, options: ElectronBridgeTrustOptions = {} ): Promise { + url = url.trim(); + const generation = requestEpgSource(url); + const clear = this.inFlightSourceClears.get(url); + if (clear) { + await clear.catch(() => undefined); + if (generation !== epgSourceGeneration(url)) { + this.sendProgressToRenderer( + url, + 'cancelled', + undefined, + undefined, + undefined, + undefined, + undefined, + generation + ); + return; + } + return this.fetchEpgFromUrl(url, options); + } // A second request for an URL that is already being fetched must not // spawn a competing worker: both would parse and write the same EPG // data, and the late one would overwrite the early one's entry in @@ -121,7 +142,7 @@ export class EpgWorkerService { url: string, options: ElectronBridgeTrustOptions ): Promise { - const generation = requestEpgSource(url); + const generation = epgSourceGeneration(url); return new Promise((resolve, reject) => { let worker: Worker; try { @@ -383,6 +404,21 @@ export class EpgWorkerService { retireEpgSource(normalizedSourceUrl); this.fetchedUrls.delete(normalizedSourceUrl); + const previous = this.inFlightSourceClears.get(normalizedSourceUrl); + const operation = previous + ? previous + .catch(() => undefined) + .then(() => this.startSourceClear(normalizedSourceUrl)) + : this.startSourceClear(normalizedSourceUrl); + const clear = operation.finally(() => { + if (this.inFlightSourceClears.get(normalizedSourceUrl) === clear) + this.inFlightSourceClears.delete(normalizedSourceUrl); + }); + this.inFlightSourceClears.set(normalizedSourceUrl, clear); + return clear; + } + + private async startSourceClear(normalizedSourceUrl: string): Promise { const runningWorker = this.workers.get(normalizedSourceUrl); if (runningWorker) { this.workers.delete(normalizedSourceUrl); 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 a661a9d23..80e85ab10 100644 --- a/apps/electron-backend/src/app/events/epg.events.spec.ts +++ b/apps/electron-backend/src/app/events/epg.events.spec.ts @@ -402,6 +402,66 @@ describe('EpgEvents', () => { expect(terminated).toBe(true); }); + it('waits for source cleanup before starting a replacement import of the same normalized URL', async () => { + const service = new EpgWorkerService('[Test EPG]', 1000); + const url = 'https://replacement.example/guide.xml'; + const clear = service.clearEpgDataForSource(url); + const clearWorker = mockWorkerInstances[0]; + const replacement = service.fetchEpgFromUrl(` ${url} `); + const workersBeforeCleanup = mockWorkerInstances.length; + clearWorker.emit('message', { type: 'READY' }); + clearWorker.emit('message', { type: 'CLEAR_COMPLETE' }); + await clear; + await flushPromises(); + const fetchWorker = mockWorkerInstances[1]; + fetchWorker.emit('message', { type: 'READY' }); + fetchWorker.emit('message', { type: 'EPG_COMPLETE' }); + await replacement; + expect(workersBeforeCleanup).toBe(1); + expect(service.hasFetchedUrl(url)).toBe(true); + }); + + it('serializes repeated clears and retires a replacement waiting for the earlier clear', async () => { + const service = new EpgWorkerService('[Test EPG]', 1000); + const url = 'https://twice-removed.example/guide.xml'; + const firstClear = service.clearEpgDataForSource(url); + const waitingFetch = service.fetchEpgFromUrl(url); + const secondClear = service.clearEpgDataForSource(url); + expect(mockWorkerInstances).toHaveLength(1); + mockWorkerInstances[0].emit('message', { type: 'CLEAR_COMPLETE' }); + await firstClear; + await waitingFetch; + await flushPromises(); + expect(mockWorkerInstances).toHaveLength(2); + mockWorkerInstances[1].emit('message', { type: 'READY' }); + expect(mockWorkerInstances[1].postMessage).toHaveBeenCalledWith({ + type: 'CLEAR_EPG_SOURCE', + sourceUrl: url, + }); + mockWorkerInstances[1].emit('message', { type: 'CLEAR_COMPLETE' }); + await secondClear; + expect(service.hasFetchedUrl(url)).toBe(false); + }); + + it('allows a replacement import after a failed source cleanup has terminated', async () => { + const service = new EpgWorkerService('[Test EPG]', 1000); + const url = 'https://retry-clear.example/guide.xml'; + const clear = service.clearEpgDataForSource(url); + const outcome = clear.catch((error: Error) => error.message); + const replacement = service.fetchEpgFromUrl(url); + mockWorkerInstances[0].emit('message', { + type: 'EPG_ERROR', + error: 'clear failed', + }); + expect(await outcome).toBe('clear failed'); + await flushPromises(); + expect(mockWorkerInstances[0].terminate).toHaveBeenCalled(); + mockWorkerInstances[1].emit('message', { type: 'READY' }); + mockWorkerInstances[1].emit('message', { type: 'EPG_COMPLETE' }); + await replacement; + expect(service.hasFetchedUrl(url)).toBe(true); + }); + it('keeps an active EPG fetch alive when worker progress keeps moving', async () => { jest.useFakeTimers(); diff --git a/apps/web/src/app/app.component.ts b/apps/web/src/app/app.component.ts index 37b23f787..a80a56979 100644 --- a/apps/web/src/app/app.component.ts +++ b/apps/web/src/app/app.component.ts @@ -192,13 +192,14 @@ export class AppComponent implements OnInit { private async fetchStaleEpgData(urls: string[]): Promise { await this.settingsStore.loadSettings(); const revision = this.epgSources.revision(); - const fetchCurrentSources = (sources: string[]) => { + const fetchCurrentSources = async (sources: string[]) => { + await this.epgSources.waitForReconciliation(); this.epgService.fetchEpg( this.epgSources.retainCurrentSources(sources, revision) ); }; if (!this.epgBridge.supportsSourceFreshness) { - fetchCurrentSources(urls); + await fetchCurrentSources(urls); return; } @@ -206,7 +207,7 @@ export class AppComponent implements OnInit { const result = await this.epgBridge.checkFreshness(urls, 12); if (!result) { - fetchCurrentSources(urls); + await fetchCurrentSources(urls); return; } @@ -228,12 +229,12 @@ export class AppComponent implements OnInit { debugAppComponent( `EPG: Fetching ${result.staleUrls.length} stale source(s)` ); - fetchCurrentSources(result.staleUrls); + await fetchCurrentSources(result.staleUrls); } } catch (error) { console.error('Error checking EPG freshness, fetching all:', error); // Fallback: fetch all URLs if freshness check fails - fetchCurrentSources(urls); + await fetchCurrentSources(urls); } } diff --git a/docs/architecture/m3u-playlist-module.md b/docs/architecture/m3u-playlist-module.md index a0ab67b29..d3115679b 100644 --- a/docs/architecture/m3u-playlist-module.md +++ b/docs/architecture/m3u-playlist-module.md @@ -1333,6 +1333,8 @@ 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. Successfully cleared request candidates are forgotten without resetting their generation fences; failed cleanups remain retryable. +Same-URL clears are serialized and replacement imports await the outstanding +clear, so an older cleanup cannot erase a newly re-added source. 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 @@ -1341,7 +1343,10 @@ unknown (`NULL`) ownership are conservatively left alone; the existing database initialization backfill handles rows whose channel still identifies their owner. There is no new schema migration. -Renderer reconciliation increments a data revision and cancels earlier lookup +Renderer reconciliation fences lookups before its first asynchronous step. +Imports wait for serialized reconciliation (including playlist migration), then +filter against its committed owner set. Completion increments the data revision again +and cancels earlier lookup subscriptions, clears program caches and the selected M3U guide, and refreshes Xtream selection and visible channel previews, plus Stalker manual mapping overrides and bulk guides. A delayed startup import is diff --git a/libs/epg/data-access/src/lib/epg.service.spec.ts b/libs/epg/data-access/src/lib/epg.service.spec.ts index 52364ae58..69d4b29e0 100644 --- a/libs/epg/data-access/src/lib/epg.service.spec.ts +++ b/libs/epg/data-access/src/lib/epg.service.spec.ts @@ -1,8 +1,12 @@ import { TestBed } from '@angular/core/testing'; import { MatSnackBar } from '@angular/material/snack-bar'; import { TranslateService } from '@ngx-translate/core'; -import { firstValueFrom, skip } from 'rxjs'; -import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services'; +import { firstValueFrom, of, skip } from 'rxjs'; +import { + EpgSourceSettingsService, + PlaylistsService, + SettingsStore, +} from '@iptvnator/services'; import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service'; import { EpgService } from './epg.service'; @@ -46,6 +50,19 @@ describe('EpgService', () => { TestBed.configureTestingModule({ providers: [ EpgService, + { + provide: PlaylistsService, + useValue: { + getAllPlaylists: () => + of([ + { + epgUrls: [ + 'https://playlist.example/guide.xml', + ], + }, + ]), + }, + }, { provide: EpgRuntimeBridgeService, useValue: epgBridge, @@ -138,6 +155,41 @@ describe('EpgService', () => { expect(epgBridge.getChannelPrograms).toHaveBeenCalledTimes(2); }); + it('waits for ongoing reconciliation and filters imports against committed owners', async () => { + epgBridge.supportsImport = true; + const original = window.electron; + let complete!: (result: { success: boolean }) => void; + window.electron = { + reconcileEpgSources: () => + new Promise((resolve) => { + complete = resolve; + }), + } as typeof window.electron; + try { + const sources = TestBed.inject(EpgSourceSettingsService); + const reconciliation = sources.synchronize([ + 'https://kept.example/guide.xml', + ]); + await Promise.resolve(); + const pending = service.fetchEpg([ + 'https://removed.example/guide.xml', + 'https://playlist.example/guide.xml', + ]); + await Promise.resolve(); + await Promise.resolve(); + expect(epgBridge.fetchEpg).not.toHaveBeenCalled(); + complete({ success: true }); + await reconciliation; + await pending; + expect(epgBridge.fetchEpg).toHaveBeenCalledWith( + ['https://playlist.example/guide.xml'], + expect.anything() + ); + } finally { + window.electron = original; + } + }); + it('does not fetch EPG when bridge import support is disabled', () => { service.fetchEpg(['https://example.com/epg.xml']); @@ -147,7 +199,7 @@ describe('EpgService', () => { it('fetches EPG through the EPG runtime bridge when import support is enabled', async () => { epgBridge.supportsImport = true; - service.fetchEpg([ + await service.fetchEpg([ 'https://example.com/epg.xml', '', 'https://example.com/other.xml', diff --git a/libs/epg/data-access/src/lib/epg.service.ts b/libs/epg/data-access/src/lib/epg.service.ts index 4cf3c6a06..d76ae9bf6 100644 --- a/libs/epg/data-access/src/lib/epg.service.ts +++ b/libs/epg/data-access/src/lib/epg.service.ts @@ -94,6 +94,7 @@ export class EpgService { // Filter out empty and duplicate URLs and send all URLs at once. const revision = this.sourceSettings.revision(); await this.settingsStore.loadSettings(); + await this.sourceSettings.waitForReconciliation(); const validUrls = this.sourceSettings.retainCurrentSources( normalizeEpgUrls(urls), revision diff --git a/libs/services/src/lib/epg-source-settings.service.spec.ts b/libs/services/src/lib/epg-source-settings.service.spec.ts index 5027be5df..0bd81c73c 100644 --- a/libs/services/src/lib/epg-source-settings.service.spec.ts +++ b/libs/services/src/lib/epg-source-settings.service.spec.ts @@ -33,6 +33,8 @@ describe('EPG source settings synchronization', () => { lookup.subscribe(observer); const synchronization = service.synchronize([' a ', '', 'a']); expect(reconcileEpgSources).not.toHaveBeenCalled(); + pending.next('response during playlist migration'); + expect(observer).not.toHaveBeenCalled(); playlists.next([]); await synchronization; pending.next('old programme'); @@ -40,7 +42,7 @@ describe('EPG source settings synchronization', () => { lookup.subscribe(observer); pending.next('late programme'); expect(observer).not.toHaveBeenCalled(); - expect(service.revision()).toBe(1); + expect(service.revision()).toBe(2); expect(reconcileEpgSources).toHaveBeenCalledWith(['a']); expect(await firstValueFrom(of('new').pipe(service.guard()))).toBe( 'new' @@ -64,9 +66,47 @@ describe('EPG source settings synchronization', () => { await expect(service.synchronize(['current'])).rejects.toThrow( 'Failed to reconcile EPG sources' ); - expect(service.revision()).toBe(1); + expect(service.revision()).toBe(2); expect(service.retainCurrentSources(['current', 'removed'], 0)).toEqual( ['current'] ); }); + + it('serializes overlapping saves and waits for the latest committed source set', async () => { + let finishFirst!: (result: { success: boolean }) => void; + const reconcileEpgSources = jest + .fn() + .mockImplementationOnce( + () => + new Promise((resolve) => { + finishFirst = resolve; + }) + ) + .mockResolvedValue({ success: true }); + window.electron = { + reconcileEpgSources, + } as unknown as typeof window.electron; + const injector = Injector.create({ + providers: [ + EpgSourceSettingsService, + { + provide: PlaylistsService, + useValue: { getAllPlaylists: () => of([]) }, + }, + ], + }); + const service = injector.get(EpgSourceSettingsService); + const first = service.synchronize(['a']); + const second = service.synchronize(['b']); + const waiter = service.waitForReconciliation(); + await Promise.resolve(); + expect(reconcileEpgSources).toHaveBeenCalledTimes(1); + finishFirst({ success: true }); + await Promise.all([first, second, waiter]); + expect(reconcileEpgSources.mock.calls.map(([urls]) => urls)).toEqual([ + ['a'], + ['b'], + ]); + expect(service.retainCurrentSources(['a', 'b'], 0)).toEqual(['b']); + }); }); diff --git a/libs/services/src/lib/epg-source-settings.service.ts b/libs/services/src/lib/epg-source-settings.service.ts index 737882142..69a43479b 100644 --- a/libs/services/src/lib/epg-source-settings.service.ts +++ b/libs/services/src/lib/epg-source-settings.service.ts @@ -20,6 +20,7 @@ export class EpgSourceReconciliationError extends Error { export class EpgSourceSettingsService { private readonly injector = inject(Injector); private activeUrls = new Set(); + private reconciliation: Promise | undefined; readonly revision = signal(0); readonly changed$ = new Subject(); @@ -38,12 +39,35 @@ export class EpgSourceSettingsService { ); } + async waitForReconciliation(): Promise { + while (this.reconciliation) { + await this.reconciliation.catch(() => undefined); + } + } + async synchronize(urls: string[] | string | undefined): Promise { if ( typeof window === 'undefined' || !window.electron?.reconcileEpgSources ) return; + // Fence existing lookups before playlist migration or IPC can yield. + this.revision.update((revision) => revision + 1); + const previous = this.reconciliation; + const operation = previous + ? previous.catch(() => undefined).then(() => this.reconcile(urls)) + : this.reconcile(urls); + const pending = operation.finally(() => { + if (this.reconciliation === pending) + this.reconciliation = undefined; + }); + this.reconciliation = pending; + return pending; + } + + private async reconcile( + urls: string[] | string | undefined + ): Promise { const normalized = [ ...new Set( (Array.isArray(urls) ? urls : [urls ?? ''])