open-brainy/src/db/types.ts
David Snelling 292e7c0406
Some checks failed
CI / Node 24 (push) Successful in 12m21s
CI / Node 22 (push) Successful in 12m32s
CI / Integration + conformance (Node 22) (push) Failing after 13m32s
CI / Bun (latest) (push) Successful in 12m21s
fix(locks): live writers are never auto-evicted; evicted writers are fenced at every commit barrier
The production dev-store split-brain (two live writers alternating a store's
id-mapper between two internally-consistent truths), cured at all three of
its roots. (1) STALENESS REQUIRES PID-DEATH: the old rule evicted on
heartbeat age alone, so a >60s event-loop stall (debugger pause, GC, heavy
sync work) handed the lock to a second opener while the first kept writing;
a live process is now never auto-evicted — a wedged-but-alive holder is the
operator's call via {force:true}, and the heartbeat stays for observability.
(2) THE CLAIM IS ATOMIC: writeFile(wx)'s open→write→close left an empty-file
window a concurrent opener could read as torn, unlink a LIVE claim, and take
the lock; the claim is now tmp-write + hard-link — the lock appears with its
full contents in one step. (3) THE FENCE: every flush commit and transact
barrier verifies lock ownership first (one small read per window) — a
forced-out or lock-deleted writer fails typed (BRAINY_WRITER_FENCED) before
a single staged byte or manifest advance, instead of writing on unaware.

Pinned: live-with-ancient-heartbeat refuses typed; dead-PID self-clears
narrated; a forced-out writer's flush and transact both fence, advancing
nothing. Requested by a downstream team as single-writer guard or loud
lockout — this is both.
2026-08-17 16:26:41 -07:00

560 lines
25 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

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

/**
* @module db/types
* @description Shared types for the 8.0 generational-MVCC Db API: the
* declarative transaction-operation union consumed by `brain.transact()`,
* the options bags for `transact()` / `compactHistory()`, the persisted
* record shapes of the generational record layer, and the narrow storage
* contract (`GenerationStorage`) the record layer needs from an adapter.
*
* Persisted layout (all paths relative to the storage root):
*
* ```
* _system/generation.json — { generation, updatedAt } (atomic tmp+rename)
* _system/manifest.json — { version, generation, committedAt, horizon }
* (atomic tmp+rename — the rename IS the commit point)
* _system/tx-log.jsonl — one JSON line per committed transact:
* { generation, timestamp, meta? } (append-only)
* _generations/<N>/tx.json — the generation-N delta: touched noun/verb ids + meta
* _generations/<N>/prev/<id>.json — immutable before-image of <id> as of commit N
* ```
*
* Design note: the manifest holds the committed watermark + compaction
* horizon rather than a full `id → latest record generation` map. The id
* mapping is distributed across the per-generation `tx.json` deltas, which
* makes each commit O(ids touched) instead of O(all ids) — see
* `docs/ADR-001-generational-mvcc.md` for the full justification.
*/
import type { AddParams, UpdateParams, RelateParams, Entity, Relation } from '../types/brainy.types.js'
// ============================================================================
// Transaction operations (brain.transact input)
// ============================================================================
/**
* @description Create an entity. Carries the same parameters as
* `brain.add()`; supply `id` to choose the entity id (otherwise a UUID v4 is
* generated and returned in the operation's planning, exactly as `add()`
* does).
*/
export interface TxAddOperation<T = any> extends AddParams<T> {
/** Discriminator. */
op: 'add'
}
/**
* @description Update an entity. Carries the same parameters as
* `brain.update()` including per-entity CAS via `ifRev` — an `ifRev`
* conflict rejects the whole transaction before anything is applied.
*/
export interface TxUpdateOperation<T = any> extends UpdateParams<T> {
/** Discriminator. */
op: 'update'
}
/**
* @description Remove an entity (and, exactly like `brain.remove()`, every
* relationship where it is source or target — the cascade is part of the
* same atomic batch).
*/
export interface TxRemoveOperation {
/** Discriminator. */
op: 'remove'
/** Id of the entity to remove. */
id: string
}
/**
* @description Create a relationship. Carries the same parameters as
* `brain.relate()` (including `bidirectional`). Duplicate relationships
* (same from/to/type) are deduplicated exactly as `relate()` does — the
* operation becomes a no-op that resolves to the existing relationship id.
*/
export interface TxRelateOperation<T = any> extends RelateParams<T> {
/** Discriminator. */
op: 'relate'
}
/**
* @description Delete a relationship by id (mirror of `brain.unrelate()`).
*/
export interface TxUnrelateOperation {
/** Discriminator. */
op: 'unrelate'
/** Id of the relationship to delete. */
id: string
}
/**
* @description One declarative operation inside a `brain.transact()` batch.
* The batch executes atomically: either every operation applies and the
* store advances exactly one generation, or none apply and the store is
* byte-identical to its pre-transaction state.
*/
export type TxOperation<T = any> =
| TxAddOperation<T>
| TxUpdateOperation<T>
| TxRemoveOperation
| TxRelateOperation<T>
| TxUnrelateOperation
/**
* @description Options for `brain.transact()`.
*/
export interface TransactOptions {
/**
* Transaction metadata, reified Datomic-style: recorded in
* `_system/tx-log.jsonl` alongside the generation number and commit
* timestamp. Use it for audit fields (author, reason, request id) —
* it replaces commit messages from the pre-8.0 versioning API.
*/
meta?: Record<string, unknown>
/**
* Whole-store compare-and-swap. When provided, the transaction commits
* only if the store's current generation equals this value; otherwise it
* throws `GenerationConflictError` (with `expected`/`actual`) before any
* record is staged.
*/
ifAtGeneration?: number
/**
* Budget (ms) for the atomic apply phase. When omitted, the budget SCALES
* with the batch: `max(30 000, opCount × 2 000)` — production imports on
* network-attached disks measure ~2 s per operation, so a flat 30 s budget
* silently capped honest bulk work at ~15 operations. A tripped budget
* rolls the whole batch back and throws a `TransactionTimeoutError` naming
* the operation it stopped at, the batch size, and the elapsed/budget
* times. That error is retryable-with-latch, never hot-retry: its
* `retryable` field says a later attempt may succeed, its
* `hotRetryUnsafe` field says an immediate identical retry re-pays the
* full cost that just timed out — callers must latch and back off, never
* loop.
*/
timeoutMs?: number
}
/**
* @description Per-operation results of a committed transaction, in input
* order. `add` resolves to the created entity id, `relate` to the created
* (or deduplicated existing) relationship id; `update`/`remove`/`unrelate`
* resolve to the id they acted on.
*/
export interface TransactReceipt {
/** The generation this transaction committed as. */
generation: number
/** Commit timestamp (ms since epoch) recorded in the tx-log. */
timestamp: number
/** Resolved id per input operation, in input order. */
ids: string[]
}
// ============================================================================
// Compaction
// ============================================================================
/**
* @description Options for `brain.compactHistory()` — explicit retention
* **caps** (Model-B `retention` knob). Each supplied field is an upper bound:
* the oldest unpinned generations are reclaimed while **ANY** supplied cap is
* exceeded (predictable ops — the more you constrain, the more is reclaimed).
* Live `Db` pins are ALWAYS exempt. With no options, every generation not
* protected by a live pin is eligible.
*
* (8.0 rename: the pre-release `retainGenerations`/`retainMs` *floors* became
* `maxGenerations`/`maxAge` *caps* — the `max*` naming signals the semantics.)
*/
export interface CompactHistoryOptions {
/**
* Keep at most this many of the most recent committed generations — reclaim
* the oldest unpinned generations while the count exceeds it (`0` = reclaim
* everything not protected by a live pin).
*/
maxGenerations?: number
/**
* Keep only generations committed within the last `maxAge` milliseconds —
* reclaim unpinned generations older than the window.
*/
maxAge?: number
/**
* Keep total generational-history bytes at or below this — reclaim the oldest
* unpinned generations while the total exceeds it. Byte accounting is the sum
* of each surviving generation's serialized record set (`GenerationDelta.bytes`).
*/
maxBytes?: number
/**
* Stop reclaiming after this many milliseconds even if caps are still
* exceeded (8.9.0). Compaction is maintenance — a bounded pass keeps
* `close()` (and any explicit maintenance window) from stalling on a large
* backlog; the next pass resumes where this one stopped (reclamation is
* oldest-first, so an early stop is always a consistent prefix). Unset =
* run to completion.
*/
timeBudgetMs?: number
}
/**
* @description Result of `brain.compactHistory()`.
*/
export interface CompactHistoryResult {
/** Number of generation record-sets reclaimed by this call. */
removedGenerations: number
/**
* The new compaction horizon — the highest generation whose record-set has
* been reclaimed. Generations BELOW the horizon are no longer reachable via
* `asOf()` (their reads would need the reclaimed before-images); the
* horizon itself stays reachable, resolved from the record-sets above it.
*/
horizon: number
}
/**
* @description Result of `brain.historyStats()` — the read-only generational
* history footprint, for fleet audits and ops doors. A pool operator runs this
* per brain to size retention exposure (how much MVCC history each brain
* carries and under which policy) without touching any data.
*/
export interface HistoryStats {
/** Committed generation record-sets currently on disk. */
generations: number
/** Total on-disk history bytes across those record-sets. */
bytes: number
/** Oldest committed generation still on disk (null when history is empty). */
oldestGeneration: number | null
/** Newest committed generation (null when history is empty). */
newestGeneration: number | null
/** Commit timestamp (ms) of the oldest on-disk generation. */
oldestTimestamp: number | null
/** Commit timestamp (ms) of the newest on-disk generation. */
newestTimestamp: number | null
/** Compaction horizon — generations below it were reclaimed. */
horizon: number
/** The effective retention mode this brain runs under. */
retentionMode: 'all' | 'adaptive' | 'explicit'
/**
* The adaptive byte budget in force (coordinator-driven or the local
* free-memory probe); null under 'all' or explicit caps.
*/
effectiveBudgetBytes: number | null
}
// ============================================================================
// Db surfaces
// ============================================================================
/**
* @description Result of `db.since(priorDb)` — the ids touched by committed
* generations in the interval `(priorDb.generation, db.generation]`. Under
* Model-B EVERY write is its own generation, so single-operation writes
* (`add`/`update`/`remove`/`relate` outside `transact()`) ARE reflected here,
* exactly like `transact()` commits.
*/
export interface ChangedIds {
/** Entity (noun) ids touched in the interval, sorted. */
nouns: string[]
/** Relationship (verb) ids touched in the interval, sorted. */
verbs: string[]
/** Exclusive lower bound of the interval. */
fromGeneration: number
/** Inclusive upper bound of the interval. */
toGeneration: number
}
/**
* @description The result of {@link Brainy.diff}: ids that came into existence,
* went away, or changed value between two states, each split by kind. Unlike the
* raw touched-id set of {@link ChangedIds}, `diff` EARNS its name — it resolves
* each touched id at BOTH endpoints and classifies by existence (added/removed)
* and stable value comparison (modified). An id touched within the interval but
* whose endpoint states are identical (touched-then-reverted) appears in NONE of
* the three buckets. Orientation is `a → b`: `added` exists at `b` not `a`.
*/
export interface DiffResult {
/** Ids absent at `a`, present at `b`. */
added: { nouns: string[]; verbs: string[] }
/** Ids present at `a`, absent at `b`. */
removed: { nouns: string[]; verbs: string[] }
/** Ids present at both endpoints with a changed stored value. */
modified: { nouns: string[]; verbs: string[] }
/** The resolved generation of endpoint `a`. */
fromGeneration: number
/** The resolved generation of endpoint `b`. */
toGeneration: number
}
/**
* @description One version of a single entity or relationship in its
* {@link EntityHistory} — the materialized state at a committed generation that
* changed it. `value` is `null` when the id did not exist at that version
* (created-after / removed). Granularity is per-write (every `transact()` AND
* single-op write is its own generation — Model-B).
*/
export interface HistoryVersion<T = any> {
/** The committed generation that produced this version. */
generation: number
/** Commit timestamp of that generation (ms since epoch). */
timestamp: number
/** Materialized state at this version, or `null` if the id did not exist then. */
value: Entity<T> | Relation<T> | null
}
/**
* @description The result of {@link Brainy.history}: every distinct version of
* ONE id over a generation range, oldest → newest. Consecutive identical states
* are collapsed (one entry per distinct state). `kind` is auto-detected from the
* id's presence in the noun vs verb spaces.
*/
export interface EntityHistory<T = any> {
/** The id whose history this is. */
id: string
/** Whether `id` is an entity (`'noun'`) or relationship (`'verb'`). */
kind: 'noun' | 'verb'
/** Distinct versions, oldest first. */
versions: HistoryVersion<T>[]
}
// ============================================================================
// Persisted record shapes
// ============================================================================
/**
* @description An immutable before-image of one entity or relationship,
* persisted under `_generations/<N>/prev/<id>.json` — the state of the id
* immediately before commit `N` touched it. `metadata`/`vector` hold the
* *raw stored objects* (exact bytes that live at the entity's canonical
* storage paths); `null` means the corresponding file did not exist (for
* before-images of created ids, both are `null`). Before-images are both
* the crash-recovery undo log and the point-in-time read source: the state
* of id X at pinned generation G is the before-image of the first committed
* generation after G that touched X.
*/
export interface GenerationRecord {
/** Whether this record images an entity (noun) or a relationship (verb). */
kind: 'noun' | 'verb'
/** Raw stored metadata object, or `null` if the metadata file was absent. */
metadata: unknown | null
/** Raw stored vector object, or `null` if the vector file was absent. */
vector: unknown | null
}
/**
* @description The per-generation delta persisted at
* `_generations/<N>/tx.json`. For a `transact()` generation it is written and
* fsynced *before* any operation of the transaction executes, so crash
* recovery always knows exactly which ids to restore from before-images. For a
* single-operation (Model-B) generation it is written at group-commit flush
* time with `groupCommit: true` — see that flag.
*/
export interface GenerationDelta {
/** The generation this delta belongs to. */
generation: number
/** Commit timestamp (ms since epoch). */
timestamp: number
/** Transaction metadata (mirrored into the tx-log on commit). */
meta?: Record<string, unknown>
/** Entity ids touched by this generation. */
nouns: string[]
/** Relationship ids touched by this generation. */
verbs: string[]
/**
* Content-blob hashes referenced by this generation's before-image records
* (a MULTISET — one entry per referencing record occurrence), captured at
* persist time so compaction can release the exact history references it
* reclaims without re-reading the records. Always present (possibly empty)
* on deltas written under the temporal-blob contract; absent only on
* pre-contract generations, for which compaction falls back to reading the
* record-set itself.
*/
blobHashes?: string[]
/**
* `true` for a Model-B single-operation generation persisted by the async
* group-commit flush (`GenerationStore.flushPendingSingleOps`). It marks the
* **drop-without-restore** crash-recovery contract: the live write already
* mutated canonical storage and was acknowledged to the caller BEFORE the
* flush, so an uncommitted (crashed-mid-flush) generation with this flag set
* must be DISCARDED — its before-images must NEVER be restored, or recovery
* would silently revert acknowledged user writes. Absent/`false` on a
* `transact()` generation, whose before-images ARE the durable undo log and
* are restored on rollback (the before-image + execute are one atomic unit).
*/
groupCommit?: boolean
/**
* Serialized byte size of this generation's record set (the `prev/<id>.json`
* before-images plus this delta), recorded at write time so `maxBytes` /
* adaptive retention can bound total history without a storage size API. Read
* back lazily by `GenerationStore.historyBytes()`. Absent on generations
* written before byte accounting existed (treated as 0).
*/
bytes?: number
}
/**
* @description Shape of `_system/manifest.json`. Replaced via atomic
* tmp+rename on every commit — the rename is the commit point: a generation
* directory is committed if and only if its number is ≤ `generation`.
*/
export interface GenerationManifest {
/** Schema version of the manifest file. */
version: 1
/** Committed-transaction watermark. */
generation: number
/** ISO timestamp of the most recent commit (or compaction). */
committedAt: string
/** Compaction horizon — record-sets for generations ≤ this are reclaimed. */
horizon: number
}
/**
* @description One line of `_system/tx-log.jsonl`.
*/
export interface TxLogEntry {
/** The committed generation. */
generation: number
/** Commit timestamp (ms since epoch). */
timestamp: number
/** Transaction metadata, when supplied to `transact()`. */
meta?: Record<string, unknown>
/**
* WHO committed. Absent = a user write (every pre-existing consumer's
* reading stays exact). Engine-originated commits stamp themselves —
* `'system:embed-landing'` (the deferred vector landing),
* `'system:adoption-backfill'` (baseline re-commits), `'system:reconcile'`
* (the attested per-id divergence door) — so activity feeds can filter on
* fact instead of collapsing near-in-time entries (a consumer refused that
* heuristic as a quiet loss, correctly; this field is the honest cure).
* The same stamp rides the commit fact's meta, so log and tx-log agree.
*/
origin?: string
}
// ============================================================================
// Storage contract
// ============================================================================
/**
* @description The narrow storage contract the generational record layer
* needs. `BaseStorage` implements it structurally (both shipped adapters —
* filesystem and memory — inherit the implementation), so any
* `StorageAdapter` produced by the 8.0 storage factory satisfies it.
*
* Raw-object methods bypass branch scoping and operate on storage-root-
* relative paths (the record layer owns `_system/` + `_generations/`);
* entity-raw methods read/write the canonical entity files (branch-scoped,
* write-cache coherent) so before/after images capture exactly the bytes
* the live read paths see.
*/
export interface GenerationStorage {
/** Read a raw object at a storage-root-relative path (`null` if absent). */
readRawObject(path: string): Promise<any | null>
/** Write a raw object at a storage-root-relative path (atomic tmp+rename on disk). */
writeRawObject(path: string, data: any): Promise<void>
/** Delete a raw object (no-op if absent). */
deleteRawObject(path: string): Promise<void>
/** List raw object paths under a prefix (normalized, `.gz`-stripped). */
listRawObjects(prefix: string): Promise<string[]>
/** Remove every object under a prefix (and the directory itself on disk). */
removeRawPrefix(prefix: string): Promise<void>
/** Durability barrier: fsync the given object paths (no-op in memory). */
syncRawObjects(paths: string[]): Promise<void>
/**
* OPTIONAL transaction durability barrier. `commitTransaction` calls
* {@link beginWriteBarrier} immediately before running the planned operations
* and {@link flushWriteBarrier} immediately after — BEFORE the generation
* counter and manifest are advanced. An adapter whose canonical writes are
* not synchronously durable (the filesystem adapter's tmp+rename lands in the
* page cache) MUST implement these so a transaction reported "committed" is
* durable before its generation stamp: otherwise a hard kill can leave the
* fsync'd counter ahead of the still-buffered entity bytes (phantom progress
* for any generation-based consumer). Adapters whose writes are already
* durable per-call (cloud object PUT) may leave these undefined — the store
* treats them as no-ops. `beginWriteBarrier` also resets any tracking left by
* an aborted prior transaction.
*/
beginWriteBarrier?(): void
/** @see beginWriteBarrier — fsync every canonical write since begin. */
flushWriteBarrier?(): Promise<void>
/**
* OPTIONAL fold-checkpoint durability barrier: make the listed entities'
* CANONICAL live objects durable — fsync each present metadata/vector file
* AND the parent directory entry of each absent one (so a delete is as
* durable as a write). The generation store may only advance the fold
* checkpoint (`_system/fold-checkpoint.json`) after this resolves; the
* checkpoint bounds crash recovery's log fold to `(checkpoint, head]`.
* Adapters whose writes are durable per-call may leave this undefined —
* the store then treats canonical durability as immediate.
*/
syncEntityCanonical?(nouns: string[], verbs: string[]): Promise<void>
/**
* OPTIONAL writer fence: throw `BRAINY_WRITER_FENCED` when this instance
* no longer owns the store's writer lock (an operator force-takeover or a
* removed lock file). Called at every flush commit and transact barrier —
* one small read per commit window — so an evicted writer fails loudly on
* its next commit instead of split-braining the store. Adapters without a
* cross-process lock model omit it.
*/
assertWriterFenceHeld?(): Promise<void>
/** Read an entity's raw stored metadata+vector objects. */
readNounRaw(id: string): Promise<{ metadata: any | null; vector: any | null }>
/** Restore an entity's raw stored objects (`null` part ⇒ delete that file). */
writeNounRaw(id: string, record: { metadata: any | null; vector: any | null }): Promise<void>
/** Read a relationship's raw stored metadata+vector objects. */
readVerbRaw(id: string): Promise<{ metadata: any | null; vector: any | null }>
/** Restore a relationship's raw stored objects (`null` part ⇒ delete that file). */
writeVerbRaw(id: string, record: { metadata: any | null; vector: any | null }): Promise<void>
/** Append one line to `_system/tx-log.jsonl`. */
appendTxLogLine(line: string): Promise<void>
/** Read all lines of `_system/tx-log.jsonl` (empty array if absent). */
readTxLogLines(): Promise<string[]>
/**
* OPTIONAL binary raw-byte primitives — the substrate for the generation
* fact log's append-only CRC-framed segments. Feature-detected: a storage
* layer that omits them hosts no fact log (dual-write is skipped; readers
* fall back to canonical enumeration). Paths are used VERBATIM (no
* suffixing). Append durability rides `syncRawObjects` at the commit
* barrier, exactly like the staged history files.
*/
appendRawBytes?(path: string, bytes: Uint8Array): Promise<void>
/** Read a raw binary file whole; absent → null; a real fault throws. */
readRawBytes?(path: string): Promise<Uint8Array | null>
/** Replace a raw binary file atomically (tmp → fsync → rename). */
writeRawBytes?(path: string, bytes: Uint8Array): Promise<void>
/** Byte size of a raw binary file, or null when absent. */
rawByteSize?(path: string): Promise<number | null>
/**
* OPTIONAL temporal-blob contract (implemented by blob-aware storage; the
* generation store treats the hashes as opaque strings). Extract the
* content-blob hashes a record-set references — a pure multiset extraction,
* no side effects.
*/
extractBlobHashesFromRecords?(records: GenerationRecord[]): string[]
/**
* OPTIONAL: record history references for the given hashes (one increment
* per occurrence). Called BEFORE the referencing record-set is persisted —
* a crash between the two can only over-count (a leak the scrub repairs),
* never under-count (which would risk reclaiming bytes history still needs).
*/
recordHistoryBlobReferences?(hashes: string[]): Promise<void>
/**
* OPTIONAL: release history references for the given hashes (one decrement
* per occurrence) and physically reclaim any hash left with zero live AND
* zero history references. Called AFTER the referencing record-set is
* deleted by compaction — same over-count-only crash ordering.
*/
releaseHistoryBlobReferences?(hashes: string[]): Promise<void>
/**
* Register the generation-bump hook invoked on every entity-visible
* single-operation write (see `BaseStorage.setGenerationBumpHook`).
*/
setGenerationBumpHook(hook: (() => void) | undefined): void
/** Snapshot the entire store to a directory (hard-link farm on disk). */
snapshotToDirectory(targetPath: string): Promise<void>
/** Replace the entire store's contents from a snapshot directory. */
restoreFromDirectory(sourcePath: string): Promise<void>
}