From 9fc487dcb3d429ec85f7e7b175deae66368cb942 Mon Sep 17 00:00:00 2001 From: 4gray Date: Sat, 19 Sep 2026 18:21:43 +0200 Subject: [PATCH] feat(dashboard): portal EPG on the live rails, loaded lazily per visible card Xtream and Stalker live cards on the dashboard showed only the LIVE chip: they carry no XMLTV key and the rails never asked the portal. They now get their programme from the portal one card at a time, only once the card is on screen, with a bounded queue (2 in flight, 200 ms apart, 60 s TTL, 30 s failure cooldown) and every answer published as it lands, so the page never waits and a slow portal delays no other card. The rail reports visible cards through an IntersectionObserver; the presenter turns them plus the pinned hero into the wanted set; the service resolves each through StreamResolverService.loadEpgForItems, the same call the See all pages make. A shimmer placeholder shows only before a card's first answer. Follow-up to #1637. Co-Authored-By: Claude Fable 5.1 --- .changes/dashboard-portal-live-epg.md | 9 + CLAUDE.md | 1 + .../src/dashboard-activation.e2e.ts | 45 +++ docs/architecture/workspace-dashboard.md | 37 ++- .../dashboard/data-access/src/index.ts | 2 + .../dashboard-portal-live-epg.service.spec.ts | 243 ++++++++++++++++ .../lib/dashboard-portal-live-epg.service.ts | 263 ++++++++++++++++++ .../dashboard-portal-live-epg.util.spec.ts | 160 +++++++++++ .../src/lib/dashboard-portal-live-epg.util.ts | 136 +++++++++ .../lib/rails/dashboard-live-epg.presenter.ts | 60 +++- ...ashboard-portal-live-epg.presenter.spec.ts | 152 ++++++++++ .../dashboard-portal-live-epg.presenter.ts | 129 +++++++++ .../lib/rails/dashboard-rail.component.html | 20 +- .../lib/rails/dashboard-rail.component.scss | 40 ++- .../rails/dashboard-rail.component.spec.ts | 208 ++++++++++++++ .../src/lib/rails/dashboard-rail.component.ts | 98 +++++++ .../workspace-dashboard-rails.component.html | 21 +- .../workspace-dashboard-rails.component.ts | 60 +++- 18 files changed, 1653 insertions(+), 31 deletions(-) create mode 100644 .changes/dashboard-portal-live-epg.md create mode 100644 libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.spec.ts create mode 100644 libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.ts create mode 100644 libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.spec.ts create mode 100644 libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.ts create mode 100644 libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.spec.ts create mode 100644 libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.ts diff --git a/.changes/dashboard-portal-live-epg.md b/.changes/dashboard-portal-live-epg.md new file mode 100644 index 000000000..0ae5d21f3 --- /dev/null +++ b/.changes/dashboard-portal-live-epg.md @@ -0,0 +1,9 @@ +--- +type: feature +area: dashboard +--- + +Favourite and recently watched Xtream and Stalker channels on the dashboard +now show what is on air, with the programme's time and progress. Each card +asks its portal only once it scrolls into view and fills in on its own, so a +slow portal never holds up the page or the other cards. diff --git a/CLAUDE.md b/CLAUDE.md index 2a2c2b99c..e4c9f8722 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1607,6 +1607,7 @@ stream_id`); it drops `series_id`/`movie_id`, so the builder pins the - 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) +- Dashboard live rails, portal side: an Xtream or Stalker card carries no XMLTV key, so its programme comes from the portal instead, lazily and per card — `lib-dashboard-rail` reports the cards inside its viewport (`visibleCardsChanged`, `IntersectionObserver` on the track), `DashboardLiveEpgPresenter` forwards them plus the pinned hero to `DashboardPortalLiveEpgPresenter`, and `DashboardPortalLiveEpgService` (data-access, root) runs the bounded queue (2 in flight / 200 ms, 60 s TTL for a programme, 30 s for an empty answer) through `StreamResolverService.loadEpgForItems`, publishing each answer as it lands. A completion captures the display offset AND `EpgSourceSettingsService.revision()` and requeues itself when either moved; the presenter hands its wanted set back on destroy; desktop only, since the resolver is gated on `supportsProgramLookup`. The one shared facade the rails talk to is `DashboardLiveEpgPresenter`: a portal answer wins, an XMLTV title match is the fallback, and a shimmer placeholder shows only before a card's first portal answer. Contract: `docs/architecture/workspace-dashboard.md` (Data Flow, item 3) - 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/dashboard-activation.e2e.ts b/apps/electron-backend-e2e/src/dashboard-activation.e2e.ts index e78ad2d26..86a19d0d0 100644 --- a/apps/electron-backend-e2e/src/dashboard-activation.e2e.ts +++ b/apps/electron-backend-e2e/src/dashboard-activation.e2e.ts @@ -18,6 +18,7 @@ import { waitForXtreamWorkspaceReady, } from './electron-test-fixtures'; import { + fetchXtreamEpgFixture, fetchXtreamLiveFixture, fetchXtreamSeriesFixture, fetchXtreamVodFixture, @@ -47,6 +48,33 @@ test.describe('Dashboard Activation', () => { liveFixture.items, getXtreamTitle ); + // The mock's guide for the channel favourited below (its first live + // stream). The programme on air is read off the full guide at + // assertion time: the mock cuts its slots from the second the guide + // was generated, so the "current" listing of a fixture fetched in + // that same second is the slot that ended just then. + const epgFixture = await fetchXtreamEpgFixture( + request, + xtreamCredentials + ); + expect(getXtreamTitle(epgFixture.stream)).toBe(liveTitle); + const liveNowTitle = () => { + const nowSeconds = Math.floor(Date.now() / 1000); + const index = epgFixture.fullEpg.findIndex( + (listing) => + listing.startTimestamp <= nowSeconds && + nowSeconds < listing.stopTimestamp + ); + expect(index).toBeGreaterThanOrEqual(0); + // Tolerate a slot boundary passing between the app's answer and + // this assertion. + const titles = epgFixture.fullEpg + .slice(index, index + 2) + .map((listing) => + listing.title.replace(/[.*+?^${}()|[\]\\]/g, '\\$&') + ); + return new RegExp(titles.join('|')); + }; const app = await launchElectronApp(dataDir); try { @@ -139,6 +167,15 @@ test.describe('Dashboard Activation', () => { app.mainWindow, 'dashboard-live-favorites-rail' ); + // An Xtream card has no XMLTV key: its "now on air" line comes + // from the portal, asked for lazily once the card is on screen. + await expect( + dashboardRailCardByTitle( + app.mainWindow, + 'dashboard-live-favorites-rail', + liveTitle + ).locator('.rail__channel-now') + ).toContainText(liveNowTitle(), { timeout: 30000 }); await dashboardRailCardByTitle( app.mainWindow, 'dashboard-live-favorites-rail', @@ -167,6 +204,14 @@ test.describe('Dashboard Activation', () => { app.mainWindow, 'dashboard-recent-live-rail' ); + // Same channel, same key: the recent card shares the answer. + await expect( + dashboardRailCardByTitle( + app.mainWindow, + 'dashboard-recent-live-rail', + liveTitle + ).locator('.rail__channel-now') + ).toContainText(liveNowTitle(), { timeout: 30000 }); await dashboardRailCardByTitle( app.mainWindow, 'dashboard-recent-live-rail', diff --git a/docs/architecture/workspace-dashboard.md b/docs/architecture/workspace-dashboard.md index 81266e7e8..daf113338 100644 --- a/docs/architecture/workspace-dashboard.md +++ b/docs/architecture/workspace-dashboard.md @@ -129,7 +129,42 @@ Render rules: `dashboard-recent-live-rail`); there is no fallback from one to the other. M3U cards carry an `epg_lookup_key` using the app-wide XMLTV fallback order (`tvg-id` -> `tvg-name` -> channel name); EPG enrichment - must use that key before falling back to the card title. + must use that key before falling back to the card title. That XMLTV + lookup is one batched local query for every live card and re-runs on + the 30 s tick. + Xtream and Stalker cards have no XMLTV key of their own; their "now on + air" line comes from the portal, **lazily and per card**: + - `buildDashboardPortalLiveEpgEntry` (dashboard data-access) turns a + live `PortalActivityItem` into the `UnifiedCollectionItem` the + collection pages hand `StreamResolverService.loadEpgForItems`, keyed + by the collection uid — favourites and recent rows of one channel + share the answer. Radio rows and rows without a usable provider id + get no entry. Cards carry that key as `liveEpgSourceKey`. + - `lib-dashboard-rail` reports the cards inside its track viewport + (plus ~one card of `rootMargin`) through `visibleCardsChanged`, from + an `IntersectionObserver` rooted at the track; without the API every + card counts as visible. Cards that leave the list are reported gone + at once. + - `DashboardPortalLiveEpgPresenter` (component-provided) unions the + visible keys of both rails with the pinned hero key and calls + `DashboardPortalLiveEpgService.sync()` with exactly those entries — + on every change, on the 30 s tick, and on a display-offset change. + - `DashboardPortalLiveEpgService` (root) owns the queue: at most two + requests in flight, 200 ms between starts (the numbers + `EpgQueueService` proved against real panels), one card per request, + each answer published the moment it lands in `programs`, so the page + never waits and a slow portal delays no other card. Only wanted keys + are dequeued, so a card scrolled past before its turn is never + requested. Answers live 60 s; a failed portal is left alone for 30 s; + a programme that ended is asked again, but not within 30 s of the + last answer (a portal may keep returning the stale row). Every + answer is "at the provider clock", so a changed display offset or a + changed XMLTV source set drops them all. + - `enrichLiveCards` prefers the portal answer, falls back to the XMLTV + title match when the portal said "nothing on air", and marks a card + `nowPlayingState: 'pending'` only before its FIRST answer — the + channel layout then shows a shimmer placeholder in the programme + slot; a refresh keeps the previous answer on screen. 4. `xtreamRecentlyAddedCards` — maps `xtreamRecentlyAddedItems()` to rail cards. Aggregates newly added VOD and series across *all* Xtream playlists via `DashboardDataService.reloadXtreamRecentlyAddedItems()`, diff --git a/libs/workspace/dashboard/data-access/src/index.ts b/libs/workspace/dashboard/data-access/src/index.ts index 29b5d17ab..7d1146a21 100644 --- a/libs/workspace/dashboard/data-access/src/index.ts +++ b/libs/workspace/dashboard/data-access/src/index.ts @@ -1,4 +1,6 @@ export * from './lib/dashboard-data.service'; +export * from './lib/dashboard-portal-live-epg.service'; +export * from './lib/dashboard-portal-live-epg.util'; export * from './lib/dashboard-recommendations.service'; export * from './lib/dashboard-recommendations.util'; export * from './lib/dashboard-source-expiry.service'; diff --git a/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.spec.ts b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.spec.ts new file mode 100644 index 000000000..515c38a62 --- /dev/null +++ b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.spec.ts @@ -0,0 +1,243 @@ +import { TestBed } from '@angular/core/testing'; +import { Subject } from 'rxjs'; +import type { EpgProgram } from '@iptvnator/shared/interfaces'; +import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services'; +import { StreamResolverService } from '@iptvnator/portal/shared/data-access'; +import { + DASHBOARD_PORTAL_LIVE_EPG_TIMING, + DashboardPortalLiveEpgService, +} from './dashboard-portal-live-epg.service'; +import type { DashboardPortalLiveEpgEntry } from './dashboard-portal-live-epg.util'; + +interface Deferred { + readonly key: string; + resolve: (program: EpgProgram | null) => void; + reject: (error: unknown) => void; +} + +const entry = (id: number): DashboardPortalLiveEpgEntry => ({ + key: `xtream::p::${id}`, + item: { + uid: `xtream::p::${id}`, + name: `Channel ${id}`, + contentType: 'live', + sourceType: 'xtream', + playlistId: 'p', + playlistName: 'Portal', + xtreamId: id, + tvgId: String(id), + }, +}); + +const program = (title: string, stopIso?: string): EpgProgram => + ({ + channel: 'x', + title, + start: '2026-05-23T10:00:00.000Z', + stop: stopIso ?? '2026-05-23T11:00:00.000Z', + }) as EpgProgram; + +describe('DashboardPortalLiveEpgService', () => { + let service: DashboardPortalLiveEpgService; + let loadEpgForItems: jest.Mock; + let deferred: Deferred[]; + let offsetMinutes: number; + let sourceChanged: Subject; + + const { delayMs, ttlMs, failureCooldownMs, endedRefetchFloorMs } = + DASHBOARD_PORTAL_LIVE_EPG_TIMING; + + /** Let the queue loop take its next step (one inter-request delay). */ + const step = async (rounds = 1): Promise => { + for (let i = 0; i < rounds; i++) { + await jest.advanceTimersByTimeAsync(delayMs); + } + }; + + /** Settle the LATEST request for `key` — a key can be asked more than once. */ + const settle = async (key: string, value: EpgProgram | null) => { + const request = [...deferred] + .reverse() + .find((request) => request.key === key); + if (!request) throw new Error(`no request in flight for ${key}`); + request.resolve(value); + await jest.advanceTimersByTimeAsync(0); + }; + + beforeEach(() => { + jest.useFakeTimers(); + jest.setSystemTime(new Date('2026-05-23T10:30:00.000Z')); + deferred = []; + offsetMinutes = 0; + sourceChanged = new Subject(); + loadEpgForItems = jest.fn((items: { tvgId?: string }[]) => { + const key = `xtream::p::${items[0].tvgId}`; + return new Promise>( + (resolve, reject) => { + deferred.push({ + key, + resolve: (value) => + resolve(new Map([[String(items[0].tvgId), value]])), + reject, + }); + } + ); + }); + + TestBed.configureTestingModule({ + providers: [ + DashboardPortalLiveEpgService, + { + provide: StreamResolverService, + useValue: { loadEpgForItems }, + }, + { + provide: SettingsStore, + useValue: { + resolvedEpgOffsetMinutes: () => offsetMinutes, + }, + }, + { + provide: EpgSourceSettingsService, + useValue: { changed$: sourceChanged }, + }, + ], + }); + service = TestBed.inject(DashboardPortalLiveEpgService); + }); + + afterEach(() => { + jest.useRealTimers(); + }); + + it('answers each card as its own request lands, never waiting for the slowest', async () => { + service.sync([entry(1), entry(2)]); + expect(service.pending()).toEqual( + new Set(['xtream::p::1', 'xtream::p::2']) + ); + await step(2); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + + await settle('xtream::p::2', program('Fast answer')); + expect(service.programs().get('xtream::p::2')?.title).toBe( + 'Fast answer' + ); + expect(service.programs().has('xtream::p::1')).toBe(false); + expect(service.pending()).toEqual(new Set(['xtream::p::1'])); + + await settle('xtream::p::1', null); + expect(service.programs().get('xtream::p::1')).toBeNull(); + expect(service.pending().size).toBe(0); + }); + + it('never has more than two requests in flight and spaces starts by the delay', async () => { + service.sync([entry(1), entry(2), entry(3), entry(4)]); + expect(loadEpgForItems).toHaveBeenCalledTimes(1); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + await step(3); + // Two in flight: the loop keeps waiting instead of starting a third. + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + + await settle('xtream::p::1', null); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(3); + }); + + it('drops a queued card that scrolled away before its turn and keeps the rest', async () => { + service.sync([entry(1), entry(2), entry(3)]); + // 1 started; 2 and 3 still queued. + service.sync([entry(1), entry(3)]); + expect(service.pending()).toEqual( + new Set(['xtream::p::1', 'xtream::p::3']) + ); + await step(2); + const requested = loadEpgForItems.mock.calls.map( + ([items]) => items[0].tvgId + ); + expect(requested).toEqual(['1', '3']); + }); + + it('serves a fresh answer from cache and asks again once the TTL has passed', async () => { + service.sync([entry(1)]); + await settle('xtream::p::1', program('Cached')); + await step(); + + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(1); + + jest.setSystemTime(Date.now() + ttlMs); + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + }); + + it('asks again when the programme on air has ended, but not more often than the floor', async () => { + service.sync([entry(1)]); + await settle( + 'xtream::p::1', + program('Ending soon', '2026-05-23T10:35:00.000Z') + ); + await step(); + + // Ended 5 min in, but the floor (30 s) already passed → refetch. + jest.setSystemTime(Date.now() + 6 * 60_000); + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + + // Portal keeps returning the ended row: not re-asked within the floor. + await settle( + 'xtream::p::1', + program('Still ended', '2026-05-23T10:35:00.000Z') + ); + await step(); + jest.setSystemTime(Date.now() + endedRefetchFloorMs / 2); + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + }); + + it('leaves a failed portal alone for the cooldown and keeps the card unanswered', async () => { + service.sync([entry(1)]); + deferred[0].reject(new Error('portal down')); + await jest.advanceTimersByTimeAsync(0); + expect(service.programs().has('xtream::p::1')).toBe(false); + expect(service.pending().size).toBe(0); + + jest.setSystemTime(Date.now() + failureCooldownMs / 2); + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(1); + + jest.setSystemTime(Date.now() + failureCooldownMs); + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + }); + + it('retires an answer evaluated under a previous display offset and asks again', async () => { + service.sync([entry(1)]); + offsetMinutes = 60; + await settle('xtream::p::1', program('Old clock')); + await step(); + expect(service.programs().has('xtream::p::1')).toBe(false); + // Still wanted → requeued under the new clock. + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + + await settle('xtream::p::1', program('New clock')); + expect(service.programs().get('xtream::p::1')?.title).toBe('New clock'); + }); + + it('drops every answer when the EPG sources change and asks the wanted cards again', async () => { + service.sync([entry(1)]); + await settle('xtream::p::1', program('Before import')); + await step(); + + sourceChanged.next(); + expect(service.programs().size).toBe(0); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + }); +}); diff --git a/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.ts b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.ts new file mode 100644 index 000000000..e54f64312 --- /dev/null +++ b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.service.ts @@ -0,0 +1,263 @@ +import { inject, Injectable, OnDestroy, signal } from '@angular/core'; +import { + epgProviderClockMs, + type EpgProgram, +} from '@iptvnator/shared/interfaces'; +import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services'; +import { StreamResolverService } from '@iptvnator/portal/shared/data-access'; +import { + dashboardPortalLiveEpgProgramStopMs, + resolveDashboardPortalLiveEpgProgram, + type DashboardPortalLiveEpgEntry, +} from './dashboard-portal-live-epg.util'; + +interface CachedProgram { + readonly program: EpgProgram | null; + readonly fetchedAt: number; +} + +/** + * Same numbers `EpgQueueService` uses against real Xtream panels: two + * requests in flight, 200 ms between starts. A card's answer is trusted for + * a minute; a portal that failed is left alone for 30 s; a programme that + * ended is asked again, but never more often than every 30 s in case the + * portal keeps returning the stale row. + */ +export const DASHBOARD_PORTAL_LIVE_EPG_TIMING = Object.freeze({ + maxConcurrency: 2, + delayMs: 200, + ttlMs: 60_000, + failureCooldownMs: 30_000, + endedRefetchFloorMs: 30_000, +}); + +/** + * Lazy, per-card "what is on air" for the dashboard's Xtream and Stalker + * live cards. `sync()` receives the cards currently worth asking for — the + * visible ones — and the service answers each through the collection + * resolver one card at a time, publishing every answer the moment it lands + * (`programs`) so a slow portal never delays a fast one. The queue holds + * only wanted keys, so a card scrolled away before its turn is never + * requested. One instance for the app: favourites and recent rows of the + * same channel share the key and therefore the answer. + */ +@Injectable({ providedIn: 'root' }) +export class DashboardPortalLiveEpgService implements OnDestroy { + private readonly streamResolver = inject(StreamResolverService); + private readonly settingsStore = inject(SettingsStore); + private readonly sourceSubscription = inject( + EpgSourceSettingsService + ).changed$.subscribe(() => this.retireAnswers()); + + private readonly cache = new Map(); + private readonly failureAt = new Map(); + private wanted = new Map(); + private queue: string[] = []; + private readonly inFlight = new Set(); + private processing = false; + /** Display offset every cached answer was evaluated under. */ + private stateOffsetMinutes = this.offsetMinutes(); + + private readonly programsState = signal< + ReadonlyMap + >(new Map()); + private readonly pendingState = signal>(new Set()); + + /** Answers by entry key; absent = not asked yet, `null` = nothing on air. */ + readonly programs = this.programsState.asReadonly(); + /** Keys queued or in flight — the cards that may show a placeholder. */ + readonly pending = this.pendingState.asReadonly(); + + /** + * Replace the wanted set. Keys without a fresh answer are queued; keys no + * longer wanted are dropped from the queue (an in-flight request is + * allowed to finish and is cached for when the card scrolls back). + */ + sync(entries: readonly DashboardPortalLiveEpgEntry[]): void { + this.retireStateOfPreviousOffset(); + this.wanted = new Map(entries.map((entry) => [entry.key, entry])); + this.queue = this.queue.filter((key) => this.wanted.has(key)); + const now = Date.now(); + for (const entry of entries) { + if ( + this.needsFetch(entry.key, now) && + !this.queue.includes(entry.key) + ) { + this.queue.push(entry.key); + } + } + this.publishPending(); + if (!this.processing && this.queue.length > 0) { + void this.processQueue(); + } + } + + ngOnDestroy(): void { + this.sourceSubscription.unsubscribe(); + } + + private offsetMinutes(): number { + return this.settingsStore.resolvedEpgOffsetMinutes(); + } + + private needsFetch(key: string, now: number): boolean { + if (this.inFlight.has(key) || this.isCoolingDown(key, now)) { + return false; + } + const cached = this.cache.get(key); + if (!cached) { + return true; + } + if (now - cached.fetchedAt >= DASHBOARD_PORTAL_LIVE_EPG_TIMING.ttlMs) { + return true; + } + const stopMs = dashboardPortalLiveEpgProgramStopMs(cached.program); + return ( + stopMs !== null && + epgProviderClockMs(now, this.stateOffsetMinutes) >= stopMs && + now - cached.fetchedAt >= + DASHBOARD_PORTAL_LIVE_EPG_TIMING.endedRefetchFloorMs + ); + } + + private isCoolingDown(key: string, now: number): boolean { + const failedAt = this.failureAt.get(key); + if (failedAt == null) { + return false; + } + if ( + now - failedAt >= + DASHBOARD_PORTAL_LIVE_EPG_TIMING.failureCooldownMs + ) { + this.failureAt.delete(key); + return false; + } + return true; + } + + private async processQueue(): Promise { + this.processing = true; + try { + while (this.queue.length > 0) { + if ( + this.inFlight.size >= + DASHBOARD_PORTAL_LIVE_EPG_TIMING.maxConcurrency + ) { + await this.delay(); + continue; + } + const key = this.queue.shift(); + if ( + key == null || + !this.wanted.has(key) || + !this.needsFetch(key, Date.now()) + ) { + this.publishPending(); + continue; + } + this.inFlight.add(key); + void this.fetch(key); + await this.delay(); + } + } finally { + this.processing = false; + } + } + + private async fetch(key: string): Promise { + const entry = this.wanted.get(key); + if (!entry) { + this.inFlight.delete(key); + this.publishPending(); + return; + } + const offsetMinutes = this.offsetMinutes(); + let program: EpgProgram | null = null; + let failed = false; + try { + const epgMap = await this.streamResolver.loadEpgForItems([ + entry.item, + ]); + program = resolveDashboardPortalLiveEpgProgram(epgMap, entry); + } catch { + failed = true; + } + this.inFlight.delete(key); + + // The setting changed while the request was on the wire: this answer + // belongs to the previous provider clock. Retire it and ask again if + // the card is still wanted. + if (offsetMinutes !== this.offsetMinutes()) { + this.retireStateOfPreviousOffset(); + this.requeueIfWanted(key); + return; + } + + if (failed) { + this.failureAt.set(key, Date.now()); + } else { + this.cache.set(key, { program, fetchedAt: Date.now() }); + this.programsState.update((programs) => { + const next = new Map(programs); + next.set(key, program); + return next; + }); + } + this.publishPending(); + } + + private requeueIfWanted(key: string): void { + if (this.wanted.has(key) && !this.queue.includes(key)) { + this.queue.push(key); + } + this.publishPending(); + if (!this.processing && this.queue.length > 0) { + void this.processQueue(); + } + } + + /** + * Every answer here is "what is on at the provider clock", so a changed + * display offset drops all of them. The caller's own `sync` (the + * presenter re-syncs on an offset change) asks the wanted cards again. + */ + private retireStateOfPreviousOffset(): void { + const current = this.offsetMinutes(); + if (current === this.stateOffsetMinutes) { + return; + } + this.stateOffsetMinutes = current; + this.dropAnswers(); + } + + /** Removed or re-imported XMLTV: drop every answer and ask again now. */ + private retireAnswers(): void { + this.dropAnswers(); + if (this.wanted.size > 0) { + this.sync(Array.from(this.wanted.values())); + } + } + + private dropAnswers(): void { + this.cache.clear(); + this.failureAt.clear(); + this.programsState.set(new Map()); + } + + private publishPending(): void { + const pending = new Set(); + for (const key of this.queue) { + if (this.wanted.has(key)) pending.add(key); + } + for (const key of this.inFlight) { + if (this.wanted.has(key)) pending.add(key); + } + this.pendingState.set(pending); + } + + private delay(): Promise { + return new Promise((resolve) => + setTimeout(resolve, DASHBOARD_PORTAL_LIVE_EPG_TIMING.delayMs) + ); + } +} diff --git a/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.spec.ts b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.spec.ts new file mode 100644 index 000000000..9f235bbce --- /dev/null +++ b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.spec.ts @@ -0,0 +1,160 @@ +import type { + EpgProgram, + PortalActivityItem, +} from '@iptvnator/shared/interfaces'; +import { + buildDashboardPortalLiveEpgEntry, + buildDashboardPortalLiveEpgKey, + dashboardPortalLiveEpgProgramStopMs, + resolveDashboardPortalLiveEpgProgram, +} from './dashboard-portal-live-epg.util'; + +const baseItem = (overrides: Partial): PortalActivityItem => + ({ + id: 'item-1', + title: 'News 24', + type: 'live', + playlist_id: 'playlist-1', + playlist_name: 'My portal', + category_id: '7', + xtream_id: 42, + poster_url: 'https://example.com/logo.png', + ...overrides, + }) as PortalActivityItem; + +describe('buildDashboardPortalLiveEpgEntry', () => { + it('builds an Xtream live entry keyed by the collection uid with the stream id as tvgId', () => { + const entry = buildDashboardPortalLiveEpgEntry( + baseItem({ source: 'xtream', xtream_id: 42 }) + ); + + expect(entry?.key).toBe('xtream::playlist-1::42'); + expect(entry?.item).toMatchObject({ + uid: 'xtream::playlist-1::42', + name: 'News 24', + contentType: 'live', + sourceType: 'xtream', + playlistId: 'playlist-1', + playlistName: 'My portal', + logo: 'https://example.com/logo.png', + xtreamId: 42, + tvgId: '42', + categoryId: '7', + }); + }); + + it('accepts a numeric string Xtream id and rejects ids that cannot name a stream', () => { + expect( + buildDashboardPortalLiveEpgEntry( + baseItem({ source: 'xtream', xtream_id: '17' }) + )?.item.xtreamId + ).toBe(17); + for (const xtream_id of ['', 'abc', 0, -3, 1.5]) { + expect( + buildDashboardPortalLiveEpgEntry( + baseItem({ source: 'xtream', xtream_id }) + ) + ).toBeNull(); + } + }); + + it('builds a Stalker live entry from the stored portal item, keyed by the extracted id', () => { + const stalkerItem = { + id: 'ch-9', + cmd: 'ffrt http://portal/9', + radio: false, + }; + const entry = buildDashboardPortalLiveEpgEntry( + baseItem({ + source: 'stalker', + id: 'ch-9', + xtream_id: 'ch-9', + stalker_item: stalkerItem, + }) + ); + + expect(entry?.key).toBe('stalker::playlist-1::ch-9'); + expect(entry?.item).toMatchObject({ + sourceType: 'stalker', + stalkerId: 'ch-9', + tvgId: 'ch-9', + stalkerCmd: 'ffrt http://portal/9', + stalkerItem, + }); + }); + + it('skips Stalker radio rows and rows without the stored item, like the collection resolver', () => { + expect( + buildDashboardPortalLiveEpgEntry( + baseItem({ + source: 'stalker', + stalker_item: { id: 'r-1', radio: 'true' }, + }) + ) + ).toBeNull(); + expect( + buildDashboardPortalLiveEpgEntry( + baseItem({ source: 'stalker', stalker_item: undefined }) + ) + ).toBeNull(); + }); + + it('leaves M3U rows and non-live rows to the XMLTV batch', () => { + expect( + buildDashboardPortalLiveEpgEntry(baseItem({ source: 'm3u' })) + ).toBeNull(); + expect( + buildDashboardPortalLiveEpgEntry( + baseItem({ source: 'xtream', type: 'movie' }) + ) + ).toBeNull(); + expect( + buildDashboardPortalLiveEpgKey(baseItem({ source: 'xtream' })) + ).toBe('xtream::playlist-1::42'); + expect( + buildDashboardPortalLiveEpgKey(baseItem({ source: 'm3u' })) + ).toBeNull(); + }); +}); + +describe('resolveDashboardPortalLiveEpgProgram', () => { + it('reads the resolver map by the entry tvgId and treats a missing key as nothing on air', () => { + const entry = buildDashboardPortalLiveEpgEntry( + baseItem({ source: 'xtream', xtream_id: 42 }) + ); + const program = { title: 'Evening news' } as EpgProgram; + if (!entry) throw new Error('expected an entry'); + + expect( + resolveDashboardPortalLiveEpgProgram( + new Map([['42', program]]), + entry + ) + ).toBe(program); + expect( + resolveDashboardPortalLiveEpgProgram(new Map(), entry) + ).toBeNull(); + }); +}); + +describe('dashboardPortalLiveEpgProgramStopMs', () => { + it('prefers the unix-seconds stop, falls back to the ISO string, and reports null otherwise', () => { + expect( + dashboardPortalLiveEpgProgramStopMs({ + stopTimestamp: 1_700_000_000, + stop: '2026-01-01T00:00:00.000Z', + } as EpgProgram) + ).toBe(1_700_000_000_000); + expect( + dashboardPortalLiveEpgProgramStopMs({ + stop: '2026-01-01T00:00:00.000Z', + } as EpgProgram) + ).toBe(Date.parse('2026-01-01T00:00:00.000Z')); + expect( + dashboardPortalLiveEpgProgramStopMs({ + stop: 'garbage', + } as EpgProgram) + ).toBeNull(); + expect(dashboardPortalLiveEpgProgramStopMs(null)).toBeNull(); + }); +}); diff --git a/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.ts b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.ts new file mode 100644 index 000000000..f38a7c44a --- /dev/null +++ b/libs/workspace/dashboard/data-access/src/lib/dashboard-portal-live-epg.util.ts @@ -0,0 +1,136 @@ +import { + isStalkerRadioItem, + type EpgProgram, + type PortalActivityItem, +} from '@iptvnator/shared/interfaces'; +import { + buildCollectionUid, + type UnifiedCollectionItem, +} from '@iptvnator/portal/shared/util'; + +/** + * One dashboard live card's portal EPG request: the key its answer is filed + * under (the collection uid, stable across favourites and recent rows of the + * same channel) and the item `StreamResolverService.loadEpgForItems` fetches + * for — the same shape the "See all" pages hand it, so the two surfaces + * resolve a channel identically. + */ +export interface DashboardPortalLiveEpgEntry { + readonly key: string; + readonly item: UnifiedCollectionItem; +} + +/** The key alone, for cards that only need to find their answer. */ +export function buildDashboardPortalLiveEpgKey( + item: PortalActivityItem +): string | null { + return buildDashboardPortalLiveEpgEntry(item)?.key ?? null; +} + +/** + * M3U channels keep the dashboard's batched XMLTV lookup; Xtream and Stalker + * live rows have no XMLTV key of their own and are answered by their portal. + * Radio rows are skipped like the collection resolver skips them. + */ +export function buildDashboardPortalLiveEpgEntry( + item: PortalActivityItem +): DashboardPortalLiveEpgEntry | null { + if (item.type !== 'live') { + return null; + } + if (item.source === 'xtream') { + return buildXtreamEntry(item); + } + if (item.source === 'stalker') { + return buildStalkerEntry(item); + } + return null; +} + +/** + * The resolver keys its map by the item's `tvgId`; the entry sets it to the + * provider id, so a missing key means "asked, nothing on air". + */ +export function resolveDashboardPortalLiveEpgProgram( + epgMap: ReadonlyMap, + entry: DashboardPortalLiveEpgEntry +): EpgProgram | null { + const key = entry.item.tvgId?.trim(); + return key ? (epgMap.get(key) ?? null) : null; +} + +/** + * When the programme on air ends, in wall-clock milliseconds; `null` when the + * row states no usable end. Reads the pre-computed unix-seconds field first, + * then the ISO string, like the dashboard's own time-range helper. + */ +export function dashboardPortalLiveEpgProgramStopMs( + program: EpgProgram | null +): number | null { + if (!program) { + return null; + } + const cached = Number(program.stopTimestamp); + if (Number.isFinite(cached) && cached > 0) { + return cached * 1000; + } + const parsed = program.stop ? Date.parse(program.stop) : NaN; + return Number.isFinite(parsed) ? parsed : null; +} + +function buildXtreamEntry( + item: PortalActivityItem +): DashboardPortalLiveEpgEntry | null { + const xtreamId = Number(item.xtream_id); + if (!Number.isInteger(xtreamId) || xtreamId <= 0) { + return null; + } + const key = buildCollectionUid('xtream', item.playlist_id, xtreamId); + return { + key, + item: { + uid: key, + name: item.title, + contentType: 'live', + sourceType: 'xtream', + playlistId: item.playlist_id, + playlistName: item.playlist_name ?? 'Xtream', + logo: item.poster_url ?? null, + xtreamId, + tvgId: String(xtreamId), + categoryId: item.category_id, + }, + }; +} + +function buildStalkerEntry( + item: PortalActivityItem +): DashboardPortalLiveEpgEntry | null { + const raw = item.stalker_item; + if (!raw || isStalkerRadioItem(raw)) { + return null; + } + // Both dashboard mappers already store the extracted portal id as `id`. + const stalkerId = String(item.id ?? '').trim(); + if (!stalkerId) { + return null; + } + const key = buildCollectionUid('stalker', item.playlist_id, stalkerId); + return { + key, + item: { + uid: key, + name: item.title, + contentType: 'live', + sourceType: 'stalker', + playlistId: item.playlist_id, + playlistName: item.playlist_name ?? 'Stalker', + logo: item.poster_url ?? null, + stalkerId, + tvgId: stalkerId, + stalkerCmd: raw.cmd, + categoryId: item.category_id, + stalkerItem: raw, + }, + }; +} 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 index 803280606..236508129 100644 --- 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 @@ -17,11 +17,15 @@ import { switchMap, } from 'rxjs'; import { EpgService } from '@iptvnator/epg/data-access'; -import type { EpgProgram } from '@iptvnator/shared/interfaces'; +import type { + EpgProgram, + PortalActivityItem, +} 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 { DashboardPortalLiveEpgPresenter } from './dashboard-portal-live-epg.presenter'; import { buildDashboardLiveEpgDetails, buildLiveEpgLookupGroups, @@ -60,12 +64,19 @@ const emptyAnswer = (scopeKey: string): ScopeAnswer => ({ * 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. + * + * It is also the one facade the rails talk to for live EPG: an Xtream or + * Stalker card has no XMLTV key of its own, so its programme comes from + * `DashboardPortalLiveEpgPresenter` and only falls back to the title match + * here. */ @Injectable() export class DashboardLiveEpgPresenter { private readonly data = inject(DashboardDataService); private readonly epgService = inject(EpgService); private readonly settingsStore = inject(SettingsStore); + /** Xtream/Stalker cards are answered by their portal, not by XMLTV. */ + private readonly portal = inject(DashboardPortalLiveEpgPresenter); private readonly cards = signal): void { + this.portal.connect(items); + } + + /** Portal keys wanted regardless of scrolling (the hero card). */ + setPinnedPortalKeys(keys: readonly (string | null | undefined)[]): void { + this.portal.setPinnedKeys(keys); + } + + /** A rail reported the cards inside its viewport. */ + setVisibleCards(railId: string, cards: readonly DashboardRailCard[]): void { + this.portal.setVisibleCards(railId, 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) - ) - ); + // A portal answer wins. Its `null` ("asked, nothing on air") and + // "not asked yet" both fall back to the XMLTV lookup, which for a + // portal card can only ever be a title match. + const program = + this.portal.programFor(card.liveEpgSourceKey) ?? + 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( @@ -140,6 +171,17 @@ export class DashboardLiveEpgPresenter { ); } + /** + * True only before a card's FIRST portal answer, so a refresh keeps the + * previous answer on screen instead of flashing a placeholder. + */ + isAwaitingFirstAnswer(card: DashboardRailCard): boolean { + return ( + this.portal.programFor(card.liveEpgSourceKey) === undefined && + this.portal.isPending(card.liveEpgSourceKey) + ); + } + /** The XMLTV scope a live card's programme must be resolved in. */ private sourceUrlsForCard(card: DashboardRailCard): string[] { const byPlaylistId = this.sourceUrlsByPlaylist(); diff --git a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.spec.ts b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.spec.ts new file mode 100644 index 000000000..409c9abb6 --- /dev/null +++ b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.spec.ts @@ -0,0 +1,152 @@ +import { signal } from '@angular/core'; +import { TestBed } from '@angular/core/testing'; +import type { + EpgProgram, + PortalActivityItem, +} from '@iptvnator/shared/interfaces'; +import { SettingsStore } from '@iptvnator/services'; +import { + DashboardPortalLiveEpgService, + type DashboardPortalLiveEpgEntry, +} from '@iptvnator/workspace/dashboard/data-access'; +import { DashboardPortalLiveEpgPresenter } from './dashboard-portal-live-epg.presenter'; + +const xtreamLive = (id: number, playlist = 'p'): PortalActivityItem => + ({ + id: `x-${id}`, + title: `Channel ${id}`, + type: 'live', + source: 'xtream', + playlist_id: playlist, + category_id: '1', + xtream_id: id, + }) as PortalActivityItem; + +describe('DashboardPortalLiveEpgPresenter', () => { + let presenter: DashboardPortalLiveEpgPresenter; + let sync: jest.Mock; + let programs: ReturnType< + typeof signal> + >; + let pending: ReturnType>>; + let offsetMinutes: ReturnType>; + + /** Keys of every sync call, sorted: the wanted set has no order. */ + const wantedKeys = (): string[][] => + sync.mock.calls.map(([entries]: [DashboardPortalLiveEpgEntry[]]) => + entries.map((entry) => entry.key).sort() + ); + + beforeEach(() => { + sync = jest.fn(); + programs = signal>(new Map()); + pending = signal>(new Set()); + offsetMinutes = signal(0); + TestBed.configureTestingModule({ + providers: [ + DashboardPortalLiveEpgPresenter, + { + provide: DashboardPortalLiveEpgService, + useValue: { sync, programs, pending }, + }, + { + provide: SettingsStore, + useValue: { resolvedEpgOffsetMinutes: offsetMinutes }, + }, + ], + }); + presenter = TestBed.inject(DashboardPortalLiveEpgPresenter); + }); + + it('asks only for pinned and visible cards, deduplicated across rails', () => { + const items = signal([ + xtreamLive(1), + xtreamLive(2), + xtreamLive(3), + // The same channel in the recent rail shares its key. + xtreamLive(2), + ]); + presenter.connect(items); + TestBed.tick(); + expect(wantedKeys().at(-1)).toEqual([]); + + presenter.setPinnedKeys(['xtream::p::1', null, undefined]); + presenter.setVisibleCards('favorites', [ + { id: 'f2', liveEpgSourceKey: 'xtream::p::2' }, + { id: 'f3', liveEpgSourceKey: 'xtream::p::3' }, + ] as never); + presenter.setVisibleCards('recent', [ + { id: 'r2', liveEpgSourceKey: 'xtream::p::2' }, + { id: 'm3u', liveEpgSourceKey: null }, + ] as never); + TestBed.tick(); + expect(wantedKeys().at(-1)).toEqual([ + 'xtream::p::1', + 'xtream::p::2', + 'xtream::p::3', + ]); + + // Scrolled away: the favourites rail now shows only channel 3. + presenter.setVisibleCards('favorites', [ + { id: 'f3', liveEpgSourceKey: 'xtream::p::3' }, + ] as never); + TestBed.tick(); + expect(wantedKeys().at(-1)).toEqual([ + 'xtream::p::1', + 'xtream::p::2', + 'xtream::p::3', + ]); + presenter.setVisibleCards('recent', []); + TestBed.tick(); + expect(wantedKeys().at(-1)).toEqual(['xtream::p::1', 'xtream::p::3']); + }); + + it('ignores visible keys whose item is no longer on the dashboard', () => { + const items = signal([xtreamLive(1)]); + presenter.connect(items); + presenter.setVisibleCards('favorites', [ + { id: 'f1', liveEpgSourceKey: 'xtream::p::1' }, + { id: 'gone', liveEpgSourceKey: 'xtream::p::9' }, + ] as never); + TestBed.tick(); + expect(wantedKeys().at(-1)).toEqual(['xtream::p::1']); + + items.set([]); + TestBed.tick(); + expect(wantedKeys().at(-1)).toEqual([]); + }); + + it('re-syncs the same wanted set when the display offset changes', () => { + presenter.connect( + signal([xtreamLive(1)]) + ); + presenter.setPinnedKeys(['xtream::p::1']); + TestBed.tick(); + const before = sync.mock.calls.length; + + offsetMinutes.set(30); + TestBed.tick(); + expect(sync.mock.calls.length).toBe(before + 1); + expect(wantedKeys().at(-1)).toEqual(['xtream::p::1']); + }); + + it('answers a card from the service: undefined until asked, null when nothing is on air', () => { + const program = { title: 'Now' } as EpgProgram; + expect(presenter.programFor('xtream::p::1')).toBeUndefined(); + expect(presenter.programFor(null)).toBeUndefined(); + expect(presenter.isPending('xtream::p::1')).toBe(false); + + pending.set(new Set(['xtream::p::1'])); + expect(presenter.isPending('xtream::p::1')).toBe(true); + expect(presenter.isPending(undefined)).toBe(false); + + programs.set( + new Map([ + ['xtream::p::1', program], + ['xtream::p::2', null], + ]) + ); + expect(presenter.programFor('xtream::p::1')).toBe(program); + expect(presenter.programFor('xtream::p::2')).toBeNull(); + }); +}); diff --git a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.ts b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.ts new file mode 100644 index 000000000..0c7d3126b --- /dev/null +++ b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-portal-live-epg.presenter.ts @@ -0,0 +1,129 @@ +import { + computed, + effect, + inject, + Injectable, + signal, + untracked, + type Signal, +} from '@angular/core'; +import { toSignal } from '@angular/core/rxjs-interop'; +import { interval, map } from 'rxjs'; +import type { + EpgProgram, + PortalActivityItem, +} from '@iptvnator/shared/interfaces'; +import { SettingsStore } from '@iptvnator/services'; +import { + buildDashboardPortalLiveEpgEntry, + DashboardPortalLiveEpgService, + type DashboardPortalLiveEpgEntry, +} from '@iptvnator/workspace/dashboard/data-access'; +import type { DashboardRailCard } from './dashboard-rail.component'; +import { LIVE_EPG_TICK_MS } from './dashboard-live-epg.utils'; + +/** + * The rails component's view of `DashboardPortalLiveEpgService`: which + * portal live cards exist, which of them are on screen, and the answer for + * a card. Component-provided, so its wanted set dies with the dashboard. + * + * "On screen" is what the rails report through their visibility output, + * plus the pinned keys (the hero card, always at the top). A card the user + * never scrolls to is never asked for. The 30 s tick re-syncs so a programme + * that ended is asked again and progress bars keep moving; a changed display + * offset re-syncs at once because every cached answer was just retired. + */ +@Injectable() +export class DashboardPortalLiveEpgPresenter { + private readonly service = inject(DashboardPortalLiveEpgService); + private readonly settingsStore = inject(SettingsStore); + + private readonly source = signal | null>(null); + private readonly visibleKeysByRail = signal< + ReadonlyMap> + >(new Map()); + private readonly pinnedKeys = signal>(new Set()); + + /** Heartbeat shared with the XMLTV batch: shifted by one so the first + * emission differs from `initialValue` and is not swallowed. */ + readonly tick = toSignal( + interval(LIVE_EPG_TICK_MS).pipe(map((tick) => tick + 1)), + { initialValue: 0 } + ); + + private readonly entries = computed< + ReadonlyMap + >(() => { + const source = this.source(); + const entries = new Map(); + for (const item of source?.() ?? []) { + const entry = buildDashboardPortalLiveEpgEntry(item); + if (entry && !entries.has(entry.key)) { + entries.set(entry.key, entry); + } + } + return entries; + }); + + private readonly wanted = computed(() => { + const entries = this.entries(); + const keys = new Set(this.pinnedKeys()); + for (const railKeys of this.visibleKeysByRail().values()) { + for (const key of railKeys) keys.add(key); + } + const wanted: DashboardPortalLiveEpgEntry[] = []; + for (const key of keys) { + const entry = entries.get(key); + if (entry) wanted.push(entry); + } + return wanted; + }); + + constructor() { + effect(() => { + const wanted = this.wanted(); + this.tick(); + this.settingsStore.resolvedEpgOffsetMinutes(); + untracked(() => this.service.sync(wanted)); + }); + } + + /** The live rows every portal card on the dashboard is built from. */ + connect(source: Signal): void { + this.source.set(source); + } + + /** Keys wanted regardless of scrolling (the hero card). */ + setPinnedKeys(keys: readonly (string | null | undefined)[]): void { + this.pinnedKeys.set( + new Set(keys.filter((key): key is string => !!key)) + ); + } + + /** A rail reported the cards inside its viewport. */ + setVisibleCards(railId: string, cards: readonly DashboardRailCard[]): void { + const keys = new Set(); + for (const card of cards) { + if (card.liveEpgSourceKey) keys.add(card.liveEpgSourceKey); + } + this.visibleKeysByRail.update((byRail) => { + const next = new Map(byRail); + next.set(railId, keys); + return next; + }); + } + + /** + * `undefined` = not asked yet or still in flight, `null` = asked and + * nothing on air, else the programme. + */ + programFor(key: string | null | undefined): EpgProgram | null | undefined { + return key ? this.service.programs().get(key) : undefined; + } + + isPending(key: string | null | undefined): boolean { + return !!key && this.service.pending().has(key); + } +} diff --git a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.html b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.html index cfe91881b..96c5584e8 100644 --- a/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.html +++ b/libs/workspace/dashboard/feature/src/lib/rails/dashboard-rail.component.html @@ -21,8 +21,7 @@ 'WORKSPACE.DASHBOARD.SEE_ALL_COUNT' | translate : { - count: - totalCount() ?? items().length, + count: totalCount() ?? items().length, } }} } @else { @@ -56,8 +55,10 @@ @for (card of items(); track card.id) { @if (layout() === 'channel') {
@if (card.nowPlayingTitle) { {{ card.nowPlayingTitle }} + } @else if ( + card.nowPlayingState === 'pending' + ) { + } @else if (card.subtitle) { {{ card.subtitle }} } @@ -132,8 +141,10 @@
} @else {