refactor(epg): shrink the two oversized files this PR grew, keep portal cards in scope

Two review follow-ups.

Greptile: the repo requires refactoring an already-oversized production file
before adding to it, and this PR had grown both files it touched. Extracted,
without behaviour change:

- `EpgProgramCache` — the 60 s "currently airing" memory, its in-flight
  request slots and the scope/offset cache keys.
- `EpgChannelMetadataLookup` — the icon/display-name lookup with the same
  playlist-first, Settings-fallback ladder.
- `DashboardLiveEpgPresenter` — the rails' XMLTV lookup and card enrichment,
  component-provided so the 30 s heartbeat dies with the page.

`epg.service.ts` 901 -> 727 and the rails component 1000 -> 842, both now
smaller than before this PR; the three new files are well under the 400-line
limit, and the generated max-lines baseline is unchanged.

Codex P2: the any-source retry now applies only to cards that carry a real
XMLTV key. An Xtream or Stalker card has none, so its lookup key is just its
display title, and searching every imported 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 flag is
part of the scope identity, so a guide-less M3U card and a portal card never
share an answer.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
4grayandClaude Opus 5 committed 2026-09-20 20:27:32 +02:00
1 parent 6579151dcf
commit 2ef0afbeca
9 files changed
+565 -402

No files matched your search

+1 -1
View File
@@ -1606,7 +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. The channel list keeps the strict scope. Contract: `docs/architecture/m3u-playlist-module.md` (the "Scoped lookups" bullet under playlist-scoped URLs)
- 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/`
+8 -1
View File
@@ -998,7 +998,14 @@ These URLs are playlist-scoped by default:
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. The channel list keeps the strict scope. Single-channel current
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,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();
}
}
+38 -212
View File
@@ -23,15 +23,10 @@ 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 { 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,18 @@ 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()
);
/** Display offset every "currently airing" decision in here is made with. */
private epgOffsetMinutes(): number {
@@ -213,10 +209,10 @@ export class EpgService {
channelId: string
): Observable<EpgProgram | null> {
// Check cache first
const cacheKey = this.createProgramCacheKey(channelId);
const cacheKey = this.programCache.keyFor(channelId);
// Fetch from backend
return this.getCachedOrFetchCurrentProgram(cacheKey, () =>
return this.programCache.getOrFetch(cacheKey, () =>
from(this.epgBridge.getChannelPrograms(channelId)).pipe(
this.sourceSettings.guard(),
map((programs) => normalizeEpgPrograms(programs ?? [])),
@@ -406,12 +402,8 @@ export class EpgService {
// 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
) {
const cached = this.programCache.getFresh(channelId, now);
if (cached && cached.offsetMinutes === offsetMinutes) {
resultMap.set(channelId, cached.program);
} else {
channelsToFetch.push(channelId);
@@ -439,11 +431,12 @@ export class EpgService {
channelsToFetch.forEach((channelId) => {
const program = batchResult?.[channelId] ?? null;
resultMap.set(channelId, program);
this.programCache.set(channelId, {
this.programCache.set(
channelId,
program,
timestamp: cacheTimestamp,
offsetMinutes,
});
cacheTimestamp
);
});
return resultMap;
}),
@@ -494,8 +487,8 @@ export class EpgService {
const channelsToFetch: string[] = [];
normalizedChannelIds.forEach((channelId) => {
const cached = this.getCachedProgram(
this.createProgramCacheKey(channelId, sourceUrls)
const cached = this.programCache.get(
this.programCache.keyFor(channelId, sourceUrls)
);
if (cached) {
resultMap.set(channelId, cached.program);
@@ -508,13 +501,12 @@ export class EpgService {
return of(resultMap);
}
const batchCacheKey = this.createProgramBatchCacheKey(
const batchCacheKey = this.programCache.batchKeyFor(
channelsToFetch,
sourceUrls,
fallbackSourceUrls
);
const existingRequest =
this.fetchingCurrentProgramBatches.get(batchCacheKey);
const existingRequest = this.programCache.batchInFlight(batchCacheKey);
// Tag entries with the offset the verdict was computed with, not the
// one current when the response lands.
const offsetMinutes = this.epgOffsetMinutes();
@@ -530,30 +522,21 @@ export class EpgService {
const cacheTimestamp = Date.now();
channelsToFetch.forEach((channelId) => {
this.programCache.set(
this.createProgramCacheKey(channelId, sourceUrls),
{
program: fetchedMap.get(channelId) ?? null,
timestamp: cacheTimestamp,
offsetMinutes,
}
this.programCache.keyFor(channelId, sourceUrls),
fetchedMap.get(channelId) ?? null,
offsetMinutes,
cacheTimestamp
);
});
}),
finalize(() => {
if (
this.fetchingCurrentProgramBatches.get(
batchCacheKey
) === request$
)
this.fetchingCurrentProgramBatches.delete(
batchCacheKey
);
}),
finalize(() =>
this.programCache.releaseBatch(batchCacheKey, request$)
),
shareReplay({ bufferSize: 1, refCount: false })
);
if (!existingRequest) {
this.fetchingCurrentProgramBatches.set(batchCacheKey, request$);
this.programCache.registerBatch(batchCacheKey, request$);
}
return request$.pipe(
@@ -647,59 +630,13 @@ export class EpgService {
}
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(
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);
})
);
})
this.normalizeSourceUrls(options)
);
}
@@ -717,123 +654,14 @@ export class EpgService {
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);
const cacheKey = this.programCache.keyFor(channelId, sourceUrls);
return this.getCachedOrFetchCurrentProgram(cacheKey, () =>
return this.programCache.getOrFetch(cacheKey, () =>
from(
this.epgBridge.getChannelPrograms(channelId, { sourceUrls })
).pipe(
@@ -895,7 +723,5 @@ export class EpgService {
*/
clearCache(): void {
this.programCache.clear();
this.fetchingCurrentPrograms.clear();
this.fetchingCurrentProgramBatches.clear();
}
}
@@ -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;
}
@@ -339,10 +339,14 @@ describe('Live rail helpers', () => {
() => []
);
expect(groups).toHaveLength(1);
expect(groups[0].lookupKeys).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([]);
expect(groups[0].scopeKey).toBe('');
});
it('asks each XMLTV source scope separately and shares one lookup between playlists on the same guide', () => {
@@ -379,6 +383,29 @@ describe('Live rail helpers', () => {
]);
});
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', () => {
const program = { title: 'Tagesschau' } as EpgProgram;
const wrongProgram = { title: 'Wrong channel' } as EpgProgram;
@@ -404,22 +431,28 @@ describe('Live rail helpers', () => {
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), 'ard.de'), fromGuideA],
[liveEpgProgramKey(liveEpgScopeKey(guideB), 'ard.de'), fromGuideB],
[
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)
liveEpgScopeKey(guideA, true)
)
).toBe(fromGuideA);
expect(
getLiveEpgProgramForCard(
channelCard({ epgLookupKey: 'ard.de', epgPlaylistId: 'b' }),
epgMap,
liveEpgScopeKey(guideB)
liveEpgScopeKey(guideB, true)
)
).toBe(fromGuideB);
// A scope with no answer stays empty instead of borrowing one.
@@ -433,10 +466,15 @@ describe('Live rail helpers', () => {
});
it('treats the same URL set as one scope whatever its order or duplicates', () => {
expect(liveEpgScopeKey(['b', 'a'])).toBe(
liveEpgScopeKey(['a', 'b', 'a'])
expect(liveEpgScopeKey(['b', 'a'], true)).toBe(
liveEpgScopeKey(['a', 'b', 'a'], true)
);
expect(liveEpgScopeKey(['a'])).not.toBe(liveEpgScopeKey(['a', 'b']));
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', () => {
@@ -7,20 +7,9 @@ import {
signal,
untracked,
} from '@angular/core';
import { toObservable, toSignal } from '@angular/core/rxjs-interop';
import { toSignal } from '@angular/core/rxjs-interop';
import { interval, map, startWith } from 'rxjs';
import {
catchError,
defaultIfEmpty,
forkJoin,
interval,
map,
of,
startWith,
switchMap,
} from 'rxjs';
import { EpgService } from '@iptvnator/epg/data-access';
import {
type EpgProgram,
isStalkerAccountPlaylist,
isXtreamAccountPlaylist,
normalizeDashboardRailsSettings,
@@ -31,7 +20,6 @@ import { MatIcon } from '@angular/material/icon';
import { MatSnackBar } from '@angular/material/snack-bar';
import { Router, RouterLink } from '@angular/router';
import { isPortalPlaybackWatched } from '@iptvnator/portal/shared/util';
import { normalizeEpgUrls } from '@iptvnator/shared/m3u-utils';
import { Store } from '@ngrx/store';
import { TranslatePipe, TranslateService } from '@ngx-translate/core';
import {
@@ -73,16 +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,
buildLiveEpgLookupGroups,
getLiveEpgProgramForCard,
liveEpgProgramKey,
liveEpgScopeKey,
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,
@@ -122,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);
@@ -140,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);
@@ -202,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' &&
@@ -279,132 +260,21 @@ export class WorkspaceDashboardRailsComponent {
})
);
// 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 liveEpgSourceUrlsByPlaylist = 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;
});
// 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.
//
// 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.
private readonly liveEpgLookupGroups = computed(() => {
const heroLiveCard = this.heroLiveCard();
return buildLiveEpgLookupGroups(
buildLiveEpgCardsForEnabledRails(
this.dashboardRails(),
heroLiveCard,
this.liveFavoriteCards(),
this.recentLiveCards()
),
(card) => this.liveEpgSourceUrlsForCard(card)
);
});
// 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.liveEpgLookupGroups).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.epgService
.getCurrentProgramsForChannels(
group.lookupKeys,
{
sourceUrls: group.sourceUrls,
anySourceFallback: true,
}
)
.pipe(
map((programs) => ({
scopeKey: group.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
// empty.
defaultIfEmpty({
scopeKey: group.scopeKey,
programs: new Map<
string,
EpgProgram | null
>(),
}),
catchError(() =>
of({
scopeKey: group.scopeKey,
programs: new Map<
string,
EpgProgram | null
>(),
})
)
)
)
).pipe(
map((answers) => {
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;
})
)
)
)
)
),
{ initialValue: new Map<string, EpgProgram | null>() }
);
readonly liveFavoriteCardsEnriched = computed<DashboardRailCard[]>(() =>
this.enrichLiveCards(this.liveFavoriteCards())
);
@@ -516,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
@@ -699,21 +571,11 @@ export class WorkspaceDashboardRailsComponent {
});
}
/** The XMLTV scope a live card's programme must be resolved in. */
private liveEpgSourceUrlsForCard(card: DashboardRailCard): string[] {
const byPlaylistId = this.liveEpgSourceUrlsByPlaylist();
return (
(card.epgPlaylistId
? byPlaylistId.get(card.epgPlaylistId)
: undefined) ?? []
);
}
private enrichLiveCards(
cards: readonly DashboardRailCard[]
): DashboardRailCard[] {
return cards.map((card) => {
const details = this.getLiveEpgDetailsForCard(card);
const details = this.liveEpg.detailsFor(card);
if (!details) {
return card;
}
@@ -721,26 +583,6 @@ export class WorkspaceDashboardRailsComponent {
});
}
private getLiveEpgDetailsForCard(
card: DashboardRailCard | null
): DashboardLiveEpgDetails | null {
if (!card) {
return null;
}
const program = getLiveEpgProgramForCard(
card,
this.liveEpgPrograms(),
liveEpgScopeKey(this.liveEpgSourceUrlsForCard(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()
);
}
private buildNonLiveSeeAllState(
cards: readonly DashboardRailCard[]
): Record<string, unknown> {