23f7624596
ADR-166 MCP Bridge Security Lock / Static-source security lock (push) Failing after 0s
ADR-166 MCP Bridge Security Lock / Compose default binds loopback + Mongo has auth (push) Failing after 2s
CodeQL Advanced / Analyze (rust) (push) Failing after 0s
ADR-166 MCP Bridge Security Lock / plugin-agent-federation bindHost default (push) Failing after 1s
ADR-166 MCP Bridge Security Lock / Runtime behavior — 401 + terminal gate + fail-closed (push) Failing after 4s
business-pods-smoke / smoke (push) Failing after 1s
all-plugins-smoke / smoke-all (push) Failing after 2s
CI/CD Pipeline / Security & Code Quality (push) Failing after 1s
CI/CD Pipeline / Test Suite (ubuntu-latest) (push) Failing after 1s
CI/CD Pipeline / Build & Package (macos-latest) (push) Has been skipped
CI/CD Pipeline / Build & Package (ubuntu-latest) (push) Has been skipped
CI/CD Pipeline / Build & Package (windows-latest) (push) Has been skipped
CI/CD Pipeline / Documentation & Examples (push) Failing after 1s
Clone Tracker (14-day rolling) / Snapshot clones for ruflo ecosystem (push) Failing after 1s
CodeQL Advanced / Analyze (actions) (push) Failing after 1s
CodeQL Advanced / Analyze (javascript-typescript) (push) Failing after 1s
federation-peer-rust / stable-noop (push) Failing after 1s
metaharness-ci / score (push) Failing after 1s
metaharness-ci / router-compat (push) Failing after 0s
metaharness-ci / similarity-tests (push) Failing after 0s
no-agentbbs-smoke / smoke-without-agentbbs (push) Failing after 1s
V3 CI/CD Pipeline / Build V3 (windows-latest) (push) Has been skipped
codex-integration-audit / Codex integration audit (push) Failing after 1s
helpers-manifest-guard / guard (push) Failing after 1s
🔗 Cross-Agent Integration Tests / 🤝 Agent Coordination Tests (push) Has been skipped
🔗 Cross-Agent Integration Tests / 🧠 Memory Sharing Integration (push) Has been skipped
🔗 Cross-Agent Integration Tests / 🛡️ Fault Tolerance Tests (push) Has been skipped
🔗 Cross-Agent Integration Tests / ⚡ Performance Integration Tests (push) Has been skipped
metaharness-ci / mcp-scan (push) Failing after 1s
metaharness-ci / eject-dryrun (push) Failing after 1s
metaharness-ci / metaharness-real-data (push) Failing after 0s
no-cli-optdep-bloat-2561 / guard (push) Failing after 1s
no-metaharness-smoke / smoke-without-metaharness (push) Failing after 1s
no-phantom-agentic-flow-subpath / guard (push) Failing after 1s
🔄 Automated Rollback Manager / 🚨 Failure Detection (push) Failing after 1s
V3 CI/CD Pipeline / Plugin hooks smoke / ubuntu-latest / Node 22 (push) Failing after 1s
V3 CI/CD Pipeline / ruflo-graph-intelligence build + test smoke (#2044, ADR-123) (push) Failing after 1s
CVE Audit Gate / Audit root (critical-blocking) (push) Failing after 2s
cost-tracker-smoke / smoke (push) Failing after 3s
oia-audit-weekly / audit (push) Failing after 2s
ruflo-agent-smoke / ruflo-agent structural smoke (push) Failing after 1s
📊 Status Badges Update / 📊 Update Status Badges (push) Failing after 1s
V3 CI/CD Pipeline / Static regression guards (#2267 YAML + (push) Failing after 1s
V3 CI/CD Pipeline / Test V3 Packages (push) Failing after 0s
V3 CI/CD Pipeline / agent_execute provider routing smoke (#2042) (push) Failing after 0s
CVE Audit Gate / Audit v3 (critical-blocking) (push) Failing after 1s
federation-peer-rust / stable-native (push) Failing after 2s
🔗 Cross-Agent Integration Tests / 🚀 Integration Test Setup (push) Failing after 2s
neural-trader-smoke / runtime-smoke (push) Failing after 1s
V3 CI/CD Pipeline / Build V3 (macos-latest) (push) Has been skipped
V3 CI/CD Pipeline / Build V3 (ubuntu-latest) (push) Has been skipped
V3 CI/CD Pipeline / Type Check V3 (push) Failing after 1s
V3 CI/CD Pipeline / Smoke (no better-sqlite3) / ubuntu-latest / Node 24 (push) Failing after 1s
V3 CI/CD Pipeline / Smoke (no better-sqlite3) / ubuntu-latest / Node 22 (push) Failing after 2s
V3 CI/CD Pipeline / browser rvf create flag smoke (#2015) (push) Failing after 0s
V3 CI/CD Pipeline / Dependency review (#2046) (push) Has been skipped
V3 CI/CD Pipeline / Supply-chain audit (#2046) (push) Failing after 0s
V3 CI/CD Pipeline / witness marker drift smoke (#2021) (push) Failing after 1s
V3 CI/CD Pipeline / neural-trader portfolio CG smoke (#2068, ADR-126 Phase 3) (push) Failing after 1s
V3 CI/CD Pipeline / neural-trader backtest signing smoke (#2068, ADR-126 Phase 4) (push) Failing after 1s
V3 CI/CD Pipeline / kg-extract type-import classification smoke (#2049) (push) Failing after 0s
V3 CI/CD Pipeline / witness verify precondition smoke (#1880) (push) Failing after 2s
V3 CI/CD Pipeline / neural-trader pipeline risk-gate smoke (#2068, ADR-126 Phase 5) (push) Failing after 0s
V3 CI/CD Pipeline / neural-trader feature attribution smoke (#2068, ADR-126 Phase 6) (push) Failing after 0s
V3 CI/CD Pipeline / plugin-registry signature verification smoke (#1922, CWE-347) (push) Failing after 4s
V3 CI/CD Pipeline / memory stats legacy-DB smoke (#2120) (push) Failing after 4s
V3 CI/CD Pipeline / github deprecated actions smoke (#2089, ADR-127 Phase 3) (push) Failing after 1s
V3 CI/CD Pipeline / graph query + pathfinder smoke (ADR-130 P2+P5) (push) Has been skipped
V3 CI/CD Pipeline / graph trajectory hooks smoke (ADR-130 P3) (push) Has been skipped
V3 CI/CD Pipeline / graph plugin adapter smoke (ADR-130 P4) (push) Has been skipped
V3 CI/CD Pipeline / graph benchmark (ADR-130 P6) (push) Has been skipped
V3 CI/CD Pipeline / statusline generator delegation smoke (#2195) (push) Failing after 1s
V3 CI/CD Pipeline / wizard init regression guard (#2206 (push) Failing after 1s
V3 CI/CD Pipeline / memory no-stray-db smoke (ADR-125 P7) (push) Failing after 1s
V3 CI/CD Pipeline / github-safe injection smoke (#2089, ADR-127 Phase 1) (push) Failing after 1s
V3 CI/CD Pipeline / github actions pin smoke (#2089, ADR-127 Phase 1) (push) Failing after 1s
V3 CI/CD Pipeline / github attribution opt-in smoke (#2089, ADR-127 Phase 4) (push) Failing after 1s
V3 CI/CD Pipeline / pre-bash hook safety smoke (#2017) (push) Failing after 1s
V3 CI/CD Pipeline / Memory import smoke / ubuntu-latest (push) Failing after 0s
V3 CI/CD Pipeline / MCP protocol smoke / ubuntu-latest (push) Failing after 2s
V3 CI/CD Pipeline / ruvllm WASM auto-init smoke (#2086) (push) Failing after 4s
V3 CI/CD Pipeline / MCP paired-tool round-trip smoke (#1889) (push) Failing after 1s
V3 CI/CD Pipeline / Plugin package install-safety (#1902/#1903/#1904) (push) Failing after 1s
V3 CI/CD Pipeline / Tool description discoverability (ADR-112) (push) Failing after 3s
V3 CI/CD Pipeline / CLI npx-install smoke (#1147 / (22) (push) Failing after 1s
V3 CI/CD Pipeline / CLI npx-install smoke (#1147 / (24) (push) Failing after 1s
V3 CI/CD Pipeline / Windows hook shim smoke (#2132) / ubuntu-latest (push) Failing after 2s
V3 CI/CD Pipeline / Windows hook execution smoke (#2132) / ubuntu-latest (push) Failing after 1s
V3 CI/CD Pipeline / Windows init hooks smoke (#2132) / ubuntu-latest (push) Failing after 1s
V3 CI/CD Pipeline / Vector-index dimension audit (#1947) (push) Failing after 0s
V3 CI/CD Pipeline / Hook-command install safety (#1921) (push) Failing after 1s
V3 CI/CD Pipeline / ToolOutputGuardrail smoke (ADR-131, (push) Failing after 1s
V3 CI/CD Pipeline / init-bundle invariants smoke (#2095, ADR-128 Phase 5) (push) Failing after 1s
V3 CI/CD Pipeline / wasm provider bridge smoke (ADR-129 P1) (push) Failing after 2s
V3 CI/CD Pipeline / wasm gallery CRUD smoke (ADR-129 P3) (push) Failing after 1s
V3 CI/CD Pipeline / wasm plugin bridge smoke (ADR-129 P4) (push) Failing after 0s
V3 CI/CD Pipeline / wasm compose smoke (ADR-129 P2) (push) Failing after 4s
V3 CI/CD Pipeline / graph schema smoke (ADR-130 P1) (push) Failing after 0s
Validate Marketplace / validate (push) Failing after 1s
🔍 Verification Pipeline / 🚀 Setup Verification (push) Failing after 1s
🔍 Verification Pipeline / 🛡️ Security Verification (push) Has been skipped
🔍 Verification Pipeline / 📝 Code Quality (push) Has been skipped
🔍 Verification Pipeline / 🧪 Test Verification (${{ matrix.os }}, Node ${{ matrix.node }}) (push) Has been skipped
🔍 Verification Pipeline / 🏗️ Build Verification (push) Has been skipped
🔍 Verification Pipeline / 📚 Documentation Verification (push) Has been skipped
CVE Audit Gate / High-severity report (warn only) (push) Has been cancelled
🔄 Automated Rollback Manager / 🔄 Execute Rollback (push) Has been cancelled
🔄 Automated Rollback Manager / ✅ Post-Rollback Verification (push) Has been cancelled
🔄 Automated Rollback Manager / 📊 Rollback Monitoring (push) Has been cancelled
V3 CI/CD Pipeline / Windows init hooks smoke (#2132) / windows-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook execution smoke (#2132) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook execution smoke (#2132) / windows-latest (push) Has been cancelled
🔄 Automated Rollback Manager / ⏳ Manual Rollback Approval (push) Has been cancelled
V3 CI/CD Pipeline / MCP protocol smoke / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Memory import smoke / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook shim smoke (#2132) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook shim smoke (#2132) / windows-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows init hooks smoke (#2132) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Witness verify (signed manifest) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Witness verify (signed manifest) / ubuntu-latest (push) Has been cancelled
V3 CI/CD Pipeline / Witness verify (signed manifest) / windows-latest (push) Has been cancelled
V3 CI/CD Pipeline / Publish to npm (alpha) (push) Has been cancelled
V3 CI/CD Pipeline / Smoke (no better-sqlite3) / macos-latest / Node 22 (push) Has been cancelled
V3 CI/CD Pipeline / Plugin hooks smoke / macos-latest / Node 22 (push) Has been cancelled
CI/CD Pipeline / Deploy & Release (push) Has been cancelled
CI/CD Pipeline / CI Status (push) Has been cancelled
🔗 Cross-Agent Integration Tests / 📊 Integration Test Report (push) Has been cancelled
🔄 Automated Rollback Manager / 🔍 Pre-Rollback Validation (push) Has been cancelled
🔍 Verification Pipeline / ⚡ Performance Verification (push) Has been cancelled
🔍 Verification Pipeline / 📊 Verification Report (push) Has been cancelled
644 lines
19 KiB
TypeScript
644 lines
19 KiB
TypeScript
/**
|
|
* Tests for RvfEventLog (ADR-057 Phase 2)
|
|
*
|
|
* Covers: initialize, append, getEvents, getAllEvents, snapshots,
|
|
* stats, persistence, close, and edge cases.
|
|
*/
|
|
|
|
import { describe, it, beforeEach, afterEach } from 'node:test';
|
|
import assert from 'node:assert/strict';
|
|
import crypto from 'node:crypto';
|
|
import { existsSync, mkdirSync, rmSync, readFileSync, writeFileSync, appendFileSync } from 'node:fs';
|
|
import { join } from 'node:path';
|
|
import { tmpdir } from 'node:os';
|
|
|
|
// Use dynamic import — tsx handles TS source directly.
|
|
const { RvfEventLog } = await import(
|
|
'../v3/@claude-flow/shared/src/events/rvf-event-log.ts'
|
|
);
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Helpers
|
|
// ---------------------------------------------------------------------------
|
|
|
|
type DomainEvent = {
|
|
id: string;
|
|
type: string;
|
|
aggregateId: string;
|
|
aggregateType: string;
|
|
version: number;
|
|
timestamp: number;
|
|
source: string;
|
|
payload: Record<string, unknown>;
|
|
metadata?: Record<string, unknown>;
|
|
causationId?: string;
|
|
correlationId?: string;
|
|
};
|
|
|
|
type EventSnapshot = {
|
|
aggregateId: string;
|
|
aggregateType: string;
|
|
version: number;
|
|
state: Record<string, unknown>;
|
|
timestamp: number;
|
|
};
|
|
|
|
let tmpCounter = 0;
|
|
|
|
function makeTmpDir(): string {
|
|
const dir = join(tmpdir(), `rvf-test-${Date.now()}-${++tmpCounter}`);
|
|
mkdirSync(dir, { recursive: true });
|
|
return dir;
|
|
}
|
|
|
|
function makeEvent(
|
|
aggregateId: string,
|
|
type: string,
|
|
timestampOverride?: number,
|
|
): DomainEvent {
|
|
return {
|
|
id: crypto.randomUUID(),
|
|
type,
|
|
aggregateId,
|
|
aggregateType: 'test' as any,
|
|
timestamp: timestampOverride ?? Date.now(),
|
|
version: 0,
|
|
source: 'swarm',
|
|
payload: {},
|
|
metadata: { correlationId: crypto.randomUUID(), causationId: '' },
|
|
};
|
|
}
|
|
|
|
function makeSnapshot(
|
|
aggregateId: string,
|
|
version: number,
|
|
state: Record<string, unknown> = {},
|
|
): EventSnapshot {
|
|
return {
|
|
aggregateId,
|
|
aggregateType: 'test' as any,
|
|
version,
|
|
state,
|
|
timestamp: Date.now(),
|
|
};
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Test state
|
|
// ---------------------------------------------------------------------------
|
|
|
|
let dir: string;
|
|
let logPath: string;
|
|
|
|
beforeEach(() => {
|
|
dir = makeTmpDir();
|
|
logPath = join(dir, 'events.rvf');
|
|
});
|
|
|
|
afterEach(() => {
|
|
if (existsSync(dir)) {
|
|
rmSync(dir, { recursive: true, force: true });
|
|
}
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 1. Initialize
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog#initialize', () => {
|
|
it('creates event file with magic header', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
assert.ok(existsSync(logPath));
|
|
const buf = readFileSync(logPath);
|
|
assert.equal(buf.subarray(0, 4).toString(), 'RVFL');
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('creates snapshot file alongside event file', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const snapPath = logPath.replace(/\.rvf$/, '.snap.rvf');
|
|
assert.ok(existsSync(snapPath));
|
|
const buf = readFileSync(snapPath);
|
|
assert.equal(buf.subarray(0, 4).toString(), 'RVFL');
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('rebuilds indexes on re-open', async () => {
|
|
const log1 = new RvfEventLog({ logPath });
|
|
await log1.initialize();
|
|
await log1.append(makeEvent('agg-1', 'created'));
|
|
await log1.append(makeEvent('agg-1', 'updated'));
|
|
await log1.close();
|
|
|
|
const log2 = new RvfEventLog({ logPath });
|
|
await log2.initialize();
|
|
|
|
const events = await log2.getEvents('agg-1');
|
|
assert.equal(events.length, 2);
|
|
assert.equal(events[0].version, 1);
|
|
assert.equal(events[1].version, 2);
|
|
|
|
await log2.close();
|
|
});
|
|
|
|
it('emits initialized event', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
let emitted = false;
|
|
log.on('initialized', () => { emitted = true; });
|
|
await log.initialize();
|
|
|
|
assert.ok(emitted);
|
|
await log.close();
|
|
});
|
|
|
|
it('is idempotent when called twice', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
await log.initialize();
|
|
|
|
assert.ok(existsSync(logPath));
|
|
await log.close();
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 2. Append
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog#append', () => {
|
|
it('assigns incrementing versions per aggregate', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const e1 = makeEvent('agg-A', 'step');
|
|
const e2 = makeEvent('agg-A', 'step');
|
|
const e3 = makeEvent('agg-B', 'step');
|
|
|
|
await log.append(e1);
|
|
await log.append(e2);
|
|
await log.append(e3);
|
|
|
|
assert.equal(e1.version, 1);
|
|
assert.equal(e2.version, 2);
|
|
assert.equal(e3.version, 1); // separate aggregate
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('persists to disk immediately', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
await log.append(makeEvent('agg-1', 'op'));
|
|
|
|
// File should be larger than just the 4-byte magic header.
|
|
const size = readFileSync(logPath).length;
|
|
assert.ok(size > 4);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('emits event:appended', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const captured: any[] = [];
|
|
log.on('event:appended', (e: any) => captured.push(e));
|
|
|
|
const ev = makeEvent('agg-1', 'op');
|
|
await log.append(ev);
|
|
|
|
assert.equal(captured.length, 1);
|
|
assert.equal(captured[0].id, ev.id);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('throws when not initialized', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await assert.rejects(
|
|
() => log.append(makeEvent('x', 'y')),
|
|
/not initialized/i,
|
|
);
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 3. GetEvents
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog#getEvents', () => {
|
|
it('filters by aggregateId', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
await log.append(makeEvent('agg-A', 'x'));
|
|
await log.append(makeEvent('agg-B', 'x'));
|
|
await log.append(makeEvent('agg-A', 'y'));
|
|
|
|
const result = await log.getEvents('agg-A');
|
|
assert.equal(result.length, 2);
|
|
assert.ok(result.every((e: any) => e.aggregateId === 'agg-A'));
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('filters by fromVersion', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
await log.append(makeEvent('agg-A', 'a'));
|
|
await log.append(makeEvent('agg-A', 'b'));
|
|
await log.append(makeEvent('agg-A', 'c'));
|
|
|
|
const result = await log.getEvents('agg-A', 2);
|
|
assert.equal(result.length, 2);
|
|
assert.equal(result[0].version, 2);
|
|
assert.equal(result[1].version, 3);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('returns empty array for unknown aggregate', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
await log.append(makeEvent('agg-A', 'x'));
|
|
|
|
const result = await log.getEvents('nonexistent');
|
|
assert.deepEqual(result, []);
|
|
|
|
await log.close();
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 4. GetAllEvents
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog#getAllEvents', () => {
|
|
it('returns all events sorted by timestamp when no filter', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const base = Date.now();
|
|
await log.append(makeEvent('a', 'x', base + 20));
|
|
await log.append(makeEvent('b', 'y', base + 10));
|
|
await log.append(makeEvent('c', 'z', base + 30));
|
|
|
|
const result = await log.getAllEvents();
|
|
assert.equal(result.length, 3);
|
|
assert.ok(result[0].timestamp <= result[1].timestamp);
|
|
assert.ok(result[1].timestamp <= result[2].timestamp);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('filters by eventTypes', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
await log.append(makeEvent('a', 'created'));
|
|
await log.append(makeEvent('a', 'updated'));
|
|
await log.append(makeEvent('a', 'deleted'));
|
|
|
|
const result = await log.getAllEvents({ eventTypes: ['created', 'deleted'] });
|
|
assert.equal(result.length, 2);
|
|
const types = result.map((e: any) => e.type).sort();
|
|
assert.deepEqual(types, ['created', 'deleted']);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('filters by timestamps (after/before)', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const base = 1000000;
|
|
await log.append(makeEvent('a', 'x', base + 100));
|
|
await log.append(makeEvent('a', 'y', base + 200));
|
|
await log.append(makeEvent('a', 'z', base + 300));
|
|
|
|
const result = await log.getAllEvents({
|
|
afterTimestamp: base + 100,
|
|
beforeTimestamp: base + 300,
|
|
});
|
|
assert.equal(result.length, 1);
|
|
assert.equal(result[0].timestamp, base + 200);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('supports pagination (offset + limit)', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const base = Date.now();
|
|
for (let i = 0; i < 10; i++) {
|
|
await log.append(makeEvent('a', 'step', base + i));
|
|
}
|
|
|
|
const page = await log.getAllEvents({ offset: 3, limit: 4 });
|
|
assert.equal(page.length, 4);
|
|
|
|
await log.close();
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 5. Snapshots
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog#snapshots', () => {
|
|
it('saveSnapshot persists and getSnapshot retrieves', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const snap = makeSnapshot('agg-1', 5, { counter: 42 });
|
|
await log.saveSnapshot(snap);
|
|
|
|
const retrieved = await log.getSnapshot('agg-1');
|
|
assert.ok(retrieved);
|
|
assert.equal(retrieved!.aggregateId, 'agg-1');
|
|
assert.equal(retrieved!.version, 5);
|
|
assert.deepEqual(retrieved!.state, { counter: 42 });
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('returns null for missing snapshot', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const result = await log.getSnapshot('nonexistent');
|
|
assert.equal(result, null);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('multiple snapshots: latest wins', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
await log.saveSnapshot(makeSnapshot('agg-1', 3, { v: 'old' }));
|
|
await log.saveSnapshot(makeSnapshot('agg-1', 7, { v: 'new' }));
|
|
|
|
const snap = await log.getSnapshot('agg-1');
|
|
assert.equal(snap!.version, 7);
|
|
assert.deepEqual(snap!.state, { v: 'new' });
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('snapshots survive close and re-initialize', async () => {
|
|
const log1 = new RvfEventLog({ logPath });
|
|
await log1.initialize();
|
|
await log1.saveSnapshot(makeSnapshot('agg-1', 10, { val: 99 }));
|
|
await log1.close();
|
|
|
|
const log2 = new RvfEventLog({ logPath });
|
|
await log2.initialize();
|
|
const snap = await log2.getSnapshot('agg-1');
|
|
assert.ok(snap);
|
|
assert.equal(snap!.version, 10);
|
|
assert.deepEqual(snap!.state, { val: 99 });
|
|
|
|
await log2.close();
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 6. Stats
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog#getStats', () => {
|
|
it('returns correct statistics', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const base = 1000000;
|
|
await log.append(makeEvent('agg-1', 'created', base + 10));
|
|
await log.append(makeEvent('agg-1', 'updated', base + 20));
|
|
await log.append(makeEvent('agg-2', 'created', base + 30));
|
|
await log.saveSnapshot(makeSnapshot('agg-1', 2));
|
|
|
|
const stats = await log.getStats();
|
|
assert.equal(stats.totalEvents, 3);
|
|
assert.equal(stats.eventsByType['created'], 2);
|
|
assert.equal(stats.eventsByType['updated'], 1);
|
|
assert.equal(stats.eventsByAggregate['agg-1'], 2);
|
|
assert.equal(stats.eventsByAggregate['agg-2'], 1);
|
|
assert.equal(stats.oldestEvent, base + 10);
|
|
assert.equal(stats.newestEvent, base + 30);
|
|
assert.equal(stats.snapshotCount, 1);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('returns nulls for empty log', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const stats = await log.getStats();
|
|
assert.equal(stats.totalEvents, 0);
|
|
assert.equal(stats.oldestEvent, null);
|
|
assert.equal(stats.newestEvent, null);
|
|
|
|
await log.close();
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 7. Persistence
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog persistence', () => {
|
|
it('data survives close and re-initialize cycle', async () => {
|
|
const log1 = new RvfEventLog({ logPath });
|
|
await log1.initialize();
|
|
await log1.append(makeEvent('agg-1', 'created'));
|
|
await log1.append(makeEvent('agg-1', 'updated'));
|
|
await log1.append(makeEvent('agg-2', 'created'));
|
|
await log1.close();
|
|
|
|
const log2 = new RvfEventLog({ logPath });
|
|
await log2.initialize();
|
|
|
|
const all = await log2.getAllEvents();
|
|
assert.equal(all.length, 3);
|
|
|
|
// Versions should continue from where they left off.
|
|
const e = makeEvent('agg-1', 'deleted');
|
|
await log2.append(e);
|
|
assert.equal(e.version, 3);
|
|
|
|
await log2.close();
|
|
});
|
|
|
|
it('truncated records are handled gracefully', async () => {
|
|
const log1 = new RvfEventLog({ logPath });
|
|
await log1.initialize();
|
|
await log1.append(makeEvent('agg-1', 'created'));
|
|
await log1.close();
|
|
|
|
// Simulate a truncated write: append a length prefix that claims
|
|
// more bytes than actually exist.
|
|
const lengthBuf = Buffer.allocUnsafe(4);
|
|
lengthBuf.writeUInt32BE(99999, 0);
|
|
appendFileSync(logPath, Buffer.concat([lengthBuf, Buffer.from('partial')]));
|
|
|
|
// Re-open should recover the valid event and skip the truncated one.
|
|
const log2 = new RvfEventLog({ logPath });
|
|
await log2.initialize();
|
|
|
|
const events = await log2.getEvents('agg-1');
|
|
assert.equal(events.length, 1);
|
|
assert.equal(events[0].type, 'created');
|
|
|
|
await log2.close();
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 8. Close
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog#close', () => {
|
|
it('clears in-memory state', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
await log.append(makeEvent('agg-1', 'x'));
|
|
await log.close();
|
|
|
|
// After close, operations should fail because initialized is false.
|
|
await assert.rejects(() => log.getEvents('agg-1'), /not initialized/i);
|
|
await assert.rejects(() => log.getAllEvents(), /not initialized/i);
|
|
await assert.rejects(() => log.getStats(), /not initialized/i);
|
|
});
|
|
|
|
it('prevents append after close', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
await log.close();
|
|
|
|
await assert.rejects(
|
|
() => log.append(makeEvent('x', 'y')),
|
|
/not initialized/i,
|
|
);
|
|
});
|
|
|
|
it('emits shutdown event', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
let emitted = false;
|
|
log.on('shutdown', () => { emitted = true; });
|
|
await log.close();
|
|
|
|
assert.ok(emitted);
|
|
});
|
|
|
|
it('is idempotent', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
await log.close();
|
|
await log.close(); // should not throw
|
|
});
|
|
});
|
|
|
|
// ===========================================================================
|
|
// 9. Edge Cases
|
|
// ===========================================================================
|
|
|
|
describe('RvfEventLog edge cases', () => {
|
|
it('empty log returns empty arrays', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
assert.deepEqual(await log.getEvents('anything'), []);
|
|
assert.deepEqual(await log.getAllEvents(), []);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('corrupt JSON records are skipped', async () => {
|
|
// Create a valid file with magic header, then inject bad JSON.
|
|
writeFileSync(logPath, Buffer.from('RVFL'));
|
|
const badPayload = Buffer.from('{{{invalid json');
|
|
const lenBuf = Buffer.allocUnsafe(4);
|
|
lenBuf.writeUInt32BE(badPayload.length, 0);
|
|
appendFileSync(logPath, Buffer.concat([lenBuf, badPayload]));
|
|
|
|
// Now append a valid record manually.
|
|
const validEvent = makeEvent('agg-1', 'ok');
|
|
validEvent.version = 1;
|
|
const validJson = Buffer.from(JSON.stringify(validEvent), 'utf8');
|
|
const validLen = Buffer.allocUnsafe(4);
|
|
validLen.writeUInt32BE(validJson.length, 0);
|
|
appendFileSync(logPath, Buffer.concat([validLen, validJson]));
|
|
|
|
// Also create the snapshot file so initialize does not fail.
|
|
const snapPath = logPath.replace(/\.rvf$/, '.snap.rvf');
|
|
writeFileSync(snapPath, Buffer.from('RVFL'));
|
|
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const events = await log.getEvents('agg-1');
|
|
assert.equal(events.length, 1);
|
|
assert.equal(events[0].type, 'ok');
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('handles large number of events', async () => {
|
|
const log = new RvfEventLog({ logPath });
|
|
await log.initialize();
|
|
|
|
const count = 500;
|
|
for (let i = 0; i < count; i++) {
|
|
await log.append(makeEvent('bulk', 'tick', Date.now() + i));
|
|
}
|
|
|
|
const all = await log.getAllEvents();
|
|
assert.equal(all.length, count);
|
|
|
|
const stats = await log.getStats();
|
|
assert.equal(stats.totalEvents, count);
|
|
|
|
await log.close();
|
|
});
|
|
|
|
it('invalid file header throws', async () => {
|
|
writeFileSync(logPath, Buffer.from('BAAD'));
|
|
const snapPath = logPath.replace(/\.rvf$/, '.snap.rvf');
|
|
writeFileSync(snapPath, Buffer.from('RVFL'));
|
|
|
|
const log = new RvfEventLog({ logPath });
|
|
await assert.rejects(() => log.initialize(), /Invalid file header/);
|
|
});
|
|
|
|
it('snapshot:recommended emitted at threshold', async () => {
|
|
const log = new RvfEventLog({ logPath, snapshotThreshold: 3 });
|
|
await log.initialize();
|
|
|
|
const recommended: any[] = [];
|
|
log.on('snapshot:recommended', (info: any) => recommended.push(info));
|
|
|
|
for (let i = 0; i < 6; i++) {
|
|
await log.append(makeEvent('agg-1', 'tick'));
|
|
}
|
|
|
|
// Versions 3 and 6 should trigger the recommendation.
|
|
assert.equal(recommended.length, 2);
|
|
assert.equal(recommended[0].version, 3);
|
|
assert.equal(recommended[1].version, 6);
|
|
|
|
await log.close();
|
|
});
|
|
});
|