Instancia transcoder t1: cambios locales
This commit is contained in:
+2
-1
@@ -67,7 +67,8 @@ function getDefaultSettings() {
|
||||
forceVideoTranscode: false, // Force Video Transcode
|
||||
forceRemux: false,
|
||||
autoTranscode: true,
|
||||
streamFormat: 'm3u8',
|
||||
// Dispatcharr (y muchos Xtream) sirven MPEG-TS aunque la extensión sea .m3u8; Hls.js falla. Usar ts.
|
||||
streamFormat: 'ts',
|
||||
epgRefreshInterval: '24',
|
||||
// User-Agent settings
|
||||
userAgentPreset: 'chrome', // chrome | vlc | tivimate | custom
|
||||
|
||||
+9
-2
@@ -27,7 +27,14 @@ app.use(session({
|
||||
app.use(passport.initialize());
|
||||
app.use(passport.session());
|
||||
|
||||
app.use(express.static(path.join(__dirname, '..', 'public')));
|
||||
// Serve static files with no-cache for JS/CSS so browser always loads the latest version
|
||||
app.use(express.static(path.join(__dirname, '..', 'public'), {
|
||||
setHeaders(res, filePath) {
|
||||
if (filePath.endsWith('.js') || filePath.endsWith('.css')) {
|
||||
res.setHeader('Cache-Control', 'no-cache, must-revalidate');
|
||||
}
|
||||
}
|
||||
}));
|
||||
|
||||
// FFMPEG Configuration (optional - for transcoding support)
|
||||
// Priority: 1. System FFmpeg (better Docker DNS support), 2. ffmpeg-static npm package
|
||||
@@ -199,7 +206,7 @@ app.use((err, req, res, next) => {
|
||||
});
|
||||
|
||||
app.listen(PORT, async () => {
|
||||
console.log(`NodeCast TV server running on http://localhost:${PORT}`);
|
||||
console.log(`KiraTV server running on http://localhost:${PORT}`);
|
||||
|
||||
// Load plugins
|
||||
await loadPlugins().catch(err => {
|
||||
|
||||
+66
-24
@@ -1,6 +1,14 @@
|
||||
const express = require('express');
|
||||
const router = express.Router();
|
||||
const { getDb } = require('../db/sqlite');
|
||||
const { sources } = require('../db');
|
||||
const { requireAuth } = require('../auth');
|
||||
const {
|
||||
getAccessibleSourceIds,
|
||||
assertCanUseSourceById
|
||||
} = require('../sourceAccess');
|
||||
|
||||
router.use(requireAuth);
|
||||
|
||||
// Helper to map API item types to DB types and tables
|
||||
function mapItemType(apiType) {
|
||||
@@ -21,27 +29,44 @@ router.get('/hidden', async (req, res) => {
|
||||
const { sourceId } = req.query;
|
||||
const db = getDb();
|
||||
|
||||
let hidden = [];
|
||||
let allowedIds = null;
|
||||
if (req.user.role !== 'admin') {
|
||||
allowedIds = await getAccessibleSourceIds(req, sources);
|
||||
if (allowedIds.length === 0) {
|
||||
return res.json([]);
|
||||
}
|
||||
}
|
||||
|
||||
const resultFormat = (row, itemType) => ({
|
||||
source_id: row.source_id,
|
||||
item_type: itemType,
|
||||
item_id: itemType.includes('category') || itemType === 'group' ? row.category_id : row.item_id
|
||||
});
|
||||
|
||||
// Query Categories
|
||||
let catQuery = `SELECT source_id, category_id, type FROM categories WHERE is_hidden = 1`;
|
||||
let itemQuery = `SELECT source_id, item_id, type FROM playlist_items WHERE is_hidden = 1`;
|
||||
let catParams = [];
|
||||
let itemParams = [];
|
||||
|
||||
const params = [];
|
||||
if (sourceId) {
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sourceId))) return;
|
||||
const sid = parseInt(sourceId, 10);
|
||||
catQuery += ` AND source_id = ?`;
|
||||
itemQuery += ` AND source_id = ?`;
|
||||
const sid = parseInt(sourceId);
|
||||
params.push(sid);
|
||||
catParams.push(sid);
|
||||
itemParams.push(sid);
|
||||
} else if (allowedIds) {
|
||||
const ph = allowedIds.map(() => '?').join(',');
|
||||
catQuery += ` AND source_id IN (${ph})`;
|
||||
itemQuery += ` AND source_id IN (${ph})`;
|
||||
catParams = [...allowedIds];
|
||||
itemParams = [...allowedIds];
|
||||
}
|
||||
|
||||
const hiddenCats = db.prepare(catQuery).all(...params);
|
||||
const hiddenItems = db.prepare(itemQuery).all(...params);
|
||||
const hiddenCats = db.prepare(catQuery).all(...catParams);
|
||||
const hiddenItems = db.prepare(itemQuery).all(...itemParams);
|
||||
|
||||
const hidden = [];
|
||||
|
||||
hiddenCats.forEach(row => {
|
||||
let apiType;
|
||||
@@ -75,6 +100,7 @@ router.post('/hide', async (req, res) => {
|
||||
const mapping = mapItemType(itemType);
|
||||
|
||||
if (!mapping) return res.status(400).json({ error: 'Invalid item type' });
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sourceId))) return;
|
||||
|
||||
const db = getDb();
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
@@ -101,6 +127,7 @@ router.post('/show', async (req, res) => {
|
||||
const mapping = mapItemType(itemType);
|
||||
|
||||
if (!mapping) return res.status(400).json({ error: 'Invalid item type' });
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sourceId))) return;
|
||||
|
||||
const db = getDb();
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
@@ -127,6 +154,8 @@ router.get('/hidden/check', async (req, res) => {
|
||||
const mapping = mapItemType(itemType);
|
||||
if (!mapping) return res.json({ hidden: false });
|
||||
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sourceId))) return;
|
||||
|
||||
const db = getDb();
|
||||
const idCol = mapping.table === 'categories' ? 'category_id' : 'item_id';
|
||||
|
||||
@@ -148,13 +177,16 @@ router.post('/hide/bulk', async (req, res) => {
|
||||
const { items } = req.body;
|
||||
if (!Array.isArray(items)) return res.status(400).json({ error: 'items array required' });
|
||||
|
||||
const uniqueIds = [...new Set(items.map(i => i.sourceId))];
|
||||
for (const sid of uniqueIds) {
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sid))) return;
|
||||
}
|
||||
|
||||
const db = getDb();
|
||||
|
||||
// Prepare statements once
|
||||
const hideCat = db.prepare('UPDATE categories SET is_hidden = 1 WHERE source_id = ? AND type = ? AND category_id = ?');
|
||||
const hideItem = db.prepare('UPDATE playlist_items SET is_hidden = 1 WHERE source_id = ? AND type = ? AND item_id = ?');
|
||||
|
||||
// Cascading statements (hide all children of a category)
|
||||
const hideCatChildren = db.prepare('UPDATE playlist_items SET is_hidden = 1 WHERE source_id = ? AND type = ? AND category_id = ?');
|
||||
|
||||
const runBulk = db.transaction((list) => {
|
||||
@@ -162,12 +194,9 @@ router.post('/hide/bulk', async (req, res) => {
|
||||
const mapping = mapItemType(item.itemType);
|
||||
if (mapping) {
|
||||
if (mapping.table === 'categories') {
|
||||
// Hide the category
|
||||
hideCat.run(item.sourceId, mapping.type, item.itemId);
|
||||
// Cascade to children
|
||||
hideCatChildren.run(item.sourceId, mapping.type, item.itemId);
|
||||
} else {
|
||||
// Hide individual item
|
||||
hideItem.run(item.sourceId, mapping.type, item.itemId);
|
||||
}
|
||||
}
|
||||
@@ -191,13 +220,16 @@ router.post('/show/bulk', async (req, res) => {
|
||||
const { items } = req.body;
|
||||
if (!Array.isArray(items)) return res.status(400).json({ error: 'items array required' });
|
||||
|
||||
const uniqueIds = [...new Set(items.map(i => i.sourceId))];
|
||||
for (const sid of uniqueIds) {
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sid))) return;
|
||||
}
|
||||
|
||||
const db = getDb();
|
||||
|
||||
// Prepare statements once
|
||||
const showCat = db.prepare('UPDATE categories SET is_hidden = 0 WHERE source_id = ? AND type = ? AND category_id = ?');
|
||||
const showItem = db.prepare('UPDATE playlist_items SET is_hidden = 0 WHERE source_id = ? AND type = ? AND item_id = ?');
|
||||
|
||||
// Cascading statements (show all children of a category)
|
||||
const showCatChildren = db.prepare('UPDATE playlist_items SET is_hidden = 0 WHERE source_id = ? AND type = ? AND category_id = ?');
|
||||
|
||||
const runBulk = db.transaction((list) => {
|
||||
@@ -205,12 +237,9 @@ router.post('/show/bulk', async (req, res) => {
|
||||
const mapping = mapItemType(item.itemType);
|
||||
if (mapping) {
|
||||
if (mapping.table === 'categories') {
|
||||
// Show the category
|
||||
showCat.run(item.sourceId, mapping.type, item.itemId);
|
||||
// Cascade to children
|
||||
showCatChildren.run(item.sourceId, mapping.type, item.itemId);
|
||||
} else {
|
||||
// Show individual item
|
||||
showItem.run(item.sourceId, mapping.type, item.itemId);
|
||||
}
|
||||
}
|
||||
@@ -234,14 +263,15 @@ router.post('/show/all', async (req, res) => {
|
||||
const { sourceId, contentType } = req.body;
|
||||
if (!sourceId) return res.status(400).json({ error: 'sourceId required' });
|
||||
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sourceId))) return;
|
||||
|
||||
const db = getDb();
|
||||
let catCount = 0;
|
||||
let itemCount = 0;
|
||||
|
||||
// Determine which types to update based on contentType
|
||||
const types = contentType === 'movies' ? ['movie']
|
||||
: contentType === 'series' ? ['series']
|
||||
: ['live']; // default to channels
|
||||
: ['live'];
|
||||
|
||||
for (const type of types) {
|
||||
const catResult = db.prepare(`UPDATE categories SET is_hidden = 0 WHERE source_id = ? AND type = ?`).run(sourceId, type);
|
||||
@@ -264,14 +294,15 @@ router.post('/hide/all', async (req, res) => {
|
||||
const { sourceId, contentType } = req.body;
|
||||
if (!sourceId) return res.status(400).json({ error: 'sourceId required' });
|
||||
|
||||
if (!(await assertCanUseSourceById(req, res, sources, sourceId))) return;
|
||||
|
||||
const db = getDb();
|
||||
let catCount = 0;
|
||||
let itemCount = 0;
|
||||
|
||||
// Determine which types to update based on contentType
|
||||
const types = contentType === 'movies' ? ['movie']
|
||||
: contentType === 'series' ? ['series']
|
||||
: ['live']; // default to channels
|
||||
: ['live'];
|
||||
|
||||
for (const type of types) {
|
||||
const catResult = db.prepare(`UPDATE categories SET is_hidden = 1 WHERE source_id = ? AND type = ?`).run(sourceId, type);
|
||||
@@ -296,6 +327,18 @@ router.get('/recent', async (req, res) => {
|
||||
return res.status(400).json({ error: 'Valid type (movie or series) is required' });
|
||||
}
|
||||
|
||||
let accessibleFilter = '';
|
||||
const params = [type];
|
||||
if (req.user.role !== 'admin') {
|
||||
const ids = await getAccessibleSourceIds(req, sources);
|
||||
if (ids.length === 0) {
|
||||
return res.json([]);
|
||||
}
|
||||
accessibleFilter = ` AND p.source_id IN (${ids.map(() => '?').join(',')})`;
|
||||
params.push(...ids);
|
||||
}
|
||||
params.push(parseInt(limit, 10));
|
||||
|
||||
const db = getDb();
|
||||
const recentItems = db.prepare(`
|
||||
SELECT * FROM playlist_items p
|
||||
@@ -308,11 +351,11 @@ router.get('/recent', async (req, res) => {
|
||||
AND c.type = p.type
|
||||
AND c.is_hidden = 1
|
||||
)
|
||||
${accessibleFilter}
|
||||
ORDER BY p.added_at DESC
|
||||
LIMIT ?
|
||||
`).all(type, parseInt(limit));
|
||||
`).all(...params);
|
||||
|
||||
// Parse JSON data for each item
|
||||
const formatted = recentItems.map(item => ({
|
||||
...item,
|
||||
data: JSON.parse(item.data)
|
||||
@@ -326,4 +369,3 @@ router.get('/recent', async (req, res) => {
|
||||
});
|
||||
|
||||
module.exports = router;
|
||||
|
||||
|
||||
+44
-21
@@ -170,32 +170,55 @@ router.get('/', async (req, res) => {
|
||||
|
||||
console.log(`[Probe] Probing: ${url.substring(0, 80)}... ${ua ? `(UA: ${ua})` : ''}`);
|
||||
|
||||
try {
|
||||
const probeResult = await probeStream(url, ffprobePath, ua);
|
||||
const analysis = analyzeProbeResult(probeResult, url);
|
||||
// Retry probe up to 3 times on 5XX errors (provider slot not released yet)
|
||||
const MAX_PROBE_RETRIES = 3;
|
||||
const PROBE_RETRY_DELAY_MS = 2500;
|
||||
let lastErr = null;
|
||||
|
||||
// Cache result
|
||||
probeCache.set(cacheKey, { result: analysis, timestamp: Date.now() });
|
||||
for (let attempt = 1; attempt <= MAX_PROBE_RETRIES; attempt++) {
|
||||
try {
|
||||
const probeResult = await probeStream(url, ffprobePath, ua);
|
||||
const analysis = analyzeProbeResult(probeResult, url);
|
||||
|
||||
console.log(`[Probe] Result: video=${analysis.video}, audio=${analysis.audio}, ` +
|
||||
`container=${analysis.container}, compatible=${analysis.compatible}, ` +
|
||||
`needsRemux=${analysis.needsRemux}, needsTranscode=${analysis.needsTranscode}`);
|
||||
// Cache successful result
|
||||
probeCache.set(cacheKey, { result: analysis, timestamp: Date.now() });
|
||||
|
||||
res.json(analysis);
|
||||
} catch (err) {
|
||||
console.error('[Probe] Failed:', err.message);
|
||||
console.log(`[Probe] Result (attempt ${attempt}): video=${analysis.video}, audio=${analysis.audio}, ` +
|
||||
`container=${analysis.container}, compatible=${analysis.compatible}, ` +
|
||||
`needsRemux=${analysis.needsRemux}, needsTranscode=${analysis.needsTranscode}`);
|
||||
|
||||
// On error, assume transcode needed to be safe
|
||||
res.json({
|
||||
video: 'unknown',
|
||||
audio: 'unknown',
|
||||
container: 'unknown',
|
||||
compatible: false,
|
||||
needsRemux: false,
|
||||
needsTranscode: true,
|
||||
error: err.message
|
||||
});
|
||||
return res.json(analysis);
|
||||
} catch (err) {
|
||||
lastErr = err;
|
||||
const is5xx = /5[0-9]{2}|Server Error/i.test(err.message);
|
||||
if (is5xx && attempt < MAX_PROBE_RETRIES) {
|
||||
console.warn(`[Probe] 5XX from provider (attempt ${attempt}/${MAX_PROBE_RETRIES}), retrying in ${PROBE_RETRY_DELAY_MS}ms...`);
|
||||
await new Promise(r => setTimeout(r, PROBE_RETRY_DELAY_MS));
|
||||
continue;
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
console.error('[Probe] Failed after retries:', lastErr.message);
|
||||
|
||||
// If there's a stale cached result, use it rather than assuming needsTranscode
|
||||
const stale = probeCache.get(cacheKey);
|
||||
if (stale) {
|
||||
console.warn('[Probe] Using stale cached result as fallback');
|
||||
return res.json({ ...stale.result, stale: true });
|
||||
}
|
||||
|
||||
// Final fallback: assume transcode needed
|
||||
res.json({
|
||||
video: 'unknown',
|
||||
audio: 'unknown',
|
||||
container: 'unknown',
|
||||
compatible: false,
|
||||
needsRemux: false,
|
||||
needsTranscode: true,
|
||||
error: lastErr.message
|
||||
});
|
||||
});
|
||||
|
||||
module.exports = router;
|
||||
|
||||
+45
-18
@@ -12,6 +12,35 @@ const https = require('https');
|
||||
const { spawn } = require('child_process');
|
||||
const ffmpegPath = require('ffmpeg-static');
|
||||
const { Readable } = require('stream');
|
||||
const auth = require('../auth');
|
||||
const { canUserAccessSource } = require('../sourceAccess');
|
||||
|
||||
// Rutas que el navegador pide sin Bearer (HLS, <img>, etc.) deben quedar fuera de JWT
|
||||
router.use((req, res, next) => {
|
||||
const p = req.path || (req.url || '').split('?')[0];
|
||||
if (p === '/stream' || p === '/image') return next();
|
||||
return auth.requireAuth(req, res, next);
|
||||
});
|
||||
|
||||
router.param('sourceId', async (req, res, next, id) => {
|
||||
try {
|
||||
const sid = parseInt(id, 10);
|
||||
if (Number.isNaN(sid)) {
|
||||
return res.status(400).json({ error: 'Invalid source id' });
|
||||
}
|
||||
const source = await sources.getById(sid);
|
||||
if (!source) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
if (!canUserAccessSource(req.user, source)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
req.grantedSource = source;
|
||||
next();
|
||||
} catch (err) {
|
||||
next(err);
|
||||
}
|
||||
});
|
||||
|
||||
// Default cache max age in hours
|
||||
const DEFAULT_MAX_AGE_HOURS = 24;
|
||||
@@ -83,8 +112,8 @@ function getStreamsFromDb(sourceId, type, categoryId = null, includeHidden = fal
|
||||
// Login / Authenticate
|
||||
router.get('/xtream/:sourceId', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source || source.type !== 'xtream') return res.status(404).send('Source not found');
|
||||
const source = req.grantedSource;
|
||||
if (source.type !== 'xtream') return res.status(404).send('Source not found');
|
||||
|
||||
// Proxy auth check to upstream to ensure credentials are still valid
|
||||
|
||||
@@ -185,8 +214,7 @@ router.get('/xtream/:sourceId/series', async (req, res) => {
|
||||
// Proxy series info request
|
||||
router.get('/xtream/:sourceId/series_info', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source) return res.status(404).send('Source not found');
|
||||
const source = req.grantedSource;
|
||||
|
||||
const seriesId = req.query.series_id;
|
||||
if (!seriesId) return res.status(400).send('series_id required');
|
||||
@@ -207,8 +235,7 @@ router.get('/xtream/:sourceId/series_info', async (req, res) => {
|
||||
// VOD Info
|
||||
router.get('/xtream/:sourceId/vod_info', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source) return res.status(404).send('Source not found');
|
||||
const source = req.grantedSource;
|
||||
|
||||
const vodId = req.query.vod_id;
|
||||
if (!vodId) return res.status(400).send('vod_id required');
|
||||
@@ -230,14 +257,14 @@ router.get('/xtream/:sourceId/vod_info', async (req, res) => {
|
||||
// Returns the direct stream URL for a given stream ID
|
||||
router.get('/xtream/:sourceId/stream/:streamId/:type', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source || source.type !== 'xtream') {
|
||||
const source = req.grantedSource;
|
||||
if (source.type !== 'xtream') {
|
||||
return res.status(404).json({ error: 'Xtream source not found' });
|
||||
}
|
||||
|
||||
const streamId = req.params.streamId;
|
||||
const type = req.params.type || 'live';
|
||||
const container = req.query.container || 'm3u8';
|
||||
const container = req.query.container || 'ts';
|
||||
|
||||
// Construct the Xtream stream URL
|
||||
// Format: http://server:port/live/username/password/streamId.container (for live)
|
||||
@@ -393,8 +420,8 @@ router.delete('/cache/:sourceId', (req, res) => {
|
||||
router.get('/xtream/:sourceId/:action', async (req, res) => {
|
||||
try {
|
||||
const sourceId = req.params.sourceId;
|
||||
const source = await sources.getById(sourceId);
|
||||
if (!source || source.type !== 'xtream') {
|
||||
const source = req.grantedSource;
|
||||
if (source.type !== 'xtream') {
|
||||
return res.status(404).json({ error: 'Xtream source not found' });
|
||||
}
|
||||
|
||||
@@ -478,14 +505,14 @@ router.get('/xtream/:sourceId/:action', async (req, res) => {
|
||||
*/
|
||||
router.get('/xtream/:sourceId/stream/:streamId/:type?', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source || source.type !== 'xtream') {
|
||||
const source = req.grantedSource;
|
||||
if (source.type !== 'xtream') {
|
||||
return res.status(404).json({ error: 'Xtream source not found' });
|
||||
}
|
||||
|
||||
const api = xtreamApi.createFromSource(source);
|
||||
const { streamId, type = 'live' } = req.params;
|
||||
const { container = 'm3u8' } = req.query;
|
||||
const { container = 'ts' } = req.query;
|
||||
|
||||
const url = api.buildStreamUrl(streamId, type, container);
|
||||
res.json({ url });
|
||||
@@ -505,8 +532,8 @@ router.get('/xtream/:sourceId/stream/:streamId/:type?', async (req, res) => {
|
||||
router.get('/epg/:sourceId', async (req, res) => {
|
||||
try {
|
||||
const sourceId = req.params.sourceId;
|
||||
const source = await sources.getById(sourceId);
|
||||
if (!source || (source.type !== 'epg' && source.type !== 'xtream')) {
|
||||
const source = req.grantedSource;
|
||||
if (source.type !== 'epg' && source.type !== 'xtream') {
|
||||
return res.status(404).json({ error: 'Valid EPG source not found' });
|
||||
}
|
||||
|
||||
@@ -567,8 +594,8 @@ router.delete('/epg/:sourceId/cache', (req, res) => {
|
||||
*/
|
||||
router.post('/epg/:sourceId/channels', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.sourceId);
|
||||
if (!source || source.type !== 'epg') {
|
||||
const source = req.grantedSource;
|
||||
if (source.type !== 'epg') {
|
||||
return res.status(404).json({ error: 'EPG source not found' });
|
||||
}
|
||||
|
||||
|
||||
+137
-97
@@ -5,29 +5,84 @@ const { getDb } = require('../db/sqlite');
|
||||
const xtreamApi = require('../services/xtreamApi');
|
||||
const syncService = require('../services/syncService');
|
||||
const m3uParser = require('../services/m3uParser');
|
||||
const { requireAuth, requireAdmin } = require('../auth');
|
||||
const { canUserAccessSource, getAccessibleSourceIds } = require('../sourceAccess');
|
||||
|
||||
// Get all sources
|
||||
router.use(requireAuth);
|
||||
|
||||
function sanitizeSourceList(list) {
|
||||
return list.map((s) => ({
|
||||
...s,
|
||||
password: s.password ? '••••••••' : null
|
||||
}));
|
||||
}
|
||||
|
||||
function parseOwnerId(body) {
|
||||
if (body.ownerId === undefined || body.ownerId === '') return null;
|
||||
const n = parseInt(body.ownerId, 10);
|
||||
return Number.isNaN(n) ? null : n;
|
||||
}
|
||||
|
||||
// Estimate by URL (must stay before /:id routes)
|
||||
const M3U_LARGE_THRESHOLD = 50000;
|
||||
|
||||
router.post('/estimate', async (req, res) => {
|
||||
try {
|
||||
const { url, type } = req.body;
|
||||
|
||||
if (!url) {
|
||||
return res.status(400).json({ error: 'URL is required' });
|
||||
}
|
||||
|
||||
if (type !== 'm3u') {
|
||||
return res.json({ count: 0, needsWarning: false, threshold: M3U_LARGE_THRESHOLD });
|
||||
}
|
||||
|
||||
const count = await m3uParser.countEntries(url);
|
||||
|
||||
res.json({
|
||||
count,
|
||||
needsWarning: count > M3U_LARGE_THRESHOLD,
|
||||
threshold: M3U_LARGE_THRESHOLD
|
||||
});
|
||||
} catch (err) {
|
||||
console.error('Error estimating M3U size:', err);
|
||||
res.status(500).json({ error: 'Failed to estimate playlist size', message: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
// Global sync — admin only
|
||||
router.post('/sync-all', requireAdmin, async (req, res) => {
|
||||
try {
|
||||
syncService.syncAll().catch(console.error);
|
||||
res.json({ success: true, message: 'Global sync started' });
|
||||
} catch (err) {
|
||||
console.error('Error starting global sync:', err);
|
||||
res.status(500).json({ error: 'Failed to start global sync' });
|
||||
}
|
||||
});
|
||||
|
||||
// Get all sources (filtered by access)
|
||||
router.get('/', async (req, res) => {
|
||||
try {
|
||||
const allSources = await sources.getAll();
|
||||
// Don't expose passwords in list view
|
||||
const sanitized = allSources.map(s => ({
|
||||
...s,
|
||||
password: s.password ? '••••••••' : null
|
||||
}));
|
||||
res.json(sanitized);
|
||||
const filtered = allSources.filter((s) => canUserAccessSource(req.user, s));
|
||||
res.json(sanitizeSourceList(filtered));
|
||||
} catch (err) {
|
||||
console.error('Error getting sources:', err);
|
||||
res.status(500).json({ error: 'Failed to get sources' });
|
||||
}
|
||||
});
|
||||
|
||||
// Get sync status for all sources
|
||||
// Get sync status (filtered for non-admins)
|
||||
router.get('/status', async (req, res) => {
|
||||
try {
|
||||
const { getDb } = require('../db/sqlite');
|
||||
const db = getDb();
|
||||
const statuses = db.prepare('SELECT * FROM sync_status').all();
|
||||
let statuses = db.prepare('SELECT * FROM sync_status').all();
|
||||
if (req.user.role !== 'admin') {
|
||||
const allowed = new Set(await getAccessibleSourceIds(req, sources));
|
||||
statuses = statuses.filter((row) => allowed.has(row.source_id));
|
||||
}
|
||||
res.json(statuses);
|
||||
} catch (err) {
|
||||
console.error('Error getting sync status:', err);
|
||||
@@ -39,13 +94,42 @@ router.get('/status', async (req, res) => {
|
||||
router.get('/type/:type', async (req, res) => {
|
||||
try {
|
||||
const typeSources = await sources.getByType(req.params.type);
|
||||
res.json(typeSources);
|
||||
const filtered = typeSources.filter((s) => canUserAccessSource(req.user, s));
|
||||
res.json(sanitizeSourceList(filtered));
|
||||
} catch (err) {
|
||||
console.error('Error getting sources by type:', err);
|
||||
res.status(500).json({ error: 'Failed to get sources' });
|
||||
}
|
||||
});
|
||||
|
||||
// Estimate by source ID (before bare GET /:id — same base path length)
|
||||
router.get('/:id/estimate', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.id);
|
||||
if (!source) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
if (!canUserAccessSource(req.user, source)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
|
||||
if (source.type !== 'm3u') {
|
||||
return res.json({ count: 0, needsWarning: false, threshold: M3U_LARGE_THRESHOLD });
|
||||
}
|
||||
|
||||
const count = await m3uParser.countEntries(source.url);
|
||||
|
||||
res.json({
|
||||
count,
|
||||
needsWarning: count > M3U_LARGE_THRESHOLD,
|
||||
threshold: M3U_LARGE_THRESHOLD
|
||||
});
|
||||
} catch (err) {
|
||||
console.error('Error estimating M3U size:', err);
|
||||
res.status(500).json({ error: 'Failed to estimate playlist size', message: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
// Get single source
|
||||
router.get('/:id', async (req, res) => {
|
||||
try {
|
||||
@@ -53,6 +137,9 @@ router.get('/:id', async (req, res) => {
|
||||
if (!source) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
if (!canUserAccessSource(req.user, source)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
res.json(source);
|
||||
} catch (err) {
|
||||
console.error('Error getting source:', err);
|
||||
@@ -73,8 +160,14 @@ router.post('/', async (req, res) => {
|
||||
return res.status(400).json({ error: 'Invalid source type' });
|
||||
}
|
||||
|
||||
const source = await sources.create({ type, name, url, username, password });
|
||||
// Trigger Sync
|
||||
let ownerId = null;
|
||||
if (req.user.role === 'admin') {
|
||||
ownerId = parseOwnerId(req.body);
|
||||
} else {
|
||||
ownerId = req.user.id;
|
||||
}
|
||||
|
||||
const source = await sources.create({ type, name, url, username, password, ownerId });
|
||||
syncService.syncSource(source.id).catch(console.error);
|
||||
res.status(201).json(source);
|
||||
} catch (err) {
|
||||
@@ -90,16 +183,23 @@ router.put('/:id', async (req, res) => {
|
||||
if (!existing) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
if (!canUserAccessSource(req.user, existing)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
|
||||
const { name, url, username, password } = req.body;
|
||||
const updated = await sources.update(req.params.id, {
|
||||
const updates = {
|
||||
name: name || existing.name,
|
||||
url: url || existing.url,
|
||||
username: username !== undefined ? username : existing.username,
|
||||
password: password !== undefined ? password : existing.password
|
||||
});
|
||||
// Trigger Sync (if critical fields changed? safely just trigger it)
|
||||
syncService.syncSource(parseInt(req.params.id)).catch(console.error);
|
||||
};
|
||||
if (req.user.role === 'admin' && req.body.ownerId !== undefined) {
|
||||
updates.ownerId = parseOwnerId(req.body);
|
||||
}
|
||||
|
||||
const updated = await sources.update(req.params.id, updates);
|
||||
syncService.syncSource(parseInt(req.params.id, 10)).catch(console.error);
|
||||
res.json(updated);
|
||||
} catch (err) {
|
||||
console.error('Error updating source:', err);
|
||||
@@ -110,13 +210,15 @@ router.put('/:id', async (req, res) => {
|
||||
// Delete source
|
||||
router.delete('/:id', async (req, res) => {
|
||||
try {
|
||||
const sourceId = parseInt(req.params.id);
|
||||
const sourceId = parseInt(req.params.id, 10);
|
||||
const existing = await sources.getById(sourceId);
|
||||
if (!existing) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
if (!canUserAccessSource(req.user, existing)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
|
||||
// Cascade delete: Clean up SQLite data for this source
|
||||
const db = getDb();
|
||||
const deleteCategories = db.prepare('DELETE FROM categories WHERE source_id = ?');
|
||||
const deleteItems = db.prepare('DELETE FROM playlist_items WHERE source_id = ?');
|
||||
@@ -130,7 +232,6 @@ router.delete('/:id', async (req, res) => {
|
||||
|
||||
console.log(`[Source] Cascade delete for source ${sourceId}: ${catResult.changes} categories, ${itemResult.changes} items, ${epgResult.changes} EPG programs`);
|
||||
|
||||
// Delete source config and related hidden items (favorites handled by db.js)
|
||||
await sources.delete(sourceId);
|
||||
|
||||
res.json({ success: true });
|
||||
@@ -143,14 +244,21 @@ router.delete('/:id', async (req, res) => {
|
||||
// Toggle source enabled/disabled
|
||||
router.post('/:id/toggle', async (req, res) => {
|
||||
try {
|
||||
const existing = await sources.getById(req.params.id);
|
||||
if (!existing) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
if (!canUserAccessSource(req.user, existing)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
|
||||
const updated = await sources.toggleEnabled(req.params.id);
|
||||
if (!updated) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
|
||||
// If enabled, trigger sync
|
||||
if (updated.enabled) {
|
||||
syncService.syncSource(parseInt(req.params.id)).catch(console.error);
|
||||
syncService.syncSource(parseInt(req.params.id, 10)).catch(console.error);
|
||||
}
|
||||
|
||||
res.json(updated);
|
||||
@@ -163,11 +271,13 @@ router.post('/:id/toggle', async (req, res) => {
|
||||
// Manual Sync
|
||||
router.post('/:id/sync', async (req, res) => {
|
||||
try {
|
||||
const id = parseInt(req.params.id);
|
||||
const id = parseInt(req.params.id, 10);
|
||||
const source = await sources.getById(id);
|
||||
if (!source) return res.status(404).json({ error: 'Source not found' });
|
||||
if (!canUserAccessSource(req.user, source)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
|
||||
// Trigger sync (async)
|
||||
syncService.syncSource(id).catch(console.error);
|
||||
|
||||
res.json({ success: true, message: 'Sync started' });
|
||||
@@ -184,6 +294,9 @@ router.post('/:id/test', async (req, res) => {
|
||||
if (!source) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
if (!canUserAccessSource(req.user, source)) {
|
||||
return res.status(403).json({ error: 'No access to this source' });
|
||||
}
|
||||
|
||||
if (source.type === 'xtream') {
|
||||
const result = await xtreamApi.authenticate(source.url, source.username, source.password);
|
||||
@@ -205,77 +318,4 @@ router.post('/:id/test', async (req, res) => {
|
||||
}
|
||||
});
|
||||
|
||||
// Estimate M3U playlist size (for large playlist warning)
|
||||
const M3U_LARGE_THRESHOLD = 50000;
|
||||
|
||||
// Estimate by URL (for new sources before creation)
|
||||
router.post('/estimate', async (req, res) => {
|
||||
try {
|
||||
const { url, type } = req.body;
|
||||
|
||||
if (!url) {
|
||||
return res.status(400).json({ error: 'URL is required' });
|
||||
}
|
||||
|
||||
// Only M3U sources need estimation
|
||||
if (type !== 'm3u') {
|
||||
return res.json({ count: 0, needsWarning: false, threshold: M3U_LARGE_THRESHOLD });
|
||||
}
|
||||
|
||||
console.log(`[Sources] Estimating M3U size for URL...`);
|
||||
const count = await m3uParser.countEntries(url);
|
||||
console.log(`[Sources] M3U estimate: ${count} entries`);
|
||||
|
||||
res.json({
|
||||
count,
|
||||
needsWarning: count > M3U_LARGE_THRESHOLD,
|
||||
threshold: M3U_LARGE_THRESHOLD
|
||||
});
|
||||
} catch (err) {
|
||||
console.error('Error estimating M3U size:', err);
|
||||
res.status(500).json({ error: 'Failed to estimate playlist size', message: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
// Estimate by source ID (for existing sources)
|
||||
router.get('/:id/estimate', async (req, res) => {
|
||||
try {
|
||||
const source = await sources.getById(req.params.id);
|
||||
if (!source) {
|
||||
return res.status(404).json({ error: 'Source not found' });
|
||||
}
|
||||
|
||||
// Only M3U sources need estimation
|
||||
if (source.type !== 'm3u') {
|
||||
return res.json({ count: 0, needsWarning: false, threshold: M3U_LARGE_THRESHOLD });
|
||||
}
|
||||
|
||||
console.log(`[Sources] Estimating M3U size for ${source.name}...`);
|
||||
const count = await m3uParser.countEntries(source.url);
|
||||
console.log(`[Sources] M3U estimate: ${count} entries`);
|
||||
|
||||
res.json({
|
||||
count,
|
||||
needsWarning: count > M3U_LARGE_THRESHOLD,
|
||||
threshold: M3U_LARGE_THRESHOLD
|
||||
});
|
||||
} catch (err) {
|
||||
console.error('Error estimating M3U size:', err);
|
||||
res.status(500).json({ error: 'Failed to estimate playlist size', message: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
// Global Sync - sync all enabled sources
|
||||
router.post('/sync-all', async (req, res) => {
|
||||
try {
|
||||
// Trigger global sync (async - don't wait for completion)
|
||||
syncService.syncAll().catch(console.error);
|
||||
res.json({ success: true, message: 'Global sync started' });
|
||||
} catch (err) {
|
||||
console.error('Error starting global sync:', err);
|
||||
res.status(500).json({ error: 'Failed to start global sync' });
|
||||
}
|
||||
});
|
||||
|
||||
module.exports = router;
|
||||
|
||||
|
||||
+134
-26
@@ -1,10 +1,45 @@
|
||||
const express = require('express');
|
||||
const router = express.Router();
|
||||
const { spawn } = require('child_process');
|
||||
const { spawn, execFile } = require('child_process');
|
||||
const path = require('path');
|
||||
const fs = require('fs').promises;
|
||||
const db = require('../db');
|
||||
const transcodeSession = require('../services/transcodeSession');
|
||||
const streamInput = require('../utils/streamInput');
|
||||
|
||||
/**
|
||||
* Probe a stream URL with FFprobe to get its total duration in seconds.
|
||||
* VOD pesado usa probes mayores y timeout más alto (MKV Jellyfin/remux).
|
||||
*/
|
||||
function probeDuration(url, ffprobePath, userAgent) {
|
||||
return new Promise((resolve) => {
|
||||
const sizing = streamInput.probeSizingForUrl(url);
|
||||
const referer = streamInput.deriveRefererFromUrl(url);
|
||||
const args = [
|
||||
'-v', 'quiet',
|
||||
'-print_format', 'json',
|
||||
'-show_entries', 'format=duration',
|
||||
'-user_agent', userAgent || 'Mozilla/5.0',
|
||||
'-probesize', sizing.probesize,
|
||||
'-analyzeduration', sizing.analyzeduration,
|
||||
];
|
||||
if (referer) {
|
||||
args.push('-referer', referer);
|
||||
}
|
||||
args.push(url);
|
||||
const timeoutMs = sizing.heavy ? 45000 : 8000;
|
||||
const proc = execFile(ffprobePath || 'ffprobe', args, { timeout: timeoutMs }, (err, stdout) => {
|
||||
if (err) { resolve(0); return; }
|
||||
try {
|
||||
const data = JSON.parse(stdout);
|
||||
const dur = parseFloat(data?.format?.duration);
|
||||
resolve(Number.isFinite(dur) && dur > 0 ? Math.round(dur) : 0);
|
||||
} catch { resolve(0); }
|
||||
});
|
||||
// Ensure the process is killed if execFile timeout fires
|
||||
proc.on('error', () => resolve(0));
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Transcode Routes
|
||||
@@ -40,38 +75,78 @@ router.post('/session', async (req, res) => {
|
||||
const userAgent = db.getUserAgent(settings);
|
||||
|
||||
try {
|
||||
const session = await transcodeSession.createSession(url, {
|
||||
ffmpegPath,
|
||||
userAgent,
|
||||
seekOffset: seekOffset || 0,
|
||||
hwEncoder: settings.hwEncoder || 'software',
|
||||
maxResolution: settings.maxResolution || '1080p',
|
||||
quality: settings.quality || 'medium',
|
||||
audioMixPreset: settings.audioMixPreset || 'auto', // Audio downmix preset
|
||||
// Upscaling options
|
||||
upscaleEnabled: settings.upscaleEnabled || false,
|
||||
upscaleMethod: settings.upscaleMethod || 'hardware',
|
||||
upscaleTarget: settings.upscaleTarget || '1080p',
|
||||
videoMode: videoMode, // 'copy' or 'encode'
|
||||
videoCodec: videoCodec, // 'h264', 'hevc', etc.
|
||||
audioCodec: audioCodec, // 'aac', 'ac3', etc.
|
||||
audioChannels: audioChannels // number of channels (2=stereo)
|
||||
});
|
||||
const ffprobePath = (ffmpegPath || 'ffmpeg').replace(/ffmpeg$/, 'ffprobe');
|
||||
|
||||
await session.start();
|
||||
// Probe duration and start session in parallel (don't block on duration)
|
||||
const [session, durationSec] = await Promise.all([
|
||||
transcodeSession.createSession(url, {
|
||||
ffmpegPath,
|
||||
userAgent,
|
||||
seekOffset: seekOffset || 0,
|
||||
hwEncoder: settings.hwEncoder || 'software',
|
||||
maxResolution: settings.maxResolution || '1080p',
|
||||
quality: settings.quality || 'medium',
|
||||
audioMixPreset: settings.audioMixPreset || 'auto',
|
||||
upscaleEnabled: settings.upscaleEnabled || false,
|
||||
upscaleMethod: settings.upscaleMethod || 'hardware',
|
||||
upscaleTarget: settings.upscaleTarget || '1080p',
|
||||
videoMode: videoMode,
|
||||
videoCodec: videoCodec,
|
||||
audioCodec: audioCodec,
|
||||
audioChannels: audioChannels
|
||||
}).then(async s => { await s.start(); return s; }),
|
||||
probeDuration(url, ffprobePath, userAgent)
|
||||
]);
|
||||
|
||||
// Wait for playlist to be ready (first segments generated)
|
||||
const ready = await session.waitForPlaylist(15000);
|
||||
// VOD pesado (4K HEVC/remux) puede tardar >15 s en generar el primer segmento
|
||||
const waitTimeout = streamInput.isHeavyVodUrl(url) ? 60000 : 15000;
|
||||
let ready = await session.waitForPlaylist(waitTimeout);
|
||||
|
||||
// Fallback a software si el encoder HW falla inmediatamente (error rápido)
|
||||
if (!ready && session.status === 'error' && (settings.hwEncoder || 'software') !== 'software') {
|
||||
console.warn(`[Transcode] HW encoder (${settings.hwEncoder}) falló, reintentando con software libx264...`);
|
||||
await transcodeSession.removeSession(session.id);
|
||||
const swSession = await transcodeSession.createSession(url, {
|
||||
ffmpegPath,
|
||||
userAgent,
|
||||
seekOffset: seekOffset || 0,
|
||||
hwEncoder: 'software',
|
||||
maxResolution: settings.maxResolution || '1080p',
|
||||
quality: settings.quality || 'medium',
|
||||
audioMixPreset: settings.audioMixPreset || 'auto',
|
||||
upscaleEnabled: false,
|
||||
videoMode: videoMode,
|
||||
videoCodec: videoCodec,
|
||||
audioCodec: audioCodec,
|
||||
audioChannels: audioChannels
|
||||
});
|
||||
await swSession.start();
|
||||
ready = await swSession.waitForPlaylist(waitTimeout);
|
||||
if (ready) {
|
||||
console.log(`[Transcode] Software fallback OK → session ${swSession.id}`);
|
||||
return res.json({
|
||||
sessionId: swSession.id,
|
||||
playlistUrl: `/api/transcode/${swSession.id}/stream.m3u8`,
|
||||
status: swSession.status,
|
||||
durationSec: durationSec || 0
|
||||
});
|
||||
}
|
||||
await transcodeSession.removeSession(swSession.id);
|
||||
}
|
||||
|
||||
if (!ready) {
|
||||
await transcodeSession.removeSession(session.id);
|
||||
return res.status(500).json({ error: 'Transcoding failed to start', reason: 'Playlist not generated in time' });
|
||||
}
|
||||
|
||||
console.log(`[Transcode] Session ${session.id} ready. Source duration: ${durationSec}s`);
|
||||
|
||||
res.json({
|
||||
sessionId: session.id,
|
||||
playlistUrl: `/api/transcode/${session.id}/stream.m3u8`,
|
||||
status: session.status
|
||||
status: session.status,
|
||||
durationSec: durationSec || 0
|
||||
});
|
||||
|
||||
} catch (err) {
|
||||
@@ -144,6 +219,30 @@ router.delete('/:sessionId', async (req, res) => {
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Heartbeat — keeps a session alive while the player is active.
|
||||
* POST /api/transcode/:sessionId/heartbeat
|
||||
*/
|
||||
router.post('/:sessionId/heartbeat', (req, res) => {
|
||||
const session = transcodeSession.getSession(req.params.sessionId);
|
||||
if (!session) return res.status(404).json({ error: 'Session not found' });
|
||||
session.lastAccess = Date.now();
|
||||
res.json({ ok: true });
|
||||
});
|
||||
|
||||
/**
|
||||
* Stop session via POST (for sendBeacon on page close)
|
||||
* POST /api/transcode/:sessionId/stop
|
||||
*/
|
||||
router.post('/:sessionId/stop', async (req, res) => {
|
||||
try {
|
||||
await transcodeSession.removeSession(req.params.sessionId);
|
||||
res.json({ success: true });
|
||||
} catch (err) {
|
||||
res.status(500).json({ error: 'Failed to remove session', details: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* List all active sessions (for debugging)
|
||||
* GET /api/transcode/sessions
|
||||
@@ -175,6 +274,9 @@ router.get('/', async (req, res) => {
|
||||
console.log(`[Transcode] Using User-Agent: ${settings.userAgentPreset}`);
|
||||
console.log(`[Transcode] Using binary: ${ffmpegPath}`);
|
||||
|
||||
const sizing = streamInput.probeSizingForUrl(url);
|
||||
const referer = streamInput.deriveRefererFromUrl(url);
|
||||
|
||||
// FFmpeg arguments for transcoding
|
||||
// Optimized for VOD content with incompatible audio (Dolby/AC3/EAC3)
|
||||
// Also works for live streams with ad stitching (Pluto TV, etc.)
|
||||
@@ -182,9 +284,13 @@ router.get('/', async (req, res) => {
|
||||
'-hide_banner',
|
||||
'-loglevel', 'warning',
|
||||
'-user_agent', userAgent,
|
||||
// Faster startup - reduced probe/analyze for quicker first bytes
|
||||
'-probesize', '2000000', // 2MB (reduced from 5MB)
|
||||
'-analyzeduration', '3000000', // 3 seconds (reduced from 10s)
|
||||
];
|
||||
if (referer) {
|
||||
args.push('-referer', referer);
|
||||
}
|
||||
args.push(
|
||||
'-probesize', sizing.probesize,
|
||||
'-analyzeduration', sizing.analyzeduration,
|
||||
// Error resilience: generate timestamps, discard corrupt packets
|
||||
'-fflags', '+genpts+discardcorrupt+nobuffer',
|
||||
// Ignore errors in stream and continue
|
||||
@@ -197,7 +303,9 @@ router.get('/', async (req, res) => {
|
||||
'-reconnect_delay_max', '3',
|
||||
// Prevent Range/HEAD requests that some providers reject with 405
|
||||
'-seekable', '0',
|
||||
'-i', url,
|
||||
'-i', url
|
||||
);
|
||||
args.push(
|
||||
// Map only first video and audio stream (avoid subtitle streams causing issues)
|
||||
'-map', '0:v:0',
|
||||
'-map', '0:a:0?', // ? makes audio optional if not present
|
||||
@@ -218,7 +326,7 @@ router.get('/', async (req, res) => {
|
||||
'-movflags', 'frag_keyframe+empty_moov+default_base_moof+faststart',
|
||||
'-flush_packets', '1', // Send data immediately
|
||||
'-' // Output to stdout
|
||||
];
|
||||
);
|
||||
|
||||
console.log(`[Transcode] Full command: ${ffmpegPath} ${args.join(' ')}`);
|
||||
|
||||
|
||||
@@ -228,16 +228,18 @@ async function detect() {
|
||||
detectAMF()
|
||||
]);
|
||||
|
||||
// Determine recommended encoder (priority: NVENC > AMF > QSV > VAAPI > Software)
|
||||
// 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 (qsv.available) {
|
||||
recommended = 'qsv';
|
||||
} else if (vaapi.available) {
|
||||
recommended = 'vaapi';
|
||||
} else if (qsv.available) {
|
||||
recommended = 'qsv';
|
||||
}
|
||||
|
||||
hwCapabilities = {
|
||||
|
||||
@@ -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
|
||||
};
|
||||
@@ -18,6 +18,7 @@ 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();
|
||||
@@ -26,9 +27,9 @@ const sessions = new Map();
|
||||
const CACHE_DIR = path.join(process.cwd(), 'transcode-cache');
|
||||
|
||||
// Session settings
|
||||
const SESSION_TIMEOUT_MS = 30 * 60 * 1000; // 30 minutes idle timeout
|
||||
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 = 5 * 60 * 1000; // Check every 5 minutes
|
||||
const CLEANUP_INTERVAL_MS = 60 * 1000; // Check every 60 seconds
|
||||
|
||||
/**
|
||||
* Generate a unique session ID
|
||||
@@ -181,11 +182,17 @@ class TranscodeSession extends EventEmitter {
|
||||
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') {
|
||||
@@ -194,8 +201,8 @@ class TranscodeSession extends EventEmitter {
|
||||
|
||||
// Input options (common)
|
||||
args.push(
|
||||
'-probesize', '5000000',
|
||||
'-analyzeduration', '5000000',
|
||||
'-probesize', sizing.probesize,
|
||||
'-analyzeduration', sizing.analyzeduration,
|
||||
'-fflags', '+genpts+discardcorrupt',
|
||||
'-err_detect', 'ignore_err',
|
||||
'-reconnect', '1',
|
||||
@@ -222,7 +229,10 @@ class TranscodeSession extends EventEmitter {
|
||||
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') {
|
||||
args.push('-bsf:v', 'h264_mp4toannexb');
|
||||
// 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');
|
||||
@@ -296,12 +306,12 @@ class TranscodeSession extends EventEmitter {
|
||||
);
|
||||
break;
|
||||
case 'vaapi':
|
||||
// VAAPI hardware decoding (Linux)
|
||||
args.push(
|
||||
'-hwaccel', 'vaapi',
|
||||
'-hwaccel_device', '/dev/dri/renderD128',
|
||||
'-hwaccel_output_format', '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
|
||||
@@ -404,7 +414,9 @@ class TranscodeSession extends EventEmitter {
|
||||
// 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`;
|
||||
// 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':
|
||||
@@ -466,14 +478,11 @@ class TranscodeSession extends EventEmitter {
|
||||
// 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
|
||||
// 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',
|
||||
'-profile:v', 'main', // Use main profile for compatibility
|
||||
'-global_quality', String(qp),
|
||||
'-bf', '3',
|
||||
'-pix_fmt', 'yuv420p' // Force 8-bit output for compatibility
|
||||
'-global_quality', String(qp)
|
||||
);
|
||||
}
|
||||
|
||||
@@ -551,6 +560,7 @@ class TranscodeSession extends EventEmitter {
|
||||
|
||||
/**
|
||||
* 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();
|
||||
@@ -558,6 +568,11 @@ class TranscodeSession extends EventEmitter {
|
||||
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;
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
/**
|
||||
* Who may use a source for playback and API access.
|
||||
* ownerId null/undefined: legacy shared — any authenticated user may use it.
|
||||
* ownerId N: only that user id (and admins) may use it.
|
||||
*/
|
||||
function canUserAccessSource(user, source) {
|
||||
if (!user || !source) return false;
|
||||
if (user.role === 'admin') return true;
|
||||
const oid = source.ownerId;
|
||||
if (oid == null || oid === undefined) return true;
|
||||
return Number(oid) === Number(user.id);
|
||||
}
|
||||
|
||||
async function getAccessibleSourceIds(req, sourcesApi) {
|
||||
const all = await sourcesApi.getAll();
|
||||
return all.filter((s) => canUserAccessSource(req.user, s)).map((s) => s.id);
|
||||
}
|
||||
|
||||
/** Returns false after sending res if missing, forbidden, or invalid. */
|
||||
async function assertCanUseSourceById(req, res, sourcesApi, rawSourceId) {
|
||||
const sid = typeof rawSourceId === 'number' ? rawSourceId : parseInt(rawSourceId, 10);
|
||||
if (Number.isNaN(sid)) {
|
||||
res.status(400).json({ error: 'Invalid source id' });
|
||||
return false;
|
||||
}
|
||||
const source = await sourcesApi.getById(sid);
|
||||
if (!source) {
|
||||
res.status(404).json({ error: 'Source not found' });
|
||||
return false;
|
||||
}
|
||||
if (!canUserAccessSource(req.user, source)) {
|
||||
res.status(403).json({ error: 'No access to this source' });
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
canUserAccessSource,
|
||||
getAccessibleSourceIds,
|
||||
assertCanUseSourceById
|
||||
};
|
||||
@@ -0,0 +1,38 @@
|
||||
'use strict';
|
||||
|
||||
/**
|
||||
* Opciones de entrada HTTP para ffprobe/ffmpeg con orígenes tipo Xtream / Jellyfin (Dispatcharr).
|
||||
* Muchos backends exigen Referer/origin coherentes además del User-Agent.
|
||||
*/
|
||||
|
||||
function deriveRefererFromUrl(inputUrl) {
|
||||
try {
|
||||
const u = new URL(inputUrl);
|
||||
return `${u.origin}/`;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** VOD pesado (Remux Jellyfin MKV/MP4, GOP largo) — probes más grandes. */
|
||||
function isHeavyVodUrl(url) {
|
||||
if (!url || typeof url !== 'string') return false;
|
||||
const u = url.toLowerCase();
|
||||
if (u.includes('/movie/') || u.includes('/series/')) return true;
|
||||
if (/\.mkv\b/i.test(url) || /\.mp4\b/i.test(url)) return true;
|
||||
return false;
|
||||
}
|
||||
|
||||
function probeSizingForUrl(url) {
|
||||
if (isHeavyVodUrl(url)) {
|
||||
/** Microseconds; VOD Jellyfin/remux pueden necesitar leer bastante antes del analyze */
|
||||
return { probesize: '32000000', analyzeduration: '45000000', heavy: true };
|
||||
}
|
||||
return { probesize: '8000000', analyzeduration: '8000000', heavy: false };
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
deriveRefererFromUrl,
|
||||
isHeavyVodUrl,
|
||||
probeSizingForUrl
|
||||
};
|
||||
Reference in New Issue
Block a user