feat: implement GPU encoder pipelines
This commit is contained in:
+142
-1
@@ -1,10 +1,150 @@
|
||||
const express = require('express');
|
||||
const router = express.Router();
|
||||
const { spawn } = require('child_process');
|
||||
const path = require('path');
|
||||
const fs = require('fs').promises;
|
||||
const db = require('../db');
|
||||
const transcodeSession = require('../services/transcodeSession');
|
||||
|
||||
/**
|
||||
* Transcode stream
|
||||
* Transcode Routes
|
||||
*
|
||||
* Direct streaming (backward compatible):
|
||||
* GET /api/transcode?url=...
|
||||
*
|
||||
* HLS session-based (new, supports seeking):
|
||||
* POST /api/transcode/session - Create new session
|
||||
* GET /api/transcode/:id/stream.m3u8 - Get HLS playlist
|
||||
* GET /api/transcode/:id/:segment.ts - Get segment file
|
||||
* DELETE /api/transcode/:id - Stop and cleanup session
|
||||
* GET /api/transcode/sessions - List all sessions (debug)
|
||||
*/
|
||||
|
||||
// Start session cleanup interval
|
||||
transcodeSession.startCleanupInterval();
|
||||
|
||||
/**
|
||||
* Create a new transcode session
|
||||
* POST /api/transcode/session
|
||||
* Body: { url: string, seekOffset?: number }
|
||||
*/
|
||||
router.post('/session', async (req, res) => {
|
||||
const { url, seekOffset } = req.body;
|
||||
|
||||
if (!url) {
|
||||
return res.status(400).json({ error: 'URL is required' });
|
||||
}
|
||||
|
||||
const ffmpegPath = req.app.locals.ffmpegPath || 'ffmpeg';
|
||||
const settings = await db.settings.get();
|
||||
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'
|
||||
});
|
||||
|
||||
await session.start();
|
||||
|
||||
// Wait for playlist to be ready (first segments generated)
|
||||
const ready = await session.waitForPlaylist(15000);
|
||||
|
||||
if (!ready) {
|
||||
await transcodeSession.removeSession(session.id);
|
||||
return res.status(500).json({ error: 'Transcoding failed to start', reason: 'Playlist not generated in time' });
|
||||
}
|
||||
|
||||
res.json({
|
||||
sessionId: session.id,
|
||||
playlistUrl: `/api/transcode/${session.id}/stream.m3u8`,
|
||||
status: session.status
|
||||
});
|
||||
|
||||
} catch (err) {
|
||||
console.error('[Transcode] Session creation failed:', err);
|
||||
res.status(500).json({ error: 'Failed to create session', details: err.message });
|
||||
}
|
||||
});
|
||||
|
||||
/**
|
||||
* Get HLS playlist for a session
|
||||
* GET /api/transcode/:sessionId/stream.m3u8
|
||||
*/
|
||||
router.get('/:sessionId/stream.m3u8', async (req, res) => {
|
||||
const { sessionId } = req.params;
|
||||
const session = transcodeSession.getSession(sessionId);
|
||||
|
||||
if (!session) {
|
||||
return res.status(404).json({ error: 'Session not found' });
|
||||
}
|
||||
|
||||
const playlist = await session.getPlaylist();
|
||||
if (!playlist) {
|
||||
return res.status(404).json({ error: 'Playlist not ready' });
|
||||
}
|
||||
|
||||
res.setHeader('Content-Type', 'application/vnd.apple.mpegurl');
|
||||
res.setHeader('Cache-Control', 'no-cache');
|
||||
res.send(playlist);
|
||||
});
|
||||
|
||||
/**
|
||||
* Get a segment file for a session
|
||||
* GET /api/transcode/:sessionId/:segment.ts
|
||||
*/
|
||||
router.get('/:sessionId/:segment', async (req, res) => {
|
||||
const { sessionId, segment } = req.params;
|
||||
|
||||
// Only handle .ts files
|
||||
if (!segment.endsWith('.ts')) {
|
||||
return res.status(404).json({ error: 'Invalid segment' });
|
||||
}
|
||||
|
||||
const session = transcodeSession.getSession(sessionId);
|
||||
if (!session) {
|
||||
return res.status(404).json({ error: 'Session not found' });
|
||||
}
|
||||
|
||||
const segmentPath = await session.getSegment(segment);
|
||||
if (!segmentPath) {
|
||||
return res.status(404).json({ error: 'Segment not found' });
|
||||
}
|
||||
|
||||
res.setHeader('Content-Type', 'video/MP2T');
|
||||
res.setHeader('Cache-Control', 'public, max-age=31536000'); // Cache forever (immutable)
|
||||
res.sendFile(segmentPath);
|
||||
});
|
||||
|
||||
/**
|
||||
* Stop and cleanup a session
|
||||
* DELETE /api/transcode/:sessionId
|
||||
*/
|
||||
router.delete('/:sessionId', async (req, res) => {
|
||||
const { sessionId } = req.params;
|
||||
|
||||
try {
|
||||
await transcodeSession.removeSession(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
|
||||
*/
|
||||
router.get('/sessions', (req, res) => {
|
||||
res.json(transcodeSession.getAllSessions());
|
||||
});
|
||||
|
||||
/**
|
||||
* Direct transcode stream (backward compatible, no seeking)
|
||||
* GET /api/transcode?url=...
|
||||
*
|
||||
* Transcodes audio to AAC for browser compatibility while passing video through.
|
||||
@@ -121,3 +261,4 @@ router.get('/', async (req, res) => {
|
||||
});
|
||||
|
||||
module.exports = router;
|
||||
|
||||
|
||||
@@ -0,0 +1,658 @@
|
||||
/**
|
||||
* 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');
|
||||
|
||||
// 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',
|
||||
...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.ts');
|
||||
const encoder = this.options.hwEncoder || 'software';
|
||||
|
||||
const args = [
|
||||
'-hide_banner',
|
||||
'-loglevel', 'warning',
|
||||
'-user_agent', this.options.userAgent,
|
||||
];
|
||||
|
||||
// Add hardware acceleration input options based on encoder
|
||||
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',
|
||||
'-seekable', '0'
|
||||
);
|
||||
|
||||
// Add seek offset if specified
|
||||
if (this.options.seekOffset > 0) {
|
||||
args.push('-ss', String(this.options.seekOffset));
|
||||
}
|
||||
|
||||
args.push('-i', this.url);
|
||||
|
||||
// Map streams
|
||||
args.push('-map', '0:v:0');
|
||||
args.push('-map', '0:a:0?');
|
||||
|
||||
// Add video encoder and filters based on selected encoder
|
||||
this.addVideoEncoderArgs(args, encoder);
|
||||
|
||||
// Audio: Transcode to AAC
|
||||
args.push(
|
||||
'-c:a', 'aac',
|
||||
'-ar', '48000',
|
||||
'-b:a', '192k'
|
||||
);
|
||||
|
||||
// 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', segmentPattern,
|
||||
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 setting
|
||||
*/
|
||||
getTargetHeight() {
|
||||
const resolutionMap = {
|
||||
'4k': 2160,
|
||||
'1080p': 1080,
|
||||
'720p': 720,
|
||||
'480p': 480
|
||||
};
|
||||
return resolutionMap[this.options.maxResolution] || 1080;
|
||||
}
|
||||
|
||||
/**
|
||||
* NVIDIA NVENC encoder arguments
|
||||
*/
|
||||
addNvencEncoderArgs(args, height, qp) {
|
||||
// Video filter for scaling on GPU
|
||||
args.push('-vf', `scale_cuda=-2:${height}:interp_algo=lanczos`);
|
||||
|
||||
args.push(
|
||||
'-c:v', 'h264_nvenc',
|
||||
'-preset', 'p4', // Balanced preset
|
||||
'-tune', 'hq', // High quality tuning
|
||||
'-rc', 'constqp', // Constant QP mode
|
||||
'-qp', String(qp),
|
||||
'-rc-lookahead', '32',
|
||||
'-bf', '3', // B-frames
|
||||
'-b_ref_mode', 'middle',
|
||||
'-spatial-aq', '1',
|
||||
'-temporal-aq', '1'
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* AMD AMF encoder arguments
|
||||
*/
|
||||
addAmfEncoderArgs(args, height, qp) {
|
||||
// CPU decoding + software scale + AMF encode
|
||||
args.push('-vf', `scale=-2:${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)
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* VAAPI encoder arguments (Linux)
|
||||
*/
|
||||
addVaapiEncoderArgs(args, height, qp) {
|
||||
// Scale on VAAPI
|
||||
args.push('-vf', `scale_vaapi=w=-2:h=${height}`);
|
||||
|
||||
args.push(
|
||||
'-c:v', 'h264_vaapi',
|
||||
'-rc_mode', 'CQP',
|
||||
'-qp', String(qp),
|
||||
'-bf', '3'
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Intel QuickSync encoder arguments
|
||||
*/
|
||||
addQsvEncoderArgs(args, height, qp) {
|
||||
// Scale on QSV
|
||||
args.push('-vf', `scale_qsv=w=-2:h=${height}`);
|
||||
|
||||
args.push(
|
||||
'-c:v', 'h264_qsv',
|
||||
'-preset', 'medium',
|
||||
'-global_quality', String(qp),
|
||||
'-look_ahead', '1',
|
||||
'-look_ahead_depth', '40'
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Software encoder arguments (fallback)
|
||||
*/
|
||||
addSoftwareEncoderArgs(args, height, crf) {
|
||||
// Software scaling
|
||||
args.push('-vf', `scale=-2:${height}`);
|
||||
|
||||
args.push(
|
||||
'-c:v', 'libx264',
|
||||
'-preset', 'veryfast', // Fast for real-time
|
||||
'-crf', String(crf),
|
||||
'-profile:v', 'high',
|
||||
'-level', '4.1'
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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)
|
||||
*/
|
||||
async waitForPlaylist(timeoutMs = 10000) {
|
||||
const startTime = Date.now();
|
||||
while (Date.now() - startTime < timeoutMs) {
|
||||
if (await this.isPlaylistReady()) {
|
||||
return true;
|
||||
}
|
||||
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
|
||||
};
|
||||
Reference in New Issue
Block a user