refactor: remove expired program cleanup functionality and adjust EPG program queries

This commit is contained in:
4gray committed 2026-02-01 22:54:00 +01:00
1 parent bb1f3af082
commit c3deb56bff
3 files changed
+117 -312

No files matched your search

@@ -128,12 +128,8 @@ ipcMain.handle(
async (_event, channelId: string, fromTime?: string, toTime?: string) => {
try {
const db = await getDatabase();
const now = new Date().toISOString();
let query = db
.select()
.from(schema.epgPrograms)
.where(eq(schema.epgPrograms.channelId, channelId));
let query;
// Apply time filters if provided
if (fromTime && toTime) {
@@ -148,16 +144,11 @@ ipcMain.handle(
)
);
} else {
// Default: from now onwards
// Return all programs for this channel
query = db
.select()
.from(schema.epgPrograms)
.where(
and(
eq(schema.epgPrograms.channelId, channelId),
gte(schema.epgPrograms.stop, now)
)
);
.where(eq(schema.epgPrograms.channelId, channelId));
}
const results = await query.orderBy(schema.epgPrograms.start);
@@ -176,29 +167,32 @@ ipcMain.handle(
/**
* Get current program for a channel (what's on now)
*/
ipcMain.handle('EPG_DB_GET_CURRENT_PROGRAM', async (_event, channelId: string) => {
try {
const db = await getDatabase();
const now = new Date().toISOString();
ipcMain.handle(
'EPG_DB_GET_CURRENT_PROGRAM',
async (_event, channelId: string) => {
try {
const db = await getDatabase();
const now = new Date().toISOString();
const result = await db
.select()
.from(schema.epgPrograms)
.where(
and(
eq(schema.epgPrograms.channelId, channelId),
lte(schema.epgPrograms.start, now),
gte(schema.epgPrograms.stop, now)
const result = await db
.select()
.from(schema.epgPrograms)
.where(
and(
eq(schema.epgPrograms.channelId, channelId),
lte(schema.epgPrograms.start, now),
gte(schema.epgPrograms.stop, now)
)
)
)
.limit(1);
.limit(1);
return result[0] || null;
} catch (error) {
console.error(loggerLabel, 'Error getting current program:', error);
throw error;
return result[0] || null;
} catch (error) {
console.error(loggerLabel, 'Error getting current program:', error);
throw error;
}
}
});
);
/**
* Full-text search EPG programs using FTS5 with LIKE fallback
@@ -210,7 +204,6 @@ ipcMain.handle(
async (_event, searchTerm: string, limit = 50) => {
try {
const db = await getDatabase();
const now = new Date().toISOString();
const trimmedTerm = searchTerm.trim();
if (!trimmedTerm) {
@@ -222,6 +215,7 @@ ipcMain.handle(
const likePattern = `%${trimmedTerm}%`;
// JOIN with epg_channels to get channel display name
// Include all programs (past and future) for catchup/archive feature
const results = await db.all(sql`
SELECT
p.*,
@@ -233,7 +227,6 @@ ipcMain.handle(
OR p.description LIKE ${likePattern}
OR p.category LIKE ${likePattern}
)
AND p.stop >= ${now}
ORDER BY p.start
LIMIT ${limit}
`);
@@ -266,59 +259,37 @@ ipcMain.handle('EPG_DB_GET_CHANNELS', async () => {
/**
* Get channel by ID or display name
*/
ipcMain.handle('EPG_DB_GET_CHANNEL', async (_event, channelIdOrName: string) => {
try {
const db = await getDatabase();
ipcMain.handle(
'EPG_DB_GET_CHANNEL',
async (_event, channelIdOrName: string) => {
try {
const db = await getDatabase();
// Try exact ID match first
let result = await db
.select()
.from(schema.epgChannels)
.where(eq(schema.epgChannels.id, channelIdOrName))
.limit(1);
// Try exact ID match first
let result = await db
.select()
.from(schema.epgChannels)
.where(eq(schema.epgChannels.id, channelIdOrName))
.limit(1);
if (result.length > 0) return result[0];
if (result.length > 0) return result[0];
// Try display name match (case-insensitive)
result = await db
.select()
.from(schema.epgChannels)
.where(
sql`LOWER(${schema.epgChannels.displayName}) = LOWER(${channelIdOrName})`
)
.limit(1);
// Try display name match (case-insensitive)
result = await db
.select()
.from(schema.epgChannels)
.where(
sql`LOWER(${schema.epgChannels.displayName}) = LOWER(${channelIdOrName})`
)
.limit(1);
return result[0] || null;
} catch (error) {
console.error(loggerLabel, 'Error getting EPG channel:', error);
throw error;
return result[0] || null;
} catch (error) {
console.error(loggerLabel, 'Error getting EPG channel:', error);
throw error;
}
}
});
/**
* Cleanup expired programs (older than specified hours)
*/
ipcMain.handle('EPG_DB_CLEANUP_EXPIRED', async (_event, hoursToKeep = 24) => {
try {
const db = await getDatabase();
const cutoff = new Date(
Date.now() - hoursToKeep * 60 * 60 * 1000
).toISOString();
const result = await db
.delete(schema.epgPrograms)
.where(lte(schema.epgPrograms.stop, cutoff));
console.log(
loggerLabel,
`Cleaned up expired programs (older than ${hoursToKeep}h)`
);
return { success: true };
} catch (error) {
console.error(loggerLabel, 'Error cleaning up EPG programs:', error);
throw error;
}
});
);
/**
* Clear all EPG data for a specific source URL
@@ -1,4 +1,4 @@
import { and, eq, gte, sql } from 'drizzle-orm';
import { eq, sql } from 'drizzle-orm';
import { app, BrowserWindow, ipcMain } from 'electron';
import * as path from 'path';
import { EpgProgram } from 'shared-interfaces';
@@ -75,14 +75,6 @@ export default class EpgEvents {
return await this.handleFetchEpg([url]);
});
// Cleanup expired programs (uses worker thread)
ipcMain.handle(
'EPG_CLEANUP_EXPIRED',
async (_event, hoursToKeep?: number) => {
return this.runCleanupInWorker(hoursToKeep ?? 24);
}
);
// Clear all EPG data
ipcMain.handle('EPG_CLEAR_ALL', async () => {
await this.clearEpgData();
@@ -280,7 +272,11 @@ export default class EpgEvents {
// In packaged app, native modules are in app.asar.unpacked/node_modules
// which is separate from the worker location in extraResources
const nativeModulesPath = app.isPackaged
? path.join(path.dirname(app.getAppPath()), 'app.asar.unpacked', 'node_modules')
? path.join(
path.dirname(app.getAppPath()),
'app.asar.unpacked',
'node_modules'
)
: undefined;
worker = new Worker(workerURL, {
resourceLimits: {
@@ -348,18 +344,6 @@ export default class EpgEvents {
message.stats
);
this.fetchedUrls.add(url);
// Trigger cleanup in worker thread (non-blocking)
worker.postMessage({
type: 'CLEANUP_EXPIRED',
hoursToKeep: 24,
});
break;
case 'CLEANUP_COMPLETE':
console.log(
this.loggerLabel,
'Cleanup complete, terminating worker'
);
worker.terminate();
this.workers.delete(url);
resolve();
@@ -450,20 +434,14 @@ export default class EpgEvents {
): Promise<EpgProgram[]> {
try {
const db = await getDatabase();
const now = new Date().toISOString();
// Try exact channel ID match first
let results = await db
.select()
.from(schema.epgPrograms)
.where(
and(
eq(schema.epgPrograms.channelId, channelId),
gte(schema.epgPrograms.stop, now)
)
)
.where(eq(schema.epgPrograms.channelId, channelId))
.orderBy(schema.epgPrograms.start)
.limit(100);
.limit(500);
if (results.length > 0) {
return results.map(this.transformDbRowToEpgProgram);
@@ -482,14 +460,9 @@ export default class EpgEvents {
results = await db
.select()
.from(schema.epgPrograms)
.where(
and(
eq(schema.epgPrograms.channelId, channel[0].id),
gte(schema.epgPrograms.stop, now)
)
)
.where(eq(schema.epgPrograms.channelId, channel[0].id))
.orderBy(schema.epgPrograms.start)
.limit(100);
.limit(500);
return results.map(this.transformDbRowToEpgProgram);
}
@@ -587,88 +560,6 @@ export default class EpgEvents {
}
}
/**
* Run cleanup in worker thread to avoid blocking main thread
*/
private static async runCleanupInWorker(
hoursToKeep: number
): Promise<{ success: boolean }> {
return new Promise((resolve, reject) => {
let workerPath: string;
if (app.isPackaged) {
const resourcesPath = path.dirname(app.getAppPath());
workerPath = path.join(
resourcesPath,
'dist',
'apps',
'electron-backend',
'workers',
'epg-parser.worker.js'
);
} else {
workerPath = path.join(
__dirname,
'workers',
'epg-parser.worker.js'
);
}
let worker: Worker;
try {
const workerURL = pathToFileURL(workerPath);
// In packaged app, native modules are in app.asar.unpacked/node_modules
const nativeModulesPath = app.isPackaged
? path.join(path.dirname(app.getAppPath()), 'app.asar.unpacked', 'node_modules')
: undefined;
worker = new Worker(workerURL, {
workerData: { nativeModulesPath },
});
} catch (error) {
console.error(
this.loggerLabel,
'Failed to create worker for cleanup:',
error
);
reject(error);
return;
}
worker.on(
'message',
(message: { type: string; error?: string }) => {
if (message.type === 'READY') {
worker.postMessage({
type: 'CLEANUP_EXPIRED',
hoursToKeep,
});
} else if (message.type === 'CLEANUP_COMPLETE') {
worker.terminate();
resolve({ success: true });
} else if (message.type === 'EPG_ERROR') {
console.error(
this.loggerLabel,
'Worker cleanup error:',
message.error
);
worker.terminate();
reject(new Error(message.error || 'Cleanup failed'));
}
}
);
worker.on('error', (error) => {
console.error(
this.loggerLabel,
'Worker error during cleanup:',
error
);
worker.terminate();
reject(error);
});
});
}
/**
* Clear all EPG data using worker thread to avoid blocking main thread
*/
@@ -699,7 +590,11 @@ export default class EpgEvents {
const workerURL = pathToFileURL(workerPath);
// In packaged app, native modules are in app.asar.unpacked/node_modules
const nativeModulesPath = app.isPackaged
? path.join(path.dirname(app.getAppPath()), 'app.asar.unpacked', 'node_modules')
? path.join(
path.dirname(app.getAppPath()),
'app.asar.unpacked',
'node_modules'
)
: undefined;
worker = new Worker(workerURL, {
workerData: { nativeModulesPath },
@@ -1,12 +1,12 @@
import { parentPort, workerData } from 'worker_threads';
import { createGunzip } from 'zlib';
import type BetterSqlite3 from 'better-sqlite3';
import { existsSync, mkdirSync } from 'fs';
import { createRequire } from 'module';
import { homedir } from 'os';
import { join } from 'path';
import { SaxesParser, SaxesTagPlain } from 'saxes';
import { Readable } from 'stream';
import { existsSync, mkdirSync } from 'fs';
import { homedir } from 'os';
import { join, dirname } from 'path';
import { createRequire } from 'module';
import type BetterSqlite3 from 'better-sqlite3';
import { parentPort, workerData } from 'worker_threads';
import { createGunzip } from 'zlib';
// In packaged app, native modules are in app.asar.unpacked/node_modules
// which is separate from the worker location in extraResources
@@ -14,25 +14,46 @@ let Database: typeof BetterSqlite3;
function loadBetterSqlite3(): typeof BetterSqlite3 {
// Try workerData path first (passed from main process)
if (workerData?.nativeModulesPath && existsSync(workerData.nativeModulesPath)) {
if (
workerData?.nativeModulesPath &&
existsSync(workerData.nativeModulesPath)
) {
try {
const nativeRequire = createRequire(join(workerData.nativeModulesPath, 'index.js'));
const nativeRequire = createRequire(
join(workerData.nativeModulesPath, 'index.js')
);
return nativeRequire('better-sqlite3');
} catch (e) {
console.error('[EPG Worker] Failed to load from workerData path:', e);
console.error(
'[EPG Worker] Failed to load from workerData path:',
e
);
}
}
// Try process.resourcesPath (available in packaged Electron apps)
if ((process as NodeJS.Process & { resourcesPath?: string }).resourcesPath) {
const resourcesPath = (process as NodeJS.Process & { resourcesPath?: string }).resourcesPath!;
const unpackedPath = join(resourcesPath, 'app.asar.unpacked', 'node_modules');
if (
(process as NodeJS.Process & { resourcesPath?: string }).resourcesPath
) {
const resourcesPath = (
process as NodeJS.Process & { resourcesPath?: string }
).resourcesPath!;
const unpackedPath = join(
resourcesPath,
'app.asar.unpacked',
'node_modules'
);
if (existsSync(unpackedPath)) {
try {
const nativeRequire = createRequire(join(unpackedPath, 'index.js'));
const nativeRequire = createRequire(
join(unpackedPath, 'index.js')
);
return nativeRequire('better-sqlite3');
} catch (e) {
console.error('[EPG Worker] Failed to load from resourcesPath:', e);
console.error(
'[EPG Worker] Failed to load from resourcesPath:',
e
);
}
}
}
@@ -97,9 +118,8 @@ interface ParsedProgram {
*/
interface WorkerMessage {
type: 'FETCH_EPG' | 'FORCE_FETCH' | 'CLEAR_EPG' | 'CLEANUP_EXPIRED';
type: 'FETCH_EPG' | 'FORCE_FETCH' | 'CLEAR_EPG';
url?: string;
hoursToKeep?: number;
}
interface WorkerResponse {
@@ -108,7 +128,6 @@ interface WorkerResponse {
| 'EPG_ERROR'
| 'EPG_PROGRESS'
| 'CLEAR_COMPLETE'
| 'CLEANUP_COMPLETE'
| 'READY';
error?: string;
url?: string;
@@ -124,34 +143,6 @@ const loggerLabel = '[EPG Worker]';
const CHANNEL_BATCH_SIZE = 100;
const PROGRAM_BATCH_SIZE = 1000;
// Skip programs that ended more than this many hours ago
const SKIP_PROGRAMS_OLDER_THAN_HOURS = 2;
/**
* Calculate the cutoff time for filtering old programs
* Programs that ended before this time will be skipped
*/
function getProgramCutoffTime(): Date {
return new Date(Date.now() - SKIP_PROGRAMS_OLDER_THAN_HOURS * 60 * 60 * 1000);
}
/**
* Check if a program has already ended (is expired)
* Handles timezone-aware ISO date strings
*/
function isProgramExpired(stopTime: string, cutoff: Date): boolean {
if (!stopTime) return true; // Skip programs without stop time
try {
const stopDate = new Date(stopTime);
// Check if the date is valid
if (isNaN(stopDate.getTime())) return false; // Don't skip if we can't parse
return stopDate < cutoff;
} catch {
return false; // Don't skip if parsing fails
}
}
/**
* Get database file path (same as main app)
*/
@@ -220,7 +211,8 @@ class EpgDatabase {
insertChannels(channels: ParsedChannel[], sourceUrl: string): void {
const insertMany = this.db.transaction((channels: ParsedChannel[]) => {
for (const channel of channels) {
const displayName = channel.displayName?.[0]?.value || channel.id;
const displayName =
channel.displayName?.[0]?.value || channel.id;
const iconUrl = channel.icon?.[0]?.src || null;
const url = channel.url?.[0] || null;
@@ -323,10 +315,6 @@ class StreamingEpgParser {
private programs: ParsedProgram[] = [];
private totalChannels = 0;
private totalPrograms = 0;
private skippedPrograms = 0;
// Cutoff time for filtering old programs (calculated once at start)
private readonly cutoffTime: Date;
// Current element being parsed
private currentChannel: Partial<ParsedChannel> | null = null;
@@ -343,7 +331,6 @@ class StreamingEpgParser {
private onProgress: (channels: number, programs: number) => void
) {
this.parser = new SaxesParser();
this.cutoffTime = getProgramCutoffTime();
this.setupParser();
}
@@ -452,7 +439,9 @@ class StreamingEpgParser {
if (text) this.currentChannel.url!.push(text);
break;
case 'channel':
this.channels.push(this.currentChannel as ParsedChannel);
this.channels.push(
this.currentChannel as ParsedChannel
);
this.totalChannels++;
this.currentChannel = null;
@@ -504,16 +493,13 @@ class StreamingEpgParser {
}
break;
case 'programme':
// Skip programs that have already ended
if (isProgramExpired(this.currentProgram.stop!, this.cutoffTime)) {
this.skippedPrograms++;
} else {
this.programs.push(this.currentProgram as ParsedProgram);
this.totalPrograms++;
this.programs.push(
this.currentProgram as ParsedProgram
);
this.totalPrograms++;
if (this.programs.length >= PROGRAM_BATCH_SIZE) {
this.flushPrograms();
}
if (this.programs.length >= PROGRAM_BATCH_SIZE) {
this.flushPrograms();
}
this.currentProgram = null;
break;
@@ -549,22 +535,14 @@ class StreamingEpgParser {
this.parser.write(chunk);
}
finish(): { totalChannels: number; totalPrograms: number; skippedPrograms: number } {
finish(): { totalChannels: number; totalPrograms: number } {
this.parser.close();
this.flushChannels();
this.flushPrograms();
if (this.skippedPrograms > 0) {
console.log(
loggerLabel,
`Skipped ${this.skippedPrograms} expired programs (ended more than ${SKIP_PROGRAMS_OLDER_THAN_HOURS}h ago)`
);
}
return {
totalChannels: this.totalChannels,
totalPrograms: this.totalPrograms,
skippedPrograms: this.skippedPrograms,
};
}
}
@@ -650,10 +628,7 @@ async function fetchAndParseEpgStreaming(url: string): Promise<void> {
const stats = parser.finish();
console.log(
loggerLabel,
`Parsing complete: ${stats.totalChannels} channels, ${stats.totalPrograms} programs` +
(stats.skippedPrograms > 0
? ` (skipped ${stats.skippedPrograms} expired)`
: '')
`Parsing complete: ${stats.totalChannels} channels, ${stats.totalPrograms} programs`
);
// Close database connection
@@ -719,40 +694,6 @@ function clearAllEpgData(): void {
}
}
/**
* Cleans up expired EPG programs from the database
* Runs in worker thread to avoid blocking main thread
*/
function cleanupExpiredPrograms(hoursToKeep: number): void {
const dbPath = getDatabasePath();
const db = new Database(dbPath);
try {
const cutoff = new Date(
Date.now() - hoursToKeep * 60 * 60 * 1000
).toISOString();
console.log(loggerLabel, `Cleaning up programs older than ${hoursToKeep}h (before ${cutoff})...`);
const stmt = db.prepare('DELETE FROM epg_programs WHERE stop <= ?');
const result = stmt.run(cutoff);
console.log(loggerLabel, `Cleaned up ${result.changes} expired programs`);
const response: WorkerResponse = { type: 'CLEANUP_COMPLETE' };
parentPort?.postMessage(response);
} catch (error) {
console.error(loggerLabel, 'Error cleaning up expired programs:', error);
const errorResponse: WorkerResponse = {
type: 'EPG_ERROR',
error: error instanceof Error ? error.message : String(error),
};
parentPort?.postMessage(errorResponse);
} finally {
db.close();
}
}
/**
* Worker message handler
*/
@@ -766,8 +707,6 @@ if (parentPort) {
await fetchAndParseEpgStreaming(message.url!);
} else if (message.type === 'CLEAR_EPG') {
clearAllEpgData();
} else if (message.type === 'CLEANUP_EXPIRED') {
cleanupExpiredPrograms(message.hoursToKeep ?? 24);
}
} catch (error) {
console.error(loggerLabel, 'Worker error:', error);