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, }); }); });