Instancia transcoder t3: snapshot
This commit is contained in:
@@ -0,0 +1,151 @@
|
||||
/**
|
||||
* File-based Cache Service
|
||||
* Stores cached data as JSON files in data/cache/
|
||||
*/
|
||||
|
||||
const fs = require('fs');
|
||||
const path = require('path');
|
||||
|
||||
// Cache directory
|
||||
const cacheDir = path.join(__dirname, '..', '..', 'data', 'cache');
|
||||
|
||||
// Ensure cache directories exist
|
||||
function ensureCacheDir(type, sourceId) {
|
||||
const dir = path.join(cacheDir, type, String(sourceId));
|
||||
if (!fs.existsSync(dir)) {
|
||||
fs.mkdirSync(dir, { recursive: true });
|
||||
}
|
||||
return dir;
|
||||
}
|
||||
|
||||
// Get cache file path
|
||||
function getCachePath(type, sourceId, key) {
|
||||
const dir = ensureCacheDir(type, sourceId);
|
||||
// Sanitize key for filename
|
||||
const safeKey = String(key || 'default').replace(/[^a-zA-Z0-9_-]/g, '_');
|
||||
return path.join(dir, `${safeKey}.json`);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get cached data if not expired
|
||||
* @param {string} type - Cache type (epg, m3u, xtream)
|
||||
* @param {number|string} sourceId - Source ID
|
||||
* @param {string} key - Cache key (e.g., action name)
|
||||
* @param {number} maxAgeMs - Maximum age in milliseconds
|
||||
* @returns {any|null} - Cached data or null if expired/missing
|
||||
*/
|
||||
function get(type, sourceId, key, maxAgeMs) {
|
||||
try {
|
||||
const cachePath = getCachePath(type, sourceId, key);
|
||||
|
||||
if (!fs.existsSync(cachePath)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const cached = JSON.parse(fs.readFileSync(cachePath, 'utf-8'));
|
||||
const age = Date.now() - cached.timestamp;
|
||||
|
||||
if (age > maxAgeMs) {
|
||||
return null; // Expired
|
||||
}
|
||||
|
||||
return cached.data;
|
||||
} catch (err) {
|
||||
console.warn(`Cache read error for ${type}/${sourceId}/${key}:`, err.message);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Store data in cache
|
||||
* @param {string} type - Cache type
|
||||
* @param {number|string} sourceId - Source ID
|
||||
* @param {string} key - Cache key
|
||||
* @param {any} data - Data to cache
|
||||
*/
|
||||
function set(type, sourceId, key, data) {
|
||||
try {
|
||||
const cachePath = getCachePath(type, sourceId, key);
|
||||
const cached = {
|
||||
timestamp: Date.now(),
|
||||
data: data
|
||||
};
|
||||
fs.writeFileSync(cachePath, JSON.stringify(cached));
|
||||
} catch (err) {
|
||||
console.error(`Cache write error for ${type}/${sourceId}/${key}:`, err.message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Clear specific cache entry
|
||||
*/
|
||||
function clear(type, sourceId, key) {
|
||||
try {
|
||||
const cachePath = getCachePath(type, sourceId, key);
|
||||
if (fs.existsSync(cachePath)) {
|
||||
fs.unlinkSync(cachePath);
|
||||
}
|
||||
} catch (err) {
|
||||
console.warn(`Cache clear error:`, err.message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Clear all cache for a source
|
||||
*/
|
||||
function clearSource(sourceId) {
|
||||
try {
|
||||
const types = ['epg', 'm3u', 'xtream'];
|
||||
for (const type of types) {
|
||||
const dir = path.join(cacheDir, type, String(sourceId));
|
||||
if (fs.existsSync(dir)) {
|
||||
fs.rmSync(dir, { recursive: true });
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
console.warn(`Cache clear source error:`, err.message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Clear all cache
|
||||
*/
|
||||
function clearAll() {
|
||||
try {
|
||||
if (fs.existsSync(cacheDir)) {
|
||||
fs.rmSync(cacheDir, { recursive: true });
|
||||
}
|
||||
} catch (err) {
|
||||
console.warn(`Cache clear all error:`, err.message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get cache info for debugging
|
||||
*/
|
||||
function getInfo(type, sourceId, key) {
|
||||
try {
|
||||
const cachePath = getCachePath(type, sourceId, key);
|
||||
if (!fs.existsSync(cachePath)) {
|
||||
return null;
|
||||
}
|
||||
const cached = JSON.parse(fs.readFileSync(cachePath, 'utf-8'));
|
||||
const stats = fs.statSync(cachePath);
|
||||
return {
|
||||
timestamp: cached.timestamp,
|
||||
age: Date.now() - cached.timestamp,
|
||||
size: stats.size
|
||||
};
|
||||
} catch (err) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
get,
|
||||
set,
|
||||
clear,
|
||||
clearSource,
|
||||
clearAll,
|
||||
getInfo
|
||||
};
|
||||
@@ -0,0 +1,448 @@
|
||||
/**
|
||||
* EPG (XMLTV) Parser (Streaming)
|
||||
* Parses XMLTV format EPG data and extracts channel/programme information using streaming XML parser
|
||||
*/
|
||||
|
||||
const sax = require('sax');
|
||||
const zlib = require('zlib');
|
||||
const { Readable } = require('stream');
|
||||
|
||||
/**
|
||||
* Parse XMLTV date format (YYYYMMDDHHmmss +ZZZZ)
|
||||
* @param {string} dateStr - XMLTV format date string
|
||||
* @returns {Date}
|
||||
*/
|
||||
function parseXmltvDate(dateStr) {
|
||||
if (!dateStr) return null;
|
||||
|
||||
// Format: 20231225120000 +0000
|
||||
const match = dateStr.match(/^(\d{4})(\d{2})(\d{2})(\d{2})(\d{2})(\d{2})\s*([+-]\d{4})?$/);
|
||||
if (!match) {
|
||||
// Try ISO format fallback
|
||||
return new Date(dateStr);
|
||||
}
|
||||
|
||||
const [, year, month, day, hour, minute, second, tz] = match;
|
||||
let isoStr = `${year}-${month}-${day}T${hour}:${minute}:${second}`;
|
||||
|
||||
if (tz) {
|
||||
const tzHours = tz.substring(0, 3);
|
||||
const tzMins = tz.substring(3);
|
||||
isoStr += `${tzHours}:${tzMins}`;
|
||||
} else {
|
||||
isoStr += 'Z';
|
||||
}
|
||||
|
||||
return new Date(isoStr);
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse XMLTV content (Stream or String)
|
||||
* @param {Readable|string} input - XMLTV content as Stream or String
|
||||
* @returns {Promise<{ channels: Array, programmes: Array }>}
|
||||
*/
|
||||
function parse(input) {
|
||||
return new Promise((resolve, reject) => {
|
||||
const channels = [];
|
||||
const programmes = [];
|
||||
|
||||
const saxStream = sax.createStream(true, { trim: true, normalize: true }); // strict mode
|
||||
|
||||
let currentTag = null;
|
||||
let currentObject = null;
|
||||
let textBuffer = '';
|
||||
|
||||
saxStream.on('error', function (e) {
|
||||
// clear the error
|
||||
this._parser.error = null;
|
||||
this._parser.resume();
|
||||
console.warn('XML Parse Warning:', e.message);
|
||||
});
|
||||
|
||||
saxStream.on('opentag', function (node) {
|
||||
currentTag = node.name;
|
||||
const attr = node.attributes;
|
||||
|
||||
if (currentTag === 'channel') {
|
||||
currentObject = {
|
||||
id: attr.id,
|
||||
name: null, // 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 = '';
|
||||
});
|
||||
|
||||
saxStream.on('text', function (text) {
|
||||
textBuffer += text;
|
||||
});
|
||||
|
||||
saxStream.on('cdata', function (text) {
|
||||
textBuffer += text;
|
||||
});
|
||||
|
||||
saxStream.on('closetag', function (tagName) {
|
||||
if (tagName === 'channel') {
|
||||
if (currentObject) channels.push(currentObject);
|
||||
currentObject = null;
|
||||
} else if (tagName === 'programme') {
|
||||
if (currentObject) 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) 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;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
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);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Get programmes for a specific channel
|
||||
*/
|
||||
function getProgrammesForChannel(programmes, channelId) {
|
||||
return programmes.filter(p => p.channelId === channelId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get current and upcoming programmes for a channel
|
||||
*/
|
||||
function getCurrentAndUpcoming(programmes, channelId, count = 5) {
|
||||
const now = new Date();
|
||||
const channelProgrammes = getProgrammesForChannel(programmes, channelId);
|
||||
|
||||
// Sort by start time
|
||||
channelProgrammes.sort((a, b) => a.start - b.start);
|
||||
|
||||
// Find current and upcoming
|
||||
const current = channelProgrammes.find(p => p.start <= now && p.stop > now);
|
||||
const upcoming = channelProgrammes
|
||||
.filter(p => p.start > now)
|
||||
.slice(0, count);
|
||||
|
||||
return { current, upcoming };
|
||||
}
|
||||
|
||||
/**
|
||||
* Fetch and parse XMLTV from URL
|
||||
*/
|
||||
async function fetchAndParse(url) {
|
||||
const response = await fetch(url);
|
||||
if (!response.ok) {
|
||||
throw new Error(`Failed to fetch EPG: ${response.status} ${response.statusText}`);
|
||||
}
|
||||
|
||||
let stream;
|
||||
if (response.body && typeof response.body.pipe === 'function') {
|
||||
stream = response.body;
|
||||
} else if (response.body) {
|
||||
stream = Readable.fromWeb(response.body);
|
||||
} else {
|
||||
stream = Readable.from([]);
|
||||
}
|
||||
|
||||
// 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);
|
||||
}
|
||||
|
||||
/**
|
||||
* Streaming EPG parser that yields batches of programmes (memory-efficient)
|
||||
* Channels are collected and returned with the first batch, then programmes are yielded in batches.
|
||||
*
|
||||
* @param {string} url - XMLTV URL
|
||||
* @param {number} batchSize - Number of programmes per batch (default: 1000)
|
||||
* @yields {{ channels: Array|null, programmes: Array, isLast: boolean }}
|
||||
*/
|
||||
async function* fetchAndParseStreaming(url, batchSize = 1000) {
|
||||
const response = await fetch(url);
|
||||
if (!response.ok) {
|
||||
throw new Error(`Failed to fetch EPG: ${response.status} ${response.statusText}`);
|
||||
}
|
||||
|
||||
let stream;
|
||||
if (response.body && typeof response.body.pipe === 'function') {
|
||||
stream = response.body;
|
||||
} else if (response.body) {
|
||||
stream = Readable.fromWeb(response.body);
|
||||
} else {
|
||||
stream = Readable.from([]);
|
||||
}
|
||||
|
||||
const isGzipped = url.endsWith('.gz') || (response.headers.get('content-type') || '').includes('gzip');
|
||||
|
||||
if (isGzipped) {
|
||||
const gunzip = zlib.createGunzip();
|
||||
stream.pipe(gunzip);
|
||||
stream = gunzip;
|
||||
}
|
||||
|
||||
// Use async iterator pattern with SAX
|
||||
yield* parseStreaming(stream, batchSize);
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse XMLTV as streaming async generator
|
||||
* @param {Readable} input - XMLTV stream
|
||||
* @param {number} batchSize - Number of programmes per batch
|
||||
* @yields {{ channels: Array|null, programmes: Array, isLast: boolean }}
|
||||
*/
|
||||
async function* parseStreaming(input, batchSize = 1000) {
|
||||
const channels = [];
|
||||
let programmeBatch = [];
|
||||
let channelsYielded = false;
|
||||
|
||||
// We need to convert SAX events to an async iterator
|
||||
// This requires collecting events and yielding when batch is full
|
||||
|
||||
const saxStream = sax.createStream(true, { trim: true, normalize: true });
|
||||
|
||||
let currentTag = null;
|
||||
let currentObject = null;
|
||||
let textBuffer = '';
|
||||
let resolveNext = null;
|
||||
let pendingBatch = null;
|
||||
let ended = false;
|
||||
let error = null;
|
||||
|
||||
saxStream.on('error', function (e) {
|
||||
this._parser.error = null;
|
||||
this._parser.resume();
|
||||
console.warn('XML Parse Warning:', e.message);
|
||||
});
|
||||
|
||||
saxStream.on('opentag', function (node) {
|
||||
currentTag = node.name;
|
||||
const attr = node.attributes;
|
||||
|
||||
if (currentTag === 'channel') {
|
||||
currentObject = {
|
||||
id: attr.id,
|
||||
name: null,
|
||||
icon: null,
|
||||
url: null
|
||||
};
|
||||
} else if (currentTag === 'programme') {
|
||||
currentObject = {
|
||||
channelId: attr.channel,
|
||||
start: parseXmltvDate(attr.start),
|
||||
stop: parseXmltvDate(attr.stop),
|
||||
title: null,
|
||||
subtitle: null,
|
||||
description: null,
|
||||
category: [],
|
||||
icon: null,
|
||||
date: null,
|
||||
episodeNum: null
|
||||
};
|
||||
} else if (currentTag === 'icon') {
|
||||
if (currentObject) {
|
||||
currentObject.icon = attr.src;
|
||||
}
|
||||
}
|
||||
textBuffer = '';
|
||||
});
|
||||
|
||||
saxStream.on('text', function (text) {
|
||||
textBuffer += text;
|
||||
});
|
||||
|
||||
saxStream.on('cdata', function (text) {
|
||||
textBuffer += text;
|
||||
});
|
||||
|
||||
saxStream.on('closetag', function (tagName) {
|
||||
if (tagName === 'channel') {
|
||||
if (currentObject) channels.push(currentObject);
|
||||
currentObject = null;
|
||||
} else if (tagName === 'programme') {
|
||||
if (currentObject) {
|
||||
programmeBatch.push(currentObject);
|
||||
|
||||
// Check if we should yield a batch
|
||||
if (programmeBatch.length >= batchSize) {
|
||||
const batch = {
|
||||
channels: !channelsYielded ? channels : null,
|
||||
programmes: programmeBatch,
|
||||
isLast: false
|
||||
};
|
||||
channelsYielded = true;
|
||||
programmeBatch = [];
|
||||
|
||||
if (resolveNext) {
|
||||
resolveNext(batch);
|
||||
resolveNext = null;
|
||||
} else {
|
||||
pendingBatch = batch;
|
||||
}
|
||||
}
|
||||
}
|
||||
currentObject = null;
|
||||
} else if (currentObject) {
|
||||
switch (tagName) {
|
||||
case 'display-name':
|
||||
if (!currentObject.name) currentObject.name = textBuffer;
|
||||
break;
|
||||
case 'url':
|
||||
currentObject.url = textBuffer;
|
||||
break;
|
||||
case 'title':
|
||||
currentObject.title = textBuffer;
|
||||
break;
|
||||
case 'sub-title':
|
||||
currentObject.subtitle = textBuffer;
|
||||
break;
|
||||
case 'desc':
|
||||
currentObject.description = textBuffer;
|
||||
break;
|
||||
case 'category':
|
||||
if (textBuffer && currentObject.category) currentObject.category.push(textBuffer);
|
||||
break;
|
||||
case 'date':
|
||||
currentObject.date = textBuffer;
|
||||
break;
|
||||
case 'episode-num':
|
||||
currentObject.episodeNum = textBuffer;
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
saxStream.on('end', function () {
|
||||
ended = true;
|
||||
// Yield final batch
|
||||
const batch = {
|
||||
channels: !channelsYielded ? channels : null,
|
||||
programmes: programmeBatch,
|
||||
isLast: true
|
||||
};
|
||||
if (resolveNext) {
|
||||
resolveNext(batch);
|
||||
resolveNext = null;
|
||||
} else {
|
||||
pendingBatch = batch;
|
||||
}
|
||||
});
|
||||
|
||||
saxStream.on('error', function (e) {
|
||||
error = e;
|
||||
if (resolveNext) {
|
||||
resolveNext(null);
|
||||
}
|
||||
});
|
||||
|
||||
// Start piping
|
||||
input.pipe(saxStream);
|
||||
|
||||
// Yield batches as they become available
|
||||
while (!ended || pendingBatch) {
|
||||
if (pendingBatch) {
|
||||
const batch = pendingBatch;
|
||||
pendingBatch = null;
|
||||
yield batch;
|
||||
if (batch.isLast) break;
|
||||
} else if (!ended) {
|
||||
// Wait for next batch
|
||||
const batch = await new Promise(resolve => {
|
||||
resolveNext = resolve;
|
||||
});
|
||||
if (batch) {
|
||||
yield batch;
|
||||
if (batch.isLast) break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (error) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
parse,
|
||||
parseXmltvDate,
|
||||
fetchAndParse,
|
||||
fetchAndParseStreaming,
|
||||
parseStreaming,
|
||||
getProgrammesForChannel,
|
||||
getCurrentAndUpcoming
|
||||
};
|
||||
|
||||
@@ -0,0 +1,283 @@
|
||||
/**
|
||||
* Hardware Detection Service
|
||||
*
|
||||
* Detects available hardware acceleration capabilities:
|
||||
* - NVIDIA GPU (NVENC/NVDEC via nvidia-smi)
|
||||
* - VAAPI (Linux integrated GPU acceleration)
|
||||
* - QuickSync (Intel GPU acceleration)
|
||||
*
|
||||
* Results are cached at startup and exposed via API.
|
||||
*/
|
||||
|
||||
const { execSync, exec } = require('child_process');
|
||||
const os = require('os');
|
||||
|
||||
// Cache detection results
|
||||
let hwCapabilities = null;
|
||||
|
||||
/**
|
||||
* NVIDIA GPU compute capability requirements for codec support
|
||||
* Source: https://developer.nvidia.com/video-encode-and-decode-gpu-support-matrix-new
|
||||
*/
|
||||
const NVDEC_MIN_COMPUTE = {
|
||||
h264: 3.0, // Kepler+
|
||||
hevc: 5.0, // Maxwell GM206+
|
||||
av1: 8.0, // Ada Lovelace+
|
||||
};
|
||||
|
||||
/**
|
||||
* Detect NVIDIA GPU and its capabilities
|
||||
*/
|
||||
async function detectNvidia() {
|
||||
try {
|
||||
// Query GPU info via nvidia-smi
|
||||
const result = execSync(
|
||||
'nvidia-smi --query-gpu=name,compute_cap --format=csv,noheader',
|
||||
{ timeout: 5000, encoding: 'utf-8', windowsHide: true }
|
||||
);
|
||||
|
||||
const lines = result.trim().split('\n');
|
||||
if (lines.length === 0 || !lines[0]) {
|
||||
return { available: false };
|
||||
}
|
||||
|
||||
// Parse "NVIDIA GeForce RTX 3080, 8.6"
|
||||
const parts = lines[0].split(',');
|
||||
if (parts.length < 2) {
|
||||
return { available: false };
|
||||
}
|
||||
|
||||
const gpuName = parts[0].trim();
|
||||
const computeCap = parseFloat(parts[1].trim());
|
||||
|
||||
// Determine supported codecs based on compute capability
|
||||
const supportedCodecs = [];
|
||||
for (const [codec, minCap] of Object.entries(NVDEC_MIN_COMPUTE)) {
|
||||
if (computeCap >= minCap) {
|
||||
supportedCodecs.push(codec);
|
||||
}
|
||||
}
|
||||
|
||||
console.log(`[HwDetect] NVIDIA GPU detected: ${gpuName} (compute ${computeCap})`);
|
||||
console.log(`[HwDetect] NVDEC supported codecs: ${supportedCodecs.join(', ')}`);
|
||||
|
||||
return {
|
||||
available: true,
|
||||
name: gpuName,
|
||||
computeCap,
|
||||
supportedCodecs,
|
||||
encoder: 'h264_nvenc',
|
||||
decoder: 'h264_cuvid'
|
||||
};
|
||||
} catch (err) {
|
||||
// nvidia-smi not found or no GPU
|
||||
console.log('[HwDetect] No NVIDIA GPU detected');
|
||||
return { available: false };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Detect VAAPI support (Linux only)
|
||||
* Checks for /dev/dri/renderD* devices
|
||||
*/
|
||||
async function detectVAAPI() {
|
||||
if (os.platform() !== 'linux') {
|
||||
return { available: false, reason: 'VAAPI is Linux-only' };
|
||||
}
|
||||
|
||||
try {
|
||||
// Check for render nodes
|
||||
const result = execSync('ls /dev/dri/renderD* 2>/dev/null', {
|
||||
timeout: 2000,
|
||||
encoding: 'utf-8',
|
||||
shell: true
|
||||
});
|
||||
|
||||
const devices = result.trim().split('\n').filter(d => d);
|
||||
if (devices.length === 0) {
|
||||
return { available: false, reason: 'No render devices found' };
|
||||
}
|
||||
|
||||
// Use first available device
|
||||
const device = devices[0];
|
||||
console.log(`[HwDetect] VAAPI device found: ${device}`);
|
||||
|
||||
return {
|
||||
available: true,
|
||||
device,
|
||||
encoder: 'h264_vaapi',
|
||||
decoder: 'vaapi'
|
||||
};
|
||||
} catch (err) {
|
||||
console.log('[HwDetect] VAAPI not available');
|
||||
return { available: false };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Detect Intel QuickSync support
|
||||
* Checks if FFmpeg can use QSV
|
||||
*/
|
||||
async function detectQuickSync() {
|
||||
try {
|
||||
// Check if Intel GPU exists
|
||||
let hasIntelGpu = false;
|
||||
|
||||
if (os.platform() === 'win32') {
|
||||
// Windows: Check via WMIC
|
||||
const result = execSync(
|
||||
'wmic path win32_VideoController get name',
|
||||
{ timeout: 5000, encoding: 'utf-8', windowsHide: true }
|
||||
);
|
||||
hasIntelGpu = result.toLowerCase().includes('intel');
|
||||
} else if (os.platform() === 'linux') {
|
||||
// Linux: Check lspci
|
||||
try {
|
||||
const result = execSync('lspci | grep -i "vga\\|display" | grep -i intel', {
|
||||
timeout: 5000,
|
||||
encoding: 'utf-8',
|
||||
shell: true
|
||||
});
|
||||
hasIntelGpu = result.trim().length > 0;
|
||||
} catch {
|
||||
hasIntelGpu = false;
|
||||
}
|
||||
}
|
||||
|
||||
if (!hasIntelGpu) {
|
||||
return { available: false, reason: 'No Intel GPU found' };
|
||||
}
|
||||
|
||||
console.log('[HwDetect] Intel GPU detected, QSV may be available');
|
||||
|
||||
return {
|
||||
available: true,
|
||||
encoder: 'h264_qsv',
|
||||
decoder: 'h264_qsv'
|
||||
};
|
||||
} catch (err) {
|
||||
console.log('[HwDetect] QuickSync not available');
|
||||
return { available: false };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Detect AMD AMF support (Windows only)
|
||||
* Linux AMD uses VAAPI (detected separately)
|
||||
*/
|
||||
async function detectAMF() {
|
||||
if (os.platform() !== 'win32') {
|
||||
// On Linux, AMD GPUs use VAAPI which is detected separately
|
||||
return { available: false, reason: 'AMF is Windows-only (Linux uses VAAPI)' };
|
||||
}
|
||||
|
||||
try {
|
||||
// Windows: Check via WMIC for AMD/Radeon
|
||||
const result = execSync(
|
||||
'wmic path win32_VideoController get name',
|
||||
{ timeout: 5000, encoding: 'utf-8', windowsHide: true }
|
||||
);
|
||||
|
||||
const lowerResult = result.toLowerCase();
|
||||
const hasAmdGpu = lowerResult.includes('amd') || lowerResult.includes('radeon');
|
||||
|
||||
if (!hasAmdGpu) {
|
||||
return { available: false, reason: 'No AMD GPU found' };
|
||||
}
|
||||
|
||||
// Extract GPU name for display
|
||||
const lines = result.trim().split('\n').filter(l => l.trim());
|
||||
let gpuName = 'AMD GPU';
|
||||
for (const line of lines) {
|
||||
const trimmed = line.trim();
|
||||
if (trimmed.toLowerCase().includes('amd') || trimmed.toLowerCase().includes('radeon')) {
|
||||
gpuName = trimmed;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
console.log(`[HwDetect] AMD GPU detected: ${gpuName}`);
|
||||
|
||||
return {
|
||||
available: true,
|
||||
name: gpuName,
|
||||
encoder: 'h264_amf',
|
||||
decoder: 'h264' // AMF has limited decode support, often uses software
|
||||
};
|
||||
} catch (err) {
|
||||
console.log('[HwDetect] AMF not available');
|
||||
return { available: false };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Detect all hardware capabilities
|
||||
* Results are cached for performance
|
||||
*/
|
||||
async function detect() {
|
||||
if (hwCapabilities !== null) {
|
||||
return hwCapabilities;
|
||||
}
|
||||
|
||||
console.log('[HwDetect] Probing hardware acceleration capabilities...');
|
||||
|
||||
const [nvidia, vaapi, qsv, amf] = await Promise.all([
|
||||
detectNvidia(),
|
||||
detectVAAPI(),
|
||||
detectQuickSync(),
|
||||
detectAMF()
|
||||
]);
|
||||
|
||||
// Determine recommended encoder
|
||||
// On Linux: prefer VAAPI over QSV — VAAPI is more reliable in LXC containers
|
||||
// (QSV may fail to initialize the MFX device even when /dev/dri is passed through)
|
||||
let recommended = 'software';
|
||||
if (nvidia.available) {
|
||||
recommended = 'nvenc';
|
||||
} else if (amf.available) {
|
||||
recommended = 'amf';
|
||||
} else if (vaapi.available) {
|
||||
recommended = 'vaapi';
|
||||
} else if (qsv.available) {
|
||||
recommended = 'qsv';
|
||||
}
|
||||
|
||||
hwCapabilities = {
|
||||
nvidia,
|
||||
amf,
|
||||
vaapi,
|
||||
qsv,
|
||||
recommended,
|
||||
platform: os.platform(),
|
||||
detectedAt: new Date().toISOString()
|
||||
};
|
||||
|
||||
console.log(`[HwDetect] Recommended encoder: ${recommended}`);
|
||||
|
||||
return hwCapabilities;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get cached capabilities (or detect if not cached)
|
||||
*/
|
||||
function getCapabilities() {
|
||||
return hwCapabilities;
|
||||
}
|
||||
|
||||
/**
|
||||
* Force re-detection (clears cache)
|
||||
*/
|
||||
async function refresh() {
|
||||
hwCapabilities = null;
|
||||
return detect();
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
detect,
|
||||
getCapabilities,
|
||||
refresh,
|
||||
detectNvidia,
|
||||
detectAMF,
|
||||
detectVAAPI,
|
||||
detectQuickSync
|
||||
};
|
||||
@@ -0,0 +1,316 @@
|
||||
/**
|
||||
* 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
|
||||
* @param {string} group - Group title
|
||||
* @returns {string} Stable ID
|
||||
*/
|
||||
function generateStableId(name, group) {
|
||||
const str = `${name || 'unknown'}:${group || 'unknown'}`;
|
||||
// Simple hash function
|
||||
let hash = 0;
|
||||
for (let i = 0; i < str.length; i++) {
|
||||
const char = str.charCodeAt(i);
|
||||
hash = ((hash << 5) - hash) + char;
|
||||
hash = hash & hash; // Convert to 32bit integer
|
||||
}
|
||||
return `m3u_${Math.abs(hash).toString(36)}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse EXTINF line and extract attributes
|
||||
* @param {string} line - EXTINF line
|
||||
* @returns {Object} Parsed channel info
|
||||
*/
|
||||
function parseExtinf(line) {
|
||||
const info = {
|
||||
duration: -1,
|
||||
tvgId: null,
|
||||
tvgName: null,
|
||||
tvgLogo: null,
|
||||
groupTitle: null,
|
||||
name: null
|
||||
};
|
||||
|
||||
// Extract duration and rest
|
||||
const match = line.match(/#EXTINF:(-?\d+\.?\d*)\s*(.*)/);
|
||||
if (!match) return info;
|
||||
|
||||
info.duration = parseFloat(match[1]);
|
||||
const rest = match[2];
|
||||
|
||||
// Extract attributes using regex
|
||||
const attrPatterns = {
|
||||
tvgId: /tvg-id="([^"]*)"/i,
|
||||
tvgName: /tvg-name="([^"]*)"/i,
|
||||
tvgLogo: /tvg-logo="([^"]*)"/i,
|
||||
groupTitle: /group-title="([^"]*)"/i
|
||||
};
|
||||
|
||||
for (const [key, pattern] of Object.entries(attrPatterns)) {
|
||||
const attrMatch = rest.match(pattern);
|
||||
if (attrMatch) {
|
||||
info[key] = attrMatch[1];
|
||||
}
|
||||
}
|
||||
|
||||
// Extract channel name (after the comma)
|
||||
const commaIndex = rest.lastIndexOf(',');
|
||||
if (commaIndex !== -1) {
|
||||
info.name = rest.substring(commaIndex + 1).trim();
|
||||
} else {
|
||||
// Fallback: use tvg-name or the whole rest
|
||||
info.name = info.tvgName || rest.trim();
|
||||
}
|
||||
|
||||
// Generate ID if not present
|
||||
if (!info.tvgId) {
|
||||
info.tvgId = info.name ? info.name.toLowerCase().replace(/\s+/g, '_') : `channel_${Date.now()}`;
|
||||
}
|
||||
|
||||
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 lines;
|
||||
|
||||
if (typeof input === 'string') {
|
||||
// Handle string input directly
|
||||
lines = input.split(/\r?\n/);
|
||||
} else {
|
||||
// Handle stream input
|
||||
const rl = readline.createInterface({
|
||||
input: input,
|
||||
crlfDelay: Infinity
|
||||
});
|
||||
lines = rl;
|
||||
}
|
||||
|
||||
for await (const line of lines) {
|
||||
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
|
||||
* @returns {Promise<{ channels: Array, groups: Array }>}
|
||||
*/
|
||||
async function fetchAndParse(url) {
|
||||
const response = await fetch(url);
|
||||
if (!response.ok) {
|
||||
throw new Error(`Failed to fetch M3U: ${response.status} ${response.statusText}`);
|
||||
}
|
||||
|
||||
// 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);
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse M3U content as a streaming async generator (memory-efficient)
|
||||
* Yields batches of channels to avoid loading entire playlist into memory.
|
||||
*
|
||||
* @param {Readable|string} input - M3U content as Stream or String
|
||||
* @param {number} batchSize - Number of channels per batch (default: 500)
|
||||
* @yields {{ channels: Array, groups: Set, isLast: boolean }}
|
||||
*/
|
||||
async function* parseStreaming(input, batchSize = 500) {
|
||||
const groupsSet = new Set();
|
||||
let currentInfo = null;
|
||||
let currentGroup = null;
|
||||
let batch = [];
|
||||
|
||||
let lines;
|
||||
if (typeof input === 'string') {
|
||||
lines = input.split(/\r?\n/);
|
||||
} else {
|
||||
const rl = readline.createInterface({
|
||||
input: input,
|
||||
crlfDelay: Infinity
|
||||
});
|
||||
lines = rl;
|
||||
}
|
||||
|
||||
for await (const line of lines) {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) continue;
|
||||
|
||||
if (trimmed.startsWith('#EXTINF:')) {
|
||||
currentInfo = parseExtinf(trimmed);
|
||||
if (currentInfo.groupTitle) {
|
||||
groupsSet.add(currentInfo.groupTitle);
|
||||
currentGroup = currentInfo.groupTitle;
|
||||
}
|
||||
} else if (trimmed.startsWith('#EXTGRP:')) {
|
||||
currentGroup = trimmed.substring(8).trim();
|
||||
groupsSet.add(currentGroup);
|
||||
if (currentInfo) {
|
||||
currentInfo.groupTitle = currentGroup;
|
||||
}
|
||||
} else if (!trimmed.startsWith('#')) {
|
||||
if (currentInfo) {
|
||||
const groupTitle = currentInfo.groupTitle || currentGroup || 'Uncategorized';
|
||||
const stableId = currentInfo.tvgId || generateStableId(currentInfo.name, groupTitle);
|
||||
|
||||
batch.push({
|
||||
...currentInfo,
|
||||
id: stableId,
|
||||
url: trimmed,
|
||||
groupTitle: groupTitle
|
||||
});
|
||||
currentInfo = null;
|
||||
|
||||
// Yield batch when full
|
||||
if (batch.length >= batchSize) {
|
||||
yield { channels: batch, groups: groupsSet, isLast: false };
|
||||
batch = [];
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Yield remaining channels
|
||||
if (batch.length > 0) {
|
||||
yield { channels: batch, groups: groupsSet, isLast: true };
|
||||
} else {
|
||||
// Yield empty final batch with isLast=true so caller knows we're done
|
||||
yield { channels: [], groups: groupsSet, isLast: true };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fetch and parse M3U from URL as streaming async generator (memory-efficient)
|
||||
* @param {string} url - M3U playlist URL
|
||||
* @param {number} batchSize - Number of channels per batch
|
||||
* @yields {{ channels: Array, groups: Set, isLast: boolean }}
|
||||
*/
|
||||
async function* fetchAndParseStreaming(url, batchSize = 500) {
|
||||
const response = await fetch(url);
|
||||
if (!response.ok) {
|
||||
throw new Error(`Failed to fetch M3U: ${response.status} ${response.statusText}`);
|
||||
}
|
||||
|
||||
let stream;
|
||||
if (response.body && typeof response.body.pipe === 'function') {
|
||||
stream = response.body;
|
||||
} else if (response.body) {
|
||||
stream = Readable.fromWeb(response.body);
|
||||
} else {
|
||||
stream = Readable.from([]);
|
||||
}
|
||||
|
||||
yield* parseStreaming(stream, batchSize);
|
||||
}
|
||||
|
||||
/**
|
||||
* Fast count of entries in an M3U playlist (for size estimation)
|
||||
* Streams the file and counts #EXTINF lines without full parsing
|
||||
* @param {string} url - URL of the M3U playlist
|
||||
* @returns {Promise<number>} Number of entries
|
||||
*/
|
||||
async function countEntries(url) {
|
||||
const response = await fetch(url, {
|
||||
headers: {
|
||||
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
|
||||
}
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
throw new Error(`Failed to fetch playlist: ${response.status}`);
|
||||
}
|
||||
|
||||
let stream;
|
||||
if (response.body && typeof response.body.pipe === 'function') {
|
||||
stream = response.body;
|
||||
} else if (response.body) {
|
||||
stream = Readable.fromWeb(response.body);
|
||||
} else {
|
||||
return 0;
|
||||
}
|
||||
|
||||
return new Promise((resolve, reject) => {
|
||||
let count = 0;
|
||||
const rl = readline.createInterface({ input: stream, crlfDelay: Infinity });
|
||||
|
||||
rl.on('line', (line) => {
|
||||
if (line.startsWith('#EXTINF:')) {
|
||||
count++;
|
||||
}
|
||||
});
|
||||
|
||||
rl.on('close', () => resolve(count));
|
||||
rl.on('error', reject);
|
||||
});
|
||||
}
|
||||
|
||||
module.exports = { parse, parseExtinf, fetchAndParse, parseStreaming, fetchAndParseStreaming, countEntries };
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
/**
|
||||
* M3U Xtream Adapter
|
||||
*
|
||||
* Makes M3U sources respond to Xtream-style API methods.
|
||||
* Queries data from SQLite (already synced during source refresh).
|
||||
*/
|
||||
|
||||
const { getDb } = require('../db/sqlite');
|
||||
|
||||
class M3uXtreamAdapter {
|
||||
constructor(sourceId) {
|
||||
this.sourceId = sourceId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get live categories (groups) for this M3U source.
|
||||
* Returns Xtream-compatible format: [{ category_id, category_name, parent_id }]
|
||||
*/
|
||||
getLiveCategories(includeHidden = false) {
|
||||
const db = getDb();
|
||||
|
||||
// M3U stores group name in category_id field of playlist_items
|
||||
// We need to aggregate unique groups with counts
|
||||
let query = `
|
||||
SELECT
|
||||
category_id,
|
||||
category_id as category_name,
|
||||
NULL as parent_id,
|
||||
COUNT(*) as channel_count
|
||||
FROM playlist_items
|
||||
WHERE source_id = ? AND type = 'live'
|
||||
${!includeHidden ? 'AND is_hidden = 0' : ''}
|
||||
GROUP BY category_id
|
||||
ORDER BY category_id ASC
|
||||
`;
|
||||
|
||||
const rows = db.prepare(query).all(this.sourceId);
|
||||
|
||||
return rows.map(row => ({
|
||||
category_id: row.category_id || 'Uncategorized',
|
||||
category_name: row.category_id || 'Uncategorized',
|
||||
parent_id: null,
|
||||
// Bonus: include count for lazy-loading UI
|
||||
channel_count: row.channel_count
|
||||
}));
|
||||
}
|
||||
|
||||
/**
|
||||
* Get live streams (channels), optionally filtered by category.
|
||||
* Returns Xtream-compatible format: [{ stream_id, name, stream_icon, category_id, ... }]
|
||||
*/
|
||||
getLiveStreams(categoryId = null, includeHidden = false) {
|
||||
const db = getDb();
|
||||
|
||||
let query = `
|
||||
SELECT
|
||||
item_id as stream_id,
|
||||
name,
|
||||
stream_icon,
|
||||
stream_url,
|
||||
category_id,
|
||||
added_at,
|
||||
data
|
||||
FROM playlist_items
|
||||
WHERE source_id = ? AND type = 'live'
|
||||
${!includeHidden ? 'AND is_hidden = 0' : ''}
|
||||
`;
|
||||
|
||||
const params = [this.sourceId];
|
||||
|
||||
if (categoryId) {
|
||||
query += ` AND category_id = ?`;
|
||||
params.push(categoryId);
|
||||
}
|
||||
|
||||
query += ` ORDER BY name ASC`;
|
||||
|
||||
const rows = db.prepare(query).all(...params);
|
||||
|
||||
return rows.map(row => {
|
||||
// Parse any extra data stored as JSON
|
||||
let extra = {};
|
||||
if (row.data) {
|
||||
try { extra = JSON.parse(row.data); } catch (e) { }
|
||||
}
|
||||
|
||||
return {
|
||||
stream_id: row.stream_id,
|
||||
name: row.name,
|
||||
stream_icon: row.stream_icon,
|
||||
category_id: row.category_id,
|
||||
added: row.added_at,
|
||||
// M3U-specific: direct stream URL (Xtream builds URLs from credentials)
|
||||
stream_url: row.stream_url,
|
||||
// Include extra fields from parser (tvgId, etc.)
|
||||
epg_channel_id: extra.tvgId || null,
|
||||
...extra
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Build stream URL for playback.
|
||||
* For M3U, we return the stored URL directly (no credential building).
|
||||
*/
|
||||
buildStreamUrl(streamId, type = 'live', container = 'ts') {
|
||||
const db = getDb();
|
||||
|
||||
const row = db.prepare(`
|
||||
SELECT stream_url FROM playlist_items
|
||||
WHERE source_id = ? AND item_id = ?
|
||||
`).get(this.sourceId, streamId);
|
||||
|
||||
return row?.stream_url || null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get XMLTV EPG URL - M3U sources don't have built-in EPG URLs.
|
||||
* Returns null; EPG must be configured as a separate source.
|
||||
*/
|
||||
getXmltvUrl() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Factory function to create adapter from source ID
|
||||
*/
|
||||
function createFromSourceId(sourceId) {
|
||||
return new M3uXtreamAdapter(sourceId);
|
||||
}
|
||||
|
||||
module.exports = { M3uXtreamAdapter, createFromSourceId };
|
||||
@@ -0,0 +1,781 @@
|
||||
/**
|
||||
* Transcode Session Service
|
||||
*
|
||||
* Manages HLS transcoding sessions with segment caching for VOD seeking.
|
||||
* Each session transcodes a source URL to HLS segments on disk.
|
||||
*
|
||||
* Key features:
|
||||
* - Session-based transcoding with unique IDs
|
||||
* - HLS segment output for seeking support
|
||||
* - Segment caching for fast access
|
||||
* - Session persistence for recovery after restart
|
||||
* - Automatic cleanup of stale sessions
|
||||
*/
|
||||
|
||||
const { spawn } = require('child_process');
|
||||
const path = require('path');
|
||||
const fs = require('fs').promises;
|
||||
const crypto = require('crypto');
|
||||
const EventEmitter = require('events');
|
||||
const hwDetect = require('./hwDetect');
|
||||
|
||||
// Session storage
|
||||
const sessions = new Map();
|
||||
|
||||
// Cache directory for transcoded segments
|
||||
const CACHE_DIR = path.join(process.cwd(), 'transcode-cache');
|
||||
|
||||
// Session settings
|
||||
const SESSION_TIMEOUT_MS = 30 * 60 * 1000; // 30 minutes idle timeout
|
||||
const SEGMENT_DURATION = 4; // seconds per HLS segment
|
||||
const CLEANUP_INTERVAL_MS = 5 * 60 * 1000; // Check every 5 minutes
|
||||
|
||||
/**
|
||||
* Generate a unique session ID
|
||||
*/
|
||||
function generateSessionId() {
|
||||
return crypto.randomBytes(8).toString('hex');
|
||||
}
|
||||
|
||||
/**
|
||||
* Ensure cache directory exists
|
||||
*/
|
||||
async function ensureCacheDir() {
|
||||
try {
|
||||
await fs.mkdir(CACHE_DIR, { recursive: true });
|
||||
} catch (err) {
|
||||
if (err.code !== 'EEXIST') throw err;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* TranscodeSession class
|
||||
* Manages a single transcoding session from source URL to HLS segments
|
||||
*/
|
||||
class TranscodeSession extends EventEmitter {
|
||||
constructor(url, options = {}) {
|
||||
super();
|
||||
this.id = generateSessionId();
|
||||
this.url = url;
|
||||
this.dir = path.join(CACHE_DIR, this.id);
|
||||
this.playlistPath = path.join(this.dir, 'stream.m3u8');
|
||||
this.process = null;
|
||||
this.segments = new Map(); // segment index -> { ready: boolean, path: string }
|
||||
this.status = 'pending'; // pending | starting | running | stopped | error
|
||||
this.error = null;
|
||||
this.startTime = Date.now();
|
||||
this.lastAccess = Date.now();
|
||||
this.options = {
|
||||
ffmpegPath: options.ffmpegPath || 'ffmpeg',
|
||||
userAgent: options.userAgent || 'Mozilla/5.0',
|
||||
seekOffset: options.seekOffset || 0,
|
||||
hwEncoder: options.hwEncoder || 'software',
|
||||
maxResolution: options.maxResolution || '1080p',
|
||||
quality: options.quality || 'medium',
|
||||
// Upscaling options
|
||||
upscaleEnabled: options.upscaleEnabled || false,
|
||||
upscaleMethod: options.upscaleMethod || 'hardware', // 'hardware' or 'software'
|
||||
upscaleTarget: options.upscaleTarget || '1080p',
|
||||
...options
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Start the transcoding process
|
||||
*/
|
||||
async start() {
|
||||
if (this.status === 'running') {
|
||||
return;
|
||||
}
|
||||
|
||||
this.status = 'starting';
|
||||
console.log(`[TranscodeSession ${this.id}] Starting session for: ${this.url}`);
|
||||
|
||||
// Create session directory
|
||||
try {
|
||||
await fs.mkdir(this.dir, { recursive: true });
|
||||
} catch (err) {
|
||||
this.status = 'error';
|
||||
this.error = err.message;
|
||||
throw err;
|
||||
}
|
||||
|
||||
// Build FFmpeg arguments for HLS output
|
||||
const args = this.buildFFmpegArgs();
|
||||
|
||||
console.log(`[TranscodeSession ${this.id}] Command: ${this.options.ffmpegPath} ${args.join(' ')}`);
|
||||
|
||||
try {
|
||||
this.process = spawn(this.options.ffmpegPath, args, {
|
||||
cwd: this.dir,
|
||||
windowsHide: true
|
||||
});
|
||||
|
||||
this.status = 'running';
|
||||
|
||||
// Handle stdout (should be empty for file output)
|
||||
this.process.stdout.on('data', (data) => {
|
||||
console.log(`[TranscodeSession ${this.id}] stdout: ${data}`);
|
||||
});
|
||||
|
||||
// Handle stderr (FFmpeg progress/errors)
|
||||
let stderrBuffer = '';
|
||||
this.process.stderr.on('data', (data) => {
|
||||
stderrBuffer += data.toString();
|
||||
// Log periodically to avoid spam
|
||||
const lines = stderrBuffer.split('\n');
|
||||
if (lines.length > 1) {
|
||||
lines.slice(0, -1).forEach(line => {
|
||||
if (line.trim()) {
|
||||
console.log(`[FFmpeg ${this.id}] ${line}`);
|
||||
}
|
||||
});
|
||||
stderrBuffer = lines[lines.length - 1];
|
||||
}
|
||||
});
|
||||
|
||||
// Handle process exit
|
||||
this.process.on('exit', (code) => {
|
||||
if (code === 0 || code === null) {
|
||||
console.log(`[TranscodeSession ${this.id}] FFmpeg completed successfully`);
|
||||
this.status = 'stopped';
|
||||
} else if (code !== 255) { // 255 is often from SIGKILL
|
||||
console.error(`[TranscodeSession ${this.id}] FFmpeg exited with code ${code}`);
|
||||
this.status = 'error';
|
||||
this.error = `FFmpeg exited with code ${code}`;
|
||||
}
|
||||
this.process = null;
|
||||
this.emit('exit', code);
|
||||
});
|
||||
|
||||
// Handle spawn errors
|
||||
this.process.on('error', (err) => {
|
||||
console.error(`[TranscodeSession ${this.id}] FFmpeg error:`, err);
|
||||
this.status = 'error';
|
||||
this.error = err.message;
|
||||
this.emit('error', err);
|
||||
});
|
||||
|
||||
// Save session metadata
|
||||
await this.persist();
|
||||
|
||||
} catch (err) {
|
||||
this.status = 'error';
|
||||
this.error = err.message;
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Build FFmpeg arguments for HLS output with optional GPU encoding
|
||||
*/
|
||||
buildFFmpegArgs() {
|
||||
const segmentPattern = path.join(this.dir, 'seg%04d.m4s');
|
||||
const videoMode = this.options.videoMode || 'encode';
|
||||
|
||||
// Resolve 'auto' encoder to detected hardware, fallback to software
|
||||
let encoder = this.options.hwEncoder || 'software';
|
||||
if (encoder === 'auto') {
|
||||
const hwCaps = hwDetect.getCapabilities();
|
||||
encoder = hwCaps?.recommended || 'software';
|
||||
console.log(`[TranscodeSession ${this.id}] Auto encoder resolved to: ${encoder}`);
|
||||
}
|
||||
|
||||
const args = [
|
||||
'-hide_banner',
|
||||
'-loglevel', 'warning',
|
||||
'-user_agent', this.options.userAgent,
|
||||
];
|
||||
|
||||
// Add hardware acceleration input options based on encoder (only if encoding)
|
||||
if (videoMode === 'encode') {
|
||||
this.addHwAccelInputArgs(args, encoder);
|
||||
}
|
||||
|
||||
// Input options (common)
|
||||
args.push(
|
||||
'-probesize', '5000000',
|
||||
'-analyzeduration', '5000000',
|
||||
'-fflags', '+genpts+discardcorrupt',
|
||||
'-err_detect', 'ignore_err',
|
||||
'-reconnect', '1',
|
||||
'-reconnect_streamed', '1',
|
||||
'-reconnect_delay_max', '3'
|
||||
);
|
||||
|
||||
args.push('-i', this.url);
|
||||
|
||||
// Add seek offset if specified (as output option to avoid Range requests)
|
||||
if (this.options.seekOffset > 0) {
|
||||
args.push('-ss', String(this.options.seekOffset));
|
||||
}
|
||||
|
||||
// Map streams
|
||||
args.push('-map', '0:v:0');
|
||||
args.push('-map', '0:a:0?');
|
||||
|
||||
// Add video encoder and filters based on selected encoder OR copy
|
||||
if (videoMode === 'copy') {
|
||||
args.push('-c:v', 'copy');
|
||||
|
||||
// Critical for MKV/MP4 -> TS copy: Convert bitstream from AVCC/HVCC to Annex B
|
||||
if (this.options.videoCodec === 'hevc' || this.options.videoCodec === 'h265') {
|
||||
args.push('-bsf:v', 'hevc_mp4toannexb');
|
||||
} else if (this.options.videoCodec === 'h264' || this.options.videoCodec === 'avc') {
|
||||
// Keep Annex-B conversion and force parameter sets to be repeated.
|
||||
// This helps players recover on streams where SPS/PPS appear sporadically,
|
||||
// which often manifests as "audio OK, black video".
|
||||
args.push('-bsf:v', 'h264_mp4toannexb,dump_extra');
|
||||
} else {
|
||||
// Fallback (e.g. unknown codec), try strict extraction
|
||||
args.push('-bsf:v', 'dump_extra');
|
||||
}
|
||||
} else {
|
||||
this.addVideoEncoderArgs(args, encoder);
|
||||
}
|
||||
|
||||
// Audio: Apply mix preset
|
||||
const audioCodec = this.options.audioCodec?.toLowerCase() || 'unknown';
|
||||
const audioChannels = this.options.audioChannels || 0;
|
||||
const audioMixPreset = this.options.audioMixPreset || 'auto';
|
||||
const isStereoAac = audioCodec.includes('aac') && audioChannels === 2;
|
||||
|
||||
// Define pan filter presets for 5.1 -> Stereo downmix
|
||||
const AUDIO_MIX_FILTERS = {
|
||||
// ITU-R BS.775 Standard: Mathematically balanced, transparent
|
||||
itu: 'pan=stereo|FL=FL+0.707*FC+0.707*BL+0.5*LFE|FR=FR+0.707*FC+0.707*BR+0.5*LFE',
|
||||
// Night Mode: Heavy dialogue boost, reduced bass/surrounds for quiet viewing
|
||||
night: 'pan=stereo|FL=0.5*FL+1.2*FC+0.3*BL+0.1*LFE|FR=0.5*FR+1.2*FC+0.3*BR+0.1*LFE',
|
||||
// Cinematic: Wide soundstage, immersive (original "dialogue boost" mix)
|
||||
cinematic: 'pan=stereo|FL=FC+0.80*FL+0.60*BL+0.5*LFE|FR=FC+0.80*FR+0.60*BR+0.5*LFE'
|
||||
};
|
||||
|
||||
if (audioMixPreset === 'passthrough') {
|
||||
// Passthrough: Always copy audio, no processing
|
||||
console.log(`[TranscodeSession ${this.id}] Audio: Passthrough (copy)`);
|
||||
args.push('-c:a', 'copy');
|
||||
} else if (audioMixPreset === 'auto' && isStereoAac) {
|
||||
// Auto + Stereo AAC source: Smart copy
|
||||
console.log(`[TranscodeSession ${this.id}] Audio: Auto (Smart Copy) - Source is Stereo AAC`);
|
||||
args.push('-c:a', 'copy');
|
||||
} else {
|
||||
// Transcode to AAC with selected mix preset (default to ITU for 'auto')
|
||||
const mixPreset = (audioMixPreset === 'auto') ? 'itu' : audioMixPreset;
|
||||
const panFilter = AUDIO_MIX_FILTERS[mixPreset] || AUDIO_MIX_FILTERS.itu;
|
||||
|
||||
console.log(`[TranscodeSession ${this.id}] Audio: ${mixPreset.toUpperCase()} mix (${audioCodec} ${audioChannels}ch -> Stereo AAC)`);
|
||||
args.push(
|
||||
'-c:a', 'aac',
|
||||
'-ar', '48000',
|
||||
'-b:a', '192k',
|
||||
'-af', `${panFilter},aresample=async=1`
|
||||
);
|
||||
}
|
||||
|
||||
// HLS output options
|
||||
args.push(
|
||||
'-f', 'hls',
|
||||
'-hls_time', String(SEGMENT_DURATION),
|
||||
'-hls_list_size', '0', // Keep all segments in playlist
|
||||
'-hls_flags', 'independent_segments+append_list',
|
||||
'-hls_segment_type', 'mpegts',
|
||||
'-hls_segment_filename', path.join(this.dir, 'seg%04d.ts'),
|
||||
this.playlistPath
|
||||
);
|
||||
|
||||
return args;
|
||||
}
|
||||
|
||||
/**
|
||||
* Add hardware acceleration input arguments
|
||||
*/
|
||||
addHwAccelInputArgs(args, encoder) {
|
||||
switch (encoder) {
|
||||
case 'nvenc':
|
||||
// NVIDIA CUDA/NVDEC hardware decoding
|
||||
args.push(
|
||||
'-hwaccel', 'cuda',
|
||||
'-hwaccel_output_format', 'cuda'
|
||||
);
|
||||
break;
|
||||
case 'vaapi':
|
||||
// VAAPI hardware decoding (Linux)
|
||||
args.push(
|
||||
'-hwaccel', 'vaapi',
|
||||
'-hwaccel_device', '/dev/dri/renderD128',
|
||||
'-hwaccel_output_format', 'vaapi'
|
||||
);
|
||||
break;
|
||||
case 'qsv':
|
||||
// Intel QuickSync hardware decoding
|
||||
args.push(
|
||||
'-hwaccel', 'qsv',
|
||||
'-hwaccel_output_format', 'qsv'
|
||||
);
|
||||
break;
|
||||
case 'amf':
|
||||
// AMD AMF (no hwaccel input, AMF is encode-only)
|
||||
// Decode on CPU, encode on GPU
|
||||
break;
|
||||
case 'software':
|
||||
case 'auto':
|
||||
default:
|
||||
// No hardware acceleration for input
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Add video encoder arguments based on selected encoder
|
||||
*/
|
||||
addVideoEncoderArgs(args, encoder) {
|
||||
const resolution = this.getTargetHeight();
|
||||
const quality = this.options.quality || 'medium';
|
||||
|
||||
// Quality presets mapping
|
||||
const qualityPresets = {
|
||||
'high': { nvenc: 18, vaapi: 18, qsv: 18, amf: 18, software: 18 },
|
||||
'medium': { nvenc: 24, vaapi: 24, qsv: 24, amf: 24, software: 23 },
|
||||
'low': { nvenc: 30, vaapi: 30, qsv: 30, amf: 30, software: 28 }
|
||||
};
|
||||
const qp = qualityPresets[quality] || qualityPresets.medium;
|
||||
|
||||
switch (encoder) {
|
||||
case 'nvenc':
|
||||
this.addNvencEncoderArgs(args, resolution, qp.nvenc);
|
||||
break;
|
||||
case 'amf':
|
||||
this.addAmfEncoderArgs(args, resolution, qp.amf);
|
||||
break;
|
||||
case 'vaapi':
|
||||
this.addVaapiEncoderArgs(args, resolution, qp.vaapi);
|
||||
break;
|
||||
case 'qsv':
|
||||
this.addQsvEncoderArgs(args, resolution, qp.qsv);
|
||||
break;
|
||||
case 'software':
|
||||
case 'auto':
|
||||
default:
|
||||
this.addSoftwareEncoderArgs(args, resolution, qp.software);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get target height based on maxResolution or upscaleTarget setting
|
||||
* When upscaling is enabled, uses the upscaleTarget resolution.
|
||||
* Otherwise, uses maxResolution to cap the output.
|
||||
*/
|
||||
getTargetHeight() {
|
||||
const resolutionMap = {
|
||||
'4k': 2160,
|
||||
'1080p': 1080,
|
||||
'720p': 720,
|
||||
'480p': 480
|
||||
};
|
||||
|
||||
// When upscaling is enabled, use the upscale target resolution
|
||||
if (this.options.upscaleEnabled) {
|
||||
const target = resolutionMap[this.options.upscaleTarget] || 1080;
|
||||
console.log(`[TranscodeSession ${this.id}] Upscale target height: ${target}p`);
|
||||
return target;
|
||||
}
|
||||
|
||||
// Otherwise, use max resolution as the cap
|
||||
return resolutionMap[this.options.maxResolution] || 1080;
|
||||
}
|
||||
|
||||
/**
|
||||
* Build scale filter string based on encoder and upscaling settings
|
||||
* @param {string} encoder - The encoder being used
|
||||
* @param {number} height - Target height
|
||||
*/
|
||||
buildScaleFilter(encoder, height) {
|
||||
const useUpscale = this.options.upscaleEnabled;
|
||||
const upscaleMethod = this.options.upscaleMethod || 'hardware';
|
||||
|
||||
// Log upscaling status
|
||||
if (useUpscale) {
|
||||
console.log(`[TranscodeSession ${this.id}] Upscaling: ${upscaleMethod} method to ${height}p`);
|
||||
}
|
||||
|
||||
// Hardware scaling filters (for both upscale and downscale)
|
||||
if (upscaleMethod === 'hardware' || !useUpscale) {
|
||||
switch (encoder) {
|
||||
case 'nvenc':
|
||||
// NVIDIA CUDA scaling with Lanczos
|
||||
// Force nv12 (8-bit) output to handle 10-bit inputs (fixes "10 bit encode not supported")
|
||||
return `scale_cuda=-2:${height}:interp_algo=lanczos:format=nv12`;
|
||||
case 'vaapi':
|
||||
return `scale_vaapi=w=-2:h=${height}:format=nv12`;
|
||||
case 'qsv':
|
||||
return `scale_qsv=w=-2:h=${height}:format=nv12`;
|
||||
case 'amf':
|
||||
// AMF uses CPU decode, so use software scale
|
||||
return useUpscale ? `scale=-2:${height}:flags=lanczos` : `scale=-2:${height}`;
|
||||
case 'software':
|
||||
default:
|
||||
return useUpscale ? `scale=-2:${height}:flags=lanczos` : `scale=-2:${height}`;
|
||||
}
|
||||
}
|
||||
|
||||
// Software Lanczos scaling (high quality, slower)
|
||||
return `scale=-2:${height}:flags=lanczos`;
|
||||
}
|
||||
|
||||
/**
|
||||
* NVIDIA NVENC encoder arguments
|
||||
*/
|
||||
addNvencEncoderArgs(args, height, qp) {
|
||||
// Video filter for scaling on GPU
|
||||
args.push('-vf', this.buildScaleFilter('nvenc', height));
|
||||
|
||||
// NVENC encoder with quality settings
|
||||
// Using portable options that work across FFmpeg builds
|
||||
args.push(
|
||||
'-c:v', 'h264_nvenc',
|
||||
'-preset', 'p4', // Balanced preset (p1=fastest, p7=best)
|
||||
'-rc', 'constqp', // Constant QP mode
|
||||
'-qp', String(qp),
|
||||
'-bf', '3' // B-frames for better compression
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* AMD AMF encoder arguments
|
||||
*/
|
||||
addAmfEncoderArgs(args, height, qp) {
|
||||
// CPU decoding + software scale + AMF encode
|
||||
args.push('-vf', this.buildScaleFilter('amf', height));
|
||||
|
||||
args.push(
|
||||
'-c:v', 'h264_amf',
|
||||
'-quality', 'quality', // Quality preset
|
||||
'-rc', 'cqp', // Constant QP
|
||||
'-qp_i', String(qp),
|
||||
'-qp_p', String(qp + 2),
|
||||
'-qp_b', String(qp + 4),
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* VAAPI encoder arguments (Linux)
|
||||
*/
|
||||
addVaapiEncoderArgs(args, height, qp) {
|
||||
// VAAPI filter chain:
|
||||
// 1. scale_vaapi to resize on GPU
|
||||
// 2. Ensure output format is nv12 for maximum encoder compatibility
|
||||
// The format is handled automatically when using -hwaccel_output_format vaapi
|
||||
args.push('-vf', this.buildScaleFilter('vaapi', height));
|
||||
|
||||
// VAAPI encoder with quality setting
|
||||
// Note: -global_quality is the portable way to set quality for VAAPI
|
||||
args.push(
|
||||
'-c:v', 'h264_vaapi',
|
||||
'-profile:v', 'main', // Use main profile for compatibility
|
||||
'-global_quality', String(qp),
|
||||
'-bf', '3',
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Intel QuickSync encoder arguments
|
||||
*/
|
||||
addQsvEncoderArgs(args, height, qp) {
|
||||
// Scale on QSV
|
||||
args.push('-vf', this.buildScaleFilter('qsv', height));
|
||||
|
||||
args.push(
|
||||
'-c:v', 'h264_qsv',
|
||||
'-preset', 'medium',
|
||||
'-global_quality', String(qp),
|
||||
'-look_ahead', '1',
|
||||
'-look_ahead_depth', '40',
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Software encoder arguments (fallback)
|
||||
*/
|
||||
addSoftwareEncoderArgs(args, height, crf) {
|
||||
// Software scaling (use Lanczos for upscaling if enabled)
|
||||
args.push('-vf', this.buildScaleFilter('software', height));
|
||||
|
||||
args.push(
|
||||
'-c:v', 'libx264',
|
||||
'-preset', 'veryfast', // Fast for real-time
|
||||
'-crf', String(crf),
|
||||
'-profile:v', 'high',
|
||||
'-level', '4.1',
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility (fixes 10-bit input errors)
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop the transcoding process
|
||||
*/
|
||||
stop() {
|
||||
if (this.process) {
|
||||
console.log(`[TranscodeSession ${this.id}] Stopping FFmpeg process`);
|
||||
this.process.kill('SIGTERM');
|
||||
// Force kill after 2 seconds if still running
|
||||
setTimeout(() => {
|
||||
if (this.process) {
|
||||
this.process.kill('SIGKILL');
|
||||
}
|
||||
}, 2000);
|
||||
}
|
||||
this.status = 'stopped';
|
||||
}
|
||||
|
||||
/**
|
||||
* Update last access time (prevents cleanup)
|
||||
*/
|
||||
touch() {
|
||||
this.lastAccess = Date.now();
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if playlist exists and is ready
|
||||
*/
|
||||
async isPlaylistReady() {
|
||||
try {
|
||||
await fs.access(this.playlistPath);
|
||||
const content = await fs.readFile(this.playlistPath, 'utf8');
|
||||
// Check if playlist has at least one segment
|
||||
return content.includes('.ts');
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait for playlist to be ready (with timeout)
|
||||
* Fast-fails if FFmpeg exited with an error before the timeout.
|
||||
*/
|
||||
async waitForPlaylist(timeoutMs = 10000) {
|
||||
const startTime = Date.now();
|
||||
while (Date.now() - startTime < timeoutMs) {
|
||||
if (await this.isPlaylistReady()) {
|
||||
return true;
|
||||
}
|
||||
// Fast-fail: FFmpeg already exited with an error — no point waiting
|
||||
if (this.status === 'error' && !this.process) {
|
||||
console.log(`[TranscodeSession ${this.id}] Fast-fail: FFmpeg exited with error, aborting wait`);
|
||||
return false;
|
||||
}
|
||||
await new Promise(resolve => setTimeout(resolve, 200));
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get the HLS playlist content
|
||||
*/
|
||||
async getPlaylist() {
|
||||
this.touch();
|
||||
try {
|
||||
return await fs.readFile(this.playlistPath, 'utf8');
|
||||
} catch (err) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get a specific segment
|
||||
*/
|
||||
async getSegment(segmentName) {
|
||||
this.touch();
|
||||
const segmentPath = path.join(this.dir, segmentName);
|
||||
try {
|
||||
await fs.access(segmentPath);
|
||||
return segmentPath;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Save session metadata to disk for recovery
|
||||
*/
|
||||
async persist() {
|
||||
const metadata = {
|
||||
id: this.id,
|
||||
url: this.url,
|
||||
status: this.status,
|
||||
startTime: this.startTime,
|
||||
lastAccess: this.lastAccess,
|
||||
options: this.options,
|
||||
seekOffset: this.options.seekOffset
|
||||
};
|
||||
const metaPath = path.join(this.dir, 'session.json');
|
||||
await fs.writeFile(metaPath, JSON.stringify(metadata, null, 2));
|
||||
}
|
||||
|
||||
/**
|
||||
* Restore a session from disk metadata
|
||||
*/
|
||||
static async restore(sessionDir) {
|
||||
const metaPath = path.join(sessionDir, 'session.json');
|
||||
try {
|
||||
const data = await fs.readFile(metaPath, 'utf8');
|
||||
const metadata = JSON.parse(data);
|
||||
const session = new TranscodeSession(metadata.url, metadata.options);
|
||||
session.id = metadata.id;
|
||||
session.dir = sessionDir;
|
||||
session.playlistPath = path.join(sessionDir, 'stream.m3u8');
|
||||
session.startTime = metadata.startTime;
|
||||
session.lastAccess = metadata.lastAccess;
|
||||
session.status = 'stopped'; // Not running after restart
|
||||
return session;
|
||||
} catch (err) {
|
||||
console.error(`Failed to restore session from ${sessionDir}:`, err.message);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete session directory and all segments
|
||||
*/
|
||||
async cleanup() {
|
||||
this.stop();
|
||||
try {
|
||||
await fs.rm(this.dir, { recursive: true, force: true });
|
||||
console.log(`[TranscodeSession ${this.id}] Cleaned up session directory`);
|
||||
} catch (err) {
|
||||
console.error(`[TranscodeSession ${this.id}] Failed to cleanup:`, err.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Session Manager
|
||||
*/
|
||||
|
||||
/**
|
||||
* Create a new transcode session
|
||||
*/
|
||||
async function createSession(url, options = {}) {
|
||||
await ensureCacheDir();
|
||||
const session = new TranscodeSession(url, options);
|
||||
sessions.set(session.id, session);
|
||||
return session;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get an existing session by ID
|
||||
*/
|
||||
function getSession(sessionId) {
|
||||
const session = sessions.get(sessionId);
|
||||
if (session) {
|
||||
session.touch();
|
||||
}
|
||||
return session;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get or create a session for a URL (reuses existing if still valid)
|
||||
*/
|
||||
async function getOrCreateSession(url, options = {}) {
|
||||
// Check for existing session with same URL
|
||||
for (const session of sessions.values()) {
|
||||
if (session.url === url && session.status === 'running') {
|
||||
session.touch();
|
||||
return session;
|
||||
}
|
||||
}
|
||||
// Create new session
|
||||
return createSession(url, options);
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop and remove a session
|
||||
*/
|
||||
async function removeSession(sessionId) {
|
||||
const session = sessions.get(sessionId);
|
||||
if (session) {
|
||||
await session.cleanup();
|
||||
sessions.delete(sessionId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Cleanup stale sessions (idle for too long)
|
||||
*/
|
||||
async function cleanupStaleSessions() {
|
||||
const now = Date.now();
|
||||
for (const [id, session] of sessions) {
|
||||
if (now - session.lastAccess > SESSION_TIMEOUT_MS) {
|
||||
console.log(`[TranscodeSession] Cleaning up stale session ${id}`);
|
||||
await removeSession(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Recover sessions from disk after server restart
|
||||
*/
|
||||
async function recoverSessions() {
|
||||
try {
|
||||
await fs.access(CACHE_DIR);
|
||||
const dirs = await fs.readdir(CACHE_DIR, { withFileTypes: true });
|
||||
|
||||
for (const dirent of dirs) {
|
||||
if (dirent.isDirectory()) {
|
||||
const sessionDir = path.join(CACHE_DIR, dirent.name);
|
||||
const session = await TranscodeSession.restore(sessionDir);
|
||||
if (session) {
|
||||
sessions.set(session.id, session);
|
||||
console.log(`[TranscodeSession] Recovered session ${session.id}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
// Cache dir doesn't exist yet, that's fine
|
||||
if (err.code !== 'ENOENT') {
|
||||
console.error('[TranscodeSession] Error recovering sessions:', err.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Start cleanup interval
|
||||
*/
|
||||
let cleanupInterval = null;
|
||||
function startCleanupInterval() {
|
||||
if (!cleanupInterval) {
|
||||
cleanupInterval = setInterval(cleanupStaleSessions, CLEANUP_INTERVAL_MS);
|
||||
cleanupInterval.unref(); // Don't prevent process exit
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get all active sessions (for debugging/monitoring)
|
||||
*/
|
||||
function getAllSessions() {
|
||||
return Array.from(sessions.values()).map(s => ({
|
||||
id: s.id,
|
||||
url: s.url,
|
||||
status: s.status,
|
||||
startTime: s.startTime,
|
||||
lastAccess: s.lastAccess,
|
||||
idleMs: Date.now() - s.lastAccess
|
||||
}));
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
TranscodeSession,
|
||||
createSession,
|
||||
getSession,
|
||||
getOrCreateSession,
|
||||
removeSession,
|
||||
cleanupStaleSessions,
|
||||
recoverSessions,
|
||||
startCleanupInterval,
|
||||
getAllSessions,
|
||||
CACHE_DIR,
|
||||
SEGMENT_DURATION
|
||||
};
|
||||
@@ -0,0 +1,574 @@
|
||||
const { getDb } = require('../db/sqlite');
|
||||
const { sources, settings } = require('../db'); // For source config and settings
|
||||
const xtreamApi = require('./xtreamApi');
|
||||
const m3uParser = require('./m3uParser');
|
||||
const epgParser = require('./epgParser');
|
||||
|
||||
// Sync tracking
|
||||
const activeSyncs = new Set(); // sourceId
|
||||
|
||||
class SyncService {
|
||||
constructor() {
|
||||
this.lastSyncTime = null; // Track when global sync last completed
|
||||
this._syncTimer = null; // Server-side sync timer
|
||||
this._currentInterval = null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get when the last global sync completed
|
||||
*/
|
||||
getLastSyncTime() {
|
||||
return this.lastSyncTime;
|
||||
}
|
||||
|
||||
/**
|
||||
* Start the server-side sync timer based on settings
|
||||
* Should be called once on server startup after initial sync
|
||||
*/
|
||||
async startSyncTimer() {
|
||||
// Get interval from settings
|
||||
const currentSettings = await settings.get();
|
||||
const intervalHours = parseInt(currentSettings.epgRefreshInterval) || 24;
|
||||
|
||||
// If interval is 0, don't start timer (manual only mode)
|
||||
if (intervalHours <= 0) {
|
||||
console.log('[Sync] Auto-sync disabled (manual only mode)');
|
||||
this.stopSyncTimer();
|
||||
this._currentInterval = 0;
|
||||
return;
|
||||
}
|
||||
|
||||
const intervalMs = intervalHours * 60 * 60 * 1000;
|
||||
|
||||
// Don't restart if interval hasn't changed and timer exists
|
||||
if (this._currentInterval === intervalHours && this._syncTimer) {
|
||||
console.log(`[Sync] Timer already running for ${intervalHours} hours, not restarting`);
|
||||
return;
|
||||
}
|
||||
|
||||
// Clear existing timer
|
||||
this.stopSyncTimer();
|
||||
|
||||
const nextSyncTime = new Date(Date.now() + intervalMs);
|
||||
console.log(`[Sync] Starting server-side sync timer: every ${intervalHours} hours`);
|
||||
console.log(`[Sync] Next scheduled sync at: ${nextSyncTime.toLocaleString()}`);
|
||||
|
||||
this._syncTimer = setInterval(async () => {
|
||||
console.log('[Sync] Scheduled sync triggered');
|
||||
await this.syncAll();
|
||||
// Log next sync time
|
||||
const next = new Date(Date.now() + intervalMs);
|
||||
console.log(`[Sync] Next scheduled sync at: ${next.toLocaleString()}`);
|
||||
}, intervalMs);
|
||||
|
||||
this._currentInterval = intervalHours;
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop the server-side sync timer
|
||||
*/
|
||||
stopSyncTimer() {
|
||||
if (this._syncTimer) {
|
||||
clearInterval(this._syncTimer);
|
||||
this._syncTimer = null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Restart the sync timer with updated settings
|
||||
* Called when sync interval setting changes
|
||||
*/
|
||||
async restartSyncTimer() {
|
||||
await this.startSyncTimer();
|
||||
}
|
||||
|
||||
/**
|
||||
* 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);
|
||||
}
|
||||
}
|
||||
this.lastSyncTime = new Date();
|
||||
console.log('[Sync] Global sync completed at', this.lastSyncTime.toISOString());
|
||||
} 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));
|
||||
}
|
||||
});
|
||||
|
||||
// Reduced batch size for better event loop interleaving
|
||||
const BATCH_SIZE = 100;
|
||||
for (let i = 0; i < categories.length; i += BATCH_SIZE) {
|
||||
insertBatch(categories.slice(i, i + BATCH_SIZE));
|
||||
// Yield to event loop between batches to allow other requests
|
||||
await new Promise(resolve => setImmediate(resolve));
|
||||
}
|
||||
|
||||
console.log(`[Sync] Saved ${categories.length} ${type} categories`);
|
||||
}
|
||||
|
||||
/**
|
||||
* Batch save streams (channels, vod, series)
|
||||
* Also purges stale entries that no longer exist in the source (unless skipPurge is true)
|
||||
* @param {number} sourceId - Source ID
|
||||
* @param {string} type - Type of items (live, movie, series)
|
||||
* @param {Array} items - Items to save
|
||||
* @param {Object} options - Options { skipPurge: boolean }
|
||||
* @returns {Set} Set of synced IDs (for external purge if skipPurge was true)
|
||||
*/
|
||||
async saveStreams(sourceId, type, items, options = {}) {
|
||||
if (!items || items.length === 0) return new Set();
|
||||
const db = getDb();
|
||||
const { skipPurge = false } = options;
|
||||
|
||||
// Collect all IDs we're syncing
|
||||
const syncedIds = new Set();
|
||||
|
||||
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 || `Channel ${item.stream_id}`;
|
||||
catId = item.category_id;
|
||||
icon = item.stream_icon;
|
||||
added = item.added;
|
||||
} else if (type === 'movie') {
|
||||
itemId = item.stream_id;
|
||||
name = item.name || `Movie ${item.stream_id}`;
|
||||
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 || `Series ${item.series_id}`;
|
||||
catId = item.category_id;
|
||||
icon = item.cover;
|
||||
rating = item.rating;
|
||||
year = item.releaseDate;
|
||||
added = item.last_modified;
|
||||
}
|
||||
|
||||
const id = `${sourceId}:${itemId}`;
|
||||
syncedIds.add(id);
|
||||
|
||||
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)
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
// Reduced batch size for better event loop interleaving
|
||||
const BATCH_SIZE = 100;
|
||||
for (let i = 0; i < items.length; i += BATCH_SIZE) {
|
||||
insertBatch(items.slice(i, i + BATCH_SIZE));
|
||||
// Yield to event loop between batches to allow other requests
|
||||
await new Promise(resolve => setImmediate(resolve));
|
||||
}
|
||||
|
||||
// Purge stale entries (skip if doing batch sync like M3U)
|
||||
if (!skipPurge && syncedIds.size > 0) {
|
||||
await this.purgeStaleItems(sourceId, type, syncedIds);
|
||||
}
|
||||
|
||||
console.log(`[Sync] Saved ${items.length} ${type} items`);
|
||||
return syncedIds;
|
||||
}
|
||||
|
||||
/**
|
||||
* Purge stale items that are no longer in the source
|
||||
* @param {number} sourceId - Source ID
|
||||
* @param {string} type - Type of items (live, movie, series)
|
||||
* @param {Set} syncedIds - Set of IDs that should be kept
|
||||
*/
|
||||
async purgeStaleItems(sourceId, type, syncedIds) {
|
||||
if (!syncedIds || syncedIds.size === 0) return;
|
||||
|
||||
const db = getDb();
|
||||
db.exec('CREATE TEMP TABLE IF NOT EXISTS synced_ids (id TEXT PRIMARY KEY)');
|
||||
db.exec('DELETE FROM synced_ids');
|
||||
|
||||
const insertTemp = db.prepare('INSERT OR IGNORE INTO synced_ids (id) VALUES (?)');
|
||||
const insertTempBatch = db.transaction((ids) => {
|
||||
for (const id of ids) {
|
||||
insertTemp.run(id);
|
||||
}
|
||||
});
|
||||
insertTempBatch([...syncedIds]);
|
||||
|
||||
const deleteStmt = db.prepare(`
|
||||
DELETE FROM playlist_items
|
||||
WHERE source_id = ? AND type = ?
|
||||
AND id NOT IN (SELECT id FROM synced_ids)
|
||||
`);
|
||||
const deleted = deleteStmt.run(sourceId, type);
|
||||
|
||||
if (deleted.changes > 0) {
|
||||
console.log(`[Sync] Purged ${deleted.changes} stale ${type} items`);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Sync EPG from URL (Streaming - Memory Efficient)
|
||||
* Processes EPG files in batches to avoid OOM on large EPG data
|
||||
*/
|
||||
async syncEpgFromUrl(sourceId, url) {
|
||||
console.log(`[Sync] Fetching EPG from: ${url.substring(0, 60)}...`);
|
||||
|
||||
// Temporary memory logging for verification
|
||||
const logMemory = () => {
|
||||
const used = process.memoryUsage();
|
||||
console.log(`[Sync] Memory: ${Math.round(used.heapUsed / 1024 / 1024)}MB heap`);
|
||||
};
|
||||
|
||||
logMemory();
|
||||
|
||||
const db = getDb();
|
||||
let allChannels = [];
|
||||
let totalProgrammes = 0;
|
||||
let batchCount = 0;
|
||||
|
||||
// Clear old programmes first
|
||||
db.prepare('DELETE FROM epg_programs WHERE source_id = ?').run(sourceId);
|
||||
|
||||
const programmeStmt = db.prepare(`
|
||||
INSERT INTO epg_programs (channel_id, source_id, start_time, end_time, title, description, data)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||
`);
|
||||
|
||||
const insertProgrammes = db.transaction((progs) => {
|
||||
for (const p of progs) {
|
||||
programmeStmt.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)
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
// Stream and process in batches (default 1000 programmes per batch)
|
||||
for await (const batch of epgParser.fetchAndParseStreaming(url)) {
|
||||
batchCount++;
|
||||
|
||||
// Collect channels from first batch
|
||||
if (batch.channels) {
|
||||
allChannels = batch.channels;
|
||||
}
|
||||
|
||||
// Save this batch of programmes immediately
|
||||
if (batch.programmes.length > 0) {
|
||||
insertProgrammes(batch.programmes);
|
||||
totalProgrammes += batch.programmes.length;
|
||||
}
|
||||
|
||||
// Log progress every 10 batches
|
||||
if (batchCount % 10 === 0) {
|
||||
console.log(`[Sync] Processed ${totalProgrammes} programmes so far...`);
|
||||
logMemory();
|
||||
}
|
||||
|
||||
// Yield to event loop
|
||||
await new Promise(resolve => setImmediate(resolve));
|
||||
}
|
||||
|
||||
console.log(`[Sync] EPG Parsed: ${allChannels.length} channels, ${totalProgrammes} programmes`);
|
||||
logMemory();
|
||||
|
||||
// Save EPG Channels
|
||||
if (allChannels.length > 0) {
|
||||
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
|
||||
`);
|
||||
|
||||
const insertChannels = db.transaction((chanList) => {
|
||||
for (const ch of chanList) {
|
||||
const id = `${sourceId}:${ch.id}`;
|
||||
channelStmt.run(
|
||||
id,
|
||||
sourceId,
|
||||
ch.id,
|
||||
'epg_channel',
|
||||
ch.name,
|
||||
ch.icon || null,
|
||||
null,
|
||||
null,
|
||||
JSON.stringify(ch)
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
insertChannels(allChannels);
|
||||
console.log(`[Sync] Saved ${allChannels.length} EPG channels`);
|
||||
}
|
||||
|
||||
console.log(`[Sync] Saved ${totalProgrammes} programmes`);
|
||||
}
|
||||
|
||||
/**
|
||||
* M3U Sync Logic (Streaming - Memory Efficient)
|
||||
* Processes M3U files in batches to avoid OOM on large playlists
|
||||
*/
|
||||
async syncM3u(source) {
|
||||
console.log(`[Sync] Fetching M3U playlist for ${source.name}`);
|
||||
|
||||
// Temporary memory logging for verification
|
||||
const logMemory = () => {
|
||||
const used = process.memoryUsage();
|
||||
console.log(`[Sync] Memory: ${Math.round(used.heapUsed / 1024 / 1024)}MB heap`);
|
||||
};
|
||||
|
||||
logMemory();
|
||||
|
||||
const allGroups = new Set();
|
||||
const allSyncedIds = new Set(); // Collect IDs across all batches
|
||||
let totalChannels = 0;
|
||||
let batchCount = 0;
|
||||
|
||||
// Stream and process in batches (default 500 channels per batch)
|
||||
for await (const batch of m3uParser.fetchAndParseStreaming(source.url)) {
|
||||
batchCount++;
|
||||
|
||||
// Map M3U channel format to our schema
|
||||
const playlistItems = batch.channels.map(ch => ({
|
||||
stream_id: ch.id,
|
||||
name: ch.name,
|
||||
category_id: ch.groupTitle || 'Uncategorized',
|
||||
stream_icon: ch.tvgLogo,
|
||||
stream_url: ch.url,
|
||||
tvgId: ch.tvgId || null,
|
||||
}));
|
||||
|
||||
// Save this batch immediately (skip purge - we'll do it at the end)
|
||||
if (playlistItems.length > 0) {
|
||||
const batchIds = await this.saveStreams(source.id, 'live', playlistItems, { skipPurge: true });
|
||||
batchIds.forEach(id => allSyncedIds.add(id));
|
||||
totalChannels += playlistItems.length;
|
||||
}
|
||||
|
||||
// Collect groups for category creation at the end
|
||||
batch.groups.forEach(g => allGroups.add(g));
|
||||
|
||||
// Log progress every 10 batches
|
||||
if (batchCount % 10 === 0) {
|
||||
console.log(`[Sync] Processed ${totalChannels} channels so far...`);
|
||||
logMemory();
|
||||
}
|
||||
}
|
||||
|
||||
console.log(`[Sync] M3U Parsed: ${totalChannels} channels, ${allGroups.size} groups`);
|
||||
logMemory();
|
||||
|
||||
// Purge stale items after all batches are complete
|
||||
if (allSyncedIds.size > 0) {
|
||||
await this.purgeStaleItems(source.id, 'live', allSyncedIds);
|
||||
}
|
||||
|
||||
// Save Categories (Groups) at the end
|
||||
const categories = Array.from(allGroups).map(name => ({
|
||||
category_id: name,
|
||||
category_name: name,
|
||||
parent_id: null
|
||||
}));
|
||||
|
||||
await this.saveCategories(source.id, 'live', categories);
|
||||
console.log(`[Sync] M3U sync complete for ${source.name}`);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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();
|
||||
@@ -0,0 +1,787 @@
|
||||
/**
|
||||
* Transcode Session Service
|
||||
*
|
||||
* Manages HLS transcoding sessions with segment caching for VOD seeking.
|
||||
* Each session transcodes a source URL to HLS segments on disk.
|
||||
*
|
||||
* Key features:
|
||||
* - Session-based transcoding with unique IDs
|
||||
* - HLS segment output for seeking support
|
||||
* - Segment caching for fast access
|
||||
* - Session persistence for recovery after restart
|
||||
* - Automatic cleanup of stale sessions
|
||||
*/
|
||||
|
||||
const { spawn } = require('child_process');
|
||||
const path = require('path');
|
||||
const fs = require('fs').promises;
|
||||
const crypto = require('crypto');
|
||||
const EventEmitter = require('events');
|
||||
const hwDetect = require('./hwDetect');
|
||||
const streamInput = require('../utils/streamInput');
|
||||
|
||||
// Session storage
|
||||
const sessions = new Map();
|
||||
|
||||
// Cache directory for transcoded segments
|
||||
const CACHE_DIR = path.join(process.cwd(), 'transcode-cache');
|
||||
|
||||
// Session settings
|
||||
const SESSION_TIMEOUT_MS = 3 * 60 * 1000; // 3 minutes idle timeout (kills orphaned sessions quickly)
|
||||
const SEGMENT_DURATION = 4; // seconds per HLS segment
|
||||
const CLEANUP_INTERVAL_MS = 60 * 1000; // Check every 60 seconds
|
||||
|
||||
/**
|
||||
* Generate a unique session ID
|
||||
*/
|
||||
function generateSessionId() {
|
||||
return crypto.randomBytes(8).toString('hex');
|
||||
}
|
||||
|
||||
/**
|
||||
* Ensure cache directory exists
|
||||
*/
|
||||
async function ensureCacheDir() {
|
||||
try {
|
||||
await fs.mkdir(CACHE_DIR, { recursive: true });
|
||||
} catch (err) {
|
||||
if (err.code !== 'EEXIST') throw err;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* TranscodeSession class
|
||||
* Manages a single transcoding session from source URL to HLS segments
|
||||
*/
|
||||
class TranscodeSession extends EventEmitter {
|
||||
constructor(url, options = {}) {
|
||||
super();
|
||||
this.id = generateSessionId();
|
||||
this.url = url;
|
||||
this.dir = path.join(CACHE_DIR, this.id);
|
||||
this.playlistPath = path.join(this.dir, 'stream.m3u8');
|
||||
this.process = null;
|
||||
this.segments = new Map(); // segment index -> { ready: boolean, path: string }
|
||||
this.status = 'pending'; // pending | starting | running | stopped | error
|
||||
this.error = null;
|
||||
this.startTime = Date.now();
|
||||
this.lastAccess = Date.now();
|
||||
this.options = {
|
||||
ffmpegPath: options.ffmpegPath || 'ffmpeg',
|
||||
userAgent: options.userAgent || 'Mozilla/5.0',
|
||||
seekOffset: options.seekOffset || 0,
|
||||
hwEncoder: options.hwEncoder || 'software',
|
||||
maxResolution: options.maxResolution || '1080p',
|
||||
quality: options.quality || 'medium',
|
||||
// Upscaling options
|
||||
upscaleEnabled: options.upscaleEnabled || false,
|
||||
upscaleMethod: options.upscaleMethod || 'hardware', // 'hardware' or 'software'
|
||||
upscaleTarget: options.upscaleTarget || '1080p',
|
||||
...options
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Start the transcoding process
|
||||
*/
|
||||
async start() {
|
||||
if (this.status === 'running') {
|
||||
return;
|
||||
}
|
||||
|
||||
this.status = 'starting';
|
||||
console.log(`[TranscodeSession ${this.id}] Starting session for: ${this.url}`);
|
||||
|
||||
// Create session directory
|
||||
try {
|
||||
await fs.mkdir(this.dir, { recursive: true });
|
||||
} catch (err) {
|
||||
this.status = 'error';
|
||||
this.error = err.message;
|
||||
throw err;
|
||||
}
|
||||
|
||||
// Build FFmpeg arguments for HLS output
|
||||
const args = this.buildFFmpegArgs();
|
||||
|
||||
console.log(`[TranscodeSession ${this.id}] Command: ${this.options.ffmpegPath} ${args.join(' ')}`);
|
||||
|
||||
try {
|
||||
this.process = spawn(this.options.ffmpegPath, args, {
|
||||
cwd: this.dir,
|
||||
windowsHide: true
|
||||
});
|
||||
|
||||
this.status = 'running';
|
||||
|
||||
// Handle stdout (should be empty for file output)
|
||||
this.process.stdout.on('data', (data) => {
|
||||
console.log(`[TranscodeSession ${this.id}] stdout: ${data}`);
|
||||
});
|
||||
|
||||
// Handle stderr (FFmpeg progress/errors)
|
||||
let stderrBuffer = '';
|
||||
this.process.stderr.on('data', (data) => {
|
||||
stderrBuffer += data.toString();
|
||||
// Log periodically to avoid spam
|
||||
const lines = stderrBuffer.split('\n');
|
||||
if (lines.length > 1) {
|
||||
lines.slice(0, -1).forEach(line => {
|
||||
if (line.trim()) {
|
||||
console.log(`[FFmpeg ${this.id}] ${line}`);
|
||||
}
|
||||
});
|
||||
stderrBuffer = lines[lines.length - 1];
|
||||
}
|
||||
});
|
||||
|
||||
// Handle process exit
|
||||
this.process.on('exit', (code) => {
|
||||
if (code === 0 || code === null) {
|
||||
console.log(`[TranscodeSession ${this.id}] FFmpeg completed successfully`);
|
||||
this.status = 'stopped';
|
||||
} else if (code !== 255) { // 255 is often from SIGKILL
|
||||
console.error(`[TranscodeSession ${this.id}] FFmpeg exited with code ${code}`);
|
||||
this.status = 'error';
|
||||
this.error = `FFmpeg exited with code ${code}`;
|
||||
}
|
||||
this.process = null;
|
||||
this.emit('exit', code);
|
||||
});
|
||||
|
||||
// Handle spawn errors
|
||||
this.process.on('error', (err) => {
|
||||
console.error(`[TranscodeSession ${this.id}] FFmpeg error:`, err);
|
||||
this.status = 'error';
|
||||
this.error = err.message;
|
||||
this.emit('error', err);
|
||||
});
|
||||
|
||||
// Save session metadata
|
||||
await this.persist();
|
||||
|
||||
} catch (err) {
|
||||
this.status = 'error';
|
||||
this.error = err.message;
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Build FFmpeg arguments for HLS output with optional GPU encoding
|
||||
*/
|
||||
buildFFmpegArgs() {
|
||||
const segmentPattern = path.join(this.dir, 'seg%04d.m4s');
|
||||
const videoMode = this.options.videoMode || 'encode';
|
||||
|
||||
// Resolve 'auto' encoder to detected hardware, fallback to software
|
||||
let encoder = this.options.hwEncoder || 'software';
|
||||
if (encoder === 'auto') {
|
||||
const hwCaps = hwDetect.getCapabilities();
|
||||
encoder = hwCaps?.recommended || 'software';
|
||||
console.log(`[TranscodeSession ${this.id}] Auto encoder resolved to: ${encoder}`);
|
||||
}
|
||||
|
||||
const sizing = streamInput.probeSizingForUrl(this.url);
|
||||
const referer = streamInput.deriveRefererFromUrl(this.url);
|
||||
|
||||
const args = [
|
||||
'-hide_banner',
|
||||
'-loglevel', 'warning',
|
||||
'-user_agent', this.options.userAgent,
|
||||
];
|
||||
if (referer) {
|
||||
args.push('-referer', referer);
|
||||
}
|
||||
|
||||
// Add hardware acceleration input options based on encoder (only if encoding)
|
||||
if (videoMode === 'encode') {
|
||||
this.addHwAccelInputArgs(args, encoder);
|
||||
}
|
||||
|
||||
// Input options (common)
|
||||
args.push(
|
||||
'-probesize', sizing.probesize,
|
||||
'-analyzeduration', sizing.analyzeduration,
|
||||
'-fflags', '+genpts+discardcorrupt',
|
||||
'-err_detect', 'ignore_err',
|
||||
'-reconnect', '1',
|
||||
'-reconnect_streamed', '1',
|
||||
'-reconnect_delay_max', '3'
|
||||
);
|
||||
|
||||
args.push('-i', this.url);
|
||||
|
||||
// Add seek offset if specified (as output option to avoid Range requests)
|
||||
if (this.options.seekOffset > 0) {
|
||||
args.push('-ss', String(this.options.seekOffset));
|
||||
}
|
||||
|
||||
// Map streams
|
||||
args.push('-map', '0:v:0');
|
||||
args.push('-map', '0:a:0?');
|
||||
|
||||
// Add video encoder and filters based on selected encoder OR copy
|
||||
if (videoMode === 'copy') {
|
||||
args.push('-c:v', 'copy');
|
||||
|
||||
// Critical for MKV/MP4 -> TS copy: Convert bitstream from AVCC/HVCC to Annex B
|
||||
if (this.options.videoCodec === 'hevc' || this.options.videoCodec === 'h265') {
|
||||
args.push('-bsf:v', 'hevc_mp4toannexb');
|
||||
} else if (this.options.videoCodec === 'h264' || this.options.videoCodec === 'avc') {
|
||||
// Keep Annex-B conversion and force parameter sets to be repeated.
|
||||
// This helps players recover on streams where SPS/PPS appear sporadically,
|
||||
// which often manifests as "audio OK, black video".
|
||||
args.push('-bsf:v', 'h264_mp4toannexb,dump_extra');
|
||||
} else {
|
||||
// Fallback (e.g. unknown codec), try strict extraction
|
||||
args.push('-bsf:v', 'dump_extra');
|
||||
}
|
||||
} else {
|
||||
this.addVideoEncoderArgs(args, encoder);
|
||||
}
|
||||
|
||||
// Audio: Apply mix preset
|
||||
const audioCodec = this.options.audioCodec?.toLowerCase() || 'unknown';
|
||||
const audioChannels = this.options.audioChannels || 0;
|
||||
const audioMixPreset = this.options.audioMixPreset || 'auto';
|
||||
const isStereoAac = audioCodec.includes('aac') && audioChannels === 2;
|
||||
|
||||
// Define pan filter presets for 5.1 -> Stereo downmix
|
||||
const AUDIO_MIX_FILTERS = {
|
||||
// ITU-R BS.775 Standard: Mathematically balanced, transparent
|
||||
itu: 'pan=stereo|FL=FL+0.707*FC+0.707*BL+0.5*LFE|FR=FR+0.707*FC+0.707*BR+0.5*LFE',
|
||||
// Night Mode: Heavy dialogue boost, reduced bass/surrounds for quiet viewing
|
||||
night: 'pan=stereo|FL=0.5*FL+1.2*FC+0.3*BL+0.1*LFE|FR=0.5*FR+1.2*FC+0.3*BR+0.1*LFE',
|
||||
// Cinematic: Wide soundstage, immersive (original "dialogue boost" mix)
|
||||
cinematic: 'pan=stereo|FL=FC+0.80*FL+0.60*BL+0.5*LFE|FR=FC+0.80*FR+0.60*BR+0.5*LFE'
|
||||
};
|
||||
|
||||
if (audioMixPreset === 'passthrough') {
|
||||
// Passthrough: Always copy audio, no processing
|
||||
console.log(`[TranscodeSession ${this.id}] Audio: Passthrough (copy)`);
|
||||
args.push('-c:a', 'copy');
|
||||
} else if (audioMixPreset === 'auto' && isStereoAac) {
|
||||
// Auto + Stereo AAC source: Smart copy
|
||||
console.log(`[TranscodeSession ${this.id}] Audio: Auto (Smart Copy) - Source is Stereo AAC`);
|
||||
args.push('-c:a', 'copy');
|
||||
} else {
|
||||
// Transcode to AAC with selected mix preset (default to ITU for 'auto')
|
||||
const mixPreset = (audioMixPreset === 'auto') ? 'itu' : audioMixPreset;
|
||||
const panFilter = AUDIO_MIX_FILTERS[mixPreset] || AUDIO_MIX_FILTERS.itu;
|
||||
|
||||
console.log(`[TranscodeSession ${this.id}] Audio: ${mixPreset.toUpperCase()} mix (${audioCodec} ${audioChannels}ch -> Stereo AAC)`);
|
||||
args.push(
|
||||
'-c:a', 'aac',
|
||||
'-ar', '48000',
|
||||
'-b:a', '192k',
|
||||
'-af', `${panFilter},aresample=async=1`
|
||||
);
|
||||
}
|
||||
|
||||
// HLS output options
|
||||
args.push(
|
||||
'-f', 'hls',
|
||||
'-hls_time', String(SEGMENT_DURATION),
|
||||
'-hls_list_size', '0', // Keep all segments in playlist
|
||||
'-hls_flags', 'independent_segments+append_list',
|
||||
'-hls_segment_type', 'mpegts',
|
||||
'-hls_segment_filename', path.join(this.dir, 'seg%04d.ts'),
|
||||
this.playlistPath
|
||||
);
|
||||
|
||||
return args;
|
||||
}
|
||||
|
||||
/**
|
||||
* Add hardware acceleration input arguments
|
||||
*/
|
||||
addHwAccelInputArgs(args, encoder) {
|
||||
switch (encoder) {
|
||||
case 'nvenc':
|
||||
// NVIDIA CUDA/NVDEC hardware decoding
|
||||
args.push(
|
||||
'-hwaccel', 'cuda',
|
||||
'-hwaccel_output_format', 'cuda'
|
||||
);
|
||||
break;
|
||||
case 'vaapi':
|
||||
// Solo registramos el dispositivo VAAPI globalmente.
|
||||
// Se usa decodificación SOFTWARE para máxima compatibilidad con
|
||||
// cualquier codec de entrada (HEVC, VP9, etc.) — el GPU Intel
|
||||
// puede no tener soporte de decodificación para todos ellos.
|
||||
// El upload a memoria VAAPI se hace en el filtro de vídeo.
|
||||
args.push('-vaapi_device', '/dev/dri/renderD128');
|
||||
break;
|
||||
case 'qsv':
|
||||
// Intel QuickSync hardware decoding
|
||||
args.push(
|
||||
'-hwaccel', 'qsv',
|
||||
'-hwaccel_output_format', 'qsv'
|
||||
);
|
||||
break;
|
||||
case 'amf':
|
||||
// AMD AMF (no hwaccel input, AMF is encode-only)
|
||||
// Decode on CPU, encode on GPU
|
||||
break;
|
||||
case 'software':
|
||||
case 'auto':
|
||||
default:
|
||||
// No hardware acceleration for input
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Add video encoder arguments based on selected encoder
|
||||
*/
|
||||
addVideoEncoderArgs(args, encoder) {
|
||||
const resolution = this.getTargetHeight();
|
||||
const quality = this.options.quality || 'medium';
|
||||
|
||||
// Quality presets mapping
|
||||
const qualityPresets = {
|
||||
'high': { nvenc: 18, vaapi: 18, qsv: 18, amf: 18, software: 18 },
|
||||
'medium': { nvenc: 24, vaapi: 24, qsv: 24, amf: 24, software: 23 },
|
||||
'low': { nvenc: 30, vaapi: 30, qsv: 30, amf: 30, software: 28 }
|
||||
};
|
||||
const qp = qualityPresets[quality] || qualityPresets.medium;
|
||||
|
||||
switch (encoder) {
|
||||
case 'nvenc':
|
||||
this.addNvencEncoderArgs(args, resolution, qp.nvenc);
|
||||
break;
|
||||
case 'amf':
|
||||
this.addAmfEncoderArgs(args, resolution, qp.amf);
|
||||
break;
|
||||
case 'vaapi':
|
||||
this.addVaapiEncoderArgs(args, resolution, qp.vaapi);
|
||||
break;
|
||||
case 'qsv':
|
||||
this.addQsvEncoderArgs(args, resolution, qp.qsv);
|
||||
break;
|
||||
case 'software':
|
||||
case 'auto':
|
||||
default:
|
||||
this.addSoftwareEncoderArgs(args, resolution, qp.software);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get target height based on maxResolution or upscaleTarget setting
|
||||
* When upscaling is enabled, uses the upscaleTarget resolution.
|
||||
* Otherwise, uses maxResolution to cap the output.
|
||||
*/
|
||||
getTargetHeight() {
|
||||
const resolutionMap = {
|
||||
'4k': 2160,
|
||||
'1080p': 1080,
|
||||
'720p': 720,
|
||||
'480p': 480
|
||||
};
|
||||
|
||||
// When upscaling is enabled, use the upscale target resolution
|
||||
if (this.options.upscaleEnabled) {
|
||||
const target = resolutionMap[this.options.upscaleTarget] || 1080;
|
||||
console.log(`[TranscodeSession ${this.id}] Upscale target height: ${target}p`);
|
||||
return target;
|
||||
}
|
||||
|
||||
// Otherwise, use max resolution as the cap
|
||||
return resolutionMap[this.options.maxResolution] || 1080;
|
||||
}
|
||||
|
||||
/**
|
||||
* Build scale filter string based on encoder and upscaling settings
|
||||
* @param {string} encoder - The encoder being used
|
||||
* @param {number} height - Target height
|
||||
*/
|
||||
buildScaleFilter(encoder, height) {
|
||||
const useUpscale = this.options.upscaleEnabled;
|
||||
const upscaleMethod = this.options.upscaleMethod || 'hardware';
|
||||
|
||||
// Log upscaling status
|
||||
if (useUpscale) {
|
||||
console.log(`[TranscodeSession ${this.id}] Upscaling: ${upscaleMethod} method to ${height}p`);
|
||||
}
|
||||
|
||||
// Hardware scaling filters (for both upscale and downscale)
|
||||
if (upscaleMethod === 'hardware' || !useUpscale) {
|
||||
switch (encoder) {
|
||||
case 'nvenc':
|
||||
// NVIDIA CUDA scaling with Lanczos
|
||||
// Force nv12 (8-bit) output to handle 10-bit inputs (fixes "10 bit encode not supported")
|
||||
return `scale_cuda=-2:${height}:interp_algo=lanczos:format=nv12`;
|
||||
case 'vaapi':
|
||||
// Software decode → scale → nv12 → hwupload → VAAPI encode
|
||||
// format=nv12 ANTES de hwupload es obligatorio para Intel iGPU en LXC
|
||||
return `scale=-2:${height},format=nv12,hwupload`;
|
||||
case 'qsv':
|
||||
return `scale_qsv=w=-2:h=${height}:format=nv12`;
|
||||
case 'amf':
|
||||
// AMF uses CPU decode, so use software scale
|
||||
return useUpscale ? `scale=-2:${height}:flags=lanczos` : `scale=-2:${height}`;
|
||||
case 'software':
|
||||
default:
|
||||
return useUpscale ? `scale=-2:${height}:flags=lanczos` : `scale=-2:${height}`;
|
||||
}
|
||||
}
|
||||
|
||||
// Software Lanczos scaling (high quality, slower)
|
||||
return `scale=-2:${height}:flags=lanczos`;
|
||||
}
|
||||
|
||||
/**
|
||||
* NVIDIA NVENC encoder arguments
|
||||
*/
|
||||
addNvencEncoderArgs(args, height, qp) {
|
||||
// Video filter for scaling on GPU
|
||||
args.push('-vf', this.buildScaleFilter('nvenc', height));
|
||||
|
||||
// NVENC encoder with quality settings
|
||||
// Using portable options that work across FFmpeg builds
|
||||
args.push(
|
||||
'-c:v', 'h264_nvenc',
|
||||
'-preset', 'p4', // Balanced preset (p1=fastest, p7=best)
|
||||
'-rc', 'constqp', // Constant QP mode
|
||||
'-qp', String(qp),
|
||||
'-bf', '3' // B-frames for better compression
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* AMD AMF encoder arguments
|
||||
*/
|
||||
addAmfEncoderArgs(args, height, qp) {
|
||||
// CPU decoding + software scale + AMF encode
|
||||
args.push('-vf', this.buildScaleFilter('amf', height));
|
||||
|
||||
args.push(
|
||||
'-c:v', 'h264_amf',
|
||||
'-quality', 'quality', // Quality preset
|
||||
'-rc', 'cqp', // Constant QP
|
||||
'-qp_i', String(qp),
|
||||
'-qp_p', String(qp + 2),
|
||||
'-qp_b', String(qp + 4),
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* VAAPI encoder arguments (Linux)
|
||||
*/
|
||||
addVaapiEncoderArgs(args, height, qp) {
|
||||
// VAAPI filter chain:
|
||||
// 1. scale_vaapi to resize on GPU
|
||||
// 2. Ensure output format is nv12 for maximum encoder compatibility
|
||||
// The format is handled automatically when using -hwaccel_output_format vaapi
|
||||
args.push('-vf', this.buildScaleFilter('vaapi', height));
|
||||
|
||||
// VAAPI encoder. No especificar -profile:v ni -pix_fmt (los gestiona el driver).
|
||||
// -bf 3 puede causar problemas en algunos drivers Intel; se omite.
|
||||
args.push(
|
||||
'-c:v', 'h264_vaapi',
|
||||
'-global_quality', String(qp)
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Intel QuickSync encoder arguments
|
||||
*/
|
||||
addQsvEncoderArgs(args, height, qp) {
|
||||
// Scale on QSV
|
||||
args.push('-vf', this.buildScaleFilter('qsv', height));
|
||||
|
||||
args.push(
|
||||
'-c:v', 'h264_qsv',
|
||||
'-preset', 'medium',
|
||||
'-global_quality', String(qp),
|
||||
'-look_ahead', '1',
|
||||
'-look_ahead_depth', '40',
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Software encoder arguments (fallback)
|
||||
*/
|
||||
addSoftwareEncoderArgs(args, height, crf) {
|
||||
// Software scaling (use Lanczos for upscaling if enabled)
|
||||
args.push('-vf', this.buildScaleFilter('software', height));
|
||||
|
||||
args.push(
|
||||
'-c:v', 'libx264',
|
||||
'-preset', 'veryfast', // Fast for real-time
|
||||
'-crf', String(crf),
|
||||
'-profile:v', 'high',
|
||||
'-level', '4.1',
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility (fixes 10-bit input errors)
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop the transcoding process
|
||||
*/
|
||||
stop() {
|
||||
if (this.process) {
|
||||
console.log(`[TranscodeSession ${this.id}] Stopping FFmpeg process`);
|
||||
this.process.kill('SIGTERM');
|
||||
// Force kill after 2 seconds if still running
|
||||
setTimeout(() => {
|
||||
if (this.process) {
|
||||
this.process.kill('SIGKILL');
|
||||
}
|
||||
}, 2000);
|
||||
}
|
||||
this.status = 'stopped';
|
||||
}
|
||||
|
||||
/**
|
||||
* Update last access time (prevents cleanup)
|
||||
*/
|
||||
touch() {
|
||||
this.lastAccess = Date.now();
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if playlist exists and is ready
|
||||
*/
|
||||
async isPlaylistReady() {
|
||||
try {
|
||||
await fs.access(this.playlistPath);
|
||||
const content = await fs.readFile(this.playlistPath, 'utf8');
|
||||
// Check if playlist has at least one segment
|
||||
return content.includes('.ts');
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait for playlist to be ready (with timeout)
|
||||
* Fast-fails if FFmpeg exited with an error before the timeout.
|
||||
*/
|
||||
async waitForPlaylist(timeoutMs = 10000) {
|
||||
const startTime = Date.now();
|
||||
while (Date.now() - startTime < timeoutMs) {
|
||||
if (await this.isPlaylistReady()) {
|
||||
return true;
|
||||
}
|
||||
// Fast-fail: FFmpeg already exited with an error — no point waiting
|
||||
if (this.status === 'error' && !this.process) {
|
||||
console.log(`[TranscodeSession ${this.id}] Fast-fail: FFmpeg exited with error, aborting wait`);
|
||||
return false;
|
||||
}
|
||||
await new Promise(resolve => setTimeout(resolve, 200));
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get the HLS playlist content
|
||||
*/
|
||||
async getPlaylist() {
|
||||
this.touch();
|
||||
try {
|
||||
return await fs.readFile(this.playlistPath, 'utf8');
|
||||
} catch (err) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get a specific segment
|
||||
*/
|
||||
async getSegment(segmentName) {
|
||||
this.touch();
|
||||
const segmentPath = path.join(this.dir, segmentName);
|
||||
try {
|
||||
await fs.access(segmentPath);
|
||||
return segmentPath;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Save session metadata to disk for recovery
|
||||
*/
|
||||
async persist() {
|
||||
const metadata = {
|
||||
id: this.id,
|
||||
url: this.url,
|
||||
status: this.status,
|
||||
startTime: this.startTime,
|
||||
lastAccess: this.lastAccess,
|
||||
options: this.options,
|
||||
seekOffset: this.options.seekOffset
|
||||
};
|
||||
const metaPath = path.join(this.dir, 'session.json');
|
||||
await fs.writeFile(metaPath, JSON.stringify(metadata, null, 2));
|
||||
}
|
||||
|
||||
/**
|
||||
* Restore a session from disk metadata
|
||||
*/
|
||||
static async restore(sessionDir) {
|
||||
const metaPath = path.join(sessionDir, 'session.json');
|
||||
try {
|
||||
const data = await fs.readFile(metaPath, 'utf8');
|
||||
const metadata = JSON.parse(data);
|
||||
const session = new TranscodeSession(metadata.url, metadata.options);
|
||||
session.id = metadata.id;
|
||||
session.dir = sessionDir;
|
||||
session.playlistPath = path.join(sessionDir, 'stream.m3u8');
|
||||
session.startTime = metadata.startTime;
|
||||
session.lastAccess = metadata.lastAccess;
|
||||
session.status = 'stopped'; // Not running after restart
|
||||
return session;
|
||||
} catch (err) {
|
||||
console.error(`Failed to restore session from ${sessionDir}:`, err.message);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete session directory and all segments
|
||||
*/
|
||||
async cleanup() {
|
||||
this.stop();
|
||||
try {
|
||||
await fs.rm(this.dir, { recursive: true, force: true });
|
||||
console.log(`[TranscodeSession ${this.id}] Cleaned up session directory`);
|
||||
} catch (err) {
|
||||
console.error(`[TranscodeSession ${this.id}] Failed to cleanup:`, err.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Session Manager
|
||||
*/
|
||||
|
||||
/**
|
||||
* Create a new transcode session
|
||||
*/
|
||||
async function createSession(url, options = {}) {
|
||||
await ensureCacheDir();
|
||||
const session = new TranscodeSession(url, options);
|
||||
sessions.set(session.id, session);
|
||||
return session;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get an existing session by ID
|
||||
*/
|
||||
function getSession(sessionId) {
|
||||
const session = sessions.get(sessionId);
|
||||
if (session) {
|
||||
session.touch();
|
||||
}
|
||||
return session;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get or create a session for a URL (reuses existing if still valid)
|
||||
*/
|
||||
async function getOrCreateSession(url, options = {}) {
|
||||
// Check for existing session with same URL
|
||||
for (const session of sessions.values()) {
|
||||
if (session.url === url && session.status === 'running') {
|
||||
session.touch();
|
||||
return session;
|
||||
}
|
||||
}
|
||||
// Create new session
|
||||
return createSession(url, options);
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop and remove a session
|
||||
*/
|
||||
async function removeSession(sessionId) {
|
||||
const session = sessions.get(sessionId);
|
||||
if (session) {
|
||||
await session.cleanup();
|
||||
sessions.delete(sessionId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Cleanup stale sessions (idle for too long)
|
||||
*/
|
||||
async function cleanupStaleSessions() {
|
||||
const now = Date.now();
|
||||
for (const [id, session] of sessions) {
|
||||
if (now - session.lastAccess > SESSION_TIMEOUT_MS) {
|
||||
console.log(`[TranscodeSession] Cleaning up stale session ${id}`);
|
||||
await removeSession(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Recover sessions from disk after server restart
|
||||
*/
|
||||
async function recoverSessions() {
|
||||
try {
|
||||
await fs.access(CACHE_DIR);
|
||||
const dirs = await fs.readdir(CACHE_DIR, { withFileTypes: true });
|
||||
|
||||
for (const dirent of dirs) {
|
||||
if (dirent.isDirectory()) {
|
||||
const sessionDir = path.join(CACHE_DIR, dirent.name);
|
||||
const session = await TranscodeSession.restore(sessionDir);
|
||||
if (session) {
|
||||
sessions.set(session.id, session);
|
||||
console.log(`[TranscodeSession] Recovered session ${session.id}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
// Cache dir doesn't exist yet, that's fine
|
||||
if (err.code !== 'ENOENT') {
|
||||
console.error('[TranscodeSession] Error recovering sessions:', err.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Start cleanup interval
|
||||
*/
|
||||
let cleanupInterval = null;
|
||||
function startCleanupInterval() {
|
||||
if (!cleanupInterval) {
|
||||
cleanupInterval = setInterval(cleanupStaleSessions, CLEANUP_INTERVAL_MS);
|
||||
cleanupInterval.unref(); // Don't prevent process exit
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get all active sessions (for debugging/monitoring)
|
||||
*/
|
||||
function getAllSessions() {
|
||||
return Array.from(sessions.values()).map(s => ({
|
||||
id: s.id,
|
||||
url: s.url,
|
||||
status: s.status,
|
||||
startTime: s.startTime,
|
||||
lastAccess: s.lastAccess,
|
||||
idleMs: Date.now() - s.lastAccess
|
||||
}));
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
TranscodeSession,
|
||||
createSession,
|
||||
getSession,
|
||||
getOrCreateSession,
|
||||
removeSession,
|
||||
cleanupStaleSessions,
|
||||
recoverSessions,
|
||||
startCleanupInterval,
|
||||
getAllSessions,
|
||||
CACHE_DIR,
|
||||
SEGMENT_DURATION
|
||||
};
|
||||
@@ -0,0 +1,161 @@
|
||||
/**
|
||||
* Xtream Codes API v2 Client
|
||||
* Handles authentication and API calls to Xtream servers
|
||||
*/
|
||||
|
||||
class XtreamApi {
|
||||
constructor(baseUrl, username, password) {
|
||||
// Clean up base URL
|
||||
this.baseUrl = baseUrl.replace(/\/+$/, '');
|
||||
this.username = username;
|
||||
this.password = password;
|
||||
}
|
||||
|
||||
/**
|
||||
* Build API URL with authentication
|
||||
*/
|
||||
buildApiUrl(action, params = {}) {
|
||||
const url = new URL(`${this.baseUrl}/player_api.php`);
|
||||
url.searchParams.set('username', this.username);
|
||||
url.searchParams.set('password', this.password);
|
||||
if (action) {
|
||||
url.searchParams.set('action', action);
|
||||
}
|
||||
for (const [key, value] of Object.entries(params)) {
|
||||
if (value !== undefined && value !== null) {
|
||||
url.searchParams.set(key, value);
|
||||
}
|
||||
}
|
||||
return url.toString();
|
||||
}
|
||||
|
||||
/**
|
||||
* Make API request
|
||||
*/
|
||||
async request(action, params = {}) {
|
||||
const url = this.buildApiUrl(action, params);
|
||||
const response = await fetch(url);
|
||||
if (!response.ok) {
|
||||
throw new Error(`Xtream API error: ${response.status} ${response.statusText}`);
|
||||
}
|
||||
return response.json();
|
||||
}
|
||||
|
||||
/**
|
||||
* Authenticate and get server/user info
|
||||
*/
|
||||
async authenticate() {
|
||||
const data = await this.request(null);
|
||||
if (!data.user_info) {
|
||||
throw new Error('Invalid credentials or server response');
|
||||
}
|
||||
return data;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get live channel categories
|
||||
*/
|
||||
async getLiveCategories() {
|
||||
return this.request('get_live_categories');
|
||||
}
|
||||
|
||||
/**
|
||||
* Get live streams, optionally filtered by category
|
||||
*/
|
||||
async getLiveStreams(categoryId = null) {
|
||||
return this.request('get_live_streams', { category_id: categoryId });
|
||||
}
|
||||
|
||||
/**
|
||||
* Get VOD categories
|
||||
*/
|
||||
async getVodCategories() {
|
||||
return this.request('get_vod_categories');
|
||||
}
|
||||
|
||||
/**
|
||||
* Get VOD streams, optionally filtered by category
|
||||
*/
|
||||
async getVodStreams(categoryId = null) {
|
||||
return this.request('get_vod_streams', { category_id: categoryId });
|
||||
}
|
||||
|
||||
/**
|
||||
* Get VOD info
|
||||
*/
|
||||
async getVodInfo(vodId) {
|
||||
return this.request('get_vod_info', { vod_id: vodId });
|
||||
}
|
||||
|
||||
/**
|
||||
* Get series categories
|
||||
*/
|
||||
async getSeriesCategories() {
|
||||
return this.request('get_series_categories');
|
||||
}
|
||||
|
||||
/**
|
||||
* Get series, optionally filtered by category
|
||||
*/
|
||||
async getSeries(categoryId = null) {
|
||||
return this.request('get_series', { category_id: categoryId });
|
||||
}
|
||||
|
||||
/**
|
||||
* Get series info
|
||||
*/
|
||||
async getSeriesInfo(seriesId) {
|
||||
return this.request('get_series_info', { series_id: seriesId });
|
||||
}
|
||||
|
||||
/**
|
||||
* Get short EPG for a stream
|
||||
*/
|
||||
async getShortEpg(streamId, limit = 10) {
|
||||
return this.request('get_short_epg', { stream_id: streamId, limit });
|
||||
}
|
||||
|
||||
/**
|
||||
* Get full EPG for a stream
|
||||
*/
|
||||
async getSimpleDateTable(streamId) {
|
||||
return this.request('get_simple_data_table', { stream_id: streamId });
|
||||
}
|
||||
|
||||
/**
|
||||
* Build stream URL for playback
|
||||
*/
|
||||
buildStreamUrl(streamId, type = 'live', container = 'ts') {
|
||||
const typeMap = {
|
||||
live: 'live',
|
||||
vod: 'movie',
|
||||
series: 'series'
|
||||
};
|
||||
const streamType = typeMap[type] || 'live';
|
||||
return `${this.baseUrl}/${streamType}/${this.username}/${this.password}/${streamId}.${container}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Get XMLTV EPG URL
|
||||
*/
|
||||
getXmltvUrl() {
|
||||
return `${this.baseUrl}/xmltv.php?username=${this.username}&password=${this.password}`;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Factory function to create API instance from source
|
||||
*/
|
||||
function createFromSource(source) {
|
||||
return new XtreamApi(source.url, source.username, source.password);
|
||||
}
|
||||
|
||||
/**
|
||||
* Static authenticate for testing
|
||||
*/
|
||||
async function authenticate(url, username, password) {
|
||||
const api = new XtreamApi(url, username, password);
|
||||
return api.authenticate();
|
||||
}
|
||||
|
||||
module.exports = { XtreamApi, createFromSource, authenticate };
|
||||
Reference in New Issue
Block a user