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

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

* fix(epg): close source reconciliation review races

* fix(epg): serialize cleanup with replacement imports

* fix(epg): report retired worker exits as cancellations

* fix(epg): preserve source metadata through cache cleanup

* refactor(epg): separate worker runtime and import lifecycle

* fix(epg): skip cleanup for unchanged source settings

* fix(epg): cancel retired error rows and pending retries

* fix(epg): redact diagnostics and mirror committed settings after cleanup errors

* fix(epg): preserve metadata writer order independently of timestamps
This commit is contained in:
4gray authored and GitHub committed 2026-09-06 07:32:30 +02:00
1 parent baba0529ef
commit 9de480826c
55 files changed
+3538 -615

No files matched your search

@@ -3,11 +3,13 @@ import {
EpgImportProgress,
EpgRuntimeBridgeService,
} from './epg-runtime-bridge.service';
import { SettingsStore } from '@iptvnator/services';
import { SettingsStore, EpgSourceSettingsService } from '@iptvnator/services';
import { ELECTRON_BRIDGE_SECURITY_ERROR_CODES } from '@iptvnator/shared/interfaces';
import { EpgProgressService } from './epg-progress.service';
describe('EpgProgressService', () => {
let epgBridge: Partial<EpgRuntimeBridgeService>;
let sources: { waitForReconciliation: jest.Mock };
let settingsStore: {
getSettings: jest.Mock;
getTrustOptions: jest.Mock;
@@ -15,11 +17,14 @@ describe('EpgProgressService', () => {
};
beforeEach(() => {
sources = {
waitForReconciliation: jest.fn().mockResolvedValue(undefined),
};
epgBridge = {
forceFetchEpg: jest.fn().mockResolvedValue({ success: true }),
onProgress: jest.fn(),
supportsDataManagement: false,
supportsProgress: false,
supportsProgress: true,
};
settingsStore = {
getSettings: jest.fn(() => ({
@@ -35,6 +40,7 @@ describe('EpgProgressService', () => {
});
afterEach(() => {
jest.useRealTimers();
TestBed.resetTestingModule();
jest.restoreAllMocks();
});
@@ -43,6 +49,7 @@ describe('EpgProgressService', () => {
TestBed.configureTestingModule({
providers: [
EpgProgressService,
{ provide: EpgSourceSettingsService, useValue: sources },
{
provide: EpgRuntimeBridgeService,
useValue: epgBridge,
@@ -57,25 +64,132 @@ describe('EpgProgressService', () => {
return TestBed.inject(EpgProgressService);
}
const url = 'https://example.com/epg.xml';
const emit = (progress: EpgImportProgress) =>
(epgBridge.onProgress as jest.Mock).mock.calls[0][0](progress);
const errorRow = (source = url) =>
emit({
url: source,
status: 'error',
generation: 0,
errorCode:
ELECTRON_BRIDGE_SECURITY_ERROR_CODES.InvalidTlsCertificate,
});
it('removes retained actionable errors when their source is cancelled', async () => {
epgBridge.supportsDataManagement = true;
const service = configureService();
errorRow();
emit({ url, status: 'cancelled', generation: 0 });
expect(service.isVisible()).toBe(false);
await service.retry(url);
expect(epgBridge.forceFetchEpg).not.toHaveBeenCalled();
});
it.each([true, false])(
'waits for source reconciliation before retrying (removed=%s)',
async (removed) => {
epgBridge.supportsDataManagement = true;
let finish!: () => void;
sources.waitForReconciliation.mockReturnValue(
new Promise<void>((resolve) => {
finish = resolve;
})
);
const service = configureService();
errorRow();
const retry = service.retry(url);
expect(epgBridge.forceFetchEpg).not.toHaveBeenCalled();
if (removed) emit({ url, status: 'cancelled', generation: 0 });
finish();
await retry;
expect(epgBridge.forceFetchEpg).toHaveBeenCalledTimes(
removed ? 0 : 1
);
}
);
it.each(['private', 'tls'])(
'does not retry a removed/replaced row after a pending %s trust write',
async (kind) => {
epgBridge.supportsDataManagement = true;
let finish!: () => void;
settingsStore.updateSettings.mockReturnValue(
new Promise<void>((resolve) => {
finish = resolve;
})
);
const service = configureService();
errorRow();
const retry =
kind === 'private'
? service.trustPrivateNetworkSourceAndRetry(url)
: service.trustInsecureTlsHostAndRetry(url);
emit({ url, status: 'cancelled', generation: 0 });
emit({ url, status: 'loading', generation: 2 });
finish();
await retry;
expect(epgBridge.forceFetchEpg).not.toHaveBeenCalled();
expect(service.activeCount()).toBe(1);
}
);
it('does not let an old dismissal timer remove a replacement import', () => {
jest.useFakeTimers();
const service = configureService();
emit({ url, status: 'complete', generation: 0 });
emit({ url, status: 'cancelled', generation: 0 });
emit({ url, status: 'loading', generation: 2 });
jest.advanceTimersByTime(5000);
expect(service.activeCount()).toBe(1);
jest.useRealTimers();
});
it('does not subscribe to progress events when EPG progress support is disabled', () => {
epgBridge.supportsProgress = false;
configureService();
expect(epgBridge.onProgress).not.toHaveBeenCalled();
});
it('does not force retry when EPG data management is disabled', () => {
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', async () => {
const service = configureService();
service.retry('https://example.com/epg.xml');
errorRow();
await service.retry('https://example.com/epg.xml');
expect(epgBridge.forceFetchEpg).not.toHaveBeenCalled();
});
it('forces retry through the EPG runtime bridge when data management is enabled', () => {
it('forces retry through the EPG runtime bridge when data management is enabled', async () => {
epgBridge.supportsDataManagement = true;
const service = configureService();
service.retry('https://example.com/epg.xml');
errorRow();
await service.retry('https://example.com/epg.xml');
expect(epgBridge.forceFetchEpg).toHaveBeenCalledWith(
'https://example.com/epg.xml',
@@ -90,6 +204,7 @@ describe('EpgProgressService', () => {
epgBridge.supportsDataManagement = true;
const service = configureService();
errorRow('http://192.168.1.30/guide.xml');
await service.trustPrivateNetworkSourceAndRetry(
'http://192.168.1.30/guide.xml'
);
@@ -3,7 +3,7 @@ import {
ELECTRON_BRIDGE_SECURITY_ERROR_CODES,
normalizeHost,
} from '@iptvnator/shared/interfaces';
import { SettingsStore } from '@iptvnator/services';
import { SettingsStore, EpgSourceSettingsService } from '@iptvnator/services';
import {
EpgImportProgress,
EpgRuntimeBridgeService,
@@ -13,6 +13,7 @@ import {
export class EpgProgressService {
private readonly epgBridge = inject(EpgRuntimeBridgeService);
private readonly settingsStore = inject(SettingsStore);
private readonly sources = inject(EpgSourceSettingsService);
private readonly importsMap = signal<Map<string, EpgImportProgress>>(
new Map()
);
@@ -45,20 +46,27 @@ export class EpgProgressService {
this.importsMap.set(new Map());
}
retry(url: string): void {
// Clear the errored row so the backend's subsequent 'queued' event
// reappears cleanly rather than updating an existing error row.
this.removeImport(url);
if (!this.epgBridge.supportsDataManagement) {
async retry(url: string): Promise<void> {
const progress = this.importsMap().get(url);
if (
progress?.status !== 'error' ||
!this.epgBridge.supportsDataManagement
)
return;
}
void this.epgBridge.forceFetchEpg(
// A settings save may already be retiring this row. Never bypass its
// ownership reconciliation by starting a new force-fetch generation.
await this.sources.waitForReconciliation();
if (this.importsMap().get(url) !== progress) return;
this.removeImport(url);
await this.epgBridge.forceFetchEpg(
url,
this.settingsStore.getTrustOptions()
);
}
async trustPrivateNetworkSourceAndRetry(url: string): Promise<void> {
const progress = this.importsMap().get(url);
if (progress?.status !== 'error') return;
const settings = this.settingsStore.getSettings();
const trustedUrls = new Set(
settings.trustedPrivateNetworkEpgUrls ?? []
@@ -68,13 +76,15 @@ export class EpgProgressService {
await this.settingsStore.updateSettings({
trustedPrivateNetworkEpgUrls: Array.from(trustedUrls),
});
this.retry(url);
if (this.importsMap().get(url) === progress) await this.retry(url);
}
async trustInsecureTlsHostAndRetry(
url: string,
host?: string
): Promise<void> {
const progress = this.importsMap().get(url);
if (progress?.status !== 'error') return;
const trustedHost = host ?? this.getHostname(url);
if (!trustedHost) {
return;
@@ -91,7 +101,7 @@ export class EpgProgressService {
await this.settingsStore.updateSettings({
trustedInsecureTlsHosts: Array.from(trustedHosts),
});
this.retry(url);
if (this.importsMap().get(url) === progress) await this.retry(url);
}
private initializeListener(): void {
@@ -108,6 +118,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);
@@ -118,7 +134,10 @@ export class EpgProgressService {
progress.status === 'complete' ||
(progress.status === 'error' && !this.isActionableError(progress))
) {
setTimeout(() => this.removeImport(progress.url), 5000);
setTimeout(() => {
if (this.importsMap().get(progress.url) === progress)
this.removeImport(progress.url);
}, 5000);
}
}
@@ -1,8 +1,12 @@
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 { firstValueFrom, of, skip } from 'rxjs';
import {
EpgSourceSettingsService,
PlaylistsService,
SettingsStore,
} from '@iptvnator/services';
import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service';
import { EpgService } from './epg.service';
@@ -11,6 +15,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 +34,7 @@ describe('EpgService', () => {
open: jest.fn(),
};
settingsStore = {
loadSettings: jest.fn().mockResolvedValue(undefined),
getSettings: jest.fn(() => ({
epgUrl: [],
trustedPrivateNetworkEpgUrls: ['http://192.168.1.20/guide.xml'],
@@ -44,6 +50,19 @@ describe('EpgService', () => {
TestBed.configureTestingModule({
providers: [
EpgService,
{
provide: PlaylistsService,
useValue: {
getAllPlaylists: () =>
of([
{
epgUrls: [
'https://playlist.example/guide.xml',
],
},
]),
},
},
{
provide: EpgRuntimeBridgeService,
useValue: epgBridge,
@@ -68,22 +87,126 @@ 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('waits for ongoing reconciliation and filters imports against committed owners', async () => {
epgBridge.supportsImport = true;
const original = window.electron;
let complete!: (result: { success: boolean }) => void;
window.electron = {
reconcileEpgSources: () =>
new Promise((resolve) => {
complete = resolve;
}),
} as typeof window.electron;
try {
const sources = TestBed.inject(EpgSourceSettingsService);
const reconciliation = sources.synchronize([
'https://kept.example/guide.xml',
]);
await Promise.resolve();
const pending = service.fetchEpg([
'https://removed.example/guide.xml',
'https://playlist.example/guide.xml',
]);
await Promise.resolve();
await Promise.resolve();
expect(epgBridge.fetchEpg).not.toHaveBeenCalled();
complete({ success: true });
await reconciliation;
await pending;
expect(epgBridge.fetchEpg).toHaveBeenCalledWith(
['https://playlist.example/guide.xml'],
expect.anything()
);
} finally {
window.electron = original;
}
});
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([
await service.fetchEpg([
'https://example.com/epg.xml',
'',
'https://example.com/other.xml',
' https://example.com/epg.xml ',
]);
await Promise.resolve();
expect(epgBridge.fetchEpg).toHaveBeenCalledWith(
['https://example.com/epg.xml', 'https://example.com/other.xml'],
{
+43 -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,17 @@ 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();
await this.sourceSettings.waitForReconciliation();
const validUrls = this.sourceSettings.retainCurrentSources(
normalizeEpgUrls(urls),
revision
);
if (validUrls.length === 0) return;
from(
@@ -93,6 +108,7 @@ export class EpgService {
)
)
.pipe(
this.sourceSettings.guard(),
tap((result) => {
if (result === null) return;
@@ -123,6 +139,7 @@ export class EpgService {
from(this.epgBridge.getChannelPrograms(channelId))
.pipe(
this.sourceSettings.guard(),
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
catchError((err) => {
@@ -187,6 +204,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 +305,7 @@ export class EpgService {
nowMs: this.epgClockMs(),
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((batchResult) => {
const cacheTimestamp = Date.now();
@@ -311,6 +330,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 +338,7 @@ export class EpgService {
);
return forkJoin(fetchObservables).pipe(
this.sourceSettings.guard(),
map((results) => {
results.forEach((result) => {
resultMap.set(result.channelId, result.program);
@@ -372,6 +393,7 @@ export class EpgService {
sourceUrls,
fallbackSourceUrls
).pipe(
this.sourceSettings.guard(),
tap((fetchedMap) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
@@ -386,7 +408,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 +425,7 @@ export class EpgService {
}
return request$.pipe(
this.sourceSettings.guard(),
map((fetchedMap) => {
const mergedResultMap = new Map(resultMap);
channelsToFetch.forEach((channelId) => {
@@ -421,6 +451,7 @@ export class EpgService {
nowMs,
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
switchMap((scopedResult) => {
const resultMap = new Map<string, EpgProgram | null>();
@@ -448,6 +479,7 @@ export class EpgService {
nowMs,
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((globalResult) => {
fallbackChannelIds.forEach((channelId) => {
@@ -500,6 +532,7 @@ export class EpgService {
normalizedChannelIds,
effectiveSourceUrls
).pipe(
this.sourceSettings.guard(),
switchMap((metadataMap) => {
const fallbackChannelIds =
sourceUrls.length > 0 && globalSourceUrls.length > 0
@@ -516,6 +549,7 @@ export class EpgService {
fallbackChannelIds,
globalSourceUrls
).pipe(
this.sourceSettings.guard(),
map((globalMetadataMap) => {
fallbackChannelIds.forEach((channelId) => {
metadataMap.set(
@@ -561,6 +595,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 +676,7 @@ export class EpgService {
const offsetMinutes = this.epgOffsetMinutes();
const request$ = fetchProgram().pipe(
this.sourceSettings.guard(),
tap((program) => {
this.programCache.set(cacheKey, {
program,
@@ -649,7 +685,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 +705,7 @@ export class EpgService {
from(
this.epgBridge.getChannelPrograms(channelId, { sourceUrls })
).pipe(
this.sourceSettings.guard(),
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
switchMap((programs) => {
@@ -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(),
};
})
);
}
@@ -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,112 @@
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();
pending.next('response during playlist migration');
expect(observer).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(2);
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(2);
expect(service.retainCurrentSources(['current', 'removed'], 0)).toEqual(
['current']
);
});
it('serializes overlapping saves and waits for the latest committed source set', async () => {
let finishFirst!: (result: { success: boolean }) => void;
const reconcileEpgSources = jest
.fn()
.mockImplementationOnce(
() =>
new Promise((resolve) => {
finishFirst = resolve;
})
)
.mockResolvedValue({ success: true });
window.electron = {
reconcileEpgSources,
} as unknown as typeof window.electron;
const injector = Injector.create({
providers: [
EpgSourceSettingsService,
{
provide: PlaylistsService,
useValue: { getAllPlaylists: () => of([]) },
},
],
});
const service = injector.get(EpgSourceSettingsService);
const first = service.synchronize(['a']);
const second = service.synchronize(['b']);
const waiter = service.waitForReconciliation();
await Promise.resolve();
expect(reconcileEpgSources).toHaveBeenCalledTimes(1);
finishFirst({ success: true });
await Promise.all([first, second, waiter]);
expect(reconcileEpgSources.mock.calls.map(([urls]) => urls)).toEqual([
['a'],
['b'],
]);
expect(service.retainCurrentSources(['a', 'b'], 0)).toEqual(['b']);
});
});
@@ -0,0 +1,118 @@
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');
}
}
function normalizeEpgSourceUrls(urls: string[] | string | undefined): string[] {
return [
...new Set(
(Array.isArray(urls) ? urls : [urls ?? ''])
.map((url) => url.trim())
.filter(Boolean)
),
];
}
export function epgSourceUrlsChanged(
previous: string[] | string | undefined,
next: string[] | string | undefined
): boolean {
return (
next !== undefined &&
JSON.stringify(normalizeEpgSourceUrls(previous).sort()) !==
JSON.stringify(normalizeEpgSourceUrls(next).sort())
);
}
/** 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>();
private reconciliation: Promise<void> | undefined;
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 waitForReconciliation(): Promise<void> {
while (this.reconciliation) {
await this.reconciliation.catch(() => undefined);
}
}
async synchronize(urls: string[] | string | undefined): Promise<void> {
if (
typeof window === 'undefined' ||
!window.electron?.reconcileEpgSources
)
return;
// Fence existing lookups before playlist migration or IPC can yield.
this.revision.update((revision) => revision + 1);
const previous = this.reconciliation;
const operation = previous
? previous.catch(() => undefined).then(() => this.reconcile(urls))
: this.reconcile(urls);
const pending = operation.finally(() => {
if (this.reconciliation === pending)
this.reconciliation = undefined;
});
this.reconciliation = pending;
return pending;
}
private async reconcile(
urls: string[] | string | undefined
): Promise<void> {
const normalized = normalizeEpgSourceUrls(urls);
// 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,118 @@ 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.each([
{ epgUrl: ['second', 'first'] },
{ epgUrl: [' first ', 'second', 'first', ''] },
])(
'persists unrelated settings without reconciling an unchanged normalized source set: %j',
async ({ epgUrl }) => {
storedSettings = { epgUrl: ['first', 'second'] };
const reconcile = jest
.spyOn(injector.get(EpgSourceSettingsService), 'synchronize')
.mockResolvedValue(undefined);
const store = injector.get(SettingsStore);
await store.loadSettings();
reconcile
.mockClear()
.mockRejectedValue(new Error('migration unavailable'));
await expect(
store.updateSettings({
...store.getSettings(),
language: Language.FRENCH,
epgUrl,
})
).resolves.toBeUndefined();
expect(reconcile).not.toHaveBeenCalled();
expect(store.storageFailure()).toBeNull();
expect(storage.set).toHaveBeenLastCalledWith(
STORE_KEY.Settings,
expect.objectContaining({ language: Language.FRENCH })
);
}
);
it('keeps failed cleanup retryable only when an EPG save explicitly requests it', async () => {
storedSettings = { epgUrl: ['removed'] };
const reconcile = jest
.spyOn(injector.get(EpgSourceSettingsService), 'synchronize')
.mockResolvedValue(undefined);
const store = injector.get(SettingsStore);
await store.loadSettings();
reconcile.mockClear().mockRejectedValue(new Error('cleanup failed'));
await expect(store.updateSettings({ epgUrl: [] })).rejects.toThrow(
'cleanup failed'
);
reconcile.mockClear();
await expect(
store.updateSettings({
...store.getSettings(),
language: Language.FRENCH,
})
).resolves.toBeUndefined();
expect(reconcile).not.toHaveBeenCalled();
await expect(
store.updateSettings({ epgUrl: [] }, { retryEpgCleanup: true })
).rejects.toThrow('cleanup failed');
reconcile.mockResolvedValue(undefined);
await store.updateSettings({ epgUrl: [] }, { retryEpgCleanup: true });
expect(reconcile).toHaveBeenCalledTimes(2);
});
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 +562,7 @@ describe('SettingsStore storage failure reporting', () => {
injector = Injector.create({
providers: [
SettingsStore,
EpgSourceSettingsService,
{
provide: StorageMap,
useValue: storage,
@@ -1,3 +1,7 @@
import {
EpgSourceSettingsService,
epgSourceUrlsChanged,
} from './epg-source-settings.service';
import { computed, inject } from '@angular/core';
import {
patchState,
@@ -143,6 +147,7 @@ export const SettingsStore = signalStore(
),
})),
withMethods((store, storage = inject(StorageMap)) => {
const epgSources = inject(EpgSourceSettingsService);
let settingsLoadPromise: Promise<void> | undefined;
return {
@@ -190,6 +195,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);
@@ -202,7 +215,11 @@ export const SettingsStore = signalStore(
return settingsLoadPromise;
},
async updateSettings(settings: Partial<Settings>) {
async updateSettings(
settings: Partial<Settings>,
options: { retryEpgCleanup?: boolean } = {}
) {
const previousEpgUrls = store.epgUrl();
patchState(store, {
...settings,
...(settings.webPlayerSharedControls !== undefined
@@ -252,9 +269,18 @@ 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 (
epgSourceUrlsChanged(previousEpgUrls, settings.epgUrl) ||
options.retryEpgCleanup
) {
await epgSources.synchronize(completeSettings.epgUrl);
}
},
getSettings() {
@@ -287,6 +287,18 @@ const CREATE_TABLE_STATEMENTS = [
source_url TEXT NOT NULL,
updated_at TEXT DEFAULT (datetime('now'))
)`,
`CREATE TABLE IF NOT EXISTS epg_channel_sources (
channel_id TEXT NOT NULL,
source_url TEXT NOT NULL,
display_name TEXT NOT NULL,
icon_url TEXT,
url TEXT,
updated_at TEXT DEFAULT (datetime('now')),
write_order INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (channel_id, source_url),
FOREIGN KEY (channel_id) REFERENCES epg_channels(id) ON DELETE CASCADE
)`,
`CREATE INDEX IF NOT EXISTS idx_epg_channel_sources_source ON epg_channel_sources(source_url)`,
`CREATE TABLE IF NOT EXISTS epg_programs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
channel_id TEXT NOT NULL,
@@ -419,6 +431,8 @@ const COLUMN_MIGRATION_STATEMENTS = [
`ALTER TABLE content ADD COLUMN original_title TEXT`,
// v1.7.1: Scope XMLTV programs to their source URL for playlist-local EPG lookup
`ALTER TABLE epg_programs ADD COLUMN source_url TEXT`,
// Preserve writer order independently of wall-clock precision or changes.
`ALTER TABLE epg_channel_sources ADD COLUMN write_order INTEGER NOT NULL DEFAULT 0`,
// Pause/resume: entity validator (ETag/Last-Modified) sent as If-Range on resume
`ALTER TABLE downloads ADD COLUMN resume_validator TEXT`,
// Offline details: provider-neutral display metadata captured at download time
+21
View File
@@ -11,6 +11,7 @@ import { sql } from 'drizzle-orm';
import {
index,
integer,
primaryKey,
sqliteTable,
text,
uniqueIndex,
@@ -204,6 +205,26 @@ export const epgChannels = sqliteTable(
})
);
// A global XMLTV ID can be shared by sources with different channel metadata.
export const epgChannelSources = sqliteTable(
'epg_channel_sources',
{
channelId: text('channel_id')
.notNull()
.references(() => epgChannels.id, { onDelete: 'cascade' }),
sourceUrl: text('source_url').notNull(),
displayName: text('display_name').notNull(),
iconUrl: text('icon_url'),
url: text('url'),
updatedAt: text('updated_at').default(sql`CURRENT_TIMESTAMP`),
writeOrder: integer('write_order').notNull().default(0),
},
(table) => [
primaryKey({ columns: [table.channelId, table.sourceUrl] }),
index('idx_epg_channel_sources_source').on(table.sourceUrl),
]
);
// EPG Programs table
export const epgPrograms = sqliteTable(
'epg_programs',
@@ -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',
@@ -340,8 +341,7 @@ export interface ElectronBridgeTrustOptions {
trustedInsecureTlsHosts?: string[];
}
export interface ElectronBridgePlaylistFetchOptions
extends ElectronBridgeTrustOptions {
export interface ElectronBridgePlaylistFetchOptions extends ElectronBridgeTrustOptions {
userAgent?: string;
}
@@ -375,8 +375,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;
}
@@ -387,6 +386,7 @@ export interface ElectronBridgeEpgProgressStats {
export interface ElectronBridgeEpgProgress {
url: string;
generation?: number;
status: ElectronBridgeEpgProgressStatus;
stats?: ElectronBridgeEpgProgressStats;
error?: string;
@@ -812,6 +812,7 @@ export interface ElectronBridgeApi {
options?: ElectronBridgeTrustOptions
) => Promise<ElectronBridgeEpgFetchResult>;
clearEpgData: () => Promise<ElectronBridgeResult>;
reconcileEpgSources: (urls: string[]) => Promise<ElectronBridgeResult>;
clearEpgDataForSource: (sourceUrl: string) => Promise<ElectronBridgeResult>;
checkEpgFreshness: (
urls: string[],
@@ -1291,7 +1292,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>;
@@ -96,6 +96,8 @@ export class EpgProgressPanelComponent {
return 'check_circle';
case 'error':
return 'error';
case 'cancelled':
return 'cancel';
}
}