flush() is durability work: it must cost what the current window's deltas cost, never what the history backlog costs. Under adaptive retention the byte budget derives from free memory, so bulk-load pressure shrank the budget exactly at peak write volume and flush paid actual reclaim inline — a production deployment measured single writes blocked 25-191s behind reclaim-on-flush. - flush() no longer calls autoCompactHistory(); close() is THE auto-compaction site (already ran there; now alone). - Every auto pass is time-bounded (CLOSE_COMPACTION_BUDGET_MS = 5s): reclamation is oldest-first, so an early stop is a consistent prefix and the next pass resumes. Explicit compactHistory() gains an optional timeBudgetMs for caller-chosen maintenance windows. - Documented trade stated where operators read: a long-lived writer that never closes accumulates history until its next explicit compactHistory() — predictable writes, explicit maintenance. Pins: flush-never-reclaims + close-reclaims-durably (db-mvcc), bounded pass stops-then-resumes as a consistent prefix (generationStore unit).
2530 lines
114 KiB
TypeScript
2530 lines
114 KiB
TypeScript
/**
|
||
* @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'
|
||
|
||
/**
|
||
* 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
|
||
|
||
/**
|
||
* 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
|
||
}
|
||
|
||
// 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.
|
||
*/
|
||
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 {
|
||
return []
|
||
}
|
||
const records: GenerationRecord[] = []
|
||
for (const p of paths) {
|
||
const record = (await this.storage.readRawObject(p)) as GenerationRecord | null
|
||
if (record) records.push(record)
|
||
}
|
||
return records
|
||
}
|
||
|
||
/**
|
||
* @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
|
||
}
|
||
return (await this.storage.readRawObject(
|
||
`${GENERATIONS_PREFIX}/${gen}/prev/${id}.json`
|
||
)) as GenerationRecord | 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) {
|
||
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.
|
||
*/
|
||
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)
|
||
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
|
||
}
|