diff --git a/rulesdocs/refector_rules_and_tracking.md b/rulesdocs/refector_rules_and_tracking.md
index 4008df0b..08f1827b 100644
--- a/rulesdocs/refector_rules_and_tracking.md
+++ b/rulesdocs/refector_rules_and_tracking.md
@@ -48,7 +48,7 @@
- [x] home assistant service
- [ ] llm commentary service
- [x] private rover access request service
-- [ ] replay services (already partly split; reformat consistently)
+- [x] replay services (consolidated under replayEngineV2)
- [ ] room camera services (already partly split; reformat consistently)
- [ ] rover manager service (in progress: constants/state extracted)
- [x] session service
@@ -86,6 +86,8 @@
- Finished `privateRoverAccessRequestService` decomposition by extracting in-memory maps/events/constants to `privateRoverAccessRequestService/state.js`, shared keying/lookup helpers to `privateRoverAccessRequestService/helpers.js`, request/grant business logic to `privateRoverAccessRequestService/core.js`, and rover/socket event wiring to `privateRoverAccessRequestService/hooks.js`; `privateRoverAccessRequestService/index.js` is now a thin composition layer.
- Finished `homeAssistantService` decomposition by extracting shared runtime caches/constants to `homeAssistantService/state.js`, entity/trigger normalization helpers to `homeAssistantService/entityHelpers.js`, automation/state engine logic to `homeAssistantService/runtimeEngine.js`, websocket transport/reconnect lifecycle to `homeAssistantService/transport.js`, and mode/turn/socket event wiring to `homeAssistantService/hooks.js`; `homeAssistantService/index.js` is now a thin composition layer.
- Finished `discordBotService` decomposition by extracting presence rotation/state to `discordBotService/presence.js`, channel/typing transport helpers to `discordBotService/channelIO.js`, command routing and admin command handlers to `discordBotService/commandHandlers.js`, and event-bus/chat-bridge/moderation DM workflows to `discordBotService/integrations.js`; `discordBotService/index.js` is now a thin composition layer.
+- Finished `replayEngineV2` decomposition by extracting environment/path constants to `replayEngineV2/constants.js`, mutable runtime state to `replayEngineV2/state.js`, source discovery/worker arg building to `replayEngineV2/sources.js`, ffmpeg worker lifecycle to `replayEngineV2/workerManager.js`, segment indexing/retention/health snapshot logic to `replayEngineV2/segmentStore.js`, sidebar SVG/video rendering to `replayEngineV2/sidebarRenderer.js`, and replay assembly pipeline to `replayEngineV2/replayBuilder.js`; `replayEngineV2/index.js` is now a thin orchestration layer.
+- Consolidated replay-related single-file services into `replayEngineV2` by moving cooldown state (`cooldown.js`), user-facing replay source validation/defaults (`replaySources.js`), and replay socket hooks (`socketHooks.js`) into the engine folder; removed obsolete standalone services `replayBuildService`, `replayService`, `replaySourceService`, and `replaySocketService` and rewired dependents to import directly from `replayEngineV2`.
## WebUI frontend
### BIGGEST OFFENDERS
diff --git a/server/src/services/discordBotService/index.js b/server/src/services/discordBotService/index.js
index f5b3b24a..3f651f47 100644
--- a/server/src/services/discordBotService/index.js
+++ b/server/src/services/discordBotService/index.js
@@ -13,11 +13,9 @@ const roverManager = require('../roverManager');
const { getRoster, lockRover, rovers } = roverManager;
const { MODES, getMode, setMode } = require('../modeManager');
const { sendExternalMessage, sendExternalTyping } = require('../chatService');
-const { buildReplayVideo } = require('../replayBuildService');
-const { getReplaySources, getDefaultDiscordSources, validateSources } = require('../replaySourceService');
+const { buildReplayVideo, getReplaySources, getDefaultDiscordSources, validateSources, tryTriggerReplay } = require('../replayEngineV2');
const { getActiveDrivers } = require('../turnService');
const { getNickname } = require('../nicknameService');
-const { tryTriggerReplay } = require('../replayService');
const { getCommunityGoal, setCommunityGoal, clearCommunityGoal } = require('../communityGoalService');
const { getAdminReason, setAdminReason, clearAdminReason } = require('../adminReasonService');
const {
diff --git a/server/src/services/replayBuildService/index.js b/server/src/services/replayBuildService/index.js
deleted file mode 100644
index c48df878..00000000
--- a/server/src/services/replayBuildService/index.js
+++ /dev/null
@@ -1,8 +0,0 @@
-// replay Build Service
-// Purpose: Defines the replay Build Service module and the helpers/state used by this service unit.
-// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary.
-const { buildReplayVideo } = require('../replayEngineV2');
-
-module.exports = {
- buildReplayVideo,
-};
diff --git a/server/src/services/replayEngineV2/constants.js b/server/src/services/replayEngineV2/constants.js
new file mode 100644
index 00000000..7f5e8908
--- /dev/null
+++ b/server/src/services/replayEngineV2/constants.js
@@ -0,0 +1,33 @@
+// Replay Engine Constants
+// Purpose: Defines runtime constants and filesystem paths used by replay worker and build pipelines.
+// Scope: Reads env-driven tuning values once and exposes immutable configuration values.
+const path = require('path');
+const { resolveDataDir } = require('../../helpers/dataPaths');
+
+const FFMPEG_BIN = process.env.FFMPEG_BIN || 'ffmpeg';
+const SEGMENT_ROOT = path.join(resolveDataDir(), 'replay-segments');
+const SEGMENT_SECONDS = Math.max(1, Number.parseInt(process.env.REPLAY_SEGMENT_SECONDS || '1', 10));
+const BUFFER_SECONDS = Math.max(20, Number.parseInt(process.env.REPLAY_BUFFER_SECONDS || '45', 10));
+const CLEANUP_INTERVAL_MS = 10_000;
+const BUILD_DURATION_MS = Math.max(5000, Number.parseInt(process.env.REPLAY_DURATION_MS || '20000', 10));
+const BUILD_GUARD_MS = Math.max(200, Number.parseInt(process.env.REPLAY_GUARD_MS || '1200', 10));
+const TARGET_FPS = Math.max(10, Number.parseInt(process.env.REPLAY_TARGET_FPS || '30', 10));
+const MAX_WIDTH = Math.max(320, Number.parseInt(process.env.REPLAY_MAX_WIDTH || '1280', 10));
+const MAX_HEIGHT = Math.max(180, Number.parseInt(process.env.REPLAY_MAX_HEIGHT || '720', 10));
+const MAX_BYTES = Math.floor(Number.parseFloat(process.env.REPLAY_MAX_OUTPUT_MB || '9.5') * 1024 * 1024);
+const SIDEBAR_WIDTH = 190;
+
+module.exports = {
+ FFMPEG_BIN,
+ SEGMENT_ROOT,
+ SEGMENT_SECONDS,
+ BUFFER_SECONDS,
+ CLEANUP_INTERVAL_MS,
+ BUILD_DURATION_MS,
+ BUILD_GUARD_MS,
+ TARGET_FPS,
+ MAX_WIDTH,
+ MAX_HEIGHT,
+ MAX_BYTES,
+ SIDEBAR_WIDTH,
+};
diff --git a/server/src/services/replayEngineV2/cooldown.js b/server/src/services/replayEngineV2/cooldown.js
new file mode 100644
index 00000000..8de688cc
--- /dev/null
+++ b/server/src/services/replayEngineV2/cooldown.js
@@ -0,0 +1,43 @@
+// Replay Cooldown State
+// Purpose: Tracks replay trigger cooldown state and emits updates for consumers.
+// Scope: Encapsulates replay cooldown timing and trigger bookkeeping.
+const EventEmitter = require('events');
+
+const COOLDOWN_MS = 10 * 1000;
+const replayEvents = new EventEmitter();
+let lastTriggeredAt = null;
+let lastTriggeredBy = null;
+
+function getRemainingMs(now = Date.now()) {
+ if (!lastTriggeredAt) return 0;
+ return Math.max(0, COOLDOWN_MS - (now - lastTriggeredAt));
+}
+
+function getReplayState() {
+ const remainingMs = getRemainingMs();
+ return {
+ cooldownMs: COOLDOWN_MS,
+ lastTriggeredAt,
+ lastTriggeredBy,
+ remainingMs,
+ available: remainingMs === 0,
+ };
+}
+
+function tryTriggerReplay(by = null) {
+ const remainingMs = getRemainingMs();
+ if (remainingMs > 0) {
+ return { ok: false, remainingMs, state: getReplayState() };
+ }
+ lastTriggeredAt = Date.now();
+ lastTriggeredBy = by || null;
+ const state = getReplayState();
+ replayEvents.emit('update', { state, by: lastTriggeredBy });
+ return { ok: true, state };
+}
+
+module.exports = {
+ tryTriggerReplay,
+ getReplayState,
+ replayEvents,
+};
diff --git a/server/src/services/replayEngineV2/index.js b/server/src/services/replayEngineV2/index.js
index 8de0cba9..cf12b09f 100644
--- a/server/src/services/replayEngineV2/index.js
+++ b/server/src/services/replayEngineV2/index.js
@@ -1,1048 +1,84 @@
-// replay Engine V2
-// Purpose: Defines the replay Engine V2 module and the helpers/state used by this service unit.
-// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary.
-const { spawn, execFile } = require('child_process');
-const fsp = require('fs/promises');
-const fs = require('fs');
-const os = require('os');
-const path = require('path');
+// Replay Engine V2
+// Purpose: Composes replay segment capture, indexing, replay assembly, and health tick orchestration.
+// Scope: Exposes replay build/health APIs while delegating implementation to focused submodules.
+const { execFile } = require('child_process');
const { promisify } = require('util');
-const EventEmitter = require('events');
+const fsp = require('fs/promises');
const logger = require('../../globals/logger').child('replayEngineV2');
const roverManager = require('../roverManager');
const { getRoomCameras, roomCameraEvents } = require('../roomCameraService');
-const io = require('../../globals/io');
-const { getActiveDrivers } = require('../turnService');
-const { getNickname } = require('../nicknameService');
-const { getRecentMessages } = require('../chatService');
-const { resolveDataDir } = require('../../helpers/dataPaths');
-const sharp = require('sharp');
+
+const {
+ SEGMENT_ROOT,
+ SEGMENT_SECONDS,
+ BUFFER_SECONDS,
+ CLEANUP_INTERVAL_MS,
+ BUILD_DURATION_MS,
+} = require('./constants');
+const { events, runtime } = require('./state');
+const { createWorkerManager } = require('./workerManager');
+const { createSegmentStore } = require('./segmentStore');
+const { createSidebarRenderer } = require('./sidebarRenderer');
+const { createReplayBuilder } = require('./replayBuilder');
+const { tryTriggerReplay, getReplayState, replayEvents } = require('./cooldown');
+const {
+ getReplaySources,
+ validateSources,
+ getDefaultWebSources,
+ getDefaultDiscordSources,
+} = require('./replaySources');
+const { registerReplaySocketHooks } = require('./socketHooks');
const execFileAsync = promisify(execFile);
-
-const FFMPEG_BIN = process.env.FFMPEG_BIN || 'ffmpeg';
-const SEGMENT_ROOT = path.join(resolveDataDir(), 'replay-segments');
-const SEGMENT_SECONDS = Math.max(1, Number.parseInt(process.env.REPLAY_SEGMENT_SECONDS || '1', 10));
-const BUFFER_SECONDS = Math.max(20, Number.parseInt(process.env.REPLAY_BUFFER_SECONDS || '45', 10));
-const CLEANUP_INTERVAL_MS = 10_000;
-const BUILD_DURATION_MS = Math.max(5000, Number.parseInt(process.env.REPLAY_DURATION_MS || '20000', 10));
-const BUILD_GUARD_MS = Math.max(200, Number.parseInt(process.env.REPLAY_GUARD_MS || '1200', 10));
-const TARGET_FPS = Math.max(10, Number.parseInt(process.env.REPLAY_TARGET_FPS || '30', 10));
-const MAX_WIDTH = Math.max(320, Number.parseInt(process.env.REPLAY_MAX_WIDTH || '1280', 10));
-const MAX_HEIGHT = Math.max(180, Number.parseInt(process.env.REPLAY_MAX_HEIGHT || '720', 10));
-const MAX_BYTES = Math.floor(Number.parseFloat(process.env.REPLAY_MAX_OUTPUT_MB || '9.5') * 1024 * 1024);
-const SIDEBAR_WIDTH = 190;
-
-const events = new EventEmitter();
-
-const workers = new Map(); // key -> worker
-const pendingWorkerStarts = new Set(); // key
-const segmentIndex = new Map(); // key -> [{filePath,startMs,endMs,mtimeMs,size,kind,sourceType,sourceId}]
-let cleanupTimer = null;
-let activeSegmentRoot = SEGMENT_ROOT;
-let tickInFlight = false;
-
-function sourceKey(source) {
- return `${source.sourceType}__${source.kind}__${source.id}`;
-}
-
-function sourceDirForKey(key) {
- return path.join(activeSegmentRoot, key);
-}
-
-function toSrtReadPath(streamId) {
- return `srt://127.0.0.1:9000?streamid=read:${encodeURIComponent(streamId)}`;
-}
-
-function getRoomCameraStream(camera) {
- if (camera?.streamUrl) return camera.streamUrl;
- const url = String(camera?.url || '');
- if (url.includes('.mjpg') || url.includes('mjpeg') || url.includes('stream')) return url;
- return null;
-}
-
-function listDesiredSources() {
- const sources = [];
- for (const rover of roverManager.getRoster()) {
- if (!roverManager.canReplayRoverId(rover.id)) continue;
- const roverId = String(rover.id);
- sources.push({
- id: roverId,
- sourceType: 'rover',
- kind: 'video',
- label: rover.name || roverId,
- inputUrl: toSrtReadPath(roverId),
- });
- if (rover?.media?.audioPublishUrl) {
- sources.push({
- id: `${roverId}-audio`,
- sourceType: 'rover',
- roverId,
- kind: 'audio',
- label: `${rover.name || roverId} audio`,
- inputUrl: toSrtReadPath(`${roverId}-audio`),
- });
- }
- }
-
- for (const camera of getRoomCameras()) {
- const streamUrl = getRoomCameraStream(camera);
- if (!streamUrl) continue;
- sources.push({
- id: String(camera.id),
- sourceType: 'room',
- kind: 'video',
- label: camera.name || camera.id,
- inputUrl: streamUrl,
- });
- }
-
- return sources;
-}
-
-function buildWorkerArgs(source) {
- const dir = sourceDirForKey(sourceKey(source));
- const pattern = path.join(dir, 'seg-%06d.mp4');
-
- const common = [
- '-hide_banner',
- '-loglevel',
- 'warning',
- '-y',
- '-fflags',
- '+genpts',
- '-use_wallclock_as_timestamps',
- '1',
- '-i',
- source.inputUrl,
- ];
-
- if (source.kind === 'audio') {
- return [
- ...common,
- '-vn',
- '-ac',
- '1',
- '-ar',
- '48000',
- '-af',
- 'aresample=async=1:first_pts=0:min_hard_comp=0.100000,asetpts=N/SR/TB',
- '-c:a',
- 'aac',
- '-b:a',
- '96k',
- '-f',
- 'segment',
- '-segment_time',
- String(SEGMENT_SECONDS),
- '-segment_atclocktime',
- '1',
- '-reset_timestamps',
- '1',
- pattern,
- ];
- }
-
- return [
- ...common,
- '-an',
- '-vf',
- `fps=${TARGET_FPS}`,
- '-c:v',
- 'libx264',
- '-preset',
- 'veryfast',
- '-tune',
- 'zerolatency',
- '-pix_fmt',
- 'yuv420p',
- '-g',
- String(TARGET_FPS * SEGMENT_SECONDS),
- '-keyint_min',
- String(TARGET_FPS * SEGMENT_SECONDS),
- '-sc_threshold',
- '0',
- '-f',
- 'segment',
- '-segment_time',
- String(SEGMENT_SECONDS),
- '-segment_atclocktime',
- '1',
- '-reset_timestamps',
- '1',
- pattern,
- ];
-}
+runtime.activeSegmentRoot = SEGMENT_ROOT;
async function ensureDir(dir) {
await fsp.mkdir(dir, { recursive: true });
}
-function startWorker(source) {
- const key = sourceKey(source);
- if (workers.has(key) || pendingWorkerStarts.has(key)) return;
- pendingWorkerStarts.add(key);
- const dir = sourceDirForKey(key);
- ensureDir(dir)
- .then(() => {
- if (workers.has(key)) {
- pendingWorkerStarts.delete(key);
- return;
- }
- const proc = spawn(FFMPEG_BIN, buildWorkerArgs(source), { stdio: ['ignore', 'ignore', 'pipe'] });
- const worker = {
- key,
- source,
- proc,
- };
- workers.set(key, worker);
- pendingWorkerStarts.delete(key);
- proc.stderr.on('data', (chunk) => {
- const text = String(chunk || '').trim();
- if (!text) return;
- logger.warn('worker stderr', { key, text: text.slice(0, 500) });
- });
- proc.on('exit', (code, signal) => {
- const current = workers.get(key);
- if (current?.proc === proc) {
- workers.delete(key);
- }
- logger.warn('worker exited', { key, code, signal });
- setTimeout(() => {
- const desired = listDesiredSources().find((entry) => sourceKey(entry) === key);
- if (desired && !workers.has(key)) {
- startWorker(desired);
- }
- }, 1500);
- });
- })
- .catch((err) => {
- pendingWorkerStarts.delete(key);
- logger.warn('failed to start worker', { key, error: err.message });
- });
+function getActiveSegmentRoot() {
+ return runtime.activeSegmentRoot;
}
-function stopWorker(key) {
- const worker = workers.get(key);
- if (!worker) return;
- try {
- worker.proc.kill('SIGTERM');
- } catch {
- // noop
- }
- workers.delete(key);
-}
-
-async function syncWorkers() {
- const desired = listDesiredSources();
- const desiredKeys = new Set(desired.map(sourceKey));
-
- for (const source of desired) {
- const key = sourceKey(source);
- if (!workers.has(key)) {
- startWorker(source);
- }
- }
-
- for (const key of Array.from(workers.keys())) {
- if (!desiredKeys.has(key)) {
- stopWorker(key);
- }
- }
-}
-
-async function refreshIndexForWorker(worker) {
- const key = worker.key;
- const dir = sourceDirForKey(key);
- let files;
- try {
- files = await fsp.readdir(dir, { withFileTypes: true });
- } catch {
- return;
- }
-
- const segmentFiles = files
- .filter((entry) => entry.isFile() && /^seg-\d{6}\.mp4$/.test(entry.name))
- .map((entry) => entry.name)
- .sort();
- if (!segmentFiles.length) {
- segmentIndex.set(key, []);
- return;
- }
- const cutoffMs = Date.now() - BUFFER_SECONDS * 1000 - 5000;
-
- const entries = [];
- for (const filename of segmentFiles) {
- const filePath = path.join(dir, filename);
- let stat;
- try {
- stat = await fsp.stat(filePath);
- } catch {
- continue;
- }
- if (!stat.isFile() || stat.size < 4096) continue;
- const durationSec = SEGMENT_SECONDS;
- const endMs = Math.round(stat.mtimeMs);
- const startMs = Math.round(endMs - durationSec * 1000);
- if (endMs < cutoffMs) continue;
- entries.push({
- filePath,
- startMs,
- endMs,
- mtimeMs: stat.mtimeMs,
- size: stat.size,
- kind: worker.source.kind,
- sourceType: worker.source.sourceType,
- sourceId: worker.source.id,
- roverId: worker.source.roverId || null,
- });
- }
-
- segmentIndex.set(key, entries.sort((a, b) => a.startMs - b.startMs));
-}
-
-async function refreshSegmentIndex() {
- const list = Array.from(workers.values());
- for (const worker of list) {
- await refreshIndexForWorker(worker);
- }
-}
-
-async function cleanupOldFiles() {
- const cutoff = Date.now() - BUFFER_SECONDS * 1000;
- try {
- await ensureDir(activeSegmentRoot);
- const dirs = await fsp.readdir(activeSegmentRoot, { withFileTypes: true });
- for (const dirent of dirs) {
- if (!dirent.isDirectory()) continue;
- const dirPath = path.join(activeSegmentRoot, dirent.name);
- let files;
- try {
- files = await fsp.readdir(dirPath, { withFileTypes: true });
- } catch {
- continue;
- }
- for (const file of files) {
- if (!file.isFile() || !file.name.endsWith('.mp4')) continue;
- const filePath = path.join(dirPath, file.name);
- try {
- const stat = await fsp.stat(filePath);
- if (stat.mtimeMs < cutoff) {
- await fsp.unlink(filePath);
- }
- } catch {
- // noop
- }
- }
- }
- } catch (err) {
- logger.warn('cleanup failed', err.message);
- }
-}
-
-function getVideoEntriesForSource(source) {
- const key = sourceKey({ sourceType: String(source.type), kind: 'video', id: String(source.id) });
- return segmentIndex.get(key) || [];
-}
-
-function getAudioEntriesForRover(roverId) {
- const key = sourceKey({ sourceType: 'rover', kind: 'audio', id: `${String(roverId)}-audio` });
- return segmentIndex.get(key) || [];
-}
-
-function overlapping(entries, startMs, endMs) {
- return entries.filter((entry) => entry.endMs > startMs && entry.startMs < endMs);
-}
-
-function buildGridLayout(count) {
- const cols = Math.ceil(Math.sqrt(count));
- const rows = Math.ceil(count / cols);
- return { cols, rows };
-}
-
-function clampEven(value) {
- return Math.max(2, Math.floor(value / 2) * 2);
-}
-
-function escapeDrawtext(text) {
- return String(text || '')
- .replace(/\\/g, '\\\\')
- .replace(/:/g, '\\:')
- .replace(/'/g, "\\'")
- .replace(/,/g, '\\,')
- .replace(/\[/g, '\\[')
- .replace(/\]/g, '\\]')
- .replace(/%/g, '\\%');
-}
-
-function scalePadFilter(tileWidth, tileHeight, titleText = '') {
- const safeTitle = escapeDrawtext(titleText);
- return (
- `scale=${tileWidth}:${tileHeight}:force_original_aspect_ratio=decrease:flags=lanczos,` +
- `pad=${tileWidth}:${tileHeight}:(ow-iw)/2:(oh-ih)/2:color=black,` +
- `drawtext=text='${safeTitle}':x=(w-text_w)/2:y=7:fontsize=14:fontcolor=white:borderw=1:bordercolor=black@0.7:box=1:boxcolor=black@0.52:boxborderw=4,` +
- 'setsar=1'
- );
-}
-
-function sanitizeReplayTitle(title, fallback = 'Replay') {
- const value = String(title || '').trim();
- if (!value) return fallback;
- return value.slice(0, 120);
-}
-
-function resolveDefaultReplayTitle(requester = '', sources = []) {
- const requesterLabel = String(requester || 'Someone').trim() || 'Someone';
- const rover = (Array.isArray(sources) ? sources : []).find((entry) => entry?.type === 'rover');
- const roverLabel = rover?.label || rover?.id || 'a rover';
- return `${requesterLabel} driving ${roverLabel}`;
-}
-
-function buildDriverBatterySnapshot(selectedRoverIds = []) {
- const activeDrivers = getActiveDrivers();
- const byId = new Map(roverManager.getRoster().map((rover) => [String(rover.id), rover]));
- const lines = [];
- for (const roverId of selectedRoverIds) {
- const socketId = activeDrivers[String(roverId)];
- if (!socketId) continue;
- const socket = io.sockets.sockets.get(socketId);
- const nickname = getNickname(socket) || socket?.data?.user?.username || String(socketId);
- const rover = byId.get(String(roverId));
- const roverName = rover?.name || roverId;
- const percent = rover?.batteryState?.percentDisplay;
- const batteryLabel = Number.isFinite(percent) ? `${percent}%` : '--%';
- lines.push(`${nickname} driving ${roverName} (${batteryLabel})`);
- }
- return lines;
-}
-
-function buildChatEventsForWindow(startMs, endMs, limit = 22, preWindowCount = 10) {
- const all = getRecentMessages(300, { includeSystem: false });
- const normalized = all
- .filter((msg) => Number.isFinite(msg?.ts))
- .map((msg) => {
- const nickname = String(msg?.nickname || msg?.discordUserName || 'user').trim() || 'user';
- const text = String(msg?.text || '').replace(/\s+/g, ' ').trim();
- return {
- ts: Number(msg.ts),
- nickname: nickname.slice(0, 32),
- text: text.slice(0, 120),
- role: String(msg?.role || ''),
- fromDiscord: Boolean(msg?.fromDiscord),
- roverId: msg?.roverId ? String(msg.roverId) : '',
- roverColor: msg?.roverColor ? String(msg.roverColor) : '',
- };
- });
- const beforeWindow = normalized.filter((msg) => msg.ts < startMs).slice(-preWindowCount);
- const inWindow = normalized.filter((msg) => msg.ts >= startMs && msg.ts <= endMs).slice(-limit);
- return [...beforeWindow, ...inWindow].sort((a, b) => a.ts - b.ts);
-}
-
-function escapeXml(text) {
- return String(text || '')
- .replace(/&/g, '&')
- .replace(//g, '>')
- .replace(/"/g, '"')
- .replace(/'/g, ''');
-}
-
-function wrapTextLines(text, maxChars = 24) {
- const words = String(text || '').split(/\s+/).filter(Boolean);
- const lines = [];
- let current = '';
- for (const word of words) {
- const candidate = current ? `${current} ${word}` : word;
- if (candidate.length > maxChars && current) {
- lines.push(current);
- current = word;
- } else {
- current = candidate;
- }
- }
- if (current) lines.push(current);
- return lines.slice(0, 4);
-}
-
-function hexToRgb(hex) {
- const value = String(hex || '').trim();
- const match = /^#([0-9A-Fa-f]{6})$/.exec(value);
- if (!match) return null;
- const raw = match[1];
- return {
- r: parseInt(raw.slice(0, 2), 16),
- g: parseInt(raw.slice(2, 4), 16),
- b: parseInt(raw.slice(4, 6), 16),
- };
-}
-
-function roleColor(role = '') {
- switch (String(role)) {
- case 'admin':
- case 'lockdown':
- case 'lockdown-admin':
- return '#FCD34D'; // amber-300
- case 'spectator':
- return '#94A3B8'; // slate-400
- default:
- return '#7DD3FC'; // sky-300
- }
-}
-
-function renderSidebarSvg({
- width,
- height,
- title,
- driverBatteryLines,
- chatLines,
-}) {
- const textCols = Math.max(24, Math.floor((width - 16) / 6));
- const titleParts = wrapTextLines(title, Math.max(26, textCols)).slice(0, 4).map((line) => escapeXml(line));
- const statLines = driverBatteryLines.slice(0, 8).map((line) => escapeXml(line));
- const normalizedChats = chatLines.slice(-12).map((entry) => {
- const wrapped = wrapTextLines(entry.text || '', Math.max(24, textCols)).slice(0, 4).map((line) => escapeXml(line));
- const nick = escapeXml(entry.nickname || 'user');
- const roverId = escapeXml(entry.roverId || '');
- const roverRgb = hexToRgb(entry.roverColor || '');
- const roverBadgeBg = roverRgb ? `rgba(${roverRgb.r},${roverRgb.g},${roverRgb.b},0.18)` : 'rgba(30,41,59,0.70)';
- const roverBadgeBorder = roverRgb ? `rgba(${roverRgb.r},${roverRgb.g},${roverRgb.b},0.60)` : 'rgba(71,85,105,0.75)';
- return {
- nick,
- wrapped,
- nameColor: roleColor(entry.role),
- fromDiscord: Boolean(entry.fromDiscord),
- roverId,
- roverBadgeBg,
- roverBadgeBorder,
- bubbleTone: entry.fromDiscord ? 'discordBubble' : 'chatBubble',
- };
- });
-
- const shapes = [];
- const textRows = [];
- const pad = 3;
- const cardX = pad;
- const cardW = width - pad * 2;
- let y = pad;
-
- // Title card
- const titleCardH = Math.max(26, pad * 2 + titleParts.length * 15);
- shapes.push(``);
- let ty = y + pad + 11;
- for (const part of titleParts) {
- textRows.push(`${part}`);
- ty += 15;
- }
- y += titleCardH + pad;
-
- // Drivers card
- const driverLines = statLines.length ? statLines : ['No active drivers'];
- const driversCardH = pad * 2 + 11 + driverLines.length * 13;
- shapes.push(``);
- textRows.push(`Drivers`);
- let dy = y + pad + 19;
- for (const line of driverLines) {
- textRows.push(`${line}`);
- dy += 13;
- }
- y += driversCardH + pad;
-
- // Chat card
- const chatCardH = Math.max(80, height - y - pad);
- shapes.push(``);
- textRows.push(`Chat`);
-
- let cy = y + pad + 12;
- const bubbleX = cardX + pad;
- const bubbleW = cardW - pad * 2;
- if (!normalizedChats.length) {
- textRows.push(`No chat in replay window`);
- } else {
- for (let i = 0; i < normalizedChats.length; i += 1) {
- const block = normalizedChats[i];
- const nameW = Math.min(72, block.nick.length * 5.2);
- const badgeW = block.roverId ? Math.min(46, Math.max(18, block.roverId.length * 5 + 6)) : 0;
- const badgeGap = block.roverId ? 3 : 0;
- const prefixChars = Math.ceil((nameW + badgeW + badgeGap + 10) / 5.8);
- const firstLineRaw = String(block.wrapped[0] || '').trim();
- const firstLine = firstLineRaw ? firstLineRaw : '';
- const remainingRaw = block.wrapped.slice(firstLineRaw ? 1 : 0).map((line) => String(line || '').trim()).filter(Boolean);
- const fullText = (firstLine ? [firstLine, ...remainingRaw] : remainingRaw).join(' ');
- const inlineWrapped = wrapTextLines(fullText, Math.max(16, textCols - prefixChars)).slice(0, 4).map((line) => escapeXml(line));
- const bubbleH = pad * 2 + Math.max(1, inlineWrapped.length) * 12;
- if (cy + bubbleH + pad > y + chatCardH - pad) break;
- shapes.push(
- ``,
- );
- const nameX = bubbleX + pad;
- const textStartX = nameX + nameW + 3 + (block.roverId ? badgeW + badgeGap : 0);
- textRows.push(`${block.nick}`);
- if (block.fromDiscord) {
- textRows.push(`◈`);
- }
- if (block.roverId) {
- const badgeX = nameX + nameW + 3;
- const badgeTextX = badgeX + 3;
- shapes.push(
- ``,
- );
- textRows.push(`${block.roverId}`);
- }
- let by = cy + pad + 8;
- for (let lineIdx = 0; lineIdx < inlineWrapped.length; lineIdx += 1) {
- const line = inlineWrapped[lineIdx];
- const lineX = lineIdx === 0 ? textStartX : bubbleX + pad;
- textRows.push(`${line}`);
- by += 12;
- }
- cy += bubbleH + pad;
- }
- }
-
- return `
-`;
-}
-
-async function renderSidebarVideo({
- tmpDir,
- title,
- durationSec,
- height,
- windowStartMs,
- driverBatteryLines = [],
- chatEvents = [],
-}) {
- const framesDir = path.join(tmpDir, 'sidebar-frames');
- const sidebarPath = path.join(tmpDir, 'sidebar.mp4');
- await ensureDir(framesDir);
- const secondCount = Math.max(1, Math.ceil(durationSec));
- for (let second = 0; second < secondCount; second += 1) {
- const sliceEndMs = windowStartMs + (second + 1) * 1000;
- const visibleChat = chatEvents.filter((entry) => entry.ts <= sliceEndMs).slice(-10);
- const svg = renderSidebarSvg({
- width: SIDEBAR_WIDTH,
- height,
- title,
- driverBatteryLines,
- chatLines: visibleChat,
- });
- const framePath = path.join(framesDir, `frame-${String(second + 1).padStart(4, '0')}.png`);
- await sharp(Buffer.from(svg, 'utf8')).png().toFile(framePath);
- }
-
- await execFileAsync(FFMPEG_BIN, [
- '-y',
- '-hide_banner',
- '-loglevel',
- 'error',
- '-framerate',
- '1',
- '-i',
- path.join(framesDir, 'frame-%04d.png'),
- '-vf',
- `fps=${TARGET_FPS},format=yuv420p`,
- '-t',
- durationSec.toFixed(3),
- '-c:v',
- 'libx264',
- '-preset',
- 'veryfast',
- '-pix_fmt',
- 'yuv420p',
- sidebarPath,
- ]);
- return sidebarPath;
-}
-
-async function concatFiles(inputPaths, outPath) {
- const listPath = `${outPath}.concat.txt`;
- const body = inputPaths.map((file) => `file '${file.replace(/'/g, "'\\''")}'`).join('\n');
- await fsp.writeFile(listPath, `${body}\n`, 'utf8');
- await execFileAsync(FFMPEG_BIN, [
- '-y',
- '-hide_banner',
- '-loglevel',
- 'error',
- '-f',
- 'concat',
- '-safe',
- '0',
- '-i',
- listPath,
- '-c',
- 'copy',
- outPath,
- ]);
-}
-
-async function probeMaxFrameSize(paths) {
- let maxWidth = 0;
- let maxHeight = 0;
- for (const filePath of paths) {
- try {
- const { stdout } = await execFileAsync('ffprobe', [
- '-v',
- 'error',
- '-select_streams',
- 'v:0',
- '-show_entries',
- 'stream=width,height',
- '-of',
- 'csv=p=0',
- filePath,
- ]);
- const [wRaw, hRaw] = stdout.trim().split(',');
- const w = Number(wRaw);
- const h = Number(hRaw);
- if (Number.isFinite(w) && Number.isFinite(h)) {
- maxWidth = Math.max(maxWidth, w);
- maxHeight = Math.max(maxHeight, h);
- }
- } catch {
- // noop
- }
- }
- return { maxWidth, maxHeight };
-}
-
-async function buildReplayVideo({ sources = [], title = '', requester = '', includeSidebar = true } = {}) {
- if (!Array.isArray(sources) || !sources.length) {
- throw new Error('No replay sources selected');
- }
-
- const tEnd = Date.now() - BUILD_GUARD_MS;
- const tStart = tEnd - BUILD_DURATION_MS;
- const resolvedTitle = sanitizeReplayTitle(title, resolveDefaultReplayTitle(requester, sources));
- const tmpDir = await fsp.mkdtemp(path.join(os.tmpdir(), 'mrr-replay-v2-'));
-
- try {
- await refreshSegmentIndex();
-
- const usedSources = [];
- const missingSources = [];
- const normalizedVideos = [];
- const normalizedAudios = [];
-
- for (let i = 0; i < sources.length; i += 1) {
- const source = sources[i];
- const sourceId = String(source.id);
- const videoEntries = overlapping(
- getVideoEntriesForSource({ type: String(source.type), id: sourceId }),
- tStart,
- tEnd,
- );
- if (!videoEntries.length) {
- missingSources.push({ ...source, reason: 'no video coverage in replay window' });
- continue;
- }
-
- const videoConcat = path.join(tmpDir, `video-${i}.mp4`);
- await concatFiles(videoEntries.map((entry) => entry.filePath), videoConcat);
- const videoTrimmed = path.join(tmpDir, `video-${i}.trim.mp4`);
- const firstStartMs = videoEntries[0].startMs;
- const ss = Math.max(0, (tStart - firstStartMs) / 1000);
- const to = Math.max(ss + 0.1, (tEnd - firstStartMs) / 1000);
- await execFileAsync(FFMPEG_BIN, [
- '-y',
- '-hide_banner',
- '-loglevel',
- 'error',
- '-ss',
- ss.toFixed(3),
- '-to',
- to.toFixed(3),
- '-i',
- videoConcat,
- '-an',
- '-c:v',
- 'libx264',
- '-preset',
- 'veryfast',
- '-pix_fmt',
- 'yuv420p',
- '-r',
- String(TARGET_FPS),
- videoTrimmed,
- ]);
-
- normalizedVideos.push({ path: videoTrimmed, source });
- usedSources.push(source);
-
- if (source.type === 'rover') {
- const audioEntries = overlapping(getAudioEntriesForRover(sourceId), tStart, tEnd);
- if (audioEntries.length) {
- const audioConcat = path.join(tmpDir, `audio-${i}.m4a`);
- await concatFiles(audioEntries.map((entry) => entry.filePath), audioConcat);
- const audioTrimmed = path.join(tmpDir, `audio-${i}.trim.m4a`);
- const firstAudioStartMs = audioEntries[0].startMs;
- const ass = Math.max(0, (tStart - firstAudioStartMs) / 1000);
- const ato = Math.max(ass + 0.1, (tEnd - firstAudioStartMs) / 1000);
- await execFileAsync(FFMPEG_BIN, [
- '-y',
- '-hide_banner',
- '-loglevel',
- 'error',
- '-ss',
- ass.toFixed(3),
- '-to',
- ato.toFixed(3),
- '-i',
- audioConcat,
- '-vn',
- '-ac',
- '1',
- '-ar',
- '48000',
- '-c:a',
- 'aac',
- '-b:a',
- '96k',
- audioTrimmed,
- ]);
- normalizedAudios.push(audioTrimmed);
- }
- }
- }
-
- if (!normalizedVideos.length) {
- throw new Error('No replay segments available for selected sources');
- }
-
- const layout = buildGridLayout(normalizedVideos.length);
- const { maxWidth, maxHeight } = await probeMaxFrameSize(normalizedVideos.map((v) => v.path));
- let tileWidth = maxWidth || 640;
- let tileHeight = maxHeight || 360;
- let outWidth = tileWidth * layout.cols;
- let outHeight = tileHeight * layout.rows;
- if (outWidth > MAX_WIDTH || outHeight > MAX_HEIGHT) {
- const scale = Math.min(MAX_WIDTH / outWidth, MAX_HEIGHT / outHeight);
- tileWidth *= scale;
- tileHeight *= scale;
- outWidth = tileWidth * layout.cols;
- outHeight = tileHeight * layout.rows;
- }
- tileWidth = clampEven(tileWidth);
- tileHeight = clampEven(tileHeight);
- const durationSec = BUILD_DURATION_MS / 1000;
-
- const inputArgs = [];
- const filterParts = [];
- const layoutParts = [];
-
- for (let i = 0; i < normalizedVideos.length; i += 1) {
- inputArgs.push('-i', normalizedVideos[i].path);
- const sourceTitle = normalizedVideos[i]?.source?.label || normalizedVideos[i]?.source?.id || `Source ${i + 1}`;
- filterParts.push(`[${i}:v]${scalePadFilter(tileWidth, tileHeight, sourceTitle)}[v${i}]`);
- const x = (i % layout.cols) * tileWidth;
- const y = Math.floor(i / layout.cols) * tileHeight;
- layoutParts.push(`${x}_${y}`);
- }
-
- let audioInputStart = normalizedVideos.length;
- let sidebarInputIndex = -1;
- if (includeSidebar) {
- const selectedRoverIds = usedSources
- .filter((entry) => entry?.type === 'rover')
- .map((entry) => String(entry.id));
- const driverBatteryLines = buildDriverBatterySnapshot(selectedRoverIds);
- const chatEvents = buildChatEventsForWindow(tStart, tEnd);
- const sidebarPath = await renderSidebarVideo({
- tmpDir,
- title: resolvedTitle,
- durationSec,
- height: clampEven(tileHeight * layout.rows),
- windowStartMs: tStart,
- driverBatteryLines,
- chatEvents,
- });
- inputArgs.push('-i', sidebarPath);
- sidebarInputIndex = normalizedVideos.length;
- audioInputStart = normalizedVideos.length + 1;
- }
-
- for (const audioPath of normalizedAudios) {
- inputArgs.push('-i', audioPath);
- }
-
- if (normalizedVideos.length === 1) {
- filterParts.push('[v0]null[vgrid]');
- } else {
- filterParts.push(
- `${normalizedVideos.map((_, i) => `[v${i}]`).join('')}` +
- `xstack=inputs=${normalizedVideos.length}:layout=${layoutParts.join('|')}:fill=black[vgrid]`,
- );
- }
- if (includeSidebar) {
- filterParts.push(`[vgrid][${sidebarInputIndex}:v]hstack=inputs=2[vout]`);
- } else {
- filterParts.push('[vgrid]null[vout]');
- }
-
- if (normalizedAudios.length) {
- const audioRefs = normalizedAudios.map((_, idx) => `[${audioInputStart + idx}:a]`).join('');
- filterParts.push(`${audioRefs}amix=inputs=${normalizedAudios.length}:normalize=0,alimiter=limit=0.9[aout]`);
- }
- const targetBitrateKbps = Math.max(400, Math.floor((MAX_BYTES * 8) / durationSec / 1000));
- const outPath = path.join(tmpDir, 'replay.mp4');
-
- const args = [
- '-y',
- '-hide_banner',
- '-loglevel',
- 'error',
- ...inputArgs,
- '-filter_complex',
- filterParts.join(';'),
- '-map',
- '[vout]',
- '-c:v',
- 'libx264',
- '-preset',
- 'veryfast',
- '-pix_fmt',
- 'yuv420p',
- '-r',
- String(TARGET_FPS),
- '-b:v',
- `${targetBitrateKbps}k`,
- '-maxrate',
- `${Math.floor(targetBitrateKbps * 1.15)}k`,
- '-bufsize',
- `${Math.floor(targetBitrateKbps * 2)}k`,
- ];
-
- if (normalizedAudios.length) {
- args.push('-map', '[aout]', '-c:a', 'aac', '-b:a', '128k');
- }
-
- args.push(outPath);
-
- await execFileAsync(FFMPEG_BIN, args);
- const buffer = await fsp.readFile(outPath);
- return { buffer, usedSources, missingSources, title: resolvedTitle };
- } finally {
- try {
- await fsp.rm(tmpDir, { recursive: true, force: true });
- } catch {
- // noop
- }
- }
-}
+const workerManager = createWorkerManager({ getActiveSegmentRoot });
+const segmentStore = createSegmentStore({ getActiveSegmentRoot });
+const sidebarRenderer = createSidebarRenderer({ execFileAsync, ensureDir });
+const replayBuilder = createReplayBuilder({
+ execFileAsync,
+ fsp,
+ ensureDir,
+ renderSidebarVideo: sidebarRenderer.renderSidebarVideo,
+ getVideoEntriesForSource: segmentStore.getVideoEntriesForSource,
+ getAudioEntriesForRover: segmentStore.getAudioEntriesForRover,
+ overlapping: segmentStore.overlapping,
+});
+registerReplaySocketHooks({ tryTriggerReplay, validateSources, getDefaultWebSources });
function getReplayHealthSnapshot() {
- const now = Date.now();
- const neededCount = Math.max(1, Math.ceil(BUILD_DURATION_MS / 1000));
- const sources = [];
- let readyCount = 0;
-
- const replaySources = [];
- for (const rover of roverManager.getRoster().filter((entry) => roverManager.canReplayRoverId(entry.id))) {
- replaySources.push({ type: 'rover', id: String(rover.id), label: rover.name || rover.id });
- }
- for (const camera of getRoomCameras()) {
- replaySources.push({ type: 'room', id: String(camera.id), label: camera.name || camera.id });
- }
-
- for (const source of replaySources) {
- const key = sourceKey({ sourceType: source.type, kind: 'video', id: String(source.id) });
- const dir = sourceDirForKey(key);
- let recentCount = 0;
- let lastSegmentAt = null;
- try {
- const files = fs.readdirSync(dir);
- for (const name of files) {
- if (!/^seg-\d{6}\.mp4$/.test(name)) continue;
- const full = path.join(dir, name);
- let stat;
- try {
- stat = fs.statSync(full);
- } catch {
- continue;
- }
- if (!stat.isFile() || stat.size < 4096) continue;
- if (stat.mtimeMs > (lastSegmentAt || 0)) {
- lastSegmentAt = stat.mtimeMs;
- }
- if (now - stat.mtimeMs <= BUFFER_SECONDS * 1000) {
- recentCount += 1;
- }
- }
- } catch {
- // ignore missing source directory
- }
- const ready = recentCount >= neededCount;
- if (ready) readyCount += 1;
- sources.push({
- type: source.type,
- id: source.id,
- label: source.label,
- recentCount,
- neededCount,
- lastSegmentAt,
- ready,
- });
- }
-
- return {
- sources,
- readyCount,
- totalCount: sources.length,
- };
-}
-
-async function bootstrapIndexFromDisk() {
- await ensureDir(activeSegmentRoot);
- const dirs = await fsp.readdir(activeSegmentRoot, { withFileTypes: true });
- for (const dirent of dirs) {
- if (!dirent.isDirectory()) continue;
- segmentIndex.set(dirent.name, []);
- }
+ return segmentStore.getReplayHealthSnapshot({ BUILD_DURATION_MS, roverManager, getRoomCameras });
}
async function tick() {
- if (tickInFlight) return;
- tickInFlight = true;
+ if (runtime.tickInFlight) return;
+ runtime.tickInFlight = true;
try {
- await syncWorkers();
- await refreshSegmentIndex();
- await cleanupOldFiles();
- events.emit('health', getReplayHealthSnapshot());
+ await workerManager.syncWorkers();
+ await segmentStore.refreshSegmentIndex();
+ await segmentStore.cleanupOldFiles();
+ events.emit('health', getReplayHealthSnapshot());
} finally {
- tickInFlight = false;
+ runtime.tickInFlight = false;
}
}
async function start() {
await ensureDir(SEGMENT_ROOT);
- activeSegmentRoot = SEGMENT_ROOT;
- logger.info('Replay engine using segment root', { segmentRoot: activeSegmentRoot });
- await bootstrapIndexFromDisk();
+ runtime.activeSegmentRoot = SEGMENT_ROOT;
+ logger.info('Replay engine using segment root', { segmentRoot: runtime.activeSegmentRoot });
+ await segmentStore.bootstrapIndexFromDisk();
await tick();
- if (cleanupTimer) clearInterval(cleanupTimer);
- cleanupTimer = setInterval(() => {
+ if (runtime.cleanupTimer) clearInterval(runtime.cleanupTimer);
+ runtime.cleanupTimer = setInterval(() => {
tick().catch((err) => logger.warn('tick failed', err.message));
}, CLEANUP_INTERVAL_MS);
}
@@ -1062,10 +98,17 @@ start().catch((err) => {
});
module.exports = {
- buildReplayVideo,
+ buildReplayVideo: replayBuilder.buildReplayVideo,
+ tryTriggerReplay,
+ getReplayState,
+ replayEvents,
+ getReplaySources,
+ validateSources,
+ getDefaultWebSources,
+ getDefaultDiscordSources,
getReplayHealthSnapshot,
replayEngineEvents: events,
- replaySegmentRootDir: () => activeSegmentRoot,
+ replaySegmentRootDir: () => runtime.activeSegmentRoot,
replaySegmentSeconds: SEGMENT_SECONDS,
replayBufferSeconds: BUFFER_SECONDS,
};
diff --git a/server/src/services/replayEngineV2/replayBuilder.js b/server/src/services/replayEngineV2/replayBuilder.js
new file mode 100644
index 00000000..16cc6514
--- /dev/null
+++ b/server/src/services/replayEngineV2/replayBuilder.js
@@ -0,0 +1,185 @@
+// Replay Builder Pipeline
+// Purpose: Assembles selected buffered segments into final replay output video with optional sidebar.
+// Scope: Owns concat/probe/layout/transcode pipeline and returns replay buffer plus source usage metadata.
+const os = require('os');
+const path = require('path');
+const { getActiveDrivers } = require('../turnService');
+const { getNickname } = require('../nicknameService');
+const { getRecentMessages } = require('../chatService');
+const io = require('../../globals/io');
+const roverManager = require('../roverManager');
+const { FFMPEG_BIN, BUILD_DURATION_MS, BUILD_GUARD_MS, TARGET_FPS, MAX_WIDTH, MAX_HEIGHT, MAX_BYTES } = require('./constants');
+
+function buildGridLayout(count) {
+ const cols = Math.ceil(Math.sqrt(count));
+ const rows = Math.ceil(count / cols);
+ return { cols, rows };
+}
+function clampEven(value) { return Math.max(2, Math.floor(value / 2) * 2); }
+function escapeDrawtext(text) { return String(text || '').replace(/\\/g, '\\\\').replace(/:/g, '\\:').replace(/'/g, "\\'").replace(/,/g, '\\,').replace(/\[/g, '\\[').replace(/\]/g, '\\]').replace(/%/g, '\\%'); }
+function scalePadFilter(tileWidth, tileHeight, titleText = '') {
+ const safeTitle = escapeDrawtext(titleText);
+ return `scale=${tileWidth}:${tileHeight}:force_original_aspect_ratio=decrease:flags=lanczos,pad=${tileWidth}:${tileHeight}:(ow-iw)/2:(oh-ih)/2:color=black,drawtext=text='${safeTitle}':x=(w-text_w)/2:y=7:fontsize=14:fontcolor=white:borderw=1:bordercolor=black@0.7:box=1:boxcolor=black@0.52:boxborderw=4,setsar=1`;
+}
+function sanitizeReplayTitle(title, fallback = 'Replay') { const value = String(title || '').trim(); return value ? value.slice(0, 120) : fallback; }
+function resolveDefaultReplayTitle(requester = '', sources = []) {
+ const requesterLabel = String(requester || 'Someone').trim() || 'Someone';
+ const rover = (Array.isArray(sources) ? sources : []).find((entry) => entry?.type === 'rover');
+ return `${requesterLabel} driving ${rover?.label || rover?.id || 'a rover'}`;
+}
+function buildDriverBatterySnapshot(selectedRoverIds = []) {
+ const activeDrivers = getActiveDrivers();
+ const byId = new Map(roverManager.getRoster().map((rover) => [String(rover.id), rover]));
+ const lines = [];
+ for (const roverId of selectedRoverIds) {
+ const socketId = activeDrivers[String(roverId)];
+ if (!socketId) continue;
+ const socket = io.sockets.sockets.get(socketId);
+ const nickname = getNickname(socket) || socket?.data?.user?.username || String(socketId);
+ const rover = byId.get(String(roverId));
+ const roverName = rover?.name || roverId;
+ const percent = rover?.batteryState?.percentDisplay;
+ lines.push(`${nickname} driving ${roverName} (${Number.isFinite(percent) ? `${percent}%` : '--%'})`);
+ }
+ return lines;
+}
+function buildChatEventsForWindow(startMs, endMs, limit = 22, preWindowCount = 10) {
+ const all = getRecentMessages(300, { includeSystem: false });
+ const normalized = all.filter((msg) => Number.isFinite(msg?.ts)).map((msg) => ({ ts: Number(msg.ts), nickname: String(msg?.nickname || msg?.discordUserName || 'user').trim().slice(0, 32) || 'user', text: String(msg?.text || '').replace(/\s+/g, ' ').trim().slice(0, 120), role: String(msg?.role || ''), fromDiscord: Boolean(msg?.fromDiscord), roverId: msg?.roverId ? String(msg.roverId) : '', roverColor: msg?.roverColor ? String(msg.roverColor) : '' }));
+ const beforeWindow = normalized.filter((msg) => msg.ts < startMs).slice(-preWindowCount);
+ const inWindow = normalized.filter((msg) => msg.ts >= startMs && msg.ts <= endMs).slice(-limit);
+ return [...beforeWindow, ...inWindow].sort((a, b) => a.ts - b.ts);
+}
+
+function createReplayBuilder({ execFileAsync, fsp, ensureDir, renderSidebarVideo, getVideoEntriesForSource, getAudioEntriesForRover, overlapping }) {
+ async function concatFiles(inputPaths, outPath) {
+ const listPath = `${outPath}.concat.txt`;
+ const body = inputPaths.map((file) => `file '${file.replace(/'/g, "'\\''")}'`).join('\n');
+ await fsp.writeFile(listPath, `${body}\n`, 'utf8');
+ await execFileAsync(FFMPEG_BIN, ['-y','-hide_banner','-loglevel','error','-f','concat','-safe','0','-i',listPath,'-c','copy',outPath]);
+ }
+
+ async function probeMaxFrameSize(paths) {
+ let maxWidth = 0, maxHeight = 0;
+ for (const filePath of paths) {
+ try {
+ const { stdout } = await execFileAsync('ffprobe', ['-v','error','-select_streams','v:0','-show_entries','stream=width,height','-of','csv=p=0',filePath]);
+ const [wRaw, hRaw] = stdout.trim().split(',');
+ const w = Number(wRaw), h = Number(hRaw);
+ if (Number.isFinite(w) && Number.isFinite(h)) { maxWidth = Math.max(maxWidth, w); maxHeight = Math.max(maxHeight, h); }
+ } catch {}
+ }
+ return { maxWidth, maxHeight };
+ }
+
+ async function buildReplayVideo({ sources = [], title = '', requester = '', includeSidebar = true } = {}) {
+ if (!Array.isArray(sources) || !sources.length) throw new Error('No replay sources selected');
+ const tEnd = Date.now() - BUILD_GUARD_MS;
+ const tStart = tEnd - BUILD_DURATION_MS;
+ const resolvedTitle = sanitizeReplayTitle(title, resolveDefaultReplayTitle(requester, sources));
+ const tmpDir = await fsp.mkdtemp(path.join(os.tmpdir(), 'mrr-replay-v2-'));
+
+ try {
+ const usedSources = [];
+ const missingSources = [];
+ const normalizedVideos = [];
+ const normalizedAudios = [];
+
+ for (let i = 0; i < sources.length; i += 1) {
+ const source = sources[i];
+ const sourceId = String(source.id);
+ const videoEntries = overlapping(getVideoEntriesForSource({ type: String(source.type), id: sourceId }), tStart, tEnd);
+ if (!videoEntries.length) { missingSources.push({ ...source, reason: 'no video coverage in replay window' }); continue; }
+
+ const videoConcat = path.join(tmpDir, `video-${i}.mp4`);
+ await concatFiles(videoEntries.map((entry) => entry.filePath), videoConcat);
+ const videoTrimmed = path.join(tmpDir, `video-${i}.trim.mp4`);
+ const firstStartMs = videoEntries[0].startMs;
+ const ss = Math.max(0, (tStart - firstStartMs) / 1000);
+ const to = Math.max(ss + 0.1, (tEnd - firstStartMs) / 1000);
+ await execFileAsync(FFMPEG_BIN, ['-y','-hide_banner','-loglevel','error','-ss',ss.toFixed(3),'-to',to.toFixed(3),'-i',videoConcat,'-an','-c:v','libx264','-preset','veryfast','-pix_fmt','yuv420p','-r',String(TARGET_FPS),videoTrimmed]);
+
+ normalizedVideos.push({ path: videoTrimmed, source });
+ usedSources.push(source);
+
+ if (source.type === 'rover') {
+ const audioEntries = overlapping(getAudioEntriesForRover(sourceId), tStart, tEnd);
+ if (audioEntries.length) {
+ const audioConcat = path.join(tmpDir, `audio-${i}.m4a`);
+ await concatFiles(audioEntries.map((entry) => entry.filePath), audioConcat);
+ const audioTrimmed = path.join(tmpDir, `audio-${i}.trim.m4a`);
+ const firstAudioStartMs = audioEntries[0].startMs;
+ const ass = Math.max(0, (tStart - firstAudioStartMs) / 1000);
+ const ato = Math.max(ass + 0.1, (tEnd - firstAudioStartMs) / 1000);
+ await execFileAsync(FFMPEG_BIN, ['-y','-hide_banner','-loglevel','error','-ss',ass.toFixed(3),'-to',ato.toFixed(3),'-i',audioConcat,'-vn','-ac','1','-ar','48000','-c:a','aac','-b:a','96k',audioTrimmed]);
+ normalizedAudios.push(audioTrimmed);
+ }
+ }
+ }
+
+ if (!normalizedVideos.length) throw new Error('No replay segments available for selected sources');
+
+ const layout = buildGridLayout(normalizedVideos.length);
+ const { maxWidth, maxHeight } = await probeMaxFrameSize(normalizedVideos.map((v) => v.path));
+ let tileWidth = maxWidth || 640;
+ let tileHeight = maxHeight || 360;
+ let outWidth = tileWidth * layout.cols;
+ let outHeight = tileHeight * layout.rows;
+ if (outWidth > MAX_WIDTH || outHeight > MAX_HEIGHT) {
+ const scale = Math.min(MAX_WIDTH / outWidth, MAX_HEIGHT / outHeight);
+ tileWidth *= scale; tileHeight *= scale; outWidth = tileWidth * layout.cols; outHeight = tileHeight * layout.rows;
+ }
+ tileWidth = clampEven(tileWidth);
+ tileHeight = clampEven(tileHeight);
+ const durationSec = BUILD_DURATION_MS / 1000;
+
+ const inputArgs = [];
+ const filterParts = [];
+ const layoutParts = [];
+ for (let i = 0; i < normalizedVideos.length; i += 1) {
+ inputArgs.push('-i', normalizedVideos[i].path);
+ const sourceTitle = normalizedVideos[i]?.source?.label || normalizedVideos[i]?.source?.id || `Source ${i + 1}`;
+ filterParts.push(`[${i}:v]${scalePadFilter(tileWidth, tileHeight, sourceTitle)}[v${i}]`);
+ const x = (i % layout.cols) * tileWidth;
+ const y = Math.floor(i / layout.cols) * tileHeight;
+ layoutParts.push(`${x}_${y}`);
+ }
+
+ let audioInputStart = normalizedVideos.length;
+ let sidebarInputIndex = -1;
+ if (includeSidebar) {
+ const selectedRoverIds = usedSources.filter((entry) => entry?.type === 'rover').map((entry) => String(entry.id));
+ const sidebarPath = await renderSidebarVideo({ tmpDir, title: resolvedTitle, durationSec, height: clampEven(tileHeight * layout.rows), windowStartMs: tStart, driverBatteryLines: buildDriverBatterySnapshot(selectedRoverIds), chatEvents: buildChatEventsForWindow(tStart, tEnd) });
+ inputArgs.push('-i', sidebarPath);
+ sidebarInputIndex = normalizedVideos.length;
+ audioInputStart = normalizedVideos.length + 1;
+ }
+
+ for (const audioPath of normalizedAudios) inputArgs.push('-i', audioPath);
+ if (normalizedVideos.length === 1) filterParts.push('[v0]null[vgrid]');
+ else filterParts.push(`${normalizedVideos.map((_, i) => `[v${i}]`).join('')}xstack=inputs=${normalizedVideos.length}:layout=${layoutParts.join('|')}:fill=black[vgrid]`);
+ if (includeSidebar) filterParts.push(`[vgrid][${sidebarInputIndex}:v]hstack=inputs=2[vout]`);
+ else filterParts.push('[vgrid]null[vout]');
+
+ if (normalizedAudios.length) {
+ const audioRefs = normalizedAudios.map((_, idx) => `[${audioInputStart + idx}:a]`).join('');
+ filterParts.push(`${audioRefs}amix=inputs=${normalizedAudios.length}:normalize=0,alimiter=limit=0.9[aout]`);
+ }
+
+ const targetBitrateKbps = Math.max(400, Math.floor((MAX_BYTES * 8) / durationSec / 1000));
+ const outPath = path.join(tmpDir, 'replay.mp4');
+ const args = ['-y','-hide_banner','-loglevel','error',...inputArgs,'-filter_complex',filterParts.join(';'),'-map','[vout]','-c:v','libx264','-preset','veryfast','-pix_fmt','yuv420p','-r',String(TARGET_FPS),'-b:v',`${targetBitrateKbps}k`,'-maxrate',`${Math.floor(targetBitrateKbps * 1.15)}k`,'-bufsize',`${Math.floor(targetBitrateKbps * 2)}k`];
+ if (normalizedAudios.length) args.push('-map', '[aout]', '-c:a', 'aac', '-b:a', '128k');
+ args.push(outPath);
+ await execFileAsync(FFMPEG_BIN, args);
+ const buffer = await fsp.readFile(outPath);
+ return { buffer, usedSources, missingSources, title: resolvedTitle };
+ } finally {
+ try { await fsp.rm(tmpDir, { recursive: true, force: true }); } catch {}
+ }
+ }
+
+ return { buildReplayVideo };
+}
+
+module.exports = { createReplayBuilder };
diff --git a/server/src/services/replaySourceService/index.js b/server/src/services/replayEngineV2/replaySources.js
similarity index 85%
rename from server/src/services/replaySourceService/index.js
rename to server/src/services/replayEngineV2/replaySources.js
index 31dbf54c..9c9f0e1e 100644
--- a/server/src/services/replaySourceService/index.js
+++ b/server/src/services/replayEngineV2/replaySources.js
@@ -1,6 +1,6 @@
-// replay Source Service
-// Purpose: Defines the replay Source Service module and the helpers/state used by this service unit.
-// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary.
+// Replay Source Selection
+// Purpose: Exposes replay source discovery/validation for web and Discord replay requests.
+// Scope: Handles user-visible replay source catalogs and default source selection rules.
const roverManager = require('../roverManager');
const { getRoomCameras } = require('../roomCameraService');
@@ -14,11 +14,13 @@ function getReplaySources(socket = null) {
label: rover.name || rover.id,
color: rover.color || null,
}));
+
const roomSources = getRoomCameras().map((camera) => ({
type: 'room',
id: String(camera.id),
label: camera.name || camera.id,
}));
+
return [...roverSources, ...roomSources];
}
@@ -29,10 +31,8 @@ function normalizeSource(entry) {
if (!type || !id) return null;
return { type, id: String(id) };
}
- if (typeof entry === 'object') {
- if (entry.type && entry.id) {
- return { type: entry.type, id: String(entry.id) };
- }
+ if (typeof entry === 'object' && entry.type && entry.id) {
+ return { type: entry.type, id: String(entry.id) };
}
return null;
}
@@ -42,6 +42,7 @@ function validateSources(list = [], socket = null) {
getReplaySources(socket).forEach((source) => {
allowed.set(`${source.type}:${source.id}`, source);
});
+
const unique = new Map();
(Array.isArray(list) ? list : []).forEach((entry) => {
const normalized = normalizeSource(entry);
@@ -51,6 +52,7 @@ function validateSources(list = [], socket = null) {
if (!source) return;
unique.set(key, { type: source.type, id: source.id, label: source.label });
});
+
return Array.from(unique.values());
}
diff --git a/server/src/services/replayEngineV2/segmentStore.js b/server/src/services/replayEngineV2/segmentStore.js
new file mode 100644
index 00000000..5e89953a
--- /dev/null
+++ b/server/src/services/replayEngineV2/segmentStore.js
@@ -0,0 +1,163 @@
+// Replay Segment Store
+// Purpose: Maintains in-memory segment index and performs file-system refresh/cleanup for replay windows.
+// Scope: Handles segment discovery, overlap queries, and retention cleanup.
+const fsp = require('fs/promises');
+const fs = require('fs');
+const path = require('path');
+const logger = require('../../globals/logger').child('replayEngineV2');
+const { BUFFER_SECONDS, SEGMENT_SECONDS } = require('./constants');
+const { workers, segmentIndex } = require('./state');
+const { sourceKey, sourceDirForKey } = require('./sources');
+
+async function ensureDir(dir) {
+ await fsp.mkdir(dir, { recursive: true });
+}
+
+function createSegmentStore({ getActiveSegmentRoot }) {
+ async function refreshIndexForWorker(worker) {
+ const key = worker.key;
+ const dir = sourceDirForKey(getActiveSegmentRoot(), key);
+ let files;
+ try {
+ files = await fsp.readdir(dir, { withFileTypes: true });
+ } catch {
+ return;
+ }
+
+ const segmentFiles = files.filter((entry) => entry.isFile() && /^seg-\d{6}\.mp4$/.test(entry.name)).map((entry) => entry.name).sort();
+ if (!segmentFiles.length) {
+ segmentIndex.set(key, []);
+ return;
+ }
+
+ const cutoffMs = Date.now() - BUFFER_SECONDS * 1000 - 5000;
+ const entries = [];
+ for (const filename of segmentFiles) {
+ const filePath = path.join(dir, filename);
+ let stat;
+ try { stat = await fsp.stat(filePath); } catch { continue; }
+ if (!stat.isFile() || stat.size < 4096) continue;
+ const endMs = Math.round(stat.mtimeMs);
+ const startMs = Math.round(endMs - SEGMENT_SECONDS * 1000);
+ if (endMs < cutoffMs) continue;
+ entries.push({
+ filePath,
+ startMs,
+ endMs,
+ mtimeMs: stat.mtimeMs,
+ size: stat.size,
+ kind: worker.source.kind,
+ sourceType: worker.source.sourceType,
+ sourceId: worker.source.id,
+ roverId: worker.source.roverId || null,
+ });
+ }
+ segmentIndex.set(key, entries.sort((a, b) => a.startMs - b.startMs));
+ }
+
+ async function refreshSegmentIndex() {
+ for (const worker of Array.from(workers.values())) {
+ await refreshIndexForWorker(worker);
+ }
+ }
+
+ async function cleanupOldFiles() {
+ const cutoff = Date.now() - BUFFER_SECONDS * 1000;
+ const root = getActiveSegmentRoot();
+ try {
+ await ensureDir(root);
+ const dirs = await fsp.readdir(root, { withFileTypes: true });
+ for (const dirent of dirs) {
+ if (!dirent.isDirectory()) continue;
+ const dirPath = path.join(root, dirent.name);
+ let files;
+ try { files = await fsp.readdir(dirPath, { withFileTypes: true }); } catch { continue; }
+ for (const file of files) {
+ if (!file.isFile() || !file.name.endsWith('.mp4')) continue;
+ const filePath = path.join(dirPath, file.name);
+ try {
+ const stat = await fsp.stat(filePath);
+ if (stat.mtimeMs < cutoff) await fsp.unlink(filePath);
+ } catch {}
+ }
+ }
+ } catch (err) {
+ logger.warn('cleanup failed', err.message);
+ }
+ }
+
+ function getVideoEntriesForSource(source) {
+ const key = sourceKey({ sourceType: String(source.type), kind: 'video', id: String(source.id) });
+ return segmentIndex.get(key) || [];
+ }
+
+ function getAudioEntriesForRover(roverId) {
+ const key = sourceKey({ sourceType: 'rover', kind: 'audio', id: `${String(roverId)}-audio` });
+ return segmentIndex.get(key) || [];
+ }
+
+ function overlapping(entries, startMs, endMs) {
+ return entries.filter((entry) => entry.endMs > startMs && entry.startMs < endMs);
+ }
+
+ async function bootstrapIndexFromDisk() {
+ const root = getActiveSegmentRoot();
+ await ensureDir(root);
+ const dirs = await fsp.readdir(root, { withFileTypes: true });
+ for (const dirent of dirs) {
+ if (!dirent.isDirectory()) continue;
+ segmentIndex.set(dirent.name, []);
+ }
+ }
+
+ function getReplayHealthSnapshot({ BUILD_DURATION_MS, roverManager, getRoomCameras }) {
+ const now = Date.now();
+ const neededCount = Math.max(1, Math.ceil(BUILD_DURATION_MS / 1000));
+ const sources = [];
+ let readyCount = 0;
+
+ const replaySources = [];
+ for (const rover of roverManager.getRoster().filter((entry) => roverManager.canReplayRoverId(entry.id))) {
+ replaySources.push({ type: 'rover', id: String(rover.id), label: rover.name || rover.id });
+ }
+ for (const camera of getRoomCameras()) {
+ replaySources.push({ type: 'room', id: String(camera.id), label: camera.name || camera.id });
+ }
+
+ for (const source of replaySources) {
+ const key = sourceKey({ sourceType: source.type, kind: 'video', id: String(source.id) });
+ const dir = sourceDirForKey(getActiveSegmentRoot(), key);
+ let recentCount = 0;
+ let lastSegmentAt = null;
+ try {
+ const files = fs.readdirSync(dir);
+ for (const name of files) {
+ if (!/^seg-\d{6}\.mp4$/.test(name)) continue;
+ const full = path.join(dir, name);
+ let stat;
+ try { stat = fs.statSync(full); } catch { continue; }
+ if (!stat.isFile() || stat.size < 4096) continue;
+ if (stat.mtimeMs > (lastSegmentAt || 0)) lastSegmentAt = stat.mtimeMs;
+ if (now - stat.mtimeMs <= BUFFER_SECONDS * 1000) recentCount += 1;
+ }
+ } catch {}
+ const ready = recentCount >= neededCount;
+ if (ready) readyCount += 1;
+ sources.push({ type: source.type, id: source.id, label: source.label, recentCount, neededCount, lastSegmentAt, ready });
+ }
+
+ return { sources, readyCount, totalCount: sources.length };
+ }
+
+ return {
+ refreshSegmentIndex,
+ cleanupOldFiles,
+ getVideoEntriesForSource,
+ getAudioEntriesForRover,
+ overlapping,
+ bootstrapIndexFromDisk,
+ getReplayHealthSnapshot,
+ };
+}
+
+module.exports = { createSegmentStore };
diff --git a/server/src/services/replayEngineV2/sidebarRenderer.js b/server/src/services/replayEngineV2/sidebarRenderer.js
new file mode 100644
index 00000000..23f08aee
--- /dev/null
+++ b/server/src/services/replayEngineV2/sidebarRenderer.js
@@ -0,0 +1,162 @@
+// Replay Sidebar Renderer
+// Purpose: Renders replay sidebar visuals (title, drivers, chat) and encodes them as a video stream.
+// Scope: Owns SVG/frame synthesis and ffmpeg encoding for optional replay sidebars.
+const path = require('path');
+const sharp = require('sharp');
+const { FFMPEG_BIN, TARGET_FPS, SIDEBAR_WIDTH } = require('./constants');
+
+function createSidebarRenderer({ execFileAsync, ensureDir }) {
+ function escapeXml(text) {
+ return String(text || '').replace(/&/g, '&').replace(//g, '>').replace(/"/g, '"').replace(/'/g, ''');
+ }
+
+ function wrapTextLines(text, maxChars = 24) {
+ const words = String(text || '').split(/\s+/).filter(Boolean);
+ const lines = [];
+ let current = '';
+ for (const word of words) {
+ const candidate = current ? `${current} ${word}` : word;
+ if (candidate.length > maxChars && current) {
+ lines.push(current);
+ current = word;
+ } else {
+ current = candidate;
+ }
+ }
+ if (current) lines.push(current);
+ return lines.slice(0, 4);
+ }
+
+ function hexToRgb(hex) {
+ const value = String(hex || '').trim();
+ const match = /^#([0-9A-Fa-f]{6})$/.exec(value);
+ if (!match) return null;
+ const raw = match[1];
+ return { r: parseInt(raw.slice(0, 2), 16), g: parseInt(raw.slice(2, 4), 16), b: parseInt(raw.slice(4, 6), 16) };
+ }
+
+ function roleColor(role = '') {
+ switch (String(role)) {
+ case 'admin':
+ case 'lockdown':
+ case 'lockdown-admin':
+ return '#FCD34D';
+ case 'spectator':
+ return '#94A3B8';
+ default:
+ return '#7DD3FC';
+ }
+ }
+
+ function renderSidebarSvg({ width, height, title, driverBatteryLines, chatLines }) {
+ const textCols = Math.max(24, Math.floor((width - 16) / 6));
+ const titleParts = wrapTextLines(title, Math.max(26, textCols)).slice(0, 4).map((line) => escapeXml(line));
+ const statLines = driverBatteryLines.slice(0, 8).map((line) => escapeXml(line));
+ const normalizedChats = chatLines.slice(-12).map((entry) => {
+ const wrapped = wrapTextLines(entry.text || '', Math.max(24, textCols)).slice(0, 4).map((line) => escapeXml(line));
+ const nick = escapeXml(entry.nickname || 'user');
+ const roverId = escapeXml(entry.roverId || '');
+ const roverRgb = hexToRgb(entry.roverColor || '');
+ const roverBadgeBg = roverRgb ? `rgba(${roverRgb.r},${roverRgb.g},${roverRgb.b},0.18)` : 'rgba(30,41,59,0.70)';
+ const roverBadgeBorder = roverRgb ? `rgba(${roverRgb.r},${roverRgb.g},${roverRgb.b},0.60)` : 'rgba(71,85,105,0.75)';
+ return { nick, wrapped, nameColor: roleColor(entry.role), fromDiscord: Boolean(entry.fromDiscord), roverId, roverBadgeBg, roverBadgeBorder, bubbleTone: entry.fromDiscord ? 'discordBubble' : 'chatBubble' };
+ });
+
+ const shapes = [];
+ const textRows = [];
+ const pad = 3;
+ const cardX = pad;
+ const cardW = width - pad * 2;
+ let y = pad;
+
+ const titleCardH = Math.max(26, pad * 2 + titleParts.length * 15);
+ shapes.push(``);
+ let ty = y + pad + 11;
+ for (const part of titleParts) {
+ textRows.push(`${part}`);
+ ty += 15;
+ }
+ y += titleCardH + pad;
+
+ const driverLines = statLines.length ? statLines : ['No active drivers'];
+ const driversCardH = pad * 2 + 11 + driverLines.length * 13;
+ shapes.push(``);
+ textRows.push(`Drivers`);
+ let dy = y + pad + 19;
+ for (const line of driverLines) {
+ textRows.push(`${line}`);
+ dy += 13;
+ }
+ y += driversCardH + pad;
+
+ const chatCardH = Math.max(80, height - y - pad);
+ shapes.push(``);
+ textRows.push(`Chat`);
+
+ let cy = y + pad + 12;
+ const bubbleX = cardX + pad;
+ const bubbleW = cardW - pad * 2;
+ if (!normalizedChats.length) {
+ textRows.push(`No chat in replay window`);
+ } else {
+ for (let i = 0; i < normalizedChats.length; i += 1) {
+ const block = normalizedChats[i];
+ const nameW = Math.min(72, block.nick.length * 5.2);
+ const badgeW = block.roverId ? Math.min(46, Math.max(18, block.roverId.length * 5 + 6)) : 0;
+ const badgeGap = block.roverId ? 3 : 0;
+ const prefixChars = Math.ceil((nameW + badgeW + badgeGap + 10) / 5.8);
+ const firstLineRaw = String(block.wrapped[0] || '').trim();
+ const remainingRaw = block.wrapped.slice(firstLineRaw ? 1 : 0).map((line) => String(line || '').trim()).filter(Boolean);
+ const fullText = (firstLineRaw ? [firstLineRaw, ...remainingRaw] : remainingRaw).join(' ');
+ const inlineWrapped = wrapTextLines(fullText, Math.max(16, textCols - prefixChars)).slice(0, 4).map((line) => escapeXml(line));
+ const bubbleH = pad * 2 + Math.max(1, inlineWrapped.length) * 12;
+ if (cy + bubbleH + pad > y + chatCardH - pad) break;
+ shapes.push(``);
+ const nameX = bubbleX + pad;
+ const textStartX = nameX + nameW + 3 + (block.roverId ? badgeW + badgeGap : 0);
+ textRows.push(`${block.nick}`);
+ if (block.fromDiscord) textRows.push(`◈`);
+ if (block.roverId) {
+ const badgeX = nameX + nameW + 3;
+ const badgeTextX = badgeX + 3;
+ shapes.push(``);
+ textRows.push(`${block.roverId}`);
+ }
+ let by = cy + pad + 8;
+ for (let lineIdx = 0; lineIdx < inlineWrapped.length; lineIdx += 1) {
+ const line = inlineWrapped[lineIdx];
+ const lineX = lineIdx === 0 ? textStartX : bubbleX + pad;
+ textRows.push(`${line}`);
+ by += 12;
+ }
+ cy += bubbleH + pad;
+ }
+ }
+
+ return `\n`;
+ }
+
+ async function renderSidebarVideo({ tmpDir, title, durationSec, height, windowStartMs, driverBatteryLines = [], chatEvents = [] }) {
+ const framesDir = path.join(tmpDir, 'sidebar-frames');
+ const sidebarPath = path.join(tmpDir, 'sidebar.mp4');
+ await ensureDir(framesDir);
+ const secondCount = Math.max(1, Math.ceil(durationSec));
+ for (let second = 0; second < secondCount; second += 1) {
+ const sliceEndMs = windowStartMs + (second + 1) * 1000;
+ const visibleChat = chatEvents.filter((entry) => entry.ts <= sliceEndMs).slice(-10);
+ const svg = renderSidebarSvg({ width: SIDEBAR_WIDTH, height, title, driverBatteryLines, chatLines: visibleChat });
+ const framePath = path.join(framesDir, `frame-${String(second + 1).padStart(4, '0')}.png`);
+ await sharp(Buffer.from(svg, 'utf8')).png().toFile(framePath);
+ }
+
+ await execFileAsync(FFMPEG_BIN, [
+ '-y','-hide_banner','-loglevel','error','-framerate','1','-i',path.join(framesDir, 'frame-%04d.png'),
+ '-vf',`fps=${TARGET_FPS},format=yuv420p`,'-t',durationSec.toFixed(3),'-c:v','libx264','-preset','veryfast','-pix_fmt','yuv420p',sidebarPath,
+ ]);
+ return sidebarPath;
+ }
+
+ return { renderSidebarVideo };
+}
+
+module.exports = { createSidebarRenderer };
diff --git a/server/src/services/replayEngineV2/socketHooks.js b/server/src/services/replayEngineV2/socketHooks.js
new file mode 100644
index 00000000..f753f6c2
--- /dev/null
+++ b/server/src/services/replayEngineV2/socketHooks.js
@@ -0,0 +1,81 @@
+// Replay Socket Hooks
+// Purpose: Registers web socket replay-trigger handler that publishes replay requests.
+// Scope: Applies replay mode/cooldown/source validation for socket-triggered replay requests.
+const io = require('../../globals/io');
+const logger = require('../../globals/logger').child('replaySocket');
+const { getMode, MODES } = require('../modeManager');
+const { publishEvent } = require('../eventBus');
+const assignmentService = require('../assignmentService');
+const { getNickname } = require('../nicknameService');
+const { loadConfig } = require('../../helpers/configLoader');
+
+const config = loadConfig();
+const discordConfig = config.discord || {};
+
+function buildRequesterLabel(socket) {
+ return getNickname(socket) || socket?.data?.user?.username || socket?.id || 'unknown';
+}
+
+function normalizeReplayTitle(value) {
+ if (typeof value !== 'string') return '';
+ return value.trim().slice(0, 120);
+}
+
+function normalizeIncludeSidebar(value) {
+ if (typeof value === 'boolean') return value;
+ return true;
+}
+
+function registerReplaySocketHooks({ tryTriggerReplay, validateSources, getDefaultWebSources }) {
+ io.on('connection', (socket) => {
+ socket.on('replay:trigger', (payload = {}, cb = () => {}) => {
+ if (getMode() === MODES.LOCKDOWN) {
+ cb({ error: 'Replay disabled in lockdown', state: null });
+ return;
+ }
+ const channelId = discordConfig?.channels?.replay || null;
+ if (!channelId) {
+ cb({ error: 'Replay channel not configured', state: null });
+ return;
+ }
+ const requestedSources = Array.isArray(payload?.sources) ? payload.sources : null;
+ let sources = requestedSources ? validateSources(requestedSources, socket) : [];
+ if (!sources.length) {
+ const assignment = assignmentService.describeAssignment(socket.id);
+ sources = getDefaultWebSources(assignment, socket);
+ }
+ if (!sources.length) {
+ cb({ error: 'No replay sources selected', state: null });
+ return;
+ }
+
+ const requester = buildRequesterLabel(socket);
+ const title = normalizeReplayTitle(payload?.title);
+ const includeSidebar = normalizeIncludeSidebar(payload?.includeSidebar);
+ const attempt = tryTriggerReplay({ by: { source: 'web', requester } });
+ if (!attempt.ok) {
+ cb({ error: 'Replay cooldown active', remainingMs: attempt.remainingMs, state: attempt.state });
+ return;
+ }
+
+ publishEvent({
+ source: 'replaySocket',
+ type: 'replay.requested',
+ payload: {
+ channelId,
+ requester,
+ title,
+ includeSidebar,
+ sources,
+ requestedBy: { socketId: socket.id },
+ },
+ });
+ logger.info('Replay requested via web', { socketId: socket.id });
+ cb({ success: true, state: attempt.state });
+ });
+ });
+}
+
+module.exports = {
+ registerReplaySocketHooks,
+};
diff --git a/server/src/services/replayEngineV2/sources.js b/server/src/services/replayEngineV2/sources.js
new file mode 100644
index 00000000..f139b668
--- /dev/null
+++ b/server/src/services/replayEngineV2/sources.js
@@ -0,0 +1,71 @@
+// Replay Source Catalog
+// Purpose: Resolves replay-capable media sources and ffmpeg worker arguments.
+// Scope: Converts live rover/room state into stable source descriptors and stream worker config.
+const path = require('path');
+const roverManager = require('../roverManager');
+const { getRoomCameras } = require('../roomCameraService');
+const { FFMPEG_BIN, SEGMENT_SECONDS, TARGET_FPS } = require('./constants');
+
+function sourceKey(source) {
+ return `${source.sourceType}__${source.kind}__${source.id}`;
+}
+
+function sourceDirForKey(activeSegmentRoot, key) {
+ return path.join(activeSegmentRoot, key);
+}
+
+function toSrtReadPath(streamId) {
+ return `srt://127.0.0.1:9000?streamid=read:${encodeURIComponent(streamId)}`;
+}
+
+function getRoomCameraStream(camera) {
+ if (camera?.streamUrl) return camera.streamUrl;
+ const url = String(camera?.url || '');
+ if (url.includes('.mjpg') || url.includes('mjpeg') || url.includes('stream')) return url;
+ return null;
+}
+
+function listDesiredSources() {
+ const sources = [];
+ for (const rover of roverManager.getRoster()) {
+ if (!roverManager.canReplayRoverId(rover.id)) continue;
+ const roverId = String(rover.id);
+ sources.push({ id: roverId, sourceType: 'rover', kind: 'video', label: rover.name || roverId, inputUrl: toSrtReadPath(roverId) });
+ if (rover?.media?.audioPublishUrl) {
+ sources.push({ id: `${roverId}-audio`, sourceType: 'rover', roverId, kind: 'audio', label: `${rover.name || roverId} audio`, inputUrl: toSrtReadPath(`${roverId}-audio`) });
+ }
+ }
+ for (const camera of getRoomCameras()) {
+ const streamUrl = getRoomCameraStream(camera);
+ if (!streamUrl) continue;
+ sources.push({ id: String(camera.id), sourceType: 'room', kind: 'video', label: camera.name || camera.id, inputUrl: streamUrl });
+ }
+ return sources;
+}
+
+function buildWorkerArgs(activeSegmentRoot, source) {
+ const dir = sourceDirForKey(activeSegmentRoot, sourceKey(source));
+ const pattern = path.join(dir, 'seg-%06d.mp4');
+ const common = ['-hide_banner', '-loglevel', 'warning', '-y', '-fflags', '+genpts', '-use_wallclock_as_timestamps', '1', '-i', source.inputUrl];
+
+ if (source.kind === 'audio') {
+ return [
+ ...common,
+ '-vn','-ac','1','-ar','48000','-af','aresample=async=1:first_pts=0:min_hard_comp=0.100000,asetpts=N/SR/TB',
+ '-c:a','aac','-b:a','96k','-f','segment','-segment_time',String(SEGMENT_SECONDS),'-segment_atclocktime','1','-reset_timestamps','1',pattern,
+ ];
+ }
+
+ return [
+ ...common,
+ '-an','-vf',`fps=${TARGET_FPS}`,'-c:v','libx264','-preset','veryfast','-tune','zerolatency','-pix_fmt','yuv420p','-g',String(TARGET_FPS * SEGMENT_SECONDS),'-keyint_min',String(TARGET_FPS * SEGMENT_SECONDS),'-sc_threshold','0','-f','segment','-segment_time',String(SEGMENT_SECONDS),'-segment_atclocktime','1','-reset_timestamps','1',pattern,
+ ];
+}
+
+module.exports = {
+ FFMPEG_BIN,
+ sourceKey,
+ sourceDirForKey,
+ listDesiredSources,
+ buildWorkerArgs,
+};
diff --git a/server/src/services/replayEngineV2/state.js b/server/src/services/replayEngineV2/state.js
new file mode 100644
index 00000000..d8dde379
--- /dev/null
+++ b/server/src/services/replayEngineV2/state.js
@@ -0,0 +1,23 @@
+// Replay Engine State
+// Purpose: Holds mutable runtime state for workers, segment index, and scheduler lifecycle.
+// Scope: Provides shared process-local state without embedding behavior logic.
+const EventEmitter = require('events');
+
+const events = new EventEmitter();
+const workers = new Map();
+const pendingWorkerStarts = new Set();
+const segmentIndex = new Map();
+
+const runtime = {
+ cleanupTimer: null,
+ activeSegmentRoot: null,
+ tickInFlight: false,
+};
+
+module.exports = {
+ events,
+ workers,
+ pendingWorkerStarts,
+ segmentIndex,
+ runtime,
+};
diff --git a/server/src/services/replayEngineV2/workerManager.js b/server/src/services/replayEngineV2/workerManager.js
new file mode 100644
index 00000000..a3298fea
--- /dev/null
+++ b/server/src/services/replayEngineV2/workerManager.js
@@ -0,0 +1,74 @@
+// Replay Worker Manager
+// Purpose: Starts/stops ffmpeg segment workers and keeps active worker set aligned with desired sources.
+// Scope: Owns worker process lifecycle and restart behavior for replay segment capture.
+const { spawn } = require('child_process');
+const fsp = require('fs/promises');
+const logger = require('../../globals/logger').child('replayEngineV2');
+const { FFMPEG_BIN } = require('./constants');
+const { workers, pendingWorkerStarts } = require('./state');
+const { sourceKey, sourceDirForKey, listDesiredSources, buildWorkerArgs } = require('./sources');
+
+async function ensureDir(dir) {
+ await fsp.mkdir(dir, { recursive: true });
+}
+
+function createWorkerManager({ getActiveSegmentRoot }) {
+ function startWorker(source) {
+ const key = sourceKey(source);
+ if (workers.has(key) || pendingWorkerStarts.has(key)) return;
+ pendingWorkerStarts.add(key);
+ const args = buildWorkerArgs(getActiveSegmentRoot(), source);
+ const dir = sourceDirForKey(getActiveSegmentRoot(), key);
+ ensureDir(dir)
+ .then(() => {
+ if (workers.has(key)) {
+ pendingWorkerStarts.delete(key);
+ return;
+ }
+ const proc = spawn(FFMPEG_BIN, args, { stdio: ['ignore', 'ignore', 'pipe'] });
+ workers.set(key, { key, source, proc });
+ pendingWorkerStarts.delete(key);
+ proc.stderr.on('data', (chunk) => {
+ const text = String(chunk || '').trim();
+ if (!text) return;
+ logger.warn('worker stderr', { key, text: text.slice(0, 500) });
+ });
+ proc.on('exit', (code, signal) => {
+ const current = workers.get(key);
+ if (current?.proc === proc) workers.delete(key);
+ logger.warn('worker exited', { key, code, signal });
+ setTimeout(() => {
+ const desired = listDesiredSources().find((entry) => sourceKey(entry) === key);
+ if (desired && !workers.has(key)) startWorker(desired);
+ }, 1500);
+ });
+ })
+ .catch((err) => {
+ pendingWorkerStarts.delete(key);
+ logger.warn('failed to start worker', { key, error: err.message });
+ });
+ }
+
+ function stopWorker(key) {
+ const worker = workers.get(key);
+ if (!worker) return;
+ try { worker.proc.kill('SIGTERM'); } catch {}
+ workers.delete(key);
+ }
+
+ async function syncWorkers() {
+ const desired = listDesiredSources();
+ const desiredKeys = new Set(desired.map(sourceKey));
+ for (const source of desired) {
+ const key = sourceKey(source);
+ if (!workers.has(key)) startWorker(source);
+ }
+ for (const key of Array.from(workers.keys())) {
+ if (!desiredKeys.has(key)) stopWorker(key);
+ }
+ }
+
+ return { startWorker, stopWorker, syncWorkers };
+}
+
+module.exports = { createWorkerManager };
diff --git a/server/src/services/replayService/index.js b/server/src/services/replayService/index.js
deleted file mode 100644
index 21df8a08..00000000
--- a/server/src/services/replayService/index.js
+++ /dev/null
@@ -1,45 +0,0 @@
-// replay Service
-// Purpose: Defines the replay Service module and the helpers/state used by this service unit.
-// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary.
-const EventEmitter = require('events');
-
-const COOLDOWN_MS = 10 * 1000;
-
-const events = new EventEmitter();
-let lastTriggeredAt = null;
-let lastTriggeredBy = null;
-
-function getRemainingMs(now = Date.now()) {
- if (!lastTriggeredAt) return 0;
- const elapsed = now - lastTriggeredAt;
- return Math.max(0, COOLDOWN_MS - elapsed);
-}
-
-function getState() {
- const remainingMs = getRemainingMs();
- return {
- cooldownMs: COOLDOWN_MS,
- lastTriggeredAt,
- lastTriggeredBy,
- remainingMs,
- available: remainingMs === 0,
- };
-}
-
-function tryTrigger(by = null) {
- const remainingMs = getRemainingMs();
- if (remainingMs > 0) {
- return { ok: false, remainingMs, state: getState() };
- }
- lastTriggeredAt = Date.now();
- lastTriggeredBy = by || null;
- const state = getState();
- events.emit('update', { state, by: lastTriggeredBy });
- return { ok: true, state };
-}
-
-module.exports = {
- tryTriggerReplay: tryTrigger,
- getReplayState: getState,
- replayEvents: events,
-};
diff --git a/server/src/services/replaySocketService/index.js b/server/src/services/replaySocketService/index.js
deleted file mode 100644
index 5acfa0d5..00000000
--- a/server/src/services/replaySocketService/index.js
+++ /dev/null
@@ -1,75 +0,0 @@
-// replay Socket Service
-// Purpose: Defines the replay Socket Service module and the helpers/state used by this service unit.
-// Scope: Keeps runtime behavior unchanged while isolating responsibilities into a clear module boundary.
-const io = require('../../globals/io');
-const logger = require('../../globals/logger').child('replaySocket');
-const { getMode, MODES } = require('../modeManager');
-const { publishEvent } = require('../eventBus');
-const { tryTriggerReplay } = require('../replayService');
-const { validateSources, getDefaultWebSources } = require('../replaySourceService');
-const assignmentService = require('../assignmentService');
-const { getNickname } = require('../nicknameService');
-const { loadConfig } = require('../../helpers/configLoader');
-
-const config = loadConfig();
-const discordConfig = config.discord || {};
-
-function buildRequesterLabel(socket) {
- return getNickname(socket) || socket?.data?.user?.username || socket?.id || 'unknown';
-}
-
-function normalizeReplayTitle(value) {
- if (typeof value !== 'string') return '';
- return value.trim().slice(0, 120);
-}
-
-function normalizeIncludeSidebar(value) {
- if (typeof value === 'boolean') return value;
- return true;
-}
-
-io.on('connection', (socket) => {
- socket.on('replay:trigger', (payload = {}, cb = () => {}) => {
- if (getMode() === MODES.LOCKDOWN) {
- cb({ error: 'Replay disabled in lockdown', state: null });
- return;
- }
- const channelId = discordConfig?.channels?.replay || null;
- if (!channelId) {
- cb({ error: 'Replay channel not configured', state: null });
- return;
- }
- const requestedSources = Array.isArray(payload?.sources) ? payload.sources : null;
- let sources = requestedSources ? validateSources(requestedSources, socket) : [];
- if (!sources.length) {
- const assignment = assignmentService.describeAssignment(socket.id);
- sources = getDefaultWebSources(assignment, socket);
- }
- if (!sources.length) {
- cb({ error: 'No replay sources selected', state: null });
- return;
- }
- const requester = buildRequesterLabel(socket);
- const title = normalizeReplayTitle(payload?.title);
- const includeSidebar = normalizeIncludeSidebar(payload?.includeSidebar);
- const attempt = tryTriggerReplay({ by: { source: 'web', requester } });
- if (!attempt.ok) {
- cb({ error: 'Replay cooldown active', remainingMs: attempt.remainingMs, state: attempt.state });
- return;
- }
- publishEvent({
- source: 'replaySocket',
- type: 'replay.requested',
- payload: {
- channelId,
- requester,
- title,
- includeSidebar,
- sources,
- requestedBy: { socketId: socket.id },
- },
- });
- logger.info('Replay requested via web', { socketId: socket.id });
- cb({ success: true, state: attempt.state });
- });
-});
diff --git a/server/src/services/sessionService/index.js b/server/src/services/sessionService/index.js
index 9ddef770..b971808e 100644
--- a/server/src/services/sessionService/index.js
+++ b/server/src/services/sessionService/index.js
@@ -23,8 +23,7 @@ const {
getStateForSocket: getPrivateRoverAccessStateForSocket,
requestEvents: privateRoverAccessRequestEvents,
} = require('../privateRoverAccessRequestService');
-const { getReplayState, replayEvents } = require('../replayService');
-const { getReplaySources } = require('../replaySourceService');
+const { getReplayState, replayEvents, getReplaySources } = require('../replayEngineV2');
const { getHealthSnapshot } = require('../healthService');
const { getCommunityGoal } = require('../communityGoalService');
const { getAdminReason } = require('../adminReasonService');