fix(epg): serialize cleanup with replacement imports

This commit is contained in:
4gray committed 2026-09-05 22:51:44 +02:00
1 parent e799fcec32
commit a7bd81a848
9 files changed
+238 -17

No files matched your search

+7 -5
View File
@@ -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 = () =>
@@ -39,6 +39,7 @@ export class EpgWorkerService {
private readonly fetchedUrls = new Set<string>();
private readonly workers = new Map<string, Worker>();
private readonly inFlightFetches = new Map<string, Promise<void>>();
private readonly inFlightSourceClears = new Map<string, Promise<void>>();
constructor(
private readonly loggerLabel = '[EPG Events]',
@@ -86,6 +87,26 @@ export class EpgWorkerService {
url: string,
options: ElectronBridgeTrustOptions = {}
): Promise<void> {
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<void> {
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<void> {
const runningWorker = this.workers.get(normalizedSourceUrl);
if (runningWorker) {
this.workers.delete(normalizedSourceUrl);
@@ -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();
+6 -5
View File
@@ -192,13 +192,14 @@ export class AppComponent implements OnInit {
private async fetchStaleEpgData(urls: string[]): Promise<void> {
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);
}
}
+6 -1
View File
@@ -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
@@ -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',
@@ -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
@@ -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']);
});
});
@@ -20,6 +20,7 @@ export class EpgSourceReconciliationError extends Error {
export class EpgSourceSettingsService {
private readonly injector = inject(Injector);
private activeUrls = new Set<string>();
private reconciliation: Promise<void> | undefined;
readonly revision = signal(0);
readonly changed$ = new Subject<void>();
@@ -38,12 +39,35 @@ export class EpgSourceSettingsService {
);
}
async waitForReconciliation(): Promise<void> {
while (this.reconciliation) {
await this.reconciliation.catch(() => undefined);
}
}
async synchronize(urls: string[] | string | undefined): Promise<void> {
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<void> {
const normalized = [
...new Set(
(Array.isArray(urls) ? urls : [urls ?? ''])