open-brainy/tests/integration/wait-for-indexed.test.ts
David Snelling b53e6e8987
Some checks failed
CI / Node 22 (push) Has been cancelled
CI / Node 24 (push) Has been cancelled
CI / Bun (latest) (push) Has been cancelled
feat(engine): the wiring wave — stamps ride every flush, provider generations, waitForIndexed, adopt-backfill, match-all serves
- Watermark stamping fans out at flush: all three projections stamped
  with the committed generation before their flushes persist.
- waitForIndexed(path?, {generation, timeoutMs}) — the one honest read
  barrier for write-then-recall consumers; typed timeout error carries
  the pending count and names the gauge; getIndexStatus() gains
  per-projection gauges. awaitPendingEmbeds() unchanged underneath.
- adoptLogAuthority() self-backfills curable divergences (pre-log
  records, witness drift) by identity re-commit before flipping — a
  fresh brain flips clean; log-ahead divergences still refuse loudly.
- The verification oracle gains VERB legs (all four divergence classes;
  unwired = honest verbsChecked: 0, never a scope claim).
- find({where: {}}) match-all serves (was silent-empty, warm AND cold;
  same fix in count/streaming/subgraph seeding); removeMany({where:{}})
  refuses typed — a match-all bulk delete must be explicit.
- Aggregation native envelope stamped via noteSourceGeneration before
  serializeState; the native-blob restore gates through the same
  adoption verdict as caller-side state (the unconditional adopt dies).
- LC8 pinned: a wholesale directory move opens and serves identically
  across all three intelligences, with history traveling.

Gates: unit 2031/2031 (156 files) · integration 812 (91 files) ·
conformance 27/27.
2026-08-10 10:55:11 -07:00

219 lines
8.3 KiB
TypeScript

/**
* @module tests/integration/wait-for-indexed
* @description THE READ BARRIER — `brain.waitForIndexed(path?, opts?)`. A
* consumer that writes and then semantically recalls gets ONE honest barrier
* instead of guessing. The contract pinned here:
*
* 1. SEMANTIC LEG: a deferred add followed by `waitForIndexed('semantic')`
* resolves only after the vector landed — the row is vector-searchable
* the moment the barrier returns.
* 2. TYPED TIMEOUT: `timeoutMs` expiry REJECTS with
* WaitForIndexedTimeoutError carrying the leg + the pending count and
* naming the gauge — never a silent partial wait.
* 3. NO-ARG: every projection at the head; today that means the deferred
* embed backlog is drained.
* 4. SYNCHRONOUS LEGS: metadata/graph/aggregation resolve immediately by
* design today (they update inside the write path) — even while the
* semantic backlog is wedged.
* 5. GAUGES: getIndexStatus().projections carries the per-leg numbers, and
* the top-level pendingEmbeds compat field agrees with the semantic one.
* 6. GENERATION REFINEMENT: an empty backlog satisfies any generation
* immediately; a non-empty one falls back to the full drain.
*/
import { describe, it, expect, afterEach, vi } from 'vitest'
import { Brainy, WaitForIndexedTimeoutError } from '../../src/index.js'
import { NounType } from '../../src/types/graphTypes.js'
const brains: Brainy[] = []
async function memBrain(): Promise<Brainy> {
const b = new Brainy({ storage: { type: 'memory' }, requireSubtype: false })
await b.init()
brains.push(b)
return b
}
/**
* Abandon a poisoned in-flight embed run (its embed promise never resolves —
* production is covered by the worker's 60s hang guard; the test takes the
* white-box shortcut for speed), then drain so teardown never wedges.
*/
async function unwedge(brain: Brainy): Promise<void> {
;(brain as unknown as { _embedWorkerFlight: Promise<void> | null })._embedWorkerFlight = null
await brain.awaitPendingEmbeds()
}
afterEach(async () => {
vi.restoreAllMocks()
for (const b of brains.splice(0)) await b.close().catch(() => {})
})
describe('waitForIndexed — the read barrier', () => {
it("SEMANTIC LEG: deferred add → waitForIndexed('semantic') resolves and the row is vector-searchable after", async () => {
const brain = await memBrain()
const embedSpy = vi.spyOn(brain, 'embed')
const id = await brain.add({
data: 'the quarterly revenue report for the northern region',
type: NounType.Document,
deferEmbedding: true,
metadata: { kind: 'report' }
})
expect(embedSpy, 'no embed on the ack path').not.toHaveBeenCalled()
expect(brain.pendingEmbedCount()).toBeGreaterThanOrEqual(1)
await brain.waitForIndexed('semantic')
// The barrier's meaning: backlog drained, vector real, row searchable.
expect(brain.pendingEmbedCount(), 'barrier means drained').toBe(0)
const after = await brain.get(id, { includeVectors: true })
expect((after!.vector as number[]).length, 'real vector after the barrier').toBeGreaterThan(0)
const hits = await brain.find({
query: 'the quarterly revenue report for the northern region',
searchMode: 'semantic',
limit: 5
})
expect(hits.map((r) => r.id), 'vector-searchable after the barrier').toContain(id)
})
it('TYPED TIMEOUT: a hung embedder + timeoutMs rejects with the typed error naming the pending count and the gauge', async () => {
const brain = await memBrain()
const hang = vi
.spyOn(brain, 'embed')
.mockImplementation(() => new Promise<number[]>(() => {}))
await brain.add({
data: 'never lands while the embedder hangs',
type: NounType.Document,
deferEmbedding: true,
metadata: {}
})
expect(brain.pendingEmbedCount()).toBe(1)
let caught: unknown
try {
await brain.waitForIndexed('semantic', { timeoutMs: 200 })
} catch (e) {
caught = e
}
expect(caught, 'expiry REJECTS — never a silent partial wait').toBeInstanceOf(
WaitForIndexedTimeoutError
)
const err = caught as WaitForIndexedTimeoutError
expect(err.path).toBe('semantic')
expect(err.timeoutMs).toBe(200)
expect(err.pendingEmbeds).toBeGreaterThanOrEqual(1)
// The message names what was still pending and the gauge to check.
expect(err.message).toContain(`${err.pendingEmbeds} deferred embed`)
expect(err.message).toContain('getIndexStatus().projections.semantic.pendingEmbeds')
hang.mockRestore()
await unwedge(brain)
expect(brain.pendingEmbedCount()).toBe(0)
})
it('NO-ARG: waitForIndexed() waits on the pending-embed drain (every projection at the head)', async () => {
const brain = await memBrain()
await brain.add({
data: 'a deferred capture that the bare barrier must cover',
type: NounType.Document,
deferEmbedding: true,
metadata: {}
})
expect(brain.pendingEmbedCount()).toBeGreaterThanOrEqual(1)
await brain.waitForIndexed()
expect(
brain.pendingEmbedCount(),
'the bare barrier drained the only asynchronous projection'
).toBe(0)
})
it('SYNCHRONOUS LEGS: metadata/graph/aggregation resolve immediately — even while the semantic backlog is wedged', async () => {
const brain = await memBrain()
// Quiet brain first: all three legs resolve on a brain with no backlog.
await brain.add({ data: 'quiet row', type: NounType.Document, metadata: { q: 1 } })
await brain.awaitPendingEmbeds()
await brain.waitForIndexed('metadata')
await brain.waitForIndexed('graph')
await brain.waitForIndexed('aggregation')
// The stronger pin: these projections update inside the write path today,
// so their leg resolves immediately BY DESIGN — independent of a wedged
// semantic backlog. (If any of them incorrectly delegated to the embed
// drain, this test would hang.)
const hang = vi
.spyOn(brain, 'embed')
.mockImplementation(() => new Promise<number[]>(() => {}))
await brain.add({
data: 'wedged deferred row',
type: NounType.Document,
deferEmbedding: true,
metadata: {}
})
expect(brain.pendingEmbedCount()).toBe(1)
await brain.waitForIndexed('metadata')
await brain.waitForIndexed('graph')
await brain.waitForIndexed('aggregation')
hang.mockRestore()
await unwedge(brain)
})
it('GAUGES: getIndexStatus().projections carries the per-leg shape, and the compat field agrees', async () => {
const brain = await memBrain()
await brain.add({ data: 'gauge row', type: NounType.Document, metadata: { g: 1 } })
await brain.awaitPendingEmbeds()
const status = await brain.getIndexStatus()
expect(status.projections).toEqual({
semantic: { pendingEmbeds: 0 },
metadata: { synchronous: true },
graph: { synchronous: true },
aggregation: { pendingBackfills: 0, pendingCatchUps: 0 }
})
// Compat: the existing top-level gauge stays and agrees.
expect(status.pendingEmbeds).toBe(0)
// The semantic gauge is honest while a backlog exists.
const hang = vi
.spyOn(brain, 'embed')
.mockImplementation(() => new Promise<number[]>(() => {}))
await brain.add({
data: 'backlogged row',
type: NounType.Document,
deferEmbedding: true,
metadata: {}
})
const busy = await brain.getIndexStatus()
expect(busy.projections.semantic.pendingEmbeds).toBeGreaterThanOrEqual(1)
expect(busy.pendingEmbeds).toBe(busy.projections.semantic.pendingEmbeds)
hang.mockRestore()
await unwedge(brain)
})
it('GENERATION REFINEMENT: an empty backlog satisfies any generation immediately; a non-empty one falls back to the full drain', async () => {
const brain = await memBrain()
await brain.add({ data: 'generation row', type: NounType.Document, metadata: {} })
await brain.awaitPendingEmbeds()
// Empty backlog: the semantic watermark is at the head — >= any committed G.
await brain.waitForIndexed('semantic', { generation: 1 })
// Non-empty backlog: the conservative full drain (a superset of the
// requested wait, never a partial one).
await brain.add({
data: 'second generation row',
type: NounType.Document,
deferEmbedding: true,
metadata: {}
})
await brain.waitForIndexed('semantic', { generation: 1 })
expect(brain.pendingEmbedCount(), 'the fallback is the full drain').toBe(0)
})
})