mirror of
https://github.com/legop3/MultiRoombaRover.git
synced 2026-09-16 09:31:20 -04:00
low quality preview for non-drivers
This commit is contained in:
@@ -0,0 +1,94 @@
|
||||
const EventEmitter = require('events');
|
||||
const fs = require('fs/promises');
|
||||
const path = require('path');
|
||||
const logger = require('../globals/logger').child('roverSnapshot');
|
||||
const roverManager = require('./roverManager');
|
||||
|
||||
const SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/run/rover-snapshots';
|
||||
const POLL_INTERVAL_MS = 500;
|
||||
|
||||
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();
|
||||
|
||||
module.exports = {
|
||||
roverSnapshotEvents: events,
|
||||
getRoverSnapshotState: getState,
|
||||
};
|
||||
@@ -0,0 +1,168 @@
|
||||
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 = 800;
|
||||
|
||||
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 list = Array.isArray(payload?.ids)
|
||||
? payload.ids.map(String)
|
||||
: payload?.roverId || payload?.id
|
||||
? [String(payload.roverId || payload.id)]
|
||||
: roverManager.getRoster().map((rover) => rover.id);
|
||||
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(roverManager.getRoster().map((entry) => String(entry.id)));
|
||||
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