mirror of
https://github.com/legop3/MultiRoombaRover.git
synced 2026-09-16 01:21:20 -04:00
rover sapshot
This commit is contained in:
@@ -12,6 +12,7 @@ registerRoverSnapshotSocketGateway({
|
|||||||
roverManager,
|
roverManager,
|
||||||
roverSnapshotEvents: poller.roverSnapshotEvents,
|
roverSnapshotEvents: poller.roverSnapshotEvents,
|
||||||
getRoverSnapshotState: poller.getRoverSnapshotState,
|
getRoverSnapshotState: poller.getRoverSnapshotState,
|
||||||
|
fetchSnapshotNow: poller.fetchSnapshotNow,
|
||||||
});
|
});
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
|
|||||||
@@ -8,9 +8,6 @@ const logger = require('../../globals/logger').child('roverSnapshot');
|
|||||||
|
|
||||||
const SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/var/lib/rover-snapshots';
|
const SNAPSHOT_DIR = process.env.ROVER_SNAPSHOT_DIR || '/var/lib/rover-snapshots';
|
||||||
const POLL_INTERVAL_MS = 300;
|
const POLL_INTERVAL_MS = 300;
|
||||||
const FORCE_REREAD_STALE_MS = 2000;
|
|
||||||
const STALE_WARN_MS = 10000;
|
|
||||||
|
|
||||||
const roverState = new Map();
|
const roverState = new Map();
|
||||||
const events = new EventEmitter();
|
const events = new EventEmitter();
|
||||||
let pollTimer = null;
|
let pollTimer = null;
|
||||||
@@ -28,46 +25,19 @@ function getSnapshotPath(id) {
|
|||||||
|
|
||||||
function createRoverSnapshotPoller({ roverManager }) {
|
function createRoverSnapshotPoller({ roverManager }) {
|
||||||
async function fetchSnapshot(id) {
|
async function fetchSnapshot(id) {
|
||||||
const prev = roverState.get(id);
|
const state = roverState.get(id);
|
||||||
if (prev?.fetching) return;
|
if (state?.fetching) return;
|
||||||
markState(id, { fetching: true });
|
markState(id, { fetching: true });
|
||||||
try {
|
try {
|
||||||
const filePath = getSnapshotPath(id);
|
const filePath = getSnapshotPath(id);
|
||||||
const stats = await fs.stat(filePath);
|
const stats = await fs.stat(filePath);
|
||||||
const now = Date.now();
|
if (state?.mtimeMs && stats.mtimeMs <= state.mtimeMs) return;
|
||||||
const staleAgeMs = prev?.ts ? now - prev.ts : 0;
|
|
||||||
const mtimeUnchanged = Boolean(prev?.mtimeMs && stats.mtimeMs <= prev.mtimeMs);
|
|
||||||
const shouldForceRead = mtimeUnchanged && staleAgeMs >= FORCE_REREAD_STALE_MS;
|
|
||||||
if (mtimeUnchanged && !shouldForceRead) return;
|
|
||||||
|
|
||||||
const buffer = await fs.readFile(filePath);
|
const buffer = await fs.readFile(filePath);
|
||||||
const ts = stats.mtimeMs || now;
|
const ts = stats.mtimeMs || Date.now();
|
||||||
const prevFrame = prev?.frame || null;
|
markState(id, { frame: buffer, ts, error: null, failures: 0, mtimeMs: stats.mtimeMs });
|
||||||
const changed =
|
events.emit('frame', { id, buffer, ts });
|
||||||
!prevFrame ||
|
|
||||||
prevFrame.length !== buffer.length ||
|
|
||||||
!prevFrame.equals(buffer) ||
|
|
||||||
!mtimeUnchanged;
|
|
||||||
const nextState = markState(id, {
|
|
||||||
frame: changed ? buffer : prevFrame,
|
|
||||||
ts: changed ? ts : prev?.ts || ts,
|
|
||||||
error: null,
|
|
||||||
failures: 0,
|
|
||||||
mtimeMs: stats.mtimeMs,
|
|
||||||
});
|
|
||||||
if (changed) {
|
|
||||||
events.emit('frame', { id, buffer, ts: nextState.ts });
|
|
||||||
} else if (staleAgeMs >= STALE_WARN_MS) {
|
|
||||||
logger.warn('Snapshot appears stale', {
|
|
||||||
id,
|
|
||||||
snapshotDir: SNAPSHOT_DIR,
|
|
||||||
path: filePath,
|
|
||||||
ageMs: staleAgeMs,
|
|
||||||
mtimeMs: stats.mtimeMs,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
const failures = (prev?.failures || 0) + 1;
|
const failures = (state?.failures || 0) + 1;
|
||||||
const message = err.code === 'ENOENT' ? 'Snapshot missing' : err.message;
|
const message = err.code === 'ENOENT' ? 'Snapshot missing' : err.message;
|
||||||
markState(id, { error: message, failures });
|
markState(id, { error: message, failures });
|
||||||
events.emit('status', { id, error: message });
|
events.emit('status', { id, error: message });
|
||||||
@@ -96,7 +66,6 @@ function createRoverSnapshotPoller({ roverManager }) {
|
|||||||
logger.info('Starting rover snapshot polling', {
|
logger.info('Starting rover snapshot polling', {
|
||||||
snapshotDir: SNAPSHOT_DIR,
|
snapshotDir: SNAPSHOT_DIR,
|
||||||
intervalMs: POLL_INTERVAL_MS,
|
intervalMs: POLL_INTERVAL_MS,
|
||||||
forceRereadStaleMs: FORCE_REREAD_STALE_MS,
|
|
||||||
});
|
});
|
||||||
pollTimer = setInterval(() => {
|
pollTimer = setInterval(() => {
|
||||||
const roster = roverManager.getRoster();
|
const roster = roverManager.getRoster();
|
||||||
@@ -120,6 +89,7 @@ function createRoverSnapshotPoller({ roverManager }) {
|
|||||||
return {
|
return {
|
||||||
startAll,
|
startAll,
|
||||||
stopAll,
|
stopAll,
|
||||||
|
fetchSnapshotNow: fetchSnapshot,
|
||||||
roverSnapshotEvents: events,
|
roverSnapshotEvents: events,
|
||||||
getRoverSnapshotState,
|
getRoverSnapshotState,
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -20,7 +20,12 @@ function passesMode(socket) {
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
function registerRoverSnapshotSocketGateway({ roverManager, roverSnapshotEvents, getRoverSnapshotState }) {
|
function registerRoverSnapshotSocketGateway({
|
||||||
|
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();
|
||||||
@@ -113,8 +118,19 @@ function registerRoverSnapshotSocketGateway({ roverManager, roverSnapshotEvents,
|
|||||||
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.info('Rover snapshot subscribe', {
|
||||||
|
socketId: socket.id,
|
||||||
|
requestedIds: uniqueIds,
|
||||||
|
visibleIds,
|
||||||
|
subscribedIds: validIds,
|
||||||
|
});
|
||||||
validIds.forEach((roverId) => addSubscription(socket, roverId));
|
validIds.forEach((roverId) => addSubscription(socket, roverId));
|
||||||
validIds.forEach((roverId) => {
|
validIds.forEach((roverId) => {
|
||||||
|
if (typeof fetchSnapshotNow === 'function') {
|
||||||
|
fetchSnapshotNow(String(roverId)).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 });
|
||||||
|
|||||||
Reference in New Issue
Block a user