From 57532f53d03c70723a2d6a9005eaafd76c199d85 Mon Sep 17 00:00:00 2001 From: legop3 Date: Fri, 20 Mar 2026 02:25:43 -0400 Subject: [PATCH] whipwhep --- server/src/services/audioForwardService.js | 178 +-------------------- 1 file changed, 6 insertions(+), 172 deletions(-) diff --git a/server/src/services/audioForwardService.js b/server/src/services/audioForwardService.js index 3ffdfd8b..9aba4188 100644 --- a/server/src/services/audioForwardService.js +++ b/server/src/services/audioForwardService.js @@ -8,16 +8,13 @@ const { loadConfig } = require('../helpers/configLoader'); const roverManager = require('./roverManager'); const { isVerified } = require('./verificationService'); const turnService = require('./turnService'); -const videoSessions = require('./videoSessions'); const audioForwardEvents = new EventEmitter(); const config = loadConfig(); const audioForwardConfig = config.audioForward || {}; -const mediaConfig = config.media || {}; const serviceEnabled = audioForwardConfig.enabled !== false; const ffmpegBin = audioForwardConfig.ffmpegBin || 'ffmpeg'; const streamSuffix = typeof audioForwardConfig.streamSuffix === 'string' ? audioForwardConfig.streamSuffix : '-fwd'; -const micSuffix = typeof audioForwardConfig.micSuffix === 'string' ? audioForwardConfig.micSuffix : '-mic'; const runtimeDir = path.resolve(audioForwardConfig.runtimeDir || '/tmp/mrr-audio-forward'); const uploadsDir = path.join(runtimeDir, 'uploads'); const maxUploadBytes = Number.isFinite(audioForwardConfig.maxUploadBytes) @@ -146,34 +143,6 @@ function resolveForwardUrl(roverId) { return `srt://127.0.0.1:9000?streamid=#!::r=${encodeURIComponent(roverId + streamSuffix)},m=publish&latency=10&mode=caller&transtype=live&pkt_size=1316`; } -function resolveMicPathId(roverId) { - return `${roverId}${micSuffix}`; -} - -function resolveMicReadUrl(roverId) { - const pathId = resolveMicPathId(roverId); - return `srt://127.0.0.1:9000?streamid=#!::r=${encodeURIComponent(pathId)},m=read&latency=10&mode=caller&transtype=live&pkt_size=1316`; -} - -function getMediaPrefix() { - const base = mediaConfig.whepBaseUrl; - if (!base) return ''; - try { - const parsed = new URL(base); - return `${parsed.origin}${parsed.pathname}`.replace(/\/+$/, ''); - } catch { - return String(base).replace(/\/+$/, ''); - } -} - -function buildWhipUrl(pathId) { - const prefix = getMediaPrefix(); - if (!prefix) { - throw new Error('Server media base URL missing'); - } - return `${prefix}/${encodeURIComponent(pathId)}/whip`; -} - function spawnProcess(roverId, tag, args, options = {}) { const proc = spawn(ffmpegBin, args, { stdio: [options.captureStdin ? 'pipe' : 'ignore', options.captureStdout ? 'pipe' : 'ignore', 'pipe'], @@ -289,30 +258,6 @@ function buildUploadWriterArgs(filePath) { ]; } -function buildWhipRelayReaderArgs(inputUrl) { - return [ - '-hide_banner', - '-loglevel', - 'warning', - '-fflags', - 'nobuffer', - '-flags', - 'low_delay', - '-i', - inputUrl, - '-vn', - '-af', - 'aresample=16000', - '-f', - 's16le', - '-ac', - '1', - '-ar', - '16000', - 'pipe:1', - ]; -} - function attachWriterPipe(worker, proc) { const writer = fs.createWriteStream(worker.fifoPath, { flags: 'w' }); writer.on('error', (err) => { @@ -347,10 +292,6 @@ function cleanupUploadFile(worker) { function stopContentWriter(worker) { if (!worker) return; - if (worker.micWhipRestartTimer) { - clearTimeout(worker.micWhipRestartTimer); - worker.micWhipRestartTimer = null; - } if (worker.micIdleTimer) { clearTimeout(worker.micIdleTimer); worker.micIdleTimer = null; @@ -383,7 +324,6 @@ function stopContentWriter(worker) { worker.contentProc = null; worker.contentKind = null; worker.activeOwnerSocketId = null; - worker.micWhipPathId = null; } function startSilenceWriter(roverId) { @@ -505,68 +445,6 @@ function startMicWriter(roverId, ownerSocketId = null) { setState(roverId, { state: 'playing', source: 'mic', error: null, startedAt: Date.now() }); } -function startMicWhipRelay(roverId, ownerSocketId = null) { - const worker = workers.get(roverId); - if (!worker || worker.stopping) return; - if ( - worker.contentKind === 'mic_whip' && - worker.contentProc && - worker.activeOwnerSocketId === ownerSocketId - ) { - return; - } - - worker.contentKind = 'mic_whip'; - worker.activeOwnerSocketId = ownerSocketId; - worker.micWhipPathId = resolveMicPathId(roverId); - if (worker.micWhipRestartTimer) { - clearTimeout(worker.micWhipRestartTimer); - worker.micWhipRestartTimer = null; - } - stopContentWriter(worker); - worker.contentKind = 'mic_whip'; - worker.activeOwnerSocketId = ownerSocketId; - worker.micWhipPathId = resolveMicPathId(roverId); - setState(roverId, { state: 'starting', source: 'mic-whip', error: null, startedAt: Date.now() }); - - const launch = () => { - const current = workers.get(roverId); - if (!current || current.stopping) return; - if (current.contentKind !== 'mic_whip' || current.activeOwnerSocketId !== ownerSocketId) return; - - const inputUrl = resolveMicReadUrl(roverId); - const proc = spawnProcess(roverId, 'mic-whip-reader', buildWhipRelayReaderArgs(inputUrl), { - captureStdout: true, - }); - current.contentProc = proc; - const seq = ++current.writerSeq; - attachWriterPipe(current, proc); - setState(roverId, { state: 'playing', source: 'mic-whip', error: null, startedAt: Date.now() }); - - proc.on('exit', (code, signal) => { - const next = workers.get(roverId); - if (!next || next.stopping) return; - if (next.writerSeq !== seq || next.contentProc !== proc) return; - next.contentProc = null; - if (next.contentKind !== 'mic_whip' || next.activeOwnerSocketId !== ownerSocketId) return; - if (signal === 'SIGTERM') return; - - setState(roverId, { - state: 'starting', - source: 'mic-whip', - error: code != null && code !== 0 ? `relay reconnecting (last code=${code})` : null, - startedAt: Date.now(), - }); - next.micWhipRestartTimer = setTimeout(() => { - next.micWhipRestartTimer = null; - launch(); - }, 300); - }); - }; - - launch(); -} - function decodeMicChunk(payload = {}) { const binary = payload?.data; if (Buffer.isBuffer(binary)) { @@ -653,8 +531,6 @@ function ensureWorker(roverId) { activeOwnerSocketId: null, activeUploadPath: null, micWriter: null, - micWhipPathId: null, - micWhipRestartTimer: null, micLastChunkAt: 0, micIdleTimer: null, micBackpressured: false, @@ -763,18 +639,12 @@ function stopOwnedAudioIfUnauthorized(roverId, ownerSocketId, reason = 'driver_c if (!roverId || !ownerSocketId) return; const worker = workers.get(roverId); if (!worker) return; - if (worker.contentKind !== 'upload' && worker.contentKind !== 'mic' && worker.contentKind !== 'mic_whip') return; + if (worker.contentKind !== 'upload' && worker.contentKind !== 'mic') return; if (worker.activeOwnerSocketId !== ownerSocketId) return; const ownerSocket = io.sockets.sockets.get(ownerSocketId); const ownerIsDriver = ownerSocket ? roverManager.isDriver(roverId, ownerSocket) : false; const ownerCanDrive = ownerSocket ? turnService.canDrive(roverId, ownerSocket) : false; if (ownerIsDriver && ownerCanDrive) return; - if (worker.contentKind === 'mic_whip' && worker.micWhipPathId) { - const pathId = worker.micWhipPathId; - videoSessions.revokeWhere( - (info) => info?.socketId === ownerSocketId && info?.sourceType === 'roverMic' && info?.sourceId === pathId, - ); - } logger.info('Stopping audio forward due to ownership/driver change', { roverId, ownerSocketId, reason, source: worker.contentKind }); startSilenceWriter(roverId); } @@ -789,7 +659,7 @@ roverManager.managerEvents.on('driver', ({ socketId, roverId, action } = {}) => turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => { if (!roverId) return; const worker = workers.get(roverId); - if (!worker || (worker.contentKind !== 'upload' && worker.contentKind !== 'mic' && worker.contentKind !== 'mic_whip')) return; + if (!worker || (worker.contentKind !== 'upload' && worker.contentKind !== 'mic')) return; stopOwnedAudioIfUnauthorized(roverId, worker.activeOwnerSocketId, 'turn_change'); }); @@ -856,60 +726,24 @@ io.on('connection', (socket) => { }); socket.on('audio:micWhipStart', ({ roverId } = {}, cb = () => {}) => { - try { - const normalized = String(roverId || '').trim(); - ensureAudioForwardPermission(socket, normalized); - ensureWorker(normalized); - const pathId = resolveMicPathId(normalized); - videoSessions.revokeWhere( - (info) => info?.socketId === socket.id && info?.sourceType === 'roverMic' && info?.sourceId === pathId, - ); - const token = videoSessions.createSession(socket, { type: 'roverMic', id: pathId }); - const whipUrl = buildWhipUrl(pathId); - cb({ success: true, roverId: normalized, pathId, token, whipUrl }); - } catch (err) { - cb({ error: err.message }); - } + cb({ error: 'WHIP mic forwarding is temporarily disabled; using socket fallback.' }); }); socket.on('audio:micWhipReady', ({ roverId } = {}, cb = () => {}) => { - try { - const normalized = String(roverId || '').trim(); - ensureAudioForwardPermission(socket, normalized); - startMicWhipRelay(normalized, socket.id); - cb({ success: true, roverId: normalized }); - } catch (err) { - cb({ error: err.message }); - } + cb({ error: 'WHIP mic forwarding is temporarily disabled; using socket fallback.' }); }); socket.on('audio:micWhipStop', ({ roverId } = {}, cb = () => {}) => { - try { - const normalized = String(roverId || '').trim(); - ensureAudioForwardPermission(socket, normalized); - const worker = workers.get(normalized); - if (worker && worker.contentKind === 'mic_whip' && worker.activeOwnerSocketId !== socket.id) { - throw new Error('Mic forwarding is owned by another session'); - } - const pathId = resolveMicPathId(normalized); - videoSessions.revokeWhere( - (info) => info?.socketId === socket.id && info?.sourceType === 'roverMic' && info?.sourceId === pathId, - ); - stopPlayback(normalized); - cb({ success: true, roverId: normalized }); - } catch (err) { - cb({ error: err.message }); - } + cb({ success: true }); }); socket.on('disconnect', () => { workers.forEach((worker, roverId) => { if (!worker || worker.activeOwnerSocketId !== socket.id) return; - if (worker.contentKind !== 'upload' && worker.contentKind !== 'mic' && worker.contentKind !== 'mic_whip') return; + if (worker.contentKind !== 'upload' && worker.contentKind !== 'mic') return; logger.info('Stopping owned audio forward due to socket disconnect', { roverId, socketId: socket.id, source: worker.contentKind }); startSilenceWriter(roverId); }); - videoSessions.revokeWhere((info) => info?.socketId === socket.id && info?.sourceType === 'roverMic'); }); });