fix(dashboard): live rails find XMLTV programmes outside the global EPG sources (#1637)

This commit is contained in:
4gray authored and GitHub committed 2026-09-20 21:48:55 +02:00
1 parent b30c783e85
commit 954c8ad65e
22 files changed
+2019 -659

No files matched your search

@@ -0,0 +1,9 @@
---
type: fix
area: dashboard
---
The dashboard's "Now on air on favorite channels" and "Recently watched Live
TV" rails now show the current programme for channels whose guide lives in an
XMLTV another playlist imported, not only in the global EPG sources from
Settings — the same lookup the "See all" pages already used.
+1
View File
@@ -1606,6 +1606,7 @@ stream_id`); it drops `series_id`/`movie_id`, so the builder pins the
- XMLTV format support, from `http(s)` links or local files (Electron only): a `file:` URL, an absolute POSIX path, or a Windows drive/UNC path, plain `.xml` or gzip (detected by signature). A folder button beside each row opens the native picker (`EPG_OPEN_FILE_DIALOG`, `RuntimeCapabilitiesService.supportsEpgFilePicker`). Shape rules: `classifyEpgSourceReference` in `libs/shared/interfaces`; the worker opens both kinds through `openEpgSourceStream` (`workers/epg-source-stream.ts`). Only hand-chosen sources may be local: `extractM3uEpgUrls` harvests only remote links from M3U headers (legacy stored non-remote entries are dropped unless manual), and `EpgWorkerService.startFetch` asks the main-process `EpgLocalSourceAuthorizer` before a local path reaches the worker — picker results are trusted, a typed path is confirmed once in a native message box, allowed paths persist under `TRUSTED_LOCAL_EPG_SOURCES`, and the worker's local branch requires the main-set `allowLocalFile` flag (deny-all until wired). Contract: `docs/architecture/m3u-playlist-module.md` ("Local XMLTV files")
- Background parsing in worker thread; HTTP/file gzip compatibility follows `docs/architecture/m3u-playlist-module.md` ("XMLTV response compression").
- Stored in database for quick lookup
- Source scope: batch "now" lookups search the playlist's own XMLTV first, then the Settings-managed global sources, and stop there; only `EpgLookupOptions.anySourceFallback` (renderer-only) retries the still-unresolved keys against every imported source. The dashboard live rails pass it, so a favourite whose guide lives in another playlist's XMLTV gets the same programme its "See all" row already showed; they still issue one lookup per distinct playlist source scope and namespace the answers by it (`liveEpgProgramKey`), since a `tvg-id` is unique inside a guide but not across imports, and only cards carrying a real XMLTV key are widened — a portal card's key is just its title (wiring: `DashboardLiveEpgPresenter`). The channel list keeps the strict scope. Contract: `docs/architecture/m3u-playlist-module.md` (the "Scoped lookups" bullet under playlist-scoped URLs)
- Global display-time offset (`Settings.epgOffsetMinutes`, Settings → EPG, ±720 min, Electron only): display-only, provider data is never rewritten. Two equivalent forms in `libs/shared/interfaces/src/lib/epg-display-offset.util.ts` — `epgDisplayTimeMs` (shift the programme; `ui/epg` rendering via the `offsetMinutes` input, channel rows, dashboard/recording labels; the programme dialog and the programme guide read the store themselves) and `epgProviderClockMs` (shift "now"; every "currently airing" decision: the `GET_CURRENT_PROGRAMS_BATCH` lookup takes an explicit `nowMs` and `EpgService` tags its cache with the offset, Xtream/Stalker/M3U current-programme selection and previews, the unified collection resolver, dashboard progress, recording overlap). A consumer applies exactly one form per comparison. Contract: `docs/architecture/m3u-playlist-module.md` ("EPG display offset")
- Programme guide (Electron, M3U): `app-epg-guide` in `libs/ui/epg` fed by the host-provided `EPG_GUIDE_SOURCE`; the M3U host switches into guide mode (docked player strip, no sidebar/timeline, no remount) from the header action, the palette, the EPG panel's Guide button (timeline or list view) or `G`. Data: `EPG_GET_PROGRAMS_FOR_CHANNELS` / `EPG_GET_PROGRAM_COVERAGE` (keys resolved in main; manual mappings honoured). Contract: `docs/architecture/m3u-playlist-module.md` ("Programme guide").
- Manual EPG mapping (Electron only): right-click a channel in any list (M3U views, Xtream portal list, Stalker ITV sidebar, global favorites) → "Map EPG channel" attaches it to an uploaded-XMLTV channel; stored in `epg_channel_mappings` keyed by the M3U lookup key or a playlist-scoped portal key (`xtream:{playlistId}:{id}` / `stalker:{playlistId}:{id}`, helpers in `libs/shared/interfaces/src/lib/epg-mapping-key.util.ts`); resolved on every EPG path (single + batch IPC lookups, portal detail views, preview queues); dialog: `libs/ui/components/src/lib/channel-list-container/epg-mapping-dialog/`
+129
View File
@@ -801,6 +801,135 @@ test.describe('Electron EPG', () => {
}
});
test("@epg @electron dashboard live rails find a programme that only another playlist's XMLTV carries", async ({
dataDir,
}) => {
test.setTimeout(120000);
// The reported case: the favourited channel's own playlist declares
// no guide, Settings hold a global XMLTV that does not know it, and
// the programme exists only in the guide a DIFFERENT playlist
// imported. The rails searched the global scope alone and showed
// nothing, while the "See all" row — resolved unscoped — had it.
const otherPlaylistEpgServer = await createMutableTextServer(
createCurrentXmltvFixture(
'playlist-guide-news',
'Playlist Guide News',
'Other Playlist Bulletin'
),
{
contentType: 'application/xml; charset=utf-8',
resourcePath: '/guides/other-playlist.xml',
}
);
const globalEpgServer = await createMutableTextServer(
createCurrentXmltvFixture(
'global-other',
'Global Other',
'Global Other Bulletin'
),
{
contentType: 'application/xml; charset=utf-8',
resourcePath: '/guides/global-guide.xml',
}
);
// Declares the guide, so importing it is what puts those programmes
// in the database under that source.
const guideOwnerServer = await createMutableTextServer(
buildM3uContent([
{
name: 'Playlist Guide News',
tvgId: 'playlist-guide-news',
url: 'https://example.com/live/guide-owner.m3u8',
},
]).replace(
'#EXTM3U',
`#EXTM3U x-tvg-url="${otherPlaylistEpgServer.resourceUrl}"`
),
{
contentType: 'application/x-mpegurl; charset=utf-8',
resourcePath: '/guide-owner.m3u',
}
);
// Carries the same XMLTV id under its own display name and declares
// no guide at all — the playlist the dashboard card comes from.
const guidelessServer = await createMutableTextServer(
buildM3uContent([
{
name: 'Mirror News',
tvgId: 'playlist-guide-news',
url: 'https://example.com/live/mirror-news.m3u8',
},
]),
{
contentType: 'application/x-mpegurl; charset=utf-8',
resourcePath: '/guideless.m3u',
}
);
const app = await launchElectronApp(dataDir);
try {
await importM3uPlaylistFromUrl(
app.mainWindow,
guideOwnerServer.resourceUrl
);
await expect(
app.mainWindow.locator(
'.epg-progress-panel .import-item.status-complete'
)
).toHaveCount(1, { 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')
.first()
.fill(globalEpgServer.resourceUrl);
await saveSettings(app.mainWindow);
await expect
.poll(() => getEpgChannelCount(app.mainWindow), {
timeout: 30000,
})
.toBe(2);
await importM3uPlaylistFromUrl(
app.mainWindow,
guidelessServer.resourceUrl
);
await openWorkspaceSection(app.mainWindow, 'All channels');
const channelItem = channelItemByTitle(
app.mainWindow,
'Mirror News'
);
await expect(channelItem).toBeVisible({ timeout: 20000 });
await channelItem.hover();
await channelItem.locator('.favorite-button').first().click();
await expect(
channelItem.locator('.favorite-button mat-icon').first()
).toHaveText(/star/);
await goToDashboard(app.mainWindow);
const card = app.mainWindow
.locator('[data-test-id="dashboard-live-favorites-rail-card"]')
.filter({ hasText: 'Mirror News' })
.first();
await expect(card).toBeVisible({ timeout: 20000 });
await expect(card.locator('.rail__channel-now')).toContainText(
'Other Playlist Bulletin',
{ timeout: 30000 }
);
} finally {
await closeElectronApp(app);
await guidelessServer.close();
await guideOwnerServer.close();
await otherPlaylistEpgServer.close();
await globalEpgServer.close();
}
});
test('@epg @electron uses the XMLTV channel icon as a fallback when the playlist has no tvg-logo', async ({
dataDir,
}) => {
+28 -1
View File
@@ -978,7 +978,34 @@ These URLs are playlist-scoped by default:
TTL expires.
- Scoped lookups fall back only to Settings-managed EPG URLs for channels
missing from the playlist-declared source. Playlist-local sources from other
playlists are not treated as global fallback sources. Single-channel current
playlists are not treated as global fallback sources. The one opt-out is
`EpgLookupOptions.anySourceFallback` (renderer-only, never forwarded to the
bridge): after the scope — playlist sources, then the global ones — has
answered, the keys still without a programme are retried once against every
imported source through the source-less batch path and its cache. The
ladder has the same shape with or without the bridge's batch endpoint: on
an older preload the scoped pass runs as per-channel scoped lookups
(`getScopedCurrentProgramForChannel`, same scope -> fallback-scope walk,
same scoped cache key) and only the any-source retry is source-less.
Collapsing that preload straight into the source-less lookup would drop
the caller's scope, which is what the scopes exist to prevent. The
dashboard live
rails pass the option: without it a favourite whose guide only exists in
another playlist's XMLTV showed no programme on the dashboard while its
"See all" row — resolved by `StreamResolverService`, which never scopes by
source — had one. They still ask **per source scope**, one lookup per
distinct set of playlist-declared XMLTV URLs, and namespace the answers by
that scope: a `tvg-id` is unique inside a guide, not across imports, so a
single flat map keyed by lookup key alone would hand one playlist's card
the programme another playlist's guide resolved for the same id. Playlists
sharing a guide share one lookup. Only a card that carries a real XMLTV key
is widened: an Xtream or Stalker card has none, so its lookup key is just
its display title, and searching every guide by title would let a
same-named M3U channel answer for a portal channel. Those cards keep the
strict scope (their own programmes come from the portal), and the
any-source flag is part of the scope identity so the two never share an
answer. Wiring: `DashboardLiveEpgPresenter` in
`libs/workspace/dashboard/feature/src/lib/rails/`. The channel list keeps the strict scope. Single-channel current
program lookups include the source URL set in their cache and in-flight keys,
so playlist-local and global lookups deduplicate without reusing the wrong
source scope. Batch current-program lookups use the same source-scoped
@@ -0,0 +1,101 @@
import { MonoTypeOperatorFunction, Observable, from, of } from 'rxjs';
import { catchError, map, switchMap } from 'rxjs/operators';
import { EpgChannelMetadata } from '@iptvnator/shared/interfaces';
import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service';
/**
* Channel icons and display names from the imported XMLTV, resolved with the
* same playlist-first, Settings-managed-fallback strategy as the programme
* lookups: a playlist-local guide that only supplies programmes must still
* let icons come from a global source.
*
* Extracted from `EpgService` to keep that file under the repository's
* production size limit; it carries no state of its own.
*/
export class EpgChannelMetadataLookup {
constructor(
private readonly epgBridge: EpgRuntimeBridgeService,
private readonly guard: <T>() => MonoTypeOperatorFunction<T>,
private readonly globalSourceUrls: (excluding?: string[]) => string[]
) {}
forChannels(
normalizedChannelIds: string[],
sourceUrls: string[]
): Observable<Map<string, EpgChannelMetadata | null>> {
const globalSourceUrls =
sourceUrls.length > 0
? this.globalSourceUrls(sourceUrls)
: this.globalSourceUrls();
const effectiveSourceUrls =
sourceUrls.length > 0 ? sourceUrls : globalSourceUrls;
return this.forSourceUrls(
normalizedChannelIds,
effectiveSourceUrls
).pipe(
this.guard<Map<string, EpgChannelMetadata | null>>(),
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.forSourceUrls(
fallbackChannelIds,
globalSourceUrls
).pipe(
this.guard<Map<string, EpgChannelMetadata | null>>(),
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 forSourceUrls(
channelIds: string[],
sourceUrls: string[]
): Observable<Map<string, EpgChannelMetadata | null>> {
return from(
this.epgBridge.getChannelMetadata(
channelIds,
sourceUrls.length > 0 ? { sourceUrls } : undefined
)
).pipe(
this.guard<Record<string, EpgChannelMetadata | null> | null>(),
map((metadataByChannelId) => {
return new Map<string, EpgChannelMetadata | null>(
channelIds.map((channelId) => [
channelId,
metadataByChannelId?.[channelId] ?? null,
])
);
}),
catchError((err) => {
console.error('EPG get channel metadata error:', err);
return of(new Map<string, EpgChannelMetadata | null>());
})
);
}
}
@@ -0,0 +1,256 @@
import { Observable, forkJoin, from, of } from 'rxjs';
import { catchError, map, switchMap, timeout } from 'rxjs/operators';
import { EpgProgram } from '@iptvnator/shared/interfaces';
import { EpgLookupOptions } from './epg-runtime-bridge.service';
import {
normalizeLookupChannelIds,
normalizeLookupSourceUrls,
} from './epg-lookup-normalization.util';
import type { EpgLookupContext } from './epg-lookup-context';
import { EpgSingleProgramLookup } from './epg-single-program.lookup';
import { EpgScopedBatchLookup } from './epg-scoped-batch.lookup';
/**
* "What is on air right now" for a batch of channels.
*
* The source-scope ladder lives here: the caller's playlist sources first,
* then the Settings-managed global ones, and — only when the caller opts in
* with `anySourceFallback` — a final retry across every imported guide. A
* bridge without the batch endpoint walks the same ladder one channel at a
* time through `EpgSingleProgramLookup`, so both shapes answer identically.
*/
export class EpgCurrentProgramsLookup {
constructor(
private readonly ctx: EpgLookupContext,
private readonly single: EpgSingleProgramLookup,
private readonly scopedBatch: EpgScopedBatchLookup
) {}
/**
* Gets current programs for multiple channels (batch operation)
* @param channelIds Array of channel IDs
* @returns Observable of Map with channelId -> current program
*/
forChannels(
channelIds: string[],
options?: EpgLookupOptions
): Observable<Map<string, EpgProgram | null>> {
if (!this.ctx.bridge.supportsProgramLookup) {
return of(new Map());
}
if (!channelIds || channelIds.length === 0) {
return of(new Map());
}
const scoped$ = this.getSourceScopedCurrentProgramsForChannels(
channelIds,
options
);
if (!scoped$) {
// Nothing declares a scope — neither the caller nor Settings — so
// the pool of every imported source is the only answer there is.
return this.getUnscopedCurrentProgramsForChannels(channelIds);
}
if (!options?.anySourceFallback) {
return scoped$;
}
return scoped$.pipe(
switchMap((scopedMap) => this.fillFromAnySource(scopedMap))
);
}
/**
* The scoped lookup: the caller's playlist sources first (with the global
* sources as fallback), else the Settings-managed global sources alone.
* `null` when no scope applies and the unscoped pool is the only answer.
*
* The ladder is the same whether or not the bridge has the batch
* endpoint. Collapsing a legacy preload straight into the source-less
* lookup would drop the caller's scope entirely, which is exactly what
* the scopes exist to prevent: two imported guides reusing one XMLTV id
* would answer each other's channels.
*/
private getSourceScopedCurrentProgramsForChannels(
channelIds: string[],
options?: EpgLookupOptions
): Observable<Map<string, EpgProgram | null>> | null {
const sourceUrls = normalizeLookupSourceUrls(options);
if (sourceUrls.length > 0) {
return this.scopedCurrentProgramsForChannels(
channelIds,
sourceUrls,
this.ctx.globalSourceUrls(sourceUrls)
);
}
const globalSourceUrls = this.ctx.globalSourceUrls();
if (globalSourceUrls.length > 0) {
return this.scopedCurrentProgramsForChannels(
channelIds,
globalSourceUrls,
[]
);
}
return null;
}
/** One scoped batch, or its per-channel equivalent on an older preload. */
private scopedCurrentProgramsForChannels(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<Map<string, EpgProgram | null>> {
if (this.ctx.bridge.supportsCurrentProgramBatch) {
return this.scopedBatch.forChannels(
channelIds,
sourceUrls,
fallbackSourceUrls
);
}
const normalizedChannelIds = normalizeLookupChannelIds(channelIds);
if (normalizedChannelIds.length === 0) {
return of(new Map());
}
// The per-channel lookup walks the same scope -> fallback-scope
// ladder the batch query does, and caches under the same scoped key.
return forkJoin(
normalizedChannelIds.map((channelId) =>
this.single
.scoped(channelId, sourceUrls, fallbackSourceUrls)
.pipe(
this.ctx.guard(),
timeout(5000),
map((program) => ({ channelId, program })),
catchError(() => of({ channelId, program: null }))
)
)
).pipe(
this.ctx.guard(),
map(
(results) =>
new Map(
results.map((result) => [
result.channelId,
result.program,
])
)
)
);
}
/**
* `anySourceFallback`: keys the scope answered with `null` are retried
* against every imported source. A scoped miss is kept as the answer
* when the pool has nothing either, so the merged map still names every
* requested key.
*/
private fillFromAnySource(
scopedMap: Map<string, EpgProgram | null>
): Observable<Map<string, EpgProgram | null>> {
const unresolvedIds = Array.from(scopedMap.entries())
.filter(([, program]) => !program)
.map(([channelId]) => channelId);
if (unresolvedIds.length === 0) {
return of(scopedMap);
}
return this.getUnscopedCurrentProgramsForChannels(unresolvedIds).pipe(
map((anySourceMap) => {
const mergedMap = new Map(scopedMap);
anySourceMap.forEach((program, channelId) => {
if (program) {
mergedMap.set(channelId, program);
}
});
return mergedMap;
})
);
}
/** Lookup across every imported source, cached under the source-less key. */
private getUnscopedCurrentProgramsForChannels(
channelIds: string[]
): Observable<Map<string, EpgProgram | null>> {
const resultMap = new Map<string, EpgProgram | null>();
const channelsToFetch: string[] = [];
const now = Date.now();
const offsetMinutes = this.ctx.offsetMinutes();
// Check cache for each channel
channelIds.forEach((channelId) => {
const cached = this.ctx.cache.getFresh(channelId, now);
if (cached && cached.offsetMinutes === offsetMinutes) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
// If all channels were cached, return immediately
if (channelsToFetch.length === 0) {
return of(resultMap);
}
// Single batched IPC + SQL query when the backend supports it.
// Replaces the legacy N+1 forkJoin where each channel fired its own
// GET_CHANNEL_PROGRAMS round-trip.
if (this.ctx.bridge.supportsCurrentProgramBatch) {
return from(
this.ctx.bridge.getCurrentProgramsBatch(channelsToFetch, {
nowMs: this.ctx.clockMs(),
})
).pipe(
this.ctx.guard(),
timeout(5000),
map((batchResult) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
const program = batchResult?.[channelId] ?? null;
resultMap.set(channelId, program);
this.ctx.cache.set(
channelId,
program,
offsetMinutes,
cacheTimestamp
);
});
return resultMap;
}),
catchError((err) => {
console.error('EPG batch current programs error:', err);
return of(resultMap);
})
);
}
// Fallback for older preload bundles without the batch endpoint.
// Deliberately the unscoped per-channel lookup: everything reaching
// here wants the source-less pool — either nothing declared a scope,
// or this is the any-source retry after the scoped pass. The public
// `getCurrentProgramForChannel` would put the Settings scope back on
// and make the retry re-ask the question the scoped pass answered.
const fetchObservables = channelsToFetch.map((channelId) =>
this.single.unscoped(channelId).pipe(
this.ctx.guard(),
timeout(5000),
map((program) => ({ channelId, program })),
catchError(() => of({ channelId, program: null }))
)
);
return forkJoin(fetchObservables).pipe(
this.ctx.guard(),
map((results) => {
results.forEach((result) => {
resultMap.set(result.channelId, result.program);
});
return resultMap;
})
);
}
}
@@ -0,0 +1,17 @@
import { MonoTypeOperatorFunction } from 'rxjs';
import { EpgProgramCache } from './epg-program-cache';
import { EpgRuntimeBridgeService } from './epg-runtime-bridge.service';
/** Everything the lookups need from `EpgService`, and nothing else. */
export interface EpgLookupContext {
readonly bridge: EpgRuntimeBridgeService;
readonly cache: EpgProgramCache;
/** Retires work that a changed XMLTV source set has invalidated. */
guard<T>(): MonoTypeOperatorFunction<T>;
/** Current EPG display offset, in minutes. */
offsetMinutes(): number;
/** Wall-clock now, expressed in the provider's uncorrected EPG clock. */
clockMs(): number;
/** Settings-managed XMLTV URLs, optionally minus a caller's own scope. */
globalSourceUrls(excluding?: string[]): string[];
}
@@ -0,0 +1,20 @@
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
import { EpgLookupOptions } from './epg-runtime-bridge.service';
/** Trimmed, de-duplicated, non-empty channel ids, in request order. */
export function normalizeLookupChannelIds(channelIds: string[]): string[] {
return Array.from(
new Set(
channelIds
.map((channelId) => channelId.trim())
.filter((channelId) => channelId.length > 0)
)
);
}
/** The caller's declared XMLTV scope, normalized like every stored URL. */
export function normalizeLookupSourceUrls(
options?: EpgLookupOptions
): string[] {
return normalizeEpgUrls(options?.sourceUrls ?? []);
}
@@ -0,0 +1,166 @@
import { MonoTypeOperatorFunction, Observable, of } from 'rxjs';
import { finalize, shareReplay, tap } from 'rxjs/operators';
import { EpgProgram } from '@iptvnator/shared/interfaces';
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
export interface CachedProgram {
program: EpgProgram | null;
timestamp: number;
/** Display offset the "now" verdict was computed with; a changed setting invalidates the entry. */
offsetMinutes: number;
}
const CACHE_TTL_MS = 60000;
/**
* The 60 s "currently airing" memory behind every lookup in `EpgService`,
* plus the in-flight requests sharing one answer.
*
* Owning it here keeps one rule in one place: an entry answers "what is on
* at the provider clock, in this source scope", so the scope and the display
* offset are both part of its identity. A lookup issued after the offset
* changed must neither read the previous entry nor join a request still in
* flight for it.
*
* Deliberately not `@Injectable`: it has no dependencies of its own, it
* takes the two facts it cannot know (the current offset, and the operator
* that retires work across an EPG source change) from its owner.
*/
export class EpgProgramCache {
private readonly programs = new Map<string, CachedProgram>();
private readonly inFlight = new Map<
string,
Observable<EpgProgram | null>
>();
private readonly inFlightBatches = new Map<
string,
Observable<Map<string, EpgProgram | null>>
>();
constructor(
private readonly offsetMinutes: () => number,
private readonly guard: <T>() => MonoTypeOperatorFunction<T>
) {}
/** Cache and in-flight identity of a single-channel lookup. */
keyFor(channelId: string, sourceUrls: string[] = []): string {
const normalizedSourceUrls = normalizeEpgUrls(sourceUrls);
const key =
normalizedSourceUrls.length === 0
? channelId
: `source:${channelId}:${JSON.stringify(normalizedSourceUrls)}`;
const offsetMinutes = this.offsetMinutes();
return offsetMinutes === 0 ? key : `${key}|offset:${offsetMinutes}`;
}
/** Identity of a batch request; order-insensitive in the channel ids. */
batchKeyFor(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): string {
return JSON.stringify({
channelIds: [...channelIds].sort(),
sourceUrls: normalizeEpgUrls(sourceUrls),
fallbackSourceUrls: normalizeEpgUrls(fallbackSourceUrls),
offsetMinutes: this.offsetMinutes(),
});
}
/** A live entry, or `undefined` once it expired or the offset moved. */
get(cacheKey: string): CachedProgram | undefined {
const cached = this.programs.get(cacheKey);
if (!cached) {
return undefined;
}
if (
Date.now() - cached.timestamp >= CACHE_TTL_MS ||
cached.offsetMinutes !== this.offsetMinutes()
) {
this.programs.delete(cacheKey);
return undefined;
}
return cached;
}
/**
* `offsetMinutes` is the offset the verdict was computed with, not the
* one current when the answer lands.
*/
set(
cacheKey: string,
program: EpgProgram | null,
offsetMinutes: number,
timestamp = Date.now()
): void {
this.programs.set(cacheKey, { program, timestamp, offsetMinutes });
}
/** Unexpired entry without the scope/offset checks `get` applies. */
getFresh(cacheKey: string, now: number): CachedProgram | undefined {
const cached = this.programs.get(cacheKey);
return cached && now - cached.timestamp < CACHE_TTL_MS
? cached
: undefined;
}
/** The cached answer, an in-flight one, else one started from `fetch`. */
getOrFetch(
cacheKey: string,
fetchProgram: () => Observable<EpgProgram | null>
): Observable<EpgProgram | null> {
const cached = this.get(cacheKey);
if (cached) {
return of(cached.program);
}
const existingRequest = this.inFlight.get(cacheKey);
if (existingRequest) {
return existingRequest;
}
const offsetMinutes = this.offsetMinutes();
const request$ = fetchProgram().pipe(
this.guard<EpgProgram | null>(),
tap((program) => this.set(cacheKey, program, offsetMinutes)),
finalize(() => {
if (this.inFlight.get(cacheKey) === request$)
this.inFlight.delete(cacheKey);
}),
shareReplay({ bufferSize: 1, refCount: false })
);
this.inFlight.set(cacheKey, request$);
return request$;
}
batchInFlight(
batchCacheKey: string
): Observable<Map<string, EpgProgram | null>> | undefined {
return this.inFlightBatches.get(batchCacheKey);
}
registerBatch(
batchCacheKey: string,
request$: Observable<Map<string, EpgProgram | null>>
): void {
this.inFlightBatches.set(batchCacheKey, request$);
}
/** Only the request that owns the slot may release it. */
releaseBatch(
batchCacheKey: string,
request$: Observable<Map<string, EpgProgram | null>>
): void {
if (this.inFlightBatches.get(batchCacheKey) === request$) {
this.inFlightBatches.delete(batchCacheKey);
}
}
clear(): void {
this.programs.clear();
this.inFlight.clear();
this.inFlightBatches.clear();
}
}
@@ -27,7 +27,18 @@ export type EpgImportProgress = ElectronBridgeEpgProgress;
export type EpgFetchResult = ElectronBridgeEpgFetchResult;
export type EpgFreshnessResult = ElectronBridgeEpgFreshnessResult;
export type EpgClearResult = ElectronBridgeResult;
export type EpgLookupOptions = ElectronBridgeEpgLookupOptions;
export interface EpgLookupOptions extends ElectronBridgeEpgLookupOptions {
/**
* Retry the keys the scoped lookup (playlist scope, then Settings-managed
* global sources) left without a programme against every imported XMLTV
* source. Off by default: the scopes exist so a playlist's own guide wins
* over a same-named channel in another playlist's guide. Surfaces with no
* playlist context — the dashboard live rails — opt in, matching the
* collection resolver and the timeline, which never scope by source.
* Renderer-only: it is not forwarded to the desktop bridge.
*/
anySourceFallback?: boolean;
}
export type EpgCurrentProgramsOptions = ElectronBridgeCurrentProgramsOptions;
type EpgElectronBridge = Pick<
@@ -0,0 +1,180 @@
import { Observable, from, of } from 'rxjs';
import {
catchError,
finalize,
map,
shareReplay,
switchMap,
tap,
timeout,
} from 'rxjs/operators';
import { EpgProgram } from '@iptvnator/shared/interfaces';
import { normalizeLookupChannelIds } from './epg-lookup-normalization.util';
import type { EpgLookupContext } from './epg-lookup-context';
/**
* How ONE scope is asked: the per-channel cache in front of the request, the
* order-insensitive deduplication of identical batches still in flight, and
* the retry of whatever that scope could not answer against the fallback
* scope. `EpgCurrentProgramsLookup` decides WHICH scope to ask.
*/
export class EpgScopedBatchLookup {
constructor(private readonly ctx: EpgLookupContext) {}
forChannels(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls = this.ctx.globalSourceUrls(sourceUrls)
): Observable<Map<string, EpgProgram | null>> {
const normalizedChannelIds = normalizeLookupChannelIds(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.ctx.cache.get(
this.ctx.cache.keyFor(channelId, sourceUrls)
);
if (cached) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
if (channelsToFetch.length === 0) {
return of(resultMap);
}
const batchCacheKey = this.ctx.cache.batchKeyFor(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
);
const existingRequest = this.ctx.cache.batchInFlight(batchCacheKey);
// Tag entries with the offset the verdict was computed with, not the
// one current when the response lands.
const offsetMinutes = this.ctx.offsetMinutes();
const request$ =
existingRequest ??
this.fetchBatch(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
).pipe(
this.ctx.guard(),
tap((fetchedMap) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
this.ctx.cache.set(
this.ctx.cache.keyFor(channelId, sourceUrls),
fetchedMap.get(channelId) ?? null,
offsetMinutes,
cacheTimestamp
);
});
}),
finalize(() =>
this.ctx.cache.releaseBatch(batchCacheKey, request$)
),
shareReplay({ bufferSize: 1, refCount: false })
);
if (!existingRequest) {
this.ctx.cache.registerBatch(batchCacheKey, request$);
}
return request$.pipe(
this.ctx.guard(),
map((fetchedMap) => {
const mergedResultMap = new Map(resultMap);
channelsToFetch.forEach((channelId) => {
mergedResultMap.set(
channelId,
fetchedMap.get(channelId) ?? null
);
});
return mergedResultMap;
})
);
}
private fetchBatch(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<Map<string, EpgProgram | null>> {
const nowMs = this.ctx.clockMs();
return from(
this.ctx.bridge.getCurrentProgramsBatch(channelIds, {
sourceUrls,
nowMs,
})
).pipe(
this.ctx.guard(),
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.ctx.bridge.getCurrentProgramsBatch(
fallbackChannelIds,
{
sourceUrls: fallbackSourceUrls,
nowMs,
}
)
).pipe(
this.ctx.guard(),
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.nullProgramMap(channelIds));
})
);
}
private nullProgramMap(
channelIds: string[]
): Map<string, EpgProgram | null> {
return new Map(channelIds.map((channelId) => [channelId, null]));
}
}
@@ -0,0 +1,129 @@
import { Observable, from, of } from 'rxjs';
import { catchError, map, switchMap, timeout } from 'rxjs/operators';
import { EpgProgram } from '@iptvnator/shared/interfaces';
import { EpgLookupOptions } from './epg-runtime-bridge.service';
import { normalizeEpgPrograms } from './epg-program-normalization.util';
import { normalizeLookupSourceUrls } from './epg-lookup-normalization.util';
import type { EpgLookupContext } from './epg-lookup-context';
/**
* "What is on air" for a single channel, and the scope ladder one channel
* walks: the caller's sources, then the Settings-managed ones, then — for a
* caller that asked for it — the whole imported pool. The batch lookup backs
* its own per-channel work with these, so both shapes answer identically.
*/
export class EpgSingleProgramLookup {
constructor(private readonly ctx: EpgLookupContext) {}
/**
* Gets the current EPG program for a specific channel (with caching)
* @param channelId Channel ID (tvg-id or channel name)
* @returns Observable of current program or null
*/
forChannel(
channelId: string,
options?: EpgLookupOptions
): Observable<EpgProgram | null> {
if (!this.ctx.bridge.supportsProgramLookup || !channelId) {
return of(null);
}
const sourceUrls = normalizeLookupSourceUrls(options);
if (sourceUrls.length > 0) {
return this.scoped(
channelId,
sourceUrls,
this.ctx.globalSourceUrls(sourceUrls)
);
}
const globalSourceUrls = this.ctx.globalSourceUrls();
if (globalSourceUrls.length > 0) {
return this.scoped(channelId, globalSourceUrls, []);
}
return this.unscoped(channelId);
}
/**
* The lookup across every imported source, cached under the source-less
* key. Kept separate from `getCurrentProgramForChannel`, which re-applies
* the Settings-managed scope whenever global URLs exist: a caller that
* has already decided it wants the unscoped pool — the `anySourceFallback`
* retry on a preload without the batch endpoint — must not have that
* scope put back on.
*/
unscoped(channelId: string): Observable<EpgProgram | null> {
// Check cache first
const cacheKey = this.ctx.cache.keyFor(channelId);
// Fetch from backend
return this.ctx.cache.getOrFetch(cacheKey, () =>
from(this.ctx.bridge.getChannelPrograms(channelId)).pipe(
this.ctx.guard(),
map((programs) => normalizeEpgPrograms(programs ?? [])),
map((programs: EpgProgram[]) => this.findCurrent(programs)),
catchError((err) => {
console.error('EPG get current program error:', err);
return of(null);
})
)
);
}
/**
* Finds the current program from a list of programs
*/
findCurrent(programs: EpgProgram[]): EpgProgram | null {
const now = this.ctx.clockMs();
return (
programs.find((program) => {
const start = new Date(program.start).getTime();
const stop = new Date(program.stop).getTime();
return start <= now && now <= stop;
}) || null
);
}
scoped(
channelId: string,
sourceUrls: string[],
fallbackSourceUrls: string[]
): Observable<EpgProgram | null> {
const cacheKey = this.ctx.cache.keyFor(channelId, sourceUrls);
return this.ctx.cache.getOrFetch(cacheKey, () =>
from(
this.ctx.bridge.getChannelPrograms(channelId, { sourceUrls })
).pipe(
this.ctx.guard(),
timeout(3000),
map((programs) => normalizeEpgPrograms(programs ?? [])),
switchMap((programs) => {
const currentProgram = this.findCurrent(programs);
if (currentProgram) {
return of(currentProgram);
}
return this.fallbackFor(channelId, fallbackSourceUrls);
}),
catchError((err) => {
console.error('EPG scoped current program error:', err);
return this.fallbackFor(channelId, fallbackSourceUrls);
})
)
);
}
private fallbackFor(
channelId: string,
sourceUrls: string[]
): Observable<EpgProgram | null> {
if (sourceUrls.length === 0) {
return of(null);
}
return this.scoped(channelId, sourceUrls, []);
}
}
@@ -594,6 +594,258 @@ describe('EpgService', () => {
);
});
it('retries keys the global scope left unresolved against every source when anySourceFallback is set', async () => {
// The dashboard rails: no playlist scope, one global XMLTV in
// Settings, and a channel whose guide only exists in an XMLTV another
// playlist imported. The "See all" pages resolve it unscoped, so the
// rails must too — but only after the configured scope had its say.
settingsStore.getSettings.mockReturnValue({
epgUrl: ['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: 'Global 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: 'Other Playlist Sports',
},
});
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(
['guide-news', 'guide-sports'],
{ anySourceFallback: true }
)
);
expect(result.get('guide-news')?.title).toBe('Global Bulletin');
expect(result.get('guide-sports')?.title).toBe('Other Playlist Sports');
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(2);
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenNthCalledWith(
1,
['guide-news', 'guide-sports'],
expect.objectContaining({
sourceUrls: ['https://global.example.com/guide.xml'],
})
);
// The retry carries only the unresolved key and no source scope.
const [retryIds, retryOptions] = (
epgBridge.getCurrentProgramsBatch as jest.Mock
).mock.calls[1];
expect(retryIds).toEqual(['guide-sports']);
expect(retryOptions).not.toHaveProperty('sourceUrls');
});
it('walks scope then any-source per channel on a preload without the batch endpoint', async () => {
// The ladder must not change shape with the capability: the scoped
// pass runs first and only its misses reach the source-less pool.
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = false;
epgBridge.getChannelPrograms = jest
.fn()
.mockResolvedValueOnce([])
.mockResolvedValueOnce([])
.mockResolvedValueOnce([
{
channel: 'guide-sports',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Other Playlist Sports',
},
]);
jest.useFakeTimers();
jest.setSystemTime(new Date('2026-05-23T10:30:00.000Z'));
try {
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-sports'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
anySourceFallback: true,
})
);
expect(result.get('guide-sports')?.title).toBe(
'Other Playlist Sports'
);
const calls = (epgBridge.getChannelPrograms as jest.Mock).mock
.calls;
expect(calls.map(([, options]) => options?.sourceUrls)).toEqual([
['https://playlist.example.com/guide.xml'],
['https://global.example.com/guide.xml'],
undefined,
]);
} finally {
jest.useRealTimers();
}
});
it('keeps the playlist scope on a preload without the batch endpoint when no any-source retry was asked for', async () => {
// Jumping straight to the source-less lookup here would defeat the
// scopes: two guides reusing one XMLTV id would answer each other.
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = false;
epgBridge.getChannelPrograms = jest.fn().mockResolvedValue([]);
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-sports'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
})
);
expect(result.get('guide-sports')).toBeNull();
const calls = (epgBridge.getChannelPrograms as jest.Mock).mock.calls;
expect(calls.map(([, options]) => options?.sourceUrls)).toEqual([
['https://playlist.example.com/guide.xml'],
['https://global.example.com/guide.xml'],
]);
});
it('keeps the Settings scope on the legacy per-channel path when no any-source retry was asked for', async () => {
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = false;
epgBridge.getChannelPrograms = jest.fn().mockResolvedValue([]);
await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-sports'])
);
expect(epgBridge.getChannelPrograms).toHaveBeenCalledWith(
'guide-sports',
expect.objectContaining({
sourceUrls: ['https://global.example.com/guide.xml'],
})
);
});
it('keeps the scoped verdict when the any-source retry finds nothing either', async () => {
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
epgBridge.getCurrentProgramsBatch = jest
.fn()
.mockResolvedValueOnce({ 'guide-sports': null })
.mockResolvedValueOnce({ 'guide-sports': null });
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-sports'], {
anySourceFallback: true,
})
);
expect(result.size).toBe(1);
expect(result.get('guide-sports')).toBeNull();
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(2);
});
it('skips the any-source retry when every key resolved in scope, and never retries without the option', async () => {
settingsStore.getSettings.mockReturnValue({
epgUrl: ['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: 'Global Bulletin',
},
})
.mockResolvedValueOnce({ 'guide-sports': null });
const resolved = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-news'], {
anySourceFallback: true,
})
);
expect(resolved.get('guide-news')?.title).toBe('Global Bulletin');
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(1);
// Scoped callers (the channel list) keep today's contract: a miss in
// the configured scope stays a miss.
const unresolved = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-sports'])
);
expect(unresolved.get('guide-sports')).toBeNull();
expect(epgBridge.getCurrentProgramsBatch).toHaveBeenCalledTimes(2);
});
it('runs the any-source retry after the playlist scope and its global fallback', async () => {
settingsStore.getSettings.mockReturnValue({
epgUrl: ['https://global.example.com/guide.xml'],
trustedPrivateNetworkEpgUrls: [],
trustedInsecureTlsHosts: [],
});
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
epgBridge.getCurrentProgramsBatch = jest
.fn()
.mockResolvedValueOnce({ 'guide-sports': null })
.mockResolvedValueOnce({ 'guide-sports': null })
.mockResolvedValueOnce({
'guide-sports': {
channel: 'guide-sports',
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
title: 'Other Playlist Sports',
},
});
const result = await firstValueFrom(
service.getCurrentProgramsForChannels(['guide-sports'], {
sourceUrls: ['https://playlist.example.com/guide.xml'],
anySourceFallback: true,
})
);
expect(result.get('guide-sports')?.title).toBe('Other Playlist Sports');
const calls = (epgBridge.getCurrentProgramsBatch as jest.Mock).mock
.calls;
expect(calls.map(([, options]) => options.sourceUrls)).toEqual([
['https://playlist.example.com/guide.xml'],
['https://global.example.com/guide.xml'],
undefined,
]);
});
it('caches scoped batch current programs by EPG source URL scope', async () => {
epgBridge.supportsProgramLookup = true;
epgBridge.supportsCurrentProgramBatch = true;
+47 -567
View File
@@ -1,16 +1,8 @@
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,
finalize,
map,
shareReplay,
switchMap,
tap,
timeout,
} from 'rxjs/operators';
import { BehaviorSubject, from, Observable, of } from 'rxjs';
import { catchError, map, tap, timeout } from 'rxjs/operators';
import {
createDevLogger,
EpgChannelMetadata,
@@ -23,15 +15,18 @@ import {
EpgRuntimeBridgeService,
} from './epg-runtime-bridge.service';
import { normalizeEpgPrograms } from './epg-program-normalization.util';
import { EpgProgramCache } from './epg-program-cache';
import { EpgChannelMetadataLookup } from './epg-channel-metadata.lookup';
import { EpgCurrentProgramsLookup } from './epg-current-programs.lookup';
import { EpgSingleProgramLookup } from './epg-single-program.lookup';
import { EpgScopedBatchLookup } from './epg-scoped-batch.lookup';
import type { EpgLookupContext } from './epg-lookup-context';
import {
normalizeLookupChannelIds,
normalizeLookupSourceUrls,
} from './epg-lookup-normalization.util';
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
interface CachedProgram {
program: EpgProgram | null;
timestamp: number;
/** Display offset the "now" verdict was computed with; a changed setting invalidates the entry. */
offsetMinutes: number;
}
const debugEpgService = createDevLogger('EpgService');
@Injectable({
@@ -55,17 +50,36 @@ export class EpgService {
private epgAvailable = new BehaviorSubject<boolean>(false);
private currentEpgPrograms = new BehaviorSubject<EpgProgram[]>([]);
// 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
/** Channel icons/display names; same scope ladder as the programmes. */
private readonly channelMetadata = new EpgChannelMetadataLookup(
this.epgBridge,
() => this.sourceSettings.guard(),
(excluding) => this.getGlobalEpgSourceUrls(excluding)
);
/** 60 s "currently airing" memory, keyed by scope + display offset. */
private readonly programCache = new EpgProgramCache(
() => this.epgOffsetMinutes(),
() => this.sourceSettings.guard()
);
/** The whole "what is on air" scope ladder, single channel or batch. */
private readonly lookupContext: EpgLookupContext = {
bridge: this.epgBridge,
cache: this.programCache,
guard: () => this.sourceSettings.guard(),
offsetMinutes: () => this.epgOffsetMinutes(),
clockMs: () => this.epgClockMs(),
globalSourceUrls: (excluding) => this.getGlobalEpgSourceUrls(excluding),
};
private readonly singleProgram = new EpgSingleProgramLookup(
this.lookupContext
);
private readonly currentPrograms = new EpgCurrentProgramsLookup(
this.lookupContext,
this.singleProgram,
new EpgScopedBatchLookup(this.lookupContext)
);
/** Display offset every "currently airing" decision in here is made with. */
private epgOffsetMinutes(): number {
@@ -179,57 +193,7 @@ export class EpgService {
if (!this.epgBridge.supportsProgramLookup || !channelId) {
return of(null);
}
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 this.getCachedOrFetchCurrentProgram(cacheKey, () =>
from(this.epgBridge.getChannelPrograms(channelId)).pipe(
this.sourceSettings.guard(),
map((programs) => normalizeEpgPrograms(programs ?? [])),
map((programs: EpgProgram[]) =>
this.findCurrentProgram(programs)
),
catchError((err) => {
console.error('EPG get current program error:', err);
return of(null);
})
)
);
}
/**
* Finds the current program from a list of programs
*/
private findCurrentProgram(programs: EpgProgram[]): EpgProgram | null {
const now = this.epgClockMs();
return (
programs.find((program) => {
const start = new Date(program.start).getTime();
const stop = new Date(program.stop).getTime();
return start <= now && now <= stop;
}) || null
);
return this.singleProgram.forChannel(channelId, options);
}
/**
@@ -244,266 +208,10 @@ export class EpgService {
if (!this.epgBridge.supportsProgramLookup) {
return of(new Map());
}
if (!channelIds || channelIds.length === 0) {
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();
const offsetMinutes = this.epgOffsetMinutes();
// Check cache for each channel
channelIds.forEach((channelId) => {
const cached = this.programCache.get(channelId);
if (
cached &&
now - cached.timestamp < this.CACHE_TTL &&
cached.offsetMinutes === offsetMinutes
) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
}
});
// If all channels were cached, return immediately
if (channelsToFetch.length === 0) {
return of(resultMap);
}
// Single batched IPC + SQL query when the backend supports it.
// Replaces the legacy N+1 forkJoin where each channel fired its own
// GET_CHANNEL_PROGRAMS round-trip.
if (this.epgBridge.supportsCurrentProgramBatch) {
return from(
this.epgBridge.getCurrentProgramsBatch(channelsToFetch, {
nowMs: this.epgClockMs(),
})
).pipe(
this.sourceSettings.guard(),
timeout(5000),
map((batchResult) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
const program = batchResult?.[channelId] ?? null;
resultMap.set(channelId, program);
this.programCache.set(channelId, {
program,
timestamp: cacheTimestamp,
offsetMinutes,
});
});
return resultMap;
}),
catchError((err) => {
console.error('EPG batch current programs error:', err);
return of(resultMap);
})
);
}
// 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 }))
)
);
return forkJoin(fetchObservables).pipe(
this.sourceSettings.guard(),
map((results) => {
results.forEach((result) => {
resultMap.set(result.channelId, result.program);
});
return resultMap;
})
);
}
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);
// Tag entries with the offset the verdict was computed with, not the
// one current when the response lands.
const offsetMinutes = this.epgOffsetMinutes();
const request$ =
existingRequest ??
this.fetchScopedCurrentProgramsBatch(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
).pipe(
this.sourceSettings.guard(),
tap((fetchedMap) => {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
this.programCache.set(
this.createProgramCacheKey(channelId, sourceUrls),
{
program: fetchedMap.get(channelId) ?? null,
timestamp: cacheTimestamp,
offsetMinutes,
}
);
});
}),
finalize(() => {
if (
this.fetchingCurrentProgramBatches.get(
batchCacheKey
) === request$
)
this.fetchingCurrentProgramBatches.delete(
batchCacheKey
);
}),
shareReplay({ bufferSize: 1, refCount: false })
);
if (!existingRequest) {
this.fetchingCurrentProgramBatches.set(batchCacheKey, request$);
}
return request$.pipe(
this.sourceSettings.guard(),
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>> {
const nowMs = this.epgClockMs();
return from(
this.epgBridge.getCurrentProgramsBatch(channelIds, {
sourceUrls,
nowMs,
})
).pipe(
this.sourceSettings.guard(),
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,
nowMs,
})
).pipe(
this.sourceSettings.guard(),
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));
})
);
return this.currentPrograms.forChannels(channelIds, options);
}
getChannelMetadataForChannels(
@@ -514,243 +222,17 @@ export class EpgService {
return of(new Map());
}
const normalizedChannelIds = this.normalizeChannelIds(channelIds);
const normalizedChannelIds = normalizeLookupChannelIds(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(
return this.channelMetadata.forChannels(
normalizedChannelIds,
effectiveSourceUrls
).pipe(
this.sourceSettings.guard(),
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(
this.sourceSettings.guard(),
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);
})
);
})
normalizeLookupSourceUrls(options)
);
}
private normalizeChannelIds(channelIds: string[]): string[] {
return Array.from(
new Set(
channelIds
.map((channelId) => channelId.trim())
.filter((channelId) => channelId.length > 0)
)
);
}
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(
channelIds,
sourceUrls.length > 0 ? { sourceUrls } : undefined
)
).pipe(
this.sourceSettings.guard(),
map((metadataByChannelId) => {
return new Map<string, EpgChannelMetadata | null>(
channelIds.map((channelId) => [
channelId,
metadataByChannelId?.[channelId] ?? null,
])
);
}),
catchError((err) => {
console.error('EPG get channel metadata error:', err);
return of(new Map<string, EpgChannelMetadata | null>());
})
);
}
/**
* Cache and in-flight identity of a lookup. The display offset is part of
* it: a request evaluated at another provider clock answers a different
* question, so a lookup issued after the setting changed must neither
* read the previous entry nor join a batch still in flight for it.
*/
private createProgramCacheKey(
channelId: string,
sourceUrls: string[] = []
): string {
const normalizedSourceUrls = normalizeEpgUrls(sourceUrls);
const key =
normalizedSourceUrls.length === 0
? channelId
: `source:${channelId}:${JSON.stringify(normalizedSourceUrls)}`;
const offsetMinutes = this.epgOffsetMinutes();
return offsetMinutes === 0 ? key : `${key}|offset:${offsetMinutes}`;
}
private createProgramBatchCacheKey(
channelIds: string[],
sourceUrls: string[],
fallbackSourceUrls: string[]
): string {
return JSON.stringify({
channelIds: [...channelIds].sort(),
sourceUrls: normalizeEpgUrls(sourceUrls),
fallbackSourceUrls: normalizeEpgUrls(fallbackSourceUrls),
offsetMinutes: this.epgOffsetMinutes(),
});
}
private getCachedProgram(cacheKey: string): CachedProgram | undefined {
const cached = this.programCache.get(cacheKey);
if (!cached) {
return undefined;
}
if (
Date.now() - cached.timestamp >= this.CACHE_TTL ||
cached.offsetMinutes !== this.epgOffsetMinutes()
) {
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 offsetMinutes = this.epgOffsetMinutes();
const request$ = fetchProgram().pipe(
this.sourceSettings.guard(),
tap((program) => {
this.programCache.set(cacheKey, {
program,
timestamp: Date.now(),
offsetMinutes,
});
}),
finalize(() => {
if (this.fetchingCurrentPrograms.get(cacheKey) === request$)
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(
this.sourceSettings.guard(),
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(
@@ -763,7 +245,5 @@ export class EpgService {
*/
clearCache(): void {
this.programCache.clear();
this.fetchingCurrentPrograms.clear();
this.fetchingCurrentProgramBatches.clear();
}
}
@@ -9,12 +9,17 @@ export type {
export {
buildDashboardLiveEpgDetails,
buildLiveEpgCardsForEnabledRails,
buildLiveEpgLookupKeys,
buildLiveEpgLookupGroups,
calcEpgProgress,
formatEpgTimeRange,
getLiveEpgProgramForCard,
liveEpgProgramKey,
liveEpgScopeKey,
} from './lib/rails/dashboard-live-epg.utils';
export type {
DashboardLiveEpgDetails,
DashboardLiveEpgLookupGroup,
} from './lib/rails/dashboard-live-epg.utils';
export type { DashboardLiveEpgDetails } from './lib/rails/dashboard-live-epg.utils';
export {
buildPlaybackPositionReloadKey,
formatRemainingLabel,
@@ -0,0 +1,221 @@
import { signal } from '@angular/core';
import { TestBed } from '@angular/core/testing';
import { EMPTY, of, throwError } from 'rxjs';
import { EpgService } from '@iptvnator/epg/data-access';
import type { EpgProgram, PlaylistMeta } from '@iptvnator/shared/interfaces';
import { SettingsStore } from '@iptvnator/services';
import { DashboardDataService } from '@iptvnator/workspace/dashboard/data-access';
import { DashboardLiveEpgPresenter } from './dashboard-live-epg.presenter';
import type { DashboardRailCard } from './dashboard-rail.component';
const guideA = 'https://a.example/guide.xml';
const guideB = 'https://b.example/guide.xml';
const m3uPlaylist = (id: string, epgUrls: string[]): PlaylistMeta =>
({ _id: id, epgUrls }) as PlaylistMeta;
const card = (overrides: Partial<DashboardRailCard>): DashboardRailCard =>
({
id: 'card',
title: 'Das Erste HD',
icon: 'live_tv',
contentType: 'live',
link: ['/workspace'],
...overrides,
}) as DashboardRailCard;
const program = (title: string): EpgProgram =>
({
channel: 'ard.de',
title,
start: '2026-05-23T10:00:00.000Z',
stop: '2026-05-23T11:00:00.000Z',
}) as EpgProgram;
describe('DashboardLiveEpgPresenter', () => {
let presenter: DashboardLiveEpgPresenter;
let getCurrentProgramsForChannels: jest.Mock;
let playlists: ReturnType<typeof signal<PlaylistMeta[]>>;
const setup = (cards: DashboardRailCard[]) => {
presenter.connect(signal(cards));
// `toObservable` pushes through an effect, so the connected cards
// only reach the lookup once effects run.
TestBed.tick();
};
beforeEach(() => {
jest.useFakeTimers();
jest.setSystemTime(new Date('2026-05-23T10:30:00.000Z'));
getCurrentProgramsForChannels = jest.fn(() => of(new Map()));
playlists = signal<PlaylistMeta[]>([
m3uPlaylist('a', [guideA]),
m3uPlaylist('a2', [guideA]),
m3uPlaylist('b', [guideB]),
{ _id: 'portal', serverUrl: 'http://portal' } as PlaylistMeta,
]);
TestBed.configureTestingModule({
providers: [
DashboardLiveEpgPresenter,
{
provide: DashboardDataService,
useValue: { playlists },
},
{
provide: EpgService,
useValue: { getCurrentProgramsForChannels },
},
{
provide: SettingsStore,
useValue: { resolvedEpgOffsetMinutes: () => 0 },
},
],
});
presenter = TestBed.inject(DashboardLiveEpgPresenter);
});
afterEach(() => {
jest.useRealTimers();
});
it('asks each guide once and lets playlists on the same guide share the lookup', () => {
setup([
card({ id: 'a', epgLookupKey: 'ard.de', epgPlaylistId: 'a' }),
card({ id: 'a2', epgLookupKey: 'zdf.de', epgPlaylistId: 'a2' }),
card({ id: 'b', epgLookupKey: 'ard.de', epgPlaylistId: 'b' }),
]);
expect(
getCurrentProgramsForChannels.mock.calls.map(([keys, options]) => [
keys,
options.sourceUrls,
])
).toEqual([
[['ard.de', 'zdf.de'], [guideA]],
[['ard.de'], [guideB]],
]);
});
it('keeps a portal card out of the any-source retry and in its own answer space', () => {
const portalCard = card({ id: 'p', epgPlaylistId: 'portal' });
const m3uCard = card({
id: 'm',
epgLookupKey: 'Das Erste HD',
epgPlaylistId: 'guideless',
});
getCurrentProgramsForChannels.mockImplementation((keys, options) =>
of(
new Map<string, EpgProgram | null>([
[
'Das Erste HD',
program(
options.anySourceFallback
? 'From any guide'
: 'From Settings only'
),
],
])
)
);
setup([m3uCard, portalCard]);
expect(
getCurrentProgramsForChannels.mock.calls.map(
([, options]) => options.anySourceFallback
)
).toEqual([true, false]);
// Both resolve against Settings, both key on the same title, and they
// still do not share an answer.
expect(presenter.detailsFor(m3uCard)?.nowPlayingTitle).toBe(
'From any guide'
);
expect(presenter.detailsFor(portalCard)?.nowPlayingTitle).toBe(
'From Settings only'
);
});
it('never hands a card the programme another guide resolved for the same id', () => {
const fromA = card({
id: 'a',
epgLookupKey: 'ard.de',
epgPlaylistId: 'a',
});
const fromB = card({
id: 'b',
epgLookupKey: 'ard.de',
epgPlaylistId: 'b',
});
getCurrentProgramsForChannels.mockImplementation((keys, options) =>
of(
new Map<string, EpgProgram | null>([
[
'ard.de',
program(
options.sourceUrls[0] === guideA
? 'Guide A bulletin'
: 'Guide B bulletin'
),
],
])
)
);
setup([fromA, fromB]);
expect(presenter.detailsFor(fromA)?.nowPlayingTitle).toBe(
'Guide A bulletin'
);
expect(presenter.detailsFor(fromB)?.nowPlayingTitle).toBe(
'Guide B bulletin'
);
});
it('keeps the other guides when one lookup is retired or fails mid-tick', () => {
const fromA = card({
id: 'a',
epgLookupKey: 'ard.de',
epgPlaylistId: 'a',
});
const fromB = card({
id: 'b',
epgLookupKey: 'ard.de',
epgPlaylistId: 'b',
});
getCurrentProgramsForChannels.mockImplementation((keys, options) =>
options.sourceUrls[0] === guideA
? of(
new Map<string, EpgProgram | null>([
['ard.de', program('Guide A bulletin')],
])
)
: // What the EPG source-change guard does to work in flight:
// complete without ever emitting.
EMPTY
);
setup([fromA, fromB]);
expect(presenter.detailsFor(fromA)?.nowPlayingTitle).toBe(
'Guide A bulletin'
);
expect(presenter.detailsFor(fromB)).toBeNull();
getCurrentProgramsForChannels.mockImplementation((keys, options) =>
options.sourceUrls[0] === guideA
? of(
new Map<string, EpgProgram | null>([
['ard.de', program('Guide A bulletin')],
])
)
: throwError(() => new Error('lookup failed'))
);
jest.advanceTimersByTime(30_000);
expect(presenter.detailsFor(fromA)?.nowPlayingTitle).toBe(
'Guide A bulletin'
);
expect(presenter.detailsFor(fromB)).toBeNull();
});
});
@@ -0,0 +1,183 @@
import {
computed,
inject,
Injectable,
signal,
type Signal,
} from '@angular/core';
import { toObservable, toSignal } from '@angular/core/rxjs-interop';
import {
catchError,
defaultIfEmpty,
forkJoin,
interval,
map,
of,
startWith,
switchMap,
} from 'rxjs';
import { EpgService } from '@iptvnator/epg/data-access';
import type { EpgProgram } from '@iptvnator/shared/interfaces';
import { SettingsStore } from '@iptvnator/services';
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
import { DashboardDataService } from '@iptvnator/workspace/dashboard/data-access';
import type { DashboardRailCard } from './dashboard-rail.component';
import {
buildDashboardLiveEpgDetails,
buildLiveEpgLookupGroups,
getLiveEpgProgramForCard,
liveEpgAllowsAnySource,
liveEpgProgramKey,
liveEpgScopeKey,
LIVE_EPG_TICK_MS,
type DashboardLiveEpgDetails,
type DashboardLiveEpgLookupGroup,
} from './dashboard-live-epg.utils';
type ScopeAnswer = {
readonly scopeKey: string;
readonly programs: ReadonlyMap<string, EpgProgram | null>;
};
const emptyAnswer = (scopeKey: string): ScopeAnswer => ({
scopeKey,
programs: new Map<string, EpgProgram | null>(),
});
/**
* The dashboard's uploaded-XMLTV "now on air" lookup for live cards.
* Component-provided, so the 30 s heartbeat dies with the page.
*
* Best-effort, keyed by the app-wide M3U chain (tvg-id -> tvg-name -> name)
* with the card title as a final fallback. Xtream/Stalker live items often
* have no XMLTV side-channel and simply return null — the card then renders
* without the programme row.
*
* One lookup per source scope, not one for the whole page: a `tvg-id` is
* unique inside a guide, not across imports, so a card is only ever handed
* the answer resolved in ITS playlist's scope. Each lookup then opts into the
* any-source retry — Settings global sources first, then every imported
* XMLTV, which is what the "See all" collection pages resolve against.
* Without it a channel whose guide only exists in another playlist's XMLTV
* showed no programme here while its "See all" row had one.
*/
@Injectable()
export class DashboardLiveEpgPresenter {
private readonly data = inject(DashboardDataService);
private readonly epgService = inject(EpgService);
private readonly settingsStore = inject(SettingsStore);
private readonly cards = signal<Signal<
readonly DashboardRailCard[]
> | null>(null);
// The XMLTV sources each playlist declares. Same rule as
// `ChannelListContainerComponent`: only an M3U playlist carries its own
// guide; a portal playlist is answered from the Settings-managed URLs.
private readonly sourceUrlsByPlaylist = computed(() => {
const byPlaylistId = new Map<string, string[]>();
for (const playlist of this.data.playlists()) {
byPlaylistId.set(
playlist._id,
playlist.serverUrl || playlist.macAddress
? []
: normalizeEpgUrls(playlist.epgUrls ?? [])
);
}
return byPlaylistId;
});
private readonly lookupGroups = computed(() =>
buildLiveEpgLookupGroups(this.cards()?.() ?? [], (card) =>
this.sourceUrlsForCard(card)
)
);
// Re-fetch on rail change AND on a 30s heartbeat so the progress bar
// catches the boundary between programs without a full page revisit.
private readonly programs = toSignal(
toObservable(this.lookupGroups).pipe(
switchMap((groups) =>
groups.length === 0
? of(new Map<string, EpgProgram | null>())
: interval(LIVE_EPG_TICK_MS).pipe(
startWith(0),
switchMap(() =>
forkJoin(
groups.map((group) => this.askScope(group))
).pipe(map((answers) => mergeAnswers(answers)))
)
)
)
),
{ initialValue: new Map<string, EpgProgram | null>() }
);
/** The live cards whose rails are enabled, hero included. */
connect(cards: Signal<readonly DashboardRailCard[]>): void {
this.cards.set(cards);
}
/** `null` when nothing is known about the card's current programme. */
detailsFor(card: DashboardRailCard | null): DashboardLiveEpgDetails | null {
if (!card) {
return null;
}
const program = getLiveEpgProgramForCard(
card,
this.programs(),
liveEpgScopeKey(
this.sourceUrlsForCard(card),
liveEpgAllowsAnySource(card)
)
);
// Recompute the now-window each tick so progress moves between
// 30s ticks even if the program identity is unchanged.
return buildDashboardLiveEpgDetails(
program,
Date.now(),
this.settingsStore.resolvedEpgOffsetMinutes()
);
}
/** The XMLTV scope a live card's programme must be resolved in. */
private sourceUrlsForCard(card: DashboardRailCard): string[] {
const byPlaylistId = this.sourceUrlsByPlaylist();
return (
(card.epgPlaylistId
? byPlaylistId.get(card.epgPlaylistId)
: undefined) ?? []
);
}
private askScope(group: DashboardLiveEpgLookupGroup) {
const scopeKey = group.scopeKey;
return this.epgService
.getCurrentProgramsForChannels(group.lookupKeys, {
sourceUrls: group.sourceUrls,
anySourceFallback: group.anySourceFallback,
})
.pipe(
map((programs): ScopeAnswer => ({ scopeKey, programs })),
// One scope retired by an EPG source change, or failing, must
// not blank every other scope's answer for the tick:
// `forkJoin` emits nothing at all when one input completes
// without a value.
defaultIfEmpty(emptyAnswer(scopeKey)),
catchError(() => of(emptyAnswer(scopeKey)))
);
}
}
/** Namespaced by scope, so two guides sharing an id stay apart. */
function mergeAnswers(
answers: readonly ScopeAnswer[]
): Map<string, EpgProgram | null> {
const merged = new Map<string, EpgProgram | null>();
for (const answer of answers) {
for (const [lookupKey, program] of answer.programs) {
merged.set(liveEpgProgramKey(answer.scopeKey, lookupKey), program);
}
}
return merged;
}
@@ -110,18 +110,96 @@ function liveEpgLookupKeyForCard(card: DashboardRailCard): string {
return card.epgLookupKey?.trim() || card.title.trim();
}
export function buildLiveEpgLookupKeys(
cards: readonly DashboardRailCard[]
): string[] {
const seen = new Set<string>();
const keys: string[] = [];
/** One XMLTV lookup: the scope it is answered in and the keys asked for. */
export interface DashboardLiveEpgLookupGroup {
/** Identity of the scope, as `liveEpgScopeKey` builds it. */
readonly scopeKey: string;
readonly sourceUrls: string[];
readonly lookupKeys: string[];
/** Whether unresolved keys may be retried against every imported guide. */
readonly anySourceFallback: boolean;
}
/**
* Only a card whose playlist gave it a real XMLTV key may be retried against
* every imported guide. An Xtream or Stalker card carries no such key, so the
* lookup falls back to its display title, and searching every playlist's
* guide by title would let a same-named M3U channel answer for a portal
* channel. Those cards keep the strict scope; their own programmes come from
* the portal instead.
*/
export function liveEpgAllowsAnySource(card: DashboardRailCard): boolean {
return Boolean(card.epgLookupKey?.trim());
}
/**
* Cards are grouped by the XMLTV sources their playlist declares, not by the
* playlist itself: two playlists pointing at the same guide ask one question
* and share the answer, while two playlists with DIFFERENT guides never do.
* That separation is the point — a bare `tvg-id` like `ard.de` is not unique
* across imports, so one flat map keyed by lookup key alone would hand one
* playlist's card the other playlist's programme. The any-source flag is part
* of that identity too: a guide-less M3U playlist and a portal playlist both
* resolve against Settings, but only the first may widen the search, so their
* answers for one title are not interchangeable.
*/
export function liveEpgScopeKey(
sourceUrls: readonly string[],
anySourceFallback: boolean
): string {
// JSON, not a separator character: a URL may contain anything, and
// a raw control byte would classify this source file as binary.
return JSON.stringify([
Array.from(new Set(sourceUrls)).sort(),
anySourceFallback,
]);
}
/** Namespaces an answer by the scope it was resolved in. */
export function liveEpgProgramKey(scopeKey: string, lookupKey: string): string {
return JSON.stringify([scopeKey, lookupKey]);
}
export function buildLiveEpgLookupGroups(
cards: readonly DashboardRailCard[],
sourceUrlsForCard: (card: DashboardRailCard) => string[]
): DashboardLiveEpgLookupGroup[] {
const groups = new Map<
string,
{
sourceUrls: string[];
lookupKeys: string[];
anySourceFallback: boolean;
seen: Set<string>;
}
>();
for (const card of cards) {
const key = liveEpgLookupKeyForCard(card);
if (!key || seen.has(key)) continue;
seen.add(key);
keys.push(key);
const lookupKey = liveEpgLookupKeyForCard(card);
if (!lookupKey) continue;
const sourceUrls = sourceUrlsForCard(card);
const anySourceFallback = liveEpgAllowsAnySource(card);
const scopeKey = liveEpgScopeKey(sourceUrls, anySourceFallback);
const group = groups.get(scopeKey) ?? {
sourceUrls: Array.from(new Set(sourceUrls)),
lookupKeys: [],
anySourceFallback,
seen: new Set<string>(),
};
if (!group.seen.has(lookupKey)) {
group.seen.add(lookupKey);
group.lookupKeys.push(lookupKey);
}
groups.set(scopeKey, group);
}
return keys;
return Array.from(groups.entries(), ([scopeKey, group]) => ({
scopeKey,
sourceUrls: group.sourceUrls,
lookupKeys: group.lookupKeys,
anySourceFallback: group.anySourceFallback,
}));
}
type DashboardLiveEpgRailSettings = Pick<
@@ -142,16 +220,23 @@ export function buildLiveEpgCardsForEnabledRails(
];
}
/**
* `epgMap` is keyed by `liveEpgProgramKey`, so a card only ever reads the
* answer resolved in its own source scope.
*/
export function getLiveEpgProgramForCard(
card: DashboardRailCard,
epgMap: ReadonlyMap<string, EpgProgram | null>
epgMap: ReadonlyMap<string, EpgProgram | null>,
scopeKey: string
): EpgProgram | null {
const key = liveEpgLookupKeyForCard(card);
const program = epgMap.get(key);
const program = epgMap.get(liveEpgProgramKey(scopeKey, key));
if (program) {
return program;
}
const titleKey = card.title.trim();
return key !== titleKey ? (epgMap.get(titleKey) ?? null) : null;
return key !== titleKey
? (epgMap.get(liveEpgProgramKey(scopeKey, titleKey)) ?? null)
: null;
}
@@ -40,6 +40,12 @@ export interface DashboardRailCard {
state?: Record<string, unknown>;
actions?: DashboardRailAction[];
epgLookupKey?: string;
/**
* Playlist the live card belongs to. An XMLTV key is only unique inside
* the guide its playlist declares, so the dashboard's EPG lookup is
* grouped and namespaced by that playlist's source scope.
*/
epgPlaylistId?: string;
/**
* Optional EPG enrichment shown by the 'channel' rail layout. Populated
@@ -12,10 +12,12 @@ import {
import {
buildDashboardLiveEpgDetails,
buildLiveEpgCardsForEnabledRails,
buildLiveEpgLookupKeys,
buildLiveEpgLookupGroups,
calcEpgProgress,
formatEpgTimeRange,
getLiveEpgProgramForCard,
liveEpgProgramKey,
liveEpgScopeKey,
} from './dashboard-live-epg.utils';
import {
buildPlaybackPositionReloadKey,
@@ -311,9 +313,15 @@ describe('Live rail helpers', () => {
...overrides,
});
const guideA = ['https://a.example/guide.xml'];
const guideB = ['https://b.example/guide.xml'];
const scopeFor =
(byPlaylist: Record<string, string[]>) => (card: DashboardRailCard) =>
byPlaylist[card.epgPlaylistId ?? ''] ?? [];
it('uses explicit EPG lookup keys before falling back to card titles', () => {
expect(
buildLiveEpgLookupKeys([
const groups = buildLiveEpgLookupGroups(
[
channelCard({
title: 'Das Erste HD',
epgLookupKey: 'ard.de',
@@ -327,8 +335,75 @@ describe('Live rail helpers', () => {
id: 'card-3',
title: 'Fallback News',
}),
])
).toEqual(['ard.de', 'Fallback News']);
],
() => []
);
expect(groups).toHaveLength(2);
// 'ard.de' carries an explicit key, 'Fallback News' falls back to its
// title — different any-source eligibility, so different groups.
expect(groups[0].lookupKeys).toEqual(['ard.de']);
expect(groups[0].anySourceFallback).toBe(true);
expect(groups[1].lookupKeys).toEqual(['Fallback News']);
expect(groups[1].anySourceFallback).toBe(false);
expect(groups[0].sourceUrls).toEqual([]);
});
it('asks each XMLTV source scope separately and shares one lookup between playlists on the same guide', () => {
const groups = buildLiveEpgLookupGroups(
[
channelCard({ epgLookupKey: 'ard.de', epgPlaylistId: 'a' }),
channelCard({
id: 'card-2',
epgLookupKey: 'ard.de',
epgPlaylistId: 'b',
}),
// Same guide as playlist a: one question, shared answer.
channelCard({
id: 'card-3',
epgLookupKey: 'zdf.de',
epgPlaylistId: 'a2',
}),
// No guide of its own: answered from Settings.
channelCard({
id: 'card-4',
epgLookupKey: 'cnn.us',
epgPlaylistId: 'portal',
}),
],
scopeFor({ a: guideA, a2: guideA, b: guideB })
);
expect(
groups.map((group) => [group.sourceUrls, group.lookupKeys])
).toEqual([
[guideA, ['ard.de', 'zdf.de']],
[guideB, ['ard.de']],
[[], ['cnn.us']],
]);
});
it('never widens a portal card to every guide: its key is only a title', () => {
const groups = buildLiveEpgLookupGroups(
[
// M3U: a real XMLTV key from the playlist.
channelCard({ epgLookupKey: 'ard.de', epgPlaylistId: 'm3u' }),
// Xtream/Stalker: no key at all, so the title stands in.
channelCard({
id: 'card-2',
title: 'Das Erste HD',
epgPlaylistId: 'portal',
}),
],
() => []
);
expect(
groups.map((group) => [group.lookupKeys, group.anySourceFallback])
).toEqual([
[['ard.de'], true],
[['Das Erste HD'], false],
]);
});
it('reads EPG programs by explicit lookup key instead of display title', () => {
@@ -343,13 +418,65 @@ describe('Live rail helpers', () => {
getLiveEpgProgramForCard(
card,
new Map<string, EpgProgram | null>([
['Das Erste HD', wrongProgram],
['ard.de', program],
])
[liveEpgProgramKey('', 'Das Erste HD'), wrongProgram],
[liveEpgProgramKey('', 'ard.de'), program],
]),
''
)
).toBe(program);
});
it('never hands a card the programme another playlist resolved for the same XMLTV id', () => {
// `ard.de` is unique inside a guide, not across imports.
const fromGuideA = { title: 'Guide A bulletin' } as EpgProgram;
const fromGuideB = { title: 'Guide B bulletin' } as EpgProgram;
const epgMap = new Map<string, EpgProgram | null>([
[
liveEpgProgramKey(liveEpgScopeKey(guideA, true), 'ard.de'),
fromGuideA,
],
[
liveEpgProgramKey(liveEpgScopeKey(guideB, true), 'ard.de'),
fromGuideB,
],
]);
expect(
getLiveEpgProgramForCard(
channelCard({ epgLookupKey: 'ard.de', epgPlaylistId: 'a' }),
epgMap,
liveEpgScopeKey(guideA, true)
)
).toBe(fromGuideA);
expect(
getLiveEpgProgramForCard(
channelCard({ epgLookupKey: 'ard.de', epgPlaylistId: 'b' }),
epgMap,
liveEpgScopeKey(guideB, true)
)
).toBe(fromGuideB);
// A scope with no answer stays empty instead of borrowing one.
expect(
getLiveEpgProgramForCard(
channelCard({ epgLookupKey: 'ard.de' }),
epgMap,
''
)
).toBeNull();
});
it('treats the same URL set as one scope whatever its order or duplicates', () => {
expect(liveEpgScopeKey(['b', 'a'], true)).toBe(
liveEpgScopeKey(['a', 'b', 'a'], true)
);
expect(liveEpgScopeKey(['a'], true)).not.toBe(
liveEpgScopeKey(['a', 'b'], true)
);
// A portal card and a guide-less M3U card both resolve against
// Settings, but their answers must not be interchangeable.
expect(liveEpgScopeKey([], true)).not.toBe(liveEpgScopeKey([], false));
});
it('uses honest, semantically named title keys for favorite and recent live rails', () => {
expect(liveRailTitleKeyForSource('favorites')).toBe(
'WORKSPACE.DASHBOARD.LIVE_FAVORITES'
@@ -7,11 +7,9 @@ import {
signal,
untracked,
} from '@angular/core';
import { toObservable, toSignal } from '@angular/core/rxjs-interop';
import { interval, map, of, startWith, switchMap } from 'rxjs';
import { EpgService } from '@iptvnator/epg/data-access';
import { toSignal } from '@angular/core/rxjs-interop';
import { interval, map, startWith } from 'rxjs';
import {
type EpgProgram,
isStalkerAccountPlaylist,
isXtreamAccountPlaylist,
normalizeDashboardRailsSettings,
@@ -63,14 +61,8 @@ import type {
import type { PlaylistMeta } from '@iptvnator/shared/interfaces';
import type { DashboardHeroModel } from './dashboard-hero.utils';
import { resolveDashboardHeroArtwork } from './dashboard-hero.utils';
import {
buildDashboardLiveEpgDetails,
buildLiveEpgCardsForEnabledRails,
buildLiveEpgLookupKeys,
getLiveEpgProgramForCard,
LIVE_EPG_TICK_MS,
} from './dashboard-live-epg.utils';
import type { DashboardLiveEpgDetails } from './dashboard-live-epg.utils';
import { buildLiveEpgCardsForEnabledRails } from './dashboard-live-epg.utils';
import { DashboardLiveEpgPresenter } from './dashboard-live-epg.presenter';
import {
buildPlaybackPositionReloadKey,
formatRemainingLabel,
@@ -110,9 +102,11 @@ import type {
host: {
'[class.rails-page-host--empty]': 'ready() && !hasPlaylists()',
},
providers: [DashboardLiveEpgPresenter],
})
export class WorkspaceDashboardRailsComponent {
readonly data = inject(DashboardDataService);
private readonly liveEpg = inject(DashboardLiveEpgPresenter);
private readonly dialog = inject(MatDialog);
private readonly dialogService = inject(DialogService);
private readonly playlistDeleteAction = inject(PlaylistDeleteActionService);
@@ -128,7 +122,6 @@ export class WorkspaceDashboardRailsComponent {
{ initialValue: null }
);
private readonly shellActions = inject(WORKSPACE_SHELL_ACTIONS);
private readonly epgService = inject(EpgService);
private readonly runtime = inject(RuntimeCapabilitiesService);
private readonly settingsStore = inject(SettingsStore);
private readonly heroTmdb = inject(DashboardHeroTmdbService);
@@ -190,7 +183,7 @@ export class WorkspaceDashboardRailsComponent {
const position = this.data.getPlaybackPositionForItem(item);
const liveEpgDetails =
item.type === 'live'
? this.getLiveEpgDetailsForCard(this.heroLiveCard())
? this.liveEpg.detailsFor(this.heroLiveCard())
: null;
const episodeBadge =
item.type === 'series' &&
@@ -267,46 +260,21 @@ export class WorkspaceDashboardRailsComponent {
})
);
// Best-effort EPG lookup keyed by the app-wide M3U XMLTV chain
// (tvg-id -> tvg-name -> name), with the card title as a final fallback.
// Xtream/Stalker live items often have no XMLTV side-channel and will
// simply return null — the card renders without the program row.
private readonly liveChannelLookupKeys = computed(() => {
const heroLiveCard = this.heroLiveCard();
return buildLiveEpgLookupKeys(
buildLiveEpgCardsForEnabledRails(
this.dashboardRails(),
heroLiveCard,
this.liveFavoriteCards(),
this.recentLiveCards()
)
);
});
// The live cards whose rails are enabled; the presenter looks their
// programmes up per XMLTV source scope.
private readonly enabledLiveCards = computed(() =>
buildLiveEpgCardsForEnabledRails(
this.dashboardRails(),
this.heroLiveCard(),
this.liveFavoriteCards(),
this.recentLiveCards()
)
);
private readonly playbackPositionReloadKey = computed(() =>
buildPlaybackPositionReloadKey(this.data.globalRecentVodItems())
);
// Re-fetch on rail change AND on a 30s heartbeat so the progress bar
// catches the boundary between programs without a full page revisit.
private readonly liveEpgPrograms = toSignal(
toObservable(this.liveChannelLookupKeys).pipe(
switchMap((keys) =>
keys.length === 0
? of(new Map<string, EpgProgram | null>())
: interval(LIVE_EPG_TICK_MS).pipe(
startWith(0),
switchMap(() =>
this.epgService.getCurrentProgramsForChannels(
keys
)
)
)
)
),
{ initialValue: new Map<string, EpgProgram | null>() }
);
readonly liveFavoriteCardsEnriched = computed<DashboardRailCard[]>(() =>
this.enrichLiveCards(this.liveFavoriteCards())
);
@@ -418,6 +386,8 @@ export class WorkspaceDashboardRailsComponent {
void this.data.reloadGlobalRecentItems();
void this.data.reloadGlobalFavorites();
this.liveEpg.connect(this.enabledLiveCards);
// Refresh when Xtream playlist count changes so a newly added provider
// populates the rail without a manual dashboard reload. The Xtream
// recently-added query can be the slowest dashboard worker request on
@@ -605,7 +575,7 @@ export class WorkspaceDashboardRailsComponent {
cards: readonly DashboardRailCard[]
): DashboardRailCard[] {
return cards.map((card) => {
const details = this.getLiveEpgDetailsForCard(card);
const details = this.liveEpg.detailsFor(card);
if (!details) {
return card;
}
@@ -613,22 +583,6 @@ export class WorkspaceDashboardRailsComponent {
});
}
private getLiveEpgDetailsForCard(
card: DashboardRailCard | null
): DashboardLiveEpgDetails | null {
if (!card) {
return null;
}
const program = getLiveEpgProgramForCard(card, this.liveEpgPrograms());
// Recompute the now-window each tick so progress moves between
// 30s ticks even if the program identity is unchanged.
return buildDashboardLiveEpgDetails(
program,
Date.now(),
this.settingsStore.resolvedEpgOffsetMinutes()
);
}
private buildNonLiveSeeAllState(
cards: readonly DashboardRailCard[]
): Record<string, unknown> {
@@ -671,6 +625,7 @@ export class WorkspaceDashboardRailsComponent {
icon: this.typeIcon(item.type),
contentType: item.type,
epgLookupKey: item.epg_lookup_key,
epgPlaylistId: item.playlist_id,
link: this.data.getRecentItemLink(item),
// Default click is detail-only for every card — an in-progress
// series no longer auto-plays on click (issue #1441); resuming
@@ -707,6 +662,7 @@ export class WorkspaceDashboardRailsComponent {
icon: this.typeIcon(item.type),
contentType: item.type,
epgLookupKey: item.epg_lookup_key,
epgPlaylistId: item.playlist_id,
link: this.data.getGlobalFavoriteLink(item),
state: this.data.getGlobalFavoriteNavigationState(item),
};
-1
View File
@@ -23,7 +23,6 @@ export const maxLinesBaseline = [
'apps/web/src/app/services/electron.service.ts',
'apps/web/src/app/services/pwa.service.ts',
'apps/xtream-mock-server/src/app/generators/marketing.generator.ts',
'libs/epg/data-access/src/lib/epg.service.ts',
'libs/m3u-state/src/lib/effects.ts',
'libs/playlist/m3u/feature-player/src/lib/video-player/video-player.component.ts',
'libs/playlist/shared/ui/src/lib/playlist-switcher/playlist-switcher.component.ts',