mirror of
https://github.com/legop3/MultiRoombaRover.git
synced 2026-09-16 17:40:46 -04:00
khdslfjyufe
This commit is contained in:
@@ -6,6 +6,8 @@ const path = require('path');
|
|||||||
|
|
||||||
const CANONICAL_DATA_DIR = path.resolve(__dirname, '..', '..', 'data');
|
const CANONICAL_DATA_DIR = path.resolve(__dirname, '..', '..', 'data');
|
||||||
const LEGACY_DATA_DIR = path.resolve(__dirname, '..', 'data');
|
const LEGACY_DATA_DIR = path.resolve(__dirname, '..', 'data');
|
||||||
|
const CANONICAL_ROVER_SNAPSHOT_DIR = path.join(CANONICAL_DATA_DIR, 'rover-snapshots');
|
||||||
|
const LEGACY_ROVER_SNAPSHOT_DIR = path.join(LEGACY_DATA_DIR, 'rover-snapshots');
|
||||||
|
|
||||||
function pathExists(target) {
|
function pathExists(target) {
|
||||||
try {
|
try {
|
||||||
@@ -35,7 +37,16 @@ function resolveDataPath(fileName) {
|
|||||||
return canonicalPath;
|
return canonicalPath;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function resolveRoverSnapshotDir() {
|
||||||
|
const configured = String(process.env.ROVER_SNAPSHOT_DIR || '').trim();
|
||||||
|
if (configured) return path.resolve(configured);
|
||||||
|
if (pathExists(CANONICAL_ROVER_SNAPSHOT_DIR)) return CANONICAL_ROVER_SNAPSHOT_DIR;
|
||||||
|
if (pathExists(LEGACY_ROVER_SNAPSHOT_DIR)) return LEGACY_ROVER_SNAPSHOT_DIR;
|
||||||
|
return '/var/lib/rover-snapshots';
|
||||||
|
}
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
resolveDataDir,
|
resolveDataDir,
|
||||||
resolveDataPath,
|
resolveDataPath,
|
||||||
|
resolveRoverSnapshotDir,
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -7,8 +7,9 @@ const roverManager = require('../roverManager');
|
|||||||
const { getRoomCameras } = require('../roomCameraService');
|
const { getRoomCameras } = require('../roomCameraService');
|
||||||
const { getRoomCameraState } = require('../roomCameraService');
|
const { getRoomCameraState } = require('../roomCameraService');
|
||||||
const { getReplayHealthSnapshot } = require('../replayEngineV2');
|
const { getReplayHealthSnapshot } = require('../replayEngineV2');
|
||||||
|
const { resolveRoverSnapshotDir } = require('../../helpers/dataPaths');
|
||||||
|
|
||||||
const ROVER_SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/var/lib/rover-snapshots';
|
const ROVER_SNAPSHOT_DIR = resolveRoverSnapshotDir();
|
||||||
const HEALTH_INTERVAL_MS = 5000;
|
const HEALTH_INTERVAL_MS = 5000;
|
||||||
const ROOM_CAMERA_STALE_MS = 5000;
|
const ROOM_CAMERA_STALE_MS = 5000;
|
||||||
const ROVER_SNAPSHOT_STALE_MS = 5000;
|
const ROVER_SNAPSHOT_STALE_MS = 5000;
|
||||||
|
|||||||
@@ -2,7 +2,6 @@
|
|||||||
// Purpose: Composes rover snapshot polling and snapshot socket delivery in one folderized 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.
|
// Scope: Preserves existing rover snapshot events/state exports and startup side effects.
|
||||||
const roverManager = require('../roverManager');
|
const roverManager = require('../roverManager');
|
||||||
const logger = require('../../globals/logger').child('roverSnapshotService');
|
|
||||||
const { createRoverSnapshotPoller } = require('./poller');
|
const { createRoverSnapshotPoller } = require('./poller');
|
||||||
const { registerRoverSnapshotSocketGateway } = require('./socketGateway');
|
const { registerRoverSnapshotSocketGateway } = require('./socketGateway');
|
||||||
|
|
||||||
@@ -13,11 +12,8 @@ registerRoverSnapshotSocketGateway({
|
|||||||
roverManager,
|
roverManager,
|
||||||
roverSnapshotEvents: poller.roverSnapshotEvents,
|
roverSnapshotEvents: poller.roverSnapshotEvents,
|
||||||
getRoverSnapshotState: poller.getRoverSnapshotState,
|
getRoverSnapshotState: poller.getRoverSnapshotState,
|
||||||
fetchSnapshotNow: poller.fetchSnapshotNow,
|
|
||||||
});
|
});
|
||||||
|
|
||||||
logger.warn('Rover snapshot service initialized');
|
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
roverSnapshotEvents: poller.roverSnapshotEvents,
|
roverSnapshotEvents: poller.roverSnapshotEvents,
|
||||||
getRoverSnapshotState: poller.getRoverSnapshotState,
|
getRoverSnapshotState: poller.getRoverSnapshotState,
|
||||||
|
|||||||
@@ -5,12 +5,12 @@ const EventEmitter = require('events');
|
|||||||
const fs = require('fs/promises');
|
const fs = require('fs/promises');
|
||||||
const path = require('path');
|
const path = require('path');
|
||||||
const logger = require('../../globals/logger').child('roverSnapshot');
|
const logger = require('../../globals/logger').child('roverSnapshot');
|
||||||
|
const { resolveRoverSnapshotDir } = require('../../helpers/dataPaths');
|
||||||
|
|
||||||
const SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/var/lib/rover-snapshots';
|
const SNAPSHOT_DIR = resolveRoverSnapshotDir();
|
||||||
const POLL_INTERVAL_MS = 300;
|
const POLL_INTERVAL_MS = 300;
|
||||||
const roverState = new Map();
|
const roverState = new Map();
|
||||||
const events = new EventEmitter();
|
const events = new EventEmitter();
|
||||||
const readCounts = new Map();
|
|
||||||
let pollTimer = null;
|
let pollTimer = null;
|
||||||
|
|
||||||
function markState(id, updates = {}) {
|
function markState(id, updates = {}) {
|
||||||
@@ -37,11 +37,6 @@ function createRoverSnapshotPoller({ roverManager }) {
|
|||||||
const buffer = await fs.readFile(filePath);
|
const buffer = await fs.readFile(filePath);
|
||||||
const ts = stats.mtimeMs || Date.now();
|
const ts = stats.mtimeMs || Date.now();
|
||||||
markState(id, { frame: buffer, ts, error: null, failures: 0, mtimeMs: stats.mtimeMs });
|
markState(id, { frame: buffer, ts, error: null, failures: 0, mtimeMs: stats.mtimeMs });
|
||||||
const count = (readCounts.get(id) || 0) + 1;
|
|
||||||
readCounts.set(id, count);
|
|
||||||
if (count === 1 || count % 30 === 0) {
|
|
||||||
logger.warn('Snapshot read ok', { id, count, bytes: buffer.length, ts, mtimeMs: stats.mtimeMs });
|
|
||||||
}
|
|
||||||
events.emit('frame', { id, buffer, ts });
|
events.emit('frame', { id, buffer, ts });
|
||||||
return { frame: buffer, ts };
|
return { frame: buffer, ts };
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
|
|||||||
@@ -9,7 +9,6 @@ const { isAdmin, isLockdownAdmin, getRole } = require('../roleService');
|
|||||||
const SUBSCRIBE_LIMIT = 50;
|
const SUBSCRIBE_LIMIT = 50;
|
||||||
const SUBSCRIBE_WINDOW_MS = 10000;
|
const SUBSCRIBE_WINDOW_MS = 10000;
|
||||||
const STREAM_INTERVAL_MS = 333;
|
const STREAM_INTERVAL_MS = 333;
|
||||||
const FALLBACK_REFRESH_MS = 500;
|
|
||||||
|
|
||||||
function passesMode(socket) {
|
function passesMode(socket) {
|
||||||
const mode = getMode();
|
const mode = getMode();
|
||||||
@@ -21,19 +20,11 @@ function passesMode(socket) {
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
function registerRoverSnapshotSocketGateway({
|
function registerRoverSnapshotSocketGateway({ roverManager, roverSnapshotEvents, getRoverSnapshotState }) {
|
||||||
roverManager,
|
|
||||||
roverSnapshotEvents,
|
|
||||||
getRoverSnapshotState,
|
|
||||||
fetchSnapshotNow,
|
|
||||||
}) {
|
|
||||||
const roverSubscribers = new Map();
|
const roverSubscribers = new Map();
|
||||||
const socketSubscriptions = new Map();
|
const socketSubscriptions = new Map();
|
||||||
const subscribeBuckets = new Map();
|
const subscribeBuckets = new Map();
|
||||||
const lastSentBySocket = new Map();
|
const lastSentBySocket = new Map();
|
||||||
const sentCounts = new Map();
|
|
||||||
const lastForcedTsBySocket = new Map();
|
|
||||||
let fallbackTimer = null;
|
|
||||||
|
|
||||||
function addSubscription(socket, roverId) {
|
function addSubscription(socket, roverId) {
|
||||||
if (!roverSubscribers.has(roverId)) roverSubscribers.set(roverId, new Set());
|
if (!roverSubscribers.has(roverId)) roverSubscribers.set(roverId, new Set());
|
||||||
@@ -78,29 +69,6 @@ function registerRoverSnapshotSocketGateway({
|
|||||||
socket.emit('roverSnapshot:status', { id: roverId, ...status });
|
socket.emit('roverSnapshot:status', { id: roverId, ...status });
|
||||||
}
|
}
|
||||||
|
|
||||||
function ensureFallbackTimer() {
|
|
||||||
if (fallbackTimer || typeof fetchSnapshotNow !== 'function') return;
|
|
||||||
fallbackTimer = setInterval(() => {
|
|
||||||
socketSubscriptions.forEach((roverIds, socketId) => {
|
|
||||||
const socket = io.sockets.sockets.get(socketId);
|
|
||||||
if (!socket) return;
|
|
||||||
roverIds.forEach(async (roverId) => {
|
|
||||||
try {
|
|
||||||
const result = await fetchSnapshotNow(String(roverId), { force: true });
|
|
||||||
if (!result?.frame) return;
|
|
||||||
const key = `${socketId}:${roverId}`;
|
|
||||||
const prevTs = lastForcedTsBySocket.get(key) || 0;
|
|
||||||
if ((result.ts || 0) <= prevTs) return;
|
|
||||||
lastForcedTsBySocket.set(key, result.ts || Date.now());
|
|
||||||
sendFrame(socket, roverId, { ts: result.ts || Date.now() }, result.frame);
|
|
||||||
} catch (err) {
|
|
||||||
logger.warn('Fallback refresh failed', { socketId, roverId, err: err.message });
|
|
||||||
}
|
|
||||||
});
|
|
||||||
});
|
|
||||||
}, FALLBACK_REFRESH_MS);
|
|
||||||
}
|
|
||||||
|
|
||||||
roverSnapshotEvents.on('frame', ({ id, buffer, ts }) => {
|
roverSnapshotEvents.on('frame', ({ id, buffer, ts }) => {
|
||||||
const bucket = roverSubscribers.get(id);
|
const bucket = roverSubscribers.get(id);
|
||||||
if (!bucket || !buffer) return;
|
if (!bucket || !buffer) return;
|
||||||
@@ -116,12 +84,6 @@ function registerRoverSnapshotSocketGateway({
|
|||||||
const now = ts || Date.now();
|
const now = ts || Date.now();
|
||||||
if (now - lastSent < STREAM_INTERVAL_MS) return;
|
if (now - lastSent < STREAM_INTERVAL_MS) return;
|
||||||
lastMap.set(id, now);
|
lastMap.set(id, now);
|
||||||
const key = `${socketId}:${id}`;
|
|
||||||
const sentCount = (sentCounts.get(key) || 0) + 1;
|
|
||||||
sentCounts.set(key, sentCount);
|
|
||||||
if (sentCount === 1 || sentCount % 30 === 0) {
|
|
||||||
logger.warn('Snapshot frame sent', { socketId, roverId: id, sentCount, ts, bytes: buffer.length });
|
|
||||||
}
|
|
||||||
sendFrame(socket, id, { ts }, buffer);
|
sendFrame(socket, id, { ts }, buffer);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
@@ -137,7 +99,6 @@ function registerRoverSnapshotSocketGateway({
|
|||||||
});
|
});
|
||||||
|
|
||||||
io.on('connection', (socket) => {
|
io.on('connection', (socket) => {
|
||||||
logger.warn('Rover snapshot gateway saw socket connection', { socketId: socket.id });
|
|
||||||
socket.on('roverSnapshot:subscribe', (payload = {}, cb = () => {}) => {
|
socket.on('roverSnapshot:subscribe', (payload = {}, cb = () => {}) => {
|
||||||
const visibleRoster = roverManager.getRosterForSocket(socket);
|
const visibleRoster = roverManager.getRosterForSocket(socket);
|
||||||
const visibleIds = visibleRoster.map((rover) => String(rover.id));
|
const visibleIds = visibleRoster.map((rover) => String(rover.id));
|
||||||
@@ -152,22 +113,8 @@ function registerRoverSnapshotSocketGateway({
|
|||||||
if (!passesMode(socket)) throw new Error('Not authorized for rover snapshots');
|
if (!passesMode(socket)) throw new Error('Not authorized for rover snapshots');
|
||||||
const rosterIds = new Set(visibleIds);
|
const rosterIds = new Set(visibleIds);
|
||||||
const validIds = uniqueIds.filter((id) => rosterIds.has(String(id)));
|
const validIds = uniqueIds.filter((id) => rosterIds.has(String(id)));
|
||||||
logger.warn('Rover snapshot subscribe', {
|
|
||||||
socketId: socket.id,
|
|
||||||
requestedIds: uniqueIds,
|
|
||||||
visibleIds,
|
|
||||||
subscribedIds: validIds,
|
|
||||||
});
|
|
||||||
validIds.forEach((roverId) => addSubscription(socket, roverId));
|
validIds.forEach((roverId) => addSubscription(socket, roverId));
|
||||||
ensureFallbackTimer();
|
validIds.forEach((roverId) => {
|
||||||
validIds.forEach(async (roverId) => {
|
|
||||||
if (typeof fetchSnapshotNow === 'function') {
|
|
||||||
try {
|
|
||||||
await fetchSnapshotNow(String(roverId), { force: true });
|
|
||||||
} catch (err) {
|
|
||||||
logger.warn('Immediate snapshot fetch failed', { roverId, err: err.message });
|
|
||||||
}
|
|
||||||
}
|
|
||||||
const state = getRoverSnapshotState(roverId);
|
const state = getRoverSnapshotState(roverId);
|
||||||
if (state?.frame) sendFrame(socket, roverId, { ts: state.ts }, state.frame);
|
if (state?.frame) sendFrame(socket, roverId, { ts: state.ts }, state.frame);
|
||||||
sendStatus(socket, roverId, { ts: state?.ts || null, error: state?.error || null });
|
sendStatus(socket, roverId, { ts: state?.ts || null, error: state?.error || null });
|
||||||
@@ -192,9 +139,6 @@ function registerRoverSnapshotSocketGateway({
|
|||||||
removeAllSubscriptions(socket.id);
|
removeAllSubscriptions(socket.id);
|
||||||
subscribeBuckets.delete(socket.id);
|
subscribeBuckets.delete(socket.id);
|
||||||
lastSentBySocket.delete(socket.id);
|
lastSentBySocket.delete(socket.id);
|
||||||
Array.from(lastForcedTsBySocket.keys()).forEach((key) => {
|
|
||||||
if (key.startsWith(`${socket.id}:`)) lastForcedTsBySocket.delete(key);
|
|
||||||
});
|
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user