fix(xtream): prevent stale epg queue fetches

This commit is contained in:
4gray committed 2026-05-09 12:08:14 +02:00
1 parent 505e75b67e
commit 4020822b90
3 files changed
+49 -23

No files matched your search

@@ -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<string, EpgItem>) => void = () => {};
let resolveB: (v: Record<string, EpgItem>) => void = () => {};
let resolveA!: (v: Record<string, EpgItem>) => void;
let resolveB!: (v: Record<string, EpgItem>) => 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<string, EpgItem>) => void;
fallback.getCurrentProgramsBatch.mockImplementationOnce(
() =>
new Promise<Record<string, EpgItem>>((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(
@@ -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) {
@@ -63,24 +63,8 @@ function configureStore(setup: TestStoreSetup) {
const fallbackService = {
getProgramsForChannel: jest.fn<Promise<EpgItem[]>, unknown[]>(),
resolveCurrentEpg: jest.fn(
async (args: {
epgChannelId: string | null | undefined;
preferUploaded: boolean;
fetchProvider: () => Promise<EpgItem[]>;
}): Promise<EpgItem[]> => {
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 = {