Compare commits

...

4 commits

Author SHA1 Message Date
e9ede5774b test(open): pin the pending-embed checkpoint — stuck id, crash matrix, torn fallback
Some checks are pending
CI / Node 22 (push) Waiting to run
CI / Node 24 (push) Waiting to run
CI / Integration + conformance (Node 22) (push) Waiting to run
CI / Bun (latest) (push) Waiting to run
Four things, none of them a clock:

1. A brain with one permanently-stuck pending id, closed cleanly and
   reopened, scans ONLY the facts after the checkpoint — read from the
   fold's own accounting. The same fixture pins the DEFECT it cures: no
   low-water mark exists on that brain, because it never drained, so nothing
   could have shortened its fold. A second row proves the bound stays
   O(delta) across repeated opens while the id is still stuck.

2. A crash matrix in a REAL child process (detached group, SIGKILL, no
   close), following writer-lock-clean-close's pattern: killed before any
   checkpoint was written, killed after one with an embed landed and flushed
   above it, and killed after one with an UN-FLUSHED tail. The invariant in
   every row is differential — the checkpoint-bounded fold the reopened
   brain actually ran equals a full fold from generation 1 over the same
   recovered log.

3. A torn checkpoint (bytes that are neither gzip nor JSON) falls back
   loudly — the adapter's torn-record gauge and production error, plus the
   fold's own narration of the bound it used — and still recovers the marker
   from the log. A well-formed but shape-invalid checkpoint is refused
   WHOLE: trusting its generation while ignoring its list is the one shape
   that could bound a scan behind a set that was never recovered.

4. The existing low-water pins pass unchanged — the mark is still written
   and still read, now as the fallback bound beneath the checkpoint.
2026-09-02 10:11:34 -07:00
293b131fa2 perf(open): the pending-embed fold is bounded by a checkpoint of the SET, not an empty-only mark
The low-water mark shipped in 10.4.9 can only be written when the pending
set is EMPTY, because it carries no set — it means "everything at or below
G is consumed". A brain holding even one id that never lands (an embed that
keeps failing, a data-less row reaped in memory only and re-folded every
open) never drains, so it never writes a mark, so the bound never engaged on
exactly the brains whose fold is expensive: `recover-pending-embeds` re-read
the WHOLE fact log at every open, on the open's foreground.

_system/pending_embeds_checkpoint.json carries the set: { generation,
pending, writtenAt } = "as of durable generation G the pending set was
exactly this list". Open seeds the set from the list and scans from G + 1,
so the fold is O(facts since G) whether or not the set ever drains. Measured
on a 301-row brain with one stuck id: 302 facts read before, 0 after; at 601
rows, 602 before, 0 after — same pending set both ways.

THE DURABILITY LAW, by construction. A checkpoint at head H taken while the
facts up to H are still buffered would be read back after a crash that
truncated the tail: an `embed.landed` in a truncated fact would be gone from
the log while the checkpoint still recorded its id as landed, and its
landing vector went with the fact — a LOST VECTOR. So a capture is refused
unless `0 < head <= committed`, the manifest watermark below which
FactLog.open() never truncates and which the group-commit flush only
advances after fsyncing the log. The (generation, set) pair is taken in one
synchronous instant with no await between reading the generations and
snapshotting the set. The one remaining asymmetry runs the safe way: an id
enqueued in memory whose marker lands at G+1 is captured as pending at G —
one idempotent re-embed, never a loss.

Written at clean close (inside closeDurableSteps, after the generation
store's own close flushed the log and advanced the manifest), at
drain-to-empty, and on a cadence of max(64, ceil(|pending| / 64))
transitions while open — an interval that holds the mechanism's amortized
cost at <= 64 ids written per transition however large the backlog grows, so
the cure cannot reintroduce the defect class it fixes. No timer, no knob.
The debt stays armed across attempts the durability law refuses, so a write
burst does not skip a checkpoint, it defers it.

Degradation is loud and always toward a LONGER scan: a torn checkpoint
throws typed on read (the adapter's tmp+rename write means it can never
parse into a partial list) and a malformed one is refused whole, both
falling back to the low-water mark — still written, still read — and then to
generation 1. The fold narrates which bound applied and how many facts it
read, on every open, so a bound that stops engaging is visible instead of
silent.

The worker's orphan reap splits: a row that is GONE clears durably (its
tombstone is in the log, or its create never was), while a present-but
data-less row keeps clearing in memory only and is carried in the checkpoint
list, so the bounded fold and a full fold from generation 1 agree exactly.
The crash-recovery contract is unchanged: the fold stays on the open's
foreground, markers re-armed when open() returns.
2026-09-02 10:11:34 -07:00
905c267c47 fix(find): a page the metadata block already cut is not cut again
Some checks failed
Delta Gate / Delta gate — candidate vs control (push) Waiting to run
CI / Node 24 (push) Successful in 12m34s
CI / Node 22 (push) Successful in 12m40s
CI / Integration + conformance (Node 22) (push) Failing after 17m14s
CI / Bun (latest) (push) Successful in 12m29s
`find({ query, connected, where, offset })` answered [] for every page but the
first. The metadata block ranks the fused candidates and CUTS the page itself
— rows [offset, offset+limit) — and then returns early. Two shapes do not take
that early return, `connected` and `fusion`, and they fell through to the tail,
which sliced the already-cut page by `offset` a second time: a five-row page
sliced at offset five is nothing at all. Every page after the first was empty,
and the caller had no way to tell that from "no more rows".

The block now records that it consumed the offset, and the tail returns the
page it was handed instead of re-cutting it. Nothing changes at offset 0, where
the second slice was the identity.

Pinned in tests/integration/find-hybrid-filter-before-hydrate.test.ts: page two
of a `connected` hybrid find matches the pipeline oracle row for row, paging
reaches every matching neighbour exactly once, and a `fusion` find's second
page is the same page the plain find returns.
2026-09-02 09:37:27 -07:00
b1c7054467 fix(find): the hybrid legs rank inside the filter, and only the page is read
A hybrid find fuses a text leg and a semantic leg. The semantic leg already
walked only the metadata filter's universe. The text leg did not: it ranked
the WHOLE store, took the top `limit * 4`, read every one of those rows from
canonical, and only then intersected with the filter. On a large store with a
selective filter that is hundreds of rows read to return a handful — and a row
matching both the query and the filter, but sitting outside the store-wide
text prefix, was silently dropped. The same defect `find({ connected })`
carried before the graph-first law, one leg over.

Both legs now rank ids inside the universe and neither reads canonical. The
text leg goes through a new optional `getIdsForTextQueryWithin` door on
MetadataIndexProvider — the text twin of `filterIdsWithin`, so a native index
can intersect its postings before any string crosses the boundary; the
reference index implements it from its own posting-list merge, so the two
doors can never disagree, and a provider without it is served by the
whole-store answer intersected here. The fusion ranks shells, the page is cut
from them, and canonical is read once for exactly that page — with the row
rebuilt in full, so a hydrated row is indistinguishable from an eagerly-built
one (same flattened fields, same entity, same match visibility, same key
order). The eager forms of both legs stay for the search modes whose leg
output IS the answer.

Measured on the production recall shape (query + type list + `missing`
negation + excludeVFS, limit 60) the old order read 241 rows in two batches to
return one; the new order reads the page.

Pinned in tests/integration/find-hybrid-filter-before-hydrate.test.ts. The
oracle there is the pre-change pipeline itself, replayed on the same brain
through the same doors: where the filter does not truncate the text leg the
answer is identical — rows, order, scores, match visibility and row shape —
across hybrid + where, + type list + excludeVFS + a `missing` negation, +
connected, with and without offset. Where it does truncate, the correction is
held by name: the old order's text leg contributed nothing at all, the new one
returns the matching rows and paging reaches every one of them. The cost pins
read the engine's own counters: one batchGet of `limit` ids, the whole-store
text door never called, and what the text leg marshals bounded by the universe.
2026-09-02 09:37:27 -07:00
5 changed files with 1908 additions and 144 deletions

View file

@ -776,6 +776,50 @@ export class Brainy<T = any> implements BrainyInterface<T> {
private _pendingEmbedIds = new Set<string>() private _pendingEmbedIds = new Set<string>()
private _embedWorkerFlight: Promise<void> | null = null private _embedWorkerFlight: Promise<void> | null = null
/**
* Ids cleared from {@link _pendingEmbedIds} with NO durable disarming record
* behind them today exactly one case: a pending row that still EXISTS but
* carries no embeddable data, which the worker reaps in memory only. The log
* still says those ids are pending, so the pending-embed CHECKPOINT must
* carry them: the checkpoint's contract is "as of generation G the LOG's
* pending set was exactly this list", and a checkpoint that quietly dropped
* an id the log still arms would make the bounded fold disagree with a full
* fold from generation 1 the one divergence that could lose a vector.
* Bounded by the number of such rows; an id leaves when it is re-enqueued or
* durably disarmed.
*/
private _pendingEmbedUndurableClears = new Set<string>()
/**
* Pending-set transitions (enqueue/clear) since the last checkpoint attempt
* the checkpoint CADENCE. One mechanism, one hardcoded default, no knob and
* no timer (nothing to leave running after close).
*/
private _pendingEmbedCheckpointTransitions = 0
/**
* A checkpoint is OWED: the cadence came due (or the set drained) and no
* write has satisfied it yet. It stays armed across attempts the durability
* law refuses, so the next transition that CAN be checkpointed is.
*/
private _pendingEmbedCheckpointDue = false
/** Single-flight guard for the fire-and-forget checkpoint write. */
private _pendingEmbedCheckpointFlight: Promise<void> | null = null
/**
* What the last pending-embed recovery fold actually did the bound it
* used, where it started, and how many facts it read. The narration's
* source, and the accounting a pin reads instead of a clock.
*/
private _pendingEmbedFoldReport: {
bound: 'checkpoint' | 'low-water' | 'genesis'
fromGeneration: number
factsScanned: number
seeded: number
pending: number
} | null = null
// OPEN-PATH FIX: the background embedding-engine warm kicked off (never // OPEN-PATH FIX: the background embedding-engine warm kicked off (never
// awaited) by `performInit()` when `eagerEmbeddings` resolves true. Stored // awaited) by `performInit()` when `eagerEmbeddings` resolves true. Stored
// for observability only — `embed()`/`embeddingManager.embed()` already // for observability only — `embed()`/`embeddingManager.embed()` already
@ -2434,6 +2478,47 @@ export class Brainy<T = any> implements BrainyInterface<T> {
*/ */
private static readonly PENDING_EMBED_LOWWATER_PATH = '_system/pending_embeds_lowwater.json' private static readonly PENDING_EMBED_LOWWATER_PATH = '_system/pending_embeds_lowwater.json'
/**
* Storage-root-relative path of the pending-embed CHECKPOINT:
* `{ generation, pending: string[], writtenAt }` "as of durable generation
* G the pending set was exactly this list". Open seeds the set from `pending`
* and scans the log from `G + 1`, so the fold costs O(facts since G)
* REGARDLESS of whether the set ever drains.
*
* WHY IT REPLACES THE EMPTY-ONLY MARK AS THE BOUND. The low-water mark
* ({@link PENDING_EMBED_LOWWATER_PATH}) can only be written when the pending
* set is EMPTY, because it carries no set it means "everything at or below
* G is consumed". A brain holding even ONE id that never lands (an embed that
* keeps failing; a row reaped in memory only and re-folded every open) never
* drains, so it never writes a mark, so the bound never engages on exactly
* the brains whose fold is expensive: every open re-reads the whole log. The
* checkpoint carries the set, so it needs no drain.
*
* The mark is still written and still read as the FALLBACK bound (a
* checkpoint that is absent, torn, or malformed degrades to it, and then to
* generation 1). Correctness over cost in every degradation: a stale or
* missing checkpoint only lengthens the scan.
*/
private static readonly PENDING_EMBED_CHECKPOINT_PATH = '_system/pending_embeds_checkpoint.json'
/**
* Checkpoint CADENCE BASE: attempt a checkpoint every N pending-set
* transitions (enqueues + clears) while the brain is open, on top of the
* drain-to-empty and clean-close writes. Hardcoded 90th-percentile default,
* no knob, no timer: 64 transitions is far below the cost of the fold it
* bounds and far above the per-write noise floor. An attempt that cannot
* satisfy the durability law is SKIPPED, not forced the next transition
* retries.
*
* The interval ADAPTS to the one signal that matters, the backlog's own
* size, because a checkpoint writes the WHOLE pending list: the interval is
* `max(64, ceil(|pending| / 64))`, which holds the amortized cost of the
* mechanism at 64 ids written per transition NO MATTER how large the
* backlog grows. A term that scales with the store rather than with the
* work is exactly the defect class this file is fixing; it must not be
* reintroduced by the cure.
*/
private static readonly PENDING_EMBED_CHECKPOINT_EVERY = 64
/** /**
* @description Mark a deferred embed pending (MT5): the id joins the * @description Mark a deferred embed pending (MT5): the id joins the
@ -2448,6 +2533,9 @@ export class Brainy<T = any> implements BrainyInterface<T> {
*/ */
private enqueuePendingEmbed(id: string): FactMarkerRecord { private enqueuePendingEmbed(id: string): FactMarkerRecord {
this._pendingEmbedIds.add(id) this._pendingEmbedIds.add(id)
// Re-armed for real: any earlier in-memory-only clear is superseded.
this._pendingEmbedUndurableClears.delete(id)
this.noteEmbedCheckpointCadence()
return { type: 'embed.pending', id, enqueuedAt: Date.now() } return { type: 'embed.pending', id, enqueuedAt: Date.now() }
} }
@ -2458,11 +2546,27 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* fact) the recovery fold consumes those; nothing here touches storage. * fact) the recovery fold consumes those; nothing here touches storage.
* One honest residue: a pending row whose entity still exists but carries * One honest residue: a pending row whose entity still exists but carries
* no data is reaped in memory only, so it re-folds at the next open and * no data is reaped in memory only, so it re-folds at the next open and
* is re-reaped there a bounded no-op, never a lost vector. * is re-reaped there a bounded no-op, never a lost vector. That residue
* is the ONLY `durability: 'in-memory-only'` caller, and the checkpoint
* keeps carrying those ids so the bounded fold and a full fold from
* generation 1 agree exactly (see {@link _pendingEmbedUndurableClears}).
*
* @param id - The pending id to clear.
* @param durability - `'durable'` (default) when a record in the log at or
* below the current head disarms this id (an `embed.landed` riding the
* landing or unvector commit, or the row's tombstone including the row
* simply not being there any more); `'in-memory-only'` when nothing in the
* log says so.
*/ */
private clearPendingEmbed(id: string): void { private clearPendingEmbed(
id: string,
durability: 'durable' | 'in-memory-only' = 'durable'
): void {
this._pendingEmbedIds.delete(id) this._pendingEmbedIds.delete(id)
if (durability === 'in-memory-only') this._pendingEmbedUndurableClears.add(id)
else this._pendingEmbedUndurableClears.delete(id)
if (this._pendingEmbedIds.size === 0) this.maybeWriteEmbedLowWater() if (this._pendingEmbedIds.size === 0) this.maybeWriteEmbedLowWater()
this.noteEmbedCheckpointCadence()
} }
/** /**
@ -2498,6 +2602,223 @@ export class Brainy<T = any> implements BrainyInterface<T> {
} }
} }
/**
* @description Capture a pending-embed checkpoint, or refuse.
*
* THE DURABILITY LAW, satisfied by construction. The checkpoint asserts "as
* of generation G the log's pending set was exactly this list", and the next
* open TRUSTS it: it seeds the set and never reads a fact at or below G
* again. So a checkpoint may only be taken at a G whose facts are DURABLE.
* A checkpoint taken at head H while the facts up to H are still buffered
* would be read back after a crash that truncated the tail and an
* `embed.landed` in a truncated fact would be gone from the log while the
* checkpoint still recorded its id as landed. The row's landing vector went
* with the truncated fact, so nothing would ever re-arm it: A LOST VECTOR.
*
* The gate is therefore `0 < head ≤ committed`. `committed` is the
* generation manifest's watermark — the point the store's own recovery
* treats as truth, and the point below which `FactLog.open()` never
* truncates and the group-commit flush fsyncs the log BEFORE advancing it
* (see `GenerationStore.flushPendingSingleOps`). So every fact at or below
* `head` is fsynced and survives the crash exactly as the checkpoint
* describes it. Anything else (a head above the manifest, no log, no
* generation yet, a read-only or closed brain) REFUSES: skipping a
* checkpoint costs a longer scan next open, never a marker.
*
* The snapshot is taken SYNCHRONOUSLY with reading the two generations no
* `await` between them so no commit and no worker step can slip between
* "the generation I am about to claim" and "the set I claim for it".
*
* The one asymmetry, deliberately in the safe direction: an id whose
* `embed.pending` record has not been appended yet (enqueued in memory, its
* commit still in flight) is captured as pending at G although its marker
* will land at G+1 or later. Over-stating pending costs one idempotent
* re-embed attempt; under-stating it is the shape that loses a vector, and
* cannot happen every clear either rides a durable record at or below the
* head, or is carried in {@link _pendingEmbedUndurableClears}.
*
* @returns The checkpoint payload, or `null` when this instant cannot host
* one.
*/
private captureEmbedCheckpoint(): { generation: number; pending: string[] } | null {
if (this.isReadOnly || this.closed) return null
const store = this.generationStore
if (!store) return null
const log = store.getFactLog()
if (!log) return null
// --- ONE SYNCHRONOUS INSTANT: no await until the return. ---
const generation = log.headGeneration()
const committed = store.committedGeneration()
if (!(generation > 0) || generation > committed) return null
const pending = new Set(this._pendingEmbedIds)
for (const id of this._pendingEmbedUndurableClears) pending.add(id)
// --- end of the synchronous instant. ---
return { generation, pending: [...pending] }
}
/**
* @description Fire-and-forget checkpoint write, single-flight: a burst of
* transitions never stacks writes, and because each attempt captures
* immediately before it writes, the file always ends up holding the most
* recently captured (generation, set) PAIR and every such pair is
* independently true, so even an out-of-order landing is safe.
* {@link closeDurableSteps} awaits the flight before taking the final one.
*/
private maybeWriteEmbedCheckpoint(): void {
if (this._pendingEmbedCheckpointFlight) return
this._pendingEmbedCheckpointFlight = this.writeEmbedCheckpoint()
.then((wrote) => {
if (wrote) {
this._pendingEmbedCheckpointDue = false
this._pendingEmbedCheckpointTransitions = 0
}
})
.finally(() => {
this._pendingEmbedCheckpointFlight = null
})
}
/**
* The awaitable core of {@link maybeWriteEmbedCheckpoint}.
* @returns `true` when a checkpoint was actually written.
*/
private async writeEmbedCheckpoint(): Promise<boolean> {
const snapshot = this.captureEmbedCheckpoint()
if (!snapshot) return false
try {
// Atomic on disk: the filesystem adapter's writeRawObject is tmp+rename
// (see BaseStorage.writeRawObject), so a crash mid-write leaves either
// the previous checkpoint or the new one — never a spliced file. And a
// file that IS unreadable (a torn gzip, invalid JSON) throws typed on
// read and degrades to the fallback bound; it can never parse into a
// partial `pending` list.
//
// The file is NOT separately fsynced, and does not need to be: losing
// the rename to a power cut leaves the PREVIOUS checkpoint (or none),
// which only lengthens the next scan. The invariant that matters is the
// other direction — a checkpoint that IS visible names a generation
// whose facts are durable — and that is established by the capture gate
// above, not by this write.
await this.storage.writeRawObject(Brainy.PENDING_EMBED_CHECKPOINT_PATH, {
generation: snapshot.generation,
pending: snapshot.pending,
writtenAt: Date.now()
})
return true
} catch (err) {
prodLog.warn(
`[Brainy] pending-embed checkpoint write failed at generation ` +
`${snapshot.generation}: ${(err as Error).message} — the next open scans ` +
`from the previous checkpoint`
)
return false
}
}
/**
* @description The checkpoint cadence tick: count one pending-set transition
* and OWE a checkpoint every {@link PENDING_EMBED_CHECKPOINT_EVERY}
* transitions, plus on every drain to empty. The debt stays armed across
* attempts the durability law refuses during a write burst the log head
* legitimately runs ahead of the manifest, so the first attempt often cannot
* be taken and the next transition retries it. An active brain therefore
* checkpoints steadily without ever forcing a flush; an idle one relies on
* its clean close. No timer is involved, so nothing survives close().
*/
private noteEmbedCheckpointCadence(): void {
if (this.isReadOnly || this.closed) return
this._pendingEmbedCheckpointTransitions++
const listed = this._pendingEmbedIds.size + this._pendingEmbedUndurableClears.size
const every = Math.max(
Brainy.PENDING_EMBED_CHECKPOINT_EVERY,
Math.ceil(listed / Brainy.PENDING_EMBED_CHECKPOINT_EVERY)
)
if (
this._pendingEmbedIds.size === 0 ||
this._pendingEmbedCheckpointTransitions >= every
) {
this._pendingEmbedCheckpointDue = true
}
if (this._pendingEmbedCheckpointDue) this.maybeWriteEmbedCheckpoint()
}
/**
* @description Resolve the pending-embed fold's BOUND: the checkpoint first
* (a set plus a generation), then the legacy low-water mark (a generation
* only), then genesis. Every degradation is loud and lengthens the scan
* rather than shortening it a bound that could skip a marker is never
* derived from a value this method could not fully validate.
* @returns The bound's name, the first generation to scan, and the ids to
* seed the pending set with.
*/
private async readPendingEmbedBound(): Promise<{
bound: 'checkpoint' | 'low-water' | 'genesis'
fromGeneration: number
seeded: string[]
}> {
let checkpointRejected: string | null = null
try {
const raw = await this.storage.readRawObject(Brainy.PENDING_EMBED_CHECKPOINT_PATH)
if (raw !== null && raw !== undefined) {
const parsed = Brainy.parsePendingEmbedCheckpoint(raw)
if (parsed) {
return {
bound: 'checkpoint',
fromGeneration: parsed.generation + 1,
seeded: parsed.pending
}
}
checkpointRejected = 'its shape is not { generation: number > 0, pending: string[] }'
}
} catch (err) {
// A real storage fault (EIO/EACCES/…). Corruption never lands here: the
// adapter maps a torn raw object to `null` AFTER logging it as a
// production error, so a torn checkpoint arrives as "absent" — loud at
// the adapter, and bounded here by the fallback below.
checkpointRejected = `reading it failed: ${(err as Error).message}`
}
if (checkpointRejected !== null) {
prodLog.warn(
`[Brainy] pending-embed checkpoint REFUSED (${checkpointRejected}) — falling back ` +
`to the low-water mark, else a full fold from generation 1`
)
}
try {
const mark = (await this.storage.readRawObject(Brainy.PENDING_EMBED_LOWWATER_PATH)) as {
generation?: number
} | null
if (mark && typeof mark.generation === 'number' && mark.generation > 0) {
return { bound: 'low-water', fromGeneration: mark.generation + 1, seeded: [] }
}
} catch {
// No mark (or unreadable): scan from 1 — correctness over cost.
}
return { bound: 'genesis', fromGeneration: 1, seeded: [] }
}
/**
* @description Validate a raw checkpoint object STRICTLY. Anything that is
* not exactly `{ generation: integer > 0, pending: string[] }` is refused
* whole a partially-usable checkpoint is the one shape that could seed a
* short pending set behind a high bound, which is how a vector is lost.
* @param raw - The object read back from storage.
* @returns The validated checkpoint, or `null`.
*/
private static parsePendingEmbedCheckpoint(
raw: unknown
): { generation: number; pending: string[] } | null {
if (raw === null || typeof raw !== 'object' || Array.isArray(raw)) return null
const { generation, pending } = raw as { generation?: unknown; pending?: unknown }
if (typeof generation !== 'number' || !Number.isSafeInteger(generation) || generation <= 0) {
return null
}
if (!Array.isArray(pending) || pending.some((id) => typeof id !== 'string' || id === '')) {
return null
}
return { generation, pending: pending as string[] }
}
/** /**
* @description Rebuild the pending-embed set by REPLAYING the generation * @description Rebuild the pending-embed set by REPLAYING the generation
* log's marker records (recovery = replay, not listing): `embed.pending` * log's marker records (recovery = replay, not listing): `embed.pending`
@ -2506,14 +2827,22 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* survives the fold is exactly the set of acknowledged deferred writes * survives the fold is exactly the set of acknowledged deferred writes
* whose vectors have not landed. * whose vectors have not landed.
* *
* BOUND: the scan starts at the advisory low-water mark * BOUND: the scan starts after the pending-embed CHECKPOINT
* ({@link Brainy.PENDING_EMBED_LOWWATER_PATH}) the log head at which the * ({@link Brainy.PENDING_EMBED_CHECKPOINT_PATH}) "as of durable generation
* pending set last drained to empty so a settled brain reads only the * G the pending set was exactly this list" so the fold seeds the set from
* facts since then, not its whole history. Without a mark (first open * that list and reads only the facts after G. O(delta) whether or not the
* after upgrade) it scans from generation 1, once; a stale-low mark costs * set ever drains, which is the whole point: the previous bound, the
* a longer scan, never a marker. The fold stays on the open's foreground * empty-only low-water mark, could not be written at all by a brain holding
* the crash-recovery contract pins that a reopened brain has its markers * one id that never lands, so those brains re-read their whole log at every
* re-armed when open() returns and the mark is what makes that cheap. * open. The mark remains the FALLBACK bound (checkpoint absent, torn, or
* malformed), and generation 1 the fallback below that a brain opened for
* the first time after this change has neither a checkpoint nor, if it never
* drained, a mark, so it pays one full fold and writes a checkpoint on the
* way out. A stale bound costs a longer scan, never a marker. The fold stays
* on the open's foreground the crash-recovery contract pins that a
* reopened brain has its markers re-armed when open() returns and the
* bound is what makes that cheap. What it did (bound, start, facts read) is
* narrated and kept in {@link _pendingEmbedFoldReport}.
* It is SKIPPED WHOLESALE when the log has never had a v2 tail * It is SKIPPED WHOLESALE when the log has never had a v2 tail
* ({@link FactLog.hasV2History} v1 facts cannot carry marker records), * ({@link FactLog.hasV2History} v1 facts cannot carry marker records),
* so pre-cutover brains pay nothing; on a mixed log the scan still reads * so pre-cutover brains pay nothing; on a mixed log the scan still reads
@ -2526,20 +2855,13 @@ export class Brainy<T = any> implements BrainyInterface<T> {
private async recoverPendingEmbedsFromLog(): Promise<void> { private async recoverPendingEmbedsFromLog(): Promise<void> {
const log = this.generationStore.getFactLog() const log = this.generationStore.getFactLog()
if (!log || !log.hasV2History()) return if (!log || !log.hasV2History()) return
let fromGeneration = 1 const { bound, fromGeneration, seeded } = await this.readPendingEmbedBound()
try { for (const id of seeded) this._pendingEmbedIds.add(id)
const mark = (await this.storage.readRawObject(Brainy.PENDING_EMBED_LOWWATER_PATH)) as { let factsScanned = 0
generation?: number
} | null
if (mark && typeof mark.generation === 'number' && mark.generation > 0) {
fromGeneration = mark.generation + 1
}
} catch {
// No mark (or unreadable): scan from 1 — correctness over cost.
}
const scan = log.scanFacts({ fromGeneration }) const scan = log.scanFacts({ fromGeneration })
for await (const batch of scan.batches()) { for await (const batch of scan.batches()) {
for (const fact of batch.facts) { for (const fact of batch.facts) {
factsScanned++
for (const record of fact.records ?? []) { for (const record of fact.records ?? []) {
if (record.type === 'embed.pending') { if (record.type === 'embed.pending') {
this._pendingEmbedIds.add(record.id) this._pendingEmbedIds.add(record.id)
@ -2554,6 +2876,21 @@ export class Brainy<T = any> implements BrainyInterface<T> {
} }
} }
} }
this._pendingEmbedFoldReport = {
bound,
fromGeneration,
factsScanned,
seeded: seeded.length,
pending: this._pendingEmbedIds.size
}
// The narration channel: an operator is entitled to hear which bound
// applied and what it cost, on every open — that is how a bound that
// silently stopped engaging (the defect this replaced) becomes visible.
prodLog.narrate(
`[Brainy] pending-embed fold: ${bound} bound → scanned ${factsScanned} fact(s) ` +
`from generation ${fromGeneration}, seeded ${seeded.length} id(s), ` +
`${this._pendingEmbedIds.size} pending`
)
} }
/** /**
@ -2645,11 +2982,23 @@ export class Brainy<T = any> implements BrainyInterface<T> {
for (const id of batch) { for (const id of batch) {
try { try {
const entity = await this.get(id, { includeVectors: true }) const entity = await this.get(id, { includeVectors: true })
if (!entity || entity.data === undefined || entity.data === null) { if (!entity) {
// Orphan reap: a deleted row's tombstone fact durably disarms the // The row is GONE. Either it was deleted — its tombstone fact
// marker at the next recovery fold; a data-less-but-present row // durably disarms the marker, at or below the head, exactly as the
// (edge case) re-folds and re-reaps — bounded, never a lost vector. // fold reads it — or its create never became durable, in which case
this.clearPendingEmbed(id) // the log carries no `embed.pending` for it either. Both are durable
// clears: a full fold from generation 1 reaches the same answer.
this.clearPendingEmbed(id, 'durable')
continue
}
if (entity.data === undefined || entity.data === null) {
// Orphan reap, IN MEMORY ONLY: a data-less-but-present row (edge
// case) has nothing to embed, but no record in the log says so, so
// the fold would re-arm it. Cleared here and carried in the
// checkpoint (see clearPendingEmbed) — it re-folds and re-reaps at
// the next open exactly as before: bounded, never a lost vector,
// and never a checkpoint that disagrees with the log.
this.clearPendingEmbed(id, 'in-memory-only')
continue continue
} }
// Hang guard: a wedged embedder must not block every later pending // Hang guard: a wedged embedder must not block every later pending
@ -7651,6 +8000,18 @@ export class Brainy<T = any> implements BrainyInterface<T> {
const searchMode = params.searchMode || 'auto' const searchMode = params.searchMode || 'auto'
const limit = params.limit || 10 const limit = params.limit || 10
// HYDRATE LAST (the hybrid path): its legs and its fusion rank IDS, and
// canonical is read at the two page exits below — never for a row the
// metadata filter is about to discard. This closure re-applies a hybrid
// row's match visibility once its entity is in hand; it is set only by
// the hybrid branch, so every other path hydrates unchanged.
let finishHybridRow: ((row: Result<T>, pending: Result<T>) => void) | undefined
// Set once the metadata block below has already ranked and CUT the page.
// The tail must not cut it a second time: `offset` has been consumed, and
// re-slicing a `limit`-long page by `offset` returns nothing at all.
let pagedEarly = false
// Handle text-only query (user explicitly wants text search) // Handle text-only query (user explicitly wants text search)
if (searchMode === 'text' && params.query && params.query.trim() !== '') { if (searchMode === 'text' && params.query && params.query.trim() !== '') {
results = await this.executeTextSearch(params.query, limit * 2) results = await this.executeTextSearch(params.query, limit * 2)
@ -7661,20 +8022,32 @@ export class Brainy<T = any> implements BrainyInterface<T> {
} }
// Handle explicit hybrid or auto mode with query // Handle explicit hybrid or auto mode with query
else if ((searchMode === 'auto' || searchMode === 'hybrid') && params.query && params.query.trim() !== '' && !params.vector) { else if ((searchMode === 'auto' || searchMode === 'hybrid') && params.query && params.query.trim() !== '' && !params.vector) {
// Zero-config hybrid: combine text + semantic search with RRF fusion // Zero-config hybrid: combine text + semantic search with RRF fusion.
const [textResults, semanticResults] = await Promise.all([ // BOTH legs are held to the metadata filter's universe: the vector leg
this.executeTextSearch(params.query, limit * 2), // walks it as its candidate set, and the text leg ranks inside it
this.executeVectorSearch(params, preResolvedMetadataIds ?? undefined, preResolvedAllowedIds) // instead of ranking the whole store and discarding what the filter
// would drop. Neither leg reads canonical — the page does, once.
const [textScored, semanticScored] = await Promise.all([
this.executeTextSearchScored(params.query, limit * 2, preResolvedMetadataIds ?? undefined),
this.executeVectorSearchScored(params, preResolvedMetadataIds ?? undefined, preResolvedAllowedIds)
]) ])
// Use user-specified alpha or auto-detect based on query length // Use user-specified alpha or auto-detect based on query length
const alpha = params.hybridAlpha ?? this.autoAlpha(params.query) const alpha = params.hybridAlpha ?? this.autoAlpha(params.query)
// Tokenize query for match visibility // Tokenize query for match visibility. The word list needs the entity,
// so it is computed on the page, at hydration.
const queryWords = this.metadataIndex.tokenize(params.query) const queryWords = this.metadataIndex.tokenize(params.query)
const textResultIds = new Set(textScored.map((r) => r.id))
finishHybridRow = (row, pending) => {
row.textMatches = this.findMatchingWords(row.entity, queryWords, textResultIds)
row.textScore = pending.textScore
row.semanticScore = pending.semanticScore
row.matchSource = pending.matchSource
}
// RRF fusion combines both result sets with match visibility // RRF fusion combines both ranked id sets with match visibility
results = await this.rrfFusion(textResults, semanticResults, alpha, queryWords) results = this.rrfFusion(textScored, semanticScored, alpha)
} }
// Handle direct vector search (no query text) - no hybrid needed // Handle direct vector search (no query text) - no hybrid needed
else if (params.vector && !params.query) { else if (params.vector && !params.query) {
@ -7735,20 +8108,13 @@ export class Brainy<T = any> implements BrainyInterface<T> {
const k = offset + limit const k = offset + limit
const order = rankIndicesByScore(results.map(r => r.score), k, true) const order = rankIndicesByScore(results.map(r => r.score), k, true)
results = reorderByIndices(results, order).slice(offset, k) results = reorderByIndices(results, order).slice(offset, k)
pagedEarly = true
// Batch-load entities only for the paginated results (10x faster on GCS) // Batch-load entities only for the paginated results (10x faster on GCS).
const idsToLoad = results.filter(r => !r.entity).map(r => r.id) // This is the hydrate-last seam for the deferring paths: a row that
if (idsToLoad.length > 0) { // arrives as a ranked shell is rebuilt in full here — flattened
const entitiesMap = await this.batchGet(idsToLoad) // fields, entity and match visibility — never `entity` alone.
for (const result of results) { results = await this.hydrateResultPage(results, finishHybridRow)
if (!result.entity) {
const entity = entitiesMap.get(result.id)
if (entity) {
result.entity = entity
}
}
}
}
// Early return if no other processing needed // Early return if no other processing needed
if (!params.connected && !params.fusion) { if (!params.connected && !params.fusion) {
@ -7854,8 +8220,20 @@ export class Brainy<T = any> implements BrainyInterface<T> {
const finalOffset = params.offset || 0 const finalOffset = params.offset || 0
// Efficient pagination - only slice what we need (limit already defined above) // Efficient pagination - only slice what we need (limit already defined
return results.slice(finalOffset, finalOffset + limit) // above), THEN read canonical for the page. Rows that arrived hydrated
// pass straight through; a deferred path reads exactly these rows.
//
// A page the metadata block already cut is NOT cut again: it holds the
// rows at [offset, offset+limit) of the ranking, so slicing it by
// `offset` a second time drops the whole page. That is how
// `find({ query, connected, where, offset })` — the shapes that reach
// here after early paging, `connected` and `fusion` — answered [] for
// every page but the first.
return await this.hydrateResultPage(
pagedEarly ? results : results.slice(finalOffset, finalOffset + limit),
finishHybridRow
)
})() })()
// Index-integrity guard — applied ONCE here so every find() path (metadata, // Index-integrity guard — applied ONCE here so every find() path (metadata,
@ -12949,6 +13327,31 @@ export class Brainy<T = any> implements BrainyInterface<T> {
} }
} }
/**
* The text-leg twin of {@link filterIdsWithinBelted}: rank `query` INSIDE the
* candidate universe, through the provider's own posting-list merge so the
* answer can never drift from `getIdsForTextQuery`'s. A provider without the
* door is served by its whole-store answer intersected here the same rows
* in the same order, but it pays the whole-store marshal.
*
* @param query - The text query.
* @param ids - The candidate universe (the metadata filter's ids).
* @returns `{ id, matchCount }` rows inside `ids`, ranked by match count.
*/
private async textIdsWithinBelted(
query: string,
ids: readonly string[]
): Promise<Array<{ id: string; matchCount: number }>> {
this.ensureIndexesLoaded(['metadata'])
const mip = this.metadataIndex as unknown as MetadataIndexProvider
if (typeof mip.getIdsForTextQueryWithin === 'function') {
return await mip.getIdsForTextQueryWithin(query, ids)
}
const within = new Set(ids)
const all = await this.metadataIndex.getIdsForTextQuery(query)
return all.filter((m) => within.has(m.id))
}
async getIndexStatus(): Promise<{ async getIndexStatus(): Promise<{
initialized: boolean initialized: boolean
/** `true` once open()'s index-build-if-needed step has run. Named for API /** `true` once open()'s index-build-if-needed step has run. Named for API
@ -15841,6 +16244,44 @@ export class Brainy<T = any> implements BrainyInterface<T> {
candidateIds?: string[], candidateIds?: string[],
allowedIds?: OpaqueIdSet allowedIds?: OpaqueIdSet
): Promise<Result<T>[]> { ): Promise<Result<T>[]> {
const scored = await this.executeVectorSearchScored(params, candidateIds, allowedIds)
// Batch-load entities for 10-50x faster cloud storage performance
// GCS: 10 results = 1×50ms vs 10×50ms = 500ms (10x faster)
const entitiesMap = await this.batchGet(scored.map((s) => s.id))
const results: Result<T>[] = []
for (const { id, score } of scored) {
const entity = entitiesMap.get(id)
if (entity) {
results.push(this.createResult(id, score, entity))
}
}
return results
}
/**
* The semantic leg WITHOUT hydration ranked ids and their scores.
*
* The beam walk is already restricted to the candidate universe (that is what
* `candidateIds` / `allowedIds` are for), so the leg's cost is the walk. Its
* ROWS, though, are candidates for a fusion that will keep one page of them
* so the hybrid path takes them unhydrated and reads exactly the page it
* returns. {@link executeVectorSearch} is the eager form, for the search modes
* whose leg output IS the answer.
*
* @param params - Find parameters (supplies the query/vector and the limit).
* @param candidateIds - Optional pre-resolved metadata universe (see
* {@link executeVectorSearch}).
* @param allowedIds - Optional opaque predicate-pushdown universe.
* @returns Ranked `{ id, score }` rows no entity reads.
*/
private async executeVectorSearchScored(
params: FindParams<T>,
candidateIds?: string[],
allowedIds?: OpaqueIdSet
): Promise<Array<{ id: string; score: number }>> {
// Vector cold-read guard: before trusting a semantic/vector result, verify the // Vector cold-read guard: before trusting a semantic/vector result, verify the
// vector index actually SERVES a known persisted vector (one-shot per brain). // vector index actually SERVES a known persisted vector (one-shot per brain).
// A pure semantic find({ query }) has no filter, so verifyMetadataLive never // A pure semantic find({ query }) has no filter, so verifyMetadataLive never
@ -15866,21 +16307,10 @@ export class Brainy<T = any> implements BrainyInterface<T> {
// HNSW search with optional metadata-first candidate filtering // HNSW search with optional metadata-first candidate filtering
const searchResults: [string, number][] = await this.index.search(vector, limit * 2, undefined, searchOptions) const searchResults: [string, number][] = await this.index.search(vector, limit * 2, undefined, searchOptions)
// Batch-load entities for 10-50x faster cloud storage performance return searchResults.map(([id, distance]) => ({
// GCS: 10 results = 1×50ms vs 10×50ms = 500ms (10x faster) id,
const ids = searchResults.map(([id]) => id) score: Math.max(0, Math.min(1, 1 / (1 + distance)))
const entitiesMap = await this.batchGet(ids) }))
const results: Result<T>[] = []
for (const [id, distance] of searchResults) {
const entity = entitiesMap.get(id)
if (entity) {
const score = Math.max(0, Math.min(1, 1 / (1 + distance)))
results.push(this.createResult(id, score, entity))
}
}
return results
} }
/** /**
@ -16180,30 +16610,64 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* @returns Array of Results with scores based on match count * @returns Array of Results with scores based on match count
*/ */
private async executeTextSearch(query: string, limit: number): Promise<Result<T>[]> { private async executeTextSearch(query: string, limit: number): Promise<Result<T>[]> {
const textMatches = await this.metadataIndex.getIdsForTextQuery(query) const scored = await this.executeTextSearchScored(query, limit)
if (textMatches.length === 0) return [] if (scored.length === 0) return []
// Take top matches and load entities // Batch-load entities for the whole leg — this is the eager form, kept for
const topMatches = textMatches.slice(0, limit * 2) // Get more for filtering // the text-only search mode whose results ARE the answer.
const ids = topMatches.map(m => m.id) const entitiesMap = await this.batchGet(scored.map((s) => s.id))
const entitiesMap = await this.batchGet(ids)
// Create results with scores based on match count
const maxMatches = topMatches[0]?.matchCount || 1
const results: Result<T>[] = [] const results: Result<T>[] = []
for (const { id, score } of scored) {
for (const match of topMatches) { const entity = entitiesMap.get(id)
const entity = entitiesMap.get(match.id)
if (entity) { if (entity) {
// Normalize score to 0-1 range based on match count results.push(this.createResult(id, score, entity))
const score = match.matchCount / maxMatches
results.push(this.createResult(match.id, score, entity))
} }
} }
return results return results
} }
/**
* The text leg WITHOUT hydration ranked ids and their scores.
*
* FILTER BEFORE HYDRATE: when the caller already knows the candidate
* universe (the metadata filter's ids in a hybrid `find({ query, where })`),
* it is passed here and the word index ranks INSIDE that universe. The
* earlier order ranked the whole store, took the top `limit * 2`, hydrated
* every one of them, and only then intersected with the filter so a
* filtered hybrid find on a large store hydrated hundreds of rows to return
* a handful, and a matching row outside the store-wide text prefix was
* silently dropped (the same defect `find({ connected })` had before the
* graph-first law).
*
* The score is the match count normalized against the top row's, so a
* restricted call normalizes against the top row IN THE UNIVERSE the same
* rule applied to the set actually being ranked.
*
* @param query - Text query to search for.
* @param limit - Result budget; the leg keeps `limit * 2` for the fusion.
* @param candidateIds - Optional candidate universe to rank inside.
* @returns Ranked `{ id, score }` rows no entity reads.
*/
private async executeTextSearchScored(
query: string,
limit: number,
candidateIds?: readonly string[]
): Promise<Array<{ id: string; score: number }>> {
const textMatches = candidateIds
? await this.textIdsWithinBelted(query, candidateIds)
: await this.metadataIndex.getIdsForTextQuery(query)
if (textMatches.length === 0) return []
// Take top matches (more than the page, for the fusion to rank)
const topMatches = textMatches.slice(0, limit * 2)
// Normalize score to 0-1 range based on match count
const maxMatches = topMatches[0]?.matchCount || 1
return topMatches.map((m) => ({ id: m.id, score: m.matchCount / maxMatches }))
}
/** /**
* Auto-detect optimal alpha for hybrid search * Auto-detect optimal alpha for hybrid search
* *
@ -16230,55 +16694,56 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* *
* Formula: score(d) = sum(1 / (k + rank(d))) for each list * Formula: score(d) = sum(1 / (k + rank(d))) for each list
* *
* Now includes match visibility (textMatches, textScore, semanticScore, matchSource) * Now includes match visibility (textScore, semanticScore, matchSource; the
* `textMatches` word list needs the entity and is filled at hydration).
* *
* @param textResults - Results from text search * HYDRATE LAST: both legs arrive as ranked ids + scores, and the fusion ranks
* @param semanticResults - Results from semantic search * ids no entity is read here. The rows it returns are ranked SHELLS; the
* page is cut from them and only that page is read from canonical (see
* {@link hydrateResultPage}). The earlier order hydrated both legs in full
* hundreds of rows to return one page of them.
*
* @param textResults - Ranked ids + scores from text search
* @param semanticResults - Ranked ids + scores from semantic search
* @param alpha - Weight for semantic (0=text only, 1=semantic only) * @param alpha - Weight for semantic (0=text only, 1=semantic only)
* @param queryWords - Original query words for match tracking
* @param k - RRF constant (default: 60, standard in literature) * @param k - RRF constant (default: 60, standard in literature)
* @returns Fused results sorted by combined score with match visibility * @returns Fused result shells sorted by combined score with match visibility
*/ */
private async rrfFusion( private rrfFusion(
textResults: Result<T>[], textResults: ReadonlyArray<{ id: string; score: number }>,
semanticResults: Result<T>[], semanticResults: ReadonlyArray<{ id: string; score: number }>,
alpha: number, alpha: number,
queryWords: string[],
k: number = 60 k: number = 60
): Promise<Result<T>[]> { ): Result<T>[] {
// Track scores and match details per entity // Track scores and match details per entity
interface MatchData { interface MatchData {
rrf: number rrf: number
textScore?: number textScore?: number
semanticScore?: number semanticScore?: number
textMatches: string[]
hasText: boolean hasText: boolean
hasSemantic: boolean hasSemantic: boolean
} }
const matchData = new Map<string, MatchData>() const matchData = new Map<string, MatchData>()
const entityMap = new Map<string, Entity<T>>()
// Text contribution (1 - alpha weight) // Text contribution (1 - alpha weight)
const textWeight = 1 - alpha const textWeight = 1 - alpha
textResults.forEach((r, rank) => { textResults.forEach((r, rank) => {
const rrfScore = textWeight * (1 / (k + rank + 1)) const rrfScore = textWeight * (1 / (k + rank + 1))
const existing = matchData.get(r.id) || { rrf: 0, textMatches: [], hasText: false, hasSemantic: false } const existing = matchData.get(r.id) || { rrf: 0, hasText: false, hasSemantic: false }
existing.rrf += rrfScore existing.rrf += rrfScore
existing.textScore = r.score // Original text search score (0-1) existing.textScore = r.score // Original text search score (0-1)
existing.hasText = true existing.hasText = true
matchData.set(r.id, existing) matchData.set(r.id, existing)
if (r.entity) entityMap.set(r.id, r.entity)
}) })
// Semantic contribution (alpha weight) // Semantic contribution (alpha weight)
semanticResults.forEach((r, rank) => { semanticResults.forEach((r, rank) => {
const rrfScore = alpha * (1 / (k + rank + 1)) const rrfScore = alpha * (1 / (k + rank + 1))
const existing = matchData.get(r.id) || { rrf: 0, textMatches: [], hasText: false, hasSemantic: false } const existing = matchData.get(r.id) || { rrf: 0, hasText: false, hasSemantic: false }
existing.rrf += rrfScore existing.rrf += rrfScore
existing.semanticScore = r.score // Original semantic search score (0-1) existing.semanticScore = r.score // Original semantic search score (0-1)
existing.hasSemantic = true existing.hasSemantic = true
matchData.set(r.id, existing) matchData.set(r.id, existing)
if (r.entity) entityMap.set(r.id, r.entity)
}) })
// Sort by fused score // Sort by fused score
@ -16286,27 +16751,9 @@ export class Brainy<T = any> implements BrainyInterface<T> {
.sort((a, b) => b[1].rrf - a[1].rrf) .sort((a, b) => b[1].rrf - a[1].rrf)
.map(([id, data]) => ({ id, data })) .map(([id, data]) => ({ id, data }))
// Build results - need to load any missing entities // Create ranked shells with match visibility
const missingIds = sortedIds.filter(s => !entityMap.has(s.id)).map(s => s.id)
if (missingIds.length > 0) {
const loaded = await this.batchGet(missingIds)
for (const [id, entity] of loaded) {
entityMap.set(id, entity)
}
}
// Performance: Build set of text result IDs for O(1) lookup
// This avoids re-extracting text for entities that weren't in text results
const textResultIds = new Set(textResults.map(r => r.id))
// Create final results with match visibility
const results: Result<T>[] = [] const results: Result<T>[] = []
for (const { id, data } of sortedIds) { for (const { id, data } of sortedIds) {
const entity = entityMap.get(id)
if (entity) {
// Find which query words matched - uses fast path if entity wasn't in text results
const textMatches = this.findMatchingWords(entity, queryWords, textResultIds)
// Determine match source // Determine match source
let matchSource: 'text' | 'semantic' | 'both' let matchSource: 'text' | 'semantic' | 'both'
if (data.hasText && data.hasSemantic) { if (data.hasText && data.hasSemantic) {
@ -16317,20 +16764,80 @@ export class Brainy<T = any> implements BrainyInterface<T> {
matchSource = 'semantic' matchSource = 'semantic'
} }
// Create result with match visibility const result = this.pendingResult(id, data.rrf)
const result = this.createResult(id, data.rrf, entity)
result.textMatches = textMatches
result.textScore = data.textScore result.textScore = data.textScore
result.semanticScore = data.semanticScore result.semanticScore = data.semanticScore
result.matchSource = matchSource result.matchSource = matchSource
results.push(result) results.push(result)
} }
}
return results return results
} }
/**
* A ranked candidate whose entity has NOT been read yet.
*
* The shell carries everything the ranking tail needs the id, the score,
* and the match-visibility fields and nothing that requires canonical. It
* is typed `Result` so it flows through the shared dedupe / visibility /
* filter / rank / page tail unchanged; {@link hydrateResultPage} turns the
* survivors into real results before any caller sees them, and find()'s
* index-integrity guard drops any row that never gained an entity.
*
* @param id - The candidate's canonical id.
* @param score - Its rank score.
*/
private pendingResult(id: string, score: number): Result<T> {
return { id, score } as Result<T>
}
/**
* Read canonical for exactly the rows that need it the hydrate-last seam.
*
* Rows that already carry an entity (the eager legs: metadata, text-only,
* semantic-only, proximity, graph) pass through untouched, so this is a no-op
* for every path that has not deferred. Rows that are shells are read in ONE
* batch and rebuilt through {@link createResult}, so a hydrated row is
* indistinguishable from an eagerly-built one same flattened fields, same
* `entity`, same key order with `finish` re-applying the fields only the
* deferring path knows about (a hybrid row's match visibility).
*
* A shell whose id has no canonical row is dropped, exactly as the eager legs
* dropped it; find()'s index-integrity guard makes the same judgement on the
* page it returns.
*
* @param rows - The page's rows, ranked and paged already.
* @param finish - Applied to each rebuilt row, with its shell, after the
* flattened fields are set.
* @returns The page with every surviving row hydrated.
*/
private async hydrateResultPage(
rows: Result<T>[],
finish?: (row: Result<T>, pending: Result<T>) => void
): Promise<Result<T>[]> {
const pendingIds: string[] = []
for (const row of rows) {
if (!row.entity) pendingIds.push(row.id)
}
if (pendingIds.length === 0) return rows
const entitiesMap = await this.batchGet(pendingIds)
const hydrated: Result<T>[] = []
for (const row of rows) {
if (row.entity) {
hydrated.push(row)
continue
}
const entity = entitiesMap.get(row.id)
if (!entity) continue
const filled = this.createResult(row.id, row.score, entity, row.explanation)
finish?.(filled, row)
hydrated.push(filled)
}
return hydrated
}
/** /**
* Find which query words match in an entity's text content * Find which query words match in an entity's text content
* *
@ -19824,6 +20331,21 @@ export class Brainy<T = any> implements BrainyInterface<T> {
await this.stampEntityTree() await this.stampEntityTree()
} }
// Phase 1c: the pending-embed CHECKPOINT — placed HERE and not earlier
// because this is the first point in the close where the durability law it
// must satisfy actually holds: `generationStore.close()` (in Phase 1 above)
// flushed the pending single-op tier, which fsyncs the fact log and then
// advances the manifest, so `head === committed` and every fact the
// checkpoint's generation covers is durable. Taken even when the set is
// NOT empty — that is the whole difference from the low-water mark, and it
// is what makes the next open's fold O(facts since this close) on a brain
// whose pending set never drains. Awaits any in-flight cadence write first
// so the last write to the file is this one.
if (!this.isReadOnly) {
await this._pendingEmbedCheckpointFlight?.catch(() => {})
await this.writeEmbedCheckpoint()
}
// Phase 2: Close components to release resources (timers, file handles) // Phase 2: Close components to release resources (timers, file handles)
// Data is already safe on disk from Phase 1 // Data is already safe on disk from Phase 1
await Promise.all([ await Promise.all([

View file

@ -473,6 +473,28 @@ export interface MetadataIndexProvider {
graphIndex: unknown graphIndex: unknown
): Promise<{ ids: string[]; emptyAt: 'graph' | 'filter' | 'visibility' | 'none' } | null> ): Promise<{ ids: string[]; emptyAt: 'graph' | 'filter' | 'visibility' | 'none' } | null>
getIdsForTextQuery(query: string): Promise<Array<{ id: string; matchCount: number }>> getIdsForTextQuery(query: string): Promise<Array<{ id: string; matchCount: number }>>
/**
* @description OPTIONAL: score `query` over `ids` ONLY the text-leg twin of
* {@link filterIdsWithin}, and the door a hybrid `find({ query, where })`
* walks. The metadata filter's universe is the candidate set there, so the
* text leg must cost O(|ids|) membership checks and marshal at most `|ids|`
* rows, never the whole posting list of every query word. A native index
* intersects its own postings with the candidate set (membership by entity
* int) before any string crosses the boundary; the reference index answers
* from its own `getIdsForTextQuery`, so the two doors can never disagree.
* Absent Brainy intersects `getIdsForTextQuery`'s answer with `ids` itself
* (correct, and still hydrate-last, but it marshals the whole answer).
*
* The answer keeps `getIdsForTextQuery`'s contract: `{ id, matchCount }`
* sorted by `matchCount` descending, ties in the order the whole-store answer
* would have produced. Only rows in `ids` may appear.
* @param query - The same text query accepted by `getIdsForTextQuery`.
* @param ids - The candidate ids (canonical). The answer is a subset.
*/
getIdsForTextQueryWithin?(
query: string,
ids: readonly string[]
): Promise<Array<{ id: string; matchCount: number }>>
getSortedIdsForFilter(filter: any, orderBy: string, order?: 'asc' | 'desc', topK?: number): Promise<string[]> getSortedIdsForFilter(filter: any, orderBy: string, order?: 'asc' | 'desc', topK?: number): Promise<string[]>
getFilterValues(field: string): Promise<string[]> getFilterValues(field: string): Promise<string[]>
getFilterFields(): Promise<string[]> getFilterFields(): Promise<string[]>

View file

@ -1509,11 +1509,56 @@ export class MetadataIndexManager implements MetadataIndexProvider {
* @returns Array of { id, matchCount } sorted by matchCount descending * @returns Array of { id, matchCount } sorted by matchCount descending
*/ */
async getIdsForTextQuery(query: string): Promise<Array<{ id: string; matchCount: number }>> { async getIdsForTextQuery(query: string): Promise<Array<{ id: string; matchCount: number }>> {
return this.scoreTextQuery(query)
}
/**
* Score a text query over `ids` ONLY the reference implementation of the
* optional `getIdsForTextQueryWithin` door (see
* {@link import('../plugin.js').MetadataIndexProvider}). The hybrid
* `find({ query, where })` path passes the metadata filter's universe here so
* the text leg ranks INSIDE that universe instead of ranking the whole store
* and discarding the rows the filter would have dropped.
*
* It answers from the same posting-list merge as {@link getIdsForTextQuery},
* with the candidate membership applied as each word's postings are counted,
* so the two doors can never disagree: the answer is exactly the whole-store
* answer restricted to `ids`, in the same order.
*
* @param query - Text query to search for.
* @param ids - Candidate entity ids; only these may appear in the answer.
* @returns Array of { id, matchCount } sorted by matchCount descending.
*/
async getIdsForTextQueryWithin(
query: string,
ids: readonly string[]
): Promise<Array<{ id: string; matchCount: number }>> {
if (ids.length === 0) return []
return this.scoreTextQuery(query, new Set(ids))
}
/**
* The one posting-list merge behind both text doors.
*
* Each query word contributes AT MOST one match per entity (a posting list
* can name an id more than once), and entities are ranked by how many of the
* query's words they matched. `within`, when given, restricts the count to
* those candidates applied during the merge, so a restricted call never
* materializes a whole-store match map.
*
* @param query - Text query to search for.
* @param within - Optional candidate universe; absent = the whole store.
* @returns Array of { id, matchCount } sorted by matchCount descending.
*/
private async scoreTextQuery(
query: string,
within?: ReadonlySet<string>
): Promise<Array<{ id: string; matchCount: number }>> {
const queryWords = this.tokenize(query) const queryWords = this.tokenize(query)
if (queryWords.length === 0) return [] if (queryWords.length === 0) return []
// Get IDs for each word hash // Count matches per entity, one word's postings at a time.
const wordIdSets: Map<string, number>[] = [] const matchCounts = new Map<string, number>()
for (const word of queryWords) { for (const word of queryWords) {
const wordHash = this.hashWord(word) const wordHash = this.hashWord(word)
let ids: string[] let ids: string[]
@ -1529,19 +1574,12 @@ export class MetadataIndexManager implements MetadataIndexProvider {
throw err throw err
} }
} }
const idSet = new Map<string, number>() // One count per (word, entity) — dedupe this word's postings first.
const counted = new Set<string>()
for (const id of ids) { for (const id of ids) {
idSet.set(id, 1) if (counted.has(id)) continue
} counted.add(id)
wordIdSets.push(idSet) if (within && !within.has(id)) continue
}
if (wordIdSets.length === 0) return []
// Count matches per entity
const matchCounts = new Map<string, number>()
for (const idSet of wordIdSets) {
for (const [id] of idSet) {
matchCounts.set(id, (matchCounts.get(id) || 0) + 1) matchCounts.set(id, (matchCounts.get(id) || 0) + 1)
} }
} }

View file

@ -0,0 +1,635 @@
/**
* @module tests/integration/find-hybrid-filter-before-hydrate
* @description FILTER BEFORE HYDRATE, applied to the hybrid `find({ query })` path.
*
* A hybrid find fuses two legs. The semantic leg already walked only the
* metadata filter's universe (`candidateIds` / `allowedIds`). The TEXT leg did
* not: it ranked the WHOLE store, took the top `limit * 4`, read every one of
* those rows from canonical, and only then intersected with the filter so a
* filtered hybrid find on a large store read hundreds of rows to return a
* handful of them, and a matching row outside the store-wide text prefix was
* silently dropped. That is the same defect `find({ connected })` carried
* before the graph-first law, one leg over.
*
* Both halves are pinned here.
*
* THE ANSWER. Where the filter did not truncate the text leg the universe
* covers every text match, so both orders rank the same rows the new
* pipeline's answer is IDENTICAL to the old one's: same rows, same order, same
* scores, same match visibility, same row shape. The oracle below is the
* pre-change pipeline itself, replayed on the same brain through the same
* doors, so the comparison is against what actually ran, not a remembered
* expectation.
*
* THE CORRECTION. Where the filter DID truncate it the query's words are
* common outside the universe the old order let the text leg contribute
* nothing at all: every row it ranked was discarded by the filter, and the
* answer came from the semantic leg alone. The new order ranks inside the
* universe, so the text leg contributes the rows it always should have.
*
* THE COST. Canonical is read for exactly the page: one batch, `limit` rows,
* never the legs. And the text leg is asked about the universe's ids only
* what it marshals is bounded by the universe, not by the store.
*/
import { describe, it, expect, beforeAll, vi } from 'vitest'
import { Brainy } from '../../src/brainy'
import { NounType, VerbType } from '../../src/types/graphTypes'
import { rankIndicesByScore, reorderByIndices } from '../../src/utils/resultRanking'
import { resolveEntityId } from '../../src/utils/idNormalization'
/** Embedding width of the default model — the row vectors must match it. */
const DIM = 384
/**
* A deterministic, per-row-distinct unit vector. Distinct so the semantic leg
* has a real ranking to produce (identical vectors would make its order a tie
* break), deterministic so the oracle and the pipeline see the same one.
*/
function seededVector(seed: number): number[] {
const v = new Array<number>(DIM)
for (let i = 0; i < DIM; i++) {
v[i] = Math.sin((i + 1) * 0.11 + seed * 0.37) * 0.5 + Math.cos((i + 1) * 0.05 + seed * 0.13) * 0.3
}
const magnitude = Math.sqrt(v.reduce((sum, x) => sum + x * x, 0))
return v.map((x) => x / magnitude)
}
/** The fields a caller reads off a hybrid row — the whole comparable surface. */
function project(rows: any[]): any[] {
return rows.map((r) => ({
id: r.id,
score: r.score,
type: r.type,
metadata: r.metadata,
textMatches: r.textMatches,
textScore: r.textScore,
semanticScore: r.semanticScore,
matchSource: r.matchSource
}))
}
/**
* The PRE-CHANGE hybrid pipeline, replayed on a live brain through the same
* provider doors it used: whole-store text ranking with both legs hydrated in
* full, RRF fusion, then the metadata intersection, then the page.
*
* Supports the shapes these pins exercise (query + where/type/excludeVFS +
* connected + offset); `orderBy`, `fusion` and `near` are not replayed.
*/
async function legacyHybridFind(brain: any, params: any): Promise<any[]> {
const index = brain.metadataIndex
const limit = params.limit ?? 10
const offset = params.offset ?? 0
const hasFilter = Boolean(
params.where || params.type || params.subtype || params.service || params.excludeVFS
)
let preResolvedMetadataIds: string[] | null = null
let preResolvedFilter: any = null
let graphFirstIds: string[] | null = null
if (params.connected) {
// find() normalizes the anchors to canonical ids before this stage runs.
const anchored = {
...params,
connected: {
...params.connected,
...(params.connected.from && { from: resolveEntityId(params.connected.from) }),
...(params.connected.to && { to: resolveEntityId(params.connected.to) })
}
}
graphFirstIds = await brain.resolveConnectedIds(anchored)
if (graphFirstIds!.length > 0 && hasFilter) {
preResolvedFilter = brain.buildMetadataFilter(params)
graphFirstIds = await brain.filterIdsWithinBelted(preResolvedFilter, graphFirstIds)
}
if (graphFirstIds!.length === 0) return []
preResolvedMetadataIds = graphFirstIds
} else if (hasFilter) {
preResolvedFilter = brain.buildMetadataFilter(params)
preResolvedMetadataIds = await brain.filterIdsBelted(preResolvedFilter)
if (preResolvedMetadataIds!.length === 0) return []
}
// Text leg — the whole store, then the top `limit * 4`, hydrated in full.
const allTextMatches = await index.getIdsForTextQuery(params.query)
const topMatches = allTextMatches.slice(0, limit * 2 * 2)
const maxMatches = topMatches[0]?.matchCount || 1
const textEntities = await brain.batchGet(topMatches.map((m: any) => m.id))
const textResults = topMatches
.filter((m: any) => textEntities.has(m.id))
.map((m: any) => ({ id: m.id, score: m.matchCount / maxMatches }))
// Semantic leg — the beam walk over the universe, hydrated in full.
const vector = await brain.embed(params.query)
const searchOptions = preResolvedMetadataIds ? { candidateIds: preResolvedMetadataIds } : undefined
const searchResults: [string, number][] = await brain.index.search(
vector,
limit * 2,
undefined,
searchOptions
)
const semanticEntities = await brain.batchGet(searchResults.map(([id]) => id))
const semanticResults = searchResults
.filter(([id]) => semanticEntities.has(id))
.map(([id, distance]) => ({ id, score: Math.max(0, Math.min(1, 1 / (1 + distance))) }))
// RRF fusion, with the match visibility the rows carried.
const alpha = params.hybridAlpha ?? brain.autoAlpha(params.query)
const k = 60
const matchData = new Map<string, any>()
const textWeight = 1 - alpha
textResults.forEach((r: any, rank: number) => {
const existing = matchData.get(r.id) || { rrf: 0, hasText: false, hasSemantic: false }
existing.rrf += textWeight * (1 / (k + rank + 1))
existing.textScore = r.score
existing.hasText = true
matchData.set(r.id, existing)
})
semanticResults.forEach((r: any, rank: number) => {
const existing = matchData.get(r.id) || { rrf: 0, hasText: false, hasSemantic: false }
existing.rrf += alpha * (1 / (k + rank + 1))
existing.semanticScore = r.score
existing.hasSemantic = true
matchData.set(r.id, existing)
})
const queryWords: string[] = index.tokenize(params.query)
const textResultIds = new Set(textResults.map((r: any) => r.id))
const fusedIds = Array.from(matchData.entries())
.sort((a, b) => b[1].rrf - a[1].rrf)
.map(([id, data]) => ({ id, data }))
const allEntities = await brain.batchGet(fusedIds.map((f) => f.id))
let rows: any[] = []
for (const { id, data } of fusedIds) {
const entity = allEntities.get(id)
if (!entity) continue
const textContent = textResultIds.has(id)
? index.extractTextContent({ data: entity.data, metadata: entity.metadata }).toLowerCase()
: null
rows.push({
id,
score: data.rrf,
type: entity.type,
metadata: entity.metadata,
textMatches:
textContent === null ? [] : queryWords.filter((w) => textContent.includes(w.toLowerCase())),
textScore: data.textScore,
semanticScore: data.semanticScore,
matchSource: data.hasText && data.hasSemantic ? 'both' : data.hasText ? 'text' : 'semantic'
})
}
// The metadata intersection — after the legs, as it was.
if (preResolvedMetadataIds && preResolvedFilter) {
const filteredIdSet = new Set(preResolvedMetadataIds)
rows = rows.filter((r) => filteredIdSet.has(r.id))
}
if (graphFirstIds !== null) {
const neighbourSet = new Set(graphFirstIds)
rows = rows.filter((r) => neighbourSet.has(r.id))
}
// Rank to the page, then cut it.
const order = rankIndicesByScore(
rows.map((r) => r.score),
offset + limit,
true
)
return reorderByIndices(rows, order).slice(offset, offset + limit)
}
/**
* FIXTURE A the filter's universe covers every text match, so the two orders
* rank exactly the same rows and the answers must be identical.
*/
describe('hybrid find: filter before hydrate — the answer is unchanged', () => {
let brain: Brainy<any>
const QUERY = 'orbital telemetry'
const MATCHES = 24
const FILLER = 120
const OUTSIDE = 30
const VFS = 10
const RETRACTED = 6
const anchor = 'array-anchor'
const matchIds: string[] = []
beforeAll(async () => {
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
await brain.init()
let seed = 1
await brain.add({
id: anchor,
data: 'ground station anchor record',
type: NounType.Thing,
metadata: { lane: 'alpha', role: 'anchor' },
vector: seededVector(seed++)
})
// Rows the query's words actually match — all inside every filter below.
for (let i = 0; i < MATCHES; i++) {
const id = `match-${i}`
await brain.add({
id,
data: `orbital telemetry packet ${i} recorded downlink`,
type: NounType.Document,
metadata: { lane: 'alpha', rank: i },
vector: seededVector(seed++)
})
matchIds.push(resolveEntityId(id))
await brain.relate({ from: anchor, to: id, type: VerbType.RelatedTo })
}
// Rows inside the universe that the query's words do NOT match.
for (let i = 0; i < FILLER; i++) {
await brain.add({
id: `filler-${i}`,
data: `cistern ledger entry ${i} archived`,
type: NounType.Document,
metadata: { lane: 'alpha', rank: 1000 + i },
vector: seededVector(seed++)
})
}
// Rows outside the universe.
for (let i = 0; i < OUTSIDE; i++) {
await brain.add({
id: `outside-${i}`,
data: `unrelated dossier ${i}`,
type: NounType.Person,
metadata: { lane: 'beta' },
vector: seededVector(seed++)
})
}
// VFS infrastructure rows — excluded by excludeVFS.
for (let i = 0; i < VFS; i++) {
await brain.add({
id: `vfs-${i}`,
data: `mounted path ${i}`,
type: NounType.Document,
metadata: { lane: 'alpha', vfsType: 'file' },
vector: seededVector(seed++)
})
}
// Retracted rows — excluded by a `missing` negation.
for (let i = 0; i < RETRACTED; i++) {
await brain.add({
id: `retracted-${i}`,
data: `withdrawn note ${i}`,
type: NounType.Document,
metadata: { lane: 'alpha', retracted: true },
vector: seededVector(seed++)
})
}
// The reference index has no opaque-set door, so the pipeline and the
// oracle both restrict the beam walk with the materialized candidate ids.
expect(typeof (brain as any).metadataIndex.getIdSetForFilter).not.toBe('function')
})
it('the fixture does not truncate the text leg — the universe covers every text match', async () => {
const index = (brain as any).metadataIndex
const textMatches = await index.getIdsForTextQuery(QUERY)
expect(textMatches).toHaveLength(MATCHES)
const universe = await (brain as any).filterIdsBelted({ lane: 'alpha' })
const inUniverse = new Set(universe)
for (const m of textMatches) expect(inUniverse.has(m.id)).toBe(true)
})
it('hybrid + where: identical rows, identical order, identical scores', async () => {
const params = { query: QUERY, where: { lane: 'alpha' }, limit: 8 }
const expected = await legacyHybridFind(brain as any, params)
const actual = await brain.find(params as any)
expect(actual.length).toBe(expected.length)
expect(project(actual)).toEqual(expected)
})
it('hybrid + where + offset: identical page two', async () => {
const params = { query: QUERY, where: { lane: 'alpha' }, limit: 6, offset: 6 }
const expected = await legacyHybridFind(brain as any, params)
const actual = await brain.find(params as any)
expect(actual.length).toBe(expected.length)
expect(project(actual)).toEqual(expected)
})
it('hybrid + type list + excludeVFS + a `missing` negation: identical', async () => {
const params = {
query: QUERY,
type: [NounType.Document, NounType.Person],
excludeVFS: true,
where: { lane: 'alpha', retracted: { missing: true } },
limit: 8
}
const expected = await legacyHybridFind(brain as any, params)
const actual = await brain.find(params as any)
expect(actual.length).toBe(expected.length)
expect(project(actual)).toEqual(expected)
for (const r of actual) {
expect(r.metadata.retracted).toBeUndefined()
expect(r.metadata.vfsType).toBeUndefined()
}
})
it('hybrid + type list + excludeVFS + a `missing` negation, offset: identical', async () => {
const params = {
query: QUERY,
type: [NounType.Document, NounType.Person],
excludeVFS: true,
where: { lane: 'alpha', retracted: { missing: true } },
limit: 5,
offset: 5
}
const expected = await legacyHybridFind(brain as any, params)
const actual = await brain.find(params as any)
expect(actual.length).toBe(expected.length)
expect(project(actual)).toEqual(expected)
})
it('hybrid + connected: identical, and never a non-neighbour', async () => {
const params = {
query: QUERY,
connected: { from: anchor, direction: 'out' as const },
where: { lane: 'alpha' },
limit: 8
}
const expected = await legacyHybridFind(brain as any, params)
const actual = await brain.find(params as any)
expect(actual.length).toBe(expected.length)
expect(project(actual)).toEqual(expected)
const neighbours = new Set(matchIds)
for (const r of actual) expect(neighbours.has(r.id)).toBe(true)
})
it('hybrid + connected + offset: page two is the page, not an empty answer', async () => {
const params = {
query: QUERY,
connected: { from: anchor, direction: 'out' as const },
where: { lane: 'alpha' },
limit: 5,
offset: 5
}
const expected = await legacyHybridFind(brain as any, params)
expect(expected).toHaveLength(5)
const actual = await brain.find(params as any)
expect(actual.length).toBe(expected.length)
expect(project(actual)).toEqual(expected)
})
it('hybrid + connected: paging reaches every matching neighbour exactly once', async () => {
const seen = new Set<string>()
for (let offset = 0; offset < MATCHES; offset += 6) {
const page = await brain.find({
query: QUERY,
connected: { from: anchor, direction: 'out' as const },
where: { lane: 'alpha' },
limit: 6,
offset
} as any)
for (const r of page) {
expect(seen.has(r.id)).toBe(false)
seen.add(r.id)
}
}
// Every row the fused candidate set holds is reachable by paging, and the
// neighbour set is the ceiling.
expect(seen.size).toBeGreaterThanOrEqual(MATCHES)
const neighbours = new Set(matchIds)
for (const id of seen) expect(neighbours.has(id)).toBe(true)
})
it('hybrid + fusion + offset: page two is the page', async () => {
const plain = await brain.find({
query: QUERY,
where: { lane: 'alpha' },
limit: 5,
offset: 5
} as any)
const fused = await brain.find({
query: QUERY,
where: { lane: 'alpha' },
fusion: 'weighted',
limit: 5,
offset: 5
} as any)
expect(fused).toHaveLength(plain.length)
expect(fused.map((r) => r.id)).toEqual(plain.map((r) => r.id))
})
it('a hydrated hybrid row is shaped exactly as an eagerly-built one', async () => {
const rows = await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 8 } as any)
const row = rows[0]
expect(Object.keys(row)).toEqual([
'id',
'score',
'type',
'subtype',
'visibility',
'metadata',
'data',
'confidence',
'weight',
'_rev',
'entity',
'textMatches',
'textScore',
'semanticScore',
'matchSource'
])
// The flattened fields are projections of the entity, as always.
expect(row.entity).toBeDefined()
expect(row.type).toBe(row.entity.type)
expect(row.metadata).toBe(row.entity.metadata)
expect(row.data).toBe(row.entity.data)
expect(row._rev).toBe(row.entity._rev)
// The match visibility survives the deferral — every leg's fields, on the
// rows that leg contributed, exactly as the eager pipeline set them.
expect(['text', 'semantic', 'both']).toContain(row.matchSource)
for (const r of rows) {
if (r.matchSource === 'semantic') {
expect(r.textMatches).toEqual([])
expect(r.textScore).toBeUndefined()
} else {
expect(r.textMatches).toEqual(['orbital', 'telemetry'])
expect(typeof r.textScore).toBe('number')
}
if (r.matchSource === 'text') {
expect(r.semanticScore).toBeUndefined()
} else {
expect(typeof r.semanticScore).toBe('number')
}
}
})
it('reads canonical for the page only — one batch, `limit` rows', async () => {
// Warm any first-read verification before the counters are read.
await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 1 } as any)
const hydrate = vi.spyOn(brain as any, 'batchGet')
try {
const results = await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 10 } as any)
expect(results).toHaveLength(10)
expect(hydrate).toHaveBeenCalledTimes(1)
expect((hydrate.mock.calls[0][0] as string[]).length).toBe(10)
} finally {
hydrate.mockRestore()
}
})
it('asks the text index about the universe only, never the whole store', async () => {
const index = (brain as any).metadataIndex
const wholeStore = vi.spyOn(index, 'getIdsForTextQuery')
const within = vi.spyOn(index, 'getIdsForTextQueryWithin')
try {
await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 10 } as any)
expect(wholeStore).not.toHaveBeenCalled()
expect(within).toHaveBeenCalledTimes(1)
const askedIds = within.mock.calls[0][1] as string[]
const universe = await (brain as any).filterIdsBelted({ lane: 'alpha' })
expect(askedIds).toHaveLength(universe.length)
// What the text leg marshals is bounded by the universe, not the store.
const marshalled = (await within.mock.results[0].value) as unknown[]
expect(marshalled.length).toBeLessThanOrEqual(universe.length)
expect(marshalled).toHaveLength(MATCHES)
} finally {
wholeStore.mockRestore()
within.mockRestore()
}
})
it('the two text doors agree: within is the whole-store answer restricted', async () => {
const index = (brain as any).metadataIndex
const universe: string[] = await (brain as any).filterIdsBelted({
lane: 'alpha',
retracted: { missing: true }
})
const inUniverse = new Set(universe)
const whole = await index.getIdsForTextQuery(QUERY)
const within = await index.getIdsForTextQueryWithin(QUERY, universe)
expect(within).toEqual(whole.filter((m: any) => inUniverse.has(m.id)))
expect(await index.getIdsForTextQueryWithin(QUERY, [])).toEqual([])
})
})
/**
* FIXTURE B the query's words are common OUTSIDE the universe, so the old
* order's text leg was entirely consumed by rows the filter then discarded.
* This is the corrected behaviour, held by name.
*/
describe('hybrid find: the text leg ranks inside the filter, not around it', () => {
let brain: Brainy<any>
const QUERY = 'orbital telemetry drift'
const NOISE = 150
const KEEP = 15
const keepIds: string[] = []
beforeAll(async () => {
brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' } })
await brain.init()
let seed = 5000
// Added FIRST and matching one more query word, so they lead the
// store-wide text ranking outright — and none of them pass the filter.
for (let i = 0; i < NOISE; i++) {
await brain.add({
id: `noise-${i}`,
data: `orbital telemetry drift report ${i}`,
type: NounType.Document,
metadata: { lane: 'beta' },
vector: seededVector(seed++)
})
}
for (let i = 0; i < KEEP; i++) {
const id = `keep-${i}`
await brain.add({
id,
data: `orbital telemetry summary ${i}`,
type: NounType.Document,
metadata: { lane: 'alpha' },
vector: seededVector(seed++)
})
keepIds.push(resolveEntityId(id))
}
})
it('the old order let the filter consume the whole text leg', async () => {
const index = (brain as any).metadataIndex
const universe: string[] = await (brain as any).filterIdsBelted({ lane: 'alpha' })
expect(universe).toHaveLength(KEEP)
const inUniverse = new Set(universe)
// The store-wide prefix the old text leg took (limit 10 → limit * 4).
const prefix = (await index.getIdsForTextQuery(QUERY)).slice(0, 40)
expect(prefix).toHaveLength(40)
expect(prefix.filter((m: any) => inUniverse.has(m.id))).toHaveLength(0)
// Every row the old text leg ranked was then discarded by the filter, so
// the old answer carried NO text contribution at all — fifteen rows that
// match the query's words exactly, and not one of them reached the page
// through the text leg. What the old order returned was whatever the
// semantic leg alone happened to reach.
const legacy = await legacyHybridFind(brain as any, {
query: QUERY,
where: { lane: 'alpha' },
limit: 10
})
for (const r of legacy) {
expect(r.matchSource).toBe('semantic')
expect(r.textScore).toBeUndefined()
expect(r.textMatches).toEqual([])
}
})
it('the new order ranks the text leg inside the universe', async () => {
const results = await brain.find({
query: QUERY,
where: { lane: 'alpha' },
limit: 10
} as any)
expect(results).toHaveLength(10)
const keeps = new Set(keepIds)
for (const r of results) {
expect(keeps.has(r.id)).toBe(true)
expect(r.metadata.lane).toBe('alpha')
// The text leg is the contributor the old order threw away.
expect(['text', 'both']).toContain(r.matchSource)
expect(r.textScore).toBe(1)
expect(r.textMatches).toEqual(['orbital', 'telemetry'])
}
})
it('paging reaches every matching row the old order could not see', async () => {
const seen = new Set<string>()
for (let offset = 0; offset < KEEP; offset += 5) {
const page = await brain.find({
query: QUERY,
where: { lane: 'alpha' },
limit: 5,
offset
} as any)
expect(page).toHaveLength(5)
for (const r of page) {
expect(seen.has(r.id)).toBe(false)
seen.add(r.id)
}
}
expect(seen.size).toBe(KEEP)
expect([...seen].sort()).toEqual([...keepIds].sort())
})
it('reads canonical for the page only, on the truncating shape too', async () => {
await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 1 } as any)
const hydrate = vi.spyOn(brain as any, 'batchGet')
try {
const results = await brain.find({ query: QUERY, where: { lane: 'alpha' }, limit: 10 } as any)
expect(results).toHaveLength(10)
expect(hydrate).toHaveBeenCalledTimes(1)
expect((hydrate.mock.calls[0][0] as string[]).length).toBe(10)
} finally {
hydrate.mockRestore()
}
})
})

View file

@ -0,0 +1,547 @@
/**
* @module tests/integration/pending-embed-checkpoint
* @description THE PENDING-EMBED CHECKPOINT the bound that engages on the
* brains that need it.
*
* 10.4.9 bounded the open-path `recover-pending-embeds` fold with a LOW-WATER
* MARK: the log head at which the pending set last drained to EMPTY. That mark
* carries no set, so it can only be written when the set is empty and a brain
* holding even ONE id that never lands (an embed that keeps failing, a worker
* that never gets to it, a row reaped in memory only and re-folded every open)
* never drains, therefore never writes a mark, therefore re-reads its WHOLE
* fact log on every single open. The bound was absent from exactly the brains
* whose fold is expensive: a silent scaling defect.
*
* The cure is a CHECKPOINT of the pending set
* `_system/pending_embeds_checkpoint.json` = `{ generation, pending, writtenAt }`,
* meaning "as of durable generation G the pending set was exactly this list".
* Open seeds the set from `pending` and scans only from `G + 1`, so the fold is
* O(facts since G) whether or not the set ever drains.
*
* What this suite pins:
* 1. A brain with one permanently-stuck pending id, closed cleanly and
* reopened, scans ONLY the facts after the checkpoint asserted from the
* fold's own accounting, never a clock. The same fixture pins the DEFECT:
* no low-water mark exists on that brain, because it never drained.
* 2. A crash matrix in a REAL child process (SIGKILL, no close), for kills
* before a checkpoint write, after one with embeds landed and flushed
* after it, and after one with an UN-FLUSHED tail at the moment of death.
* The invariant in every row is differential: the checkpoint-bounded fold
* the reopened brain actually ran a full fold from generation 1 over the
* same recovered log.
* 3. A torn checkpoint falls back loudly (the adapter's torn-record gauge
* plus the fold's own narration of which bound applied) and correctly.
* 4. The existing low-water pins keep passing unchanged
* (`pending-embed-low-water.test.ts`): the mark is still written and is
* still read, now as the FALLBACK bound beneath the checkpoint.
*
* The crash-recovery contract is untouched: the fold runs on the open's
* foreground, so a reopened brain has its markers re-armed when open() returns.
*/
import { describe, it, expect, afterEach } from 'vitest'
import { mkdtempSync, rmSync, existsSync, readFileSync, writeFileSync } from 'node:fs'
import { spawn } from 'node:child_process'
import { gunzipSync } from 'node:zlib'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { Brainy } from '../../src/brainy.js'
import { NounType } from '../../src/types/graphTypes.js'
import { getTornRecordGauge } from '../../src/storage/tornRecordError.js'
const CHECKPOINT_PATH = '_system/pending_embeds_checkpoint.json'
const LOWWATER_PATH = '_system/pending_embeds_lowwater.json'
const REPO_ROOT = process.cwd()
const TSX = join(REPO_ROOT, 'node_modules', '.bin', 'tsx')
/** The fold's own accounting for the most recent open. */
interface FoldReport {
bound: 'checkpoint' | 'low-water' | 'genesis'
fromGeneration: number
factsScanned: number
seeded: number
pending: number
}
const roots: string[] = []
const liveBrains: Brainy<any>[] = []
function dir(): string {
const d = mkdtempSync(join(tmpdir(), 'brainy-embed-ckpt-'))
roots.push(d)
return d
}
async function open(root: string, opts?: { blockWorker?: boolean }): Promise<Brainy<any>> {
const brain = new Brainy<any>({
requireSubtype: false,
storage: { type: 'filesystem', path: root }
})
// Blocking the worker BEFORE init() is how a "permanently stuck" pending id
// is built deterministically: the state under test is "an id the fold keeps
// re-arming and nothing ever disarms", and its production causes (a failing
// embedder, a wedged model, a data-less row) all reduce to exactly that.
if (opts?.blockWorker) (brain as unknown as { kickEmbedWorker: () => void }).kickEmbedWorker = () => {}
await brain.init()
liveBrains.push(brain)
return brain
}
function foldReport(brain: Brainy<any>): FoldReport {
const report = (brain as unknown as { _pendingEmbedFoldReport: FoldReport | null })
._pendingEmbedFoldReport
if (report === null) throw new Error('the open ran no pending-embed fold')
return report
}
function pendingIds(brain: Brainy<any>): string[] {
return [
...(brain as unknown as { _pendingEmbedIds: Set<string> })._pendingEmbedIds
].sort()
}
/** Read an artifact straight off disk (the adapter gzips raw objects). */
function readArtifact(root: string, path: string): Record<string, unknown> | null {
const plain = join(root, ...path.split('/'))
const gz = `${plain}.gz`
if (existsSync(gz)) return JSON.parse(gunzipSync(readFileSync(gz)).toString('utf-8'))
if (existsSync(plain)) return JSON.parse(readFileSync(plain, 'utf-8'))
return null
}
/** The on-disk path the adapter actually used for an artifact. */
function artifactPath(root: string, path: string): string | null {
const plain = join(root, ...path.split('/'))
const gz = `${plain}.gz`
if (existsSync(gz)) return gz
if (existsSync(plain)) return plain
return null
}
/**
* THE DIFFERENTIAL ORACLE: fold the log from generation 1 with exactly the
* engine's own rules. This is what the bounded fold must agree with, and its
* fact count is what the unbounded fold used to read at every open.
*/
async function fullFold(brain: Brainy<any>): Promise<{ ids: string[]; facts: number }> {
const log = (
brain as unknown as { generationStore: { getFactLog(): any } }
).generationStore.getFactLog()
const pending = new Set<string>()
let facts = 0
const scan = log.scanFacts({ fromGeneration: 1 })
for await (const batch of scan.batches()) {
for (const fact of batch.facts) {
facts++
for (const record of fact.records ?? []) {
if (record.type === 'embed.pending') pending.add(record.id)
else if (record.type === 'embed.landed') pending.delete(record.id)
}
for (const op of fact.ops) {
if (op.kind === 'noun' && op.record === null) pending.delete(op.id)
}
}
}
return { ids: [...pending].sort(), facts }
}
/** Capture every console.warn/error line emitted while `fn` runs. */
async function captureConsole<T>(fn: () => Promise<T>): Promise<{ result: T; lines: string[] }> {
const lines: string[] = []
const origWarn = console.warn
const origError = console.error
const sink = (...args: unknown[]) => {
lines.push(args.map((a) => String(a)).join(' '))
}
console.warn = sink as typeof console.warn
console.error = sink as typeof console.error
try {
const result = await fn()
return { result, lines }
} finally {
console.warn = origWarn
console.error = origError
}
}
/**
* Run a child process that arranges a store and then waits forever, so the
* parent can SIGKILL it. A real process death is the only honest way to pin
* "no close ran, no shutdown hook ran, RAM is gone".
*
* `detached` puts the child in its own process GROUP: tsx runs the script in a
* grandchild, and only a group-wide signal reaches the process holding the
* writer lock.
*/
function spawnArranger(root: string, body: string): Promise<{
child: ReturnType<typeof spawn>
output: () => string
}> {
const scriptPath = join(root, 'arrange.mts')
writeFileSync(scriptPath, body)
const child = spawn(TSX, [scriptPath], {
cwd: REPO_ROOT,
stdio: ['ignore', 'pipe', 'pipe'],
detached: true
})
let out = ''
child.stdout!.on('data', (d) => { out += String(d) })
child.stderr!.on('data', (d) => { out += String(d) })
return new Promise((resolvePromise, rejectPromise) => {
const timer = setTimeout(
() => rejectPromise(new Error(`arranger never became READY:\n${out}`)),
180_000
)
child.stdout!.on('data', () => {
if (out.includes('READY')) {
clearTimeout(timer)
resolvePromise({ child, output: () => out })
}
})
child.on('exit', (code) => {
clearTimeout(timer)
if (!out.includes('READY')) rejectPromise(new Error(`arranger exited ${code}:\n${out}`))
})
})
}
/** Parse the `IDS:{...}` line an arranger prints supplied ids are normalised
* to canonical uuids, and the markers, checkpoint and fold all speak those. */
function childIds(output: string): Record<string, string> {
const line = output.split('\n').find((l) => l.startsWith('IDS:'))
if (!line) throw new Error(`arranger printed no IDS line:\n${output}`)
return JSON.parse(line.slice('IDS:'.length))
}
/** SIGKILL the whole group and wait for the grandchild's death to settle. */
async function sigkill(child: ReturnType<typeof spawn>): Promise<void> {
process.kill(-(child.pid as number), 'SIGKILL')
await new Promise<void>((r) => child.on('exit', () => r()))
await new Promise<void>((r) => setTimeout(r, 500))
}
/** The preamble every arranger child shares. */
function childPreamble(root: string): string {
return `
import { Brainy } from ${JSON.stringify(join(REPO_ROOT, 'src', 'brainy.ts'))}
const ROOT = ${JSON.stringify(root)}
const brain = new Brainy<any>({ requireSubtype: false, storage: { type: 'filesystem', path: ROOT } })
const block = () => { (brain as any).kickEmbedWorker = () => {} }
const settleCheckpoint = async () => {
// The cadence write is fire-and-forget; wait for the single flight.
for (let i = 0; i < 200; i++) {
if (!(brain as any)._pendingEmbedCheckpointFlight) break
await (brain as any)._pendingEmbedCheckpointFlight.catch(() => {})
}
}
`
}
afterEach(async () => {
for (const brain of liveBrains.splice(0)) {
try { await brain.close() } catch { /* already closed / crashed — teardown only */ }
}
for (const d of roots.splice(0)) rmSync(d, { recursive: true, force: true })
})
// ===========================================================================
// 1. The stuck-id brain — the defect, and the bound that now engages on it
// ===========================================================================
describe('pending-embed checkpoint — a brain whose pending set never drains', () => {
it('a permanently-stuck pending id: the reopen scans only the facts after the checkpoint', async () => {
const root = dir()
const first = await open(root, { blockWorker: true })
// add() returns the CANONICAL id (supplied ids are normalised), and that is
// the id the markers, the checkpoint and the fold all speak.
const stuck = await first.add({
id: 'stuck',
data: 'a deferred row whose embed never lands',
type: NounType.Thing,
deferEmbedding: true
})
expect(first.pendingEmbedCount()).toBe(1)
// Ordinary traffic after it — every one of these is a fact the unbounded
// fold had to re-read at every open, forever, because of that one id.
for (let i = 0; i < 12; i++) {
await first.add({ id: `row-${i}`, data: `row ${i}`, type: NounType.Thing })
}
await first.close()
liveBrains.splice(liveBrains.indexOf(first), 1)
// THE DEFECT, PINNED: the pending set never drained, so the old bound was
// never written — nothing on this brain could have shortened its fold.
expect(readArtifact(root, LOWWATER_PATH)).toBeNull()
// The checkpoint IS written at the clean close, set non-empty and all.
const checkpoint = readArtifact(root, CHECKPOINT_PATH) as {
generation: number
pending: string[]
} | null
expect(checkpoint).not.toBeNull()
expect(checkpoint!.generation).toBeGreaterThan(0)
expect(checkpoint!.pending).toEqual([stuck])
const second = await open(root, { blockWorker: true })
const report = foldReport(second)
// THE FIX, from the fold's own counter — not the clock.
expect(report.bound).toBe('checkpoint')
expect(report.fromGeneration).toBe(checkpoint!.generation + 1)
expect(report.factsScanned).toBe(0)
expect(report.seeded).toBe(1)
// The crash-recovery contract is intact: the marker is re-armed by open().
expect(pendingIds(second)).toEqual([stuck])
expect(second.pendingEmbedCount()).toBe(1)
// The differential: the bounded answer is the full-fold answer, and the
// full fold is what the previous bound would have had to read.
const full = await fullFold(second)
expect(full.ids).toEqual([stuck])
expect(full.facts).toBeGreaterThanOrEqual(13)
expect(report.factsScanned).toBeLessThan(full.facts)
}, 180_000)
it('the bound stays O(delta) across repeated opens while the id is still stuck', async () => {
const root = dir()
const first = await open(root, { blockWorker: true })
const stuck = await first.add({
id: 'stuck',
data: 'never lands',
type: NounType.Thing,
deferEmbedding: true
})
for (let i = 0; i < 6; i++) {
await first.add({ id: `a-${i}`, data: `a ${i}`, type: NounType.Thing })
}
await first.close()
liveBrains.splice(liveBrains.indexOf(first), 1)
const second = await open(root, { blockWorker: true })
expect(foldReport(second).factsScanned).toBe(0)
// More history under the same stuck id.
for (let i = 0; i < 9; i++) {
await second.add({ id: `b-${i}`, data: `b ${i}`, type: NounType.Thing })
}
await second.close()
liveBrains.splice(liveBrains.indexOf(second), 1)
const third = await open(root, { blockWorker: true })
const report = foldReport(third)
const full = await fullFold(third)
expect(report.bound).toBe('checkpoint')
expect(report.factsScanned).toBe(0)
// The unbounded fold grew with the store; the bounded one did not.
expect(full.facts).toBeGreaterThanOrEqual(16)
expect(pendingIds(third)).toEqual([stuck])
expect(full.ids).toEqual([stuck])
}, 180_000)
})
// ===========================================================================
// 2. Torn checkpoint — falls back, loudly, correctly
// ===========================================================================
describe('pending-embed checkpoint — a torn checkpoint never shortens the fold', () => {
it('an undecodable checkpoint file degrades to the next bound, loudly, with the right pending set', async () => {
const root = dir()
const first = await open(root, { blockWorker: true })
const stuck = await first.add({
id: 'stuck',
data: 'never lands',
type: NounType.Thing,
deferEmbedding: true
})
for (let i = 0; i < 5; i++) {
await first.add({ id: `row-${i}`, data: `row ${i}`, type: NounType.Thing })
}
await first.close()
liveBrains.splice(liveBrains.indexOf(first), 1)
const onDisk = artifactPath(root, CHECKPOINT_PATH)
expect(onDisk).not.toBeNull()
// Tear it: bytes that are neither valid gzip nor valid JSON. A torn file
// must THROW on read — never parse into a partial `pending` list.
writeFileSync(onDisk!, 'not a checkpoint at all {{{')
const before = getTornRecordGauge().count
const { result: second, lines } = await captureConsole(async () =>
open(root, { blockWorker: true })
)
const report = foldReport(second)
// Fell back — never to a shorter bound, and never silently.
expect(report.bound).not.toBe('checkpoint')
expect(report.seeded).toBe(0)
expect(report.fromGeneration).toBe(1) // no mark either: this brain never drained
// LOUD, two ways: the adapter's torn-record gauge and its production error…
expect(getTornRecordGauge().count).toBeGreaterThan(before)
expect(getTornRecordGauge().lastPath).toContain('pending_embeds_checkpoint')
expect(lines.some((l) => /TORN RECORD/.test(l))).toBe(true)
// …and the fold's own narration of which bound it actually used.
expect(lines.some((l) => /pending-embed fold: genesis bound/.test(l))).toBe(true)
// CORRECT: the marker is still recovered, from the log itself.
expect(pendingIds(second)).toEqual([stuck])
const full = await fullFold(second)
expect(full.ids).toEqual([stuck])
expect(report.factsScanned).toBe(full.facts)
}, 180_000)
it('a well-formed but shape-invalid checkpoint is refused whole, never partially trusted', async () => {
const root = dir()
const first = await open(root, { blockWorker: true })
const stuck = await first.add({
id: 'stuck',
data: 'never lands',
type: NounType.Thing,
deferEmbedding: true
})
await first.add({ id: 'other', data: 'ordinary row', type: NounType.Thing })
await first.close()
liveBrains.splice(liveBrains.indexOf(first), 1)
// A checkpoint with a plausible generation but a `pending` that is not a
// list of ids: trusting the generation alone would bound the scan behind a
// set that was never recovered — the exact shape that loses a vector.
const onDisk = artifactPath(root, CHECKPOINT_PATH)!
const good = readArtifact(root, CHECKPOINT_PATH) as { generation: number }
rmSync(onDisk)
writeFileSync(
join(root, '_system', 'pending_embeds_checkpoint.json'),
JSON.stringify({ generation: good.generation, pending: { stuck: true }, writtenAt: 1 })
)
const { result: second, lines } = await captureConsole(async () =>
open(root, { blockWorker: true })
)
expect(lines.some((l) => /pending-embed checkpoint REFUSED/.test(l))).toBe(true)
const report = foldReport(second)
expect(report.bound).not.toBe('checkpoint')
expect(report.seeded).toBe(0)
expect(pendingIds(second)).toEqual([stuck])
}, 180_000)
})
// ===========================================================================
// 3. The crash matrix — real processes, real SIGKILL, differential invariant
// ===========================================================================
describe('pending-embed checkpoint — crash matrix (real child process, SIGKILL)', () => {
/**
* The invariant every row shares: whatever the reopened brain's fold did with
* whatever bound survived the crash, its pending set must equal the truth a
* full fold from generation 1 derives from the SAME recovered log.
*/
async function assertDifferentialAfterCrash(root: string): Promise<{
report: FoldReport
full: { ids: string[]; facts: number }
pending: string[]
}> {
const reopened = await open(root, { blockWorker: true })
const report = foldReport(reopened)
const full = await fullFold(reopened)
const pending = pendingIds(reopened)
expect(pending).toEqual(full.ids)
return { report, full, pending }
}
it('killed BEFORE any checkpoint was written — falls back and recovers the marker from the log', async () => {
const root = dir()
const { child, output } = await spawnArranger(
root,
`${childPreamble(root)}
block()
await brain.init()
await brain.add({ id: 'landed-row', data: 'an ordinary row', type: 'thing' })
const stuck = await brain.add({ id: 'stuck-1', data: 'deferred, never lands', type: 'thing', deferEmbedding: true })
await brain.flush()
console.log('IDS:' + JSON.stringify({ stuck }))
console.log('READY')
setInterval(() => {}, 1000)
`
)
const ids = childIds(output())
// One enqueue is well under the cadence and the set never drained, so no
// checkpoint exists — this is the pre-checkpoint crash.
expect(readArtifact(root, CHECKPOINT_PATH)).toBeNull()
await sigkill(child)
const { report, pending } = await assertDifferentialAfterCrash(root)
expect(report.bound).toBe('genesis')
expect(pending).toEqual([ids.stuck])
}, 300_000)
it('killed AFTER a checkpoint, with an embed landed and flushed after it — the post-checkpoint facts carry the disarm', async () => {
const root = dir()
const { child, output } = await spawnArranger(
root,
`${childPreamble(root)}
await brain.init()
// Land one deferred embed: the drain arms the checkpoint debt.
await brain.add({ id: 'seed', data: 'lands first', type: 'thing', deferEmbedding: true })
await brain.awaitPendingEmbeds()
await brain.flush()
// A second deferred write pays the debt (the head is at the manifest now),
// then LANDS — its embed.landed rides a fact ABOVE the checkpoint.
const landsAfter = await brain.add({ id: 'lands-after', data: 'lands after the checkpoint', type: 'thing', deferEmbedding: true })
await settleCheckpoint()
await brain.awaitPendingEmbeds()
// …and one that never will.
block()
const stuck = await brain.add({ id: 'stuck-1', data: 'deferred, never lands', type: 'thing', deferEmbedding: true })
await brain.add({ id: 'plain', data: 'more history', type: 'thing' })
await brain.flush()
console.log('IDS:' + JSON.stringify({ stuck, landsAfter }))
console.log('READY')
setInterval(() => {}, 1000)
`
)
const ids = childIds(output())
const checkpoint = readArtifact(root, CHECKPOINT_PATH) as {
generation: number
pending: string[]
} | null
expect(checkpoint).not.toBeNull()
await sigkill(child)
const { report, full, pending } = await assertDifferentialAfterCrash(root)
expect(report.bound).toBe('checkpoint')
expect(report.fromGeneration).toBe(checkpoint!.generation + 1)
// The bound really bounded: fewer facts than the whole log.
expect(report.factsScanned).toBeLessThan(full.facts)
// A landed embed above the checkpoint is disarmed by the scan, not lost;
// the stuck one is re-armed.
expect(pending).toEqual([ids.stuck])
expect(pending).not.toContain(ids.landsAfter)
}, 300_000)
it('killed AFTER a checkpoint with an UN-FLUSHED tail — truncated facts and the bounded fold still agree', async () => {
const root = dir()
const { child } = await spawnArranger(
root,
`${childPreamble(root)}
await brain.init()
await brain.add({ id: 'seed', data: 'lands first', type: 'thing', deferEmbedding: true })
await brain.awaitPendingEmbeds()
await brain.flush()
await brain.add({ id: 'lands-after', data: 'lands after the checkpoint', type: 'thing', deferEmbedding: true })
await settleCheckpoint()
await brain.awaitPendingEmbeds()
await brain.flush()
// Now write PAST the manifest and never flush: these facts are the tail a
// crash truncates. Whatever survives, the two folds must agree on it.
block()
await brain.add({ id: 'stuck-tail', data: 'deferred, never lands', type: 'thing', deferEmbedding: true })
await brain.add({ id: 'plain-tail', data: 'unflushed history', type: 'thing' })
console.log('READY')
setInterval(() => {}, 1000)
`
)
const checkpoint = readArtifact(root, CHECKPOINT_PATH) as { generation: number } | null
expect(checkpoint).not.toBeNull()
await sigkill(child)
const { report } = await assertDifferentialAfterCrash(root)
// The checkpoint's generation is at or below the manifest by construction,
// so it survived the truncation and still bounds the fold.
expect(report.bound).toBe('checkpoint')
expect(report.fromGeneration).toBe(checkpoint!.generation + 1)
}, 300_000)
})