brainy/src/db/generationStore.ts
David Snelling 1201e25543 feat: two-tier history reads + the repacker + generationDigest — D1+D3 wired end-to-end
The packed tier goes live inside GenerationStore:

- Two-tier reads: getDelta / readBeforeImage / readGenerationRecords
  fall through live-tier → sealed segments (live-tier-wins: a crash
  mid-fold leaves a duplicate representation, never a gap). Cold-open
  seeds committedRanges from the segment manifest via interval merge —
  packed generations resolve without their directories existing.
- repackHistory({timeBudgetMs, batchGenerations}): folds cold
  generations (older than the newest 1024) oldest-first into sealed
  segments, deleting per-generation directories only after segment +
  manifest are durable. Public API + automatic time-bounded pass at
  close() (before compaction, so reclaim can drop whole segments);
  re-representation only — the sole history transform under the
  archival profile. Early stop = consistent prefix, next pass resumes.
- compact(): packed generations reclaim logically in the loop and
  physically at whole-segment boundaries via dropSegmentsBelow (the
  frozen partial-segments-wait rule).
- generationDigest(g) (D8): deterministic content digest through g —
  sealed-segment checksum chain + live-tier delta hashes; O(segments +
  live window); RangeError out of range, GenerationCompactedError
  below the horizon (a gate can never silently pin reclaimed history).

Four end-to-end pins: asOf answers byte-identical across fold + cold
reopen with folded dirs physically gone; repack+reclaim composition;
digest reopen-stability/divergence/loud-horizon; budget no-op+resume.
2026-07-19 16:26:10 -07:00

2741 lines
124 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

/**
* @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/<N>` 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, PendingFlushDurabilityError, StoreInconsistentError } from './errors.js'
import type { UnreconciledRecord } from './errors.js'
import { TransactionRollbackError } from '../transaction/errors.js'
import type {
ChangedIds,
CompactHistoryOptions,
CompactHistoryResult,
GenerationDelta,
GenerationManifest,
GenerationRecord,
GenerationStorage,
TxLogEntry
} from './types.js'
import { FactLog, storageSupportsFactLog, type CommitFact, type FactOp } from './factLog.js'
import { GenerationSegmentStore, type FoldGeneration } from './generationSegments.js'
import { crc32c } from '../utils/crc32c.js'
/**
* The byte-identical before-images of every id a commit touches, read UNDER
* the commit mutex immediately before the write applies. Passed to a commit's
* optional `precommit` hook so callers can enforce compare-and-swap
* preconditions (e.g. the per-entity `ifRev` check) atomically with the apply
* — the same conditional-commit idea as `ifAtGeneration`, generalized to
* arbitrary per-record predicates. An absent id maps to the create sentinel
* (`metadata: null, vector: null`).
*/
export interface CommitBeforeImages {
nouns: Map<string, GenerationRecord>
verbs: Map<string, GenerationRecord>
}
/** 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
/**
* The generation FACT LOG (dual-write transition) — an append-only,
* CRC-framed record of every committed generation as an AFTER-IMAGE fact.
* `null` when the storage layer lacks the binary raw-byte primitives.
* Appends ride the same commit protocol: a fact-append failure FAILS the
* write (loud — a silent fact gap would make the log a lie that a later
* replay discovers), and open() reconciles the log to committed truth.
*/
private factLog: FactLog | null = null
/** 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<string>; verbs: Set<string>; 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.
*
* BOUNDED-RAM design (Approach C — hot-tail window + cold LRU). The old
* unbounded `Map<id, number[]>` held ONE resident chain per id ever touched
* across retained history → O(distinct-ids-ever-touched) RAM, which defeats
* billion-scale time travel (the chains, not the data, become the heap floor).
* These three structures replace it with RAM bounded by the knobs W/L,
* INDEPENDENT of the corpus size:
*
* 1. {@link recentNounChains}/{@link recentVerbChains} — FULL per-id chains,
* but only over the HOT TAIL window `[windowLo, counter]` (the last
* {@link recentWindowGenerations} generations). Pins at-or-above
* `windowLo` (the overwhelming common case — recent history) resolve here
* with the same O(log chain) binary search as before.
* 2. {@link coldNounChains}/{@link coldVerbChains} — a bounded LRU of FULL
* per-id chains for DEEP pins (`gen < windowLo`). Reconstructed on demand
* from the persisted deltas ({@link reconstructColdChain}), cached so a
* re-read is O(log chain), evicted oldest-first past
* {@link coldChainLruMax}. Empty chains are cached too (bounded negative
* cache for untouched-in-cold-region ids).
* 3. {@link windowNounDeltas}/{@link windowVerbDeltas} — the window's INVERSE
* index (`gen → ids touched at that gen`), fed from the touched sets
* already handed to {@link extendChains}. It makes sliding `windowLo`
* forward an O(d̄) front-trim of exactly the ids leaving the window, with
* ZERO disk I/O — the write path never reads `tx.json` on a slide.
*
* INVARIANT — purely derived, NEVER persisted. All three structures (and the
* cold LRU) are reconstructable from the on-disk deltas at any time; eviction
* (cold LRU) and window slide-out drop entries that are then rebuilt on the
* next read that needs them. Correctness across eviction/reconstruction is
* guaranteed by pin-gating: a live historical `Db` pins at `g_old`, and the
* before-image a read needs lives at the FIRST generation after `g_old`
* (`> g_old ≥ minPinned`), which compaction can never reclaim while the pin is
* held — so a reconstructed chain yields the identical answer as the original.
*/
private readonly recentNounChains = new Map<string, number[]>()
private readonly recentVerbChains = new Map<string, number[]>()
/** Window inverse index (`gen → ids touched`), one map per kind — the O(d̄),
* I/O-free slide-trim driver. Holds exactly the gens in `[windowLo, counter]`. */
private readonly windowNounDeltas = new Map<number, string[]>()
private readonly windowVerbDeltas = new Map<number, string[]>()
/** Bounded LRU of reconstructed deep-pin chains (`gen < windowLo`). JS Map
* insertion order is the LRU order: a hit re-inserts (delete+set) to mark
* MRU; eviction removes `keys().next().value` (the oldest) past the cap. */
private readonly coldNounChains = new Map<string, number[]>()
private readonly coldVerbChains = new Map<string, number[]>()
/** Inclusive lower bound of the resident hot-tail window (`0` until built). */
private windowLo = 0
/** Whether the hot-tail window + its inverse index are built and current. */
private windowReady = false
/** De-dupes concurrent first-callers of {@link ensureWindow}. */
private windowBuilding: Promise<void> | null = null
/**
* Hot-tail window width W — how many of the newest generations keep a resident
* per-id chain ({@link recentNounChains}). Bounds recent-chain + window-delta
* RAM to O(W·d̄), independent of corpus size. Deep pins below the window fall to
* the cold LRU. Not `readonly` so tests can shrink it to exercise the
* cold/slide paths without building thousands of generations.
*/
private recentWindowGenerations = 4096
/**
* Cold-chain LRU capacity L — how many reconstructed deep-pin chains stay
* resident. Bounds cold-chain RAM to O(L·avg-chain-length), independent of how
* many distinct ids are deep-read. Not `readonly` so tests can shrink it to
* exercise eviction + reconstruction.
*/
private coldChainLruMax = 4096
/**
* Cap on resident {@link deltaCache} entries (Model-B RAM bound). The per-id
* {@link recentNounChains}/{@link recentVerbChains} window is 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
/**
* Running total of on-disk history bytes across committed generations —
* `null` until {@link historyBytes} pays its one seeding walk. Maintained
* incrementally at commit/reclaim so the adaptive retention check on every
* flush() is O(1), never a tail re-walk. Never updated by cache re-reads
* ({@link setDelta} inserts are cache population, not new history).
*/
private historyBytesTotal: number | null = null
/**
* The packed tier (D1+D3): sealed segments holding folded cold
* generations. Null until {@link open} wires it (and on storage adapters
* without raw-byte primitives — the live tier then carries everything,
* exactly as before the packed tier existed).
*/
private segments: GenerationSegmentStore | null = null
/**
* Live-tier window: generations newer than `committed - REPACK_LIVE_WINDOW`
* are never folded — the hot tail stays in the per-generation layout the
* write path owns. Matches the resident chain window's scale.
*/
static readonly REPACK_LIVE_WINDOW = 1024
/**
* 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<string, GenerationRecord>; verbs: Map<string, GenerationRecord>; timestamp: number }
>()
/** Pending timer-coalesced flush handle (cleared on flush/close). */
private pendingFlushTimer: ReturnType<typeof setTimeout> | 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
/**
* Latched pending-flush durability failure. Set once background flushes have
* failed {@link PENDING_FLUSH_FAILURE_THRESHOLD} consecutive times (a real,
* persistent storage fault — not a blip); while set, writes are REFUSED with
* a {@link PendingFlushDurabilityError} rather than silently accumulate
* un-durable history in memory. Cleared the moment a flush finally succeeds
* (the store self-heals). `null` = history-durability path is healthy.
*/
private pendingFlushError: Error | null = null
/** Consecutive failed background pending-flush attempts (resets on success). */
private pendingFlushFailures = 0
/**
* Consecutive failed flush attempts before the durability latch trips and
* writes start refusing. A small tolerance so a single transient fault does
* not trip the alarm, while a genuinely stuck disk stops the bleed fast.
*/
private static readonly PENDING_FLUSH_FAILURE_THRESHOLD = 3
/** Upper bound (ms) on the exponential retry backoff between failed flushes. */
private static readonly PENDING_FLUSH_RETRY_CAP_MS = 30_000
/** Live pin refcounts, keyed by pinned generation. */
private readonly pins = new Map<number, number>()
/** Serializes transact commits, compaction, and counter persistence. */
private mutexTail: Promise<unknown> = 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<GenerationStoreOpenResult> {
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<number>()
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)
}
// The history window is 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
// Generation FACT LOG (dual-write transition): when the storage layer
// exposes the binary raw-byte primitives, open the after-image fact log
// and reconcile it to committed truth — facts are appended BEFORE the
// commit point, so a crash can only leave the log AHEAD; open truncates
// any fact beyond `committed`. Storage without the primitives simply
// hosts no fact log (readers fall back to canonical enumeration).
if (storageSupportsFactLog(this.storage)) {
this.factLog = new FactLog(this.storage)
await this.factLog.open(this.committed)
} else {
this.factLog = null
}
// PACKED TIER (D1+D3): same capability gate as the fact log. Opening
// reads ONE manifest — never a listing of the packed backlog — and seeds
// committedRanges with the sealed ranges so packed generations resolve
// exactly like live ones.
if (storageSupportsFactLog(this.storage)) {
this.segments = new GenerationSegmentStore(this.storage)
await this.segments.open()
const packedRanges = this.segments
.segments()
.map((s): [number, number] => [s.firstGeneration, Math.min(s.lastGeneration, this.committed)])
.filter(([lo, hi]) => lo <= hi)
if (packedRanges.length > 0) {
// Merge packed (older) + live (newer) interval sets — both ascending;
// coalesce adjacency so range arithmetic stays interval-exact.
const merged: Array<[number, number]> = []
for (const r of [...packedRanges, ...this.committedRanges].sort((a, b) => a[0] - b[0])) {
const last = merged[merged.length - 1]
if (last && r[0] <= last[1] + 1) last[1] = Math.max(last[1], r[1])
else merged.push([r[0], r[1]])
}
this.committedRanges = merged
}
this.horizonGen = Math.max(this.horizonGen, this.segments.compactedBelow() - 1)
} else {
this.segments = null
}
// 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<void> {
// 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 Read-only history footprint for fleet audits: how much
* generational history this store holds on disk. `bytes` pays (and seeds)
* the one-time {@link historyBytes} walk on first call — subsequent calls
* are O(1). The oldest/newest timestamps come from those generations'
* deltas (cache-bounded reads).
* @returns Counts, bytes, generation range, and the compaction horizon.
*/
/**
* @description D8 (gate-to-generation provenance): a deterministic content
* digest of the generation log THROUGH `g` — identical history ⇒ identical
* digest on any machine; any divergence (different records, different
* order, reclaimed range) ⇒ different digest. Composed from the packed
* tier's sealed-segment checksum chain (O(segments)) plus the live tier's
* per-generation delta digests (O(live window at most)). Release gates pin
* {generation, digest} and verify both at execution time.
* @param g - The generation to digest through (≤ committed).
* @returns A hex digest string, stable across reopen and repacking states
* ONLY for fully-packed prefixes — repacking changes representation, so
* the composed digest is defined over CONTENT: live-tier gens hash their
* delta + record ids, packed gens hash via frame CRCs. A gate should pin
* after a repack pass for long-term stability, or re-pin on repack.
*/
async generationDigest(g: number): Promise<string> {
if (!Number.isInteger(g) || g < 1 || g > this.committed) {
throw new RangeError(
`generationDigest(): generation ${g} is out of range [1, ${this.committed}]`
)
}
if (g <= this.horizonGen) {
throw new GenerationCompactedError(g, this.horizonGen)
}
let digest = 0
const enc = new TextEncoder()
if (this.segments) {
const packed = await this.segments.digestThroughPacked(g)
if (packed !== null) digest = packed
}
// Live-tier composition: every committed gen ≤ g not covered by a sealed
// segment hashes its delta content in ascending order.
for (const gen of this.committedGensAsc()) {
if (gen > g) break
if (this.segments?.hasGeneration(gen)) continue
const delta = await this.getDelta(gen)
digest = crc32c(
enc.encode(
`${digest}:${gen}:${delta.timestamp}:${[...delta.nouns].sort().join(',')}:${[...delta.verbs].sort().join(',')}`
)
)
}
return digest.toString(16).padStart(8, '0')
}
async historyStats(): Promise<{
generations: number
bytes: number
oldestGeneration: number | null
newestGeneration: number | null
oldestTimestamp: number | null
newestTimestamp: number | null
horizon: number
}> {
let oldest: number | null = null
let newest: number | null = null
for (const gen of this.committedGensAsc()) {
if (oldest === null) oldest = gen
newest = gen
}
return {
generations: this.committedCount(),
bytes: await this.historyBytes(),
oldestGeneration: oldest,
newestGeneration: newest,
oldestTimestamp: oldest !== null ? (await this.getDelta(oldest)).timestamp : null,
newestTimestamp: newest !== null ? (await this.getDelta(newest)).timestamp : null,
horizon: this.horizonGen
}
}
/**
* @description Read one generation's persisted before-image records — the
* compaction fallback for generations written before deltas carried
* `blobHashes`. O(that generation's records).
* @param gen - The generation whose `prev/` records to read.
* @returns The record-set (empty when absent).
*/
private async readGenerationRecords(gen: number): Promise<GenerationRecord[]> {
let paths: string[] = []
try {
paths = await this.storage.listRawObjects(`${GENERATIONS_PREFIX}/${gen}/prev`)
} catch {
paths = []
}
const records: GenerationRecord[] = []
for (const p of paths) {
const record = (await this.storage.readRawObject(p)) as GenerationRecord | null
if (record) records.push(record)
}
if (records.length > 0) return records
// Two-tier: folded generations serve their record-set from the segment.
const packed = await this.segments?.readRecords(gen)
return packed ? (packed.map((r) => r.record) as GenerationRecord[]) : []
}
/**
* @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<void> {
if (!this.opened) return
await this.withMutex(() => this.persistCounterUnlocked())
}
private async persistCounterUnlocked(): Promise<void> {
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.
*/
/**
* The generation fact log, or `null` when the storage layer cannot host one.
* Consumers scan committed facts through it (`scanFacts` / `segmentPaths`).
*/
getFactLog(): FactLog | null {
return this.factLog
}
/**
* @description Build one commit's AFTER-IMAGE fact by reading canonical
* state back for every touched id — under the commit mutex, immediately
* after the operations applied, so canonical IS the after-image (and the
* reads are page-cache-warm: the operations just wrote these files). An
* absent id (both legs null) becomes a body-less TOMBSTONE — the delete
* fact needs no body, so removal never requires reading the removed thing.
* The fact's blobHashes are extracted from the AFTER records (the content
* this generation's state references), unlike the history path's
* before-image hashes.
*/
private async buildCommitFact(args: {
generation: number
timestamp: number
nouns: string[]
verbs: string[]
meta?: Record<string, unknown>
}): Promise<CommitFact> {
const ops: FactOp[] = []
const afterRecords: GenerationRecord[] = []
for (const id of args.nouns) {
const after = await this.storage.readNounRaw(id)
const absent = after.metadata === null && after.vector === null
ops.push({ kind: 'noun', id, record: absent ? null : after })
if (!absent) afterRecords.push({ kind: 'noun', metadata: after.metadata, vector: after.vector })
}
for (const id of args.verbs) {
const after = await this.storage.readVerbRaw(id)
const absent = after.metadata === null && after.vector === null
ops.push({ kind: 'verb', id, record: absent ? null : after })
if (!absent) afterRecords.push({ kind: 'verb', metadata: after.metadata, vector: after.vector })
}
const blobHashes = this.storage.extractBlobHashesFromRecords
? this.storage.extractBlobHashesFromRecords(afterRecords)
: []
return {
generation: args.generation,
timestamp: args.timestamp,
ops,
...(args.meta ? { meta: args.meta } : {}),
...(blobHashes.length > 0 ? { blobHashes } : {})
}
}
async commitTransaction(args: {
touched: TouchedIds
meta?: Record<string, unknown>
ifAtGeneration?: number
/** Optional compare-and-swap precondition over the {@link CommitBeforeImages},
* run under the commit mutex BEFORE anything is staged or applied — the
* per-record analogue of `ifAtGeneration`. A throw aborts the whole batch:
* the generation reservation is returned and no staging I/O has happened. */
precommit?: (before: CommitBeforeImages) => void
execute: () => Promise<void>
}): Promise<{ generation: number; timestamp: number }> {
return this.withMutex(async () => {
// A latched history-durability failure compromises the whole generation
// spine — refuse a transact too (advancing the manifest past stuck,
// un-durable single-op generations would be inconsistent). Same loud
// error; self-clears when the pending tier drains.
this.assertHistoryDurable()
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
}
}
// Temporal-blob contract: hashes this record-set references, recorded
// before staging; scoped outside the try so an abort can compensate.
let txBlobHashes: string[] = []
// Hoisted so the catch can reconcile canonical state against them after a
// failed rollback (the trapdoor). Empty until populated below.
const nounBefore = new Map<string, GenerationRecord>()
const verbBefore = new Map<string, GenerationRecord>()
try {
// -- 3. Before-images + delta (the durable undo log) ------------------
// Read every before-image FIRST, then run the caller's CAS
// precondition against them, and only then stage to disk — so a
// conflicting batch aborts with zero staging I/O. The maps hold the
// byte-identical records the staged files are written from.
for (const id of nouns) {
const prev = await this.storage.readNounRaw(id)
nounBefore.set(id, { kind: 'noun', metadata: prev.metadata, vector: prev.vector })
}
for (const id of verbs) {
const prev = await this.storage.readVerbRaw(id)
verbBefore.set(id, { kind: 'verb', metadata: prev.metadata, vector: prev.vector })
}
// Conditional commit: the per-record CAS analogue of the
// `ifAtGeneration` check above, but against authoritative
// before-images under the mutex. A throw lands in the catch below,
// which returns the generation reservation; nothing was applied.
args.precommit?.({ nouns: nounBefore, verbs: verbBefore })
// Temporal-blob contract: count this record-set's content-hash
// references BEFORE it is staged (over-count-only crash ordering —
// see flushPendingSingleOps' matching note).
txBlobHashes = this.storage.extractBlobHashesFromRecords
? this.storage.extractBlobHashesFromRecords([
...nounBefore.values(),
...verbBefore.values()
])
: []
if (txBlobHashes.length > 0 && this.storage.recordHistoryBlobReferences) {
await this.storage.recordHistoryBlobReferences(txBlobHashes)
}
const stagedPaths: string[] = []
let recordBytes = 0 // serialized record-set size, for retention accounting
for (const [id, record] of nounBefore) {
const recordPath = `${dir}/prev/${id}.json`
await this.storage.writeRawObject(recordPath, record)
recordBytes += serializedBytes(record)
stagedPaths.push(recordPath)
}
for (const [id, record] of verbBefore) {
const recordPath = `${dir}/prev/${id}.json`
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,
// Always present on new deltas (empty = "no blobs"), so compaction
// only falls back to record reads for pre-contract generations.
blobHashes: txBlobHashes
}
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 -------------------------------------
// Durability barrier: record every canonical write/delete the operations
// make, then fsync them BELOW (before the counter advances). Without
// this, a hard kill after commitTransaction returns could leave the
// fsync'd generation counter ahead of still-page-cached entity bytes —
// phantom progress for any generation-based consumer.
this.storage.beginWriteBarrier?.()
this.inTransact = true
try {
await args.execute()
} finally {
this.inTransact = false
}
// The transaction's entire canonical footprint is now durable, so the
// counter/manifest advance below can never outrun the entity bytes.
await this.storage.flushWriteBarrier?.()
faultPoint('after-execute')
// Fact log (dual-write): append + fsync this generation's AFTER-IMAGE
// fact BEFORE the commit point, inside the same durability window —
// so a crash can only leave the log AHEAD (open truncates), never a
// committed generation without its fact. A real abort below this point
// compensates via dropAbove in the catch. An append failure fails the
// write, loudly — a silent fact gap would be a lie a replay discovers.
if (this.factLog) {
const fact = await this.buildCommitFact({
generation: gen,
timestamp,
nouns,
verbs,
...(args.meta ? { meta: args.meta } : {})
})
await this.factLog.append(fact)
await this.factLog.sync()
}
// -- 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
})
if (this.historyBytesTotal !== null) {
this.historyBytesTotal += 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
}
// The trapdoor for a batch: if rollback FAILED to fully apply, canonical
// storage may be inconsistent. A batch is never adopted forward (its
// other ops were rolled back — partial commit would break atomicity), so
// any unreconciled record is a fail-loud StoreInconsistentError. Compute
// it here (before the staging cleanup, which is always safe to run) and
// throw it in place of the raw error at the end.
let inconsistent: UnreconciledRecord[] = []
if (err instanceof TransactionRollbackError) {
inconsistent = await this.reconcileFailedRollback(nounBefore, verbBefore)
}
// 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)`
)
}
// Compensate the pre-staging history-reference increments. Best
// effort — a failure here only over-counts (a leak the scrub
// repairs). Reclaim inside is safe on an abort: the before-image
// hashes are the entities' still-live content, so live references
// block any physical delete.
if (txBlobHashes.length > 0 && this.storage.releaseHistoryBlobReferences) {
try {
await this.storage.releaseHistoryBlobReferences(txBlobHashes)
} catch {
// over-count-safe; the scrub restores exactness
}
}
// Fact-log compensation: a real (non-crash) abort after the fact was
// appended must take the fact back out — the generation never
// committed. A crash instead reaches open(), whose truncation does the
// same reconcile from disk.
if (this.factLog && this.factLog.headGeneration() >= gen) {
await this.factLog.dropAbove(gen - 1)
}
// 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
// A failed rollback that left canonical records unreconciled surfaces as
// a loud StoreInconsistentError (brainy quarantines writes until repair);
// otherwise the raw error (clean abort, or a derived-only undo failure
// the egress guard + rebuild handle).
if (inconsistent.length > 0) {
throw new StoreInconsistentError(
inconsistent,
(err as TransactionRollbackError).originalError
)
}
throw err
}
})
}
/**
* @description After a transaction's rollback FAILED to fully apply
* (`TransactionRollbackError`), determine which touched records are now
* unreconciled by comparing current canonical state to the before-images the
* commit captured. This observes reality rather than trusting the opaque undo
* closures — the authoritative signal for adopt-forward vs fail-loud.
*
* @param nounBefore - Byte-identical entity before-images captured pre-execute.
* @param verbBefore - Byte-identical relationship before-images.
* @returns Every record whose canonical state no longer matches its
* before-image, tagged `'orphan'` (durably present when it should be gone —
* an add whose delete-undo failed) or `'loss'` (durably gone/wrong when it
* should have been restored — a remove/update whose restore-undo failed). An
* empty array means canonical storage is cleanly rolled back (only a derived
* index undo failed — the egress guard + a rebuild handle that).
*/
private async reconcileFailedRollback(
nounBefore: Map<string, GenerationRecord>,
verbBefore: Map<string, GenerationRecord>
): Promise<UnreconciledRecord[]> {
const out: UnreconciledRecord[] = []
const classify = (
before: GenerationRecord,
current: { metadata: unknown | null; vector: unknown | null }
): 'orphan' | 'loss' | null => {
const beforeEmpty = before.metadata == null && before.vector == null
const currentEmpty = current.metadata == null && current.vector == null
if (beforeEmpty && currentEmpty) return null // add cleanly undone
if (beforeEmpty) return 'orphan' // add's delete-undo failed → still present
if (currentEmpty) return 'loss' // remove's restore-undo failed → gone
// Both present: same value = reconciled; different = restore left a wrong value.
const same =
JSON.stringify({ m: before.metadata, v: before.vector }) ===
JSON.stringify({ m: current.metadata, v: current.vector })
return same ? null : 'loss'
}
for (const [id, before] of nounBefore) {
const d = classify(before, await this.storage.readNounRaw(id))
if (d) out.push({ id, kind: 'noun', disposition: d })
}
for (const [id, before] of verbBefore) {
const d = classify(before, await this.storage.readVerbRaw(id))
if (d) out.push({ id, kind: 'verb', disposition: d })
}
return out
}
// ==========================================================================
// 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}).
*
* **Durability contract (single-op vs transact).** A single-op write's live
* canonical bytes are written via tmp+rename but NOT individually fsync'd —
* they are durable at the next {@link flushPendingSingleOps} or `close()`, not
* the instant the write resolves. On a hard kill before that flush, a
* single-op write can be lost even though the call returned; the group-commit
* counter is buffered alongside, so counter and data are lost together (no
* counter-ahead-of-state torn store). {@link commitTransaction} is the
* stronger contract: it runs a write barrier that fsyncs the whole batch's
* canonical footprint BEFORE advancing the generation counter, so a committed
* transact is durable on return. Callers that need per-write durability should
* use `transact()` (or `flush()` after a single-op); Model-B group-commit
* deliberately trades single-op fsync latency (a 3-5x write regression) for
* throughput.
*
* 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.
* @param args.precommit - Optional compare-and-swap precondition, invoked
* under the commit mutex with the just-read {@link CommitBeforeImages} and
* BEFORE `execute()`. A throw aborts the commit atomically: the generation
* reservation is returned and nothing is applied or buffered. This is the
* one place a per-record CAS (e.g. `ifRev`) is race-free — any check done
* before this mutex can interleave with a concurrent writer.
* @returns The reserved generation and its timestamp.
*/
async commitSingleOp(args: {
touched: { nouns?: string[]; verbs?: string[] }
execute: () => Promise<void>
precommit?: (before: CommitBeforeImages) => void
}): Promise<{ generation: number; timestamp: number; degraded?: string[] }> {
return this.withMutex(async () => {
// Refuse to accept a write whose history we cannot make durable: if the
// pending-tier flush has latched a persistent failure, accepting more
// writes would silently pile un-durable before-images into memory. Fail
// loud instead (the latch self-clears when a flush finally succeeds).
this.assertHistoryDurable()
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<string, GenerationRecord>()
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<string, GenerationRecord>()
for (const id of verbs) {
const prev = await this.storage.readVerbRaw(id)
verbBefore.set(id, { kind: 'verb', metadata: prev.metadata, vector: prev.vector })
}
// Conditional commit: the caller's CAS precondition runs here — under
// the mutex, against the authoritative before-images, before any write.
// A throw aborts cleanly: return the generation reservation (nothing was
// applied or buffered) and surface the conflict to the caller.
if (args.precommit) {
try {
args.precommit({ nouns: nounBefore, verbs: verbBefore })
} catch (err) {
if (this.counter === gen) this.counter = gen - 1
throw err
}
}
// 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
// A failed rollback (TransactionRollbackError) may have left canonical
// storage inconsistent — the trapdoor. Reconcile against the
// before-images to decide the honest response (David's ruling:
// adopt-forward when safe, else fail loud).
if (err instanceof TransactionRollbackError) {
const records = await this.reconcileFailedRollback(nounBefore, verbBefore)
if (records.length > 0 && records.every((r) => r.disposition === 'orphan')) {
// ADOPT FORWARD: the durably-present orphan(s) are exactly what this
// single-op add wrote. Keep + buffer the generation so the record is
// legitimately committed and readable; the derived index may be
// incomplete for these ids until the next rebuild/repairIndex (the
// egress guard prevents wrong results meanwhile). Loud, honest,
// no double-write.
this.pendingBuffer.set(gen, { nouns: nounBefore, verbs: verbBefore, timestamp })
this.pendingGens.push(gen)
this.extendChains(gen, nouns, verbs)
// The adopted generation is committed — it gets its fact like any
// other (durability rides the group-commit flush, same as the
// buffered history).
if (this.factLog) {
await this.factLog.append(
await this.buildCommitFact({ generation: gen, timestamp, nouns, verbs })
)
}
prodLog.warn(
`[GenerationStore] Recovered a failed rollback FORWARD: single-op write ` +
`committed as generation ${gen} because its canonical undo could not be ` +
`applied (record(s) ${records.map((r) => r.id).join(', ')} are durable). ` +
`The derived index may be incomplete for these ids — run repairIndex() to heal.`
)
return { generation: gen, timestamp, degraded: records.map((r) => r.id) }
}
if (records.length > 0) {
// FAIL LOUD: a loss (or mixed) cannot be safely adopted. Return the
// reservation and throw — brainy quarantines writes until repair.
if (this.counter === gen) this.counter = gen - 1
throw new StoreInconsistentError(records, err.originalError)
}
// records.length === 0: canonical is cleanly rolled back (only a
// derived undo failed) — fall through to the clean-abort path below.
}
// Live write failed with canonical cleanly rolled back: nothing buffered,
// nothing durable. 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)
// Fact log (dual-write): the acked write's AFTER-IMAGE fact, appended
// now (read back warm, under the mutex — group-commit means flush-time
// canonical only holds the LATEST state, so each generation's after-image
// exists only here). Durability rides the group-commit flush, exactly
// like the buffered before-image history: a crash before the flush loses
// the fact AND the generation together — never a torn state.
if (this.factLog) {
await this.factLog.append(
await this.buildCommitFact({ generation: gen, timestamp, 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<void>): Promise<void> {
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/<N>/` record
* set (`prev/<id>.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<void> {
if (this.pendingGens.length === 0) return
// Centralize durability accounting here so EVERY failure path — the
// background scheduler AND an explicit flush()/close() — latches
// consistently, and success (from any caller) clears the latch. The
// scheduler owns only the retry cadence.
try {
await this.flushPendingSingleOpsUnlocked()
this.onPendingFlushSuccess()
} catch (err) {
this.recordPendingFlushFailure(err as Error)
throw err
}
}
/** The actual pending-tier flush under the mutex (see {@link flushPendingSingleOps},
* which wraps this with durability accounting). */
private async flushPendingSingleOpsUnlocked(): Promise<void> {
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<number, number>() // for the post-commit setDelta
for (const gen of gens) {
const buf = this.pendingBuffer.get(gen)
if (!buf) continue
const dir = `${GENERATIONS_PREFIX}/${gen}`
// Temporal-blob contract: count this record-set's content-hash
// references BEFORE persisting it (a crash between the two only
// over-counts — the scrub repairs a leak; under-counting could let
// compaction reclaim bytes a retained generation still needs).
const blobHashes = this.storage.extractBlobHashesFromRecords
? this.storage.extractBlobHashesFromRecords([
...buf.nouns.values(),
...buf.verbs.values()
])
: []
if (blobHashes.length > 0 && this.storage.recordHistoryBlobReferences) {
await this.storage.recordHistoryBlobReferences(blobHashes)
}
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,
// Always present on new deltas (empty = "no blobs"), so compaction
// only falls back to record reads for pre-contract generations.
blobHashes
}
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)
// Fact log (dual-write): make the window's buffered facts durable in the
// same batch, BEFORE the commit point below — so a crash can only leave
// the log AHEAD of the counter (open truncates), never a committed
// generation without its durable fact.
await this.factLog?.sync()
// 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
})
if (this.historyBytesTotal !== null) {
this.historyBytesTotal += genBytes.get(gen) ?? 0
}
this.pendingBuffer.delete(gen)
}
this.pendingGens = []
for (const entry of logEntries) {
await this.storage.appendTxLogLine(JSON.stringify(entry))
}
})
}
/**
* @description Reset the durability-failure state after a successful flush.
* If a failure was latched (writes were being refused), log the recovery so
* the transition out of the loud state is as visible as the transition in.
*/
private onPendingFlushSuccess(): void {
if (this.pendingFlushError) {
prodLog.info(
`[GenerationStore] pending single-op history flush recovered after ` +
`${this.pendingFlushFailures} failed attempt(s); writes resume.`
)
}
this.pendingFlushFailures = 0
this.pendingFlushError = null
}
/**
* @description Account for a failed pending-tier flush (called from
* {@link flushPendingSingleOps}' catch, so it fires on EVERY failure path —
* background scheduler and explicit flush()/close() alike). Escalates a
* transient blip (a warn, keep retrying) into a latched durability failure
* once {@link PENDING_FLUSH_FAILURE_THRESHOLD} consecutive attempts fail —
* from that point writes are refused with a {@link PendingFlushDurabilityError}
* until a flush finally succeeds. Always ensures a backoff retry is scheduled
* so a recovering disk self-heals without operator action, regardless of who
* triggered the failing flush. The pending before-images are retained (the
* failed flush never cleared them), so no history is lost — only not-yet-durable.
*/
private recordPendingFlushFailure(err: Error): void {
this.pendingFlushFailures++
const attempts = this.pendingFlushFailures
if (attempts >= GenerationStore.PENDING_FLUSH_FAILURE_THRESHOLD) {
// Latch: refuse further writes rather than pile un-durable history.
this.pendingFlushError = err
prodLog.error(
`[GenerationStore] pending single-op history flush has FAILED ${attempts} ` +
`consecutive times (${err.message}). Generation history is not durable; ` +
`writes are now REFUSED until it drains. Live canonical data is intact. ` +
`Retrying with backoff; the store resumes writes when a flush succeeds.`
)
} else {
prodLog.warn(
`[GenerationStore] pending single-op flush failed (attempt ${attempts}/` +
`${GenerationStore.PENDING_FLUSH_FAILURE_THRESHOLD}): ${err.message}; retrying with backoff.`
)
}
this.scheduleFlushRetry(attempts)
}
/**
* @description (Re)arm the pending-flush timer at a capped exponential
* backoff after a failure, so a persistent fault is retried at a decaying
* cadence instead of hot-looping. Shares the single {@link pendingFlushTimer}
* slot with {@link schedulePendingFlush} (only one flush is ever scheduled).
* The retry's own failure is re-accounted inside {@link flushPendingSingleOps}.
*/
private scheduleFlushRetry(attempt: number): void {
if (this.pendingFlushTimer !== null) return
const backoff = Math.min(
GenerationStore.PENDING_FLUSH_DELAY_MS * 2 ** Math.min(attempt, 6),
GenerationStore.PENDING_FLUSH_RETRY_CAP_MS
)
this.pendingFlushTimer = setTimeout(() => {
this.pendingFlushTimer = null
// flushPendingSingleOps records the failure + reschedules internally.
void this.flushPendingSingleOps().catch(() => {})
}, backoff)
const t = this.pendingFlushTimer as unknown as { unref?: () => void }
if (t && typeof t.unref === 'function') t.unref()
}
/**
* @description Throw if the store has a latched history-durability failure.
* Called at the top of every write commit so a broken durability path fails
* loud and immediately instead of silently accumulating un-durable history.
*/
private assertHistoryDurable(): void {
if (this.pendingFlushError) {
throw new PendingFlushDurabilityError(this.pendingFlushError, this.pendingFlushFailures)
}
}
/** 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
* failed flush is accounted + rescheduled inside {@link flushPendingSingleOps}
* (latch-on-persistent-failure), never swallowed as a bare log line. */
private schedulePendingFlush(): void {
if (this.pendingGens.length >= this.pendingFlushThreshold) {
// Failure is recorded + retried inside flushPendingSingleOps.
void Promise.resolve()
.then(() => this.flushPendingSingleOps())
.catch(() => {})
return
}
if (this.pendingFlushTimer !== null) return
this.pendingFlushTimer = setTimeout(() => {
this.pendingFlushTimer = null
void this.flushPendingSingleOps().catch(() => {})
}, 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<number> {
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<number> {
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.ensureWindow()
// Snapshot windowLo, then read the chain synchronously (no await between):
// a concurrent commit can only slide the window at the `await` above, so the
// routing decision and the chain it reads are a consistent pair.
const windowLo = this.windowLo
let candidate: number | undefined
if (gen >= windowLo) {
// HOT path — every generation that could hold the before-image
// (`> gen ≥ windowLo`) is inside the resident window, so the recent chain
// is complete for this pin. O(log chain).
const chain = (kind === 'noun' ? this.recentNounChains : this.recentVerbChains).get(id)
if (chain === undefined) return { source: 'current' }
candidate = firstGenerationAfter(chain, gen)
} else {
// COLD path — a deep pin below the window. Serve from the bounded LRU,
// reconstructing (lock-light) on a miss. Caching empty chains too keeps
// repeated deep reads of untouched-in-cold-region ids O(1).
const cold = kind === 'noun' ? this.coldNounChains : this.coldVerbChains
let chain = cold.get(id)
if (chain !== undefined) {
cold.delete(id) // re-insert to mark most-recently-used
cold.set(id, chain)
} else {
chain = await this.reconstructColdChain(kind, id)
cold.set(id, chain)
while (cold.size > this.coldChainLruMax) {
const oldest = cold.keys().next().value
if (oldest === undefined) break
cold.delete(oldest)
}
}
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 Resolve the FIRST reserved generation after `gen` that touched
* each of `ids`, by kind — the bulk, MUTEX-FREE counterpart of
* {@link resolveAt}'s chain lookup. A SINGLE ascending pass over the reserved
* (committed pending) generations records `firstAfter[id] = g` for the first
* `g > gen` whose delta touched `id`, stopping early once every id is resolved.
*
* Ids absent from the returned map were never touched after `gen` → their live
* storage state IS their state at `gen` (the {@link resolveAt} `'current'`
* answer). Ids present map to the generation whose before-image
* ({@link readGenerationRecord}) holds their state as of `gen`.
*
* WHY MUTEX-FREE (load-bearing): {@link Brainy.materializeAtGeneration} runs its
* final reconciliation inside {@link snapshotWith} (which HOLDS the commit
* mutex). Resolving ids one-by-one through {@link resolveAt} there would (a)
* re-enter the mutex via {@link ensureWindow} → DEADLOCK, and (b) regress an
* O(R) bulk copy to O(N·R). This method takes no lock and makes one pass, so a
* deep-pin materialize over N ids stays O(R·d̄ + |ids|) with O(R) `getDelta`
* calls. It iterates SNAPSHOTS of the small committed-range list (O(gaps)) and
* the pending list (O(pending)) — NOT a materialized O(R) gen array — so it is
* safe against concurrent commits and bounded in memory.
*
* @param kind - `'noun'` for entities, `'verb'` for relationships.
* @param ids - The ids to resolve.
* @param gen - The pinned generation.
* @returns `id → first reserved generation after `gen` that touched it` for the
* subset of `ids` touched after `gen`.
*/
async resolveManyAt(
kind: 'noun' | 'verb',
ids: Set<string>,
gen: number
): Promise<Map<string, number>> {
const firstAfter = new Map<string, number>()
if (ids.size === 0) return firstAfter
// Snapshot the small range list + pending list (NOT the expanded gens) so the
// pass is immune to a concurrent compaction splice / commit append and stays
// O(gaps + pending) in memory.
const ranges = this.committedRanges.map(([s, e]) => [s, e] as [number, number])
const pending = [...this.pendingGens]
const record = async (g: number): Promise<boolean> => {
const delta = await this.getDelta(g)
const touched = kind === 'noun' ? delta.nouns : delta.verbs
for (const id of touched) {
if (ids.has(id) && !firstAfter.has(id)) firstAfter.set(id, g)
}
return firstAfter.size === ids.size
}
for (const [s, e] of ranges) {
if (e <= gen) continue
for (let g = Math.max(s, gen + 1); g <= e; g++) {
if (await record(g)) return firstAfter
}
}
for (const g of pending) {
if (g <= gen) continue
if (await record(g)) return firstAfter
}
return firstAfter
}
/**
* @description Read the before-image of `id` recorded by generation `gen` — the
* public, MUTEX-FREE accessor {@link Brainy.materializeAtGeneration}'s `copyAt`
* uses after {@link resolveManyAt} hands it the generation. Pairs with
* `resolveManyAt` so the bulk historical copy never touches the commit mutex
* (no deadlock inside {@link snapshotWith}, no per-id chain rebuild).
* @param kind - `'noun'` for entities, `'verb'` for relationships.
* @param gen - The generation whose before-image to read.
* @param id - The id to read.
* @returns The stored {@link GenerationRecord}, or `null` when the record-set
* has no before-image for `id`.
*/
async readGenerationRecord(
kind: 'noun' | 'verb',
gen: number,
id: string
): Promise<GenerationRecord | null> {
return this.readBeforeImage(kind, gen, id)
}
/**
* @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/<gen>/prev/<id>.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<GenerationRecord | null> {
const pending = this.pendingBuffer.get(gen)
if (pending) {
return (kind === 'noun' ? pending.nouns : pending.verbs).get(id) ?? null
}
const live = (await this.storage.readRawObject(
`${GENERATIONS_PREFIX}/${gen}/prev/${id}.json`
)) as GenerationRecord | null
if (live) return live
// Two-tier: the packed tier serves folded generations (live-tier-wins).
if (this.segments?.hasGeneration(gen)) {
return (await this.segments.readRecord(gen, kind, id)) as GenerationRecord | null
}
return null
}
/**
* @description Build the resident hot-tail window — the per-id chains
* ({@link recentNounChains}/{@link recentVerbChains}) and the inverse index
* ({@link windowNounDeltas}/{@link windowVerbDeltas}) over the last
* {@link recentWindowGenerations} reserved generations, `[windowLo, counter]`.
* Built once, lazily, under the commit mutex (so no concurrent commit mutates
* {@link committedRanges}/pending mid-build); thereafter maintained
* incrementally by {@link extendChains} and invalidated by {@link compact} /
* reopen. Concurrent first-callers share one build. O(W) `getDelta` reads the
* first time (only the window, NOT all of history); O(1) afterwards. Deep pins
* below `windowLo` are served by {@link reconstructColdChain} on demand.
*/
private async ensureWindow(): Promise<void> {
if (this.windowReady) return
if (!this.windowBuilding) {
this.windowBuilding = this.withMutex(async () => {
if (this.windowReady) return
this.recentNounChains.clear()
this.recentVerbChains.clear()
this.windowNounDeltas.clear()
this.windowVerbDeltas.clear()
// The window is the newest W reserved generations, clamped above the
// compaction horizon (gens ≤ horizon are reclaimed — never resident).
this.windowLo = Math.max(this.horizonGen + 1, this.counter - this.recentWindowGenerations + 1)
// reservedGensAsc (committed pending) is ascending, so chains accrue in
// order and un-flushed single-ops in-window are indexed too.
for (const g of this.reservedGensAsc()) {
if (g < this.windowLo) continue
const delta = await this.getDelta(g)
const nouns = [...delta.nouns]
const verbs = [...delta.verbs]
for (const id of nouns) appendToChain(this.recentNounChains, id, g)
for (const id of verbs) appendToChain(this.recentVerbChains, id, g)
this.windowNounDeltas.set(g, nouns)
this.windowVerbDeltas.set(g, verbs)
}
this.windowReady = true
this.windowBuilding = null
})
}
await this.windowBuilding
}
/**
* @description Reconstruct the FULL per-id chain for a deep pin (`gen <
* windowLo`) — every generation that touched `id`, ascending. Built from the
* persisted deltas BELOW the window plus the resident in-window tail, then
* cached in the cold LRU by the caller.
*
* LOCK-LIGHT — deliberately does NOT hold the commit mutex (a deep read must
* never stall writers). Correctness without the lock:
* - The recent tail and `windowLo` are snapshotted SYNCHRONOUSLY up front, so
* even if the window slides during the (awaited) cold scan, the chain is the
* one consistent with the snapshot. New writes during the scan land ABOVE the
* snapshot and cannot change any deep pin's first-after answer; the cold
* entry is also dropped by {@link extendChains} the moment `id` is next
* touched, so it is never served stale.
* - The cold scan iterates a SNAPSHOT of the small committed-range list
* (O(gaps) memory, immune to a concurrent compaction splice).
* - A generation concurrently reclaimed by compaction surfaces as a missing
* delta. Such a gen is always `≤ minPinned` (compaction never reclaims above
* the lowest live pin) and therefore NEVER a live pin's answer (`first-after
* > pin ≥ minPinned`), so it is skipped, not raised. A missing delta ABOVE
* `minPinned` is genuine corruption and still throws.
*
* @param kind - `'noun'` for entities, `'verb'` for relationships.
* @param id - The id whose chain to reconstruct.
* @returns The id's full ascending generation chain (`[]` when never touched).
*/
private async reconstructColdChain(kind: 'noun' | 'verb', id: string): Promise<number[]> {
// Synchronous, consistent snapshot of the window state.
const windowLo = this.windowLo
const recentTail = (kind === 'noun' ? this.recentNounChains : this.recentVerbChains).get(id)
const recentSnapshot = recentTail ? [...recentTail] : []
const ranges = this.committedRanges.map(([s, e]) => [s, e] as [number, number])
const pendingCold = this.pendingGens.filter((g) => g < windowLo)
const cold: number[] = []
const scan = async (g: number): Promise<void> => {
let delta: { nouns: Set<string>; verbs: Set<string> }
try {
delta = await this.getDelta(g)
} catch (err) {
// Concurrent-compaction tolerance: a reclaimed gen (≤ minPinned, never a
// live pin's answer) is skipped; a missing gen above the lowest pin is
// real corruption and rethrows.
if (g <= this.minPinnedGeneration()) return
throw err
}
if ((kind === 'noun' ? delta.nouns : delta.verbs).has(id)) cold.push(g)
}
for (const [s, e] of ranges) {
if (s >= windowLo) break // ascending → nothing newer is below the window
const hi = Math.min(e, windowLo - 1)
for (let g = s; g <= hi; g++) await scan(g)
}
// Cold pending gens (only when more than W writes are buffered, i.e. tests
// with a small window) sit between the committed gens and the window.
for (const g of pendingCold) await scan(g)
// Cold gens (< windowLo) then the resident in-window tail (≥ windowLo) — both
// ascending and disjoint, so the concatenation is the full ascending chain.
return recentSnapshot.length ? cold.concat(recentSnapshot) : cold
}
/**
* Record that generation `gen` (the newest) touched these ids and advance the
* hot-tail window. Keeps the resident window current after a commit without a
* rebuild. No-op until the window is built (the eventual build reads `gen` from
* {@link reservedGensAsc}).
*
* Three steps, all in-memory and I/O-FREE:
* 1. Append `gen` to each touched id's recent chain, and DROP each touched id
* from the cold LRU — a freshly-touched id's cached deep chain is now
* incomplete, so it must be reconstructed on the next deep read (else a
* deep pin whose first-after just became `gen` would wrongly read `current`).
* 2. Record the window inverse index entry `gen → touched ids` (the slide-trim
* driver).
* 3. Slide `windowLo` forward to `max(horizon+1, gen-W+1)`, front-trimming the
* chains of exactly the ids leaving the window using the inverse index — NO
* `getDelta`, so the write path never reads `tx.json` on a slide regardless
* of W vs the delta-cache size.
*/
private extendChains(gen: number, nouns: Iterable<string>, verbs: Iterable<string>): void {
if (!this.windowReady) return
const nounArr = [...nouns]
const verbArr = [...verbs]
for (const id of nounArr) {
appendToChain(this.recentNounChains, id, gen)
this.coldNounChains.delete(id)
}
for (const id of verbArr) {
appendToChain(this.recentVerbChains, id, gen)
this.coldVerbChains.delete(id)
}
this.windowNounDeltas.set(gen, nounArr)
this.windowVerbDeltas.set(gen, verbArr)
const newLo = Math.max(this.horizonGen + 1, gen - this.recentWindowGenerations + 1)
if (newLo <= this.windowLo) return
// Collect the ids leaving the window from the inverse index (O(slid·d̄)), then
// front-trim their chains to the new low watermark (each id once).
const trimNouns = new Set<string>()
const trimVerbs = new Set<string>()
for (let g = this.windowLo; g < newLo; g++) {
const n = this.windowNounDeltas.get(g)
if (n) for (const id of n) trimNouns.add(id)
this.windowNounDeltas.delete(g)
const v = this.windowVerbDeltas.get(g)
if (v) for (const id of v) trimVerbs.add(id)
this.windowVerbDeltas.delete(g)
}
for (const id of trimNouns) dropChainFront(this.recentNounChains, id, newLo)
for (const id of trimVerbs) dropChainFront(this.recentVerbChains, id, newLo)
this.windowLo = newLo
}
/** Drop the resident window + cold LRU so the next read rebuilds from the
* surviving generations (called after compaction removes record-sets / on
* reopen). */
private invalidateChains(): void {
this.windowReady = false
this.windowBuilding = null
this.recentNounChains.clear()
this.recentVerbChains.clear()
this.windowNounDeltas.clear()
this.windowVerbDeltas.clear()
this.coldNounChains.clear()
this.coldVerbChains.clear()
this.windowLo = 0
}
/**
* @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<ChangedIds> {
const nouns = new Set<string>()
const verbs = new Set<string>()
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<number[]> {
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<TxLogEntry[]> {
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<number | null> {
// 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<string>; verbs: Set<string>; 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) {
// Two-tier read (D1+D3): not in the live tier → the packed tier.
// Live-tier-wins ordering (a crash mid-fold leaves a duplicate, never
// a gap), so the segment lookup runs only after the live miss.
const packed = await this.segments?.readDelta(gen)
if (packed) {
const d = packed.delta as GenerationDelta
const entry = {
nouns: new Set(d.nouns),
verbs: new Set(d.verbs),
timestamp: packed.timestamp,
bytes: d.bytes ?? 0
}
this.setDelta(gen, entry)
return entry
}
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<string>; verbs: Set<string>; 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. O(1) after the first call: the
* total is computed by ONE walk over committed deltas, then maintained
* incrementally at every commit (+bytes) and reclaim (bytes) and dropped on
* a wholesale state replacement (restore). Without the running total, the
* adaptive auto-compaction on every flush() re-walked the ENTIRE history —
* O(committed generations) file reads per flush past the delta-cache bound —
* which is how a 70k-generation production brain turned every write into a
* full-tail scan (SELF-GENERATIONS-GROWTH). Pending (un-flushed) generations
* are excluded (they are not on disk).
* @returns The total on-disk history byte count.
*/
async historyBytes(): Promise<number> {
if (this.historyBytesTotal !== null) {
return this.historyBytesTotal
}
let total = 0
for (const gen of this.committedGensAsc()) {
total += (await this.getDelta(gen)).bytes
}
this.historyBytesTotal = total
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.
*/
/**
* @description The REPACKER (D1+D3+repacking): fold cold live-tier
* generations into sealed segments — re-representation, never deletion.
* Every record and delta stays readable (asOf/chains unchanged); the
* per-generation directories are deleted only AFTER their segment is
* durable (crash between = duplicate representation, resolved
* live-tier-wins by every reader; never a gap). This is the transform that
* takes a 70k-file history to tens of segment files, and the ONLY history
* transform permitted under the archival profile.
*
* Folds oldest-first, contiguous from the packed boundary, in batches, and
* stops at the live window ({@link GenerationStore.REPACK_LIVE_WINDOW})
* or when `timeBudgetMs` is spent — an early stop is a consistent prefix;
* the next pass resumes.
*/
async repackHistory(options?: { timeBudgetMs?: number; batchGenerations?: number }): Promise<{
foldedGenerations: number
segmentsCreated: number
}> {
if (!this.segments) return { foldedGenerations: 0, segmentsCreated: 0 }
const segments = this.segments
return this.withMutex(async () => {
const deadline =
options?.timeBudgetMs !== undefined ? Date.now() + options.timeBudgetMs : undefined
const batchSize = options?.batchGenerations ?? 512
const coldCeiling = this.committed - GenerationStore.REPACK_LIVE_WINDOW
const packedThrough =
segments.segments().length > 0
? segments.segments()[segments.segments().length - 1].lastGeneration
: 0
// Cold, unpacked, committed generations — ascending, contiguous scan.
const eligible: number[] = []
for (const gen of this.committedGensAsc()) {
if (gen > coldCeiling) break
if (gen <= packedThrough) continue // already packed (dup fold barred)
if (this.pendingBuffer.has(gen)) continue // un-flushed = live by definition
eligible.push(gen)
}
let folded = 0
let segmentsCreated = 0
for (let i = 0; i < eligible.length; i += batchSize) {
if (deadline !== undefined && Date.now() >= deadline) break
const batch = eligible.slice(i, i + batchSize)
const foldInput: FoldGeneration[] = []
for (const gen of batch) {
const delta = (await this.storage.readRawObject(
`${GENERATIONS_PREFIX}/${gen}/tx.json`
)) as GenerationDelta | null
if (delta === null) {
// Already folded by a prior crashed pass whose dirs were removed,
// or damage — getDelta's two-tier read decides which, loudly,
// when someone asks. Skip; never fold a generation we cannot read.
continue
}
const records: FoldGeneration['records'] = []
for (const [kind, ids] of [
['noun', delta.nouns] as const,
['verb', delta.verbs] as const
]) {
for (const id of ids) {
const record = await this.storage.readRawObject(
`${GENERATIONS_PREFIX}/${gen}/prev/${id}.json`
)
if (record) records.push({ kind, id, record })
}
}
foldInput.push({ generation: gen, timestamp: delta.timestamp, delta, records })
}
if (foldInput.length === 0) continue
await segments.fold(foldInput)
segmentsCreated++
// Segment + manifest durable → the live copies retire.
for (const g of foldInput) {
await this.storage.removeRawPrefix(`${GENERATIONS_PREFIX}/${g.generation}`)
}
folded += foldInput.length
}
if (folded > 0) {
prodLog.info(
`[GenerationStore] repacked ${folded} cold generation(s) into ${segmentsCreated} segment(s) — history preserved, file count reduced`
)
}
return { foldedGenerations: folded, segmentsCreated }
})
}
async compact(options?: CompactHistoryOptions): Promise<CompactHistoryResult> {
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
// Bounded maintenance pass (8.9.0): stop reclaiming once the budget is
// spent. Safe mid-loop — reclamation is oldest-first, so an early stop
// leaves a consistent contiguous prefix and the next pass resumes.
const deadline =
options?.timeBudgetMs !== undefined ? Date.now() + options.timeBudgetMs : 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 (deadline !== undefined && Date.now() >= deadline) break // budget spent — resume next pass
const delta = await this.getDelta(gen)
if (!noCaps) {
const violatesCount = maxGenerations !== undefined && remainingCount > maxGenerations
const violatesBytes = maxBytes !== undefined && remainingBytes > maxBytes
const violatesAge = ageCutoff !== undefined && 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 ? delta.bytes : 0
// Temporal-blob contract: resolve the content-blob hashes this
// generation's record-set references BEFORE deleting it. New
// generations carry the multiset in their persisted delta (empty
// array when none — distinguishing "new format, no blobs" from a
// pre-contract delta); legacy generations fall back to reading the
// records themselves. Skipped entirely on non-blob-aware storage.
let blobHashes: string[] | undefined
if (this.storage.releaseHistoryBlobReferences) {
const rawDelta = (await this.storage.readRawObject(
`${GENERATIONS_PREFIX}/${gen}/tx.json`
)) as GenerationDelta | null
blobHashes = rawDelta?.blobHashes
if (blobHashes === undefined && this.storage.extractBlobHashesFromRecords) {
blobHashes = this.storage.extractBlobHashesFromRecords(
await this.readGenerationRecords(gen)
)
}
}
await this.storage.removeRawPrefix(`${GENERATIONS_PREFIX}/${gen}`)
this.deltaCache.delete(gen)
if (this.historyBytesTotal !== null) {
this.historyBytesTotal -= delta.bytes
}
// AFTER the record-set is gone (over-count-only crash ordering):
// release its history references and reclaim any blob left with zero
// live AND zero history references — the system's one byte-reclaim point.
if (blobHashes && blobHashes.length > 0 && this.storage.releaseHistoryBlobReferences) {
await this.storage.releaseHistoryBlobReferences(blobHashes)
}
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)
// Packed-tier reclaim (D3): a packed generation's bytes live in a
// sealed segment — removeRawPrefix above was a no-op for it. Drop
// WHOLE segments now fully below the horizon; a partially-reclaimed
// segment keeps its bytes until the boundary passes it (the frozen
// partial-segments-wait rule; logical reclamation above still holds —
// the generations left committedRanges and asOf below the horizon
// throws regardless).
if (this.segments) {
await this.segments.dropSegmentsBelow(this.horizonGen + 1)
}
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<void> {
await this.withMutex(async () => {
this.deltaCache.clear()
// The running history-byte total describes the REPLACED store — drop it;
// the next historyBytes() re-seeds with one walk over the new state.
this.historyBytesTotal = null
// 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<void> {
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<void> {
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<R>(fn: () => Promise<R>): Promise<R> {
return this.withMutex(async () => {
await this.persistCounterUnlocked()
return fn()
})
}
/** Run `fn` exclusively, chained behind every prior exclusive section. */
private withMutex<R>(fn: () => Promise<R>): Promise<R> {
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/<N>/...`. 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<string, number[]>, 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 Drop the leading entries of an id's ascending chain that fall
* below `lo` (generations that slid out of the hot-tail window), deleting the
* chain entirely when it empties. Used by the O(d̄) window slide in
* {@link GenerationStore.extendChains}. Tolerates an already-trimmed front.
*/
function dropChainFront(chains: Map<string, number[]>, id: string, lo: number): void {
const chain = chains.get(id)
if (chain === undefined) return
let i = 0
while (i < chain.length && chain[i] < lo) i++
if (i === chain.length) chains.delete(id)
else if (i > 0) chain.splice(0, i)
}
/**
* @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
}