This commit is contained in:
legop3
2026-03-17 19:35:36 -04:00
parent 618fd3e4de
commit 93204ec1f1
12 changed files with 529 additions and 137 deletions
+188 -7
View File
@@ -145,7 +145,7 @@ function resolveForwardUrl(roverId) {
function spawnProcess(roverId, tag, args, options = {}) {
const proc = spawn(ffmpegBin, args, {
stdio: ['ignore', options.captureStdout ? 'pipe' : 'ignore', 'pipe'],
stdio: [options.captureStdin ? 'pipe' : 'ignore', options.captureStdout ? 'pipe' : 'ignore', 'pipe'],
});
proc.stderr?.on('data', (chunk) => {
const text = String(chunk || '').trim();
@@ -258,6 +258,30 @@ function buildUploadWriterArgs(filePath) {
];
}
function buildMicWriterArgs() {
return [
'-hide_banner',
'-loglevel',
'warning',
'-fflags',
'nobuffer',
'-flags',
'low_delay',
'-i',
'pipe:0',
'-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) => {
@@ -292,6 +316,18 @@ function cleanupUploadFile(worker) {
function stopContentWriter(worker) {
if (!worker?.contentProc) return;
if (worker.micIdleTimer) {
clearTimeout(worker.micIdleTimer);
worker.micIdleTimer = null;
}
worker.micLastChunkAt = 0;
if (worker.contentProc.stdin && !worker.contentProc.stdin.destroyed) {
try {
worker.contentProc.stdin.destroy();
} catch {
// noop
}
}
stopProc(worker.contentProc);
worker.contentProc = null;
worker.contentKind = null;
@@ -358,6 +394,100 @@ function startUploadWriter(roverId, filePath) {
});
}
function scheduleMicIdleTimeout(roverId) {
const worker = workers.get(roverId);
if (!worker || worker.contentKind !== 'mic') return;
if (worker.micIdleTimer) {
clearTimeout(worker.micIdleTimer);
}
worker.micIdleTimer = setTimeout(() => {
const current = workers.get(roverId);
if (!current || current.contentKind !== 'mic') return;
const staleForMs = Date.now() - (current.micLastChunkAt || 0);
if (staleForMs < 2500) return;
logger.info('Stopping mic writer due to idle chunk timeout', { roverId, staleForMs });
startSilenceWriter(roverId);
}, 3000);
}
function startMicWriter(roverId, ownerSocketId = null) {
const worker = workers.get(roverId);
if (!worker || worker.stopping) return;
if (
worker.contentKind === 'mic' &&
worker.contentProc &&
worker.activeOwnerSocketId === ownerSocketId
) {
return;
}
stopContentWriter(worker);
cleanupUploadFile(worker);
const proc = spawnProcess(roverId, 'mic-writer', buildMicWriterArgs(), {
captureStdout: true,
captureStdin: true,
});
worker.contentProc = proc;
worker.contentKind = 'mic';
worker.activeOwnerSocketId = ownerSocketId;
worker.micLastChunkAt = Date.now();
scheduleMicIdleTimeout(roverId);
const seq = ++worker.writerSeq;
attachWriterPipe(worker, proc);
proc.stdin?.on('error', (err) => {
const code = err?.code || 'unknown';
if (code !== 'EPIPE') {
logger.warn('mic writer stdin error', { roverId, code, message: err?.message || String(err) });
}
});
setState(roverId, { state: 'playing', source: 'mic', 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;
current.activeOwnerSocketId = null;
if (current.micIdleTimer) {
clearTimeout(current.micIdleTimer);
current.micIdleTimer = null;
}
current.micLastChunkAt = 0;
if (code != null && code !== 0 && signal !== 'SIGTERM') {
setState(roverId, { state: 'error', source: 'mic', error: `mic writer exited code=${code} signal=${signal || 'none'}` });
}
startSilenceWriter(roverId);
});
}
function pushMicChunk(roverId, ownerSocketId, dataBase64) {
const worker = workers.get(roverId);
if (!worker) {
throw new Error('Audio forward worker unavailable');
}
if (worker.contentKind !== 'mic' || !worker.contentProc || worker.activeOwnerSocketId !== ownerSocketId) {
throw new Error('Mic forwarding is not active');
}
const encoded = typeof dataBase64 === 'string' ? dataBase64.trim() : '';
if (!encoded) {
throw new Error('Mic chunk missing');
}
const bytes = Buffer.from(encoded, 'base64');
if (!bytes.length) {
throw new Error('Mic chunk decode failed');
}
if (worker.contentProc.stdin?.writable !== true) {
throw new Error('Mic writer input is not writable');
}
worker.micLastChunkAt = Date.now();
worker.contentProc.stdin.write(bytes);
scheduleMicIdleTimeout(roverId);
}
function ensureWorker(roverId) {
ensureServiceEnabled();
if (!roverId) {
@@ -388,7 +518,10 @@ function ensureWorker(roverId) {
publisherProc: publisher,
contentProc: null,
contentKind: null,
activeOwnerSocketId: null,
activeUploadPath: null,
micLastChunkAt: 0,
micIdleTimer: null,
writerSeq: 0,
stopping: false,
};
@@ -490,31 +623,32 @@ roverManager.managerEvents.on('rover', ({ roverId, action } = {}) => {
}
});
function stopOwnedUploadIfUnauthorized(roverId, ownerSocketId, reason = 'driver_change') {
function stopOwnedAudioIfUnauthorized(roverId, ownerSocketId, reason = 'driver_change') {
if (!roverId || !ownerSocketId) return;
const worker = workers.get(roverId);
if (!worker || worker.contentKind !== 'upload') return;
if (!worker) 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;
logger.info('Stopping upload due to ownership/driver change', { roverId, ownerSocketId, reason });
logger.info('Stopping audio forward due to ownership/driver change', { roverId, ownerSocketId, reason, source: worker.contentKind });
startSilenceWriter(roverId);
}
roverManager.managerEvents.on('driver', ({ socketId, roverId, action } = {}) => {
if (!socketId || !roverId) return;
if (action === 'remove' || action === 'add') {
stopOwnedUploadIfUnauthorized(roverId, socketId, action);
stopOwnedAudioIfUnauthorized(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');
if (!worker || (worker.contentKind !== 'upload' && worker.contentKind !== 'mic')) return;
stopOwnedAudioIfUnauthorized(roverId, worker.activeOwnerSocketId, 'turn_change');
});
io.on('connection', (socket) => {
@@ -540,6 +674,53 @@ io.on('connection', (socket) => {
cb({ error: err.message });
}
});
socket.on('audio:micStart', ({ roverId } = {}, cb = () => {}) => {
try {
const normalized = String(roverId || '').trim();
ensureAudioForwardPermission(socket, normalized);
ensureWorker(normalized);
startMicWriter(normalized, socket.id);
cb({ success: true, roverId: normalized });
} catch (err) {
cb({ error: err.message });
}
});
socket.on('audio:micChunk', (payload = {}, cb) => {
try {
const normalized = String(payload?.roverId || '').trim();
ensureAudioForwardPermission(socket, normalized);
pushMicChunk(normalized, socket.id, payload?.dataBase64);
if (typeof cb === 'function') cb({ success: true });
} catch (err) {
if (typeof cb === 'function') cb({ error: err.message });
}
});
socket.on('audio:micStop', ({ roverId } = {}, cb = () => {}) => {
try {
const normalized = String(roverId || '').trim();
ensureAudioForwardPermission(socket, normalized);
const worker = workers.get(normalized);
if (worker && worker.contentKind === 'mic' && worker.activeOwnerSocketId !== socket.id) {
throw new Error('Mic forwarding 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.activeOwnerSocketId !== socket.id) 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);
});
});
});
module.exports = {