diff --git a/rulesdocs/refector_rules_and_tracking.md b/rulesdocs/refector_rules_and_tracking.md index 63d92681..ddc46b94 100644 --- a/rulesdocs/refector_rules_and_tracking.md +++ b/rulesdocs/refector_rules_and_tracking.md @@ -49,7 +49,7 @@ - [ ] llm commentary service - [x] private rover access request service - [x] replay services (consolidated under replayEngineV2) -- [ ] room camera services (already partly split; reformat consistently) +- [x] room camera services (consolidated into roomCameraService multipart folder) - [ ] rover manager service (in progress: constants/state extracted) - [x] session service - [x] turn service @@ -90,6 +90,9 @@ - Finished `replayEngineV2` decomposition by extracting environment/path constants to `replayEngineV2/constants.js`, mutable runtime state to `replayEngineV2/state.js`, source discovery/worker arg building to `replayEngineV2/sources.js`, ffmpeg worker lifecycle to `replayEngineV2/workerManager.js`, segment indexing/retention/health snapshot logic to `replayEngineV2/segmentStore.js`, sidebar SVG/video rendering to `replayEngineV2/sidebarRenderer.js`, and replay assembly pipeline to `replayEngineV2/replayBuilder.js`; `replayEngineV2/index.js` is now a thin orchestration layer. - Consolidated replay-related single-file services into `replayEngineV2` by moving cooldown state (`cooldown.js`), user-facing replay source validation/defaults (`replaySources.js`), and replay socket hooks (`socketHooks.js`) into the engine folder; removed obsolete standalone services `replayBuildService`, `replayService`, `replaySourceService`, and `replaySocketService` and rewired dependents to import directly from `replayEngineV2`. - Finished `chatService` decomposition by extracting runtime constants (`chatService/constants.js`), shared mutable state (`chatService/state.js`), content/moderation helpers (`chatService/contentFilters.js`), payload/context builders (`chatService/contextBuilders.js`), bus/history broadcast pipeline (`chatService/broadcast.js`), rover-side notification helpers (`chatService/notifications.js`), message handlers (`chatService/handlers.js`), and socket/event-bus wiring (`chatService/socketHooks.js`); `chatService/index.js` is now a thin orchestration layer. +- Consolidated room-camera services into `roomCameraService` by absorbing catalog (`roomCameraService`), snapshot polling/streaming (`roomCameraSnapshotService`), socket fan-out (`roomCameraSocketService`), and replay assembly (`roomCameraReplayService`) into one multipart folder (`roomCameraService/catalog.js`, `snapshotEngine.js`, `socketGateway.js`, `replayBuilder.js`), and updated imports/startup wiring to use the consolidated service exports. +- Follow-up: moved room-camera replay assembly module from `roomCameraService/replayBuilder.js` into `replayEngineV2/roomCameraReplayBuilder.js`; `roomCameraService` now consumes replay functionality from replay engine ownership while keeping the same exported room-camera replay API. +- Consolidated rover snapshot polling/socket services into `roverSnapshotService` by absorbing `roverSnapshotSocketService` into folder modules (`roverSnapshotService/poller.js`, `socketGateway.js`) and keeping `roverSnapshotService/index.js` as the startup composition/export layer. ## WebUI frontend ### BIGGEST OFFENDERS diff --git a/server/index.js b/server/index.js index ce515fcb..846eeecf 100644 --- a/server/index.js +++ b/server/index.js @@ -25,8 +25,8 @@ require('./src/services/serverControlService'); require('./src/services/videoSessions'); require('./src/services/videoAuthService'); require('./src/services/videoSocketService'); -require('./src/services/roomCameraSocketService'); -require('./src/services/roverSnapshotSocketService'); +require('./src/services/roomCameraService'); +require('./src/services/roverSnapshotService'); require('./src/services/humanAlertButtonService'); require('./src/services/embedHttpService'); require('./src/services/logStreamService'); diff --git a/server/src/services/embedService/index.js b/server/src/services/embedService/index.js index fef6a627..a7c15b7c 100644 --- a/server/src/services/embedService/index.js +++ b/server/src/services/embedService/index.js @@ -9,7 +9,7 @@ const { getMode } = require('../modeManager'); const roverManager = require('../roverManager'); const { getActiveDrivers, getTurnQueues } = require('../turnService'); const { getRoomCameras } = require('../roomCameraService'); -const { getRoomCameraState } = require('../roomCameraSnapshotService'); +const { getRoomCameraState } = require('../roomCameraService'); const { loadConfig } = require('../../helpers/configLoader'); const INDEX_HTML_PATH = path.join(__dirname, '..', '..', '..', 'public', 'index.html'); diff --git a/server/src/services/healthService/index.js b/server/src/services/healthService/index.js index 0f5a98a5..2bfb1820 100644 --- a/server/src/services/healthService/index.js +++ b/server/src/services/healthService/index.js @@ -5,7 +5,7 @@ const fsp = require('fs/promises'); const path = require('path'); const roverManager = require('../roverManager'); const { getRoomCameras } = require('../roomCameraService'); -const { getRoomCameraState } = require('../roomCameraSnapshotService'); +const { getRoomCameraState } = require('../roomCameraService'); const { getReplayHealthSnapshot } = require('../replayEngineV2'); const ROVER_SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/var/lib/rover-snapshots'; diff --git a/server/src/services/humanAlertButtonService/index.js b/server/src/services/humanAlertButtonService/index.js index ccbb023b..97848fe2 100644 --- a/server/src/services/humanAlertButtonService/index.js +++ b/server/src/services/humanAlertButtonService/index.js @@ -9,7 +9,7 @@ const { issueCommand } = require('../commandService'); const roverManager = require('../roverManager'); const { toggleLightsLockedOn } = require('../homeAssistantService'); const { getRoomCameras } = require('../roomCameraService'); -const { getRoomCameraState } = require('../roomCameraSnapshotService'); +const { getRoomCameraState } = require('../roomCameraService'); const HA_BUTTON_EVENT_TYPE = 'ha.button.action'; const HUMAN_ALERT_ACTION = 'humanAlert'; diff --git a/server/src/services/roomCameraReplayService/index.js b/server/src/services/replayEngineV2/roomCameraReplayBuilder.js similarity index 68% rename from server/src/services/roomCameraReplayService/index.js rename to server/src/services/replayEngineV2/roomCameraReplayBuilder.js index 332e8373..7d7fbeb6 100644 --- a/server/src/services/roomCameraReplayService/index.js +++ b/server/src/services/replayEngineV2/roomCameraReplayBuilder.js @@ -1,6 +1,6 @@ -// room Camera Replay Service -// Purpose: Defines the room Camera Replay Service module and the helpers/state used by this service unit. -// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary. +// Room Camera Replay Builder +// Purpose: Records recent room-camera frames and assembles short replay videos for operator workflows. +// Scope: Owns frame history retention and ffmpeg/ffprobe-based replay rendering. const { execFile } = require('child_process'); const EventEmitter = require('events'); const fsp = require('fs/promises'); @@ -9,7 +9,6 @@ const path = require('path'); const { promisify } = require('util'); const logger = require('../../globals/logger').child('roomCameraReplay'); -const { getRoomCamera, getRoomCameras } = require('../roomCameraService'); const execFileAsync = promisify(execFile); @@ -20,34 +19,29 @@ const MAX_REPLAY_WIDTH = 1280; const MAX_REPLAY_HEIGHT = 720; const REPLAY_MAX_BYTES = Math.floor(9.5 * 1024 * 1024); -const frameHistory = new Map(); // id -> [{ buffer, ts }] -const latestFrames = new Map(); // id -> { buffer, ts } +const frameHistory = new Map(); +const latestFrames = new Map(); const events = new EventEmitter(); -function recordFrame(id, buffer, ts = Date.now()) { +function recordRoomCameraFrame(id, buffer, ts = Date.now()) { if (!id || !buffer) return; const entry = { buffer, ts }; latestFrames.set(id, entry); const history = frameHistory.get(id) || []; history.push(entry); const cutoff = ts - HISTORY_WINDOW_MS; - while (history.length && history[0].ts < cutoff) { - history.shift(); - } + while (history.length && history[0].ts < cutoff) history.shift(); frameHistory.set(id, history); events.emit('frame', { id, ts }); } -function clearFrames() { +function clearRoomCameraReplayFrames() { frameHistory.clear(); latestFrames.clear(); } -function getReplayMetadata() { - return { - durationMs: REPLAY_DURATION_MS, - fps: REPLAY_FPS, - }; +function getRoomCameraReplayMetadata() { + return { durationMs: REPLAY_DURATION_MS, fps: REPLAY_FPS }; } function buildTimelineForCamera(id, startMs, frameCount, frameStepMs) { @@ -60,9 +54,7 @@ function buildTimelineForCamera(id, startMs, frameCount, frameStepMs) { lastBuffer = history[idx].buffer; idx += 1; } - if (!lastBuffer) { - lastBuffer = history[0]?.buffer || fallback; - } + if (!lastBuffer) lastBuffer = history[0]?.buffer || fallback; const frames = new Array(frameCount); for (let i = 0; i < frameCount; i += 1) { const slotTs = startMs + i * frameStepMs; @@ -87,15 +79,7 @@ async function probeMaxFrameSize(framePaths) { for (const framePath of framePaths) { try { const { stdout } = await execFileAsync('ffprobe', [ - '-v', - 'error', - '-select_streams', - 'v:0', - '-show_entries', - 'stream=width,height', - '-of', - 'csv=p=0', - framePath, + '-v', 'error', '-select_streams', 'v:0', '-show_entries', 'stream=width,height', '-of', 'csv=p=0', framePath, ]); const [widthRaw, heightRaw] = stdout.trim().split(','); const width = Number(widthRaw); @@ -119,27 +103,20 @@ function clampEven(value) { return Math.max(2, Math.floor(value / 2) * 2); } -async function buildReplayVideo({ cameraId = null } = {}) { +async function buildRoomCameraReplayVideo({ cameraId = null } = {}, { getRoomCamera, getRoomCameras }) { const cameras = cameraId ? [getRoomCamera(cameraId)].filter(Boolean) : getRoomCameras(); - if (!cameras.length) { - throw new Error('No room cameras configured'); - } + if (!cameras.length) throw new Error('No room cameras configured'); - const fpsValue = REPLAY_FPS; - const frameCount = Math.max(1, Math.round((REPLAY_DURATION_MS / 1000) * fpsValue)); - const frameStepMs = 1000 / fpsValue; + const frameCount = Math.max(1, Math.round((REPLAY_DURATION_MS / 1000) * REPLAY_FPS)); + const frameStepMs = 1000 / REPLAY_FPS; const startMs = Date.now() - REPLAY_DURATION_MS; const cameraEntries = []; cameras.forEach((camera) => { const frames = buildTimelineForCamera(camera.id, startMs, frameCount, frameStepMs); - if (!frames) return; - cameraEntries.push({ camera, frames }); + if (frames) cameraEntries.push({ camera, frames }); }); - - if (!cameraEntries.length) { - throw new Error('No camera frames available yet'); - } + if (!cameraEntries.length) throw new Error('No camera frames available yet'); const tmpDir = await fsp.mkdtemp(path.join(os.tmpdir(), 'rover-replay-')); try { @@ -151,14 +128,10 @@ async function buildReplayVideo({ cameraId = null } = {}) { await fsp.mkdir(camDir, { recursive: true }); for (let j = 0; j < entry.frames.length; j += 1) { const buffer = entry.frames[j]; - if (!buffer) { - throw new Error(`Camera ${entry.camera.id} missing replay frame`); - } + if (!buffer) throw new Error(`Camera ${entry.camera.id} missing replay frame`); const filename = `frame-${String(j + 1).padStart(4, '0')}.jpg`; const fullPath = path.join(camDir, filename); - if (j === 0) { - firstFramePaths.push(fullPath); - } + if (j === 0) firstFramePaths.push(fullPath); await fsp.writeFile(fullPath, buffer); } } @@ -171,8 +144,8 @@ async function buildReplayVideo({ cameraId = null } = {}) { let outputHeight = tileHeight * layout.rows; if (outputWidth > MAX_REPLAY_WIDTH || outputHeight > MAX_REPLAY_HEIGHT) { const scale = Math.min(MAX_REPLAY_WIDTH / outputWidth, MAX_REPLAY_HEIGHT / outputHeight); - tileWidth = tileWidth * scale; - tileHeight = tileHeight * scale; + tileWidth *= scale; + tileHeight *= scale; outputWidth = tileWidth * layout.cols; outputHeight = tileHeight * layout.rows; } @@ -181,8 +154,8 @@ async function buildReplayVideo({ cameraId = null } = {}) { outputWidth = clampEven(tileWidth * layout.cols); outputHeight = clampEven(tileHeight * layout.rows); - const fps = fpsValue.toFixed(3); - const durationSec = Math.max(1, frameCount / fpsValue); + const fps = REPLAY_FPS.toFixed(3); + const durationSec = Math.max(1, frameCount / REPLAY_FPS); const targetBitrateKbps = Math.max(300, Math.floor((REPLAY_MAX_BYTES * 8) / durationSec / 1000)); const maxrateKbps = Math.floor(targetBitrateKbps * 1.1); const bufsizeKbps = Math.floor(targetBitrateKbps * 2); @@ -192,9 +165,7 @@ async function buildReplayVideo({ cameraId = null } = {}) { const layoutParts = []; for (let i = 0; i < cameraEntries.length; i += 1) { inputArgs.push('-framerate', fps, '-i', path.join(cameraEntries[i].dir, 'frame-%04d.jpg')); - filterParts.push( - `[${i}:v]${buildScalePadFilter(tileWidth, tileHeight)}[v${i}]`, - ); + filterParts.push(`[${i}:v]${buildScalePadFilter(tileWidth, tileHeight)}[v${i}]`); const x = (i % layout.cols) * tileWidth; const y = Math.floor(i / layout.cols) * tileHeight; layoutParts.push(`${x}_${y}`); @@ -203,38 +174,25 @@ async function buildReplayVideo({ cameraId = null } = {}) { filterParts.push('[v0]null[v]'); } else { filterParts.push( - `${cameraEntries.map((_, i) => `[v${i}]`).join('')}` + - `xstack=inputs=${cameraEntries.length}:layout=${layoutParts.join('|')}:fill=black[v]`, + `${cameraEntries.map((_, i) => `[v${i}]`).join('')}xstack=inputs=${cameraEntries.length}:layout=${layoutParts.join('|')}:fill=black[v]`, ); } const outPath = path.join(tmpDir, 'replay.mp4'); await execFileAsync('ffmpeg', [ - '-y', - '-hide_banner', - '-loglevel', - 'error', + '-y', '-hide_banner', '-loglevel', 'error', ...inputArgs, - '-filter_complex', - filterParts.join(';'), - '-map', - '[v]', - '-r', - fps, - '-c:v', - 'libx264', - '-b:v', - `${targetBitrateKbps}k`, - '-maxrate', - `${maxrateKbps}k`, - '-bufsize', - `${bufsizeKbps}k`, - '-pix_fmt', - 'yuv420p', + '-filter_complex', filterParts.join(';'), + '-map', '[v]', + '-r', fps, + '-c:v', 'libx264', + '-b:v', `${targetBitrateKbps}k`, + '-maxrate', `${maxrateKbps}k`, + '-bufsize', `${bufsizeKbps}k`, + '-pix_fmt', 'yuv420p', outPath, ]); - const buffer = await fsp.readFile(outPath); - return buffer; + return await fsp.readFile(outPath); } finally { try { await fsp.rm(tmpDir, { recursive: true, force: true }); @@ -245,9 +203,9 @@ async function buildReplayVideo({ cameraId = null } = {}) { } module.exports = { - recordRoomCameraFrame: recordFrame, - clearRoomCameraReplayFrames: clearFrames, - getRoomCameraReplayMetadata: getReplayMetadata, - buildRoomCameraReplayVideo: buildReplayVideo, + recordRoomCameraFrame, + clearRoomCameraReplayFrames, + getRoomCameraReplayMetadata, + buildRoomCameraReplayVideo, roomCameraReplayEvents: events, }; diff --git a/server/src/services/roomCameraService/catalog.js b/server/src/services/roomCameraService/catalog.js new file mode 100644 index 00000000..b9003f66 --- /dev/null +++ b/server/src/services/roomCameraService/catalog.js @@ -0,0 +1,57 @@ +// Room Camera Catalog +// Purpose: Loads and normalizes configured room camera entries and emits update events when roster data changes. +// Scope: Owns camera identity/url normalization and read-only accessors for room camera metadata. +const EventEmitter = require('events'); +const logger = require('../../globals/logger').child('roomCameraService'); +const { loadConfig } = require('../../helpers/configLoader'); + +const events = new EventEmitter(); +const config = loadConfig(); +const cameraMap = new Map(); + +function normalizeCamera(camera) { + if (!camera) return null; + const id = camera.id || camera.name; + if (!id) { + logger.warn('Room camera missing id', camera); + return null; + } + if (!camera.url && !camera.streamUrl && !camera.mjpegUrl) { + logger.warn('Room camera missing url/streamUrl', { id, camera }); + return null; + } + return { + id: String(id), + name: camera.name || camera.id || String(id), + description: camera.description || null, + url: camera.url || null, + streamUrl: camera.streamUrl || camera.mjpegUrl || null, + }; +} + +function getRoomCameras() { + return Array.from(cameraMap.values()); +} + +function getRoomCamera(id) { + if (!id) return null; + return cameraMap.get(String(id)) || null; +} + +function loadFromConfig() { + cameraMap.clear(); + const list = Array.isArray(config.roomCameras) ? config.roomCameras : []; + list.forEach((camera) => { + const normalized = normalizeCamera(camera); + if (normalized) cameraMap.set(normalized.id, normalized); + }); + logger.info('Loaded room cameras', { count: cameraMap.size }); + events.emit('update', getRoomCameras()); +} + +module.exports = { + loadFromConfig, + getRoomCameras, + getRoomCamera, + roomCameraEvents: events, +}; diff --git a/server/src/services/roomCameraService/index.js b/server/src/services/roomCameraService/index.js index 28425c10..2d032073 100644 --- a/server/src/services/roomCameraService/index.js +++ b/server/src/services/roomCameraService/index.js @@ -1,61 +1,36 @@ -// room Camera Service -// Purpose: Defines the room Camera Service module and the helpers/state used by this service unit. -// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary. -const EventEmitter = require('events'); -const logger = require('../../globals/logger').child('roomCameraService'); -const { loadConfig } = require('../../helpers/configLoader'); +// Room Camera Service +// Purpose: Composes room camera catalog, snapshot streaming, socket delivery, and replay helpers in one service folder. +// Scope: Exposes the existing room-camera public API while preserving side-effect startup behavior. +const { loadFromConfig, getRoomCameras, getRoomCamera, roomCameraEvents } = require('./catalog'); +const { createSnapshotEngine } = require('./snapshotEngine'); +const { registerRoomCameraSocketGateway } = require('./socketGateway'); +const replay = require('../replayEngineV2/roomCameraReplayBuilder'); -const events = new EventEmitter(); -const config = loadConfig(); +const snapshotEngine = createSnapshotEngine({ getRoomCameras, roomCameraEvents }); +snapshotEngine.startAll(); -const cameraMap = new Map(); - -function normalizeCamera(camera) { - if (!camera) return null; - const id = camera.id || camera.name; - if (!id) { - logger.warn('Room camera missing id', camera); - return null; - } - if (!camera.url && !camera.streamUrl && !camera.mjpegUrl) { - logger.warn('Room camera missing url/streamUrl', { id, camera }); - return null; - } - return { - id: String(id), - name: camera.name || camera.id || String(id), - description: camera.description || null, - url: camera.url || null, - streamUrl: camera.streamUrl || camera.mjpegUrl || null, - }; -} - -function loadFromConfig() { - cameraMap.clear(); - const list = Array.isArray(config.roomCameras) ? config.roomCameras : []; - list.forEach((camera) => { - const normalized = normalizeCamera(camera); - if (normalized) { - cameraMap.set(normalized.id, normalized); - } - }); - logger.info('Loaded room cameras', { count: cameraMap.size }); - events.emit('update', getRoomCameras()); -} - -function getRoomCameras() { - return Array.from(cameraMap.values()); -} - -function getRoomCamera(id) { - if (!id) return null; - return cameraMap.get(String(id)) || null; -} +registerRoomCameraSocketGateway({ + getRoomCamera, + getRoomCameras, + getRoomCameraState: snapshotEngine.getRoomCameraState, + roomCameraStreamEvents: snapshotEngine.roomCameraStreamEvents, +}); loadFromConfig(); +function buildRoomCameraReplayVideo(options = {}) { + return replay.buildRoomCameraReplayVideo(options, { getRoomCamera, getRoomCameras }); +} + module.exports = { getRoomCameras, getRoomCamera, - roomCameraEvents: events, + roomCameraEvents, + roomCameraStreamEvents: snapshotEngine.roomCameraStreamEvents, + getRoomCameraState: snapshotEngine.getRoomCameraState, + recordRoomCameraFrame: replay.recordRoomCameraFrame, + clearRoomCameraReplayFrames: replay.clearRoomCameraReplayFrames, + getRoomCameraReplayMetadata: replay.getRoomCameraReplayMetadata, + buildRoomCameraReplayVideo, + roomCameraReplayEvents: replay.roomCameraReplayEvents, }; diff --git a/server/src/services/roomCameraService/snapshotEngine.js b/server/src/services/roomCameraService/snapshotEngine.js new file mode 100644 index 00000000..33f80c41 --- /dev/null +++ b/server/src/services/roomCameraService/snapshotEngine.js @@ -0,0 +1,221 @@ +// Room Camera Snapshot Engine +// Purpose: Maintains per-camera snapshot/stream frame state and handles polling/reconnect loops for camera feeds. +// Scope: Owns frame acquisition from MJPEG streams or snapshot URLs and emits frame/status events. +const EventEmitter = require('events'); +const http = require('http'); +const https = require('https'); +const logger = require('../../globals/logger').child('roomCameraSnapshot'); + +const POLL_INTERVAL_MS = 67; +const FETCH_TIMEOUT_MS = 2000; +const STREAM_RETRY_MS = 1500; + +const cameraState = new Map(); +const streamState = new Map(); +const events = new EventEmitter(); +let pollTimer = null; + +const JPEG_START = Buffer.from([0xff, 0xd8]); +const JPEG_END = Buffer.from([0xff, 0xd9]); + +function markState(id, updates = {}) { + const prev = cameraState.get(id) || {}; + const next = { ...prev, ...updates }; + cameraState.set(id, next); + return next; +} + +function isMjpegUrl(rawUrl) { + if (!rawUrl) return false; + const lower = String(rawUrl).toLowerCase(); + return lower.endsWith('.mjpg') || lower.endsWith('.mjpeg') || lower.includes('mjpeg'); +} + +function getStreamUrl(camera) { + if (camera.streamUrl) return camera.streamUrl; + if (isMjpegUrl(camera.url)) return camera.url; + return null; +} + +function stopStream(id) { + const entry = streamState.get(id); + if (!entry) return; + if (entry.req) entry.req.destroy(); + if (entry.reconnectTimer) clearTimeout(entry.reconnectTimer); + streamState.delete(id); +} + +function scheduleStreamReconnect(camera, startStream) { + const id = camera.id; + const entry = streamState.get(id) || {}; + if (entry.reconnectTimer) return; + entry.reconnectTimer = setTimeout(() => { + const current = streamState.get(id) || {}; + current.reconnectTimer = null; + streamState.set(id, current); + startStream(camera); + }, STREAM_RETRY_MS); + streamState.set(id, entry); +} + +function clearStreamRequest(id, req) { + const entry = streamState.get(id); + if (!entry || entry.req !== req) return; + entry.req = null; + streamState.set(id, entry); +} + +function handleStreamError(camera, err) { + const state = cameraState.get(camera.id) || {}; + const failures = (state.failures || 0) + 1; + markState(camera.id, { error: err.message || 'stream error', failures }); + events.emit('status', { id: camera.id, error: err.message || 'stream error' }); + logger.warn('Stream failed', { id: camera.id, err: err.message }); +} + +function createSnapshotEngine({ getRoomCameras, roomCameraEvents }) { + function startStream(camera) { + const streamUrl = getStreamUrl(camera); + if (!streamUrl || streamState.get(camera.id)?.req) return; + const url = new URL(streamUrl); + const client = url.protocol === 'https:' ? https : http; + const req = client.get( + { + hostname: url.hostname, + port: url.port || (url.protocol === 'https:' ? 443 : 80), + path: `${url.pathname}${url.search}`, + headers: { Accept: 'multipart/x-mixed-replace' }, + }, + (res) => { + if (res.statusCode !== 200) { + res.resume(); + clearStreamRequest(camera.id, req); + handleStreamError(camera, new Error(`HTTP ${res.statusCode}`)); + scheduleStreamReconnect(camera, startStream); + return; + } + let buffer = Buffer.alloc(0); + res.on('data', (chunk) => { + buffer = Buffer.concat([buffer, chunk]); + while (true) { + let start = buffer.indexOf(JPEG_START); + if (start === -1) { + if (buffer.length > 2 * 1024 * 1024) buffer = buffer.slice(-1024 * 1024); + break; + } + if (start > 0) { + buffer = buffer.slice(start); + start = 0; + } + const end = buffer.indexOf(JPEG_END, 2); + if (end === -1) break; + const frame = buffer.slice(0, end + 2); + buffer = buffer.slice(end + 2); + const ts = Date.now(); + markState(camera.id, { frame, ts, error: null, failures: 0 }); + events.emit('frame', { id: camera.id, buffer: frame, ts }); + } + }); + res.on('end', () => { + clearStreamRequest(camera.id, req); + scheduleStreamReconnect(camera, startStream); + }); + res.on('aborted', () => { + clearStreamRequest(camera.id, req); + handleStreamError(camera, new Error('stream aborted')); + scheduleStreamReconnect(camera, startStream); + }); + res.on('error', (err) => { + clearStreamRequest(camera.id, req); + handleStreamError(camera, err); + scheduleStreamReconnect(camera, startStream); + }); + }, + ); + req.on('error', (err) => { + clearStreamRequest(camera.id, req); + handleStreamError(camera, err); + scheduleStreamReconnect(camera, startStream); + }); + streamState.set(camera.id, { req, reconnectTimer: null }); + } + + async function fetchSnapshot(camera) { + const { id, url } = camera; + const state = cameraState.get(id); + if (!url || state?.fetching) return; + markState(id, { fetching: true }); + const abortController = new AbortController(); + const timeout = setTimeout(() => abortController.abort(), FETCH_TIMEOUT_MS); + try { + const res = await fetch(url, { signal: abortController.signal }); + if (!res.ok) throw new Error(`HTTP ${res.status}`); + const arrayBuffer = await res.arrayBuffer(); + const buffer = Buffer.from(arrayBuffer); + const ts = Date.now(); + markState(id, { frame: buffer, ts, error: null, failures: 0 }); + events.emit('frame', { id, buffer, ts }); + } catch (err) { + const failures = (state?.failures || 0) + 1; + markState(id, { error: err.message, failures }); + events.emit('status', { id, error: err.message }); + logger.warn('Snapshot fetch failed', { id, err: err.message }); + } finally { + clearTimeout(timeout); + markState(id, { fetching: false }); + } + } + + function stopAll() { + if (pollTimer) { + clearInterval(pollTimer); + pollTimer = null; + } + Array.from(streamState.keys()).forEach((id) => stopStream(id)); + cameraState.clear(); + } + + function startAll() { + stopAll(); + const cameras = getRoomCameras(); + const snapshotCameras = cameras.filter((camera) => !getStreamUrl(camera) && camera.url); + cameras.forEach((camera) => startStream(camera)); + if (snapshotCameras.length) { + pollTimer = setInterval(() => { + snapshotCameras.forEach((camera) => fetchSnapshot(camera)); + }, POLL_INTERVAL_MS); + snapshotCameras.forEach((camera) => fetchSnapshot(camera)); + } + logger.info('Started room camera feeds', { + total: cameras.length, + streaming: cameras.length - snapshotCameras.length, + snapshots: snapshotCameras.length, + }); + } + + function getRoomCameraState(id) { + const state = cameraState.get(id); + if (!state) return null; + return { + frame: state.frame || null, + ts: state.ts || null, + error: state.error || null, + }; + } + + roomCameraEvents.on('update', () => { + logger.info('Room cameras changed; restarting snapshot pollers'); + startAll(); + }); + + return { + startAll, + stopAll, + roomCameraStreamEvents: events, + getRoomCameraState, + }; +} + +module.exports = { + createSnapshotEngine, +}; diff --git a/server/src/services/roomCameraService/socketGateway.js b/server/src/services/roomCameraService/socketGateway.js new file mode 100644 index 00000000..834ae80f --- /dev/null +++ b/server/src/services/roomCameraService/socketGateway.js @@ -0,0 +1,145 @@ +// Room Camera Socket Gateway +// Purpose: Manages room-camera subscription sockets and forwards frame/status events to authorized viewers. +// Scope: Owns socket-level auth checks, subscription state, throttling, and event fan-out. +const io = require('../../globals/io'); +const logger = require('../../globals/logger').child('roomCameraSocket'); +const { getMode, MODES } = require('../modeManager'); +const { isAdmin, isLockdownAdmin, getRole } = require('../roleService'); + +const SUBSCRIBE_LIMIT = 50; +const SUBSCRIBE_WINDOW_MS = 10000; +const STREAM_INTERVAL_MS = 1000; + +function passesMode(socket) { + const mode = getMode(); + if (mode === MODES.LOCKDOWN) return isLockdownAdmin(socket); + if (mode === MODES.ADMIN) { + const role = getRole(socket); + return role === 'spectator' || isAdmin(socket); + } + return true; +} + +function registerRoomCameraSocketGateway({ getRoomCamera, getRoomCameras, getRoomCameraState, roomCameraStreamEvents }) { + const cameraSubscribers = new Map(); + const socketSubscriptions = new Map(); + const subscribeBuckets = new Map(); + const lastSentBySocket = new Map(); + + function addSubscription(socket, cameraId) { + if (!cameraSubscribers.has(cameraId)) cameraSubscribers.set(cameraId, new Set()); + cameraSubscribers.get(cameraId).add(socket.id); + if (!socketSubscriptions.has(socket.id)) socketSubscriptions.set(socket.id, new Set()); + socketSubscriptions.get(socket.id).add(cameraId); + } + + function removeSubscription(socketId, cameraId) { + const bucket = cameraSubscribers.get(cameraId); + if (bucket) { + bucket.delete(socketId); + if (bucket.size === 0) cameraSubscribers.delete(cameraId); + } + const socketBucket = socketSubscriptions.get(socketId); + if (socketBucket) { + socketBucket.delete(cameraId); + if (socketBucket.size === 0) socketSubscriptions.delete(socketId); + } + } + + function removeAllSubscriptions(socketId) { + const bucket = socketSubscriptions.get(socketId); + if (!bucket) return; + bucket.forEach((cameraId) => removeSubscription(socketId, cameraId)); + } + + function allowSubscribe(socketId) { + const now = Date.now(); + let bucket = subscribeBuckets.get(socketId); + if (!bucket || now - bucket.start >= SUBSCRIBE_WINDOW_MS) bucket = { start: now, count: 0 }; + bucket.count += 1; + subscribeBuckets.set(socketId, bucket); + return bucket.count <= SUBSCRIBE_LIMIT; + } + + function sendFrame(socket, cameraId, payload, buffer) { + socket.emit('roomCamera:frame', { id: cameraId, ...payload }, buffer); + } + + function sendStatus(socket, cameraId, status) { + socket.emit('roomCamera:status', { id: cameraId, ...status }); + } + + roomCameraStreamEvents.on('frame', ({ id, buffer, ts }) => { + const bucket = cameraSubscribers.get(id); + if (!bucket || !buffer) return; + bucket.forEach((socketId) => { + const socket = io.sockets.sockets.get(socketId); + if (!socket) return; + let lastMap = lastSentBySocket.get(socketId); + if (!lastMap) { + lastMap = new Map(); + lastSentBySocket.set(socketId, lastMap); + } + const lastSent = lastMap.get(id) || 0; + const now = ts || Date.now(); + if (now - lastSent < STREAM_INTERVAL_MS) return; + lastMap.set(id, now); + sendFrame(socket, id, { ts }, buffer); + }); + }); + + roomCameraStreamEvents.on('status', ({ id, error }) => { + const bucket = cameraSubscribers.get(id); + if (!bucket) return; + bucket.forEach((socketId) => { + const socket = io.sockets.sockets.get(socketId); + if (!socket) return; + sendStatus(socket, id, { error: error || null }); + }); + }); + + io.on('connection', (socket) => { + socket.on('roomCamera:subscribe', (payload = {}, cb = () => {}) => { + const list = Array.isArray(payload?.ids) + ? payload.ids.map(String) + : payload?.roomCameraId || payload?.id + ? [String(payload.roomCameraId || payload.id)] + : getRoomCameras().map((cam) => cam.id); + const uniqueIds = Array.from(new Set(list)); + try { + if (!allowSubscribe(socket.id)) return cb({ error: 'Rate limited' }); + if (!passesMode(socket)) throw new Error('Not authorized for room camera'); + const validIds = uniqueIds.filter((id) => !!getRoomCamera(id)); + validIds.forEach((cameraId) => addSubscription(socket, cameraId)); + validIds.forEach((cameraId) => { + const state = getRoomCameraState(cameraId); + if (state?.frame) sendFrame(socket, cameraId, { ts: state.ts }, state.frame); + sendStatus(socket, cameraId, { ts: state?.ts || null, error: state?.error || null }); + }); + cb({ ok: true, subscribed: validIds }); + } catch (err) { + logger.warn('Room camera subscribe failed', { socketId: socket.id, err: err.message }); + cb({ error: err.message }); + } + }); + + socket.on('roomCamera:unsubscribe', (payload = {}) => { + const list = Array.isArray(payload?.ids) + ? payload.ids.map(String) + : payload?.roomCameraId || payload?.id + ? [String(payload.roomCameraId || payload.id)] + : []; + list.forEach((cameraId) => removeSubscription(socket.id, cameraId)); + }); + + socket.on('disconnect', () => { + removeAllSubscriptions(socket.id); + subscribeBuckets.delete(socket.id); + lastSentBySocket.delete(socket.id); + }); + }); +} + +module.exports = { + registerRoomCameraSocketGateway, +}; diff --git a/server/src/services/roomCameraSnapshotService/index.js b/server/src/services/roomCameraSnapshotService/index.js deleted file mode 100644 index 1eafb921..00000000 --- a/server/src/services/roomCameraSnapshotService/index.js +++ /dev/null @@ -1,225 +0,0 @@ -// room Camera Snapshot Service -// Purpose: Defines the room Camera Snapshot Service module and the helpers/state used by this service unit. -// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary. -const EventEmitter = require('events'); -const http = require('http'); -const https = require('https'); -const logger = require('../../globals/logger').child('roomCameraSnapshot'); -const { getRoomCameras, roomCameraEvents } = require('../roomCameraService'); - -const POLL_INTERVAL_MS = 67; -const FETCH_TIMEOUT_MS = 2000; -const STREAM_RETRY_MS = 1500; - -const cameraState = new Map(); // id -> {frame, ts, error, failures, fetching} -const events = new EventEmitter(); // frame, status -let pollTimer = null; -const streamState = new Map(); // id -> { req, reconnectTimer } -const JPEG_START = Buffer.from([0xff, 0xd8]); -const JPEG_END = Buffer.from([0xff, 0xd9]); - -function markState(id, updates = {}) { - const prev = cameraState.get(id) || {}; - const next = { ...prev, ...updates }; - cameraState.set(id, next); - return next; -} - -function isMjpegUrl(rawUrl) { - if (!rawUrl) return false; - const lower = String(rawUrl).toLowerCase(); - return lower.endsWith('.mjpg') || lower.endsWith('.mjpeg') || lower.includes('mjpeg'); -} - -function getStreamUrl(camera) { - if (camera.streamUrl) return camera.streamUrl; - if (isMjpegUrl(camera.url)) return camera.url; - return null; -} - -function stopStream(id) { - const entry = streamState.get(id); - if (!entry) return; - if (entry.req) { - entry.req.destroy(); - } - if (entry.reconnectTimer) { - clearTimeout(entry.reconnectTimer); - } - streamState.delete(id); -} - -function scheduleStreamReconnect(camera) { - const id = camera.id; - const entry = streamState.get(id) || {}; - if (entry.reconnectTimer) return; - entry.reconnectTimer = setTimeout(() => { - const current = streamState.get(id) || {}; - current.reconnectTimer = null; - streamState.set(id, current); - startStream(camera); - }, STREAM_RETRY_MS); - streamState.set(id, entry); -} - -function clearStreamRequest(id, req) { - const entry = streamState.get(id); - if (!entry) return; - if (entry.req !== req) return; - entry.req = null; - streamState.set(id, entry); -} - -function handleStreamError(camera, err) { - const state = cameraState.get(camera.id) || {}; - const failures = (state.failures || 0) + 1; - markState(camera.id, { error: err.message || 'stream error', failures }); - events.emit('status', { id: camera.id, error: err.message || 'stream error' }); - logger.warn('Stream failed', { id: camera.id, err: err.message }); -} - -function startStream(camera) { - const streamUrl = getStreamUrl(camera); - if (!streamUrl) return; - if (streamState.get(camera.id)?.req) return; - const url = new URL(streamUrl); - const client = url.protocol === 'https:' ? https : http; - const req = client.get( - { - hostname: url.hostname, - port: url.port || (url.protocol === 'https:' ? 443 : 80), - path: `${url.pathname}${url.search}`, - headers: { Accept: 'multipart/x-mixed-replace' }, - }, - (res) => { - if (res.statusCode !== 200) { - res.resume(); - clearStreamRequest(camera.id, req); - handleStreamError(camera, new Error(`HTTP ${res.statusCode}`)); - scheduleStreamReconnect(camera); - return; - } - let buffer = Buffer.alloc(0); - res.on('data', (chunk) => { - buffer = Buffer.concat([buffer, chunk]); - while (true) { - let start = buffer.indexOf(JPEG_START); - if (start === -1) { - if (buffer.length > 2 * 1024 * 1024) { - buffer = buffer.slice(-1024 * 1024); - } - break; - } - if (start > 0) { - buffer = buffer.slice(start); - start = 0; - } - const end = buffer.indexOf(JPEG_END, 2); - if (end === -1) break; - const frame = buffer.slice(0, end + 2); - buffer = buffer.slice(end + 2); - const ts = Date.now(); - markState(camera.id, { frame, ts, error: null, failures: 0 }); - events.emit('frame', { id: camera.id, buffer: frame, ts }); - } - }); - res.on('end', () => { - clearStreamRequest(camera.id, req); - scheduleStreamReconnect(camera); - }); - res.on('aborted', () => { - clearStreamRequest(camera.id, req); - handleStreamError(camera, new Error('stream aborted')); - scheduleStreamReconnect(camera); - }); - res.on('error', (err) => { - clearStreamRequest(camera.id, req); - handleStreamError(camera, err); - scheduleStreamReconnect(camera); - }); - }, - ); - req.on('error', (err) => { - clearStreamRequest(camera.id, req); - handleStreamError(camera, err); - scheduleStreamReconnect(camera); - }); - streamState.set(camera.id, { req, reconnectTimer: null }); -} - -async function fetchSnapshot(camera) { - const { id, url } = camera; - const state = cameraState.get(id); - if (!url || state?.fetching) return; - markState(id, { fetching: true }); - const abortController = new AbortController(); - const timeout = setTimeout(() => abortController.abort(), FETCH_TIMEOUT_MS); - try { - const res = await fetch(url, { signal: abortController.signal }); - if (!res.ok) { - throw new Error(`HTTP ${res.status}`); - } - const arrayBuffer = await res.arrayBuffer(); - const buffer = Buffer.from(arrayBuffer); - const ts = Date.now(); - markState(id, { frame: buffer, ts, error: null, failures: 0 }); - events.emit('frame', { id, buffer, ts }); - } catch (err) { - const failures = (state?.failures || 0) + 1; - markState(id, { error: err.message, failures }); - events.emit('status', { id, error: err.message }); - logger.warn('Snapshot fetch failed', { id, err: err.message }); - } finally { - clearTimeout(timeout); - markState(id, { fetching: false }); - } -} - -function stopAll() { - if (pollTimer) { - clearInterval(pollTimer); - pollTimer = null; - } - Array.from(streamState.keys()).forEach((id) => stopStream(id)); - cameraState.clear(); -} - -function startAll() { - stopAll(); - const cameras = getRoomCameras(); - const snapshotCameras = cameras.filter((camera) => !getStreamUrl(camera) && camera.url); - cameras.forEach((camera) => startStream(camera)); - if (snapshotCameras.length) { - pollTimer = setInterval(() => { - snapshotCameras.forEach((camera) => fetchSnapshot(camera)); - }, POLL_INTERVAL_MS); - snapshotCameras.forEach((camera) => fetchSnapshot(camera)); - } - logger.info('Started room camera feeds', { - total: cameras.length, - streaming: cameras.length - snapshotCameras.length, - snapshots: snapshotCameras.length, - }); -} - -function getState(id) { - const state = cameraState.get(id); - if (!state) return null; - return { - frame: state.frame || null, - ts: state.ts || null, - error: state.error || null, - }; -} - -roomCameraEvents.on('update', () => { - logger.info('Room cameras changed; restarting snapshot pollers'); - startAll(); -}); - -startAll(); - -module.exports = { - roomCameraStreamEvents: events, - getRoomCameraState: getState, -}; diff --git a/server/src/services/roomCameraSocketService/index.js b/server/src/services/roomCameraSocketService/index.js deleted file mode 100644 index 073280b8..00000000 --- a/server/src/services/roomCameraSocketService/index.js +++ /dev/null @@ -1,170 +0,0 @@ -// room Camera Socket Service -// Purpose: Defines the room Camera Socket Service module and the helpers/state used by this service unit. -// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary. -const io = require('../../globals/io'); -const logger = require('../../globals/logger').child('roomCameraSocket'); -const { getMode, MODES } = require('../modeManager'); -const { isAdmin, isLockdownAdmin, getRole } = require('../roleService'); -const { getRoomCamera, getRoomCameras } = require('../roomCameraService'); -const { roomCameraStreamEvents, getRoomCameraState } = require('../roomCameraSnapshotService'); - -const SUBSCRIBE_LIMIT = 50; -const SUBSCRIBE_WINDOW_MS = 10000; -const STREAM_INTERVAL_MS = 1000; - -function passesMode(socket) { - const mode = getMode(); - if (mode === MODES.LOCKDOWN) { - return isLockdownAdmin(socket); - } - if (mode === MODES.ADMIN) { - const role = getRole(socket); - return role === 'spectator' || isAdmin(socket); - } - return true; -} - -function canViewRoomCamera(socket) { - return passesMode(socket); -} - -const cameraSubscribers = new Map(); // id -> Set(socketId) -const socketSubscriptions = new Map(); // socketId -> Set(id) -const subscribeBuckets = new Map(); // socketId -> { start, count } -const lastSentBySocket = new Map(); // socketId -> Map(cameraId -> ts) - -function addSubscription(socket, cameraId) { - if (!cameraSubscribers.has(cameraId)) { - cameraSubscribers.set(cameraId, new Set()); - } - cameraSubscribers.get(cameraId).add(socket.id); - - if (!socketSubscriptions.has(socket.id)) { - socketSubscriptions.set(socket.id, new Set()); - } - socketSubscriptions.get(socket.id).add(cameraId); -} - -function removeSubscription(socketId, cameraId) { - const bucket = cameraSubscribers.get(cameraId); - if (bucket) { - bucket.delete(socketId); - if (bucket.size === 0) { - cameraSubscribers.delete(cameraId); - } - } - const socketBucket = socketSubscriptions.get(socketId); - if (socketBucket) { - socketBucket.delete(cameraId); - if (socketBucket.size === 0) { - socketSubscriptions.delete(socketId); - } - } -} - -function removeAllSubscriptions(socketId) { - const bucket = socketSubscriptions.get(socketId); - if (!bucket) return; - bucket.forEach((cameraId) => removeSubscription(socketId, cameraId)); -} - -function allowSubscribe(socketId) { - const now = Date.now(); - let bucket = subscribeBuckets.get(socketId); - if (!bucket || now - bucket.start >= SUBSCRIBE_WINDOW_MS) { - bucket = { start: now, count: 0 }; - } - bucket.count += 1; - subscribeBuckets.set(socketId, bucket); - return bucket.count <= SUBSCRIBE_LIMIT; -} - -function sendFrame(socket, cameraId, payload, buffer) { - socket.emit('roomCamera:frame', { id: cameraId, ...payload }, buffer); -} - -function sendStatus(socket, cameraId, status) { - socket.emit('roomCamera:status', { id: cameraId, ...status }); -} - -roomCameraStreamEvents.on('frame', ({ id, buffer, ts }) => { - const bucket = cameraSubscribers.get(id); - if (!bucket || !buffer) return; - bucket.forEach((socketId) => { - const socket = io.sockets.sockets.get(socketId); - if (!socket) return; - let lastMap = lastSentBySocket.get(socketId); - if (!lastMap) { - lastMap = new Map(); - lastSentBySocket.set(socketId, lastMap); - } - const lastSent = lastMap.get(id) || 0; - const now = ts || Date.now(); - if (now - lastSent < STREAM_INTERVAL_MS) { - return; - } - lastMap.set(id, now); - sendFrame(socket, id, { ts }, buffer); - }); -}); - -roomCameraStreamEvents.on('status', ({ id, error }) => { - const bucket = cameraSubscribers.get(id); - if (!bucket) return; - bucket.forEach((socketId) => { - const socket = io.sockets.sockets.get(socketId); - if (!socket) return; - sendStatus(socket, id, { error: error || null }); - }); -}); - -io.on('connection', (socket) => { - socket.on('roomCamera:subscribe', (payload = {}, cb = () => {}) => { - const list = Array.isArray(payload?.ids) - ? payload.ids.map(String) - : payload?.roomCameraId || payload?.id - ? [String(payload.roomCameraId || payload.id)] - : getRoomCameras().map((cam) => cam.id); - const uniqueIds = Array.from(new Set(list)); - try { - if (!allowSubscribe(socket.id)) { - cb({ error: 'Rate limited' }); - return; - } - if (!canViewRoomCamera(socket)) { - throw new Error('Not authorized for room camera'); - } - const validIds = uniqueIds.filter((id) => !!getRoomCamera(id)); - validIds.forEach((cameraId) => addSubscription(socket, cameraId)); - validIds.forEach((cameraId) => { - const state = getRoomCameraState(cameraId); - if (state?.frame) { - sendFrame(socket, cameraId, { ts: state.ts }, state.frame); - } - sendStatus(socket, cameraId, { - ts: state?.ts || null, - error: state?.error || null, - }); - }); - cb({ ok: true, subscribed: validIds }); - } catch (err) { - logger.warn('Room camera subscribe failed', { socketId: socket.id, err: err.message }); - cb({ error: err.message }); - } - }); - - socket.on('roomCamera:unsubscribe', (payload = {}) => { - const list = Array.isArray(payload?.ids) - ? payload.ids.map(String) - : payload?.roomCameraId || payload?.id - ? [String(payload.roomCameraId || payload.id)] - : []; - list.forEach((cameraId) => removeSubscription(socket.id, cameraId)); - }); - - socket.on('disconnect', () => { - removeAllSubscriptions(socket.id); - subscribeBuckets.delete(socket.id); - lastSentBySocket.delete(socket.id); - }); -}); diff --git a/server/src/services/roverSnapshotService/index.js b/server/src/services/roverSnapshotService/index.js index f3c63264..ac3711ec 100644 --- a/server/src/services/roverSnapshotService/index.js +++ b/server/src/services/roverSnapshotService/index.js @@ -1,97 +1,20 @@ -// rover Snapshot Service -// Purpose: Defines the rover Snapshot Service module and the helpers/state used by this service unit. -// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary. -const EventEmitter = require('events'); -const fs = require('fs/promises'); -const path = require('path'); -const logger = require('../../globals/logger').child('roverSnapshot'); +// Rover Snapshot Service +// Purpose: Composes rover snapshot polling and snapshot socket delivery in one folderized service. +// Scope: Preserves existing rover snapshot events/state exports and startup side effects. const roverManager = require('../roverManager'); +const { createRoverSnapshotPoller } = require('./poller'); +const { registerRoverSnapshotSocketGateway } = require('./socketGateway'); -const SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/var/lib/rover-snapshots'; -const POLL_INTERVAL_MS = 300; +const poller = createRoverSnapshotPoller({ roverManager }); +poller.startAll(); -const roverState = new Map(); // id -> { frame, ts, error, failures, fetching, mtimeMs } -const events = new EventEmitter(); // frame, status -let pollTimer = null; - -function markState(id, updates = {}) { - const prev = roverState.get(id) || {}; - const next = { ...prev, ...updates }; - roverState.set(id, next); - return next; -} - -function getSnapshotPath(id) { - return path.join(SNAPSHOT_DIR, `${id}.jpg`); -} - -async function fetchSnapshot(id) { - const state = roverState.get(id); - if (state?.fetching) return; - markState(id, { fetching: true }); - try { - const filePath = getSnapshotPath(id); - const stats = await fs.stat(filePath); - if (state?.mtimeMs && stats.mtimeMs <= state.mtimeMs) { - return; - } - const buffer = await fs.readFile(filePath); - const ts = stats.mtimeMs || Date.now(); - markState(id, { frame: buffer, ts, error: null, failures: 0, mtimeMs: stats.mtimeMs }); - events.emit('frame', { id, buffer, ts }); - } catch (err) { - const failures = (state?.failures || 0) + 1; - const message = err.code === 'ENOENT' ? 'Snapshot missing' : err.message; - markState(id, { error: message, failures }); - events.emit('status', { id, error: message }); - if (failures % 20 === 1) { - logger.warn('Snapshot read failed', { id, err: message }); - } - } finally { - markState(id, { fetching: false }); - } -} - -function cleanupInactive(activeIds) { - roverState.forEach((_, id) => { - if (!activeIds.has(id)) { - roverState.delete(id); - } - }); -} - -function stopAll() { - if (pollTimer) { - clearInterval(pollTimer); - pollTimer = null; - } - roverState.clear(); -} - -function startAll() { - stopAll(); - pollTimer = setInterval(() => { - const roster = roverManager.getRoster(); - const activeIds = new Set(roster.map((entry) => String(entry.id))); - cleanupInactive(activeIds); - roster.forEach((entry) => fetchSnapshot(String(entry.id))); - }, POLL_INTERVAL_MS); - logger.info('Started rover snapshot polling'); -} - -function getState(id) { - const state = roverState.get(id); - if (!state) return null; - return { - frame: state.frame || null, - ts: state.ts || null, - error: state.error || null, - }; -} - -startAll(); +registerRoverSnapshotSocketGateway({ + roverManager, + roverSnapshotEvents: poller.roverSnapshotEvents, + getRoverSnapshotState: poller.getRoverSnapshotState, +}); module.exports = { - roverSnapshotEvents: events, - getRoverSnapshotState: getState, + roverSnapshotEvents: poller.roverSnapshotEvents, + getRoverSnapshotState: poller.getRoverSnapshotState, }; diff --git a/server/src/services/roverSnapshotService/poller.js b/server/src/services/roverSnapshotService/poller.js new file mode 100644 index 00000000..ac0d6182 --- /dev/null +++ b/server/src/services/roverSnapshotService/poller.js @@ -0,0 +1,96 @@ +// Rover Snapshot Poller +// Purpose: Polls on-disk rover snapshot images and emits frame/status updates as rover files change. +// Scope: Owns snapshot file IO, stale-rover cleanup, and in-memory state tracking. +const EventEmitter = require('events'); +const fs = require('fs/promises'); +const path = require('path'); +const logger = require('../../globals/logger').child('roverSnapshot'); + +const SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/var/lib/rover-snapshots'; +const POLL_INTERVAL_MS = 300; + +const roverState = new Map(); +const events = new EventEmitter(); +let pollTimer = null; + +function markState(id, updates = {}) { + const prev = roverState.get(id) || {}; + const next = { ...prev, ...updates }; + roverState.set(id, next); + return next; +} + +function getSnapshotPath(id) { + return path.join(SNAPSHOT_DIR, `${id}.jpg`); +} + +function createRoverSnapshotPoller({ roverManager }) { + async function fetchSnapshot(id) { + const state = roverState.get(id); + if (state?.fetching) return; + markState(id, { fetching: true }); + try { + const filePath = getSnapshotPath(id); + const stats = await fs.stat(filePath); + if (state?.mtimeMs && stats.mtimeMs <= state.mtimeMs) return; + const buffer = await fs.readFile(filePath); + const ts = stats.mtimeMs || Date.now(); + markState(id, { frame: buffer, ts, error: null, failures: 0, mtimeMs: stats.mtimeMs }); + events.emit('frame', { id, buffer, ts }); + } catch (err) { + const failures = (state?.failures || 0) + 1; + const message = err.code === 'ENOENT' ? 'Snapshot missing' : err.message; + markState(id, { error: message, failures }); + events.emit('status', { id, error: message }); + if (failures % 20 === 1) logger.warn('Snapshot read failed', { id, err: message }); + } finally { + markState(id, { fetching: false }); + } + } + + function cleanupInactive(activeIds) { + roverState.forEach((_, id) => { + if (!activeIds.has(id)) roverState.delete(id); + }); + } + + function stopAll() { + if (pollTimer) { + clearInterval(pollTimer); + pollTimer = null; + } + roverState.clear(); + } + + function startAll() { + stopAll(); + pollTimer = setInterval(() => { + const roster = roverManager.getRoster(); + const activeIds = new Set(roster.map((entry) => String(entry.id))); + cleanupInactive(activeIds); + roster.forEach((entry) => fetchSnapshot(String(entry.id))); + }, POLL_INTERVAL_MS); + logger.info('Started rover snapshot polling'); + } + + function getRoverSnapshotState(id) { + const state = roverState.get(id); + if (!state) return null; + return { + frame: state.frame || null, + ts: state.ts || null, + error: state.error || null, + }; + } + + return { + startAll, + stopAll, + roverSnapshotEvents: events, + getRoverSnapshotState, + }; +} + +module.exports = { + createRoverSnapshotPoller, +}; diff --git a/server/src/services/roverSnapshotService/socketGateway.js b/server/src/services/roverSnapshotService/socketGateway.js new file mode 100644 index 00000000..5d18626e --- /dev/null +++ b/server/src/services/roverSnapshotService/socketGateway.js @@ -0,0 +1,148 @@ +// Rover Snapshot Socket Gateway +// Purpose: Handles rover snapshot socket subscriptions and streams rover frame/status updates to clients. +// Scope: Owns auth/rate-limit checks and per-socket subscription fan-out for rover snapshot events. +const io = require('../../globals/io'); +const logger = require('../../globals/logger').child('roverSnapshotSocket'); +const { getMode, MODES } = require('../modeManager'); +const { isAdmin, isLockdownAdmin, getRole } = require('../roleService'); + +const SUBSCRIBE_LIMIT = 50; +const SUBSCRIBE_WINDOW_MS = 10000; +const STREAM_INTERVAL_MS = 333; + +function passesMode(socket) { + const mode = getMode(); + if (mode === MODES.LOCKDOWN) return isLockdownAdmin(socket); + if (mode === MODES.ADMIN) { + const role = getRole(socket); + return role === 'spectator' || isAdmin(socket); + } + return true; +} + +function registerRoverSnapshotSocketGateway({ roverManager, roverSnapshotEvents, getRoverSnapshotState }) { + const roverSubscribers = new Map(); + const socketSubscriptions = new Map(); + const subscribeBuckets = new Map(); + const lastSentBySocket = new Map(); + + function addSubscription(socket, roverId) { + if (!roverSubscribers.has(roverId)) roverSubscribers.set(roverId, new Set()); + roverSubscribers.get(roverId).add(socket.id); + if (!socketSubscriptions.has(socket.id)) socketSubscriptions.set(socket.id, new Set()); + socketSubscriptions.get(socket.id).add(roverId); + } + + function removeSubscription(socketId, roverId) { + const bucket = roverSubscribers.get(roverId); + if (bucket) { + bucket.delete(socketId); + if (bucket.size === 0) roverSubscribers.delete(roverId); + } + const socketBucket = socketSubscriptions.get(socketId); + if (socketBucket) { + socketBucket.delete(roverId); + if (socketBucket.size === 0) socketSubscriptions.delete(socketId); + } + } + + function removeAllSubscriptions(socketId) { + const bucket = socketSubscriptions.get(socketId); + if (!bucket) return; + bucket.forEach((roverId) => removeSubscription(socketId, roverId)); + } + + function allowSubscribe(socketId) { + const now = Date.now(); + let bucket = subscribeBuckets.get(socketId); + if (!bucket || now - bucket.start >= SUBSCRIBE_WINDOW_MS) bucket = { start: now, count: 0 }; + bucket.count += 1; + subscribeBuckets.set(socketId, bucket); + return bucket.count <= SUBSCRIBE_LIMIT; + } + + function sendFrame(socket, roverId, payload, buffer) { + socket.emit('roverSnapshot:frame', { id: roverId, ...payload }, buffer); + } + + function sendStatus(socket, roverId, status) { + socket.emit('roverSnapshot:status', { id: roverId, ...status }); + } + + roverSnapshotEvents.on('frame', ({ id, buffer, ts }) => { + const bucket = roverSubscribers.get(id); + if (!bucket || !buffer) return; + bucket.forEach((socketId) => { + const socket = io.sockets.sockets.get(socketId); + if (!socket) return; + let lastMap = lastSentBySocket.get(socketId); + if (!lastMap) { + lastMap = new Map(); + lastSentBySocket.set(socketId, lastMap); + } + const lastSent = lastMap.get(id) || 0; + const now = ts || Date.now(); + if (now - lastSent < STREAM_INTERVAL_MS) return; + lastMap.set(id, now); + sendFrame(socket, id, { ts }, buffer); + }); + }); + + roverSnapshotEvents.on('status', ({ id, error }) => { + const bucket = roverSubscribers.get(id); + if (!bucket) return; + bucket.forEach((socketId) => { + const socket = io.sockets.sockets.get(socketId); + if (!socket) return; + sendStatus(socket, id, { error: error || null }); + }); + }); + + io.on('connection', (socket) => { + socket.on('roverSnapshot:subscribe', (payload = {}, cb = () => {}) => { + const visibleRoster = roverManager.getRosterForSocket(socket); + const visibleIds = visibleRoster.map((rover) => String(rover.id)); + const list = Array.isArray(payload?.ids) + ? payload.ids.map(String) + : payload?.roverId || payload?.id + ? [String(payload.roverId || payload.id)] + : visibleIds; + const uniqueIds = Array.from(new Set(list)); + try { + if (!allowSubscribe(socket.id)) return cb({ error: 'Rate limited' }); + if (!passesMode(socket)) throw new Error('Not authorized for rover snapshots'); + const rosterIds = new Set(visibleIds); + const validIds = uniqueIds.filter((id) => rosterIds.has(String(id))); + validIds.forEach((roverId) => addSubscription(socket, roverId)); + validIds.forEach((roverId) => { + const state = getRoverSnapshotState(roverId); + if (state?.frame) sendFrame(socket, roverId, { ts: state.ts }, state.frame); + sendStatus(socket, roverId, { ts: state?.ts || null, error: state?.error || null }); + }); + cb({ ok: true, subscribed: validIds }); + } catch (err) { + logger.warn('Rover snapshot subscribe failed', { socketId: socket.id, err: err.message }); + cb({ error: err.message }); + } + }); + + socket.on('roverSnapshot:unsubscribe', (payload = {}) => { + const list = Array.isArray(payload?.ids) + ? payload.ids.map(String) + : payload?.roverId || payload?.id + ? [String(payload.roverId || payload.id)] + : []; + list.forEach((roverId) => removeSubscription(socket.id, roverId)); + }); + + socket.on('disconnect', () => { + removeAllSubscriptions(socket.id); + subscribeBuckets.delete(socket.id); + lastSentBySocket.delete(socket.id); + }); + }); +} + +module.exports = { + registerRoverSnapshotSocketGateway, +}; diff --git a/server/src/services/roverSnapshotSocketService/index.js b/server/src/services/roverSnapshotSocketService/index.js deleted file mode 100644 index 27a3f93e..00000000 --- a/server/src/services/roverSnapshotSocketService/index.js +++ /dev/null @@ -1,173 +0,0 @@ -// rover Snapshot Socket Service -// Purpose: Defines the rover Snapshot Socket Service module and the helpers/state used by this service unit. -// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary. -const io = require('../../globals/io'); -const logger = require('../../globals/logger').child('roverSnapshotSocket'); -const { getMode, MODES } = require('../modeManager'); -const { isAdmin, isLockdownAdmin, getRole } = require('../roleService'); -const roverManager = require('../roverManager'); -const { roverSnapshotEvents, getRoverSnapshotState } = require('../roverSnapshotService'); - -const SUBSCRIBE_LIMIT = 50; -const SUBSCRIBE_WINDOW_MS = 10000; -const STREAM_INTERVAL_MS = 333; - -function passesMode(socket) { - const mode = getMode(); - if (mode === MODES.LOCKDOWN) { - return isLockdownAdmin(socket); - } - if (mode === MODES.ADMIN) { - const role = getRole(socket); - return role === 'spectator' || isAdmin(socket); - } - return true; -} - -function canViewSnapshots(socket) { - return passesMode(socket); -} - -const roverSubscribers = new Map(); // id -> Set(socketId) -const socketSubscriptions = new Map(); // socketId -> Set(id) -const subscribeBuckets = new Map(); // socketId -> { start, count } -const lastSentBySocket = new Map(); // socketId -> Map(roverId -> ts) - -function addSubscription(socket, roverId) { - if (!roverSubscribers.has(roverId)) { - roverSubscribers.set(roverId, new Set()); - } - roverSubscribers.get(roverId).add(socket.id); - - if (!socketSubscriptions.has(socket.id)) { - socketSubscriptions.set(socket.id, new Set()); - } - socketSubscriptions.get(socket.id).add(roverId); -} - -function removeSubscription(socketId, roverId) { - const bucket = roverSubscribers.get(roverId); - if (bucket) { - bucket.delete(socketId); - if (bucket.size === 0) { - roverSubscribers.delete(roverId); - } - } - const socketBucket = socketSubscriptions.get(socketId); - if (socketBucket) { - socketBucket.delete(roverId); - if (socketBucket.size === 0) { - socketSubscriptions.delete(socketId); - } - } -} - -function removeAllSubscriptions(socketId) { - const bucket = socketSubscriptions.get(socketId); - if (!bucket) return; - bucket.forEach((roverId) => removeSubscription(socketId, roverId)); -} - -function allowSubscribe(socketId) { - const now = Date.now(); - let bucket = subscribeBuckets.get(socketId); - if (!bucket || now - bucket.start >= SUBSCRIBE_WINDOW_MS) { - bucket = { start: now, count: 0 }; - } - bucket.count += 1; - subscribeBuckets.set(socketId, bucket); - return bucket.count <= SUBSCRIBE_LIMIT; -} - -function sendFrame(socket, roverId, payload, buffer) { - socket.emit('roverSnapshot:frame', { id: roverId, ...payload }, buffer); -} - -function sendStatus(socket, roverId, status) { - socket.emit('roverSnapshot:status', { id: roverId, ...status }); -} - -roverSnapshotEvents.on('frame', ({ id, buffer, ts }) => { - const bucket = roverSubscribers.get(id); - if (!bucket || !buffer) return; - bucket.forEach((socketId) => { - const socket = io.sockets.sockets.get(socketId); - if (!socket) return; - let lastMap = lastSentBySocket.get(socketId); - if (!lastMap) { - lastMap = new Map(); - lastSentBySocket.set(socketId, lastMap); - } - const lastSent = lastMap.get(id) || 0; - const now = ts || Date.now(); - if (now - lastSent < STREAM_INTERVAL_MS) { - return; - } - lastMap.set(id, now); - sendFrame(socket, id, { ts }, buffer); - }); -}); - -roverSnapshotEvents.on('status', ({ id, error }) => { - const bucket = roverSubscribers.get(id); - if (!bucket) return; - bucket.forEach((socketId) => { - const socket = io.sockets.sockets.get(socketId); - if (!socket) return; - sendStatus(socket, id, { error: error || null }); - }); -}); - -io.on('connection', (socket) => { - socket.on('roverSnapshot:subscribe', (payload = {}, cb = () => {}) => { - const visibleRoster = roverManager.getRosterForSocket(socket); - const visibleIds = visibleRoster.map((rover) => String(rover.id)); - const list = Array.isArray(payload?.ids) - ? payload.ids.map(String) - : payload?.roverId || payload?.id - ? [String(payload.roverId || payload.id)] - : visibleIds; - const uniqueIds = Array.from(new Set(list)); - try { - if (!allowSubscribe(socket.id)) { - cb({ error: 'Rate limited' }); - return; - } - if (!canViewSnapshots(socket)) { - throw new Error('Not authorized for rover snapshots'); - } - const rosterIds = new Set(visibleIds); - const validIds = uniqueIds.filter((id) => rosterIds.has(String(id))); - validIds.forEach((roverId) => addSubscription(socket, roverId)); - validIds.forEach((roverId) => { - const state = getRoverSnapshotState(roverId); - if (state?.frame) { - sendFrame(socket, roverId, { ts: state.ts }, state.frame); - } - sendStatus(socket, roverId, { - ts: state?.ts || null, - error: state?.error || null, - }); - }); - cb({ ok: true, subscribed: validIds }); - } catch (err) { - logger.warn('Rover snapshot subscribe failed', { socketId: socket.id, err: err.message }); - cb({ error: err.message }); - } - }); - - socket.on('roverSnapshot:unsubscribe', (payload = {}) => { - const list = Array.isArray(payload?.ids) - ? payload.ids.map(String) - : payload?.roverId || payload?.id - ? [String(payload.roverId || payload.id)] - : []; - list.forEach((roverId) => removeSubscription(socket.id, roverId)); - }); - - socket.on('disconnect', () => { - removeAllSubscriptions(socket.id); - subscribeBuckets.delete(socket.id); - lastSentBySocket.delete(socket.id); - }); -});