From d8a01f9b60e1a587ddd1b1ae5a75d6876e82cef7 Mon Sep 17 00:00:00 2001 From: legop3 Date: Wed, 29 Apr 2026 14:30:24 -0400 Subject: [PATCH] khdslfjyufe --- server/src/helpers/dataPaths.js | 11 ++++ server/src/services/healthService/index.js | 3 +- .../services/roverSnapshotService/index.js | 4 -- .../services/roverSnapshotService/poller.js | 9 +-- .../roverSnapshotService/socketGateway.js | 60 +------------------ 5 files changed, 17 insertions(+), 70 deletions(-) diff --git a/server/src/helpers/dataPaths.js b/server/src/helpers/dataPaths.js index 0dc60afb..2bfe1291 100644 --- a/server/src/helpers/dataPaths.js +++ b/server/src/helpers/dataPaths.js @@ -6,6 +6,8 @@ const path = require('path'); const CANONICAL_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) { try { @@ -35,7 +37,16 @@ function resolveDataPath(fileName) { 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 = { resolveDataDir, resolveDataPath, + resolveRoverSnapshotDir, }; diff --git a/server/src/services/healthService/index.js b/server/src/services/healthService/index.js index 2bfb1820..938ccde0 100644 --- a/server/src/services/healthService/index.js +++ b/server/src/services/healthService/index.js @@ -7,8 +7,9 @@ const roverManager = require('../roverManager'); const { getRoomCameras } = require('../roomCameraService'); const { getRoomCameraState } = require('../roomCameraService'); 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 ROOM_CAMERA_STALE_MS = 5000; const ROVER_SNAPSHOT_STALE_MS = 5000; diff --git a/server/src/services/roverSnapshotService/index.js b/server/src/services/roverSnapshotService/index.js index f4af4bae..ac3711ec 100644 --- a/server/src/services/roverSnapshotService/index.js +++ b/server/src/services/roverSnapshotService/index.js @@ -2,7 +2,6 @@ // 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 logger = require('../../globals/logger').child('roverSnapshotService'); const { createRoverSnapshotPoller } = require('./poller'); const { registerRoverSnapshotSocketGateway } = require('./socketGateway'); @@ -13,11 +12,8 @@ registerRoverSnapshotSocketGateway({ roverManager, roverSnapshotEvents: poller.roverSnapshotEvents, getRoverSnapshotState: poller.getRoverSnapshotState, - fetchSnapshotNow: poller.fetchSnapshotNow, }); -logger.warn('Rover snapshot service initialized'); - module.exports = { roverSnapshotEvents: poller.roverSnapshotEvents, getRoverSnapshotState: poller.getRoverSnapshotState, diff --git a/server/src/services/roverSnapshotService/poller.js b/server/src/services/roverSnapshotService/poller.js index c592ef08..56e4de3c 100644 --- a/server/src/services/roverSnapshotService/poller.js +++ b/server/src/services/roverSnapshotService/poller.js @@ -5,12 +5,12 @@ const EventEmitter = require('events'); const fs = require('fs/promises'); const path = require('path'); 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 roverState = new Map(); const events = new EventEmitter(); -const readCounts = new Map(); let pollTimer = null; function markState(id, updates = {}) { @@ -37,11 +37,6 @@ function createRoverSnapshotPoller({ roverManager }) { const buffer = await fs.readFile(filePath); const ts = stats.mtimeMs || Date.now(); 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 }); return { frame: buffer, ts }; } catch (err) { diff --git a/server/src/services/roverSnapshotService/socketGateway.js b/server/src/services/roverSnapshotService/socketGateway.js index 31ae18ce..5d18626e 100644 --- a/server/src/services/roverSnapshotService/socketGateway.js +++ b/server/src/services/roverSnapshotService/socketGateway.js @@ -9,7 +9,6 @@ const { isAdmin, isLockdownAdmin, getRole } = require('../roleService'); const SUBSCRIBE_LIMIT = 50; const SUBSCRIBE_WINDOW_MS = 10000; const STREAM_INTERVAL_MS = 333; -const FALLBACK_REFRESH_MS = 500; function passesMode(socket) { const mode = getMode(); @@ -21,19 +20,11 @@ function passesMode(socket) { return true; } -function registerRoverSnapshotSocketGateway({ - roverManager, - roverSnapshotEvents, - getRoverSnapshotState, - fetchSnapshotNow, -}) { +function registerRoverSnapshotSocketGateway({ roverManager, roverSnapshotEvents, getRoverSnapshotState }) { const roverSubscribers = new Map(); const socketSubscriptions = new Map(); const subscribeBuckets = new Map(); const lastSentBySocket = new Map(); - const sentCounts = new Map(); - const lastForcedTsBySocket = new Map(); - let fallbackTimer = null; function addSubscription(socket, roverId) { if (!roverSubscribers.has(roverId)) roverSubscribers.set(roverId, new Set()); @@ -78,29 +69,6 @@ function registerRoverSnapshotSocketGateway({ 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 }) => { const bucket = roverSubscribers.get(id); if (!bucket || !buffer) return; @@ -116,12 +84,6 @@ function registerRoverSnapshotSocketGateway({ const now = ts || Date.now(); if (now - lastSent < STREAM_INTERVAL_MS) return; 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); }); }); @@ -137,7 +99,6 @@ function registerRoverSnapshotSocketGateway({ }); io.on('connection', (socket) => { - logger.warn('Rover snapshot gateway saw socket connection', { socketId: socket.id }); socket.on('roverSnapshot:subscribe', (payload = {}, cb = () => {}) => { const visibleRoster = roverManager.getRosterForSocket(socket); 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'); const rosterIds = new Set(visibleIds); 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)); - ensureFallbackTimer(); - 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 }); - } - } + 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 }); @@ -192,9 +139,6 @@ function registerRoverSnapshotSocketGateway({ removeAllSubscriptions(socket.id); subscribeBuckets.delete(socket.id); lastSentBySocket.delete(socket.id); - Array.from(lastForcedTsBySocket.keys()).forEach((key) => { - if (key.startsWith(`${socket.id}:`)) lastForcedTsBySocket.delete(key); - }); }); }); }