diff --git a/.changes/dashboard-live-rails-epg-any-source.md b/.changes/dashboard-live-rails-epg-any-source.md new file mode 100644 index 000000000..e0dce1935 --- /dev/null +++ b/.changes/dashboard-live-rails-epg-any-source.md @@ -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. diff --git a/CLAUDE.md b/CLAUDE.md index 908f3a906..2a2c2b99c 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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/` diff --git a/apps/electron-backend-e2e/src/epg.e2e.ts b/apps/electron-backend-e2e/src/epg.e2e.ts index 58525a4af..a1453e3bb 100644 --- a/apps/electron-backend-e2e/src/epg.e2e.ts +++ b/apps/electron-backend-e2e/src/epg.e2e.ts @@ -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, }) => { diff --git a/docs/architecture/m3u-playlist-module.md b/docs/architecture/m3u-playlist-module.md index ca85817f1..75d698dc9 100644 --- a/docs/architecture/m3u-playlist-module.md +++ b/docs/architecture/m3u-playlist-module.md @@ -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 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-current-programs.lookup.ts b/libs/epg/data-access/src/lib/epg-current-programs.lookup.ts new file mode 100644 index 000000000..38c296596 --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-current-programs.lookup.ts @@ -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> { + 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> | 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> { + 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 + ): Observable> { + 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> { + const resultMap = new Map(); + 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; + }) + ); + } +} diff --git a/libs/epg/data-access/src/lib/epg-lookup-context.ts b/libs/epg/data-access/src/lib/epg-lookup-context.ts new file mode 100644 index 000000000..c32c8e245 --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-lookup-context.ts @@ -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(): MonoTypeOperatorFunction; + /** 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[]; +} diff --git a/libs/epg/data-access/src/lib/epg-lookup-normalization.util.ts b/libs/epg/data-access/src/lib/epg-lookup-normalization.util.ts new file mode 100644 index 000000000..41eb39a26 --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-lookup-normalization.util.ts @@ -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 ?? []); +} 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-runtime-bridge.service.ts b/libs/epg/data-access/src/lib/epg-runtime-bridge.service.ts index 9b4193bfd..a73e49dc9 100644 --- a/libs/epg/data-access/src/lib/epg-runtime-bridge.service.ts +++ b/libs/epg/data-access/src/lib/epg-runtime-bridge.service.ts @@ -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< diff --git a/libs/epg/data-access/src/lib/epg-scoped-batch.lookup.ts b/libs/epg/data-access/src/lib/epg-scoped-batch.lookup.ts new file mode 100644 index 000000000..9080f9036 --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-scoped-batch.lookup.ts @@ -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> { + const normalizedChannelIds = normalizeLookupChannelIds(channelIds); + if (normalizedChannelIds.length === 0) { + return of(new Map()); + } + + const resultMap = new Map(); + 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> { + 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(); + 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 { + return new Map(channelIds.map((channelId) => [channelId, null])); + } +} diff --git a/libs/epg/data-access/src/lib/epg-single-program.lookup.ts b/libs/epg/data-access/src/lib/epg-single-program.lookup.ts new file mode 100644 index 000000000..036db380c --- /dev/null +++ b/libs/epg/data-access/src/lib/epg-single-program.lookup.ts @@ -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 { + 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 { + // 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 { + 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 { + if (sourceUrls.length === 0) { + return of(null); + } + + return this.scoped(channelId, sourceUrls, []); + } +} diff --git a/libs/epg/data-access/src/lib/epg.service.spec.ts b/libs/epg/data-access/src/lib/epg.service.spec.ts index 69d4b29e0..4a0fbcb12 100644 --- a/libs/epg/data-access/src/lib/epg.service.spec.ts +++ b/libs/epg/data-access/src/lib/epg.service.spec.ts @@ -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; diff --git a/libs/epg/data-access/src/lib/epg.service.ts b/libs/epg/data-access/src/lib/epg.service.ts index d76ae9bf6..525e308da 100644 --- a/libs/epg/data-access/src/lib/epg.service.ts +++ b/libs/epg/data-access/src/lib/epg.service.ts @@ -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(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() + ); + + /** 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(); - 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> { - const normalizedChannelIds = this.normalizeChannelIds(channelIds); - if (normalizedChannelIds.length === 0) { - return of(new Map()); - } - - const resultMap = new Map(); - 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> { - const nowMs = this.epgClockMs(); - return from( - this.epgBridge.getCurrentProgramsBatch(channelIds, { - sourceUrls, - nowMs, - }) - ).pipe( - this.sourceSettings.guard(), - timeout(5000), - switchMap((scopedResult) => { - const resultMap = new Map(); - 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> { - 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); - - 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 { - if (sourceUrls.length === 0) { - return of(null); - } - - return this.getScopedCurrentProgramForChannel( - channelId, - sourceUrls, - [] - ); - } - - private createNullProgramMap( - channelIds: string[] - ): Map { - 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(); } } diff --git a/libs/workspace/dashboard/feature/src/index.ts b/libs/workspace/dashboard/feature/src/index.ts index 363b10a9a..767acd00c 100644 --- a/libs/workspace/dashboard/feature/src/index.ts +++ b/libs/workspace/dashboard/feature/src/index.ts @@ -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, diff --git a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.presenter.spec.ts b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.presenter.spec.ts new file mode 100644 index 000000000..fa4f72d3a --- /dev/null +++ b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-live-epg.presenter.spec.ts @@ -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 => + ({ + 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>; + + 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([ + 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([ + [ + '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([ + [ + '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([ + ['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([ + ['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(); + }); +}); 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 99d22d4aa..98526cabb 100644 --- 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 @@ -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(); - 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; + } + >(); + 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(), + }; + 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 + epgMap: ReadonlyMap, + 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; } diff --git a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.ts b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.ts index e3eaaa589..60d01a07a 100644 --- a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.ts +++ b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.ts @@ -40,6 +40,12 @@ export interface DashboardRailCard { state?: Record; 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 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 3d9d265b5..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 @@ -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) => (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([ - ['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([ + [ + 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' 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 1b851edf3..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,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()) - : interval(LIVE_EPG_TICK_MS).pipe( - startWith(0), - switchMap(() => - this.epgService.getCurrentProgramsForChannels( - keys - ) - ) - ) - ) - ), - { initialValue: new Map() } - ); - readonly liveFavoriteCardsEnriched = computed(() => 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 { @@ -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), }; diff --git a/tools/eslint/max-lines-baseline.mjs b/tools/eslint/max-lines-baseline.mjs index 26f1032a9..1afe8ffae 100644 --- a/tools/eslint/max-lines-baseline.mjs +++ b/tools/eslint/max-lines-baseline.mjs @@ -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',