diff --git a/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.spec.ts b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.spec.ts index efbc31936..b0e199ebe 100644 --- a/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.spec.ts @@ -33,7 +33,7 @@ describe('EpgQueueService', () => { beforeEach(() => { jest.useFakeTimers(); - xtreamApi = { getShortEpg: jest.fn() }; + xtreamApi = { getShortEpg: jest.fn().mockResolvedValue([]) }; fallback = { getProgramsForChannel: jest.fn().mockResolvedValue([]), getCurrentProgramsBatch: jest.fn().mockResolvedValue({}), @@ -165,8 +165,8 @@ describe('EpgQueueService', () => { it('only the latest enqueue commits queue state when prefetches overlap', async () => { // Defer batch resolutions so we can interleave them. - let resolveA: (v: Record) => void = () => {}; - let resolveB: (v: Record) => void = () => {}; + let resolveA!: (v: Record) => void; + let resolveB!: (v: Record) => void; fallback.getCurrentProgramsBatch .mockImplementationOnce( () => @@ -229,6 +229,44 @@ describe('EpgQueueService', () => { sub.unsubscribe(); }); + it('drops stale queued provider fetches while a newer XMLTV prefetch is pending', async () => { + let resolveLatestBatch!: (v: Record) => void; + fallback.getCurrentProgramsBatch.mockImplementationOnce( + () => + new Promise>((resolve) => { + resolveLatestBatch = resolve; + }) + ); + xtreamApi.getShortEpg.mockResolvedValue([]); + + await service.enqueue([1, 2], new Set([1, 2]), credentials); + expect(xtreamApi.getShortEpg).toHaveBeenCalledWith( + credentials, + 1, + 3, + { suppressErrorLog: true } + ); + + const latestEnqueue = service.enqueue( + [{ streamId: 3, epgChannelId: 'three.epg' }], + new Set([3]), + credentials + ); + + jest.advanceTimersByTime(201); + await Promise.resolve(); + + expect(xtreamApi.getShortEpg).not.toHaveBeenCalledWith( + credentials, + 2, + 3, + { suppressErrorLog: true } + ); + + resolveLatestBatch({}); + await latestEnqueue; + }); + it('clones the caller visibleIds Set so external mutation is harmless', async () => { const visible = new Set([1, 2, 3]); await service.enqueue( diff --git a/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts index 23c22ee7e..b26e3ec2a 100644 --- a/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts +++ b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts @@ -123,6 +123,12 @@ export class EpgQueueService implements OnDestroy { streamsByEpgId.set(id, list); } + // Make the latest viewport visible to any currently running queue + // before the async XMLTV prefetch returns, so stale queued provider + // requests are dropped immediately on fast scroll. + this.visibleSet = new Set(visibleIds); + this.queue = []; + const batchResult = await this.fetchXmltvCurrentPure( Array.from(streamsByEpgId.keys()) ); @@ -130,8 +136,6 @@ export class EpgQueueService implements OnDestroy { if (generation !== this.enqueueGeneration) return; // Atomic commit block — no awaits below. - // Clone the caller's Set so external mutation cannot shift our state. - this.visibleSet = new Set(visibleIds); this.pruneEphemeralMaps(this.visibleSet); for (const entry of normalized) { diff --git a/libs/portal/xtream/data-access/src/lib/stores/features/with-epg.feature.spec.ts b/libs/portal/xtream/data-access/src/lib/stores/features/with-epg.feature.spec.ts index d1fbf1e3e..c27223632 100644 --- a/libs/portal/xtream/data-access/src/lib/stores/features/with-epg.feature.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/stores/features/with-epg.feature.spec.ts @@ -63,24 +63,8 @@ function configureStore(setup: TestStoreSetup) { const fallbackService = { getProgramsForChannel: jest.fn, unknown[]>(), - resolveCurrentEpg: jest.fn( - async (args: { - epgChannelId: string | null | undefined; - preferUploaded: boolean; - fetchProvider: () => Promise; - }): Promise => { - const id = (args.epgChannelId ?? '').trim(); - if (args.preferUploaded && id) { - const xmltv = - await fallbackService.getProgramsForChannel(id); - if (xmltv.length > 0) return xmltv; - return args.fetchProvider(); - } - const provider = await args.fetchProvider(); - if (provider.length > 0 || !id) return provider; - return fallbackService.getProgramsForChannel(id); - } - ), + resolveCurrentEpg: + XtreamXmltvFallbackService.prototype.resolveCurrentEpg, }; const settingsStore = {