diff --git a/dist/dummy1.yml b/dist/dummy1.yml index 568c6668..3ffd50e8 100644 --- a/dist/dummy1.yml +++ b/dist/dummy1.yml @@ -14,9 +14,8 @@ battery: urgent: 1650 maxWheelSpeed: 350 media: - whepUrl: https://mediaserver.local/whep/roomba-alpha - streamKey: roomba-alpha + publishUrl: https://control-server.local/whip/roomba-alpha manage: false service: mediamtx.service - healthUrl: http://127.0.0.1:9997/v3/paths/list/rovercam - healthInterval: 30s \ No newline at end of file + healthUrl: http://127.0.0.1:9997/v3/paths/list + healthInterval: 30s diff --git a/dist/dummy2.yml b/dist/dummy2.yml index a3b559f7..89d64536 100644 --- a/dist/dummy2.yml +++ b/dist/dummy2.yml @@ -14,9 +14,8 @@ battery: urgent: 1650 maxWheelSpeed: 350 media: - whepUrl: https://mediaserver.local/whep/roomba-alpha - streamKey: roomba-alpha + publishUrl: https://control-server.local/whip/roomba-alpha manage: false service: mediamtx.service - healthUrl: http://127.0.0.1:9997/v3/paths/list/rovercam - healthInterval: 30s \ No newline at end of file + healthUrl: http://127.0.0.1:9997/v3/paths/list + healthInterval: 30s diff --git a/dist/roverd b/dist/roverd index 65406101..d34196db 100755 Binary files a/dist/roverd and b/dist/roverd differ diff --git a/docs/pi-deployment.md b/docs/pi-deployment.md index 93f7ab29..6b3398a7 100644 --- a/docs/pi-deployment.md +++ b/docs/pi-deployment.md @@ -55,7 +55,7 @@ Flags: | `--mediamtx` | download/install mediaMTX plus the provided config + unit | | `--mediamtx-version X.Y.Z` | override the mediaMTX release tag (default `1.15.3`) | -If the script installs the sample config, it will remind you to edit `/etc/roverd.yaml` before manually restarting the service: set `name`, `serverUrl`, serial device, BRC pin, battery thresholds, and the WHEP URL that the central server should expose. +If the script installs the sample config, it will remind you to edit `/etc/roverd.yaml` before manually restarting the service: set `name`, `serverUrl`, serial device, BRC pin, battery thresholds, and the `media.publishUrl` that points at your central mediaMTX WHIP endpoint (for example `https://control-server.local/whip/roomba-alpha`). ## Manual installation @@ -65,7 +65,7 @@ If the script installs the sample config, it will remind you to edit `/etc/rover sudo install -o roverd -g roverd -m 0755 dist/roverd /usr/local/bin/roverd sudo install -o roverd -g roverd -m 0640 pi/roverd/roverd.sample.yaml /etc/roverd.yaml ``` - Adjust `/etc/roverd.yaml` for each rover: `name`, `serverUrl` (e.g. `ws://control-server:8080/rover`), serial port path, battery thresholds, GPIO pin for BRC, and the media WHEP URL that points at the central distribution server. + Adjust `/etc/roverd.yaml` for each rover: `name`, `serverUrl` (e.g. `ws://control-server:8080/rover`), serial port path, battery thresholds, GPIO pin for BRC, and the media `publishUrl` that points at the central media server’s WHIP endpoint for that rover. 2. Install the systemd unit: ```bash @@ -92,8 +92,9 @@ If you set `media.manage: true` in `/etc/roverd.yaml`, make sure the `roverd` se sudo systemctl enable --now mediamtx.service ``` -The sample config uses the Raspberry Pi camera module as the source and exposes a WHEP endpoint at `http://:8889/whep/rovercam`. Point `media.whepUrl` in `roverd.yaml` at this URL so the central server can list it. -Expose the mediaMTX HTTP API locally (default `http://127.0.0.1:9997`) and set `media.healthUrl` so `roverd` can monitor the pipeline; `media.service` should match the systemd unit name (default `mediamtx.service`). +The sample config uses the Raspberry Pi camera module as the source and enables the local mediaMTX HTTP API so `roverd` can health-check the service. Edit `/etc/mediamtx/mediamtx.yml` to point its WHIP client at your central media server (see `mediamtx_server_integration.md` for details). +Expose the mediaMTX HTTP API locally (default `http://127.0.0.1:9997`) and set `media.healthUrl` so `roverd` can monitor the pipeline; `media.service` should match the systemd unit name (default `mediamtx.service`). +Finally, edit `/etc/mediamtx/mediamtx.yml` so the Pi publishes the camera feed to your central control server via WHIP (for example by pointing it at `https://control.example.com/whip/`). The Pi never serves viewers directly—drivers and spectators always watch through the control server. ## Server + UI diff --git a/mediamtx_server_integration.md b/mediamtx_server_integration.md new file mode 100644 index 00000000..0045aa91 --- /dev/null +++ b/mediamtx_server_integration.md @@ -0,0 +1,19 @@ +# mediaMTX Integration + +Each rover runs mediaMTX locally to capture the Pi camera and publish it upstream, while the control server hosts a central mediaMTX instance that fans video out to drivers and spectators. + +## Pi (publisher) + +- mediaMTX samples the Pi camera (`paths.rovercam.source: rpiCamera`) and exposes the HTTP API on `http://127.0.0.1:9997`. +- `/etc/roverd.yaml` contains `media.publishUrl`, pointing at the control server’s WHIP endpoint for that rover (e.g. `https://control.example.com/whip/roomba-alpha`). No auth is required when the Pi network is trusted. +- The `media.manage` flag keeps the local service alive via `systemctl` and hits the API for health checks (`media.healthUrl`, defaults to `http://127.0.0.1:9997/v3/paths/list`). + +## Control server (viewer) + +- The central mediaMTX instance accepts WHIP ingest at `/whip/` and serves WHEP playback at `/whep/`. +- The Node server issues viewer sessions (one per socket) and mediamtx calls back into the Node server to authorize new WHEP connections. Lockdown mode simply stops minting sessions for non-lockdown admins. + +## Driver / spectator UIs + +- When a user requests video the browser asks the Node server for a session token, and if the current mode/role allows it, the server returns the WHEP URL plus the session id. +- Spectator mode and lockdown are enforced purely on the Node server: no direct Pi URLs are ever exposed, and once a user has a session it remains valid until they disconnect or lockdown revokes it. diff --git a/pi/mediamtx/mediamtx.yml b/pi/mediamtx/mediamtx.yml index 4bbcd4d5..9b741517 100644 --- a/pi/mediamtx/mediamtx.yml +++ b/pi/mediamtx/mediamtx.yml @@ -1,4 +1,7 @@ logLevel: info +api: yes +apiAddress: 127.0.0.1:9997 +apiEncryption: no webrtc: yes webrtcLocalUDPAddress: :8189 webrtcLocalTCPAddress: :8189 diff --git a/pi/roverd/config.go b/pi/roverd/config.go index 1a236b08..c0e3d4bb 100644 --- a/pi/roverd/config.go +++ b/pi/roverd/config.go @@ -53,8 +53,7 @@ type BatteryConfig struct { } type MediaConfig struct { - WhepURL string `yaml:"whepUrl"` - StreamKey string `yaml:"streamKey"` + PublishURL string `yaml:"publishUrl"` Manage bool `yaml:"manage"` Service string `yaml:"service"` HealthURL string `yaml:"healthUrl"` diff --git a/pi/roverd/media_supervisor.go b/pi/roverd/media_supervisor.go index bd14d1c3..d70c5d86 100644 --- a/pi/roverd/media_supervisor.go +++ b/pi/roverd/media_supervisor.go @@ -22,6 +22,9 @@ func NewMediaSupervisor(cfg MediaConfig, logger *log.Logger) *MediaSupervisor { if !cfg.Manage || cfg.Service == "" { return nil } + if cfg.HealthURL == "" { + cfg.HealthURL = "http://127.0.0.1:9997/v3/paths/list" + } interval := cfg.HealthInterval.Duration if interval <= 0 { interval = 30 * time.Second diff --git a/pi/roverd/roverd b/pi/roverd/roverd new file mode 100755 index 00000000..8fbabf72 Binary files /dev/null and b/pi/roverd/roverd differ diff --git a/pi/roverd/roverd.sample.yaml b/pi/roverd/roverd.sample.yaml index 975b5105..958e8988 100644 --- a/pi/roverd/roverd.sample.yaml +++ b/pi/roverd/roverd.sample.yaml @@ -15,9 +15,8 @@ battery: urgent: 1650 maxWheelSpeed: 350 media: - whepUrl: https://mediaserver.local/whep/roomba-alpha - streamKey: roomba-alpha + publishUrl: https://control-server.local/whip/roomba-alpha manage: false service: mediamtx.service - healthUrl: http://127.0.0.1:9997/v3/paths/list/rovercam + healthUrl: http://127.0.0.1:9997/v3/paths/list healthInterval: 30s diff --git a/pi/roverd/roverd.yaml b/pi/roverd/roverd.yaml index 975b5105..958e8988 100644 --- a/pi/roverd/roverd.yaml +++ b/pi/roverd/roverd.yaml @@ -15,9 +15,8 @@ battery: urgent: 1650 maxWheelSpeed: 350 media: - whepUrl: https://mediaserver.local/whep/roomba-alpha - streamKey: roomba-alpha + publishUrl: https://control-server.local/whip/roomba-alpha manage: false service: mediamtx.service - healthUrl: http://127.0.0.1:9997/v3/paths/list/rovercam + healthUrl: http://127.0.0.1:9997/v3/paths/list healthInterval: 30s diff --git a/server/index.js b/server/index.js index 77cf4c28..b32f11c5 100644 --- a/server/index.js +++ b/server/index.js @@ -13,4 +13,5 @@ require('./src/services/lockdownGuard'); require('./src/services/roverManager'); require('./src/services/commandService'); require('./src/services/roverConnectionService'); +require('./src/services/assignmentService'); require('./src/services/httpServer'); diff --git a/server/public/index.html b/server/public/index.html index 9c4f64c8..5027ad5b 100644 --- a/server/public/index.html +++ b/server/public/index.html @@ -24,17 +24,25 @@ - Role: user
-
- +
+
+
@@ -51,7 +59,8 @@

Sensor Frames

-
Active driver: --
+
Active rover: --
+
Driver: --

         

Use WASD keys to drive the selected rover. Shift increases speed.

@@ -67,5 +76,6 @@ +

Open spectator view

diff --git a/server/public/spectate/index.html b/server/public/spectate/index.html new file mode 100644 index 00000000..c0601dd9 --- /dev/null +++ b/server/public/spectate/index.html @@ -0,0 +1,29 @@ + + + + + Rover Spectator + + + +

Rover Spectator

+
Connecting…
+
+ + + + + + + + diff --git a/server/public/spectate/spectator.js b/server/public/spectate/spectator.js new file mode 100644 index 00000000..a26b1c3d --- /dev/null +++ b/server/public/spectate/spectator.js @@ -0,0 +1,63 @@ +registerModule('spectator/main', (require, exports) => { + const { socket } = require('globals/socket'); + + const statusEl = document.getElementById('status'); + const grid = document.getElementById('grid'); + const cards = new Map(); + + socket.on('connect', () => { + statusEl.textContent = 'Connected'; + }); + + socket.on('disconnect', () => { + statusEl.textContent = 'Disconnected'; + }); + + socket.on('lockdown', ({ message }) => { + statusEl.textContent = message || 'Spectator mode disabled'; + }); + + socket.on('rovers', (list) => { + const seen = new Set(); + list.forEach((rover) => { + const card = ensureCard(rover.id); + card.querySelector('.meta').textContent = `Locked: ${rover.locked ? 'yes' : 'no'} | Last seen: ${new Date(rover.lastSeen).toLocaleTimeString()}`; + }); + cards.forEach((_, id) => { + if (!list.find((r) => r.id === id)) { + removeCard(id); + } + }); + }); + + socket.on('sensorFrame', ({ roverId, sensors }) => { + const card = ensureCard(roverId); + card.querySelector('.sensor').textContent = sensors ? JSON.stringify(sensors, null, 2) : 'No data'; + }); + + function ensureCard(roverId) { + if (cards.has(roverId)) { + return cards.get(roverId); + } + const card = document.createElement('div'); + card.className = 'rover-card'; + card.innerHTML = ` +

${roverId}

+
+
Waiting for data...
+ `; + grid.appendChild(card); + cards.set(roverId, card); + return card; + } + + function removeCard(roverId) { + const card = cards.get(roverId); + if (card) { + card.remove(); + cards.delete(roverId); + } + } +}); + +requireModule('spectator/main'); diff --git a/server/public/src/globals/socket.js b/server/public/src/globals/socket.js index bbe01878..a3f50f70 100644 --- a/server/public/src/globals/socket.js +++ b/server/public/src/globals/socket.js @@ -1,4 +1,8 @@ registerModule('globals/socket', (require, exports) => { - const socket = io(); + const query = {}; + if (window.__desiredRole) { + query.role = window.__desiredRole; + } + const socket = io(query.role ? { query } : undefined); exports.socket = socket; }); diff --git a/server/public/src/services/authControls.js b/server/public/src/services/authControls.js index e4d3c680..b3f730f0 100644 --- a/server/public/src/services/authControls.js +++ b/server/public/src/services/authControls.js @@ -6,7 +6,6 @@ registerModule('services/authControls', (require, exports) => { const passwordInput = document.getElementById('password'); const loginBtn = document.getElementById('loginBtn'); const roleStatus = document.getElementById('roleStatus'); - const spectatorBtn = document.getElementById('spectatorBtn'); function updateRole(role) { if (roleStatus) { @@ -33,9 +32,4 @@ registerModule('services/authControls', (require, exports) => { }); // initialize from default role updateRole(state.getRole()); - - spectatorBtn?.addEventListener('click', () => { - socket.emit('role:set', { role: 'spectator' }); - updateRole('spectator'); - }); }); diff --git a/server/public/src/services/roverUI.js b/server/public/src/services/roverUI.js index 0f3467e0..f326c5f5 100644 --- a/server/public/src/services/roverUI.js +++ b/server/public/src/services/roverUI.js @@ -6,16 +6,28 @@ registerModule('services/roverUI', (require, exports) => { const statusEl = document.getElementById('status'); const roverSelect = document.getElementById('roverSelect'); const sensorOutput = document.getElementById('sensorOutput'); - const requestBtn = document.getElementById('requestControl'); const lockBtn = document.getElementById('lockToggle'); + const forceCheckbox = document.getElementById('forceControl'); + const adminControls = document.getElementById('adminControls'); + const modeControls = document.getElementById('modeControls'); + const modeSelect = document.getElementById('modeSelect'); const activeDriverEl = document.getElementById('activeDriver'); + const currentRoverEl = document.getElementById('currentRover'); let roster = []; + let lastSelection = null; state.onRoleChange((role) => { const isAdmin = role === 'admin' || role === 'lockdown'; - if (requestBtn) requestBtn.disabled = !isAdmin; if (lockBtn) lockBtn.disabled = !isAdmin; + if (roverSelect) roverSelect.disabled = !isAdmin; + if (adminControls) { + adminControls.style.display = isAdmin ? 'flex' : 'none'; + } + if (modeControls) { + modeControls.style.display = isAdmin ? 'block' : 'none'; + if (modeSelect) modeSelect.disabled = !isAdmin; + } }); socket.on('rovers', (list) => { @@ -27,19 +39,28 @@ registerModule('services/roverUI', (require, exports) => { option.textContent = `${rover.name}${rover.locked ? ' (locked)' : ''}`; roverSelect.appendChild(option); }); - if (!state.getSelected() && list.length) { + const selected = state.getSelected(); + if ((!selected || !list.find((r) => r.id === selected)) && list.length) { state.setSelected(list[0].id); - roverSelect.value = list[0].id; + } else if (!list.length) { + state.setSelected(null); } + renderSelection(); }); roverSelect.addEventListener('change', () => { - state.setSelected(roverSelect.value); + const roverId = roverSelect.value; + const role = state.getRole(); + state.setSelected(roverId); + renderSelection(true); + if (role === 'admin' || role === 'lockdown') { + socket.emit('requestControl', { roverId, force: !!forceCheckbox?.checked }); + } }); socket.on('controlGranted', ({ roverId }) => { state.setSelected(roverId); - roverSelect.value = roverId; + renderSelection(true); statusEl.textContent = `Driving ${roverId}`; }); @@ -62,10 +83,6 @@ registerModule('services/roverUI', (require, exports) => { } }); - requestBtn?.addEventListener('click', () => { - socket.emit('requestControl', { roverId: state.getSelected() }); - }); - lockBtn?.addEventListener('click', () => { const roverId = state.getSelected(); const rover = roster.find((r) => r.id === roverId); @@ -73,10 +90,36 @@ registerModule('services/roverUI', (require, exports) => { socket.emit('lockRover', { roverId, locked: !rover.locked }); }); + modeSelect?.addEventListener('change', () => { + socket.emit('setMode', { mode: modeSelect.value }); + }); + + socket.on('mode', ({ mode }) => { + if (modeSelect) { + modeSelect.value = mode; + } + }); + socket.on('connect', () => { statusEl.textContent = 'Connected'; }); socket.on('disconnect', () => { statusEl.textContent = 'Disconnected'; }); + + function renderSelection(clearOutput) { + const roverId = state.getSelected(); + if (roverSelect && roverId) { + roverSelect.value = roverId; + } else if (roverSelect && !roverId) { + roverSelect.value = ''; + } + if (currentRoverEl) { + currentRoverEl.textContent = `Active rover: ${roverId || '--'}`; + } + if (sensorOutput && (clearOutput || roverId !== lastSelection)) { + sensorOutput.textContent = roverId ? 'Waiting for sensor data...' : ''; + } + lastSelection = roverId; + } }); diff --git a/server/src/globals/logger.js b/server/src/globals/logger.js index 137960b0..3f9ffe12 100644 --- a/server/src/globals/logger.js +++ b/server/src/globals/logger.js @@ -1,9 +1,20 @@ -function stamp(level, args) { - return [new Date().toISOString(), `[${level}]`, ...args]; +function stamp(level, label, args) { + const fields = [new Date().toISOString(), `[${level}]`]; + if (label) { + fields.push(`[${label}]`); + } + return [...fields, ...args]; +} + +function baseLogger(label) { + return { + info: (...args) => console.log(...stamp('INFO', label, args)), + warn: (...args) => console.warn(...stamp('WARN', label, args)), + error: (...args) => console.error(...stamp('ERROR', label, args)), + }; } module.exports = { - info: (...args) => console.log(...stamp('INFO', args)), - warn: (...args) => console.warn(...stamp('WARN', args)), - error: (...args) => console.error(...stamp('ERROR', args)), + ...baseLogger(), + child: (label) => baseLogger(label), }; diff --git a/server/src/services/assignmentService.js b/server/src/services/assignmentService.js new file mode 100644 index 00000000..d7436851 --- /dev/null +++ b/server/src/services/assignmentService.js @@ -0,0 +1,135 @@ +const io = require('../globals/io'); +const logger = require('../globals/logger').child('assignment'); +const { MODES, getMode, modeEvents } = require('./modeManager'); +const { roleEvents, getRole, isAdmin } = require('./roleService'); +const roverManager = require('./roverManager'); + +const socketRefs = new Map(); // socketId -> socket +const assignments = new Map(); // socketId -> roverId +const waiting = new Set(); // socketIds waiting for placement + +io.on('connection', (socket) => { + socketRefs.set(socket.id, socket); + socket.on('disconnect', () => { + socketRefs.delete(socket.id); + unassignSocket(socket); + }); +}); + +roleEvents.on('change', ({ socket, role }) => { + if (!socket || !socket.id) return; + if (role === 'user') { + assignSocket(socket); + } else { + unassignSocket(socket); + } +}); + +modeEvents.on('change', (mode) => { + if (mode === MODES.ADMIN || mode === MODES.LOCKDOWN) { + // release non-admin drivers + for (const [socketId, roverId] of assignments.entries()) { + const socket = socketRefs.get(socketId); + if (socket && !isAdmin(socket)) { + releaseAssignment(socket, roverId); + } + } + } + reassignWaiting(); +}); + +roverManager.managerEvents.on('lock', ({ roverId, locked }) => { + if (locked) { + reassignFromRover(roverId); + } else { + reassignWaiting(); + } +}); + +roverManager.managerEvents.on('rover', ({ action }) => { + if (action === 'removed' || action === 'upsert') { + reassignWaiting(); + } +}); + +function assignSocket(socket) { + if (!socket || isAdmin(socket) || getRole(socket) !== 'user') { + return; + } + // avoid double assignment + if (assignments.has(socket.id)) { + return; + } + const target = pickRover(); + if (!target) { + waiting.add(socket.id); + logger.info('No rover available, user waiting', socket.id); + return; + } + try { + roverManager.requestControl(target.id, socket, { allowUser: true }); + assignments.set(socket.id, target.id); + waiting.delete(socket.id); + logger.info('Assigned user to rover', socket.id, target.id); + } catch (err) { + logger.warn('Failed to assign user', err.message); + waiting.add(socket.id); + } +} + +function unassignSocket(socket) { + if (!socket) return; + waiting.delete(socket.id); + const roverId = assignments.get(socket.id); + if (roverId) { + roverManager.releaseControl(roverId, socket); + assignments.delete(socket.id); + } +} + +function reassignFromRover(roverId) { + for (const [socketId, rid] of assignments.entries()) { + if (rid !== roverId) continue; + const socket = socketRefs.get(socketId); + if (!socket) { + assignments.delete(socketId); + continue; + } + roverManager.releaseControl(rid, socket); + assignments.delete(socketId); + assignSocket(socket); + } +} + +function reassignWaiting() { + for (const socketId of Array.from(waiting)) { + const socket = socketRefs.get(socketId); + if (socket) { + assignSocket(socket); + } else { + waiting.delete(socketId); + } + } +} + +function releaseAssignment(socket, roverId) { + roverManager.releaseControl(roverId, socket); + assignments.delete(socket.id); + waiting.add(socket.id); +} + +function pickRover() { + const mode = getMode(); + if (mode === MODES.ADMIN || mode === MODES.LOCKDOWN) { + return null; + } + const candidates = Array.from(roverManager.rovers.values()).filter((rover) => { + if (!rover || rover.locked) return false; + return true; + }); + if (candidates.length === 0) { + return null; + } + candidates.sort((a, b) => a.drivers.size - b.drivers.size); + return candidates[0]; +} diff --git a/server/src/services/authService.js b/server/src/services/authService.js index 854b005d..98e0677f 100644 --- a/server/src/services/authService.js +++ b/server/src/services/authService.js @@ -1,8 +1,9 @@ const bcrypt = require('bcrypt'); const io = require('../globals/io'); +const logger = require('../globals/logger').child('authService'); const { loadConfig } = require('../helpers/configLoader'); const { clearLockdownTimer } = require('./lockdownGuard'); -const { setRole, roleEvents } = require('./roleService'); +const { setRole } = require('./roleService'); const config = loadConfig(); const admins = config.admins || []; @@ -32,8 +33,10 @@ function isLockdownAdmin(socket) { } io.on('connection', (socket) => { - setRole(socket, 'user'); - socket.emit('auth:role', { role: 'user' }); + const requestedRole = socket.handshake?.query?.role; + const initialRole = requestedRole === 'spectator' ? 'spectator' : 'user'; + setRole(socket, initialRole); + socket.emit('auth:role', { role: initialRole }); socket.on('auth:login', async ({ username, password }, cb = () => {}) => { try { const admin = await authenticate(username, password); diff --git a/server/src/services/httpServer.js b/server/src/services/httpServer.js index 88d1a1bf..c98251ab 100644 --- a/server/src/services/httpServer.js +++ b/server/src/services/httpServer.js @@ -1,6 +1,6 @@ const { httpServer } = require('../globals/http'); const config = require('../globals/config'); -const logger = require('../globals/logger'); +const logger = require('../globals/logger').child('httpServer'); httpServer.listen(config.port, () => { logger.info(`Server listening on :${config.port}`); diff --git a/server/src/services/modeManager.js b/server/src/services/modeManager.js index 1a8876da..c45aa352 100644 --- a/server/src/services/modeManager.js +++ b/server/src/services/modeManager.js @@ -37,6 +37,7 @@ function setMode(nextMode, socket) { message: `Server mode set to ${nextMode}`, }); modeEvents.emit('change', currentMode); + io.emit('mode', { mode: currentMode }); return currentMode; } @@ -52,6 +53,7 @@ module.exports = { }; io.on('connection', (socket) => { + socket.emit('mode', { mode: currentMode }); socket.on('setMode', ({ mode }) => { try { setMode(mode, socket); diff --git a/server/src/services/roverConnectionService.js b/server/src/services/roverConnectionService.js index 0562c4e5..665da354 100644 --- a/server/src/services/roverConnectionService.js +++ b/server/src/services/roverConnectionService.js @@ -1,5 +1,5 @@ const roverWSS = require('../globals/ws'); -const logger = require('../globals/logger'); +const logger = require('../globals/logger').child('roverConnection'); const roverManager = require('./roverManager'); const { sendAlert, COLORS } = require('./alertService'); const { handleAck } = require('./commandService'); diff --git a/server/src/services/roverManager.js b/server/src/services/roverManager.js index 756944f8..fff5a215 100644 --- a/server/src/services/roverManager.js +++ b/server/src/services/roverManager.js @@ -1,5 +1,6 @@ +const EventEmitter = require('events'); const io = require('../globals/io'); -const logger = require('../globals/logger'); +const logger = require('../globals/logger').child('roverManager'); const { sendAlert, COLORS } = require('./alertService'); const { parseSensorFrame } = require('../helpers/sensorDecoder'); const { MODES, getMode } = require('./modeManager'); @@ -9,6 +10,7 @@ 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(); function ensureRecord(id) { if (!rovers.has(id)) { @@ -37,6 +39,7 @@ function upsertRover(meta, ws) { const sock = io.sockets.sockets.get(socketId); sock?.join(record.room); }); + managerEvents.emit('rover', { roverId: id, action: 'upsert', record }); broadcastRoster(); return record; } @@ -51,6 +54,7 @@ function removeRover(id) { sock?.leave(record.room); }); broadcastRoster(); + managerEvents.emit('rover', { roverId: id, action: 'removed' }); } function lockRover(id, locked, actorSocket) { @@ -66,6 +70,7 @@ function lockRover(id, locked, actorSocket) { sendAlert({ color: COLORS.success, title: 'Rover Unlocked', message: `${id} unlocked.` }); } broadcastRoster(); + managerEvents.emit('lock', { roverId: id, locked: record.locked }); return record.locked; } @@ -115,22 +120,23 @@ function removeSocket(socket) { disableSpectator(socket); } -function requestControl(roverId, 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 (!isAdmin(socket)) { + if (!allowUser && !isAdmin(socket)) { throw new Error('Only admins can request control'); } - if (record.locked && !isAdmin(socket)) { + if (record.locked && !isAdmin(socket) && !allowUser) { throw new Error('Rover locked'); } const mode = getMode(); - if (mode === MODES.ADMIN && !isAdmin(socket)) { + if (!allowUser && mode === MODES.ADMIN && !isAdmin(socket)) { throw new Error('Admins only'); } - if (mode === MODES.LOCKDOWN && !isAdmin(socket)) { + if (!allowUser && mode === MODES.LOCKDOWN && !isAdmin(socket)) { throw new Error('Server in lockdown'); } record.drivers.add(socket.id); @@ -139,7 +145,8 @@ function requestControl(roverId, socket) { } socketToRovers.get(socket.id).add(roverId); socket.join(record.room); - turnService.driverAdded(roverId, socket.id); + turnService.driverAdded(roverId, socket.id, force && isAdmin(socket)); + socket.emit('controlGranted', { roverId }); sendAlert({ color: COLORS.success, title: 'Control Granted', @@ -188,6 +195,7 @@ module.exports = { enableSpectator, disableSpectator, rovers, + managerEvents, }; roleEvents.on('change', ({ socket, role }) => { @@ -204,13 +212,13 @@ io.on('connection', (socket) => { enableSpectator(socket); } - socket.on('requestControl', ({ roverId } = {}) => { + socket.on('requestControl', ({ roverId, force } = {}) => { try { const targetId = roverId || Array.from(rovers.keys())[0]; if (!targetId) { throw new Error('No rovers available'); } - requestControl(targetId, socket); + requestControl(targetId, socket, { force: Boolean(force) }); socket.emit('controlGranted', { roverId: targetId }); } catch (err) { sendAlert({ color: COLORS.warn, title: 'Control denied', message: err.message }); diff --git a/server/src/services/turnService.js b/server/src/services/turnService.js index d8cb5d32..ff08659c 100644 --- a/server/src/services/turnService.js +++ b/server/src/services/turnService.js @@ -6,12 +6,12 @@ const driverQueues = new Map(); // roverId -> { queue: [], current: socketId, ti const activeDrivers = new Map(); const TURN_DURATION_MS = 60 * 1000; -function driverAdded(roverId, socketId) { +function driverAdded(roverId, socketId, force) { const queue = ensureQueue(roverId); if (!queue.queue.includes(socketId)) { queue.queue.push(socketId); } - if (!queue.current) { + if (!queue.current || force) { queue.current = socketId; } syncState(roverId);