Compare commits

..

1 commit

Author SHA1 Message Date
aab501d4a2 ci(test): perf and scale benchmarks leave the correctness gate
Some checks are pending
CI / Node 22 (push) Waiting to run
CI / Node 24 (push) Waiting to run
CI / Integration + conformance (Node 22) (push) Waiting to run
CI / Bun (latest) (push) Waiting to run
Delta Gate / Delta gate — candidate vs control (push) Waiting to run
2026-09-02 10:17:05 -07:00
4 changed files with 47 additions and 1399 deletions

View file

@ -776,50 +776,6 @@ export class Brainy<T = any> implements BrainyInterface<T> {
private _pendingEmbedIds = new Set<string>()
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
// awaited) by `performInit()` when `eagerEmbeddings` resolves true. Stored
// for observability only — `embed()`/`embeddingManager.embed()` already
@ -2478,47 +2434,6 @@ export class Brainy<T = any> implements BrainyInterface<T> {
*/
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
@ -2533,9 +2448,6 @@ export class Brainy<T = any> implements BrainyInterface<T> {
*/
private enqueuePendingEmbed(id: string): FactMarkerRecord {
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() }
}
@ -2546,27 +2458,11 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* fact) the recovery fold consumes those; nothing here touches storage.
* 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
* 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.
* is re-reaped there a bounded no-op, never a lost vector.
*/
private clearPendingEmbed(
id: string,
durability: 'durable' | 'in-memory-only' = 'durable'
): void {
private clearPendingEmbed(id: string): void {
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()
this.noteEmbedCheckpointCadence()
}
/**
@ -2602,223 +2498,6 @@ 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
* log's marker records (recovery = replay, not listing): `embed.pending`
@ -2827,22 +2506,14 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* survives the fold is exactly the set of acknowledged deferred writes
* whose vectors have not landed.
*
* BOUND: the scan starts after the pending-embed CHECKPOINT
* ({@link Brainy.PENDING_EMBED_CHECKPOINT_PATH}) "as of durable generation
* G the pending set was exactly this list" so the fold seeds the set from
* that list and reads only the facts after G. O(delta) whether or not the
* set ever drains, which is the whole point: the previous bound, the
* empty-only low-water mark, could not be written at all by a brain holding
* one id that never lands, so those brains re-read their whole log at every
* 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}.
* BOUND: the scan starts at the advisory low-water mark
* ({@link Brainy.PENDING_EMBED_LOWWATER_PATH}) the log head at which the
* pending set last drained to empty so a settled brain reads only the
* facts since then, not its whole history. Without a mark (first open
* after upgrade) it scans from generation 1, once; a stale-low mark 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 mark is what makes that cheap.
* It is SKIPPED WHOLESALE when the log has never had a v2 tail
* ({@link FactLog.hasV2History} v1 facts cannot carry marker records),
* so pre-cutover brains pay nothing; on a mixed log the scan still reads
@ -2855,13 +2526,20 @@ export class Brainy<T = any> implements BrainyInterface<T> {
private async recoverPendingEmbedsFromLog(): Promise<void> {
const log = this.generationStore.getFactLog()
if (!log || !log.hasV2History()) return
const { bound, fromGeneration, seeded } = await this.readPendingEmbedBound()
for (const id of seeded) this._pendingEmbedIds.add(id)
let factsScanned = 0
let fromGeneration = 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) {
fromGeneration = mark.generation + 1
}
} catch {
// No mark (or unreadable): scan from 1 — correctness over cost.
}
const scan = log.scanFacts({ fromGeneration })
for await (const batch of scan.batches()) {
for (const fact of batch.facts) {
factsScanned++
for (const record of fact.records ?? []) {
if (record.type === 'embed.pending') {
this._pendingEmbedIds.add(record.id)
@ -2876,21 +2554,6 @@ 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`
)
}
/**
@ -2982,23 +2645,11 @@ export class Brainy<T = any> implements BrainyInterface<T> {
for (const id of batch) {
try {
const entity = await this.get(id, { includeVectors: true })
if (!entity) {
// The row is GONE. Either it was deleted — its tombstone fact
// durably disarms the marker, at or below the head, exactly as the
// fold reads it — or its create never became durable, in which case
// 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')
if (!entity || entity.data === undefined || entity.data === null) {
// Orphan reap: a deleted row's tombstone fact durably disarms the
// marker at the next recovery fold; a data-less-but-present row
// (edge case) re-folds and re-reaps — bounded, never a lost vector.
this.clearPendingEmbed(id)
continue
}
// Hang guard: a wedged embedder must not block every later pending
@ -20331,21 +19982,6 @@ export class Brainy<T = any> implements BrainyInterface<T> {
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)
// Data is already safe on disk from Phase 1
await Promise.all([

View file

@ -40,10 +40,7 @@
* The manifest (`_generations/facts/manifest.json`, JSON forensics stay
* terminal-readable) is the single source of truth for the segment SET;
* rotation flips it atomically (write-new fsync rename) BEFORE the new
* tail's first byte exists, so no segment file is ever unaccounted for. Its
* per-segment `firstGeneration`/`lastGeneration` are LOAD-BEARING at open: a
* recovery pass looking for facts above a bound reads only the segments those
* bounds cannot rule out (the prune law see `segmentsHoldingFactsAbove`).
* tail's first byte exists, so no segment file is ever unaccounted for.
*
* ## Mixed-version logs (the v2 live-write cutover)
*
@ -692,74 +689,6 @@ function parseSegment(
return { facts, validBytes: offset, formatVersion: FACT_LOG_FORMAT_V1 }
}
/**
* THE PRUNE LAW which segment files a pass looking for facts ABOVE
* `committedGeneration` actually has to read, and how many the manifest's own
* recorded bounds took off the table.
*
* A sealed segment's `lastGeneration` is written at SEAL time and never
* mutated upward afterwards ({@link FactLog.rotate}, unchanged since the log
* was introduced): the tail's bytes are fsynced FIRST (`await this.sync()`
* "sealed segments are always fully durable"), the entry is then built from
* the content that fsync covered, and only then does the manifest flip
* atomically (tmp+rename) and fsynced which in the SAME write re-points
* `tailSegment` at a new file, so the sealed file is never appended to again.
* A crash anywhere in that order is safe in the pruning direction: crash
* before the manifest write and the segment is still the TAIL (read whole);
* crash after it and the entry describes bytes that were already durable. The
* only later mutation of a sealed segment is `open()`'s straddle truncation,
* which REMOVES facts and re-derives the entry from the actual bytes so a
* recorded bound can drift DOWN with its file, never up.
*
* Therefore: `lastGeneration = L` proves the file holds no fact above L, and
* a pass above `committedGeneration >= L` can skip it whole no read, no
* CRC decode, no msgpack. What the manifest cannot PROVE is never pruned: an
* entry with no numeric `lastGeneration` (a legacy or hand-repaired manifest)
* is read, and the unsealed tail is always read.
*
* This is the difference between an open that costs O(whole fact log) and one
* that costs O(the facts that could matter). MEASURED in production: a 16k-row
* brain at generation ~478,819 paid 34-37s of segment reads and CRC decoding
* in `generation-store-open-fold` on EVERY open to answer a question whose
* answer, after a clean close, is always "nothing".
*/
function segmentsHoldingFactsAbove(
stored: FactsManifest,
committedGeneration: number
): { files: string[]; pruned: number } {
const files: string[] = []
let pruned = 0
for (const entry of stored.segments) {
const last = (entry as Partial<SegmentEntry>).lastGeneration
if (typeof last === 'number' && Number.isFinite(last) && last <= committedGeneration) {
pruned++
continue
}
files.push(entry.file)
}
if (stored.tailSegment) files.push(stored.tailSegment)
return { files, pruned }
}
/**
* Say what the open actually read. One line, and only when the log holds more
* than one segment (a single-segment log has nothing to prune and nothing to
* report) the operator's receipt that the open is paying for the tail, not
* for the whole history.
*/
function narrateAboveScan(
pass: string,
committedGeneration: number,
read: number,
pruned: number
): void {
if (read + pruned <= 1) return
prodLog.narrate(
`[FactLog] ${pass} above generation ${committedGeneration}: ${read} segment(s) read, ` +
`${pruned} pruned of ${read + pruned} (sealed at or below the bound)`
)
}
/**
* The generation fact log. One instance per open store; every method assumes
* the single-writer discipline the generation store already enforces (calls
@ -825,6 +754,22 @@ export class FactLog {
return this.manifest.brainId !== undefined || this.tailVersion === FACT_LOG_FORMAT_V2
}
/**
* Open the log and reconcile it to committed truth: read the manifest,
* establish the tail's intact content (torn-tail scan), then TRUNCATE any
* fact with `generation > committedGeneration` those never committed (a
* crash between fact-append and the commit point). After open, the log is
* exactly the committed prefix.
*/
/**
* Read (without truncating) every intact fact ABOVE a generation the
* log-authority recovery surface: after a crash, facts beyond the
* manifest watermark that survived with valid CRCs are ACKED writes in
* durable-at-ack mode, and the owner REPLAYS them instead of letting
* open() truncate them. Must be called BEFORE open() (it reads the raw
* segments directly; the torn tail's invalid suffix is ignored exactly
* like open() would).
*/
/**
* STREAMING twin of {@link FactLog.peekFactsAbove} for the recovery fold:
* yields facts above the bound one SEGMENT at a time, ascending, without
@ -834,18 +779,13 @@ export class FactLog {
* Works manifest-direct (safe before {@link FactLog.open}). Ordering is
* structural (segments rotate in order; appends are ordered within one) and
* ASSERTED a violation aborts loudly, never a silent misordered replay.
*
* Reads only the segments that CAN hold a fact above the bound see
* {@link segmentsHoldingFactsAbove}. A bounded fold above a high checkpoint
* therefore reads its own tail, not the whole history it already proved
* durable.
*/
async *streamFactsAbove(committedGeneration: number): AsyncGenerator<CommitFact[], void> {
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return
const { files, pruned } = segmentsHoldingFactsAbove(stored, committedGeneration)
narrateAboveScan('recovery fold', committedGeneration, files.length, pruned)
const files = [...stored.segments.map((s) => s.file)]
if (stored.tailSegment) files.push(stored.tailSegment)
let lastGen = committedGeneration
for (const file of files) {
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
@ -867,27 +807,13 @@ export class FactLog {
}
}
/**
* Read (without truncating) every intact fact ABOVE a generation the
* log-authority recovery surface: after a crash, facts beyond the
* manifest watermark that survived with valid CRCs are ACKED writes in
* durable-at-ack mode, and the owner REPLAYS them instead of letting
* open() truncate them. Must be called BEFORE open() (it reads the raw
* segments directly; the torn tail's invalid suffix is ignored exactly
* like open() would).
*
* Reads only the segments that CAN hold such a fact see
* {@link segmentsHoldingFactsAbove}. This runs on EVERY log-authority open,
* including the clean one where the answer is always empty, so the segments
* the manifest already proves irrelevant are never opened at all.
*/
async peekFactsAbove(committedGeneration: number): Promise<CommitFact[]> {
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
if (!stored || typeof stored !== 'object' || !Array.isArray(stored.segments)) return []
if (stored.formatVersion !== FACTS_FORMAT_VERSION) return []
const out: CommitFact[] = []
const { files, pruned } = segmentsHoldingFactsAbove(stored, committedGeneration)
narrateAboveScan('above-manifest peek', committedGeneration, files.length, pruned)
const files = [...stored.segments.map((s) => s.file)]
if (stored.tailSegment) files.push(stored.tailSegment)
for (const file of files) {
const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`)
if (bytes === null) continue
@ -900,13 +826,6 @@ export class FactLog {
return out
}
/**
* Open the log and reconcile it to committed truth: read the manifest,
* establish the tail's intact content (torn-tail scan), then TRUNCATE any
* fact with `generation > committedGeneration` those never committed (a
* crash between fact-append and the commit point). After open, the log is
* exactly the committed prefix.
*/
async open(committedGeneration: number): Promise<void> {
const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null
if (stored && typeof stored === 'object' && Array.isArray(stored.segments)) {

View file

@ -1,360 +0,0 @@
/**
* @module tests/integration/factlog-open-prune
* @description THE OPEN READS THE TAIL, NOT THE HISTORY.
*
* Every log-authority open asks the fact log one question "is there a fact
* above the committed pointer?" and until this lane existed it answered by
* reading and CRC-decoding EVERY segment file the manifest names. MEASURED in
* production on a 16k-row brain at generation ~478,819: 34-37 seconds inside
* the `generation-store-open-fold` phase, on every open, including the clean
* one where the answer is always "nothing".
*
* The manifest already records each sealed segment's `lastGeneration`, written
* at seal time AFTER the segment's bytes are fsynced and into a manifest that
* is itself written atomically and fsynced and a sealed file is never
* appended to again (the same manifest flip re-points `tailSegment`). So an
* entry recording `lastGeneration ≤ committed` PROVES its file holds nothing
* above the bound, and the open can skip it whole.
*
* Pinned here, from the log's own counters (the narration line), never a clock:
*
* 1. A clean close and reopen on a log with 4 sealed segments reads
* EXACTLY the tail (1 of 6), prunes the rest, and finds nothing.
* 2. A real SIGKILLed process that sealed segments holding facts ABOVE the
* committed pointer: the reopen READS those sealed segments and recovers
* byte-identically to an unpruned open (differential the same store,
* with the provable field stripped from its manifest, takes the full-scan
* path and must agree fact for fact, before and after `open()`).
* 3. A manifest entry with no `lastGeneration` (legacy, or hand-repaired) is
* READ. Never prune what the manifest cannot prove.
*/
import { describe, it, expect, afterEach } from 'vitest'
import * as fs from 'node:fs'
import * as os from 'node:os'
import * as path from 'node:path'
import { spawn } from 'node:child_process'
import {
FactLog,
FACTS_MANIFEST_PATH,
type CommitFact,
type FactLogStorage
} from '../../src/db/factLog.js'
import { FileSystemStorage } from '../../src/storage/adapters/fileSystemStorage.js'
const REPO_ROOT = process.cwd()
const TSX = path.join(REPO_ROOT, 'node_modules', '.bin', 'tsx')
/** ~1KB frames against a 4KB rotation threshold: ~5 facts per segment. */
const ROTATE_BYTES = 4096
const tmpDirs: string[] = []
function makeTempDir(): string {
const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-factlog-prune-'))
tmpDirs.push(dir)
return dir
}
afterEach(() => {
for (const dir of tmpDirs.splice(0)) {
try {
fs.rmSync(dir, { recursive: true, force: true })
} catch {
/* best effort */
}
try {
fs.rmSync(`${dir}.ready.json`, { force: true })
} catch {
/* best effort */
}
}
})
const UUID = (n: number): string => `00000000-0000-4000-8000-${String(n).padStart(12, '0')}`
/** One ~1KB fact — the padding is what makes rotation cheap to provoke. */
function fact(generation: number): CommitFact {
return {
generation,
timestamp: 1_700_000_000_000 + generation,
ops: [
{
kind: 'noun',
id: UUID(generation),
record: {
metadata: { noun: 'document', pad: 'x'.repeat(900), g: generation },
vector: null
}
}
]
}
}
/**
* A deterministic int minter so the log writes the V2 format production
* writes (the prune is a manifest-level decision and never touches segment
* bytes but the pins should run against the bytes the fleet actually has).
*/
function makeMinter(): (kind: 'noun' | 'verb', id: string) => bigint {
const ints = new Map<string, bigint>()
return (kind, id) => {
const key = `${kind}:${id}`
let minted = ints.get(key)
if (minted === undefined) {
minted = BigInt(ints.size + 1)
ints.set(key, minted)
}
return minted
}
}
/** Open a fact log over a store directory (a fresh adapter each time this is
* what a reopen actually does). */
async function openStore(dir: string): Promise<{ storage: any; log: FactLog }> {
const storage: any = new FileSystemStorage(dir)
await storage.init()
const log = new FactLog(storage as FactLogStorage, { rotateBytes: ROTATE_BYTES })
log.setIntMinter(makeMinter())
return { storage, log }
}
/** Build a log of `count` facts (rotating every ~5), left durable, not closed. */
async function buildLog(dir: string, count: number): Promise<number> {
const { log } = await openStore(dir)
await log.open(0)
for (let g = 1; g <= count; g++) await log.append(fact(g))
await log.sync()
return log.headGeneration()
}
/** Capture the narration channel (`prodLog.narrate` → console.warn). */
async function captureNarration<T>(
fn: () => Promise<T>
): Promise<{ result: T; lines: string[] }> {
const lines: string[] = []
const original = console.warn
console.warn = ((...args: unknown[]) => {
lines.push(args.map((a) => String(a)).join(' '))
}) as typeof console.warn
try {
return { result: await fn(), lines }
} finally {
console.warn = original
}
}
/** The counters the open narrated the pin's only source of truth for what
* was read (a wall-clock assertion could pass on a warm page cache). */
function scanCounts(lines: string[]): { read: number; pruned: number; total: number } {
const line = lines.find((l) => l.includes('[FactLog] above-manifest peek above generation'))
if (!line) {
throw new Error(`no peek narration in:\n${lines.join('\n')}`)
}
const match = /(\d+) segment\(s\) read, (\d+) pruned of (\d+)/.exec(line)
if (!match) throw new Error(`unparsable peek narration: ${line}`)
return { read: Number(match[1]), pruned: Number(match[2]), total: Number(match[3]) }
}
interface SegmentEntryOnDisk {
file: string
firstGeneration: number
lastGeneration?: number
facts: number
bytes: number
}
async function readManifest(dir: string): Promise<{
segments: SegmentEntryOnDisk[]
tailSegment: string | null
}> {
const storage: any = new FileSystemStorage(dir)
await storage.init()
return (await storage.readRawObject(FACTS_MANIFEST_PATH)) as any
}
async function rewriteManifest(
dir: string,
mutate: (manifest: any) => void
): Promise<void> {
const storage: any = new FileSystemStorage(dir)
await storage.init()
const manifest = await storage.readRawObject(FACTS_MANIFEST_PATH)
mutate(manifest)
await storage.writeRawObject(FACTS_MANIFEST_PATH, manifest)
await storage.syncRawObjects([FACTS_MANIFEST_PATH])
}
/** Every fact the log holds, in order — the recovered state, read back. */
async function allFacts(log: FactLog): Promise<CommitFact[]> {
const out: CommitFact[] = []
const handle = log.scanFacts()
for await (const batch of handle.batches()) out.push(...batch.facts)
return out
}
describe('fact log — the open reads only the segments that can hold facts above the bound', () => {
it('a clean close + reopen over ≥4 sealed segments reads exactly the tail and finds nothing', async () => {
const dir = makeTempDir()
const head = await buildLog(dir, 30)
const manifest = await readManifest(dir)
expect(manifest.segments.length).toBeGreaterThanOrEqual(4) // the fixture is real
expect(manifest.tailSegment).not.toBeNull()
// The reopen: a clean close means committed === the log's head.
const { log } = await openStore(dir)
const { result: orphans, lines } = await captureNarration(() => log.peekFactsAbove(head))
expect(orphans).toEqual([]) // the fold finds nothing, as it always does after a clean close
const counts = scanCounts(lines)
expect(counts.read).toBe(1) // EXACTLY the tail
expect(counts.total).toBe(manifest.segments.length + 1)
expect(counts.pruned).toBe(manifest.segments.length)
// And the reconciling open still lands on the same committed prefix.
await log.open(head)
expect(log.headGeneration()).toBe(head)
expect((await allFacts(log)).map((f) => f.generation)).toEqual(
Array.from({ length: head }, (_, i) => i + 1)
)
})
it('a manifest entry with no lastGeneration is READ — never prune what you cannot prove', async () => {
const dir = makeTempDir()
const head = await buildLog(dir, 30)
const before = await readManifest(dir)
expect(before.segments.length).toBeGreaterThanOrEqual(4)
// A legacy/hand-repaired entry: the field the prune needs is simply absent.
await rewriteManifest(dir, (m) => {
delete m.segments[0].lastGeneration
})
const { log } = await openStore(dir)
const { result: orphans, lines } = await captureNarration(() => log.peekFactsAbove(head))
expect(orphans).toEqual([]) // still nothing above the bound — it was READ to find out
const counts = scanCounts(lines)
expect(counts.read).toBe(2) // the unprovable entry + the tail
expect(counts.pruned).toBe(before.segments.length - 1)
expect(counts.total).toBe(before.segments.length + 1)
})
it(
'a SIGKILLed writer that sealed segments above the committed pointer recovers identically to an unpruned open',
async () => {
const dir = makeTempDir()
const readyPath = `${dir}.ready.json`
// A real process death: the child fsyncs its segments, records what it
// reached, and SIGKILLs ITSELF — no close, no unwind, no chance to tidy.
const script = `
import * as fs from 'node:fs'
import { FactLog } from ${JSON.stringify(path.join(REPO_ROOT, 'src', 'db', 'factLog.ts'))}
import { FileSystemStorage } from ${JSON.stringify(path.join(REPO_ROOT, 'src', 'storage', 'adapters', 'fileSystemStorage.ts'))}
const UUID = (n) => '00000000-0000-4000-8000-' + String(n).padStart(12, '0')
const fact = (g) => ({
generation: g,
timestamp: 1700000000000 + g,
ops: [{ kind: 'noun', id: UUID(g), record: { metadata: { noun: 'document', pad: 'x'.repeat(900), g }, vector: null } }]
})
const ints = new Map()
const storage = new FileSystemStorage(${JSON.stringify(dir)})
await storage.init()
const log = new FactLog(storage, { rotateBytes: ${ROTATE_BYTES} })
log.setIntMinter((kind, id) => {
const key = kind + ':' + id
if (!ints.has(key)) ints.set(key, BigInt(ints.size + 1))
return ints.get(key)
})
await log.open(0)
for (let g = 1; g <= 30; g++) await log.append(fact(g))
await log.sync()
fs.writeFileSync(${JSON.stringify(readyPath)}, JSON.stringify({ head: log.headGeneration() }))
process.kill(process.pid, 'SIGKILL')
`
const scriptPath = path.join(dir, 'crash-writer.mts')
fs.writeFileSync(scriptPath, script)
const child = spawn(TSX, [scriptPath], { cwd: REPO_ROOT, stdio: ['ignore', 'pipe', 'pipe'] })
let output = ''
child.stdout.on('data', (d) => { output += String(d) })
child.stderr.on('data', (d) => { output += String(d) })
const exit = await new Promise<{ code: number | null; signal: string | null }>((resolve) =>
child.on('exit', (code, signal) => resolve({ code, signal }))
)
if (!fs.existsSync(readyPath)) {
throw new Error(`the crash writer never reached its kill point:\n${output}`)
}
// Death, not a shutdown: no close(), no unwind, no orderly exit code.
expect(exit.signal ?? `code ${exit.code}`).not.toBe('code 0')
const head = JSON.parse(fs.readFileSync(readyPath, 'utf8')).head as number
expect(head).toBe(30)
// The committed pointer the survivor comes back on: mid-log, so sealed
// segments hold facts ABOVE it — the exact shape the prune must not skip.
const committed = 12
const manifest = await readManifest(dir)
const straddling = manifest.segments.filter(
(s) => s.firstGeneration <= committed && (s.lastGeneration ?? 0) > committed
)
const entirelyAbove = manifest.segments.filter((s) => s.firstGeneration > committed)
expect(straddling.length).toBeGreaterThanOrEqual(1)
expect(entirelyAbove.length).toBeGreaterThanOrEqual(1)
// THE DIFFERENTIAL. The unpruned answer, through the SAME code on the
// SAME bytes: a peek above generation 0 can prune nothing (no sealed
// segment ends at or below 0), so it reads every segment file and
// decodes every frame — exactly what this open used to do — and its
// facts above the pointer are what the fold is entitled to replay.
const { log } = await openStore(dir)
const { result: fullScan, lines: fullLines } = await captureNarration(() =>
log.peekFactsAbove(0)
)
expect(scanCounts(fullLines)).toEqual({
read: manifest.segments.length + 1,
pruned: 0,
total: manifest.segments.length + 1
})
const unprunedAnswer = fullScan.filter((f) => f.generation > committed)
const { result: prunedAnswer, lines } = await captureNarration(() =>
log.peekFactsAbove(committed)
)
// The sealed segments above the bound were READ, not skipped.
const counts = scanCounts(lines)
expect(counts.read).toBe(straddling.length + entirelyAbove.length + 1)
expect(counts.pruned).toBe(manifest.segments.length - straddling.length - entirelyAbove.length)
expect(counts.pruned).toBeGreaterThan(0) // the prune did engage, and was still right
expect(prunedAnswer.map((f) => f.generation)).toEqual(
Array.from({ length: head - committed }, (_, i) => committed + 1 + i)
)
// Facts that live in a SEALED segment (not the tail) came back.
expect(prunedAnswer.some((f) => f.generation <= (straddling[0].lastGeneration ?? 0))).toBe(
true
)
// Fact for fact, the pruned answer IS the unpruned answer — so whatever
// the recovery replays, it replays identically.
expect(prunedAnswer).toEqual(unprunedAnswer)
// The fold's streaming twin (the unclean-open path) agrees too.
const streamed: CommitFact[] = []
for await (const batch of log.streamFactsAbove(committed)) streamed.push(...batch)
expect(streamed).toEqual(unprunedAnswer)
// And the reconciling open rolls back exactly as it always did: the two
// never-committed sealed segments dropped, the straddling one cut, the
// tail truncated — the log left as the committed prefix.
await log.open(committed)
expect(log.headGeneration()).toBe(committed)
expect((await allFacts(log)).map((f) => f.generation)).toEqual(
Array.from({ length: committed }, (_, i) => i + 1)
)
const after = await readManifest(dir)
expect(after.segments.map((s) => s.file)).toEqual(
manifest.segments
.filter((s) => s.firstGeneration <= committed)
.map((s) => s.file)
)
expect(after.segments[after.segments.length - 1].lastGeneration).toBe(committed)
},
120_000
)
})

View file

@ -1,547 +0,0 @@
/**
* @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)
})