diff --git a/.changes/xtream-epg-current-program-refresh.md b/.changes/xtream-epg-current-program-refresh.md new file mode 100644 index 000000000..aac681249 --- /dev/null +++ b/.changes/xtream-epg-current-program-refresh.md @@ -0,0 +1,7 @@ +--- +type: fix +area: epg +issues: [767] +--- + +The current program shown under each channel in the Xtream Live TV list now moves on by itself once that program ends, and its progress bar keeps advancing. Previously both stayed frozen on the old program until you left the category and came back in. diff --git a/CLAUDE.md b/CLAUDE.md index 3e01da88c..c2022dc4e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -488,6 +488,10 @@ libs/portal/xtream/ └── feature/src/lib/ # Routed components ├── xtream-feature.routes.ts # createXtreamRoutes(): /workspace/xtreams/:id tree ├── live-stream-layout/, vod-details/, serial-details/, ... + ├── portal-channels-list/ + │ ├── epg-preview-program.ts # Pure "which programme is current" rules + │ ├── epg-refill-limiter.service.ts # Floor on refetching an exhausted guide + │ └── epg-refresh-coordinator.service.ts # One refresh timer, merged queue requests └── global-search-results/ # Global search (Electron-only route) ``` @@ -1609,6 +1613,7 @@ stream_id`); it drops `series_id`/`movie_id`, so the builder pins the - 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") +- Xtream channel-row refresh (#767): the "current programme" under each Live TV row is re-checked once a minute, since `applyProgram()` otherwise runs only on scroll-into-view, a new EPG result or an offset change. A programme still on air only has its progress bar advanced (no cache read, no request); once it ends the row is re-picked and only an on-air or upcoming programme may replace it, because the earliest-item fallback that fills a blank row on first paint would otherwise move an advanced row backwards; a cached guide whose programmes have ALL ended is invalidated and refetched (the queue skips any stream still holding an answer) while an EMPTY answer means the provider has no guide for that channel and is left alone. A programme occupies `[start, stop)` in every comparison, so "has it ended" and "what is on air" cannot disagree on the boundary. Pure rules: `epg-preview-program.ts`. Two root services exist because a live layout mounts the list more than once over one `EpgQueueService` — `EpgRefillLimiter` (floor on dropping an exhausted cache, keyed by playlist + stream since ids are provider-local and the service outlives a playlist switch, expiring by age not viewport membership) and `EpgRefreshCoordinator` (owns the single timer and merges every mounted list's request, because `enqueue` is latest-wins and separate timers would cancel each other on every programme boundary). Contract: `docs/architecture/m3u-playlist-module.md` ("Xtream channel-row programme refresh") - Programme guide (Electron, M3U): `app-epg-guide` in `libs/ui/epg` fed by the host-provided `EPG_GUIDE_SOURCE`; the M3U host switches into guide mode (docked player strip, no sidebar/timeline, no remount) from the header action, the palette, the EPG panel's Guide button (timeline or list view) or `G`. Data: `EPG_GET_PROGRAMS_FOR_CHANNELS` / `EPG_GET_PROGRAM_COVERAGE` (keys resolved in main; manual mappings honoured). Contract: `docs/architecture/m3u-playlist-module.md` ("Programme guide"). - Manual EPG mapping (Electron only): right-click a channel in any list (M3U views, Xtream portal list, Stalker ITV sidebar, global favorites) → "Map EPG channel" attaches it to an uploaded-XMLTV channel; stored in `epg_channel_mappings` keyed by the M3U lookup key or a playlist-scoped portal key (`xtream:{playlistId}:{id}` / `stalker:{playlistId}:{id}`, helpers in `libs/shared/interfaces/src/lib/epg-mapping-key.util.ts`); resolved on every EPG path (single + batch IPC lookups, portal detail views, preview queues); dialog: `libs/ui/components/src/lib/channel-list-container/epg-mapping-dialog/` diff --git a/docs/architecture/m3u-playlist-module.md b/docs/architecture/m3u-playlist-module.md index 75d698dc9..47181aeec 100644 --- a/docs/architecture/m3u-playlist-module.md +++ b/docs/architecture/m3u-playlist-module.md @@ -1090,6 +1090,62 @@ These URLs are playlist-scoped by default: ## EPG Integration +### Xtream channel-row programme refresh + +The "current programme" line under each Xtream Live TV channel used to be +written only when a row scrolled into view, when an EPG result arrived, or +when the display offset changed. Nothing re-evaluated it as wall-clock time +passed, so once a programme ended the row stayed on it until the category was +left and re-entered (#767). + +`PortalChannelsListComponent` re-checks the rows on screen once a minute. The +rules, each of which exists to avoid a specific failure: + +- A programme still on air only has its **progress bar** advanced. No cache + read, no request: a row with nothing to learn must not cost traffic. +- Once it ends the row is re-picked from the queue's cache, and only a + programme that is **on air or upcoming** may replace it. A finished + programme is never re-applied — it would keep presenting itself as current, + and the earliest-item fallback that fills a blank row on first paint would + move an advanced row *backwards*. +- A cached guide whose programmes have **all** ended is dropped + (`EpgQueueService.invalidate`) and refetched, because the queue skips any + stream that still holds a cached answer. An **empty** answer means the + provider has no guide for that channel and is left alone; re-asking would + put one call per EPG-less visible row on the wire every minute. +- What is on screen stays there until a replacement arrives, so a refreshing + row never blanks out. + +A programme occupies `[start, stop)` in every one of these comparisons, so +"has it ended" and "what is on air" cannot disagree on the boundary instant. +The selection rules are pure functions in +`libs/portal/xtream/feature/src/lib/portal-channels-list/epg-preview-program.ts` +and take an explicit `nowMs` in the PROVIDER's clock (`epgProviderClockMs`), +never `Date.now()`. + +Two root-provided services exist because a live layout mounts the channel list +more than once — the sidebar and the fullscreen channel panel render side by +side — over one shared `EpgQueueService`: + +- `EpgRefillLimiter` is the floor on dropping an exhausted cache, the one + place that overrides the queue's own throttling. A provider whose guide has + genuinely run out answers the refill with the same finished programmes, so + without a floor the row would ask again on the very next tick. Records carry + the owning playlist, since a stream id is provider-local and the service + outlives a playlist switch, and they expire by age rather than by viewport + membership — a claim dropped when its row scrolled away would be handed back + the moment the user scrolled to it again. +- `EpgRefreshCoordinator` owns the single timer and merges what every mounted + list needs into one queue request. `EpgQueueService.enqueue` is latest-wins: + it bumps one generation, replaces the queue and the visible set, and drops an + earlier caller's entries after its XMLTV await. Separate timers would cancel + each other whenever both lists had rows to fill, which is exactly what + happens on a programme boundary. Each list still decides for itself what is + stale (that reads only its own state) and contributes its whole visible + slice, since the queue drops anything outside the visible set it was last + handed; a channel both lists show is fetched once, and playlists stay apart + because their credentials differ. + ### XMLTV response compression The Electron EPG worker decodes HTTP `Content-Encoding` layers in reverse diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-preview-program.spec.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-preview-program.spec.ts new file mode 100644 index 000000000..49b098ab5 --- /dev/null +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-preview-program.spec.ts @@ -0,0 +1,153 @@ +import { EpgItem, EpgProgram } from '@iptvnator/shared/interfaces'; +import { + epgProgramProgressPercent, + hasEpgProgramEnded, + pickAiringOrUpcomingEpgItem, + pickEpgPreviewItem, + toSharedEpgProgram, +} from './epg-preview-program'; + +function item(title: string, startIso: string, stopIso: string): EpgItem { + return { + id: title, + epg_id: `epg-${title}`, + title, + lang: 'en', + start: startIso, + end: stopIso, + stop: stopIso, + description: '', + channel_id: 'channel-1', + start_timestamp: String(Math.floor(Date.parse(startIso) / 1000)), + stop_timestamp: String(Math.floor(Date.parse(stopIso) / 1000)), + }; +} + +const at = (iso: string) => Date.parse(iso); + +describe('pickAiringOrUpcomingEpgItem', () => { + const early = item('Early', '2026-04-05T05:30:00Z', '2026-04-05T06:00:00Z'); + const late = item('Late', '2026-04-05T06:00:00Z', '2026-04-05T06:30:00Z'); + + it('returns the program on air', () => { + expect( + pickAiringOrUpcomingEpgItem( + [late, early], + at('2026-04-05T05:45:00Z') + ) + ).toBe(early); + }); + + it('returns the next program to start when nothing is on air yet', () => { + expect( + pickAiringOrUpcomingEpgItem([late], at('2026-04-05T05:45:00Z')) + ).toBe(late); + }); + + it('returns null once every program has ended', () => { + // The caller must request fresh data rather than present a finished + // program as the current one (#767). + expect( + pickAiringOrUpcomingEpgItem( + [early, late], + at('2026-04-05T07:00:00Z') + ) + ).toBeNull(); + }); + + it('returns null for an empty guide', () => { + expect( + pickAiringOrUpcomingEpgItem([], at('2026-04-05T07:00:00Z')) + ).toBeNull(); + }); + + it('hands the boundary instant to the program that starts on it', () => { + // A program occupies [start, stop), the same rule hasEpgProgramEnded + // applies -- otherwise the two disagree at the boundary and the row + // re-applies the finished program for another minute. + expect( + pickAiringOrUpcomingEpgItem( + [early, late], + at('2026-04-05T06:00:00Z') + ) + ).toBe(late); + }); +}); + +describe('pickEpgPreviewItem', () => { + it('falls back to the earliest item so a first paint is never blank', () => { + const early = item( + 'Early', + '2026-04-05T05:00:00Z', + '2026-04-05T05:30:00Z' + ); + const late = item( + 'Late', + '2026-04-05T05:30:00Z', + '2026-04-05T06:00:00Z' + ); + + expect( + pickEpgPreviewItem([late, early], at('2026-04-05T07:00:00Z')) + ).toBe(early); + }); +}); + +describe('hasEpgProgramEnded', () => { + const shown = toSharedEpgProgram( + item('Show', '2026-04-05T05:30:00Z', '2026-04-05T06:00:00Z') + ); + + it('is false while the program runs and true at its stop time', () => { + expect(hasEpgProgramEnded(shown, at('2026-04-05T05:59:00Z'))).toBe( + false + ); + expect(hasEpgProgramEnded(shown, at('2026-04-05T06:00:00Z'))).toBe( + true + ); + }); + + it('keeps a program with an unreadable end time on screen', () => { + // Answering "ended" would re-request this row once a minute forever. + const broken: EpgProgram = { + ...shown, + stop: 'not-a-date', + stopTimestamp: null, + }; + + expect(hasEpgProgramEnded(broken, at('2026-04-05T09:00:00Z'))).toBe( + false + ); + }); +}); + +describe('epgProgramProgressPercent', () => { + const shown = toSharedEpgProgram( + item('Show', '2026-04-05T05:00:00Z', '2026-04-05T07:00:00Z') + ); + + it('reports the elapsed share of the running program', () => { + expect( + epgProgramProgressPercent(shown, at('2026-04-05T06:00:00Z')) + ).toBeCloseTo(50, 5); + }); + + it('reports nothing outside the program', () => { + expect( + epgProgramProgressPercent(shown, at('2026-04-05T04:00:00Z')) + ).toBeNull(); + expect( + epgProgramProgressPercent(shown, at('2026-04-05T08:00:00Z')) + ).toBeNull(); + }); + + it('reports nothing for a zero-length program', () => { + const instant = toSharedEpgProgram( + item('Instant', '2026-04-05T06:00:00Z', '2026-04-05T06:00:00Z') + ); + + expect( + epgProgramProgressPercent(instant, at('2026-04-05T06:00:00Z')) + ).toBeNull(); + }); +}); diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-preview-program.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-preview-program.ts new file mode 100644 index 000000000..e5f6167c1 --- /dev/null +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-preview-program.ts @@ -0,0 +1,168 @@ +import { EpgItem, EpgProgram } from '@iptvnator/shared/interfaces'; + +/** + * Selection rules for the "current program" line under a channel row, kept + * out of the component so both the scroll-driven fill and the periodic + * refresh (#767) read the guide the same way. + * + * Every comparison here takes an explicit `nowMs` in the PROVIDER's clock + * (`epgProviderClockMs`), never `Date.now()` — the EPG display offset moves + * "now", not the programmes. + * + * A programme occupies `[start, stop)`: at exactly its stop time it is over + * and its successor has begun. One rule for every function below, so the + * "has it ended" and "what is on air" answers cannot disagree on the + * boundary and leave the row a minute behind the rest of the EPG surfaces. + */ + +/** + * Milliseconds for one guide boundary. The provider's unix timestamp wins + * when it is present; the ISO string is the fallback. `NaN` when neither + * form parses, so callers must check before comparing. + */ +export function epgBoundaryMs( + dateValue: string | undefined, + unixTimestampValue: string | undefined +): number { + const unixTimestamp = Number(unixTimestampValue); + if (Number.isFinite(unixTimestamp) && unixTimestamp > 0) { + return unixTimestamp * 1000; + } + + return new Date(dateValue ?? '').getTime(); +} + +/** Same boundary in whole seconds, or `null` when neither form parses. */ +export function epgBoundarySeconds( + dateValue: string | undefined, + unixTimestampValue: string | undefined +): number | null { + const boundaryMs = epgBoundaryMs(dateValue, unixTimestampValue); + return Number.isFinite(boundaryMs) ? Math.floor(boundaryMs / 1000) : null; +} + +function epgItemStartMs(item: EpgItem): number { + return epgBoundaryMs(item.start, item.start_timestamp); +} + +function epgItemEndMs(item: EpgItem): number { + return epgBoundaryMs(item.stop ?? item.end, item.stop_timestamp); +} + +/** + * The programme on air at `nowMs`, else the next one to start. + * + * `null` when every item has already ended: a finished programme must never + * be presented as the current one, which is the whole point of the periodic + * refresh. Returning the newest ended item instead would be just as wrong — + * the row would still name a programme that is over. + */ +export function pickAiringOrUpcomingEpgItem( + items: readonly EpgItem[], + nowMs: number +): EpgItem | null { + if (items.length === 0) { + return null; + } + + const byStart = [...items].sort( + (left, right) => epgItemStartMs(left) - epgItemStartMs(right) + ); + + const airing = byStart.find( + (item) => nowMs >= epgItemStartMs(item) && nowMs < epgItemEndMs(item) + ); + + return ( + airing ?? byStart.find((item) => epgItemStartMs(item) > nowMs) ?? null + ); +} + +/** + * The programme to show for a row that has none yet. Falls back to the + * earliest known item so a first paint shows what the guide holds rather + * than an empty line; the refresh path deliberately does NOT use this + * fallback, because re-applying it to a row that already advanced could + * move it backwards to an older programme. + */ +export function pickEpgPreviewItem( + items: readonly EpgItem[], + nowMs: number +): EpgItem | null { + if (items.length === 0) { + return null; + } + + const airingOrUpcoming = pickAiringOrUpcomingEpgItem(items, nowMs); + if (airingOrUpcoming) { + return airingOrUpcoming; + } + + return [...items].sort( + (left, right) => epgItemStartMs(left) - epgItemStartMs(right) + )[0]; +} + +function epgProgramStartMs(program: EpgProgram): number { + return program.startTimestamp != null + ? program.startTimestamp * 1000 + : new Date(program.start ?? '').getTime(); +} + +function epgProgramEndMs(program: EpgProgram): number { + return program.stopTimestamp != null + ? program.stopTimestamp * 1000 + : new Date(program.stop ?? '').getTime(); +} + +/** + * Whether the programme already on screen has run out. + * + * An unreadable end time answers `false`: the row keeps what it has instead + * of asking the provider for a replacement once a minute, forever. + */ +export function hasEpgProgramEnded( + program: EpgProgram, + nowMs: number +): boolean { + const endMs = epgProgramEndMs(program); + return Number.isFinite(endMs) ? endMs <= nowMs : false; +} + +/** Elapsed percentage, or `null` while `nowMs` sits outside the programme. */ +export function epgProgramProgressPercent( + program: EpgProgram, + nowMs: number +): number | null { + const startMs = epgProgramStartMs(program); + const endMs = epgProgramEndMs(program); + if (!Number.isFinite(startMs) || !Number.isFinite(endMs)) { + return null; + } + + if (nowMs < startMs || nowMs >= endMs) { + return null; + } + + return ((nowMs - startMs) / (endMs - startMs)) * 100; +} + +/** Guide item in the shape the shared channel-row components render. */ +export function toSharedEpgProgram(program: EpgItem): EpgProgram { + return { + start: program.start, + stop: program.stop ?? program.end, + channel: program.channel_id ?? program.id, + title: program.title, + desc: program.description ?? null, + category: null, + startTimestamp: epgBoundarySeconds( + program.start, + program.start_timestamp + ), + stopTimestamp: epgBoundarySeconds( + program.stop ?? program.end, + program.stop_timestamp + ), + }; +} diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refill-limiter.service.spec.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refill-limiter.service.spec.ts new file mode 100644 index 000000000..882a6be0d --- /dev/null +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refill-limiter.service.spec.ts @@ -0,0 +1,78 @@ +import { TestBed } from '@angular/core/testing'; +import { + EPG_REFILL_MIN_INTERVAL_MS, + EpgRefillLimiter, +} from './epg-refill-limiter.service'; + +describe('EpgRefillLimiter', () => { + let limiter: EpgRefillLimiter; + + beforeEach(() => { + TestBed.configureTestingModule({}); + limiter = TestBed.inject(EpgRefillLimiter); + }); + + it('is one shared instance for every mounted channel list', () => { + // The sidebar list and the fullscreen channel panel are mounted at + // once over a single EPG queue; a record per component would give + // each copy its own allowance and divide the floor between them. + expect(TestBed.inject(EpgRefillLimiter)).toBe(limiter); + }); + + it('allows one claim per interval and then blocks', () => { + expect(limiter.claim('playlist-1', 1, 0)).toBe(true); + expect( + limiter.claim('playlist-1', 1, EPG_REFILL_MIN_INTERVAL_MS - 1) + ).toBe(false); + expect(limiter.claim('playlist-1', 1, EPG_REFILL_MIN_INTERVAL_MS)).toBe( + true + ); + }); + + it('tracks each stream separately', () => { + expect(limiter.claim('playlist-1', 1, 0)).toBe(true); + expect(limiter.claim('playlist-1', 2, 0)).toBe(true); + }); + + it('lets a released stream claim again straight away', () => { + limiter.claim('playlist-1', 1, 0); + limiter.release('playlist-1', 1); + + expect(limiter.claim('playlist-1', 1, 1)).toBe(true); + }); + + it('forgets records only once they have expired', () => { + limiter.claim('playlist-1', 1, 0); + + limiter.forgetExpired(EPG_REFILL_MIN_INTERVAL_MS - 1); + expect( + limiter.claim('playlist-1', 1, EPG_REFILL_MIN_INTERVAL_MS - 1) + ).toBe(false); + + limiter.forgetExpired(EPG_REFILL_MIN_INTERVAL_MS); + expect(limiter.claim('playlist-1', 1, EPG_REFILL_MIN_INTERVAL_MS)).toBe( + true + ); + }); + + it('does not let one playlist hold back another playlist same stream id', () => { + // Stream ids are provider-local and the root-provided service outlives + // a playlist switch, so two accounts numbering a channel alike must + // not share an allowance. + expect(limiter.claim('playlist-1', 1, 0)).toBe(true); + + expect(limiter.claim('playlist-2', 1, 0)).toBe(true); + }); + + it('holds a claim across a row scrolling out of view and back', () => { + // Housekeeping runs on every tick while the row is off screen; the + // floor must survive it, or scrolling up and down the list would hand + // out a fresh request each pass. + limiter.claim('playlist-1', 1, 0); + + limiter.forgetExpired(60_000); + limiter.forgetExpired(120_000); + + expect(limiter.claim('playlist-1', 1, 180_000)).toBe(false); + }); +}); diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refill-limiter.service.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refill-limiter.service.ts new file mode 100644 index 000000000..5418a6334 --- /dev/null +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refill-limiter.service.ts @@ -0,0 +1,82 @@ +import { Injectable } from '@angular/core'; + +/** + * Floor between two refills of the same channel's exhausted guide. Matches + * the EPG queue's own cache lifetime, so a provider with no fresh data is + * asked no more often than its cached answer would have expired anyway. + */ +export const EPG_REFILL_MIN_INTERVAL_MS = 5 * 60_000; + +/** Stream ids are provider-local, so a record belongs to one playlist. */ +function refillKey( + playlistId: string | null | undefined, + streamId: number +): string { + return `${playlistId ?? ''}:${streamId}`; +} + +/** + * Rate limit for the one place that overrides the EPG queue's own throttling: + * dropping a cached guide whose programmes have all ended so that it can be + * fetched again. + * + * A provider whose guide has genuinely run out answers that refill with the + * same finished programmes, which would leave the row stale and ask again on + * the very next tick — one request per visible channel per minute, against the + * queue that exists to keep providers from banning the client. + * + * Root-provided on purpose: a live layout mounts the channel list more than + * once (the sidebar and the fullscreen channel panel render side by side) + * over one shared queue, so a per-component record would hand every mounted + * copy its own allowance and divide the floor between them. + * + * Because it is root-provided it also outlives a playlist switch, so records + * carry the owning playlist: a stream id is provider-local, and two Xtream + * accounts routinely number their channels alike. A bare id would let one + * account's claim hold back a channel of the account now on screen. + */ +@Injectable({ providedIn: 'root' }) +export class EpgRefillLimiter { + private readonly requestedAt = new Map(); + + /** True when this stream may be refilled now, recording the attempt. */ + claim( + playlistId: string | null | undefined, + streamId: number, + wallClockMs: number + ): boolean { + const key = refillKey(playlistId, streamId); + const previous = this.requestedAt.get(key); + if ( + previous !== undefined && + wallClockMs - previous < EPG_REFILL_MIN_INTERVAL_MS + ) { + return false; + } + + this.requestedAt.set(key, wallClockMs); + return true; + } + + /** The guide is flowing again, so the next gap may refill immediately. */ + release(playlistId: string | null | undefined, streamId: number): void { + this.requestedAt.delete(refillKey(playlistId, streamId)); + } + + /** + * Drops records that have outlived the interval and so no longer hold + * anything back, which is what keeps the map bounded. + * + * Deliberately not keyed on the viewport: a claim dropped when its row + * scrolls out of sight would be handed back the moment the user scrolled + * to it again, and a few passes up and down the list would bypass the + * floor entirely. + */ + forgetExpired(wallClockMs: number): void { + for (const [streamId, requestedAt] of this.requestedAt) { + if (wallClockMs - requestedAt >= EPG_REFILL_MIN_INTERVAL_MS) { + this.requestedAt.delete(streamId); + } + } + } +} diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refresh-coordinator.service.spec.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refresh-coordinator.service.spec.ts new file mode 100644 index 000000000..e2a0ffd67 --- /dev/null +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refresh-coordinator.service.spec.ts @@ -0,0 +1,132 @@ +import { TestBed } from '@angular/core/testing'; +import { EpgQueueService } from '@iptvnator/portal/xtream/data-access'; +import { + EPG_REFRESH_INTERVAL_MS, + EpgRefreshContribution, + EpgRefreshCoordinator, +} from './epg-refresh-coordinator.service'; + +const credentials = { + serverUrl: 'http://demo.example', + username: 'demo', + password: 'secret', +}; + +function contribution( + overrides: Partial = {} +): EpgRefreshContribution { + return { + playlistId: 'playlist-1', + credentials, + visibleStreamIds: [1], + staleEntries: [{ streamId: 1 }], + ...overrides, + }; +} + +describe('EpgRefreshCoordinator', () => { + const epgQueueService = { + enqueue: jest.fn().mockResolvedValue(undefined), + }; + let coordinator: EpgRefreshCoordinator; + + beforeEach(() => { + jest.useFakeTimers(); + epgQueueService.enqueue.mockClear(); + TestBed.configureTestingModule({ + providers: [ + { provide: EpgQueueService, useValue: epgQueueService }, + ], + }); + coordinator = TestBed.inject(EpgRefreshCoordinator); + }); + + afterEach(() => { + jest.useRealTimers(); + }); + + it('merges what every mounted list needs into one request', () => { + // Separate requests would cancel one another: the queue is latest-wins + // and replaces its visible set and pending queue on every call. + coordinator.register(() => + contribution({ + visibleStreamIds: [1, 2], + staleEntries: [{ streamId: 1 }], + }) + ); + coordinator.register(() => + contribution({ + visibleStreamIds: [3, 4], + staleEntries: [{ streamId: 3 }], + }) + ); + + jest.advanceTimersByTime(EPG_REFRESH_INTERVAL_MS); + + expect(epgQueueService.enqueue).toHaveBeenCalledTimes(1); + const [entries, visibleIds] = epgQueueService.enqueue.mock.calls[0]; + expect(entries.map((entry: { streamId: number }) => entry.streamId)) // + .toEqual([1, 3]); + // The whole of both slices, so neither list's rows are dropped while + // the queue works through them. + expect([...(visibleIds as Set)].sort()).toEqual([1, 2, 3, 4]); + }); + + it('fetches a channel both lists show only once', () => { + coordinator.register(() => contribution()); + coordinator.register(() => contribution()); + + jest.advanceTimersByTime(EPG_REFRESH_INTERVAL_MS); + + const [entries] = epgQueueService.enqueue.mock.calls[0]; + expect(entries).toHaveLength(1); + }); + + it('keeps playlists apart, since their credentials differ', () => { + coordinator.register(() => contribution({ playlistId: 'playlist-1' })); + coordinator.register(() => + contribution({ + playlistId: 'playlist-2', + staleEntries: [{ streamId: 9 }], + }) + ); + + jest.advanceTimersByTime(EPG_REFRESH_INTERVAL_MS); + + expect(epgQueueService.enqueue).toHaveBeenCalledTimes(2); + }); + + it('asks for nothing when no list needs anything', () => { + coordinator.register(() => null); + coordinator.register(() => contribution({ staleEntries: [] })); + + jest.advanceTimersByTime(EPG_REFRESH_INTERVAL_MS); + + expect(epgQueueService.enqueue).not.toHaveBeenCalled(); + }); + + it('stops ticking once the last list has left', () => { + const leaveFirst = coordinator.register(() => contribution()); + const leaveSecond = coordinator.register(() => contribution()); + + leaveFirst(); + jest.advanceTimersByTime(EPG_REFRESH_INTERVAL_MS); + expect(epgQueueService.enqueue).toHaveBeenCalledTimes(1); + + leaveSecond(); + jest.advanceTimersByTime(5 * EPG_REFRESH_INTERVAL_MS); + expect(epgQueueService.enqueue).toHaveBeenCalledTimes(1); + }); + + it('survives a rejected enqueue', async () => { + epgQueueService.enqueue.mockRejectedValueOnce(new Error('offline')); + coordinator.register(() => contribution()); + + jest.advanceTimersByTime(EPG_REFRESH_INTERVAL_MS); + await Promise.resolve(); + + epgQueueService.enqueue.mockResolvedValue(undefined); + jest.advanceTimersByTime(EPG_REFRESH_INTERVAL_MS); + expect(epgQueueService.enqueue).toHaveBeenCalledTimes(2); + }); +}); diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refresh-coordinator.service.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refresh-coordinator.service.ts new file mode 100644 index 000000000..c65b1a4ae --- /dev/null +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/epg-refresh-coordinator.service.ts @@ -0,0 +1,148 @@ +import { Injectable, inject } from '@angular/core'; +import { + EpgQueueService, + XtreamCredentials, +} from '@iptvnator/portal/xtream/data-access'; +import { createLogger } from '@iptvnator/portal/shared/util'; + +/** How often the rows on screen re-check the programme they are showing. */ +export const EPG_REFRESH_INTERVAL_MS = 60_000; + +/** What one mounted channel list needs fetched, and the rows it is showing. */ +export interface EpgRefreshContribution { + readonly playlistId: string | null | undefined; + readonly credentials: XtreamCredentials; + /** Every row on screen, whether or not it needs anything. */ + readonly visibleStreamIds: readonly number[]; + readonly staleEntries: readonly { + streamId: number; + epgChannelId?: string | null; + playlistId?: string | null; + }[]; +} + +/** Credentials for the queue, from the playlist row the store holds. */ +export function xtreamCredentialsOf(playlist: { + serverUrl: string; + username: string; + password: string; + serverTimezone?: string | null; +}): XtreamCredentials { + return { + serverUrl: playlist.serverUrl, + username: playlist.username, + password: playlist.password, + serverTimezone: playlist.serverTimezone ?? undefined, + }; +} + +/** One queue entry for a channel row. */ +export function epgQueueEntryFor( + channel: { xtream_id: number; epg_channel_id?: string | null }, + playlistId: string | null | undefined +) { + return { + streamId: channel.xtream_id, + epgChannelId: channel.epg_channel_id ?? null, + playlistId: playlistId ?? null, + }; +} + +/** Returns what this list needs on a tick, or `null` when it needs nothing. */ +export type EpgRefreshParticipant = () => EpgRefreshContribution | null; + +/** + * Owns the one timer behind the channel rows' periodic EPG refresh (#767) and + * merges what every mounted list needs into a single queue request. + * + * A live layout mounts the channel list more than once — the sidebar and the + * fullscreen channel panel render side by side — and `EpgQueueService.enqueue` + * is latest-wins: it bumps one generation, replaces the queue and the visible + * set, and drops an earlier caller's entries after its XMLTV await. Two lists + * ticking on their own timers would therefore cancel each other whenever both + * had rows to fill, which is exactly what happens on a programme boundary, + * leaving one of them a minute behind. Rows that share a timer and one merged + * request cannot. + * + * Each list still decides for itself what is stale — that part reads only its + * own state and costs nothing — and contributes its whole visible slice, since + * the queue drops anything outside the visible set it was last handed. + */ +@Injectable({ providedIn: 'root' }) +export class EpgRefreshCoordinator { + private readonly epgQueueService = inject(EpgQueueService); + private readonly logger = createLogger('EpgRefreshCoordinator'); + private readonly participants = new Set(); + private intervalId?: number; + + /** Joins the shared tick. Call the returned function on teardown. */ + register(participant: EpgRefreshParticipant): () => void { + this.participants.add(participant); + if (this.intervalId === undefined) { + this.intervalId = window.setInterval( + () => this.tick(), + EPG_REFRESH_INTERVAL_MS + ); + } + + return () => { + this.participants.delete(participant); + if (this.participants.size === 0 && this.intervalId !== undefined) { + clearInterval(this.intervalId); + this.intervalId = undefined; + } + }; + } + + private tick(): void { + const contributions: EpgRefreshContribution[] = []; + for (const participant of this.participants) { + const contribution = participant(); + if (contribution) { + contributions.push(contribution); + } + } + + // Lists of different playlists never share a queue request: the + // credentials differ, and a stream id only means anything inside its + // own account. In practice one live layout means one playlist. + const byPlaylist = new Map(); + for (const contribution of contributions) { + const key = contribution.playlistId ?? ''; + byPlaylist.set(key, [...(byPlaylist.get(key) ?? []), contribution]); + } + + for (const group of byPlaylist.values()) { + this.requestMerged(group); + } + } + + private requestMerged(group: EpgRefreshContribution[]): void { + const visibleIds = new Set(); + const entries = new Map< + number, + EpgRefreshContribution['staleEntries'][number] + >(); + for (const contribution of group) { + for (const streamId of contribution.visibleStreamIds) { + visibleIds.add(streamId); + } + for (const entry of contribution.staleEntries) { + // Both lists can show the same channel; it is fetched once. + if (!entries.has(entry.streamId)) { + entries.set(entry.streamId, entry); + } + } + } + + if (entries.size === 0) { + return; + } + + this.epgQueueService + .enqueue([...entries.values()], visibleIds, group[0].credentials) + .catch((error) => { + this.logger.warn('EPG refresh enqueue failed', error); + }); + } +} diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.spec.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.spec.ts index 1c40d00b3..7afa25423 100644 --- a/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.spec.ts +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.spec.ts @@ -70,7 +70,8 @@ describe('PortalChannelsListComponent', () => { const epgQueueService = { epgResult$: epgResults$, getCached: jest.fn().mockReturnValue(null), - enqueue: jest.fn(), + enqueue: jest.fn().mockResolvedValue(undefined), + invalidate: jest.fn(), }; beforeEach(async () => { @@ -91,6 +92,7 @@ describe('PortalChannelsListComponent', () => { favoritesService.getFavorites.mockReturnValue(of([] as FavoriteItem[])); epgQueueService.getCached.mockReturnValue(null); epgQueueService.enqueue.mockClear(); + epgQueueService.invalidate.mockClear(); await TestBed.configureTestingModule({ imports: [PortalChannelsListComponent, NoopAnimationsModule], @@ -272,6 +274,504 @@ describe('PortalChannelsListComponent', () => { expect(component.currentProgramsProgress.get(50)).toBeCloseTo(50, 1); }); + /** Sizes the viewport so its rendered range covers the given channels and + * `lastVisibleChannels` gets populated, matching what a real scroll does. */ + async function renderViewport( + fixture: ComponentFixture + ) { + const viewport = fixture.componentInstance.viewport(); + if (!viewport) { + throw new Error('Expected channel viewport'); + } + Object.defineProperty( + viewport.elementRef.nativeElement, + 'clientHeight', + { + configurable: true, + value: 520, + } + ); + viewport.checkViewportSize(); + fixture.detectChanges(); + await fixture.whenStable(); + jest.advanceTimersByTime(300); // renderedRangeStream's debounceTime + fixture.detectChanges(); + } + + it('re-picks a channel current program once its end time passes, without a scroll or re-entry (#767)', async () => { + jest.useFakeTimers(); + const firstStartTimestamp = Math.floor( + Date.parse('2026-04-05T05:30:00.000Z') / 1000 + ); + const firstStopTimestamp = Math.floor( + Date.parse('2026-04-05T06:00:00.000Z') / 1000 + ); + const secondStartTimestamp = firstStopTimestamp; + const secondStopTimestamp = Math.floor( + Date.parse('2026-04-05T06:30:00.000Z') / 1000 + ); + + jest.setSystemTime(new Date('2026-04-05T05:45:00.000Z')); + + selectedTypeContentLoading.set(false); + selectedChannels.set([ + { + title: 'Cartoon Network', + xtream_id: 50, + }, + ]); + currentPlaylist.set({ + id: 'playlist-1', + password: 'secret', + serverUrl: 'http://demo.example', + username: 'demo', + }); + + fixture.detectChanges(); + await renderViewport(fixture); + + const firstProgram = buildEpgItem({ + id: 'first', + title: 'Current Show', + start: '2026-04-05T05:30:00.000Z', + stop: '2026-04-05T06:00:00.000Z', + startTimestamp: firstStartTimestamp, + stopTimestamp: firstStopTimestamp, + }); + const secondProgram = buildEpgItem({ + id: 'second', + title: 'Next Show', + start: '2026-04-05T06:00:00.000Z', + stop: '2026-04-05T06:30:00.000Z', + startTimestamp: secondStartTimestamp, + stopTimestamp: secondStopTimestamp, + }); + + epgResults$.next({ + streamId: 50, + items: [firstProgram, secondProgram], + }); + fixture.detectChanges(); + + const component = fixture.componentInstance; + expect(component.epgPrograms.get(50)?.title).toBe('Current Show'); + + // The queue's cache holds both programmes for this channel, as it + // would once the EPG result above had arrived through the real + // service. + epgQueueService.getCached.mockImplementation((streamId: number) => + streamId === 50 ? [firstProgram, secondProgram] : null + ); + epgQueueService.enqueue.mockClear(); + + // Wall-clock time moves past the first programme's end with no + // scroll-out/in and no EPG settings change: the reporter's exact + // scenario from #767. The channel stays on screen throughout; the + // 60s interval fires several times along the way, landing on the + // next programme by the time 20 minutes have passed. + jest.advanceTimersByTime(20 * 60 * 1000); + + expect(component.epgPrograms.get(50)?.title).toBe('Next Show'); + // The cache was already warm, so the refresh re-picked from it + // instead of enqueuing a new fetch. + expect(epgQueueService.enqueue).not.toHaveBeenCalled(); + }); + + it('re-fetches EPG on the refresh tick once the cache has expired, without blanking the row (#767)', async () => { + // The queue's own cache is short-lived (a 5 minute TTL); a channel + // left on screen without a scroll event needs the periodic refresh + // to re-fetch it, not just re-pick from a cache entry that is gone. + jest.useFakeTimers(); + jest.setSystemTime(new Date('2026-04-05T05:45:00.000Z')); + + selectedTypeContentLoading.set(false); + selectedChannels.set([ + { + title: 'Cartoon Network', + xtream_id: 50, + }, + ]); + currentPlaylist.set({ + id: 'playlist-1', + password: 'secret', + serverUrl: 'http://demo.example', + username: 'demo', + }); + + fixture.detectChanges(); + await renderViewport(fixture); + + epgResults$.next({ + streamId: 50, + items: [ + buildEpgItem({ + id: 'first', + title: 'Current Show', + start: '2026-04-05T05:30:00.000Z', + stop: '2026-04-05T06:00:00.000Z', + startTimestamp: Math.floor( + Date.parse('2026-04-05T05:30:00.000Z') / 1000 + ), + stopTimestamp: Math.floor( + Date.parse('2026-04-05T06:00:00.000Z') / 1000 + ), + }), + ], + }); + fixture.detectChanges(); + const component = fixture.componentInstance; + expect(component.epgPrograms.get(50)?.title).toBe('Current Show'); + + // Simulate the cache having aged out: getCached now returns null. + epgQueueService.getCached.mockReturnValue(null); + epgQueueService.enqueue.mockClear(); + + jest.advanceTimersByTime(20 * 60 * 1000); + + expect(epgQueueService.enqueue).toHaveBeenCalledWith( + expect.arrayContaining([expect.objectContaining({ streamId: 50 })]), + expect.any(Set), + expect.objectContaining({ serverUrl: 'http://demo.example' }) + ); + // The row is never cleared while the fetch is in flight: it keeps + // showing the last-known program instead of going blank. + expect(component.epgPrograms.get(50)?.title).toBe('Current Show'); + + // The fetch resolves with the next programme; the row updates. + epgResults$.next({ + streamId: 50, + items: [ + buildEpgItem({ + id: 'second', + title: 'Next Show', + start: '2026-04-05T06:00:00.000Z', + stop: '2026-04-05T06:30:00.000Z', + startTimestamp: Math.floor( + Date.parse('2026-04-05T06:00:00.000Z') / 1000 + ), + stopTimestamp: Math.floor( + Date.parse('2026-04-05T06:30:00.000Z') / 1000 + ), + }), + ], + }); + fixture.detectChanges(); + expect(component.epgPrograms.get(50)?.title).toBe('Next Show'); + }); + + it('does not fetch off-screen channels on the refresh tick when nothing has been scrolled into view', () => { + // No renderViewport() here: lastVisibleChannels stays empty, as it + // would before the first scroll event settles. + jest.useFakeTimers(); + selectedTypeContentLoading.set(false); + selectedChannels.set( + Array.from({ length: 60 }, (_, index) => ({ + title: `Channel ${index + 1}`, + xtream_id: index + 1, + })) + ); + currentPlaylist.set({ + id: 'playlist-1', + password: 'secret', + serverUrl: 'http://demo.example', + username: 'demo', + }); + + fixture.detectChanges(); + epgQueueService.enqueue.mockClear(); + + jest.advanceTimersByTime(60_000); + + expect(epgQueueService.enqueue).not.toHaveBeenCalled(); + }); + + /** + * One channel on screen, the clock parked at `nowIso`. + * + * `keepTimers` is for a second list mounted beside the first: re-installing + * the fake clock would drop the intervals the first one had registered, and + * the test would silently observe a single component. + */ + async function renderSingleChannel( + fixture: ComponentFixture, + nowIso: string, + options: { keepTimers?: boolean } = {} + ) { + if (!options.keepTimers) { + jest.useFakeTimers(); + } + jest.setSystemTime(new Date(nowIso)); + selectedTypeContentLoading.set(false); + selectedChannels.set([{ title: 'Cartoon Network', xtream_id: 50 }]); + currentPlaylist.set({ + id: 'playlist-1', + password: 'secret', + serverUrl: 'http://demo.example', + username: 'demo', + }); + + fixture.detectChanges(); + await renderViewport(fixture); + epgQueueService.enqueue.mockClear(); + epgQueueService.invalidate.mockClear(); + } + + function buildProgram(title: string, startIso: string, stopIso: string) { + return buildEpgItem({ + id: title, + title, + start: startIso, + stop: stopIso, + startTimestamp: Math.floor(Date.parse(startIso) / 1000), + stopTimestamp: Math.floor(Date.parse(stopIso) / 1000), + }); + } + + it('never falls back to a finished program once the cached guide runs out (#767)', async () => { + // The queue caches a short EPG for five minutes, so a channel whose + // listings are shorter than that reaches a state where every cached + // program has already ended -- the shape an uploaded-XMLTV fallback + // always has, since it caches the single current program. + await renderSingleChannel(fixture, '2026-04-05T05:45:00.000Z'); + + const early = buildProgram( + 'Early Show', + '2026-04-05T05:30:00.000Z', + '2026-04-05T06:00:00.000Z' + ); + const short = buildProgram( + 'Short Show', + '2026-04-05T06:00:00.000Z', + '2026-04-05T06:02:00.000Z' + ); + epgQueueService.getCached.mockImplementation((streamId: number) => + streamId === 50 ? [early, short] : null + ); + + epgResults$.next({ streamId: 50, items: [early, short] }); + fixture.detectChanges(); + const component = fixture.componentInstance; + expect(component.epgPrograms.get(50)?.title).toBe('Early Show'); + + // 05:45 -> 06:03: both cached programs are now over. The row must not + // rewind to the oldest one, which is what the preview pick returns + // when nothing is on air. + jest.advanceTimersByTime(18 * 60 * 1000); + + expect(component.epgPrograms.get(50)?.title).toBe('Short Show'); + // The exhausted entry is dropped, because the queue skips any stream + // that still has a cached answer, and fresh data is requested. + expect(epgQueueService.invalidate).toHaveBeenCalledWith(50); + expect(epgQueueService.enqueue).toHaveBeenCalledWith( + expect.arrayContaining([expect.objectContaining({ streamId: 50 })]), + expect.any(Set), + expect.objectContaining({ serverUrl: 'http://demo.example' }) + ); + + epgResults$.next({ + streamId: 50, + items: [ + buildProgram( + 'Live Show', + '2026-04-05T06:02:00.000Z', + '2026-04-05T06:40:00.000Z' + ), + ], + }); + fixture.detectChanges(); + expect(component.epgPrograms.get(50)?.title).toBe('Live Show'); + }); + + it('keeps the last program when a refill answers with the same stale guide', async () => { + // The provider has simply stopped publishing for this channel, so the + // refill returns the window that is already over. Feeding that through + // the first-paint fallback would walk the row back to the oldest entry. + await renderSingleChannel(fixture, '2026-04-05T05:45:00.000Z'); + + const early = buildProgram( + 'Early Show', + '2026-04-05T05:30:00.000Z', + '2026-04-05T06:00:00.000Z' + ); + const short = buildProgram( + 'Short Show', + '2026-04-05T06:00:00.000Z', + '2026-04-05T06:02:00.000Z' + ); + epgQueueService.getCached.mockImplementation((streamId: number) => + streamId === 50 ? [early, short] : null + ); + epgResults$.next({ streamId: 50, items: [early, short] }); + fixture.detectChanges(); + + jest.advanceTimersByTime(18 * 60 * 1000); + const component = fixture.componentInstance; + expect(component.epgPrograms.get(50)?.title).toBe('Short Show'); + + epgResults$.next({ streamId: 50, items: [early, short] }); + fixture.detectChanges(); + + expect(component.epgPrograms.get(50)?.title).toBe('Short Show'); + }); + + it('refills an exhausted guide at most once per cache lifetime', async () => { + // A provider whose guide has genuinely run out answers the refill with + // the same finished program, so without a floor the row would drop and + // re-request its cache on every single tick. + await renderSingleChannel(fixture, '2026-04-05T06:03:00.000Z'); + + const finished = buildProgram( + 'Finished Show', + '2026-04-05T05:30:00.000Z', + '2026-04-05T06:00:00.000Z' + ); + epgQueueService.getCached.mockImplementation((streamId: number) => + streamId === 50 ? [finished] : null + ); + epgResults$.next({ streamId: 50, items: [finished] }); + fixture.detectChanges(); + + jest.advanceTimersByTime(4 * 60 * 1000); + expect(epgQueueService.invalidate).toHaveBeenCalledTimes(1); + + jest.advanceTimersByTime(2 * 60 * 1000); + expect(epgQueueService.invalidate).toHaveBeenCalledTimes(2); + }); + + it('shares the refill floor between the sidebar list and its fullscreen copy', async () => { + // A live layout mounts this component more than once over one EPG + // queue. A refill record per component would hand each copy its own + // allowance, so the same exhausted guide would be dropped and + // refetched once per mounted list every minute. + await renderSingleChannel(fixture, '2026-04-05T06:03:00.000Z'); + const fullscreenCopy = TestBed.createComponent( + PortalChannelsListComponent + ); + await renderSingleChannel(fullscreenCopy, '2026-04-05T06:03:00.000Z', { + keepTimers: true, + }); + + const finished = buildProgram( + 'Finished Show', + '2026-04-05T05:30:00.000Z', + '2026-04-05T06:00:00.000Z' + ); + epgQueueService.getCached.mockImplementation((streamId: number) => + streamId === 50 ? [finished] : null + ); + epgResults$.next({ streamId: 50, items: [finished] }); + fixture.detectChanges(); + fullscreenCopy.detectChanges(); + + jest.advanceTimersByTime(4 * 60 * 1000); + + expect(epgQueueService.invalidate).toHaveBeenCalledTimes(1); + }); + + it('refreshes the rows on screen now, not the ones a stale snapshot holds', async () => { + // The viewport's rendered range only emits when the INDEX range + // changes, so swapping a one-channel category for another leaves the + // snapshot pointing at the previous list while the indices still + // describe what is rendered. + await renderSingleChannel(fixture, '2026-04-05T06:03:00.000Z'); + + selectedChannels.set([{ title: 'Boomerang', xtream_id: 77 }]); + fixture.detectChanges(); + epgQueueService.getCached.mockReturnValue(null); + epgQueueService.enqueue.mockClear(); + + jest.advanceTimersByTime(60_000); + + expect(epgQueueService.enqueue).toHaveBeenCalledWith( + [expect.objectContaining({ streamId: 77 })], + new Set([77]), + expect.objectContaining({ serverUrl: 'http://demo.example' }) + ); + }); + + it('refreshes nothing while a search shows no results', async () => { + // A search with no matches destroys the viewport, so there is no + // rendered range at all; naming the rows it used to hold would keep + // fetching channels that are off screen. + await renderSingleChannel(fixture, '2026-04-05T06:03:00.000Z'); + epgQueueService.getCached.mockReturnValue(null); + + fixture.componentRef.setInput('searchTermInput', 'no-such-channel'); + fixture.detectChanges(); + epgQueueService.enqueue.mockClear(); + + jest.advanceTimersByTime(5 * 60 * 1000); + + expect(epgQueueService.enqueue).not.toHaveBeenCalled(); + expect(epgQueueService.invalidate).not.toHaveBeenCalled(); + }); + + it('does not walk a row back when the selected channel guide has run out', async () => { + // Selecting a row loads the provider's full guide. When that guide + // holds only finished programmes, the first-paint fallback would pick + // the earliest of them and undo what the refresh advanced to. + await renderSingleChannel(fixture, '2026-04-05T06:03:00.000Z'); + + const early = buildProgram( + 'Early Show', + '2026-04-05T05:30:00.000Z', + '2026-04-05T06:00:00.000Z' + ); + const short = buildProgram( + 'Short Show', + '2026-04-05T06:00:00.000Z', + '2026-04-05T06:02:00.000Z' + ); + epgResults$.next({ streamId: 50, items: [short] }); + fixture.detectChanges(); + const component = fixture.componentInstance; + expect(component.epgPrograms.get(50)?.title).toBe('Short Show'); + + selectedItem.set({ xtream_id: 50 }); + epgItems.set([early, short]); + fixture.detectChanges(); + + expect(component.epgPrograms.get(50)?.title).toBe('Short Show'); + }); + + it('leaves a channel the provider has no EPG for alone', async () => { + // An empty answer is cached deliberately; re-requesting it would put + // one call per EPG-less visible row on the wire every minute. + await renderSingleChannel(fixture, '2026-04-05T06:03:00.000Z'); + epgQueueService.getCached.mockReturnValue([]); + + jest.advanceTimersByTime(5 * 60 * 1000); + + expect(epgQueueService.enqueue).not.toHaveBeenCalled(); + expect(epgQueueService.invalidate).not.toHaveBeenCalled(); + }); + + it('advances the progress bar of a running program without touching the queue', async () => { + await renderSingleChannel(fixture, '2026-04-05T06:00:00.000Z'); + epgQueueService.getCached.mockReturnValue(null); + + epgResults$.next({ + streamId: 50, + items: [ + buildProgram( + 'Long Show', + '2026-04-05T05:00:00.000Z', + '2026-04-05T07:00:00.000Z' + ), + ], + }); + fixture.detectChanges(); + const component = fixture.componentInstance; + expect(component.currentProgramsProgress.get(50)).toBeCloseTo(50, 1); + epgQueueService.getCached.mockClear(); + + jest.advanceTimersByTime(30 * 60 * 1000); + + expect(component.currentProgramsProgress.get(50)).toBeCloseTo(75, 1); + expect(epgQueueService.getCached).not.toHaveBeenCalled(); + expect(epgQueueService.enqueue).not.toHaveBeenCalled(); + }); + it('does not derive or subscribe to row EPG previews in browser/PWA mode', async () => { Object.defineProperty(window, 'electron', { configurable: true, diff --git a/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.ts b/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.ts index 6df4a38cf..a8723a295 100644 --- a/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.ts +++ b/libs/portal/xtream/feature/src/lib/portal-channels-list/portal-channels-list.component.ts @@ -42,13 +42,13 @@ import { EpgMappingDialogComponent, } from '@iptvnator/ui/components'; import { + createLogger, getXtreamCatchupDays, isXtreamCatchupAvailable, PortalChannelSortMode, sortPortalChannelItems, } from '@iptvnator/portal/shared/util'; import { EpgQueueService } from '@iptvnator/portal/xtream/data-access'; -import { XtreamCredentials } from '@iptvnator/portal/xtream/data-access'; import { FavoritesService } from '@iptvnator/portal/xtream/data-access'; import { XtreamStore } from '@iptvnator/portal/xtream/data-access'; import { @@ -56,6 +56,20 @@ import { RuntimeCapabilitiesService, SettingsStore, } from '@iptvnator/services'; +import { + epgProgramProgressPercent, + hasEpgProgramEnded, + pickAiringOrUpcomingEpgItem, + pickEpgPreviewItem, + toSharedEpgProgram, +} from './epg-preview-program'; +import { EpgRefillLimiter } from './epg-refill-limiter.service'; +import { + epgQueueEntryFor, + EpgRefreshContribution, + EpgRefreshCoordinator, + xtreamCredentialsOf, +} from './epg-refresh-coordinator.service'; import { XtreamFavoriteMarksService } from './xtream-favorite-marks.service'; export interface XtreamChannelListItem { @@ -165,8 +179,18 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { epgPrograms = new Map(); currentProgramsProgress = new Map(); - /** Last viewport slice, reused to refresh previews after a mapping change. */ - private lastVisibleChannels: XtreamChannelListItem[] = []; + /** Leaves the shared refresh tick; called from `ngOnDestroy`. */ + private leaveEpgRefresh?: () => void; + + private readonly epgRefreshCoordinator = inject(EpgRefreshCoordinator); + private readonly logger = createLogger('XtreamPortalChannelsList'); + + /** + * Shared, not per-instance: a live layout mounts this list more than once + * (sidebar plus the fullscreen channel panel) over one EPG queue, so a + * local record would give each copy its own refill allowance. + */ + private readonly epgRefill = inject(EpgRefillLimiter); readonly viewport = viewChild(CdkVirtualScrollViewport); @@ -275,12 +299,15 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { return; } - const previewProgram = this.pickPreviewProgram(epgItems); - if (!previewProgram) { - return; - } + // Same rule as a queue result: selecting a row must not walk it + // back when the provider's guide holds only finished programmes. + const program = this.pickProgramFor( + selectedItem.xtream_id, + epgItems + ); + if (!program) return; - this.applyProgram(selectedItem.xtream_id, previewProgram); + this.applyProgram(selectedItem.xtream_id, program); }); } @@ -330,13 +357,26 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { this.subscriptions.add( this.epgQueueService.epgResult$.subscribe( ({ streamId, items }) => { - const previewProgram = this.pickPreviewProgram(items); - if (previewProgram) { - this.applyProgram(streamId, previewProgram); - } + // The earliest-item fallback belongs to a row that has + // nothing to show yet. Once a row holds a programme, + // only one on air or upcoming may replace it: a refill + // answered with the same finished window would + // otherwise walk the row back to an older programme, + // the very thing the refresh exists to prevent (#767). + const program = this.pickProgramFor(streamId, items); + if (program) this.applyProgram(streamId, program); } ) ); + + // Nothing else re-evaluates the shown "current program" as + // wall-clock time passes (#767): applyProgram() only runs on + // scroll-into-view, a new EPG result, or an offset change. The + // tick is shared with every other mounted list so their requests + // merge instead of cancelling one another. + this.leaveEpgRefresh = this.epgRefreshCoordinator.register(() => + this.collectEpgRefresh() + ); } } @@ -355,13 +395,46 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { range.start, range.end ); - this.lastVisibleChannels = visibleChannels; this.loadEpgForVisibleChannels(visibleChannels); }) ); } } + /** + * The rows actually on screen now, read from the viewport rather than + * remembered. `renderedRangeStream` only emits when the index RANGE + * changes, so a new category, search term or sort order leaves a snapshot + * holding the previous list's rows; and a search with no results destroys + * the viewport altogether, where a snapshot would keep naming rows that + * are gone. Empty means nothing is on screen, which is the honest answer. + */ + private visibleChannelsNow(): XtreamChannelListItem[] { + const range = this.viewport()?.getRenderedRange(); + return range + ? this.filteredChannels().slice(range.start, range.end) + : []; + } + + /** Rows on screen, or a head of the list while nothing is rendered. */ + private channelsToPreview(): XtreamChannelListItem[] { + const visible = this.visibleChannelsNow(); + return visible.length ? visible : this.filteredChannels().slice(0, 50); + } + + /** + * The programme to show for `streamId`. A row with nothing yet may take + * the earliest known item so it is not blank; a row that already holds one + * accepts only an airing or upcoming programme, or a guide that has run + * out would walk it backwards (#767). + */ + private pickProgramFor(streamId: number, items: EpgItem[]): EpgItem | null { + const now = this.epgClockMs(); + return this.epgPrograms.has(streamId) + ? pickAiringOrUpcomingEpgItem(items, now) + : pickEpgPreviewItem(items, now); + } + /** Wall-clock now in the provider's EPG clock (`epg-display-offset.util.ts`). */ private epgClockMs(): number { return epgProviderClockMs( @@ -373,10 +446,7 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { private repickPreviewsForOffsetChange(): void { this.epgPrograms.clear(); this.currentProgramsProgress.clear(); - const visible = this.lastVisibleChannels.length - ? this.lastVisibleChannels - : this.filteredChannels().slice(0, 50); - this.loadEpgForVisibleChannels(visible); + this.loadEpgForVisibleChannels(this.channelsToPreview()); this.cdr.markForCheck(); } @@ -385,76 +455,164 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { return; } - const playlist = this.xtreamStore.currentPlaylist(); - if (!playlist) return; + if (!this.xtreamStore.currentPlaylist()) return; - const credentials: XtreamCredentials = { - serverUrl: playlist.serverUrl, - username: playlist.username, - password: playlist.password, - serverTimezone: playlist.serverTimezone, - }; - - const visibleIds = new Set(channels.map((ch) => ch.xtream_id)); - const uncachedEntries: { - streamId: number; - epgChannelId?: string | null; - playlistId?: string | null; - }[] = []; + const uncachedChannels: XtreamChannelListItem[] = []; // Apply cached results immediately for (const channel of channels) { const cached = this.epgQueueService.getCached(channel.xtream_id); if (cached !== null) { - const previewProgram = this.pickPreviewProgram(cached); - if (previewProgram) { - if (!this.epgPrograms.has(channel.xtream_id)) { - this.applyProgram(channel.xtream_id, previewProgram); - } + const previewProgram = pickEpgPreviewItem( + cached, + this.epgClockMs() + ); + if ( + previewProgram && + !this.epgPrograms.has(channel.xtream_id) + ) { + this.applyProgram(channel.xtream_id, previewProgram); } continue; } if (!this.epgPrograms.has(channel.xtream_id)) { - uncachedEntries.push({ - streamId: channel.xtream_id, - epgChannelId: channel.epg_channel_id ?? null, - playlistId: playlist.id ?? null, - }); + uncachedChannels.push(channel); } } - if (uncachedEntries.length > 0) { - this.epgQueueService - .enqueue(uncachedEntries, visibleIds, credentials) - .catch((error) => { - console.warn('EPG enqueue failed', error); - }); - } + this.requestEpgFor(uncachedChannels, channels); } - private updateProgramProgress(streamId: number, program: EpgItem) { + /** + * Advances the programme shown under the rows on screen as wall-clock + * time passes (#767) — nothing else re-evaluates it, so a row stayed + * pinned to a finished programme until the category was left and + * re-entered. + * + * A programme that is still on air only has its progress bar moved: no + * cache read, no request. Once it ends the row is re-picked from the + * queue's cache, and a cache holding nothing on air or upcoming is + * dropped so the next fetch can refill it. A finished programme is never + * re-applied — it would keep presenting itself as current, and the + * earliest-item fallback used for a first paint could even move the row + * backwards. What is on screen stays there until a replacement arrives, + * so a refreshing row never blanks out. + * + * Deciding all this costs nothing but local state, so every mounted list + * does it for itself; what it cannot do alone is ASK, which is why the + * rows to fetch are handed back to the coordinator to merge. + */ + private collectEpgRefresh(): EpgRefreshContribution | null { + const channels = this.visibleChannelsNow(); + if (!this.supportsEpg || channels.length === 0) return null; + + const playlist = this.xtreamStore.currentPlaylist(); + if (!playlist) return null; + const now = this.epgClockMs(); - const start = this.getProgramTimestampMs( - program.start, - program.start_timestamp - ); - const end = this.getProgramTimestampMs( - program.stop ?? program.end, - program.stop_timestamp - ); + const wallClockNow = Date.now(); + const staleChannels: XtreamChannelListItem[] = []; + let movedProgress = false; - if (now >= start && now <= end) { - const duration = end - start; - const elapsed = now - start; - const progress = (elapsed / duration) * 100; + for (const channel of channels) { + const streamId = channel.xtream_id; + const shown = this.epgPrograms.get(streamId); + if (shown && !hasEpgProgramEnded(shown, now)) { + this.updateProgramProgress(streamId, shown); + movedProgress = true; + continue; + } - this.currentProgramsProgress.set(streamId, progress); + const cached = this.epgQueueService.getCached(streamId); + if (cached === null) { + staleChannels.push(channel); + continue; + } + + const replacement = pickAiringOrUpcomingEpgItem(cached, now); + if (replacement) { + this.applyProgram(streamId, replacement); + continue; + } + + if (cached.length === 0) { + // The provider has nothing for this channel and said so; that + // empty answer is cached on purpose, and asking again before + // it expires would put one request per EPG-less visible row + // on the wire every minute. + continue; + } + + if (!this.epgRefill.claim(playlist.id, streamId, wallClockNow)) { + continue; + } + + // Every cached programme has ended. The entry stays valid for + // minutes and the queue skips a stream that still has one, so it + // has to go before the refill below can reach the provider. + this.epgQueueService.invalidate(streamId); + staleChannels.push(channel); + } + + this.epgRefill.forgetExpired(wallClockNow); + // A replaced row already rendered through applyProgram(); this is for + // the progress bars that advanced without one. + if (movedProgress) this.cdr.markForCheck(); + + if (staleChannels.length === 0) { + return null; + } + + return { + playlistId: playlist.id, + credentials: xtreamCredentialsOf(playlist), + // The whole slice, not just what is being fetched: the queue drops + // anything outside the visible set it was last handed. + visibleStreamIds: channels.map((channel) => channel.xtream_id), + staleEntries: staleChannels.map((channel) => + epgQueueEntryFor(channel, playlist.id) + ), + }; + } + + /** + * Queues EPG for `channels`. `visibleChannels` is the full viewport slice: + * the queue drops anything outside the visible set it was last handed, so + * passing only the subset being fetched would strand entries queued for + * the other rows on screen. + */ + private requestEpgFor( + channels: XtreamChannelListItem[], + visibleChannels: XtreamChannelListItem[] + ): void { + const playlist = this.xtreamStore.currentPlaylist(); + if (!playlist || channels.length === 0) return; + + this.epgQueueService + .enqueue( + channels.map((channel) => + epgQueueEntryFor(channel, playlist.id) + ), + new Set(visibleChannels.map((channel) => channel.xtream_id)), + xtreamCredentialsOf(playlist) + ) + .catch((error) => { + // An Xtream failure carries the stream URL, which is built out + // of the username and password. + this.logger.warn('EPG enqueue failed', error); + }); + } + + private updateProgramProgress(streamId: number, program: EpgProgram) { + const progress = epgProgramProgressPercent(program, this.epgClockMs()); + if (progress === null) { + this.currentProgramsProgress.delete(streamId); return; } - this.currentProgramsProgress.delete(streamId); + this.currentProgramsProgress.set(streamId, progress); } isSelected(item: XtreamCategory | XtreamCategoryLike): boolean { @@ -528,98 +686,22 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { ngOnDestroy(): void { this.subscriptions.unsubscribe(); + this.leaveEpgRefresh?.(); } private applyProgram(streamId: number, program: EpgItem): void { - this.epgPrograms.set(streamId, this.toSharedEpgProgram(program)); - this.updateProgramProgress(streamId, program); + const shownProgram = toSharedEpgProgram(program); + this.epgPrograms.set(streamId, shownProgram); + this.updateProgramProgress(streamId, shownProgram); + if (!hasEpgProgramEnded(shownProgram, this.epgClockMs())) { + // A programme that has not run out proves the guide is flowing, + // so the next gap on this row may be refilled straight away. + const { id } = this.xtreamStore.currentPlaylist() ?? {}; + this.epgRefill.release(id, streamId); + } this.cdr.detectChanges(); } - private pickPreviewProgram(items: EpgItem[]): EpgItem | null { - if (!items.length) { - return null; - } - - const now = this.epgClockMs(); - const normalizedItems = [...items].sort( - (a, b) => - this.getProgramTimestampMs(a.start, a.start_timestamp) - - this.getProgramTimestampMs(b.start, b.start_timestamp) - ); - - const currentProgram = normalizedItems.find((item) => { - const start = this.getProgramTimestampMs( - item.start, - item.start_timestamp - ); - const end = this.getProgramTimestampMs( - item.stop ?? item.end, - item.stop_timestamp - ); - return now >= start && now <= end; - }); - - if (currentProgram) { - return currentProgram; - } - - const nextProgram = normalizedItems.find((item) => { - return ( - this.getProgramTimestampMs(item.start, item.start_timestamp) > - now - ); - }); - - return nextProgram ?? normalizedItems[0]; - } - - private getProgramTimestampMs( - dateValue: string | undefined, - unixTimestampValue: string | undefined - ): number { - const unixTimestamp = Number(unixTimestampValue); - if (Number.isFinite(unixTimestamp) && unixTimestamp > 0) { - return unixTimestamp * 1000; - } - - return new Date(dateValue ?? '').getTime(); - } - - private toSharedEpgProgram(program: EpgItem): EpgProgram { - return { - start: program.start, - stop: program.stop ?? program.end, - channel: program.channel_id ?? program.id, - title: program.title, - desc: program.description ?? null, - category: null, - startTimestamp: this.getProgramTimestampSeconds( - program.start, - program.start_timestamp - ), - stopTimestamp: this.getProgramTimestampSeconds( - program.stop ?? program.end, - program.stop_timestamp - ), - }; - } - - private getProgramTimestampSeconds( - dateValue: string | undefined, - unixTimestampValue: string | undefined - ): number | null { - const unixTimestamp = Number(unixTimestampValue); - if (Number.isFinite(unixTimestamp) && unixTimestamp > 0) { - return unixTimestamp; - } - - const parsedDate = new Date(dateValue ?? '').getTime(); - return Number.isFinite(parsedDate) - ? Math.floor(parsedDate / 1000) - : null; - } - // ── Context menu ──────────────────────────────────────────── onChannelContextMenu( @@ -689,10 +771,7 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { this.epgQueueService.invalidate(streamId); this.epgPrograms.delete(streamId); this.currentProgramsProgress.delete(streamId); - const visible = this.lastVisibleChannels.length - ? this.lastVisibleChannels - : this.filteredChannels().slice(0, 50); - this.loadEpgForVisibleChannels(visible); + this.loadEpgForVisibleChannels(this.channelsToPreview()); }); }