mirror of
https://github.com/legop3/MultiRoombaRover.git
synced 2026-09-16 09:31:20 -04:00
room camera and rover snapshots
This commit is contained in:
@@ -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');
|
||||
|
||||
@@ -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';
|
||||
|
||||
@@ -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';
|
||||
|
||||
+40
-82
@@ -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,
|
||||
};
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user