Compare commits
4 commits
97df58109b
...
367ca721a5
| Author | SHA1 | Date | |
|---|---|---|---|
| 367ca721a5 | |||
| a79db434ac | |||
| da9519903a | |||
| ec644bde56 |
7 changed files with 1323 additions and 103 deletions
381
src/brainy.ts
381
src/brainy.ts
|
|
@ -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.
|
||||
|
|
@ -20316,9 +20486,22 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
|||
await this._aggregationIndex.flush()
|
||||
}
|
||||
})(),
|
||||
// 8.0 MVCC: detach the generation-bump hook and persist the counter
|
||||
// 8.0 MVCC: detach the generation-bump hook and persist the counter.
|
||||
// READ-ONLY GUARD: a reader's open() never sets the bump hook, never
|
||||
// buffers pending single-ops, and — since generationStore.open() also
|
||||
// leaves the clean-shutdown marker untouched for a reader — never
|
||||
// consumes it either, so there is nothing of a writer's to persist or
|
||||
// release here. Calling close() anyway would still WRITE: it
|
||||
// unconditionally re-stamps `_system/clean-shutdown.json` (and can
|
||||
// advance the fold checkpoint / counter files) at the generation this
|
||||
// session merely observed — a reader vouching for a commit it never
|
||||
// made. The marker is the writer's own evidence about the writer's own
|
||||
// process; a read-only brain must leave `_system/` exactly as it found
|
||||
// it. (Mirrors the same guard already applied to every other Phase-1
|
||||
// step below, and to the signal-path shutdown in
|
||||
// registerShutdownHooks().)
|
||||
(async () => {
|
||||
if (this.generationStore) {
|
||||
if (this.generationStore && !this.isReadOnly) {
|
||||
await this.generationStore.close()
|
||||
}
|
||||
})()
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
@ -799,7 +805,16 @@ export class GenerationStore {
|
|||
if (uncleanOpen) await this.advanceFoldCheckpointUnlocked()
|
||||
// The marker is consumed: any session that can write invalidates it
|
||||
// at first commit (see the commit paths); a clean close re-writes it.
|
||||
await this.clearCleanShutdownMarker()
|
||||
// A READER NEVER CONSUMES IT. The marker is the writer's own evidence
|
||||
// about the writer's own process — clearing it here exists so that
|
||||
// if THIS session goes on to write and then dies before its next
|
||||
// clean close, the marker's absence correctly reads as unclean. A
|
||||
// reader can never write, so it can never leave the store in a state
|
||||
// its own crash would mis-describe; clearing the marker for it would
|
||||
// only cost the store's actual writer a needless whole-log fold on
|
||||
// its next open, for a generation the reader merely observed. Leave
|
||||
// `_system/` exactly as found.
|
||||
if (!options?.readOnly) await this.clearCleanShutdownMarker()
|
||||
}
|
||||
await this.factLog.open(this.committed)
|
||||
} else {
|
||||
|
|
@ -889,7 +904,11 @@ export class GenerationStore {
|
|||
}
|
||||
}
|
||||
|
||||
/** Consume the clean-shutdown marker (every open; a clean close re-writes it). */
|
||||
/**
|
||||
* Consume the clean-shutdown marker (every WRITER open; a clean close
|
||||
* re-writes it). Callers must gate this on `!options.readOnly` — a reader
|
||||
* never consumes the marker, see the call site in {@link open}.
|
||||
*/
|
||||
private async clearCleanShutdownMarker(): Promise<void> {
|
||||
try {
|
||||
await this.storage.deleteRawObject(CLEAN_SHUTDOWN_PATH)
|
||||
|
|
@ -1351,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.
|
||||
|
|
@ -1425,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
|
||||
|
|
@ -2294,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
|
||||
|
|
@ -2377,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()
|
||||
|
|
|
|||
|
|
@ -231,7 +231,8 @@ export {
|
|||
GenerationCompactedError,
|
||||
StoreInconsistentError,
|
||||
PendingFlushDurabilityError,
|
||||
CanonicalEnumerationUnavailableError
|
||||
CanonicalEnumerationUnavailableError,
|
||||
PendingSingleOpsUnflushedError
|
||||
} from './db/errors.js'
|
||||
export type { UnreconciledRecord } from './db/errors.js'
|
||||
export type {
|
||||
|
|
|
|||
250
tests/integration/readonly-close-no-marker.test.ts
Normal file
250
tests/integration/readonly-close-no-marker.test.ts
Normal file
|
|
@ -0,0 +1,250 @@
|
|||
/**
|
||||
* @module tests/integration/readonly-close-no-marker
|
||||
* @description A READ-ONLY BRAIN WRITES NO CLEAN-SHUTDOWN EVIDENCE.
|
||||
*
|
||||
* `_system/clean-shutdown.json` is the WRITER's own word about the writer's
|
||||
* own process: "everything above this line, from THIS session, is durable."
|
||||
* Two call sites treated a reader exactly like a writer:
|
||||
*
|
||||
* 1. `Brainy#closeDurableSteps()` called `generationStore.close()`
|
||||
* unconditionally — a reader's close re-stamped the marker at the
|
||||
* generation the reader merely OBSERVED, never committed.
|
||||
* 2. `GenerationStore#open()` consumed (deleted) the marker on every open,
|
||||
* reader or writer alike, so a reader that never got to a matching
|
||||
* close left the store looking crashed to the next writer.
|
||||
*
|
||||
* Both are fixed by making a read-only brain leave `_system/` exactly as it
|
||||
* found it — at open AND at close. Pinned here:
|
||||
*
|
||||
* 1. `_system/` is byte-for-byte identical (file set + contents) before and
|
||||
* after a reader opens a cleanly-closed store, reads it, and closes.
|
||||
* 2. After the reader's close, the next WRITER open adopts the marker as
|
||||
* clean — no recovery fold narrates.
|
||||
* 3. A reader creates no file under `_system/` merely by opening (before it
|
||||
* ever closes).
|
||||
* 4. A reader that opens and is then abandoned (crash-style, no close) does
|
||||
* not force the next writer to pay a recovery fold — the concrete harm
|
||||
* the fix closes.
|
||||
*/
|
||||
|
||||
import { describe, it, expect, beforeEach, afterEach } from 'vitest'
|
||||
import { mkdtempSync, rmSync, readdirSync, readFileSync, statSync } from 'node:fs'
|
||||
import { createHash } from 'node:crypto'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { Brainy } from '../../src/brainy.js'
|
||||
import { NounType } from '../../src/types/graphTypes.js'
|
||||
import { abandonAsCrashed } from '../helpers/durabilityKillMatrix.js'
|
||||
|
||||
function makeTempDir(): string {
|
||||
return mkdtempSync(join(tmpdir(), 'brainy-readonly-close-'))
|
||||
}
|
||||
|
||||
/** Recursively hash every regular file under `dir`, keyed by its path relative to `dir`. */
|
||||
function snapshotDir(dir: string): Map<string, string> {
|
||||
const out = new Map<string, string>()
|
||||
const walk = (rel: string): void => {
|
||||
const abs = rel ? join(dir, rel) : dir
|
||||
let entries: string[]
|
||||
try {
|
||||
entries = readdirSync(abs)
|
||||
} catch {
|
||||
return
|
||||
}
|
||||
for (const name of entries) {
|
||||
const childRel = rel ? join(rel, name) : name
|
||||
const childAbs = join(dir, childRel)
|
||||
const st = statSync(childAbs)
|
||||
if (st.isDirectory()) {
|
||||
walk(childRel)
|
||||
} else if (st.isFile()) {
|
||||
const hash = createHash('sha256').update(readFileSync(childAbs)).digest('hex')
|
||||
out.set(childRel, hash)
|
||||
}
|
||||
}
|
||||
}
|
||||
walk('')
|
||||
return out
|
||||
}
|
||||
|
||||
/** Capture console.warn lines (the narration channel — see `prodLog.narrate`) while `fn` runs. */
|
||||
async function captureWarn<T>(fn: () => Promise<T>): Promise<{ result: T; lines: string[] }> {
|
||||
const lines: string[] = []
|
||||
const orig = console.warn
|
||||
console.warn = ((...args: unknown[]) => {
|
||||
lines.push(args.map((a) => String(a)).join(' '))
|
||||
}) as typeof console.warn
|
||||
try {
|
||||
return { result: await fn(), lines }
|
||||
} finally {
|
||||
console.warn = orig
|
||||
}
|
||||
}
|
||||
|
||||
describe('a read-only brain writes no clean-shutdown evidence', () => {
|
||||
let dir: string
|
||||
let brain: Brainy | null = null
|
||||
|
||||
beforeEach(() => {
|
||||
dir = makeTempDir()
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
if (brain) {
|
||||
try {
|
||||
await brain.close()
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
brain = null
|
||||
}
|
||||
try {
|
||||
rmSync(dir, { recursive: true, force: true })
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
})
|
||||
|
||||
const systemDir = () => join(dir, '_system')
|
||||
/**
|
||||
* The marker file's actual on-disk name — `clean-shutdown.json` or, under
|
||||
* FileSystemStorage's default gzip compression, `clean-shutdown.json.gz`.
|
||||
* Returns null when absent.
|
||||
*/
|
||||
const findMarkerPath = (): string | null => {
|
||||
let entries: string[]
|
||||
try {
|
||||
entries = readdirSync(systemDir())
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
const name = entries.find((n) => n.startsWith('clean-shutdown.json'))
|
||||
return name ? join(systemDir(), name) : null
|
||||
}
|
||||
|
||||
it('leaves `_system/`\'s file set and the clean-shutdown marker\'s bytes identical across a reader open → read → close', async () => {
|
||||
// A writer opens, writes, and closes cleanly — the marker lands at
|
||||
// whatever generation the writer actually committed.
|
||||
const writer = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dir } })
|
||||
await writer.init()
|
||||
await writer.add({ data: 'seed entity', type: NounType.Concept })
|
||||
await writer.add({ data: 'second entity', type: NounType.Concept })
|
||||
await writer.flush()
|
||||
await writer.close()
|
||||
|
||||
const markerBeforePath = findMarkerPath()
|
||||
expect(markerBeforePath, 'the writer left a clean-shutdown marker').not.toBeNull()
|
||||
const before = snapshotDir(systemDir())
|
||||
expect(before.size).toBeGreaterThan(0)
|
||||
const markerBeforeHash = before.get(
|
||||
(markerBeforePath as string).slice(systemDir().length + 1)
|
||||
)
|
||||
expect(markerBeforeHash).toBeTruthy()
|
||||
|
||||
// A reader opens the same store, reads, and closes.
|
||||
brain = await Brainy.openReadOnly({ storage: { type: 'filesystem', path: dir } })
|
||||
expect(brain.isReadOnly).toBe(true)
|
||||
await brain.stats()
|
||||
await brain.close()
|
||||
brain = null
|
||||
|
||||
// The FILE SET under `_system/` is unchanged — a reader creates and
|
||||
// removes nothing. (Other files under `_system/` — e.g. the metadata
|
||||
// field registry, which stamps its own `lastUpdated` on every persist —
|
||||
// are a pre-existing, separate concern outside this fix's scope: this
|
||||
// pin is specifically about the generation store's clean-shutdown
|
||||
// evidence, not about every subsystem's close() being a true no-op for
|
||||
// a reader.)
|
||||
const after = snapshotDir(systemDir())
|
||||
expect([...after.keys()].sort()).toEqual([...before.keys()].sort())
|
||||
|
||||
// The MARKER's bytes are byte-for-byte identical — the reader neither
|
||||
// consumed it at open nor re-stamped it at close.
|
||||
const markerAfterPath = findMarkerPath()
|
||||
expect(markerAfterPath, 'the marker must still exist, under the same name').toBe(markerBeforePath)
|
||||
const markerAfterHash = after.get((markerAfterPath as string).slice(systemDir().length + 1))
|
||||
expect(markerAfterHash).toBe(markerBeforeHash)
|
||||
}, 120_000)
|
||||
|
||||
it('creates no file under `_system/` merely by opening read-only', async () => {
|
||||
const writer = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dir } })
|
||||
await writer.init()
|
||||
await writer.add({ data: 'seed entity', type: NounType.Concept })
|
||||
await writer.flush()
|
||||
await writer.close()
|
||||
|
||||
const baselineNames = [...snapshotDir(systemDir()).keys()].sort()
|
||||
expect(baselineNames.length).toBeGreaterThan(0)
|
||||
|
||||
// Open the reader and inspect `_system/` BEFORE it ever closes — open()
|
||||
// alone must create nothing.
|
||||
brain = await Brainy.openReadOnly({ storage: { type: 'filesystem', path: dir } })
|
||||
const whileOpenNames = [...snapshotDir(systemDir()).keys()].sort()
|
||||
expect(whileOpenNames).toEqual(baselineNames)
|
||||
|
||||
await brain.close()
|
||||
brain = null
|
||||
}, 120_000)
|
||||
|
||||
it('a writer reopening after the reader closes adopts the marker — no recovery fold', async () => {
|
||||
const writer1 = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dir } })
|
||||
await writer1.init()
|
||||
await writer1.add({ data: 'seed entity', type: NounType.Concept })
|
||||
await writer1.flush()
|
||||
await writer1.close()
|
||||
|
||||
// A reader opens and closes in between — must not disturb the marker.
|
||||
const reader = await Brainy.openReadOnly({ storage: { type: 'filesystem', path: dir } })
|
||||
await reader.stats()
|
||||
await reader.close()
|
||||
|
||||
// The next writer open must be a clean, no-fold open: no
|
||||
// "log-authority recovery" / "WHOLE-LOG fold" narration line.
|
||||
const { result: writer2, lines } = await captureWarn(async () => {
|
||||
const w = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dir } })
|
||||
await w.init()
|
||||
return w
|
||||
})
|
||||
brain = writer2
|
||||
|
||||
const foldLines = lines.filter((l) => /log-authority recovery|WHOLE-LOG fold|recovery fold/i.test(l))
|
||||
expect(foldLines, `unexpected recovery narration:\n${foldLines.join('\n')}`).toEqual([])
|
||||
|
||||
// And the store is exactly what the first writer left — the seed row is
|
||||
// still there, nothing was rolled back or re-derived.
|
||||
const found = await writer2.find({ where: {} } as any)
|
||||
expect(found.length).toBeGreaterThanOrEqual(1)
|
||||
}, 120_000)
|
||||
|
||||
it('a reader that opens and is then abandoned (never closes) does not force the next writer to fold', async () => {
|
||||
// This is the concrete harm the fix closes: pre-fix, a reader's open()
|
||||
// unconditionally DELETED the marker (consuming it as if it were the
|
||||
// writer). A reader that opened and then died — no close, exactly like
|
||||
// a killed process — left the marker gone, so the actual writer's next
|
||||
// open read the store as crashed and paid a full recovery fold for a
|
||||
// "crash" that was really just a reader that came and went.
|
||||
const writer1 = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dir } })
|
||||
await writer1.init()
|
||||
await writer1.add({ data: 'seed entity', type: NounType.Concept })
|
||||
await writer1.flush()
|
||||
await writer1.close()
|
||||
|
||||
const reader = await Brainy.openReadOnly({ storage: { type: 'filesystem', path: dir } })
|
||||
await reader.stats()
|
||||
// NEVER calls reader.close() — abandon it exactly like a killed process.
|
||||
await abandonAsCrashed(reader)
|
||||
|
||||
const { result: writer2, lines } = await captureWarn(async () => {
|
||||
const w = new Brainy({ requireSubtype: false, storage: { type: 'filesystem', path: dir } })
|
||||
await w.init()
|
||||
return w
|
||||
})
|
||||
brain = writer2
|
||||
|
||||
const foldLines = lines.filter((l) => /log-authority recovery|WHOLE-LOG fold|recovery fold/i.test(l))
|
||||
expect(
|
||||
foldLines,
|
||||
`an abandoned READER forced a recovery fold on the next writer open:\n${foldLines.join('\n')}`
|
||||
).toEqual([])
|
||||
}, 120_000)
|
||||
})
|
||||
405
tests/integration/shutdown-single-owner.test.ts
Normal file
405
tests/integration/shutdown-single-owner.test.ts
Normal 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)
|
||||
})
|
||||
254
tests/unit/db/generationStore-commit-guard.test.ts
Normal file
254
tests/unit/db/generationStore-commit-guard.test.ts
Normal 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()
|
||||
})
|
||||
})
|
||||
Reference in a new issue