mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-11 11:06:16 -08:00
fix(dashboard): guard portal EPG answers by source revision, release them on teardown
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 <noreply@anthropic.com>
This commit is contained in:
1 parent
9fc487dcb3
commit
15d7709934
6 files changed
+197
-64
No files matched your search
@@ -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.
|
||||
@@ -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
|
||||
|
||||
+92
-11
@@ -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<void>;
|
||||
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<void> => {
|
||||
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<void>();
|
||||
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);
|
||||
});
|
||||
});
|
||||
@@ -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<string, CachedProgram>();
|
||||
private readonly failureAt = new Map<string, number>();
|
||||
private wanted = new Map<string, DashboardPortalLiveEpgEntry>();
|
||||
private queue: string[] = [];
|
||||
private readonly inFlight = new Set<string>();
|
||||
@@ -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<void> {
|
||||
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());
|
||||
}
|
||||
|
||||
|
||||
+14
@@ -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<readonly PortalActivityItem[]>([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();
|
||||
|
||||
@@ -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. */
|
||||
|
||||
Reference in new issue
Block a user