fix(epg): close source reconciliation review races

This commit is contained in:
4gray committed 2026-09-05 22:37:11 +02:00
1 parent ffcbee8140
commit e799fcec32
17 files changed
+396 -41

No files matched your search

+1 -1
View File
@@ -182,7 +182,7 @@ settings load and playlist migration; failed settings reads and incomplete
playlist migration never authorize pruning. Removed sources retire queued and
running imports before worker-owned deletion. Shared channel IDs survive while
another source has programmes; manual mappings remain user preferences, but no
longer resolve deleted data. Renderer lookup generations and Xtream preview
longer resolve deleted data. Renderer lookup generations, Xtream previews and Stalker mapping-cache
invalidation prevent late results from restoring removed programmes. Provider
EPG is independent. See `docs/architecture/m3u-playlist-module.md`
("XMLTV source lifecycle").
+1 -1
View File
@@ -1667,7 +1667,7 @@ settings load and playlist migration; failed settings reads and incomplete
playlist migration never authorize pruning. Removed sources retire queued and
running imports before worker-owned deletion. Shared channel IDs survive while
another source has programmes; manual mappings remain user preferences, but no
longer resolve deleted data. Renderer lookup generations and Xtream preview
longer resolve deleted data. Renderer lookup generations, Xtream previews and Stalker mapping-cache
invalidation prevent late results from restoring removed programmes. Provider
EPG is independent. See `docs/architecture/m3u-playlist-module.md`
("XMLTV source lifecycle").
@@ -116,6 +116,102 @@ test('@epg @xtream @electron removes uploaded guide data and restores provider E
}
});
test('@epg @stalker @electron invalidates a loaded manual XMLTV mapping after source removal', async ({
dataDir,
request,
}) => {
test.setTimeout(120000);
await resetMockServers(request, ['stalker']);
const fixture = await fetchStalkerCategoryFixture(request, 'itv');
const item = fixture.items[0];
const stamp = (date: Date) =>
date.toISOString().replace(/[-:T]/g, '').slice(0, 14);
const source = await createMutableTextServer(
`<tv><channel id="stalker-mapped"><display-name>Mapped Guide</display-name></channel>
<programme channel="stalker-mapped" start="${stamp(new Date(Date.now() - 600000))} +0000" stop="${stamp(new Date(Date.now() + 3600000))} +0000"><title>Retired Stalker Bulletin</title></programme></tv>`,
{
contentType: 'application/xml',
resourcePath: '/stalker.xml',
}
);
const app = await launchElectronApp(dataDir);
try {
await app.mainWindow.route('https://test-streams.mux.dev/**', () => {
// Keep external media pending while the local guide is exercised.
});
await addStalkerPortal(app.mainWindow, {
name: 'Stalker XMLTV Removal',
});
await waitForStalkerCatalog(app.mainWindow);
const playlistId = new URL(app.mainWindow.url()).pathname.match(
/\/stalker\/([^/]+)/
)?.[1];
expect(playlistId).toBeTruthy();
const key = `stalker:${decodeURIComponent(playlistId!)}:${String(item.id).trim()}`;
await app.mainWindow.evaluate(
(mappingKey) =>
window.electron.setEpgMapping(mappingKey, 'stalker-mapped'),
key
);
await openSettings(app.mainWindow);
await openSettingsSection(app.mainWindow, 'epg');
await app.mainWindow
.getByRole('button', { name: 'Add EPG source' })
.click();
await app.mainWindow
.locator('.epg-source-row input')
.fill(source.resourceUrl);
await saveSettings(app.mainWindow);
await expect
.poll(
() =>
app.mainWindow.evaluate(async () =>
(
await window.electron.getChannelPrograms(
'stalker-mapped'
)
).map((p) => p.title)
),
{ timeout: 30000 }
)
.toContain('Retired Stalker Bulletin');
await openWorkspaceSection(app.mainWindow, 'Live TV');
await clickCategoryByNameExact(app.mainWindow, fixture.categoryName);
const row = channelItemByTitle(
app.mainWindow,
item.o_name || item.name || ''
).first();
await expect(row.locator('.epg-title')).toHaveText(
'Retired Stalker Bulletin'
);
await row.click();
await expect
.poll(() => timelineBlockTitles(app.mainWindow))
.toContain('Retired Stalker Bulletin');
await openSettings(app.mainWindow);
await openSettingsSection(app.mainWindow, 'epg');
await app.mainWindow.locator('.epg-source-row button').nth(1).click();
await saveSettings(app.mainWindow);
await openWorkspaceSection(app.mainWindow, 'Live TV');
await clickCategoryByNameExact(app.mainWindow, fixture.categoryName);
await expect(row).toBeVisible();
await row.click();
await expect
.poll(() => timelineBlockTitles(app.mainWindow))
.not.toContain('Retired Stalker Bulletin');
await expect(row.locator('.epg-title')).toHaveCount(0);
expect(
await app.mainWindow.evaluate(
(mappingKey) => window.electron.getEpgMapping(mappingKey),
key
)
).toMatchObject({ epgChannelId: 'stalker-mapped' });
} finally {
await closeElectronApp(app);
await source.close();
}
});
for (const timeZone of ['UTC', 'Europe/Berlin'] as const) {
test(`@epg @xtream @electron renders Xtream EPG previews and the timeline schedule in ${timeZone}`, async ({
dataDir,
@@ -1,4 +1,4 @@
import { epgSourceGeneration } from './epg-source-generation';
import { epgSourceGeneration, requestEpgSource } from './epg-source-generation';
import { eq } from 'drizzle-orm';
import { ElectronBridgeTrustOptions } from '@iptvnator/shared/interfaces';
import { getDatabase } from '../database/connection';
@@ -102,7 +102,7 @@ export async function handleFetchEpg(
.filter((url) => url?.trim())
.map((url) => url.trim());
const generations = new Map(
validUrls.map((url) => [url, epgSourceGeneration(url)])
validUrls.map((url) => [url, requestEpgSource(url)])
);
if (validUrls.length === 0) {
@@ -152,7 +152,19 @@ export async function handleFetchEpg(
const errors: string[] = [];
for (const url of urlsToFetch) {
try {
if (generations.get(url) !== epgSourceGeneration(url)) continue;
if (generations.get(url) !== epgSourceGeneration(url)) {
epgWorkerService.sendProgressToRenderer(
url,
'cancelled',
undefined,
undefined,
undefined,
undefined,
undefined,
generations.get(url)
);
continue;
}
await epgWorkerService.fetchEpgFromUrl(url, options);
} catch (error) {
console.error(
@@ -1,13 +1,23 @@
/** Retires queued imports as well as workers already parsing a removed URL. */
const generations = new Map<string, number>();
const requests = new Map<string, symbol>();
export function epgSourceGeneration(url: string): number {
const key = url.trim();
if (!generations.has(key)) generations.set(key, 0);
return generations.get(key)!;
return generations.get(url.trim()) ?? 0;
}
export function requestEpgSource(url: string): number {
requests.set(url.trim(), Symbol());
return epgSourceGeneration(url);
}
export function retireEpgSource(url: string): void {
generations.set(url.trim(), epgSourceGeneration(url) + 1);
}
export function requestedEpgSources(): string[] {
return [...generations.keys()];
export function requestedEpgSources(): Map<string, symbol> {
return new Map(requests);
}
export function forgetEpgSourceRequest(
url: string,
request: symbol | undefined
): void {
// A request received during cleanup still needs reconciliation next time.
if (requests.get(url.trim()) === request) requests.delete(url.trim());
}
@@ -6,7 +6,7 @@ import {
} from '../database/schema';
import { getDatabase } from '../database/connection';
import { epgWorkerService } from './epg-worker.service';
import { epgSourceGeneration } from './epg-source-generation';
import { epgSourceGeneration, requestEpgSource } from './epg-source-generation';
import { reconcileEpgSources } from './epg-source-settings.service';
jest.mock('../database/connection', () => ({ getDatabase: jest.fn() }));
@@ -65,6 +65,45 @@ describe('committed EPG source reconciliation', () => {
}
});
it('does not clear historical request keys again after successful cleanup', async () => {
requestEpgSource('removed');
await reconcileEpgSources([]);
rows.set(epgChannels, []);
rows.set(epgPrograms, []);
jest.clearAllMocks();
await reconcileEpgSources([]);
expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalled();
});
it('retries a failed cleanup even when the queued source had no database rows', async () => {
rows.set(epgChannels, []);
rows.set(epgPrograms, []);
requestEpgSource('queued-only');
(
epgWorkerService.clearEpgDataForSource as jest.Mock
).mockRejectedValueOnce(new Error('worker failure'));
await expect(reconcileEpgSources([])).rejects.toThrow('worker failure');
await reconcileEpgSources([]);
expect(epgWorkerService.clearEpgDataForSource).toHaveBeenCalledTimes(2);
});
it('preserves a new request received while the previous request is being cleared', async () => {
rows.set(epgChannels, []);
rows.set(epgPrograms, []);
requestEpgSource('requested-again');
(
epgWorkerService.clearEpgDataForSource as jest.Mock
).mockImplementationOnce(async () => {
requestEpgSource('requested-again');
});
await reconcileEpgSources([]);
await reconcileEpgSources([]);
expect(epgWorkerService.clearEpgDataForSource).toHaveBeenCalledTimes(2);
jest.clearAllMocks();
await reconcileEpgSources([]);
expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalled();
});
it('does not prune sources before playlist migration succeeds', async () => {
rows.set(appState, []);
await expect(reconcileEpgSources([])).rejects.toThrow(
@@ -7,7 +7,11 @@ import {
playlists,
} from '../database/schema';
import { epgWorkerService } from './epg-worker.service';
import { requestedEpgSources, retireEpgSource } from './epg-source-generation';
import {
forgetEpgSourceRequest,
requestedEpgSources,
retireEpgSource,
} from './epg-source-generation';
let reconciliation = Promise.resolve();
@@ -53,12 +57,13 @@ export function reconcileEpgSources(globalUrls: string[]): Promise<void> {
const programs = await db
.selectDistinct({ url: epgPrograms.sourceUrl })
.from(epgPrograms);
const requested = requestedEpgSources();
const removed = [
...new Set(
[
...channels.map((row) => row.url),
...programs.map((row) => row.url),
...requestedEpgSources(),
...requested.keys(),
].filter(
(url): url is string =>
!!url?.trim() && !active.has(url.trim())
@@ -67,8 +72,10 @@ export function reconcileEpgSources(globalUrls: string[]): Promise<void> {
];
// Fence the whole obsolete set before awaiting the first worker exit.
removed.forEach(retireEpgSource);
for (const url of removed)
for (const url of removed) {
await epgWorkerService.clearEpgDataForSource(url);
forgetEpgSourceRequest(url, requested.get(url.trim()));
}
});
reconciliation = next;
return next;
@@ -1,4 +1,8 @@
import { epgSourceGeneration, retireEpgSource } from './epg-source-generation';
import {
epgSourceGeneration,
requestEpgSource,
retireEpgSource,
} from './epg-source-generation';
import { app, BrowserWindow } from 'electron';
import * as path from 'path';
import { pathToFileURL } from 'url';
@@ -9,7 +13,8 @@ import {
} from '@iptvnator/shared/interfaces';
import { resolveWorkerRuntimeBootstrap } from '../workers/worker-runtime-paths';
export type EpgProgressStatus = 'queued' | 'loading' | 'complete' | 'error';
export type EpgProgressStatus =
'queued' | 'loading' | 'complete' | 'error' | 'cancelled';
export interface EpgProgressStats {
totalChannels: number;
@@ -59,12 +64,14 @@ export class EpgWorkerService {
error?: string,
queuePosition?: number,
errorCode?: ElectronBridgeSecurityErrorCode,
errorHost?: string
errorHost?: string,
generation = epgSourceGeneration(url)
): void {
const windows = BrowserWindow.getAllWindows();
windows.forEach((win) => {
win.webContents.send('EPG_PROGRESS_UPDATE', {
url,
generation,
status,
stats,
error,
@@ -114,7 +121,7 @@ export class EpgWorkerService {
url: string,
options: ElectronBridgeTrustOptions
): Promise<void> {
const generation = epgSourceGeneration(url);
const generation = requestEpgSource(url);
return new Promise((resolve, reject) => {
let worker: Worker;
try {
@@ -319,6 +319,7 @@ describe('EpgEvents', () => {
const { handleFetchEpg } = await import('./epg-fetch.service');
const { retireEpgSource } = await import('./epg-source-generation');
const { epgWorkerService } = await import('./epg-worker.service');
const progress = jest.spyOn(epgWorkerService, 'sendProgressToRenderer');
let finishFirst!: () => void;
const fetch = jest
.spyOn(epgWorkerService, 'fetchEpgFromUrl')
@@ -338,6 +339,17 @@ describe('EpgEvents', () => {
finishFirst();
await request;
expect(fetch).toHaveBeenCalledTimes(1);
expect(progress).toHaveBeenCalledWith(
'https://queued.example/guide.xml',
'cancelled',
undefined,
undefined,
undefined,
undefined,
undefined,
expect.any(Number)
);
progress.mockRestore();
fetch.mockRestore();
});
+5 -2
View File
@@ -1331,7 +1331,9 @@ not global XMLTV owners. No source-discovery or provider matching policy changes
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. Programmes are deleted by source; a globally keyed
existing source-clear worker. Successfully cleared request candidates are forgotten
without resetting their generation fences; failed cleanups remain retryable.
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
resolve another retained source sharing that channel ID. Legacy programmes with
@@ -1341,7 +1343,8 @@ There is no new schema migration.
Renderer reconciliation increments a data revision and cancels earlier lookup
subscriptions, clears program caches and the selected M3U guide, and refreshes
Xtream selection and visible channel previews. A delayed startup import is
Xtream selection and visible channel previews, plus Stalker manual mapping
overrides and bulk guides. A delayed startup import is
started only if its source still belongs to the reconciled configuration; its
completion observer is installed after settings initialization. Provider EPG
continues through its existing APIs. Playlist refresh is not EPG cache cleanup.
+17 -13
View File
@@ -24,13 +24,13 @@ Stalker now uses two EPG paths with different purposes:
settled bulk guide cannot answer fall back to throttled per-channel
`get_short_epg` through `StalkerEpgPreviewQueue` (see "Channel row preview
flow").
- Effect ordering matters: the eager-EPG effect is registered **after** the
playlist-change effect that calls `clearBulkItvEpgCache()`. On a portal
switch the cache is cleared first and then refilled; if the order is
reversed the clear clobbers the just-loaded bulk EPG on initial render.
- `ensureBulkItvEpg` de-duplicates (via `isLoadingBulkItvEpg` /
`bulkItvEpgLoaded` + matching playlist/period), so the eager trigger and the
play-time `loadEpgForChannel` path never double-fetch.
- Effect ordering matters: the eager-EPG effect is registered **after** the
playlist-change effect that calls `clearBulkItvEpgCache()`. On a portal
switch the cache is cleared first and then refilled; if the order is
reversed the clear clobbers the just-loaded bulk EPG on initial render.
- `ensureBulkItvEpg` de-duplicates (via `isLoadingBulkItvEpg` /
`bulkItvEpgLoaded` + matching playlist/period), so the eager trigger and the
play-time `loadEpgForChannel` path never double-fetch.
- If a portal does not return usable bulk data for the selected channel, the
active panel falls back to `get_short_epg`.
@@ -196,12 +196,12 @@ Normalization rules:
### Key files
| File | Purpose |
| -------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------- |
| `libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts` | bulk cache and fallback handling |
| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.ts` | active-channel EPG loading and controlled `app-epg-timeline` wiring |
| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.html` | active panel template |
| `libs/ui/epg/src/lib/epg-timeline/epg-timeline.component.ts` | shared controlled EPG timeline with date navigator |
| File | Purpose |
| ---------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------- |
| `libs/portal/stalker/data-access/src/lib/stores/features/with-stalker-epg.feature.ts` | bulk cache and fallback handling |
| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.ts` | active-channel EPG loading and controlled `app-epg-timeline` wiring |
| `libs/portal/stalker/feature/src/lib/stalker-live-stream-layout/stalker-live-stream-layout.component.html` | active panel template |
| `libs/ui/epg/src/lib/epg-timeline/epg-timeline.component.ts` | shared controlled EPG timeline with date navigator |
### Store API
@@ -283,6 +283,10 @@ rows stop consuming portal request capacity.
- Bulk EPG is fetched once per playlist session
- Channel switches only read from `bulkItvEpgByChannel`
- The cache is cleared when the Stalker playlist changes
- Committed XMLTV source reconciliation also clears mapping overrides, checked
IDs and bulk data, reloads the guide and selected mapping, and fences pending
mapping/bulk replies. Saved mappings remain authoritative even when removal
leaves their guide empty; portal EPG is not mixed into that empty override.
- This implementation does not add TTL-based refresh or background polling
## Authentication
@@ -63,6 +63,29 @@ describe('EpgProgressService', () => {
expect(epgBridge.onProgress).not.toHaveBeenCalled();
});
it('removes a retired queued import immediately on cancellation', () => {
epgBridge.supportsProgress = true;
const service = configureService();
const listener = (epgBridge.onProgress as jest.Mock).mock.calls[0][0];
const url = 'https://removed.example/guide.xml';
listener({ url, status: 'queued' });
expect(service.queuedCount()).toBe(1);
listener({ url, status: 'cancelled' });
expect(service.queuedCount()).toBe(0);
expect(service.isVisible()).toBe(false);
});
it('keeps a replacement import when an older queued batch finally cancels', () => {
epgBridge.supportsProgress = true;
const service = configureService();
const listener = (epgBridge.onProgress as jest.Mock).mock.calls[0][0];
const url = 'https://readded.example/guide.xml';
listener({ url, status: 'queued', generation: 0 });
listener({ url, status: 'loading', generation: 2 });
listener({ url, status: 'cancelled', generation: 0 });
expect(service.activeCount()).toBe(1);
});
it('does not force retry when EPG data management is disabled', () => {
const service = configureService();
@@ -108,6 +108,12 @@ export class EpgProgressService {
}
private updateProgress(progress: EpgImportProgress): void {
const current = this.importsMap().get(progress.url);
if ((progress.generation ?? 0) < (current?.generation ?? 0)) return;
if (progress.status === 'cancelled') {
this.removeImport(progress.url);
return;
}
this.importsMap.update((current) => {
const updated = new Map(current);
updated.set(progress.url, progress);
@@ -1,6 +1,10 @@
import { TestBed } from '@angular/core/testing';
import { signalStore, withState } from '@ngrx/signals';
import { DataService, RuntimeCapabilitiesService } from '@iptvnator/services';
import {
DataService,
EpgSourceSettingsService,
RuntimeCapabilitiesService,
} from '@iptvnator/services';
import { EpgRuntimeBridgeService } from '@iptvnator/epg/data-access';
import { EpgItem, Playlist } from '@iptvnator/shared/interfaces';
import { StalkerSessionService } from '../../stalker-session.service';
@@ -232,6 +236,44 @@ describe('withStalkerEpg', () => {
expect(epgBridge.getEpgMappingsBatch).toHaveBeenCalledTimes(1);
});
it('drops removed-source overrides and reloads the selected mapping without a portal reset', async () => {
epgBridge.getEpgMappingsBatch.mockResolvedValue({
'stalker:playlist-1:10001': 'mapped.channel.id',
});
epgBridge.getChannelPrograms.mockResolvedValue([MAPPED_PROGRAM]);
await store.applyMappedItvEpg(['10001']);
epgBridge.getChannelPrograms.mockResolvedValue([]);
const sources = TestBed.inject(EpgSourceSettingsService);
sources.revision.update((value) => value + 1);
sources.changed$.next();
expect(store.selectedItvEpgPrograms()).toEqual([]);
await store.applyMappedItvEpg(['10001']);
expect(store.selectedItvEpgPrograms()).toEqual([]);
// Saved mappings remain authoritative even when their source is gone.
expect(store.hasItvEpgMappingOverride('10001')).toBe(true);
});
it('ignores a mapped-program response that completes after source invalidation', async () => {
epgBridge.getEpgMappingsBatch.mockResolvedValue({
'stalker:playlist-1:10001': 'mapped.channel.id',
});
let resolvePrograms!: (value: unknown) => void;
epgBridge.getChannelPrograms.mockImplementationOnce(
() =>
new Promise((resolve) => {
resolvePrograms = resolve;
})
);
const pending = store.applyMappedItvEpg(['10001']);
await Promise.resolve();
const sources = TestBed.inject(EpgSourceSettingsService);
sources.revision.update((value) => value + 1);
sources.changed$.next();
resolvePrograms([MAPPED_PROGRAM]);
await pending;
expect(store.selectedItvEpgPrograms()).toEqual([]);
});
it('keeps overrides when ensureBulkItvEpg replaces the bulk record', async () => {
epgBridge.getEpgMappingsBatch.mockResolvedValue({
'stalker:playlist-1:10001': 'mapped.channel.id',
@@ -259,6 +301,63 @@ describe('withStalkerEpg', () => {
expect(record['10002']?.length).toBeGreaterThan(0);
});
it('does not mark later IDs checked after a stale mapped lookup rejects', async () => {
epgBridge.getEpgMappingsBatch.mockResolvedValue({
'stalker:playlist-1:10001': 'mapped.channel.id',
});
let rejectPrograms!: (error: Error) => void;
epgBridge.getChannelPrograms.mockImplementationOnce(
() =>
new Promise((_, reject) => {
rejectPrograms = reject;
})
);
const pending = store.applyMappedItvEpg(['10001', '10002']);
await Promise.resolve();
store.clearBulkItvEpgCache();
await store.applyMappedItvEpg(['10003']);
rejectPrograms(new Error('old request failed'));
await pending;
epgBridge.getEpgMappingsBatch.mockResolvedValue({
'stalker:playlist-1:10002': 'new.channel.id',
});
epgBridge.getChannelPrograms.mockResolvedValue([MAPPED_PROGRAM]);
await store.applyMappedItvEpg(['10002']);
expect(store.bulkItvEpgByChannel()['10002']).toEqual([
{ ...MAPPED_PROGRAM, channel: '10002' },
]);
});
it('allows initial bulk loading and mapping lookup to finish concurrently', async () => {
let finishBulk!: (value: unknown) => void;
const response = new Promise((resolve) => {
finishBulk = resolve;
});
dataService.sendIpcEvent.mockReturnValueOnce(response);
const bulk = store.ensureBulkItvEpg();
epgBridge.getEpgMappingsBatch.mockResolvedValue({
'stalker:playlist-1:10001': 'mapped.channel.id',
});
epgBridge.getChannelPrograms.mockResolvedValue([]);
await store.applyMappedItvEpg(['10001']);
finishBulk({
js: {
'10001': [
buildEntry(
'10001',
'Portal Show',
1744365600,
1744367400
),
],
},
});
await bulk;
expect(store.isLoadingBulkItvEpg()).toBe(false);
expect(store.bulkItvEpgLoaded()).toBe(true);
expect(store.selectedItvEpgPrograms()).toEqual([]);
});
it('does nothing when the mapping bridge is unsupported', async () => {
epgBridge.supportsEpgMapping = false;
@@ -3,12 +3,17 @@ import {
patchState,
signalStoreFeature,
withComputed,
withHooks,
withMethods,
withState,
} from '@ngrx/signals';
import { EpgRuntimeBridgeService } from '@iptvnator/epg/data-access';
import { createLogger } from '@iptvnator/portal/shared/util';
import { DataService, RuntimeCapabilitiesService } from '@iptvnator/services';
import {
DataService,
EpgSourceSettingsService,
RuntimeCapabilitiesService,
} from '@iptvnator/services';
import {
buildStalkerEpgMappingKey,
EpgItem,
@@ -105,7 +110,8 @@ export function withStalkerEpg() {
stalkerSession = inject(StalkerSessionService),
portalRepair = inject(StalkerPortalRepairService),
runtime = inject(RuntimeCapabilitiesService),
epgBridge = inject(EpgRuntimeBridgeService)
epgBridge = inject(EpgRuntimeBridgeService),
sources = inject(EpgSourceSettingsService)
) => {
const storeContext = store as typeof store &
StalkerEpgFeatureStoreContract;
@@ -127,6 +133,7 @@ export function withStalkerEpg() {
// an empty mapped guide must still keep the portal EPG out.
const mappingOwnedIds = new Set<string>();
let mappingPlaylistId: string | null = null;
let cacheGeneration = 0;
const resetMappingOverrides = (): void => {
mappingOverridesById.clear();
@@ -140,8 +147,8 @@ export function withStalkerEpg() {
EpgProgram[]
> => {
const record: Record<string, EpgProgram[]> = {};
for (const [id, programs] of mappingOverridesById) {
record[id] = programs;
for (const id of mappingOwnedIds) {
record[id] = mappingOverridesById.get(id) ?? [];
}
return record;
};
@@ -222,6 +229,7 @@ export function withStalkerEpg() {
}
const playlistId = String(playlist._id);
const generation = cacheGeneration;
if (!supportsEpg()) {
patchState(store, {
bulkItvEpgByChannel: {},
@@ -258,6 +266,7 @@ export function withStalkerEpg() {
type: 'itv',
period: String(periodHours),
});
if (generation !== cacheGeneration) return;
const selectedChannelId =
storeContext.selectedItvId() ?? null;
const bulkPrograms = extractBulkEpgByChannel(
@@ -274,6 +283,7 @@ export function withStalkerEpg() {
isLoadingBulkItvEpg: false,
});
} catch (error) {
if (generation !== cacheGeneration) return;
logger.warn('Bulk Stalker EPG unavailable', error);
patchState(store, {
bulkItvEpgByChannel: mappingOverridesRecord(),
@@ -312,8 +322,7 @@ export function withStalkerEpg() {
channelIds
.map((id) => normalizeStalkerEntityId(id))
.filter(
(id) =>
id && !mappingCheckedIds.has(id)
(id) => id && !mappingCheckedIds.has(id)
)
),
];
@@ -325,7 +334,11 @@ export function withStalkerEpg() {
// call can outlive a portal switch — bail out after
// every await instead of writing portal A's data
// into portal B's state.
const revision = sources.revision();
const generation = cacheGeneration;
const isStale = (): boolean =>
revision !== sources.revision() ||
generation !== cacheGeneration ||
mappingPlaylistId !== playlistId ||
String(
storeContext.currentPlaylist()?._id ?? ''
@@ -365,6 +378,7 @@ export function withStalkerEpg() {
let changed = false;
let ownershipChanged = false;
for (const [channelId, key] of keyById) {
if (isStale()) return;
const mappedEpgId = mappings[key]?.trim();
if (!mappedEpgId) {
// No mapping for this channel — a stable
@@ -443,12 +457,31 @@ export function withStalkerEpg() {
},
clearBulkItvEpgCache(): void {
cacheGeneration++;
resetMappingOverrides();
patchState(store, initialEpgState);
},
};
}
)
),
withHooks((store) => {
const sources = inject(EpgSourceSettingsService);
let subscription: { unsubscribe(): void } | undefined;
return {
onInit: () => {
subscription = sources.changed$.subscribe(() => {
const context = store as typeof store &
StalkerEpgFeatureStoreContract;
const selectedId = context.selectedItvId();
store.clearBulkItvEpgCache();
void store.ensureBulkItvEpg();
if (selectedId)
void store.applyMappedItvEpg([selectedId]);
});
},
onDestroy: () => subscription?.unsubscribe(),
};
})
);
}
@@ -85,6 +85,7 @@ export type ElectronBridgePlaylistType =
(typeof ELECTRON_BRIDGE_PLAYLIST_TYPES)[keyof typeof ELECTRON_BRIDGE_PLAYLIST_TYPES];
export const ELECTRON_BRIDGE_EPG_PROGRESS_STATUSES = {
Cancelled: 'cancelled',
Complete: 'complete',
Error: 'error',
Loading: 'loading',
@@ -385,6 +386,7 @@ export interface ElectronBridgeEpgProgressStats {
export interface ElectronBridgeEpgProgress {
url: string;
generation?: number;
status: ElectronBridgeEpgProgressStatus;
stats?: ElectronBridgeEpgProgressStats;
error?: string;
@@ -96,6 +96,8 @@ export class EpgProgressPanelComponent {
return 'check_circle';
case 'error':
return 'error';
case 'cancelled':
return 'cancel';
}
}