const EventEmitter = require('events'); const io = require('../globals/io'); const logger = require('../globals/logger').child('roverManager'); const { sendAlert } = require('./alertService'); const ALERT_COLOR = '#8bc34a'; const { parseSensorFrame } = require('../helpers/sensorDecoder'); const { MODES, getMode } = require('./modeManager'); const { isAdmin, roleEvents } = require('./roleService'); const { publishEvent } = require('./eventBus'); const videoSessions = require('./videoSessions'); const rovers = new Map(); // roverId -> record const socketToRovers = new Map(); // socketId -> Set(roverId) const spectatorSockets = new Set(); const turnService = require('./turnService'); const managerEvents = new EventEmitter(); const DOCK_GUARD_WINDOW_MS = 2 * 1000; const IDLE_UNDOCKED_MS = 2 * 60 * 1000; const PASSIVE_UNDOCKED_MS = 60 * 1000; const DOCK_GUARD_RETRY_MS = 10 * 1000; const DOCK_COMMAND_BASE64 = Buffer.from([143]).toString('base64'); const BACKOFF_MS = 500; const BACKOFF_SPEED = 300; const backoffTimers = new Map(); // roverId -> Timeout const dockGuardStates = new Map(); // roverId -> guard state function ensureRecord(id) { if (!rovers.has(id)) { rovers.set(id, { id, meta: null, ws: null, lastSensor: null, docked: null, lastBumpAt: null, drivers: new Set(), locked: false, lockReason: null, batteryState: null, room: `rover:${id}`, lastSeen: Date.now(), lastMovementAt: Date.now(), }); } return rovers.get(id); } function upsertRover(meta, ws) { const id = meta.name || meta.id; if (!meta.cameraServo) { logger.info('Rover hello missing camera servo block', { id, keys: Object.keys(meta || {}) }); } else { logger.info('Rover hello camera servo', { id, servo: meta.cameraServo }); } const isNew = !rovers.has(id); const record = ensureRecord(id); record.meta = meta; record.ws = ws; record.lastSeen = Date.now(); rovers.set(id, record); spectatorSockets.forEach((socketId) => { const sock = io.sockets.sockets.get(socketId); sock?.join(record.room); }); managerEvents.emit('rover', { roverId: id, action: 'upsert', record }); if (isNew) { publishEvent({ source: 'roverManager', type: 'rover.online', payload: { roverId: id } }); } broadcastRoster(); return record; } function removeRover(id) { const record = rovers.get(id); if (!record) return; rovers.delete(id); stopDockGuard(id); turnService.cleanupRover(id); spectatorSockets.forEach((socketId) => { const sock = io.sockets.sockets.get(socketId); sock?.leave(record.room); }); broadcastRoster(); managerEvents.emit('rover', { roverId: id, action: 'removed' }); publishEvent({ source: 'roverManager', type: 'rover.offline', payload: { roverId: id } }); } function lockRover(id, locked, options = {}) { const record = rovers.get(id); if (!record) { throw new Error('Unknown rover'); } const reason = locked ? options.reason || 'manual' : null; const silent = Boolean(options.silent); if (locked) { record.locked = true; record.lockReason = reason; if (!silent) { sendAlert({ color: ALERT_COLOR, title: 'Rover Locked', message: `${id} locked${record.lockReason ? ` (${record.lockReason})` : ''}.`, }); } publishEvent({ source: 'roverManager', type: 'rover.locked', payload: { roverId: id, reason: record.lockReason }, }); } else { record.locked = false; record.lockReason = null; if (!silent) { sendAlert({ color: ALERT_COLOR, title: 'Rover Unlocked', message: `${id} unlocked.` }); } publishEvent({ source: 'roverManager', type: 'rover.unlocked', payload: { roverId: id }, }); } broadcastRoster(); managerEvents.emit('lock', { roverId: id, locked: record.locked, reason: record.lockReason }); return record.locked; } function getRoster() { return Array.from(rovers.values()).map((record) => ({ id: record.id, name: record.meta?.name || record.id, battery: record.meta?.battery, batteryState: record.batteryState, maxWheelSpeed: record.meta?.maxWheelSpeed, media: record.meta?.media, cameraServo: record.meta?.cameraServo, audio: record.meta?.audio, nightVision: record.meta?.nightVision, locked: record.locked, lockReason: record.lockReason, lastSeen: record.lastSeen, })); } function broadcastRoster() { io.emit('rovers', getRoster()); } function computeBatteryState(record, sensors) { if (!record) return null; if (!sensors) return record.batteryState; const config = record.meta?.battery || null; const charge = sensors?.batteryChargeMah ?? null; const capacity = sensors?.batteryCapacityMah ?? null; const full = typeof config?.Full === 'number' ? config.Full : null; const warn = typeof config?.Warn === 'number' ? config.Warn : null; const urgent = typeof config?.Urgent === 'number' ? config.Urgent : null; let percent = null; if (charge != null && warn != null && full != null && full > warn) { const span = full - warn; percent = (charge - warn) / span; } else if (charge != null && capacity != null && capacity > 0) { percent = charge / capacity; } if (percent != null) { percent = Math.max(0, Math.min(1, percent)); } return { charge, capacity, full, warn, urgent, percent, percentDisplay: percent == null ? null : Math.round(percent * 100), warnActive: Boolean(warn != null && charge != null && charge <= warn), urgentActive: Boolean(urgent != null && charge != null && charge <= urgent), updatedAt: Date.now(), }; } function handleSensorFrame(roverId, frame) { const record = rovers.get(roverId); if (!record) return; record.lastSeen = Date.now(); const decoded = parseSensorFrame(frame.data); record.lastSensor = { raw: frame, decoded }; record.batteryState = computeBatteryState(record, decoded); updateMovement(record, decoded); const hasDockInfo = decoded?.chargingSources != null; if (hasDockInfo) { const prevDocked = record.docked; const docked = Boolean(decoded?.chargingSources?.homeBase); record.docked = docked; if (prevDocked === true && docked === false) { handleIdleUndock(record); } } const bumps = decoded?.bumpsAndWheelDrops; if (bumps?.bumpLeft || bumps?.bumpRight) { record.lastBumpAt = Date.now(); } io.to(record.room).volatile.emit('sensorFrame', { roverId, frame, sensors: decoded, }); managerEvents.emit('sensor', { roverId, sensors: decoded, batteryState: record.batteryState }); evaluateDockGuard(record, decoded); } function updateMovement(record, sensors) { if (!record || !sensors) return; const distance = Math.abs(sensors.distanceMm ?? 0); const angle = Math.abs(sensors.angleDeg ?? 0); const requested = Math.abs(sensors.requestedVelocity ?? 0); const requestedLeft = Math.abs(sensors.requestedLeftVelocity ?? 0); const requestedRight = Math.abs(sensors.requestedRightVelocity ?? 0); const moving = distance > 0 || angle > 0 || requested > 0 || requestedLeft > 0 || requestedRight > 0; if (moving) { record.lastMovementAt = Date.now(); } } function getDockGuardState(roverId) { if (!dockGuardStates.has(roverId)) { dockGuardStates.set(roverId, { idleUndockedSince: null, passiveUndockedSince: null, active: false, reason: null, startedAt: null, timer: null, }); } return dockGuardStates.get(roverId); } function evaluateDockGuard(record, sensors) { if (!record || !sensors) return; const state = getDockGuardState(record.id); const now = Date.now(); const docked = Boolean(sensors?.chargingSources?.homeBase); const oiMode = sensors?.oiMode?.label || null; const idleMs = record.lastMovementAt ? now - record.lastMovementAt : 0; const isIdle = idleMs >= 1000; if (docked || record.drivers.size > 0 || !isIdle) { state.idleUndockedSince = null; } else if (!state.idleUndockedSince) { state.idleUndockedSince = now; } if (docked || oiMode !== 'passive' || !isIdle) { state.passiveUndockedSince = null; } else if (!state.passiveUndockedSince) { state.passiveUndockedSince = now; } if (state.active) { const shouldStop = docked || !isIdle || (state.reason === 'idle' && record.drivers.size > 0) || (state.reason === 'passive' && oiMode !== 'passive'); if (shouldStop) { stopDockGuard(record.id); } return; } const idleReady = state.idleUndockedSince && now - state.idleUndockedSince >= IDLE_UNDOCKED_MS; const passiveReady = state.passiveUndockedSince && now - state.passiveUndockedSince >= PASSIVE_UNDOCKED_MS; if (passiveReady) { startDockGuard(record, 'passive', idleMs); } else if (idleReady) { startDockGuard(record, 'idle', idleMs); } } function startDockGuard(record, reason, idleMs) { if (!record) return; const state = getDockGuardState(record.id); if (state.active) return; state.active = true; state.reason = reason; state.startedAt = Date.now(); const reasonText = reason === 'passive' ? 'passive mode' : 'idle and undocked'; sendAlert({ color: ALERT_COLOR, title: 'Dock Guard Triggered', message: `${record.id} ${reasonText}. Seeking dock and restarting sensors until movement.`, }); publishEvent({ source: 'roverManager', type: 'rover.dockGuard', payload: { roverId: record.id, reason, reasonText, idleMs, }, }); attemptDockGuard(record.id); state.timer = setInterval(() => attemptDockGuard(record.id), DOCK_GUARD_RETRY_MS); } function stopDockGuard(roverId) { const state = dockGuardStates.get(roverId); if (!state) return; if (state.timer) { clearInterval(state.timer); } state.active = false; state.reason = null; state.startedAt = null; state.timer = null; state.idleUndockedSince = null; state.passiveUndockedSince = null; } function attemptDockGuard(roverId) { const record = rovers.get(roverId); if (!record) { stopDockGuard(roverId); return; } const { issueCommand } = require('./commandService'); try { issueCommand(roverId, { type: 'sensorStream', sensorStream: { enable: true } }); issueCommand(roverId, { type: 'raw', raw: DOCK_COMMAND_BASE64 }); } catch (err) { logger.warn('Dock guard command failed', roverId, err.message); } } function handleIdleUndock(undockedRecord) { if (!undockedRecord || undockedRecord.drivers.size > 0) return; const now = Date.now(); const { getRecentDriveActivity, setDriveCooldown, issueCommand } = require('./commandService'); const candidates = getRecentDriveActivity(DOCK_GUARD_WINDOW_MS, { excludeAdmins: true }) .filter((candidate) => candidate.roverId !== undockedRecord.id); if (candidates.length === 0) return; const candidatesWithBump = candidates.filter((candidate) => { const record = rovers.get(candidate.roverId); return record?.lastBumpAt && now - record.lastBumpAt <= DOCK_GUARD_WINDOW_MS; }); const pool = candidatesWithBump.length > 0 ? candidatesWithBump : candidates; pool.sort((a, b) => b.ts - a.ts); const suspect = pool[0]; if (!suspect) return; const suspectRecord = rovers.get(suspect.roverId); if (!suspectRecord) return; const bumpRecent = suspectRecord.lastBumpAt && now - suspectRecord.lastBumpAt <= DOCK_GUARD_WINDOW_MS; sendAlert({ color: ALERT_COLOR, title: 'Dock protection', message: `${undockedRecord.id} undocked while idle; stopping ${suspect.roverId}.`, }); try { issueCommand(suspect.roverId, { type: 'drive', driveDirect: { left: 0, right: 0 } }); issueCommand(suspect.roverId, { type: 'motors', motorPwm: { main: 0, side: 0, vacuum: 0 } }); } catch (err) { logger.warn('Dock protection stop failed', suspect.roverId, err.message); } setDriveCooldown(suspect.roverId, DOCK_GUARD_WINDOW_MS); if (bumpRecent) { nudgeRover(suspect.roverId, 'backward'); } else { nudgeRover(suspect.roverId, 'forward'); } } function nudgeRover(roverId, direction = 'backward') { if (!roverId) return; const { issueCommand } = require('./commandService'); clearTimeout(backoffTimers.get(roverId)); const speed = direction === 'forward' ? BACKOFF_SPEED : -BACKOFF_SPEED; try { issueCommand(roverId, { type: 'drive', driveDirect: { left: speed, right: speed }, }); } catch (err) { logger.warn('Dock protection nudge failed', roverId, err.message); return; } backoffTimers.set( roverId, setTimeout(() => { try { issueCommand(roverId, { type: 'drive', driveDirect: { left: 0, right: 0 } }); } catch (err) { logger.warn('Dock protection backoff stop failed', roverId, err.message); } }, BACKOFF_MS), ); } function removeSocket(socket) { const joined = socketToRovers.get(socket.id); if (!joined) { disableSpectator(socket); return; } for (const roverId of joined) { const record = rovers.get(roverId); if (record) { record.drivers.delete(socket.id); } turnService.driverRemoved(roverId, socket.id); managerEvents.emit('driver', { socketId: socket.id, roverId, action: 'remove' }); } socketToRovers.delete(socket.id); disableSpectator(socket); } function requestControl(roverId, socket, options = {}) { const { force = false, allowUser = false } = options; const record = rovers.get(roverId); if (!record) { throw new Error('Unknown rover'); } if (!allowUser && !isAdmin(socket)) { throw new Error('Only admins can request control'); } if (record.locked && !isAdmin(socket)) { throw new Error('Rover locked'); } const mode = getMode(); if (!allowUser && mode === MODES.ADMIN && !isAdmin(socket)) { throw new Error('Admins only'); } if (!allowUser && mode === MODES.LOCKDOWN && !isAdmin(socket)) { throw new Error('Server in lockdown'); } record.drivers.add(socket.id); if (!socketToRovers.has(socket.id)) { socketToRovers.set(socket.id, new Set()); } socketToRovers.get(socket.id).add(roverId); socket.join(record.room); turnService.driverAdded(roverId, socket.id, force && isAdmin(socket)); socket.emit('controlGranted', { roverId }); managerEvents.emit('driver', { socketId: socket.id, roverId, action: 'add' }); sendAlert({ color: ALERT_COLOR, title: 'Control Granted', message: `${socket.id} now driving ${roverId}`, }); return { roverId, room: record.room }; } function releaseControl(roverId, socket) { const record = rovers.get(roverId); if (!record) return; record.drivers.delete(socket.id); const joined = socketToRovers.get(socket.id); if (joined) { joined.delete(roverId); if (joined.size === 0) { socketToRovers.delete(socket.id); } } socket.leave(record.room); turnService.driverRemoved(roverId, socket.id); managerEvents.emit('driver', { socketId: socket.id, roverId, action: 'remove' }); } function isDriver(roverId, socket) { const record = rovers.get(roverId); if (!record) return false; return record.drivers.has(socket.id); } function canDrive(roverId, socket) { return turnService.canDrive(roverId, socket) || isAdmin(socket); } function getRoversForSocket(socketId) { const joined = socketToRovers.get(socketId); if (!joined || joined.size === 0) { return []; } return Array.from(joined); } function getPrimaryRoverForSocket(socketId) { const joined = socketToRovers.get(socketId); if (!joined || joined.size === 0) { return null; } const iterator = joined.values(); const first = iterator.next(); return first.done ? null : first.value; } function isDockedAndCharging(record) { const sensors = record?.lastSensor?.decoded || record?.lastSensor?.sensors || null; if (!sensors) return false; const docked = Boolean(sensors.chargingSources?.homeBase); const code = sensors.chargingState?.code; const charging = code === 2 || code === 3 || code === 4; return docked && charging; } function hasOtherDrivers(record, socketId) { if (!record) return false; for (const driver of record.drivers) { if (driver !== socketId) return true; } return false; } function canSwitchRover(socket, targetRoverId) { const currentId = getPrimaryRoverForSocket(socket.id); if (!currentId || currentId === targetRoverId) { return { ok: true, currentId }; } const currentRecord = rovers.get(currentId); if (!currentRecord) { return { ok: true, currentId }; } if (hasOtherDrivers(currentRecord, socket.id)) { return { ok: true, currentId }; } if (isDockedAndCharging(currentRecord)) { return { ok: true, currentId }; } return { ok: false, currentId, message: 'Dock and charge your current rover before switching.' }; } module.exports = { upsertRover, removeRover, lockRover, getRoster, broadcastRoster, handleSensorFrame, requestControl, releaseControl, removeSocket, isDriver, canDrive, enableSpectator, disableSpectator, rovers, managerEvents, getRoversForSocket, getPrimaryRoverForSocket, }; roleEvents.on('change', ({ socket, role }) => { if (role === 'spectator') { enableSpectator(socket); } else { disableSpectator(socket); } }); io.on('connection', (socket) => { socket.emit('rovers', getRoster()); if (socket.data?.role === 'spectator') { enableSpectator(socket); } function handleRequestControl({ roverId, force } = {}, cb = () => {}) { try { if (socket.data?.role === 'spectator') { throw new Error('Spectators cannot drive'); } const targetId = roverId || Array.from(rovers.keys())[0]; if (!targetId) { throw new Error('No rovers available'); } const previousJoined = getRoversForSocket(socket.id); if (!isAdmin(socket)) { const { ok, message } = canSwitchRover(socket, targetId); if (!ok) { throw new Error(message || 'Switch denied'); } } logger.info('Request control', socket.id, targetId, { force }); requestControl(targetId, socket, { force: Boolean(force), allowUser: true }); previousJoined.forEach((rid) => { if (rid !== targetId) { releaseControl(rid, socket); } }); videoSessions.revokeWhere( (info) => info.socketId === socket.id && info.sourceType === 'rover' && info.sourceId !== targetId && info.sourceId !== `${targetId}-audio`, ); managerEvents.emit('switch', { socketId: socket.id, roverId: targetId }); socket.emit('controlGranted', { roverId: targetId }); cb({ success: true, roverId: targetId }); } catch (err) { logger.warn('Request control failed', socket.id, err.message); sendAlert({ color: ALERT_COLOR, title: 'Control denied', message: err.message }); cb({ error: err.message }); } } function handleReleaseControl({ roverId } = {}, cb = () => {}) { if (!roverId) { cb({ error: 'roverId required' }); return; } logger.info('Release control', socket.id, roverId); releaseControl(roverId, socket); cb({ success: true, roverId }); } function handleLockToggle({ roverId, locked } = {}, cb = () => {}) { if (!isAdmin(socket)) { cb({ error: 'Not authorized' }); return; } try { lockRover(roverId, locked, { reason: 'manual' }); logger.info('Lock state changed', roverId, locked); cb({ success: true }); } catch (err) { logger.warn('Lock change failed', roverId, err.message); sendAlert({ color: ALERT_COLOR, title: 'Lock failed', message: err.message }); cb({ error: err.message }); } } function handleSubscribeAll(_, cb = () => {}) { if (socket.data?.role !== 'spectator') { cb({ error: 'Spectator role required' }); return; } if (getMode() === MODES.LOCKDOWN) { cb({ error: 'Spectating disabled in lockdown' }); return; } logger.info('Spectator subscribing to all rovers', socket.id); for (const record of rovers.values()) { socket.join(record.room); } cb({ success: true }); } socket.on('requestControl', handleRequestControl); socket.on('session:requestControl', handleRequestControl); socket.on('releaseControl', handleReleaseControl); socket.on('session:releaseControl', handleReleaseControl); socket.on('lockRover', handleLockToggle); socket.on('session:lockRover', handleLockToggle); socket.on('subscribeAll', handleSubscribeAll); socket.on('session:subscribeAll', handleSubscribeAll); socket.on('disconnecting', () => { logger.info('Socket disconnecting', socket.id); removeSocket(socket); }); socket.on('disconnect', () => { logger.info('Socket disconnected', socket.id); removeSocket(socket); }); }); function enableSpectator(socket) { if (!socket?.id || spectatorSockets.has(socket.id)) return; spectatorSockets.add(socket.id); for (const record of rovers.values()) { socket.join(record.room); } } function disableSpectator(socket) { if (!socket?.id || !spectatorSockets.has(socket.id)) return; spectatorSockets.delete(socket.id); for (const record of rovers.values()) { socket.leave(record.room); } }