This commit is contained in:
legop3
2026-03-20 23:12:09 -04:00
parent 1b7885e74b
commit cb97e73987
7 changed files with 375 additions and 24 deletions
+231 -8
View File
@@ -2,9 +2,12 @@ const fs = require('fs');
const path = require('path');
const { spawn, spawnSync } = require('child_process');
const EventEmitter = require('events');
const io = require('../globals/io');
const logger = require('../globals/logger').child('audioForwardService');
const { loadConfig } = require('../helpers/configLoader');
const roverManager = require('./roverManager');
const turnService = require('./turnService');
const { isVerified } = require('./verificationService');
const audioForwardEvents = new EventEmitter();
const config = loadConfig();
@@ -16,8 +19,12 @@ const streamSuffix =
? audioForwardConfig.streamSuffix.trim()
: '-fwd';
const runtimeDir = path.resolve(audioForwardConfig.runtimeDir || '/tmp/mrr-audio-forward');
const uploadsDir = path.join(runtimeDir, 'uploads');
const maxUploadBytes = Number.isFinite(audioForwardConfig.maxUploadBytes)
? Math.max(256 * 1024, Math.floor(audioForwardConfig.maxUploadBytes))
: 8 * 1024 * 1024;
const states = new Map(); // roverId -> { state, source, error, updatedAt }
const states = new Map(); // roverId -> { state, source, error, startedAt, updatedAt }
const workers = new Map(); // roverId -> worker
function publishStateChange(roverId) {
@@ -30,6 +37,7 @@ function setState(roverId, next = {}) {
state: next.state || prev.state || 'idle',
source: Object.prototype.hasOwnProperty.call(next, 'source') ? next.source : prev.source || 'silence',
error: Object.prototype.hasOwnProperty.call(next, 'error') ? next.error : prev.error || null,
startedAt: Object.prototype.hasOwnProperty.call(next, 'startedAt') ? next.startedAt : prev.startedAt || null,
updatedAt: Date.now(),
};
states.set(roverId, merged);
@@ -46,12 +54,46 @@ function getAudioForwardState() {
function ensureRuntimeDir() {
fs.mkdirSync(runtimeDir, { recursive: true });
fs.mkdirSync(uploadsDir, { recursive: true });
}
function sanitizeRoverId(roverId) {
return String(roverId || '').replace(/[^a-zA-Z0-9_-]+/g, '_');
}
function sanitizeFileStem(name) {
return String(name || 'upload')
.toLowerCase()
.replace(/[^a-z0-9._-]+/g, '_')
.replace(/^_+|_+$/g, '')
.slice(0, 64);
}
function extFromUpload(name, mime) {
const lowerName = String(name || '').toLowerCase();
const lowerMime = String(mime || '').toLowerCase();
if (lowerName.endsWith('.mp3') || lowerMime === 'audio/mpeg' || lowerMime === 'audio/mp3') return '.mp3';
if (lowerName.endsWith('.wav') || lowerMime === 'audio/wav' || lowerMime === 'audio/x-wav') return '.wav';
if (lowerName.endsWith('.ogg') || lowerMime === 'audio/ogg') return '.ogg';
throw new Error('Unsupported upload format (allowed: mp3, wav, ogg)');
}
function ensureVipVerified(socket) {
if (!isVerified(socket)) {
throw new Error('VIP verification required');
}
}
function ensureAudioForwardPermission(socket, roverId) {
ensureVipVerified(socket);
if (!roverManager.isDriver(roverId, socket)) {
throw new Error('Audio forwarding is only allowed on your own rover');
}
if (!turnService.canDrive(roverId, socket)) {
throw new Error('Only the current driver can play audio');
}
}
function ensureFifo(fifoPath) {
try {
const stat = fs.statSync(fifoPath);
@@ -178,16 +220,37 @@ function buildSilenceWriterArgs() {
];
}
function buildUploadWriterArgs(filePath) {
return [
'-hide_banner',
'-loglevel',
'warning',
'-re',
'-i',
filePath,
'-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) => {
const code = err?.code || 'unknown';
if (code !== 'EPIPE') {
logger.warn('silence writer pipe error', { roverId: worker?.roverId, code, message: err?.message || String(err) });
logger.warn('writer pipe error', { roverId: worker?.roverId, code, message: err?.message || String(err) });
}
});
proc.stdout.on('error', (err) => {
logger.warn('silence writer stdout error', {
logger.warn('writer stdout error', {
roverId: worker?.roverId,
code: err?.code || 'unknown',
message: err?.message || String(err),
@@ -199,12 +262,23 @@ function attachWriterPipe(worker, proc) {
});
}
function cleanupUploadFile(worker) {
if (!worker?.activeUploadPath) return;
try {
fs.unlinkSync(worker.activeUploadPath);
} catch {
// noop
}
worker.activeUploadPath = null;
}
function stopContentProc(worker) {
if (!worker) return;
if (worker.contentProc) {
stopProc(worker.contentProc);
}
worker.contentProc = null;
worker.contentKind = null;
}
function startSilenceWriter(roverId) {
@@ -212,8 +286,11 @@ function startSilenceWriter(roverId) {
if (!worker || worker.stopping) return;
stopContentProc(worker);
cleanupUploadFile(worker);
worker.activeOwnerSocketId = null;
const proc = spawnFfmpeg(roverId, 'silence-writer', buildSilenceWriterArgs(), { captureStdout: true });
worker.contentProc = proc;
worker.contentKind = 'silence';
const seq = ++worker.writerSeq;
attachWriterPipe(worker, proc);
@@ -222,14 +299,55 @@ function startSilenceWriter(roverId) {
if (!current || current.stopping) return;
if (current.writerSeq !== seq || current.contentProc !== proc) return;
current.contentProc = null;
current.contentKind = null;
if (code === 0 || signal === 'SIGTERM') return;
setState(roverId, { state: 'error', source: 'silence', error: `silence writer exited code=${code} signal=${signal || 'none'}` });
setState(roverId, {
state: 'error',
source: 'silence',
error: `silence writer exited code=${code} signal=${signal || 'none'}`,
startedAt: null,
});
setTimeout(() => {
if (workers.has(roverId)) startSilenceWriter(roverId);
}, 300);
});
setState(roverId, { state: 'idle', source: 'silence', error: null });
setState(roverId, { state: 'idle', source: 'silence', error: null, startedAt: null });
}
function startUploadWriter(roverId, filePath, ownerSocketId) {
const worker = workers.get(roverId);
if (!worker || worker.stopping) return;
stopContentProc(worker);
cleanupUploadFile(worker);
worker.activeUploadPath = filePath;
worker.activeOwnerSocketId = ownerSocketId || null;
const proc = spawnFfmpeg(roverId, 'upload-writer', buildUploadWriterArgs(filePath), { captureStdout: true });
worker.contentProc = proc;
worker.contentKind = 'upload';
const seq = ++worker.writerSeq;
attachWriterPipe(worker, proc);
setState(roverId, { state: 'playing', source: 'upload', error: null, startedAt: Date.now() });
proc.on('exit', (code, signal) => {
const current = workers.get(roverId);
if (!current || current.stopping) return;
if (current.writerSeq !== seq || current.contentProc !== proc) return;
current.contentProc = null;
current.contentKind = null;
if (code != null && code !== 0 && signal !== 'SIGTERM') {
setState(roverId, {
state: 'error',
source: 'upload',
error: `upload writer exited code=${code} signal=${signal || 'none'}`,
startedAt: null,
});
}
startSilenceWriter(roverId);
});
}
function ensureWorker(roverId) {
@@ -256,7 +374,10 @@ function ensureWorker(roverId) {
outputUrl,
publisherProc: publisher,
contentProc: null,
contentKind: null,
writerSeq: 0,
activeOwnerSocketId: null,
activeUploadPath: null,
stopping: false,
};
workers.set(roverId, worker);
@@ -266,8 +387,9 @@ function ensureWorker(roverId) {
if (!current || current.publisherProc !== publisher || current.stopping) return;
setState(roverId, {
state: 'error',
source: 'publish',
source: current.contentKind || 'publish',
error: `publisher exited code=${code} signal=${signal || 'none'}`,
startedAt: null,
});
});
@@ -282,6 +404,7 @@ function stopWorker(roverId) {
worker.stopping = true;
stopContentProc(worker);
cleanupUploadFile(worker);
stopProc(worker.publisherProc);
try {
@@ -296,7 +419,57 @@ function stopWorker(roverId) {
}
workers.delete(roverId);
setState(roverId, { state: 'offline', source: 'none', error: null });
setState(roverId, { state: 'offline', source: 'none', error: null, startedAt: null });
}
function writeUploadFile(roverId, payload = {}) {
const { name, mime, dataBase64 } = payload || {};
const ext = extFromUpload(name, mime);
const encoded = typeof dataBase64 === 'string' ? dataBase64.trim() : '';
if (!encoded) {
throw new Error('Upload payload missing');
}
const bytes = Buffer.from(encoded, 'base64');
if (!bytes.length) {
throw new Error('Upload decode failed');
}
if (bytes.length > maxUploadBytes) {
throw new Error(`Upload too large (max ${maxUploadBytes} bytes)`);
}
ensureRuntimeDir();
const stem = sanitizeFileStem(name || `upload-${Date.now()}`);
const filePath = path.join(uploadsDir, `${sanitizeRoverId(roverId)}-${Date.now()}-${stem}${ext}`);
fs.writeFileSync(filePath, bytes);
return filePath;
}
function playUploadedAudio(roverId, payload = {}, ownerSocketId = null) {
ensureWorker(roverId);
const uploadPath = writeUploadFile(roverId, payload);
startUploadWriter(roverId, uploadPath, ownerSocketId);
}
function stopPlayback(roverId) {
ensureWorker(roverId);
startSilenceWriter(roverId);
}
function stopOwnedUploadIfUnauthorized(roverId, ownerSocketId, reason = 'driver_change') {
if (!roverId || !ownerSocketId) return;
const worker = workers.get(roverId);
if (!worker || worker.contentKind !== 'upload' || 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;
logger.info('Stopping upload audio due to ownership/driver change', {
roverId,
ownerSocketId,
reason,
});
startSilenceWriter(roverId);
}
roverManager.managerEvents.on('rover', ({ roverId, action } = {}) => {
@@ -309,11 +482,61 @@ roverManager.managerEvents.on('rover', ({ roverId, action } = {}) => {
try {
ensureWorker(roverId);
} catch (err) {
setState(roverId, { state: 'error', source: 'init', error: err.message });
setState(roverId, { state: 'error', source: 'init', error: err.message, startedAt: null });
}
}
});
roverManager.managerEvents.on('driver', ({ socketId, roverId, action } = {}) => {
if (!socketId || !roverId) return;
if (action === 'remove' || action === 'add') {
stopOwnedUploadIfUnauthorized(roverId, socketId, action);
}
});
turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => {
if (!roverId) return;
const worker = workers.get(roverId);
if (!worker || worker.contentKind !== 'upload') return;
stopOwnedUploadIfUnauthorized(roverId, worker.activeOwnerSocketId, 'turn_change');
});
io.on('connection', (socket) => {
socket.on('audio:uploadPlay', (payload = {}, cb = () => {}) => {
try {
const roverId = String(payload?.roverId || '').trim();
ensureAudioForwardPermission(socket, roverId);
playUploadedAudio(roverId, payload, socket.id);
cb({ success: true, roverId });
} catch (err) {
cb({ error: err.message });
}
});
socket.on('audio:uploadStop', ({ roverId } = {}, cb = () => {}) => {
try {
const normalized = String(roverId || '').trim();
ensureAudioForwardPermission(socket, normalized);
const worker = workers.get(normalized);
if (worker && worker.contentKind === 'upload' && worker.activeOwnerSocketId !== socket.id) {
throw new Error('Upload playback is owned by another session');
}
stopPlayback(normalized);
cb({ success: true, roverId: normalized });
} catch (err) {
cb({ error: err.message });
}
});
socket.on('disconnect', () => {
workers.forEach((worker, roverId) => {
if (!worker || worker.contentKind !== 'upload' || worker.activeOwnerSocketId !== socket.id) return;
logger.info('Stopping owned upload audio due to socket disconnect', { roverId, socketId: socket.id });
startSilenceWriter(roverId);
});
});
});
module.exports = {
getAudioForwardState,
audioForwardEvents,