From 15d7709934ab4059e9f996f1cd33fd36e9c439f5 Mon Sep 17 00:00:00 2001 From: 4gray Date: Sun, 20 Sep 2026 13:31:41 +0200 Subject: [PATCH] fix(dashboard): guard portal EPG answers by source revision, release them on teardown MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-ups on the lazy portal EPG queue. - A completion now captures `EpgSourceSettingsService.revision()` next to the display offset and requeues itself when either moved. The retire pass cannot queue a replacement while the key is in flight, so without this the answer computed against the removed guide was published and trusted for a full TTL — the repo's late-result invalidation contract (Greptile P1, Codex P2). - The presenter hands its wanted set back on destroy: the queue lives in the root service and kept asking for cards on a page the user had left (Greptile P2, Codex P2). - The failure cooldown was dead code: the collection resolver catches per-channel portal failures and files them as `null`, so the outer catch never fired. Answers with no programme now expire after 30 s instead of 60 s, which is what lets a transient outage recover, and a resolver-level throw takes the same path (Codex P2). - `sync()` returns immediately without the local EPG program-lookup capability: the resolver is gated on it and answers nothing, so the PWA queued work that could never produce an answer (Codex P2). The release note now says desktop. Co-Authored-By: Claude Opus 5 --- .changes/dashboard-portal-live-epg.md | 8 +- docs/architecture/workspace-dashboard.md | 30 ++++- .../dashboard-portal-live-epg.service.spec.ts | 103 ++++++++++++++++-- .../lib/dashboard-portal-live-epg.service.ts | 100 +++++++++-------- ...ashboard-portal-live-epg.presenter.spec.ts | 14 +++ .../dashboard-portal-live-epg.presenter.ts | 6 + 6 files changed, 197 insertions(+), 64 deletions(-) diff --git a/.changes/dashboard-portal-live-epg.md b/.changes/dashboard-portal-live-epg.md index 0ae5d21f3..6b4547213 100644 --- a/.changes/dashboard-portal-live-epg.md +++ b/.changes/dashboard-portal-live-epg.md @@ -3,7 +3,7 @@ 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. +On the desktop app, 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/docs/architecture/workspace-dashboard.md b/docs/architecture/workspace-dashboard.md index daf113338..f239b9273 100644 --- a/docs/architecture/workspace-dashboard.md +++ b/docs/architecture/workspace-dashboard.md @@ -149,17 +149,37 @@ Render rules: 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. + The queue lives in the root service, so leaving the dashboard hands + the wanted set back (`sync([])` on destroy); otherwise the queue + would keep asking for cards on a page that is gone. - `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. + requested. A programme lives 60 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). An answer with **no** programme lives only + 30 s, because the resolver reports a failed portal and a guide-less + channel identically (it files per-channel failures as `null`), so + there is no failure cooldown to keep and the short TTL is what lets + an outage recover on the next tick. + - Every answer is "at the provider clock" and against one XMLTV source + set. A request captures both the display offset and + `EpgSourceSettingsService.revision()` — the same fence + `EpgService.guard()` uses — and a completion whose either fact moved + is discarded and requeued instead of published. That requeue has to + happen in the completion: while the key is in flight the retire pass + cannot queue a replacement, and without it the pre-change answer + would be trusted for a full TTL (the repo's late-result + invalidation contract). + - Desktop only in practice: the shared collection resolver is gated on + the local XMLTV bridge (`supportsProgramLookup`) and answers nothing + without it, so `sync()` returns immediately in the PWA rather than + filing an empty answer for every card. Lifting that gate for portal + lookups would change the collection pages too and is deliberately + out of scope here. - `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 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 index 515c38a62..9e9df0559 100644 --- 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 @@ -1,7 +1,11 @@ import { TestBed } from '@angular/core/testing'; import { Subject } from 'rxjs'; import type { EpgProgram } from '@iptvnator/shared/interfaces'; -import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services'; +import { + EpgSourceSettingsService, + RuntimeCapabilitiesService, + SettingsStore, +} from '@iptvnator/services'; import { StreamResolverService } from '@iptvnator/portal/shared/data-access'; import { DASHBOARD_PORTAL_LIVE_EPG_TIMING, @@ -43,10 +47,18 @@ describe('DashboardPortalLiveEpgService', () => { let deferred: Deferred[]; let offsetMinutes: number; let sourceChanged: Subject; + let sourceRevision: number; + let supportsEpgProgramLookup: boolean; - const { delayMs, ttlMs, failureCooldownMs, endedRefetchFloorMs } = + const { delayMs, ttlMs, emptyTtlMs, endedRefetchFloorMs } = DASHBOARD_PORTAL_LIVE_EPG_TIMING; + /** What a reconciliation does: bump the fence, then announce it. */ + const changeEpgSources = () => { + sourceRevision++; + sourceChanged.next(); + }; + /** 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++) { @@ -69,6 +81,8 @@ describe('DashboardPortalLiveEpgService', () => { jest.setSystemTime(new Date('2026-05-23T10:30:00.000Z')); deferred = []; offsetMinutes = 0; + sourceRevision = 0; + supportsEpgProgramLookup = true; sourceChanged = new Subject(); loadEpgForItems = jest.fn((items: { tvgId?: string }[]) => { const key = `xtream::p::${items[0].tvgId}`; @@ -99,7 +113,18 @@ describe('DashboardPortalLiveEpgService', () => { }, { provide: EpgSourceSettingsService, - useValue: { changed$: sourceChanged }, + useValue: { + changed$: sourceChanged, + revision: () => sourceRevision, + }, + }, + { + provide: RuntimeCapabilitiesService, + useValue: { + get supportsEpgProgramLookup() { + return supportsEpgProgramLookup; + }, + }, }, ], }); @@ -199,19 +224,42 @@ describe('DashboardPortalLiveEpgService', () => { expect(loadEpgForItems).toHaveBeenCalledTimes(2); }); - it('leaves a failed portal alone for the cooldown and keeps the card unanswered', async () => { + it('expires an answer with no programme sooner than one with a programme', async () => { + // The resolver reports a failed portal and a guide-less channel the + // same way, so the short TTL is what lets an outage recover. 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); + await settle('xtream::p::1', null); + await step(); + expect(service.programs().get('xtream::p::1')).toBeNull(); - jest.setSystemTime(Date.now() + failureCooldownMs / 2); + jest.setSystemTime(Date.now() + emptyTtlMs / 2); service.sync([entry(1)]); await step(); expect(loadEpgForItems).toHaveBeenCalledTimes(1); - jest.setSystemTime(Date.now() + failureCooldownMs); + jest.setSystemTime(Date.now() + emptyTtlMs / 2); + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + + // A programme keeps the full TTL, which outlives the empty one. + await settle('xtream::p::1', program('On air')); + await step(); + jest.setSystemTime(Date.now() + emptyTtlMs); + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + }); + + it('treats a rejected resolver call as an answer with no programme', async () => { + service.sync([entry(1)]); + deferred[0].reject(new Error('portal down')); + await jest.advanceTimersByTimeAsync(0); + + expect(service.programs().get('xtream::p::1')).toBeNull(); + expect(service.pending().size).toBe(0); + + jest.setSystemTime(Date.now() + emptyTtlMs); service.sync([entry(1)]); await step(); expect(loadEpgForItems).toHaveBeenCalledTimes(2); @@ -235,9 +283,42 @@ describe('DashboardPortalLiveEpgService', () => { await settle('xtream::p::1', program('Before import')); await step(); - sourceChanged.next(); + changeEpgSources(); expect(service.programs().size).toBe(0); await step(); expect(loadEpgForItems).toHaveBeenCalledTimes(2); }); + + it('discards an answer computed before an EPG source change and asks again', async () => { + // The key is in flight when the sources change, so the retire pass + // cannot requeue it; the completion must not publish the old guide's + // answer and must ask again itself. + service.sync([entry(1)]); + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(1); + + changeEpgSources(); + await settle('xtream::p::1', program('Removed guide')); + expect(service.programs().has('xtream::p::1')).toBe(false); + + await step(); + expect(loadEpgForItems).toHaveBeenCalledTimes(2); + await settle('xtream::p::1', program('Current guide')); + expect(service.programs().get('xtream::p::1')?.title).toBe( + 'Current guide' + ); + }); + + it('does nothing at all without the local EPG program-lookup capability', async () => { + // PWA: the collection resolver is gated on the desktop XMLTV bridge + // and answers nothing, so no request is worth queuing. + supportsEpgProgramLookup = false; + + service.sync([entry(1), entry(2)]); + await step(3); + + expect(loadEpgForItems).not.toHaveBeenCalled(); + expect(service.pending().size).toBe(0); + expect(service.programs().size).toBe(0); + }); }); 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 index e54f64312..b4a9e158d 100644 --- 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 @@ -3,7 +3,11 @@ import { epgProviderClockMs, type EpgProgram, } from '@iptvnator/shared/interfaces'; -import { EpgSourceSettingsService, SettingsStore } from '@iptvnator/services'; +import { + EpgSourceSettingsService, + RuntimeCapabilitiesService, + SettingsStore, +} from '@iptvnator/services'; import { StreamResolverService } from '@iptvnator/portal/shared/data-access'; import { dashboardPortalLiveEpgProgramStopMs, @@ -18,16 +22,22 @@ interface CachedProgram { /** * 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. + * requests in flight, 200 ms between starts. A programme is trusted for a + * minute; a programme that ended is asked again, but never more often than + * every 30 s in case the portal keeps returning the stale row. + * + * An answer with no programme expires sooner (`emptyTtlMs`) because the + * collection resolver reports a failed portal and a channel with no guide + * identically — it catches per-channel failures and files them as `null`. + * There is therefore no failure cooldown to keep: the short TTL is what lets + * a transient outage recover on the next tick, at the price of re-asking a + * genuinely guide-less channel while its card stays on screen. */ export const DASHBOARD_PORTAL_LIVE_EPG_TIMING = Object.freeze({ maxConcurrency: 2, delayMs: 200, ttlMs: 60_000, - failureCooldownMs: 30_000, + emptyTtlMs: 30_000, endedRefetchFloorMs: 30_000, }); @@ -45,12 +55,12 @@ export const DASHBOARD_PORTAL_LIVE_EPG_TIMING = Object.freeze({ 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 runtime = inject(RuntimeCapabilitiesService); + private readonly sourceSettings = inject(EpgSourceSettingsService); + private readonly sourceSubscription = + this.sourceSettings.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(); @@ -74,6 +84,13 @@ export class DashboardPortalLiveEpgService implements OnDestroy { * allowed to finish and is cached for when the card scrolls back). */ sync(entries: readonly DashboardPortalLiveEpgEntry[]): void { + // The collection resolver this queue asks through is gated on the + // local XMLTV bridge (`supportsProgramLookup`, desktop only) and + // answers nothing without it, so the PWA never queues at all rather + // than filing an empty answer for every card. + if (!this.runtime.supportsEpgProgramLookup) { + return; + } this.retireStateOfPreviousOffset(); this.wanted = new Map(entries.map((entry) => [entry.key, entry])); this.queue = this.queue.filter((key) => this.wanted.has(key)); @@ -101,14 +118,17 @@ export class DashboardPortalLiveEpgService implements OnDestroy { } private needsFetch(key: string, now: number): boolean { - if (this.inFlight.has(key) || this.isCoolingDown(key, now)) { + if (this.inFlight.has(key)) { return false; } const cached = this.cache.get(key); if (!cached) { return true; } - if (now - cached.fetchedAt >= DASHBOARD_PORTAL_LIVE_EPG_TIMING.ttlMs) { + const ttlMs = cached.program + ? DASHBOARD_PORTAL_LIVE_EPG_TIMING.ttlMs + : DASHBOARD_PORTAL_LIVE_EPG_TIMING.emptyTtlMs; + if (now - cached.fetchedAt >= ttlMs) { return true; } const stopMs = dashboardPortalLiveEpgProgramStopMs(cached.program); @@ -120,21 +140,6 @@ export class DashboardPortalLiveEpgService implements OnDestroy { ); } - 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 { @@ -171,38 +176,46 @@ export class DashboardPortalLiveEpgService implements OnDestroy { this.publishPending(); return; } + // Both facts this answer is evaluated against. The revision is the + // same fence `EpgService.guard()` uses: a reconciliation bumps it, so + // a result computed against the previous XMLTV source set can be told + // apart from one computed against the current one. const offsetMinutes = this.offsetMinutes(); + const revision = this.sourceSettings.revision(); let program: EpgProgram | null = null; - let failed = false; try { const epgMap = await this.streamResolver.loadEpgForItems([ entry.item, ]); program = resolveDashboardPortalLiveEpgProgram(epgMap, entry); } catch { - failed = true; + // The resolver files a failed portal as `null` itself, so this + // only catches a resolver-level throw. Same answer either way, + // and `emptyTtlMs` is what makes it recoverable. + program = null; } 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()) { + // A setting or the source set changed while the request was on the + // wire: this answer belongs to the previous provider clock or the + // previous guide. `retireAnswers()` could not requeue the key while + // it was in flight, so the requeue happens here — otherwise the stale + // answer would be published and trusted for a full TTL. + if ( + offsetMinutes !== this.offsetMinutes() || + revision !== this.sourceSettings.revision() + ) { 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.cache.set(key, { program, fetchedAt: Date.now() }); + this.programsState.update((programs) => { + const next = new Map(programs); + next.set(key, program); + return next; + }); this.publishPending(); } @@ -240,7 +253,6 @@ export class DashboardPortalLiveEpgService implements OnDestroy { private dropAnswers(): void { this.cache.clear(); - this.failureAt.clear(); this.programsState.set(new Map()); } 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 index 409c9abb6..9ac14530e 100644 --- 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 @@ -130,6 +130,20 @@ describe('DashboardPortalLiveEpgPresenter', () => { expect(wantedKeys().at(-1)).toEqual(['xtream::p::1']); }); + it('hands its wanted set back to the root service when the dashboard is destroyed', () => { + presenter.connect( + signal([xtreamLive(1)]) + ); + presenter.setPinnedKeys(['xtream::p::1']); + TestBed.tick(); + expect(wantedKeys().at(-1)).toEqual(['xtream::p::1']); + + // The queue lives in the root service and would otherwise keep + // asking for cards on a page the user has left. + TestBed.resetTestingModule(); + expect(wantedKeys().at(-1)).toEqual([]); + }); + 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(); 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 index 0c7d3126b..8d883fc96 100644 --- 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 @@ -1,5 +1,6 @@ import { computed, + DestroyRef, effect, inject, Injectable, @@ -88,6 +89,11 @@ export class DashboardPortalLiveEpgPresenter { this.settingsStore.resolvedEpgOffsetMinutes(); untracked(() => this.service.sync(wanted)); }); + // The queue lives in the root service; this presenter owns what it + // wants. Leaving the dashboard must hand that back, or the queue + // would keep asking for cards on a page that is gone — and a later + // source change would ask for them again. + inject(DestroyRef).onDestroy(() => this.service.sync([])); } /** The live rows every portal card on the dashboard is built from. */