diff --git a/CLAUDE.md b/CLAUDE.md index dbb4ebbb1..dd6015467 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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/` diff --git a/docs/architecture/m3u-playlist-module.md b/docs/architecture/m3u-playlist-module.md index 5522ece47..6f8898880 100644 --- a/docs/architecture/m3u-playlist-module.md +++ b/docs/architecture/m3u-playlist-module.md @@ -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 diff --git a/libs/epg/data-access/src/lib/epg-channel-metadata.lookup.ts b/libs/epg/data-access/src/lib/epg-channel-metadata.lookup.ts new file mode 100644 index 000000000..648a6a6b0 --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-channel-metadata.lookup.ts @@ -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: () => MonoTypeOperatorFunction, + private readonly globalSourceUrls: (excluding?: string[]) => string[] + ) {} + + forChannels( + normalizedChannelIds: string[], + sourceUrls: string[] + ): Observable> { + 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>(), + 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((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> { + return from( + this.epgBridge.getChannelMetadata( + channelIds, + sourceUrls.length > 0 ? { sourceUrls } : undefined + ) + ).pipe( + this.guard | null>(), + map((metadataByChannelId) => { + return new Map( + channelIds.map((channelId) => [ + channelId, + metadataByChannelId?.[channelId] ?? null, + ]) + ); + }), + catchError((err) => { + console.error('EPG get channel metadata error:', err); + return of(new Map()); + }) + ); + } +} diff --git a/libs/epg/data-access/src/lib/epg-program-cache.ts b/libs/epg/data-access/src/lib/epg-program-cache.ts new file mode 100644 index 000000000..51f851a68 --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-program-cache.ts @@ -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(); + private readonly inFlight = new Map< + string, + Observable + >(); + private readonly inFlightBatches = new Map< + string, + Observable> + >(); + + constructor( + private readonly offsetMinutes: () => number, + private readonly guard: () => MonoTypeOperatorFunction + ) {} + + /** 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 + ): Observable { + 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(), + 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> | undefined { + return this.inFlightBatches.get(batchCacheKey); + } + + registerBatch( + batchCacheKey: string, + request$: Observable> + ): void { + this.inFlightBatches.set(batchCacheKey, request$); + } + + /** Only the request that owns the slot may release it. */ + releaseBatch( + batchCacheKey: string, + request$: Observable> + ): void { + if (this.inFlightBatches.get(batchCacheKey) === request$) { + this.inFlightBatches.delete(batchCacheKey); + } + } + + clear(): void { + this.programs.clear(); + this.inFlight.clear(); + this.inFlightBatches.clear(); + } +} diff --git a/libs/epg/data-access/src/lib/epg.service.ts b/libs/epg/data-access/src/lib/epg.service.ts index 95a40113b..d99b914c0 100644 --- a/libs/epg/data-access/src/lib/epg.service.ts +++ b/libs/epg/data-access/src/lib/epg.service.ts @@ -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(false); private currentEpgPrograms = new BehaviorSubject([]); - // Cache for channel programs with 60-second TTL - private programCache = new Map(); - private fetchingCurrentPrograms = new Map< - string, - Observable - >(); - private fetchingCurrentProgramBatches = new Map< - string, - Observable> - >(); - 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 { // 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> { - return from( - this.epgBridge.getChannelMetadata( - channelIds, - sourceUrls.length > 0 ? { sourceUrls } : undefined - ) - ).pipe( - this.sourceSettings.guard(), - map((metadataByChannelId) => { - return new Map( - channelIds.map((channelId) => [ - channelId, - metadataByChannelId?.[channelId] ?? null, - ]) - ); - }), - catchError((err) => { - console.error('EPG get channel metadata error:', err); - return of(new Map()); - }) - ); - } - - /** - * 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 - ): Observable { - 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 { - 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(); } } diff --git a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.presenter.ts b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.presenter.ts new file mode 100644 index 000000000..803280606 --- /dev/null +++ b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.presenter.ts @@ -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; +}; + +const emptyAnswer = (scopeKey: string): ScopeAnswer => ({ + scopeKey, + programs: new Map(), +}); + +/** + * 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 | 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(); + 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()) + : interval(LIVE_EPG_TICK_MS).pipe( + startWith(0), + switchMap(() => + forkJoin( + groups.map((group) => this.askScope(group)) + ).pipe(map((answers) => mergeAnswers(answers))) + ) + ) + ) + ), + { initialValue: new Map() } + ); + + /** The live cards whose rails are enabled, hero included. */ + connect(cards: Signal): 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 { + const merged = new Map(); + for (const answer of answers) { + for (const [lookupKey, program] of answer.programs) { + merged.set(liveEpgProgramKey(answer.scopeKey, lookupKey), program); + } + } + return merged; +} diff --git a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.utils.ts b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.utils.ts index ed0f41a45..80cce1ede 100644 Binary files a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.utils.ts and b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.utils.ts differ diff --git a/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.spec.ts b/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.spec.ts index 4417505f1..74feeb38a 100644 --- a/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.spec.ts +++ b/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.spec.ts @@ -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([ - [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', () => { diff --git a/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.ts b/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.ts index 549d245d8..118fe9ffc 100644 --- a/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.ts +++ b/libs/workspace/dashboard/feature/src/lib/rails/workspace-dashboard-rails.component.ts @@ -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(); - 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()) - : 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() } - ); - readonly liveFavoriteCardsEnriched = computed(() => 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 {