From bc6e7018e0d9e9662bd1080d70f3d9d11dc118ed Mon Sep 17 00:00:00 2001 From: 4gray <4gray@users.noreply.github.com> Date: Sun, 19 Jul 2026 20:21:40 +0200 Subject: [PATCH] fix(epg): harden three latent edges from the #1165 manual-mapping review (#1214) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(epg): harden three latent edges from the #1165 manual-mapping review Follow-up to #1165 (manual EPG-to-channel mapping). Three minor but real issues flagged by the bot reviews on #1173, all in already-merged #1165 code rather than the Stalker delta: 1. EPG program dedup ignored source_url. The unique index and upsert key (channel_id, start, title) collapsed programmes imported from different XMLTV sources that shared those columns, and the upsert reassigned source_url to the last importer — so source-scoped queries could miss a programme and source-scoped deletes could drop another source's row. The key and index now include source_url (migrated via a _v2 index that drops the old source-blind one); the upsert no longer overwrites source_url. 2. Xtream getMapping fallback capped candidate streams at an unordered first five, so a mapping saved under a later stream sharing the provider epg_channel_id was silently ignored. Replaced the two-step fetch-then-lookup with a single content⋈categories⋈mappings join that finds a mapping under any matching stream, with no arbitrary cap. 3. The Xtream mapping dialog did not refresh after closing, unlike the Stalker path, so a remapped visible/selected channel kept its stale preview until the 5-minute TTL or a rescroll. Added EpgQueueService.invalidate(streamId) and a before/after mapping compare in the channel-list dialog flow that invalidates the cache and refetches the current viewport when the mapping actually changed. Tests: source-aware dedup index/upsert assertions; join-based getMapping resolution regardless of stream position; EpgQueueService.invalidate. Co-Authored-By: Claude Fable 5 * fix(epg): address review feedback on the #1165 follow-up - portal-channels-list: forward the already-validated playlistId into the mapping dialog instead of re-reading currentPlaylist() after the async getEpgMapping roundtrip (which could return undefined on navigation) - EpgQueueService.invalidate(): bump a per-stream invalidation epoch and clear inFlight so a request already running when the mapping changes has its (pre-change) result discarded via an epoch check in fetchEpg, and the immediate re-enqueue can schedule a fresh mapping-aware fetch - getMapping Xtream fallback: order the join deterministically before limit(1) so the resolved mapping is stable; documented that this backend layer has no caller playlist context and is a best-effort net behind the renderer's playlist-scoped resolution Tests: in-flight staleness discard for invalidate(). Co-Authored-By: Claude Fable 5 * test(epg): split EpgQueueService invalidation specs under the max-lines limit The added invalidate() tests pushed epg-queue.service.spec.ts to 415 lines, over the 400-line ESLint cap (and the baseline must not grow). Moved them to a focused epg-queue-invalidation.spec.ts; both files are now under the limit. Co-Authored-By: Claude Fable 5 * fix(epg): resolve second-order review findings on the mapping follow-up Two P2 issues Codex raised on the previous fixes: - Unscoped program lookups could return the same programme twice now that the dedup index preserves per-source rows: the M3U timeline calls getChannelPrograms without sourceUrls, so two sources sharing channel/start/title both surfaced. Added toEpgProgams(), which collapses duplicate channel|start|title slots after mapping/validation, applied at every getChannelPrograms return. - EpgQueueService.fetchEpg unconditionally cleared the in-flight marker in its finally, which could drop a marker a re-enqueued request took over after invalidate(). It now only releases the marker when the completing request still owns it (epoch unchanged), preserving per-stream dedup and concurrency accounting. Tests: unscoped duplicate-slot collapse; stale request preserving a fresh in-flight marker. Co-Authored-By: Claude Fable 5 * fix(epg): deduplicate program rows in SQL before applying row limits Codex follow-up: the JS-level slot dedup ran after the SQL LIMIT, so duplicate cross-source rows consumed the cap and truncated real data. Moved the dedup into SQL with GROUP BY, applied before the limits: - selectChannelPrograms / selectLegacyChannelPrograms: GROUP BY (channel_id, start, title) before ORDER BY start LIMIT 500, so the timeline cap counts distinct programmes rather than duplicate rows - selectCurrentProgramsForChannelIds: GROUP BY channel_id before LIMIT channelIds.length, so duplicate cross-source current slots can't starve other channels of their current-programme preview The JS toEpgPrograms() dedup stays as a safety net (e.g. legacy NULL-source rows the unique index treats as distinct). Test query-chain mocks updated for the new groupBy link. Co-Authored-By: Claude Fable 5 --------- Co-authored-by: Claude Fable 5 --- .../src/app/events/epg-query.service.spec.ts | 122 ++++++++++++++- .../src/app/events/epg-query.service.ts | 109 +++++++++----- .../src/app/events/epg.events.spec.ts | 5 +- .../src/app/workers/epg-database.spec.ts | 23 ++- .../src/app/workers/epg-database.ts | 44 ++++-- .../services/epg-queue-invalidation.spec.ts | 141 ++++++++++++++++++ .../src/lib/services/epg-queue.service.ts | 45 +++++- .../portal-channels-list.component.ts | 66 +++++++- 8 files changed, 495 insertions(+), 60 deletions(-) create mode 100644 libs/portal/xtream/data-access/src/lib/services/epg-queue-invalidation.spec.ts diff --git a/apps/electron-backend/src/app/events/epg-query.service.spec.ts b/apps/electron-backend/src/app/events/epg-query.service.spec.ts index 33350e748..577096476 100644 --- a/apps/electron-backend/src/app/events/epg-query.service.spec.ts +++ b/apps/electron-backend/src/app/events/epg-query.service.spec.ts @@ -66,16 +66,35 @@ function createLimitedSelectChain( limitResult: unknown, whereCalls: unknown[] ): { from: jest.Mock } { - const limit = jest.fn().mockResolvedValue(limitResult); + const chain: Record = { + limit: jest.fn().mockResolvedValue(limitResult), + }; + // groupBy/orderBy are optional links in the various query shapes; each + // returns the chain so where().groupBy().orderBy().limit() (and any + // subset) resolves to the same data. + chain.groupBy = jest.fn(() => chain); + chain.orderBy = jest.fn(() => chain); const where = jest.fn((condition: unknown) => { whereCalls.push(condition); - return { limit }; + return chain; }); return { from: jest.fn(() => ({ where })), }; } +/** Program-query chain: from().where().groupBy().orderBy().limit(). */ +function createProgramChain(rows: unknown): { from: jest.Mock } { + const chain: Record = { + limit: jest.fn().mockResolvedValue(rows), + }; + chain.groupBy = jest.fn(() => chain); + chain.orderBy = jest.fn(() => chain); + return { + from: jest.fn(() => ({ where: jest.fn(() => chain) })), + }; +} + describe('EpgQueryService', () => { let service: EpgQueryService; @@ -403,4 +422,103 @@ describe('EpgQueryService', () => { ) ).toBe(true); }); + + it('resolves a manual mapping saved under any stream sharing the epg_channel_id via a single join', async () => { + const whereCalls: unknown[] = []; + const joinLimit = jest + .fn() + .mockResolvedValue([{ epgChannelId: 'BBC.MAPPED' }]); + let innerJoinCount = 0; + const joinChain: Record = {}; + joinChain.innerJoin = jest.fn(() => { + innerJoinCount += 1; + return joinChain; + }); + joinChain.where = jest.fn((condition: unknown) => { + whereCalls.push(condition); + return { orderBy: jest.fn(() => ({ limit: joinLimit })) }; + }); + + // Program lookup for the resolved (mapped) channel id. + const programRow = { + id: 1, + channelId: 'BBC.MAPPED', + start: '2026-06-21T17:00:00.000Z', + stop: '2026-06-21T18:00:00.000Z', + title: 'Mapped Programme', + description: null, + category: null, + iconUrl: null, + rating: null, + episodeNum: null, + sourceUrl: null, + }; + + const select = jest + .fn() + // getMapping direct lookup — miss. + .mockReturnValueOnce(createLimitedSelectChain([], whereCalls)) + // getMapping Xtream fallback — single join across content / + // categories / mappings, no arbitrary candidate cap. + .mockReturnValueOnce({ from: jest.fn(() => joinChain) }) + // selectChannelPrograms for the mapped id. + .mockReturnValueOnce(createProgramChain([programRow])); + + getDatabase.mockResolvedValue({ select }); + + const result = await service.getChannelPrograms('provider.epg.id'); + + expect(innerJoinCount).toBe(2); + expect(joinLimit).toHaveBeenCalledWith(1); + expect(result).toHaveLength(1); + expect(result[0].title).toBe('Mapped Programme'); + }); + + it('collapses duplicate programme slots from multiple sources in an unscoped lookup', async () => { + const whereCalls: unknown[] = []; + const dupRow = { + id: 1, + channelId: 'bbc.one.uk', + start: '2026-06-21T20:00:00.000Z', + stop: '2026-06-21T21:00:00.000Z', + title: 'News', + description: null, + category: null, + iconUrl: null, + rating: null, + episodeNum: null, + }; + const select = jest + .fn() + // getMapping direct + Xtream fallback — both miss. + .mockReturnValueOnce(createLimitedSelectChain([], whereCalls)) + .mockReturnValueOnce({ + from: jest.fn(() => ({ + innerJoin: jest.fn(function inner() { + return this; + }), + where: jest.fn(() => ({ + orderBy: jest.fn(() => ({ + limit: jest.fn().mockResolvedValue([]), + })), + })), + })), + }) + // selectChannelPrograms — two sources, same channel/start/title. + // The SQL GROUP BY would collapse these; the mock returns both to + // prove the JS safety-net dedup also holds. + .mockReturnValueOnce( + createProgramChain([ + { ...dupRow, sourceUrl: 'https://a.example/g.xml' }, + { ...dupRow, id: 2, sourceUrl: 'https://b.example/g.xml' }, + ]) + ); + + getDatabase.mockResolvedValue({ select }); + + const result = await service.getChannelPrograms('bbc.one.uk'); + + expect(result).toHaveLength(1); + expect(result[0].title).toBe('News'); + }); }); diff --git a/apps/electron-backend/src/app/events/epg-query.service.ts b/apps/electron-backend/src/app/events/epg-query.service.ts index 75c8a23c4..8ed16023e 100644 --- a/apps/electron-backend/src/app/events/epg-query.service.ts +++ b/apps/electron-backend/src/app/events/epg-query.service.ts @@ -1,6 +1,5 @@ import { and, eq, inArray, isNull, or, sql, type SQL } from 'drizzle-orm'; import { - buildXtreamEpgMappingKey, EpgChannelMetadata, EpgProgram, } from '@iptvnator/shared/interfaces'; @@ -69,9 +68,7 @@ export class EpgQueryService { } if (results.length > 0) { - return results - .map(this.transformDbRowToEpgProgram) - .filter(this.isValidEpgProgram); + return this.toEpgPrograms(results); } let channel = await this.selectChannelById( @@ -102,9 +99,7 @@ export class EpgQueryService { } if (results.length > 0) { - return results - .map(this.transformDbRowToEpgProgram) - .filter(this.isValidEpgProgram); + return this.toEpgPrograms(results); } } @@ -154,9 +149,7 @@ export class EpgQueryService { ); } - return results - .map(this.transformDbRowToEpgProgram) - .filter(this.isValidEpgProgram); + return this.toEpgPrograms(results); } return []; @@ -553,6 +546,13 @@ export class EpgQueryService { sourceUrls ) ) + // Collapse identical slots from multiple sources in SQL so the + // 500-row cap counts distinct programmes, not duplicate rows. + .groupBy( + schema.epgPrograms.channelId, + schema.epgPrograms.start, + schema.epgPrograms.title + ) .orderBy(schema.epgPrograms.start) .limit(500); } @@ -576,6 +576,11 @@ export class EpgQueryService { { legacyOnly: true } ) ) + .groupBy( + schema.epgPrograms.channelId, + schema.epgPrograms.start, + schema.epgPrograms.title + ) .orderBy(schema.epgPrograms.start) .limit(500); } @@ -607,6 +612,9 @@ export class EpgQueryService { options ) ) + // One current row per channel so duplicate cross-source slots + // don't consume the per-channel cap and starve other channels. + .groupBy(schema.epgPrograms.channelId) .limit(channelIds.length); } @@ -895,6 +903,31 @@ export class EpgQueryService { ); } + /** + * Map program rows to EpgProgram and drop invalid ones, then collapse + * duplicate programme slots. The dedup index is source-aware, so an + * unscoped lookup (e.g. the M3U timeline, which queries without + * `sourceUrls`) can return the same channel/start/title from two EPG + * sources; keep the first so a programme is never shown twice. + */ + private toEpgPrograms(rows: EpgProgramRow[]): EpgProgram[] { + const seen = new Set(); + const programs: EpgProgram[] = []; + for (const row of rows) { + const program = this.transformDbRowToEpgProgram(row); + if (!this.isValidEpgProgram(program)) { + continue; + } + const key = `${program.channel}|${program.start}|${program.title}`; + if (seen.has(key)) { + continue; + } + seen.add(key); + programs.push(program); + } + return programs; + } + private transformDbRowToEpgProgram(row: EpgProgramRow): EpgProgram { return { start: row.start, @@ -940,39 +973,45 @@ export class EpgQueryService { // Fallback: the key may be an Xtream provider epg_channel_id // while the mapping was saved under the playlist-scoped Xtream // key. Resolve the owning streams (joined through categories for - // the playlist ID) and check their mapping keys. The same - // epg_channel_id can appear in several playlists, so check a few. - const contentRows = await db + // the playlist ID) and look up their mapping keys in a single + // join, so a mapping saved under any matching stream is found — + // an earlier "check the first few" cap could silently miss a + // mapping saved under a later stream sharing this epg_channel_id. + // + // The join key is built in SQL to mirror + // buildXtreamEpgMappingKey(playlistId, xtreamId) → + // `xtream:{playlistId}:{xtreamId}`; keep the two in sync. + // + // This layer only receives the provider epg_channel_id, not the + // caller's playlist, so when the same id exists in several + // playlists it cannot scope to the caller's own mapping — the + // renderer resolves the playlist-scoped key directly and this is + // a best-effort fallback. Order deterministically so the chosen + // mapping is at least stable rather than storage-order dependent. + const mapped = await db .select({ - xtreamId: schema.content.xtreamId, - playlistId: schema.categories.playlistId, + epgChannelId: schema.epgChannelMappings.epgChannelId, }) .from(schema.content) .innerJoin( schema.categories, eq(schema.content.categoryId, schema.categories.id) ) - .where(eq(schema.content.epgChannelId, channelKey)) - .limit(5); - if (contentRows.length > 0) { - const candidateKeys = contentRows.map((row) => - buildXtreamEpgMappingKey(row.playlistId, row.xtreamId) - ); - const mapped = await db - .select({ - epgChannelId: schema.epgChannelMappings.epgChannelId, - }) - .from(schema.epgChannelMappings) - .where( - inArray( - schema.epgChannelMappings.channelKey, - candidateKeys - ) + .innerJoin( + schema.epgChannelMappings, + eq( + schema.epgChannelMappings.channelKey, + sql`'xtream:' || ${schema.categories.playlistId} || ':' || ${schema.content.xtreamId}` ) - .limit(1); - if (mapped.length > 0) { - return mapped[0].epgChannelId; - } + ) + .where(eq(schema.content.epgChannelId, channelKey)) + .orderBy( + schema.categories.playlistId, + schema.content.xtreamId + ) + .limit(1); + if (mapped.length > 0) { + return mapped[0].epgChannelId; } return null; diff --git a/apps/electron-backend/src/app/events/epg.events.spec.ts b/apps/electron-backend/src/app/events/epg.events.spec.ts index 6f6799a9a..1f2977a66 100644 --- a/apps/electron-backend/src/app/events/epg.events.spec.ts +++ b/apps/electron-backend/src/app/events/epg.events.spec.ts @@ -384,6 +384,7 @@ describe('EpgEvents', () => { const chain: Record = {} as Record; chain.where = jest.fn().mockReturnValue(chain); chain.innerJoin = jest.fn().mockReturnValue(chain); + chain.groupBy = jest.fn().mockReturnValue(chain); chain.orderBy = jest.fn().mockReturnValue(chain); chain.limit = jest.fn().mockResolvedValue(data); return chain; @@ -470,12 +471,14 @@ describe('EpgEvents', () => { const select = jest.fn(); const from = jest.fn(); const where = jest.fn(); + const groupBy = jest.fn(); const orderBy = jest.fn(); const limit = jest.fn(); select.mockImplementation(() => ({ from })); from.mockReturnValue({ where }); - where.mockReturnValue({ orderBy }); + where.mockReturnValue({ groupBy }); + groupBy.mockReturnValue({ orderBy }); orderBy.mockReturnValue({ limit }); limit.mockResolvedValue([ { diff --git a/apps/electron-backend/src/app/workers/epg-database.spec.ts b/apps/electron-backend/src/app/workers/epg-database.spec.ts index ed7c5e1a2..d8c5519c2 100644 --- a/apps/electron-backend/src/app/workers/epg-database.spec.ts +++ b/apps/electron-backend/src/app/workers/epg-database.spec.ts @@ -141,19 +141,36 @@ describe('EpgDatabase', () => { sql.startsWith('DELETE FROM epg_programs WHERE id NOT IN') ); const createIndexSql = preparedSql.find((sql) => - sql.startsWith('CREATE UNIQUE INDEX idx_epg_programs_dedup') + sql.startsWith('CREATE UNIQUE INDEX idx_epg_programs_dedup_v2') + ); + const dropOldIndexSql = preparedSql.find((sql) => + sql.startsWith('DROP INDEX IF EXISTS idx_epg_programs_dedup') ); expect(dedupDeleteSql).toBeDefined(); expect(createIndexSql).toBeDefined(); + expect(dropOldIndexSql).toBeDefined(); expect(statements.get(dedupDeleteSql!)).toHaveBeenCalled(); expect(statements.get(createIndexSql!)).toHaveBeenCalled(); + expect(statements.get(dropOldIndexSql!)).toHaveBeenCalled(); + + // The dedup key and index are source-aware so programmes imported + // from different EPG sources are not collapsed. + expect(dedupDeleteSql).toContain( + 'GROUP BY channel_id, start, title, source_url' + ); + expect(createIndexSql).toContain( + 'epg_programs(channel_id, start, title, source_url)' + ); const insertProgramSql = preparedSql.find((sql) => sql.startsWith('INSERT INTO epg_programs') ); expect(insertProgramSql).toContain( - 'ON CONFLICT(channel_id, start, title) DO UPDATE SET' + 'ON CONFLICT(channel_id, start, title, source_url) DO UPDATE SET' + ); + expect(insertProgramSql).not.toContain( + 'source_url = excluded.source_url' ); }); @@ -180,7 +197,7 @@ describe('EpgDatabase', () => { sql.startsWith('INSERT INTO epg_programs') ); expect(insertProgramSql).toContain( - 'ON CONFLICT(channel_id, start, title) DO UPDATE SET' + 'ON CONFLICT(channel_id, start, title, source_url) DO UPDATE SET' ); }); diff --git a/apps/electron-backend/src/app/workers/epg-database.ts b/apps/electron-backend/src/app/workers/epg-database.ts index 669faaddb..3865b4bad 100644 --- a/apps/electron-backend/src/app/workers/epg-database.ts +++ b/apps/electron-backend/src/app/workers/epg-database.ts @@ -31,9 +31,13 @@ export class EpgDatabase { `); // Guard against duplicate entries when the clearFirst logic misses old - // rows (e.g. because the source URL changed between imports). The same - // channel + start + title is treated as the same programme — a later - // import with a corrected stop time simply updates the earlier row. + // rows. The same channel + start + title + source is treated as the + // same programme — a later import with a corrected stop time simply + // updates the earlier row. `source_url` is part of the key so that + // programmes imported from different EPG sources that happen to share + // channel_id/start/title stay isolated: the query layer scopes by + // source_url, and source-scoped deletes must not clobber another + // source's rows. // The upsert (instead of INSERT OR REPLACE) keeps the epg_programs_fts // triggers consistent: REPLACE deletes rows without firing the delete // trigger unless recursive_triggers is enabled, leaving ghost FTS rows. @@ -43,14 +47,13 @@ export class EpgDatabase { dedupIndexReady ? `INSERT INTO epg_programs (channel_id, start, stop, title, description, category, icon_url, rating, episode_num, source_url) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - ON CONFLICT(channel_id, start, title) DO UPDATE SET + ON CONFLICT(channel_id, start, title, source_url) DO UPDATE SET stop = excluded.stop, description = excluded.description, category = excluded.category, icon_url = excluded.icon_url, rating = excluded.rating, - episode_num = excluded.episode_num, - source_url = excluded.source_url` + episode_num = excluded.episode_num` : `INSERT INTO epg_programs (channel_id, start, stop, title, description, category, icon_url, rating, episode_num, source_url) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)` ); @@ -156,11 +159,19 @@ export class EpgDatabase { } /** - * Create the unique (channel_id, start, title) dedup index. Databases - * that predate the index may already contain duplicate rows — exactly - * the situation the index is meant to prevent — so those are removed - * first; a plain `CREATE UNIQUE INDEX` would otherwise throw in this - * constructor and permanently break EPG imports for upgrading users. + * Create the unique (channel_id, start, title, source_url) dedup index. + * + * Databases that predate the index may already contain duplicate rows — + * exactly the situation the index is meant to prevent — so those are + * removed first; a plain `CREATE UNIQUE INDEX` would otherwise throw in + * this constructor and permanently break EPG imports for upgrading users. + * + * The `_v2` name migrates users off the earlier source-blind index + * (`idx_epg_programs_dedup` on `channel_id, start, title`), which + * collapsed programmes imported from different EPG sources that shared + * those three columns. The old index is dropped so it can no longer + * enforce the source-blind uniqueness. + * * Returns false when the index could not be created, in which case the * insert statement falls back to the previous plain-INSERT behaviour. */ @@ -169,23 +180,26 @@ export class EpgDatabase { const exists = this.db .prepare( `SELECT 1 FROM sqlite_master - WHERE type = 'index' AND name = 'idx_epg_programs_dedup'` + WHERE type = 'index' AND name = 'idx_epg_programs_dedup_v2'` ) .get(); if (!exists) { + this.db + .prepare(`DROP INDEX IF EXISTS idx_epg_programs_dedup`) + .run(); this.db .prepare( `DELETE FROM epg_programs WHERE id NOT IN ( SELECT MIN(id) FROM epg_programs - GROUP BY channel_id, start, title + GROUP BY channel_id, start, title, source_url )` ) .run(); this.db .prepare( - `CREATE UNIQUE INDEX idx_epg_programs_dedup - ON epg_programs(channel_id, start, title)` + `CREATE UNIQUE INDEX idx_epg_programs_dedup_v2 + ON epg_programs(channel_id, start, title, source_url)` ) .run(); } diff --git a/libs/portal/xtream/data-access/src/lib/services/epg-queue-invalidation.spec.ts b/libs/portal/xtream/data-access/src/lib/services/epg-queue-invalidation.spec.ts new file mode 100644 index 000000000..1318b7729 --- /dev/null +++ b/libs/portal/xtream/data-access/src/lib/services/epg-queue-invalidation.spec.ts @@ -0,0 +1,141 @@ +import { TestBed } from '@angular/core/testing'; +import { SettingsStore } from '@iptvnator/services'; +import { EpgQueueService } from './epg-queue.service'; +import { XtreamApiService } from './xtream-api.service'; +import { XtreamXmltvFallbackService } from './xtream-xmltv-fallback.service'; +import type { EpgItem } from '@iptvnator/shared/interfaces'; + +/** + * Covers EpgQueueService.invalidate(), used when a manual EPG mapping for a + * stream changes. Split from epg-queue.service.spec.ts to keep both files + * under the max-lines limit. + */ +describe('EpgQueueService invalidation', () => { + let service: EpgQueueService; + let xtreamApi: { getShortEpg: jest.Mock }; + + const credentials = { + serverUrl: 'https://xtream.example.com', + username: 'user', + password: 'pass', + }; + + type ServicePrivates = { + fetchEpg: ( + credentials: typeof credentials, + streamId: number + ) => Promise; + shouldFetch: (streamId: number) => boolean; + xmltvPreviewByStreamId: Map; + epgChannelByStreamId: Map; + inFlight: Set; + }; + + beforeEach(() => { + xtreamApi = { getShortEpg: jest.fn().mockResolvedValue([]) }; + + TestBed.configureTestingModule({ + providers: [ + EpgQueueService, + { provide: XtreamApiService, useValue: xtreamApi }, + { + provide: XtreamXmltvFallbackService, + useValue: { + getProgramsForChannel: jest.fn().mockResolvedValue([]), + getCurrentProgramsBatch: jest.fn().mockResolvedValue({}), + }, + }, + { + provide: SettingsStore, + useValue: { + preferUploadedEpgOverXtream: jest.fn(() => false), + }, + }, + ], + }); + + service = TestBed.inject(EpgQueueService); + }); + + function priv(): ServicePrivates { + return service as unknown as ServicePrivates; + } + + function makeItem(channelId: string, title: string): EpgItem { + return { + id: `${channelId}|x`, + epg_id: '', + title, + lang: '', + start: '2026-05-07T08:00:00Z', + end: '2026-05-07T09:00:00Z', + stop: '2026-05-07T09:00:00Z', + description: '', + channel_id: channelId, + start_timestamp: '0', + stop_timestamp: '0', + }; + } + + it('drops every cached artifact so the stream refetches', async () => { + xtreamApi.getShortEpg.mockResolvedValue([makeItem('rtl.de', 'Now')]); + await priv().fetchEpg(credentials, 555); + priv().epgChannelByStreamId.set(555, 'rtl.de'); + priv().xmltvPreviewByStreamId.set(555, makeItem('rtl.de', 'now')); + + expect(service.getCached(555)).not.toBeNull(); + expect(priv().shouldFetch(555)).toBe(false); + + service.invalidate(555); + + expect(service.getCached(555)).toBeNull(); + expect(priv().shouldFetch(555)).toBe(true); + expect(priv().epgChannelByStreamId.has(555)).toBe(false); + expect(priv().xmltvPreviewByStreamId.has(555)).toBe(false); + }); + + it('discards an in-flight result when invalidated mid-request', async () => { + let resolveEpg!: (items: EpgItem[]) => void; + xtreamApi.getShortEpg.mockReturnValue( + new Promise((resolve) => { + resolveEpg = resolve; + }) + ); + const emitted: number[] = []; + const sub = service.epgResult$.subscribe(({ streamId }) => + emitted.push(streamId) + ); + + const fetchPromise = priv().fetchEpg(credentials, 707); + // The user changes the mapping while the request is still running. + service.invalidate(707); + // The provider now returns the pre-change (stale) result. + resolveEpg([makeItem('rtl.de', 'Stale')]); + await fetchPromise; + + expect(service.getCached(707)).toBeNull(); + expect(emitted).not.toContain(707); + sub.unsubscribe(); + }); + + it('leaves a fresh in-flight marker intact when a stale request completes', async () => { + let resolveA!: (items: EpgItem[]) => void; + xtreamApi.getShortEpg.mockReturnValue( + new Promise((resolve) => { + resolveA = resolve; + }) + ); + + const aPromise = priv().fetchEpg(credentials, 909); + // Mapping changes mid-flight; a re-enqueue (request B) then starts and + // takes ownership of the in-flight marker. + service.invalidate(909); + priv().inFlight.add(909); + + resolveA([makeItem('rtl.de', 'Stale')]); + await aPromise; + + // Request A was stale, so its finally must not clear B's marker. + expect(priv().inFlight.has(909)).toBe(true); + }); +}); diff --git a/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts index b1f173749..3e3b71e92 100644 --- a/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts +++ b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.ts @@ -55,6 +55,8 @@ export class EpgQueueService implements OnDestroy { private readonly cache = new Map(); private queue: number[] = []; private readonly inFlight = new Set(); + /** Bumped by invalidate() so a stale in-flight result is discarded. */ + private readonly invalidationEpoch = new Map(); private readonly epgChannelByStreamId = new Map(); private readonly xmltvPreviewByStreamId = new Map(); private visibleSet = new Set(); @@ -80,6 +82,29 @@ export class EpgQueueService implements OnDestroy { return entry.data; } + /** + * Drop every cached artifact for a stream so the next enqueue refetches + * it. Used when a manual EPG mapping for the stream changes, since the + * cached preview/resolution was computed for the previous mapping. + * + * Also bumps an invalidation epoch and clears `inFlight`: a request that + * was already running when the mapping changed carries the pre-change + * resolution, so its result is discarded (epoch mismatch in `fetchEpg`) + * and clearing `inFlight` lets the immediate re-enqueue schedule a fresh + * fetch through the mapping-aware `enqueue()` path. + */ + invalidate(streamId: number): void { + this.cache.delete(streamId); + this.failureTimestamps.delete(streamId); + this.epgChannelByStreamId.delete(streamId); + this.xmltvPreviewByStreamId.delete(streamId); + this.inFlight.delete(streamId); + this.invalidationEpoch.set( + streamId, + (this.invalidationEpoch.get(streamId) ?? 0) + 1 + ); + } + private isFailureCoolingDown(streamId: number): boolean { const timestamp = this.failureTimestamps.get(streamId); if (timestamp == null) { @@ -305,6 +330,9 @@ export class EpgQueueService implements OnDestroy { credentials: XtreamCredentials, streamId: number ): Promise { + const startEpoch = this.invalidationEpoch.get(streamId) ?? 0; + const isStale = (): boolean => + (this.invalidationEpoch.get(streamId) ?? 0) !== startEpoch; try { const apiItems = await this.apiService.getShortEpg( credentials, @@ -313,6 +341,11 @@ export class EpgQueueService implements OnDestroy { { suppressErrorLog: true } ); + // A mapping change during the request invalidated this result. + if (isStale()) { + return; + } + if (apiItems.length > 0) { this.recordSuccess(streamId, apiItems); return; @@ -321,13 +354,23 @@ export class EpgQueueService implements OnDestroy { const xmltv = this.xmltvPreviewByStreamId.get(streamId); this.recordSuccess(streamId, xmltv ? [xmltv] : []); } catch (error) { + if (isStale()) { + return; + } this.failureTimestamps.set(streamId, Date.now()); this.logger.error( `Failed to load EPG for stream ${streamId}`, error ); } finally { - this.inFlight.delete(streamId); + // Only release the in-flight marker if this request still owns it. + // When invalidate() cleared it mid-flight, a later re-enqueue may + // already have started a new request for the same stream; an + // unconditional delete here would drop that request's marker and + // let a third concurrent fetch start. + if (!isStale()) { + this.inFlight.delete(streamId); + } } } 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 56c71104b..8a2076a32 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 @@ -137,6 +137,9 @@ 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[] = []; + readonly viewport = viewChild(CdkVirtualScrollViewport); private subscriptions = new Subscription(); @@ -237,6 +240,7 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { range.start, range.end ); + this.lastVisibleChannels = visibleChannels; this.loadEpgForVisibleChannels(visibleChannels); }) ); @@ -509,10 +513,66 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { return; } + const channelKey = buildXtreamEpgMappingKey(playlistId, xtreamId); + void this.openEpgMappingDialog( + channelKey, + channel, + xtreamId, + playlistId + ); + } + + private async openEpgMappingDialog( + channelKey: string, + channel: XtreamChannelListItem, + streamId: number, + playlistId: string + ): Promise { + const mappingBefore = await this.readEpgMapping(channelKey); + EpgMappingDialogComponent.open(this.dialog, { - channelKey: buildXtreamEpgMappingKey(playlistId, xtreamId), - channelName: channel.title ?? channel.name ?? String(xtreamId), + channelKey, + channelName: channel.title ?? channel.name ?? String(streamId), playlistId, - }); + }) + .afterClosed() + .subscribe(async () => { + const mappingAfter = await this.readEpgMapping(channelKey); + if (mappingAfter === mappingBefore) { + return; + } + // The mapping changed (saved or removed) — drop the cached + // preview/resolution and refetch so the row updates now + // instead of after the 5-minute TTL or the next scroll. + 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); + }); + } + + /** Read the current mapped EPG channel id, or null (PWA / no mapping). */ + private async readEpgMapping(channelKey: string): Promise { + if (!this.supportsEpgMapping) { + return null; + } + const bridge = ( + window as unknown as { + electron?: { + getEpgMapping?: ( + key: string + ) => Promise<{ epgChannelId?: string } | null>; + }; + } + ).electron; + try { + const mapping = await bridge?.getEpgMapping?.(channelKey); + return mapping?.epgChannelId?.trim() || null; + } catch { + return null; + } } }