Compare commits
4 commits
0488547af9
...
e9ede5774b
| Author | SHA1 | Date | |
|---|---|---|---|
| e9ede5774b | |||
| 293b131fa2 | |||
| 905c267c47 | |||
| b1c7054467 |
5 changed files with 1908 additions and 144 deletions
754
src/brainy.ts
754
src/brainy.ts
|
|
@ -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([
|
||||||
|
|
|
||||||
|
|
@ -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[]>
|
||||||
|
|
|
||||||
|
|
@ -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)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
635
tests/integration/find-hybrid-filter-before-hydrate.test.ts
Normal file
635
tests/integration/find-hybrid-filter-before-hydrate.test.ts
Normal 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()
|
||||||
|
}
|
||||||
|
})
|
||||||
|
})
|
||||||
547
tests/integration/pending-embed-checkpoint.test.ts
Normal file
547
tests/integration/pending-embed-checkpoint.test.ts
Normal 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)
|
||||||
|
})
|
||||||
Reference in a new issue