This commit is contained in:
legop3
2026-04-29 14:03:56 -04:00
parent 0003338ecd
commit 5ebf5fa580
@@ -8,6 +8,8 @@ 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();
@@ -26,19 +28,46 @@ function getSnapshotPath(id) {
function createRoverSnapshotPoller({ roverManager }) { function createRoverSnapshotPoller({ roverManager }) {
async function fetchSnapshot(id) { async function fetchSnapshot(id) {
const state = roverState.get(id); const prev = roverState.get(id);
if (state?.fetching) return; if (prev?.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);
if (state?.mtimeMs && stats.mtimeMs <= state.mtimeMs) return; const now = Date.now();
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 || Date.now(); const ts = stats.mtimeMs || now;
markState(id, { frame: buffer, ts, error: null, failures: 0, mtimeMs: stats.mtimeMs }); const prevFrame = prev?.frame || null;
events.emit('frame', { id, buffer, ts }); const changed =
!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 = (state?.failures || 0) + 1; const failures = (prev?.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 });
@@ -64,6 +93,11 @@ function createRoverSnapshotPoller({ roverManager }) {
function startAll() { function startAll() {
stopAll(); stopAll();
logger.info('Starting rover snapshot polling', {
snapshotDir: SNAPSHOT_DIR,
intervalMs: POLL_INTERVAL_MS,
forceRereadStaleMs: FORCE_REREAD_STALE_MS,
});
pollTimer = setInterval(() => { pollTimer = setInterval(() => {
const roster = roverManager.getRoster(); const roster = roverManager.getRoster();
const activeIds = new Set(roster.map((entry) => String(entry.id))); const activeIds = new Set(roster.map((entry) => String(entry.id)));