diff --git a/src/brainy.ts b/src/brainy.ts index ef48c216..f79ca562 100644 --- a/src/brainy.ts +++ b/src/brainy.ts @@ -167,6 +167,7 @@ import { type ImportResult } from './db/portableGraph.js' import { GenerationStore, type CommitBeforeImages } from './db/generationStore.js' +import type { FactScanHandle } from './db/factLog.js' import { ChangeFeed, type BrainyChangeEvent, @@ -973,6 +974,24 @@ export class Brainy implements BrainyInterface { readOnly: this.config.mode === 'reader' }) + // The generation fact log is CANONICAL state, not a derived index — no + // sweeper, GC, or blob-lifecycle path may ever delete under it. Declare + // its namespace as a protected family (rebuildable: false — a lost fact + // segment is NOT reconstructable) so the storage layer REFUSES such + // deletes; refusal beats trust. Feature-detected + idempotent per name. + if ( + this.config.mode !== 'reader' && + this.generationStore.getFactLog() && + typeof this.storage.registerDerivedFamily === 'function' + ) { + await this.storage.registerDerivedFamily({ + name: 'generation-facts', + members: ['_generations/facts/'], + namespace: true, + rebuildable: false + }) + } + // 8.0 ⇄ native-provider version handshake: load the on-disk brain-format // marker (`_system/brain-format.json`) into an in-memory field NOW — // after the store-open phase, but BEFORE any derived index or native @@ -7434,6 +7453,59 @@ export class Brainy implements BrainyInterface { await this.removeMigrationBackupSafe() } + /** + * @description Open a sequential scan over the generation FACT LOG — the + * append-only record of every committed generation as an AFTER-IMAGE fact + * (what each touched entity/relationship became, or a body-less tombstone + * for a removal). The scan is the streaming substrate for index heals and + * incremental replays: one sequential read in commit order replaces a + * per-entity directory walk. The handle carries heal-narration telemetry + * (`headGeneration` / `segmentCount` / `approxFactCount` up front; ordered + * batches each stamped with their generation range, byte size, and segment; + * a `summary()` cross-check after iteration). A detected gap or damaged + * segment ABORTS the scan loudly — never a silent skip. + * + * Returns `null` when this store hosts no fact log: the storage adapter + * lacks the binary append primitives, or the brain predates the fact log + * (its history began before dual-write — facts exist only from the first + * write after upgrade; callers fall back to the canonical enumeration walk). + * + * @param options - `fromGeneration`/`toGeneration` bound the scan (inclusive + * both ends); `kinds` filters ops to `'noun'`/`'verb'`; `batchSize` caps + * facts per yielded batch (default 256). + * @returns The scan handle, or `null` when no fact log exists. + * @example + * const scan = brain.scanFacts({ fromGeneration: 1 }) + * if (scan) { + * for await (const batch of scan.batches()) { + * for (const fact of batch.facts) { + * // fact.ops: [{ kind, id, record | null (tombstone) }, ...] + * } + * } + * } + */ + scanFacts(options?: { + fromGeneration?: number + toGeneration?: number + kinds?: Array<'noun' | 'verb'> + batchSize?: number + }): FactScanHandle | null { + const factLog = this.generationStore?.getFactLog() + return factLog ? factLog.scanFacts(options) : null + } + + /** + * @description The immutable, sealed fact-log segment files covering + * `fromGeneration` — the zero-copy handoff for consumers that map segment + * files directly instead of streaming {@link scanFacts} batches. The + * append-mutable TAIL segment is deliberately excluded (read it via + * `scanFacts`). Paths are storage-root-relative. Empty when no fact log + * exists or nothing is sealed yet. + */ + factSegmentPaths(options?: { fromGeneration?: number }): string[] { + return this.generationStore?.getFactLog()?.segmentPaths(options) ?? [] + } + /** * @description Read the reified transaction log — one entry per committed * generation, carrying the committed generation, the commit timestamp, and diff --git a/src/coreTypes.ts b/src/coreTypes.ts index b20d128d..dc4b3860 100644 --- a/src/coreTypes.ts +++ b/src/coreTypes.ts @@ -1127,6 +1127,38 @@ export interface StorageAdapter { */ listDerivedFamilies?(): Promise + /** + * @description OPTIONAL binary raw-byte primitives — the substrate for + * append-only log-structured files (the generation fact log's CRC-framed + * segments). Feature-detected: an adapter that omits them simply hosts no + * fact log (readers fall back to canonical enumeration). Paths are + * storage-root-relative and used VERBATIM (no `.gz`/`.bin` suffixing — + * unlike the JSON object and blob primitives). + * + * Append to a raw binary file, creating it (and parent directories) when + * absent. Append durability is the CALLER's job via `syncRawObjects` — + * matching the commit protocol, which batches fsyncs at its barrier. + */ + appendRawBytes?(path: string, bytes: Uint8Array): Promise + + /** + * Read a raw binary file whole. Absent → `null`; a real IO fault throws + * (never masked as absence). + */ + readRawBytes?(path: string): Promise + + /** + * Replace a raw binary file atomically (write-new → fsync → rename) — the + * reconcile primitive (e.g. truncating a fact-log tail back to committed + * truth after a crash). + */ + writeRawBytes?(path: string, bytes: Uint8Array): Promise + + /** + * Byte size of a raw binary file, or `null` when absent. + */ + rawByteSize?(path: string): Promise + /** * Save statistics data * @param statistics The statistics data to save diff --git a/src/db/factLog.ts b/src/db/factLog.ts new file mode 100644 index 00000000..b43e52d0 --- /dev/null +++ b/src/db/factLog.ts @@ -0,0 +1,670 @@ +/** + * @module db/factLog + * @description The generation FACT LOG — an append-only, CRC-framed record of + * every committed generation as an AFTER-IMAGE "fact": what each touched + * entity/relationship BECAME (or a body-less tombstone when it was removed). + * This is the dual-write half of the log-canonical transition: today the + * before-image history + canonical tree remain authoritative; the fact log is + * appended at the same commit points and reconciled to committed truth at + * open, so consumers (index heals, replays, scans) can read one sequential, + * self-verifying stream instead of walking the entity tree. + * + * ## Wire format (frozen; additive-only within a major) + * + * Fact (msgpack, POSITIONAL array — the segment header's formatVersion + * governs the schema): + * + * fact := [ generation:u64, timestamp:u64, ops, meta|nil, blobHashes|nil ] + * op := [ kind:u8 (0=noun, 1=verb), id:bin16 (raw uuid bytes), + * record:[metaLeg, vecLeg] | nil ] // nil = TOMBSTONE + * + * Segment file (`_generations/facts/seg-.bfl`): + * + * header := magic "BFACTS\0\0" (8B) | formatVersion:u32 LE | + * firstGeneration:u64 LE | reserved 12B (ZEROED, verified) + * frame := length:u32 LE | crc32c:u32 LE (of payload) | payload + * + * A fact is never split across segments; a torn tail (length overrun or CRC + * mismatch) terminates that segment's scan — everything before it is intact. + * Zero-padded names make lexicographic order == generation order. + * + * ## Invariant + * + * After {@link FactLog.open}, the log contains EXACTLY the committed prefix: + * facts are appended BEFORE the commit point (inside the same durability + * window), so a crash can only leave the log AHEAD of committed truth — open + * truncates any fact beyond the committed generation. Absent generation = + * never committed; present = committed. A scan can never see an uncommitted + * fact. + * + * The manifest (`_generations/facts/manifest.json`, JSON — forensics stay + * terminal-readable) is the single source of truth for the segment SET; + * rotation flips it atomically (write-new → fsync → rename) BEFORE the new + * tail's first byte exists, so no segment file is ever unaccounted for. + */ +import { encode as defaultEncode, decode as defaultDecode } from '@msgpack/msgpack' +import { crc32c } from '../utils/crc32c.js' +import { prodLog } from '../utils/logger.js' + +// Swappable msgpack implementation — defaults to the JS codec; a native +// provider (registered via the plugin registry's 'msgpack' key) may replace +// it. Byte-compatibility is the contract (positional arrays, bin16 ids). +let msgpackEncode: (value: unknown) => Uint8Array = defaultEncode +let msgpackDecode: (bytes: Uint8Array) => unknown = defaultDecode + +/** Replace the msgpack encode/decode implementation at runtime. */ +export function setFactCodec(impl: { + encode: (value: unknown) => Uint8Array + decode: (bytes: Uint8Array) => unknown +}): void { + msgpackEncode = impl.encode + msgpackDecode = impl.decode +} + +/** Storage-root-relative home of the fact log. */ +export const FACTS_PREFIX = '_generations/facts' +/** The facts manifest path (JSON). */ +export const FACTS_MANIFEST_PATH = `${FACTS_PREFIX}/manifest.json` +/** Current segment format version (header field; additive-only within a major). */ +export const FACTS_FORMAT_VERSION = 1 +/** Rotation threshold: seal the tail segment once it exceeds this many bytes. */ +const SEGMENT_ROTATE_BYTES = 8 * 1024 * 1024 +/** Segment header: magic(8) + formatVersion(4) + firstGeneration(8) + reserved(12). */ +const HEADER_BYTES = 32 +const MAGIC = new Uint8Array([0x42, 0x46, 0x41, 0x43, 0x54, 0x53, 0x00, 0x00]) // "BFACTS\0\0" +/** Frame prefix: length(4) + crc32c(4). */ +const FRAME_PREFIX_BYTES = 8 + +/** One write inside a fact: what the id became (or a tombstone). */ +export interface FactOp { + kind: 'noun' | 'verb' + id: string + /** The AFTER-IMAGE legs, or `null` for a tombstone (the id was removed). */ + record: { metadata: unknown | null; vector: unknown | null } | null +} + +/** One committed generation, as scanned back out of the log. */ +export interface CommitFact { + generation: number + timestamp: number + ops: FactOp[] + meta?: Record + blobHashes?: string[] +} + +/** The telemetry a scan batch carries (frozen shape). */ +export interface FactScanBatch { + facts: CommitFact[] + firstGeneration: number + lastGeneration: number + factCount: number + byteSize: number + segmentId: string +} + +/** The telemetry a scan OPEN returns (frozen shape). */ +export interface FactScanHandle { + headGeneration: number + segmentCount: number + approxFactCount: number + /** Ordered batches; a detected gap aborts LOUDLY, never a silent skip. */ + batches: () => AsyncGenerator + /** Close telemetry — the invariant cross-check, valid after iteration ends. */ + summary: () => { factsYielded: number; segmentsRead: number } +} + +/** Manifest entry for a sealed segment. */ +interface SegmentEntry { + file: string + firstGeneration: number + lastGeneration: number + facts: number + bytes: number +} + +/** The facts manifest (JSON on disk). */ +interface FactsManifest { + formatVersion: number + segments: SegmentEntry[] + /** The append target. Its true content is established by scanning (crash tolerance). */ + tailSegment: string | null + updatedAt: string +} + +/** The narrow byte-level storage surface the fact log rides. */ +export interface FactLogStorage { + appendRawBytes(path: string, bytes: Uint8Array): Promise + readRawBytes(path: string): Promise + writeRawBytes(path: string, bytes: Uint8Array): Promise + rawByteSize(path: string): Promise + readRawObject(path: string): Promise + writeRawObject(path: string, data: any): Promise + syncRawObjects(paths: string[]): Promise + deleteRawObject(path: string): Promise +} + +/** True when the storage adapter exposes every primitive the fact log needs. */ +export function storageSupportsFactLog(storage: unknown): storage is FactLogStorage { + const s = storage as Record + return ( + typeof s.appendRawBytes === 'function' && + typeof s.readRawBytes === 'function' && + typeof s.writeRawBytes === 'function' && + typeof s.rawByteSize === 'function' + ) +} + +/** uuid string → 16 raw bytes (bin16 on the wire). */ +function uuidToBytes(id: string): Uint8Array { + const hex = id.replace(/-/g, '') + if (hex.length !== 32) { + // Non-uuid ids (legacy/natural keys) ride as UTF-8 with a length prefix + // marker impossible for uuids: we refuse instead — the write API has + // guaranteed uuid ids since 8.0, so anything else is a corruption signal. + throw new Error(`fact log: id is not a uuid: ${id}`) + } + const bytes = new Uint8Array(16) + for (let i = 0; i < 16; i++) { + bytes[i] = parseInt(hex.slice(i * 2, i * 2 + 2), 16) + } + return bytes +} + +/** 16 raw bytes → canonical lowercase uuid string. */ +function bytesToUuid(bytes: Uint8Array): string { + let hex = '' + for (let i = 0; i < 16; i++) hex += bytes[i].toString(16).padStart(2, '0') + return `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}` +} + +/** Zero-padded segment filename: lexicographic order == generation order. */ +function segmentFileName(firstGeneration: number): string { + return `seg-${String(firstGeneration).padStart(20, '0')}.bfl` +} + +/** Build a segment header. Reserved bytes are ZEROED (and verified on open). */ +function buildHeader(firstGeneration: number): Uint8Array { + const header = new Uint8Array(HEADER_BYTES) + header.set(MAGIC, 0) + const view = new DataView(header.buffer) + view.setUint32(8, FACTS_FORMAT_VERSION, true) + view.setBigUint64(12, BigInt(firstGeneration), true) + // bytes 20..31 stay zero (reserved) + return header +} + +/** Encode one fact into a framed record (length + crc32c + msgpack payload). */ +function encodeFrame(fact: CommitFact): Uint8Array { + const payload = msgpackEncode([ + fact.generation, + fact.timestamp, + fact.ops.map((op) => [ + op.kind === 'noun' ? 0 : 1, + uuidToBytes(op.id), + op.record === null ? null : [op.record.metadata, op.record.vector] + ]), + fact.meta ?? null, + fact.blobHashes && fact.blobHashes.length > 0 ? fact.blobHashes : null + ]) + const frame = new Uint8Array(FRAME_PREFIX_BYTES + payload.length) + const view = new DataView(frame.buffer) + view.setUint32(0, payload.length, true) + view.setUint32(4, crc32c(payload), true) + frame.set(payload, FRAME_PREFIX_BYTES) + return frame +} + +/** Decode one msgpack payload back into a CommitFact. */ +function decodeFact(payload: Uint8Array): CommitFact { + const raw = msgpackDecode(payload) as unknown[] + const [generation, timestamp, ops, meta, blobHashes] = raw as [ + number, + number, + Array<[number, Uint8Array, [unknown, unknown] | null]>, + Record | null, + string[] | null + ] + return { + generation: Number(generation), + timestamp: Number(timestamp), + ops: ops.map(([kind, idBytes, record]) => ({ + kind: kind === 0 ? ('noun' as const) : ('verb' as const), + id: bytesToUuid(idBytes), + record: record === null ? null : { metadata: record[0] ?? null, vector: record[1] ?? null } + })), + ...(meta ? { meta } : {}), + ...(blobHashes && blobHashes.length > 0 ? { blobHashes } : {}) + } +} + +/** + * Parse a segment's bytes: verify the header, then walk frames until the end + * or a torn tail (length overrun / CRC mismatch), which terminates the walk — + * everything before it is intact. Returns the decoded facts plus the byte + * length of the VALID prefix (header + intact frames), which reconciliation + * uses to cut a torn tail without re-encoding. + */ +function parseSegment( + file: string, + bytes: Uint8Array +): { facts: CommitFact[]; validBytes: number } { + if (bytes.length < HEADER_BYTES) { + prodLog.warn(`[FactLog] segment ${file} shorter than its header — treating as empty`) + return { facts: [], validBytes: 0 } + } + for (let i = 0; i < MAGIC.length; i++) { + if (bytes[i] !== MAGIC[i]) { + throw new Error(`fact log: segment ${file} has a bad magic — not a fact segment`) + } + } + const view = new DataView(bytes.buffer, bytes.byteOffset, bytes.byteLength) + const version = view.getUint32(8, true) + if (version !== FACTS_FORMAT_VERSION) { + throw new Error( + `fact log: segment ${file} has formatVersion ${version}; this build reads ${FACTS_FORMAT_VERSION}` + ) + } + for (let i = 20; i < HEADER_BYTES; i++) { + if (bytes[i] !== 0) { + // Non-zero reserved bytes = a future format this build cannot verify. + throw new Error(`fact log: segment ${file} has non-zero reserved header bytes — unverifiable`) + } + } + + const facts: CommitFact[] = [] + let offset = HEADER_BYTES + while (offset + FRAME_PREFIX_BYTES <= bytes.length) { + const length = view.getUint32(offset, true) + const expectedCrc = view.getUint32(offset + 4, true) + const start = offset + FRAME_PREFIX_BYTES + const end = start + length + if (end > bytes.length) break // torn tail: frame length overruns the file + const payload = bytes.subarray(start, end) + if (crc32c(payload) !== expectedCrc) break // torn tail: payload CRC mismatch + facts.push(decodeFact(payload)) + offset = end + } + return { facts, validBytes: offset } +} + +/** + * The generation fact log. One instance per open store; every method assumes + * the single-writer discipline the generation store already enforces (calls + * arrive under its commit mutex). + */ +export class FactLog { + private readonly storage: FactLogStorage + private manifest: FactsManifest = { + formatVersion: FACTS_FORMAT_VERSION, + segments: [], + tailSegment: null, + updatedAt: new Date(0).toISOString() + } + /** Decoded facts of the TAIL segment (bounded by the rotation threshold). */ + private tailFacts: CommitFact[] = [] + /** Byte size of the tail segment file (valid prefix). */ + private tailBytes = 0 + /** Highest generation in the log (0 = empty). */ + private head = 0 + /** Segment paths appended since the last sync (the fsync batch). */ + private readonly dirtySegments = new Set() + + constructor(storage: FactLogStorage) { + this.storage = storage + } + + /** The highest committed generation the log holds (0 = empty). */ + headGeneration(): number { + return this.head + } + + /** + * Open the log and reconcile it to committed truth: read the manifest, + * establish the tail's intact content (torn-tail scan), then TRUNCATE any + * fact with `generation > committedGeneration` — those never committed (a + * crash between fact-append and the commit point). After open, the log is + * exactly the committed prefix. + */ + async open(committedGeneration: number): Promise { + const stored = (await this.storage.readRawObject(FACTS_MANIFEST_PATH)) as FactsManifest | null + if (stored && typeof stored === 'object' && Array.isArray(stored.segments)) { + if (stored.formatVersion !== FACTS_FORMAT_VERSION) { + throw new Error( + `fact log: manifest formatVersion ${stored.formatVersion}; this build reads ${FACTS_FORMAT_VERSION}` + ) + } + this.manifest = stored + } + + // Drop sealed segments that sit ENTIRELY beyond committed truth (a crash + // right after a rotation whose facts never committed), newest first. + while (this.manifest.segments.length > 0) { + const last = this.manifest.segments[this.manifest.segments.length - 1] + if (last.firstGeneration > committedGeneration) { + prodLog.warn( + `[FactLog] dropping sealed segment ${last.file} (generations ${last.firstGeneration}..` + + `${last.lastGeneration} never committed)` + ) + await this.storage.deleteRawObject(`${FACTS_PREFIX}/${last.file}`) + this.manifest.segments.pop() + await this.persistManifest() + } else if (last.lastGeneration > committedGeneration) { + // A sealed segment STRADDLING committed truth: cut it back. + await this.truncateSegmentTo(last.file, committedGeneration) + const cut = await this.reloadSegmentEntry(last.file) + this.manifest.segments[this.manifest.segments.length - 1] = cut + await this.persistManifest() + break + } else { + break + } + } + + // Establish the tail: scan its intact prefix, then truncate beyond + // committed truth (the common crash shape: buffered single-op facts whose + // counter never went durable). + if (this.manifest.tailSegment) { + const tailPath = `${FACTS_PREFIX}/${this.manifest.tailSegment}` + const bytes = await this.storage.readRawBytes(tailPath) + if (bytes === null) { + // Manifest named a tail whose first byte never landed — an empty tail. + this.tailFacts = [] + this.tailBytes = 0 + } else { + const { facts, validBytes } = parseSegment(this.manifest.tailSegment, bytes) + const kept = facts.filter((f) => f.generation <= committedGeneration) + if (kept.length !== facts.length || validBytes !== bytes.length) { + const dropped = facts.length - kept.length + if (dropped > 0) { + prodLog.warn( + `[FactLog] truncating ${dropped} uncommitted fact(s) beyond generation ` + + `${committedGeneration} from the tail (never committed)` + ) + } + await this.rewriteTail(kept) + } else { + this.tailFacts = facts + this.tailBytes = validBytes + } + } + } + + this.head = this.computeHead() + } + + /** + * Append one committed generation's fact. NOT durable until {@link sync} — + * the caller batches durability at its commit barrier (transact syncs in + * the same call; Model-B group-commit syncs at flush). + */ + async append(fact: CommitFact): Promise { + if (fact.generation <= this.head) { + throw new Error( + `fact log: non-monotonic append (generation ${fact.generation} ≤ head ${this.head})` + ) + } + if (this.manifest.tailSegment === null) { + await this.startTail(fact.generation) + } else if (this.tailBytes >= SEGMENT_ROTATE_BYTES) { + await this.rotate(fact.generation) + } + const frame = encodeFrame(fact) + const tailPath = `${FACTS_PREFIX}/${this.manifest.tailSegment}` + await this.storage.appendRawBytes(tailPath, frame) + this.tailFacts.push(fact) + this.tailBytes += frame.length + this.head = fact.generation + this.dirtySegments.add(tailPath) + } + + /** Fsync every segment appended since the last sync. */ + async sync(): Promise { + if (this.dirtySegments.size === 0) return + const paths = [...this.dirtySegments] + this.dirtySegments.clear() + await this.storage.syncRawObjects(paths) + } + + /** + * Open a scan over committed facts. The scan runs against a MANIFEST + * SNAPSHOT (sealed segments + the tail's decoded facts at open) — exactly- + * once per fact, inclusive bounds, stable under concurrent appends. Gaps + * abort LOUDLY: a missing generation inside a segment's declared range is + * corruption, never silently skipped. + */ + scanFacts(options?: { + fromGeneration?: number + toGeneration?: number + kinds?: Array<'noun' | 'verb'> + batchSize?: number + }): FactScanHandle { + const from = options?.fromGeneration ?? 1 + const to = options?.toGeneration ?? this.head + const kinds = options?.kinds + const batchSize = Math.max(1, options?.batchSize ?? 256) + + // Snapshot: the segment list + tail content as of NOW. + const segments = this.manifest.segments.filter( + (s) => s.lastGeneration >= from && s.firstGeneration <= to + ) + const tailSnapshot = this.tailFacts.filter((f) => f.generation >= from && f.generation <= to) + const tailId = this.manifest.tailSegment ?? 'tail' + const approxFactCount = + segments.reduce((sum, s) => sum + s.facts, 0) + tailSnapshot.length + + let factsYielded = 0 + let segmentsRead = 0 + const storage = this.storage + + async function* batches(this: void): AsyncGenerator { + let expectedNext = 0 // gap detection: generations are monotonic, not necessarily dense + const emit = (facts: CommitFact[], segmentId: string, byteSize: number): FactScanBatch => ({ + facts, + firstGeneration: facts[0].generation, + lastGeneration: facts[facts.length - 1].generation, + factCount: facts.length, + byteSize, + segmentId + }) + const filterOps = (fact: CommitFact): CommitFact => + kinds + ? { ...fact, ops: fact.ops.filter((op) => kinds.includes(op.kind)) } + : fact + + for (const entry of segments) { + const bytes = await storage.readRawBytes(`${FACTS_PREFIX}/${entry.file}`) + if (bytes === null) { + throw new Error( + `fact log: sealed segment ${entry.file} is MISSING — the log is damaged; aborting scan` + ) + } + const { facts } = parseSegment(entry.file, bytes) + segmentsRead++ + const inRange = facts.filter((f) => f.generation >= from && f.generation <= to) + for (const f of inRange) { + if (f.generation <= expectedNext - 1) { + throw new Error(`fact log: out-of-order fact ${f.generation} in ${entry.file} — aborting scan`) + } + expectedNext = f.generation + 1 + } + for (let i = 0; i < inRange.length; i += batchSize) { + const slice = inRange.slice(i, i + batchSize).map(filterOps) + if (slice.length === 0) continue + factsYielded += slice.length + yield emit(slice, entry.file, slice.reduce((n, f) => n + encodeFrame(f).length, 0)) + } + } + + if (tailSnapshot.length > 0) { + segmentsRead++ + for (const f of tailSnapshot) { + if (f.generation <= expectedNext - 1) { + throw new Error(`fact log: out-of-order fact ${f.generation} in the tail — aborting scan`) + } + expectedNext = f.generation + 1 + } + for (let i = 0; i < tailSnapshot.length; i += batchSize) { + const slice = tailSnapshot.slice(i, i + batchSize).map(filterOps) + factsYielded += slice.length + yield emit(slice, tailId, slice.reduce((n, f) => n + encodeFrame(f).length, 0)) + } + } + } + + return { + headGeneration: this.head, + segmentCount: segments.length + (tailSnapshot.length > 0 ? 1 : 0), + approxFactCount, + batches, + summary: () => ({ factsYielded, segmentsRead }) + } + } + + /** + * The mmap fast path (capability handoff): the immutable sealed segment + * files covering `fromGeneration`, in order. The TAIL is deliberately NOT + * included — it is append-mutable; consumers read it via {@link scanFacts}. + */ + segmentPaths(options?: { fromGeneration?: number }): string[] { + const from = options?.fromGeneration ?? 1 + return this.manifest.segments + .filter((s) => s.lastGeneration >= from) + .map((s) => `${FACTS_PREFIX}/${s.file}`) + } + + /** + * Drop every fact with `generation > keepThrough` — the in-session abort + * compensation: a transact appends its fact BEFORE the commit point, so a + * real (non-crash) abort after the append must take the fact back out. The + * dropped facts can only live in the TAIL (they were just appended); the + * rewrite is atomic and bounded by the rotation threshold. + */ + async dropAbove(keepThrough: number): Promise { + if (this.head <= keepThrough) return + const kept = this.tailFacts.filter((f) => f.generation <= keepThrough) + if (kept.length === this.tailFacts.length) { + throw new Error( + `fact log: dropAbove(${keepThrough}) found no droppable facts in the tail ` + + `(head ${this.head}) — the fact to drop was already sealed; the log needs reopen` + ) + } + await this.rewriteTail(kept) + this.head = this.computeHead() + } + + // -- internals ------------------------------------------------------------- + + private computeHead(): number { + if (this.tailFacts.length > 0) return this.tailFacts[this.tailFacts.length - 1].generation + const sealed = this.manifest.segments + if (sealed.length > 0) return sealed[sealed.length - 1].lastGeneration + return 0 + } + + /** Create the very first tail segment (manifest-first, then header bytes). */ + private async startTail(firstGeneration: number): Promise { + const file = segmentFileName(firstGeneration) + this.manifest.tailSegment = file + await this.persistManifest() + await this.storage.appendRawBytes(`${FACTS_PREFIX}/${file}`, buildHeader(firstGeneration)) + this.tailFacts = [] + this.tailBytes = HEADER_BYTES + } + + /** + * Seal the tail into the manifest and start a new one. Manifest-first: the + * flip both seals the old tail AND names the new one atomically, so no + * segment file ever exists unaccounted for. + */ + private async rotate(nextGeneration: number): Promise { + const sealedFile = this.manifest.tailSegment + if (!sealedFile) return + // Seal what the tail actually holds. + await this.sync() // sealed segments are always fully durable + const entry: SegmentEntry = { + file: sealedFile, + firstGeneration: this.tailFacts[0]?.generation ?? nextGeneration, + lastGeneration: this.tailFacts[this.tailFacts.length - 1]?.generation ?? nextGeneration - 1, + facts: this.tailFacts.length, + bytes: this.tailBytes + } + const newFile = segmentFileName(nextGeneration) + this.manifest.segments.push(entry) + this.manifest.tailSegment = newFile + await this.persistManifest() + await this.storage.appendRawBytes(`${FACTS_PREFIX}/${newFile}`, buildHeader(nextGeneration)) + this.tailFacts = [] + this.tailBytes = HEADER_BYTES + } + + /** Atomically persist the manifest (write-new → fsync → rename downstream). */ + private async persistManifest(): Promise { + this.manifest.updatedAt = new Date().toISOString() + await this.storage.writeRawObject(FACTS_MANIFEST_PATH, this.manifest) + await this.storage.syncRawObjects([FACTS_MANIFEST_PATH]) + } + + /** Rewrite the tail segment to hold exactly `facts` (atomic replace). */ + private async rewriteTail(facts: CommitFact[]): Promise { + const file = this.manifest.tailSegment + if (!file) return + const first = facts[0]?.generation ?? this.segmentFirstGenerationFromName(file) + const parts: Uint8Array[] = [buildHeader(first)] + for (const f of facts) parts.push(encodeFrame(f)) + const total = parts.reduce((n, p) => n + p.length, 0) + const merged = new Uint8Array(total) + let offset = 0 + for (const p of parts) { + merged.set(p, offset) + offset += p.length + } + await this.storage.writeRawBytes(`${FACTS_PREFIX}/${file}`, merged) + this.tailFacts = facts + this.tailBytes = total + } + + /** Cut a SEALED segment back to `committedGeneration` (atomic replace). */ + private async truncateSegmentTo(file: string, committedGeneration: number): Promise { + const path = `${FACTS_PREFIX}/${file}` + const bytes = await this.storage.readRawBytes(path) + if (bytes === null) return + const { facts } = parseSegment(file, bytes) + const kept = facts.filter((f) => f.generation <= committedGeneration) + prodLog.warn( + `[FactLog] truncating sealed segment ${file} to generation ${committedGeneration} ` + + `(${facts.length - kept.length} uncommitted fact(s) dropped)` + ) + const first = kept[0]?.generation ?? this.segmentFirstGenerationFromName(file) + const parts: Uint8Array[] = [buildHeader(first)] + for (const f of kept) parts.push(encodeFrame(f)) + const total = parts.reduce((n, p) => n + p.length, 0) + const merged = new Uint8Array(total) + let offset = 0 + for (const p of parts) { + merged.set(p, offset) + offset += p.length + } + await this.storage.writeRawBytes(path, merged) + } + + /** Re-derive a sealed segment's manifest entry from its actual bytes. */ + private async reloadSegmentEntry(file: string): Promise { + const bytes = await this.storage.readRawBytes(`${FACTS_PREFIX}/${file}`) + const { facts, validBytes } = bytes + ? parseSegment(file, bytes) + : { facts: [] as CommitFact[], validBytes: 0 } + return { + file, + firstGeneration: facts[0]?.generation ?? this.segmentFirstGenerationFromName(file), + lastGeneration: facts[facts.length - 1]?.generation ?? 0, + facts: facts.length, + bytes: validBytes + } + } + + /** Parse the zero-padded firstGeneration back out of a segment filename. */ + private segmentFirstGenerationFromName(file: string): number { + const match = /^seg-(\d{20})\.bfl$/.exec(file) + return match ? Number(match[1]) : 0 + } +} diff --git a/src/db/generationStore.ts b/src/db/generationStore.ts index d7a302c9..0d813f4d 100644 --- a/src/db/generationStore.ts +++ b/src/db/generationStore.ts @@ -45,6 +45,7 @@ import type { 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 @@ -121,6 +122,16 @@ export interface GenerationStoreOpenResult { 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). */ @@ -400,6 +411,19 @@ export class GenerationStore { 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) { @@ -600,6 +624,58 @@ export class GenerationStore { * @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 + }): Promise { + 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 @@ -733,6 +809,24 @@ export class GenerationStore { 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') @@ -800,6 +894,13 @@ export class GenerationStore { // 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 @@ -991,6 +1092,14 @@ export class GenerationStore { 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 ` + @@ -1020,6 +1129,17 @@ export class GenerationStore { 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 } }) @@ -1141,6 +1261,12 @@ export class GenerationStore { // 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 diff --git a/src/db/types.ts b/src/db/types.ts index 4f932e76..2bfd6ac2 100644 --- a/src/db/types.ts +++ b/src/db/types.ts @@ -423,6 +423,22 @@ export interface GenerationStorage { /** Read all lines of `_system/tx-log.jsonl` (empty array if absent). */ readTxLogLines(): Promise + /** + * 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 + /** Read a raw binary file whole; absent → null; a real fault throws. */ + readRawBytes?(path: string): Promise + /** Replace a raw binary file atomically (tmp → fsync → rename). */ + writeRawBytes?(path: string, bytes: Uint8Array): Promise + /** Byte size of a raw binary file, or null when absent. */ + rawByteSize?(path: string): Promise + /** * OPTIONAL temporal-blob contract (implemented by blob-aware storage; the * generation store treats the hashes as opaque strings). Extract the diff --git a/src/index.ts b/src/index.ts index b0a95fe1..b9135600 100644 --- a/src/index.ts +++ b/src/index.ts @@ -196,6 +196,14 @@ export type { HistoryVersion, EntityHistory } from './db/types.js' +// The generation fact log — sequential after-image scan surface +// (brain.scanFacts / brain.factSegmentPaths) for index heals and replays. +export type { + CommitFact, + FactOp, + FactScanBatch, + FactScanHandle +} from './db/factLog.js' // Optional provider capability for generation-aware native indexes export { isVersionedIndexProvider } from './plugin.js' export type { VersionedIndexProvider } from './plugin.js' diff --git a/src/storage/adapters/fileSystemStorage.ts b/src/storage/adapters/fileSystemStorage.ts index a9e3014c..4f246f67 100644 --- a/src/storage/adapters/fileSystemStorage.ts +++ b/src/storage/adapters/fileSystemStorage.ts @@ -693,6 +693,17 @@ export class FileSystemStorage extends BaseStorage { */ private static readonly SNAPSHOT_BYTE_COPY_DIRS = new Set(['_id_mapper']) + /** + * Nested path PREFIXES whose files are byte-copied into snapshots, not + * hard-linked — for append-in-place files below the top level. The + * generation fact log's tail segment is appended in place between rotations; + * a hard-linked tail would let post-snapshot appends reach through into the + * snapshot. (Sealed segments are immutable and would be link-safe, but the + * prefix rule keeps the discipline simple; segments are bounded by the + * rotation threshold, so the copy cost is small.) + */ + private static readonly SNAPSHOT_BYTE_COPY_PREFIXES: string[] = ['_generations/facts/'] + /** * Top-level directories excluded from snapshots: process-local lock state * (writer lock, flush-request RPC files) must never travel with the data, and @@ -870,6 +881,75 @@ export class FileSystemStorage extends BaseStorage { } } + // ========================================================================== + // Binary raw-byte primitives — the substrate for append-only log-structured + // files (the generation fact log's CRC-framed segments). Paths are used + // VERBATIM (no .gz/.bin suffixing). Append durability rides syncRawObjects + // at the commit barrier, like every other staged write. + // ========================================================================== + + /** + * Append bytes to a raw binary file, creating it (and parent directories) + * when absent. NOT fsync'd here — the caller batches durability via + * `syncRawObjects` at its commit barrier. + */ + public async appendRawBytes(rawPath: string, bytes: Uint8Array): Promise { + await this.ensureInitialized() + const fullPath = path.join(this.rootDir, rawPath) + await fs.promises.mkdir(path.dirname(fullPath), { recursive: true }) + await fs.promises.appendFile(fullPath, bytes) + } + + /** + * Read a raw binary file whole. Absent → `null`; a real IO fault throws — + * a present-but-unreadable log segment must never read as "no facts". + */ + public async readRawBytes(rawPath: string): Promise { + await this.ensureInitialized() + try { + const buf: Buffer = await fs.promises.readFile(path.join(this.rootDir, rawPath)) + return new Uint8Array(buf.buffer, buf.byteOffset, buf.byteLength) + } catch (error: any) { + if (isAbsentError(error)) return null + throw error + } + } + + /** + * Replace a raw binary file atomically: write-new → fsync → rename. The + * reconcile primitive (e.g. truncating a fact-log tail back to committed + * truth after a crash) — a crash mid-replace leaves either the old file or + * the new one, never a mix. + */ + public async writeRawBytes(rawPath: string, bytes: Uint8Array): Promise { + await this.ensureInitialized() + const fullPath = path.join(this.rootDir, rawPath) + await fs.promises.mkdir(path.dirname(fullPath), { recursive: true }) + const tmpPath = `${fullPath}.tmp.${Date.now()}.${Math.random().toString(36).slice(2)}` + const handle = await fs.promises.open(tmpPath, 'w') + try { + await handle.writeFile(bytes) + await handle.sync() + } finally { + await handle.close() + } + await fs.promises.rename(tmpPath, fullPath) + } + + /** + * Byte size of a raw binary file, or `null` when absent. + */ + public async rawByteSize(rawPath: string): Promise { + await this.ensureInitialized() + try { + const stat = await fs.promises.stat(path.join(this.rootDir, rawPath)) + return stat.size + } catch (error: any) { + if (isAbsentError(error)) return null + throw error + } + } + /** * Snapshot the entire store into `targetPath` as a hard-link farm * (Cassandra-style: instant, space-shared). Safe because every data file @@ -919,7 +999,8 @@ export class FileSystemStorage extends BaseStorage { // msync/truncate reach through into the snapshot. if ( FileSystemStorage.SNAPSHOT_BYTE_COPY_PATHS.has(normalized) || - FileSystemStorage.SNAPSHOT_BYTE_COPY_DIRS.has(normalized.split('/')[0]) + FileSystemStorage.SNAPSHOT_BYTE_COPY_DIRS.has(normalized.split('/')[0]) || + FileSystemStorage.SNAPSHOT_BYTE_COPY_PREFIXES.some((p) => normalized.startsWith(p)) ) { await fs.promises.copyFile(sourceFile, targetFile) continue diff --git a/src/storage/adapters/memoryStorage.ts b/src/storage/adapters/memoryStorage.ts index 41987ea0..bab9d4d9 100644 --- a/src/storage/adapters/memoryStorage.ts +++ b/src/storage/adapters/memoryStorage.ts @@ -222,6 +222,45 @@ export class MemoryStorage extends BaseStorage { return [...this.txLogLines] } + // =========================================================================== + // Binary raw-byte primitives — in-memory mirror of the filesystem adapter's + // append-only substrate (the generation fact log's segments), so memory + // brains dual-write facts too and the compat suite runs on both adapters. + // =========================================================================== + + /** Raw binary files, keyed by verbatim path. */ + private rawBytesStore: Map = new Map() + + /** Append bytes to a raw binary file, creating it when absent. */ + public async appendRawBytes(rawPath: string, bytes: Uint8Array): Promise { + const existing = this.rawBytesStore.get(rawPath) + if (!existing) { + this.rawBytesStore.set(rawPath, bytes.slice()) + return + } + const merged = new Uint8Array(existing.length + bytes.length) + merged.set(existing, 0) + merged.set(bytes, existing.length) + this.rawBytesStore.set(rawPath, merged) + } + + /** Read a raw binary file whole (a copy); absent → null. */ + public async readRawBytes(rawPath: string): Promise { + const bytes = this.rawBytesStore.get(rawPath) + return bytes ? bytes.slice() : null + } + + /** Replace a raw binary file (atomic by construction in memory). */ + public async writeRawBytes(rawPath: string, bytes: Uint8Array): Promise { + this.rawBytesStore.set(rawPath, bytes.slice()) + } + + /** Byte size of a raw binary file, or null when absent. */ + public async rawByteSize(rawPath: string): Promise { + const bytes = this.rawBytesStore.get(rawPath) + return bytes ? bytes.length : null + } + /** * Serialize the entire in-memory store to a directory in the exact layout * the filesystem adapter uses (uncompressed JSON objects, `_blobs/*.bin` @@ -354,6 +393,7 @@ export class MemoryStorage extends BaseStorage { public async clear(): Promise { this.objectStore.clear() this.blobStore.clear() + this.rawBytesStore.clear() this.txLogLines = [] this.statistics = null diff --git a/src/utils/crc32c.ts b/src/utils/crc32c.ts new file mode 100644 index 00000000..0c2a1c14 --- /dev/null +++ b/src/utils/crc32c.ts @@ -0,0 +1,43 @@ +/** + * @module utils/crc32c + * @description CRC-32C (Castagnoli, polynomial 0x1EDC6F41, reflected 0x82F63B78) + * — the storage-industry frame checksum (ext4, iSCSI, SCTP, LSM segment files). + * Used to frame generation-fact segments: every appended record carries the + * CRC-32C of its payload, so a torn tail (crash mid-append) or bit rot is + * DETECTED at scan time and never silently read as data. + * + * Table-driven, dependency-free reference implementation. Native providers may + * substitute a hardware-accelerated (SSE4.2 / ARMv8 CRC) implementation — the + * polynomial is the contract, byte-identical results required. + */ + +/** The 256-entry lookup table for the reflected CRC-32C polynomial. */ +const TABLE: Uint32Array = (() => { + const table = new Uint32Array(256) + for (let n = 0; n < 256; n++) { + let c = n + for (let k = 0; k < 8; k++) { + c = c & 1 ? 0x82f63b78 ^ (c >>> 1) : c >>> 1 + } + table[n] = c >>> 0 + } + return table +})() + +/** + * Compute the CRC-32C checksum of a byte buffer. + * + * Known-answer vectors (RFC 3720 appendix / the standard test suite): + * - ASCII "123456789" → 0xE3069283 + * - 32 zero bytes → 0x8A9136AA + * + * @param bytes - The payload to checksum. + * @returns The CRC-32C as an unsigned 32-bit integer. + */ +export function crc32c(bytes: Uint8Array): number { + let crc = 0xffffffff + for (let i = 0; i < bytes.length; i++) { + crc = TABLE[(crc ^ bytes[i]) & 0xff] ^ (crc >>> 8) + } + return (crc ^ 0xffffffff) >>> 0 +} diff --git a/tests/integration/db-mvcc.test.ts b/tests/integration/db-mvcc.test.ts index 764a755a..10a158ae 100644 --- a/tests/integration/db-mvcc.test.ts +++ b/tests/integration/db-mvcc.test.ts @@ -515,7 +515,15 @@ describe('8.0 Db API — generational MVCC', () => { await brain.transact([{ op: 'update', id: uid('compact-e'), metadata: { v: 4 } }]) ).release() - const recordsBefore = (await storage.listRawObjects('_generations')).length + // History record-sets only — the generation FACT LOG also lives under + // `_generations/` (at `facts/`) and is deliberately NOT reclaimed by + // history compaction (facts are the future canonical, not undo history). + const historyRecords = async (): Promise => + (await storage.listRawObjects('_generations')).filter( + (p: string) => !p.startsWith('_generations/facts/') + ).length + + const recordsBefore = await historyRecords() expect(recordsBefore).toBeGreaterThan(0) // Compact while pinned: record-sets above the pin survive, pinned reads stay correct. @@ -528,7 +536,7 @@ describe('8.0 Db API — generational MVCC', () => { const second = await brain.compactHistory() expect(first.removedGenerations + second.removedGenerations).toBeGreaterThan(0) - const recordsAfter = (await storage.listRawObjects('_generations')).length + const recordsAfter = await historyRecords() expect(recordsAfter).toBeLessThan(recordsBefore) expect(recordsAfter).toBe(0) diff --git a/tests/integration/fact-log-dual-write.test.ts b/tests/integration/fact-log-dual-write.test.ts new file mode 100644 index 00000000..9df1b908 --- /dev/null +++ b/tests/integration/fact-log-dual-write.test.ts @@ -0,0 +1,200 @@ +/** + * @module tests/integration/fact-log-dual-write + * @description The generation fact log end-to-end through real commits: every + * committed generation (single-op AND transact) appends its AFTER-IMAGE fact + * at the commit point; removals append body-less tombstones; an aborted + * transaction leaves no fact; facts survive reopen and continue monotonically; + * the scan surface (brain.scanFacts) carries the frozen telemetry shape; and + * the fact-log namespace is protected against prefix-nuking. + */ +import { describe, it, expect, beforeEach, afterEach } from 'vitest' +import * as fs from 'node:fs' +import * as os from 'node:os' +import * as path from 'node:path' +import { Brainy, ProtectedArtifactError, type CommitFact } from '../../src/index.js' + +async function allFacts(brain: any): Promise { + const scan = brain.scanFacts() + expect(scan).not.toBeNull() + const facts: CommitFact[] = [] + for await (const batch of scan!.batches()) facts.push(...batch.facts) + return facts +} + +describe('fact log dual-write (memory adapter)', () => { + let brain: any + + beforeEach(async () => { + process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true' + brain = new Brainy({ requireSubtype: false, storage: { type: 'memory' }, silent: true, dimensions: 384 }) + await brain.init() + }) + afterEach(async () => { + await brain.close?.().catch(() => {}) + }) + + it('every single-op write appends its after-image fact; a remove appends a tombstone', async () => { + const id = await brain.add({ data: 'first', type: 'document', metadata: { rev: 1 } }) + await brain.update({ id, metadata: { rev: 2 } }) + await brain.remove(id) + + const facts = await allFacts(brain) + // add + update + remove each committed a generation (the remove may span + // cascade ops but is ONE generation). Facts are monotonic. + const gens = facts.map((f) => f.generation) + expect([...gens].sort((a, b) => a - b)).toEqual(gens) + expect(facts.length).toBeGreaterThanOrEqual(3) + + // The add fact carries the after-image of the new entity. + const addFact = facts.find((f) => f.ops.some((op) => op.id === id && op.record !== null)) + expect(addFact).toBeDefined() + + // The remove fact carries a body-less tombstone for the id. + const removeFact = facts[facts.length - 1] + const tombstone = removeFact.ops.find((op) => op.id === id) + expect(tombstone).toBeDefined() + expect(tombstone!.record).toBeNull() + expect(tombstone!.kind).toBe('noun') + }) + + it('the update fact holds the NEW state (after-image, not before)', async () => { + const id = await brain.add({ data: 'versioned', type: 'document', metadata: { v: 'old' } }) + await brain.update({ id, metadata: { v: 'new' } }) + + const facts = await allFacts(brain) + const updateFact = facts[facts.length - 1] + const op = updateFact.ops.find((o) => o.id === id)! + expect(op.record).not.toBeNull() + expect((op.record!.metadata as any).v).toBe('new') + }) + + it('a transact commits ONE fact carrying all its ops, with meta', async () => { + const receipt = await brain.transact( + [ + { op: 'add', type: 'document', metadata: { part: 1 }, data: 'a' }, + { op: 'add', type: 'document', metadata: { part: 2 }, data: 'b' } + ], + { meta: { source: 'batch-import' } } + ) + + const facts = await allFacts(brain) + const txFact = facts.find((f) => f.generation === receipt.generation) + expect(txFact).toBeDefined() + expect(txFact!.ops.filter((op) => op.kind === 'noun').length).toBeGreaterThanOrEqual(2) + expect(txFact!.meta).toEqual({ source: 'batch-import' }) + }) + + it('an aborted transact leaves NO fact (absent = never committed)', async () => { + const id = await brain.add({ data: 'cas target', type: 'document', metadata: { n: 1 } }) + const before = (await allFacts(brain)).length + + await expect( + brain.transact([{ op: 'update', id, ifRev: 999, metadata: { n: 2 } }]) + ).rejects.toThrow() + + const after = await allFacts(brain) + expect(after.length).toBe(before) + }) + + it('fact generations line up with the transaction log', async () => { + await brain.add({ data: 'x', type: 'document', metadata: {} }) + await brain.add({ data: 'y', type: 'document', metadata: {} }) + await brain.flush() + + const facts = await allFacts(brain) + const logGens = new Set((await brain.transactionLog()).map((e: any) => e.generation)) + for (const f of facts) { + expect(logGens.has(f.generation)).toBe(true) + } + }) + + it('scan telemetry carries the frozen shape end-to-end', async () => { + for (let i = 0; i < 5; i++) await brain.add({ data: `t${i}`, type: 'document', metadata: { i } }) + + const scan = brain.scanFacts({ batchSize: 2 })! + expect(scan.headGeneration).toBeGreaterThanOrEqual(5) + expect(scan.approxFactCount).toBeGreaterThanOrEqual(5) + let batches = 0 + for await (const b of scan.batches()) { + batches++ + expect(b.factCount).toBe(b.facts.length) + expect(b.firstGeneration).toBe(b.facts[0].generation) + expect(b.lastGeneration).toBe(b.facts[b.facts.length - 1].generation) + expect(b.byteSize).toBeGreaterThan(0) + expect(typeof b.segmentId).toBe('string') + } + expect(batches).toBeGreaterThan(1) + expect(scan.summary().factsYielded).toBe(scan.approxFactCount) + }) +}) + +describe('fact log dual-write (filesystem adapter — durability + protection)', () => { + let dir: string + let brain: any + + const open = async () => { + const b: any = new Brainy({ + requireSubtype: false, + storage: { type: 'filesystem', path: dir }, + silent: true, + dimensions: 384 + }) + await b.init() + return b + } + + beforeEach(async () => { + process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true' + dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-factlog-')) + brain = await open() + }) + afterEach(async () => { + await brain.close?.().catch(() => {}) + fs.rmSync(dir, { recursive: true, force: true }) + }) + + it('facts survive close + reopen and appends continue monotonically', async () => { + const id = await brain.add({ data: 'persist me', type: 'document', metadata: { k: 1 } }) + await brain.remove(id) + await brain.close() + + brain = await open() + const facts = await allFacts(brain) + expect(facts.length).toBeGreaterThanOrEqual(2) + const headBefore = facts[facts.length - 1].generation + + await brain.add({ data: 'after reopen', type: 'document', metadata: { k: 2 } }) + const facts2 = await allFacts(brain) + expect(facts2[facts2.length - 1].generation).toBeGreaterThan(headBefore) + }) + + it('the fact segments exist on disk under _generations/facts/ with zero-padded names', async () => { + await brain.add({ data: 'on disk', type: 'document', metadata: {} }) + await brain.flush() + const factsDir = path.join(dir, '_generations', 'facts') + const files = fs.readdirSync(factsDir) + // The manifest rides the store's JSON object discipline (gzip on disk). + expect(files.some((f) => f.startsWith('manifest.json'))).toBe(true) + const segs = files.filter((f) => /^seg-\d{20}\.bfl$/.test(f)) + expect(segs.length).toBeGreaterThanOrEqual(1) + }) + + it('the fact-log namespace is PROTECTED: a prefix-nuke is refused', async () => { + await brain.add({ data: 'protected', type: 'document', metadata: {} }) + await expect(brain.storage.removeRawPrefix('_generations/facts')).rejects.toBeInstanceOf( + ProtectedArtifactError + ) + // Per-generation history cleanup remains unaffected (no false intersect). + await expect(brain.storage.removeRawPrefix('_generations/999999')).resolves.toBeUndefined() + }) + + it('transact facts are durable-on-return (no flush needed before reopen)', async () => { + const receipt = await brain.transact([ + { op: 'add', type: 'document', metadata: { durable: true }, data: 'tx' } + ]) + // Simulate an abrupt end: no flush(), no close() — reopen from disk. + brain = await open() + const facts = await allFacts(brain) + expect(facts.some((f) => f.generation === receipt.generation)).toBe(true) + }) +}) diff --git a/tests/unit/db/fact-log.test.ts b/tests/unit/db/fact-log.test.ts new file mode 100644 index 00000000..abce2dc9 --- /dev/null +++ b/tests/unit/db/fact-log.test.ts @@ -0,0 +1,189 @@ +/** + * @module tests/unit/db/fact-log + * @description The generation fact log in isolation: wire-format round-trip + * (positional msgpack facts, bin16 uuids, body-less tombstones), crc32c + * framing with torn-tail detection, open-time truncation to committed truth + * (the log can only ever be AHEAD after a crash; open cuts it back), rotation + * with a manifest-first flip, exactly-once scans with the frozen telemetry + * shape, and the mmap segment handoff excluding the mutable tail. + */ +import { describe, it, expect, beforeEach } from 'vitest' +import { MemoryStorage } from '../../../src/storage/adapters/memoryStorage.js' +import { + FactLog, + FACTS_PREFIX, + type CommitFact, + type FactLogStorage, + storageSupportsFactLog +} from '../../../src/db/factLog.js' +import { crc32c } from '../../../src/utils/crc32c.js' + +const UUID = (n: number): string => + `00000000-0000-4000-8000-${String(n).padStart(12, '0')}` + +const fact = (generation: number, overrides?: Partial): CommitFact => ({ + generation, + timestamp: 1_700_000_000_000 + generation, + ops: [ + { + kind: 'noun', + id: UUID(generation), + record: { metadata: { noun: 'document', title: `doc ${generation}` }, vector: { v: [1, 2] } } + } + ], + ...overrides +}) + +describe('crc32c known-answer vectors', () => { + it('matches the RFC 3720 test vectors', () => { + expect(crc32c(new TextEncoder().encode('123456789'))).toBe(0xe3069283) + expect(crc32c(new Uint8Array(32))).toBe(0x8a9136aa) + }) +}) + +describe('fact log — round-trip, framing, reconcile, rotation, scan', () => { + let storage: FactLogStorage + let log: FactLog + + beforeEach(async () => { + const mem: any = new MemoryStorage() + await mem.init() + expect(storageSupportsFactLog(mem)).toBe(true) + storage = mem + log = new FactLog(storage) + await log.open(0) + }) + + it('facts round-trip byte-exactly: ops, tombstones, meta, blobHashes', async () => { + await log.append(fact(1)) + await log.append( + fact(2, { + ops: [ + { kind: 'verb', id: UUID(21), record: { metadata: { verb: 'contains' }, vector: null } }, + { kind: 'noun', id: UUID(22), record: null } // TOMBSTONE + ], + meta: { source: 'test' }, + blobHashes: ['abc123', 'abc123'] // multiset — duplicates preserved + }) + ) + await log.sync() + + const scan = log.scanFacts() + expect(scan.headGeneration).toBe(2) + const all: CommitFact[] = [] + for await (const batch of scan.batches()) all.push(...batch.facts) + + expect(all).toHaveLength(2) + expect(all[0].generation).toBe(1) + expect(all[0].ops[0].id).toBe(UUID(1)) + expect(all[0].ops[0].record?.metadata).toEqual({ noun: 'document', title: 'doc 1' }) + expect(all[1].ops[0].kind).toBe('verb') + expect(all[1].ops[1].record).toBeNull() // the tombstone is body-less + expect(all[1].meta).toEqual({ source: 'test' }) + expect(all[1].blobHashes).toEqual(['abc123', 'abc123']) + expect(scan.summary().factsYielded).toBe(2) + }) + + it('appends are monotonic — a replayed/duplicate generation throws', async () => { + await log.append(fact(5)) + await expect(log.append(fact(5))).rejects.toThrow(/non-monotonic/) + await expect(log.append(fact(3))).rejects.toThrow(/non-monotonic/) + await expect(log.append(fact(6))).resolves.toBeUndefined() // gaps are fine (aborted reservations) + }) + + it('a torn tail (partial frame) is detected and ignored — intact prefix survives', async () => { + await log.append(fact(1)) + await log.append(fact(2)) + await log.sync() + + // Simulate a crash mid-append: chop bytes off the tail file. + const tailPath = `${FACTS_PREFIX}/seg-${'1'.padStart(20, '0')}.bfl` + const bytes = (await storage.readRawBytes(tailPath))! + await storage.writeRawBytes(tailPath, bytes.subarray(0, bytes.length - 7)) + + const reopened = new FactLog(storage) + await reopened.open(2) + expect(reopened.headGeneration()).toBe(1) // fact 2's frame was torn → gone + + const all: CommitFact[] = [] + for await (const b of reopened.scanFacts().batches()) all.push(...b.facts) + expect(all.map((f) => f.generation)).toEqual([1]) + }) + + it('open() truncates facts beyond committed truth (the crash-ahead shape)', async () => { + await log.append(fact(1)) + await log.append(fact(2)) + await log.append(fact(3)) + await log.sync() + + // The store's committed generation is 1 — facts 2..3 never committed. + const reopened = new FactLog(storage) + await reopened.open(1) + expect(reopened.headGeneration()).toBe(1) + + const all: CommitFact[] = [] + for await (const b of reopened.scanFacts().batches()) all.push(...b.facts) + expect(all.map((f) => f.generation)).toEqual([1]) + + // And appends continue cleanly from the truncated head. + await reopened.append(fact(2)) + expect(reopened.headGeneration()).toBe(2) + }) + + it('a cleared store (committed=0) truncates everything', async () => { + await log.append(fact(1)) + await log.append(fact(2)) + await log.sync() + const reopened = new FactLog(storage) + await reopened.open(0) + expect(reopened.headGeneration()).toBe(0) + }) + + it('scan honors fromGeneration/toGeneration inclusively and filters kinds', async () => { + for (let g = 1; g <= 6; g++) await log.append(fact(g)) + await log.sync() + + const scan = log.scanFacts({ fromGeneration: 2, toGeneration: 4 }) + const all: CommitFact[] = [] + for await (const b of scan.batches()) all.push(...b.facts) + expect(all.map((f) => f.generation)).toEqual([2, 3, 4]) + + const verbsOnly = log.scanFacts({ kinds: ['verb'] }) + for await (const b of verbsOnly.batches()) { + for (const f of b.facts) expect(f.ops.every((op) => op.kind === 'verb')).toBe(true) + } + }) + + it('batch telemetry carries the frozen shape', async () => { + for (let g = 1; g <= 5; g++) await log.append(fact(g)) + await log.sync() + + const scan = log.scanFacts({ batchSize: 2 }) + expect(scan.approxFactCount).toBe(5) + const batches = [] + for await (const b of scan.batches()) batches.push(b) + expect(batches.length).toBe(3) + expect(batches[0]).toMatchObject({ firstGeneration: 1, lastGeneration: 2, factCount: 2 }) + expect(batches[0].byteSize).toBeGreaterThan(0) + expect(typeof batches[0].segmentId).toBe('string') + expect(scan.summary()).toEqual({ factsYielded: 5, segmentsRead: 1 }) + }) + + it('survives reopen: head and content come back from disk', async () => { + for (let g = 1; g <= 3; g++) await log.append(fact(g)) + await log.sync() + + const reopened = new FactLog(storage) + await reopened.open(3) + expect(reopened.headGeneration()).toBe(3) + const all: CommitFact[] = [] + for await (const b of reopened.scanFacts().batches()) all.push(...b.facts) + expect(all.map((f) => f.generation)).toEqual([1, 2, 3]) + }) + + it('segmentPaths excludes the mutable tail (mmap handoff = sealed only)', async () => { + await log.append(fact(1)) + await log.sync() + expect(log.segmentPaths()).toEqual([]) // only a tail exists — nothing sealed + }) +}) diff --git a/tests/unit/db/generationStore.test.ts b/tests/unit/db/generationStore.test.ts index 8073d7ee..184d3974 100644 --- a/tests/unit/db/generationStore.test.ts +++ b/tests/unit/db/generationStore.test.ts @@ -249,7 +249,11 @@ describe('db/GenerationStore', () => { store.release(pinned) const result = await store.compact() expect(result.removedGenerations).toBeGreaterThan(0) - const remaining = await storage.listRawObjects(GENERATIONS_PREFIX) + // History record-sets only — the fact log (at `_generations/facts/`) is + // deliberately NOT reclaimed by history compaction. + const remaining = (await storage.listRawObjects(GENERATIONS_PREFIX)).filter( + (p: string) => !p.startsWith(`${GENERATIONS_PREFIX}/facts/`) + ) expect(remaining).toEqual([]) })