diff --git a/server/routes/transcode.js b/server/routes/transcode.js index f1e25d9..2ffbb05 100644 --- a/server/routes/transcode.js +++ b/server/routes/transcode.js @@ -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; + diff --git a/server/services/transcodeSession.js b/server/services/transcodeSession.js new file mode 100644 index 0000000..6f11c93 --- /dev/null +++ b/server/services/transcodeSession.js @@ -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 +};