fix(epg): remove cached XMLTV data after source deletion

This commit is contained in:
4gray committed 2026-09-05 22:15:20 +02:00
1 parent b55a3cc322
commit ffcbee8140
32 files changed
+1051 -82

No files matched your search

+6
View File
@@ -0,0 +1,6 @@
---
type: fix
area: epg
---
Removing and saving an EPG source now clears its cached programmes, including data left by previously removed sources on restart. Other configured sources and playlist guides are preserved, and Live TV stops showing programmes from the removed source.
+13
View File
@@ -174,6 +174,19 @@ Button states use the matching group; "Total selected" counts the whole catalog.
Save persists the complete draft, Close discards it, and refresh restores hidden
categories by provider ID and type. See `docs/architecture/category-management.md`.
## XMLTV Source Removal
Saving Settings → EPG reconciles cached XMLTV with committed global URLs and
all enabled M3U playlist sources. Startup runs the same reconciliation after
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
invalidation prevent late results from restoring removed programmes. Provider
EPG is independent. See `docs/architecture/m3u-playlist-module.md`
("XMLTV source lifecycle").
## Portal Connectivity Preference
- Desktop Settings > General > Portal connections exposes default-on
+13
View File
@@ -1659,6 +1659,19 @@ No formal migration system yet. Schema changes are applied via raw SQL in the `c
<!-- nx configuration end-->
## XMLTV Source Removal
Saving Settings → EPG reconciles cached XMLTV with committed global URLs and
all enabled M3U playlist sources. Startup runs the same reconciliation after
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
invalidation prevent late results from restoring removed programmes. Provider
EPG is independent. See `docs/architecture/m3u-playlist-module.md`
("XMLTV source lifecycle").
## Portal Connectivity Preference
- Desktop Settings > General > Portal connections exposes default-on
+114 -13
View File
@@ -65,7 +65,7 @@ function formatXmltvDate(date: Date): string {
}
test.describe('Electron EPG', () => {
test('@epg @electron adds an EPG source, fetches guide data, removes the source row, and clears stored EPG data', async ({
test('@epg @electron adds an EPG source, fetches guide data, removes its stored EPG data on save', async ({
dataDir,
}) => {
const epgServer = await createMutableTextServer(epgFixtureXml, {
@@ -103,6 +103,8 @@ test.describe('Electron EPG', () => {
.first()
).toBeVisible();
await saveSettings(app.mainWindow);
await app.mainWindow
.locator('.epg-source-row button')
.nth(1)
@@ -116,18 +118,6 @@ test.describe('Electron EPG', () => {
// would block the close and leak the Electron process.
await saveSettings(app.mainWindow);
await app.mainWindow
.getByRole('button', { name: 'Clear EPG data' })
.click();
const dialog = app.mainWindow.locator('mat-dialog-container');
await expect(dialog).toBeVisible();
await dialog
.getByRole('button', { name: 'Yes', exact: true })
.click();
await app.mainWindow.waitForSelector('mat-dialog-container', {
state: 'detached',
});
await expect
.poll(() => getEpgChannelCount(app.mainWindow), {
timeout: 20000,
@@ -139,9 +129,99 @@ test.describe('Electron EPG', () => {
}
});
test('@epg @electron removes only the saved source, preserves shared IDs and mappings, and repairs old orphan data on restart', async ({
dataDir,
}) => {
test.setTimeout(120000);
const first = await createMutableTextServer(
createCurrentXmltvFixture(
'shared-news',
'Shared News',
'Removed Bulletin'
),
{ contentType: 'application/xml', resourcePath: '/first.xml' }
);
const second = await createMutableTextServer(
createCurrentXmltvFixture(
'shared-news',
'Shared News',
'Retained Bulletin'
),
{ contentType: 'application/xml', resourcePath: '/second.xml' }
);
let app = await launchElectronApp(dataDir);
try {
await openSettings(app.mainWindow);
await openSettingsSection(app.mainWindow, 'epg');
for (const url of [first.resourceUrl, second.resourceUrl]) {
await app.mainWindow
.getByRole('button', { name: 'Add EPG source' })
.click();
await app.mainWindow
.locator('.epg-source-row input')
.last()
.fill(url);
}
await saveSettings(app.mainWindow);
const programs = () =>
app.mainWindow.evaluate(async () =>
(
await window.electron.getChannelPrograms('mapped-news')
).map((p) => p.title)
);
await app.mainWindow.evaluate(() =>
window.electron.setEpgMapping('mapped-news', 'shared-news')
);
await expect
.poll(programs, { timeout: 30000 })
.toEqual(['Removed Bulletin', 'Retained Bulletin']);
await app.mainWindow
.locator('.epg-source-row')
.first()
.locator('button')
.nth(1)
.click();
// A staged removal must not delete data before Save.
expect(await programs()).toContain('Removed Bulletin');
await saveSettings(app.mainWindow);
await expect.poll(programs).toEqual(['Retained Bulletin']);
await closeElectronApp(app);
app = await launchElectronApp(dataDir);
await expect.poll(programs).toEqual(['Retained Bulletin']);
// Simulate a cache left by 0.23: import a source absent from settings.
await app.mainWindow.evaluate(
(url) => window.electron.forceFetchEpg(url),
first.resourceUrl
);
await expect.poll(programs).toContain('Removed Bulletin');
await closeElectronApp(app);
app = await launchElectronApp(dataDir);
await expect.poll(programs).toEqual(['Retained 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 expect.poll(programs).toEqual([]);
await expect.poll(() => getEpgChannelCount(app.mainWindow)).toBe(0);
expect(
await app.mainWindow.evaluate(() =>
window.electron.getEpgMapping('mapped-news')
)
).toMatchObject({ epgChannelId: 'shared-news' });
} finally {
await closeElectronApp(app);
await first.close();
await second.close();
}
});
test('@epg @electron imports and renders an EPG source declared by an M3U playlist header', async ({
dataDir,
}) => {
test.setTimeout(90000);
const epgServer = await createMutableTextServer(
createCurrentXmltvFixture(
'playlist-guide-news',
@@ -212,6 +292,27 @@ test.describe('Electron EPG', () => {
'Playlist Scoped Bulletin',
{ timeout: 30000 }
);
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(epgServer.resourceUrl);
await saveSettings(app.mainWindow);
await app.mainWindow
.locator('.epg-source-row button')
.nth(1)
.click();
await saveSettings(app.mainWindow);
// The same URL still belongs to the saved M3U playlist.
const retained = await app.mainWindow.evaluate(async () =>
window.electron.getChannelPrograms('playlist-guide-news')
);
expect(retained.map((program) => program.title)).toContain(
'Playlist Scoped Bulletin'
);
} finally {
await closeElectronApp(app);
await playlistServer.close();
+93 -10
View File
@@ -5,6 +5,7 @@ import {
channelItemByTitle,
clickCategoryByNameExact,
closeElectronApp,
createMutableTextServer,
expect,
goToDashboard,
launchElectronApp,
@@ -28,6 +29,93 @@ const epgCredentials = {
password: 'epg',
};
test('@epg @xtream @electron removes uploaded guide data and restores provider EPG without restarting', async ({
dataDir,
request,
}) => {
test.setTimeout(120000);
await resetMockServers(request, ['xtream']);
const fixture = await fetchXtreamEpgFixture(request, epgCredentials);
const id = fixture.stream.epg_channel_id;
expect(id).toBeTruthy();
const stamp = (date: Date) =>
date.toISOString().replace(/[-:T]/g, '').slice(0, 14);
const xml = `<tv><channel id="${id}"><display-name>Uploaded Guide</display-name></channel>
<programme channel="${id}" start="${stamp(new Date(Date.now() - 600000))} +0000" stop="${stamp(new Date(Date.now() + 3600000))} +0000"><title>Temporary XMLTV Bulletin</title></programme></tv>`;
const source = await createMutableTextServer(xml, {
contentType: 'application/xml',
resourcePath: '/temporary.xml',
});
const app = await launchElectronApp(dataDir);
try {
await addXtreamPortal(app.mainWindow, {
name: epgPortalName,
...epgCredentials,
});
await waitForXtreamWorkspaceReady(app.mainWindow);
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 app.mainWindow
.getByTestId('toggle-prefer-uploaded-epg')
.locator('input')
.check();
await saveSettings(app.mainWindow);
await expect
.poll(
() =>
app.mainWindow.evaluate(
async (channelId) =>
(
await window.electron.getChannelPrograms(
channelId!
)
).map((p) => p.title),
id
),
{ timeout: 30000 }
)
.toContain('Temporary XMLTV Bulletin');
await openWorkspaceSection(app.mainWindow, 'Live TV');
await clickCategoryByNameExact(app.mainWindow, fixture.categoryName);
const row = channelItemByTitle(
app.mainWindow,
fixture.stream.name ?? ''
).first();
await expect(row.locator('.epg-title')).toHaveText(
'Temporary XMLTV Bulletin'
);
await row.click();
await expect
.poll(() => timelineBlockTitles(app.mainWindow))
.toContain('Temporary XMLTV 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.locator('.epg-title')).toHaveText(
fixture.shortEpg[0].title
);
await row.click();
await expect
.poll(() => timelineBlockTitles(app.mainWindow))
.not.toContain('Temporary XMLTV Bulletin');
await expect
.poll(() => timelineBlockTitles(app.mainWindow))
.toContain(fixture.fullEpg[0].title);
} 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,
@@ -66,9 +154,7 @@ for (const timeZone of ['UTC', 'Europe/Berlin'] as const) {
// Provider-declared catch-up (tv_archive=1 in the fixture) is
// surfaced as a badge on the sidebar row (#1128).
await expect(
channelRow.getByTestId('catchup-badge')
).toBeVisible();
await expect(channelRow.getByTestId('catchup-badge')).toBeVisible();
// Sidebar channel list shows the per-channel "now" programme line.
await expect
@@ -263,10 +349,7 @@ test('@radio @stalker @electron keeps radio rows compact at narrow widths', asyn
await expect(categoryButton).toBeVisible();
await categoryButton.click();
const radioRow = channelItemByTitle(
app.mainWindow,
firstTitle
).first();
const radioRow = channelItemByTitle(app.mainWindow, firstTitle).first();
await expect(radioRow).toBeVisible({ timeout: 20000 });
await expect(radioRow).toHaveClass(/compact/);
await expect(radioRow.locator('.epg-placeholder')).toHaveCount(0);
@@ -301,7 +384,8 @@ test('@epg @xtream @electron shifts the sidebar preview and the timeline "now" b
// The guide runs an hour ahead of the real schedule, so the programme
// actually on air is the one the provider files under "an hour ago" —
// and every surface must agree on that same programme.
const providerNowSeconds = Math.floor(Date.now() / 1000) - offsetMinutes * 60;
const providerNowSeconds =
Math.floor(Date.now() / 1000) - offsetMinutes * 60;
const expectedProgram = fixture.fullEpg.find(
(listing) =>
listing.startTimestamp <= providerNowSeconds &&
@@ -341,8 +425,7 @@ test('@epg @xtream @electron shifts the sidebar preview and the timeline "now" b
await expect
.poll(async () =>
(
(await channelRow.locator('.epg-title').textContent()) ??
''
(await channelRow.locator('.epg-title').textContent()) ?? ''
).trim()
)
.toBe(expectedProgram.title);
@@ -524,6 +524,12 @@ export const epgPreloadCases: PreloadInvokeCase[] = [
channel: 'EPG_CLEAR_ALL',
forwardedArgs: [],
},
{
method: 'reconcileEpgSources',
args: [epgUrls],
channel: 'EPG_RECONCILE_SOURCES',
forwardedArgs: [{ urls: epgUrls }],
},
{
method: 'clearEpgDataForSource',
args: ['https://example.com/guide.xml'],
@@ -615,11 +615,7 @@ const electronApi: ElectronBridgeApi = {
sessionId: string,
style: EmbeddedMpvSubtitleStyle
): Promise<EmbeddedMpvSession | null> =>
ipcRenderer.invoke(
'EMBEDDED_MPV_SET_SUBTITLE_STYLE',
sessionId,
style
),
ipcRenderer.invoke('EMBEDDED_MPV_SET_SUBTITLE_STYLE', sessionId, style),
selectEmbeddedMpvSubtitleFile: (): Promise<string | null> =>
ipcRenderer.invoke('EMBEDDED_MPV_SELECT_SUBTITLE_FILE'),
setEmbeddedMpvSpeed: (
@@ -681,6 +677,8 @@ const electronApi: ElectronBridgeApi = {
forceFetchEpg: (url: string, options?: ElectronBridgeTrustOptions) =>
ipcRenderer.invoke('EPG_FORCE_FETCH', { url, options }),
clearEpgData: () => ipcRenderer.invoke('EPG_CLEAR_ALL'),
reconcileEpgSources: (urls: string[]) =>
ipcRenderer.invoke('EPG_RECONCILE_SOURCES', { urls }),
clearEpgDataForSource: (sourceUrl: string) =>
ipcRenderer.invoke('EPG_CLEAR_SOURCE', { sourceUrl }),
checkEpgFreshness: (urls: string[], maxAgeHours?: number) =>
@@ -1114,8 +1112,7 @@ const electronApi: ElectronBridgeApi = {
recordingsUpdatePrograms: (
targetPath: string,
programs: RecordingProgramSnapshot[]
) =>
ipcRenderer.invoke('RECORDINGS_UPDATE_PROGRAMS', targetPath, programs),
) => ipcRenderer.invoke('RECORDINGS_UPDATE_PROGRAMS', targetPath, programs),
recordingsRevealFile: (filePath: string) =>
ipcRenderer.invoke('RECORDINGS_REVEAL_FILE', filePath),
recordingsPlayFile: (filePath: string) =>
@@ -1,3 +1,4 @@
import { epgSourceGeneration } from './epg-source-generation';
import { eq } from 'drizzle-orm';
import { ElectronBridgeTrustOptions } from '@iptvnator/shared/interfaces';
import { getDatabase } from '../database/connection';
@@ -45,6 +46,7 @@ export async function checkEpgFreshness(
const db = await getDatabase();
for (const url of urls) {
const generation = epgSourceGeneration(url);
if (!url?.trim()) continue;
const result = await db
@@ -58,6 +60,7 @@ export async function checkEpgFreshness(
result[0].updatedAt &&
result[0].updatedAt >= cutoffTime;
if (generation !== epgSourceGeneration(url)) continue;
if (isFresh) {
freshUrls.push(url);
epgWorkerService.markFetchedUrl(url);
@@ -95,7 +98,12 @@ export async function handleFetchEpg(
urls: string[],
options: ElectronBridgeTrustOptions = {}
): Promise<EpgFetchResult> {
const validUrls = urls.filter((url) => url?.trim());
const validUrls = urls
.filter((url) => url?.trim())
.map((url) => url.trim());
const generations = new Map(
validUrls.map((url) => [url, epgSourceGeneration(url)])
);
if (validUrls.length === 0) {
return { success: false, message: 'No valid URLs provided' };
@@ -118,7 +126,9 @@ export async function handleFetchEpg(
// a 'queued' status, then fetchEpgFromUrl silently skips the URL and no
// completion update ever arrives, leaving the UI stuck at "queued".
const urlsToFetch = staleUrls.filter(
(url) => !epgWorkerService.hasFetchedUrl(url)
(url) =>
generations.get(url) === epgSourceGeneration(url) &&
!epgWorkerService.hasFetchedUrl(url)
);
if (urlsToFetch.length === 0) {
@@ -142,6 +152,7 @@ export async function handleFetchEpg(
const errors: string[] = [];
for (const url of urlsToFetch) {
try {
if (generations.get(url) !== epgSourceGeneration(url)) continue;
await epgWorkerService.fetchEpgFromUrl(url, options);
} catch (error) {
console.error(
@@ -149,9 +160,7 @@ export async function handleFetchEpg(
`Error fetching EPG from ${url}:`,
error
);
errors.push(
error instanceof Error ? error.message : String(error)
);
errors.push(error instanceof Error ? error.message : String(error));
}
}
@@ -0,0 +1,13 @@
/** Retires queued imports as well as workers already parsing a removed URL. */
const generations = new Map<string, number>();
export function epgSourceGeneration(url: string): number {
const key = url.trim();
if (!generations.has(key)) generations.set(key, 0);
return generations.get(key)!;
}
export function retireEpgSource(url: string): void {
generations.set(url.trim(), epgSourceGeneration(url) + 1);
}
export function requestedEpgSources(): string[] {
return [...generations.keys()];
}
@@ -0,0 +1,81 @@
import {
appState,
epgChannels,
epgPrograms,
playlists,
} from '../database/schema';
import { getDatabase } from '../database/connection';
import { epgWorkerService } from './epg-worker.service';
import { epgSourceGeneration } from './epg-source-generation';
import { reconcileEpgSources } from './epg-source-settings.service';
jest.mock('../database/connection', () => ({ getDatabase: jest.fn() }));
jest.mock('./epg-worker.service', () => ({
epgWorkerService: {
clearEpgDataForSource: jest.fn().mockResolvedValue(undefined),
},
}));
describe('committed EPG source reconciliation', () => {
let rows: Map<unknown, unknown[]>;
beforeEach(() => {
jest.clearAllMocks();
rows = new Map<unknown, unknown[]>([
[appState, [{ value: '1' }]],
[playlists, []],
[epgChannels, [{ url: 'removed' }, { url: 'shared' }]],
[epgPrograms, [{ url: 'removed' }, { url: 'second' }]],
]);
const select = () => ({
from: (table: unknown) => {
const result = rows.get(table) ?? [];
return Object.assign(Promise.resolve(result), {
where: () => Promise.resolve(result),
});
},
});
(getDatabase as jest.Mock).mockResolvedValue({
select,
selectDistinct: select,
});
});
it('preserves global and every enabled M3U source, including an overlapping channel owner', async () => {
rows.set(playlists, [{ type: 'm3u-url', urls: '["shared"]' }]);
const generation = epgSourceGeneration('removed');
await reconcileEpgSources([' second ']);
expect(epgWorkerService.clearEpgDataForSource).toHaveBeenCalledWith(
'removed'
);
expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalledWith(
'shared'
);
expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalledWith(
'second'
);
expect(epgSourceGeneration('removed')).toBeGreaterThan(generation);
});
it('cleans already orphaned sources when the saved global list is empty', async () => {
await reconcileEpgSources([]);
for (const url of ['removed', 'shared', 'second']) {
expect(epgWorkerService.clearEpgDataForSource).toHaveBeenCalledWith(
url
);
}
});
it('does not prune sources before playlist migration succeeds', async () => {
rows.set(appState, []);
await expect(reconcileEpgSources([])).rejects.toThrow(
'migrated playlists'
);
expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalled();
});
it('does not guess ownership when persisted playlist source metadata is invalid', async () => {
rows.set(playlists, [{ type: 'm3u-url', urls: '{broken' }]);
await expect(reconcileEpgSources([])).rejects.toThrow();
expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalled();
});
});
@@ -0,0 +1,75 @@
import { eq } from 'drizzle-orm';
import { getDatabase } from '../database/connection';
import {
appState,
epgChannels,
epgPrograms,
playlists,
} from '../database/schema';
import { epgWorkerService } from './epg-worker.service';
import { requestedEpgSources, retireEpgSource } from './epg-source-generation';
let reconciliation = Promise.resolve();
/** Called only with settings successfully read from or written to IndexedDB. */
export function reconcileEpgSources(globalUrls: string[]): Promise<void> {
const next = reconciliation
.catch(() => undefined)
.then(async () => {
const db = await getDatabase();
const migration = await db
.select()
.from(appState)
.where(
eq(appState.key, 'm3u-playlists-indexeddb-to-sqlite-v1')
);
// A failed/incomplete migration cannot establish all playlist owners.
if (migration[0]?.value !== '1') {
throw new Error(
'EPG source reconciliation requires migrated playlists'
);
}
const active = new Set(
globalUrls.map((url) => url.trim()).filter(Boolean)
);
const savedPlaylists = await db
.select({ type: playlists.type, urls: playlists.epgUrls })
.from(playlists);
for (const playlist of savedPlaylists) {
if (!playlist.type.startsWith('m3u-')) continue;
const urls: unknown = JSON.parse(playlist.urls || '[]');
if (
!Array.isArray(urls) ||
urls.some((url) => typeof url !== 'string')
) {
throw new Error('Invalid saved playlist EPG ownership');
}
for (const url of urls as string[])
if (url.trim()) active.add(url.trim());
}
const channels = await db
.selectDistinct({ url: epgChannels.sourceUrl })
.from(epgChannels);
const programs = await db
.selectDistinct({ url: epgPrograms.sourceUrl })
.from(epgPrograms);
const removed = [
...new Set(
[
...channels.map((row) => row.url),
...programs.map((row) => row.url),
...requestedEpgSources(),
].filter(
(url): url is string =>
!!url?.trim() && !active.has(url.trim())
)
),
];
// Fence the whole obsolete set before awaiting the first worker exit.
removed.forEach(retireEpgSource);
for (const url of removed)
await epgWorkerService.clearEpgDataForSource(url);
});
reconciliation = next;
return next;
}
@@ -1,3 +1,4 @@
import { epgSourceGeneration, retireEpgSource } from './epg-source-generation';
import { app, BrowserWindow } from 'electron';
import * as path from 'path';
import { pathToFileURL } from 'url';
@@ -113,6 +114,7 @@ export class EpgWorkerService {
url: string,
options: ElectronBridgeTrustOptions
): Promise<void> {
const generation = epgSourceGeneration(url);
return new Promise((resolve, reject) => {
let worker: Worker;
try {
@@ -201,6 +203,7 @@ export class EpgWorkerService {
scheduleFetchTimeout();
worker.on('message', async (message: EpgWorkerMessage) => {
if (settled || generation !== epgSourceGeneration(url)) return;
try {
switch (message.type) {
case 'READY':
@@ -371,6 +374,8 @@ export class EpgWorkerService {
return;
}
retireEpgSource(normalizedSourceUrl);
this.fetchedUrls.delete(normalizedSourceUrl);
const runningWorker = this.workers.get(normalizedSourceUrl);
if (runningWorker) {
this.workers.delete(normalizedSourceUrl);
@@ -287,6 +287,60 @@ describe('EpgEvents', () => {
await flushPromises();
});
it('does not revive a retired source from late worker READY or COMPLETE messages', async () => {
const service = new EpgWorkerService('[Test EPG]', 1000);
const url = 'https://removed.example/guide.xml';
const fetch = service.fetchEpgFromUrl(url).catch(() => undefined);
const worker = mockWorkerInstances[0];
let finishTermination!: () => void;
worker.terminate.mockReturnValue(
new Promise<void>((resolve) => {
finishTermination = resolve;
})
);
const clear = service.clearEpgDataForSource(url);
worker.emit('message', { type: 'READY' });
worker.emit('message', { type: 'EPG_COMPLETE' });
expect(worker.postMessage).not.toHaveBeenCalled();
expect(service.hasFetchedUrl(url)).toBe(false);
expect(mockWorkerInstances).toHaveLength(1);
worker.emit('exit', 1);
finishTermination();
await flushPromises();
const clearWorker = mockWorkerInstances[1];
clearWorker.emit('message', { type: 'READY' });
clearWorker.emit('message', { type: 'CLEAR_COMPLETE' });
await clear;
await fetch;
});
it('does not start a queued source removed while an earlier source imports', async () => {
getDatabase.mockRejectedValue(new Error('force stale for test'));
const { handleFetchEpg } = await import('./epg-fetch.service');
const { retireEpgSource } = await import('./epg-source-generation');
const { epgWorkerService } = await import('./epg-worker.service');
let finishFirst!: () => void;
const fetch = jest
.spyOn(epgWorkerService, 'fetchEpgFromUrl')
.mockImplementationOnce(
() =>
new Promise<void>((resolve) => {
finishFirst = resolve;
})
)
.mockResolvedValue(undefined);
const request = handleFetchEpg([
'https://first.example/guide.xml',
'https://queued.example/guide.xml',
]);
await flushPromises();
retireEpgSource('https://queued.example/guide.xml');
finishFirst();
await request;
expect(fetch).toHaveBeenCalledTimes(1);
fetch.mockRestore();
});
it('clears one EPG source through a worker and allows it to be fetched again', async () => {
const workerService = new EpgWorkerService('[Test EPG]', 1000);
const sourceUrl = 'https://playlist.example.com/guide.xml';
@@ -381,7 +435,10 @@ describe('EpgEvents', () => {
* inserts extra queries that the original test didn't anticipate.
*/
function queryChain<T>(data: T) {
const chain: Record<string, jest.Mock> = {} as Record<string, jest.Mock>;
const chain: Record<string, jest.Mock> = {} as Record<
string,
jest.Mock
>;
chain.where = jest.fn().mockReturnValue(chain);
chain.innerJoin = jest.fn().mockReturnValue(chain);
chain.groupBy = jest.fn().mockReturnValue(chain);
@@ -393,25 +450,30 @@ describe('EpgEvents', () => {
// getMapping queries (must come first — they return empty)
const from = jest
.fn()
.mockReturnValueOnce(queryChain([])) // 1. getMapping: epgChannelMappings
.mockReturnValueOnce(queryChain([])) // 2. getMapping: content table
.mockReturnValueOnce(queryChain([])) // 1. getMapping: epgChannelMappings
.mockReturnValueOnce(queryChain([])) // 2. getMapping: content table
// Test expectations below
.mockReturnValueOnce(queryChain([])) // 3. selectChannelPrograms (bbc.one.uk → empty)
.mockReturnValueOnce(queryChain([{ id: 'BBC.ONE.UK', displayName: 'BBC One' }])) // 4. selectChannelById → channel found
.mockReturnValueOnce(queryChain([ // 5. selectChannelPrograms (BBC.ONE.UK → programs)
{
id: 1,
channelId: 'BBC.ONE.UK',
start: '2026-04-14T10:00:00Z',
stop: '2026-04-14T11:00:00Z',
title: 'News',
description: null,
category: null,
iconUrl: null,
rating: null,
episodeNum: null,
},
]));
.mockReturnValueOnce(queryChain([])) // 3. selectChannelPrograms (bbc.one.uk → empty)
.mockReturnValueOnce(
queryChain([{ id: 'BBC.ONE.UK', displayName: 'BBC One' }])
) // 4. selectChannelById → channel found
.mockReturnValueOnce(
queryChain([
// 5. selectChannelPrograms (BBC.ONE.UK → programs)
{
id: 1,
channelId: 'BBC.ONE.UK',
start: '2026-04-14T10:00:00Z',
stop: '2026-04-14T11:00:00Z',
title: 'News',
description: null,
category: null,
iconUrl: null,
rating: null,
episodeNum: null,
},
])
);
select.mockImplementation(() => ({ from }));
@@ -1,3 +1,4 @@
import { reconcileEpgSources } from './epg-source-settings.service';
import { ipcMain } from 'electron';
import {
ElectronBridgeCurrentProgramsOptions,
@@ -29,6 +30,14 @@ export default class EpgEvents {
* Bootstrap EPG events
*/
static bootstrapEpgEvents(): Electron.IpcMain {
ipcMain.handle(
'EPG_RECONCILE_SOURCES',
async (_event, args: { urls: string[] }) => {
await reconcileEpgSources(args.urls);
return { success: true };
}
);
ipcMain.handle(
'FETCH_EPG',
async (
@@ -179,10 +188,7 @@ export default class EpgEvents {
ipcMain.handle(
'EPG_CHANNEL_SEARCH',
async (
_event,
args: { searchTerm: string; limit?: number }
) => {
async (_event, args: { searchTerm: string; limit?: number }) => {
return handleSearchEpgChannels(args.searchTerm, args.limit);
}
);
@@ -267,6 +267,20 @@ export class EpgDatabaseSourceClearOperation {
const clearSource = this.db.transaction((url: string) => {
this.deleteProgramsForSourceStmt.run(url);
// Channels have one legacy global ID. Transfer their ownership
// to a surviving source before a subsequent source removal.
this.db
.prepare(
`UPDATE epg_channels SET source_url = (
SELECT source_url FROM epg_programs
WHERE channel_id = epg_channels.id AND source_url IS NOT NULL
ORDER BY source_url LIMIT 1
) WHERE source_url = ? AND EXISTS (
SELECT 1 FROM epg_programs
WHERE channel_id = epg_channels.id AND source_url IS NOT NULL
)`
)
.run(url);
this.deleteOrphanChannelsForSourceStmt.run(url);
});
+37 -1
View File
@@ -14,7 +14,12 @@ import {
} from '@iptvnator/workspace/shell/util';
import { MockProvider } from 'ng-mocks';
import { EMPTY, of } from 'rxjs';
import { DataService, RuntimeCapabilitiesService } from '@iptvnator/services';
import {
DataService,
EpgSourceSettingsService,
SettingsStore,
RuntimeCapabilitiesService,
} from '@iptvnator/services';
import {
Language,
Settings,
@@ -235,6 +240,37 @@ describe('AppComponent', () => {
expect(snackBar.open).not.toHaveBeenCalled();
});
it('does not reimport a deleted source from a late startup freshness response', async () => {
await TestBed.inject(SettingsStore).loadSettings();
let finishFreshness!: (value: {
freshUrls: string[];
staleUrls: string[];
}) => void;
epgBridge.checkFreshness = jest.fn(
() =>
new Promise((resolve) => {
finishFreshness = resolve;
})
);
const pending = (
component as unknown as {
fetchStaleEpgData(urls: string[]): Promise<void>;
}
).fetchStaleEpgData(['https://removed.example/guide.xml']);
await Promise.resolve();
const sources = TestBed.inject(EpgSourceSettingsService);
sources.revision.update((value) => value + 1);
sources.changed$.next();
finishFreshness({
freshUrls: [],
staleUrls: ['https://removed.example/guide.xml'],
});
await pending;
expect(epgService.fetchEpg).not.toHaveBeenCalledWith([
'https://removed.example/guide.xml',
]);
});
it('does not fetch EPG settings when the EPG bridge cannot import EPG', async () => {
const settings: Settings = {
...DEFAULT_SETTINGS,
+13 -4
View File
@@ -20,6 +20,7 @@ import {
DataService,
RuntimeCapabilitiesService,
SettingsStore,
EpgSourceSettingsService,
} from '@iptvnator/services';
import {
AUTO_UPDATE_PLAYLISTS,
@@ -63,6 +64,7 @@ export class AppComponent implements OnInit {
private translate = inject(TranslateService);
private settingsService = inject(SettingsService);
private settingsStore = inject(SettingsStore);
private readonly epgSources = inject(EpgSourceSettingsService);
private playbackKeepAwake = inject(PlaybackKeepAwakeService);
private playlistOpenRequests = inject(PlaylistOpenRequestService);
private runtime = inject(RuntimeCapabilitiesService);
@@ -188,8 +190,15 @@ export class AppComponent implements OnInit {
* Data is considered fresh if updated within the last 12 hours.
*/
private async fetchStaleEpgData(urls: string[]): Promise<void> {
await this.settingsStore.loadSettings();
const revision = this.epgSources.revision();
const fetchCurrentSources = (sources: string[]) => {
this.epgService.fetchEpg(
this.epgSources.retainCurrentSources(sources, revision)
);
};
if (!this.epgBridge.supportsSourceFreshness) {
this.epgService.fetchEpg(urls);
fetchCurrentSources(urls);
return;
}
@@ -197,7 +206,7 @@ export class AppComponent implements OnInit {
const result = await this.epgBridge.checkFreshness(urls, 12);
if (!result) {
this.epgService.fetchEpg(urls);
fetchCurrentSources(urls);
return;
}
@@ -219,12 +228,12 @@ export class AppComponent implements OnInit {
debugAppComponent(
`EPG: Fetching ${result.staleUrls.length} stale source(s)`
);
this.epgService.fetchEpg(result.staleUrls);
fetchCurrentSources(result.staleUrls);
}
} catch (error) {
console.error('Error checking EPG freshness, fetching all:', error);
// Fallback: fetch all URLs if freshness check fails
this.epgService.fetchEpg(urls);
fetchCurrentSources(urls);
}
}
@@ -16,7 +16,10 @@ import { MatIconModule } from '@angular/material/icon';
import { ActivatedRoute, Router } from '@angular/router';
import { SettingsContextService } from '@iptvnator/workspace/shell/util';
import { TranslateModule, TranslateService } from '@ngx-translate/core';
import { RuntimeCapabilitiesService } from '@iptvnator/services';
import {
EpgSourceReconciliationError,
RuntimeCapabilitiesService,
} from '@iptvnator/services';
import { VodSourceDiscoveryService } from '@iptvnator/portal/shared/data-access';
import { Language, StreamFormat } from '@iptvnator/shared/interfaces';
import { firstValueFrom, map } from 'rxjs';
@@ -307,7 +310,13 @@ export class SettingsComponent
try {
await this.form.save(() => this.applyChangedSettings());
return true;
} catch {
} catch (error) {
if (error instanceof EpgSourceReconciliationError) {
this.settingsSnackbar.open(
this.translate.instant('SETTINGS.EPG_DATA_CLEAR_FAILED')
);
return false;
}
// The store already applied the change in memory, so without
// this the save looks successful until the next restart. The
// unsaved-changes bar stays visible so it can be retried.
+34
View File
@@ -1311,3 +1311,37 @@ Routes live in `libs/playlist/m3u/feature-player/src/lib/m3u-workspace.routes.ts
3. Dispatch `FavoritesActions.hydrateFavorites` only when copying values that
were already read from persistence into NgRx
4. Effects persist the two user-mutation actions; hydration is reducer-only
### XMLTV source lifecycle
Electron treats `Settings.epgUrl` as the committed global source list. Removing
an input is a draft edit; only a successful IndexedDB write authorizes source
reconciliation. A storage write failure restores the previous in-memory EPG
URLs. If subsequent cache cleanup fails, the saved URLs remain authoritative,
the settings form remains retryable and shows the existing EPG cleanup failure
message rather than claiming settings storage failed.
`EpgSourceSettingsService` waits for `PlaylistsService.getAllPlaylists()` (which
performs the legacy playlist migration) before invoking `EPG_RECONCILE_SOURCES`.
Startup includes an empty global list and does not prune after a failed settings
read. Main verifies the completed migration flag, reads every enabled M3U
`epg_urls` list and unions them with the saved globals. Invalid ownership metadata
aborts pruning. Detected-but-disabled sources and Xtream/Stalker provider EPG are
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
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
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
subscriptions, clears program caches and the selected M3U guide, and refreshes
Xtream selection and visible channel previews. 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.
@@ -2,7 +2,7 @@ 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 { SettingsStore } from '@iptvnator/services';
import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services';
import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service';
import { EpgService } from './epg.service';
@@ -11,6 +11,7 @@ describe('EpgService', () => {
let epgBridge: Partial<EpgRuntimeBridgeService>;
let snackBar: { open: jest.Mock };
let settingsStore: {
loadSettings: jest.Mock;
getSettings: jest.Mock;
getTrustOptions: jest.Mock;
resolvedEpgOffsetMinutes: jest.Mock;
@@ -29,6 +30,7 @@ describe('EpgService', () => {
open: jest.fn(),
};
settingsStore = {
loadSettings: jest.fn().mockResolvedValue(undefined),
getSettings: jest.fn(() => ({
epgUrl: [],
trustedPrivateNetworkEpgUrls: ['http://192.168.1.20/guide.xml'],
@@ -68,13 +70,81 @@ describe('EpgService', () => {
service = TestBed.inject(EpgService);
});
it('observes startup import completion after initial source reconciliation', async () => {
epgBridge.supportsImport = true;
let loaded!: () => void;
settingsStore.loadSettings.mockReturnValue(
new Promise<void>((resolve) => {
loaded = resolve;
})
);
const sources = TestBed.inject(EpgSourceSettingsService);
jest.spyOn(sources, 'retainCurrentSources').mockImplementation(
(urls) => urls
);
const pending = service.fetchEpg([
'https://configured.example/guide.xml',
]);
sources.revision.update((revision) => revision + 1);
sources.changed$.next();
const availability: boolean[] = [];
const subscription = service.epgAvailable$.subscribe((value) =>
availability.push(value)
);
loaded();
await pending;
await Promise.resolve();
expect(epgBridge.fetchEpg).toHaveBeenCalledTimes(1);
expect(availability.filter(Boolean)).toHaveLength(2);
subscription.unsubscribe();
});
it('does not launch an obsolete import deferred behind settings initialization', async () => {
epgBridge.supportsImport = true;
let loaded!: () => void;
settingsStore.loadSettings.mockReturnValue(
new Promise<void>((resolve) => {
loaded = resolve;
})
);
const pending = service.fetchEpg(['https://removed.example/guide.xml']);
const sources = TestBed.inject(EpgSourceSettingsService);
sources.revision.update((revision) => revision + 1);
sources.changed$.next();
loaded();
await pending;
expect(epgBridge.fetchEpg).not.toHaveBeenCalled();
});
it('prevents pending programmes from repopulating a cache after source deletion', async () => {
epgBridge.supportsProgramLookup = true;
let resolveOld!: (programs: unknown[]) => void;
(epgBridge.getChannelPrograms as jest.Mock).mockReturnValueOnce(
new Promise((resolve) => {
resolveOld = resolve;
})
);
const stale = jest.fn();
service.getCurrentProgramForChannel('deleted-channel').subscribe(stale);
const sources = TestBed.inject(EpgSourceSettingsService);
sources.revision.update((revision) => revision + 1);
sources.changed$.next();
resolveOld([]);
await Promise.resolve();
expect(stale).not.toHaveBeenCalled();
await firstValueFrom(
service.getCurrentProgramForChannel('deleted-channel')
);
expect(epgBridge.getChannelPrograms).toHaveBeenCalledTimes(2);
});
it('does not fetch EPG when bridge import support is disabled', () => {
service.fetchEpg(['https://example.com/epg.xml']);
expect(epgBridge.fetchEpg).not.toHaveBeenCalled();
});
it('fetches EPG through the EPG runtime bridge when import support is enabled', () => {
it('fetches EPG through the EPG runtime bridge when import support is enabled', async () => {
epgBridge.supportsImport = true;
service.fetchEpg([
@@ -84,6 +154,7 @@ describe('EpgService', () => {
' https://example.com/epg.xml ',
]);
await Promise.resolve();
expect(epgBridge.fetchEpg).toHaveBeenCalledWith(
['https://example.com/epg.xml', 'https://example.com/other.xml'],
{
+42 -5
View File
@@ -17,7 +17,7 @@ import {
EpgProgram,
epgProviderClockMs,
} from '@iptvnator/shared/interfaces';
import { SettingsStore } from '@iptvnator/services';
import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services';
import {
EpgLookupOptions,
EpgRuntimeBridgeService,
@@ -42,6 +42,15 @@ export class EpgService {
private translate = inject(TranslateService);
private readonly epgBridge = inject(EpgRuntimeBridgeService);
private readonly settingsStore = inject(SettingsStore);
private readonly sourceSettings = inject(EpgSourceSettingsService);
constructor() {
this.sourceSettings.changed$.subscribe(() => {
this.clearCache();
this.currentEpgPrograms.next([]);
this.epgAvailable.next(true);
});
}
private epgAvailable = new BehaviorSubject<boolean>(false);
private currentEpgPrograms = new BehaviorSubject<EpgProgram[]>([]);
@@ -79,11 +88,16 @@ export class EpgService {
/**
* Fetches EPG from the given URLs
*/
fetchEpg(urls: string[]): void {
async fetchEpg(urls: string[]): Promise<void> {
if (!this.epgBridge.supportsImport) return;
// Filter out empty and duplicate URLs and send all URLs at once.
const validUrls = normalizeEpgUrls(urls);
const revision = this.sourceSettings.revision();
await this.settingsStore.loadSettings();
const validUrls = this.sourceSettings.retainCurrentSources(
normalizeEpgUrls(urls),
revision
);
if (validUrls.length === 0) return;
from(
@@ -93,6 +107,7 @@ export class EpgService {
)
)
.pipe(
this.sourceSettings.guard(),
tap((result) => {
if (result === null) return;
@@ -123,6 +138,7 @@ export class EpgService {
from(this.epgBridge.getChannelPrograms(channelId))
.pipe(
this.sourceSettings.guard(),
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
catchError((err) => {
@@ -187,6 +203,7 @@ export class EpgService {
// Fetch from backend
return this.getCachedOrFetchCurrentProgram(cacheKey, () =>
from(this.epgBridge.getChannelPrograms(channelId)).pipe(
this.sourceSettings.guard(),
map((programs) => normalizeEpgPrograms(programs ?? [])),
map((programs: EpgProgram[]) =>
this.findCurrentProgram(programs)
@@ -287,6 +304,7 @@ export class EpgService {
nowMs: this.epgClockMs(),
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((batchResult) => {
const cacheTimestamp = Date.now();
@@ -311,6 +329,7 @@ export class EpgService {
// Fallback for older preload bundles without the batch endpoint.
const fetchObservables = channelsToFetch.map((channelId) =>
this.getCurrentProgramForChannel(channelId).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((program) => ({ channelId, program })),
catchError(() => of({ channelId, program: null }))
@@ -318,6 +337,7 @@ export class EpgService {
);
return forkJoin(fetchObservables).pipe(
this.sourceSettings.guard(),
map((results) => {
results.forEach((result) => {
resultMap.set(result.channelId, result.program);
@@ -372,6 +392,7 @@ export class EpgService {
sourceUrls,
fallbackSourceUrls
).pipe(
this.sourceSettings.guard(),
tap((fetchedMap) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
@@ -386,7 +407,14 @@ export class EpgService {
});
}),
finalize(() => {
this.fetchingCurrentProgramBatches.delete(batchCacheKey);
if (
this.fetchingCurrentProgramBatches.get(
batchCacheKey
) === request$
)
this.fetchingCurrentProgramBatches.delete(
batchCacheKey
);
}),
shareReplay({ bufferSize: 1, refCount: false })
);
@@ -396,6 +424,7 @@ export class EpgService {
}
return request$.pipe(
this.sourceSettings.guard(),
map((fetchedMap) => {
const mergedResultMap = new Map(resultMap);
channelsToFetch.forEach((channelId) => {
@@ -421,6 +450,7 @@ export class EpgService {
nowMs,
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
switchMap((scopedResult) => {
const resultMap = new Map<string, EpgProgram | null>();
@@ -448,6 +478,7 @@ export class EpgService {
nowMs,
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((globalResult) => {
fallbackChannelIds.forEach((channelId) => {
@@ -500,6 +531,7 @@ export class EpgService {
normalizedChannelIds,
effectiveSourceUrls
).pipe(
this.sourceSettings.guard(),
switchMap((metadataMap) => {
const fallbackChannelIds =
sourceUrls.length > 0 && globalSourceUrls.length > 0
@@ -516,6 +548,7 @@ export class EpgService {
fallbackChannelIds,
globalSourceUrls
).pipe(
this.sourceSettings.guard(),
map((globalMetadataMap) => {
fallbackChannelIds.forEach((channelId) => {
metadataMap.set(
@@ -561,6 +594,7 @@ export class EpgService {
sourceUrls.length > 0 ? { sourceUrls } : undefined
)
).pipe(
this.sourceSettings.guard(),
map((metadataByChannelId) => {
return new Map<string, EpgChannelMetadata | null>(
channelIds.map((channelId) => [
@@ -641,6 +675,7 @@ export class EpgService {
const offsetMinutes = this.epgOffsetMinutes();
const request$ = fetchProgram().pipe(
this.sourceSettings.guard(),
tap((program) => {
this.programCache.set(cacheKey, {
program,
@@ -649,7 +684,8 @@ export class EpgService {
});
}),
finalize(() => {
this.fetchingCurrentPrograms.delete(cacheKey);
if (this.fetchingCurrentPrograms.get(cacheKey) === request$)
this.fetchingCurrentPrograms.delete(cacheKey);
}),
shareReplay({ bufferSize: 1, refCount: false })
);
@@ -668,6 +704,7 @@ export class EpgService {
from(
this.epgBridge.getChannelPrograms(channelId, { sourceUrls })
).pipe(
this.sourceSettings.guard(),
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
switchMap((programs) => {
@@ -5,7 +5,7 @@ import {
EpgItem,
windowEpgItemsAtProviderClock,
} from '@iptvnator/shared/interfaces';
import { SettingsStore } from '@iptvnator/services';
import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services';
import { XtreamApiService, XtreamCredentials } from './xtream-api.service';
import { XtreamXmltvFallbackService } from './xtream-xmltv-fallback.service';
import { createLogger } from '@iptvnator/portal/shared/util';
@@ -50,6 +50,20 @@ export class EpgQueueService implements OnDestroy {
private readonly apiService = inject(XtreamApiService);
private readonly fallbackService = inject(XtreamXmltvFallbackService);
private readonly settingsStore = inject(SettingsStore);
private readonly sourceSubscription = inject(
EpgSourceSettingsService
).changed$.subscribe(() => {
this.enqueueGeneration++;
this.queue = [];
for (const id of new Set([
...this.cache.keys(),
...this.inFlight,
...this.visibleSet,
])) {
this.invalidate(id);
this.epgResult$.next({ streamId: id, items: [] });
}
});
private readonly logger = createLogger('EpgQueueService');
private readonly previewLimit = 3;
@@ -494,6 +508,7 @@ export class EpgQueueService implements OnDestroy {
}
ngOnDestroy(): void {
this.sourceSubscription.unsubscribe();
this.epgResult$.complete();
}
}
@@ -5,7 +5,7 @@ import {
epgProviderClockMs,
} from '@iptvnator/shared/interfaces';
import { createLogger } from '@iptvnator/portal/shared/util';
import { SettingsStore } from '@iptvnator/services';
import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services';
type ElectronEpgBridge = {
getChannelPrograms?: (channelId: string) => Promise<EpgProgram[]>;
@@ -19,6 +19,7 @@ type ElectronEpgBridge = {
export class XtreamXmltvFallbackService {
private readonly logger = createLogger('XtreamXmltvFallback');
private readonly settingsStore = inject(SettingsStore);
private readonly sources = inject(EpgSourceSettingsService);
/**
* `DataService.isElectron` is intentionally not consulted here: it
@@ -44,6 +45,7 @@ export class XtreamXmltvFallbackService {
async getProgramsForChannel(
epgChannelId: string | null | undefined
): Promise<EpgItem[]> {
const revision = this.sources.revision();
const id = (epgChannelId ?? '').trim();
if (!id) return [];
@@ -52,6 +54,7 @@ export class XtreamXmltvFallbackService {
try {
const programs = await fn.call(this.bridge, id);
if (revision !== this.sources.revision()) return [];
return (programs ?? []).map((p) => mapEpgProgramToEpgItem(p, id));
} catch (error) {
this.logger.error(`Failed to load XMLTV programs for ${id}`, error);
@@ -62,6 +65,7 @@ export class XtreamXmltvFallbackService {
async getCurrentProgramsBatch(
epgChannelIds: ReadonlyArray<string | null | undefined>
): Promise<Record<string, EpgItem>> {
const revision = this.sources.revision();
const fn = this.bridge?.getCurrentProgramsBatch;
if (typeof fn !== 'function') return {};
@@ -83,6 +87,7 @@ export class XtreamXmltvFallbackService {
this.settingsStore.resolvedEpgOffsetMinutes()
),
});
if (revision !== this.sources.revision()) return {};
const out: Record<string, EpgItem> = {};
for (const id of ids) {
const row = rows?.[id];
@@ -4,6 +4,7 @@ import {
signalStoreFeature,
withComputed,
withMethods,
withHooks,
withState,
} from '@ngrx/signals';
import {
@@ -11,7 +12,11 @@ import {
EpgItem,
epgProviderClockMs,
} from '@iptvnator/shared/interfaces';
import { RuntimeCapabilitiesService, SettingsStore } from '@iptvnator/services';
import {
EpgSourceSettingsService,
RuntimeCapabilitiesService,
SettingsStore,
} from '@iptvnator/services';
import {
XtreamApiService,
XtreamCredentials,
@@ -104,6 +109,7 @@ export function withEpg() {
const fallbackService = inject(XtreamXmltvFallbackService);
const runtime = inject(RuntimeCapabilitiesService);
const settingsStore = inject(SettingsStore);
const sources = inject(EpgSourceSettingsService);
const supportsEpg = (): boolean => runtime.supportsEpg;
@@ -146,6 +152,7 @@ export function withEpg() {
* sets `preferUploadedEpgOverXtream`.
*/
async loadEpg(): Promise<EpgItem[]> {
const sourceRevision = sources.revision();
if (!supportsEpg()) {
patchState(store, {
epgItems: [],
@@ -203,6 +210,7 @@ export function withEpg() {
fetchFullProvider(credentials, xtreamId),
});
if (sourceRevision !== sources.revision()) return [];
patchState(store, {
epgItems,
isLoadingEpg: false,
@@ -210,6 +218,7 @@ export function withEpg() {
return epgItems;
} catch (error) {
if (sourceRevision !== sources.revision()) return [];
logger.error('Error loading EPG', error);
patchState(store, {
epgItems: [],
@@ -255,6 +264,19 @@ export function withEpg() {
patchState(store, initialEpgState);
},
};
}),
withHooks((store) => {
const sources = inject(EpgSourceSettingsService);
let subscription: { unsubscribe(): void } | undefined;
return {
onInit: () => {
subscription = sources.changed$.subscribe(() => {
store.clearEpg();
void store.loadEpg();
});
},
onDestroy: () => subscription?.unsubscribe(),
};
})
);
}
@@ -1,3 +1,4 @@
import { EpgSourceSettingsService } from '@iptvnator/services';
import { CdkFixedSizeVirtualScroll } from '@angular/cdk/scrolling';
import { signal } from '@angular/core';
import { ComponentFixture, TestBed } from '@angular/core/testing';
@@ -170,6 +171,15 @@ describe('PortalChannelsListComponent', () => {
});
});
it('clears visible cached XMLTV previews and re-enqueues after a source changes', () => {
const component = fixture.componentInstance;
component.epgPrograms.set(50, { title: 'Removed programme' } as never);
component.currentProgramsProgress.set(50, 20);
TestBed.inject(EpgSourceSettingsService).changed$.next();
expect(component.epgPrograms.size).toBe(0);
expect(component.currentProgramsProgress.size).toBe(0);
});
it('renders a loading placeholder instead of the empty state while xtream live content is still loading', () => {
fixture.detectChanges();
@@ -49,7 +49,11 @@ import { EpgQueueService } from '@iptvnator/portal/xtream/data-access';
import { XtreamCredentials } from '@iptvnator/portal/xtream/data-access';
import { FavoritesService } from '@iptvnator/portal/xtream/data-access';
import { XtreamStore } from '@iptvnator/portal/xtream/data-access';
import { RuntimeCapabilitiesService, SettingsStore } from '@iptvnator/services';
import {
EpgSourceSettingsService,
RuntimeCapabilitiesService,
SettingsStore,
} from '@iptvnator/services';
import { XtreamFavoriteMarksService } from './xtream-favorite-marks.service';
export interface XtreamChannelListItem {
@@ -164,6 +168,11 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy {
private readonly settingsStore = inject(SettingsStore);
constructor(private cdr: ChangeDetectorRef) {
this.subscriptions.add(
inject(EpgSourceSettingsService).changed$.subscribe(() => {
this.repickPreviewsForOffsetChange();
})
);
// A changed display offset moves "now" in the provider's clock, so
// the visible previews are re-picked from the cached EPG. The first
// run only records the initial value.
+2
View File
@@ -21,3 +21,5 @@ export * from './lib/tmdb';
export * from './lib/xtream-pending-restore.service';
export * from './lib/stream-probe.service';
export * from './lib/vod-source-pin.service';
export * from './lib/epg-source-settings.service';
@@ -0,0 +1,72 @@
import { Injector } from '@angular/core';
import { firstValueFrom, of, Subject } from 'rxjs';
import { EpgSourceSettingsService } from './epg-source-settings.service';
import { PlaylistsService } from './playlists.service';
describe('EPG source settings synchronization', () => {
const original = window.electron;
afterEach(() => {
window.electron = original;
});
it('waits for playlist migration, normalizes URLs and invalidates pending lookups after success', async () => {
const playlists = new Subject<never[]>();
const reconcileEpgSources = jest
.fn()
.mockResolvedValue({ success: true });
window.electron = {
reconcileEpgSources,
} as unknown as typeof window.electron;
const injector = Injector.create({
providers: [
EpgSourceSettingsService,
{
provide: PlaylistsService,
useValue: { getAllPlaylists: () => playlists },
},
],
});
const service = injector.get(EpgSourceSettingsService);
const pending = new Subject<string>();
const observer = jest.fn();
const lookup = pending.pipe(service.guard());
lookup.subscribe(observer);
const synchronization = service.synchronize([' a ', '', 'a']);
expect(reconcileEpgSources).not.toHaveBeenCalled();
playlists.next([]);
await synchronization;
pending.next('old programme');
// Subscribing to a pre-deletion request later also cannot revive it.
lookup.subscribe(observer);
pending.next('late programme');
expect(observer).not.toHaveBeenCalled();
expect(service.revision()).toBe(1);
expect(reconcileEpgSources).toHaveBeenCalledWith(['a']);
expect(await firstValueFrom(of('new').pipe(service.guard()))).toBe(
'new'
);
});
it('invalidates possibly partially deleted data and reports reconciliation failure', async () => {
window.electron = {
reconcileEpgSources: jest.fn().mockRejectedValue(new Error('disk')),
} as unknown as typeof window.electron;
const injector = Injector.create({
providers: [
EpgSourceSettingsService,
{
provide: PlaylistsService,
useValue: { getAllPlaylists: () => of([]) },
},
],
});
const service = injector.get(EpgSourceSettingsService);
await expect(service.synchronize(['current'])).rejects.toThrow(
'Failed to reconcile EPG sources'
);
expect(service.revision()).toBe(1);
expect(service.retainCurrentSources(['current', 'removed'], 0)).toEqual(
['current']
);
});
});
@@ -0,0 +1,79 @@
import { inject, Injectable, Injector, signal } from '@angular/core';
import {
firstValueFrom,
Subject,
Observable,
MonoTypeOperatorFunction,
filter,
takeUntil,
} from 'rxjs';
import { PlaylistsService } from './playlists.service';
export class EpgSourceReconciliationError extends Error {
constructor() {
super('Failed to reconcile EPG sources');
}
}
/** Synchronizes committed global XMLTV settings, never unsaved form edits. */
@Injectable({ providedIn: 'root' })
export class EpgSourceSettingsService {
private readonly injector = inject(Injector);
private activeUrls = new Set<string>();
readonly revision = signal(0);
readonly changed$ = new Subject<void>();
retainCurrentSources(urls: string[], requestedRevision: number): string[] {
return requestedRevision === this.revision()
? urls
: urls.filter((url) => this.activeUrls.has(url));
}
guard<T>(): MonoTypeOperatorFunction<T> {
const revision = this.revision();
return (source: Observable<T>) =>
source.pipe(
takeUntil(this.changed$),
filter(() => revision === this.revision())
);
}
async synchronize(urls: string[] | string | undefined): Promise<void> {
if (
typeof window === 'undefined' ||
!window.electron?.reconcileEpgSources
)
return;
const normalized = [
...new Set(
(Array.isArray(urls) ? urls : [urls ?? ''])
.map((url) => url.trim())
.filter(Boolean)
),
];
// These globals are committed even if playlist ownership or cleanup
// cannot be read. Never keep the previous global list on failure.
this.activeUrls = new Set(normalized);
try {
// This includes the legacy IndexedDB → SQLite playlist migration.
const playlists = await firstValueFrom(
this.injector.get(PlaylistsService).getAllPlaylists()
);
for (const playlist of playlists) {
if (playlist.serverUrl || playlist.macAddress) continue;
for (const url of playlist.epgUrls ?? []) {
if (url.trim()) this.activeUrls.add(url.trim());
}
}
const result =
await window.electron.reconcileEpgSources(normalized);
if (!result.success)
throw new Error('EPG source reconciliation failed');
} catch {
throw new EpgSourceReconciliationError();
} finally {
this.revision.update((revision) => revision + 1);
this.changed$.next();
}
}
}
@@ -1,3 +1,4 @@
import { EpgSourceSettingsService } from './epg-source-settings.service';
import { Injector } from '@angular/core';
import { StorageMap } from '@ngx-pwa/local-storage';
import { of, Subject } from 'rxjs';
@@ -46,6 +47,7 @@ describe('SettingsStore dashboard rail settings', () => {
injector = Injector.create({
providers: [
SettingsStore,
EpgSourceSettingsService,
{
provide: StorageMap,
useValue: storage,
@@ -54,6 +56,60 @@ describe('SettingsStore dashboard rail settings', () => {
});
});
it('reconciles only after persistence and restores the previous EPG list on a failed save', async () => {
storedSettings = { epgUrl: ['https://old.example/guide.xml'] };
const sources = injector.get(EpgSourceSettingsService);
const reconcile = jest
.spyOn(sources, 'synchronize')
.mockResolvedValue(undefined);
const store = injector.get(SettingsStore);
await store.loadSettings();
reconcile.mockClear();
const write = new Subject<void>();
storage.set.mockReturnValue(write);
const saving = store.updateSettings({ epgUrl: [] });
expect(reconcile).not.toHaveBeenCalled();
write.error(new Error('storage unavailable'));
await expect(saving).rejects.toThrow('storage unavailable');
expect(reconcile).not.toHaveBeenCalled();
expect(store.epgUrl()).toEqual(['https://old.example/guide.xml']);
storage.set.mockReturnValue(of(undefined));
await store.updateSettings({ epgUrl: [] });
expect(reconcile).toHaveBeenCalledWith([]);
});
it('never reconciles empty defaults after settings storage fails to load', async () => {
const pending = new Subject<unknown>();
storage.get.mockReturnValue(pending);
const reconcile = jest.spyOn(
injector.get(EpgSourceSettingsService),
'synchronize'
);
const store = injector.get(SettingsStore);
pending.error(new Error('cannot read settings'));
await store.loadSettings();
expect(reconcile).not.toHaveBeenCalled();
});
it('keeps persisted URLs authoritative when subsequent EPG cleanup fails', async () => {
const sources = injector.get(EpgSourceSettingsService);
const reconcile = jest
.spyOn(sources, 'synchronize')
.mockResolvedValue(undefined);
const store = injector.get(SettingsStore);
await store.loadSettings();
reconcile.mockRejectedValue(new Error('cleanup failed'));
await expect(
store.updateSettings({ epgUrl: ['new-source'] })
).rejects.toThrow('cleanup failed');
expect(store.epgUrl()).toEqual(['new-source']);
expect(store.storageFailure()).toBeNull();
expect(storage.set).toHaveBeenLastCalledWith(
STORE_KEY.Settings,
expect.objectContaining({ epgUrl: ['new-source'] })
);
});
it('defaults portal request pauses on and persists an explicit opt-out', async () => {
const store = injector.get(SettingsStore);
expect(store.getSettings().portalConnectivityGuard).toBe(true);
@@ -448,6 +504,7 @@ describe('SettingsStore storage failure reporting', () => {
injector = Injector.create({
providers: [
SettingsStore,
EpgSourceSettingsService,
{
provide: StorageMap,
useValue: storage,
@@ -1,3 +1,4 @@
import { EpgSourceSettingsService } from './epg-source-settings.service';
import { computed, inject } from '@angular/core';
import {
patchState,
@@ -143,6 +144,7 @@ export const SettingsStore = signalStore(
),
})),
withMethods((store, storage = inject(StorageMap)) => {
const epgSources = inject(EpgSourceSettingsService);
let settingsLoadPromise: Promise<void> | undefined;
return {
@@ -190,6 +192,14 @@ export const SettingsStore = signalStore(
}
);
}
await epgSources
.synchronize(this.getSettings().epgUrl)
.catch((error) => {
console.warn(
'Could not reconcile cached EPG sources on startup.',
error
);
});
})().catch((error) => {
settingsLoadPromise = undefined;
console.error('Failed to load settings:', error);
@@ -203,6 +213,7 @@ export const SettingsStore = signalStore(
},
async updateSettings(settings: Partial<Settings>) {
const previousEpgUrls = store.epgUrl();
patchState(store, {
...settings,
...(settings.webPlayerSharedControls !== undefined
@@ -252,9 +263,15 @@ export const SettingsStore = signalStore(
console.error('Failed to save settings:', error);
// The in-memory patch above already applied, so without
// this flag the change looks saved until the next restart.
patchState(store, { storageFailure: 'save' });
patchState(store, {
storageFailure: 'save',
epgUrl: previousEpgUrls,
});
throw error;
}
if (settings.epgUrl !== undefined) {
await epgSources.synchronize(completeSettings.epgUrl);
}
},
getSettings() {
@@ -340,8 +340,7 @@ export interface ElectronBridgeTrustOptions {
trustedInsecureTlsHosts?: string[];
}
export interface ElectronBridgePlaylistFetchOptions
extends ElectronBridgeTrustOptions {
export interface ElectronBridgePlaylistFetchOptions extends ElectronBridgeTrustOptions {
userAgent?: string;
}
@@ -375,8 +374,7 @@ export interface ElectronBridgeEpgLookupOptions {
* renderer will also present as "now". Absent, the main process uses its own
* clock.
*/
export interface ElectronBridgeCurrentProgramsOptions
extends ElectronBridgeEpgLookupOptions {
export interface ElectronBridgeCurrentProgramsOptions extends ElectronBridgeEpgLookupOptions {
nowMs?: number;
}
@@ -812,6 +810,7 @@ export interface ElectronBridgeApi {
options?: ElectronBridgeTrustOptions
) => Promise<ElectronBridgeEpgFetchResult>;
clearEpgData: () => Promise<ElectronBridgeResult>;
reconcileEpgSources: (urls: string[]) => Promise<ElectronBridgeResult>;
clearEpgDataForSource: (sourceUrl: string) => Promise<ElectronBridgeResult>;
checkEpgFreshness: (
urls: string[],
@@ -1287,7 +1286,9 @@ export interface ElectronBridgeApi {
recordingsGet?: (
recordingId: number
) => Promise<ElectronRecordingItem | null>;
recordingsStop?: (recordingId: number) => Promise<ElectronBridgeErrorResult>;
recordingsStop?: (
recordingId: number
) => Promise<ElectronBridgeErrorResult>;
recordingsRemove?: (
recordingId: number
) => Promise<ElectronBridgeErrorResult>;