From bfc80738f4faaa1aaaa750ce1d78d386f02c8337 Mon Sep 17 00:00:00 2001 From: avandeweghe <1706945+avandeweghe@users.noreply.github.com> Date: Sat, 31 Jan 2026 19:31:03 -0500 Subject: [PATCH] Refactor parseStreaming to use pendingBatches array Was running into an issue where epg wouldn't process passed the first batch. Refactored the batching to use an array and added backpressure to keep memory usage low. --- server/services/epgParser.js | 44 ++++++++++++++++++++++++++++-------- 1 file changed, 35 insertions(+), 9 deletions(-) diff --git a/server/services/epgParser.js b/server/services/epgParser.js index 9b75fff..53ec2df 100644 --- a/server/services/epgParser.js +++ b/server/services/epgParser.js @@ -266,9 +266,11 @@ async function* fetchAndParseStreaming(url, batchSize = 1000) { * @yields {{ channels: Array|null, programmes: Array, isLast: boolean }} */ async function* parseStreaming(input, batchSize = 1000) { - const channels = []; + let channels = []; let programmeBatch = []; let channelsYielded = false; + const maxPenmdingBatches = 4; + let paused = false; // We need to convert SAX events to an async iterator // This requires collecting events and yielding when batch is full @@ -279,7 +281,7 @@ async function* parseStreaming(input, batchSize = 1000) { let currentObject = null; let textBuffer = ''; let resolveNext = null; - let pendingBatch = null; + const pendingBatches = []; let ended = false; let error = null; @@ -345,13 +347,27 @@ async function* parseStreaming(input, batchSize = 1000) { isLast: false }; channelsYielded = true; + + // allow channels array to be GC'd after included in first batch + try { if (channels && channels.length) channels = null; } catch (e) {} + programmeBatch = []; if (resolveNext) { resolveNext(batch); resolveNext = null; } else { - pendingBatch = batch; + pendingBatches.push(batch); + + // Apply backpressure when queue grows too large + try { + if (!paused && pendingBatches.length >= maxPenmdingBatches) { + if (input && typeof input.pause === 'function') { + input.pause(); + paused = true; + } + } + } catch (e) {} } } } @@ -398,7 +414,7 @@ async function* parseStreaming(input, batchSize = 1000) { resolveNext(batch); resolveNext = null; } else { - pendingBatch = batch; + pendingBatches.push(batch); } }); @@ -413,10 +429,21 @@ async function* parseStreaming(input, batchSize = 1000) { input.pipe(saxStream); // Yield batches as they become available - while (!ended || pendingBatch) { - if (pendingBatch) { - const batch = pendingBatch; - pendingBatch = null; + + while (!ended || pendingBatches.length > 0) { + if (pendingBatches.length > 0) { + const batch = pendingBatches.shift(); + + // If paused and queue drained below threshold, resume + try { + if (paused && pendingBatches.length < maxPenmdingBatches) { + if (input && typeof input.resume === 'function') { + input.resume(); + paused = false; + } + } + } catch (e) {} + yield batch; if (batch.isLast) break; } else if (!ended) { @@ -445,4 +472,3 @@ module.exports = { getProgrammesForChannel, getCurrentAndUpcoming }; -