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 e9b6cbf4c..fd2fffa62 100644 --- a/apps/electron-backend/src/app/workers/epg-parser.worker.ts +++ b/apps/electron-backend/src/app/workers/epg-parser.worker.ts @@ -1,10 +1,14 @@ import type BetterSqlite3 from 'better-sqlite3'; import { existsSync, mkdirSync } from 'fs'; import { getIptvnatorDatabasePath } from 'database-path-utils'; -import { SaxesParser, SaxesTagPlain } from 'saxes'; import { Readable } from 'stream'; import { parentPort, workerData } from 'worker_threads'; import { createGunzip } from 'zlib'; +import { + ParsedChannel, + ParsedProgram, + StreamingEpgParser, +} from './epg-streaming-parser'; import { shouldGunzipEpgResponse } from './epg-response-utils'; import { getNativeModuleSearchPaths, @@ -39,51 +43,6 @@ function loadBetterSqlite3(): typeof BetterSqlite3 { Database = loadBetterSqlite3(); -/** - * Internal parsing types with arrays for XML parsing - * These are different from the flat EpgProgram interface used by the frontend - */ -interface ParsedTextValue { - lang: string; - value: string; -} - -interface ParsedIcon { - src: string; - width?: number; - height?: number; -} - -interface ParsedRating { - system: string; - value: string; -} - -interface ParsedEpisodeNum { - system: string; - value: string; -} - -interface ParsedChannel { - id: string; - displayName: ParsedTextValue[]; - icon: ParsedIcon[]; - url: string[]; -} - -interface ParsedProgram { - start: string; - stop: string; - channel: string; - title: ParsedTextValue[]; - desc: ParsedTextValue[]; - category: ParsedTextValue[]; - date: string; - episodeNum: ParsedEpisodeNum[]; - icon: ParsedIcon[]; - rating: ParsedRating[]; -} - /** * Streaming EPG Parser Worker * Uses SAX parsing to process XML incrementally without loading entire file into memory. @@ -244,274 +203,6 @@ class EpgDatabase { } } -/** - * Parse XMLTV datetime format to ISO string - * Format: YYYYMMDDHHmmss +HHMM or YYYYMMDDHHmmss - */ -function parseXmltvDate(dateStr: string): string { - if (!dateStr) return ''; - - const match = dateStr.match( - /^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})\s*([+-]\d{4})?$/ - ); - - if (!match) return dateStr; - - const [, year, month, day, hour, minute, second, tz] = match; - - let isoString = `${year}-${month}-${day}T${hour}:${minute}:${second}`; - - if (tz) { - isoString += `${tz.slice(0, 3)}:${tz.slice(3)}`; - } else { - isoString += 'Z'; - } - - return isoString; -} - -/** - * Streaming EPG parser using SAX - */ -class StreamingEpgParser { - private parser: SaxesParser; - private channels: ParsedChannel[] = []; - private programs: ParsedProgram[] = []; - private totalChannels = 0; - private totalPrograms = 0; - - // Current element being parsed - private currentChannel: Partial | null = null; - private currentProgram: Partial | null = null; - private currentTextContent = ''; - private currentLang = ''; - - // For nested elements - private elementStack: string[] = []; - - constructor( - private onChannelsBatch: (channels: ParsedChannel[]) => void, - private onProgramsBatch: (programs: ParsedProgram[]) => void, - private onProgress: (channels: number, programs: number) => void - ) { - this.parser = new SaxesParser(); - this.setupParser(); - } - - private setupParser(): void { - this.parser.on('opentag', (tag: SaxesTagPlain) => { - this.elementStack.push(tag.name); - this.currentTextContent = ''; - - switch (tag.name) { - case 'channel': - this.currentChannel = { - id: (tag.attributes['id'] as string) || '', - displayName: [], - icon: [], - url: [], - }; - break; - - case 'programme': - this.currentProgram = { - start: parseXmltvDate( - (tag.attributes['start'] as string) || '' - ), - stop: parseXmltvDate( - (tag.attributes['stop'] as string) || '' - ), - channel: (tag.attributes['channel'] as string) || '', - title: [], - desc: [], - category: [], - date: '', - episodeNum: [], - icon: [], - rating: [], - }; - break; - - case 'icon': - if (this.currentChannel) { - this.currentChannel.icon!.push({ - src: (tag.attributes['src'] as string) || '', - width: tag.attributes['width'] - ? parseInt(tag.attributes['width'] as string) - : undefined, - height: tag.attributes['height'] - ? parseInt(tag.attributes['height'] as string) - : undefined, - }); - } else if (this.currentProgram) { - this.currentProgram.icon!.push({ - src: (tag.attributes['src'] as string) || '', - width: tag.attributes['width'] - ? parseInt(tag.attributes['width'] as string) - : undefined, - height: tag.attributes['height'] - ? parseInt(tag.attributes['height'] as string) - : undefined, - }); - } - break; - - case 'display-name': - case 'title': - case 'desc': - case 'category': - this.currentLang = (tag.attributes['lang'] as string) || ''; - break; - - case 'rating': - if (this.currentProgram) { - const system = - (tag.attributes['system'] as string) || ''; - this.currentProgram.rating!.push({ system, value: '' }); - } - break; - - case 'episode-num': - if (this.currentProgram) { - const system = - (tag.attributes['system'] as string) || ''; - this.currentProgram.episodeNum!.push({ - system, - value: '', - }); - } - break; - } - }); - - this.parser.on('text', (text: string) => { - this.currentTextContent += text; - }); - - this.parser.on('closetag', (tag: SaxesTagPlain) => { - const text = this.currentTextContent.trim(); - - if (this.currentChannel) { - switch (tag.name) { - case 'display-name': - this.currentChannel.displayName!.push({ - lang: this.currentLang, - value: text, - }); - break; - case 'url': - if (text) this.currentChannel.url!.push(text); - break; - case 'channel': - this.channels.push( - this.currentChannel as ParsedChannel - ); - this.totalChannels++; - this.currentChannel = null; - - if (this.channels.length >= CHANNEL_BATCH_SIZE) { - this.flushChannels(); - } - break; - } - } - - if (this.currentProgram) { - switch (tag.name) { - case 'title': - this.currentProgram.title!.push({ - lang: this.currentLang, - value: text, - }); - break; - case 'desc': - this.currentProgram.desc!.push({ - lang: this.currentLang, - value: text, - }); - break; - case 'category': - this.currentProgram.category!.push({ - lang: this.currentLang, - value: text, - }); - break; - case 'date': - this.currentProgram.date = text; - break; - case 'value': - if ( - this.elementStack.includes('rating') && - this.currentProgram.rating!.length > 0 - ) { - this.currentProgram.rating![ - this.currentProgram.rating!.length - 1 - ].value = text; - } - break; - case 'episode-num': - if (this.currentProgram.episodeNum!.length > 0) { - this.currentProgram.episodeNum![ - this.currentProgram.episodeNum!.length - 1 - ].value = text; - } - break; - case 'programme': - this.programs.push( - this.currentProgram as ParsedProgram - ); - this.totalPrograms++; - - if (this.programs.length >= PROGRAM_BATCH_SIZE) { - this.flushChannels(); - this.flushPrograms(); - } - this.currentProgram = null; - break; - } - } - - this.elementStack.pop(); - this.currentTextContent = ''; - }); - - this.parser.on('error', (err: Error) => { - console.error(loggerLabel, 'Parser error:', err.message); - }); - } - - private flushChannels(): void { - if (this.channels.length > 0) { - this.onChannelsBatch([...this.channels]); - this.channels = []; - this.onProgress(this.totalChannels, this.totalPrograms); - } - } - - private flushPrograms(): void { - if (this.programs.length > 0) { - this.onProgramsBatch([...this.programs]); - this.programs = []; - this.onProgress(this.totalChannels, this.totalPrograms); - } - } - - write(chunk: string): void { - this.parser.write(chunk); - } - - finish(): { totalChannels: number; totalPrograms: number } { - this.parser.close(); - this.flushChannels(); - this.flushPrograms(); - - return { - totalChannels: this.totalChannels, - totalPrograms: this.totalPrograms, - }; - } -} - /** * Fetches and parses EPG data from URL using streaming * Inserts directly into SQLite to avoid blocking main thread @@ -566,7 +257,9 @@ async function fetchAndParseEpgStreaming(url: string): Promise { stats: { totalChannels, totalPrograms }, }; parentPort?.postMessage(response); - } + }, + CHANNEL_BATCH_SIZE, + PROGRAM_BATCH_SIZE ); // Convert web stream to Node.js stream diff --git a/apps/electron-backend/src/app/workers/epg-streaming-parser.spec.ts b/apps/electron-backend/src/app/workers/epg-streaming-parser.spec.ts new file mode 100644 index 000000000..a36a3320f --- /dev/null +++ b/apps/electron-backend/src/app/workers/epg-streaming-parser.spec.ts @@ -0,0 +1,57 @@ +import { + parseXmltvDate, + StreamingEpgParser, + type ParsedChannel, + type ParsedProgram, +} from './epg-streaming-parser'; + +describe('parseXmltvDate', () => { + it('normalizes XMLTV timestamps with timezone offsets', () => { + expect(parseXmltvDate('20260415053700 +0000')).toBe( + '2026-04-15T05:37:00+00:00' + ); + }); +}); + +describe('StreamingEpgParser', () => { + it('flushes pending channels before the first programme batch', () => { + const callbackOrder: string[] = []; + const channelIdsByBatch: string[][] = []; + const programmeChannelsByBatch: string[][] = []; + + const parser = new StreamingEpgParser( + (channels: ParsedChannel[]) => { + callbackOrder.push('channels'); + channelIdsByBatch.push(channels.map((channel) => channel.id)); + }, + (programs: ParsedProgram[]) => { + callbackOrder.push('programs'); + programmeChannelsByBatch.push( + programs.map((program) => program.channel) + ); + }, + () => undefined + ); + + const xml = [ + '', + '', + ...Array.from({ length: 101 }, (_, index) => { + const id = `channel-${index + 1}`; + return `${id}`; + }), + 'Late channel', + '', + ].join(''); + + parser.write(xml); + parser.finish(); + + expect(callbackOrder).toEqual(['channels', 'channels', 'programs']); + expect(channelIdsByBatch).toEqual([ + Array.from({ length: 100 }, (_, index) => `channel-${index + 1}`), + ['channel-101'], + ]); + expect(programmeChannelsByBatch).toEqual([['channel-101']]); + }); +}); diff --git a/apps/electron-backend/src/app/workers/epg-streaming-parser.ts b/apps/electron-backend/src/app/workers/epg-streaming-parser.ts new file mode 100644 index 000000000..7c82e7d00 --- /dev/null +++ b/apps/electron-backend/src/app/workers/epg-streaming-parser.ts @@ -0,0 +1,309 @@ +import { SaxesParser, SaxesTagPlain } from 'saxes'; + +export interface ParsedTextValue { + lang: string; + value: string; +} + +export interface ParsedIcon { + src: string; + width?: number; + height?: number; +} + +export interface ParsedRating { + system: string; + value: string; +} + +export interface ParsedEpisodeNum { + system: string; + value: string; +} + +export interface ParsedChannel { + id: string; + displayName: ParsedTextValue[]; + icon: ParsedIcon[]; + url: string[]; +} + +export interface ParsedProgram { + start: string; + stop: string; + channel: string; + title: ParsedTextValue[]; + desc: ParsedTextValue[]; + category: ParsedTextValue[]; + date: string; + episodeNum: ParsedEpisodeNum[]; + icon: ParsedIcon[]; + rating: ParsedRating[]; +} + +/** + * Parse XMLTV datetime format to ISO string + * Format: YYYYMMDDHHmmss +HHMM or YYYYMMDDHHmmss + */ +export function parseXmltvDate(dateStr: string): string { + if (!dateStr) return ''; + + const match = dateStr.match( + /^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})\s*([+-]\d{4})?$/ + ); + + if (!match) return dateStr; + + const [, year, month, day, hour, minute, second, tz] = match; + + let isoString = `${year}-${month}-${day}T${hour}:${minute}:${second}`; + + if (tz) { + isoString += `${tz.slice(0, 3)}:${tz.slice(3)}`; + } else { + isoString += 'Z'; + } + + return isoString; +} + +/** + * Streaming EPG parser using SAX + */ +export class StreamingEpgParser { + private parser: SaxesParser; + private channels: ParsedChannel[] = []; + private programs: ParsedProgram[] = []; + private totalChannels = 0; + private totalPrograms = 0; + + // Current element being parsed + private currentChannel: Partial | null = null; + private currentProgram: Partial | null = null; + private currentTextContent = ''; + private currentLang = ''; + + // For nested elements + private elementStack: string[] = []; + + constructor( + private readonly onChannelsBatch: (channels: ParsedChannel[]) => void, + private readonly onProgramsBatch: (programs: ParsedProgram[]) => void, + private readonly onProgress: (channels: number, programs: number) => void, + private readonly channelBatchSize = 100, + private readonly programBatchSize = 1000 + ) { + this.parser = new SaxesParser(); + this.setupParser(); + } + + private setupParser(): void { + this.parser.on('opentag', (tag: SaxesTagPlain) => { + this.elementStack.push(tag.name); + this.currentTextContent = ''; + + switch (tag.name) { + case 'channel': + this.currentChannel = { + id: (tag.attributes['id'] as string) || '', + displayName: [], + icon: [], + url: [], + }; + break; + + case 'programme': + // XMLTV feeds typically emit all entries before + // the first . Flush any pending channel batch + // so downstream program consumers can resolve late-channel + // IDs immediately instead of dropping the first rows. + this.flushChannels(); + this.currentProgram = { + start: parseXmltvDate( + (tag.attributes['start'] as string) || '' + ), + stop: parseXmltvDate( + (tag.attributes['stop'] as string) || '' + ), + channel: (tag.attributes['channel'] as string) || '', + title: [], + desc: [], + category: [], + date: '', + episodeNum: [], + icon: [], + rating: [], + }; + break; + + case 'icon': + if (this.currentChannel) { + this.currentChannel.icon!.push({ + src: (tag.attributes['src'] as string) || '', + width: tag.attributes['width'] + ? parseInt(tag.attributes['width'] as string) + : undefined, + height: tag.attributes['height'] + ? parseInt(tag.attributes['height'] as string) + : undefined, + }); + } else if (this.currentProgram) { + this.currentProgram.icon!.push({ + src: (tag.attributes['src'] as string) || '', + width: tag.attributes['width'] + ? parseInt(tag.attributes['width'] as string) + : undefined, + height: tag.attributes['height'] + ? parseInt(tag.attributes['height'] as string) + : undefined, + }); + } + break; + + case 'display-name': + case 'title': + case 'desc': + case 'category': + this.currentLang = (tag.attributes['lang'] as string) || ''; + break; + + case 'rating': + if (this.currentProgram) { + const system = + (tag.attributes['system'] as string) || ''; + this.currentProgram.rating!.push({ system, value: '' }); + } + break; + + case 'episode-num': + if (this.currentProgram) { + const system = + (tag.attributes['system'] as string) || ''; + this.currentProgram.episodeNum!.push({ + system, + value: '', + }); + } + break; + } + }); + + this.parser.on('text', (text: string) => { + this.currentTextContent += text; + }); + + this.parser.on('closetag', (tag: SaxesTagPlain) => { + const text = this.currentTextContent.trim(); + + if (this.currentChannel) { + switch (tag.name) { + case 'display-name': + this.currentChannel.displayName!.push({ + lang: this.currentLang, + value: text, + }); + break; + case 'url': + if (text) this.currentChannel.url!.push(text); + break; + case 'channel': + this.channels.push(this.currentChannel as ParsedChannel); + this.totalChannels++; + this.currentChannel = null; + + if (this.channels.length >= this.channelBatchSize) { + this.flushChannels(); + } + break; + } + } + + if (this.currentProgram) { + switch (tag.name) { + case 'title': + this.currentProgram.title!.push({ + lang: this.currentLang, + value: text, + }); + break; + case 'desc': + this.currentProgram.desc!.push({ + lang: this.currentLang, + value: text, + }); + break; + case 'category': + this.currentProgram.category!.push({ + lang: this.currentLang, + value: text, + }); + break; + case 'date': + this.currentProgram.date = text; + break; + case 'value': + if ( + this.elementStack.includes('rating') && + this.currentProgram.rating!.length > 0 + ) { + this.currentProgram.rating![ + this.currentProgram.rating!.length - 1 + ].value = text; + } + break; + case 'episode-num': + if (this.currentProgram.episodeNum!.length > 0) { + this.currentProgram.episodeNum![ + this.currentProgram.episodeNum!.length - 1 + ].value = text; + } + break; + case 'programme': + this.programs.push(this.currentProgram as ParsedProgram); + this.totalPrograms++; + + if (this.programs.length >= this.programBatchSize) { + this.flushChannels(); + this.flushPrograms(); + } + this.currentProgram = null; + break; + } + } + + this.elementStack.pop(); + this.currentTextContent = ''; + }); + } + + private flushChannels(): void { + if (this.channels.length > 0) { + this.onChannelsBatch([...this.channels]); + this.channels = []; + this.onProgress(this.totalChannels, this.totalPrograms); + } + } + + private flushPrograms(): void { + if (this.programs.length > 0) { + this.onProgramsBatch([...this.programs]); + this.programs = []; + this.onProgress(this.totalChannels, this.totalPrograms); + } + } + + write(chunk: string): void { + this.parser.write(chunk); + } + + finish(): { totalChannels: number; totalPrograms: number } { + this.parser.close(); + this.flushChannels(); + this.flushPrograms(); + + return { + totalChannels: this.totalChannels, + totalPrograms: this.totalPrograms, + }; + } +}