Compare commits

..

4 commits

Author SHA1 Message Date
367ca721a5 fix(close): a read-only brain writes no clean-shutdown evidence — the marker is the writer's word about itself
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:59:29 -07:00
a79db434ac fix(generation-store): commitTransaction refuses while single-ops are pending — the order invariant is enforced, not assumed
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
reservedGensAsc() documented but never enforced that pending single-op
generations must sort above every committed one. A direct
commitTransaction() call bypassing Brainy.transact()'s flush-first step
could commit a fresh generation into committedRanges above lower,
still-pending ones, unsorting the committed-then-pending walk
resolveManyAt relies on and returning a wrong before-image for a
point-in-time read — silently. commitTransaction() now refuses via a new
PendingSingleOpsUnflushedError when the pending tier is non-empty, before
any staging I/O. Behavior-neutral: both sanctioned callers
(Brainy.transact(), Brainy.compactHistory()) already flush first.
2026-09-02 10:55:40 -07:00
da9519903a test(shutdown): pin one owner per brain — real processes, real signals
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
Four pins in real child processes under real SIGTERM, following the
writer-lock-clean-close spawn pattern:

(a) A host owner registered on SIGTERM closes two brains while the engine's
    hooks are live: exactly one close entered and one close body run per brain,
    the writer lock given up exactly ONCE per brain, the handler announcing
    that it stepped aside, no "Writer fence lost", no failed instance, both
    durability markers written, exit 0, and both reopens adopting rather than
    folding. The release count is the discriminating assertion — against the
    old handler it reads {a: 2, b: 2}, one release from the owner's close and
    one from the handler's own finally.
(b) No host owner: the engine's handler closes every instance by the same
    path — one close each, markers written, clean exit, clean reopen.
(c) Two concurrent close() callers share one promise (by identity) and one
    execution; a third call after they settle runs nothing.
(d) Eight kicks during a running flush — five through the cadence door, three
    direct — arm exactly ONE follow-up: two flush bodies total, and the
    concurrency high-water mark stays at 1.

The counts come out of the child through a file written synchronously on the
way out: the engine calls process.exit(0) when it is the sole shutdown owner,
and a console.log to a pipe can be dropped by that exit.
2026-09-02 10:55:09 -07:00
ec644bde56 fix(shutdown): one owner per brain — the signal handler defers to close(), and flush is single-flight
MEASURED IN PRODUCTION. A host that owns its own shutdown — one SIGTERM
listener calling close() on every pooled store — ran head-on into the engine's
own signal handler, which iterated every live instance, flushed its components
in parallel, and released its writer lock in its own finally. Two teardowns of
the same brain at the same moment: "Shutdown signal received - flushing pending
data...", 148s of silence, "Flushed successfully (1 instance)", and the host's
pool close of that same store returning 1s later — 149s against 24s for the six
stores with no engine work in flight. The same race reproduced locally as
"Failed to flush one Brainy instance on shutdown: Writer fence lost … the lock
file is gone": the handler observing a lock the close it was racing had already
released.

Three changes, one law — a brain's teardown belongs to whoever started it.

1. close() is idempotent and re-entrant. The first call stores its promise
   synchronously in _closeInFlight and every later or concurrent caller gets
   that same promise back; the teardown runs once. close() is no longer async
   so the promise is shared by identity, not just outcome. The state is
   observable: isClosing (begun) and isClosed (finished).

2. The signal handler defers one macrotask, then per instance either steps
   aside (a close has begun or finished — its owner owns the flush, the markers
   and the lock) or awaits instance.close(): the same settle/flush/attest/
   marker/lock path any caller gets. Its old parallel per-component flush and
   separate lock release are gone; the three laws that block carried are each
   satisfied by close(), verified line by line and recorded in the new comment.
   Per-instance isolation stays here, in the loop's try/catch.

   Sole-owner exit now reads the listener count WHEN THE SIGNAL ARRIVES.
   Asking afterwards reads a process that has already torn itself down —
   closing the last brain deregisters the engine's own listeners, so a host's
   single remaining listener would look like "<= 1" and be force-exited out of
   its own graceful shutdown.

3. Flush is single-flight with a queue one deep. It did not coalesce: the
   cadence's guard covered only the flushes the cadence started, so a
   cross-process flush request or an application flush() overlapped it freely —
   production showed two "Flushing Brainy indexes…" runs 3s apart, walls
   growing 295ms to 4.9s. The gate now lives in flush() itself and covers every
   caller: run, or join the ONE queued follow-up. A follow-up rather than
   joining the running flush, because a caller flushes to make ITS writes
   durable and those may have landed after the running flush read its state; it
   costs nothing when there is nothing new. close() drains that chain too.

The idle law is untouched: a clean brain's flush still returns immediately, and
an idle brain still flushes zero times.
2026-09-02 10:55:09 -07:00
6 changed files with 1043 additions and 99 deletions

View file

@ -767,6 +767,34 @@ export class Brainy<T = any> implements BrainyInterface<T> {
private _persistIdleTimer: ReturnType<typeof setTimeout> | null = null
private _persistBackgroundFlight: Promise<void> | null = null
/**
* FLUSH IS SINGLE-FLIGHT, AND THE QUEUE IS ONE DEEP. `_flushInFlight` is the
* flush body actually running; `_flushFollowUp` is the AT MOST ONE flush
* queued behind it. Every caller the write cadence, the cross-process
* flush-request watcher, an application calling `flush()` directly either
* runs (nothing in flight), or joins the single queued follow-up.
*
* WHY A FOLLOW-UP RATHER THAN JOINING THE RUNNING FLUSH: a caller flushes to
* make ITS writes durable, and those writes may have landed after the
* running flush read its state. Joining would return "flushed" over data
* that was never persisted. Chaining one follow-up costs nothing when there
* is nothing new (a clean brain's flush returns immediately see
* `_dirtySinceLastFlush`) and is correct when there is.
*
* MEASURED, in the production shutdown this was written for: two
* "Flushing Brainy indexes and caches to disk..." runs overlapping 3s
* apart on one brain, their walls growing 295ms 4.9s as they contended
* for the same providers.
*/
private _flushInFlight: Promise<void> | null = null
private _flushFollowUp: Promise<void> | null = null
/** Flush bodies that got past the single-flight gate (pinned by tests). */
private _flushBodyRuns = 0
/** Flush bodies running right now, and the high-water mark which the
* single-flight law requires to stay at 1 (pinned by tests). */
private _flushBodiesActive = 0
private _flushConcurrencyPeak = 0
// DEFERRED EMBEDDING (MT5): pending markers are LOG RECORDS — an
// embed.pending record rides the deferred write's own commit fact and
// embed.landed rides the landing commit; this set is the in-memory
@ -889,6 +917,24 @@ export class Brainy<T = any> implements BrainyInterface<T> {
// applies only to instances that were never closed.
private closed = false
/**
* THE ONE CLOSE. Set SYNCHRONOUSLY by the first `close()` call, before that
* call yields, and never cleared close is terminal. Every later or
* concurrent caller receives this same promise, so a shutdown with two
* callers (a host's pool close and the engine's own signal handler) runs
* ONE teardown, not two.
*
* MEASURED, the day this was added: a host that owns shutdown called
* `close()` on every pooled store at SIGTERM while the engine's signal
* handler flushed the same instances in parallel and released their writer
* locks in its own `finally`. One store took 149s to close (148s of it
* silent) against 24s for its idle siblings, and the same race in a local
* reproduction printed `Writer fence lost … the lock file is gone` the
* handler observing a lock the close it was racing had already released.
* Two owners of one shutdown; now there is one, whoever calls first.
*/
private _closeInFlight: Promise<void> | null = null
// Index-build-at-open state. `lazyRebuildCompleted` predates the health-gate
// law (it named a first-QUERY lazy rebuild) and stays for `getIndexStatus()`
// API compatibility, but its truth changed: a needed rebuild now runs
@ -2076,105 +2122,88 @@ export class Brainy<T = any> implements BrainyInterface<T> {
*/
private registerShutdownHooks(): void {
/**
* The signal-path shutdown. THREE LAWS, each written by a production
* shutdown that looked clean and wasn't:
* The signal-path shutdown. ONE OWNER PER BRAIN, AND THE PATH IS `close()`.
*
* 1. PER-INSTANCE ISOLATION. This used to be one `try` around a loop over
* every open brain: the first instance whose flush rejected aborted the
* loop, so every remaining brain kept its writer lock and its unwritten
* markers and the process still exited 0. A pool of brains failed in
* a batch, not one at a time.
* 2. THE MARKER IS PART OF SHUTDOWN. Flushing the indexes without closing
* the generation store leaves the clean-shutdown marker unwritten, so
* the NEXT open reads the store as crashed and folds the whole
* generation log measured in tens of seconds on a real store, paid on
* every restart, after a shutdown the operator saw exit 0.
* 3. THE LOCK IS ALWAYS GIVEN UP. In a `finally`, per instance: a process
* on its way out holds nothing.
* WHAT THIS REPLACED, and why. The handler used to run its own shutdown
* a parallel per-component flush, the generation store's close, a second
* parallel round of component closes, and a `finally` that stopped the
* flush-request watcher and released the writer lock. That is a SECOND
* teardown of the same brain, and a host application with its own SIGTERM
* handler (the shape every pooled deployment has) ran the FIRST one at the
* same moment. MEASURED in production the day this changed: a host closing
* seven pooled stores at SIGTERM printed "Shutdown signal received -
* flushing pending data...", went silent for 148s, printed "Flushed
* successfully (1 instance)", and the host's own close of that same store
* returned 1s later 149s, against 24s for the six stores with no engine
* work in flight. The same race reproduced locally as
* `Failed to flush one Brainy instance on shutdown: Writer fence lost …
* the lock file is gone`: this handler observing a lock that the close it
* was racing had already released.
*
* SO: defer one macrotask, then per instance either STEP ASIDE (a close
* has begun or finished its owner owns the flush, the markers and the
* lock) or `await instance.close()` the one durable path, identical to
* what any caller gets. The three laws the old block carried are all
* satisfied by `close()`, each verified against its code:
*
* 1. PER-INSTANCE ISOLATION kept HERE, in the per-instance try/catch
* below: one brain's failed close never aborts the loop over the rest.
* (`close()` itself is per-instance by construction.)
* 2. THE MARKER IS PART OF SHUTDOWN `close()` `closeDurableSteps()`
* Phase 1 awaits `this.generationStore.close()`, which persists the
* counter, advances the fold checkpoint and stamps the clean-shutdown
* marker LAST. That is the step that decides adopt-vs-fold at the next
* open, and it is the same call the old block made.
* 3. THE LOCK IS ALWAYS GIVEN UP `close()`'s terminal releases run
* whether the durable steps threw or not (its contract: "TWO PARTS, AND
* THE SECOND IS UNCONDITIONAL"): `stopFlushRequestWatcher()` then
* `releaseWriterLock()`, then the VFS shutdown and the terminal
* `closed` flag, and only then is the original failure rethrown.
* `close()` releases the lock in MORE cases than the old block did it
* also drains the metadata write buffer first, so no pending write can
* land after a successor writer claims the lock.
*/
const flushOnShutdown = async () => {
const closeOnShutdown = async () => {
console.log('Shutdown signal received - flushing pending data...')
let flushedCount = 0
// DEFER ONE MACROTASK. A host application registers its own listener on
// the same signal, and Node runs listeners in registration order — ours
// is usually first, because the brain was opened before the host wired
// its shutdown. Yielding once lets every other listener for this signal
// run its synchronous prologue, so a host that calls close() gets to be
// the owner. It is only a courtesy, never the safety: close()'s own
// single-flight gate is what makes a lost race harmless.
await new Promise<void>((resolve) => setImmediate(resolve))
let closedCount = 0
let deferredCount = 0
let failedCount = 0
// Snapshot: close() splices Brainy.instances while we iterate.
for (const instance of [...Brainy.instances]) {
if (!instance.initialized) continue
// SOMEONE ELSE OWNS THIS ONE. Not a flush, not a lock release, not a
// component close — nothing. Touching a brain whose close is running
// is the whole defect this handler was rewritten for.
if (instance.closed || instance._closeInFlight !== null) {
deferredCount++
continue
}
try {
// Flush all buffered data (parallel across components, this brain only).
await Promise.all([
(async () => {
if (instance.storage && typeof instance.storage.flushCounts === 'function') {
await instance.storage.flushCounts()
}
})(),
(async () => {
if (instance.metadataIndex && typeof instance.metadataIndex.flush === 'function') {
await instance.metadataIndex.flush()
}
})(),
(async () => {
if (instance.graphIndex && typeof instance.graphIndex.flush === 'function') {
await instance.graphIndex.flush()
}
})(),
(async () => {
if (instance.index && typeof instance.index.flush === 'function') {
await instance.index.flush()
}
})()
])
// Close the generation store: persists the counter, advances the
// fold checkpoint, and stamps the clean-shutdown marker LAST — the
// one step that decides whether the next open adopts or folds. Law 2.
if (instance.generationStore && !instance.isReadOnly) {
await instance.generationStore.close()
}
// Close components to stop timers that would prevent clean process exit
await Promise.all([
(async () => {
if (instance.graphIndex && typeof instance.graphIndex.close === 'function') {
await instance.graphIndex.close()
}
})(),
(async () => {
const index = instance.index as JsHnswVectorIndex & VectorIndexOptionalHooks
if (index && typeof index.close === 'function') {
await index.close()
}
})(),
(async () => {
const metadataIndex = instance.metadataIndex as MetadataIndexManager & MetadataIndexOptionalHooks
if (metadataIndex && typeof metadataIndex.close === 'function') {
await metadataIndex.close()
}
})()
])
flushedCount++
// Law 1: this try/catch is the isolation — the loop continues.
await instance.close()
closedCount++
} catch (error) {
failedCount++
console.error('Failed to flush one Brainy instance on shutdown:', error)
} finally {
// Law 3 — the lock and the watcher go regardless.
try {
if (instance.storage && typeof instance.storage.stopFlushRequestWatcher === 'function') {
instance.storage.stopFlushRequestWatcher()
}
} catch (error) {
console.error('Failed to stop the flush-request watcher on shutdown:', error)
}
try {
if (instance.storage && typeof instance.storage.releaseWriterLock === 'function') {
await instance.storage.releaseWriterLock()
}
} catch (error) {
console.error('Failed to release the writer lock on shutdown:', error)
}
console.error('Failed to close one Brainy instance on shutdown:', error)
}
}
if (flushedCount > 0) {
console.log(`Flushed successfully (${flushedCount} instance${flushedCount > 1 ? 's' : ''})`)
if (closedCount > 0) {
console.log(`Flushed successfully (${closedCount} instance${closedCount > 1 ? 's' : ''})`)
}
if (deferredCount > 0) {
console.log(
`${deferredCount} Brainy instance${deferredCount > 1 ? 's are' : ' is'} already ` +
`closing — left to the caller that owns that close.`
)
}
if (failedCount > 0) {
console.error(
@ -2201,19 +2230,29 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* markers unwritten. When the host has its own handler (listener count
* above our own), the host owns the exit; Brainy only makes its data
* durable and steps aside.
*
* THE COUNT IS TAKEN WHEN THE SIGNAL ARRIVES, not after the shutdown ran.
* "Is anyone else handling this signal?" is a question about the moment
* the signal landed. Asking afterwards reads a process that has already
* torn itself down: the handler now CLOSES its instances, and closing the
* last brain deregisters Brainy's own listeners — so a host application's
* single remaining listener would look like `<= 1` and get force-exited
* out of its own graceful shutdown, precisely the failure above.
*/
const exitIfSoleShutdownOwner = (signal: 'SIGTERM' | 'SIGINT'): void => {
if (process.listenerCount(signal) <= 1) {
const exitIfSoleShutdownOwner = (ownersWhenSignalled: number): void => {
if (ownersWhenSignalled <= 1) {
process.exit(0)
}
}
Brainy.sigtermListener = async () => {
await flushOnShutdown()
exitIfSoleShutdownOwner('SIGTERM')
const owners = process.listenerCount('SIGTERM')
await closeOnShutdown()
exitIfSoleShutdownOwner(owners)
}
Brainy.sigintListener = async () => {
await flushOnShutdown()
exitIfSoleShutdownOwner('SIGINT')
const owners = process.listenerCount('SIGINT')
await closeOnShutdown()
exitIfSoleShutdownOwner(owners)
}
Brainy.beforeExitListener = async () => {
// Self-deregister FIRST: Node re-emits 'beforeExit' after every event-
@ -2225,7 +2264,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
process.off('beforeExit', Brainy.beforeExitListener)
Brainy.beforeExitListener = undefined
}
await flushOnShutdown()
await closeOnShutdown()
}
process.on('SIGTERM', Brainy.sigtermListener)
process.on('SIGINT', Brainy.sigintListener)
@ -2298,6 +2337,33 @@ export class Brainy<T = any> implements BrainyInterface<T> {
return this.initialized
}
/**
* @description Whether `close()` has BEGUN on this instance in flight or
* already finished. The question a shutdown owner asks: this brain's
* teardown belongs to whoever started it, and a second party must not flush
* its components or release its writer lock underneath it.
*
* True from the synchronous moment `close()` is entered, so a listener that
* yields a tick and comes back reads the truth, not a stale "not yet".
* @returns `true` once a close has started.
*/
get isClosing(): boolean {
return this._closeInFlight !== null
}
/**
* @description Whether `close()` has FINISHED tearing this instance down
* durable steps attempted, writer lock released, instance terminal. A
* closed brain never re-initializes; every operation on it throws.
*
* True after a close that FAILED partway, too: such a brain still holds no
* writer lock and still serves nothing (see {@link close}).
* @returns `true` once the teardown has completed.
*/
get isClosed(): boolean {
return this.closed
}
/**
* Promise that resolves when Brainy is fully initialized and ready to use
*
@ -3271,9 +3337,18 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* toward the next trigger. A failure is LOUD and leaves the writes counted
* again silence is not an option, and neither is a retry storm (the next
* trigger re-attempts).
*
* COALESCING LIVES IN {@link flush}, NOT HERE. A kick that arrives while a
* flush is running used to return without doing anything the writes it
* counted waited for some LATER trigger, and this method's guard also could
* not coalesce the flushes it does not start (the cross-process
* flush-request watcher and application `flush()` calls both go straight to
* `flush()`; two of those overlapping is exactly what production showed).
* The gate in `flush()` covers every caller: this kick now either runs the
* flush or joins the single queued follow-up, so the writes it counted are
* always someone's work, and there is still never a second concurrent run.
*/
private kickBackgroundFlush(reason: 'threshold' | 'idle'): void {
if (this._persistBackgroundFlight) return
const counted = this._persistDirtyWrites
this._persistDirtyWrites = 0
this._persistLastFlushAt = Date.now()
@ -12903,7 +12978,58 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* process.exit(0)
* })
*/
async flush(): Promise<void> {
flush(): Promise<void> {
// ---- THE SINGLE-FLIGHT GATE ----
// One flush body runs at a time, with at most ONE queued behind it. See
// `_flushInFlight` / `_flushFollowUp` for the measurement that required
// this. NOT `async`: the gate hands back the very promise the work is on,
// so joining callers share identity, not just an outcome. The gate is
// crossed BEFORE any await, so two callers in the same tick cannot both
// find the field empty.
if (this._flushInFlight) {
if (!this._flushFollowUp) {
// The running flush's failure is not this follow-up's failure: it is
// reported to ITS caller, and the queued work still gets its turn.
this._flushFollowUp = this._flushInFlight
.catch(() => {})
.then(() => {
this._flushFollowUp = null
return this.flush()
})
}
return this._flushFollowUp
}
const run = this._runFlush()
// `finally` and not `then`: a failed flush must still open the gate, or
// one rejection would wedge every later flush behind a promise nobody
// will ever settle.
const gated = run.finally(() => {
if (this._flushInFlight === gated) this._flushInFlight = null
})
this._flushInFlight = gated
return gated
}
/**
* @description The flush body everything {@link flush} promises, run
* exactly once at a time by that method's single-flight gate. Private
* because non-overlap is part of the contract: there is no supported way to
* run two of these at once, and the counters here witness that.
* @returns Nothing.
*/
private async _runFlush(): Promise<void> {
this._flushBodyRuns++
this._flushBodiesActive++
this._flushConcurrencyPeak = Math.max(this._flushConcurrencyPeak, this._flushBodiesActive)
try {
await this._flushSteps()
} finally {
this._flushBodiesActive--
}
}
/** @description The flush steps themselves. See {@link flush}. */
private async _flushSteps(): Promise<void> {
await this.ensureInitialized()
// Read-only instances have no buffered writes to flush. close() may call
@ -20150,11 +20276,42 @@ export class Brainy<T = any> implements BrainyInterface<T> {
*
* The original failure is never swallowed: it is narrated with what it costs
* the next open, then rethrown to the caller.
*
* IDEMPOTENT AND RE-ENTRANT. The teardown below runs ONCE. Concurrent
* callers share the one in-flight promise and settle together; a caller
* arriving after it finished gets that same settled promise (close is
* terminal there is nothing left to redo, and a failed close has already
* released the lock and set `closed`). This is what makes the shutdown
* ownership question answerable at all: whoever calls first owns the close,
* everyone else including the engine's own signal handler joins it or
* steps aside. See `_closeInFlight`.
* @returns Nothing.
* @throws The first failure from the durable close steps, after the
* terminal releases have run.
*/
async close(): Promise<void> {
close(): Promise<void> {
// NOT `async`: an async wrapper allocates a FRESH promise per call, so
// callers would hold different handles to the same work. Returning the
// stored promise itself makes "one close" observable identity, not just
// observable behaviour. The gate is crossed with NO await before it, so
// two callers in the same tick — and a signal handler resuming mid-close
// — always see the same answer; `isClosing` is true from this assignment
// onward. (`_closeOnce()` is async, so a failure is always a rejection,
// never a synchronous throw out of this method.)
if (this._closeInFlight) return this._closeInFlight
const run = this._closeOnce()
this._closeInFlight = run
return run
}
/**
* @description The close body everything {@link close} promises, run
* exactly once by that method's gate.
* @returns Nothing.
* @throws The first failure from the durable close steps, after the
* terminal releases have run.
*/
private async _closeOnce(): Promise<void> {
if (this._pendingEmbedIds.size === 0) await this.writeEmbedLowWater()
let closeFailure: unknown = null
try {
@ -20243,6 +20400,19 @@ export class Brainy<T = any> implements BrainyInterface<T> {
if (this._persistBackgroundFlight) {
await this._persistBackgroundFlight.catch(() => {})
}
// Drain the flush chain itself: the running flush AND the single follow-up
// queued behind it. The cadence's own handle above covers only the flushes
// the cadence started — a flush-request from another process, or an
// application's own flush() racing this close, is on the chain and nowhere
// else, and a flush landing mid-close writes behind the close's work.
// Bounded by construction: at most one follow-up exists, and awaiting it
// awaits its leader too, so the second pass is a no-op unless a writer
// raced this close.
for (let pass = 0; pass < 2; pass++) {
const chain = this._flushFollowUp ?? this._flushInFlight
if (!chain) break
await chain.catch(() => {})
}
// Cancel any pending post-import background deduplication FIRST — it is a
// writer (merge-deletes), and no delete pass may start mid- or post-close.

View file

@ -351,3 +351,63 @@ export class PendingFlushDurabilityError extends Error {
this.failedAttempts = failedAttempts
}
}
/**
* @description Thrown by {@link GenerationStore.commitTransaction} when the
* PENDING single-op tier is non-empty i.e. one or more `commitSingleOp()`
* generations are buffered in memory, not yet flushed to
* `committedRanges` via `flushPendingSingleOps()`.
*
* The invariant `reservedGensAsc()` (and everything built on it
* `resolveManyAt`, `resolveAt`, `changedBetween`, the hot-tail window) relies
* on is documented, not enforced by types: pending generations must always be
* numerically greater than every committed one, because the ONLY sanctioned
* callers of `commitTransaction()` `Brainy.transact()` and
* `Brainy.compactHistory()` flush the pending tier FIRST. A caller that
* invokes `commitTransaction()` directly while single-ops are still pending
* breaks that invariant: the new commit lands in `committedRanges` ABOVE
* generations still sitting in `pendingGens`, so the committed-then-pending
* concatenation `reservedGensAsc()` yields is no longer ascending. The
* concrete failure this produces is silent, not a crash: `resolveManyAt`
* walks committed ranges before pending ones, so it can report a NEWER
* generation as the "first after" a pin than an older, still-pending one that
* actually touched the id first a wrong before-image at a point-in-time
* read, without a compensating error to warn a caller anything went wrong.
*
* This error refuses the commit outright, before any staging I/O: nothing is
* written, the generation counter reservation is untouched, and
* `committedRanges`/`pendingGens` are exactly as they were. Call
* `flushPendingSingleOps()` first (or go through `Brainy.transact()`, which
* already does).
*
* @example
* try {
* await generationStore.commitTransaction({ touched, execute })
* } catch (err) {
* if (err instanceof PendingSingleOpsUnflushedError) {
* await generationStore.flushPendingSingleOps()
* await generationStore.commitTransaction({ touched, execute }) // now safe
* }
* }
*/
export class PendingSingleOpsUnflushedError extends Error {
/** How many un-flushed single-op generations were buffered at refusal time. */
public readonly pendingCount: number
/**
* @param pendingCount - `pendingGens.length` at the moment of refusal (always 1).
*/
constructor(pendingCount: number) {
super(
`commitTransaction() refused: ${pendingCount} pending single-op generation(s) ` +
`are still buffered and un-flushed. Flush the pending single-op tier before ` +
`committing a transaction — Brainy.transact() does this automatically; a ` +
`direct commitTransaction() call with pending generations would leave the ` +
`generation order unsorted (committed generations landing above lower, ` +
`still-pending ones) and make point-in-time reads (resolveManyAt/resolveAt) ` +
`return the wrong before-image. Call flushPendingSingleOps() first, then retry.`
)
this.name = 'PendingSingleOpsUnflushedError'
this.pendingCount = pendingCount
}
}

View file

@ -32,7 +32,13 @@
*/
import { prodLog } from '../utils/logger.js'
import { GenerationCompactedError, GenerationConflictError, PendingFlushDurabilityError, StoreInconsistentError } from './errors.js'
import {
GenerationCompactedError,
GenerationConflictError,
PendingFlushDurabilityError,
PendingSingleOpsUnflushedError,
StoreInconsistentError
} from './errors.js'
import type { UnreconciledRecord } from './errors.js'
import { TransactionRollbackError } from '../transaction/errors.js'
import type {
@ -1364,6 +1370,9 @@ export class GenerationStore {
* @param args.execute - Runs the planned operation batch atomically.
* @returns The committed generation and its commit timestamp.
* @throws GenerationConflictError when the CAS expectation fails.
* @throws PendingSingleOpsUnflushedError when the pending single-op tier is
* non-empty call `flushPendingSingleOps()` first (both `Brainy.transact()`
* and `Brainy.compactHistory()` already do).
*/
/**
* The generation fact log, or `null` when the storage layer cannot host one.
@ -1438,6 +1447,13 @@ export class GenerationStore {
execute: () => Promise<void>
}): Promise<{ generation: number; timestamp: number }> {
return this.withMutex(async () => {
// The generation-order guard (see assertPendingSingleOpsFlushed): a
// direct commitTransaction() call while single-ops are still pending
// would commit above them, unsorting reservedGensAsc() and corrupting
// point-in-time reads. Both sanctioned callers (Brainy.transact(),
// Brainy.compactHistory()) already flush first, so this is
// behavior-neutral on every real path.
this.assertPendingSingleOpsFlushed()
// A latched history-durability failure compromises the whole generation
// chain — refuse a transact too (advancing the manifest past stuck,
// un-durable single-op generations would be inconsistent). Same loud
@ -2307,6 +2323,37 @@ export class GenerationStore {
}
}
/**
* @description Throw if the pending single-op tier is non-empty. Called at
* the top of {@link commitTransaction} (the ONLY method that appends a
* fresh commit directly into {@link committedRanges} outside recovery) so
* the ordering invariant {@link reservedGensAsc}'s own doc comment states
* "pending generations are always greater than every committed one" is
* ENFORCED there rather than merely assumed.
*
* That invariant holds today only because both sanctioned callers flush the
* pending tier before committing: `Brainy.transact()` (src/brainy.ts,
* `await this.generationStore.flushPendingSingleOps()` immediately before
* its `commitTransaction()` call) and `Brainy.compactHistory()`
* (src/brainy.ts, the same flush immediately before its `compact()` call
* `compact()` itself only ever RECLAIMS an existing committed prefix, so it
* cannot land a commit out of order and needs no guard of its own). A
* caller that reaches `commitTransaction()` by any other path bypassing
* that flush would commit a new generation into `committedRanges` ABOVE
* generations still sitting in `pendingGens`, breaking `reservedGensAsc`'s
* "committed-then-pending is already sorted" assumption and making
* `resolveManyAt`'s single ascending pass (and `resolveAt`'s consumers)
* return the WRONG before-image for a point-in-time read silently, no
* compensating error. Refusing here, before any staging I/O, keeps the
* store untouched (nothing committed, nothing staged, the generation
* counter reservation unaffected) on every path that already flushes.
*/
private assertPendingSingleOpsFlushed(): void {
if (this.pendingGens.length > 0) {
throw new PendingSingleOpsUnflushedError(this.pendingGens.length)
}
}
/** Schedule a coalesced pending-tier flush (size trigger fires immediately on
* the next microtask; otherwise a {@link PENDING_FLUSH_DELAY_MS} timer). Both
* defer outside the current mutex section so the flush can re-acquire it. A
@ -2390,6 +2437,13 @@ export class GenerationStore {
* committed-then-pending concatenation is already sorted identical to the old
* `[...committedGens, ...pendingGens]`. This is the union historical reads
* resolve over so un-flushed single-ops are visible to pins/`asOf`.
*
* The "flush first" half of that invariant is ENFORCED, not just documented:
* {@link commitTransaction} the only method that lands a fresh commit into
* {@link committedRanges} outside crash recovery refuses via
* {@link assertPendingSingleOpsFlushed} whenever {@link pendingGens} is
* non-empty, so a committed generation can never land above a still-pending
* one and break this ordering.
*/
private *reservedGensAsc(): IterableIterator<number> {
yield* this.committedGensAsc()

View file

@ -231,7 +231,8 @@ export {
GenerationCompactedError,
StoreInconsistentError,
PendingFlushDurabilityError,
CanonicalEnumerationUnavailableError
CanonicalEnumerationUnavailableError,
PendingSingleOpsUnflushedError
} from './db/errors.js'
export type { UnreconciledRecord } from './db/errors.js'
export type {

View file

@ -0,0 +1,405 @@
/**
* @module tests/integration/shutdown-single-owner
* @description ONE SHUTDOWN, ONE OWNER.
*
* MEASURED IN PRODUCTION. A host that owns its own shutdown one SIGTERM
* listener calling `close()` on every pooled store ran head-on into the
* engine's own signal handler, which iterated every live instance, flushed its
* components in parallel, and released its writer lock in a `finally`. Two
* teardowns of the same brain at the same moment. The log shape:
*
* "Shutdown signal received - flushing pending data..." (SIGTERM)
* ...148 seconds of silence...
* "Flushed successfully (1 instance)"
* ...the host's pool close of that same store returns 1s later
*
* 149s for the one store with engine work in flight, against 24s for its six
* idle siblings. The same race in a local reproduction printed
* `Failed to flush one Brainy instance on shutdown: Writer fence lost … the
* lock file is gone` — the handler observing a lock the close it was racing
* had already released.
*
* The contract pinned here:
* (a) A host owner and the engine's hooks both live: EXACTLY ONE close runs
* per brain, no fence is lost, both durability markers are written, the
* process exits 0, and the reopen adopts rather than folding.
* (b) No host owner: the engine's handler closes every instance by the same
* `close()` path markers written, clean exit.
* (c) `close()` is idempotent and re-entrant: concurrent callers share ONE
* execution and all of them settle.
* (d) Flush is single-flight: N kicks during a running flush arm exactly one
* follow-up, and two flush bodies never overlap.
*/
import { describe, it, expect, beforeEach, afterEach } from 'vitest'
import { mkdtempSync, rmSync, existsSync, readFileSync, writeFileSync } from 'node:fs'
import { spawn } from 'node:child_process'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { Brainy } from '../../src/brainy.js'
import { NounType } from '../../src/types/graphTypes.js'
const REPO_ROOT = process.cwd()
const TSX = join(REPO_ROOT, 'node_modules', '.bin', 'tsx')
const BRAINY_SRC = join(REPO_ROOT, 'src', 'brainy.ts')
function makeTempDir(prefix: string): string {
return mkdtempSync(join(tmpdir(), prefix))
}
/** The writer lock's clean-close record — written by `releaseWriterLock()`. */
const closeRecordPath = (dir: string) => join(dir, 'locks', '_writer.close')
/**
* The generation store's clean-shutdown marker the adopt-vs-fold gate.
* (`FileSystemStorage` gzips raw objects, so the file on disk carries `.gz`;
* both spellings are accepted so the pin survives a compression change.)
*/
const cleanShutdownWritten = (dir: string) =>
existsSync(join(dir, '_system', 'clean-shutdown.json.gz')) ||
existsSync(join(dir, '_system', 'clean-shutdown.json'))
/**
* Write a child script and start it under tsx, in its OWN process group so a
* group-wide signal reaches the grandchild that actually holds the writer
* lock. (A file, not `tsx -e`: the eval form compiles to CommonJS, which has
* no top-level await.)
*/
function startChild(scriptDir: string, body: string): ReturnType<typeof spawn> {
const scriptPath = join(scriptDir, 'child-process.mts')
writeFileSync(scriptPath, body)
return spawn(TSX, [scriptPath], {
cwd: REPO_ROOT,
stdio: ['ignore', 'pipe', 'pipe'],
detached: true
})
}
/** Start a child and resolve once it prints READY, collecting all its output. */
function startAndAwaitReady(
scriptDir: string,
body: string
): Promise<{ child: ReturnType<typeof spawn>; output: () => string }> {
const child = startChild(scriptDir, body)
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(`child never became READY:\n${out}`)),
120_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(`child exited ${code} before READY:\n${out}`))
})
})
}
/** Capture console.warn/error/log lines emitted while `fn` runs. */
async function captureConsole<T>(fn: () => Promise<T>): Promise<{ result: T; lines: string[] }> {
const lines: string[] = []
const orig = { log: console.log, warn: console.warn, error: console.error }
const sink = (...args: unknown[]) => { lines.push(args.map((a) => String(a)).join(' ')) }
console.log = sink as typeof console.log
console.warn = sink as typeof console.warn
console.error = sink as typeof console.error
try {
return { result: await fn(), lines }
} finally {
console.log = orig.log
console.warn = orig.warn
console.error = orig.error
}
}
/**
* Reopen a store and assert the open ADOPTED: no crash-recovery fold, no
* stale-lock verdict. This is the whole point of a close having run exactly
* once a fold is measured in tens of seconds on a real store.
*/
async function expectCleanReopen(dir: string): Promise<void> {
const { result, lines } = await captureConsole(async () => {
const next = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dir } })
await next.init()
return next
})
try {
expect(lines.filter((l) => /log-authority recovery|unclean shutdown detected/i.test(l))).toEqual([])
expect(lines.filter((l) => /Overwriting stale writer lock|appears dead/i.test(l))).toEqual([])
} finally {
await result.close()
}
}
/** The child's counts of closes entered and close bodies run, per brain. */
function readResult(
resultPath: string,
out: string
): { entries: Record<string, number>; bodies: Record<string, number>; releases: Record<string, number> } {
if (!existsSync(resultPath)) throw new Error(`child wrote no result file:\n${out}`)
return JSON.parse(readFileSync(resultPath, 'utf-8'))
}
/**
* The child-side instrumentation, shared by (a) and (b): count how many times
* `close()` is ENTERED per brain and how many times its body actually RUNS.
* The counting wrapper is an OWN property, so it shadows the prototype for
* every caller including the engine's own signal handler, which calls
* `instance.close()`.
*
* `report()` writes SYNCHRONOUSLY to a file: it runs on the way out of the
* process (the engine's handler calls `process.exit(0)` when it is the sole
* shutdown owner), and a `console.log` to a pipe is asynchronous and can be
* dropped by that exit.
*/
function childCounters(resultPath: string): string {
return `
const entries = {}
const bodies = {}
const releases = {}
function instrument(name, brain) {
entries[name] = 0
bodies[name] = 0
releases[name] = 0
const enter = brain.close.bind(brain)
brain.close = () => { entries[name]++; return enter() }
const durable = brain.closeDurableSteps.bind(brain)
brain.closeDurableSteps = () => { bodies[name]++; return durable() }
// The writer lock is the ownership witness: the old handler released it
// in its own finally, on top of the owner's close doing the same.
const storage = brain.storage
const release = storage.releaseWriterLock.bind(storage)
storage.releaseWriterLock = () => { releases[name]++; return release() }
}
const report = () => {
__writeFileSync(${JSON.stringify(resultPath)}, JSON.stringify({ entries, bodies, releases }))
}
`
}
describe('shutdown has exactly one owner', () => {
let dirA: string
let dirB: string
let scriptDir: string
let resultPath: string
beforeEach(() => {
dirA = makeTempDir('brainy-shutdown-owner-a-')
dirB = makeTempDir('brainy-shutdown-owner-b-')
scriptDir = makeTempDir('brainy-shutdown-owner-script-')
resultPath = join(scriptDir, 'result.json')
})
afterEach(() => {
for (const d of [dirA, dirB, scriptDir]) {
try { rmSync(d, { recursive: true, force: true }) } catch { /* ignore */ }
}
})
it('(a) a host owner closes both brains and the engine handler steps aside', async () => {
const script = `
import { writeFileSync as __writeFileSync } from 'node:fs'
import { Brainy } from ${JSON.stringify(BRAINY_SRC)}
const a = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: ${JSON.stringify(dirA)} } })
const b = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: ${JSON.stringify(dirB)} } })
await a.init()
await b.init()
await a.add({ data: 'row in brain a', type: 'concept' })
await b.add({ data: 'row in brain b', type: 'concept' })
${childCounters(resultPath)}
instrument('a', a)
instrument('b', b)
// THE HOST'S OWN SHUTDOWN OWNER, registered after the engine's hooks —
// the ordinary shape: the pool was built before the signal wiring.
process.on('SIGTERM', async () => {
await Promise.all([a.close(), b.close()])
// Stay alive a beat so the engine's deferred handler gets its turn and
// has to decide what to do about two already-closed brains.
await new Promise((r) => setTimeout(r, 1500))
report()
process.exit(0)
})
console.log('READY')
setInterval(() => {}, 1000)
`
const { child, output } = await startAndAwaitReady(scriptDir, script)
process.kill(-(child.pid as number), 'SIGTERM')
const code = await new Promise<number | null>((r) => child.on('exit', (c) => r(c)))
// The tsx wrapper's exit event and the grandchild that actually held the
// locks are asynchronous with each other — let its last writes land.
await new Promise<void>((r) => setTimeout(r, 750))
const out = output()
// The process shut down cleanly.
expect(code, `child output:\n${out}`).toBe(0)
// EXACTLY ONE close per brain — entered once, body run once. A second
// entry would mean the engine's handler closed a brain its owner was
// already closing; a second body would mean close() is not single-flight.
const { entries, bodies, releases } = readResult(resultPath, out)
expect(entries).toEqual({ a: 1, b: 1 })
expect(bodies).toEqual({ a: 1, b: 1 })
// ...and the writer lock was given up exactly once per brain. This is the
// assertion that fails on the old handler, which released the lock in its
// own `finally` on top of the owner's close doing the same — two owners.
expect(releases).toEqual({ a: 1, b: 1 })
// The engine's handler ran (it announced the signal) and stepped aside for
// both brains rather than touching them. setImmediate lands in the check
// phase of the same loop turn, so a close that has begun cannot have
// finished — it is still in flight when the handler looks.
expect(out).toContain('Shutdown signal received')
expect(out).toMatch(/2 Brainy instances are already closing/)
// Nothing was taken out from under the owner, and nothing failed.
expect(out).not.toMatch(/Writer fence lost/i)
expect(out).not.toMatch(/Failed to (flush|close) one Brainy instance/i)
// Both durability markers, both brains: the writer lock's clean-close
// record and the generation store's clean-shutdown marker.
for (const dir of [dirA, dirB]) {
expect(existsSync(closeRecordPath(dir)), `clean-close record missing in ${dir}`).toBe(true)
expect(cleanShutdownWritten(dir), `clean-shutdown marker missing in ${dir}`).toBe(true)
}
// And the next open adopts instead of folding.
await expectCleanReopen(dirA)
await expectCleanReopen(dirB)
}, 240_000)
it('(b) with no host owner the engine closes every instance the same way', async () => {
const script = `
import { writeFileSync as __writeFileSync } from 'node:fs'
import { Brainy } from ${JSON.stringify(BRAINY_SRC)}
const a = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: ${JSON.stringify(dirA)} } })
const b = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: ${JSON.stringify(dirB)} } })
await a.init()
await b.init()
await a.add({ data: 'row in brain a', type: 'concept' })
await b.add({ data: 'row in brain b', type: 'concept' })
${childCounters(resultPath)}
instrument('a', a)
instrument('b', b)
process.on('exit', report)
console.log('READY')
setInterval(() => {}, 1000)
`
const { child, output } = await startAndAwaitReady(scriptDir, script)
process.kill(-(child.pid as number), 'SIGTERM')
const code = await new Promise<number | null>((r) => child.on('exit', (c) => r(c)))
// The tsx wrapper's exit event and the grandchild that actually held the
// locks are asynchronous with each other — let its last writes land.
await new Promise<void>((r) => setTimeout(r, 750))
const out = output()
expect(code, `child output:\n${out}`).toBe(0)
// The engine owned this shutdown: one close per brain, through close().
const { entries, bodies, releases } = readResult(resultPath, out)
expect(entries).toEqual({ a: 1, b: 1 })
expect(bodies).toEqual({ a: 1, b: 1 })
expect(releases).toEqual({ a: 1, b: 1 })
expect(out).toContain('Shutdown signal received')
expect(out).toMatch(/Flushed successfully \(2 instances\)/)
expect(out).not.toMatch(/Writer fence lost/i)
expect(out).not.toMatch(/Failed to (flush|close) one Brainy instance/i)
for (const dir of [dirA, dirB]) {
expect(existsSync(closeRecordPath(dir)), `clean-close record missing in ${dir}`).toBe(true)
expect(cleanShutdownWritten(dir), `clean-shutdown marker missing in ${dir}`).toBe(true)
}
await expectCleanReopen(dirA)
await expectCleanReopen(dirB)
}, 240_000)
it('(c) two concurrent close() callers share ONE execution, and both settle', async () => {
const brain = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dirA } })
await brain.init()
await brain.add({ data: 'one row', type: NounType.Concept })
const inner = brain as unknown as { closeDurableSteps: () => Promise<void> }
const durable = inner.closeDurableSteps.bind(inner)
let bodies = 0
inner.closeDurableSteps = () => { bodies++; return durable() }
expect(brain.isClosing).toBe(false)
expect(brain.isClosed).toBe(false)
const first = brain.close()
// The state is observable IMMEDIATELY — a signal handler that yields a
// tick and comes back must not read a stale "not yet".
expect(brain.isClosing).toBe(true)
const second = brain.close()
expect(first === second, 'concurrent callers must share the one promise').toBe(true)
await Promise.all([first, second])
expect(bodies).toBe(1)
expect(brain.isClosed).toBe(true)
// A caller arriving after the close finished gets the same settled answer,
// and nothing runs again.
await brain.close()
expect(bodies).toBe(1)
expect(existsSync(closeRecordPath(dirA))).toBe(true)
expect(cleanShutdownWritten(dirA)).toBe(true)
}, 120_000)
it('(d) N kicks during a running flush arm exactly one follow-up, never a second flush', async () => {
const brain = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dirA } })
await brain.init()
const inner = brain as unknown as {
_flushBodyRuns: number
_flushConcurrencyPeak: number
_flushInFlight: Promise<void> | null
_flushFollowUp: Promise<void> | null
_persistBackgroundFlight: Promise<void> | null
metadataIndex: { flush: () => Promise<void> }
kickBackgroundFlush: (reason: 'threshold' | 'idle') => void
}
// Widen the flush body's window so the kicks land INSIDE it — the
// production shape, where two flushes overlapped 3s apart.
const metaFlush = inner.metadataIndex.flush.bind(inner.metadataIndex)
inner.metadataIndex.flush = async () => {
await new Promise((r) => setTimeout(r, 400))
return metaFlush()
}
await brain.add({ data: 'a write to flush', type: NounType.Concept })
const runsBefore = inner._flushBodyRuns
const leader = brain.flush()
await new Promise((r) => setTimeout(r, 50)) // the leader is inside its body
expect(inner._flushInFlight, 'a flush is running').not.toBeNull()
// The cadence kicks — the door named in the defect — plus direct callers
// (an application flush, the cross-process flush-request watcher).
for (let i = 0; i < 5; i++) inner.kickBackgroundFlush('threshold')
const direct = [brain.flush(), brain.flush(), brain.flush()]
// EXACTLY ONE follow-up is armed, however many callers arrived.
expect(inner._flushFollowUp, 'the eight kicks armed one follow-up').not.toBeNull()
await Promise.all([leader, ...direct, inner._persistBackgroundFlight ?? Promise.resolve()])
// One leader + one follow-up. Not nine, and never two at once.
expect(inner._flushBodyRuns - runsBefore).toBe(2)
expect(inner._flushConcurrencyPeak).toBe(1)
expect(inner._flushInFlight).toBeNull()
expect(inner._flushFollowUp).toBeNull()
inner.metadataIndex.flush = metaFlush
await brain.close()
}, 120_000)
})

View file

@ -0,0 +1,254 @@
/**
* @module tests/unit/db/generationStore-commit-guard
* @description Pins the commit-order guard on
* `GenerationStore.commitTransaction()` (`src/db/generationStore.ts`).
*
* `reservedGensAsc()`'s own doc comment states an invariant it never
* enforced: pending single-op generations are always greater than every
* committed one, because the store's only two sanctioned callers
* `Brainy.transact()` and `Brainy.compactHistory()` flush the pending tier
* before committing. Nothing stopped a caller from invoking
* `commitTransaction()` directly while single-ops were still buffered: the
* fresh commit would land in `committedRanges` ABOVE those lower,
* still-pending generations, so the committed-then-pending concatenation
* `reservedGensAsc()` yields is no longer ascending and `resolveManyAt`
* (which walks committed ranges before pending ones) would silently report a
* WRONG before-image for a point-in-time read. `commitTransaction()` now
* refuses loudly (`PendingSingleOpsUnflushedError`) instead of assuming.
*
* Four pins:
* 1. A direct `commitTransaction()` call while single-ops are pending throws
* and commits NOTHING.
* 2. The same commit succeeds once the pending tier is flushed first.
* 3. `Brainy.transact()` which already flushes first is unaffected
* (mirrors `tests/unit/db/generation-chain.test.ts`'s `seedX()`/`bumpX()`
* transact pin: add, then transact-update, generation advances by one
* each time, the update lands).
* 4. `reservedGensAsc()` stays ascending across a real add+transact+delete
* workload proven by point-in-time reads (`asOf`) staying correct
* throughout, which is exactly what an ordering break would corrupt.
*/
import { describe, it, expect, beforeEach, afterEach } from 'vitest'
import { MemoryStorage } from '../../../src/storage/adapters/memoryStorage.js'
import {
GenerationStore,
GENERATIONS_PREFIX,
MANIFEST_PATH
} from '../../../src/db/generationStore.js'
import { PendingSingleOpsUnflushedError } from '../../../src/db/errors.js'
import { Brainy } from '../../../src/index.js'
import { NounType } from '../../../src/types/graphTypes.js'
import { createTestConfig, generateTestVector } from '../../helpers/test-factory.js'
/** Precomputed embedding so Brainy-level adds skip the (slow) embedding model
* these tests exercise the generation layer, not semantics. */
const VEC = generateTestVector()
// Entity ids must be UUID-shaped (the sharded storage layout derives the
// shard from the UUID hex) — same fixture convention as generationStore.test.ts.
const ID_A = '00000000-0000-4000-8000-0000000000aa'
const ID_B = '00000000-0000-4000-8000-0000000000bb'
/** Stored-metadata fixture in the canonical shape the live write paths use
* (matches generationStore.test.ts's fixture exactly). */
function metadataFixture(version: number): Record<string, unknown> {
return {
noun: NounType.Document,
subtype: 'note',
data: `payload-v${version}`,
version,
createdAt: 1000,
updatedAt: 1000 + version,
_rev: version
}
}
describe('db/GenerationStore — commitTransaction pending-tier guard (store level)', () => {
let storage: MemoryStorage
let store: GenerationStore
beforeEach(async () => {
storage = new MemoryStorage()
await storage.init()
store = new GenerationStore(storage)
await store.open()
})
/** Buffer one single-op generation via commitSingleOp WITHOUT flushing
* the pending tier that must be drained before commitTransaction(). */
async function pendingSingleOp(id: string, version: number): Promise<number> {
const { generation } = await store.commitSingleOp({
touched: { nouns: [id] },
execute: async () => {
await storage.saveNounMetadata(id, metadataFixture(version))
}
})
return generation
}
/** A direct transact commit exactly what a caller bypassing
* Brainy.transact()'s flush-first step would issue. */
function directCommit(id: string, version: number): Promise<{ generation: number; timestamp: number }> {
return store.commitTransaction({
touched: { nouns: [id], verbs: [] },
execute: async () => {
await storage.saveNounMetadata(id, metadataFixture(version))
}
})
}
it('PIN 1: refuses a direct commitTransaction() while single-ops are pending, and commits NOTHING', async () => {
const g1 = await pendingSingleOp(ID_A, 1)
expect(g1).toBe(1)
expect(store.committedGeneration()).toBe(0) // nothing flushed to disk yet
let caught: unknown
try {
await directCommit(ID_B, 1)
expect.unreachable('should have thrown PendingSingleOpsUnflushedError')
} catch (err) {
caught = err
}
expect(caught).toBeInstanceOf(PendingSingleOpsUnflushedError)
expect((caught as PendingSingleOpsUnflushedError).pendingCount).toBe(1)
// Nothing committed: the head + committed ranges are unchanged, and the
// counter never advanced for the refused attempt (the guard fires before
// a generation is even reserved).
expect(store.committedGeneration()).toBe(0)
expect(store.generation()).toBe(1) // still just the pending single-op's gen
expect(await storage.readRawObject(MANIFEST_PATH)).toBeNull()
// The guard fires BEFORE a generation is reserved (`gen = ++this.counter`
// never runs), so the refused attempt's would-be directory (generation 2,
// the next number after the pending single-op's 1) was never created.
expect(await storage.listRawObjects(`${GENERATIONS_PREFIX}/2`)).toEqual([])
// The refused write never touched canonical storage.
expect((await storage.readNounRaw(ID_B)).metadata).toBeNull()
// The pending tier itself is untouched by the refused attempt — flushing
// now still commits the ORIGINAL single-op cleanly.
await store.flushPendingSingleOps()
expect(store.committedGeneration()).toBe(1)
const atG0 = await store.resolveAt('noun', ID_A, 0)
expect(atG0).toEqual({ source: 'absent' }) // the create sentinel before g1's write
})
it('PIN 2: the same commit succeeds once the pending tier is flushed first', async () => {
await pendingSingleOp(ID_A, 1)
await expect(directCommit(ID_B, 1)).rejects.toBeInstanceOf(PendingSingleOpsUnflushedError)
await store.flushPendingSingleOps()
expect(store.committedGeneration()).toBe(1)
const { generation } = await directCommit(ID_B, 1)
expect(generation).toBe(2)
expect(store.committedGeneration()).toBe(2)
expect((await storage.readNounRaw(ID_B)).metadata).toMatchObject({ version: 1 })
})
})
describe('Brainy public API — commitTransaction pending-tier guard is behavior-neutral', () => {
let brain: Brainy
beforeEach(async () => {
brain = new Brainy(createTestConfig())
await brain.init()
})
afterEach(async () => {
await brain.close()
})
it('PIN 3: Brainy.transact() still commits normally over pending single-ops (mirrors generation-chain.test.ts\'s seedX()/bumpX() transact pin)', async () => {
const store = (brain as any).generationStore as GenerationStore
// Relative, not absolute: under the adopt-at-open default the open-time
// baseline backfill takes a generation of its own (see
// bounded-chains.test.ts's identical note), so the first user add is not
// necessarily generation 1.
const baseGen = brain.generation()
const baseCommitted = store.committedGeneration()
const id = await brain.add({
data: 'x',
type: NounType.Document,
subtype: 'note',
metadata: { v: 1 },
vector: VEC
})
// The add is a pending single-op generation — NOT yet flushed.
expect(brain.generation()).toBe(baseGen + 1)
expect(store.committedGeneration()).toBe(baseCommitted)
// Brainy.transact() flushes the pending tier FIRST (src/brainy.ts:
// `await this.generationStore.flushPendingSingleOps()`, immediately
// before its `generationStore.commitTransaction()` call), so the guard
// never fires on this path — same shape as generation-chain.test.ts's
// seedX() (add) → bumpX() (transact update) → generation advances by one.
const db = await brain.transact([{ op: 'update', id, metadata: { v: 2 } }])
await db.release()
expect(brain.generation()).toBe(baseGen + 2)
expect(store.committedGeneration()).toBe(baseGen + 2) // the flushed add + the transact update
const entity = (await brain.get(id)) as any
expect(entity.metadata.v).toBe(2)
})
it('PIN 4: reservedGensAsc() stays ascending across a real add+transact+delete workload — point-in-time reads stay correct', async () => {
const store = (brain as any).generationStore as GenerationStore
const baseGen = brain.generation()
const baseCommitted = store.committedGeneration()
const idX = await brain.add({
data: 'x',
type: NounType.Document,
subtype: 'note',
metadata: { v: 1 },
vector: VEC
})
expect(brain.generation()).toBe(baseGen + 1) // pending (un-flushed)
const idY = await brain.add({
data: 'y',
type: NounType.Document,
subtype: 'note',
metadata: { v: 1 },
vector: VEC
})
// Pin right after BOTH adds — before the transact update — so X reads v1
// and Y still exists at this pin, unlike the live head after the rest of
// the workload runs.
const pinAfterBothAdds = brain.generation()
expect(pinAfterBothAdds).toBe(baseGen + 2) // ALSO pending — two un-flushed single-ops
expect(store.committedGeneration()).toBe(baseCommitted)
// A transact() flushes baseGen+1 and baseGen+2 first, then commits its
// own update as baseGen+3. If committed-vs-pending ordering ever broke,
// this is exactly the step that would land a commit ABOVE still-pending
// generations.
const db = await brain.transact([{ op: 'update', id: idX, metadata: { v: 3 } }])
await db.release()
expect(brain.generation()).toBe(baseGen + 3)
expect(store.committedGeneration()).toBe(baseGen + 3)
// A single-op delete, pending again (un-flushed).
await brain.remove(idY)
expect(brain.generation()).toBe(baseGen + 4)
// A point-in-time read pinned right after the two adds (before the
// transact update) must see X's PRE-update value and Y still present.
// This is precisely what resolveManyAt/resolveAt get WRONG if committed
// and pending generations were ever interleaved out of ascending order.
const past = await brain.asOf(pinAfterBothAdds)
const xAtPin = (await past.get(idX)) as any
expect(xAtPin?.metadata?.v).toBe(1)
const yAtPin = (await past.get(idY)) as any
expect(yAtPin?.metadata?.v).toBe(1) // not yet removed, as of this pin
await past.release()
// Live state reflects every later write, in the right order.
const xNow = (await brain.get(idX)) as any
expect(xNow.metadata.v).toBe(3)
expect(await brain.get(idY)).toBeNull()
})
})