diff --git a/apps/electron-backend/src/app/events/epg.events.ts b/apps/electron-backend/src/app/events/epg.events.ts index 01523102d..a74c89744 100644 --- a/apps/electron-backend/src/app/events/epg.events.ts +++ b/apps/electron-backend/src/app/events/epg.events.ts @@ -1,11 +1,11 @@ import { and, eq, gte, sql } from 'drizzle-orm'; import { app, BrowserWindow, ipcMain } from 'electron'; import * as path from 'path'; -import { Worker } from 'worker_threads'; +import { EpgProgram } from 'shared-interfaces'; import { pathToFileURL } from 'url'; +import { Worker } from 'worker_threads'; import { getDatabase } from '../database/connection'; import * as schema from '../database/schema'; -import { EpgProgram } from 'shared-interfaces'; /** * EPG Events Handler @@ -96,7 +96,10 @@ export default class EpgEvents { _event, args: { urls: string[]; maxAgeHours?: number } ): Promise<{ staleUrls: string[]; freshUrls: string[] }> => { - return this.checkEpgFreshness(args.urls, args.maxAgeHours ?? 12); + return this.checkEpgFreshness( + args.urls, + args.maxAgeHours ?? 12 + ); } ); @@ -143,16 +146,26 @@ export default class EpgEvents { } } } catch (error) { - console.error(this.loggerLabel, 'Error checking EPG freshness:', error); + console.error( + this.loggerLabel, + 'Error checking EPG freshness:', + error + ); return { staleUrls: urls, freshUrls: [] }; } // Log summary if (freshUrls.length > 0) { - console.log(this.loggerLabel, `EPG fresh (skipping): ${freshUrls.length} source(s)`); + console.log( + this.loggerLabel, + `EPG fresh (skipping): ${freshUrls.length} source(s)` + ); } if (staleUrls.length > 0) { - console.log(this.loggerLabel, `EPG stale (will fetch): ${staleUrls.length} source(s)`); + console.log( + this.loggerLabel, + `EPG stale (will fetch): ${staleUrls.length} source(s)` + ); } return { staleUrls, freshUrls }; @@ -173,15 +186,28 @@ export default class EpgEvents { } // Check which URLs have fresh data and can be skipped - const { staleUrls, freshUrls } = await this.checkEpgFreshness(validUrls, 12); + const { staleUrls, freshUrls } = await this.checkEpgFreshness( + validUrls, + 12 + ); if (staleUrls.length === 0) { - return { success: true, message: 'All EPG data is fresh', skipped: freshUrls }; + return { + success: true, + message: 'All EPG data is fresh', + skipped: freshUrls, + }; } // Send queued status for all stale URLs first staleUrls.forEach((url, index) => { - this.sendProgressToRenderer(url, 'queued', undefined, undefined, index + 1); + this.sendProgressToRenderer( + url, + 'queued', + undefined, + undefined, + index + 1 + ); }); // Process only stale URLs sequentially to avoid database locking @@ -191,8 +217,14 @@ export default class EpgEvents { try { await this.fetchEpgFromUrl(url); } catch (error) { - console.error(this.loggerLabel, `Error fetching EPG from ${url}:`, error); - errors.push(error instanceof Error ? error.message : String(error)); + console.error( + this.loggerLabel, + `Error fetching EPG from ${url}:`, + error + ); + errors.push( + error instanceof Error ? error.message : String(error) + ); } } @@ -214,7 +246,10 @@ export default class EpgEvents { private static async fetchEpgFromUrl(url: string): Promise { // Skip if already fetched this session if (this.fetchedUrls.has(url)) { - console.log(this.loggerLabel, `Skipping already fetched URL: ${url}`); + console.log( + this.loggerLabel, + `Skipping already fetched URL: ${url}` + ); return; } @@ -232,7 +267,11 @@ export default class EpgEvents { 'epg-parser.worker.js' ); } else { - workerPath = path.join(__dirname, 'workers', 'epg-parser.worker.js'); + workerPath = path.join( + __dirname, + 'workers', + 'epg-parser.worker.js' + ); } let worker: Worker; @@ -245,7 +284,11 @@ export default class EpgEvents { }, }); } catch (error) { - console.error(this.loggerLabel, 'Failed to create worker:', error); + console.error( + this.loggerLabel, + 'Failed to create worker:', + error + ); reject(error); return; } @@ -264,7 +307,10 @@ export default class EpgEvents { switch (message.type) { case 'READY': // Notify renderer that loading started - this.sendProgressToRenderer(url, 'loading', { totalChannels: 0, totalPrograms: 0 }); + this.sendProgressToRenderer(url, 'loading', { + totalChannels: 0, + totalPrograms: 0, + }); worker.postMessage({ type: 'FETCH_EPG', url }); break; @@ -275,7 +321,11 @@ export default class EpgEvents { `Progress: ${message.stats.totalChannels} channels, ${message.stats.totalPrograms} programs` ); // Forward progress to renderer - this.sendProgressToRenderer(url, 'loading', message.stats); + this.sendProgressToRenderer( + url, + 'loading', + message.stats + ); } break; @@ -286,10 +336,17 @@ export default class EpgEvents { message.stats ); // Notify renderer of completion - this.sendProgressToRenderer(url, 'complete', message.stats); + this.sendProgressToRenderer( + url, + 'complete', + message.stats + ); this.fetchedUrls.add(url); // Trigger cleanup in worker thread (non-blocking) - worker.postMessage({ type: 'CLEANUP_EXPIRED', hoursToKeep: 24 }); + worker.postMessage({ + type: 'CLEANUP_EXPIRED', + hoursToKeep: 24, + }); break; case 'CLEANUP_COMPLETE': @@ -309,14 +366,25 @@ export default class EpgEvents { message.error ); // Notify renderer of error - this.sendProgressToRenderer(url, 'error', undefined, message.error); + this.sendProgressToRenderer( + url, + 'error', + undefined, + message.error + ); worker.terminate(); this.workers.delete(url); - reject(new Error(message.error || 'Unknown error')); + reject( + new Error(message.error || 'Unknown error') + ); break; } } catch (err) { - console.error(this.loggerLabel, 'Error handling message:', err); + console.error( + this.loggerLabel, + 'Error handling message:', + err + ); reject(err); } } @@ -371,7 +439,9 @@ export default class EpgEvents { /** * Get programs for a specific channel from database */ - private static async handleGetChannelPrograms(channelId: string): Promise { + private static async handleGetChannelPrograms( + channelId: string + ): Promise { try { const db = await getDatabase(); const now = new Date().toISOString(); @@ -420,7 +490,11 @@ export default class EpgEvents { return []; } catch (error) { - console.error(this.loggerLabel, 'Error getting channel programs:', error); + console.error( + this.loggerLabel, + 'Error getting channel programs:', + error + ); return []; } } @@ -444,7 +518,11 @@ export default class EpgEvents { return { channels, programs: [] }; } catch (error) { - console.error(this.loggerLabel, 'Error getting all channels:', error); + console.error( + this.loggerLabel, + 'Error getting all channels:', + error + ); return { channels: [], programs: [] }; } } @@ -455,7 +533,14 @@ export default class EpgEvents { private static async handleGetChannelsByRange( skip: number, limit: number - ): Promise> { + ): Promise< + Array<{ + id: string; + displayName: string; + iconUrl: string | null; + programs: EpgProgram[]; + }> + > { try { const db = await getDatabase(); const channels = await db @@ -487,7 +572,11 @@ export default class EpgEvents { return channelsWithPrograms; } catch (error) { - console.error(this.loggerLabel, 'Error getting channels by range:', error); + console.error( + this.loggerLabel, + 'Error getting channels by range:', + error + ); return []; } } @@ -495,7 +584,9 @@ export default class EpgEvents { /** * Run cleanup in worker thread to avoid blocking main thread */ - private static async runCleanupInWorker(hoursToKeep: number): Promise<{ success: boolean }> { + private static async runCleanupInWorker( + hoursToKeep: number + ): Promise<{ success: boolean }> { return new Promise((resolve, reject) => { let workerPath: string; @@ -510,7 +601,11 @@ export default class EpgEvents { 'epg-parser.worker.js' ); } else { - workerPath = path.join(__dirname, 'workers', 'epg-parser.worker.js'); + workerPath = path.join( + __dirname, + 'workers', + 'epg-parser.worker.js' + ); } let worker: Worker; @@ -518,26 +613,44 @@ export default class EpgEvents { const workerURL = pathToFileURL(workerPath); worker = new Worker(workerURL); } catch (error) { - console.error(this.loggerLabel, 'Failed to create worker for cleanup:', 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( + '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); + console.error( + this.loggerLabel, + 'Worker error during cleanup:', + error + ); worker.terminate(); reject(error); }); @@ -562,7 +675,11 @@ export default class EpgEvents { 'epg-parser.worker.js' ); } else { - workerPath = path.join(__dirname, 'workers', 'epg-parser.worker.js'); + workerPath = path.join( + __dirname, + 'workers', + 'epg-parser.worker.js' + ); } let worker: Worker; @@ -570,31 +687,49 @@ export default class EpgEvents { const workerURL = pathToFileURL(workerPath); worker = new Worker(workerURL); } catch (error) { - console.error(this.loggerLabel, 'Failed to create worker for clear:', error); + console.error( + this.loggerLabel, + 'Failed to create worker for clear:', + error + ); reject(error); return; } - worker.on('message', (message: { type: string; error?: string }) => { - if (message.type === 'READY') { - worker.postMessage({ type: 'CLEAR_EPG' }); - } else if (message.type === 'CLEAR_COMPLETE') { - console.log(this.loggerLabel, 'EPG data cleared via worker'); - this.fetchedUrls.clear(); - // Terminate any running fetch workers - this.workers.forEach((w) => w.terminate()); - this.workers.clear(); - worker.terminate(); - resolve(); - } else if (message.type === 'EPG_ERROR') { - console.error(this.loggerLabel, 'Worker clear error:', message.error); - worker.terminate(); - reject(new Error(message.error || 'Clear failed')); + worker.on( + 'message', + (message: { type: string; error?: string }) => { + if (message.type === 'READY') { + worker.postMessage({ type: 'CLEAR_EPG' }); + } else if (message.type === 'CLEAR_COMPLETE') { + console.log( + this.loggerLabel, + 'EPG data cleared via worker' + ); + this.fetchedUrls.clear(); + // Terminate any running fetch workers + this.workers.forEach((w) => w.terminate()); + this.workers.clear(); + worker.terminate(); + resolve(); + } else if (message.type === 'EPG_ERROR') { + console.error( + this.loggerLabel, + 'Worker clear error:', + message.error + ); + worker.terminate(); + reject(new Error(message.error || 'Clear failed')); + } } - }); + ); worker.on('error', (error) => { - console.error(this.loggerLabel, 'Worker error during clear:', error); + console.error( + this.loggerLabel, + 'Worker error during clear:', + error + ); worker.terminate(); reject(error); }); diff --git a/libs/services/src/lib/epg.service.ts b/libs/services/src/lib/epg.service.ts index 7b1bda605..0e857645d 100644 --- a/libs/services/src/lib/epg.service.ts +++ b/libs/services/src/lib/epg.service.ts @@ -34,7 +34,6 @@ export class EpgService { */ fetchEpg(urls: string[]): void { if (!this.isDesktop) return; - this.showFetchSnackbar(); // Filter out empty URLs and send all URLs at once const validUrls = urls.filter((url) => url?.trim()); @@ -45,7 +44,6 @@ export class EpgService { tap((result) => { if (result.success) { this.epgAvailable.next(true); - this.showSuccessSnackbar(); } else { this.epgAvailable.next(false); this.showErrorSnackbar(result.message); @@ -91,30 +89,6 @@ export class EpgService { }); } - /** - * Shows fetch in progress snackbar - */ - showFetchSnackbar(): void { - this.snackBar.open(this.translate.instant('EPG.FETCH_EPG'), undefined, { - duration: 2000, - horizontalPosition: 'start', - }); - } - - /** - * Shows success snackbar - */ - private showSuccessSnackbar(): void { - this.snackBar.open( - this.translate.instant('EPG.FETCH_SUCCESS'), - undefined, - { - duration: 2000, - horizontalPosition: 'start', - } - ); - } - /** * Shows error snackbar */ @@ -131,7 +105,9 @@ export class EpgService { * @param channelId Channel ID (tvg-id or channel name) * @returns Observable of current program or null */ - getCurrentProgramForChannel(channelId: string): Observable { + getCurrentProgramForChannel( + channelId: string + ): Observable { if (!this.isDesktop || !channelId) { return of(null); } @@ -140,7 +116,7 @@ export class EpgService { const cached = this.programCache.get(channelId); const now = Date.now(); - if (cached && (now - cached.timestamp) < this.CACHE_TTL) { + if (cached && now - cached.timestamp < this.CACHE_TTL) { return of(cached.program); } @@ -148,7 +124,10 @@ export class EpgService { return from(window.electron.getChannelPrograms(channelId)).pipe( map((programs: EpgProgram[]) => { if (!programs || programs.length === 0) { - this.programCache.set(channelId, { program: null, timestamp: now }); + this.programCache.set(channelId, { + program: null, + timestamp: now, + }); return null; } @@ -160,19 +139,23 @@ export class EpgService { })); // Find current program from transformed programs - const currentProgram = this.findCurrentProgram(transformedPrograms); + const currentProgram = + this.findCurrentProgram(transformedPrograms); // Cache the result this.programCache.set(channelId, { program: currentProgram, - timestamp: now + timestamp: now, }); return currentProgram; }), catchError((err) => { console.error('EPG get current program error:', err); - this.programCache.set(channelId, { program: null, timestamp: now }); + this.programCache.set(channelId, { + program: null, + timestamp: now, + }); return of(null); }) ); @@ -184,11 +167,13 @@ export class EpgService { private findCurrentProgram(programs: EpgProgram[]): EpgProgram | null { const now = new Date(); - return programs.find(program => { - const start = new Date(program.start); - const stop = new Date(program.stop); - return start <= now && now <= stop; - }) || null; + return ( + programs.find((program) => { + const start = new Date(program.start); + const stop = new Date(program.stop); + return start <= now && now <= stop; + }) || null + ); } /** @@ -196,7 +181,9 @@ export class EpgService { * @param channelIds Array of channel IDs * @returns Observable of Map with channelId -> current program */ - getCurrentProgramsForChannels(channelIds: string[]): Observable> { + getCurrentProgramsForChannels( + channelIds: string[] + ): Observable> { if (!this.isDesktop) { return of(new Map()); } @@ -210,9 +197,9 @@ export class EpgService { const now = Date.now(); // Check cache for each channel - channelIds.forEach(channelId => { + channelIds.forEach((channelId) => { const cached = this.programCache.get(channelId); - if (cached && (now - cached.timestamp) < this.CACHE_TTL) { + if (cached && now - cached.timestamp < this.CACHE_TTL) { resultMap.set(channelId, cached.program); } else { channelsToFetch.push(channelId); @@ -225,18 +212,18 @@ export class EpgService { } // Fetch uncached channels with timeout and error handling per request - const fetchObservables = channelsToFetch.map(channelId => + const fetchObservables = channelsToFetch.map((channelId) => this.getCurrentProgramForChannel(channelId).pipe( timeout(5000), // 5 second timeout per request - map(program => ({ channelId, program })), + map((program) => ({ channelId, program })), catchError(() => of({ channelId, program: null })) ) ); // Combine all fetches using forkJoin return forkJoin(fetchObservables).pipe( - map(results => { - results.forEach(result => { + map((results) => { + results.forEach((result) => { resultMap.set(result.channelId, result.program); }); return resultMap;