Compare commits

..

No commits in common. "da9519903a3de52b2e0aeeb6e33a1257c6f5b749" and "65493ba2de09e291972b0539bf49c221e3f18ce7" have entirely different histories.

2 changed files with 97 additions and 672 deletions

View file

@ -767,34 +767,6 @@ 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
@ -917,24 +889,6 @@ 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
@ -2122,88 +2076,105 @@ export class Brainy<T = any> implements BrainyInterface<T> {
*/
private registerShutdownHooks(): void {
/**
* The signal-path shutdown. ONE OWNER PER BRAIN, AND THE PATH IS `close()`.
* The signal-path shutdown. THREE LAWS, each written by a production
* shutdown that looked clean and wasn't:
*
* 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.
* 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.
*/
const closeOnShutdown = async () => {
const flushOnShutdown = async () => {
console.log('Shutdown signal received - flushing pending data...')
// 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 flushedCount = 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 {
// Law 1: this try/catch is the isolation — the loop continues.
await instance.close()
closedCount++
// 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++
} catch (error) {
failedCount++
console.error('Failed to close one Brainy instance on shutdown:', error)
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)
}
}
}
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 (flushedCount > 0) {
console.log(`Flushed successfully (${flushedCount} instance${flushedCount > 1 ? 's' : ''})`)
}
if (failedCount > 0) {
console.error(
@ -2230,29 +2201,19 @@ 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 = (ownersWhenSignalled: number): void => {
if (ownersWhenSignalled <= 1) {
const exitIfSoleShutdownOwner = (signal: 'SIGTERM' | 'SIGINT'): void => {
if (process.listenerCount(signal) <= 1) {
process.exit(0)
}
}
Brainy.sigtermListener = async () => {
const owners = process.listenerCount('SIGTERM')
await closeOnShutdown()
exitIfSoleShutdownOwner(owners)
await flushOnShutdown()
exitIfSoleShutdownOwner('SIGTERM')
}
Brainy.sigintListener = async () => {
const owners = process.listenerCount('SIGINT')
await closeOnShutdown()
exitIfSoleShutdownOwner(owners)
await flushOnShutdown()
exitIfSoleShutdownOwner('SIGINT')
}
Brainy.beforeExitListener = async () => {
// Self-deregister FIRST: Node re-emits 'beforeExit' after every event-
@ -2264,7 +2225,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
process.off('beforeExit', Brainy.beforeExitListener)
Brainy.beforeExitListener = undefined
}
await closeOnShutdown()
await flushOnShutdown()
}
process.on('SIGTERM', Brainy.sigtermListener)
process.on('SIGINT', Brainy.sigintListener)
@ -2337,33 +2298,6 @@ 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
*
@ -3337,18 +3271,9 @@ 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()
@ -12978,58 +12903,7 @@ export class Brainy<T = any> implements BrainyInterface<T> {
* process.exit(0)
* })
*/
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> {
async flush(): Promise<void> {
await this.ensureInitialized()
// Read-only instances have no buffered writes to flush. close() may call
@ -20276,42 +20150,11 @@ 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.
*/
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> {
async close(): Promise<void> {
if (this._pendingEmbedIds.size === 0) await this.writeEmbedLowWater()
let closeFailure: unknown = null
try {
@ -20400,19 +20243,6 @@ 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

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