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 <noreply@anthropic.com>
This commit is contained in:
4grayandClaude Fable 5 committed 2026-07-19 17:59:47 +02:00
1 parent b1fc23abbb
commit ef1cff25ed
7 files changed
+216 -48

No files matched your search

@@ -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<string, jest.Mock> = {};
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');
});
});
@@ -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;
@@ -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'
);
});
@@ -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();
}
@@ -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 {
@@ -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) {
@@ -137,6 +137,9 @@ export class PortalChannelsListComponent implements AfterViewInit, OnDestroy {
epgPrograms = new Map<number, EpgProgram>();
currentProgramsProgress = new Map<number, number>();
/** 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<void> {
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<string | null> {
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;
}
}
}