/** * @module db/generationStore * @description The generational record layer behind Brainy 8.0's MVCC Db API. * * One `GenerationStore` is attached to every writable `Brainy` instance. It * owns: * * - the **monotonic generation counter** (`_system/generation.json`), * reserved per `transact()` commit and bumped per single-operation write * batch via the storage-layer hook; * - the **commit protocol** for `brain.transact()`: capture before-images → * write the generation delta → fsync → execute the operation batch → * atomically rename `_system/manifest.json` (the rename IS the commit * point) → append the tx-log line; * - **crash recovery**: on open, any `_generations/` directory with * `N > manifest.generation` is an uncommitted transaction — its * before-images are restored to the canonical entity paths and the * directory is removed, so a crash anywhere before the manifest rename * leaves the store byte-identical to its pre-transaction state; * - **pin/release refcounting** for live `Db` values, which is what makes * `compactHistory()` safe: a generation record-set is reclaimed only when * no live pin could ever need it; * - **point-in-time resolution**: the state of id X at pinned generation G * is the before-image stored by the *first* committed generation after G * that touched X — or the live storage state when nothing after G touched * it (so current-generation reads stay on the existing fast paths, * untouched). * * Durability and atomicity guarantees, failure modes, and the intellectual * lineage (Datomic database-as-value, LMDB reader pins, LSM immutable * segments) are specified in `docs/ADR-001-generational-mvcc.md`. */ import { prodLog } from '../utils/logger.js' import { GenerationCompactedError, GenerationConflictError } from './errors.js' import type { ChangedIds, CompactHistoryOptions, CompactHistoryResult, GenerationDelta, GenerationManifest, GenerationRecord, GenerationStorage, TxLogEntry } from './types.js' /** Storage-root-relative path of the persisted generation counter. */ export const GENERATION_COUNTER_PATH = '_system/generation.json' /** Storage-root-relative path of the commit manifest. */ export const MANIFEST_PATH = '_system/manifest.json' /** Storage-root-relative prefix of the per-generation record directories. */ export const GENERATIONS_PREFIX = '_generations' /** * @description Phases of the {@link GenerationStore.commitTransaction} commit * protocol at which a test-only fault injector can simulate a process crash. * The phases map 1:1 onto the protocol steps documented on * `commitTransaction`: * * - `'after-staging'` — before-images + the `tx.json` delta are written and * fsynced; the operation batch has NOT executed yet. * - `'after-execute'` — the operation batch has been applied to canonical * storage and indexes; the manifest still points at the prior generation. * - `'before-manifest-rename'` — the batch has executed and the counter is * persisted; the atomic manifest rename — the commit point — has NOT * happened. A crash here is the canonical "fully staged, never committed" * case that recovery must roll back byte-identically. * - `'after-manifest-rename'` — the manifest rename landed (the transaction * IS committed); the tx-log append has NOT happened yet. A crash here must * keep the transaction (the tx-log is advisory metadata, not the source of * commit truth). */ export type CommitFaultPhase = | 'after-staging' | 'after-execute' | 'before-manifest-rename' | 'after-manifest-rename' /** * @description Identifies which ids a transaction touches, split by kind. */ export interface TouchedIds { /** Entity ids the transaction writes or deletes. */ nouns: string[] /** Relationship ids the transaction writes or deletes. */ verbs: string[] } /** * @description Result of {@link GenerationStore.open} — how many uncommitted * generations crash recovery rolled back (a non-zero value tells `Brainy` to * force an index reconciliation, since derived index state may reference the * rolled-back writes). */ export interface GenerationStoreOpenResult { /** Count of uncommitted generation directories rolled back and removed. */ rolledBackGenerations: number } /** * @description The generational MVCC record layer. See the module * documentation for the full protocol; every public method documents its * own contract. */ export class GenerationStore { private readonly storage: GenerationStorage /** Latest reserved/observed generation (≥ {@link committed}). */ private counter = 0 /** Committed-transaction watermark (manifest generation). */ private committed = 0 /** Compaction horizon — record-sets ≤ this are reclaimed. */ private horizonGen = 0 /** * Committed generations whose record dirs exist, stored as a SORTED, DISJOINT, * ascending list of INCLUSIVE `[start, end]` intervals (a run-length set). * * Committed generations come from the monotonic `counter++`, so they are dense * contiguous integers EXCEPT where {@link compact} punches a hole (reclaiming * an oldest contiguous prefix). Storing them as intervals makes the resident * size O(number-of-compaction-gaps) — typically a SINGLE interval * `[firstGen, lastGen]` — instead of one array element per generation. A * billion-insert corpus (every single-op reserves a distinct generation) would * otherwise hold a ~1B-element array resident for the process lifetime; the * interval form collapses that to O(1). All access goes through the helpers * ({@link appendCommittedGen}, {@link committedGensAsc}, {@link committedCount}, * {@link removeCommittedUpTo}, {@link lastCommittedGen}, * {@link largestReservedAtOrBefore}) so the multiset and ordering are identical * to the old flat array at every point. */ private committedRanges: Array<[number, number]> = [] /** Delta cache, keyed by generation (lazily loaded from `tx.json`). `bytes` * is the serialized record-set size, summed by {@link historyBytes} for the * `maxBytes` / adaptive retention caps (0 for generations written before * byte accounting, or for un-flushed pending generations). */ private readonly deltaCache = new Map< number, { nouns: Set; verbs: Set; timestamp: number; bytes: number } >() /** * Per-id inverted history index (Model-B scalability) — `id → ascending * generations that touched it`, one map per kind. {@link resolveAt} binary- * searches the id's OWN chain (O(log chain)) for the first generation after * the pin, instead of linearly scanning the global committed-generation set * ({@link committedRanges}) * (which is O(database-age) — confirmed by the scalability spike: a read of an * unchanged entity at an old pin scaled 11.9x for 10x history depth). This * mirrors cor's `delta_history` chain. Built lazily once (under the mutex), * maintained on commit, invalidated on compaction. */ private readonly nounChains = new Map() private readonly verbChains = new Map() /** Whether {@link nounChains}/{@link verbChains} reflect every committed generation. */ private chainsReady = false /** De-dupes concurrent first-callers of {@link ensureChains}. */ private chainsBuilding: Promise | null = null /** * Cap on resident {@link deltaCache} entries (Model-B RAM bound). The per-id * {@link nounChains}/{@link verbChains} are the hot read structure; the raw * deltas are only needed for range queries ({@link changedBetween}, `diff`, * `since`) and chain (re)builds, and {@link getDelta} transparently re-reads an * evicted delta from storage. Bounding this keeps a long-lived high-write * process's heap O(cap) instead of O(generations) on the disk-backed path. * Not `readonly` so tests can lower it to exercise the eviction/re-read path * without building thousands of generations. */ private deltaCacheMax = 4096 /** * Model-B per-write group-commit — the in-memory PENDING tier. * * Each single-operation write reserves its OWN generation (the locked * cross-project contract: "every write versioned; a distinct generation per * single-op"), captures byte-identical before-images, and buffers them here * instead of paying a synchronous per-write manifest fsync (the 3-5x write * regression). {@link flushPendingSingleOps} persists the whole buffer to * disk in ONE fsync (durability batching — NOT a generation collapse). * * Crucially, these pending generations participate in point-in-time * resolution EXACTLY like committed ones ({@link resolveAt}, the per-id * chains, {@link changedBetween}, {@link hasCommittedAfter} all read the * union via {@link reservedGensAsc}; before-images resolve from this buffer * while pending, from disk once flushed). That is what lets the SYNCHRONOUS * `now()` pin freeze against un-flushed single-ops with no forced flush. * * Durability boundary: a hard crash before a flush loses only the buffered * *history* of un-flushed writes (live data is already in canonical storage); * a crash mid-flush is recovered by drop-without-restore (see * {@link GenerationDelta.groupCommit}). */ private pendingGens: number[] = [] private readonly pendingBuffer = new Map< number, { nouns: Map; verbs: Map; timestamp: number } >() /** Pending timer-coalesced flush handle (cleared on flush/close). */ private pendingFlushTimer: ReturnType | null = null /** * Flush the pending tier once it holds this many generations (size trigger), * complementing the {@link PENDING_FLUSH_DELAY_MS} timer trigger. Not * `readonly` so tests can drive the boundary deterministically. */ private pendingFlushThreshold = 256 /** Coalescing window (ms) before an idle pending tier is flushed. */ private static readonly PENDING_FLUSH_DELAY_MS = 50 /** Live pin refcounts, keyed by pinned generation. */ private readonly pins = new Map() /** Serializes transact commits, compaction, and counter persistence. */ private mutexTail: Promise = Promise.resolve() /** True while a transact batch is executing (suppresses the bump hook). */ private inTransact = false /** Coalescing state for lazy persistence of single-op counter bumps. */ private counterPersistScheduled = false /** * Test-only fault injector for the commit protocol (see * {@link GenerationStore.setCommitFaultInjector}). `undefined` in production. */ private commitFaultInjector?: (phase: CommitFaultPhase) => void private opened = false /** * @param storage - The storage adapter, narrowed to the * {@link GenerationStorage} contract (`BaseStorage` implements it). */ constructor(storage: GenerationStorage) { this.storage = storage } // ========================================================================== // Lifecycle // ========================================================================== /** * @description Load the persisted counter + manifest and run crash * recovery. Must be called once, before indexes are built, because * recovery may rewrite canonical entity files (restoring before-images of * an uncommitted transaction). * * @param options.readOnly - When true (reader-mode brains), recovery is * skipped: readers never write. An un-repaired crashed directory is * repaired by the next *writer* open. * @returns How many uncommitted generations were rolled back. */ async open(options?: { readOnly?: boolean }): Promise { const counterFile = (await this.storage.readRawObject(GENERATION_COUNTER_PATH)) as | { generation?: number } | null const manifest = (await this.storage.readRawObject(MANIFEST_PATH)) as GenerationManifest | null this.committed = manifest?.generation ?? 0 this.horizonGen = manifest?.horizon ?? 0 this.counter = Math.max(counterFile?.generation ?? 0, this.committed) // Discover existing generation record directories. const recordPaths = await this.storage.listRawObjects(GENERATIONS_PREFIX) const seenGens = new Set() for (const p of recordPaths) { const gen = parseGenerationFromPath(p) if (gen !== null) seenGens.add(gen) } let rolledBack = 0 // Coalesce the ascending on-disk committed gens into interval form: each // contiguous run becomes a single `[start, end]` range (appendCommittedGen // extends the last range on a +1 step, opens a new one across a gap). this.committedRanges = [] for (const gen of [...seenGens].sort((a, b) => a - b)) { if (gen <= this.committed) { this.appendCommittedGen(gen) } else if (!options?.readOnly) { // Uncommitted (crashed) generation. A transact generation is rolled // back by RESTORING its before-images (its before-image + execute are // one atomic unit). A single-op group-commit generation is DROPPED // WITHOUT RESTORE — its live write already landed and was acknowledged // before the flush, so restoring would silently revert it. The branch // is decided by the persisted delta's `groupCommit` flag. await this.recoverUncommittedGeneration(gen) rolledBack++ } // Reader mode: leave orphan dirs alone; resolution ignores them because // committedRanges only includes generations ≤ the manifest watermark. this.counter = Math.max(this.counter, gen) } // Chains are rebuilt lazily from the (possibly changed) generation set on // the next historical read — covers reopen and reopen-after-restore. this.invalidateChains() if (rolledBack > 0) { prodLog.warn( `[GenerationStore] Crash recovery rolled back ${rolledBack} uncommitted ` + `transaction(s); store restored to generation ${this.committed}.` ) } this.opened = true // Hook single-op write batches so generation() is always meaningful. // Suppressed while a transact batch executes (the batch is ONE generation). if (!options?.readOnly) { this.storage.setGenerationBumpHook(() => this.noteSingleOpWrite()) } return { rolledBackGenerations: rolledBack } } /** * @description Detach the storage hook and persist the counter. Called * from `brain.close()`. */ async close(): Promise { // Persist any buffered single-op history before detaching, so a clean close // never loses generations the caller already observed. this.clearPendingFlushTimer() await this.flushPendingSingleOps() this.storage.setGenerationBumpHook(undefined) await this.persistCounterNow() } /** * @description TEST-ONLY: install (or clear, with `undefined`) a fault * injector that is invoked at each {@link CommitFaultPhase} of the commit * protocol. A **throwing** injector simulates a process crash at that * phase: the abort cleanup (staging-directory removal, generation-counter * reservation return) is deliberately skipped — exactly as if the process * had died — so crash-consistency tests exercise the REAL recovery path * ({@link GenerationStore.open} rolling back the uncommitted generation on * the next open). Never installed in production code; the injector exists * solely so the durability protocol can be proven, not trusted. * * @param injector - Callback invoked synchronously at each phase, or * `undefined` to detach. */ setCommitFaultInjector(injector: ((phase: CommitFaultPhase) => void) | undefined): void { this.commitFaultInjector = injector } // ========================================================================== // Generation counter // ========================================================================== /** @returns The store's current generation (latest write watermark). */ generation(): number { return this.counter } /** @returns The committed-transaction watermark (manifest generation). */ committedGeneration(): number { return this.committed } /** @returns The compaction horizon (generations ≤ horizon are reclaimed). */ horizon(): number { return this.horizonGen } /** * @description Single-operation write hook (registered with the storage * layer in {@link open}). Bumps the in-memory counter and schedules a * coalesced persist of `_system/generation.json` — durable artifacts * (records, manifests, snapshots) always persist the counter * synchronously at their own commit points, so a crash inside the * coalescing window can never move the counter backwards relative to * anything durable. */ private noteSingleOpWrite(): void { if (this.inTransact) return this.counter++ this.schedulePersistCounter() } /** Schedule one coalesced counter persist for the current write burst. */ private schedulePersistCounter(): void { if (this.counterPersistScheduled) return this.counterPersistScheduled = true setImmediate(() => { this.counterPersistScheduled = false void this.withMutex(() => this.persistCounterUnlocked()).catch((err) => { prodLog.warn(`[GenerationStore] generation counter persist failed: ${(err as Error).message}`) }) }) } /** * @description Persist the generation counter now (atomic tmp+rename). * Called from `flush()`, `close()`, snapshotting, and the commit protocol. */ async persistCounterNow(): Promise { if (!this.opened) return await this.withMutex(() => this.persistCounterUnlocked()) } private async persistCounterUnlocked(): Promise { await this.storage.writeRawObject(GENERATION_COUNTER_PATH, { generation: this.counter, updatedAt: new Date().toISOString() }) } // ========================================================================== // Pins // ========================================================================== /** * @description Pin a generation (refcounted). Compaction never reclaims a * record-set a live pin could need (i.e. any generation directory `N` * with `N > G` for some pinned `G`). * @param gen - The generation to pin. */ pin(gen: number): void { this.pins.set(gen, (this.pins.get(gen) ?? 0) + 1) } /** * @description Release one pin on a generation. Safe to call for a * generation with no live pins (no-op with a warning). * @param gen - The generation to release. */ release(gen: number): void { const count = this.pins.get(gen) if (count === undefined) { prodLog.warn(`[GenerationStore] release(${gen}) called with no live pin at that generation`) return } if (count <= 1) this.pins.delete(gen) else this.pins.set(gen, count - 1) } /** @returns Total number of live pins across all generations. */ activePinCount(): number { let total = 0 for (const count of this.pins.values()) total += count return total } /** @returns The smallest pinned generation, or `Infinity` with no pins. */ private minPinnedGeneration(): number { let min = Infinity for (const gen of this.pins.keys()) min = Math.min(min, gen) return min } // ========================================================================== // Commit protocol // ========================================================================== /** * @description Commit one transaction. The caller (Brainy.transact) plans * the batch up front — resolving every touched id, including generated * entity/relationship ids and delete cascades — and hands this method the * touched-id sets plus an `execute` closure that runs the already-built * operation batch through the TransactionManager. * * Protocol (under the commit mutex): * 1. `ifAtGeneration` CAS check (rejects before anything is staged). * 2. Reserve generation `N` (counter increment). * 3. Capture before-images of every touched id and write them, plus the * `tx.json` delta, under `_generations/N/`; fsync. From this point a * crash anywhere is recoverable to the exact pre-transaction bytes. * 4. Run `execute()`. On failure the TransactionManager has already * rolled back its applied operations; the staging directory is removed * and the reservation is returned (when no concurrent bump consumed a * later number), so a failed transaction leaves the generation * unchanged. * 5. Persist the counter, then atomically rename the manifest to * generation `N` — **the rename is the commit point** — and fsync it. * 6. Append the tx-log line and update in-memory bookkeeping. * * @param args.touched - Every id the batch writes or deletes. * @param args.meta - Transaction metadata for the tx-log. * @param args.ifAtGeneration - Whole-store CAS expectation. * @param args.execute - Runs the planned operation batch atomically. * @returns The committed generation and its commit timestamp. * @throws GenerationConflictError when the CAS expectation fails. */ async commitTransaction(args: { touched: TouchedIds meta?: Record ifAtGeneration?: number execute: () => Promise }): Promise<{ generation: number; timestamp: number }> { return this.withMutex(async () => { if (args.ifAtGeneration !== undefined && args.ifAtGeneration !== this.counter) { throw new GenerationConflictError(args.ifAtGeneration, this.counter) } const nouns = [...new Set(args.touched.nouns)] const verbs = [...new Set(args.touched.verbs)] const gen = ++this.counter const dir = `${GENERATIONS_PREFIX}/${gen}` const timestamp = Date.now() // Test-only crash simulation (see setCommitFaultInjector): a throwing // injector must bypass the abort cleanup below, exactly as a real // process death would, so recovery-on-open is what restores the store. let crashSimulated = false const faultPoint = (phase: CommitFaultPhase): void => { if (!this.commitFaultInjector) return try { this.commitFaultInjector(phase) } catch (err) { crashSimulated = true throw err } } try { // -- 3. Before-images + delta (the durable undo log) ------------------ const stagedPaths: string[] = [] let recordBytes = 0 // serialized record-set size, for retention accounting for (const id of nouns) { const prev = await this.storage.readNounRaw(id) const recordPath = `${dir}/prev/${id}.json` const record: GenerationRecord = { kind: 'noun', metadata: prev.metadata, vector: prev.vector } await this.storage.writeRawObject(recordPath, record) recordBytes += serializedBytes(record) stagedPaths.push(recordPath) } for (const id of verbs) { const prev = await this.storage.readVerbRaw(id) const recordPath = `${dir}/prev/${id}.json` const record: GenerationRecord = { kind: 'verb', metadata: prev.metadata, vector: prev.vector } await this.storage.writeRawObject(recordPath, record) recordBytes += serializedBytes(record) stagedPaths.push(recordPath) } const delta: GenerationDelta = { generation: gen, timestamp, ...(args.meta && { meta: args.meta }), nouns, verbs } delta.bytes = recordBytes + serializedBytes(delta) const deltaPath = `${dir}/tx.json` await this.storage.writeRawObject(deltaPath, delta) stagedPaths.push(deltaPath) await this.storage.syncRawObjects(stagedPaths) faultPoint('after-staging') // -- 4. Execute the planned batch ------------------------------------- this.inTransact = true try { await args.execute() } finally { this.inTransact = false } faultPoint('after-execute') // -- 5. Counter + manifest rename (COMMIT POINT) ---------------------- await this.persistCounterUnlocked() faultPoint('before-manifest-rename') const manifest: GenerationManifest = { version: 1, generation: gen, committedAt: new Date(timestamp).toISOString(), horizon: this.horizonGen } await this.storage.writeRawObject(MANIFEST_PATH, manifest) await this.storage.syncRawObjects([MANIFEST_PATH]) faultPoint('after-manifest-rename') // -- 6. Post-commit bookkeeping --------------------------------------- this.committed = gen this.appendCommittedGen(gen) this.setDelta(gen, { nouns: new Set(nouns), verbs: new Set(verbs), timestamp, bytes: delta.bytes ?? 0 }) this.extendChains(gen, nouns, verbs) const logEntry: TxLogEntry = { generation: gen, timestamp, ...(args.meta && { meta: args.meta }) } await this.storage.appendTxLogLine(JSON.stringify(logEntry)) return { generation: gen, timestamp } } catch (err) { // Simulated crash (test-only fault injector): skip the abort cleanup // entirely — a dead process cannot clean up. Recovery on the next // open() is what must (and does) restore the store. if (crashSimulated) { throw err } // Failed before the manifest rename: nothing is committed. Remove the // staging directory; the TransactionManager already restored any // applied operation byte-identically. try { await this.storage.removeRawPrefix(dir) } catch (cleanupErr) { prodLog.warn( `[GenerationStore] failed to remove staging dir ${dir} after aborted transaction: ` + `${(cleanupErr as Error).message} (recovery will remove it on next open)` ) } // Return the reservation when no concurrent bump consumed a later // number, so a failed transaction leaves generation() unchanged. if (this.counter === gen) this.counter = gen - 1 throw err } }) } // ========================================================================== // Single-operation commit (Model-B per-write group-commit) // ========================================================================== /** * @description Commit ONE single-operation write (`add`/`update`/`remove`/ * `relate`/`updateRelation`/`unrelate` and the per-item calls of `addMany`/ * `updateMany`/`relateMany`) as its own immutable generation — the Model-B * "every write versioned" contract. (`removeMany` commits each delete *chunk* * as a single generation for bulk efficiency — the one non-per-item path.) Structurally a one-operation {@link commitTransaction} with * **deferred durability**: * * - Under the commit mutex it reserves generation `N`, captures byte-identical * before-images of every touched id (so the pin/`asOf` read source is * exact), runs `execute()` (the existing single-op `executeTransaction` * batch) with the storage bump hook suppressed, then BUFFERS the * before-images in the in-memory pending tier and returns — **no per-write * manifest fsync** (that synchronous fsync is the 3-5x write regression this * design avoids). * - The generation is immediately visible to point-in-time reads (chains + * {@link reservedGensAsc}), so a synchronous `now()` pin freezes against it * with no forced flush. * - {@link flushPendingSingleOps} later persists the buffer to disk in one * fsync (async group-commit). A crash before that flush loses only the * buffered *history*; live data is already in canonical storage. A crash * mid-flush is recovered by drop-without-restore * (see {@link GenerationDelta.groupCommit}). * * The per-write generation is read by graph-index ops through the same * `generation()` watermark `transact()` uses — because this method, like * {@link commitTransaction}, holds the mutex across the counter bump AND * `execute()`, so the watermark equals `N` for the duration of the write (a * distinct generation threaded to cor per single-op, the locked contract). * * @param args.touched - The ids this write creates/updates/deletes, by kind. * @param args.execute - Runs the single-op's existing operation batch. * @returns The reserved generation and its timestamp. */ async commitSingleOp(args: { touched: { nouns?: string[]; verbs?: string[] } execute: () => Promise }): Promise<{ generation: number; timestamp: number }> { return this.withMutex(async () => { const nouns = args.touched.nouns ? [...new Set(args.touched.nouns)] : [] const verbs = args.touched.verbs ? [...new Set(args.touched.verbs)] : [] const gen = ++this.counter const timestamp = Date.now() // Before-images BEFORE execute overwrites live storage (mirrors // commitTransaction's staging, but buffered in memory — durability is // deferred to the flush). readNounRaw/readVerbRaw on an absent id return // {metadata:null, vector:null} = the create sentinel. const nounBefore = new Map() for (const id of nouns) { const prev = await this.storage.readNounRaw(id) nounBefore.set(id, { kind: 'noun', metadata: prev.metadata, vector: prev.vector }) } const verbBefore = new Map() for (const id of verbs) { const prev = await this.storage.readVerbRaw(id) verbBefore.set(id, { kind: 'verb', metadata: prev.metadata, vector: prev.vector }) } // Execute the live write. inTransact suppresses the storage bump hook so // the counter is not double-advanced (this generation already reserved // it) — the same suppression transact() uses. this.inTransact = true try { await args.execute() } catch (err) { this.inTransact = false // Live write failed: nothing buffered, nothing on disk. Return the // reservation (when no concurrent bump consumed a later number) so // generation() is unchanged. if (this.counter === gen) this.counter = gen - 1 throw err } this.inTransact = false // Buffer the pending generation + make it instantly visible to reads. this.pendingBuffer.set(gen, { nouns: nounBefore, verbs: verbBefore, timestamp }) this.pendingGens.push(gen) this.extendChains(gen, nouns, verbs) this.schedulePendingFlush() return { generation: gen, timestamp } }) } /** * @description Run a write WITHOUT creating a generation — for init-time / * infrastructure writes (e.g. the VFS root directory) that constitute the * generation-0 baseline of a freshly-materialized brain rather than a user * mutation. The storage bump hook is suppressed (`inTransact`) so the * generation counter does not advance, so a brand-new brain reports * `generation() === 0` and an empty `transactionLog()`, and the first USER * write is generation 1. Mirrors Datomic semantics where the empty database * value is the starting point and bootstrap is not a transaction. * @param execute - The infrastructure write to apply to the baseline. */ async runWithoutGeneration(execute: () => Promise): Promise { this.inTransact = true try { await execute() } finally { this.inTransact = false } } /** * @description Persist the in-memory pending single-op generations to disk in * ONE fsync (async group-commit) and advance the committed watermark. Called * on the size/timer triggers and forced by `flush()`/`close()`/`transact()`/ * `compactHistory()` (so generations stay strictly ordered and durable before * any of those). A no-op when the pending tier is empty. * * Each pending generation is written as a normal `_generations//` record * set (`prev/.json` before-images + `tx.json` delta) but with the delta's * `groupCommit: true` flag set — the drop-without-restore recovery marker. * The single manifest rename to the highest generation commits the whole * batch atomically (committed iff `gen ≤ manifest.generation`). */ async flushPendingSingleOps(): Promise { if (this.pendingGens.length === 0) return return this.withMutex(async () => { if (this.pendingGens.length === 0) return this.clearPendingFlushTimer() const gens = [...this.pendingGens].sort((a, b) => a - b) const stagedPaths: string[] = [] const logEntries: TxLogEntry[] = [] const genBytes = new Map() // for the post-commit setDelta for (const gen of gens) { const buf = this.pendingBuffer.get(gen) if (!buf) continue const dir = `${GENERATIONS_PREFIX}/${gen}` const nounIds: string[] = [] const verbIds: string[] = [] let recordBytes = 0 for (const [id, record] of buf.nouns) { const path = `${dir}/prev/${id}.json` await this.storage.writeRawObject(path, record) recordBytes += serializedBytes(record) stagedPaths.push(path) nounIds.push(id) } for (const [id, record] of buf.verbs) { const path = `${dir}/prev/${id}.json` await this.storage.writeRawObject(path, record) recordBytes += serializedBytes(record) stagedPaths.push(path) verbIds.push(id) } const delta: GenerationDelta = { generation: gen, timestamp: buf.timestamp, groupCommit: true, nouns: nounIds, verbs: verbIds } delta.bytes = recordBytes + serializedBytes(delta) genBytes.set(gen, delta.bytes) const deltaPath = `${dir}/tx.json` await this.storage.writeRawObject(deltaPath, delta) stagedPaths.push(deltaPath) logEntries.push({ generation: gen, timestamp: buf.timestamp }) } // ONE fsync for the whole window — the durability-batching win. await this.storage.syncRawObjects(stagedPaths) // Test-only crash simulation: a throwing injector here leaves the staged // group-commit generation dirs on disk with NO manifest advance — the // exact "crashed mid-flush" state recovery must DROP-WITHOUT-RESTORE // (the pending in-memory state would be lost in a real crash; we abandon // this store and recover from disk on the next open()). if (this.commitFaultInjector) this.commitFaultInjector('before-manifest-rename') // Commit point: persist the counter, then atomically rename the manifest // to the highest pending generation — that commits ALL of them at once. const highest = gens[gens.length - 1] const highestTs = this.pendingBuffer.get(highest)?.timestamp ?? Date.now() await this.persistCounterUnlocked() const manifest: GenerationManifest = { version: 1, generation: highest, committedAt: new Date(highestTs).toISOString(), horizon: this.horizonGen } await this.storage.writeRawObject(MANIFEST_PATH, manifest) await this.storage.syncRawObjects([MANIFEST_PATH]) // Move pending → committed. Chains already carry these gens; seed the // bounded delta cache from the buffer so the next read avoids a re-read, // then drop the in-memory before-images. this.committed = highest for (const gen of gens) { const buf = this.pendingBuffer.get(gen) if (!buf) continue this.appendCommittedGen(gen) this.setDelta(gen, { nouns: new Set(buf.nouns.keys()), verbs: new Set(buf.verbs.keys()), timestamp: buf.timestamp, bytes: genBytes.get(gen) ?? 0 }) this.pendingBuffer.delete(gen) } this.pendingGens = [] for (const entry of logEntries) { await this.storage.appendTxLogLine(JSON.stringify(entry)) } }) } /** 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. */ private schedulePendingFlush(): void { if (this.pendingGens.length >= this.pendingFlushThreshold) { void Promise.resolve() .then(() => this.flushPendingSingleOps()) .catch((err) => prodLog.warn(`[GenerationStore] pending single-op flush failed: ${(err as Error).message}`) ) return } if (this.pendingFlushTimer !== null) return this.pendingFlushTimer = setTimeout(() => { this.pendingFlushTimer = null void this.flushPendingSingleOps().catch((err) => prodLog.warn(`[GenerationStore] pending single-op flush failed: ${(err as Error).message}`) ) }, GenerationStore.PENDING_FLUSH_DELAY_MS) // A background flush must not keep the process alive (Node only). const t = this.pendingFlushTimer as unknown as { unref?: () => void } if (t && typeof t.unref === 'function') t.unref() } /** Cancel any scheduled pending-tier flush timer. */ private clearPendingFlushTimer(): void { if (this.pendingFlushTimer !== null) { clearTimeout(this.pendingFlushTimer) this.pendingFlushTimer = null } } // ========================================================================== // Committed-generation interval set (run-length representation) // ========================================================================== /** * @description Record a newly committed generation. `gen` is ALWAYS greater * than every committed generation (the monotonic `counter++`), so it either * extends the last interval (the normal contiguous append, `gen === * lastEnd + 1`) or opens a fresh single-element interval (the first * generation, or the first commit after a {@link compact} gap). * @param gen - The just-committed generation. */ private appendCommittedGen(gen: number): void { const ranges = this.committedRanges const last = ranges[ranges.length - 1] if (last !== undefined && gen === last[1] + 1) { last[1] = gen } else { ranges.push([gen, gen]) } } /** @returns The newest committed generation, or `undefined` when none exist. */ private lastCommittedGen(): number | undefined { const ranges = this.committedRanges return ranges.length ? ranges[ranges.length - 1][1] : undefined } /** @returns How many committed generations the interval set represents. */ private committedCount(): number { let total = 0 for (const [start, end] of this.committedRanges) total += end - start + 1 return total } /** * @description Yield every committed generation ascending — the exact multiset * (and order) the old flat `committedGens` array held. */ private *committedGensAsc(): IterableIterator { for (const [start, end] of this.committedRanges) { for (let g = start; g <= end; g++) yield g } } /** * @description Yield committed ∪ pending generations, ascending. Pending * generations are always greater than every committed one (reserved after the * last commit; `transact()`/`compact()` flush the pending tier first), so the * 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`. */ private *reservedGensAsc(): IterableIterator { yield* this.committedGensAsc() for (const g of this.pendingGens) yield g } /** * @description Remove every committed generation `<= maxGen`. {@link compact} * always reclaims an oldest CONTIGUOUS prefix of the committed gens, so this * walk over the front of {@link committedRanges} is exact: drop whole ranges * that end at or below `maxGen`, then clip the first surviving range up to * `maxGen + 1` if `maxGen` falls inside it, and stop. * @param maxGen - Inclusive upper bound of the reclaimed prefix. */ private removeCommittedUpTo(maxGen: number): void { const ranges = this.committedRanges let drop = 0 // count of leading ranges fully reclaimed while (drop < ranges.length) { const range = ranges[drop] if (range[1] <= maxGen) { drop++ // whole range is at or below the cut → reclaim it } else { if (range[0] <= maxGen) range[0] = maxGen + 1 // clip the straddling range break // this range survives; nothing newer is below the cut either } } if (drop > 0) ranges.splice(0, drop) } /** * @description The largest reserved (committed ∪ pending) generation * `<= gen` — exactly what `largestAtOrBefore([...committedGens, ...pendingGens], * gen)` returned, in O(log ranges + scanned pending) instead of materializing * the union. Pending generations are the newest reserved gens (all greater than * every committed one), so the largest pending `<= gen`, when one exists, * dominates every committed gen and is the answer; otherwise the answer is the * largest committed gen `<= gen`, found by binary-searching {@link committedRanges} * for the rightmost range whose `start <= gen` and clamping `gen` into it. * @param gen - The upper bound (inclusive). * @returns The largest reserved generation `<= gen`, or `undefined` when every * reserved generation is greater than `gen`. */ private largestReservedAtOrBefore(gen: number): number | undefined { // Pending (ascending, all > every committed): the largest one ≤ gen wins. let bestPending: number | undefined for (const p of this.pendingGens) { if (p <= gen) bestPending = p else break } if (bestPending !== undefined) return bestPending // No pending ≤ gen → the largest committed gen ≤ gen. Binary-search for the // rightmost range whose start ≤ gen; the answer is min(gen, that range.end). const ranges = this.committedRanges let lo = 0 let hi = ranges.length // first range index whose start is > gen while (lo < hi) { const mid = (lo + hi) >>> 1 if (ranges[mid][0] <= gen) lo = mid + 1 else hi = mid } if (lo === 0) return undefined // every committed range starts above gen const [, end] = ranges[lo - 1] return Math.min(gen, end) } // ========================================================================== // Point-in-time resolution // ========================================================================== /** * @description Whether any transaction committed after generation `gen`. * O(1) — this is the fast-path check that lets a `Db` pinned at the * current generation delegate every read to the live brain untouched. * @param gen - The pinned generation. */ hasCommittedAfter(gen: number): boolean { // Pending single-op generations count: a now() pin must report itself // historical the moment any later write (committed OR un-flushed) lands, // else the Model-A "pin doesn't freeze against single-ops" hole reopens. const last = this.pendingGens.length > 0 ? this.pendingGens[this.pendingGens.length - 1] : this.lastCommittedGen() return last !== undefined && last > gen } /** * @description Resolve the state of an entity or relationship at pinned * generation `gen`: the before-image stored by the first committed * generation after `gen` that touched the id — or `{ source: 'current' }` * when nothing after `gen` touched it (the live storage state *is* the * state at `gen`). * * @param kind - `'noun'` for entities, `'verb'` for relationships. * @param id - The id to resolve. * @param gen - The pinned generation. * @returns `{ source: 'current' }`, `{ source: 'absent' }` (did not exist * at `gen`), or `{ source: 'record', metadata, vector }` with the raw * stored objects as of `gen`. */ async resolveAt( kind: 'noun' | 'verb', id: string, gen: number ): Promise< | { source: 'current' } | { source: 'absent' } | { source: 'record'; metadata: any; vector: any | null } > { await this.ensureChains() const chain = (kind === 'noun' ? this.nounChains : this.verbChains).get(id) if (chain === undefined) { // The id was never touched by any committed generation → the live // storage state IS its state at `gen`. return { source: 'current' } } // The first committed generation strictly greater than `gen` that touched // `id` holds the before-image of its state AS OF `gen` — O(log chain). const candidate = firstGenerationAfter(chain, gen) if (candidate === undefined) { // Nothing after `gen` touched `id`; the live state is its state at `gen`. return { source: 'current' } } const record = await this.readBeforeImage(kind, candidate, id) if (record === null) { // The record-set exists (candidate is committed) but the id's // before-image is missing — only possible through external tampering. throw new Error( `Generation record missing: ${GENERATIONS_PREFIX}/${candidate}/prev/${id}.json ` + `(store corrupted or records removed outside compactHistory())` ) } if (record.metadata === null && record.vector === null) { return { source: 'absent' } } return { source: 'record', metadata: record.metadata, vector: record.vector } } /** * @description Read the before-image of `id` stored by generation `gen` from * the right tier: the in-memory pending buffer while the generation is * un-flushed, or `_generations//prev/.json` on disk once committed. * Returns `null` when the record-set has no before-image for `id`. */ private async readBeforeImage( kind: 'noun' | 'verb', gen: number, id: string ): Promise { const pending = this.pendingBuffer.get(gen) if (pending) { return (kind === 'noun' ? pending.nouns : pending.verbs).get(id) ?? null } return (await this.storage.readRawObject( `${GENERATIONS_PREFIX}/${gen}/prev/${id}.json` )) as GenerationRecord | null } /** * @description Ensure the per-id history chains ({@link nounChains} / * {@link verbChains}) reflect every committed generation. Built once, lazily, * under the commit mutex (so no concurrent commit mutates {@link committedRanges} * mid-build); thereafter maintained incrementally by {@link commit} and * invalidated by {@link compact}. Concurrent first-callers share one build. * O(committed generations) the first time; O(1) afterwards. */ private async ensureChains(): Promise { if (this.chainsReady) return if (!this.chainsBuilding) { this.chainsBuilding = this.withMutex(async () => { if (this.chainsReady) return this.nounChains.clear() this.verbChains.clear() // reservedGensAsc (committed ∪ pending) is sorted ascending, so chains // accrue in ascending order and un-flushed single-ops are indexed too. for (const g of this.reservedGensAsc()) { const delta = await this.getDelta(g) for (const nounId of delta.nouns) appendToChain(this.nounChains, nounId, g) for (const verbId of delta.verbs) appendToChain(this.verbChains, verbId, g) } this.chainsReady = true this.chainsBuilding = null }) } await this.chainsBuilding } /** Record that generation `gen` (the newest) touched these ids — keeps chains * current after a commit without a full rebuild. No-op until chains are built * (the eventual build reads `gen` from {@link committedRanges}). */ private extendChains(gen: number, nouns: Iterable, verbs: Iterable): void { if (!this.chainsReady) return for (const nounId of nouns) appendToChain(this.nounChains, nounId, gen) for (const verbId of verbs) appendToChain(this.verbChains, verbId, gen) } /** Drop the chains so the next read rebuilds them from the surviving * generations (called after compaction removes record-sets). */ private invalidateChains(): void { this.chainsReady = false this.chainsBuilding = null this.nounChains.clear() this.verbChains.clear() } /** * @description The ids touched by committed transactions in the interval * `(fromGen, toGen]`. Used by `db.since()` and by historical `find()` to * build its overlay of changed entities. * @param fromGen - Exclusive lower bound. * @param toGen - Inclusive upper bound. */ async changedBetween(fromGen: number, toGen: number): Promise { const nouns = new Set() const verbs = new Set() for (const gen of this.reservedGensAsc()) { if (gen <= fromGen || gen > toGen) continue const delta = await this.getDelta(gen) for (const id of delta.nouns) nouns.add(id) for (const id of delta.verbs) verbs.add(id) } return { nouns: [...nouns].sort(), verbs: [...verbs].sort(), fromGeneration: fromGen, toGeneration: toGen } } /** * @description The committed generations in `(fromGen, toGen]` that touched a * single id, ascending. The per-id companion to {@link changedBetween} — used * by `history()` to walk one entity's change-points without scanning every id. * No all-ids index: it consults the same cached per-generation deltas. * @param kind - `'noun'` for entities, `'verb'` for relationships. * @param id - The id to track. * @param fromGen - Exclusive lower bound. * @param toGen - Inclusive upper bound. */ async generationsTouching( kind: 'noun' | 'verb', id: string, fromGen: number, toGen: number ): Promise { const result: number[] = [] for (const gen of this.reservedGensAsc()) { if (gen <= fromGen || gen > toGen) continue const delta = await this.getDelta(gen) const touched = kind === 'noun' ? delta.nouns : delta.verbs if (touched.has(id)) result.push(gen) } return result } /** * @description Assert that generation `gen` is reachable (not compacted * away) and within the committed range, then return it. Used by `asOf()`. * Reachability boundary: reclaiming record-set `N` removes the * before-images that reads at generations BELOW `N` depend on, so * generations `< horizon` are unreachable while the horizon itself remains * servable from the record-sets above it. * @param gen - The requested generation. * @throws GenerationCompactedError when `gen` is below the horizon. * @throws RangeError when `gen` is negative or beyond the current counter. */ assertReachable(gen: number): number { if (!Number.isInteger(gen) || gen < 0) { throw new RangeError(`asOf(): generation must be a non-negative integer (got ${gen})`) } if (gen > this.counter) { throw new RangeError( `asOf(): generation ${gen} is in the future (store is at generation ${this.counter})` ) } if (gen < this.horizonGen) { throw new GenerationCompactedError(gen, this.horizonGen) } return gen } /** * @description Resolve a wall-clock timestamp to the generation that was * committed at-or-before it, via the tx-log. Under Model-B every write — * `transact()` AND single-op — has its own tx-log entry, so resolution is * per-write (the tx-log includes un-flushed single-op generations too, so * `asOf(Date)` is consistent with the rest of the temporal surface). * @param timestampMs - Milliseconds since epoch. * @returns The resolved generation (`0` when `timestampMs` predates the * first commit) and the matched tx-log entry when one exists. */ async resolveTimestamp( timestampMs: number ): Promise<{ generation: number; entry: TxLogEntry | null }> { let best: TxLogEntry | null = null for (const entry of await this.txLog()) { if (entry.timestamp <= timestampMs) { if (best === null || entry.generation > best.generation) best = entry } } return { generation: best?.generation ?? 0, entry: best } } /** * @description Read all committed tx-log entries, oldest first. Torn * trailing lines from a crashed append are tolerated and skipped, and * entries beyond the committed watermark are excluded — the tx-log is * advisory metadata; the manifest rename is the source of commit truth. * Compaction never rewrites the log, so entries may reference generations * whose record-sets were already reclaimed. * @returns Every committed {@link TxLogEntry}, in commit order. */ async txLog(): Promise { const lines = await this.storage.readTxLogLines() const entries: TxLogEntry[] = [] for (const line of lines) { let entry: TxLogEntry try { entry = JSON.parse(line) as TxLogEntry } catch { continue // tolerate a torn trailing line from a crashed append } if (entry.generation <= this.committed) entries.push(entry) } // Include un-flushed single-op generations so the tx-log is consistent with // the rest of the temporal surface (`since`/`diff`/`asOf`/`changedBetween` // already resolve over the pending tier). Reads are thus flush-agnostic — // whether the async group-commit has run yet never changes results, only // crash durability. Pending single-ops are always newer than every // committed generation, so they extend the ascending list. for (const gen of this.pendingGens) { const buf = this.pendingBuffer.get(gen) if (buf) entries.push({ generation: gen, timestamp: buf.timestamp }) } return entries } /** * @description The commit timestamp of the newest committed generation at * or below `gen`, when its record-set is still resident. Used to populate * `Db.timestamp` for `asOf()` values. * @param gen - The pinned generation. */ async commitTimestampAtOrBefore(gen: number): Promise { // The largest reserved (committed ∪ pending) generation ≤ gen, found in // O(log ranges) over the interval set (not an O(database-age) backward // scan). Pending single-op generations carry timestamps too. const candidate = this.largestReservedAtOrBefore(gen) if (candidate === undefined) return null const delta = await this.getDelta(candidate) return delta.timestamp } private async getDelta( gen: number ): Promise<{ nouns: Set; verbs: Set; timestamp: number; bytes: number }> { const cached = this.deltaCache.get(gen) if (cached) return cached // A pending (un-flushed) generation has no tx.json on disk yet — derive its // delta from the in-memory before-image buffer. Not cached: the bounded // cache is for the disk path; the buffer is the source of truth here. // bytes = 0: pending generations are not on disk, so they do not count // toward the on-disk history byte budget ({@link historyBytes}). const pending = this.pendingBuffer.get(gen) if (pending) { return { nouns: new Set(pending.nouns.keys()), verbs: new Set(pending.verbs.keys()), timestamp: pending.timestamp, bytes: 0 } } const delta = (await this.storage.readRawObject( `${GENERATIONS_PREFIX}/${gen}/tx.json` )) as GenerationDelta | null if (delta === null) { throw new Error( `Generation delta missing: ${GENERATIONS_PREFIX}/${gen}/tx.json ` + `(store corrupted or records removed outside compactHistory())` ) } const entry = { nouns: new Set(delta.nouns), verbs: new Set(delta.verbs), timestamp: delta.timestamp, bytes: delta.bytes ?? 0 } this.setDelta(gen, entry) return entry } /** * @description Insert a delta into {@link deltaCache}, evicting the oldest * entries (FIFO, by insertion order) past {@link deltaCacheMax}. Eviction is * safe: {@link getDelta} re-reads any evicted delta from storage on demand. */ private setDelta( gen: number, entry: { nouns: Set; verbs: Set; timestamp: number; bytes: number } ): void { this.deltaCache.set(gen, entry) while (this.deltaCache.size > this.deltaCacheMax) { const oldest = this.deltaCache.keys().next().value if (oldest === undefined) break this.deltaCache.delete(oldest) } } // ========================================================================== // Compaction // ========================================================================== /** * @description Total serialized bytes of the ON-DISK generational history — * the sum of every committed generation's recorded `bytes`. Backs the * `maxBytes` and adaptive retention caps. Reads each committed generation's * delta (cached; a re-read only for cache-evicted ones) — O(committed * generations), bounded by retention itself and invoked only at compaction * time. Pending (un-flushed) generations are excluded (they are not on disk). * @returns The total on-disk history byte count. */ async historyBytes(): Promise { let total = 0 for (const gen of this.committedGensAsc()) { total += (await this.getDelta(gen)).bytes } return total } /** * @description Reclaim generation record-sets per the retention CAPS. A * generation is removed only when it is at or below **every** live pin * (deleting `N` can only break readers pinned *below* `N`, because resolution * reads before-images from generations strictly greater than the pin) — pins * are ALWAYS exempt. Among the unpinned generations, the OLDEST are reclaimed * while **ANY** supplied cap is exceeded: * * - `maxGenerations` — reclaim while the committed-generation count exceeds it. * - `maxAge` — reclaim generations older than `now − maxAge`. * - `maxBytes` — reclaim while total on-disk history bytes exceed it. * * Because reclamation is oldest-first and every cap is monotone in age, once * the oldest unpinned generation violates no cap, no newer one does either, so * the scan stops. With no options, every unpinned generation is reclaimed. * * @param options - Retention caps (see {@link CompactHistoryOptions}). * @returns Count of removed record-sets and the new horizon. */ async compact(options?: CompactHistoryOptions): Promise { return this.withMutex(async () => { const minPinned = this.minPinnedGeneration() const maxGenerations = options?.maxGenerations const maxAge = options?.maxAge const maxBytes = options?.maxBytes const ageCutoff = maxAge !== undefined ? Date.now() - maxAge : undefined const noCaps = maxGenerations === undefined && maxAge === undefined && maxBytes === undefined // Running totals the caps are evaluated against; updated as we reclaim. let remainingCount = this.committedCount() let remainingBytes = maxBytes !== undefined ? await this.historyBytes() : 0 // Snapshot the committed gens ascending so the reclaim loop iterates safely // while removeCommittedUpTo mutates the interval set afterwards. const removed: number[] = [] for (const gen of [...this.committedGensAsc()]) { // Pins are always exempt: never reclaim a generation a live pin needs. if (gen > minPinned) break // committedGensAsc ascending → nothing newer is eligible either if (!noCaps) { const violatesCount = maxGenerations !== undefined && remainingCount > maxGenerations const violatesBytes = maxBytes !== undefined && remainingBytes > maxBytes let violatesAge = false if (ageCutoff !== undefined) { const delta = await this.getDelta(gen) violatesAge = delta.timestamp < ageCutoff } // Oldest-first: once the oldest unpinned gen trips no cap, none newer do. if (!violatesCount && !violatesBytes && !violatesAge) break } const genBytes = maxBytes !== undefined ? (await this.getDelta(gen)).bytes : 0 await this.storage.removeRawPrefix(`${GENERATIONS_PREFIX}/${gen}`) this.deltaCache.delete(gen) remainingCount-- remainingBytes -= genBytes removed.push(gen) } if (removed.length > 0) { // The reclaim loop walks committedGensAsc() oldest-first and stops at the // first survivor, so `removed` is the oldest CONTIGUOUS prefix of the // committed gens — every committed gen ≤ max(removed) was reclaimed, which // is exactly what removeCommittedUpTo strips. Guard that prefix invariant. const highestRemoved = Math.max(...removed) const countBefore = this.committedCount() this.removeCommittedUpTo(highestRemoved) if (this.committedCount() !== countBefore - removed.length) { throw new Error( `[GenerationStore] compaction invariant violated: reclaimed ${removed.length} ` + `generation(s) but committedCount dropped by ${countBefore - this.committedCount()} ` + `(the reclaimed set is not the oldest contiguous prefix)` ) } // Reclaimed generations leave the per-id chains stale → rebuild on next read. this.invalidateChains() this.horizonGen = Math.max(this.horizonGen, highestRemoved) const manifest: GenerationManifest = { version: 1, generation: this.committed, committedAt: new Date().toISOString(), horizon: this.horizonGen } await this.storage.writeRawObject(MANIFEST_PATH, manifest) await this.storage.syncRawObjects([MANIFEST_PATH]) } return { removedGenerations: removed.length, horizon: this.horizonGen } }) } // ========================================================================== // Restore support // ========================================================================== /** * @description Re-read all persisted state after `brain.restore()` * replaced the store's contents from a snapshot, enforcing counter * monotonicity: the counter never moves below `floorGeneration` (the * pre-restore counter), so generation numbers observed before a restore * are never reissued. * @param floorGeneration - The counter value before the restore. */ async reopenAfterRestore(floorGeneration: number): Promise { await this.withMutex(async () => { this.deltaCache.clear() // A wholesale state replacement invalidates any buffered single-op // history — discard the pending tier (its live writes are gone with the // replaced store). this.clearPendingFlushTimer() this.pendingGens = [] this.pendingBuffer.clear() this.opened = false // open() re-reads counter/manifest and re-registers the bump hook. await this.open() if (this.counter < floorGeneration) { this.counter = floorGeneration await this.persistCounterUnlocked() } }) } // ========================================================================== // Internals // ========================================================================== /** * Recover one uncommitted (crashed) generation, branching on whether it is a * single-op group-commit generation (drop-without-restore) or a `transact()` * generation (restore before-images). The decision reads the generation's * `tx.json` `groupCommit` flag: * * - **`groupCommit: true` → DROP** the directory, never restore. The single-op * write already mutated canonical storage and was acknowledged to the caller * before the flush; restoring its before-images would silently revert an * acknowledged user write (the data-corruption trap). Only the *history* of * the crashed flush window is lost — live data is correct as-is. * - **absent/`false` (a transact generation) → RESTORE** then drop (the * existing atomic-rollback path). * - **`tx.json` unreadable/missing → DROP** (safe for both): a transact * fsyncs `tx.json` BEFORE it executes, so an unreadable delta means the * execute never ran → live storage is unchanged → dropping is a no-op on * live state, while a partial group-commit must be dropped anyway. */ private async recoverUncommittedGeneration(gen: number): Promise { const dir = `${GENERATIONS_PREFIX}/${gen}` const delta = (await this.storage.readRawObject(`${dir}/tx.json`)) as GenerationDelta | null if (delta === null || delta.groupCommit === true) { // Drop-without-restore (group-commit, or an indeterminate partial dir). await this.storage.removeRawPrefix(dir) prodLog.warn( `[GenerationStore] dropped uncommitted ${ delta === null ? 'indeterminate' : 'group-commit' } generation ${gen} (live data left intact; history of that flush window discarded)` ) return } await this.rollBackUncommittedGeneration(gen) } /** * Restore the before-images of an uncommitted `transact()` generation to the * canonical entity paths, then remove its record directory. Restores are * idempotent (writing identical bytes / deleting already-absent files), so * recovery is safe at every crash point of the transact commit protocol. */ private async rollBackUncommittedGeneration(gen: number): Promise { const dir = `${GENERATIONS_PREFIX}/${gen}` const prevPaths = await this.storage.listRawObjects(`${dir}/prev`) for (const recordPath of prevPaths) { const id = recordIdFromPath(recordPath) if (id === null) continue const record = (await this.storage.readRawObject(recordPath)) as GenerationRecord | null if (record === null) continue const image = { metadata: record.metadata, vector: record.vector } if (record.kind === 'verb') { await this.storage.writeVerbRaw(id, image) } else { await this.storage.writeNounRaw(id, image) } } await this.storage.removeRawPrefix(dir) prodLog.warn( `[GenerationStore] rolled back uncommitted generation ${gen} ` + `(${prevPaths.length} record(s) restored)` ) } /** * @description Run a snapshot section serialized against transact commits, * compaction, and counter persistence (the store's single commit mutex), * with the generation counter durably persisted first. Used by * `db.persist()`: the lock guarantees a snapshot can never interleave with * a commit's record-write/manifest-rename sequence, and the up-front * counter persist guarantees the snapshot's `_system/generation.json` * reflects every single-operation bump (whose persistence is otherwise * coalesced and may still be scheduled). * @param fn - The exclusive snapshot section. * @returns `fn`'s result. */ snapshotWith(fn: () => Promise): Promise { return this.withMutex(async () => { await this.persistCounterUnlocked() return fn() }) } /** Run `fn` exclusively, chained behind every prior exclusive section. */ private withMutex(fn: () => Promise): Promise { const run = this.mutexTail.then(fn, fn) // Keep the chain alive regardless of fn's outcome. this.mutexTail = run.catch(() => {}) return run } } /** * @description Approximate serialized byte size of a record or delta, for the * Model-B retention byte caps (`maxBytes` / adaptive). The JSON string length * (chars ≈ bytes) is a deliberately cheap soft estimate — retention is a budget, * not exact accounting — and avoids a storage size API. Returns 0 on any * serialization failure (e.g. a cyclic structure, which records never are). * @param obj - The record or delta about to be persisted. */ function serializedBytes(obj: unknown): number { try { return JSON.stringify(obj)?.length ?? 0 } catch { return 0 } } /** * @description Extract the generation number from a record path of the form * `_generations//...`. Returns `null` for paths that don't match. * @param path - A storage-root-relative object path. */ function parseGenerationFromPath(path: string): number | null { const match = /^_generations[/\\](\d+)[/\\]/.exec(path) if (!match) return null const gen = Number(match[1]) return Number.isSafeInteger(gen) ? gen : null } /** * @description Extract the record id (file basename without `.json`) from a * `prev/` before-image record path. Returns `null` for non-record paths. * @param path - A storage-root-relative object path. */ function recordIdFromPath(path: string): string | null { const match = /[/\\]prev[/\\]([^/\\]+)\.json$/.exec(path) return match ? match[1] : null } /** * @description Append `gen` to an id's ascending history chain in `chains`, * creating the chain on first touch. Callers append in ascending generation * order, so the chain stays sorted without a re-sort. */ function appendToChain(chains: Map, id: string, gen: number): void { const chain = chains.get(id) if (chain === undefined) { chains.set(id, [gen]) } else if (chain[chain.length - 1] !== gen) { chain.push(gen) } } /** * @description Binary-search an ascending generation chain for the FIRST entry * strictly greater than `gen` (the generation whose before-image holds the * value as of `gen`). Returns `undefined` when every entry is ≤ `gen`. O(log n). */ function firstGenerationAfter(chain: number[], gen: number): number | undefined { let lo = 0 let hi = chain.length // upper-bound search over [lo, hi) while (lo < hi) { const mid = (lo + hi) >>> 1 if (chain[mid] > gen) hi = mid else lo = mid + 1 } return lo < chain.length ? chain[lo] : undefined }