From c3deb56bff66b2f8aea3248d4c02b6b078e12715 Mon Sep 17 00:00:00 2001 From: 4gray Date: Sun, 1 Feb 2026 22:54:00 +0100 Subject: [PATCH] refactor: remove expired program cleanup functionality and adjust EPG program queries --- .../src/app/events/database/epg-db.events.ts | 133 ++++++--------- .../src/app/events/epg.events.ts | 135 ++------------- .../src/app/workers/epg-parser.worker.ts | 161 ++++++------------ 3 files changed, 117 insertions(+), 312 deletions(-) diff --git a/apps/electron-backend/src/app/events/database/epg-db.events.ts b/apps/electron-backend/src/app/events/database/epg-db.events.ts index 2dfedb00e..1b7edb4fb 100644 --- a/apps/electron-backend/src/app/events/database/epg-db.events.ts +++ b/apps/electron-backend/src/app/events/database/epg-db.events.ts @@ -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 diff --git a/apps/electron-backend/src/app/events/epg.events.ts b/apps/electron-backend/src/app/events/epg.events.ts index 763a3ae64..ad53ee5db 100644 --- a/apps/electron-backend/src/app/events/epg.events.ts +++ b/apps/electron-backend/src/app/events/epg.events.ts @@ -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 { 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 }, diff --git a/apps/electron-backend/src/app/workers/epg-parser.worker.ts b/apps/electron-backend/src/app/workers/epg-parser.worker.ts index e20115c9c..5601b526f 100644 --- a/apps/electron-backend/src/app/workers/epg-parser.worker.ts +++ b/apps/electron-backend/src/app/workers/epg-parser.worker.ts @@ -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 | 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 { 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);