mirror of
https://github.com/legop3/MultiRoombaRover.git
synced 2026-09-16 17:40:46 -04:00
whipwhep
This commit is contained in:
@@ -8,16 +8,13 @@ const { loadConfig } = require('../helpers/configLoader');
|
|||||||
const roverManager = require('./roverManager');
|
const roverManager = require('./roverManager');
|
||||||
const { isVerified } = require('./verificationService');
|
const { isVerified } = require('./verificationService');
|
||||||
const turnService = require('./turnService');
|
const turnService = require('./turnService');
|
||||||
const videoSessions = require('./videoSessions');
|
|
||||||
|
|
||||||
const audioForwardEvents = new EventEmitter();
|
const audioForwardEvents = new EventEmitter();
|
||||||
const config = loadConfig();
|
const config = loadConfig();
|
||||||
const audioForwardConfig = config.audioForward || {};
|
const audioForwardConfig = config.audioForward || {};
|
||||||
const mediaConfig = config.media || {};
|
|
||||||
const serviceEnabled = audioForwardConfig.enabled !== false;
|
const serviceEnabled = audioForwardConfig.enabled !== false;
|
||||||
const ffmpegBin = audioForwardConfig.ffmpegBin || 'ffmpeg';
|
const ffmpegBin = audioForwardConfig.ffmpegBin || 'ffmpeg';
|
||||||
const streamSuffix = typeof audioForwardConfig.streamSuffix === 'string' ? audioForwardConfig.streamSuffix : '-fwd';
|
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 runtimeDir = path.resolve(audioForwardConfig.runtimeDir || '/tmp/mrr-audio-forward');
|
||||||
const uploadsDir = path.join(runtimeDir, 'uploads');
|
const uploadsDir = path.join(runtimeDir, 'uploads');
|
||||||
const maxUploadBytes = Number.isFinite(audioForwardConfig.maxUploadBytes)
|
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`;
|
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 = {}) {
|
function spawnProcess(roverId, tag, args, options = {}) {
|
||||||
const proc = spawn(ffmpegBin, args, {
|
const proc = spawn(ffmpegBin, args, {
|
||||||
stdio: [options.captureStdin ? 'pipe' : 'ignore', options.captureStdout ? 'pipe' : 'ignore', 'pipe'],
|
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) {
|
function attachWriterPipe(worker, proc) {
|
||||||
const writer = fs.createWriteStream(worker.fifoPath, { flags: 'w' });
|
const writer = fs.createWriteStream(worker.fifoPath, { flags: 'w' });
|
||||||
writer.on('error', (err) => {
|
writer.on('error', (err) => {
|
||||||
@@ -347,10 +292,6 @@ function cleanupUploadFile(worker) {
|
|||||||
|
|
||||||
function stopContentWriter(worker) {
|
function stopContentWriter(worker) {
|
||||||
if (!worker) return;
|
if (!worker) return;
|
||||||
if (worker.micWhipRestartTimer) {
|
|
||||||
clearTimeout(worker.micWhipRestartTimer);
|
|
||||||
worker.micWhipRestartTimer = null;
|
|
||||||
}
|
|
||||||
if (worker.micIdleTimer) {
|
if (worker.micIdleTimer) {
|
||||||
clearTimeout(worker.micIdleTimer);
|
clearTimeout(worker.micIdleTimer);
|
||||||
worker.micIdleTimer = null;
|
worker.micIdleTimer = null;
|
||||||
@@ -383,7 +324,6 @@ function stopContentWriter(worker) {
|
|||||||
worker.contentProc = null;
|
worker.contentProc = null;
|
||||||
worker.contentKind = null;
|
worker.contentKind = null;
|
||||||
worker.activeOwnerSocketId = null;
|
worker.activeOwnerSocketId = null;
|
||||||
worker.micWhipPathId = null;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
function startSilenceWriter(roverId) {
|
function startSilenceWriter(roverId) {
|
||||||
@@ -505,68 +445,6 @@ function startMicWriter(roverId, ownerSocketId = null) {
|
|||||||
setState(roverId, { state: 'playing', source: 'mic', error: null, startedAt: Date.now() });
|
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 = {}) {
|
function decodeMicChunk(payload = {}) {
|
||||||
const binary = payload?.data;
|
const binary = payload?.data;
|
||||||
if (Buffer.isBuffer(binary)) {
|
if (Buffer.isBuffer(binary)) {
|
||||||
@@ -653,8 +531,6 @@ function ensureWorker(roverId) {
|
|||||||
activeOwnerSocketId: null,
|
activeOwnerSocketId: null,
|
||||||
activeUploadPath: null,
|
activeUploadPath: null,
|
||||||
micWriter: null,
|
micWriter: null,
|
||||||
micWhipPathId: null,
|
|
||||||
micWhipRestartTimer: null,
|
|
||||||
micLastChunkAt: 0,
|
micLastChunkAt: 0,
|
||||||
micIdleTimer: null,
|
micIdleTimer: null,
|
||||||
micBackpressured: false,
|
micBackpressured: false,
|
||||||
@@ -763,18 +639,12 @@ function stopOwnedAudioIfUnauthorized(roverId, ownerSocketId, reason = 'driver_c
|
|||||||
if (!roverId || !ownerSocketId) return;
|
if (!roverId || !ownerSocketId) return;
|
||||||
const worker = workers.get(roverId);
|
const worker = workers.get(roverId);
|
||||||
if (!worker) return;
|
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;
|
if (worker.activeOwnerSocketId !== ownerSocketId) return;
|
||||||
const ownerSocket = io.sockets.sockets.get(ownerSocketId);
|
const ownerSocket = io.sockets.sockets.get(ownerSocketId);
|
||||||
const ownerIsDriver = ownerSocket ? roverManager.isDriver(roverId, ownerSocket) : false;
|
const ownerIsDriver = ownerSocket ? roverManager.isDriver(roverId, ownerSocket) : false;
|
||||||
const ownerCanDrive = ownerSocket ? turnService.canDrive(roverId, ownerSocket) : false;
|
const ownerCanDrive = ownerSocket ? turnService.canDrive(roverId, ownerSocket) : false;
|
||||||
if (ownerIsDriver && ownerCanDrive) return;
|
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 });
|
logger.info('Stopping audio forward due to ownership/driver change', { roverId, ownerSocketId, reason, source: worker.contentKind });
|
||||||
startSilenceWriter(roverId);
|
startSilenceWriter(roverId);
|
||||||
}
|
}
|
||||||
@@ -789,7 +659,7 @@ roverManager.managerEvents.on('driver', ({ socketId, roverId, action } = {}) =>
|
|||||||
turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => {
|
turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => {
|
||||||
if (!roverId) return;
|
if (!roverId) return;
|
||||||
const worker = workers.get(roverId);
|
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');
|
stopOwnedAudioIfUnauthorized(roverId, worker.activeOwnerSocketId, 'turn_change');
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -856,60 +726,24 @@ io.on('connection', (socket) => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
socket.on('audio:micWhipStart', ({ roverId } = {}, cb = () => {}) => {
|
socket.on('audio:micWhipStart', ({ roverId } = {}, cb = () => {}) => {
|
||||||
try {
|
cb({ error: 'WHIP mic forwarding is temporarily disabled; using socket fallback.' });
|
||||||
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 });
|
|
||||||
}
|
|
||||||
});
|
});
|
||||||
|
|
||||||
socket.on('audio:micWhipReady', ({ roverId } = {}, cb = () => {}) => {
|
socket.on('audio:micWhipReady', ({ roverId } = {}, cb = () => {}) => {
|
||||||
try {
|
cb({ error: 'WHIP mic forwarding is temporarily disabled; using socket fallback.' });
|
||||||
const normalized = String(roverId || '').trim();
|
|
||||||
ensureAudioForwardPermission(socket, normalized);
|
|
||||||
startMicWhipRelay(normalized, socket.id);
|
|
||||||
cb({ success: true, roverId: normalized });
|
|
||||||
} catch (err) {
|
|
||||||
cb({ error: err.message });
|
|
||||||
}
|
|
||||||
});
|
});
|
||||||
|
|
||||||
socket.on('audio:micWhipStop', ({ roverId } = {}, cb = () => {}) => {
|
socket.on('audio:micWhipStop', ({ roverId } = {}, cb = () => {}) => {
|
||||||
try {
|
cb({ success: true });
|
||||||
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 });
|
|
||||||
}
|
|
||||||
});
|
});
|
||||||
|
|
||||||
socket.on('disconnect', () => {
|
socket.on('disconnect', () => {
|
||||||
workers.forEach((worker, roverId) => {
|
workers.forEach((worker, roverId) => {
|
||||||
if (!worker || worker.activeOwnerSocketId !== socket.id) return;
|
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 });
|
logger.info('Stopping owned audio forward due to socket disconnect', { roverId, socketId: socket.id, source: worker.contentKind });
|
||||||
startSilenceWriter(roverId);
|
startSilenceWriter(roverId);
|
||||||
});
|
});
|
||||||
videoSessions.revokeWhere((info) => info?.socketId === socket.id && info?.sourceType === 'roverMic');
|
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user