From d63c9c37d7604eca0fd162f5d81e209fe981235a Mon Sep 17 00:00:00 2001 From: Trevor Mears Date: Tue, 30 Dec 2025 00:34:55 -0800 Subject: [PATCH] Refactored server/db.js and all routes to use asynchronous fs operations, refactored m3uParser and epgParser to use streaming (via readline and sax) to prevent errors with large files, updated server/routes/proxy.js to use native stream piping for images and video, reducing latency and memory overhead --- package-lock.json | 1 + package.json | 3 +- server/db.js | 134 ++++++++++--------- server/routes/channels.js | 25 ++-- server/routes/favorites.js | 16 +-- server/routes/proxy.js | 68 ++++------ server/routes/sources.js | 34 ++--- server/services/epgParser.js | 241 ++++++++++++++++++++--------------- server/services/m3uParser.js | 155 ++++++++++++---------- 9 files changed, 366 insertions(+), 311 deletions(-) diff --git a/package-lock.json b/package-lock.json index 2d5ed73..c62c745 100644 --- a/package-lock.json +++ b/package-lock.json @@ -10,6 +10,7 @@ "license": "GPL-3.0-only", "dependencies": { "express": "^4.18.2", + "sax": "^1.4.3", "xml2js": "^0.6.2" }, "engines": { diff --git a/package.json b/package.json index 0f23544..799d410 100644 --- a/package.json +++ b/package.json @@ -10,6 +10,7 @@ }, "dependencies": { "express": "^4.18.2", + "sax": "^1.4.3", "xml2js": "^0.6.2" }, "optionalDependencies": { @@ -18,4 +19,4 @@ "engines": { "node": ">=18.0.0" } -} \ No newline at end of file +} diff --git a/server/db.js b/server/db.js index 03cfd11..7a50ce6 100644 --- a/server/db.js +++ b/server/db.js @@ -1,61 +1,75 @@ -const fs = require('fs'); +const fs = require('fs/promises'); const path = require('path'); +const { existsSync, mkdirSync } = require('fs'); -// Ensure data directory exists +// Ensure data directory exists (sync is fine for startup) const dataDir = path.join(__dirname, '..', 'data'); -if (!fs.existsSync(dataDir)) { - fs.mkdirSync(dataDir, { recursive: true }); +if (!existsSync(dataDir)) { + mkdirSync(dataDir, { recursive: true }); } const dbPath = path.join(dataDir, 'db.json'); // Initialize database structure -function loadDb() { +async function loadDb() { try { - if (fs.existsSync(dbPath)) { - const data = JSON.parse(fs.readFileSync(dbPath, 'utf-8')); - // Ensure all keys exist + // Check if file exists (using fs.access is better for async, but we can catch ENOENT) + try { + const fileContent = await fs.readFile(dbPath, 'utf-8'); + const data = JSON.parse(fileContent); return { sources: data.sources || [], hiddenItems: data.hiddenItems || [], favorites: data.favorites || [], nextId: data.nextId || 1 }; + } catch (error) { + if (error.code === 'ENOENT') { + // File doesn't exist, return default + return { + sources: [], + hiddenItems: [], + favorites: [], + nextId: 1 + }; + } + throw error; } } catch (err) { console.error('Error loading database:', err); + // Return safe default on error to prevent crashing, but log it + return { + sources: [], + hiddenItems: [], + favorites: [], + nextId: 1 + }; } - return { - sources: [], - hiddenItems: [], - favorites: [], - nextId: 1 - }; } -function saveDb(data) { - fs.writeFileSync(dbPath, JSON.stringify(data, null, 2)); +async function saveDb(data) { + await fs.writeFile(dbPath, JSON.stringify(data, null, 2)); } // Source CRUD operations const sources = { - getAll() { - const db = loadDb(); + async getAll() { + const db = await loadDb(); return db.sources; }, - getById(id) { - const db = loadDb(); + async getById(id) { + const db = await loadDb(); return db.sources.find(s => s.id === parseInt(id)); }, - getByType(type) { - const db = loadDb(); + async getByType(type) { + const db = await loadDb(); return db.sources.filter(s => s.type === type && s.enabled); }, - create(source) { - const db = loadDb(); + async create(source) { + const db = await loadDb(); const newSource = { id: db.nextId++, ...source, @@ -64,12 +78,12 @@ const sources = { updated_at: new Date().toISOString() }; db.sources.push(newSource); - saveDb(db); + await saveDb(db); return newSource; }, - update(id, updates) { - const db = loadDb(); + async update(id, updates) { + const db = await loadDb(); const index = db.sources.findIndex(s => s.id === parseInt(id)); if (index === -1) return null; @@ -78,25 +92,25 @@ const sources = { ...updates, updated_at: new Date().toISOString() }; - saveDb(db); + await saveDb(db); return db.sources[index]; }, - delete(id) { - const db = loadDb(); + async delete(id) { + const db = await loadDb(); db.sources = db.sources.filter(s => s.id !== parseInt(id)); // Also delete related hidden items db.hiddenItems = db.hiddenItems.filter(h => h.source_id !== parseInt(id)); - saveDb(db); + await saveDb(db); }, - toggleEnabled(id) { - const db = loadDb(); + async toggleEnabled(id) { + const db = await loadDb(); const source = db.sources.find(s => s.id === parseInt(id)); if (source) { source.enabled = !source.enabled; source.updated_at = new Date().toISOString(); - saveDb(db); + await saveDb(db); } return source; } @@ -104,16 +118,16 @@ const sources = { // Hidden items operations const hiddenItems = { - getAll(sourceId = null) { - const db = loadDb(); + async getAll(sourceId = null) { + const db = await loadDb(); if (sourceId) { return db.hiddenItems.filter(h => h.source_id === parseInt(sourceId)); } return db.hiddenItems; }, - hide(sourceId, itemType, itemId) { - const db = loadDb(); + async hide(sourceId, itemType, itemId) { + const db = await loadDb(); // Check if already hidden const exists = db.hiddenItems.find( h => h.source_id === parseInt(sourceId) && h.item_type === itemType && h.item_id === itemId @@ -125,27 +139,27 @@ const hiddenItems = { item_type: itemType, item_id: itemId }); - saveDb(db); + await saveDb(db); } }, - show(sourceId, itemType, itemId) { - const db = loadDb(); + async show(sourceId, itemType, itemId) { + const db = await loadDb(); db.hiddenItems = db.hiddenItems.filter( h => !(h.source_id === parseInt(sourceId) && h.item_type === itemType && h.item_id === itemId) ); - saveDb(db); + await saveDb(db); }, - isHidden(sourceId, itemType, itemId) { - const db = loadDb(); + async isHidden(sourceId, itemType, itemId) { + const db = await loadDb(); return db.hiddenItems.some( h => h.source_id === parseInt(sourceId) && h.item_type === itemType && h.item_id === itemId ); }, - bulkHide(items) { - const db = loadDb(); + async bulkHide(items) { + const db = await loadDb(); let modified = false; items.forEach(item => { @@ -166,13 +180,13 @@ const hiddenItems = { }); if (modified) { - saveDb(db); + await saveDb(db); } return true; }, - bulkShow(items) { - const db = loadDb(); + async bulkShow(items) { + const db = await loadDb(); const initialLength = db.hiddenItems.length; // Create a set of "signatures" for O(1) lookup of items to remove @@ -183,7 +197,7 @@ const hiddenItems = { ); if (db.hiddenItems.length !== initialLength) { - saveDb(db); + await saveDb(db); } return true; } @@ -191,8 +205,8 @@ const hiddenItems = { // Favorites operations const favorites = { - getAll(sourceId = null, itemType = null) { - const db = loadDb(); + async getAll(sourceId = null, itemType = null) { + const db = await loadDb(); let results = db.favorites; if (sourceId) { results = results.filter(f => f.source_id === parseInt(sourceId)); @@ -203,8 +217,8 @@ const favorites = { return results; }, - add(sourceId, itemId, itemType = 'channel') { - const db = loadDb(); + async add(sourceId, itemId, itemType = 'channel') { + const db = await loadDb(); // Check if already favorited const exists = db.favorites.find( f => f.source_id === parseInt(sourceId) && f.item_id === String(itemId) && f.item_type === itemType @@ -217,22 +231,22 @@ const favorites = { item_type: itemType, // 'channel', 'movie', 'series' created_at: new Date().toISOString() }); - saveDb(db); + await saveDb(db); } return true; }, - remove(sourceId, itemId, itemType = 'channel') { - const db = loadDb(); + async remove(sourceId, itemId, itemType = 'channel') { + const db = await loadDb(); db.favorites = db.favorites.filter( f => !(f.source_id === parseInt(sourceId) && f.item_id === String(itemId) && f.item_type === itemType) ); - saveDb(db); + await saveDb(db); return true; }, - isFavorite(sourceId, itemId, itemType = 'channel') { - const db = loadDb(); + async isFavorite(sourceId, itemId, itemType = 'channel') { + const db = await loadDb(); return db.favorites.some( f => f.source_id === parseInt(sourceId) && f.item_id === String(itemId) && f.item_type === itemType ); diff --git a/server/routes/channels.js b/server/routes/channels.js index 488915f..ebbda1c 100644 --- a/server/routes/channels.js +++ b/server/routes/channels.js @@ -3,10 +3,10 @@ const router = express.Router(); const { hiddenItems } = require('../db'); // Get all hidden items -router.get('/hidden', (req, res) => { +router.get('/hidden', async (req, res) => { try { const { sourceId } = req.query; - const items = hiddenItems.getAll(sourceId ? parseInt(sourceId) : null); + const items = await hiddenItems.getAll(sourceId ? parseInt(sourceId) : null); res.json(items); } catch (err) { console.error('Error getting hidden items:', err); @@ -15,7 +15,7 @@ router.get('/hidden', (req, res) => { }); // Hide a channel, group, or category -router.post('/hide', (req, res) => { +router.post('/hide', async (req, res) => { try { const { sourceId, itemType, itemId } = req.body; @@ -28,7 +28,7 @@ router.post('/hide', (req, res) => { return res.status(400).json({ error: `itemType must be one of: ${validTypes.join(', ')}` }); } - hiddenItems.hide(sourceId, itemType, itemId); + await hiddenItems.hide(sourceId, itemType, itemId); res.json({ success: true }); } catch (err) { console.error('Error hiding item:', err); @@ -37,7 +37,7 @@ router.post('/hide', (req, res) => { }); // Show (unhide) a channel or group -router.post('/show', (req, res) => { +router.post('/show', async (req, res) => { try { const { sourceId, itemType, itemId } = req.body; @@ -45,7 +45,7 @@ router.post('/show', (req, res) => { return res.status(400).json({ error: 'sourceId, itemType, and itemId are required' }); } - hiddenItems.show(sourceId, itemType, itemId); + await hiddenItems.show(sourceId, itemType, itemId); res.json({ success: true }); } catch (err) { console.error('Error showing item:', err); @@ -54,7 +54,7 @@ router.post('/show', (req, res) => { }); // Check if item is hidden -router.get('/hidden/check', (req, res) => { +router.get('/hidden/check', async (req, res) => { try { const { sourceId, itemType, itemId } = req.query; @@ -62,7 +62,8 @@ router.get('/hidden/check', (req, res) => { return res.status(400).json({ error: 'sourceId, itemType, and itemId are required' }); } - const isHidden = hiddenItems.isHidden(parseInt(sourceId), itemType, itemId); + // isHidden is now async + const isHidden = await hiddenItems.isHidden(parseInt(sourceId), itemType, itemId); res.json({ hidden: isHidden }); } catch (err) { console.error('Error checking hidden status:', err); @@ -71,7 +72,7 @@ router.get('/hidden/check', (req, res) => { }); // Bulk hide channels and groups -router.post('/hide/bulk', (req, res) => { +router.post('/hide/bulk', async (req, res) => { try { const { items } = req.body; @@ -79,7 +80,7 @@ router.post('/hide/bulk', (req, res) => { return res.status(400).json({ error: 'items array is required' }); } - hiddenItems.bulkHide(items); + await hiddenItems.bulkHide(items); res.json({ success: true, count: items.length }); } catch (err) { console.error('Error bulk hiding items:', err); @@ -88,7 +89,7 @@ router.post('/hide/bulk', (req, res) => { }); // Bulk show channels and groups -router.post('/show/bulk', (req, res) => { +router.post('/show/bulk', async (req, res) => { try { const { items } = req.body; @@ -96,7 +97,7 @@ router.post('/show/bulk', (req, res) => { return res.status(400).json({ error: 'items array is required' }); } - hiddenItems.bulkShow(items); + await hiddenItems.bulkShow(items); res.json({ success: true, count: items.length }); } catch (err) { console.error('Error bulk showing items:', err); diff --git a/server/routes/favorites.js b/server/routes/favorites.js index 34b415b..8533d68 100644 --- a/server/routes/favorites.js +++ b/server/routes/favorites.js @@ -3,10 +3,10 @@ const router = express.Router(); const { favorites } = require('../db'); // Get all favorites -router.get('/', (req, res) => { +router.get('/', async (req, res) => { try { const { sourceId, itemType } = req.query; - const items = favorites.getAll(sourceId, itemType); + const items = await favorites.getAll(sourceId, itemType); res.json(items); } catch (err) { res.status(500).json({ error: err.message }); @@ -14,14 +14,14 @@ router.get('/', (req, res) => { }); // Add favorite -router.post('/', (req, res) => { +router.post('/', async (req, res) => { try { const { sourceId, itemId, itemType = 'channel' } = req.body; if (!sourceId || !itemId) { return res.status(400).json({ error: 'Source ID and Item ID are required' }); } - favorites.add(sourceId, itemId, itemType); + await favorites.add(sourceId, itemId, itemType); res.json({ success: true }); } catch (err) { res.status(500).json({ error: err.message }); @@ -29,14 +29,14 @@ router.post('/', (req, res) => { }); // Remove favorite -router.delete('/', (req, res) => { +router.delete('/', async (req, res) => { try { const { sourceId, itemId, itemType = 'channel' } = req.body; if (!sourceId || !itemId) { return res.status(400).json({ error: 'Source ID and Item ID are required' }); } - favorites.remove(sourceId, itemId, itemType); + await favorites.remove(sourceId, itemId, itemType); res.json({ success: true }); } catch (err) { res.status(500).json({ error: err.message }); @@ -44,14 +44,14 @@ router.delete('/', (req, res) => { }); // Check if item is favorited -router.get('/check', (req, res) => { +router.get('/check', async (req, res) => { try { const { sourceId, itemId, itemType = 'channel' } = req.query; if (!sourceId || !itemId) { return res.status(400).json({ error: 'Source ID and Item ID are required' }); } - const isFav = favorites.isFavorite(sourceId, itemId, itemType); + const isFav = await favorites.isFavorite(sourceId, itemId, itemType); res.json({ isFavorite: isFav }); } catch (err) { res.status(500).json({ error: err.message }); diff --git a/server/routes/proxy.js b/server/routes/proxy.js index d3e269f..1008c17 100644 --- a/server/routes/proxy.js +++ b/server/routes/proxy.js @@ -5,6 +5,7 @@ const xtreamApi = require('../services/xtreamApi'); const m3uParser = require('../services/m3uParser'); const epgParser = require('../services/epgParser'); const cache = require('../services/cache'); +const { Readable } = require('stream'); // Default cache TTL: 24 hours const DEFAULT_MAX_AGE_HOURS = 24; @@ -16,7 +17,7 @@ const DEFAULT_MAX_AGE_HOURS = 24; router.get('/xtream/:sourceId/:action', async (req, res) => { try { const sourceId = req.params.sourceId; - const source = sources.getById(sourceId); + const source = await sources.getById(sourceId); if (!source || source.type !== 'xtream') { return res.status(404).json({ error: 'Xtream source not found' }); } @@ -99,9 +100,9 @@ router.get('/xtream/:sourceId/:action', async (req, res) => { * Get Xtream stream URL * GET /api/proxy/xtream/:sourceId/stream/:streamId */ -router.get('/xtream/:sourceId/stream/:streamId/:type?', (req, res) => { +router.get('/xtream/:sourceId/stream/:streamId/:type?', async (req, res) => { try { - const source = sources.getById(req.params.sourceId); + const source = await sources.getById(req.params.sourceId); if (!source || source.type !== 'xtream') { return res.status(404).json({ error: 'Xtream source not found' }); } @@ -125,7 +126,7 @@ router.get('/xtream/:sourceId/stream/:streamId/:type?', (req, res) => { router.get('/m3u/:sourceId', async (req, res) => { try { const sourceId = req.params.sourceId; - const source = sources.getById(sourceId); + const source = await sources.getById(sourceId); if (!source || source.type !== 'm3u') { return res.status(404).json({ error: 'M3U source not found' }); } @@ -164,7 +165,7 @@ router.get('/m3u/:sourceId', async (req, res) => { router.get('/epg/:sourceId', async (req, res) => { try { const sourceId = req.params.sourceId; - const source = sources.getById(sourceId); + const source = await sources.getById(sourceId); if (!source || (source.type !== 'epg' && source.type !== 'xtream')) { return res.status(404).json({ error: 'Valid EPG source not found' }); } @@ -226,7 +227,7 @@ router.delete('/epg/:sourceId/cache', (req, res) => { */ router.post('/epg/:sourceId/channels', async (req, res) => { try { - const source = sources.getById(req.params.sourceId); + const source = await sources.getById(req.params.sourceId); if (!source || source.type !== 'epg') { return res.status(404).json({ error: 'EPG source not found' }); } @@ -277,7 +278,6 @@ router.get('/stream', async (req, res) => { const response = await fetch(url, { headers }); if (!response.ok) { console.error(`Upstream error for ${url.substring(0, 80)}...: ${response.status} ${response.statusText}`); - // Log response body for debugging 403s if (response.status === 403) { const errorBody = await response.text().catch(() => 'N/A'); console.error(`403 Response body: ${errorBody.substring(0, 200)}`); @@ -289,25 +289,22 @@ router.get('/stream', async (req, res) => { res.set('Access-Control-Allow-Origin', '*'); // Create an async iterator for the response body - // Note: undici fetch returns an iterable body const iterator = response.body[Symbol.asyncIterator](); const first = await iterator.next(); if (first.done) { - // Empty response res.set('Content-Type', contentType || 'application/octet-stream'); return res.end(); } const firstChunk = Buffer.from(first.value); - // Peek at first bytes to check for HLS manifest (#EXTM3U) + // Peek at first bytes to check for HLS manifest ({ #EXTM3U }) const textPrefix = firstChunk.subarray(0, 7).toString('utf8'); const contentLooksLikeHls = textPrefix === '#EXTM3U'; if (contentLooksLikeHls) { - // HLS Manifest: We must read the WHOLE stream to rewrite it - // This is fine because manifests are small text files + // HLS Manifest: We must read the WHOLE manifest to rewrite it const chunks = [firstChunk]; // Consume the rest of the stream @@ -318,8 +315,6 @@ router.get('/stream', async (req, res) => { } const buffer = Buffer.concat(chunks); - - // Use the final URL after redirects for base URL calculation const finalUrl = response.url || url; console.log(`[Proxy] Processing HLS manifest from: ${finalUrl.substring(0, 80)}...`); res.set('Content-Type', 'application/vnd.apple.mpegurl'); @@ -328,11 +323,9 @@ router.get('/stream', async (req, res) => { const finalUrlObj = new URL(finalUrl); const baseUrl = finalUrlObj.origin + finalUrlObj.pathname.substring(0, finalUrlObj.pathname.lastIndexOf('/') + 1); - console.log(`[Proxy] Base URL for rewriting: ${baseUrl}`); manifest = manifest.split('\n').map(line => { const trimmed = line.trim(); - // ... same rewrite logic as before ... if (trimmed === '' || trimmed.startsWith('#')) { if (trimmed.includes('URI="')) { return line.replace(/URI="([^"]+)"/g, (match, p1) => { @@ -360,34 +353,19 @@ router.get('/stream', async (req, res) => { return res.send(manifest); } - // Binary content (Video Segment): STREAM IT! - // This is the critical fix for "SocketError: other side closed" + // Binary content (Video Segment): Efficient Pipe console.log(`[Proxy] Piping binary stream (${contentType})`); res.set('Content-Type', contentType || 'application/octet-stream'); - // Write the first chunk we already peeked + // Write the chunk we peeked res.write(firstChunk); - // Pipe the rest of the iterator to the response - try { - let result = await iterator.next(); - while (!result.done) { - // If client disconnects, stop reading - if (res.writableEnded || res.closed) break; + // Stream the rest + // Create a readable stream from the iterator + const restOfStream = Readable.from(iterator); - const canWrite = res.write(Buffer.from(result.value)); - if (!canWrite) { - // Handle backpressure - await new Promise(resolve => res.once('drain', resolve)); - } - result = await iterator.next(); - } - res.end(); - } catch (streamErr) { - console.error('Stream pipe error:', streamErr); - if (!res.headersSent) res.status(500).end(); - else res.end(); - } + // Pipe to response + restOfStream.pipe(res); } catch (err) { console.error('Stream proxy error:', err); @@ -420,15 +398,21 @@ router.get('/image', async (req, res) => { return res.status(response.status).send('Failed to fetch image'); } - // Forward content type const contentType = response.headers.get('content-type') || 'image/png'; res.set('Content-Type', contentType); res.set('Access-Control-Allow-Origin', '*'); res.set('Cache-Control', 'public, max-age=86400'); // Cache for 24 hours - // Pipe the image data - const buffer = await response.arrayBuffer(); - res.send(Buffer.from(buffer)); + // Efficiently pipe the response body + if (response.body) { + // response.body is an AsyncIterable in standard fetch/undici + // Readable.from converts it to a Node.js Readable stream + const stream = Readable.from(response.body); + stream.pipe(res); + } else { + res.end(); + } + } catch (err) { console.error('Image proxy error:', err.message); res.status(500).send('Image proxy error'); diff --git a/server/routes/sources.js b/server/routes/sources.js index 996d6d4..6d00601 100644 --- a/server/routes/sources.js +++ b/server/routes/sources.js @@ -4,9 +4,9 @@ const { sources } = require('../db'); const xtreamApi = require('../services/xtreamApi'); // Get all sources -router.get('/', (req, res) => { +router.get('/', async (req, res) => { try { - const allSources = sources.getAll(); + const allSources = await sources.getAll(); // Don't expose passwords in list view const sanitized = allSources.map(s => ({ ...s, @@ -20,9 +20,9 @@ router.get('/', (req, res) => { }); // Get sources by type -router.get('/type/:type', (req, res) => { +router.get('/type/:type', async (req, res) => { try { - const typeSources = sources.getByType(req.params.type); + const typeSources = await sources.getByType(req.params.type); res.json(typeSources); } catch (err) { console.error('Error getting sources by type:', err); @@ -31,9 +31,9 @@ router.get('/type/:type', (req, res) => { }); // Get single source -router.get('/:id', (req, res) => { +router.get('/:id', async (req, res) => { try { - const source = sources.getById(req.params.id); + const source = await sources.getById(req.params.id); if (!source) { return res.status(404).json({ error: 'Source not found' }); } @@ -45,7 +45,7 @@ router.get('/:id', (req, res) => { }); // Create source -router.post('/', (req, res) => { +router.post('/', async (req, res) => { try { const { type, name, url, username, password } = req.body; @@ -57,7 +57,7 @@ router.post('/', (req, res) => { return res.status(400).json({ error: 'Invalid source type' }); } - const source = sources.create({ type, name, url, username, password }); + const source = await sources.create({ type, name, url, username, password }); res.status(201).json(source); } catch (err) { console.error('Error creating source:', err); @@ -66,15 +66,15 @@ router.post('/', (req, res) => { }); // Update source -router.put('/:id', (req, res) => { +router.put('/:id', async (req, res) => { try { - const existing = sources.getById(req.params.id); + const existing = await sources.getById(req.params.id); if (!existing) { return res.status(404).json({ error: 'Source not found' }); } const { name, url, username, password } = req.body; - const updated = sources.update(req.params.id, { + const updated = await sources.update(req.params.id, { name: name || existing.name, url: url || existing.url, username: username !== undefined ? username : existing.username, @@ -88,13 +88,13 @@ router.put('/:id', (req, res) => { }); // Delete source -router.delete('/:id', (req, res) => { +router.delete('/:id', async (req, res) => { try { - const existing = sources.getById(req.params.id); + const existing = await sources.getById(req.params.id); if (!existing) { return res.status(404).json({ error: 'Source not found' }); } - sources.delete(req.params.id); + await sources.delete(req.params.id); res.json({ success: true }); } catch (err) { console.error('Error deleting source:', err); @@ -103,9 +103,9 @@ router.delete('/:id', (req, res) => { }); // Toggle source enabled/disabled -router.post('/:id/toggle', (req, res) => { +router.post('/:id/toggle', async (req, res) => { try { - const updated = sources.toggleEnabled(req.params.id); + const updated = await sources.toggleEnabled(req.params.id); if (!updated) { return res.status(404).json({ error: 'Source not found' }); } @@ -119,7 +119,7 @@ router.post('/:id/toggle', (req, res) => { // Test source connection router.post('/:id/test', async (req, res) => { try { - const source = sources.getById(req.params.id); + const source = await sources.getById(req.params.id); if (!source) { return res.status(404).json({ error: 'Source not found' }); } diff --git a/server/services/epgParser.js b/server/services/epgParser.js index c031a8f..c0cde58 100644 --- a/server/services/epgParser.js +++ b/server/services/epgParser.js @@ -1,14 +1,11 @@ /** - * EPG (XMLTV) Parser - * Parses XMLTV format EPG data and extracts channel/programme information + * EPG (XMLTV) Parser (Streaming) + * Parses XMLTV format EPG data and extracts channel/programme information using streaming XML parser */ -const { parseString } = require('xml2js'); -const { promisify } = require('util'); +const sax = require('sax'); const zlib = require('zlib'); - -const parseXml = promisify(parseString); -const gunzip = promisify(zlib.gunzip); +const { Readable } = require('stream'); /** * Parse XMLTV date format (YYYYMMDDHHmmss +ZZZZ) @@ -40,96 +37,120 @@ function parseXmltvDate(dateStr) { } /** - * Parse XMLTV content - * @param {string} content - Raw XMLTV content + * Parse XMLTV content (Stream or String) + * @param {Readable|string} input - XMLTV content as Stream or String * @returns {Promise<{ channels: Array, programmes: Array }>} */ -async function parse(content) { - const result = await parseXml(content, { - explicitArray: false, - mergeAttrs: true - }); +function parse(input) { + return new Promise((resolve, reject) => { + const channels = []; + const programmes = []; - if (!result.tv) { - throw new Error('Invalid XMLTV format: missing root element'); - } + const saxStream = sax.createStream(true, { trim: true, normalize: true }); // strict mode - const tv = result.tv; - const channels = []; - const programmes = []; + let currentTag = null; + let currentObject = null; + let textBuffer = ''; - // Parse channels - const channelList = Array.isArray(tv.channel) ? tv.channel : (tv.channel ? [tv.channel] : []); - for (const ch of channelList) { - const channel = { - id: ch.id, - name: extractText(ch['display-name']), - icon: ch.icon ? (ch.icon.src || ch.icon) : null, - url: extractText(ch.url) - }; - channels.push(channel); - } + saxStream.on('error', function (e) { + // clear the error + this._parser.error = null; + this._parser.resume(); + console.warn('XML Parse Warning:', e.message); + }); - // Parse programmes - const programmeList = Array.isArray(tv.programme) ? tv.programme : (tv.programme ? [tv.programme] : []); - for (const prog of programmeList) { - const programme = { - channelId: prog.channel, - start: parseXmltvDate(prog.start), - stop: parseXmltvDate(prog.stop), - title: extractText(prog.title), - subtitle: extractText(prog['sub-title']), - description: extractText(prog.desc), - category: extractCategories(prog.category), - icon: prog.icon ? (prog.icon.src || prog.icon) : null, - date: extractText(prog.date), - episodeNum: extractEpisodeNum(prog['episode-num']) - }; - programmes.push(programme); - } + saxStream.on('opentag', function (node) { + currentTag = node.name; + const attr = node.attributes; - return { channels, programmes }; -} + if (currentTag === 'channel') { + currentObject = { + id: attr.id, + name: null, // Will be populated by display-name tag + 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 = ''; + }); -/** - * Extract text from XMLTV element (handles both string and object formats) - */ -function extractText(element) { - if (!element) return null; - if (typeof element === 'string') return element; - if (Array.isArray(element)) { - // Prefer English or first item - const en = element.find(e => e.lang === 'en' || !e.lang); - return extractText(en || element[0]); - } - if (element._) return element._; - if (element['#text']) return element['#text']; - return String(element); -} + saxStream.on('text', function (text) { + textBuffer += text; + }); -/** - * Extract categories array - */ -function extractCategories(category) { - if (!category) return []; - const cats = Array.isArray(category) ? category : [category]; - return cats.map(c => extractText(c)).filter(Boolean); -} + saxStream.on('cdata', function (text) { + textBuffer += text; + }); -/** - * Extract episode number - */ -function extractEpisodeNum(episodeNum) { - if (!episodeNum) return null; - const nums = Array.isArray(episodeNum) ? episodeNum : [episodeNum]; + saxStream.on('closetag', function (tagName) { + if (tagName === 'channel') { + if (currentObject) channels.push(currentObject); + currentObject = null; + } else if (tagName === 'programme') { + if (currentObject) programmes.push(currentObject); + currentObject = null; + } else if (currentObject) { + // Handle properties within objects + switch (tagName) { + case 'display-name': // channel name + if (!currentObject.name) currentObject.name = textBuffer; + break; + case 'url': // channel 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': + // Prefer system "xmltv_ns" or just take text + // Complex episode parsing logic can go here if needed + currentObject.episodeNum = textBuffer; + break; + } + } + }); - for (const num of nums) { - if (typeof num === 'string') return num; - if (num._ || num['#text']) { - return num._ || num['#text']; + saxStream.on('end', function () { + resolve({ channels, programmes }); + }); + + // Handle input type + if (typeof input === 'string') { + const inputStream = Readable.from([input]); + inputStream.pipe(saxStream); + } else { + input.pipe(saxStream); } - } - return null; + }); } /** @@ -167,28 +188,40 @@ async function fetchAndParse(url) { throw new Error(`Failed to fetch EPG: ${response.status} ${response.statusText}`); } - let content; - const buffer = await response.arrayBuffer(); - const bytes = new Uint8Array(buffer); - - // Check for gzip magic bytes (1f 8b) to detect actual gzip content - // This handles cases where the server auto-decompresses or the URL doesn't reflect actual encoding - const isActuallyGzipped = bytes.length >= 2 && bytes[0] === 0x1f && bytes[1] === 0x8b; - - if (isActuallyGzipped) { - // Handle gzipped EPG files (.xml.gz) - try { - const decompressed = await gunzip(Buffer.from(buffer)); - content = decompressed.toString('utf-8'); - } catch (err) { - throw new Error(`Failed to decompress gzipped EPG: ${err.message}`); - } + let stream; + if (response.body && typeof response.body.pipe === 'function') { + stream = response.body; + } else if (response.body) { + stream = Readable.fromWeb(response.body); } else { - // Already plain text (or auto-decompressed by fetch) - content = Buffer.from(buffer).toString('utf-8'); + stream = Readable.from([]); } - return parse(content); + // Check for GZIP + // Note: We can't easily check for magic bytes on a stream without buffering. + // We'll rely on response headers or file extension mostly, or try to peek. + // For now, let's assume if content-encoding is gzip OR url ends in .gz + + // However, undici/fetch usually handles 'Content-Encoding: gzip' automatically transparently. + // We only need to manually gunzip if the server serves it as application/octet-stream but it's actually gzipped, + // or if it's a .gz file download. + + // A robust way for streams is checking magic bytes, but that requires peeking. + // Simplified approach: try to pipe through gunzip if the URL indicates it. + + const isGzipped = url.endsWith('.gz') || (response.headers.get('content-type') || '').includes('gzip'); + + if (isGzipped) { + const gunzip = zlib.createGunzip(); + stream.pipe(gunzip); + return parse(gunzip); + } + + // In the previous version we read magic bytes. + // To support that with streams we'd need a peek stream. + // For now let's trust the transparent decompression of fetch or the URL. + + return parse(stream); } module.exports = { diff --git a/server/services/m3uParser.js b/server/services/m3uParser.js index 72415e6..0708763 100644 --- a/server/services/m3uParser.js +++ b/server/services/m3uParser.js @@ -1,8 +1,11 @@ /** - * M3U Playlist Parser - * Parses EXTM3U format playlists and extracts channel information + * M3U Playlist Parser (Streaming) + * Parses EXTM3U format playlists and extracts channel information line-by-line */ +const readline = require('readline'); +const { Readable } = require('stream'); + /** * Generate a simple stable ID from name and group * @param {string} name - Channel name @@ -21,69 +24,6 @@ function generateStableId(name, group) { return `m3u_${Math.abs(hash).toString(36)}`; } -/** - * Parse M3U playlist content - * @param {string} content - Raw M3U playlist content - * @returns {{ channels: Array, groups: Array }} - */ -function parse(content) { - const lines = content.split('\n').map(line => line.trim()); - const channels = []; - const groupsSet = new Set(); - - // Verify it's a valid M3U file - if (!lines[0] || !lines[0].startsWith('#EXTM3U')) { - throw new Error('Invalid M3U format: missing #EXTM3U header'); - } - - let currentInfo = null; - let currentGroup = null; - - for (let i = 1; i < lines.length; i++) { - const line = lines[i]; - - if (line.startsWith('#EXTINF:')) { - // Parse EXTINF line - currentInfo = parseExtinf(line); - if (currentInfo.groupTitle) { - groupsSet.add(currentInfo.groupTitle); - currentGroup = currentInfo.groupTitle; - } - } else if (line.startsWith('#EXTGRP:')) { - // Parse EXTGRP line (alternative group specification) - currentGroup = line.substring(8).trim(); - groupsSet.add(currentGroup); - if (currentInfo) { - currentInfo.groupTitle = currentGroup; - } - } else if (line && !line.startsWith('#')) { - // This is a stream URL - if (currentInfo) { - const groupTitle = currentInfo.groupTitle || currentGroup || 'Uncategorized'; - // Generate a stable ID: use tvgId if present, otherwise hash name+group - const stableId = currentInfo.tvgId || generateStableId(currentInfo.name, groupTitle); - - channels.push({ - ...currentInfo, - id: stableId, - url: line, - groupTitle: groupTitle - }); - currentInfo = null; - } - } - } - - // Convert groups to array of objects - const groups = Array.from(groupsSet).map((name, index) => ({ - id: `group_${index}`, - name, - channelCount: channels.filter(c => c.groupTitle === name).length - })); - - return { channels, groups }; -} - /** * Parse EXTINF line and extract attributes * @param {string} line - EXTINF line @@ -138,6 +78,75 @@ function parseExtinf(line) { return info; } +/** + * Parse M3U content (Stream or String) + * @param {Readable|string} input - M3U content as Stream or String + * @returns {Promise<{ channels: Array, groups: Array }>} + */ +async function parse(input) { + const channels = []; + const groupsSet = new Set(); + let currentInfo = null; + let currentGroup = null; + + let inputStream; + if (typeof input === 'string') { + inputStream = Readable.from([input]); + } else { + inputStream = input; + } + + const rl = readline.createInterface({ + input: inputStream, + crlfDelay: Infinity + }); + + for await (const line of rl) { + const trimmed = line.trim(); + if (!trimmed) continue; + + if (trimmed.startsWith('#EXTINF:')) { + // Parse EXTINF line + currentInfo = parseExtinf(trimmed); + if (currentInfo.groupTitle) { + groupsSet.add(currentInfo.groupTitle); + currentGroup = currentInfo.groupTitle; + } + } else if (trimmed.startsWith('#EXTGRP:')) { + // Parse EXTGRP line (alternative group specification) + currentGroup = trimmed.substring(8).trim(); + groupsSet.add(currentGroup); + if (currentInfo) { + currentInfo.groupTitle = currentGroup; + } + } else if (!trimmed.startsWith('#')) { + // This is a stream URL + if (currentInfo) { + const groupTitle = currentInfo.groupTitle || currentGroup || 'Uncategorized'; + // Generate a stable ID: use tvgId if present, otherwise hash name+group + const stableId = currentInfo.tvgId || generateStableId(currentInfo.name, groupTitle); + + channels.push({ + ...currentInfo, + id: stableId, + url: trimmed, + groupTitle: groupTitle + }); + currentInfo = null; + } + } + } + + // Convert groups to array of objects + const groups = Array.from(groupsSet).map((name, index) => ({ + id: `group_${index}`, + name, + channelCount: channels.filter(c => c.groupTitle === name).length + })); + + return { channels, groups }; +} + /** * Fetch and parse M3U from URL * @param {string} url - M3U playlist URL @@ -148,8 +157,20 @@ async function fetchAndParse(url) { if (!response.ok) { throw new Error(`Failed to fetch M3U: ${response.status} ${response.statusText}`); } - const content = await response.text(); - return parse(content); + + // Check if body is a Node.js stream (undici/node-fetch) or web stream + let stream; + if (response.body && typeof response.body.pipe === 'function') { + stream = response.body; + } else if (response.body) { + // Convert Web Stream to Node Readable for readline + stream = Readable.fromWeb(response.body); + } else { + // Fallback for empty body + stream = Readable.from([]); + } + + return parse(stream); } module.exports = { parse, parseExtinf, fetchAndParse };