brainy/tests/integration/aggregation-state-persistence.test.ts
David Snelling da55be7520 fix: aggregation state adoption on reopen + single-flight backfill + query-cap ratchet removal
- AggregationIndex: defineAggregate before init no longer forces a backfill.
  init reconciles instead of clobbering: the app definition wins, persisted
  state is adopted on hash match, a write landing pre-adoption forces an exact
  rescan, and a successful state load clears the backfill flag. New ready()
  settles every adoption decision before query paths consult backfill state.
- brainy: backfills are single-flight and batched. Concurrent queries share one
  store walk and every pending aggregate fills from that same walk; the old
  behavior let each concurrent query wipe the others' partial state and start
  its own full walk, so a store under steady aggregate traffic never converged.
  queryAggregate also waits for persisted definitions before deciding an
  aggregate does not exist.
- paramValidation: recordQuery is telemetry-only. The duration-based cap
  ratchet (x0.8 per recorded query while lifetime-average exceeded 1s, floored
  at 1000 - below the documented 10000 auto floor, reported under a stale
  basis label) is removed; the cap comes from its construction-time basis or
  explicit overrides alone.
- storage: __aggregation_* and singleton system keys (brainy:entityIdMapper)
  are recognized before the unknown-key warning fires; routing is unchanged.
- docs: find-limits cap-immutability note, aggregation reopen/adoption
  semantics, RELEASES.md 8.5.1 entry.
2026-07-17 09:20:09 -07:00

235 lines
8 KiB
TypeScript

/**
* @module tests/integration/aggregation-state-persistence
* @description The boot-order contract for the aggregation engine. Five laws:
* (1) STATE ADOPTION — a reopen with an unchanged defineAggregate() adopts the
* persisted state and performs NO store walk. (The pre-fix behavior: the
* synchronous define always beat the async init, flagged a backfill, and
* the first query wiped the just-loaded state and re-walked the whole
* store — every restart, forever.)
* (2) SINGLE-FLIGHT + BATCH — concurrent cold queries across multiple pending
* aggregates share exactly ONE store walk; a query never wipes another's
* partial progress and M pending aggregates cost one enumeration, not M.
* (3) CHANGED DEFINITION — a real definition change still backfills, exactly.
* (4) PERSISTED-ONLY DEFINITIONS — an app that does not re-define at boot can
* query a persisted aggregate without racing a spurious "not defined".
* (5) QUIET KEYS — the engine's persistence keys (__aggregation_*) are
* recognized system keys: no "Unknown key format" warning at boot.
*/
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'
import * as fs from 'node:fs'
import * as os from 'node:os'
import * as path from 'node:path'
import { Brainy } from '../../src/index.js'
import { NounType } from '../../src/types/graphTypes.js'
import type { AggregateDefinition } from '../../src/types/brainy.types.js'
import { prodLog } from '../../src/utils/logger.js'
const SPENDING: AggregateDefinition = {
name: 'spending',
source: { type: NounType.Event, where: { domain: 'financial' } },
groupBy: ['category'],
metrics: {
total: { op: 'sum', field: 'amount' },
count: { op: 'count' }
}
}
/** Same name, different metrics — a REAL definition change (hash differs). */
const SPENDING_CHANGED: AggregateDefinition = {
...SPENDING,
metrics: {
total: { op: 'sum', field: 'amount' },
count: { op: 'count' },
average: { op: 'avg', field: 'amount' }
}
}
describe('aggregation state persistence — boot-order contract', () => {
let dir: string
beforeEach(() => {
dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-agg-persist-'))
})
afterEach(() => {
fs.rmSync(dir, { recursive: true, force: true })
})
const open = async (): Promise<any> => {
const b: any = new Brainy({
requireSubtype: false,
storage: { type: 'filesystem', path: dir },
silent: true
})
await b.init()
return b
}
/** Count store walks by intercepting the storage adapter's getNouns. */
const countWalks = (brain: any): { count: () => number } => {
const storage = brain.storage
const orig = storage.getNouns.bind(storage)
let calls = 0
storage.getNouns = async (opts: unknown) => {
calls++
return orig(opts)
}
return { count: () => calls }
}
const seed = async (brain: any): Promise<void> => {
for (let i = 0; i < 12; i++) {
await brain.add({
data: `tx ${i}`,
type: NounType.Event,
metadata: {
domain: 'financial',
category: i % 2 === 0 ? 'food' : 'transport',
amount: 10 + i
}
})
}
}
it('adopts persisted state on reopen with an unchanged definition — zero walks', async () => {
const brain1 = await open()
brain1.defineAggregate(SPENDING)
await seed(brain1)
const before = await brain1.queryAggregate('spending')
expect(before.length).toBe(2)
await brain1.close()
const brain2 = await open()
brain2.defineAggregate(SPENDING) // the standard declarative boot pattern
await brain2.getNounCount() // settle init paths before counting walks
const walks = countWalks(brain2)
const after = await brain2.queryAggregate('spending')
expect(walks.count()).toBe(0)
const key = (r: any) => r.groupKey.category
expect(
after.map((r: any) => [key(r), r.metrics.total, r.metrics.count]).sort()
).toEqual(
before.map((r: any) => [key(r), r.metrics.total, r.metrics.count]).sort()
)
await brain2.close()
})
it('adopted state keeps accumulating: post-reopen writes land on top of it', async () => {
const brain1 = await open()
brain1.defineAggregate(SPENDING)
await seed(brain1)
await brain1.queryAggregate('spending')
await brain1.close()
const brain2 = await open()
brain2.defineAggregate(SPENDING)
// A write BEFORE the first query: if it lands before state adoption the
// engine must choose an exact rescan over adoption — either way the
// result must include all 13 entities.
await brain2.add({
data: 'late tx',
type: NounType.Event,
metadata: { domain: 'financial', category: 'food', amount: 100 }
})
const rows = await brain2.queryAggregate('spending')
const food = rows.find((r: any) => r.groupKey.category === 'food')
expect(food.metrics.count).toBe(7) // 6 seeded + 1 late
await brain2.close()
})
it('a changed definition still backfills — exactly once, with correct results', async () => {
const brain1 = await open()
brain1.defineAggregate(SPENDING)
await seed(brain1)
await brain1.queryAggregate('spending')
await brain1.close()
const brain2 = await open()
brain2.defineAggregate(SPENDING_CHANGED)
await brain2.getNounCount()
const walks = countWalks(brain2)
const rows = await brain2.queryAggregate('spending')
expect(walks.count()).toBe(1) // 12 entities = one page = one getNouns call
const food = rows.find((r: any) => r.groupKey.category === 'food')
expect(food.metrics.count).toBe(6)
expect(food.metrics.average).toBeCloseTo(food.metrics.total / 6)
await brain2.close()
})
it('concurrent cold queries across two aggregates share exactly ONE walk', async () => {
const brain = await open()
brain.defineAggregate(SPENDING)
brain.defineAggregate({
...SPENDING,
name: 'by_category_count',
metrics: { count: { op: 'count' } }
})
await seed(brain)
await brain.getNounCount()
const walks = countWalks(brain)
const results = await Promise.all([
brain.queryAggregate('spending'),
brain.queryAggregate('by_category_count'),
brain.queryAggregate('spending'),
brain.queryAggregate('by_category_count'),
brain.queryAggregate('spending'),
brain.queryAggregate('by_category_count')
])
expect(walks.count()).toBe(1) // 12 entities = one page; one walk fills both
for (const rows of results) {
const total = rows.reduce((s: number, r: any) => s + r.metrics.count, 0)
expect(total).toBe(12)
}
// Warm re-query: converged, no further walks.
await brain.queryAggregate('spending')
expect(walks.count()).toBe(1)
await brain.close()
})
it('persisted-only definitions are queryable without re-defining at boot', async () => {
const brain1 = await open()
brain1.defineAggregate(SPENDING)
await seed(brain1)
await brain1.queryAggregate('spending')
await brain1.close()
const brain2 = await open()
// NO defineAggregate — the app relies on the persisted definition.
const walks = countWalks(brain2)
const rows = await brain2.queryAggregate('spending')
expect(walks.count()).toBe(0) // persisted state adopted here too
expect(rows.length).toBe(2)
await brain2.close()
})
it('aggregation persistence keys never log "Unknown key format"', async () => {
const warnSpy = vi.spyOn(prodLog, 'warn')
const brain1 = await open()
brain1.defineAggregate(SPENDING)
await seed(brain1)
await brain1.queryAggregate('spending')
await brain1.close()
const brain2 = await open()
brain2.defineAggregate(SPENDING)
await brain2.queryAggregate('spending')
await brain2.close()
const offenders = warnSpy.mock.calls
.map(args => String(args[0]))
.filter(
msg =>
msg.includes('Unknown key format') &&
(msg.includes('__aggregation_') || msg.includes('brainy:entityIdMapper'))
)
expect(offenders).toEqual([])
warnSpy.mockRestore()
})
})