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 }; -