SLOP ANALYTICS!!

This commit is contained in:
legop3
2026-07-19 22:43:56 -04:00
parent b261a4bfa2
commit 90d4f778a5
32 changed files with 2552 additions and 19 deletions
+6
View File
@@ -52,6 +52,7 @@ function buildFeatureFlags(config = loadConfig()) {
const interInstanceConfig = config.interInstance || {};
const ptzCameraConfig = config.ptzCamera || {};
const discordConfig = config.discord || {};
const fleetReportsConfig = config.fleetReports || {};
const homeAssistant = Boolean(
asBoolean(homeAssistantConfig.enabled) &&
asTrimmedString(homeAssistantConfig.url) &&
@@ -100,6 +101,11 @@ function buildFeatureFlags(config = loadConfig()) {
to run without the integration.
*/
discord: Boolean(asBoolean(discordConfig.enabled) && asTrimmedString(discordConfig.token)),
// Fleet reports are deliberately controlled by one explicit server switch.
// Storage contents, Discord availability, or historical database files must
// never cause the reporting UI to appear on an installation that has not
// opted into the collector.
fleetReports: asBoolean(fleetReportsConfig.enabled),
};
}
@@ -2,6 +2,7 @@
// Purpose: Defines the command 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 { v4: uuidv4 } = require('uuid');
const EventEmitter = require('events');
const io = require('../../globals/io');
const roverManager = require('../roverManager');
const { isAdmin, isLockdownAdmin } = require('../roleService');
@@ -11,6 +12,12 @@ const { isHeadlightBlocked } = require('../../rewards/definitions/darkness');
const homeAssistantService = require('../homeAssistantService');
const overcurrentProtectionService = require('../overcurrentProtectionService');
// Command observations are intentionally separate from the global event bus.
// Drive and motor commands can run at control-loop frequency, and publishing
// every packet onto the logging event bus would manufacture noise. Optional
// observers can aggregate this emitter without changing command delivery.
const commandEvents = new EventEmitter();
const pendingCommands = new Map(); // id -> { roverId }
const lastDriveActivity = new Map(); // roverId -> { ts, socketId, direction, speed, isAdmin }
const driveCooldowns = new Map(); // roverId -> blockedUntil
@@ -130,6 +137,16 @@ function handleAck(msg) {
status: msg.status || 'ok',
error: msg.error,
});
commandEvents.emit('observation', {
ts: Date.now(),
roverId: pending.roverId,
type: pending.type,
commandId: msg.id,
outcome: msg.error ? 'failed' : 'acknowledged',
latencyMs: Date.now() - pending.ts,
status: msg.status || 'ok',
error: msg.error || null,
});
}
function issueUpdateToAllRovers() {
@@ -237,6 +254,7 @@ module.exports = {
handleAck,
getRecentDriveActivity,
setDriveCooldown,
commandEvents,
};
io.on('connection', (socket) => {
@@ -331,6 +349,17 @@ io.on('connection', (socket) => {
});
}
const id = issueCommand(roverId, { type, ...payload });
commandEvents.emit('observation', {
ts: Date.now(),
roverId: String(roverId),
type,
commandId: id,
outcome: 'issued',
socketId: socket.id,
// Payloads are omitted deliberately: raw OI, TTS, and maintenance
// commands can carry arbitrary content. Their structured type/outcome
// supplies analytics without accidentally persisting secret material.
});
logger.info('Queued command', socket.id, roverId, type);
if (shouldRecordTurnActivity(type, payload)) {
try {
@@ -343,6 +372,14 @@ io.on('connection', (socket) => {
reply({ id });
} catch (err) {
logger.warn('Command rejected', socket.id, err.message);
commandEvents.emit('observation', {
ts: Date.now(),
roverId: roverId ? String(roverId) : null,
type: type || 'unknown',
outcome: 'rejected',
socketId: socket.id,
error: err.message,
});
reply({ error: err.message });
}
}
@@ -0,0 +1,147 @@
// Discord Fleet Daily Reports
// Purpose: Schedules and delivers completed-day fleet summaries to the existing admin alert channel.
// Scope: Discord owns timing/formatting/delivery; the fleet service owns evidence, analysis, and durable delivery state.
const { DateTime } = require('luxon');
const { AttachmentBuilder, EmbedBuilder } = require('discord.js');
function parseSendTime(value) {
const match = /^(\d{1,2}):(\d{2})$/.exec(String(value || '').trim());
if (!match) return { hour: 8, minute: 0 };
return {
hour: Math.max(0, Math.min(23, Number(match[1]))),
minute: Math.max(0, Math.min(59, Number(match[2]))),
};
}
function nextRunAt({ zone, hour, minute }) {
const now = DateTime.now().setZone(zone);
let next = now.set({ hour, minute, second: 0, millisecond: 0 });
if (next <= now) next = next.plus({ days: 1 });
return next;
}
function formatNumber(value, digits = 1) {
return Number(value || 0).toLocaleString(undefined, { maximumFractionDigits: digits });
}
function createFleetDailyReports({ logger, discordConfig, fleetConfig, fleetReportService, roverManager, sendToChannel }) {
let timer = null;
const reportConfig = fleetConfig?.discord || {};
const enabled = fleetReportService?.enabled && reportConfig.enabled !== false;
const channelId = discordConfig?.channels?.adminAlerts;
const zone = String(reportConfig.timezone || 'America/New_York');
const { hour, minute } = parseSendTime(reportConfig.sendAt);
function publicRoverIds() {
// Discord's shared admin-alert channel does not provide a per-viewer socket
// against which private-rover grants can be checked. Excluding private
// rovers here preserves the existing privacy boundary instead of assuming
// every channel reader has every private grant.
return roverManager.getRoster()
.filter((rover) => !rover?.private?.enabled)
.map((rover) => String(rover.id));
}
function completedDayRange() {
const end = DateTime.now().setZone(zone).startOf('day');
const start = end.minus({ days: 1 });
return {
reportDate: start.toISODate(),
since: start.toMillis(),
until: end.toMillis(),
};
}
function buildEmbed(reportDate, report) {
const totals = report.totals;
const attention = report.findings.slice(0, 12).map((finding) =>
`${finding.roverId ? `${finding.roverId}: ` : ''}${finding.title} (${finding.severity}, ${finding.confidence} confidence)`,
).join('\n') || 'No report findings.';
const roverLines = report.rovers.map((rover) =>
`${rover.name}: ${formatNumber(rover.dischargedMah)} mAh used, ${formatNumber(rover.chargedMah)} mAh charged, ${formatNumber(rover.sampleCount, 0)} samples, ${formatNumber(rover.gapCount, 0)} gaps`,
).join('\n') || 'No public rover telemetry.';
return new EmbedBuilder()
.setTitle(`Daily fleet report — ${reportDate}`)
.setColor(totals.criticalFindingCount ? 0xe53935 : totals.warningFindingCount ? 0xf0b651 : 0x4caf50)
.addFields(
{
name: 'Fleet totals',
value: `${totals.onlineRoverCount}/${totals.roverCount} online · ${formatNumber(totals.sampleCount, 0)} samples · ${formatNumber(totals.dischargedMah)} mAh used · ${formatNumber(totals.chargedMah)} mAh charged · ${formatNumber(totals.telemetryGapCount, 0)} gaps`,
},
{ name: 'Needs attention', value: attention.slice(0, 1024) },
{ name: 'Rovers', value: roverLines.slice(0, 1024) },
)
.setFooter({ text: 'Detailed read-only evidence is available on the server reports page.' });
}
async function deliverPreviousDay() {
if (!enabled || !channelId) return;
const range = completedDayRange();
const existing = fleetReportService.storage.getDailyReport(range.reportDate);
if (existing?.discordDeliveredAt) return;
const report = existing?.report || fleetReportService.getDailyReport({
since: range.since,
until: range.until,
roverIds: publicRoverIds(),
});
if (!report) return;
// Lockdown-only records and raw event payloads do not belong in a shared
// Discord attachment. The interactive server UI applies per-socket access
// and remains the place for complete event evidence.
const attachmentReport = {
...report,
events: report.events.filter((event) => event.visibility !== 'lockdown').map(({ payload, ...event }) => ({
...event,
payload,
})),
};
fleetReportService.storage.saveDailyReport(range.reportDate, attachmentReport);
const attachment = new AttachmentBuilder(
Buffer.from(JSON.stringify(attachmentReport, null, 2)),
{ name: `fleet-report-${range.reportDate}.json` },
);
const sent = await sendToChannel(
channelId,
`Daily fleet report for ${range.reportDate}`,
{ embeds: [buildEmbed(range.reportDate, report)], files: [attachment] },
{ parse: [] },
);
if (sent) {
fleetReportService.storage.markDailyReportDelivery(range.reportDate, { deliveredAt: Date.now(), error: null });
} else {
fleetReportService.storage.markDailyReportDelivery(range.reportDate, { error: 'Discord delivery returned no message' });
}
}
function scheduleNext() {
if (!enabled || !channelId) return;
const next = nextRunAt({ zone, hour, minute });
const delay = Math.max(1000, next.toMillis() - Date.now());
timer = setTimeout(async () => {
try {
await deliverPreviousDay();
} catch (err) {
logger.warn('Daily fleet report delivery failed', { error: err.message });
} finally {
scheduleNext();
}
}, delay);
timer.unref?.();
logger.info('Scheduled daily fleet report', { nextRunAt: next.toISO(), channelId });
}
function start() {
scheduleNext();
}
function stop() {
if (timer) clearTimeout(timer);
timer = null;
}
return { start, stop, deliverPreviousDay };
}
module.exports = {
createFleetDailyReports,
};
@@ -56,6 +56,8 @@ const { createChannelIO } = require('./channelIO');
const { createCommandHandlers } = require('../operatorCommandService');
const { createDiscordTransportHandlers, createDiscordCommandRequest } = require('./commandAdapter');
const { createIntegrations } = require('./integrations');
const { createFleetDailyReports } = require('./fleetDailyReports');
const fleetReportService = require('../fleetReportService');
const { registerPreferredDeliveryProvider } = require('../replayDeliveryService');
const {
DEFAULT_ALLOWED_MENTIONS,
@@ -347,6 +349,17 @@ client.on('messageCreate', async (message) => {
client.once('ready', () => {
logger.info('Discord bot logged in', { tag: client.user?.tag });
presence.schedulePresenceRotation();
// Discord is only a delivery consumer. Starting its scheduler after the bot
// is ready avoids failed sends during login while the collector continues to
// operate independently of Discord availability.
createFleetDailyReports({
logger,
discordConfig,
fleetConfig: config.fleetReports || {},
fleetReportService,
roverManager,
sendToChannel: channelIO.sendToChannel,
}).start();
});
client.login(discordConfig.token).catch((err) => {
@@ -9,7 +9,7 @@ const { renderIndexHtml, renderOgImage } = require('../embedService');
entry point. Including /ptz here lets direct loads and browser refreshes
receive the same rendered index document as navigation from the driver page.
*/
app.get(['/', '/spectate', '/mini', '/display', '/scanner', '/database', '/ptz'], async (req, res) => {
app.get(['/', '/spectate', '/mini', '/display', '/scanner', '/database', '/ptz', '/reports'], async (req, res) => {
try {
const html = await renderIndexHtml(req);
res.type('html').send(html);
@@ -0,0 +1,581 @@
// Fleet Report Collector
// Purpose: Converts existing server events and high-rate rover sensor frames into bounded historical evidence.
// Scope: Performs passive normalization, battery-current integration, minute aggregation, and battery-session classification.
const crypto = require('crypto');
const MINUTE_MS = 60 * 1000;
const FULL_WAIT_MS = 5 * 60 * 1000;
const SESSION_KIND_CONFIRM_SAMPLES = 3;
function finite(value) {
const number = Number(value);
return Number.isFinite(number) ? number : null;
}
function minimum(previous, value) {
if (value == null) return previous;
return previous == null ? value : Math.min(previous, value);
}
function maximum(previous, value) {
if (value == null) return previous;
return previous == null ? value : Math.max(previous, value);
}
function eventRoverId(event) {
const payload = event?.payload || {};
// Generic payload `id` fields are commonly message, request, or job IDs and
// must not be mistaken for rover identities. Producers use roverId (or an
// explicit rover object) whenever the existing privacy resolver should scope
// an event to a physical rover.
return payload.roverId || payload.rover?.id || null;
}
function inferVisibility(event) {
const payload = event?.payload || {};
// Producers that know a stricter visibility scope may attach it explicitly.
// Otherwise rover-scoped events are filtered later against the same visible
// roster that drives the live UI, while verification/auth details retain the
// existing lockdown-only boundary.
if (payload.visibility) return String(payload.visibility);
if (event?.source === 'verification' || event?.source === 'identity' || event?.source === 'auth') {
return 'lockdown';
}
return eventRoverId(event) ? 'rover' : 'global';
}
function inferSeverity(type = '') {
const value = String(type).toLowerCase();
if (/fault|critical|urgent|failed|failure/.test(value)) return 'critical';
if (/warn|offline|rejected|stopped|removed|denied/.test(value)) return 'warning';
if (/started|completed|online|resolved|updated/.test(value)) return 'notice';
return 'informational';
}
function batteryKey(roverId) {
// Until an admin registers a physical battery, the stable fallback keeps all
// observations attached to the rover without pretending the hardware has a
// serial number exposed by OI.
return `unregistered:${roverId}`;
}
function makeMinute(roverId, now) {
return {
roverId,
bucketTs: Math.floor(now / MINUTE_MS) * MINUTE_MS,
sampleCount: 0,
coverageMs: 0,
gapCount: 0,
chargedMah: 0,
dischargedMah: 0,
minVoltageMv: null,
maxVoltageMv: null,
voltageTotal: 0,
voltageCount: 0,
minCurrentMa: null,
maxCurrentMa: null,
currentTotal: 0,
currentCount: 0,
minTemperatureC: null,
maxTemperatureC: null,
temperatureTotal: 0,
temperatureCount: 0,
minChargeMah: null,
maxChargeMah: null,
lastChargeMah: null,
reportedCapacityMah: null,
dockedSamples: 0,
chargingSamples: 0,
commandCount: 0,
driveCommandCount: 0,
rejectedCommandCount: 0,
distanceMm: 0,
bumpCount: 0,
cliffCount: 0,
wheelDropCount: 0,
virtualWallCount: 0,
overcurrentEpisodeCount: 0,
};
}
function persistedMinute(minute) {
return {
roverId: minute.roverId,
bucketTs: minute.bucketTs,
sampleCount: minute.sampleCount,
coverageMs: Math.round(minute.coverageMs),
gapCount: minute.gapCount,
chargedMah: minute.chargedMah,
dischargedMah: minute.dischargedMah,
minVoltageMv: minute.minVoltageMv,
maxVoltageMv: minute.maxVoltageMv,
avgVoltageMv: minute.voltageCount ? minute.voltageTotal / minute.voltageCount : null,
minCurrentMa: minute.minCurrentMa,
maxCurrentMa: minute.maxCurrentMa,
avgCurrentMa: minute.currentCount ? minute.currentTotal / minute.currentCount : null,
minTemperatureC: minute.minTemperatureC,
maxTemperatureC: minute.maxTemperatureC,
avgTemperatureC: minute.temperatureCount ? minute.temperatureTotal / minute.temperatureCount : null,
minChargeMah: minute.minChargeMah,
maxChargeMah: minute.maxChargeMah,
lastChargeMah: minute.lastChargeMah,
reportedCapacityMah: minute.reportedCapacityMah,
dockedSamples: minute.dockedSamples,
chargingSamples: minute.chargingSamples,
commandCount: minute.commandCount,
driveCommandCount: minute.driveCommandCount,
rejectedCommandCount: minute.rejectedCommandCount,
distanceMm: minute.distanceMm,
bumpCount: minute.bumpCount,
cliffCount: minute.cliffCount,
wheelDropCount: minute.wheelDropCount,
virtualWallCount: minute.virtualWallCount,
overcurrentEpisodeCount: minute.overcurrentEpisodeCount,
};
}
function newBatterySession(roverId, kind, now, sensors, state) {
return {
roverId,
batteryKey: state.batteryKey || batteryKey(roverId),
kind,
startedAt: now,
endedAt: null,
startChargeMah: finite(sensors?.batteryChargeMah),
endChargeMah: null,
chargedMah: 0,
dischargedMah: 0,
minVoltageMv: finite(sensors?.voltageMv),
maxVoltageMv: finite(sensors?.voltageMv),
minTemperatureC: finite(sensors?.batteryTemperatureC),
maxTemperatureC: finite(sensors?.batteryTemperatureC),
sampleCount: 0,
gapCount: 0,
status: 'open',
confidence: 'low',
qualificationReason: 'session is still open',
details: {
startedFromQualifiedFull: Boolean(state.fullQualifiedAt),
fullQualifiedAt: state.fullQualifiedAt,
warnMah: finite(state.lastBatteryState?.warn),
urgentMah: finite(state.lastBatteryState?.urgent),
configuredFullMah: finite(state.lastBatteryState?.full),
},
};
}
function observedSessionKind(sensors) {
const code = finite(sensors?.chargingState?.code);
const docked = Boolean(sensors?.chargingSources?.homeBase || sensors?.chargingSources?.internalCharger);
const current = finite(sensors?.currentMa);
if (docked && (code === 1 || code === 2 || code === 3 || (current != null && current > 25))) return 'charging';
if (!docked && current != null && current < -25) return 'discharging';
return 'idle';
}
function createCollector({ storage, logger, maximumIntegrationGapMs, minimumCapacityTestDepthPercent }) {
const roverStates = new Map();
const lastManagerSampleAt = new Map();
const diagnostics = {
startedAt: Date.now(),
eventsObserved: 0,
eventsStored: 0,
sensorFramesObserved: 0,
validBatteryFrames: 0,
integrationGaps: 0,
minuteWrites: 0,
sessionsCompleted: 0,
lastEventAt: null,
lastSensorAt: null,
lastError: null,
};
function stateFor(roverId, now) {
if (!roverStates.has(roverId)) {
roverStates.set(roverId, {
roverId,
lastAt: null,
minute: makeMinute(roverId, now),
candidateKind: null,
candidateCount: 0,
sessionKind: 'idle',
session: null,
waitingSince: null,
fullQualifiedAt: null,
lastBatteryState: null,
batteryKey: storage.getActiveBattery?.(roverId)?.batteryKey || batteryKey(roverId),
lastOdometerTotalMm: null,
safety: {
bump: false,
cliff: false,
wheelDrop: false,
virtualWall: false,
overcurrent: false,
overcurrentStartedAt: null,
overcurrentSamples: 0,
},
});
}
return roverStates.get(roverId);
}
function collectEvent(event = {}) {
diagnostics.eventsObserved += 1;
diagnostics.lastEventAt = Date.now();
try {
const normalized = {
ts: finite(event.ts) || Date.now(),
source: String(event.source || 'unknown'),
type: String(event.type || 'unknown'),
roverId: eventRoverId(event),
visibility: inferVisibility(event),
severity: inferSeverity(event.type),
correlationId: event?.payload?.correlationId || event?.payload?.jobId || event?.payload?.sessionId || null,
payload: event.payload ?? null,
};
if (storage.insertEvent(normalized)) diagnostics.eventsStored += 1;
} catch (err) {
diagnostics.lastError = err.message;
logger.warn('Fleet collector ignored malformed domain event', { error: err.message });
}
}
function updateMinute(minute, sensors, elapsedMs, chargedMah, dischargedMah, gap) {
const voltage = finite(sensors?.voltageMv);
const current = finite(sensors?.currentMa);
const temperature = finite(sensors?.batteryTemperatureC);
const charge = finite(sensors?.batteryChargeMah);
const capacity = finite(sensors?.batteryCapacityMah);
minute.sampleCount += 1;
minute.coverageMs += elapsedMs;
minute.gapCount += gap ? 1 : 0;
minute.chargedMah += chargedMah;
minute.dischargedMah += dischargedMah;
minute.minVoltageMv = minimum(minute.minVoltageMv, voltage);
minute.maxVoltageMv = maximum(minute.maxVoltageMv, voltage);
if (voltage != null) { minute.voltageTotal += voltage; minute.voltageCount += 1; }
minute.minCurrentMa = minimum(minute.minCurrentMa, current);
minute.maxCurrentMa = maximum(minute.maxCurrentMa, current);
if (current != null) { minute.currentTotal += current; minute.currentCount += 1; }
minute.minTemperatureC = minimum(minute.minTemperatureC, temperature);
minute.maxTemperatureC = maximum(minute.maxTemperatureC, temperature);
if (temperature != null) { minute.temperatureTotal += temperature; minute.temperatureCount += 1; }
minute.minChargeMah = minimum(minute.minChargeMah, charge);
minute.maxChargeMah = maximum(minute.maxChargeMah, charge);
minute.lastChargeMah = charge;
minute.reportedCapacityMah = capacity;
if (sensors?.chargingSources?.homeBase) minute.dockedSamples += 1;
if (observedSessionKind(sensors) === 'charging') minute.chargingSamples += 1;
}
function finishSession(state, now, sensors, reason) {
const session = state.session;
if (!session) return;
session.endedAt = now;
session.endChargeMah = finite(sensors?.batteryChargeMah);
session.status = 'completed';
if (session.kind === 'discharging') {
const configuredFull = finite(session.details.configuredFullMah);
const startCharge = finite(session.startChargeMah);
const endCharge = finite(session.endChargeMah);
const reference = configuredFull || startCharge;
const observedDepth = reference && startCharge != null && endCharge != null
? Math.max(0, ((startCharge - endCharge) / reference) * 100)
: 0;
const reachedLowEndpoint = Boolean(
state.lastBatteryState?.urgentActive ||
(finite(session.details.urgentMah) != null && endCharge != null && endCharge <= session.details.urgentMah),
);
const qualified = Boolean(
session.details.startedFromQualifiedFull &&
reachedLowEndpoint &&
observedDepth >= minimumCapacityTestDepthPercent &&
session.gapCount === 0,
);
session.details.observedDepthPercent = observedDepth;
session.details.reachedLowEndpoint = reachedLowEndpoint;
session.details.capacityTestQualified = qualified;
if (qualified) {
session.confidence = 'high';
session.qualificationReason = 'continuous qualified-full to low-endpoint discharge';
} else if (observedDepth >= 30 && session.gapCount <= 1) {
session.confidence = 'medium';
session.qualificationReason = reason || 'useful partial discharge; not a full capacity test';
} else {
session.confidence = 'low';
session.qualificationReason = reason || 'insufficient depth, endpoint, or telemetry coverage';
}
} else {
session.confidence = session.gapCount === 0 ? 'high' : session.gapCount <= 1 ? 'medium' : 'low';
session.qualificationReason = reason || 'charging session completed';
}
storage.insertBatterySession(session);
diagnostics.sessionsCompleted += 1;
state.session = null;
}
function applySessionKind(state, nextKind, now, sensors) {
if (nextKind === state.sessionKind) {
state.candidateKind = null;
state.candidateCount = 0;
return;
}
if (state.candidateKind !== nextKind) {
state.candidateKind = nextKind;
state.candidateCount = 1;
return;
}
state.candidateCount += 1;
if (state.candidateCount < SESSION_KIND_CONFIRM_SAMPLES) return;
finishSession(state, now, sensors, `state changed from ${state.sessionKind} to ${nextKind}`);
state.sessionKind = nextKind;
state.candidateKind = null;
state.candidateCount = 0;
if (nextKind === 'charging' || nextKind === 'discharging') {
state.session = newBatterySession(state.roverId, nextKind, now, sensors, state);
// A qualified-full marker is consumed by the next discharge. Leaving it
// set during the open session records the evidence in session details,
// while clearing it prevents later partial sessions from inheriting it.
if (nextKind === 'discharging') state.fullQualifiedAt = null;
}
}
function updateFullQualification(state, now, sensors) {
const waiting = finite(sensors?.chargingState?.code) === 4;
if (waiting) {
if (state.waitingSince == null) state.waitingSince = now;
if (now - state.waitingSince >= FULL_WAIT_MS && state.fullQualifiedAt == null) {
state.fullQualifiedAt = state.waitingSince + FULL_WAIT_MS;
}
} else {
state.waitingSince = null;
}
}
function collectSensor({ roverId, sensors, batteryState } = {}) {
diagnostics.sensorFramesObserved += 1;
diagnostics.lastSensorAt = Date.now();
if (!roverId || !sensors || finite(sensors.currentMa) == null) return;
diagnostics.validBatteryFrames += 1;
const now = Date.now();
const state = stateFor(String(roverId), now);
state.lastBatteryState = batteryState || state.lastBatteryState;
const elapsedMs = state.lastAt == null ? 0 : Math.max(0, now - state.lastAt);
const gap = elapsedMs > maximumIntegrationGapMs;
const validElapsedMs = gap ? 0 : elapsedMs;
if (gap) diagnostics.integrationGaps += 1;
const currentMa = finite(sensors.currentMa) || 0;
const deltaMah = currentMa * validElapsedMs / 3600000;
const chargedMah = Math.max(0, deltaMah);
const dischargedMah = Math.max(0, -deltaMah);
const bucketTs = Math.floor(now / MINUTE_MS) * MINUTE_MS;
if (state.minute.bucketTs !== bucketTs) {
storage.upsertMinute(persistedMinute(state.minute));
diagnostics.minuteWrites += 1;
state.minute = makeMinute(state.roverId, now);
}
updateMinute(state.minute, sensors, validElapsedMs, chargedMah, dischargedMah, gap);
const bumps = sensors?.bumpsAndWheelDrops || {};
const nextSafety = {
bump: Boolean(bumps.bumpLeft || bumps.bumpRight),
cliff: Boolean(sensors.cliffLeft || sensors.cliffFrontLeft || sensors.cliffFrontRight || sensors.cliffRight),
wheelDrop: Boolean(bumps.wheelDropLeft || bumps.wheelDropRight),
virtualWall: Boolean(sensors.virtualWall),
overcurrent: Boolean(
sensors?.wheelOvercurrents?.leftWheel || sensors?.wheelOvercurrents?.rightWheel ||
sensors?.wheelOvercurrents?.mainBrush || sensors?.wheelOvercurrents?.sideBrush,
),
};
if (nextSafety.bump && !state.safety.bump) state.minute.bumpCount += 1;
if (nextSafety.cliff && !state.safety.cliff) state.minute.cliffCount += 1;
if (nextSafety.wheelDrop && !state.safety.wheelDrop) state.minute.wheelDropCount += 1;
if (nextSafety.virtualWall && !state.safety.virtualWall) state.minute.virtualWallCount += 1;
if (nextSafety.overcurrent) state.safety.overcurrentSamples += 1;
if (nextSafety.overcurrent && !state.safety.overcurrent) {
state.minute.overcurrentEpisodeCount += 1;
state.safety.overcurrentStartedAt = now;
state.safety.overcurrentSamples = 1;
collectEvent({
source: 'fleetReportService',
type: 'overcurrent.episode.started',
ts: now,
payload: { roverId: state.roverId, motors: sensors.wheelOvercurrents },
});
} else if (!nextSafety.overcurrent && state.safety.overcurrent) {
collectEvent({
source: 'fleetReportService',
type: 'overcurrent.episode.resolved',
ts: now,
payload: {
roverId: state.roverId,
startedAt: state.safety.overcurrentStartedAt,
durationMs: Math.max(0, now - (state.safety.overcurrentStartedAt || now)),
sampleCount: state.safety.overcurrentSamples,
},
});
state.safety.overcurrentStartedAt = null;
state.safety.overcurrentSamples = 0;
}
Object.assign(state.safety, nextSafety);
updateFullQualification(state, now, sensors);
applySessionKind(state, observedSessionKind(sensors), now, sensors);
if (state.session) {
const session = state.session;
session.sampleCount += 1;
session.gapCount += gap ? 1 : 0;
session.chargedMah += chargedMah;
session.dischargedMah += dischargedMah;
session.minVoltageMv = minimum(session.minVoltageMv, finite(sensors.voltageMv));
session.maxVoltageMv = maximum(session.maxVoltageMv, finite(sensors.voltageMv));
session.minTemperatureC = minimum(session.minTemperatureC, finite(sensors.batteryTemperatureC));
session.maxTemperatureC = maximum(session.maxTemperatureC, finite(sensors.batteryTemperatureC));
}
state.lastAt = now;
}
function collectCommand(command = {}) {
if (!command.roverId) return;
const now = finite(command.ts) || Date.now();
const state = stateFor(String(command.roverId), now);
const bucketTs = Math.floor(now / MINUTE_MS) * MINUTE_MS;
if (state.minute.bucketTs !== bucketTs) {
if (state.minute.sampleCount || state.minute.commandCount) {
storage.upsertMinute(persistedMinute(state.minute));
diagnostics.minuteWrites += 1;
}
state.minute = makeMinute(state.roverId, now);
}
state.minute.commandCount += 1;
if (command.type === 'drive' || command.type === 'motors') state.minute.driveCommandCount += 1;
if (command.outcome === 'rejected') state.minute.rejectedCommandCount += 1;
// Drive/motor commands can arrive at control-loop frequency. Their exact
// volume belongs in minute counters, while rejections and low-frequency
// actions remain individually inspectable. This preserves operational
// depth without turning normal held movement into an event-timeline flood.
if ((command.type !== 'drive' && command.type !== 'motors') || command.outcome === 'rejected') {
collectEvent({
source: 'commandService',
type: `command.${command.outcome || 'observed'}`,
ts: now,
payload: command,
});
}
}
function collectManagerEvent(kind, event = {}) {
const roverId = event.roverId ? String(event.roverId) : null;
if (kind === 'hostStats') {
const key = `${kind}:${roverId || 'unknown'}`;
const now = Date.now();
// Host statistics arrive periodically and change gradually. One exact
// sample every five minutes retains long-term diagnostic evidence while
// avoiding a timeline row for every routine host heartbeat.
if (now - (lastManagerSampleAt.get(key) || 0) < 5 * 60 * 1000) return;
lastManagerSampleAt.set(key, now);
collectEvent({
source: 'roverHost',
type: 'host.sample',
ts: event.receivedAt || now,
payload: { roverId, stats: event.stats || null },
});
return;
}
const payload = { ...event };
// Live rover records contain websocket handles, sets, and other runtime
// objects. The lifecycle facts are sufficient evidence and serialize
// predictably without copying those control-owned objects into storage.
delete payload.record;
collectEvent({
source: 'roverManager',
type: `roverManager.${kind}`,
payload,
});
}
function collectOdometer({ roverId, odometer } = {}) {
if (!roverId || !odometer) return;
const now = finite(odometer.updatedAt) || Date.now();
const state = stateFor(String(roverId), now);
const totalMm = finite(odometer.totalMm);
if (totalMm == null) return;
const bucketTs = Math.floor(now / MINUTE_MS) * MINUTE_MS;
if (state.minute.bucketTs !== bucketTs) {
if (state.minute.sampleCount || state.minute.commandCount || state.minute.distanceMm) {
storage.upsertMinute(persistedMinute(state.minute));
diagnostics.minuteWrites += 1;
}
state.minute = makeMinute(state.roverId, now);
}
if (state.lastOdometerTotalMm != null && totalMm >= state.lastOdometerTotalMm) {
// Odometer total is already rollover-corrected and sanity-filtered by its
// owning service. Only non-negative increments belong in this report;
// resets establish a new baseline instead of subtracting fleet distance.
state.minute.distanceMm += totalMm - state.lastOdometerTotalMm;
}
state.lastOdometerTotalMm = totalMm;
}
function flushMinutes() {
roverStates.forEach((state) => {
if (!state.minute.sampleCount && !state.minute.commandCount && !state.minute.distanceMm) return;
storage.upsertMinute(persistedMinute(state.minute));
diagnostics.minuteWrites += 1;
});
}
function getLiveState() {
return Array.from(roverStates.values()).map((state) => ({
roverId: state.roverId,
lastAt: state.lastAt,
sessionKind: state.sessionKind,
waitingSince: state.waitingSince,
fullQualifiedAt: state.fullQualifiedAt,
minute: persistedMinute(state.minute),
openSession: state.session ? {
...state.session,
// The database serializer is not involved in this live response, so a
// defensive copy prevents UI consumers from mutating collector state.
details: { ...state.session.details },
} : null,
}));
}
function getDiagnostics() {
return {
...diagnostics,
activeRovers: roverStates.size,
openSessions: Array.from(roverStates.values()).filter((state) => state.session).length,
instanceId: crypto.createHash('sha1').update(String(diagnostics.startedAt)).digest('hex').slice(0, 10),
};
}
function refreshBatteryIdentity(roverId) {
const id = String(roverId || '');
if (!id) return;
const state = roverStates.get(id);
if (!state) return;
state.batteryKey = storage.getActiveBattery?.(id)?.batteryKey || batteryKey(id);
}
return {
collectEvent,
collectSensor,
collectCommand,
collectManagerEvent,
collectOdometer,
flushMinutes,
getLiveState,
getDiagnostics,
refreshBatteryIdentity,
};
}
module.exports = {
createCollector,
};
@@ -0,0 +1,93 @@
// Fleet Report Collector Tests
// Purpose: Verifies signed-current integration, gap rejection, and high-rate command noise reduction independently of SQLite.
// Scope: Uses an in-memory storage double so tests exercise collection policy without touching development data files.
const test = require('node:test');
const assert = require('node:assert/strict');
const { createCollector } = require('./collector');
function makeHarness() {
const writes = { events: [], minutes: [], sessions: [] };
const storage = {
insertEvent(event) { writes.events.push(event); return { changes: 1 }; },
upsertMinute(minute) { writes.minutes.push({ ...minute }); return { changes: 1 }; },
insertBatterySession(session) { writes.sessions.push({ ...session }); return { changes: 1 }; },
};
const collector = createCollector({
storage,
logger: { warn() {} },
maximumIntegrationGapMs: 5000,
minimumCapacityTestDepthPercent: 60,
});
return { collector, writes };
}
function sensors(overrides = {}) {
return {
currentMa: -3600,
voltageMv: 14500,
batteryTemperatureC: 25,
batteryChargeMah: 2000,
batteryCapacityMah: 3000,
chargingState: { code: 0, label: 'not charging' },
chargingSources: { homeBase: false, internalCharger: false },
...overrides,
};
}
test('integrates signed battery current while excluding long telemetry gaps', () => {
const { collector } = makeHarness();
const originalNow = Date.now;
let now = 1_000_000;
Date.now = () => now;
try {
collector.collectSensor({ roverId: 'alpha', sensors: sensors() });
now += 1000;
collector.collectSensor({ roverId: 'alpha', sensors: sensors() });
now += 6000;
collector.collectSensor({ roverId: 'alpha', sensors: sensors() });
const live = collector.getLiveState()[0].minute;
// -3600 mA for one valid second is exactly one discharged mAh. The six
// second interval exceeds the configured integration gap and adds no
// fictional throughput.
assert.equal(live.dischargedMah, 1);
assert.equal(live.chargedMah, 0);
assert.equal(live.gapCount, 1);
assert.equal(live.coverageMs, 1000);
} finally {
Date.now = originalNow;
}
});
test('aggregates drive commands into minute counters instead of event noise', () => {
const { collector, writes } = makeHarness();
const originalNow = Date.now;
Date.now = () => 2_000_000;
try {
for (let index = 0; index < 100; index += 1) {
collector.collectCommand({ roverId: 'alpha', type: 'drive', outcome: 'issued', ts: 2_000_000 + index });
}
collector.collectCommand({ roverId: 'alpha', type: 'drive', outcome: 'rejected', error: 'safety cooldown' });
collector.collectCommand({ roverId: 'alpha', type: 'horn', outcome: 'issued' });
const live = collector.getLiveState()[0].minute;
assert.equal(live.commandCount, 102);
assert.equal(live.driveCommandCount, 101);
assert.equal(live.rejectedCommandCount, 1);
assert.equal(writes.events.length, 2);
assert.deepEqual(writes.events.map((event) => event.type), ['command.rejected', 'command.issued']);
} finally {
Date.now = originalNow;
}
});
test('preserves chat content in structured global events', () => {
const { collector, writes } = makeHarness();
collector.collectEvent({
source: 'chat',
type: 'chat:message',
ts: 3_000_000,
payload: { text: 'hello fleet history', nickname: 'Otter' },
});
assert.equal(writes.events.length, 1);
assert.equal(writes.events[0].payload.text, 'hello fleet history');
assert.equal(writes.events[0].visibility, 'global');
});
@@ -0,0 +1,117 @@
// Fleet Report Service
// Purpose: Composes optional passive collection, storage, analysis, retention, and read-only transport.
// Scope: This is the sole feature boundary; disabled installations register no collectors, timers, database, or sockets.
const { loadConfig } = require('../../helpers/configLoader');
const { isFeatureEnabled } = require('../../helpers/features');
const logger = require('../../globals/logger').child('fleetReportService');
if (!isFeatureEnabled('fleetReports')) {
module.exports = {
enabled: false,
getDailyReport: () => null,
};
} else {
const { subscribeAll } = require('../eventBus');
const roverManager = require('../roverManager');
const { commandEvents } = require('../commandService');
const { odometerEvents } = require('../odometerService');
const { createStorage } = require('./storage');
const { createCollector } = require('./collector');
const { createReportBuilder } = require('./reportBuilder');
const { registerSocketGateway } = require('./socketGateway');
const config = loadConfig().fleetReports || {};
const batteryConfig = config.battery || {};
const retentionConfig = config.retention || {};
const maximumIntegrationGapMs = Math.max(
250,
(Number(batteryConfig.maximumIntegrationGapSeconds) || 5) * 1000,
);
const minimumCapacityTestDepthPercent = Math.max(
10,
Math.min(100, Number(batteryConfig.minimumCapacityTestDepthPercent) || 60),
);
const batteryEnabled = batteryConfig.enabled !== false;
const storage = createStorage({ logger });
const collector = createCollector({
storage,
logger,
maximumIntegrationGapMs,
minimumCapacityTestDepthPercent,
});
const reportBuilder = createReportBuilder({ storage, collector, roverManager });
storage.open();
const unsubscribeEvents = subscribeAll(collector.collectEvent);
if (batteryEnabled) roverManager.managerEvents.on('sensor', collector.collectSensor);
commandEvents.on('observation', collector.collectCommand);
odometerEvents.on('update', collector.collectOdometer);
const managerEventKinds = ['rover', 'hostStats', 'driver', 'switch', 'lock', 'private', 'privateSafety'];
const managerEventHandlers = new Map(managerEventKinds.map((kind) => {
const handler = (event) => collector.collectManagerEvent(kind, event);
roverManager.managerEvents.on(kind, handler);
return [kind, handler];
}));
registerSocketGateway({ roverManager, reportBuilder, storage, collector, logger });
// Periodic upserts bound data-loss on an unclean shutdown while still
// avoiding writes at the 20 Hz sensor-frame rate.
const flushTimer = setInterval(() => collector.flushMinutes(), 30 * 1000);
flushTimer.unref?.();
function retentionDays(value, fallback) {
const number = Number(value);
return Number.isFinite(number) && number >= 0 ? number : fallback;
}
function pruneNow() {
const now = Date.now();
const detailedDays = retentionDays(retentionConfig.detailedDays, 0);
const minuteDays = retentionDays(retentionConfig.minuteSamplesDays, 0);
storage.prune({
detailedBefore: detailedDays === 0 ? 0 : now - detailedDays * 86400000,
minuteBefore: minuteDays === 0 ? 0 : now - minuteDays * 86400000,
});
}
pruneNow();
const retentionTimer = setInterval(pruneNow, 6 * 60 * 60 * 1000);
retentionTimer.unref?.();
function getDailyReport({ since, until, roverIds } = {}) {
const end = Number(until) || Date.now();
return reportBuilder.build({
since: Number(since) || end - 24 * 60 * 60 * 1000,
until: end,
roverIds: Array.isArray(roverIds) ? roverIds : undefined,
includeEvents: true,
eventLimit: 1000,
});
}
logger.info('Fleet reporting enabled', {
databaseAvailable: storage.getDiagnostics().available,
maximumIntegrationGapMs,
minimumCapacityTestDepthPercent,
batteryEnabled,
});
module.exports = {
enabled: true,
getDailyReport,
collector,
storage,
reportBuilder,
// Exposed for controlled tests and graceful future shutdown wiring. Normal
// runtime leaves subscriptions active for the lifetime of the server.
stop() {
unsubscribeEvents();
if (batteryEnabled) roverManager.managerEvents.off('sensor', collector.collectSensor);
commandEvents.off('observation', collector.collectCommand);
odometerEvents.off('update', collector.collectOdometer);
managerEventHandlers.forEach((handler, kind) => roverManager.managerEvents.off(kind, handler));
clearInterval(flushTimer);
clearInterval(retentionTimer);
collector.flushMinutes();
},
};
}
@@ -0,0 +1,247 @@
// Fleet Report Builder
// Purpose: Produces dense read models from stored evidence without leaking presentation logic into collection.
// Scope: Owns totals, rover comparisons, attention findings, and exact supporting datasets for UI/Discord consumers.
function sum(rows, key) {
return rows.reduce((total, row) => total + (Number(row?.[key]) || 0), 0);
}
function groupByRover(minutes, roster = []) {
const rosterById = new Map(roster.map((rover) => [String(rover.id), rover]));
const grouped = new Map();
minutes.forEach((minute) => {
const roverId = String(minute.roverId);
if (!grouped.has(roverId)) grouped.set(roverId, []);
grouped.get(roverId).push(minute);
});
return Array.from(new Set([...rosterById.keys(), ...grouped.keys()])).map((roverId) => {
const rows = grouped.get(roverId) || [];
const samples = sum(rows, 'sampleCount');
const voltageWeighted = rows.reduce(
(total, row) => total + (Number(row.avgVoltageMv) || 0) * (Number(row.sampleCount) || 0),
0,
);
const temperatureWeighted = rows.reduce(
(total, row) => total + (Number(row.avgTemperatureC) || 0) * (Number(row.sampleCount) || 0),
0,
);
const latest = rows[rows.length - 1] || null;
return {
roverId,
name: rosterById.get(roverId)?.name || roverId,
color: rosterById.get(roverId)?.color || null,
online: Boolean(rosterById.get(roverId)),
sampleCount: samples,
coverageMs: sum(rows, 'coverageMs'),
gapCount: sum(rows, 'gapCount'),
commandCount: sum(rows, 'commandCount'),
driveCommandCount: sum(rows, 'driveCommandCount'),
rejectedCommandCount: sum(rows, 'rejectedCommandCount'),
distanceMm: sum(rows, 'distanceMm'),
bumpCount: sum(rows, 'bumpCount'),
cliffCount: sum(rows, 'cliffCount'),
wheelDropCount: sum(rows, 'wheelDropCount'),
virtualWallCount: sum(rows, 'virtualWallCount'),
overcurrentEpisodeCount: sum(rows, 'overcurrentEpisodeCount'),
chargedMah: sum(rows, 'chargedMah'),
dischargedMah: sum(rows, 'dischargedMah'),
averageVoltageMv: samples ? voltageWeighted / samples : null,
minimumVoltageMv: rows.reduce((value, row) => row.minVoltageMv == null ? value : Math.min(value ?? Infinity, row.minVoltageMv), null),
maximumVoltageMv: rows.reduce((value, row) => row.maxVoltageMv == null ? value : Math.max(value ?? -Infinity, row.maxVoltageMv), null),
averageTemperatureC: samples ? temperatureWeighted / samples : null,
minimumTemperatureC: rows.reduce((value, row) => row.minTemperatureC == null ? value : Math.min(value ?? Infinity, row.minTemperatureC), null),
maximumTemperatureC: rows.reduce((value, row) => row.maxTemperatureC == null ? value : Math.max(value ?? -Infinity, row.maxTemperatureC), null),
latestChargeMah: latest?.lastChargeMah ?? null,
reportedCapacityMah: latest?.reportedCapacityMah ?? null,
lastSampleAt: latest ? latest.bucketTs + 60000 : null,
};
}).sort((a, b) => a.name.localeCompare(b.name));
}
function buildFindings({ roverRows, events, now }) {
const findings = [];
roverRows.forEach((rover) => {
if (!rover.sampleCount) {
findings.push({
key: `no-telemetry:${rover.roverId}`,
roverId: rover.roverId,
severity: rover.online ? 'warning' : 'notice',
confidence: 'high',
status: 'ongoing',
title: rover.online ? 'No battery telemetry in selected range' : 'Rover offline or absent',
evidence: { sampleCount: 0 },
});
return;
}
if (rover.maximumTemperatureC != null && rover.maximumTemperatureC >= 45) {
findings.push({
key: `battery-temperature:${rover.roverId}`,
roverId: rover.roverId,
severity: rover.maximumTemperatureC >= 50 ? 'critical' : 'warning',
confidence: 'high',
status: 'observed',
title: 'High battery temperature observed',
evidence: { maximumTemperatureC: rover.maximumTemperatureC, sampleCount: rover.sampleCount },
});
}
if (rover.gapCount > 0) {
findings.push({
key: `telemetry-gaps:${rover.roverId}`,
roverId: rover.roverId,
severity: 'notice',
confidence: 'high',
status: 'observed',
title: 'Battery integration contains telemetry gaps',
evidence: { gapCount: rover.gapCount, coverageMs: rover.coverageMs },
});
}
if (rover.lastSampleAt && now - rover.lastSampleAt > 5 * 60 * 1000 && rover.online) {
findings.push({
key: `stale-telemetry:${rover.roverId}`,
roverId: rover.roverId,
severity: 'warning',
confidence: 'high',
status: 'ongoing',
title: 'Telemetry is stale while rover is online',
evidence: { lastSampleAt: rover.lastSampleAt },
});
}
});
const criticalEvents = events.filter((event) => event.severity === 'critical');
if (criticalEvents.length) {
findings.push({
key: 'critical-events',
roverId: null,
severity: 'critical',
confidence: 'high',
status: 'observed',
title: `${criticalEvents.length} critical event${criticalEvents.length === 1 ? '' : 's'} in selected range`,
evidence: { eventIds: criticalEvents.slice(0, 25).map((event) => event.id) },
});
}
const rank = { critical: 0, warning: 1, notice: 2, informational: 3 };
return findings.sort((a, b) => (rank[a.severity] ?? 9) - (rank[b.severity] ?? 9));
}
function median(values) {
if (!values.length) return null;
const sorted = [...values].sort((a, b) => a - b);
const middle = Math.floor(sorted.length / 2);
return sorted.length % 2 ? sorted[middle] : (sorted[middle - 1] + sorted[middle]) / 2;
}
function buildBatteryHealth({ sessions, roverRows, batteryRegistry }) {
return roverRows.map((rover) => {
const roverSessions = sessions.filter((session) => String(session.roverId) === String(rover.roverId));
const qualified = roverSessions
.filter((session) => session.kind === 'discharging' && session.details?.capacityTestQualified)
.sort((a, b) => a.startedAt - b.startedAt);
// The first three qualified tests establish the learned healthy baseline.
// A median resists one unusually light/heavy run while remaining auditable
// in the session table. Battery replacement identity will start a separate
// key, so only sessions for the current key should contribute once one is
// registered.
const activeRegistryEntry = batteryRegistry.find((battery) =>
String(battery.roverId) === String(rover.roverId) && battery.retiredAt == null,
);
const currentKey = activeRegistryEntry?.batteryKey || qualified[qualified.length - 1]?.batteryKey || roverSessions[0]?.batteryKey || `unregistered:${rover.roverId}`;
const sameBattery = qualified.filter((session) => session.batteryKey === currentKey);
const baselineTests = sameBattery.slice(0, 3);
const baselineMah = median(baselineTests.map((session) => Number(session.dischargedMah)));
const latest = sameBattery[sameBattery.length - 1] || null;
const measuredUsableMah = latest ? Number(latest.dischargedMah) : null;
const capacityRetentionPercent = baselineMah && measuredUsableMah != null
? measuredUsableMah / baselineMah * 100
: null;
const throughputMah = roverSessions.reduce((total, session) => total + (Number(session.dischargedMah) || 0), 0);
const cycleReferenceMah = baselineMah || Number(rover.reportedCapacityMah) || null;
return {
roverId: rover.roverId,
batteryKey: currentKey,
qualifiedTestCount: sameBattery.length,
baselineTestCount: baselineTests.length,
baselineMah,
measuredUsableMah,
capacityRetentionPercent,
equivalentFullCycles: cycleReferenceMah ? throughputMah / cycleReferenceMah : null,
dischargedThroughputMah: throughputMah,
latestQualifiedTestAt: latest?.endedAt || null,
confidence: sameBattery.length >= 3 ? 'high' : sameBattery.length >= 1 ? 'medium' : 'low',
confidenceReason: sameBattery.length >= 3
? 'at least three qualified full-to-low tests'
: sameBattery.length >= 1
? 'fewer than three qualified tests'
: 'no qualified full-to-low capacity test',
};
});
}
function createReportBuilder({ storage, collector, roverManager }) {
function build({ since, until, roverIds, includeEvents = true, eventLimit = 500 }) {
const visibleRoster = Array.from(roverManager.rovers?.values?.() || []).map((record) => ({
id: record.id,
name: record.name || record.id,
color: record.color || record.meta?.color || null,
}));
const requestedIds = Array.isArray(roverIds) ? roverIds.map(String) : null;
const roster = requestedIds
? visibleRoster.filter((rover) => requestedIds.includes(String(rover.id)))
: visibleRoster;
const effectiveIds = requestedIds || roster.map((rover) => String(rover.id));
const minutes = storage.listMinutes({ since, until, roverIds: effectiveIds });
const events = includeEvents
? storage.listEvents({ since, until, roverIds: effectiveIds, limit: eventLimit })
: [];
const batterySessions = storage.listBatterySessions({ since, until, roverIds: effectiveIds, limit: 500 });
const roverRows = groupByRover(minutes, roster);
const batteryRegistry = storage.listBatteries(effectiveIds);
const batteryHealth = buildBatteryHealth({ sessions: batterySessions, roverRows, batteryRegistry });
const findings = buildFindings({ roverRows, events, now: Date.now() });
return {
generatedAt: Date.now(),
range: { since, until },
totals: {
roverCount: roverRows.length,
onlineRoverCount: roverRows.filter((rover) => rover.online).length,
sampleCount: sum(minutes, 'sampleCount'),
coverageMs: sum(minutes, 'coverageMs'),
telemetryGapCount: sum(minutes, 'gapCount'),
commandCount: sum(minutes, 'commandCount'),
driveCommandCount: sum(minutes, 'driveCommandCount'),
rejectedCommandCount: sum(minutes, 'rejectedCommandCount'),
distanceMm: sum(minutes, 'distanceMm'),
bumpCount: sum(minutes, 'bumpCount'),
cliffCount: sum(minutes, 'cliffCount'),
wheelDropCount: sum(minutes, 'wheelDropCount'),
virtualWallCount: sum(minutes, 'virtualWallCount'),
overcurrentEpisodeCount: sum(minutes, 'overcurrentEpisodeCount'),
chargedMah: sum(minutes, 'chargedMah'),
dischargedMah: sum(minutes, 'dischargedMah'),
eventCountReturned: events.length,
batterySessionCount: batterySessions.length,
criticalFindingCount: findings.filter((finding) => finding.severity === 'critical').length,
warningFindingCount: findings.filter((finding) => finding.severity === 'warning').length,
},
rovers: roverRows,
findings,
minutes,
batterySessions,
batteryHealth,
batteryRegistry,
dailyReportHistory: storage.listDailyReports(365),
events,
live: collector.getLiveState(),
diagnostics: {
collector: collector.getDiagnostics(),
storage: storage.getDiagnostics(),
},
};
}
return { build };
}
module.exports = {
createReportBuilder,
};
@@ -0,0 +1,103 @@
// Fleet Report Socket Gateway
// Purpose: Exposes read-only, visibility-filtered fleet history to browser clients.
// Scope: Owns query validation and existing private/lockdown access boundaries; it performs no collection or analysis.
const io = require('../../globals/io');
const crypto = require('crypto');
const { isAdmin, isLockdownAdmin } = require('../roleService');
const MAX_RANGE_MS = 366 * 24 * 60 * 60 * 1000;
function normalizeRange(payload = {}) {
const now = Date.now();
const until = Number.isFinite(Number(payload.until)) ? Number(payload.until) : now;
const requestedSince = Number.isFinite(Number(payload.since))
? Number(payload.since)
: until - 24 * 60 * 60 * 1000;
const since = Math.max(0, Math.max(requestedSince, until - MAX_RANGE_MS));
return { since, until: Math.max(since + 1, until) };
}
function registerSocketGateway({ roverManager, reportBuilder, storage, collector, logger }) {
io.on('connection', (socket) => {
socket.on('fleetReports:get', (payload = {}, cb = () => {}) => {
try {
const { since, until } = normalizeRange(payload);
// getRosterForSocket is the canonical live private-rover visibility
// resolver. Historical queries use precisely those currently visible
// rover IDs so fleet totals cannot indirectly disclose a private rover.
const visibleRoverIds = roverManager.getRosterForSocket(socket).map((rover) => String(rover.id));
const requestedIds = Array.isArray(payload.roverIds)
? payload.roverIds.map(String).filter((id) => visibleRoverIds.includes(id))
: visibleRoverIds;
const report = reportBuilder.build({
since,
until,
roverIds: requestedIds,
includeEvents: payload.includeEvents !== false,
eventLimit: payload.eventLimit,
});
// Lockdown-only events are deliberately removed after query assembly.
// They are global rather than rover-scoped, so rover filtering alone is
// insufficient to preserve the pre-existing lockdown privacy boundary.
if (!isLockdownAdmin(socket)) {
report.events = report.events.filter((event) => event.visibility !== 'lockdown');
report.totals.eventCountReturned = report.events.length;
}
if (payload.compact === true) {
// The Activities card needs current totals and findings, not the
// underlying time series. Removing bulky evidence here keeps the
// always-visible surface cheap while `/reports` retains full depth.
report.minutes = [];
report.events = [];
report.batterySessions = [];
report.dailyReportHistory = [];
report.live = report.live.map((entry) => ({
roverId: entry.roverId,
lastAt: entry.lastAt,
sessionKind: entry.sessionKind,
}));
}
cb({ ok: true, report });
} catch (err) {
logger.warn('Fleet report query failed', { socketId: socket.id, error: err.message });
cb({ error: 'Fleet report query failed' });
}
});
socket.on('fleetReports:replaceBattery', (payload = {}, cb = () => {}) => {
try {
if (!isAdmin(socket)) throw new Error('Admin access required');
const roverId = String(payload.roverId || '').trim();
if (!roverId || !roverManager.rovers.has(roverId)) throw new Error('Known online rover required');
const ratedCapacityMah = Number(payload.ratedCapacityMah);
if (!Number.isFinite(ratedCapacityMah) || ratedCapacityMah <= 0 || ratedCapacityMah > 65535) {
throw new Error('Rated capacity must be between 1 and 65535 mAh');
}
const installedAt = Number.isFinite(Number(payload.installedAt)) ? Number(payload.installedAt) : Date.now();
const entry = storage.replaceBattery({
roverId,
batteryKey: `battery:${roverId}:${installedAt}:${crypto.randomUUID().slice(0, 8)}`,
chemistry: String(payload.chemistry || '').trim() || null,
ratedCapacityMah: Math.round(ratedCapacityMah),
installedAt,
notes: String(payload.notes || '').trim() || null,
});
if (!entry) throw new Error('Battery registry write failed');
collector.refreshBatteryIdentity(roverId);
collector.collectEvent({
source: 'fleetReportService',
type: 'battery.replaced',
payload: { roverId, battery: entry },
});
cb({ ok: true, battery: entry });
} catch (err) {
logger.warn('Fleet battery replacement rejected', { socketId: socket.id, error: err.message });
cb({ error: err.message });
}
});
});
}
module.exports = {
registerSocketGateway,
};
@@ -0,0 +1,506 @@
// Fleet Report Storage
// Purpose: Owns the reporting database, schema, bounded writes, and read queries.
// Scope: Keeps SQLite details out of telemetry collection, analysis, UI transport, and Discord delivery.
const fs = require('fs');
const path = require('path');
const Database = require('better-sqlite3');
const { resolveDataPath } = require('../../helpers/dataPaths');
const DB_PATH = resolveDataPath('fleet-reports.sqlite');
function safeJson(value) {
try {
return JSON.stringify(value ?? null);
} catch (_err) {
// An unusual circular payload must not break the event subscriber. The
// placeholder still records that an event occurred and explains why its
// supporting payload is unavailable.
return JSON.stringify({ serializationError: true });
}
}
function parseJson(value, fallback = null) {
try {
return value == null ? fallback : JSON.parse(value);
} catch (_err) {
return fallback;
}
}
function createStorage({ logger }) {
let db = null;
let statements = null;
function open() {
if (db) return true;
try {
fs.mkdirSync(path.dirname(DB_PATH), { recursive: true });
db = new Database(DB_PATH);
db.pragma('journal_mode = WAL');
db.pragma('synchronous = NORMAL');
db.pragma('foreign_keys = ON');
db.exec(`
CREATE TABLE IF NOT EXISTS fleet_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ts INTEGER NOT NULL,
source TEXT NOT NULL,
type TEXT NOT NULL,
rover_id TEXT,
visibility TEXT NOT NULL DEFAULT 'global',
severity TEXT NOT NULL DEFAULT 'informational',
correlation_id TEXT,
payload_json TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_fleet_events_ts ON fleet_events(ts DESC);
CREATE INDEX IF NOT EXISTS idx_fleet_events_rover_ts ON fleet_events(rover_id, ts DESC);
CREATE INDEX IF NOT EXISTS idx_fleet_events_type_ts ON fleet_events(type, ts DESC);
CREATE TABLE IF NOT EXISTS fleet_minute_samples (
rover_id TEXT NOT NULL,
bucket_ts INTEGER NOT NULL,
sample_count INTEGER NOT NULL,
coverage_ms INTEGER NOT NULL,
gap_count INTEGER NOT NULL,
charged_mah REAL NOT NULL,
discharged_mah REAL NOT NULL,
min_voltage_mv INTEGER,
max_voltage_mv INTEGER,
avg_voltage_mv REAL,
min_current_ma INTEGER,
max_current_ma INTEGER,
avg_current_ma REAL,
min_temperature_c INTEGER,
max_temperature_c INTEGER,
avg_temperature_c REAL,
min_charge_mah INTEGER,
max_charge_mah INTEGER,
last_charge_mah INTEGER,
reported_capacity_mah INTEGER,
docked_samples INTEGER NOT NULL,
charging_samples INTEGER NOT NULL,
command_count INTEGER NOT NULL DEFAULT 0,
drive_command_count INTEGER NOT NULL DEFAULT 0,
rejected_command_count INTEGER NOT NULL DEFAULT 0,
distance_mm REAL NOT NULL DEFAULT 0,
bump_count INTEGER NOT NULL DEFAULT 0,
cliff_count INTEGER NOT NULL DEFAULT 0,
wheel_drop_count INTEGER NOT NULL DEFAULT 0,
virtual_wall_count INTEGER NOT NULL DEFAULT 0,
overcurrent_episode_count INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (rover_id, bucket_ts)
);
CREATE INDEX IF NOT EXISTS idx_fleet_minutes_ts ON fleet_minute_samples(bucket_ts DESC);
CREATE TABLE IF NOT EXISTS fleet_battery_sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
rover_id TEXT NOT NULL,
battery_key TEXT NOT NULL,
kind TEXT NOT NULL,
started_at INTEGER NOT NULL,
ended_at INTEGER,
start_charge_mah INTEGER,
end_charge_mah INTEGER,
charged_mah REAL NOT NULL DEFAULT 0,
discharged_mah REAL NOT NULL DEFAULT 0,
min_voltage_mv INTEGER,
max_voltage_mv INTEGER,
min_temperature_c INTEGER,
max_temperature_c INTEGER,
sample_count INTEGER NOT NULL DEFAULT 0,
gap_count INTEGER NOT NULL DEFAULT 0,
status TEXT NOT NULL DEFAULT 'open',
confidence TEXT NOT NULL DEFAULT 'low',
qualification_reason TEXT,
details_json TEXT NOT NULL DEFAULT '{}'
);
CREATE INDEX IF NOT EXISTS idx_battery_sessions_rover_time ON fleet_battery_sessions(rover_id, started_at DESC);
CREATE TABLE IF NOT EXISTS fleet_batteries (
battery_key TEXT PRIMARY KEY,
rover_id TEXT NOT NULL,
chemistry TEXT,
rated_capacity_mah INTEGER,
installed_at INTEGER,
retired_at INTEGER,
healthy_baseline_mah REAL,
notes TEXT,
updated_at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_fleet_batteries_rover ON fleet_batteries(rover_id, installed_at DESC);
CREATE TABLE IF NOT EXISTS fleet_daily_reports (
report_date TEXT PRIMARY KEY,
generated_at INTEGER NOT NULL,
report_json TEXT NOT NULL,
discord_delivered_at INTEGER,
discord_error TEXT
);
`);
// SQLite's CREATE TABLE IF NOT EXISTS does not add columns to an older
// reporting database. These narrow additive migrations keep development
// databases usable as collection coverage expands without coupling this
// optional feature to the identity database's migration history.
const minuteColumns = new Set(db.prepare('PRAGMA table_info(fleet_minute_samples)').all().map((column) => column.name));
[
['command_count', 'INTEGER NOT NULL DEFAULT 0'],
['drive_command_count', 'INTEGER NOT NULL DEFAULT 0'],
['rejected_command_count', 'INTEGER NOT NULL DEFAULT 0'],
['distance_mm', 'REAL NOT NULL DEFAULT 0'],
['bump_count', 'INTEGER NOT NULL DEFAULT 0'],
['cliff_count', 'INTEGER NOT NULL DEFAULT 0'],
['wheel_drop_count', 'INTEGER NOT NULL DEFAULT 0'],
['virtual_wall_count', 'INTEGER NOT NULL DEFAULT 0'],
['overcurrent_episode_count', 'INTEGER NOT NULL DEFAULT 0'],
].forEach(([name, definition]) => {
if (!minuteColumns.has(name)) db.exec(`ALTER TABLE fleet_minute_samples ADD COLUMN ${name} ${definition}`);
});
statements = {
insertEvent: db.prepare(`
INSERT INTO fleet_events (ts, source, type, rover_id, visibility, severity, correlation_id, payload_json)
VALUES (@ts, @source, @type, @roverId, @visibility, @severity, @correlationId, @payloadJson)
`),
upsertMinute: db.prepare(`
INSERT INTO fleet_minute_samples (
rover_id, bucket_ts, sample_count, coverage_ms, gap_count,
charged_mah, discharged_mah, min_voltage_mv, max_voltage_mv, avg_voltage_mv,
min_current_ma, max_current_ma, avg_current_ma, min_temperature_c,
max_temperature_c, avg_temperature_c, min_charge_mah, max_charge_mah,
last_charge_mah, reported_capacity_mah, docked_samples, charging_samples,
command_count, drive_command_count, rejected_command_count,
distance_mm, bump_count, cliff_count, wheel_drop_count,
virtual_wall_count, overcurrent_episode_count
) VALUES (
@roverId, @bucketTs, @sampleCount, @coverageMs, @gapCount,
@chargedMah, @dischargedMah, @minVoltageMv, @maxVoltageMv, @avgVoltageMv,
@minCurrentMa, @maxCurrentMa, @avgCurrentMa, @minTemperatureC,
@maxTemperatureC, @avgTemperatureC, @minChargeMah, @maxChargeMah,
@lastChargeMah, @reportedCapacityMah, @dockedSamples, @chargingSamples,
@commandCount, @driveCommandCount, @rejectedCommandCount,
@distanceMm, @bumpCount, @cliffCount, @wheelDropCount,
@virtualWallCount, @overcurrentEpisodeCount
)
ON CONFLICT(rover_id, bucket_ts) DO UPDATE SET
sample_count = excluded.sample_count,
coverage_ms = excluded.coverage_ms,
gap_count = excluded.gap_count,
charged_mah = excluded.charged_mah,
discharged_mah = excluded.discharged_mah,
min_voltage_mv = excluded.min_voltage_mv,
max_voltage_mv = excluded.max_voltage_mv,
avg_voltage_mv = excluded.avg_voltage_mv,
min_current_ma = excluded.min_current_ma,
max_current_ma = excluded.max_current_ma,
avg_current_ma = excluded.avg_current_ma,
min_temperature_c = excluded.min_temperature_c,
max_temperature_c = excluded.max_temperature_c,
avg_temperature_c = excluded.avg_temperature_c,
min_charge_mah = excluded.min_charge_mah,
max_charge_mah = excluded.max_charge_mah,
last_charge_mah = excluded.last_charge_mah,
reported_capacity_mah = excluded.reported_capacity_mah,
docked_samples = excluded.docked_samples,
charging_samples = excluded.charging_samples,
command_count = excluded.command_count,
drive_command_count = excluded.drive_command_count,
rejected_command_count = excluded.rejected_command_count
,distance_mm = excluded.distance_mm
,bump_count = excluded.bump_count
,cliff_count = excluded.cliff_count
,wheel_drop_count = excluded.wheel_drop_count
,virtual_wall_count = excluded.virtual_wall_count
,overcurrent_episode_count = excluded.overcurrent_episode_count
`),
insertSession: db.prepare(`
INSERT INTO fleet_battery_sessions (
rover_id, battery_key, kind, started_at, ended_at, start_charge_mah,
end_charge_mah, charged_mah, discharged_mah, min_voltage_mv,
max_voltage_mv, min_temperature_c, max_temperature_c, sample_count,
gap_count, status, confidence, qualification_reason, details_json
) VALUES (
@roverId, @batteryKey, @kind, @startedAt, @endedAt, @startChargeMah,
@endChargeMah, @chargedMah, @dischargedMah, @minVoltageMv,
@maxVoltageMv, @minTemperatureC, @maxTemperatureC, @sampleCount,
@gapCount, @status, @confidence, @qualificationReason, @detailsJson
)
`),
};
return true;
} catch (err) {
logger.error('Failed to open fleet report database; reporting will remain fail-open', {
path: DB_PATH,
error: err.message,
});
if (db) {
try { db.close(); } catch (_closeErr) { /* Best-effort cleanup after failed initialization. */ }
}
db = null;
statements = null;
return false;
}
}
function runSafely(label, operation, fallback = null) {
if (!open()) return fallback;
try {
return operation();
} catch (err) {
// Reporting is observer-only. A failed write/query is visible in logs but
// is never allowed to propagate into a rover sensor or command callback.
logger.warn(`Fleet report storage ${label} failed`, { error: err.message });
return fallback;
}
}
function insertEvent(event) {
return runSafely('event write', () => statements.insertEvent.run({
ts: event.ts,
source: event.source,
type: event.type,
roverId: event.roverId || null,
visibility: event.visibility || 'global',
severity: event.severity || 'informational',
correlationId: event.correlationId || null,
payloadJson: safeJson(event.payload),
}));
}
function upsertMinute(sample) {
return runSafely('minute write', () => statements.upsertMinute.run(sample));
}
function insertBatterySession(session) {
return runSafely('battery session write', () => statements.insertSession.run({
...session,
detailsJson: safeJson(session.details || {}),
}));
}
function listEvents({ since, until, roverIds, limit = 500, offset = 0, type = null }) {
return runSafely('event query', () => {
const clauses = ['ts >= ?', 'ts < ?'];
const params = [since, until];
if (type) {
clauses.push('type = ?');
params.push(type);
}
if (Array.isArray(roverIds)) {
if (roverIds.length === 0) clauses.push('rover_id IS NULL');
else {
clauses.push(`(rover_id IS NULL OR rover_id IN (${roverIds.map(() => '?').join(',')}))`);
params.push(...roverIds);
}
}
params.push(Math.max(1, Math.min(2000, Number(limit) || 500)), Math.max(0, Number(offset) || 0));
const rows = db.prepare(`
SELECT id, ts, source, type, rover_id AS roverId, visibility, severity,
correlation_id AS correlationId, payload_json AS payloadJson
FROM fleet_events
WHERE ${clauses.join(' AND ')}
ORDER BY ts DESC
LIMIT ? OFFSET ?
`).all(...params);
return rows.map(({ payloadJson, ...row }) => ({ ...row, payload: parseJson(payloadJson, {}) }));
}, []);
}
function listMinutes({ since, until, roverIds }) {
return runSafely('minute query', () => {
const clauses = ['bucket_ts >= ?', 'bucket_ts < ?'];
const params = [since, until];
if (Array.isArray(roverIds)) {
if (roverIds.length === 0) return [];
clauses.push(`rover_id IN (${roverIds.map(() => '?').join(',')})`);
params.push(...roverIds);
}
return db.prepare(`
SELECT rover_id AS roverId, bucket_ts AS bucketTs, sample_count AS sampleCount,
coverage_ms AS coverageMs, gap_count AS gapCount, charged_mah AS chargedMah,
discharged_mah AS dischargedMah, min_voltage_mv AS minVoltageMv,
max_voltage_mv AS maxVoltageMv, avg_voltage_mv AS avgVoltageMv,
min_current_ma AS minCurrentMa, max_current_ma AS maxCurrentMa,
avg_current_ma AS avgCurrentMa, min_temperature_c AS minTemperatureC,
max_temperature_c AS maxTemperatureC, avg_temperature_c AS avgTemperatureC,
min_charge_mah AS minChargeMah, max_charge_mah AS maxChargeMah,
last_charge_mah AS lastChargeMah, reported_capacity_mah AS reportedCapacityMah,
docked_samples AS dockedSamples, charging_samples AS chargingSamples,
command_count AS commandCount, drive_command_count AS driveCommandCount,
rejected_command_count AS rejectedCommandCount
,distance_mm AS distanceMm, bump_count AS bumpCount,
cliff_count AS cliffCount, wheel_drop_count AS wheelDropCount,
virtual_wall_count AS virtualWallCount,
overcurrent_episode_count AS overcurrentEpisodeCount
FROM fleet_minute_samples
WHERE ${clauses.join(' AND ')}
ORDER BY bucket_ts ASC, rover_id ASC
`).all(...params);
}, []);
}
function listBatterySessions({ since, until, roverIds, limit = 500 }) {
return runSafely('battery session query', () => {
if (Array.isArray(roverIds) && roverIds.length === 0) return [];
const roverClause = Array.isArray(roverIds)
? `AND rover_id IN (${roverIds.map(() => '?').join(',')})`
: '';
const params = [since, until, ...(roverIds || []), Math.max(1, Math.min(2000, Number(limit) || 500))];
return db.prepare(`
SELECT id, rover_id AS roverId, battery_key AS batteryKey, kind,
started_at AS startedAt, ended_at AS endedAt,
start_charge_mah AS startChargeMah, end_charge_mah AS endChargeMah,
charged_mah AS chargedMah, discharged_mah AS dischargedMah,
min_voltage_mv AS minVoltageMv, max_voltage_mv AS maxVoltageMv,
min_temperature_c AS minTemperatureC, max_temperature_c AS maxTemperatureC,
sample_count AS sampleCount, gap_count AS gapCount, status,
confidence, qualification_reason AS qualificationReason,
details_json AS detailsJson
FROM fleet_battery_sessions
WHERE started_at < ? AND COALESCE(ended_at, started_at) >= ? ${roverClause}
ORDER BY started_at DESC
LIMIT ?
`).all(until, since, ...(roverIds || []), params[params.length - 1]).map(({ detailsJson, ...row }) => ({
...row,
details: parseJson(detailsJson, {}),
}));
}, []);
}
function prune({ detailedBefore, minuteBefore }) {
return runSafely('retention prune', () => db.transaction(() => {
const events = db.prepare('DELETE FROM fleet_events WHERE ts < ?').run(detailedBefore).changes;
const minutes = db.prepare('DELETE FROM fleet_minute_samples WHERE bucket_ts < ?').run(minuteBefore).changes;
return { events, minutes };
})());
}
function getDailyReport(reportDate) {
return runSafely('daily report query', () => {
const row = db.prepare(`
SELECT report_date AS reportDate, generated_at AS generatedAt,
report_json AS reportJson, discord_delivered_at AS discordDeliveredAt,
discord_error AS discordError
FROM fleet_daily_reports WHERE report_date = ?
`).get(reportDate);
if (!row) return null;
const { reportJson, ...metadata } = row;
return { ...metadata, report: parseJson(reportJson, null) };
});
}
function listBatteries(roverIds = null) {
return runSafely('battery registry query', () => {
if (Array.isArray(roverIds) && roverIds.length === 0) return [];
const clause = Array.isArray(roverIds)
? `WHERE rover_id IN (${roverIds.map(() => '?').join(',')})`
: '';
return db.prepare(`
SELECT battery_key AS batteryKey, rover_id AS roverId, chemistry,
rated_capacity_mah AS ratedCapacityMah, installed_at AS installedAt,
retired_at AS retiredAt, healthy_baseline_mah AS healthyBaselineMah,
notes, updated_at AS updatedAt
FROM fleet_batteries ${clause}
ORDER BY rover_id ASC, installed_at DESC
`).all(...(roverIds || []));
}, []);
}
function getActiveBattery(roverId) {
return runSafely('active battery query', () => db.prepare(`
SELECT battery_key AS batteryKey, rover_id AS roverId, chemistry,
rated_capacity_mah AS ratedCapacityMah, installed_at AS installedAt,
healthy_baseline_mah AS healthyBaselineMah, notes, updated_at AS updatedAt
FROM fleet_batteries
WHERE rover_id = ? AND retired_at IS NULL
ORDER BY installed_at DESC
LIMIT 1
`).get(String(roverId)) || null);
}
function replaceBattery(entry) {
return runSafely('battery replacement write', () => db.transaction(() => {
const now = Date.now();
db.prepare('UPDATE fleet_batteries SET retired_at = ?, updated_at = ? WHERE rover_id = ? AND retired_at IS NULL')
.run(entry.installedAt || now, now, entry.roverId);
db.prepare(`
INSERT INTO fleet_batteries (
battery_key, rover_id, chemistry, rated_capacity_mah, installed_at,
retired_at, healthy_baseline_mah, notes, updated_at
) VALUES (?, ?, ?, ?, ?, NULL, ?, ?, ?)
`).run(
entry.batteryKey,
entry.roverId,
entry.chemistry || null,
entry.ratedCapacityMah || null,
entry.installedAt || now,
entry.healthyBaselineMah || null,
entry.notes || null,
now,
);
return getActiveBattery(entry.roverId);
})());
}
function saveDailyReport(reportDate, report) {
return runSafely('daily report write', () => db.prepare(`
INSERT INTO fleet_daily_reports (report_date, generated_at, report_json)
VALUES (?, ?, ?)
ON CONFLICT(report_date) DO UPDATE SET
generated_at = excluded.generated_at,
report_json = excluded.report_json
`).run(reportDate, Date.now(), safeJson(report)));
}
function listDailyReports(limit = 90) {
return runSafely('daily report history query', () => db.prepare(`
SELECT report_date AS reportDate, generated_at AS generatedAt,
discord_delivered_at AS discordDeliveredAt, discord_error AS discordError,
length(report_json) AS reportBytes
FROM fleet_daily_reports
ORDER BY report_date DESC
LIMIT ?
`).all(Math.max(1, Math.min(1000, Number(limit) || 90))), []);
}
function markDailyReportDelivery(reportDate, { deliveredAt = null, error = null } = {}) {
return runSafely('daily delivery update', () => db.prepare(`
UPDATE fleet_daily_reports
SET discord_delivered_at = ?, discord_error = ?
WHERE report_date = ?
`).run(deliveredAt, error, reportDate));
}
function getDiagnostics() {
return runSafely('diagnostics query', () => ({
available: true,
path: DB_PATH,
bytes: fs.statSync(DB_PATH).size,
eventCount: db.prepare('SELECT COUNT(*) AS count FROM fleet_events').get().count,
minuteCount: db.prepare('SELECT COUNT(*) AS count FROM fleet_minute_samples').get().count,
sessionCount: db.prepare('SELECT COUNT(*) AS count FROM fleet_battery_sessions').get().count,
}), { available: false, path: DB_PATH });
}
return {
open,
insertEvent,
upsertMinute,
insertBatterySession,
listEvents,
listMinutes,
listBatterySessions,
prune,
getDailyReport,
saveDailyReport,
listDailyReports,
markDailyReportDelivery,
listBatteries,
getActiveBattery,
replaceBattery,
getDiagnostics,
};
}
module.exports = {
createStorage,
};