This commit is contained in:
legop3
2026-03-20 22:09:59 -04:00
parent 57532f53d0
commit f364e6ac60
2 changed files with 116 additions and 5 deletions
+105 -3
View File
@@ -8,10 +8,12 @@ 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';
@@ -23,6 +25,7 @@ const maxUploadBytes = Number.isFinite(audioForwardConfig.maxUploadBytes)
const states = new Map(); // roverId -> { state, source, error, startedAt, updatedAt } const states = new Map(); // roverId -> { state, source, error, startedAt, updatedAt }
const workers = new Map(); // roverId -> worker const workers = new Map(); // roverId -> worker
const whipOwners = new Map(); // roverId -> socketId
function publishStateChange(roverId) { function publishStateChange(roverId) {
audioForwardEvents.emit('change', { roverId, state: states.get(roverId) || null }); audioForwardEvents.emit('change', { roverId, state: states.get(roverId) || null });
@@ -143,6 +146,29 @@ 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 resolveForwardPathId(roverId) {
return `${roverId}${streamSuffix}`;
}
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'],
@@ -403,6 +429,7 @@ function scheduleMicIdleTimeout(roverId) {
} }
function startMicWriter(roverId, ownerSocketId = null) { function startMicWriter(roverId, ownerSocketId = null) {
stopWhipForRover(roverId, 'socket_mic_override');
const worker = workers.get(roverId); const worker = workers.get(roverId);
if (!worker || worker.stopping) return; if (!worker || worker.stopping) return;
if ( if (
@@ -578,6 +605,7 @@ function writeUploadFile(roverId, payload = {}) {
} }
function playUploadedAudio(roverId, payload = {}) { function playUploadedAudio(roverId, payload = {}) {
stopWhipForRover(roverId, 'upload_override');
const ownerSocketId = typeof payload?.ownerSocketId === 'string' ? payload.ownerSocketId : null; const ownerSocketId = typeof payload?.ownerSocketId === 'string' ? payload.ownerSocketId : null;
const uploadPath = writeUploadFile(roverId, payload); const uploadPath = writeUploadFile(roverId, payload);
ensureWorker(roverId); ensureWorker(roverId);
@@ -592,11 +620,35 @@ function playUploadedAudio(roverId, payload = {}) {
} }
function stopPlayback(roverId) { function stopPlayback(roverId) {
stopWhipForRover(roverId, 'stop_playback');
ensureWorker(roverId); ensureWorker(roverId);
startSilenceWriter(roverId); startSilenceWriter(roverId);
} }
function revokeWhipSessionForRover(roverId, ownerSocketId) {
if (!roverId || !ownerSocketId) return;
const pathId = resolveForwardPathId(roverId);
videoSessions.revokeWhere(
(info) => info?.socketId === ownerSocketId && info?.sourceType === 'roverMic' && info?.sourceId === pathId,
);
}
function stopWhipForRover(roverId, reason = 'unknown') {
const ownerSocketId = whipOwners.get(roverId);
if (!ownerSocketId) return;
whipOwners.delete(roverId);
revokeWhipSessionForRover(roverId, ownerSocketId);
logger.info('Stopping WHIP mic session', { roverId, ownerSocketId, reason });
try {
ensureWorker(roverId);
startSilenceWriter(roverId);
} catch (err) {
setState(roverId, { state: 'error', source: 'whip', error: err?.message || String(err), startedAt: null });
}
}
function stopWorker(roverId) { function stopWorker(roverId) {
stopWhipForRover(roverId, 'worker_stop');
const worker = workers.get(roverId); const worker = workers.get(roverId);
if (!worker) return; if (!worker) return;
worker.stopping = true; worker.stopping = true;
@@ -637,6 +689,14 @@ roverManager.managerEvents.on('rover', ({ roverId, action } = {}) => {
function stopOwnedAudioIfUnauthorized(roverId, ownerSocketId, reason = 'driver_change') { function stopOwnedAudioIfUnauthorized(roverId, ownerSocketId, reason = 'driver_change') {
if (!roverId || !ownerSocketId) return; if (!roverId || !ownerSocketId) return;
if (whipOwners.get(roverId) === ownerSocketId) {
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) {
stopWhipForRover(roverId, reason);
}
}
const worker = workers.get(roverId); const worker = workers.get(roverId);
if (!worker) return; if (!worker) return;
if (worker.contentKind !== 'upload' && worker.contentKind !== 'mic') return; if (worker.contentKind !== 'upload' && worker.contentKind !== 'mic') return;
@@ -658,6 +718,10 @@ roverManager.managerEvents.on('driver', ({ socketId, roverId, action } = {}) =>
turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => { turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => {
if (!roverId) return; if (!roverId) return;
const whipOwner = whipOwners.get(roverId);
if (whipOwner) {
stopOwnedAudioIfUnauthorized(roverId, whipOwner, 'turn_change');
}
const worker = workers.get(roverId); const worker = workers.get(roverId);
if (!worker || (worker.contentKind !== 'upload' && worker.contentKind !== 'mic')) return; if (!worker || (worker.contentKind !== 'upload' && worker.contentKind !== 'mic')) return;
stopOwnedAudioIfUnauthorized(roverId, worker.activeOwnerSocketId, 'turn_change'); stopOwnedAudioIfUnauthorized(roverId, worker.activeOwnerSocketId, 'turn_change');
@@ -726,15 +790,49 @@ io.on('connection', (socket) => {
}); });
socket.on('audio:micWhipStart', ({ roverId } = {}, cb = () => {}) => { socket.on('audio:micWhipStart', ({ roverId } = {}, cb = () => {}) => {
cb({ error: 'WHIP mic forwarding is temporarily disabled; using socket fallback.' }); try {
const normalized = String(roverId || '').trim();
ensureAudioForwardPermission(socket, normalized);
// WHIP publishes directly to the same forward path; stop local publisher to avoid path conflicts.
stopWorker(normalized);
whipOwners.set(normalized, socket.id);
const pathId = resolveForwardPathId(normalized);
revokeWhipSessionForRover(normalized, socket.id);
const token = videoSessions.createSession(socket, { type: 'roverMic', id: pathId });
const whipUrl = buildWhipUrl(pathId);
setState(normalized, { state: 'starting', source: 'mic-whip', error: null, startedAt: Date.now() });
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 = () => {}) => {
cb({ error: 'WHIP mic forwarding is temporarily disabled; using socket fallback.' }); try {
const normalized = String(roverId || '').trim();
ensureAudioForwardPermission(socket, normalized);
if (whipOwners.get(normalized) !== socket.id) {
throw new Error('WHIP session not owned by this client');
}
setState(normalized, { state: 'playing', source: 'mic-whip', error: null, startedAt: Date.now() });
cb({ success: true, roverId: normalized });
} catch (err) {
cb({ error: err.message });
}
}); });
socket.on('audio:micWhipStop', ({ roverId } = {}, cb = () => {}) => { socket.on('audio:micWhipStop', ({ roverId } = {}, cb = () => {}) => {
cb({ success: true }); try {
const normalized = String(roverId || '').trim();
ensureAudioForwardPermission(socket, normalized);
if (whipOwners.get(normalized) && whipOwners.get(normalized) !== socket.id) {
throw new Error('Mic forwarding is owned by another session');
}
stopWhipForRover(normalized, 'client_stop');
cb({ success: true, roverId: normalized });
} catch (err) {
cb({ error: err.message });
}
}); });
socket.on('disconnect', () => { socket.on('disconnect', () => {
@@ -744,6 +842,10 @@ io.on('connection', (socket) => {
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);
}); });
for (const [roverId, ownerSocketId] of whipOwners.entries()) {
if (ownerSocketId !== socket.id) continue;
stopWhipForRover(roverId, 'socket_disconnect');
}
}); });
}); });
+11 -2
View File
@@ -50,6 +50,9 @@ function extractStreamInfo(path) {
const remaining = segments.slice(start, end); const remaining = segments.slice(start, end);
if (remaining.length === 1) { if (remaining.length === 1) {
const rawId = remaining[0] || ''; const rawId = remaining[0] || '';
if (rawId.endsWith('-fwd')) {
return { type: 'rover', id: rawId, baseId: rawId.slice(0, -4) };
}
if (rawId.endsWith('-mic')) { if (rawId.endsWith('-mic')) {
return { type: 'roverMic', id: rawId, baseId: rawId.slice(0, -4) }; return { type: 'roverMic', id: rawId, baseId: rawId.slice(0, -4) };
} }
@@ -92,6 +95,9 @@ function extractStreamInfoFromBody(body = {}) {
if (srtId.endsWith('-mic')) { if (srtId.endsWith('-mic')) {
return { type: 'roverMic', id: srtId, baseId: srtId.slice(0, -4) }; return { type: 'roverMic', id: srtId, baseId: srtId.slice(0, -4) };
} }
if (srtId.endsWith('-fwd')) {
return { type: 'rover', id: srtId, baseId: srtId.slice(0, -4) };
}
const baseId = srtId.endsWith('-audio') ? srtId.slice(0, -6) : srtId; const baseId = srtId.endsWith('-audio') ? srtId.slice(0, -6) : srtId;
return { type: 'rover', id: srtId, baseId }; return { type: 'rover', id: srtId, baseId };
} }
@@ -147,7 +153,10 @@ app.post('/mediamtx/auth', (req, res) => {
} }
const info = videoSessions.getSession(sessionId); const info = videoSessions.getSession(sessionId);
if (!info || info.sourceType !== streamInfo.type || info.sourceId !== streamInfo.id) { const streamTypeMatches =
info &&
(info.sourceType === streamInfo.type || (info.sourceType === 'roverMic' && streamInfo.type === 'rover'));
if (!info || !streamTypeMatches || info.sourceId !== streamInfo.id) {
logger.warn('invalid session %s for stream %s:%s', sessionId, streamInfo.type, streamInfo.id); logger.warn('invalid session %s for stream %s:%s', sessionId, streamInfo.type, streamInfo.id);
return res.status(401).end(); return res.status(401).end();
} }
@@ -159,7 +168,7 @@ app.post('/mediamtx/auth', (req, res) => {
if (!canView(socket)) { if (!canView(socket)) {
return res.status(401).end(); return res.status(401).end();
} }
if (streamInfo.type === 'roverMic' && action === 'publish') { if (info.sourceType === 'roverMic' && action === 'publish') {
const roverId = streamInfo.baseId || streamInfo.id; const roverId = streamInfo.baseId || streamInfo.id;
if (!isVerified(socket)) { if (!isVerified(socket)) {
return res.status(401).end(); return res.status(401).end();