Files
harbor-net/test/server/device-traffic.test.js
T

432 lines
20 KiB
JavaScript

import assert from 'node:assert/strict';
import fs from 'node:fs';
import path from 'node:path';
import test from 'node:test';
import {
buildTrafficRestore,
createDeviceTrafficService,
parseTrafficCounters,
selectTrafficDevices,
} from '../../dist/server/services/deviceTrafficService.js';
const uploadChain = 'VPN_PROXY_TRAFFIC_UP';
const downloadChain = 'VPN_PROXY_TRAFFIC_DOWN';
const routeAccounting = {
directChain: 'VPN_PROXY_DIRECT',
directMark: '0x40000000',
tproxyMark: '1',
gatewayClientCidrs: ['10.0.0.0/8', '172.16.0.0/12', '192.168.0.0/16'],
};
const dataplaneSource = fs.readFileSync(
path.resolve(import.meta.dirname, '../../src/server/dataplane.ts'),
'utf8',
);
const observation = (ip, mac = '00:11:22:33:44:55', deviceInterface = 'eth0') => ({
ip,
mac,
interface: deviceInterface,
});
test('dataplane exposes cached traffic snapshots without making accounting a readiness dependency', () => {
assert.match(dataplaneSource, /req\.method === 'GET' && req\.url === '\/device-traffic'[\s\S]*traffic\.snapshot\(\)/);
assert.match(dataplaneSource, /async function refreshDeviceTraffic\(\)[\s\S]*traffic\.refresh\(\)/);
assert.match(dataplaneSource, /ready = true;[\s\S]*deviceTrafficAccountingEnabled[\s\S]*setImmediate[\s\S]*refreshDeviceTraffic\(\)/);
assert.match(dataplaneSource, /refreshDeviceTraffic\(\)\.catch/);
});
test('traffic selection keeps only unambiguous IPv4 neighbors', () => {
const selected = selectTrafficDevices([
observation('192.168.50.7'),
observation('192.168.50.7'),
observation('192.168.50.8', '00:11:22:33:44:66'),
observation('192.168.50.9', '00:11:22:33:44:66'),
observation('192.168.50.10', '00:11:22:33:44:77', 'bad interface'),
observation('192.168.50.11', '00:11:22:33:44:88', 'br-docker0'),
observation('192.168.50.12', '00:11:22:33:44:99', 'eth+'),
observation('172.17.0.2', '00:11:22:33:44:aa', 'docker0'),
observation('172.18.0.2', '00:11:22:33:44:bb', 'veth1234'),
observation('2001:db8::7', '00:11:22:33:44:88'),
]);
assert.equal(selected.length, 1);
assert.deepEqual(
{ ip: selected[0].ip, mac: selected[0].mac, interface: selected[0].interface },
observation('192.168.50.7'),
);
assert.match(selected[0].key, /^[a-f0-9]{16}$/);
});
test('traffic rules split local proxy traffic from public Gateway traffic in one restore batch', () => {
const [device] = selectTrafficDevices([observation('192.168.50.7')]);
const restore = buildTrafficRestore({
devices: [device],
bypassCidrs: ['10.0.0.0/8', '192.168.0.0/16'],
uploadChain,
downloadChain,
...routeAccounting,
slot: 'A',
proxyPort: 8080,
});
const lines = restore.trim().split('\n');
const proxyUploadJump = '-A VPN_PROXY_TRAFFIC_UP_A -p tcp --dport 8080 -m addrtype --dst-type LOCAL -j VPN_PROXY_TRAFFIC_UP_A_P';
const proxyUploadReturn = '-A VPN_PROXY_TRAFFIC_UP_A -p tcp --dport 8080 -m addrtype --dst-type LOCAL -j RETURN';
const proxyDownloadJump = '-A VPN_PROXY_TRAFFIC_DOWN_A -p tcp --sport 8080 -m addrtype --src-type LOCAL -j VPN_PROXY_TRAFFIC_DOWN_A_P';
const proxyDownloadReturn = '-A VPN_PROXY_TRAFFIC_DOWN_A -p tcp --sport 8080 -m addrtype --src-type LOCAL -j RETURN';
assert.ok(lines.includes(proxyUploadJump));
assert.ok(lines.includes(proxyUploadReturn));
assert.ok(lines.indexOf(proxyUploadJump) < lines.indexOf(proxyUploadReturn));
assert.ok(lines.includes(proxyDownloadJump));
assert.ok(lines.includes(proxyDownloadReturn));
assert.ok(lines.indexOf(proxyDownloadJump) < lines.indexOf(proxyDownloadReturn));
assert.match(restore, /-A VPN_PROXY_TRAFFIC_UP_A_P .*--mac-source 00:11:22:33:44:55 .*harbor-traffic:[a-f0-9]{16}:proxy-upload/);
assert.match(restore, /-A VPN_PROXY_TRAFFIC_DOWN_A_P .*192\.168\.50\.7 .*harbor-traffic:[a-f0-9]{16}:proxy-download/);
assert.ok(lines.indexOf('-A VPN_PROXY_TRAFFIC_UP_A -d 10.0.0.0/8 -j RETURN')
< lines.findIndex((line) => line.endsWith(':upload -j RETURN')));
assert.ok(lines.indexOf('-A VPN_PROXY_TRAFFIC_DOWN_A -s 10.0.0.0/8 -j RETURN')
< lines.findIndex((line) => line.endsWith(':download -j RETURN')));
assert.match(restore, /-A VPN_PROXY_DIRECT_A -m addrtype --dst-type LOCAL -j RETURN/);
assert.match(restore, /-A VPN_PROXY_DIRECT_A -m mark --mark 1\/1 -j RETURN/);
assert.match(restore, /-A VPN_PROXY_DIRECT_A -i br-\+ -j RETURN/);
assert.match(restore, /-A VPN_PROXY_DIRECT_A -d 10\.0\.0\.0\/8 -j RETURN/);
assert.match(restore, /-A VPN_PROXY_DIRECT_A -s 192\.168\.0\.0\/16 .*harbor-traffic:global:direct-upload -j CONNMARK --set-xmark 0x40000000\/0x40000000/);
assert.match(restore, /-A VPN_PROXY_TRAFFIC_DOWN_A -d 192\.168\.0\.0\/16 .*harbor-traffic:global:direct-download$/m);
assert.doesNotMatch(restore, /-A VPN_PROXY_TRAFFIC_DOWN_A -m connmark .*harbor-traffic:global:direct-download$/m);
assert.match(restore, /-A VPN_PROXY_DIRECT_A .*--mark 0x40000000\/0x40000000 .*harbor-traffic:[a-f0-9]{16}:direct-upload$/m);
assert.match(restore, /-A VPN_PROXY_TRAFFIC_DOWN_A .*--mark 0x40000000\/0x40000000 .*harbor-traffic:[a-f0-9]{16}:direct-download$/m);
assert.doesNotMatch(restore, /TPROXY|DNAT|SNAT|REDIRECT/);
assert.throws(() => buildTrafficRestore({
devices: [],
bypassCidrs: [],
uploadChain,
downloadChain,
...routeAccounting,
gatewayClientCidrs: ['10.0.0.0/8', '10.1.0.0/16'],
slot: 'A',
proxyPort: 8080,
}), /traffic accounting/);
});
test('counter parser preserves exact uint64 strings and sums TCP plus UDP proxy rules', () => {
const counters = parseTrafficCounters(
'[7:9007199254740993] -A VPN_PROXY_TRAFFIC_UP_A_P -p tcp -m comment --comment "harbor-traffic:0123456789abcdef:proxy-upload" -j RETURN\n'
+ '[3:9] -A VPN_PROXY_TRAFFIC_UP_A_P -p udp -m comment --comment "harbor-traffic:0123456789abcdef:proxy-upload" -j RETURN\n',
'VPN_PROXY_TRAFFIC_UP_A_P',
);
assert.equal(counters.get('0123456789abcdef:proxy-upload'), '9007199254741002');
const direct = parseTrafficCounters(
'[5:77] -A VPN_PROXY_TRAFFIC_DOWN_A -m connmark --mark 0x40000000/0x40000000 -m comment --comment "harbor-traffic:0123456789abcdef:direct-download"\n',
'VPN_PROXY_TRAFFIC_DOWN_A',
);
assert.equal(direct.get('0123456789abcdef:direct-download'), '77');
});
test('traffic service preserves active rules and snapshot when replacement fails', async () => {
let observed = {
observedAt: '2026-08-07T12:00:00.000Z',
observations: [observation('192.168.50.7')],
error: null,
};
const firstDevice = selectTrafficDevices(observed.observations)[0];
const calls = [];
let failReplacement = false;
let failCounters = false;
const generations = ['boot', 'rules-a', 'rules-b'];
const run = (command, args, options) => {
calls.push([command, args, options]);
if (command === 'iptables-save') {
if (failCounters) return { status: null, stdout: '', stderr: '', error: new Error('counter read timed out') };
const direction = args.includes('raw') ? 'upload' : 'download';
const chain = direction === 'upload' ? `${uploadChain}_A` : `${downloadChain}_A`;
const bytes = direction === 'upload' ? '1200' : '3400';
const proxyDirection = direction === 'upload' ? 'proxy-upload' : 'proxy-download';
const proxyBytes = direction === 'upload' ? ['100', '10'] : ['200', '20'];
return {
status: 0,
stdout: `[1:${bytes}] -A ${chain} -m comment --comment "harbor-traffic:${firstDevice.key}:${direction}" -j RETURN\n`
+ proxyBytes.map((value) => `[1:${value}] -A ${chain}_P -m comment --comment "harbor-traffic:${firstDevice.key}:${proxyDirection}" -j RETURN`).join('\n')
+ (direction === 'download'
? `\n[1:300] -A ${routeAccounting.directChain}_A -m comment --comment "harbor-traffic:${firstDevice.key}:direct-upload" -j CONNMARK\n`
+ `[1:500] -A ${routeAccounting.directChain}_A -m comment --comment "harbor-traffic:global:direct-upload" -j CONNMARK\n`
+ `[1:400] -A ${chain} -m comment --comment "harbor-traffic:${firstDevice.key}:direct-download"\n`
+ `[1:600] -A ${chain} -m comment --comment "harbor-traffic:global:direct-download"\n`
: ''),
stderr: '',
};
}
if (failReplacement && command === 'iptables-restore'
&& options.input.includes(`-R ${downloadChain}`)) {
return { status: 1, stdout: '', stderr: 'cannot switch download rules' };
}
return { status: 0, stdout: '', stderr: '' };
};
const service = createDeviceTrafficService({
observe: async () => observed,
uploadChain,
downloadChain,
bypassCidrs: ['10.0.0.0/8'],
proxyPort: 8080,
...routeAccounting,
run,
nextGeneration: () => generations.shift(),
});
const first = await service.refresh();
assert.equal(first.epoch, 'boot');
assert.equal(first.generation, 'rules-a');
assert.deepEqual(first.direct, { uploadBytes: '500', downloadBytes: '600' });
assert.deepEqual(first.devices, [{
ip: '192.168.50.7',
mac: '00:11:22:33:44:55',
interface: 'eth0',
uploadBytes: '1200',
downloadBytes: '3400',
proxyUploadBytes: '110',
proxyDownloadBytes: '220',
directUploadBytes: '300',
directDownloadBytes: '400',
}]);
const switchCallsBefore = calls.filter(([, args]) => args.includes('-R') || args.includes('-A')).length;
await service.refresh();
const switchCallsAfter = calls.filter(([, args]) => args.includes('-R') || args.includes('-A')).length;
assert.equal(switchCallsAfter, switchCallsBefore);
observed = {
observedAt: '2026-08-07T12:01:00.000Z',
observations: [observation('192.168.50.8')],
error: null,
};
failReplacement = true;
const failed = await service.refresh();
assert.equal(failed.generation, 'rules-a');
assert.match(failed.source.error, /cannot switch download rules/);
assert.deepEqual(failed.devices, first.devices);
const uploadSwitches = calls
.filter(([, args]) => args.includes('-R') && args.includes(uploadChain))
.map(([, args]) => args.at(-1));
assert.deepEqual(uploadSwitches, [`${uploadChain}_B`, `${uploadChain}_A`]);
const failedMangleSwitch = calls.find(([command, , options]) => (
command === 'iptables-restore' && options.input.includes(`-R ${downloadChain}`)
));
assert.match(failedMangleSwitch[2].input, new RegExp(`-R ${downloadChain} 1 -j ${downloadChain}_B`));
assert.match(failedMangleSwitch[2].input, new RegExp(`-R ${routeAccounting.directChain} 1 -j ${routeAccounting.directChain}_B`));
failReplacement = false;
observed = {
observedAt: '2026-08-07T12:02:00.000Z',
observations: [],
error: 'neighbors unavailable',
};
const stale = await service.refresh();
assert.equal(stale.generation, 'rules-a');
assert.equal(stale.source.error, 'neighbors unavailable');
assert.deepEqual(stale.devices, first.devices);
failCounters = true;
observed = {
observedAt: '2026-08-07T12:03:00.000Z',
observations: [observation('192.168.50.7')],
error: null,
};
const timedOut = await service.refresh();
assert.equal(timedOut.observedAt, '2026-08-07T12:02:00.000Z');
assert.match(timedOut.source.error, /counter read timed out/);
assert.deepEqual(timedOut.devices, first.devices);
assert.ok(calls.every(([, , options]) => options.timeout === 2_000));
});
test('an initial accounting failure does not publish fresh zero counters', async () => {
const service = createDeviceTrafficService({
observe: () => ({
observedAt: '2026-08-07T12:00:00.000Z',
observations: [observation('192.168.50.7')],
error: null,
}),
uploadChain,
downloadChain,
bypassCidrs: [],
proxyPort: 8080,
...routeAccounting,
run: (command) => ({
status: command === 'iptables-restore' ? 1 : 0,
stdout: '',
stderr: command === 'iptables-restore' ? 'prepare failed' : '',
}),
});
const snapshot = await service.refresh();
assert.equal(snapshot.observedAt, null);
assert.deepEqual(snapshot.direct, { uploadBytes: '0', downloadBytes: '0' });
assert.deepEqual(snapshot.devices, []);
assert.match(snapshot.source.error, /prepare failed/);
});
test('a 512-device refresh keeps a fixed subprocess count and a cached snapshot', async () => {
const observations = Array.from({ length: 512 }, (_, index) => observation(
`10.${Math.floor(index / 254)}.${Math.floor((index % 254) / 254)}.${(index % 254) + 1}`,
`02:00:${Math.floor(index / 256).toString(16).padStart(2, '0')}:${Math.floor(index / 16).toString(16).padStart(2, '0')}:${(index % 16).toString(16).padStart(2, '0')}:01`,
));
const calls = [];
let releaseRestore;
let blockedRestore = false;
const run = (command, args, options) => {
calls.push([command, args, options]);
if (command === 'iptables-restore' && !blockedRestore) {
blockedRestore = true;
return new Promise((resolve) => {
releaseRestore = () => resolve({ status: 0, stdout: '', stderr: '' });
});
}
return { status: 0, stdout: '', stderr: '' };
};
const generations = ['epoch', 'rules-a'];
const service = createDeviceTrafficService({
observe: () => ({ observedAt: '2026-08-07T12:00:00.000Z', observations, error: null }),
uploadChain,
downloadChain,
bypassCidrs: [],
proxyPort: 8080,
...routeAccounting,
run,
nextGeneration: () => generations.shift(),
});
const pending = service.refresh();
await new Promise((resolve) => setImmediate(resolve));
assert.equal(service.snapshot().observedAt, null);
assert.equal(calls.length, 1);
assert.equal(calls[0][0], 'iptables-restore');
assert.match(calls[0][2].input, /proxy-upload/);
releaseRestore();
const snapshot = await pending;
assert.equal(snapshot.devices.length, 512);
assert.deepEqual(
calls.map(([command]) => command),
['iptables-restore', 'iptables', 'iptables-restore', 'iptables-save', 'iptables-save'],
);
assert.match(calls[2][2].input, new RegExp(`-A ${downloadChain} -j ${downloadChain}_A`));
assert.match(calls[2][2].input, new RegExp(`-A ${routeAccounting.directChain} -j ${routeAccounting.directChain}_A`));
});
test('traffic service finalizes a detached slot once and keeps epoch totals monotonic', async () => {
const firstObservation = observation('192.168.50.7', '00:11:22:33:44:55');
const secondObservation = observation('192.168.50.8', '00:11:22:33:44:66');
const firstDevice = selectTrafficDevices([firstObservation])[0];
const secondDevice = selectTrafficDevices([secondObservation])[0];
let observed = {
observedAt: '2026-08-07T12:00:00.000Z',
observations: [firstObservation],
error: null,
};
const values = {
A: {
upload: '100', download: '200', proxyUpload: '30', proxyDownload: '40',
directUpload: '11', directDownload: '22', globalDirectUpload: '33', globalDirectDownload: '44',
},
B: {
upload: '5', download: '7', proxyUpload: '2', proxyDownload: '3',
directUpload: '1', directDownload: '2', globalDirectUpload: '3', globalDirectDownload: '4',
},
};
const keys = { A: firstDevice.key, B: secondDevice.key };
let failNextCounterRead = false;
const run = (command, args) => {
if (command !== 'iptables-save') return { status: 0, stdout: '', stderr: '' };
if (failNextCounterRead) {
failNextCounterRead = false;
return { status: null, stdout: '', stderr: '', error: new Error('retired slot read failed') };
}
const direction = args.includes('raw') ? 'upload' : 'download';
const proxyDirection = direction === 'upload' ? 'proxyUpload' : 'proxyDownload';
const proxyKind = direction === 'upload' ? 'proxy-upload' : 'proxy-download';
const tableChain = args.includes('raw') ? uploadChain : downloadChain;
return {
status: 0,
stdout: ['A', 'B'].flatMap((slot) => [
`[1:${values[slot][direction]}] -A ${tableChain}_${slot} -m comment --comment "harbor-traffic:${keys[slot]}:${direction}" -j RETURN`,
`[1:${values[slot][proxyDirection]}] -A ${tableChain}_${slot}_P -m comment --comment "harbor-traffic:${keys[slot]}:${proxyKind}" -j RETURN`,
...(direction === 'download' ? [
`[1:${values[slot].directUpload}] -A ${routeAccounting.directChain}_${slot} -m comment --comment "harbor-traffic:${keys[slot]}:direct-upload"`,
`[1:${values[slot].directDownload}] -A ${tableChain}_${slot} -m comment --comment "harbor-traffic:${keys[slot]}:direct-download"`,
`[1:${values[slot].globalDirectUpload}] -A ${routeAccounting.directChain}_${slot} -m comment --comment "harbor-traffic:global:direct-upload"`,
`[1:${values[slot].globalDirectDownload}] -A ${tableChain}_${slot} -m comment --comment "harbor-traffic:global:direct-download"`,
] : []),
]).join('\n'),
stderr: '',
};
};
const generations = ['epoch-1', 'rules-a', 'rules-b'];
const service = createDeviceTrafficService({
observe: async () => observed,
uploadChain,
downloadChain,
bypassCidrs: [],
proxyPort: 8080,
...routeAccounting,
run,
nextGeneration: () => generations.shift(),
});
const first = await service.refresh();
assert.equal(first.epoch, 'epoch-1');
assert.equal(first.generation, 'rules-a');
assert.deepEqual(first.direct, { uploadBytes: '33', downloadBytes: '44' });
assert.deepEqual(first.devices.map(({ mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes }) => ({
mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes,
})), [{
mac: firstObservation.mac,
uploadBytes: '100',
downloadBytes: '200',
proxyUploadBytes: '30',
proxyDownloadBytes: '40',
directUploadBytes: '11',
directDownloadBytes: '22',
}]);
values.A = {
upload: '130', download: '240', proxyUpload: '35', proxyDownload: '48',
directUpload: '15', directDownload: '28', globalDirectUpload: '39', globalDirectDownload: '52',
};
observed = {
observedAt: '2026-08-07T12:01:00.000Z',
observations: [secondObservation],
error: null,
};
failNextCounterRead = true;
const pending = await service.refresh();
assert.equal(pending.epoch, 'epoch-1');
assert.equal(pending.generation, 'rules-b');
assert.match(pending.source.error, /retired slot read failed/);
assert.deepEqual(pending.direct, { uploadBytes: '36', downloadBytes: '48' });
assert.deepEqual(pending.devices.map(({ mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes }) => ({
mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes,
})), [
{ mac: firstObservation.mac, uploadBytes: '100', downloadBytes: '200', proxyUploadBytes: '30', proxyDownloadBytes: '40', directUploadBytes: '11', directDownloadBytes: '22' },
{ mac: secondObservation.mac, uploadBytes: '5', downloadBytes: '7', proxyUploadBytes: '2', proxyDownloadBytes: '3', directUploadBytes: '1', directDownloadBytes: '2' },
]);
const finalized = await service.refresh();
assert.equal(finalized.generation, 'rules-b');
assert.equal(finalized.source.error, null);
assert.deepEqual(finalized.direct, { uploadBytes: '42', downloadBytes: '56' });
assert.deepEqual(finalized.devices.map(({ mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes }) => ({
mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes,
})), [
{ mac: firstObservation.mac, uploadBytes: '130', downloadBytes: '240', proxyUploadBytes: '35', proxyDownloadBytes: '48', directUploadBytes: '15', directDownloadBytes: '28' },
{ mac: secondObservation.mac, uploadBytes: '5', downloadBytes: '7', proxyUploadBytes: '2', proxyDownloadBytes: '3', directUploadBytes: '1', directDownloadBytes: '2' },
]);
values.B = {
upload: '15', download: '17', proxyUpload: '4', proxyDownload: '6',
directUpload: '4', directDownload: '6', globalDirectUpload: '8', globalDirectDownload: '10',
};
const polled = await service.refresh();
assert.deepEqual(polled.direct, { uploadBytes: '47', downloadBytes: '62' });
assert.deepEqual(polled.devices.map(({ mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes }) => ({
mac, uploadBytes, downloadBytes, proxyUploadBytes, proxyDownloadBytes, directUploadBytes, directDownloadBytes,
})), [
{ mac: firstObservation.mac, uploadBytes: '130', downloadBytes: '240', proxyUploadBytes: '35', proxyDownloadBytes: '48', directUploadBytes: '15', directDownloadBytes: '28' },
{ mac: secondObservation.mac, uploadBytes: '15', downloadBytes: '17', proxyUploadBytes: '4', proxyDownloadBytes: '6', directUploadBytes: '4', directDownloadBytes: '6' },
]);
});