mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-09 01:16:15 -08:00
fix(epg): preserve source metadata through cache cleanup
This commit is contained in:
1 parent
6d406c85e2
commit
dc0b7581cb
16 files changed
+514
-59
No files matched your search
@@ -181,7 +181,10 @@ all enabled M3U playlist sources. Startup runs the same reconciliation after
|
||||
settings load and playlist migration; failed settings reads and incomplete
|
||||
playlist migration never authorize pruning. Removed sources retire queued and
|
||||
running imports before worker-owned deletion. Shared channel IDs survive while
|
||||
another source has programmes; manual mappings remain user preferences, but no
|
||||
another source has programmes or per-source channel metadata. The additive
|
||||
`epg_channel_sources` table preserves each imported source's name, logo, URL and
|
||||
timestamp so removal restores a surviving snapshot; ambiguous legacy metadata
|
||||
falls back to the XMLTV ID until reimport. Manual mappings remain user preferences, but no
|
||||
longer resolve deleted data. Renderer lookup generations, Xtream previews and Stalker mapping-cache
|
||||
invalidation prevent late results from restoring removed programmes. Provider
|
||||
EPG is independent. See `docs/architecture/m3u-playlist-module.md`
|
||||
|
||||
@@ -1666,7 +1666,10 @@ all enabled M3U playlist sources. Startup runs the same reconciliation after
|
||||
settings load and playlist migration; failed settings reads and incomplete
|
||||
playlist migration never authorize pruning. Removed sources retire queued and
|
||||
running imports before worker-owned deletion. Shared channel IDs survive while
|
||||
another source has programmes; manual mappings remain user preferences, but no
|
||||
another source has programmes or per-source channel metadata. The additive
|
||||
`epg_channel_sources` table preserves each imported source's name, logo, URL and
|
||||
timestamp so removal restores a surviving snapshot; ambiguous legacy metadata
|
||||
falls back to the XMLTV ID until reimport. Manual mappings remain user preferences, but no
|
||||
longer resolve deleted data. Renderer lookup generations, Xtream previews and Stalker mapping-cache
|
||||
invalidation prevent late results from restoring removed programmes. Provider
|
||||
EPG is independent. See `docs/architecture/m3u-playlist-module.md`
|
||||
|
||||
@@ -33,7 +33,8 @@ const epgFixtureXml = `<?xml version="1.0" encoding="UTF-8"?>
|
||||
function createCurrentXmltvFixture(
|
||||
channelId: string,
|
||||
channelName: string,
|
||||
programTitle: string
|
||||
programTitle: string,
|
||||
iconUrl?: string
|
||||
): string {
|
||||
const start = new Date(Date.now() - 15 * 60 * 1000);
|
||||
const stop = new Date(Date.now() + 45 * 60 * 1000);
|
||||
@@ -42,6 +43,7 @@ function createCurrentXmltvFixture(
|
||||
<tv>
|
||||
<channel id="${channelId}">
|
||||
<display-name>${channelName}</display-name>
|
||||
${iconUrl ? `<icon src="${iconUrl}"/>` : ''}
|
||||
</channel>
|
||||
<programme start="${formatXmltvDate(start)} +0000" stop="${formatXmltvDate(stop)} +0000" channel="${channelId}">
|
||||
<title>${programTitle}</title>
|
||||
@@ -220,6 +222,82 @@ test.describe('Electron EPG', () => {
|
||||
}
|
||||
});
|
||||
|
||||
test('@epg @electron restores the first source name and logo after removing the last metadata writer', async ({
|
||||
dataDir,
|
||||
}) => {
|
||||
const first = await createMutableTextServer(
|
||||
createCurrentXmltvFixture(
|
||||
'shared-logo',
|
||||
'First News',
|
||||
'First Bulletin',
|
||||
'https://example.com/first.png'
|
||||
),
|
||||
{ contentType: 'application/xml', resourcePath: '/first.xml' }
|
||||
);
|
||||
const second = await createMutableTextServer(
|
||||
createCurrentXmltvFixture(
|
||||
'shared-logo',
|
||||
'Second News',
|
||||
'Second Bulletin',
|
||||
'https://example.com/second.png'
|
||||
),
|
||||
{ contentType: 'application/xml', resourcePath: '/second.xml' }
|
||||
);
|
||||
let app = await launchElectronApp(dataDir);
|
||||
const metadata = () =>
|
||||
app.mainWindow.evaluate(
|
||||
async () =>
|
||||
(
|
||||
await window.electron.getEpgChannelMetadata([
|
||||
'shared-logo',
|
||||
])
|
||||
)['shared-logo']
|
||||
);
|
||||
try {
|
||||
await openSettings(app.mainWindow);
|
||||
await openSettingsSection(app.mainWindow, 'epg');
|
||||
for (const [index, source] of [first, second].entries()) {
|
||||
await app.mainWindow
|
||||
.getByRole('button', { name: 'Add EPG source' })
|
||||
.click();
|
||||
const inputs = app.mainWindow.locator('.epg-source-row input');
|
||||
await expect(inputs).toHaveCount(index + 1);
|
||||
await inputs.nth(index).fill(source.resourceUrl);
|
||||
await saveSettings(app.mainWindow);
|
||||
await expect.poll(metadata, { timeout: 30000 }).toMatchObject({
|
||||
displayName: index === 0 ? 'First News' : 'Second News',
|
||||
iconUrl: `https://example.com/${index === 0 ? 'first' : 'second'}.png`,
|
||||
});
|
||||
}
|
||||
await app.mainWindow
|
||||
.locator('.epg-source-row')
|
||||
.nth(1)
|
||||
.locator('button')
|
||||
.nth(1)
|
||||
.click();
|
||||
await saveSettings(app.mainWindow);
|
||||
const firstMetadata = {
|
||||
displayName: 'First News',
|
||||
iconUrl: 'https://example.com/first.png',
|
||||
};
|
||||
await expect.poll(metadata).toMatchObject(firstMetadata);
|
||||
await closeElectronApp(app);
|
||||
app = await launchElectronApp(dataDir);
|
||||
await expect.poll(metadata).toMatchObject(firstMetadata);
|
||||
expect(
|
||||
await app.mainWindow.evaluate(async () =>
|
||||
(
|
||||
await window.electron.getChannelPrograms('shared-logo')
|
||||
).map((p) => p.title)
|
||||
)
|
||||
).toEqual(['First Bulletin']);
|
||||
} finally {
|
||||
await closeElectronApp(app);
|
||||
await first.close();
|
||||
await second.close();
|
||||
}
|
||||
});
|
||||
|
||||
test('@epg @electron imports and renders an EPG source declared by an M3U playlist header', async ({
|
||||
dataDir,
|
||||
}) => {
|
||||
|
||||
@@ -4,40 +4,41 @@
|
||||
*/
|
||||
|
||||
export {
|
||||
// Tables
|
||||
playlists,
|
||||
categories,
|
||||
content,
|
||||
recentlyViewed,
|
||||
favorites,
|
||||
epgChannels,
|
||||
epgChannelMappings,
|
||||
epgPrograms,
|
||||
playbackPositions,
|
||||
downloads,
|
||||
recordings,
|
||||
appState,
|
||||
// Types
|
||||
type Playlist,
|
||||
type NewPlaylist,
|
||||
type AppState,
|
||||
type NewAppState,
|
||||
type Category,
|
||||
type NewCategory,
|
||||
type Content,
|
||||
type NewContent,
|
||||
type RecentlyViewed,
|
||||
type NewRecentlyViewed,
|
||||
type Favorite,
|
||||
type NewFavorite,
|
||||
type EpgChannel,
|
||||
type NewEpgChannel,
|
||||
type EpgProgramDb,
|
||||
type NewEpgProgramDb,
|
||||
type PlaybackPosition,
|
||||
type NewPlaybackPosition,
|
||||
type Download,
|
||||
type NewDownload,
|
||||
type Recording,
|
||||
type NewRecording,
|
||||
// Tables
|
||||
playlists,
|
||||
categories,
|
||||
content,
|
||||
recentlyViewed,
|
||||
favorites,
|
||||
epgChannels,
|
||||
epgChannelSources,
|
||||
epgChannelMappings,
|
||||
epgPrograms,
|
||||
playbackPositions,
|
||||
downloads,
|
||||
recordings,
|
||||
appState,
|
||||
// Types
|
||||
type Playlist,
|
||||
type NewPlaylist,
|
||||
type AppState,
|
||||
type NewAppState,
|
||||
type Category,
|
||||
type NewCategory,
|
||||
type Content,
|
||||
type NewContent,
|
||||
type RecentlyViewed,
|
||||
type NewRecentlyViewed,
|
||||
type Favorite,
|
||||
type NewFavorite,
|
||||
type EpgChannel,
|
||||
type NewEpgChannel,
|
||||
type EpgProgramDb,
|
||||
type NewEpgProgramDb,
|
||||
type PlaybackPosition,
|
||||
type NewPlaybackPosition,
|
||||
type Download,
|
||||
type NewDownload,
|
||||
type Recording,
|
||||
type NewRecording,
|
||||
} from '@iptvnator/shared/database';
|
||||
@@ -50,9 +50,9 @@ export async function checkEpgFreshness(
|
||||
if (!url?.trim()) continue;
|
||||
|
||||
const result = await db
|
||||
.select({ updatedAt: schema.epgChannels.updatedAt })
|
||||
.from(schema.epgChannels)
|
||||
.where(eq(schema.epgChannels.sourceUrl, url))
|
||||
.select({ updatedAt: schema.epgChannelSources.updatedAt })
|
||||
.from(schema.epgChannelSources)
|
||||
.where(eq(schema.epgChannelSources.sourceUrl, url))
|
||||
.limit(1);
|
||||
|
||||
const isFresh =
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import {
|
||||
appState,
|
||||
epgChannels,
|
||||
epgChannelSources,
|
||||
epgPrograms,
|
||||
playlists,
|
||||
} from '../database/schema';
|
||||
@@ -65,6 +66,22 @@ describe('committed EPG source reconciliation', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('finds metadata-only source owners after restart', async () => {
|
||||
rows.set(epgChannels, [{ url: 'active' }]);
|
||||
rows.set(epgPrograms, []);
|
||||
rows.set(epgChannelSources, [
|
||||
{ url: 'metadata-only' },
|
||||
{ url: 'active' },
|
||||
]);
|
||||
await reconcileEpgSources(['active']);
|
||||
expect(epgWorkerService.clearEpgDataForSource).toHaveBeenCalledWith(
|
||||
'metadata-only'
|
||||
);
|
||||
expect(epgWorkerService.clearEpgDataForSource).not.toHaveBeenCalledWith(
|
||||
'active'
|
||||
);
|
||||
});
|
||||
|
||||
it('does not clear historical request keys again after successful cleanup', async () => {
|
||||
requestEpgSource('removed');
|
||||
await reconcileEpgSources([]);
|
||||
|
||||
@@ -3,6 +3,7 @@ import { getDatabase } from '../database/connection';
|
||||
import {
|
||||
appState,
|
||||
epgChannels,
|
||||
epgChannelSources,
|
||||
epgPrograms,
|
||||
playlists,
|
||||
} from '../database/schema';
|
||||
@@ -57,12 +58,16 @@ export function reconcileEpgSources(globalUrls: string[]): Promise<void> {
|
||||
const programs = await db
|
||||
.selectDistinct({ url: epgPrograms.sourceUrl })
|
||||
.from(epgPrograms);
|
||||
const metadata = await db
|
||||
.selectDistinct({ url: epgChannelSources.sourceUrl })
|
||||
.from(epgChannelSources);
|
||||
const requested = requestedEpgSources();
|
||||
const removed = [
|
||||
...new Set(
|
||||
[
|
||||
...channels.map((row) => row.url),
|
||||
...programs.map((row) => row.url),
|
||||
...metadata.map((row) => row.url),
|
||||
...requested.keys(),
|
||||
].filter(
|
||||
(url): url is string =>
|
||||
|
||||
@@ -445,11 +445,15 @@ export class EpgWorkerService {
|
||||
}
|
||||
|
||||
private async startSourceClear(normalizedSourceUrl: string): Promise<void> {
|
||||
const inFlight = this.inFlightFetches.get(normalizedSourceUrl);
|
||||
const runningWorker = this.workers.get(normalizedSourceUrl);
|
||||
if (runningWorker) {
|
||||
this.workers.delete(normalizedSourceUrl);
|
||||
await this.terminateWorker(runningWorker, 'source clear');
|
||||
}
|
||||
// Error/timeout handlers may already have removed the worker from the
|
||||
// lookup, but its fetch promise still owns asynchronous termination.
|
||||
if (inFlight) await inFlight.catch(() => undefined);
|
||||
|
||||
return this.runClearWorker({
|
||||
timeoutLabel: 'EPG source clear',
|
||||
|
||||
@@ -78,6 +78,37 @@ describe('EpgEvents', () => {
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
}
|
||||
|
||||
it.each(['missing', 'stale', 'fresh'])(
|
||||
'uses source-owned freshness for %s metadata',
|
||||
async (state) => {
|
||||
const { epgChannelSources } = await import('../database/schema');
|
||||
const { checkEpgFreshness } = await import('./epg-fetch.service');
|
||||
const fresh = [{ updatedAt: new Date().toISOString() }];
|
||||
const snapshots =
|
||||
state === 'missing'
|
||||
? []
|
||||
: state === 'fresh'
|
||||
? fresh
|
||||
: [{ updatedAt: '2000-01-01T00:00:00Z' }];
|
||||
getDatabase.mockResolvedValue({
|
||||
select: () => ({
|
||||
from: (table: unknown) => ({
|
||||
where: () => ({
|
||||
limit: async () =>
|
||||
table === epgChannelSources ? snapshots : fresh,
|
||||
}),
|
||||
}),
|
||||
}),
|
||||
});
|
||||
const result = await checkEpgFreshness(['shared-source'], 12);
|
||||
expect(result).toEqual(
|
||||
state === 'fresh'
|
||||
? { freshUrls: ['shared-source'], staleUrls: [] }
|
||||
: { freshUrls: [], staleUrls: ['shared-source'] }
|
||||
);
|
||||
}
|
||||
);
|
||||
|
||||
it('uses the shared worker bootstrap and passes native module search paths to the EPG worker', async () => {
|
||||
const fetchPromise = (EpgEvents as unknown as Record<string, any>)[
|
||||
'fetchEpgFromUrl'
|
||||
@@ -366,6 +397,41 @@ describe('EpgEvents', () => {
|
||||
fetch.mockRestore();
|
||||
});
|
||||
|
||||
it.each([true, false])(
|
||||
'awaits a failing worker termination before source cleanup (already retired: %s)',
|
||||
async (alreadyRetired) => {
|
||||
const service = new EpgWorkerService('[Test EPG]', 1000);
|
||||
const url = `https://terminating-${alreadyRetired}.example/guide.xml`;
|
||||
const fetch = service.fetchEpgFromUrl(url).catch(() => undefined);
|
||||
const worker = mockWorkerInstances[0];
|
||||
let finishTermination!: () => void;
|
||||
worker.terminate.mockReturnValue(
|
||||
new Promise<void>((resolve) => {
|
||||
finishTermination = resolve;
|
||||
})
|
||||
);
|
||||
if (alreadyRetired) {
|
||||
const { retireEpgSource } =
|
||||
await import('./epg-source-generation');
|
||||
retireEpgSource(url);
|
||||
}
|
||||
worker.emit(
|
||||
'error',
|
||||
new Error('worker failed while another source clears')
|
||||
);
|
||||
const clear = service.clearEpgDataForSource(url);
|
||||
const workersBeforeTermination = mockWorkerInstances.length;
|
||||
finishTermination();
|
||||
await fetch;
|
||||
await flushPromises();
|
||||
const clearWorker = mockWorkerInstances[1];
|
||||
clearWorker.emit('message', { type: 'READY' });
|
||||
clearWorker.emit('message', { type: 'CLEAR_COMPLETE' });
|
||||
await clear;
|
||||
expect(workersBeforeTermination).toBe(1);
|
||||
}
|
||||
);
|
||||
|
||||
it('clears one EPG source through a worker and allows it to be fetched again', async () => {
|
||||
const workerService = new EpgWorkerService('[Test EPG]', 1000);
|
||||
const sourceUrl = 'https://playlist.example.com/guide.xml';
|
||||
|
||||
@@ -93,6 +93,10 @@ describe('EpgDatabase', () => {
|
||||
FROM epg_programs
|
||||
WHERE epg_programs.channel_id = epg_channels.id
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM epg_channel_sources
|
||||
WHERE epg_channel_sources.channel_id = epg_channels.id
|
||||
)
|
||||
`);
|
||||
|
||||
expect(preparedSql).toContain(deleteTodayAndFutureSql);
|
||||
@@ -230,6 +234,7 @@ describe('EpgDatabaseClearOperation', () => {
|
||||
expect(exec.mock.calls.map(([statement]) => statement)).toEqual([
|
||||
'BEGIN',
|
||||
'DELETE FROM epg_programs',
|
||||
'DELETE FROM epg_channel_sources',
|
||||
'DELETE FROM epg_channels',
|
||||
'COMMIT',
|
||||
]);
|
||||
@@ -250,6 +255,7 @@ describe('EpgDatabaseClearOperation', () => {
|
||||
expect(exec.mock.calls.map(([statement]) => statement)).toEqual([
|
||||
'BEGIN',
|
||||
'DELETE FROM epg_programs',
|
||||
'DELETE FROM epg_channel_sources',
|
||||
'DELETE FROM epg_channels',
|
||||
'ROLLBACK',
|
||||
]);
|
||||
@@ -276,6 +282,10 @@ describe('EpgDatabaseSourceClearOperation', () => {
|
||||
FROM epg_programs
|
||||
WHERE epg_programs.channel_id = epg_channels.id
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM epg_channel_sources
|
||||
WHERE epg_channel_sources.channel_id = epg_channels.id
|
||||
)
|
||||
`);
|
||||
|
||||
expect(preparedSql).toContain(deleteProgramsSql);
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import type BetterSqlite3 from 'better-sqlite3';
|
||||
import { getIptvnatorDatabasePath } from '@iptvnator/shared/database/path-utils';
|
||||
import type { ParsedChannel, ParsedProgram } from './epg-streaming-parser';
|
||||
import { restoreSurvivingChannelMetadata } from './epg-source-metadata';
|
||||
|
||||
/**
|
||||
* Database helper for worker-owned EPG operations.
|
||||
@@ -10,6 +11,8 @@ export class EpgDatabase {
|
||||
private readonly db: BetterSqlite3.Database;
|
||||
private readonly knownChannelIds = new Set<string>();
|
||||
private readonly insertChannelStmt: BetterSqlite3.Statement;
|
||||
private readonly insertChannelSourceStmt: BetterSqlite3.Statement;
|
||||
private readonly deleteChannelSourcesStmt: BetterSqlite3.Statement;
|
||||
private readonly insertProgramStmt: BetterSqlite3.Statement;
|
||||
private readonly deleteOrphanChannelsForSourceStmt: BetterSqlite3.Statement;
|
||||
private readonly deleteTodayAndFutureStmt: BetterSqlite3.Statement;
|
||||
@@ -30,6 +33,20 @@ export class EpgDatabase {
|
||||
updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
|
||||
`);
|
||||
|
||||
this.insertChannelSourceStmt = this.db.prepare(`
|
||||
INSERT INTO epg_channel_sources
|
||||
(channel_id, display_name, icon_url, url, source_url, updated_at)
|
||||
VALUES (?, ?, ?, ?, ?, strftime('%Y-%m-%dT%H:%M:%SZ', 'now'))
|
||||
ON CONFLICT(channel_id, source_url) DO UPDATE SET
|
||||
display_name = excluded.display_name,
|
||||
icon_url = excluded.icon_url,
|
||||
url = excluded.url,
|
||||
updated_at = excluded.updated_at
|
||||
`);
|
||||
this.deleteChannelSourcesStmt = this.db.prepare(
|
||||
'DELETE FROM epg_channel_sources WHERE source_url = ?'
|
||||
);
|
||||
|
||||
// Guard against duplicate entries when the clearFirst logic misses old
|
||||
// rows. The same channel + start + title + source is treated as the
|
||||
// same programme — a later import with a corrected stop time simply
|
||||
@@ -66,6 +83,10 @@ export class EpgDatabase {
|
||||
FROM epg_programs
|
||||
WHERE epg_programs.channel_id = epg_channels.id
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM epg_channel_sources
|
||||
WHERE epg_channel_sources.channel_id = epg_channels.id
|
||||
)
|
||||
`);
|
||||
|
||||
this.deleteTodayAndFutureStmt = this.db.prepare(`
|
||||
@@ -88,7 +109,9 @@ export class EpgDatabase {
|
||||
): void {
|
||||
const insertMany = this.db.transaction((channels: ParsedChannel[]) => {
|
||||
if (clearTodayAndFuture) {
|
||||
restoreSurvivingChannelMetadata(this.db, sourceUrl);
|
||||
this.deleteTodayAndFutureStmt.run(sourceUrl);
|
||||
this.deleteChannelSourcesStmt.run(sourceUrl);
|
||||
this.deleteOrphanChannelsForSourceStmt.run(sourceUrl);
|
||||
this.knownChannelIds.clear();
|
||||
}
|
||||
@@ -106,6 +129,13 @@ export class EpgDatabase {
|
||||
url,
|
||||
sourceUrl
|
||||
);
|
||||
this.insertChannelSourceStmt.run(
|
||||
channel.id,
|
||||
displayName,
|
||||
iconUrl,
|
||||
url,
|
||||
sourceUrl
|
||||
);
|
||||
this.knownChannelIds.add(channel.id);
|
||||
}
|
||||
});
|
||||
@@ -222,6 +252,7 @@ export class EpgDatabaseClearOperation {
|
||||
this.db.exec('BEGIN');
|
||||
try {
|
||||
this.db.exec('DELETE FROM epg_programs');
|
||||
this.db.exec('DELETE FROM epg_channel_sources');
|
||||
this.db.exec('DELETE FROM epg_channels');
|
||||
this.db.exec('COMMIT');
|
||||
} catch (error) {
|
||||
@@ -256,6 +287,10 @@ export class EpgDatabaseSourceClearOperation {
|
||||
FROM epg_programs
|
||||
WHERE epg_programs.channel_id = epg_channels.id
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM epg_channel_sources
|
||||
WHERE epg_channel_sources.channel_id = epg_channels.id
|
||||
)
|
||||
`);
|
||||
}
|
||||
|
||||
@@ -266,20 +301,11 @@ export class EpgDatabaseSourceClearOperation {
|
||||
}
|
||||
|
||||
const clearSource = this.db.transaction((url: string) => {
|
||||
// Capture affected IDs while removed-source provenance still exists.
|
||||
restoreSurvivingChannelMetadata(this.db, url);
|
||||
this.deleteProgramsForSourceStmt.run(url);
|
||||
// Channels have one legacy global ID. Transfer their ownership
|
||||
// to a surviving source before a subsequent source removal.
|
||||
this.db
|
||||
.prepare(
|
||||
`UPDATE epg_channels SET source_url = (
|
||||
SELECT source_url FROM epg_programs
|
||||
WHERE channel_id = epg_channels.id AND source_url IS NOT NULL
|
||||
ORDER BY source_url LIMIT 1
|
||||
) WHERE source_url = ? AND EXISTS (
|
||||
SELECT 1 FROM epg_programs
|
||||
WHERE channel_id = epg_channels.id AND source_url IS NOT NULL
|
||||
)`
|
||||
)
|
||||
.prepare('DELETE FROM epg_channel_sources WHERE source_url = ?')
|
||||
.run(url);
|
||||
this.deleteOrphanChannelsForSourceStmt.run(url);
|
||||
});
|
||||
|
||||
@@ -0,0 +1,164 @@
|
||||
import { execFileSync } from 'node:child_process';
|
||||
import { createRequire } from 'node:module';
|
||||
import { resolve } from 'node:path';
|
||||
import { pathToFileURL } from 'node:url';
|
||||
|
||||
/** Exercise worker SQL with Electron's native SQLite binding, including upgrades. */
|
||||
function runDatabaseScenario(scenario: string): unknown {
|
||||
const moduleUrl = pathToFileURL(resolve(__dirname, 'epg-database.ts')).href;
|
||||
const connectionUrl = pathToFileURL(
|
||||
resolve(process.cwd(), 'libs/shared/database/src/lib/connection.ts')
|
||||
).href;
|
||||
const script = `
|
||||
const { default: Database } = await import('better-sqlite3');
|
||||
const { EpgDatabase, EpgDatabaseSourceClearOperation, EpgDatabaseClearOperation }
|
||||
= await import(${JSON.stringify(moduleUrl)});
|
||||
const { __databaseConnectionTestHooks } = await import(${JSON.stringify(connectionUrl)});
|
||||
const sqlite = new Database(':memory:');
|
||||
const statements = __databaseConnectionTestHooks.createTableStatements
|
||||
.filter(sql => /(?:TABLE|INDEX|TRIGGER) IF NOT EXISTS (?:idx_)?epg_/.test(sql));
|
||||
for (const sql of statements) sqlite.exec(sql);
|
||||
class SharedDatabase { constructor() { return sqlite; } }
|
||||
const worker = new EpgDatabase(SharedDatabase);
|
||||
const clear = new EpgDatabaseSourceClearOperation(SharedDatabase);
|
||||
const importSource = (source, programs = true, clearFirst = false) => {
|
||||
worker.insertChannels([{
|
||||
id: 'shared', displayName: [{ value: source + ' News' }],
|
||||
icon: [{ src: source + '.png' }], url: [source + '.website']
|
||||
}], source, clearFirst);
|
||||
if (programs) worker.insertPrograms([{
|
||||
channel: 'shared', start: '2099-01-01T00:00:00Z',
|
||||
stop: '2099-01-01T01:00:00Z', title: [{value: source}]
|
||||
}], source);
|
||||
};
|
||||
const channel = () => sqlite.prepare('SELECT * FROM epg_channels WHERE id = ?').get('shared');
|
||||
${scenario}
|
||||
sqlite.close();
|
||||
`;
|
||||
return JSON.parse(
|
||||
execFileSync(
|
||||
createRequire(__filename)('electron'),
|
||||
['--import', 'tsx', '--eval', script],
|
||||
{
|
||||
cwd: process.cwd(),
|
||||
encoding: 'utf8',
|
||||
env: {
|
||||
...process.env,
|
||||
ELECTRON_RUN_AS_NODE: '1',
|
||||
TSX_TSCONFIG_PATH: resolve(
|
||||
process.cwd(),
|
||||
'tsconfig.base.json'
|
||||
),
|
||||
},
|
||||
}
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
describe('source-owned XMLTV channel metadata', () => {
|
||||
it.each(['A', 'B'])(
|
||||
'restores the survivor when removing %s, then removes the final owner',
|
||||
(removed) => {
|
||||
const survivor = removed === 'A' ? 'B' : 'A';
|
||||
const result = runDatabaseScenario(`
|
||||
importSource('A'); importSource('B');
|
||||
clear.run('${removed}');
|
||||
clear.run('${removed}');
|
||||
const surviving = channel();
|
||||
const programs = sqlite.prepare('SELECT source_url FROM epg_programs').all();
|
||||
clear.run('${survivor}');
|
||||
process.stdout.write(JSON.stringify({ surviving, programs, remaining: channel() ?? null }));
|
||||
`);
|
||||
expect(result).toEqual({
|
||||
surviving: expect.objectContaining({
|
||||
display_name: `${survivor} News`,
|
||||
icon_url: `${survivor}.png`,
|
||||
url: `${survivor}.website`,
|
||||
source_url: survivor,
|
||||
}),
|
||||
programs: [{ source_url: survivor }],
|
||||
remaining: null,
|
||||
});
|
||||
}
|
||||
);
|
||||
|
||||
it('retains metadata-only owners during another source refresh and removal', () => {
|
||||
expect(
|
||||
runDatabaseScenario(`
|
||||
importSource('A', false); importSource('B', false);
|
||||
worker.insertChannels([], 'A', true);
|
||||
clear.run('A');
|
||||
const surviving = channel();
|
||||
clear.run('B');
|
||||
process.stdout.write(JSON.stringify({surviving, remaining: channel() ?? null}));
|
||||
`)
|
||||
).toEqual({
|
||||
surviving: expect.objectContaining({
|
||||
display_name: 'B News',
|
||||
source_url: 'B',
|
||||
}),
|
||||
remaining: null,
|
||||
});
|
||||
});
|
||||
|
||||
it('removes an orphan whose original owner disappeared during refresh', () => {
|
||||
expect(
|
||||
runDatabaseScenario(`
|
||||
importSource('A', false); importSource('B', false);
|
||||
worker.insertChannels([], 'A', true);
|
||||
clear.run('B');
|
||||
process.stdout.write(JSON.stringify(channel() ?? null));
|
||||
`)
|
||||
).toBeNull();
|
||||
});
|
||||
|
||||
it('restores metadata before a refresh removes the last writer provenance', () => {
|
||||
expect(
|
||||
runDatabaseScenario(`
|
||||
importSource('A'); importSource('B');
|
||||
worker.insertChannels([], 'B', true);
|
||||
clear.run('B');
|
||||
process.stdout.write(JSON.stringify(channel()));
|
||||
`)
|
||||
).toEqual(
|
||||
expect.objectContaining({
|
||||
display_name: 'A News',
|
||||
icon_url: 'A.png',
|
||||
source_url: 'A',
|
||||
})
|
||||
);
|
||||
});
|
||||
|
||||
it('neutralizes ambiguous legacy metadata while preserving surviving programmes', () => {
|
||||
expect(
|
||||
runDatabaseScenario(`
|
||||
importSource('A'); importSource('B');
|
||||
sqlite.exec('DROP TABLE IF EXISTS epg_channel_sources');
|
||||
for (const sql of statements) sqlite.exec(sql);
|
||||
clear.run('B');
|
||||
process.stdout.write(JSON.stringify(channel()));
|
||||
`)
|
||||
).toEqual(
|
||||
expect.objectContaining({
|
||||
display_name: 'shared',
|
||||
icon_url: null,
|
||||
url: null,
|
||||
source_url: 'A',
|
||||
updated_at: null,
|
||||
})
|
||||
);
|
||||
});
|
||||
|
||||
it('clears all metadata and permits a clean reimport', () => {
|
||||
expect(
|
||||
runDatabaseScenario(`
|
||||
importSource('A'); importSource('B');
|
||||
new EpgDatabaseClearOperation(SharedDatabase).run();
|
||||
const count = sqlite.prepare('SELECT count(*) AS count FROM epg_channel_sources').get().count;
|
||||
importSource('C');
|
||||
clear.run('C');
|
||||
process.stdout.write(JSON.stringify({ count, remaining: channel() ?? null }));
|
||||
`)
|
||||
).toEqual({ count: 0, remaining: null });
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,32 @@
|
||||
import type BetterSqlite3 from 'better-sqlite3';
|
||||
|
||||
/** Called inside the source-clear transaction, before deleting its provenance. */
|
||||
export function restoreSurvivingChannelMetadata(
|
||||
db: BetterSqlite3.Database,
|
||||
sourceUrl: string
|
||||
): void {
|
||||
// Existing global rows cannot be backfilled: their last metadata writer may
|
||||
// differ from source_url. Only imports recorded in the ledger prove origin.
|
||||
const survivingValue = (column: string) => `(
|
||||
SELECT ${column} FROM epg_channel_sources
|
||||
WHERE channel_id = epg_channels.id AND source_url != @sourceUrl
|
||||
ORDER BY updated_at DESC, source_url LIMIT 1
|
||||
)`;
|
||||
db.prepare(
|
||||
`
|
||||
UPDATE epg_channels SET
|
||||
display_name = COALESCE(${survivingValue('display_name')}, id),
|
||||
icon_url = ${survivingValue('icon_url')},
|
||||
url = ${survivingValue('url')},
|
||||
updated_at = ${survivingValue('updated_at')},
|
||||
source_url = COALESCE(${survivingValue('source_url')}, (
|
||||
SELECT source_url FROM epg_programs
|
||||
WHERE channel_id = epg_channels.id AND source_url != @sourceUrl
|
||||
ORDER BY source_url LIMIT 1
|
||||
), @sourceUrl)
|
||||
WHERE source_url = @sourceUrl
|
||||
OR id IN (SELECT channel_id FROM epg_channel_sources WHERE source_url = @sourceUrl)
|
||||
OR id IN (SELECT channel_id FROM epg_programs WHERE source_url = @sourceUrl)
|
||||
`
|
||||
).run({ sourceUrl });
|
||||
}
|
||||
@@ -1329,20 +1329,35 @@ read. Main verifies the completed migration flag, reads every enabled M3U
|
||||
aborts pruning. Detected-but-disabled sources and Xtream/Stalker provider EPG are
|
||||
not global XMLTV owners. No source-discovery or provider matching policy changes.
|
||||
|
||||
Reconciliation finds old sources in both XMLTV tables and queued imports. It
|
||||
Reconciliation finds old sources in XMLTV channel, programme and per-source
|
||||
metadata tables, plus queued imports. It
|
||||
retires their generations before waiting for workers to exit, then uses the
|
||||
existing source-clear worker. Successfully cleared request candidates are forgotten
|
||||
without resetting their generation fences; failed cleanups remain retryable.
|
||||
Same-URL clears are serialized and replacement imports await the outstanding
|
||||
clear, so an older cleanup cannot erase a newly re-added source.
|
||||
clear, so an older cleanup cannot erase a newly re-added source. Source cleanup
|
||||
also awaits the in-flight fetch promise when an error/timeout has already removed
|
||||
the worker from the lookup but its termination is still pending.
|
||||
Retired queued and running imports emit generation-scoped cancellation so progress
|
||||
rows disappear without reporting routine worker termination as an import failure. Programmes are deleted by source; a globally keyed
|
||||
channel is retained while another source still has programmes, transferring its
|
||||
legacy owner to that remaining source. Manual mappings are preserved and can
|
||||
channel is retained while another source still has programmes or channel metadata.
|
||||
The additive `epg_channel_sources` table is created during database initialization
|
||||
and records each imported source's name, logo, URL and timestamp. Before removing
|
||||
source provenance, cleanup restores affected global channels from a surviving
|
||||
snapshot, including when the legacy channel owner differs from the removed
|
||||
metadata writer. Metadata-only owners survive another source's refresh or removal.
|
||||
Refresh restores affected metadata before discarding that source's old snapshots;
|
||||
clear-all deletes the snapshots too. Freshness reads source-specific snapshot
|
||||
timestamps, so legacy sources without snapshots are refreshed on their next import
|
||||
check. No legacy snapshot is guessed from global
|
||||
channel metadata: its last writer may differ from its recorded owner. Affected
|
||||
legacy channels without a surviving snapshot retain programmes with a neutral
|
||||
XMLTV ID label, no logo/URL and no freshness timestamp until reimport. Historical
|
||||
metadata with no remaining source provenance cannot be selectively reconstructed. Manual mappings are preserved and can
|
||||
resolve another retained source sharing that channel ID. Legacy programmes with
|
||||
unknown (`NULL`) ownership are conservatively left alone; the existing database
|
||||
initialization backfill handles rows whose channel still identifies their owner.
|
||||
There is no new schema migration.
|
||||
|
||||
|
||||
Renderer reconciliation fences lookups before its first asynchronous step.
|
||||
Imports wait for serialized reconciliation (including playlist migration), then
|
||||
|
||||
@@ -288,6 +288,17 @@ const CREATE_TABLE_STATEMENTS = [
|
||||
source_url TEXT NOT NULL,
|
||||
updated_at TEXT DEFAULT (datetime('now'))
|
||||
)`,
|
||||
`CREATE TABLE IF NOT EXISTS epg_channel_sources (
|
||||
channel_id TEXT NOT NULL,
|
||||
source_url TEXT NOT NULL,
|
||||
display_name TEXT NOT NULL,
|
||||
icon_url TEXT,
|
||||
url TEXT,
|
||||
updated_at TEXT DEFAULT (datetime('now')),
|
||||
PRIMARY KEY (channel_id, source_url),
|
||||
FOREIGN KEY (channel_id) REFERENCES epg_channels(id) ON DELETE CASCADE
|
||||
)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_epg_channel_sources_source ON epg_channel_sources(source_url)`,
|
||||
`CREATE TABLE IF NOT EXISTS epg_programs (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
channel_id TEXT NOT NULL,
|
||||
|
||||
@@ -11,6 +11,7 @@ import { sql } from 'drizzle-orm';
|
||||
import {
|
||||
index,
|
||||
integer,
|
||||
primaryKey,
|
||||
sqliteTable,
|
||||
text,
|
||||
uniqueIndex,
|
||||
@@ -204,6 +205,25 @@ export const epgChannels = sqliteTable(
|
||||
})
|
||||
);
|
||||
|
||||
// A global XMLTV ID can be shared by sources with different channel metadata.
|
||||
export const epgChannelSources = sqliteTable(
|
||||
'epg_channel_sources',
|
||||
{
|
||||
channelId: text('channel_id')
|
||||
.notNull()
|
||||
.references(() => epgChannels.id, { onDelete: 'cascade' }),
|
||||
sourceUrl: text('source_url').notNull(),
|
||||
displayName: text('display_name').notNull(),
|
||||
iconUrl: text('icon_url'),
|
||||
url: text('url'),
|
||||
updatedAt: text('updated_at').default(sql`CURRENT_TIMESTAMP`),
|
||||
},
|
||||
(table) => [
|
||||
primaryKey({ columns: [table.channelId, table.sourceUrl] }),
|
||||
index('idx_epg_channel_sources_source').on(table.sourceUrl),
|
||||
]
|
||||
);
|
||||
|
||||
// EPG Programs table
|
||||
export const epgPrograms = sqliteTable(
|
||||
'epg_programs',
|
||||
|
||||
Reference in new issue
Block a user