Refactored code to use sqlite for channel, epg and vod data
This commit is contained in:
@@ -0,0 +1,111 @@
|
||||
const Database = require('better-sqlite3');
|
||||
const path = require('path');
|
||||
const fs = require('fs');
|
||||
|
||||
const dataDir = path.join(__dirname, '..', '..', 'data');
|
||||
const dbPath = path.join(dataDir, 'content.db');
|
||||
|
||||
// Ensure data directory exists
|
||||
if (!fs.existsSync(dataDir)) {
|
||||
fs.mkdirSync(dataDir, { recursive: true });
|
||||
}
|
||||
|
||||
let db;
|
||||
|
||||
function getDb() {
|
||||
if (!db) {
|
||||
console.log('[SQLite] Opening database at', dbPath);
|
||||
db = new Database(dbPath);
|
||||
// Optimize performance
|
||||
db.pragma('journal_mode = WAL');
|
||||
db.pragma('synchronous = NORMAL');
|
||||
initSchema();
|
||||
}
|
||||
return db;
|
||||
}
|
||||
|
||||
function initSchema() {
|
||||
if (!db) throw new Error('Database not initialized');
|
||||
|
||||
// Categories (Groups)
|
||||
db.exec(`
|
||||
CREATE TABLE IF NOT EXISTS categories (
|
||||
id TEXT PRIMARY KEY, -- Composite key: sourceId:categoryId
|
||||
source_id INTEGER NOT NULL,
|
||||
category_id TEXT NOT NULL,
|
||||
type TEXT NOT NULL, -- 'live', 'movie', 'series'
|
||||
name TEXT NOT NULL,
|
||||
parent_id TEXT, -- For nested categories
|
||||
is_hidden INTEGER DEFAULT 0,
|
||||
data JSON -- Extra provider data
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_categories_source_type ON categories(source_id, type);
|
||||
`);
|
||||
|
||||
// Playlist Items (Channels, Movies, Series, Episodes)
|
||||
db.exec(`
|
||||
CREATE TABLE IF NOT EXISTS playlist_items (
|
||||
id TEXT PRIMARY KEY, -- Composite key: sourceId:itemId
|
||||
source_id INTEGER NOT NULL,
|
||||
item_id TEXT NOT NULL, -- Original ID from provider
|
||||
type TEXT NOT NULL, -- 'live', 'movie', 'series', 'episode'
|
||||
name TEXT NOT NULL,
|
||||
category_id TEXT, -- maps to categories.category_id (not our composite id)
|
||||
parent_id TEXT, -- For episodes -> series_id
|
||||
|
||||
-- Common Media Fields
|
||||
stream_icon TEXT,
|
||||
stream_url TEXT, -- Direct link if available
|
||||
container_extension TEXT,
|
||||
|
||||
-- VOD/Series Specific
|
||||
rating REAL,
|
||||
year TEXT,
|
||||
added_at TEXT,
|
||||
|
||||
-- App State
|
||||
is_hidden INTEGER DEFAULT 0,
|
||||
is_favorite INTEGER DEFAULT 0,
|
||||
|
||||
data JSON -- Full original JSON object
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_items_source_type ON playlist_items(source_id, type);
|
||||
CREATE INDEX IF NOT EXISTS idx_items_category ON playlist_items(source_id, category_id);
|
||||
`);
|
||||
|
||||
// EPG Programs
|
||||
// Optimized for range queries
|
||||
db.exec(`
|
||||
CREATE TABLE IF NOT EXISTS epg_programs (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
channel_id TEXT NOT NULL, -- matches playlist_items.id if possible, or mapping key
|
||||
source_id INTEGER NOT NULL,
|
||||
start_time INTEGER NOT NULL, -- Unix timestamp (ms)
|
||||
end_time INTEGER NOT NULL, -- Unix timestamp (ms)
|
||||
title TEXT,
|
||||
description TEXT,
|
||||
data JSON
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_epg_channel_time ON epg_programs(channel_id, start_time, end_time);
|
||||
CREATE INDEX IF NOT EXISTS idx_epg_cleanup ON epg_programs(end_time); -- For deleting old programs
|
||||
`);
|
||||
|
||||
// Sync Status
|
||||
db.exec(`
|
||||
CREATE TABLE IF NOT EXISTS sync_status (
|
||||
source_id INTEGER NOT NULL,
|
||||
type TEXT NOT NULL, -- 'live', 'vod', 'series', 'epg'
|
||||
last_sync INTEGER NOT NULL,
|
||||
status TEXT, -- 'success', 'error', 'syncing'
|
||||
error TEXT,
|
||||
PRIMARY KEY (source_id, type)
|
||||
);
|
||||
`);
|
||||
|
||||
console.log('[SQLite] Schema initialized');
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
getDb,
|
||||
initSchema
|
||||
};
|
||||
+11
-1
@@ -1,6 +1,8 @@
|
||||
const express = require('express');
|
||||
const path = require('path');
|
||||
const passport = require('passport');
|
||||
const { migrateUserData } = require('./services/migrationService');
|
||||
const syncService = require('./services/syncService');
|
||||
|
||||
// Initialize database
|
||||
require('./db');
|
||||
@@ -74,6 +76,14 @@ app.use((err, req, res, next) => {
|
||||
res.status(500).json({ error: 'Internal server error' });
|
||||
});
|
||||
|
||||
app.listen(PORT, () => {
|
||||
app.listen(PORT, async () => {
|
||||
console.log(`NodeCast TV server running on http://localhost:${PORT}`);
|
||||
|
||||
// Run migration after startup
|
||||
// await migrateUserData();
|
||||
|
||||
// Trigger background sync with delay to allow server to settle
|
||||
setTimeout(() => {
|
||||
syncService.syncAll().catch(console.error);
|
||||
}, 5000);
|
||||
});
|
||||
|
||||
+201
-41
@@ -1,34 +1,92 @@
|
||||
const express = require('express');
|
||||
const router = express.Router();
|
||||
const { hiddenItems } = require('../db');
|
||||
const { getDb } = require('../db/sqlite');
|
||||
|
||||
// Get all hidden items
|
||||
// Helper to map API item types to DB types and tables
|
||||
function mapItemType(apiType) {
|
||||
switch (apiType) {
|
||||
case 'channel': return { table: 'playlist_items', type: 'live' };
|
||||
case 'group': return { table: 'categories', type: 'live' };
|
||||
case 'vod_category': return { table: 'categories', type: 'movie' };
|
||||
case 'series_category': return { table: 'categories', type: 'series' };
|
||||
case 'movie': return { table: 'playlist_items', type: 'movie' };
|
||||
case 'series': return { table: 'playlist_items', type: 'series' };
|
||||
default: return null;
|
||||
}
|
||||
}
|
||||
|
||||
// Get all hidden items (formatted like db.json for frontend compatibility)
|
||||
router.get('/hidden', async (req, res) => {
|
||||
try {
|
||||
const { sourceId } = req.query;
|
||||
const items = await hiddenItems.getAll(sourceId ? parseInt(sourceId) : null);
|
||||
res.json(items);
|
||||
const db = getDb();
|
||||
|
||||
let hidden = [];
|
||||
const resultFormat = (row, itemType) => ({
|
||||
source_id: row.source_id,
|
||||
item_type: itemType,
|
||||
item_id: itemType.includes('category') || itemType === 'group' ? row.category_id : row.item_id
|
||||
});
|
||||
|
||||
// Query Categories
|
||||
let catQuery = `SELECT source_id, category_id, type FROM categories WHERE is_hidden = 1`;
|
||||
let itemQuery = `SELECT source_id, item_id, type FROM playlist_items WHERE is_hidden = 1`;
|
||||
|
||||
const params = [];
|
||||
if (sourceId) {
|
||||
catQuery += ` AND source_id = ?`;
|
||||
itemQuery += ` AND source_id = ?`;
|
||||
const sid = parseInt(sourceId);
|
||||
params.push(sid);
|
||||
}
|
||||
|
||||
const hiddenCats = db.prepare(catQuery).all(...params);
|
||||
const hiddenItems = db.prepare(itemQuery).all(...params);
|
||||
|
||||
hiddenCats.forEach(row => {
|
||||
let apiType;
|
||||
if (row.type === 'live') apiType = 'group';
|
||||
else if (row.type === 'movie') apiType = 'vod_category';
|
||||
else if (row.type === 'series') apiType = 'series_category';
|
||||
|
||||
if (apiType) hidden.push(resultFormat(row, apiType));
|
||||
});
|
||||
|
||||
hiddenItems.forEach(row => {
|
||||
let apiType;
|
||||
if (row.type === 'live') apiType = 'channel';
|
||||
else if (row.type === 'movie') apiType = 'movie';
|
||||
else if (row.type === 'series') apiType = 'series';
|
||||
|
||||
if (apiType) hidden.push(resultFormat(row, apiType));
|
||||
});
|
||||
|
||||
res.json(hidden);
|
||||
} catch (err) {
|
||||
console.error('Error getting hidden items:', err);
|
||||
res.status(500).json({ error: 'Failed to get hidden items' });
|
||||
}
|
||||
});
|
||||
|
||||
// Hide a channel, group, or category
|
||||
// Hide item
|
||||
router.post('/hide', async (req, res) => {
|
||||
try {
|
||||
const { sourceId, itemType, itemId } = req.body;
|
||||
const mapping = mapItemType(itemType);
|
||||
|
||||
if (!sourceId || !itemType || !itemId) {
|
||||
return res.status(400).json({ error: 'sourceId, itemType, and itemId are required' });
|
||||
}
|
||||
if (!mapping) return res.status(400).json({ error: 'Invalid item type' });
|
||||
|
||||
const validTypes = ['channel', 'group', 'vod_category', 'series_category'];
|
||||
if (!validTypes.includes(itemType)) {
|
||||
return res.status(400).json({ error: `itemType must be one of: ${validTypes.join(', ')}` });
|
||||
}
|
||||
const db = getDb();
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
|
||||
const stmt = db.prepare(`
|
||||
UPDATE ${mapping.table}
|
||||
SET is_hidden = 1
|
||||
WHERE source_id = ? AND type = ? AND ${idCol} = ?
|
||||
`);
|
||||
|
||||
stmt.run(sourceId, mapping.type, itemId);
|
||||
|
||||
await hiddenItems.hide(sourceId, itemType, itemId);
|
||||
res.json({ success: true });
|
||||
} catch (err) {
|
||||
console.error('Error hiding item:', err);
|
||||
@@ -36,16 +94,25 @@ router.post('/hide', async (req, res) => {
|
||||
}
|
||||
});
|
||||
|
||||
// Show (unhide) a channel or group
|
||||
// Show item
|
||||
router.post('/show', async (req, res) => {
|
||||
try {
|
||||
const { sourceId, itemType, itemId } = req.body;
|
||||
const mapping = mapItemType(itemType);
|
||||
|
||||
if (!sourceId || !itemType || !itemId) {
|
||||
return res.status(400).json({ error: 'sourceId, itemType, and itemId are required' });
|
||||
}
|
||||
if (!mapping) return res.status(400).json({ error: 'Invalid item type' });
|
||||
|
||||
const db = getDb();
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
|
||||
const stmt = db.prepare(`
|
||||
UPDATE ${mapping.table}
|
||||
SET is_hidden = 0
|
||||
WHERE source_id = ? AND type = ? AND ${idCol} = ?
|
||||
`);
|
||||
|
||||
stmt.run(sourceId, mapping.type, itemId);
|
||||
|
||||
await hiddenItems.show(sourceId, itemType, itemId);
|
||||
res.json({ success: true });
|
||||
} catch (err) {
|
||||
console.error('Error showing item:', err);
|
||||
@@ -53,56 +120,149 @@ router.post('/show', async (req, res) => {
|
||||
}
|
||||
});
|
||||
|
||||
// Check if item is hidden
|
||||
// Check hidden status
|
||||
router.get('/hidden/check', async (req, res) => {
|
||||
try {
|
||||
const { sourceId, itemType, itemId } = req.query;
|
||||
const mapping = mapItemType(itemType);
|
||||
if (!mapping) return res.json({ hidden: false });
|
||||
|
||||
if (!sourceId || !itemType || !itemId) {
|
||||
return res.status(400).json({ error: 'sourceId, itemType, and itemId are required' });
|
||||
}
|
||||
const db = getDb();
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
|
||||
// isHidden is now async
|
||||
const isHidden = await hiddenItems.isHidden(parseInt(sourceId), itemType, itemId);
|
||||
res.json({ hidden: isHidden });
|
||||
const row = db.prepare(`
|
||||
SELECT is_hidden FROM ${mapping.table}
|
||||
WHERE source_id = ? AND type = ? AND ${idCol} = ?
|
||||
`).get(sourceId, mapping.type, itemId);
|
||||
|
||||
res.json({ hidden: !!(row && row.is_hidden) });
|
||||
} catch (err) {
|
||||
console.error('Error checking hidden status:', err);
|
||||
res.status(500).json({ error: 'Failed to check hidden status' });
|
||||
console.error('Error checking hidden:', err);
|
||||
res.status(500).json({ error: 'Failed to check status' });
|
||||
}
|
||||
});
|
||||
|
||||
// Bulk hide channels and groups
|
||||
// Bulk Hide
|
||||
router.post('/hide/bulk', async (req, res) => {
|
||||
try {
|
||||
const { items } = req.body;
|
||||
if (!Array.isArray(items)) return res.status(400).json({ error: 'items array required' });
|
||||
|
||||
if (!items || !Array.isArray(items) || items.length === 0) {
|
||||
return res.status(400).json({ error: 'items array is required' });
|
||||
}
|
||||
const db = getDb();
|
||||
const runBulk = db.transaction((list) => {
|
||||
for (const item of list) {
|
||||
const mapping = mapItemType(item.itemType);
|
||||
if (mapping) {
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
db.prepare(`
|
||||
UPDATE ${mapping.table} SET is_hidden = 1
|
||||
WHERE source_id = ? AND type = ? AND ${idCol} = ?
|
||||
`).run(item.sourceId, mapping.type, item.itemId);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
await hiddenItems.bulkHide(items);
|
||||
runBulk(items);
|
||||
res.json({ success: true, count: items.length });
|
||||
} catch (err) {
|
||||
console.error('Error bulk hiding items:', err);
|
||||
res.status(500).json({ error: 'Failed to bulk hide items' });
|
||||
if (err.code === 'SQLITE_BUSY') {
|
||||
return res.status(503).json({ error: 'Database is busy, please try again' });
|
||||
}
|
||||
console.error('Error bulk hide:', err);
|
||||
res.status(500).json({ error: 'Failed' });
|
||||
}
|
||||
});
|
||||
|
||||
// Bulk show channels and groups
|
||||
// Bulk Show
|
||||
router.post('/show/bulk', async (req, res) => {
|
||||
try {
|
||||
const { items } = req.body;
|
||||
if (!Array.isArray(items)) return res.status(400).json({ error: 'items array required' });
|
||||
|
||||
if (!items || !Array.isArray(items) || items.length === 0) {
|
||||
return res.status(400).json({ error: 'items array is required' });
|
||||
}
|
||||
const db = getDb();
|
||||
const runBulk = db.transaction((list) => {
|
||||
for (const item of list) {
|
||||
const mapping = mapItemType(item.itemType);
|
||||
if (mapping) {
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
db.prepare(`
|
||||
UPDATE ${mapping.table} SET is_hidden = 0
|
||||
WHERE source_id = ? AND type = ? AND ${idCol} = ?
|
||||
`).run(item.sourceId, mapping.type, item.itemId);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
await hiddenItems.bulkShow(items);
|
||||
runBulk(items);
|
||||
res.json({ success: true, count: items.length });
|
||||
} catch (err) {
|
||||
console.error('Error bulk showing items:', err);
|
||||
res.status(500).json({ error: 'Failed to bulk show items' });
|
||||
if (err.code === 'SQLITE_BUSY') {
|
||||
return res.status(503).json({ error: 'Database is busy, please try again' });
|
||||
}
|
||||
console.error('Error bulk show:', err);
|
||||
res.status(500).json({ error: 'Failed' });
|
||||
}
|
||||
});
|
||||
|
||||
// Show ALL items for a source (single SQL statement - much faster than bulk)
|
||||
router.post('/show/all', async (req, res) => {
|
||||
try {
|
||||
const { sourceId, contentType } = req.body;
|
||||
if (!sourceId) return res.status(400).json({ error: 'sourceId required' });
|
||||
|
||||
const db = getDb();
|
||||
let catCount = 0;
|
||||
let itemCount = 0;
|
||||
|
||||
// Determine which types to update based on contentType
|
||||
const types = contentType === 'movies' ? ['movie']
|
||||
: contentType === 'series' ? ['series']
|
||||
: ['live']; // default to channels
|
||||
|
||||
for (const type of types) {
|
||||
const catResult = db.prepare(`UPDATE categories SET is_hidden = 0 WHERE source_id = ? AND type = ?`).run(sourceId, type);
|
||||
const itemResult = db.prepare(`UPDATE playlist_items SET is_hidden = 0 WHERE source_id = ? AND type = ?`).run(sourceId, type);
|
||||
catCount += catResult.changes;
|
||||
itemCount += itemResult.changes;
|
||||
}
|
||||
|
||||
console.log(`[Channels] Show all for source ${sourceId} (${contentType}): ${catCount} categories, ${itemCount} items`);
|
||||
res.json({ success: true, categoriesUpdated: catCount, itemsUpdated: itemCount });
|
||||
} catch (err) {
|
||||
console.error('Error show all:', err);
|
||||
res.status(500).json({ error: 'Failed to show all' });
|
||||
}
|
||||
});
|
||||
|
||||
// Hide ALL items for a source (single SQL statement - much faster than bulk)
|
||||
router.post('/hide/all', async (req, res) => {
|
||||
try {
|
||||
const { sourceId, contentType } = req.body;
|
||||
if (!sourceId) return res.status(400).json({ error: 'sourceId required' });
|
||||
|
||||
const db = getDb();
|
||||
let catCount = 0;
|
||||
let itemCount = 0;
|
||||
|
||||
// Determine which types to update based on contentType
|
||||
const types = contentType === 'movies' ? ['movie']
|
||||
: contentType === 'series' ? ['series']
|
||||
: ['live']; // default to channels
|
||||
|
||||
for (const type of types) {
|
||||
const catResult = db.prepare(`UPDATE categories SET is_hidden = 1 WHERE source_id = ? AND type = ?`).run(sourceId, type);
|
||||
const itemResult = db.prepare(`UPDATE playlist_items SET is_hidden = 1 WHERE source_id = ? AND type = ?`).run(sourceId, type);
|
||||
catCount += catResult.changes;
|
||||
itemCount += itemResult.changes;
|
||||
}
|
||||
|
||||
console.log(`[Channels] Hide all for source ${sourceId} (${contentType}): ${catCount} categories, ${itemCount} items`);
|
||||
res.json({ success: true, categoriesUpdated: catCount, itemsUpdated: itemCount });
|
||||
} catch (err) {
|
||||
console.error('Error hide all:', err);
|
||||
res.status(500).json({ error: 'Failed to hide all' });
|
||||
}
|
||||
});
|
||||
|
||||
module.exports = router;
|
||||
|
||||
|
||||
+427
-399
@@ -1,472 +1,500 @@
|
||||
const express = require('express');
|
||||
const router = express.Router();
|
||||
const { sources } = require('../db');
|
||||
const { getDb } = require('../db/sqlite'); // Import SQLite
|
||||
const xtreamApi = require('../services/xtreamApi');
|
||||
const m3uParser = require('../services/m3uParser');
|
||||
const epgParser = require('../services/epgParser');
|
||||
const cache = require('../services/cache');
|
||||
const path = require('path');
|
||||
const fs = require('fs');
|
||||
const http = require('http');
|
||||
const https = require('https');
|
||||
const { spawn } = require('child_process');
|
||||
const ffmpegPath = require('ffmpeg-static');
|
||||
const { Readable } = require('stream');
|
||||
|
||||
// Default cache TTL: 24 hours
|
||||
const DEFAULT_MAX_AGE_HOURS = 24;
|
||||
// Helper to get formatted category list from DB
|
||||
function getCategoriesFromDb(sourceId, type, includeHidden = false) {
|
||||
const db = getDb();
|
||||
let query = `
|
||||
SELECT category_id, name as category_name, parent_id
|
||||
FROM categories
|
||||
WHERE source_id = ? AND type = ?
|
||||
`;
|
||||
if (!includeHidden) {
|
||||
query += ` AND is_hidden = 0`;
|
||||
}
|
||||
query += ` ORDER BY name ASC`;
|
||||
const cats = db.prepare(query).all(sourceId, type);
|
||||
return cats;
|
||||
}
|
||||
|
||||
/**
|
||||
* Proxy Xtream API calls
|
||||
* GET /api/proxy/xtream/:sourceId/:action
|
||||
*/
|
||||
router.get('/xtream/:sourceId/:action', async (req, res) => {
|
||||
// Helper to get formatted streams from DB
|
||||
function getStreamsFromDb(sourceId, type, categoryId = null, includeHidden = false) {
|
||||
const db = getDb();
|
||||
let query = `
|
||||
SELECT item_id, name, stream_icon, added_at, rating, container_extension, year, category_id, data
|
||||
FROM playlist_items
|
||||
WHERE source_id = ? AND type = ?
|
||||
`;
|
||||
if (!includeHidden) {
|
||||
query += ` AND is_hidden = 0`;
|
||||
}
|
||||
const params = [sourceId, type];
|
||||
|
||||
if (categoryId) {
|
||||
query += ` AND category_id = ?`;
|
||||
params.push(categoryId);
|
||||
}
|
||||
|
||||
// Default sorting
|
||||
// query += ` ORDER BY name ASC`; // Sorting usually handled by client
|
||||
|
||||
const items = db.prepare(query).all(...params);
|
||||
|
||||
// Map to Xtream format
|
||||
return items.map(item => {
|
||||
const data = JSON.parse(item.data || '{}');
|
||||
// Override with our local fields if needed, or just return the mixed object
|
||||
// We should ensure critical fields are present
|
||||
return {
|
||||
...data,
|
||||
stream_id: item.item_id, // ensure ID matches what client expects
|
||||
series_id: type === 'series' ? item.item_id : undefined,
|
||||
name: item.name,
|
||||
stream_icon: item.stream_icon,
|
||||
cover: item.stream_icon, // series/vod often use cover
|
||||
added: item.added_at,
|
||||
rating: item.rating,
|
||||
container_extension: item.container_extension,
|
||||
category_id: item.category_id
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
// --- Xtream Codes Proxy API --- //
|
||||
|
||||
// Login / Authenticate
|
||||
router.get('/xtream/:sourceId', async (req, res) => {
|
||||
try {
|
||||
const sourceId = req.params.sourceId;
|
||||
const source = await sources.getById(sourceId);
|
||||
if (!source || source.type !== 'xtream') {
|
||||
return res.status(404).json({ error: 'Xtream source not found' });
|
||||
}
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source || source.type !== 'xtream') return res.status(404).send('Source not found');
|
||||
|
||||
const { action } = req.params;
|
||||
const { category_id, stream_id, vod_id, series_id, limit, refresh, maxAge } = req.query;
|
||||
const forceRefresh = refresh === '1';
|
||||
const maxAgeHours = parseInt(maxAge) || DEFAULT_MAX_AGE_HOURS;
|
||||
const maxAgeMs = maxAgeHours * 60 * 60 * 1000;
|
||||
// Proxy auth check to upstream to ensure credentials are still valid
|
||||
|
||||
// Actions that should be cached
|
||||
const cacheableActions = [
|
||||
'live_categories', 'live_streams',
|
||||
'vod_categories', 'vod_streams',
|
||||
'series_categories', 'series'
|
||||
];
|
||||
|
||||
// Build cache key (include category_id if present)
|
||||
const cacheKey = category_id ? `${action}_${category_id}` : action;
|
||||
const cached = cache.get(`xtream:${source.id}:auth`);
|
||||
if (cached) return res.json(cached);
|
||||
|
||||
// Check cache for cacheable actions
|
||||
if (!forceRefresh && cacheableActions.includes(action)) {
|
||||
const cached = cache.get('xtream', sourceId, cacheKey, maxAgeMs);
|
||||
if (cached) {
|
||||
return res.json(cached);
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch fresh data
|
||||
const api = xtreamApi.createFromSource(source);
|
||||
let data;
|
||||
switch (action) {
|
||||
case 'auth':
|
||||
data = await api.authenticate();
|
||||
break;
|
||||
case 'live_categories':
|
||||
data = await api.getLiveCategories();
|
||||
break;
|
||||
case 'live_streams':
|
||||
data = await api.getLiveStreams(category_id);
|
||||
break;
|
||||
case 'vod_categories':
|
||||
data = await api.getVodCategories();
|
||||
break;
|
||||
case 'vod_streams':
|
||||
data = await api.getVodStreams(category_id);
|
||||
break;
|
||||
case 'vod_info':
|
||||
data = await api.getVodInfo(vod_id);
|
||||
break;
|
||||
case 'series_categories':
|
||||
data = await api.getSeriesCategories();
|
||||
break;
|
||||
case 'series':
|
||||
data = await api.getSeries(category_id);
|
||||
break;
|
||||
case 'series_info':
|
||||
data = await api.getSeriesInfo(series_id);
|
||||
break;
|
||||
case 'short_epg':
|
||||
data = await api.getShortEpg(stream_id, limit);
|
||||
break;
|
||||
default:
|
||||
return res.status(400).json({ error: 'Unknown action' });
|
||||
}
|
||||
|
||||
// Cache the result for cacheable actions
|
||||
if (cacheableActions.includes(action)) {
|
||||
cache.set('xtream', sourceId, cacheKey, data);
|
||||
}
|
||||
|
||||
const data = await api.authenticate();
|
||||
cache.set(`xtream:${source.id}:auth`, data, 300); // 5 min cache
|
||||
res.json(data);
|
||||
} catch (err) {
|
||||
console.error('Xtream proxy error:', err);
|
||||
res.status(500).json({ error: err.message });
|
||||
res.status(502).json({ error: 'Upstream error', details: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Get Xtream stream URL
|
||||
* GET /api/proxy/xtream/:sourceId/stream/:streamId
|
||||
*/
|
||||
router.get('/xtream/:sourceId/stream/:streamId/:type?', async (req, res) => {
|
||||
// Live Categories
|
||||
router.get('/xtream/:sourceId/live_categories', async (req, res) => {
|
||||
try {
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const includeHidden = req.query.includeHidden === 'true';
|
||||
const cats = getCategoriesFromDb(sourceId, 'live', includeHidden);
|
||||
res.json(cats);
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
// Live Streams
|
||||
router.get('/xtream/:sourceId/live_streams', async (req, res) => {
|
||||
try {
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const categoryId = req.query.category_id;
|
||||
const includeHidden = req.query.includeHidden === 'true';
|
||||
const streams = getStreamsFromDb(sourceId, 'live', categoryId, includeHidden);
|
||||
res.json(streams);
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
// VOD Categories
|
||||
router.get('/xtream/:sourceId/vod_categories', async (req, res) => {
|
||||
try {
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const includeHidden = req.query.includeHidden === 'true';
|
||||
const cats = getCategoriesFromDb(sourceId, 'movie', includeHidden);
|
||||
res.json(cats);
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
// VOD Streams
|
||||
router.get('/xtream/:sourceId/vod_streams', async (req, res) => {
|
||||
try {
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const categoryId = req.query.category_id;
|
||||
const includeHidden = req.query.includeHidden === 'true';
|
||||
const streams = getStreamsFromDb(sourceId, 'movie', categoryId, includeHidden);
|
||||
res.json(streams);
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
// Series Categories
|
||||
router.get('/xtream/:sourceId/series_categories', async (req, res) => {
|
||||
try {
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const includeHidden = req.query.includeHidden === 'true';
|
||||
const cats = getCategoriesFromDb(sourceId, 'series', includeHidden);
|
||||
res.json(cats);
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
// Series
|
||||
router.get('/xtream/:sourceId/series', async (req, res) => {
|
||||
try {
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const categoryId = req.query.category_id;
|
||||
const includeHidden = req.query.includeHidden === 'true';
|
||||
const streams = getStreamsFromDb(sourceId, 'series', categoryId, includeHidden);
|
||||
res.json(streams);
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
// Series Info (Episodes)
|
||||
// Proxy series info request
|
||||
router.get('/xtream/:sourceId/series_info', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source) return res.status(404).send('Source not found');
|
||||
|
||||
const seriesId = req.query.series_id;
|
||||
if (!seriesId) return res.status(400).send('series_id required');
|
||||
|
||||
const cacheKey = `xtream:${source.id}:series_info:${seriesId}`;
|
||||
const cached = cache.get(cacheKey);
|
||||
if (cached) return res.json(cached);
|
||||
|
||||
const api = xtreamApi.createFromSource(source);
|
||||
const data = await api.getSeriesInfo(seriesId);
|
||||
cache.set(cacheKey, data, 3600); // 1 hour
|
||||
res.json(data);
|
||||
} catch (err) {
|
||||
res.status(502).json({ error: 'Upstream error', details: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
// VOD Info
|
||||
router.get('/xtream/:sourceId/vod_info', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source) return res.status(404).send('Source not found');
|
||||
|
||||
const vodId = req.query.vod_id;
|
||||
if (!vodId) return res.status(400).send('vod_id required');
|
||||
|
||||
const cacheKey = `xtream:${source.id}:vod_info:${vodId}`;
|
||||
const cached = cache.get(cacheKey);
|
||||
if (cached) return res.json(cached);
|
||||
|
||||
const api = xtreamApi.createFromSource(source);
|
||||
const data = await api.getVodInfo(vodId);
|
||||
cache.set(cacheKey, data, 3600); // 1 hour
|
||||
res.json(data);
|
||||
} catch (err) {
|
||||
res.status(502).json({ error: 'Upstream error', details: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
// Get Stream URL for playback
|
||||
// Returns the direct stream URL for a given stream ID
|
||||
router.get('/xtream/:sourceId/stream/:streamId/:type', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source || source.type !== 'xtream') {
|
||||
return res.status(404).json({ error: 'Xtream source not found' });
|
||||
}
|
||||
|
||||
const api = xtreamApi.createFromSource(source);
|
||||
const { streamId, type = 'live' } = req.params;
|
||||
const { container = 'm3u8' } = req.query;
|
||||
const streamId = req.params.streamId;
|
||||
const type = req.params.type || 'live';
|
||||
const container = req.query.container || 'm3u8';
|
||||
|
||||
const url = api.buildStreamUrl(streamId, type, container);
|
||||
res.json({ url });
|
||||
// Construct the Xtream stream URL
|
||||
// Format: http://server:port/live/username/password/streamId.container (for live)
|
||||
// Format: http://server:port/movie/username/password/streamId.container (for movie)
|
||||
// Format: http://server:port/series/username/password/streamId.container (for series)
|
||||
|
||||
let streamUrl;
|
||||
const baseUrl = source.url.replace(/\/$/, ''); // Remove trailing slash
|
||||
|
||||
if (type === 'live') {
|
||||
streamUrl = `${baseUrl}/live/${source.username}/${source.password}/${streamId}.${container}`;
|
||||
} else if (type === 'movie') {
|
||||
streamUrl = `${baseUrl}/movie/${source.username}/${source.password}/${streamId}.${container}`;
|
||||
} else if (type === 'series') {
|
||||
streamUrl = `${baseUrl}/series/${source.username}/${source.password}/${streamId}.${container}`;
|
||||
} else {
|
||||
return res.status(400).json({ error: 'Invalid stream type' });
|
||||
}
|
||||
|
||||
res.json({ url: streamUrl });
|
||||
} catch (err) {
|
||||
console.error('Stream URL error:', err);
|
||||
res.status(500).json({ error: err.message });
|
||||
console.error('Error getting stream URL:', err);
|
||||
res.status(500).json({ error: 'Failed to get stream URL' });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Fetch and parse M3U playlist
|
||||
* GET /api/proxy/m3u/:sourceId
|
||||
*/
|
||||
|
||||
// --- Other Proxy Routes --- //
|
||||
|
||||
// M3U Playlist
|
||||
// (For M3U sources, we now have data in DB. We can reconstruct M3U or return JSON)
|
||||
// Frontend ChannelList.js for M3U sources calls `API.proxy.m3u.get(sourceId)`
|
||||
// which points here. It expects { channels, groups }.
|
||||
router.get('/m3u/:sourceId', async (req, res) => {
|
||||
try {
|
||||
const sourceId = req.params.sourceId;
|
||||
const source = await sources.getById(sourceId);
|
||||
if (!source || source.type !== 'm3u') {
|
||||
return res.status(404).json({ error: 'M3U source not found' });
|
||||
}
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const includeHidden = req.query.includeHidden === 'true';
|
||||
|
||||
const forceRefresh = req.query.refresh === '1';
|
||||
const maxAgeHours = parseInt(req.query.maxAge) || DEFAULT_MAX_AGE_HOURS;
|
||||
const maxAgeMs = maxAgeHours * 60 * 60 * 1000;
|
||||
// Fetch from DB
|
||||
const channels = getStreamsFromDb(sourceId, 'live', null, includeHidden);
|
||||
const groups = getCategoriesFromDb(sourceId, 'live', includeHidden);
|
||||
|
||||
// Check cache
|
||||
if (!forceRefresh) {
|
||||
const cached = cache.get('m3u', sourceId, 'playlist', maxAgeMs);
|
||||
if (cached) {
|
||||
return res.json(cached);
|
||||
}
|
||||
}
|
||||
// Format for frontend helper
|
||||
// ChannelList expects:
|
||||
// {
|
||||
// channels: [ { id, name, groupTitle, url, tvgLogo, ... } ],
|
||||
// groups: [ { id, name, channelCount } ]
|
||||
// }
|
||||
// Note: DB `live` items from M3U sync have `category_id` as their group name usually.
|
||||
|
||||
const data = await m3uParser.fetchAndParse(source.url);
|
||||
const reformattedChannels = channels.map(c => ({
|
||||
...c,
|
||||
id: c.stream_id,
|
||||
groupTitle: c.category_id || 'Uncategorized',
|
||||
url: c.stream_url || c.url,
|
||||
tvgLogo: c.stream_icon
|
||||
}));
|
||||
|
||||
// Store in cache
|
||||
cache.set('m3u', sourceId, 'playlist', data);
|
||||
const reformattedGroups = groups.map(g => ({
|
||||
id: g.category_id,
|
||||
name: g.category_name,
|
||||
channelCount: 0 // Frontend calculates this or we can
|
||||
}));
|
||||
|
||||
// Add implicit groups check?
|
||||
// The frontend M3U parser generates groups from the channels if explicit groups missing.
|
||||
// Our SyncService `saveCategories` handles explicit groups.
|
||||
|
||||
res.json({ channels: reformattedChannels, groups: reformattedGroups });
|
||||
|
||||
res.json(data);
|
||||
} catch (err) {
|
||||
console.error('M3U proxy error:', err);
|
||||
res.status(500).json({ error: err.message });
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Fetch and parse EPG (with file-based caching)
|
||||
* GET /api/proxy/epg/:sourceId
|
||||
* Query params:
|
||||
* - refresh=1 Force refresh, bypass cache
|
||||
* - maxAge=N Max cache age in hours (default 24)
|
||||
*/
|
||||
// EPG
|
||||
router.get('/epg/:sourceId', async (req, res) => {
|
||||
try {
|
||||
const sourceId = req.params.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' });
|
||||
const sourceId = parseInt(req.params.sourceId);
|
||||
const db = getDb();
|
||||
|
||||
// Return channels + programmes
|
||||
// Frontend EpgGuide.js expects { channels: [], programmes: [] }
|
||||
// It fetches ALL data. We should support timestamps later.
|
||||
|
||||
// EPG Channels (from playlist_items where type='live')
|
||||
// Ideally we only return channels that HAVE epg data?
|
||||
// Or just all channels for this source?
|
||||
// The parser returns `channels` (xmltv channel definitions).
|
||||
// We didn't strictly store XMLTV channel defs in `epg_programs`.
|
||||
// We relied on mapping `channel_id` in `epg_programs`.
|
||||
// Let's return the `playlist_items` as "channels" but mapped to EPG expectations?
|
||||
// EpgGuide joins on `id === tvgId` or `name === name`.
|
||||
|
||||
// Fetch programs
|
||||
|
||||
|
||||
let programsQuery = `SELECT channel_id as channelId, start_time, end_time, title, description, data FROM epg_programs WHERE source_id = ?`;
|
||||
const params = [sourceId];
|
||||
|
||||
// Only valid programs?
|
||||
// programsQuery += ` AND end_time > ?`;
|
||||
// params.push(Date.now() - 86400000); // last 24h
|
||||
|
||||
const programs = db.prepare(programsQuery).all(...params);
|
||||
|
||||
const formattedPrograms = programs.map(p => ({
|
||||
channelId: p.channelId,
|
||||
start: new Date(p.start_time).toISOString(), // EpgGuide parse this back
|
||||
stop: new Date(p.end_time).toISOString(),
|
||||
title: p.title,
|
||||
desc: p.description
|
||||
}));
|
||||
|
||||
// Fetch EPG channels from playlist_items (type='epg_channel')
|
||||
|
||||
|
||||
let epgChannels = [];
|
||||
|
||||
// Try getting stored channels first
|
||||
const storedChannels = db.prepare(`
|
||||
SELECT item_id as id, name, stream_icon as icon, data
|
||||
FROM playlist_items
|
||||
WHERE source_id = ? AND type = 'epg_channel'
|
||||
`).all(sourceId);
|
||||
|
||||
if (storedChannels.length > 0) {
|
||||
epgChannels = storedChannels;
|
||||
} else {
|
||||
// Fallback: Build from unique channelIds in programmes (Legacy behavior)
|
||||
const uniqueChannelIds = [...new Set(programs.map(p => p.channelId))];
|
||||
epgChannels = uniqueChannelIds.map(id => ({
|
||||
id: id,
|
||||
name: id // Use channelId as name (fallback)
|
||||
}));
|
||||
}
|
||||
|
||||
const forceRefresh = req.query.refresh === '1';
|
||||
const maxAgeHours = parseInt(req.query.maxAge) || DEFAULT_MAX_AGE_HOURS;
|
||||
const maxAgeMs = maxAgeHours * 60 * 60 * 1000;
|
||||
res.json({
|
||||
channels: epgChannels,
|
||||
programmes: formattedPrograms
|
||||
});
|
||||
|
||||
// Check file cache (unless force refresh)
|
||||
if (!forceRefresh) {
|
||||
const cached = cache.get('epg', sourceId, 'data', maxAgeMs);
|
||||
if (cached) {
|
||||
return res.json(cached);
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch fresh data
|
||||
let url = source.url;
|
||||
if (source.type === 'xtream') {
|
||||
const api = xtreamApi.createFromSource(source);
|
||||
url = api.getXmltvUrl();
|
||||
}
|
||||
|
||||
const data = await epgParser.fetchAndParse(url);
|
||||
|
||||
// Store in file cache
|
||||
cache.set('epg', sourceId, 'data', data);
|
||||
|
||||
res.json(data);
|
||||
} catch (err) {
|
||||
console.error('EPG proxy error:', err);
|
||||
res.status(500).json({ error: err.message });
|
||||
console.error(err);
|
||||
res.status(500).json({ error: 'Database error' });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Clear cache for a source
|
||||
* DELETE /api/proxy/cache/:sourceId
|
||||
*/
|
||||
// Clear cache (kept for compatibility)
|
||||
router.delete('/cache/:sourceId', (req, res) => {
|
||||
const sourceId = req.params.sourceId;
|
||||
cache.clearSource(sourceId);
|
||||
res.json({ success: true });
|
||||
});
|
||||
|
||||
/**
|
||||
* Clear EPG cache for a source (legacy endpoint, calls clearSource)
|
||||
* DELETE /api/proxy/epg/:sourceId/cache
|
||||
*/
|
||||
router.delete('/epg/:sourceId/cache', (req, res) => {
|
||||
const sourceId = req.params.sourceId;
|
||||
cache.clear('epg', sourceId, 'data');
|
||||
res.json({ success: true });
|
||||
});
|
||||
|
||||
/**
|
||||
* Get EPG for specific channels
|
||||
* POST /api/proxy/epg/:sourceId/channels
|
||||
*/
|
||||
router.post('/epg/:sourceId/channels', async (req, res) => {
|
||||
// --- Stream Proxy (Unchanged mostly) --- //
|
||||
|
||||
// Rewrite M3U8 for proxying
|
||||
async function rewriteM3u8(m3u8Url, baseUrl) {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source || source.type !== 'epg') {
|
||||
return res.status(404).json({ error: 'EPG source not found' });
|
||||
}
|
||||
const response = await fetch(m3u8Url);
|
||||
if (!response.ok) throw new Error(`Fetch failed: ${response.status}`);
|
||||
let content = await response.text();
|
||||
|
||||
const { channelIds } = req.body;
|
||||
if (!channelIds || !Array.isArray(channelIds)) {
|
||||
return res.status(400).json({ error: 'channelIds array required' });
|
||||
}
|
||||
// Resolve relative URLs
|
||||
const m3u8Base = m3u8Url.substring(0, m3u8Url.lastIndexOf('/') + 1);
|
||||
|
||||
const data = await epgParser.fetchAndParse(source.url);
|
||||
|
||||
// Filter programmes for requested channels
|
||||
const result = {};
|
||||
for (const channelId of channelIds) {
|
||||
result[channelId] = epgParser.getCurrentAndUpcoming(data.programmes, channelId);
|
||||
}
|
||||
|
||||
res.json(result);
|
||||
} catch (err) {
|
||||
console.error('EPG channels error:', err);
|
||||
res.status(500).json({ error: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Proxy stream for playback
|
||||
* This handles CORS for streams that don't allow cross-origin
|
||||
* Supports HTTP Range requests for video seeking
|
||||
*/
|
||||
router.get('/stream', async (req, res) => {
|
||||
const maxRetries = 2;
|
||||
let lastError = null;
|
||||
|
||||
for (let attempt = 1; attempt <= maxRetries; attempt++) {
|
||||
try {
|
||||
let { url } = req.query;
|
||||
if (!url) {
|
||||
return res.status(400).json({ error: 'URL required' });
|
||||
}
|
||||
|
||||
// Forward some headers to be more "transparent" back to the origin
|
||||
const isPluto = url.includes('pluto.tv');
|
||||
|
||||
const headers = {
|
||||
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||
'Accept': '*/*',
|
||||
'Accept-Language': 'en-US,en;q=0.9',
|
||||
// Using https and matching the origin of the request
|
||||
'Origin': isPluto ? 'https://pluto.tv' : new URL(url).origin,
|
||||
'Referer': isPluto ? 'https://pluto.tv/' : new URL(url).origin + '/'
|
||||
};
|
||||
|
||||
// Forward Range header for video seeking support
|
||||
const rangeHeader = req.get('range');
|
||||
if (rangeHeader) {
|
||||
headers['Range'] = rangeHeader;
|
||||
}
|
||||
|
||||
const response = await fetch(url, { headers });
|
||||
|
||||
// Retry on 5xx errors (transient upstream issues)
|
||||
if (response.status >= 500 && attempt < maxRetries) {
|
||||
console.log(`[Proxy] Upstream 5xx error (attempt ${attempt}/${maxRetries}), retrying in 500ms...`);
|
||||
await new Promise(r => setTimeout(r, 500));
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!response.ok) {
|
||||
console.error(`Upstream error for ${url.substring(0, 80)}...: ${response.status} ${response.statusText}`);
|
||||
if (response.status === 403) {
|
||||
const errorBody = await response.text().catch(() => 'N/A');
|
||||
console.error(`403 Response body: ${errorBody.substring(0, 200)}`);
|
||||
}
|
||||
return res.status(response.status).send(`Failed to fetch stream: ${response.statusText}`);
|
||||
}
|
||||
|
||||
const contentType = response.headers.get('content-type') || '';
|
||||
res.set('Access-Control-Allow-Origin', '*');
|
||||
|
||||
// Forward range-related headers for video seeking support
|
||||
const contentLength = response.headers.get('content-length');
|
||||
const contentRange = response.headers.get('content-range');
|
||||
const acceptRanges = response.headers.get('accept-ranges');
|
||||
|
||||
if (contentLength) {
|
||||
res.set('Content-Length', contentLength);
|
||||
}
|
||||
if (contentRange) {
|
||||
res.set('Content-Range', contentRange);
|
||||
}
|
||||
if (acceptRanges) {
|
||||
res.set('Accept-Ranges', acceptRanges);
|
||||
} else if (contentLength && !contentRange) {
|
||||
// If server supports content-length but didn't explicitly state accept-ranges,
|
||||
// we can safely assume it supports byte ranges
|
||||
res.set('Accept-Ranges', 'bytes');
|
||||
}
|
||||
|
||||
// Set status code (206 for partial content when range request was made)
|
||||
res.status(response.status);
|
||||
|
||||
// Create an async iterator for the response body
|
||||
const iterator = response.body[Symbol.asyncIterator]();
|
||||
const first = await iterator.next();
|
||||
|
||||
if (first.done) {
|
||||
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 })
|
||||
const textPrefix = firstChunk.subarray(0, 7).toString('utf8');
|
||||
const contentLooksLikeHls = textPrefix === '#EXTM3U';
|
||||
|
||||
if (contentLooksLikeHls) {
|
||||
// HLS Manifest: We must read the WHOLE manifest to rewrite it
|
||||
const chunks = [firstChunk];
|
||||
|
||||
// Consume the rest of the stream
|
||||
let result = await iterator.next();
|
||||
while (!result.done) {
|
||||
chunks.push(Buffer.from(result.value));
|
||||
result = await iterator.next();
|
||||
}
|
||||
|
||||
const buffer = Buffer.concat(chunks);
|
||||
const finalUrl = response.url || url;
|
||||
console.log(`[Proxy] Processing HLS manifest from: ${finalUrl.substring(0, 80)}...`);
|
||||
res.set('Content-Type', 'application/vnd.apple.mpegurl');
|
||||
|
||||
let manifest = buffer.toString('utf-8');
|
||||
|
||||
const finalUrlObj = new URL(finalUrl);
|
||||
const baseUrl = finalUrlObj.origin + finalUrlObj.pathname.substring(0, finalUrlObj.pathname.lastIndexOf('/') + 1);
|
||||
|
||||
manifest = manifest.split('\n').map(line => {
|
||||
const trimmed = line.trim();
|
||||
if (trimmed === '' || trimmed.startsWith('#')) {
|
||||
if (trimmed.includes('URI="')) {
|
||||
return line.replace(/URI="([^"]+)"/g, (match, p1) => {
|
||||
try {
|
||||
const absoluteUrl = new URL(p1, baseUrl).href;
|
||||
return `URI="${req.protocol}://${req.get('host')}${req.baseUrl}/stream?url=${encodeURIComponent(absoluteUrl)}"`;
|
||||
} catch (e) { return match; }
|
||||
});
|
||||
}
|
||||
return line;
|
||||
}
|
||||
|
||||
// Stream URL handling
|
||||
try {
|
||||
let absoluteUrl;
|
||||
if (trimmed.startsWith('http://') || trimmed.startsWith('https://')) {
|
||||
absoluteUrl = trimmed;
|
||||
} else {
|
||||
absoluteUrl = new URL(trimmed, baseUrl).href;
|
||||
}
|
||||
return `${req.protocol}://${req.get('host')}${req.baseUrl}/stream?url=${encodeURIComponent(absoluteUrl)}`;
|
||||
} catch (e) { return line; }
|
||||
}).join('\n');
|
||||
|
||||
return res.send(manifest);
|
||||
}
|
||||
|
||||
// Binary content (Video Segment): Efficient Pipe
|
||||
console.log(`[Proxy] Piping binary stream (${contentType})`);
|
||||
res.set('Content-Type', contentType || 'application/octet-stream');
|
||||
|
||||
// Write the chunk we peeked
|
||||
res.write(firstChunk);
|
||||
|
||||
// Stream the rest
|
||||
// Create a readable stream from the iterator
|
||||
const restOfStream = Readable.from(iterator);
|
||||
|
||||
// Pipe to response
|
||||
restOfStream.pipe(res);
|
||||
return; // Success - exit the retry loop
|
||||
|
||||
} catch (err) {
|
||||
lastError = err;
|
||||
console.error(`Stream proxy error (attempt ${attempt}/${maxRetries}):`, err.message);
|
||||
if (attempt < maxRetries) {
|
||||
console.log('[Proxy] Retrying after error...');
|
||||
await new Promise(r => setTimeout(r, 500));
|
||||
continue;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// All retries failed
|
||||
if (!res.headersSent) {
|
||||
res.status(500).json({ error: lastError?.message || 'Stream proxy failed after retries' });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Proxy images (channel logos, posters)
|
||||
* Fixes mixed content errors when loading HTTP images on HTTPS pages
|
||||
* GET /api/proxy/image?url=...
|
||||
*/
|
||||
router.get('/image', async (req, res) => {
|
||||
try {
|
||||
const { url } = req.query;
|
||||
if (!url) {
|
||||
return res.status(400).json({ error: 'URL required' });
|
||||
}
|
||||
|
||||
const response = await fetch(url, {
|
||||
headers: {
|
||||
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||
'Accept': 'image/*,*/*;q=0.8'
|
||||
content = content.replace(/^(?!#)(.+)$/gm, (match) => {
|
||||
let chunkUrl = match.trim();
|
||||
if (!chunkUrl.startsWith('http')) {
|
||||
chunkUrl = m3u8Base + chunkUrl;
|
||||
}
|
||||
return `${baseUrl}?url=${encodeURIComponent(chunkUrl)}`;
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
return res.status(response.status).send('Failed to fetch image');
|
||||
return content;
|
||||
} catch (e) {
|
||||
console.error('M3U8 Rewrite error:', e);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
router.get('/stream', async (req, res) => {
|
||||
const { url } = req.query;
|
||||
if (!url) return res.status(400).send('URL required');
|
||||
|
||||
try {
|
||||
// Handle M3U8 rewrite
|
||||
if (url.includes('.m3u8')) {
|
||||
const proxyBase = `${req.protocol}://${req.get('host')}/api/proxy/stream`;
|
||||
const manifest = await rewriteM3u8(url, proxyBase);
|
||||
if (manifest) {
|
||||
res.setHeader('Content-Type', 'application/vnd.apple.mpegurl');
|
||||
return res.send(manifest);
|
||||
}
|
||||
}
|
||||
|
||||
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
|
||||
// Native Proxy
|
||||
const range = req.headers.range;
|
||||
const options = {
|
||||
headers: range ? { Range: range } : {}
|
||||
};
|
||||
|
||||
// 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();
|
||||
}
|
||||
// Handle different protocols
|
||||
const lib = url.startsWith('https') ? https : http;
|
||||
|
||||
const proxyReq = lib.get(url, options, (proxyRes) => {
|
||||
// Forward headers
|
||||
res.status(proxyRes.statusCode);
|
||||
for (const [key, value] of Object.entries(proxyRes.headers)) {
|
||||
res.setHeader(key, value);
|
||||
}
|
||||
// Pipe data
|
||||
proxyRes.pipe(res);
|
||||
});
|
||||
|
||||
proxyReq.on('error', (err) => {
|
||||
console.error('Stream proxy error:', err.message);
|
||||
if (!res.headersSent) res.status(502).end();
|
||||
});
|
||||
|
||||
// Handle aborts
|
||||
req.on('close', () => {
|
||||
proxyReq.destroy();
|
||||
});
|
||||
|
||||
} catch (err) {
|
||||
console.error('Image proxy error:', err.message);
|
||||
res.status(500).send('Image proxy error');
|
||||
console.error('Stream handler error:', err);
|
||||
if (!res.headersSent) res.status(500).end();
|
||||
}
|
||||
});
|
||||
|
||||
// Image Proxy
|
||||
router.get('/image', async (req, res) => {
|
||||
const { url } = req.query;
|
||||
if (!url) return res.status(400).send('URL required');
|
||||
|
||||
// Valid check
|
||||
|
||||
if (!url.startsWith('http')) return res.redirect(url); // already local or invalid
|
||||
|
||||
try {
|
||||
const lib = url.startsWith('https') ? https : http;
|
||||
lib.get(url, (proxyRes) => {
|
||||
res.status(proxyRes.statusCode);
|
||||
for (const [key, value] of Object.entries(proxyRes.headers)) {
|
||||
// cors
|
||||
if (key === 'access-control-allow-origin') continue;
|
||||
res.setHeader(key, value);
|
||||
}
|
||||
res.setHeader('Access-Control-Allow-Origin', '*');
|
||||
res.setHeader('Cache-Control', 'public, max-age=86400');
|
||||
proxyRes.pipe(res);
|
||||
}).on('error', err => {
|
||||
res.status(404).end();
|
||||
});
|
||||
} catch (e) {
|
||||
res.status(500).end();
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ const express = require('express');
|
||||
const router = express.Router();
|
||||
const { sources } = require('../db');
|
||||
const xtreamApi = require('../services/xtreamApi');
|
||||
const syncService = require('../services/syncService');
|
||||
|
||||
// Get all sources
|
||||
router.get('/', async (req, res) => {
|
||||
@@ -19,6 +20,19 @@ router.get('/', async (req, res) => {
|
||||
}
|
||||
});
|
||||
|
||||
// Get sync status for all sources
|
||||
router.get('/status', async (req, res) => {
|
||||
try {
|
||||
const { getDb } = require('../db/sqlite');
|
||||
const db = getDb();
|
||||
const statuses = db.prepare('SELECT * FROM sync_status').all();
|
||||
res.json(statuses);
|
||||
} catch (err) {
|
||||
console.error('Error getting sync status:', err);
|
||||
res.status(500).json({ error: 'Failed to get sync status' });
|
||||
}
|
||||
});
|
||||
|
||||
// Get sources by type
|
||||
router.get('/type/:type', async (req, res) => {
|
||||
try {
|
||||
@@ -58,6 +72,8 @@ router.post('/', async (req, res) => {
|
||||
}
|
||||
|
||||
const source = await sources.create({ type, name, url, username, password });
|
||||
// Trigger Sync
|
||||
syncService.syncSource(source.id).catch(console.error);
|
||||
res.status(201).json(source);
|
||||
} catch (err) {
|
||||
console.error('Error creating source:', err);
|
||||
@@ -80,6 +96,8 @@ router.put('/:id', async (req, res) => {
|
||||
username: username !== undefined ? username : existing.username,
|
||||
password: password !== undefined ? password : existing.password
|
||||
});
|
||||
// Trigger Sync (if critical fields changed? safely just trigger it)
|
||||
syncService.syncSource(parseInt(req.params.id)).catch(console.error);
|
||||
res.json(updated);
|
||||
} catch (err) {
|
||||
console.error('Error updating source:', err);
|
||||
@@ -109,6 +127,12 @@ router.post('/:id/toggle', async (req, res) => {
|
||||
if (!updated) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
|
||||
// If enabled, trigger sync
|
||||
if (updated.enabled) {
|
||||
syncService.syncSource(parseInt(req.params.id)).catch(console.error);
|
||||
}
|
||||
|
||||
res.json(updated);
|
||||
} catch (err) {
|
||||
console.error('Error toggling source:', err);
|
||||
@@ -116,6 +140,23 @@ router.post('/:id/toggle', async (req, res) => {
|
||||
}
|
||||
});
|
||||
|
||||
// Manual Sync
|
||||
router.post('/:id/sync', async (req, res) => {
|
||||
try {
|
||||
const id = parseInt(req.params.id);
|
||||
const source = await sources.getById(id);
|
||||
if (!source) return res.status(404).json({ error: 'Source not found' });
|
||||
|
||||
// Trigger sync (async)
|
||||
syncService.syncSource(id).catch(console.error);
|
||||
|
||||
res.json({ success: true, message: 'Sync started' });
|
||||
} catch (err) {
|
||||
console.error('Error starting sync:', err);
|
||||
res.status(500).json({ error: 'Failed to start sync' });
|
||||
}
|
||||
});
|
||||
|
||||
// Test source connection
|
||||
router.post('/:id/test', async (req, res) => {
|
||||
try {
|
||||
|
||||
@@ -25,9 +25,10 @@ router.get('/', (req, res) => {
|
||||
const args = [
|
||||
'-hide_banner',
|
||||
'-loglevel', 'warning',
|
||||
// Low-latency startup: reduce probe/analyze time for faster first bytes
|
||||
'-probesize', '32768',
|
||||
'-analyzeduration', '500000', // 0.5 seconds - enough to detect audio
|
||||
// Low-latency startup vs Reliability trade-off
|
||||
// Increased to 5MB/10s to handle slow HLS manifests and streams with large headers
|
||||
'-probesize', '5000000', // 5MB
|
||||
'-analyzeduration', '10000000', // 10 seconds
|
||||
// Error resilience: discard corrupt packets, generate timestamps, ignore DTS, no buffering
|
||||
'-fflags', '+genpts+discardcorrupt+igndts+nobuffer',
|
||||
// Ignore errors in stream and continue
|
||||
|
||||
@@ -89,19 +89,21 @@ async function parse(input) {
|
||||
let currentInfo = null;
|
||||
let currentGroup = null;
|
||||
|
||||
let inputStream;
|
||||
let lines;
|
||||
|
||||
if (typeof input === 'string') {
|
||||
inputStream = Readable.from([input]);
|
||||
// Handle string input directly
|
||||
lines = input.split(/\r?\n/);
|
||||
} else {
|
||||
inputStream = input;
|
||||
// Handle stream input
|
||||
const rl = readline.createInterface({
|
||||
input: input,
|
||||
crlfDelay: Infinity
|
||||
});
|
||||
lines = rl;
|
||||
}
|
||||
|
||||
const rl = readline.createInterface({
|
||||
input: inputStream,
|
||||
crlfDelay: Infinity
|
||||
});
|
||||
|
||||
for await (const line of rl) {
|
||||
for await (const line of lines) {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) continue;
|
||||
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
const { getDb } = require('../db/sqlite');
|
||||
const { hiddenItems, favorites } = require('../db');
|
||||
|
||||
async function migrateUserData() {
|
||||
console.log('[Migration] Starting user data migration (Hidden Items & Favorites)...');
|
||||
const db = getDb();
|
||||
|
||||
try {
|
||||
// 1. Migrate Hidden Items
|
||||
const allHidden = await hiddenItems.getAll(); // returns array of { source_id, item_type, item_id }
|
||||
|
||||
const hidStmt = db.prepare(`
|
||||
UPDATE playlist_items
|
||||
SET is_hidden = 1
|
||||
WHERE source_id = ? AND type = ? AND item_id = ?
|
||||
`);
|
||||
|
||||
const catHidStmt = db.prepare(`
|
||||
UPDATE categories
|
||||
SET is_hidden = 1
|
||||
WHERE source_id = ? AND type = ? AND category_id = ?
|
||||
`);
|
||||
|
||||
let hidCount = 0;
|
||||
const migrateHidden = db.transaction((items) => {
|
||||
for (const item of items) {
|
||||
// Map db.json types to SQLite types
|
||||
// db.json types: 'channel', 'group', 'vod_category', 'series_category'
|
||||
// SQLite types: 'live', 'movie', 'series' (for categories/items)
|
||||
|
||||
if (item.item_type === 'channel') {
|
||||
// Hidden channel -> playlist_items (type='live')
|
||||
const res = hidStmt.run(item.source_id, 'live', item.item_id);
|
||||
if (res.changes > 0) hidCount++;
|
||||
} else if (item.item_type === 'group') {
|
||||
// Hidden group -> categories (type='live')
|
||||
const res = catHidStmt.run(item.source_id, 'live', item.item_id); // item_id is group name/id
|
||||
if (res.changes > 0) hidCount++;
|
||||
} else if (item.item_type === 'vod_category') {
|
||||
const res = catHidStmt.run(item.source_id, 'movie', item.item_id);
|
||||
if (res.changes > 0) hidCount++;
|
||||
} else if (item.item_type === 'series_category') {
|
||||
const res = catHidStmt.run(item.source_id, 'series', item.item_id);
|
||||
if (res.changes > 0) hidCount++;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
migrateHidden(allHidden);
|
||||
console.log(`[Migration] Migrated ${hidCount} hidden items/categories.`);
|
||||
|
||||
|
||||
// 2. Migrate Favorites
|
||||
const allFavorites = await favorites.getAll();
|
||||
|
||||
const favStmt = db.prepare(`
|
||||
UPDATE playlist_items
|
||||
SET is_favorite = 1
|
||||
WHERE source_id = ? AND type = ? AND item_id = ?
|
||||
`);
|
||||
|
||||
let favCount = 0;
|
||||
const migrateFavorites = db.transaction((favs) => {
|
||||
for (const fav of favs) {
|
||||
// Map types
|
||||
// db.json: 'channel', 'movie', 'series'
|
||||
// SQLite: 'live', 'movie', 'series'
|
||||
|
||||
let type = fav.item_type;
|
||||
if (type === 'channel') type = 'live';
|
||||
|
||||
const res = favStmt.run(fav.source_id, type, fav.item_id);
|
||||
if (res.changes > 0) favCount++;
|
||||
}
|
||||
});
|
||||
|
||||
migrateFavorites(allFavorites);
|
||||
console.log(`[Migration] Migrated ${favCount} favorites.`);
|
||||
|
||||
} catch (e) {
|
||||
console.error('[Migration] Failed:', e);
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { migrateUserData };
|
||||
@@ -0,0 +1,385 @@
|
||||
const { getDb } = require('../db/sqlite');
|
||||
const { sources } = require('../db'); // For source config
|
||||
const xtreamApi = require('./xtreamApi');
|
||||
const m3uParser = require('./m3uParser');
|
||||
const epgParser = require('./epgParser');
|
||||
|
||||
// Sync tracking
|
||||
const activeSyncs = new Set(); // sourceId
|
||||
|
||||
class SyncService {
|
||||
/**
|
||||
* Sync all enabled sources
|
||||
*/
|
||||
async syncAll() {
|
||||
console.log('[Sync] Starting global sync...');
|
||||
try {
|
||||
const allSources = await sources.getAll();
|
||||
for (const source of allSources) {
|
||||
if (source.enabled) {
|
||||
// Run sequentially to not overload
|
||||
await this.syncSource(source.id);
|
||||
}
|
||||
}
|
||||
console.log('[Sync] Global sync completed');
|
||||
} catch (err) {
|
||||
console.error('[Sync] Global sync failed:', err);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Start sync for a source
|
||||
*/
|
||||
async syncSource(sourceId) {
|
||||
if (activeSyncs.has(sourceId)) {
|
||||
console.log(`[Sync] Source ${sourceId} is already syncing`);
|
||||
return;
|
||||
}
|
||||
|
||||
activeSyncs.add(sourceId);
|
||||
|
||||
try {
|
||||
const db = getDb();
|
||||
const source = await sources.getById(sourceId);
|
||||
|
||||
if (!source) {
|
||||
throw new Error(`Source ${sourceId} not found`);
|
||||
}
|
||||
|
||||
console.log(`[Sync] Starting sync for source ${source.name} (ID: ${sourceId})`);
|
||||
|
||||
if (!source.enabled) {
|
||||
console.log(`[Sync] Skipping disabled source ${source.name}`);
|
||||
activeSyncs.delete(sourceId);
|
||||
return;
|
||||
}
|
||||
|
||||
// Update status
|
||||
this.updateSyncStatus(sourceId, 'all', 'syncing');
|
||||
|
||||
if (source.type === 'xtream') {
|
||||
await this.syncXtream(source);
|
||||
} else if (source.type === 'm3u') {
|
||||
await this.syncM3u(source);
|
||||
} else if (source.type === 'epg') {
|
||||
await this.syncEpg(source);
|
||||
}
|
||||
|
||||
this.updateSyncStatus(sourceId, 'all', 'success');
|
||||
console.log(`[Sync] Completed sync for source ${source.name}`);
|
||||
|
||||
} catch (err) {
|
||||
console.error(`[Sync] Failed sync for source ${sourceId}:`, err);
|
||||
this.updateSyncStatus(sourceId, 'all', 'error', err.message);
|
||||
} finally {
|
||||
activeSyncs.delete(sourceId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Update sync status in DB
|
||||
*/
|
||||
updateSyncStatus(sourceId, type, status, error = null) {
|
||||
const db = getDb();
|
||||
const stmt = db.prepare(`
|
||||
INSERT INTO sync_status (source_id, type, last_sync, status, error)
|
||||
VALUES (?, ?, ?, ?, ?)
|
||||
ON CONFLICT(source_id, type) DO UPDATE SET
|
||||
last_sync = excluded.last_sync,
|
||||
status = excluded.status,
|
||||
error = excluded.error
|
||||
`);
|
||||
stmt.run(sourceId, type, Date.now(), status, error);
|
||||
}
|
||||
|
||||
/**
|
||||
* Xtream Sync Logic
|
||||
*/
|
||||
async syncXtream(source) {
|
||||
const api = xtreamApi.createFromSource(source);
|
||||
const db = getDb();
|
||||
|
||||
// 1. Live Categories
|
||||
console.log(`[Sync] Fetching Live Categories for ${source.name}`);
|
||||
const liveCats = await api.getLiveCategories();
|
||||
await this.saveCategories(source.id, 'live', liveCats);
|
||||
|
||||
// 2. Live Streams
|
||||
console.log(`[Sync] Fetching Live Streams for ${source.name}`);
|
||||
const liveStreams = await api.getLiveStreams();
|
||||
await this.saveStreams(source.id, 'live', liveStreams);
|
||||
|
||||
// 3. VOD Categories
|
||||
console.log(`[Sync] Fetching VOD Categories for ${source.name}`);
|
||||
const vodCats = await api.getVodCategories();
|
||||
await this.saveCategories(source.id, 'movie', vodCats);
|
||||
|
||||
// 4. VOD Streams
|
||||
console.log(`[Sync] Fetching VOD Streams for ${source.name}`);
|
||||
const vodStreams = await api.getVodStreams();
|
||||
await this.saveStreams(source.id, 'movie', vodStreams);
|
||||
|
||||
// 5. Series Categories
|
||||
console.log(`[Sync] Fetching Series Categories for ${source.name}`);
|
||||
const seriesCats = await api.getSeriesCategories();
|
||||
await this.saveCategories(source.id, 'series', seriesCats);
|
||||
|
||||
// 6. Series
|
||||
console.log(`[Sync] Fetching Series for ${source.name}`);
|
||||
const series = await api.getSeries();
|
||||
await this.saveStreams(source.id, 'series', series);
|
||||
|
||||
// 7. EPG (Xmltv)
|
||||
// Try to fetch XMLTV if available
|
||||
console.log(`[Sync] Fetching EPG for ${source.name}`);
|
||||
try {
|
||||
const xmltvUrl = api.getXmltvUrl();
|
||||
await this.syncEpgFromUrl(source.id, xmltvUrl);
|
||||
} catch (e) {
|
||||
console.warn('[Sync] XMLTV fetch failed, skipping EPG sync for now:', e.message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Batch save categories
|
||||
*/
|
||||
async saveCategories(sourceId, type, categories) {
|
||||
if (!categories || categories.length === 0) return;
|
||||
console.log(`[Sync] Saving ${categories.length} ${type} categories for source ${sourceId}...`);
|
||||
const db = getDb();
|
||||
const stmt = db.prepare(`
|
||||
INSERT INTO categories (id, source_id, category_id, type, name, parent_id, data)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(id) DO UPDATE SET
|
||||
name = excluded.name,
|
||||
data = excluded.data
|
||||
`);
|
||||
|
||||
const insertBatch = db.transaction((batch) => {
|
||||
for (const cat of batch) {
|
||||
const catId = cat.category_id; // standard xtream field
|
||||
const name = cat.category_name;
|
||||
const id = `${sourceId}:${catId}`;
|
||||
stmt.run(id, sourceId, String(catId), type, name, cat.parent_id || null, JSON.stringify(cat));
|
||||
}
|
||||
});
|
||||
|
||||
const BATCH_SIZE = 500;
|
||||
for (let i = 0; i < categories.length; i += BATCH_SIZE) {
|
||||
insertBatch(categories.slice(i, i + BATCH_SIZE));
|
||||
await new Promise(resolve => setImmediate(resolve));
|
||||
}
|
||||
|
||||
console.log(`[Sync] Saved ${categories.length} ${type} categories`);
|
||||
}
|
||||
|
||||
/**
|
||||
* Batch save streams (channels, vod, series)
|
||||
*/
|
||||
async saveStreams(sourceId, type, items) {
|
||||
if (!items || items.length === 0) return;
|
||||
const db = getDb();
|
||||
const stmt = db.prepare(`
|
||||
INSERT INTO playlist_items (
|
||||
id, source_id, item_id, type, name, category_id,
|
||||
stream_icon, stream_url, container_extension,
|
||||
rating, year, added_at, data
|
||||
)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(id) DO UPDATE SET
|
||||
name = excluded.name,
|
||||
category_id = excluded.category_id,
|
||||
stream_icon = excluded.stream_icon,
|
||||
container_extension = excluded.container_extension,
|
||||
data = excluded.data
|
||||
`);
|
||||
|
||||
const insertBatch = db.transaction((batch) => {
|
||||
for (const item of batch) {
|
||||
// Map fields based on type
|
||||
let itemId, name, catId, icon, container;
|
||||
let rating = null, year = null, added = null;
|
||||
|
||||
if (type === 'live') {
|
||||
itemId = item.stream_id;
|
||||
name = item.name;
|
||||
catId = item.category_id;
|
||||
icon = item.stream_icon;
|
||||
added = item.added;
|
||||
} else if (type === 'movie') {
|
||||
itemId = item.stream_id;
|
||||
name = item.name;
|
||||
catId = item.category_id;
|
||||
icon = item.stream_icon; // or cover
|
||||
container = item.container_extension;
|
||||
rating = item.rating;
|
||||
added = item.added;
|
||||
} else if (type === 'series') {
|
||||
itemId = item.series_id;
|
||||
name = item.name;
|
||||
catId = item.category_id;
|
||||
icon = item.cover;
|
||||
rating = item.rating;
|
||||
year = item.releaseDate;
|
||||
added = item.last_modified;
|
||||
}
|
||||
|
||||
const id = `${sourceId}:${itemId}`;
|
||||
|
||||
stmt.run(
|
||||
id,
|
||||
sourceId,
|
||||
String(itemId),
|
||||
type,
|
||||
name,
|
||||
String(catId),
|
||||
icon,
|
||||
null, // Direct URL not stored for Xtream usually, built on fly
|
||||
container,
|
||||
rating,
|
||||
year,
|
||||
added,
|
||||
JSON.stringify(item)
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
const BATCH_SIZE = 500;
|
||||
for (let i = 0; i < items.length; i += BATCH_SIZE) {
|
||||
insertBatch(items.slice(i, i + BATCH_SIZE));
|
||||
await new Promise(resolve => setImmediate(resolve));
|
||||
}
|
||||
|
||||
console.log(`[Sync] Saved ${items.length} ${type} items`);
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Sync EPG from URL
|
||||
*/
|
||||
async syncEpgFromUrl(sourceId, url) {
|
||||
// Use our streaming parser
|
||||
const { channels, programmes } = await epgParser.fetchAndParse(url);
|
||||
|
||||
console.log(`[Sync] EPG Parsed: ${channels.length} channels, ${programmes.length} programs`);
|
||||
|
||||
const db = getDb();
|
||||
|
||||
// 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
|
||||
db.prepare('DELETE FROM epg_programs WHERE source_id = ?').run(sourceId);
|
||||
|
||||
const stmt = db.prepare(`
|
||||
INSERT INTO epg_programs (channel_id, source_id, start_time, end_time, title, description, data)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||
`);
|
||||
|
||||
const insertMany = db.transaction((progs) => {
|
||||
for (const p of progs) {
|
||||
stmt.run(
|
||||
p.channelId,
|
||||
sourceId,
|
||||
p.start ? p.start.getTime() : 0,
|
||||
p.stop ? p.stop.getTime() : 0,
|
||||
p.title,
|
||||
p.description || p.desc,
|
||||
JSON.stringify(p)
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
insertMany(programmes);
|
||||
console.log(`[Sync] Saved ${programmes.length} programs`);
|
||||
}
|
||||
|
||||
/**
|
||||
* M3U Sync Logic
|
||||
*/
|
||||
/**
|
||||
* M3U Sync Logic
|
||||
*/
|
||||
async syncM3u(source) {
|
||||
console.log(`[Sync] Fetching M3U playlist for ${source.name}`);
|
||||
|
||||
// Fetch and parse using existing parser
|
||||
// We use fetch directly here to get the stream/text
|
||||
const response = await fetch(source.url);
|
||||
if (!response.ok) throw new Error(`Failed to fetch M3U: ${response.status}`);
|
||||
|
||||
const text = await response.text();
|
||||
const { channels, groups } = await m3uParser.parse(text);
|
||||
|
||||
console.log(`[Sync] M3U Parsed: ${channels.length} channels, ${groups.length} groups`);
|
||||
|
||||
// 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,
|
||||
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);
|
||||
}
|
||||
|
||||
/**
|
||||
* EPG Source Sync Logic
|
||||
*/
|
||||
async syncEpg(source) {
|
||||
console.log(`[Sync] Fetching standalone EPG for ${source.name}`);
|
||||
await this.syncEpgFromUrl(source.id, source.url);
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = new SyncService();
|
||||
Reference in New Issue
Block a user