const assert = require('assert');
|
const { MetricAggregator } = require('../metric-aggregator');
|
|
describe('MetricAggregator', () => {
|
it('merges latest metrics by device and always includes n', async () => {
|
const flushed = [];
|
const aggregator = new MetricAggregator({
|
mode: 'batch',
|
onFlush: async (device, payload) => {
|
flushed.push({ device, payload });
|
return true;
|
},
|
});
|
|
aggregator.ingest({ deviceId: 'JH-001', ip: '127.0.0.1' }, { F: 24.6 });
|
aggregator.ingest({ deviceId: 'JH-001', ip: '127.0.0.1' }, { A: 0, F: 25.1 });
|
await aggregator.flush({ reason: 'test' });
|
|
assert.strictEqual(flushed.length, 1);
|
assert.deepStrictEqual(flushed[0].payload, {
|
n: 'JH-001',
|
F: 25.1,
|
A: 0,
|
});
|
});
|
|
it('flushes only dirty devices', async () => {
|
const flushed = [];
|
const aggregator = new MetricAggregator({
|
mode: 'batch',
|
onFlush: async (device, payload) => {
|
flushed.push({ deviceId: device.deviceId, payload });
|
return true;
|
},
|
});
|
|
aggregator.ingest({ deviceId: 'JH-001' }, { F: 24.6 });
|
await aggregator.flush({ reason: 'test-1' });
|
await aggregator.flush({ reason: 'test-2' });
|
|
assert.strictEqual(flushed.length, 1);
|
assert.strictEqual(flushed[0].deviceId, 'JH-001');
|
});
|
|
it('keeps device dirty when flush callback reports failure', async () => {
|
const flushed = [];
|
const aggregator = new MetricAggregator({
|
mode: 'batch',
|
onFlush: async (_device, payload) => {
|
flushed.push(payload);
|
return false;
|
},
|
});
|
|
aggregator.ingest({ deviceId: 'JH-001' }, { F: 24.6 });
|
await aggregator.flush({ reason: 'test' });
|
|
const snapshot = aggregator.getSnapshot('JH-001');
|
assert.strictEqual(flushed.length, 1);
|
assert.strictEqual(snapshot.dirty, true);
|
});
|
|
it('supports immediate mode', async () => {
|
const flushed = [];
|
const aggregator = new MetricAggregator({
|
mode: 'immediate',
|
onFlush: async (_device, payload) => {
|
flushed.push(payload);
|
return true;
|
},
|
});
|
|
await aggregator.ingest({ deviceId: 'JH-001' }, { F: 24.6 });
|
|
assert.strictEqual(flushed.length, 1);
|
assert.deepStrictEqual(flushed[0], {
|
n: 'JH-001',
|
F: 24.6,
|
});
|
});
|
|
it('waits a full flush interval before timer flushes newly cached data', async () => {
|
const originalNow = Date.now;
|
let now = 1_000;
|
Date.now = () => now;
|
|
try {
|
const flushed = [];
|
const aggregator = new MetricAggregator({
|
mode: 'batch',
|
flushIntervalMs: 60_000,
|
onFlush: async (_device, payload) => {
|
flushed.push(payload);
|
return true;
|
},
|
});
|
|
aggregator.ingest({ deviceId: 'JH-001' }, { F: 24.6 });
|
|
now = 30_000;
|
await aggregator.flush({ reason: 'timer' });
|
assert.strictEqual(flushed.length, 0);
|
|
now = 61_000;
|
await aggregator.flush({ reason: 'timer' });
|
assert.strictEqual(flushed.length, 1);
|
assert.deepStrictEqual(flushed[0], {
|
n: 'JH-001',
|
F: 24.6,
|
});
|
} finally {
|
Date.now = originalNow;
|
}
|
});
|
|
it('preserves immediate flush intent when a second flush is queued during an ongoing flush', async () => {
|
const flushed = [];
|
let releaseFirstFlush = null;
|
const firstFlushDone = new Promise((resolve) => {
|
releaseFirstFlush = resolve;
|
});
|
|
const aggregator = new MetricAggregator({
|
mode: 'batch',
|
onFlush: async (_device, payload) => {
|
flushed.push(payload);
|
|
if (flushed.length === 1) {
|
await firstFlushDone;
|
}
|
|
return true;
|
},
|
});
|
|
aggregator.ingest({ deviceId: 'JH-001' }, { F: 24.6 });
|
const firstFlushPromise = aggregator.flush({ reason: 'manual' });
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
aggregator.ingest({ deviceId: 'JH-001' }, { N: 120, O: 80, P: 89 });
|
const queuedFlushResult = await aggregator.flush({ reason: 'blood-pressure', deviceId: 'JH-001' });
|
|
assert.strictEqual(queuedFlushResult, false);
|
|
releaseFirstFlush();
|
await firstFlushPromise;
|
|
assert.strictEqual(flushed.length, 2);
|
assert.deepStrictEqual(flushed[0], {
|
n: 'JH-001',
|
F: 24.6,
|
});
|
assert.deepStrictEqual(flushed[1], {
|
n: 'JH-001',
|
F: 24.6,
|
N: 120,
|
O: 80,
|
P: 89,
|
});
|
});
|
});
|