From e8e6ba27d7591af4330b6839331befebc4b2732e Mon Sep 17 00:00:00 2001 From: Trevor Mears Date: Wed, 7 Jan 2026 23:09:23 -0800 Subject: [PATCH] Implemented streaming parsers for M3U and EPG to reduce memory usage --- server/services/epgParser.js | 215 +++++++++++++++++++++++++++++++++ server/services/m3uParser.js | 99 ++++++++++++++- server/services/syncService.js | 205 +++++++++++++++++++------------ 3 files changed, 443 insertions(+), 76 deletions(-) diff --git a/server/services/epgParser.js b/server/services/epgParser.js index c0cde58..181b022 100644 --- a/server/services/epgParser.js +++ b/server/services/epgParser.js @@ -224,10 +224,225 @@ async function fetchAndParse(url) { return parse(stream); } +/** + * Streaming EPG parser that yields batches of programmes (memory-efficient) + * Channels are collected and returned with the first batch, then programmes are yielded in batches. + * + * @param {string} url - XMLTV URL + * @param {number} batchSize - Number of programmes per batch (default: 1000) + * @yields {{ channels: Array|null, programmes: Array, isLast: boolean }} + */ +async function* fetchAndParseStreaming(url, batchSize = 1000) { + const response = await fetch(url); + if (!response.ok) { + throw new Error(`Failed to fetch EPG: ${response.status} ${response.statusText}`); + } + + let stream; + if (response.body && typeof response.body.pipe === 'function') { + stream = response.body; + } else if (response.body) { + stream = Readable.fromWeb(response.body); + } else { + stream = Readable.from([]); + } + + const isGzipped = url.endsWith('.gz') || (response.headers.get('content-type') || '').includes('gzip'); + + if (isGzipped) { + const gunzip = zlib.createGunzip(); + stream.pipe(gunzip); + stream = gunzip; + } + + // Use async iterator pattern with SAX + yield* parseStreaming(stream, batchSize); +} + +/** + * Parse XMLTV as streaming async generator + * @param {Readable} input - XMLTV stream + * @param {number} batchSize - Number of programmes per batch + * @yields {{ channels: Array|null, programmes: Array, isLast: boolean }} + */ +async function* parseStreaming(input, batchSize = 1000) { + const channels = []; + let programmeBatch = []; + let channelsYielded = false; + + // We need to convert SAX events to an async iterator + // This requires collecting events and yielding when batch is full + + const saxStream = sax.createStream(true, { trim: true, normalize: true }); + + let currentTag = null; + let currentObject = null; + let textBuffer = ''; + let resolveNext = null; + let pendingBatch = null; + let ended = false; + let error = null; + + saxStream.on('error', function (e) { + this._parser.error = null; + this._parser.resume(); + console.warn('XML Parse Warning:', e.message); + }); + + saxStream.on('opentag', function (node) { + currentTag = node.name; + const attr = node.attributes; + + if (currentTag === 'channel') { + currentObject = { + id: attr.id, + name: null, + icon: null, + url: null + }; + } else if (currentTag === 'programme') { + currentObject = { + channelId: attr.channel, + start: parseXmltvDate(attr.start), + stop: parseXmltvDate(attr.stop), + title: null, + subtitle: null, + description: null, + category: [], + icon: null, + date: null, + episodeNum: null + }; + } else if (currentTag === 'icon') { + if (currentObject) { + currentObject.icon = attr.src; + } + } + textBuffer = ''; + }); + + saxStream.on('text', function (text) { + textBuffer += text; + }); + + saxStream.on('cdata', function (text) { + textBuffer += text; + }); + + saxStream.on('closetag', function (tagName) { + if (tagName === 'channel') { + if (currentObject) channels.push(currentObject); + currentObject = null; + } else if (tagName === 'programme') { + if (currentObject) { + programmeBatch.push(currentObject); + + // Check if we should yield a batch + if (programmeBatch.length >= batchSize) { + const batch = { + channels: !channelsYielded ? channels : null, + programmes: programmeBatch, + isLast: false + }; + channelsYielded = true; + programmeBatch = []; + + if (resolveNext) { + resolveNext(batch); + resolveNext = null; + } else { + pendingBatch = batch; + } + } + } + currentObject = null; + } else if (currentObject) { + switch (tagName) { + case 'display-name': + if (!currentObject.name) currentObject.name = textBuffer; + break; + case 'url': + currentObject.url = textBuffer; + break; + case 'title': + currentObject.title = textBuffer; + break; + case 'sub-title': + currentObject.subtitle = textBuffer; + break; + case 'desc': + currentObject.description = textBuffer; + break; + case 'category': + if (textBuffer) currentObject.category.push(textBuffer); + break; + case 'date': + currentObject.date = textBuffer; + break; + case 'episode-num': + currentObject.episodeNum = textBuffer; + break; + } + } + }); + + saxStream.on('end', function () { + ended = true; + // Yield final batch + const batch = { + channels: !channelsYielded ? channels : null, + programmes: programmeBatch, + isLast: true + }; + if (resolveNext) { + resolveNext(batch); + resolveNext = null; + } else { + pendingBatch = batch; + } + }); + + saxStream.on('error', function (e) { + error = e; + if (resolveNext) { + resolveNext(null); + } + }); + + // Start piping + input.pipe(saxStream); + + // Yield batches as they become available + while (!ended || pendingBatch) { + if (pendingBatch) { + const batch = pendingBatch; + pendingBatch = null; + yield batch; + if (batch.isLast) break; + } else if (!ended) { + // Wait for next batch + const batch = await new Promise(resolve => { + resolveNext = resolve; + }); + if (batch) { + yield batch; + if (batch.isLast) break; + } + } + } + + if (error) { + throw error; + } +} + module.exports = { parse, parseXmltvDate, fetchAndParse, + fetchAndParseStreaming, + parseStreaming, getProgrammesForChannel, getCurrentAndUpcoming }; + diff --git a/server/services/m3uParser.js b/server/services/m3uParser.js index ac11026..433004a 100644 --- a/server/services/m3uParser.js +++ b/server/services/m3uParser.js @@ -175,4 +175,101 @@ async function fetchAndParse(url) { return parse(stream); } -module.exports = { parse, parseExtinf, fetchAndParse }; +/** + * Parse M3U content as a streaming async generator (memory-efficient) + * Yields batches of channels to avoid loading entire playlist into memory. + * + * @param {Readable|string} input - M3U content as Stream or String + * @param {number} batchSize - Number of channels per batch (default: 500) + * @yields {{ channels: Array, groups: Set, isLast: boolean }} + */ +async function* parseStreaming(input, batchSize = 500) { + const groupsSet = new Set(); + let currentInfo = null; + let currentGroup = null; + let batch = []; + + let lines; + if (typeof input === 'string') { + lines = input.split(/\r?\n/); + } else { + const rl = readline.createInterface({ + input: input, + crlfDelay: Infinity + }); + lines = rl; + } + + for await (const line of lines) { + const trimmed = line.trim(); + if (!trimmed) continue; + + if (trimmed.startsWith('#EXTINF:')) { + currentInfo = parseExtinf(trimmed); + if (currentInfo.groupTitle) { + groupsSet.add(currentInfo.groupTitle); + currentGroup = currentInfo.groupTitle; + } + } else if (trimmed.startsWith('#EXTGRP:')) { + currentGroup = trimmed.substring(8).trim(); + groupsSet.add(currentGroup); + if (currentInfo) { + currentInfo.groupTitle = currentGroup; + } + } else if (!trimmed.startsWith('#')) { + if (currentInfo) { + const groupTitle = currentInfo.groupTitle || currentGroup || 'Uncategorized'; + const stableId = currentInfo.tvgId || generateStableId(currentInfo.name, groupTitle); + + batch.push({ + ...currentInfo, + id: stableId, + url: trimmed, + groupTitle: groupTitle + }); + currentInfo = null; + + // Yield batch when full + if (batch.length >= batchSize) { + yield { channels: batch, groups: groupsSet, isLast: false }; + batch = []; + } + } + } + } + + // Yield remaining channels + if (batch.length > 0) { + yield { channels: batch, groups: groupsSet, isLast: true }; + } else { + // Yield empty final batch with isLast=true so caller knows we're done + yield { channels: [], groups: groupsSet, isLast: true }; + } +} + +/** + * Fetch and parse M3U from URL as streaming async generator (memory-efficient) + * @param {string} url - M3U playlist URL + * @param {number} batchSize - Number of channels per batch + * @yields {{ channels: Array, groups: Set, isLast: boolean }} + */ +async function* fetchAndParseStreaming(url, batchSize = 500) { + const response = await fetch(url); + if (!response.ok) { + throw new Error(`Failed to fetch M3U: ${response.status} ${response.statusText}`); + } + + let stream; + if (response.body && typeof response.body.pipe === 'function') { + stream = response.body; + } else if (response.body) { + stream = Readable.fromWeb(response.body); + } else { + stream = Readable.from([]); + } + + yield* parseStreaming(stream, batchSize); +} + +module.exports = { parse, parseExtinf, fetchAndParse, parseStreaming, fetchAndParseStreaming }; + diff --git a/server/services/syncService.js b/server/services/syncService.js index bf84df4..d95dac9 100644 --- a/server/services/syncService.js +++ b/server/services/syncService.js @@ -334,63 +334,36 @@ class SyncService { /** - * Sync EPG from URL + * Sync EPG from URL (Streaming - Memory Efficient) + * Processes EPG files in batches to avoid OOM on large EPG data */ async syncEpgFromUrl(sourceId, url) { - // Use our streaming parser - const { channels, programmes } = await epgParser.fetchAndParse(url); + console.log(`[Sync] Fetching EPG from: ${url.substring(0, 60)}...`); - console.log(`[Sync] EPG Parsed: ${channels.length} channels, ${programmes.length} programs`); + // Temporary memory logging for verification + const logMemory = () => { + const used = process.memoryUsage(); + console.log(`[Sync] Memory: ${Math.round(used.heapUsed / 1024 / 1024)}MB heap`); + }; + + logMemory(); const db = getDb(); + let allChannels = []; + let totalProgrammes = 0; + let batchCount = 0; - // 1. Save EPG Channels to playlist_items (for Name/Icon matching) - - const channelStmt = db.prepare(` - INSERT INTO playlist_items ( - id, source_id, item_id, type, name, stream_icon, - stream_url, category_id, data - ) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) - ON CONFLICT(id) DO UPDATE SET - name = excluded.name, - stream_icon = excluded.stream_icon, - data = excluded.data - `); - - // Use transaction for channels - const insertChannels = db.transaction((chanList) => { - for (const ch of chanList) { - const id = `${sourceId}:${ch.id}`; - channelStmt.run( - id, - sourceId, - ch.id, // XMLTV ID - 'epg_channel', - ch.name, - ch.icon || null, - null, // No URL - null, // No Category - JSON.stringify(ch) - ); - } - }); - - insertChannels(channels); - console.log(`[Sync] Saved ${channels.length} EPG channels`); - - // 2. Save Programs - // First delete old programs for this source + // Clear old programmes first db.prepare('DELETE FROM epg_programs WHERE source_id = ?').run(sourceId); - const stmt = db.prepare(` + const programmeStmt = db.prepare(` INSERT INTO epg_programs (channel_id, source_id, start_time, end_time, title, description, data) VALUES (?, ?, ?, ?, ?, ?, ?) `); - const insertMany = db.transaction((progs) => { + const insertProgrammes = db.transaction((progs) => { for (const p of progs) { - stmt.run( + programmeStmt.run( p.channelId, sourceId, p.start ? p.start.getTime() : 0, @@ -402,50 +375,132 @@ class SyncService { } }); - insertMany(programmes); - console.log(`[Sync] Saved ${programmes.length} programs`); + // Stream and process in batches (default 1000 programmes per batch) + for await (const batch of epgParser.fetchAndParseStreaming(url)) { + batchCount++; + + // Collect channels from first batch + if (batch.channels) { + allChannels = batch.channels; + } + + // Save this batch of programmes immediately + if (batch.programmes.length > 0) { + insertProgrammes(batch.programmes); + totalProgrammes += batch.programmes.length; + } + + // Log progress every 10 batches + if (batchCount % 10 === 0) { + console.log(`[Sync] Processed ${totalProgrammes} programmes so far...`); + logMemory(); + } + + // Yield to event loop + await new Promise(resolve => setImmediate(resolve)); + } + + console.log(`[Sync] EPG Parsed: ${allChannels.length} channels, ${totalProgrammes} programmes`); + logMemory(); + + // Save EPG Channels + if (allChannels.length > 0) { + const channelStmt = db.prepare(` + INSERT INTO playlist_items ( + id, source_id, item_id, type, name, stream_icon, + stream_url, category_id, data + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(id) DO UPDATE SET + name = excluded.name, + stream_icon = excluded.stream_icon, + data = excluded.data + `); + + const insertChannels = db.transaction((chanList) => { + for (const ch of chanList) { + const id = `${sourceId}:${ch.id}`; + channelStmt.run( + id, + sourceId, + ch.id, + 'epg_channel', + ch.name, + ch.icon || null, + null, + null, + JSON.stringify(ch) + ); + } + }); + + insertChannels(allChannels); + console.log(`[Sync] Saved ${allChannels.length} EPG channels`); + } + + console.log(`[Sync] Saved ${totalProgrammes} programmes`); } /** - * M3U Sync Logic - */ - /** - * M3U Sync Logic + * M3U Sync Logic (Streaming - Memory Efficient) + * Processes M3U files in batches to avoid OOM on large playlists */ async syncM3u(source) { console.log(`[Sync] Fetching M3U playlist for ${source.name}`); - // Use the streaming parser directly to avoid loading entire file into memory - // This prevents OOM crashes on large playlists (100MB+) - const { channels, groups } = await m3uParser.fetchAndParse(source.url); + // Temporary memory logging for verification + const logMemory = () => { + const used = process.memoryUsage(); + console.log(`[Sync] Memory: ${Math.round(used.heapUsed / 1024 / 1024)}MB heap`); + }; - console.log(`[Sync] M3U Parsed: ${channels.length} channels, ${groups.length} groups`); + logMemory(); - // Save Categories (Groups) - // M3U groups are just strings usually, we need to normalize them - const categories = groups.map(g => ({ - category_id: g.name, // use name as ID for M3U groups - category_name: g.name, + const allGroups = new Set(); + let totalChannels = 0; + let batchCount = 0; + + // Stream and process in batches (default 500 channels per batch) + for await (const batch of m3uParser.fetchAndParseStreaming(source.url)) { + batchCount++; + + // Map M3U channel format to our schema + const playlistItems = batch.channels.map(ch => ({ + stream_id: ch.id, + name: ch.name, + category_id: ch.groupTitle || 'Uncategorized', + stream_icon: ch.tvgLogo, + stream_url: ch.url, + })); + + // Save this batch immediately + if (playlistItems.length > 0) { + await this.saveStreams(source.id, 'live', playlistItems); + totalChannels += playlistItems.length; + } + + // Collect groups for category creation at the end + batch.groups.forEach(g => allGroups.add(g)); + + // Log progress every 10 batches + if (batchCount % 10 === 0) { + console.log(`[Sync] Processed ${totalChannels} channels so far...`); + logMemory(); + } + } + + console.log(`[Sync] M3U Parsed: ${totalChannels} channels, ${allGroups.size} groups`); + logMemory(); + + // Save Categories (Groups) at the end + const categories = Array.from(allGroups).map(name => ({ + category_id: name, + category_name: name, parent_id: null })); await this.saveCategories(source.id, 'live', categories); - - // Save Channels - // Map M3U channel format to our schema - const playlistItems = channels.map(ch => ({ - stream_id: ch.id, // parser generates a stable-ish ID - name: ch.name, - category_id: ch.groupTitle || 'Uncategorized', - stream_icon: ch.tvgLogo, - stream_url: ch.url, - // M3U doesn't usually have VOD metadata like rating/year easily accessible unless extended tags used - // We assume 'live' for now, but could detect VOD from URL extension? - // For now, treat all as type='live' for M3U or maybe check info? - // The parser doesn't differentiate types well yet. - })); - - await this.saveStreams(source.id, 'live', playlistItems); + console.log(`[Sync] M3U sync complete for ${source.name}`); } /**