Track per-device traffic totals and recover inventory state
Build and Deploy Gateway / build-and-push (push) Successful in 15s
Build and Deploy Gateway / deploy (push) Successful in 13s

This commit is contained in:
2026-08-07 15:24:54 +03:00
parent e774486b99
commit 307ad02cd7
13 changed files with 842 additions and 96 deletions
+20 -2
View File
@@ -39,13 +39,18 @@ import {
import { HarborError, normalizeHarborError } from '../shared/errors.js';
import { normalizeRouteRules } from '../shared/routingRules.js';
import { createJsonStore, createStateStore } from './services/stateStore.js';
import { createDeviceInventoryService, createVendorLookup } from './services/deviceInventoryService.js';
import {
createDeviceInventoryService,
createVendorLookup,
DEVICE_INVENTORY_SCHEMA_VERSION,
migrateDeviceInventoryState,
} from './services/deviceInventoryService.js';
import { buildGatewayVersionInfo, buildVersionInfo } from './version.js';
const MAX_BODY_BYTES = 1_000_000;
const SUBSCRIPTION_REFRESH_INTERVAL_MS = 15 * 60 * 1000;
const GATEWAY_DISCOVERY_INTERVAL_MS = 5_000;
const DEVICE_DISCOVERY_INTERVAL_MS = 60_000;
const DEVICE_DISCOVERY_INTERVAL_MS = 15_000;
const TERMINAL_SUBSCRIPTION_CODES = new Set([
'SUBSCRIPTION_EXPIRED',
'SUBSCRIPTION_DISABLED',
@@ -62,7 +67,17 @@ const subscriptionCacheStore = createJsonStore({
const deviceStore = createJsonStore({
filePath: settings.deviceStatePath,
defaultValue: {},
migrate: migrateDeviceInventoryState,
initializeMissing: true,
backupWhen: () => true,
});
deviceStore.read();
if (deviceStore.migration) {
console.log(`[storage] devices migrated to v${DEVICE_INVENTORY_SCHEMA_VERSION}; backup: ${deviceStore.migration.backupPath}`);
}
if (deviceStore.recovery) {
console.warn(`[storage] corrupt devices recovered; backup: ${deviceStore.recovery.backupPath}`);
}
let cacheRecoveryLogged = false;
function readSubscriptionCache() {
@@ -99,6 +114,9 @@ const deviceInventory = settings.appMode === 'gateway'
observe: remoteDataplane
? () => singboxRuntime.observeDevices()
: () => readNeighborSnapshot(),
observeTraffic: remoteDataplane
? () => singboxRuntime.observeTraffic()
: null,
vendor: createVendorLookup(),
})
: null;
+201 -20
View File
@@ -3,22 +3,40 @@ import fs from 'node:fs';
import net from 'node:net';
import { HarborError } from '../../shared/errors.js';
const SCHEMA_VERSION = 1;
export const DEVICE_INVENTORY_SCHEMA_VERSION = 2;
const ONLINE_MS = 2 * 60 * 1000;
const RECENT_MS = 24 * 60 * 60 * 1000;
const RETENTION_MS = 30 * 24 * 60 * 60 * 1000;
const COUNTER_PATTERN = /^\d+$/;
const MAC_PATTERN = /^[0-9a-f]{2}(?::[0-9a-f]{2}){5}$/;
const DEFAULT_STATE = {
schemaVersion: SCHEMA_VERSION,
schemaVersion: DEVICE_INVENTORY_SCHEMA_VERSION,
revision: 0,
lastObservedAt: null,
lastError: null,
traffic: {
epoch: null,
generation: null,
lastObservedAt: null,
lastError: null,
baselinesByMac: {},
totalsByMac: {},
rebaselineMacs: [],
},
devices: [],
};
const normalizeMac = (value) => String(value || '').trim().toLowerCase();
const deviceId = (mac) => `dev_${crypto.createHash('sha256').update(mac).digest('hex').slice(0, 16)}`;
const isPrivateMac = (mac) => (Number.parseInt(mac.slice(0, 2), 16) & 2) !== 0;
const recordEntries = (value) => value && typeof value === 'object' && !Array.isArray(value)
? Object.entries(value)
: [];
const parseStoredCounter = (value) => {
const counter = String(value ?? '');
return COUNTER_PATTERN.test(counter) ? BigInt(counter).toString() : null;
};
export function parseOuiVendors(text) {
const vendors = new Map();
@@ -44,18 +62,77 @@ export function createVendorLookup(filePath = '/usr/share/ieee-data/oui.txt') {
};
}
function migrate(value) {
export function migrateDeviceInventoryState(value) {
const state = value && typeof value === 'object' && !Array.isArray(value) ? value : {};
const version = Number.isSafeInteger(state.schemaVersion) ? state.schemaVersion : 0;
if (version < 0 || version > SCHEMA_VERSION) {
if (version < 0 || version > DEVICE_INVENTORY_SCHEMA_VERSION) {
throw new Error(`Unsupported device inventory schemaVersion: ${version}`);
}
const traffic = state.traffic && typeof state.traffic === 'object' && !Array.isArray(state.traffic)
? state.traffic
: {};
const devices = Array.isArray(state.devices) ? state.devices : [];
const rebaselineMacs = new Set((Array.isArray(traffic.rebaselineMacs) ? traffic.rebaselineMacs : [])
.map(normalizeMac).filter((mac) => MAC_PATTERN.test(mac)));
let recoveredTraffic = version >= 2 && (
traffic !== state.traffic
|| !traffic.baselinesByMac || typeof traffic.baselinesByMac !== 'object' || Array.isArray(traffic.baselinesByMac)
|| !traffic.totalsByMac || typeof traffic.totalsByMac !== 'object' || Array.isArray(traffic.totalsByMac)
);
if (recoveredTraffic) {
for (const device of devices) rebaselineMacs.add(normalizeMac(device.mac));
}
const baselinesByMac = {};
for (const [rawMac, baseline] of recordEntries(traffic.baselinesByMac)) {
const mac = normalizeMac(rawMac);
const uploadBytes = parseStoredCounter(baseline?.uploadBytes);
const downloadBytes = parseStoredCounter(baseline?.downloadBytes);
if (!MAC_PATTERN.test(mac) || typeof baseline?.epoch !== 'string' || !baseline.epoch
|| uploadBytes == null || downloadBytes == null) {
if (MAC_PATTERN.test(mac)) rebaselineMacs.add(mac);
recoveredTraffic = true;
continue;
}
baselinesByMac[mac] = { epoch: baseline.epoch, uploadBytes, downloadBytes };
}
const totalsByMac = {};
for (const [rawMac, total] of recordEntries(traffic.totalsByMac)) {
const mac = normalizeMac(rawMac);
const uploadBytes = parseStoredCounter(total?.uploadBytes);
const downloadBytes = parseStoredCounter(total?.downloadBytes);
if (!MAC_PATTERN.test(mac) || uploadBytes == null || downloadBytes == null) {
if (MAC_PATTERN.test(mac)) rebaselineMacs.add(mac);
recoveredTraffic = true;
continue;
}
totalsByMac[mac] = {
uploadBytes,
downloadBytes,
observedAt: typeof total?.observedAt === 'string' ? total.observedAt : null,
};
}
for (const mac of new Set([...Object.keys(baselinesByMac), ...Object.keys(totalsByMac)])) {
if (!Object.hasOwn(baselinesByMac, mac) || !Object.hasOwn(totalsByMac, mac)) {
rebaselineMacs.add(mac);
recoveredTraffic = true;
}
}
return {
...DEFAULT_STATE,
...state,
schemaVersion: SCHEMA_VERSION,
schemaVersion: DEVICE_INVENTORY_SCHEMA_VERSION,
revision: Number.isSafeInteger(state.revision) ? state.revision : 0,
devices: Array.isArray(state.devices) ? state.devices : [],
traffic: {
...DEFAULT_STATE.traffic,
...traffic,
lastError: recoveredTraffic
? 'Повреждённый traffic checkpoint восстановлен из корректных данных'
: traffic.lastError || null,
baselinesByMac,
totalsByMac,
rebaselineMacs: [...rebaselineMacs].filter((mac) => MAC_PATTERN.test(mac)),
},
devices,
};
}
@@ -66,17 +143,29 @@ function deviceStatus(lastSeenAt, now) {
return 'offline';
}
export function createDeviceInventoryService({ store, observe, vendor = () => null, now = () => new Date() }) {
export function createDeviceInventoryService({
store,
observe,
observeTraffic = null,
vendor = () => null,
now = () => new Date(),
}) {
let refreshPromise = null;
function snapshot() {
const state = migrate(store.read());
const state = migrateDeviceInventoryState(store.read());
const current = now();
const rank = { online: 0, recent: 1, offline: 2 };
const devices = state.devices.map((device) => ({
...device,
status: deviceStatus(device.lastSeenAt, current),
})).sort((left, right) => (
const devices = state.devices.map((device) => {
const traffic = state.traffic.totalsByMac[device.mac];
return {
...device,
status: deviceStatus(device.lastSeenAt, current),
uploadBytes: traffic?.uploadBytes || '0',
downloadBytes: traffic?.downloadBytes || '0',
trafficObservedAt: traffic?.observedAt || null,
};
}).sort((left, right) => (
Number(right.pinned) - Number(left.pinned)
|| rank[left.status] - rank[right.status]
|| String(right.lastSeenAt).localeCompare(String(left.lastSeenAt))
@@ -87,18 +176,27 @@ export function createDeviceInventoryService({ store, observe, vendor = () => nu
kind: 'neighbor',
lastObservedAt: state.lastObservedAt,
error: state.lastError,
traffic: {
lastObservedAt: state.traffic.lastObservedAt,
error: state.traffic.lastError,
},
},
devices,
};
}
async function performRefresh() {
let result;
try {
result = await observe();
} catch (error) {
result = { observedAt: now().toISOString(), observations: [], error: error.message || String(error) };
}
const [result, trafficResult] = await Promise.all([
Promise.resolve().then(() => observe()).catch((error) => ({
observedAt: now().toISOString(),
observations: [],
error: error.message || String(error),
})),
observeTraffic
? Promise.resolve().then(() => observeTraffic())
.catch((error) => ({ transportError: error.message || String(error) }))
: null,
]);
const observedAt = result?.observedAt || now().toISOString();
const observations = Array.isArray(result?.observations) ? result.observations : [];
const ipsByMac = new Map();
@@ -109,7 +207,7 @@ export function createDeviceInventoryService({ store, observe, vendor = () => nu
ipsByMac.get(mac).add(String(observation.ip));
}
store.update((stored) => {
const state = migrate(stored);
const state = migrateDeviceInventoryState(stored);
const byMac = new Map(state.devices.map((device) => [device.mac, device]));
for (const observation of observations) {
const mac = normalizeMac(observation.mac);
@@ -139,11 +237,94 @@ export function createDeviceInventoryService({ store, observe, vendor = () => nu
const devices = [...byMac.values()].filter((device) => (
device.pinned || device.alias || new Date(device.lastSeenAt).getTime() >= cutoff
));
let traffic = state.traffic;
if (trafficResult) {
if (trafficResult.transportError) {
traffic = { ...traffic, lastError: trafficResult.transportError };
} else {
try {
if (typeof trafficResult.epoch !== 'string' || !trafficResult.epoch) {
throw new Error('Dataplane не вернул traffic epoch');
}
const processByMac = new Map();
for (const row of Array.isArray(trafficResult.devices) ? trafficResult.devices : []) {
const mac = normalizeMac(row?.mac);
const upload = String(row?.uploadBytes ?? '');
const download = String(row?.downloadBytes ?? '');
if (!MAC_PATTERN.test(mac) || !COUNTER_PATTERN.test(upload) || !COUNTER_PATTERN.test(download)) {
throw new Error('Dataplane вернул невалидный traffic counter');
}
const previous = processByMac.get(mac) || { upload: 0n, download: 0n };
processByMac.set(mac, {
upload: previous.upload + BigInt(upload),
download: previous.download + BigInt(download),
});
}
const knownMacs = new Set(devices.map((device) => device.mac));
const baselinesByMac = { ...traffic.baselinesByMac };
const totalsByMac = { ...traffic.totalsByMac };
const rebaselineMacs = new Set(traffic.rebaselineMacs);
for (const [mac, processTotal] of processByMac) {
if (!knownMacs.has(mac)) continue;
const baseline = baselinesByMac[mac];
const recovering = rebaselineMacs.has(mac);
const sameEpoch = !recovering && baseline?.epoch === trafficResult.epoch;
const baselineUpload = sameEpoch ? BigInt(baseline.uploadBytes) : 0n;
const baselineDownload = sameEpoch ? BigInt(baseline.downloadBytes) : 0n;
if (processTotal.upload < baselineUpload || processTotal.download < baselineDownload) {
throw new Error('Dataplane traffic counter уменьшился внутри одного epoch');
}
const total = totalsByMac[mac] || { uploadBytes: '0', downloadBytes: '0' };
totalsByMac[mac] = {
uploadBytes: (BigInt(total.uploadBytes)
+ (recovering ? 0n : processTotal.upload - baselineUpload)).toString(),
downloadBytes: (BigInt(total.downloadBytes)
+ (recovering ? 0n : processTotal.download - baselineDownload)).toString(),
observedAt: trafficResult.observedAt || traffic.lastObservedAt,
};
baselinesByMac[mac] = {
epoch: trafficResult.epoch,
uploadBytes: processTotal.upload.toString(),
downloadBytes: processTotal.download.toString(),
};
rebaselineMacs.delete(mac);
}
for (const mac of Object.keys(totalsByMac)) {
if (!knownMacs.has(mac)) {
delete totalsByMac[mac];
delete baselinesByMac[mac];
rebaselineMacs.delete(mac);
}
}
for (const mac of rebaselineMacs) {
if (!knownMacs.has(mac)) {
delete totalsByMac[mac];
delete baselinesByMac[mac];
rebaselineMacs.delete(mac);
}
}
traffic = {
...traffic,
epoch: trafficResult.epoch,
generation: trafficResult.generation || traffic.generation,
lastObservedAt: trafficResult.observedAt || traffic.lastObservedAt,
lastError: trafficResult.source?.error
|| (rebaselineMacs.size ? traffic.lastError : null),
baselinesByMac,
totalsByMac,
rebaselineMacs: [...rebaselineMacs],
};
} catch (error) {
traffic = { ...traffic, lastError: error.message || String(error) };
}
}
}
return {
...state,
revision: state.revision + 1,
lastObservedAt: observedAt,
lastError: result?.error || null,
traffic,
devices,
};
});
@@ -172,7 +353,7 @@ export function createDeviceInventoryService({ store, observe, vendor = () => nu
throw new HarborError('REQUEST_INVALID');
}
store.update((stored) => {
const state = migrate(stored);
const state = migrateDeviceInventoryState(stored);
if (state.revision !== expectedRevision) throw new HarborError('STATE_CONFLICT');
const index = state.devices.findIndex((device) => device.id === id);
if (index < 0) throw new HarborError('DEVICE_NOT_FOUND');
+89 -15
View File
@@ -112,12 +112,18 @@ export function createDeviceTrafficService({
run = spawnSync,
nextGeneration = () => crypto.randomUUID(),
}) {
const epoch = nextGeneration();
let activeSlot = null;
let activeDevices = [];
let activeSignature = '';
let activeCounters = new Map();
let pendingRetired = null;
let refreshPromise = null;
const finalized = new Map();
const devicesByKey = new Map();
let current = {
generation: nextGeneration(),
epoch,
generation: epoch,
observedAt: null,
source: { error: null },
devices: [],
@@ -162,21 +168,72 @@ export function createDeviceTrafficService({
}
}
function readCounters(devices) {
if (!activeSlot) return [];
function readCounters(devices, slot) {
if (!slot) return new Map();
const upload = parseTrafficCounters(
execute('iptables-save', ['-c', '-t', 'raw']),
childChain(uploadChain, activeSlot),
childChain(uploadChain, slot),
);
const download = parseTrafficCounters(
execute('iptables-save', ['-c', '-t', 'mangle']),
childChain(downloadChain, activeSlot),
childChain(downloadChain, slot),
);
return devices.map(({ key, ...device }) => ({
...device,
uploadBytes: upload.get(`${key}:upload`) || '0',
downloadBytes: download.get(`${key}:download`) || '0',
}));
const counters = new Map();
for (const { key } of devices) {
counters.set(`${key}:upload`, upload.get(`${key}:upload`) || '0');
counters.set(`${key}:download`, download.get(`${key}:download`) || '0');
}
return counters;
}
function counter(counters, key, direction) {
return BigInt(counters.get(`${key}:${direction}`) || '0');
}
function remember(devices) {
for (const device of devices) devicesByKey.set(device.key, device);
}
function finalizeRetired() {
if (!pendingRetired) return false;
const counters = readCounters(pendingRetired.devices, pendingRetired.slot);
for (const { key } of pendingRetired.devices) {
const previous = finalized.get(key) || { upload: 0n, download: 0n };
finalized.set(key, {
upload: previous.upload + counter(counters, key, 'upload'),
download: previous.download + counter(counters, key, 'download'),
});
}
pendingRetired = null;
return true;
}
function processTotals() {
const activeByMac = new Map(activeDevices.map((device) => [device.mac, device]));
const totalsByMac = new Map();
for (const [key, remembered] of devicesByKey) {
const base = finalized.get(key) || { upload: 0n, download: 0n };
const pending = pendingRetired?.counters || new Map();
const upload = base.upload
+ counter(pending, key, 'upload')
+ counter(activeCounters, key, 'upload');
const download = base.download
+ counter(pending, key, 'download')
+ counter(activeCounters, key, 'download');
const previous = totalsByMac.get(remembered.mac) || { upload: 0n, download: 0n };
totalsByMac.set(remembered.mac, {
...(activeByMac.get(remembered.mac) || remembered),
upload: previous.upload + upload,
download: previous.download + download,
});
}
return [...totalsByMac.values()]
.map(({ key: _key, upload, download, ...device }) => ({
...device,
uploadBytes: upload.toString(),
downloadBytes: download.toString(),
}))
.sort((left, right) => left.mac.localeCompare(right.mac));
}
async function performRefresh() {
@@ -192,34 +249,51 @@ export function createDeviceTrafficService({
? activeDevices
: selectTrafficDevices(observed?.observations);
const nextSignature = JSON.stringify(nextDevices);
let countersRead = false;
if (!sourceError && nextSignature !== activeSignature) {
if (pendingRetired) {
try {
countersRead = finalizeRetired() || countersRead;
} catch (error) {
sourceError = sourceError || error.message || String(error);
}
}
if (!pendingRetired && !sourceError && nextSignature !== activeSignature) {
const nextSlot = activeSlot === 'A' ? 'B' : 'A';
try {
prepare(nextSlot, nextDevices);
switchTo(nextSlot);
const retired = activeSlot ? {
slot: activeSlot,
devices: activeDevices,
counters: activeCounters,
} : null;
activeSlot = nextSlot;
activeDevices = nextDevices;
activeSignature = nextSignature;
activeCounters = new Map();
pendingRetired = retired;
remember(nextDevices);
current.generation = nextGeneration();
if (pendingRetired) countersRead = finalizeRetired() || countersRead;
} catch (error) {
sourceError = error.message || String(error);
}
}
let devices = current.devices;
let countersRead = false;
try {
devices = readCounters(activeDevices);
activeCounters = readCounters(activeDevices, activeSlot);
countersRead = true;
} catch (error) {
sourceError = sourceError || error.message || String(error);
}
current = {
epoch,
generation: current.generation,
observedAt: countersRead ? observed?.observedAt || current.observedAt : current.observedAt,
source: { error: sourceError },
devices,
devices: countersRead ? processTotals() : current.devices,
};
return structuredClone(current);
}