From ccd502c0f9113b58cd71cf07cd2bd581383f6a1a Mon Sep 17 00:00:00 2001 From: legop3 Date: Tue, 28 Apr 2026 21:03:08 -0400 Subject: [PATCH] audio forward --- rulesdocs/refector_rules_and_tracking.md | 1 + .../src/services/audioForwardService/hooks.js | 152 +++++++++++++ .../src/services/audioForwardService/index.js | 210 +++--------------- .../services/audioForwardService/policy.js | 81 +++++++ 4 files changed, 270 insertions(+), 174 deletions(-) create mode 100644 server/src/services/audioForwardService/hooks.js create mode 100644 server/src/services/audioForwardService/policy.js diff --git a/rulesdocs/refector_rules_and_tracking.md b/rulesdocs/refector_rules_and_tracking.md index 40ae2970..21ca8f0a 100644 --- a/rulesdocs/refector_rules_and_tracking.md +++ b/rulesdocs/refector_rules_and_tracking.md @@ -78,6 +78,7 @@ - Continued `llmCommentaryService` decomposition by extracting sensor activity aggregation and snapshot assembly to `llmCommentaryService/snapshotEngine.js`; rewired commentary tick/event flow to use the new engine. - Continued `llmCommentaryService` decomposition by extracting socket/role/rover event wiring into `llmCommentaryService/hooks.js` and keeping `index.js` focused on orchestration. - Finished major `llmCommentaryService` decomposition by extracting tick scheduling, run-loop orchestration, and history-reset behavior into `llmCommentaryService/runner.js`; `llmCommentaryService/index.js` is now a thin composition layer. +- Began `audioForwardService` decomposition by extracting permission/path policy helpers to `audioForwardService/policy.js` and rover/turn/socket event wiring to `audioForwardService/hooks.js`; rewired service entrypoint to use extracted modules. ## WebUI frontend ### BIGGEST OFFENDERS diff --git a/server/src/services/audioForwardService/hooks.js b/server/src/services/audioForwardService/hooks.js new file mode 100644 index 00000000..0f6dc97b --- /dev/null +++ b/server/src/services/audioForwardService/hooks.js @@ -0,0 +1,152 @@ +// audio Forward Service hooks +// Purpose: Registers rover/turn lifecycle listeners and socket handlers for audio forwarding controls. +// Scope: Keeps runtime behavior unchanged by delegating to injected core operations. +function registerAudioForwardHooks(deps) { + const { + io, + roverManager, + turnService, + logger, + serviceEnabled, + workers, + whipOwners, + ensureWorker, + stopWorker, + setState, + stopOwnedAudioIfUnauthorized, + stopWhipForRover, + ensureAudioForwardPermission, + playUploadedAudio, + stopPlayback, + resolveForwardPathId, + revokeWhipSessionForRover, + buildWhipUrl, + videoSessions, + startSilenceWriter, + } = deps; + + roverManager.managerEvents.on('rover', ({ roverId, action } = {}) => { + if (!roverId) return; + if (action === 'removed') { + stopWorker(roverId); + return; + } + if (action === 'upsert' && serviceEnabled) { + if (whipOwners.has(roverId)) { + return; + } + try { + ensureWorker(roverId); + } catch (err) { + 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') { + stopOwnedAudioIfUnauthorized(roverId, socketId, action); + } + }); + + turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => { + if (!roverId) return; + const whipOwner = whipOwners.get(roverId); + if (whipOwner) { + stopOwnedAudioIfUnauthorized(roverId, whipOwner, 'turn_change'); + } + const worker = workers.get(roverId); + if (!worker || worker.contentKind !== 'upload') return; + stopOwnedAudioIfUnauthorized(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('audio:micWhipStart', ({ roverId } = {}, cb = () => {}) => { + try { + const normalized = String(roverId || '').trim(); + ensureAudioForwardPermission(socket, normalized); + 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 = () => {}) => { + 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 = () => {}) => { + 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', () => { + 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); + }); + for (const [roverId, ownerSocketId] of whipOwners.entries()) { + if (ownerSocketId !== socket.id) continue; + stopWhipForRover(roverId, 'socket_disconnect'); + } + }); + }); +} + +module.exports = { + registerAudioForwardHooks, +}; diff --git a/server/src/services/audioForwardService/index.js b/server/src/services/audioForwardService/index.js index 0282652f..abc72738 100644 --- a/server/src/services/audioForwardService/index.js +++ b/server/src/services/audioForwardService/index.js @@ -12,6 +12,8 @@ const roverManager = require('../roverManager'); const turnService = require('../turnService'); const { isVerified } = require('../verificationService'); const videoSessions = require('../videoSessions'); +const { createAudioForwardPolicy } = require('./policy'); +const { registerAudioForwardHooks } = require('./hooks'); const audioForwardEvents = new EventEmitter(); const config = loadConfig(); @@ -84,21 +86,19 @@ function extFromUpload(name, mime) { 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'); - } -} +const audioForwardPolicy = createAudioForwardPolicy({ + isVerified, + roverManager, + turnService, + streamSuffix, + mediaConfig, +}); +const { + ensureAudioForwardPermission, + resolveForwardUrl, + resolveForwardPathId, + buildWhipUrl, +} = audioForwardPolicy; function ensureFifo(fifoPath) { try { @@ -114,46 +114,6 @@ function ensureFifo(fifoPath) { } } -function forcePublishStreamMode(rawUrl) { - const value = String(rawUrl || '').trim(); - if (!value) return ''; - if (!/[?&]streamid=#!::/.test(value)) return value; - if (/,m=publish\b/.test(value)) return value; - if (/,m=[a-zA-Z]+\b/.test(value)) return value.replace(/,m=[a-zA-Z]+\b/, ',m=publish'); - return value.replace(/([?&]streamid=#!::[^&]*)/, '$1,m=publish'); -} - -function resolveForwardUrl(roverId) { - const record = roverManager.rovers.get(roverId); - const configured = record?.meta?.media?.audioForwardUrl; - if (configured) return forcePublishStreamMode(configured); - 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 spawnFfmpeg(roverId, tag, args, options = {}) { const proc = spawn(ffmpegBin, args, { @@ -533,125 +493,27 @@ function stopOwnedAudioIfUnauthorized(roverId, ownerSocketId, reason = 'driver_c startSilenceWriter(roverId); } -roverManager.managerEvents.on('rover', ({ roverId, action } = {}) => { - if (!roverId) return; - if (action === 'removed') { - stopWorker(roverId); - return; - } - if (action === 'upsert' && serviceEnabled) { - if (whipOwners.has(roverId)) { - return; - } - try { - ensureWorker(roverId); - } catch (err) { - 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') { - stopOwnedAudioIfUnauthorized(roverId, socketId, action); - } -}); - -turnService.turnEvents.on('activeDriver', ({ roverId } = {}) => { - if (!roverId) return; - const whipOwner = whipOwners.get(roverId); - if (whipOwner) { - stopOwnedAudioIfUnauthorized(roverId, whipOwner, 'turn_change'); - } - const worker = workers.get(roverId); - if (!worker || worker.contentKind !== 'upload') return; - stopOwnedAudioIfUnauthorized(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('audio:micWhipStart', ({ roverId } = {}, cb = () => {}) => { - try { - const normalized = String(roverId || '').trim(); - ensureAudioForwardPermission(socket, normalized); - 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 = () => {}) => { - 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 = () => {}) => { - 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', () => { - 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); - }); - for (const [roverId, ownerSocketId] of whipOwners.entries()) { - if (ownerSocketId !== socket.id) continue; - stopWhipForRover(roverId, 'socket_disconnect'); - } - }); +registerAudioForwardHooks({ + io, + roverManager, + turnService, + logger, + serviceEnabled, + workers, + whipOwners, + ensureWorker, + stopWorker, + setState, + stopOwnedAudioIfUnauthorized, + stopWhipForRover, + ensureAudioForwardPermission, + playUploadedAudio, + stopPlayback, + resolveForwardPathId, + revokeWhipSessionForRover, + buildWhipUrl, + videoSessions, + startSilenceWriter, }); module.exports = { diff --git a/server/src/services/audioForwardService/policy.js b/server/src/services/audioForwardService/policy.js new file mode 100644 index 00000000..bb72a80d --- /dev/null +++ b/server/src/services/audioForwardService/policy.js @@ -0,0 +1,81 @@ +// audio Forward Service policy +// Purpose: Encapsulates permission checks and media path/url derivation helpers. +// Scope: Keeps runtime behavior unchanged while isolating validation and path-construction logic. +function createAudioForwardPolicy(deps) { + const { + isVerified, + roverManager, + turnService, + streamSuffix, + mediaConfig, + } = deps; + + 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 forcePublishStreamMode(rawUrl) { + const value = String(rawUrl || '').trim(); + if (!value) return ''; + if (!/[?&]streamid=#!::/.test(value)) return value; + if (/,m=publish\b/.test(value)) return value; + if (/,m=[a-zA-Z]+\b/.test(value)) return value.replace(/,m=[a-zA-Z]+\b/, ',m=publish'); + return value.replace(/([?&]streamid=#!::[^&]*)/, '$1,m=publish'); + } + + function resolveForwardUrl(roverId) { + const record = roverManager.rovers.get(roverId); + const configured = record?.meta?.media?.audioForwardUrl; + if (configured) return forcePublishStreamMode(configured); + 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`; + } + + return { + ensureVipVerified, + ensureAudioForwardPermission, + resolveForwardUrl, + resolveForwardPathId, + buildWhipUrl, + }; +} + +module.exports = { + createAudioForwardPolicy, +};