feat(epg): support playlist-scoped sources

Add playlist-scoped EPG source support for M3U playlists.
This commit is contained in:
4gray authored and GitHub committed 2026-06-21 22:53:07 +02:00
1 parent 36a7ce9f3d
commit e801028005
71 files changed
+4876 -373

No files matched your search

@@ -66,11 +66,15 @@ describe('EpgRuntimeBridgeService', () => {
const fetchEpg = jest.fn().mockResolvedValue({ success: true });
const forceFetchEpg = jest.fn().mockResolvedValue({ success: true });
const clearEpgData = jest.fn().mockResolvedValue({ success: true });
const clearEpgDataForSource = jest
.fn()
.mockResolvedValue({ success: true });
window.electron = {
...window.electron,
fetchEpg,
forceFetchEpg,
clearEpgData,
clearEpgDataForSource,
} as unknown as typeof window.electron;
runtimeCapabilities.supportsEpgImport = true;
runtimeCapabilities.supportsEpgDataManagement = true;
@@ -84,6 +88,13 @@ describe('EpgRuntimeBridgeService', () => {
await expect(service.clearEpgData()).resolves.toEqual({
success: true,
});
await expect(
service.clearEpgDataForSource(
' https://playlist.example.com/guide.xml '
)
).resolves.toEqual({
success: true,
});
expect(fetchEpg).toHaveBeenCalledWith(
['https://example.com/epg.xml'],
@@ -94,6 +105,9 @@ describe('EpgRuntimeBridgeService', () => {
undefined
);
expect(clearEpgData).toHaveBeenCalledTimes(1);
expect(clearEpgDataForSource).toHaveBeenCalledWith(
'https://playlist.example.com/guide.xml'
);
});
it('delegates read-side EPG calls through the typed Electron bridge', async () => {
@@ -127,15 +141,23 @@ describe('EpgRuntimeBridgeService', () => {
runtimeCapabilities.supportsEpgProgramSearch = true;
await service.getChannelPrograms('channel-1');
await service.getCurrentProgramsBatch(['channel-1']);
await service.getChannelMetadata(['channel-1']);
await service.getCurrentProgramsBatch(['channel-1'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
});
await service.getChannelMetadata(['channel-1'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
});
await service.checkFreshness(['https://example.com/epg.xml'], 12);
await service.getChannelsByRange(0, 20);
await service.searchPrograms('news', 20);
expect(getChannelPrograms).toHaveBeenCalledWith('channel-1');
expect(getCurrentProgramsBatch).toHaveBeenCalledWith(['channel-1']);
expect(getEpgChannelMetadata).toHaveBeenCalledWith(['channel-1']);
expect(getCurrentProgramsBatch).toHaveBeenCalledWith(['channel-1'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
});
expect(getEpgChannelMetadata).toHaveBeenCalledWith(['channel-1'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
});
expect(checkEpgFreshness).toHaveBeenCalledWith(
['https://example.com/epg.xml'],
12
@@ -7,6 +7,7 @@ import {
ElectronBridgeEpgProgressStatus,
ElectronBridgeEpgChannelWithPrograms,
ElectronBridgeEpgFreshnessResult,
ElectronBridgeEpgLookupOptions,
ELECTRON_BRIDGE_EPG_PROGRESS_STATUSES,
ElectronBridgeResult,
ElectronBridgeTrustOptions,
@@ -22,11 +23,13 @@ export type EpgImportProgress = ElectronBridgeEpgProgress;
export type EpgFetchResult = ElectronBridgeEpgFetchResult;
export type EpgFreshnessResult = ElectronBridgeEpgFreshnessResult;
export type EpgClearResult = ElectronBridgeResult;
export type EpgLookupOptions = ElectronBridgeEpgLookupOptions;
type EpgElectronBridge = Pick<
Partial<ElectronBridgeApi>,
| 'checkEpgFreshness'
| 'clearEpgData'
| 'clearEpgDataForSource'
| 'fetchEpg'
| 'forceFetchEpg'
| 'getChannelPrograms'
@@ -109,41 +112,71 @@ export class EpgRuntimeBridgeService {
return this.bridge?.clearEpgData?.() ?? Promise.resolve(null);
}
getChannelPrograms(channelId: string): Promise<EpgProgram[] | null> {
if (!this.supportsProgramLookup) {
clearEpgDataForSource(sourceUrl: string): Promise<EpgClearResult | null> {
if (!this.supportsDataManagement) {
return Promise.resolve(null);
}
const normalizedSourceUrl = sourceUrl.trim();
if (!normalizedSourceUrl) {
return Promise.resolve(null);
}
return (
this.bridge?.getChannelPrograms?.(channelId) ??
this.bridge?.clearEpgDataForSource?.(normalizedSourceUrl) ??
Promise.resolve(null)
);
}
getChannelPrograms(
channelId: string,
options?: EpgLookupOptions
): Promise<EpgProgram[] | null> {
if (!this.supportsProgramLookup) {
return Promise.resolve(null);
}
if (!this.bridge?.getChannelPrograms) {
return Promise.resolve(null);
}
return options
? this.bridge.getChannelPrograms(channelId, options)
: this.bridge.getChannelPrograms(channelId);
}
getCurrentProgramsBatch(
channelIds: string[]
channelIds: string[],
options?: EpgLookupOptions
): Promise<Record<string, EpgProgram | null> | null> {
if (!this.supportsCurrentProgramBatch) {
return Promise.resolve(null);
}
return (
this.bridge?.getCurrentProgramsBatch?.(channelIds) ??
Promise.resolve(null)
);
if (!this.bridge?.getCurrentProgramsBatch) {
return Promise.resolve(null);
}
return options
? this.bridge.getCurrentProgramsBatch(channelIds, options)
: this.bridge.getCurrentProgramsBatch(channelIds);
}
getChannelMetadata(
channelIds: string[]
channelIds: string[],
options?: EpgLookupOptions
): Promise<Record<string, EpgChannelMetadata | null> | null> {
if (!this.supportsChannelMetadata) {
return Promise.resolve(null);
}
return (
this.bridge?.getEpgChannelMetadata?.(channelIds) ??
Promise.resolve(null)
);
if (!this.bridge?.getEpgChannelMetadata) {
return Promise.resolve(null);
}
return options
? this.bridge.getEpgChannelMetadata(channelIds, options)
: this.bridge.getEpgChannelMetadata(channelIds);
}
checkFreshness(
@@ -1,7 +1,7 @@
import { TestBed } from '@angular/core/testing';
import { MatSnackBar } from '@angular/material/snack-bar';
import { TranslateService } from '@ngx-translate/core';
import { firstValueFrom } from 'rxjs';
import { firstValueFrom, skip } from 'rxjs';
import { SettingsStore } from '@iptvnator/services';
import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service';
import { EpgService } from './epg.service';
@@ -26,6 +26,7 @@ describe('EpgService', () => {
};
settingsStore = {
getSettings: jest.fn(() => ({
epgUrl: [],
trustedPrivateNetworkEpgUrls: ['http://192.168.1.20/guide.xml'],
trustedInsecureTlsHosts: ['playlist.local'],
})),
@@ -75,6 +76,7 @@ describe('EpgService', () => {
'https://example.com/epg.xml',
'',
'https://example.com/other.xml',
' https://example.com/epg.xml ',
]);
expect(epgBridge.fetchEpg).toHaveBeenCalledWith(
@@ -132,4 +134,484 @@ describe('EpgService', () => {
expect(epgBridge.getChannelPrograms).toHaveBeenCalledWith('channel-1');
jest.useRealTimers();
});
it('caches current program lookups separately by EPG source URL scope', async () => {
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.getChannelPrograms = jest
.fn()
.mockResolvedValueOnce([
{
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Guide Bulletin',
},
])
.mockResolvedValueOnce([
{
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Global News Bulletin',
},
])
.mockResolvedValue([]);
jest.useFakeTimers();
jest.setSystemTime(new Date('2026-05-23T10:30:00.000Z'));
try {
const playlistResult = await firstValueFrom(
service.getCurrentProgramForChannel('guide-news', {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
const globalResult = await firstValueFrom(
service.getCurrentProgramForChannel('guide-news')
);
const cachedPlaylistResult = await firstValueFrom(
service.getCurrentProgramForChannel('guide-news', {
sourceUrls: [' https://playlist.example.com/guide.xml '],
})
);
const cachedGlobalResult = await firstValueFrom(
service.getCurrentProgramForChannel('guide-news')
);
expect(playlistResult?.title).toBe('Playlist Guide Bulletin');
expect(globalResult?.title).toBe('Global News Bulletin');
expect(cachedPlaylistResult?.title).toBe('Playlist Guide Bulletin');
expect(cachedGlobalResult?.title).toBe('Global News Bulletin');
expect(epgBridge.getChannelPrograms).toHaveBeenCalledTimes(2);
expect(epgBridge.getChannelPrograms).toHaveBeenNthCalledWith(
1,
'guide-news',
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
);
expect(epgBridge.getChannelPrograms).toHaveBeenNthCalledWith(
2,
'guide-news',
{ sourceUrls: ['https://global.example.com/guide.xml'] }
);
} finally {
jest.useRealTimers();
}
});
it('deduplicates concurrent scoped current program lookups for the same source scope', async () => {
epgBridge.supportsProgramLookup = true;
let resolvePrograms:
| ((
programs: {
channel: string;
start: string;
stop: string;
title: string;
}[]
) => void)
| undefined;
epgBridge.getChannelPrograms = jest.fn(
() =>
new Promise((resolve) => {
resolvePrograms = resolve;
})
);
jest.useFakeTimers();
jest.setSystemTime(new Date('2026-05-23T10:30:00.000Z'));
try {
const firstLookup = firstValueFrom(
service.getCurrentProgramForChannel('guide-news', {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
const secondLookup = firstValueFrom(
service.getCurrentProgramForChannel('guide-news', {
sourceUrls: [' https://playlist.example.com/guide.xml '],
})
);
expect(epgBridge.getChannelPrograms).toHaveBeenCalledTimes(1);
resolvePrograms?.([
{
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Guide Bulletin',
},
]);
await expect(
Promise.all([firstLookup, secondLookup])
).resolves.toEqual([
expect.objectContaining({
title: 'Playlist Guide Bulletin',
}),
expect.objectContaining({
title: 'Playlist Guide Bulletin',
}),
]);
} finally {
jest.useRealTimers();
}
});
it('queries playlist-scoped current programs first and falls back to global EPG for missing channels', async () => {
settingsStore.getSettings.mockReturnValue({
epgUrl: [
'https://global.example.com/guide.xml',
' https://global.example.com/guide.xml ',
],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
epgBridge.getCurrentProgramsBatch = jest
.fn()
.mockResolvedValueOnce({
'guide-news': {
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Guide Bulletin',
},
'guide-sports': null,
})
.mockResolvedValueOnce({
'guide-sports': {
channel: 'guide-sports',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Global Sports Bulletin',
},
});
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(
['guide-news', 'guide-sports'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
)
);
expect(result.get('guide-news')?.title).toBe('Playlist Guide Bulletin');
expect(result.get('guide-sports')?.title).toBe(
'Global Sports Bulletin'
);
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenNthCalledWith(
1,
['guide-news', 'guide-sports'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
);
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenNthCalledWith(
2,
['guide-sports'],
{ sourceUrls: ['https://global.example.com/guide.xml'] }
);
});
it('does not fall back to the unscoped EPG pool when no global EPG URLs are configured', async () => {
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
epgBridge.getCurrentProgramsBatch = jest.fn().mockResolvedValue({
'guide-sports': null,
});
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-sports'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
expect(result.get('guide-sports')).toBeNull();
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(1);
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledWith(
['guide-sports'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
);
});
it('caches scoped batch current programs by EPG source URL scope', async () => {
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
epgBridge.getCurrentProgramsBatch = jest
.fn()
.mockResolvedValue({})
.mockResolvedValueOnce({
'guide-news': {
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Guide Bulletin',
},
});
const firstResult = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
const secondResult = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
sourceUrls: [' https://playlist.example.com/guide.xml '],
})
);
expect(firstResult.get('guide-news')?.title).toBe(
'Playlist Guide Bulletin'
);
expect(secondResult.get('guide-news')?.title).toBe(
'Playlist Guide Bulletin'
);
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(1);
});
it('clears cached null current programs after a successful EPG import', async () => {
epgBridge.supportsImport = true;
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
epgBridge.getCurrentProgramsBatch = jest
.fn()
.mockResolvedValueOnce({
'guide-news': null,
})
.mockResolvedValueOnce({
'guide-news': {
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Guide Bulletin',
},
});
const beforeImport = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
expect(beforeImport.get('guide-news')).toBeNull();
const availability = firstValueFrom(service.epgAvailable$.pipe(skip(1)));
service.fetchEpg(['https://playlist.example.com/guide.xml']);
await expect(availability).resolves.toBe(true);
const afterImport = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
expect(afterImport.get('guide-news')?.title).toBe(
'Playlist Guide Bulletin'
);
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(2);
});
it('deduplicates concurrent scoped batch current program lookups for the same source scope', async () => {
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
const batchResolvers: Array<
(programs: Record<string, unknown>) => void
> = [];
epgBridge.getCurrentProgramsBatch = jest.fn(
() =>
new Promise((resolve) => {
batchResolvers.push(resolve);
})
);
const firstLookup = firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
const secondLookup = firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
sourceUrls: [' https://playlist.example.com/guide.xml '],
})
);
try {
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(1);
batchResolvers.forEach((resolve) =>
resolve({
'guide-news': {
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Guide Bulletin',
},
})
);
const [firstResult, secondResult] = await Promise.all([
firstLookup,
secondLookup,
]);
expect(firstResult.get('guide-news')?.title).toBe(
'Playlist Guide Bulletin'
);
expect(secondResult.get('guide-news')?.title).toBe(
'Playlist Guide Bulletin'
);
} finally {
batchResolvers.forEach((resolve) => resolve({}));
}
});
it('deduplicates concurrent scoped batch current program lookups regardless of channel order', async () => {
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
let resolveBatch:
| ((programs: Record<string, unknown>) => void)
| undefined;
epgBridge.getCurrentProgramsBatch = jest.fn(
() =>
new Promise((resolve) => {
resolveBatch = resolve;
})
);
const firstLookup = firstValueFrom(
service.getCurrentProgramsForChannels(
['guide-news', 'guide-sports'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
)
);
const secondLookup = firstValueFrom(
service.getCurrentProgramsForChannels(
['guide-sports', 'guide-news'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
)
);
try {
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(1);
resolveBatch?.({
'guide-news': {
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Guide Bulletin',
},
'guide-sports': {
channel: 'guide-sports',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Playlist Sports Bulletin',
},
});
const [firstResult, secondResult] = await Promise.all([
firstLookup,
secondLookup,
]);
expect(firstResult.get('guide-news')?.title).toBe(
'Playlist Guide Bulletin'
);
expect(secondResult.get('guide-sports')?.title).toBe(
'Playlist Sports Bulletin'
);
} finally {
resolveBatch?.({});
}
});
it('returns null scoped current programs instead of re-entering global lookup when scoped batch lookup fails', async () => {
const consoleError = jest
.spyOn(console, 'error')
.mockImplementation(() => undefined);
try {
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
epgBridge.getCurrentProgramsBatch = jest
.fn()
.mockRejectedValueOnce(new Error('ipc down'))
.mockResolvedValueOnce({
'guide-news': {
channel: 'guide-news',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Global News Bulletin',
},
});
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
expect(result.get('guide-news')).toBeNull();
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(1);
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledWith(
['guide-news'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
);
} finally {
consoleError.mockRestore();
}
});
it('falls back to global EPG metadata for channels missing from playlist-scoped sources', async () => {
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsChannelMetadata = true;
epgBridge.getChannelMetadata = jest
.fn()
.mockResolvedValueOnce({
'guide-news': {
id: 'guide-news',
displayName: 'Playlist News',
iconUrl: 'https://playlist.example.com/news.png',
},
'guide-sports': null,
})
.mockResolvedValueOnce({
'guide-sports': {
id: 'guide-sports',
displayName: 'Global Sports',
iconUrl: 'https://global.example.com/sports.png',
},
});
const result = await firstValueFrom(
service.getChannelMetadataForChannels(
['guide-news', 'guide-sports'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
)
);
expect(result.get('guide-news')).toMatchObject({
displayName: 'Playlist News',
});
expect(result.get('guide-sports')).toMatchObject({
displayName: 'Global Sports',
iconUrl: 'https://global.example.com/sports.png',
});
expect(epgBridge.getChannelMetadata).toHaveBeenNthCalledWith(
1,
['guide-news', 'guide-sports'],
{ sourceUrls: ['https://playlist.example.com/guide.xml'] }
);
expect(epgBridge.getChannelMetadata).toHaveBeenNthCalledWith(
2,
['guide-sports'],
{ sourceUrls: ['https://global.example.com/guide.xml'] }
);
});
});
+427 -48
View File
@@ -2,15 +2,27 @@ import { inject, Injectable } from '@angular/core';
import { MatSnackBar } from '@angular/material/snack-bar';
import { TranslateService } from '@ngx-translate/core';
import { BehaviorSubject, forkJoin, from, Observable, of } from 'rxjs';
import { catchError, map, tap, timeout } from 'rxjs/operators';
import {
catchError,
finalize,
map,
shareReplay,
switchMap,
tap,
timeout,
} from 'rxjs/operators';
import {
createDevLogger,
EpgChannelMetadata,
EpgProgram,
} from '@iptvnator/shared/interfaces';
import { SettingsStore } from '@iptvnator/services';
import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service';
import {
EpgLookupOptions,
EpgRuntimeBridgeService,
} from './epg-runtime-bridge.service';
import { normalizeEpgPrograms } from './epg-program-normalization.util';
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
interface CachedProgram {
program: EpgProgram | null;
@@ -33,6 +45,14 @@ export class EpgService {
// Cache for channel programs with 60-second TTL
private programCache = new Map<string, CachedProgram>();
private fetchingCurrentPrograms = new Map<
string,
Observable<EpgProgram | null>
>();
private fetchingCurrentProgramBatches = new Map<
string,
Observable<Map<string, EpgProgram | null>>
>();
private readonly CACHE_TTL = 60000; // 60 seconds
readonly epgAvailable$ = this.epgAvailable.asObservable();
@@ -44,8 +64,8 @@ export class EpgService {
fetchEpg(urls: string[]): void {
if (!this.epgBridge.supportsImport) return;
// Filter out empty URLs and send all URLs at once
const validUrls = urls.filter((url) => url?.trim());
// Filter out empty and duplicate URLs and send all URLs at once.
const validUrls = normalizeEpgUrls(urls);
if (validUrls.length === 0) return;
from(
@@ -59,6 +79,7 @@ export class EpgService {
if (result === null) return;
if (result.success) {
this.clearCache();
this.epgAvailable.next(true);
} else {
this.epgAvailable.next(false);
@@ -117,50 +138,46 @@ export class EpgService {
* @returns Observable of current program or null
*/
getCurrentProgramForChannel(
channelId: string
channelId: string,
options?: EpgLookupOptions
): Observable<EpgProgram | null> {
if (!this.epgBridge.supportsProgramLookup || !channelId) {
return of(null);
}
// Check cache first
const cached = this.programCache.get(channelId);
const now = Date.now();
if (cached && now - cached.timestamp < this.CACHE_TTL) {
return of(cached.program);
const sourceUrls = this.normalizeSourceUrls(options);
if (sourceUrls.length > 0) {
return this.getScopedCurrentProgramForChannel(
channelId,
sourceUrls,
this.getGlobalEpgSourceUrls(sourceUrls)
);
}
const globalSourceUrls = this.getGlobalEpgSourceUrls();
if (globalSourceUrls.length > 0) {
return this.getScopedCurrentProgramForChannel(
channelId,
globalSourceUrls,
[]
);
}
// Check cache first
const cacheKey = this.createProgramCacheKey(channelId);
// Fetch from backend
return from(this.epgBridge.getChannelPrograms(channelId)).pipe(
map((programs) => normalizeEpgPrograms(programs ?? [])),
map((programs: EpgProgram[]) => {
if (!programs.length) {
this.programCache.set(channelId, {
program: null,
timestamp: now,
});
return null;
}
const currentProgram = this.findCurrentProgram(programs);
// Cache the result
this.programCache.set(channelId, {
program: currentProgram,
timestamp: now,
});
return currentProgram;
}),
catchError((err) => {
console.error('EPG get current program error:', err);
this.programCache.set(channelId, {
program: null,
timestamp: now,
});
return of(null);
})
return this.getCachedOrFetchCurrentProgram(cacheKey, () =>
from(this.epgBridge.getChannelPrograms(channelId)).pipe(
map((programs) => normalizeEpgPrograms(programs ?? [])),
map((programs: EpgProgram[]) =>
this.findCurrentProgram(programs)
),
catchError((err) => {
console.error('EPG get current program error:', err);
return of(null);
})
)
);
}
@@ -185,7 +202,8 @@ export class EpgService {
* @returns Observable of Map with channelId -> current program
*/
getCurrentProgramsForChannels(
channelIds: string[]
channelIds: string[],
options?: EpgLookupOptions
): Observable<Map<string, EpgProgram | null>> {
if (!this.epgBridge.supportsProgramLookup) {
return of(new Map());
@@ -195,6 +213,29 @@ export class EpgService {
return of(new Map());
}
const sourceUrls = this.normalizeSourceUrls(options);
if (
sourceUrls.length > 0 &&
this.epgBridge.supportsCurrentProgramBatch
) {
return this.getScopedCurrentProgramsForChannels(
channelIds,
sourceUrls
);
}
const globalSourceUrls = this.getGlobalEpgSourceUrls();
if (
globalSourceUrls.length > 0 &&
this.epgBridge.supportsCurrentProgramBatch
) {
return this.getScopedCurrentProgramsForChannels(
channelIds,
globalSourceUrls,
[]
);
}
const resultMap = new Map<string, EpgProgram | null>();
const channelsToFetch: string[] = [];
const now = Date.now();
@@ -260,31 +301,236 @@ export class EpgService {
);
}
private getScopedCurrentProgramsForChannels(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls = this.getGlobalEpgSourceUrls(sourceUrls)
): Observable<Map<string, EpgProgram | null>> {
const normalizedChannelIds = this.normalizeChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
const resultMap = new Map<string, EpgProgram | null>();
const channelsToFetch: string[] = [];
normalizedChannelIds.forEach((channelId) => {
const cached = this.getCachedProgram(
this.createProgramCacheKey(channelId, sourceUrls)
);
if (cached) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
if (channelsToFetch.length === 0) {
return of(resultMap);
}
const batchCacheKey = this.createProgramBatchCacheKey(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
);
const existingRequest =
this.fetchingCurrentProgramBatches.get(batchCacheKey);
const request$ =
existingRequest ??
this.fetchScopedCurrentProgramsBatch(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
).pipe(
tap((fetchedMap) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
this.programCache.set(
this.createProgramCacheKey(channelId, sourceUrls),
{
program: fetchedMap.get(channelId) ?? null,
timestamp: cacheTimestamp,
}
);
});
}),
finalize(() => {
this.fetchingCurrentProgramBatches.delete(batchCacheKey);
}),
shareReplay({ bufferSize: 1, refCount: false })
);
if (!existingRequest) {
this.fetchingCurrentProgramBatches.set(batchCacheKey, request$);
}
return request$.pipe(
map((fetchedMap) => {
const mergedResultMap = new Map(resultMap);
channelsToFetch.forEach((channelId) => {
mergedResultMap.set(
channelId,
fetchedMap.get(channelId) ?? null
);
});
return mergedResultMap;
})
);
}
private fetchScopedCurrentProgramsBatch(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<Map<string, EpgProgram | null>> {
return from(
this.epgBridge.getCurrentProgramsBatch(channelIds, {
sourceUrls,
})
).pipe(
timeout(5000),
switchMap((scopedResult) => {
const resultMap = new Map<string, EpgProgram | null>();
const fallbackChannelIds: string[] = [];
channelIds.forEach((channelId) => {
const program = scopedResult?.[channelId] ?? null;
resultMap.set(channelId, program);
if (!program) {
fallbackChannelIds.push(channelId);
}
});
if (fallbackChannelIds.length === 0) {
return of(resultMap);
}
if (fallbackSourceUrls.length === 0) {
return of(resultMap);
}
return from(
this.epgBridge.getCurrentProgramsBatch(fallbackChannelIds, {
sourceUrls: fallbackSourceUrls,
})
).pipe(
timeout(5000),
map((globalResult) => {
fallbackChannelIds.forEach((channelId) => {
resultMap.set(
channelId,
globalResult?.[channelId] ?? null
);
});
return resultMap;
}),
catchError((err) => {
console.error(
'EPG global fallback current programs error:',
err
);
return of(resultMap);
})
);
}),
catchError((err) => {
console.error('EPG scoped batch current programs error:', err);
return of(this.createNullProgramMap(channelIds));
})
);
}
getChannelMetadataForChannels(
channelIds: string[]
channelIds: string[],
options?: EpgLookupOptions
): Observable<Map<string, EpgChannelMetadata | null>> {
if (!this.epgBridge.supportsChannelMetadata) {
return of(new Map());
}
const normalizedChannelIds = Array.from(
const normalizedChannelIds = this.normalizeChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
const sourceUrls = this.normalizeSourceUrls(options);
const globalSourceUrls =
sourceUrls.length > 0
? this.getGlobalEpgSourceUrls(sourceUrls)
: this.getGlobalEpgSourceUrls();
const effectiveSourceUrls =
sourceUrls.length > 0 ? sourceUrls : globalSourceUrls;
return this.getChannelMetadataMapForSourceUrls(
normalizedChannelIds,
effectiveSourceUrls
).pipe(
switchMap((metadataMap) => {
const fallbackChannelIds =
sourceUrls.length > 0 && globalSourceUrls.length > 0
? normalizedChannelIds.filter(
(channelId) => !metadataMap.get(channelId)
)
: [];
if (fallbackChannelIds.length === 0) {
return of(metadataMap);
}
return this.getChannelMetadataMapForSourceUrls(
fallbackChannelIds,
globalSourceUrls
).pipe(
map((globalMetadataMap) => {
fallbackChannelIds.forEach((channelId) => {
metadataMap.set(
channelId,
globalMetadataMap.get(channelId) ?? null
);
});
return metadataMap;
}),
catchError((err) => {
console.error(
'EPG global fallback channel metadata error:',
err
);
return of(metadataMap);
})
);
})
);
}
private normalizeChannelIds(channelIds: string[]): string[] {
return Array.from(
new Set(
channelIds
.map((channelId) => channelId.trim())
.filter((channelId) => channelId.length > 0)
)
);
}
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
private normalizeSourceUrls(options?: EpgLookupOptions): string[] {
return normalizeEpgUrls(options?.sourceUrls ?? []);
}
private getChannelMetadataMapForSourceUrls(
channelIds: string[],
sourceUrls: string[]
): Observable<Map<string, EpgChannelMetadata | null>> {
return from(
this.epgBridge.getChannelMetadata(normalizedChannelIds)
this.epgBridge.getChannelMetadata(
channelIds,
sourceUrls.length > 0 ? { sourceUrls } : undefined
)
).pipe(
map((metadataByChannelId) => {
return new Map<string, EpgChannelMetadata | null>(
normalizedChannelIds.map((channelId) => [
channelIds.map((channelId) => [
channelId,
metadataByChannelId?.[channelId] ?? null,
])
@@ -297,10 +543,143 @@ export class EpgService {
);
}
private createProgramCacheKey(
channelId: string,
sourceUrls: string[] = []
): string {
const normalizedSourceUrls = normalizeEpgUrls(sourceUrls);
if (normalizedSourceUrls.length === 0) {
return channelId;
}
return `source:${channelId}:${JSON.stringify(normalizedSourceUrls)}`;
}
private createProgramBatchCacheKey(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): string {
return JSON.stringify({
channelIds: [...channelIds].sort(),
sourceUrls: normalizeEpgUrls(sourceUrls),
fallbackSourceUrls: normalizeEpgUrls(fallbackSourceUrls),
});
}
private getCachedProgram(cacheKey: string): CachedProgram | undefined {
const cached = this.programCache.get(cacheKey);
if (!cached) {
return undefined;
}
if (Date.now() - cached.timestamp >= this.CACHE_TTL) {
this.programCache.delete(cacheKey);
return undefined;
}
return cached;
}
private getCachedOrFetchCurrentProgram(
cacheKey: string,
fetchProgram: () => Observable<EpgProgram | null>
): Observable<EpgProgram | null> {
const cached = this.getCachedProgram(cacheKey);
if (cached) {
return of(cached.program);
}
const existingRequest = this.fetchingCurrentPrograms.get(cacheKey);
if (existingRequest) {
return existingRequest;
}
const request$ = fetchProgram().pipe(
tap((program) => {
this.programCache.set(cacheKey, {
program,
timestamp: Date.now(),
});
}),
finalize(() => {
this.fetchingCurrentPrograms.delete(cacheKey);
}),
shareReplay({ bufferSize: 1, refCount: false })
);
this.fetchingCurrentPrograms.set(cacheKey, request$);
return request$;
}
private getScopedCurrentProgramForChannel(
channelId: string,
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<EpgProgram | null> {
const cacheKey = this.createProgramCacheKey(channelId, sourceUrls);
return this.getCachedOrFetchCurrentProgram(cacheKey, () =>
from(
this.epgBridge.getChannelPrograms(channelId, { sourceUrls })
).pipe(
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
switchMap((programs) => {
const currentProgram = this.findCurrentProgram(programs);
if (currentProgram) {
return of(currentProgram);
}
return this.getFallbackCurrentProgramForChannel(
channelId,
fallbackSourceUrls
);
}),
catchError((err) => {
console.error('EPG scoped current program error:', err);
return this.getFallbackCurrentProgramForChannel(
channelId,
fallbackSourceUrls
);
})
)
);
}
private getFallbackCurrentProgramForChannel(
channelId: string,
sourceUrls: string[]
): Observable<EpgProgram | null> {
if (sourceUrls.length === 0) {
return of(null);
}
return this.getScopedCurrentProgramForChannel(
channelId,
sourceUrls,
[]
);
}
private createNullProgramMap(
channelIds: string[]
): Map<string, EpgProgram | null> {
return new Map(channelIds.map((channelId) => [channelId, null]));
}
private getGlobalEpgSourceUrls(excluding: string[] = []): string[] {
const excludedUrls = new Set(excluding);
return normalizeEpgUrls(
this.settingsStore.getSettings().epgUrl ?? []
).filter((url) => !excludedUrls.has(url));
}
/**
* Clears the program cache (useful when EPG is refreshed)
*/
clearCache(): void {
this.programCache.clear();
this.fetchingCurrentPrograms.clear();
this.fetchingCurrentProgramBatches.clear();
}
}