From adb76a3fca92315b98437ceb9355d621efb06548 Mon Sep 17 00:00:00 2001 From: 4gray <4gray@users.noreply.github.com> Date: Wed, 10 Jun 2026 12:06:10 +0200 Subject: [PATCH] fix(epg): dedupe concurrent fetches and await worker termination (#1040) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(epg): dedupe concurrent fetches and await worker termination Two concurrent fetchEpgFromUrl calls for the same URL spawned two workers parsing and writing the same EPG data, with the second one overwriting the first one's entry in the workers map and leaking that worker. Share the in-flight promise instead of spawning a competitor. worker.terminate() was also fired without awaiting it in every settle path. A terminated-but-still-running worker can keep holding the SQLite lock, blocking the next EPG operation. All settle paths now resolve or reject only after the worker thread has really exited; the settle guard runs first so the worker's own exit event cannot hijack the outcome. Co-Authored-By: Claude Fable 5 * fix(epg): close review gaps in fetch dedupe and clear sequencing - Check the in-flight map before the fetched-URL shortcut: a completed fetch is added to fetchedUrls while its worker is still terminating, and a concurrent request must keep awaiting that window instead of resolving early. - clearEpgData now resolves only after every interrupted fetch worker has terminated too, not just the clear worker — they may still hold the SQLite lock the caller expects to be free. Addresses Codex/Greptile review feedback on #1040. Co-Authored-By: Claude Fable 5 --------- Co-authored-by: 4gray Co-authored-by: Claude Fable 5 --- .../src/app/events/epg-worker.service.ts | 136 +++++++++++++--- .../src/app/events/epg.events.spec.ts | 146 ++++++++++++++++++ 2 files changed, 258 insertions(+), 24 deletions(-) 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 cb44cd1ce..f725b3a61 100644 --- a/apps/electron-backend/src/app/events/epg-worker.service.ts +++ b/apps/electron-backend/src/app/events/epg-worker.service.ts @@ -21,6 +21,7 @@ interface EpgWorkerMessage { export class EpgWorkerService { private readonly fetchedUrls = new Set(); private readonly workers = new Map(); + private readonly inFlightFetches = new Map>(); constructor( private readonly loggerLabel = '[EPG Events]', @@ -59,6 +60,22 @@ export class EpgWorkerService { } async fetchEpgFromUrl(url: string): Promise { + // 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 + // `workers`, leaking that worker. Share the in-flight promise instead. + // Checked before the fetched-URL shortcut: a completed fetch is added + // to `fetchedUrls` while its worker is still terminating, and callers + // must keep awaiting that termination window. + const inFlight = this.inFlightFetches.get(url); + if (inFlight) { + console.log( + this.loggerLabel, + `Reusing in-flight EPG fetch: ${url}` + ); + return inFlight; + } + if (this.fetchedUrls.has(url)) { console.log( this.loggerLabel, @@ -67,6 +84,14 @@ export class EpgWorkerService { return; } + const fetchPromise = this.startFetch(url).finally(() => { + this.inFlightFetches.delete(url); + }); + this.inFlightFetches.set(url, fetchPromise); + return fetchPromise; + } + + private startFetch(url: string): Promise { return new Promise((resolve, reject) => { let worker: Worker; try { @@ -104,9 +129,15 @@ export class EpgWorkerService { undefined, errorMessage ); - worker.terminate(); this.workers.delete(url); - settle(() => reject(new Error(errorMessage))); + // Settle only after the worker thread is really gone: a + // terminated-but-still-running worker can keep holding the + // SQLite lock and block the next EPG fetch. + settle(() => { + void this.terminateWorker(worker, 'timed out fetch').then( + () => reject(new Error(errorMessage)) + ); + }); }, this.fetchTimeoutMs); worker.on('message', async (message: EpgWorkerMessage) => { @@ -142,9 +173,13 @@ export class EpgWorkerService { message.stats ); this.fetchedUrls.add(url); - worker.terminate(); this.workers.delete(url); - settle(() => resolve()); + settle(() => { + void this.terminateWorker( + worker, + 'completed fetch' + ).then(() => resolve()); + }); break; case 'EPG_ERROR': @@ -159,13 +194,19 @@ export class EpgWorkerService { undefined, message.error ); - worker.terminate(); this.workers.delete(url); - settle(() => - reject( - new Error(message.error || 'Unknown error') - ) - ); + settle(() => { + void this.terminateWorker( + worker, + 'failed fetch' + ).then(() => + reject( + new Error( + message.error || 'Unknown error' + ) + ) + ); + }); break; } } catch (err) { @@ -180,9 +221,13 @@ export class EpgWorkerService { undefined, err instanceof Error ? err.message : String(err) ); - worker.terminate(); this.workers.delete(url); - settle(() => reject(err)); + settle(() => { + void this.terminateWorker( + worker, + 'failed message handling' + ).then(() => reject(err)); + }); } }); @@ -194,9 +239,12 @@ export class EpgWorkerService { undefined, error.message ); - worker.terminate(); this.workers.delete(url); - settle(() => reject(error)); + settle(() => { + void this.terminateWorker(worker, 'errored fetch').then( + () => reject(error) + ); + }); }); worker.on('exit', (code) => { @@ -244,8 +292,9 @@ export class EpgWorkerService { }s`; console.error(this.loggerLabel, errorMessage); settle(() => { - worker.terminate(); - reject(new Error(errorMessage)); + void this.terminateWorker(worker, 'timed out clear').then( + () => reject(new Error(errorMessage)) + ); }); }, this.fetchTimeoutMs); @@ -261,12 +310,24 @@ export class EpgWorkerService { 'EPG data cleared via worker' ); this.fetchedUrls.clear(); - this.workers.forEach((runningWorker) => - runningWorker.terminate() + // Resolve only after every interrupted fetch + // worker has exited too — they may still hold the + // SQLite lock the caller expects to be free. + const terminations = [ + ...this.workers.values(), + ].map((runningWorker) => + this.terminateWorker( + runningWorker, + 'fetch during clear' + ) ); this.workers.clear(); - worker.terminate(); - resolve(); + terminations.push( + this.terminateWorker(worker, 'completed clear') + ); + void Promise.all(terminations).then(() => + resolve() + ); }); } else if (message.type === 'EPG_ERROR') { console.error( @@ -275,8 +336,14 @@ export class EpgWorkerService { message.error ); settle(() => { - worker.terminate(); - reject(new Error(message.error || 'Clear failed')); + void this.terminateWorker( + worker, + 'failed clear' + ).then(() => + reject( + new Error(message.error || 'Clear failed') + ) + ); }); } } @@ -289,8 +356,9 @@ export class EpgWorkerService { error ); settle(() => { - worker.terminate(); - reject(error); + void this.terminateWorker(worker, 'errored clear').then( + () => reject(error) + ); }); }); @@ -303,6 +371,26 @@ export class EpgWorkerService { }); } + /** + * Awaits worker shutdown so callers can sequence work (e.g. the next DB + * access) after the thread has really exited. Termination failures are + * logged and swallowed — there is nothing actionable left to do. + */ + private async terminateWorker( + worker: Worker, + context: string + ): Promise { + try { + await worker.terminate(); + } catch (error) { + console.error( + this.loggerLabel, + `Failed to terminate ${context} worker:`, + error + ); + } + } + private createEpgWorker(): Worker { const bootstrap = resolveWorkerRuntimeBootstrap({ isPackaged: app.isPackaged, 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 8d61ab6ac..349b51d46 100644 --- a/apps/electron-backend/src/app/events/epg.events.spec.ts +++ b/apps/electron-backend/src/app/events/epg.events.spec.ts @@ -160,6 +160,152 @@ describe('EpgEvents', () => { expect(worker.terminate).toHaveBeenCalled(); }); + it('reuses the in-flight worker when the same EPG URL is fetched concurrently', async () => { + const workerService = new EpgWorkerService('[Test EPG]', 1000); + const url = 'https://example.com/guide.xml'; + + const firstPromise = workerService.fetchEpgFromUrl(url); + const secondPromise = workerService.fetchEpgFromUrl(url); + + expect(mockWorkerInstances).toHaveLength(1); + + const worker = mockWorkerInstances[0]; + worker.emit('message', { type: 'READY' }); + await flushPromises(); + worker.emit('message', { + type: 'EPG_COMPLETE', + stats: { totalChannels: 1, totalPrograms: 2 }, + }); + + await expect(firstPromise).resolves.toBeUndefined(); + await expect(secondPromise).resolves.toBeUndefined(); + expect(mockWorkerInstances).toHaveLength(1); + }); + + it('starts a fresh worker once a failed fetch for the same URL has settled', async () => { + const workerService = new EpgWorkerService('[Test EPG]', 1000); + const url = 'https://example.com/guide.xml'; + + const firstPromise = workerService.fetchEpgFromUrl(url); + mockWorkerInstances[0].emit('message', { + type: 'EPG_ERROR', + error: 'parse failed', + }); + await expect(firstPromise).rejects.toThrow('parse failed'); + + const secondPromise = workerService.fetchEpgFromUrl(url); + expect(mockWorkerInstances).toHaveLength(2); + + const retryWorker = mockWorkerInstances[1]; + retryWorker.emit('message', { type: 'READY' }); + await flushPromises(); + retryWorker.emit('message', { + type: 'EPG_COMPLETE', + stats: { totalChannels: 1, totalPrograms: 2 }, + }); + + await expect(secondPromise).resolves.toBeUndefined(); + }); + + it('shares the in-flight promise while a completed fetch is still terminating', async () => { + const workerService = new EpgWorkerService('[Test EPG]', 1000); + const url = 'https://example.com/guide.xml'; + + const firstPromise = workerService.fetchEpgFromUrl(url); + const worker = mockWorkerInstances[0]; + + let releaseTerminate!: () => void; + worker.terminate.mockReturnValue( + new Promise((resolve) => { + releaseTerminate = resolve; + }) + ); + + worker.emit('message', { type: 'READY' }); + await flushPromises(); + worker.emit('message', { + type: 'EPG_COMPLETE', + stats: { totalChannels: 1, totalPrograms: 2 }, + }); + await flushPromises(); + + // The URL is already marked as fetched, but the worker is still + // terminating: a new request must keep awaiting that window instead + // of resolving early via the fetched-URL shortcut. + const secondPromise = workerService.fetchEpgFromUrl(url); + expect(mockWorkerInstances).toHaveLength(1); + + let secondResolved = false; + void secondPromise.then(() => { + secondResolved = true; + }); + await flushPromises(); + expect(secondResolved).toBe(false); + + releaseTerminate(); + await expect(firstPromise).resolves.toBeUndefined(); + await flushPromises(); + expect(secondResolved).toBe(true); + }); + + it('does not resolve clearEpgData until interrupted fetch workers have terminated', async () => { + const workerService = new EpgWorkerService('[Test EPG]', 1000); + + void workerService + .fetchEpgFromUrl('https://example.com/guide.xml') + .catch(() => undefined); + const fetchWorker = mockWorkerInstances[0]; + + let releaseFetchTerminate!: () => void; + fetchWorker.terminate.mockReturnValue( + new Promise((resolve) => { + releaseFetchTerminate = resolve; + }) + ); + + const clearPromise = workerService.clearEpgData(); + const clearWorker = mockWorkerInstances[1]; + clearWorker.emit('message', { type: 'READY' }); + clearWorker.emit('message', { type: 'CLEAR_COMPLETE' }); + + let cleared = false; + void clearPromise.then(() => { + cleared = true; + }); + await flushPromises(); + expect(fetchWorker.terminate).toHaveBeenCalled(); + expect(cleared).toBe(false); + + releaseFetchTerminate(); + await flushPromises(); + expect(cleared).toBe(true); + }); + + it('rejects a timed-out fetch with the timeout error after the worker has terminated', async () => { + jest.useFakeTimers(); + + const workerService = new EpgWorkerService('[Test EPG]', 25); + const fetchPromise = workerService.fetchEpgFromUrl( + 'https://example.com/guide.xml' + ); + const worker = mockWorkerInstances[0]; + + let terminated = false; + worker.terminate.mockImplementation(() => { + terminated = true; + // Worker threads emit 'exit' as part of termination; the timeout + // rejection must win over the generic exit-handler rejection. + worker.emit('exit', 1); + return Promise.resolve(0); + }); + + worker.emit('message', { type: 'READY' }); + jest.advanceTimersByTime(25); + + await expect(fetchPromise).rejects.toThrow('EPG fetch timed out after'); + expect(terminated).toBe(true); + }); + it('falls back to case-insensitive channel id lookup for EPG programs', async () => { const select = jest.fn(); const programLimitExact = jest.fn().mockResolvedValue([]);