Compare commits
1 commit
a79db434ac
...
30af405b7b
| Author | SHA1 | Date | |
|---|---|---|---|
| 30af405b7b |
2 changed files with 97 additions and 672 deletions
364
src/brainy.ts
364
src/brainy.ts
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
})
|
||||
Reference in a new issue