Implemented streaming parsers for M3U and EPG to reduce memory usage

This commit is contained in:
Trevor Mears
2026-01-07 23:09:23 -08:00
parent 67d9cb26d2
commit e8e6ba27d7
3 changed files with 443 additions and 76 deletions
+215
View File
@@ -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
};