From ef1cff25ed1ae8138b597b9ddb408f5ef6bb9d5b Mon Sep 17 00:00:00 2001 From: 4gray Date: Sun, 19 Jul 2026 17:59:47 +0200 Subject: [PATCH] fix(epg): harden three latent edges from the #1165 manual-mapping review MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../src/app/events/epg-query.service.spec.ts | 60 ++++++++++++++++++ .../src/app/events/epg-query.service.ts | 46 ++++++-------- .../src/app/workers/epg-database.spec.ts | 23 ++++++- .../src/app/workers/epg-database.ts | 44 ++++++++----- .../lib/services/epg-queue.service.spec.ts | 17 +++++ .../src/lib/services/epg-queue.service.ts | 12 ++++ .../portal-channels-list.component.ts | 62 +++++++++++++++++-- 7 files changed, 216 insertions(+), 48 deletions(-) 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..e4603bad0 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 @@ -403,4 +403,64 @@ 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 { 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 programChain = { + from: jest.fn(() => ({ + where: jest.fn(() => ({ + orderBy: jest.fn(() => ({ + limit: jest.fn().mockResolvedValue([programRow]), + })), + })), + })), + }; + + 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(programChain); + + 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'); + }); }); 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..68739d6a4 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'; @@ -940,39 +939,34 @@ 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. + 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)) + .limit(1); + if (mapped.length > 0) { + return mapped[0].epgChannelId; } return null; 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.service.spec.ts b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.spec.ts index 3aa315f41..3001cc0e0 100644 --- a/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.spec.ts +++ b/libs/portal/xtream/data-access/src/lib/services/epg-queue.service.spec.ts @@ -355,6 +355,23 @@ describe('EpgQueueService', () => { expect(events).toEqual([808]); sub.unsubscribe(); }); + + it('invalidate() drops every cached artifact so the stream refetches', async () => { + xtreamApi.getShortEpg.mockResolvedValue([makeEpgItem('rtl.de', 'Now')]); + await priv().fetchEpg(credentials, 555); + priv().epgChannelByStreamId.set(555, 'rtl.de'); + seedXmltvPreview(555, 'rtl.de'); + + 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); + }); }); function makeEpgItem(channelId: string, title: string): EpgItem { 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..2cc41a476 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 @@ -80,6 +80,18 @@ 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. + */ + invalidate(streamId: number): void { + this.cache.delete(streamId); + this.failureTimestamps.delete(streamId); + this.epgChannelByStreamId.delete(streamId); + this.xmltvPreviewByStreamId.delete(streamId); + } + private isFailureCoolingDown(streamId: number): boolean { const timestamp = this.failureTimestamps.get(streamId); if (timestamp == null) { 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..3443689a8 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,60 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy { return; } + const channelKey = buildXtreamEpgMappingKey(playlistId, xtreamId); + void this.openEpgMappingDialog(channelKey, channel, xtreamId); + } + + private async openEpgMappingDialog( + channelKey: string, + channel: XtreamChannelListItem, + streamId: number + ): Promise { + const mappingBefore = await this.readEpgMapping(channelKey); + EpgMappingDialogComponent.open(this.dialog, { - channelKey: buildXtreamEpgMappingKey(playlistId, xtreamId), - channelName: channel.title ?? channel.name ?? String(xtreamId), - playlistId, - }); + channelKey, + channelName: channel.title ?? channel.name ?? String(streamId), + playlistId: this.xtreamStore.currentPlaylist()?.id, + }) + .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; + } } }