From d1651f986c5d235f580daf93ad70e1381c066021 Mon Sep 17 00:00:00 2001 From: David Snelling Date: Mon, 10 Aug 2026 11:39:27 -0700 Subject: [PATCH] =?UTF-8?q?feat(reprojection):=20the=20one=20doors-open=20?= =?UTF-8?q?machinery=20=E2=80=94=20budget-capped,=20yielding,=20foreground?= =?UTF-8?q?-preempted,=20atomic-swap;=20poison=20records=20quarantine=20ty?= =?UTF-8?q?ped?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The generic reprojection engine (pure TS; the twin of the native implementation — same frozen contract, one shared conformance intent): register any ProjectionAdapter; advance(family, {budgetMs}) folds facts from the adapter's own watermark to the head in installments ≤50ms with real macrotask yields; foreground door traffic bumps the DoorSignal and an in-flight advance yields within one installment ('preempted'); advanceAll round-robins families fairly. swap(family, buildAdapter) is the doors-open migration primitive: the OLD projection keeps serving while the new one builds beside it, the flip is atomic at parity, and a concurrent second swap refuses typed. A fact the fold cannot apply (typed ProjectionApplyError) is QUARANTINED — skipped, ledgered, narrated per-doubling, exposed for refuse-affected-reads — the service class law's fourth answer: never a wedged rebuild, never a silent skip. The engine never writes stamps: each adapter owns its durability and its stamp-after-data discipline. Upgrade, heal, and rebuild are now the same machinery behind open doors. FactLogSource wires any host's fact scan in one line (factSourceFromHost(brain)); window-contract violations are loud. Pins: 23 unit (budget resume without refold · preemption within one installment · round-robin fairness under a skewed backlog · build-beside visibility mid-swap · atomic flip · single-flight refusal · quarantine skip/ledger/doubling · non-typed throw aborts · losing adapter discarded) + 3 integration on a real brain (fold matches ground truth · doors answer mid-fold with the preemption path exercised · crash mid-fold resumes from the stamp, never refolds). Gates: unit 2054/2054 (157 files) · integration 820 (93 files) · conformance 31/31. --- src/reprojection/factLogSource.ts | 141 ++++ src/reprojection/reprojectionEngine.ts | 648 ++++++++++++++++++ .../reprojection-doors-open.test.ts | 257 +++++++ .../reprojection/reprojection-engine.test.ts | 590 ++++++++++++++++ 4 files changed, 1636 insertions(+) create mode 100644 src/reprojection/factLogSource.ts create mode 100644 src/reprojection/reprojectionEngine.ts create mode 100644 tests/integration/reprojection-doors-open.test.ts create mode 100644 tests/unit/reprojection/reprojection-engine.test.ts diff --git a/src/reprojection/factLogSource.ts b/src/reprojection/factLogSource.ts new file mode 100644 index 00000000..796fb343 --- /dev/null +++ b/src/reprojection/factLogSource.ts @@ -0,0 +1,141 @@ +/** + * @module reprojection/factLogSource + * @description The production {@link FactSource}: adapts the database's + * committed-fact scan to the reprojection engine's `scan(from, limit)` + * window contract. + * + * DEPENDENCY-CLEAN BY DESIGN: this module never imports the database class. + * It wraps a host-owned scan callback `(from, limit) => Promise` + * injected at construction, so the host wires itself in one line — either by + * handing {@link FactLogSource} a callback built on its own scan API, or via + * {@link factSourceFromHost}, which builds that callback from any object + * structurally exposing `scanFacts` (the batch-handle shape the fact log + * serves). + * + * CONTRACT ENFORCEMENT — loud, never quiet: every `scan` return is checked + * (≤ limit facts, strictly ascending generations, all strictly above `from`); + * a violating callback throws instead of silently corrupting a fold. A host + * with NO fact log throws too — reporting "caught up" against an unscannable + * store would be a silent lie. + */ + +import type { CommitFact } from '../db/factLog.js' +import type { FactSource } from './reprojectionEngine.js' + +/** + * The host-owned scan callback: return up to `limit` committed facts with + * generation strictly greater than `from`, in ascending generation order; + * empty means caught up to the head as of the call. + */ +export type FactScanCallback = (from: number, limit: number) => Promise + +/** + * The minimal structural surface of a fact-scanning host — matches the + * database's `scanFacts` shape without importing it. `scanFacts` returns a + * handle whose `batches()` yields ordered, non-empty fact batches, or `null` + * when the store hosts no fact log. + */ +export interface FactScanHost { + scanFacts(options?: { fromGeneration?: number; batchSize?: number }): { + batches: () => AsyncGenerator<{ facts: CommitFact[] }> + } | null +} + +/** + * The production {@link FactSource}: wraps an injected scan callback and + * enforces the window contract on every return. + * + * COST NOTE: each `scan` call is stateless (a fresh window above the caller's + * watermark), which is exactly what resumable, crash-tolerant folds need — + * at the price of the host re-opening its scan per call. Fine for + * budget-capped maintenance; not a hot-path read primitive. + */ +export class FactLogSource implements FactSource { + private readonly scanCallback: FactScanCallback + + /** @param scanCallback - The host-owned scan (see {@link FactScanCallback}). */ + constructor(scanCallback: FactScanCallback) { + if (typeof scanCallback !== 'function') { + throw new Error('FactLogSource: a scan callback (from, limit) => Promise is required') + } + this.scanCallback = scanCallback + } + + /** + * Fetch up to `limit` committed facts strictly above generation `from`, + * verifying the callback honored the window contract. + * @param from - Exclusive lower bound generation (≥ 0 integer). + * @param limit - Maximum facts to return (≥ 1 integer). + */ + async scan(from: number, limit: number): Promise { + if (!Number.isInteger(from) || from < 0) { + throw new Error(`FactLogSource.scan: 'from' must be a non-negative integer (got ${from})`) + } + if (!Number.isInteger(limit) || limit < 1) { + throw new Error(`FactLogSource.scan: 'limit' must be a positive integer (got ${limit})`) + } + const facts = await this.scanCallback(from, limit) + if (!Array.isArray(facts)) { + throw new Error('FactLogSource.scan: the scan callback must resolve to an array of facts') + } + if (facts.length > limit) { + throw new Error( + `FactLogSource.scan: the scan callback returned ${facts.length} facts for limit ${limit} — ` + + `contract violation; refusing to fold an oversized window` + ) + } + let prev = from + for (const fact of facts) { + const g = fact?.generation + if (typeof g !== 'number' || !Number.isFinite(g) || g <= prev) { + throw new Error( + `FactLogSource.scan: the scan callback violated the window contract — generation ` + + `${String(g)} is not strictly ascending above ${prev} (from=${from}); refusing to fold` + ) + } + prev = g + } + return facts + } +} + +/** + * Build the production source from any host structurally exposing + * `scanFacts` — the one-line wiring for the database side: + * + * ```ts + * const source = factSourceFromHost(brain) + * ``` + * + * Each `scan(from, limit)` opens `scanFacts({ fromGeneration: from + 1, + * batchSize: limit })` (the engine's `from` is exclusive; `scanFacts` bounds + * are inclusive) and returns the FIRST batch, closing the handle — short + * batches at segment boundaries are legal under the source contract (only + * EMPTY means caught up). A host with no fact log throws loudly. + * + * @param host - Any object with the `scanFacts` batch-handle shape. + */ +export function factSourceFromHost(host: FactScanHost): FactLogSource { + if (!host || typeof host.scanFacts !== 'function') { + throw new Error('factSourceFromHost: the host must expose scanFacts(options)') + } + return new FactLogSource(async (from, limit) => { + const scan = host.scanFacts({ fromGeneration: from + 1, batchSize: limit }) + if (scan === null) { + throw new Error( + 'reprojection: this store hosts no fact log — reprojection folds committed facts, ' + + 'and reporting a caught-up fold against an unscannable store would be a silent lie' + ) + } + const iterator = scan.batches() + try { + const first = await iterator.next() + return first.done ? [] : first.value.facts + } finally { + // Close the abandoned generator so its cleanup (timers) runs. + if (typeof iterator.return === 'function') { + await iterator.return(undefined) + } + } + }) +} diff --git a/src/reprojection/reprojectionEngine.ts b/src/reprojection/reprojectionEngine.ts new file mode 100644 index 00000000..1465b1f6 --- /dev/null +++ b/src/reprojection/reprojectionEngine.ts @@ -0,0 +1,648 @@ +/** + * @module reprojection/reprojectionEngine + * @description The pure-TS reprojection engine — the ONE machinery for + * rebuilding, healing, and migrating persisted projections from the committed + * fact log on the JS side. It is the TypeScript twin of the native engine's + * reprojection core: the same frozen contract (names AND semantics), so a + * single shared conformance suite runs against both implementations and + * TS-only deployments green the same rows without native code. + * + * THE AVAILABILITY LAW — maintenance never holds the doors: + * + * - Work proceeds in INSTALLMENTS of at most {@link MAX_INSTALLMENT_MS} (50ms) + * of wall time each. Between installments the loop awaits a REAL macrotask + * boundary (never a busy loop, never a bare microtask), so foreground I/O + * and timers always interleave with a running fold. + * - Foreground door traffic announces itself via {@link DoorSignal.bump}. An + * in-flight {@link ReprojectionEngine.advance} yields at the next + * installment boundary and returns `{ status: 'preempted' }` — the doors + * never wait for maintenance to finish. + * - Budgets are honored: `advance` stops once `budgetMs` is spent and reports + * exactly how far it got; a later call RESUMES from the adapter's own + * watermark. Nothing ever refolds from zero because a budget ran out. + * + * WATERMARK DISCIPLINE — the engine NEVER writes stamps. Each adapter's + * `applyBatch` owns its own durability and its own stamp (stamp-after-data, + * the law stated in src/utils/projectionWatermark.ts); the engine only READS + * `watermark()` to decide the next scan window. Delivery is therefore + * at-least-once: an adapter that crashed between data and stamp is re-served + * the same facts on resume and MUST apply idempotently. + * + * THE FOUR ANSWER CLASSES of an advance: `'caught-up'` (folded to the head of + * the requested window, ledger clean), `'preempted'` (a door bumped), + * `'budget-exhausted'` (time ran out mid-stream), and `'quarantined'` (folded + * to the head, but this family's quarantine ledger is non-empty — one or more + * poison facts are being skipped and reads touching them are suspect). + */ + +import type { CommitFact } from '../db/factLog.js' +import { prodLog } from '../utils/logger.js' + +/** + * The hard ceiling on one installment of fold work, in wall-clock ms. An + * advance loop that has run this long without yielding closes the installment + * and awaits a macrotask boundary so foreground traffic interleaves. Frozen by + * the shared contract — both engines install the same ceiling. + */ +export const MAX_INSTALLMENT_MS = 50 + +/** Default facts-per-batch pulled from the {@link FactSource} per step. */ +export const DEFAULT_REPROJECTION_BATCH_SIZE = 256 + +/** + * One registered projection family: a named consumer that folds committed + * facts into its own persisted artifact and stamps its own watermark. + * + * OWNERSHIP: the adapter owns durability AND the stamp. `applyBatch` must + * persist its data first and stamp `upTo` after (stamp-after-data), and must + * tolerate at-least-once delivery — on resume after a crash between data and + * stamp, the same facts arrive again. + */ +export interface ProjectionAdapter { + /** Unique family name — the registry key; one adapter serves a family at a time. */ + family: string + /** + * The highest generation this projection's persisted state reflects, or + * `null` when the projection is unbuilt/unstamped. The engine reads this to + * open the next scan window; it never writes it. + */ + watermark(): number | null + /** + * Fold `facts` (ascending generations, all strictly above the current + * watermark) into the projection, then stamp `watermark = upTo`. + * + * `facts` MAY be empty while `upTo` is above the current watermark: that is + * a pure watermark advance past quarantined generations — the adapter must + * still stamp, or the fold cannot make progress past the poison. + * + * FAILURE CONTRACT: throw a {@link ProjectionApplyError} to name exactly one + * poison fact (the engine quarantines it and continues). ANY other throw + * aborts the advance loudly — an unknown failure is never treated as a + * poison record. + */ + applyBatch(facts: CommitFact[], upTo: number): Promise + /** + * Destroy this adapter's persisted artifact(s). The engine calls this on + * the LOSING adapter after a successful {@link ReprojectionEngine.swap}, + * and on a partially-built replacement whose build aborted. + */ + discard(): Promise +} + +/** + * The committed-fact scan the engine folds from. `from` is an EXCLUSIVE lower + * bound generation; the source returns at most `limit` facts in ascending + * generation order, and an empty array means caught up to the head as of this + * call. Short non-empty returns are legal (e.g. a segment boundary) — only + * empty means done. + */ +export interface FactSource { + scan(from: number, limit: number): Promise +} + +/** + * The foreground-preemption signal. Door traffic (foreground reads/writes) + * calls {@link DoorSignal.bump}; an in-flight `advance` observes the bump at + * its next installment boundary, yields a macrotask, and returns + * `{ status: 'preempted' }`. Bumps are edge-triggered per advance: only bumps + * that arrive AFTER an advance began preempt it. + */ +export class DoorSignal { + private count = 0 + + /** Announce foreground door traffic — an in-flight advance will yield. */ + bump(): void { + this.count++ + } + + /** + * The current bump epoch — the engine snapshots this at advance entry and + * compares at installment boundaries. + * @internal + */ + epoch(): number { + return this.count + } +} + +/** + * The TYPED poison-record failure an adapter throws from `applyBatch` to name + * exactly one unfoldable fact. The engine quarantines that generation for + * that family (skips it, ledgers it, narrates per-doubling) and keeps + * folding. Any OTHER throw from `applyBatch` aborts the advance loudly. + */ +export class ProjectionApplyError extends Error { + /** The generation of the fact that cannot be applied. */ + readonly generation: number + /** Optional index of the offending record within the fact's ops. */ + readonly recordIndex?: number + /** The underlying failure. */ + override readonly cause: unknown + + /** + * @param args - `generation` names the poison fact; `recordIndex` + * optionally narrows to one record inside it; `cause` carries the + * underlying failure. + */ + constructor(args: { generation: number; recordIndex?: number; cause: unknown }) { + super( + `projection apply failed at generation ${args.generation}` + + (args.recordIndex !== undefined ? ` (record ${args.recordIndex})` : '') + ) + this.name = 'ProjectionApplyError' + this.generation = args.generation + if (args.recordIndex !== undefined) this.recordIndex = args.recordIndex + this.cause = args.cause + } +} + +/** + * The TYPED single-flight refusal: a second concurrent + * {@link ReprojectionEngine.swap} on a family whose replacement is still + * building. The caller retries after the in-flight swap settles. + */ +export class SwapInFlightError extends Error { + /** The family whose swap is already in flight. */ + readonly family: string + + /** @param family - The family whose swap is already in flight. */ + constructor(family: string) { + super( + `reprojection: a swap is already in flight for family '${family}' — ` + + `swaps are single-flight per family; retry after the current build settles` + ) + this.name = 'SwapInFlightError' + this.family = family + } +} + +/** One quarantined fact in a family's ledger. */ +export interface QuarantineEntry { + /** The generation being skipped for this family. */ + generation: number + /** The typed apply failure that condemned it. */ + error: ProjectionApplyError + /** Wall-clock ms when it was quarantined (diagnostic). */ + at: number +} + +/** How an advance ended — the four answer classes (see the module header). */ +export type AdvanceStatus = 'caught-up' | 'preempted' | 'budget-exhausted' | 'quarantined' + +/** The result of one advance over one family. */ +export interface AdvanceResult { + /** The answer class. */ + status: AdvanceStatus + /** The family's watermark as stamped by its own adapter, after this advance. */ + watermark: number | null + /** + * Facts delivered in SUCCESSFUL `applyBatch` calls during this advance. + * At-least-once delivery means retried facts (after a quarantine or a + * resume) count again; this is delivered work, not distinct generations. + */ + applied: number +} + +/** The result of a completed {@link ReprojectionEngine.swap}. */ +export interface SwapResult { + /** The NEW adapter's watermark at the flip (parity with the head). */ + watermark: number | null + /** Facts delivered to the replacement during its beside-build. */ + applied: number +} + +/** Constructor options for {@link ReprojectionEngine}. */ +export interface ReprojectionEngineOptions { + /** The committed-fact scan every family folds from. */ + source: FactSource + /** The preemption signal; a fresh one is created when omitted. */ + doorSignal?: DoorSignal + /** + * Installment ceiling in ms, `(0, MAX_INSTALLMENT_MS]`. Out-of-range values + * throw — the 50ms law is a ceiling, never a suggestion. + */ + installmentMs?: number + /** Facts per {@link FactSource.scan} pull (default {@link DEFAULT_REPROJECTION_BATCH_SIZE}). */ + batchSize?: number +} + +/** The fold-side state shared by a serving family and a swap's beside-build. */ +interface FoldState { + adapter: ProjectionAdapter + /** The quarantine ledger, in condemnation order. */ + quarantine: QuarantineEntry[] + /** Generations filtered out of every batch served to this adapter. */ + skip: Set + /** Next ledger size that triggers a narration (1, 2, 4, 8, …). */ + nextWarnAt: number +} + +/** A registered family: fold state plus the single-flight swap latch. */ +interface FamilyState extends FoldState { + swapInFlight: boolean +} + +/** One real macrotask boundary — foreground I/O and timers run before resume. */ +function yieldToDoors(): Promise { + return new Promise((resolve) => { + if (typeof setImmediate === 'function') { + setImmediate(resolve) + } else { + setTimeout(resolve, 0) + } + }) +} + +/** + * The reprojection engine: registry of projection families, budget-capped + * yielding advances, round-robin `advanceAll`, atomic build-beside `swap`, + * and the per-family quarantine ledger. Pure TS, no storage dependencies — + * everything durable lives behind the injected {@link FactSource} and the + * registered {@link ProjectionAdapter}s. + */ +export class ReprojectionEngine { + /** The preemption signal foreground door traffic bumps. */ + readonly doorSignal: DoorSignal + + private readonly source: FactSource + private readonly installmentMs: number + private readonly batchSize: number + private readonly registry = new Map() + /** Rotates the family that leads each `advanceAll`, so repeated tiny-budget calls stay fair. */ + private roundRobinCursor = 0 + + /** @param options - See {@link ReprojectionEngineOptions}. */ + constructor(options: ReprojectionEngineOptions) { + if (!options || typeof options.source?.scan !== 'function') { + throw new Error('reprojection: a FactSource with scan(from, limit) is required') + } + const installmentMs = options.installmentMs ?? MAX_INSTALLMENT_MS + if (!(installmentMs > 0) || installmentMs > MAX_INSTALLMENT_MS) { + throw new Error( + `reprojection: installmentMs must be in (0, ${MAX_INSTALLMENT_MS}] — ` + + `${installmentMs} would let maintenance hold the doors` + ) + } + const batchSize = options.batchSize ?? DEFAULT_REPROJECTION_BATCH_SIZE + if (!Number.isInteger(batchSize) || batchSize < 1) { + throw new Error(`reprojection: batchSize must be a positive integer (got ${batchSize})`) + } + this.source = options.source + this.doorSignal = options.doorSignal ?? new DoorSignal() + this.installmentMs = installmentMs + this.batchSize = batchSize + } + + /** + * Register a projection family. Refuses a duplicate family loudly — the + * sanctioned way to replace a serving adapter is {@link swap}, never + * re-registration. + * @param adapter - The adapter that will serve this family. + */ + register(adapter: ProjectionAdapter): void { + if (!adapter || typeof adapter.family !== 'string' || adapter.family.length === 0) { + throw new Error('reprojection: adapter.family must be a non-empty string') + } + if (this.registry.has(adapter.family)) { + throw new Error( + `reprojection: family '${adapter.family}' is already registered — ` + + `replace a serving adapter via swap(), never by re-registering` + ) + } + this.registry.set(adapter.family, { + adapter, + quarantine: [], + skip: new Set(), + nextWarnAt: 1, + swapInFlight: false + }) + } + + /** + * The adapter currently serving `family` (observability — e.g. asserting + * the old adapter still serves during a swap's beside-build), or undefined + * when the family is not registered. + * @param family - The family name. + */ + getAdapter(family: string): ProjectionAdapter | undefined { + return this.registry.get(family)?.adapter + } + + /** + * This family's quarantine ledger (a defensive copy, condemnation order). + * Non-empty means one or more generations are being skipped for this + * family — the projection owner should refuse reads the skipped facts + * would have affected. + * @param family - The family name (must be registered). + */ + quarantined(family: string): QuarantineEntry[] { + return [...this.mustGet(family).quarantine] + } + + /** + * Advance one family toward the head of the fact log (or toward `upTo`), + * in installments, under a wall-clock budget, preemptible by the door + * signal. Always makes at least ONE step of progress before any budget + * check, so a zero budget still advances. + * + * @param family - The registered family to advance. + * @param options - `budgetMs` caps this call's wall time (≥ 0); `upTo` + * optionally caps the fold at a generation (inclusive). + * @returns The answer class with the adapter-stamped watermark and the + * count of facts delivered in successful applyBatch calls. + */ + async advance(family: string, options: { budgetMs: number; upTo?: number }): Promise { + const state = this.mustGet(family) + const budgetMs = options?.budgetMs + if (typeof budgetMs !== 'number' || !(budgetMs >= 0)) { + throw new Error(`reprojection: advance('${family}') requires budgetMs >= 0 (got ${budgetMs})`) + } + const start = Date.now() + const entryEpoch = this.doorSignal.epoch() + let installmentStart = start + let applied = 0 + + for (;;) { + const stepResult = await this.step(state, options.upTo) + applied += stepResult.applied + if (stepResult.done) { + return this.completed(state, applied) + } + // A bump ends the current installment immediately: yield a macrotask so + // the foreground work runs, then answer 'preempted'. + if (this.doorSignal.epoch() !== entryEpoch) { + await yieldToDoors() + return { status: 'preempted', watermark: state.adapter.watermark(), applied } + } + const t = Date.now() + if (t - start >= budgetMs) { + return { status: 'budget-exhausted', watermark: state.adapter.watermark(), applied } + } + if (t - installmentStart >= this.installmentMs) { + await yieldToDoors() + installmentStart = Date.now() + } + } + } + + /** + * Advance EVERY registered family toward the head under one shared budget, + * round-robin at batch granularity — one batch per family per turn — so no + * family starves behind another's backlog. The leading family rotates + * across calls, keeping repeated tiny-budget calls fair too. + * + * @param options - `budgetMs` caps this call's total wall time (≥ 0). + * @returns Per-family results. Families still mid-stream when the budget + * ran out (or a door bumped) report `'budget-exhausted'` (or + * `'preempted'`) at their current watermark. + */ + async advanceAll(options: { budgetMs: number }): Promise> { + const budgetMs = options?.budgetMs + if (typeof budgetMs !== 'number' || !(budgetMs >= 0)) { + throw new Error(`reprojection: advanceAll requires budgetMs >= 0 (got ${budgetMs})`) + } + const start = Date.now() + const entryEpoch = this.doorSignal.epoch() + let installmentStart = start + + const all = [...this.registry.values()] + const results: Record = {} + const appliedBy = new Map() + if (all.length === 0) return results + + // Rotate the leader across calls (fairness across repeated small budgets). + const offset = this.roundRobinCursor % all.length + this.roundRobinCursor = (this.roundRobinCursor + 1) % all.length + let queue = [...all.slice(offset), ...all.slice(0, offset)] + for (const s of queue) appliedBy.set(s.adapter.family, 0) + + const finish = ( + status: 'preempted' | 'budget-exhausted', + remaining: FamilyState[] + ): Record => { + for (const s of remaining) { + results[s.adapter.family] = { + status, + watermark: s.adapter.watermark(), + applied: appliedBy.get(s.adapter.family) ?? 0 + } + } + return results + } + + while (queue.length > 0) { + const survivors: FamilyState[] = [] + for (let i = 0; i < queue.length; i++) { + const s = queue[i] + const fam = s.adapter.family + const stepResult = await this.step(s, undefined) + appliedBy.set(fam, (appliedBy.get(fam) ?? 0) + stepResult.applied) + if (stepResult.done) { + results[fam] = this.completed(s, appliedBy.get(fam) ?? 0) + } else { + survivors.push(s) + } + const remaining = [...survivors, ...queue.slice(i + 1)] + if (this.doorSignal.epoch() !== entryEpoch) { + await yieldToDoors() + return finish('preempted', remaining) + } + const t = Date.now() + if (t - start >= budgetMs && remaining.length > 0) { + return finish('budget-exhausted', remaining) + } + if (t - installmentStart >= this.installmentMs) { + await yieldToDoors() + installmentStart = Date.now() + } + } + queue = survivors + } + return results + } + + /** + * Replace a family's adapter by BUILD-BESIDE: the old adapter keeps serving + * (stays registered, its watermark untouched) while the replacement folds + * from its own watermark (null/0 for a fresh build) to parity with the head + * of the fact log. The flip is ATOMIC — a single registry pointer swap with + * no await between the parity check and the assignment — and the losing + * adapter's `discard()` is called after the flip. + * + * SINGLE-FLIGHT: a second concurrent swap on the same family throws a + * typed {@link SwapInFlightError}. The build yields at installment + * boundaries like any fold (doors interleave), but it is never + * preemption-aborted — a swap under steady foreground traffic still + * completes. + * + * On a build failure the partially-built replacement is discarded + * (best-effort, narrated if that also fails) and the error propagates; the + * old adapter keeps serving untouched. + * + * @param family - The registered family to replace. + * @param buildAdapter - Factory for the replacement adapter (same family). + * @returns The new adapter's watermark at the flip and the facts delivered + * during the build. + */ + async swap(family: string, buildAdapter: () => Promise): Promise { + const state = this.mustGet(family) + if (state.swapInFlight) throw new SwapInFlightError(family) + state.swapInFlight = true + try { + const next = await buildAdapter() + if (!next || next.family !== family) { + throw new Error( + `reprojection: swap('${family}') built an adapter for family ` + + `'${next?.family}' — the replacement must serve the same family` + ) + } + const build: FoldState = { adapter: next, quarantine: [], skip: new Set(), nextWarnAt: 1 } + let applied = 0 + let installmentStart = Date.now() + let stalledDoneAt: number | null = null + + try { + for (;;) { + const stepResult = await this.step(build, undefined) + applied += stepResult.applied + if (stepResult.applied > 0) stalledDoneAt = null + if (stepResult.done) { + // Parity: the build just saw an empty scan (caught up to the head + // as of that call). The serving adapter can never be beyond the + // head, so newWm >= oldWm holds — verified loudly, never assumed. + const oldWm = state.adapter.watermark() ?? 0 + const newWm = next.watermark() ?? 0 + if (newWm >= oldWm) break + if (stalledDoneAt === newWm) { + throw new Error( + `reprojection: swap('${family}') build is caught up to the head at ` + + `generation ${newWm} but the serving adapter claims watermark ${oldWm} — ` + + `the serving stamp is beyond the fact log; refusing to flip` + ) + } + // The head moved past our scan (a concurrent fold advanced the + // serving adapter) — keep folding to the new head. + stalledDoneAt = newWm + } + if (Date.now() - installmentStart >= this.installmentMs) { + await yieldToDoors() + installmentStart = Date.now() + } + } + } catch (err) { + await next.discard().catch((cleanupErr) => { + prodLog.warn( + `reprojection: swap('${family}') build failed AND the failed build's discard() ` + + `also failed — its artifact may be orphaned`, + cleanupErr + ) + }) + throw err + } + + // THE FLIP — atomic by construction: no await between the parity check + // above and this pointer swap; readers see the old adapter until this + // line and the new one from it. + const losing = state.adapter + state.adapter = next + state.quarantine = build.quarantine + state.skip = build.skip + state.nextWarnAt = build.nextWarnAt + + try { + await losing.discard() + } catch (discardErr) { + // The flip already happened and the new adapter serves; the only loss + // is the loser's orphaned artifact — said out loud, never rethrown as + // a false swap failure. + prodLog.warn( + `reprojection: swap('${family}') completed but the losing adapter's discard() ` + + `failed — its artifact may be orphaned`, + discardErr + ) + } + return { watermark: next.watermark(), applied } + } finally { + state.swapInFlight = false + } + } + + /** One fold step: scan a batch above the watermark, filter quarantined generations, apply. */ + private async step(state: FoldState, upTo: number | undefined): Promise<{ done: boolean; applied: number }> { + const from = state.adapter.watermark() ?? 0 + if (upTo !== undefined && from >= upTo) return { done: true, applied: 0 } + let facts = await this.source.scan(from, this.batchSize) + if (facts.length === 0) return { done: true, applied: 0 } + if (upTo !== undefined) { + facts = facts.filter((f) => f.generation <= upTo) + if (facts.length === 0) return { done: true, applied: 0 } + } + const batchUpTo = facts[facts.length - 1].generation + const toApply = state.skip.size > 0 ? facts.filter((f) => !state.skip.has(f.generation)) : facts + try { + await state.adapter.applyBatch(toApply, batchUpTo) + } catch (err) { + if (err instanceof ProjectionApplyError) { + this.recordQuarantine(state, err) + return { done: false, applied: 0 } + } + throw err // unknown failure ≠ poison record — abort the advance loudly + } + // Anti-spin guard: a successful applyBatch that never advances the stamp + // would re-serve the same window forever. Refuse loudly instead. + const after = state.adapter.watermark() ?? 0 + if (after <= from) { + throw new Error( + `reprojection: family '${state.adapter.family}' applyBatch succeeded up to ` + + `generation ${batchUpTo} but the watermark did not advance past ${from} — ` + + `the adapter is not stamping; refusing to spin` + ) + } + return { done: false, applied: toApply.length } + } + + /** Ledger a typed apply failure, skip its generation, narrate per-doubling. */ + private recordQuarantine(state: FoldState, err: ProjectionApplyError): void { + if (!Number.isFinite(err.generation)) { + throw new Error( + `reprojection: family '${state.adapter.family}' threw ProjectionApplyError with a ` + + `non-finite generation (${err.generation}) — cannot quarantine; aborting the advance` + ) + } + if (state.skip.has(err.generation)) { + throw new Error( + `reprojection: family '${state.adapter.family}' threw ProjectionApplyError for ` + + `generation ${err.generation}, which is ALREADY quarantined and was not in the ` + + `batch — the adapter is misreporting; aborting the advance` + ) + } + state.skip.add(err.generation) + state.quarantine.push({ generation: err.generation, error: err, at: Date.now() }) + const n = state.quarantine.length + if (n === state.nextWarnAt) { + state.nextWarnAt *= 2 + prodLog.warn( + `reprojection: family '${state.adapter.family}' quarantined generation ` + + `${err.generation} (${n} quarantined total) — the fact is skipped for this family ` + + `and ledgered; reads it would have affected should be refused by the owner`, + err.cause + ) + } + } + + /** A window completed: 'caught-up' with a clean ledger, 'quarantined' otherwise. */ + private completed(state: FoldState, applied: number): AdvanceResult { + return { + status: state.quarantine.length > 0 ? 'quarantined' : 'caught-up', + watermark: state.adapter.watermark(), + applied + } + } + + /** The registered family state, or a loud refusal. */ + private mustGet(family: string): FamilyState { + const state = this.registry.get(family) + if (!state) throw new Error(`reprojection: family '${family}' is not registered`) + return state + } +} diff --git a/tests/integration/reprojection-doors-open.test.ts b/tests/integration/reprojection-doors-open.test.ts new file mode 100644 index 00000000..343bc536 --- /dev/null +++ b/tests/integration/reprojection-doors-open.test.ts @@ -0,0 +1,257 @@ +/** + * @module tests/integration/reprojection-doors-open + * @description The reprojection engine against a REAL brain on filesystem + * storage: a toy secondary projection (bucket counts with its own watermark + * artifact, stamp-after-data per src/utils/projectionWatermark.ts) folds the + * brain's committed facts through the engine, wired with the callback-form + * {@link FactLogSource} over `brain.scanFacts`. + * + * Proves the three doors-open rows: + * (i) folding to caught-up matches ground-truth counts; + * (ii) mid-fold, `find()` and `get()` still answer, and a door bump + * preempts the advance at the next boundary (mechanism-pinned via + * batch counts, not wall-clock); + * (iii) a crash mid-fold (abandon; reopen; re-advance) resumes from the + * durable stamp — never refolds from zero. + */ +import { describe, it, expect, beforeAll, afterAll } from 'vitest' +import { mkdtempSync, mkdirSync, readFileSync, writeFileSync, rmSync, existsSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { Brainy } from '../../src/index.js' +import { + ReprojectionEngine, + type ProjectionAdapter +} from '../../src/reprojection/reprojectionEngine.js' +import { FactLogSource } from '../../src/reprojection/factLogSource.js' +import { makeProjectionStamp, readStampedWatermark } from '../../src/utils/projectionWatermark.js' +import type { CommitFact } from '../../src/db/factLog.js' + +/** 50 rows, 5 buckets, 10 each. */ +const ROWS = 50 +const BUCKETS = 5 +const GROUND_TRUTH: Record = { b0: 10, b1: 10, b2: 10, b3: 10, b4: 10 } + +/** + * The toy secondary projection: latest bucket per entity id, persisted as a + * data file plus a SEPARATE stamp artifact written stamp-after-data via the + * shared projectionWatermark helpers. Idempotent by construction (latest- + * state per id), so at-least-once redelivery on resume is harmless. + */ +class BucketCountProjection implements ProjectionAdapter { + readonly family = 'bucket-counts' + /** Every generation this INSTANCE applied — the refold detector for (iii). */ + readonly appliedGenerations: number[] = [] + private latest: Map + private wm: number | null + + private constructor( + private readonly dir: string, + wm: number | null, + latest: Map + ) { + this.wm = wm + this.latest = latest + } + + /** Load from the artifact dir — data is trusted only under a valid stamp. */ + static async open(dir: string): Promise { + mkdirSync(dir, { recursive: true }) + const stampPath = join(dir, 'stamp.json') + const dataPath = join(dir, 'data.json') + let wm: number | null = null + if (existsSync(stampPath)) { + wm = readStampedWatermark(JSON.parse(readFileSync(stampPath, 'utf8'))) + } + const latest = new Map( + wm !== null && existsSync(dataPath) + ? (JSON.parse(readFileSync(dataPath, 'utf8')) as Array<[string, string | null]>) + : [] + ) + return new BucketCountProjection(dir, wm, latest) + } + + /** Non-null bucket tallies from the latest-state map. */ + counts(): Record { + const out: Record = {} + for (const bucket of this.latest.values()) { + if (bucket !== null) out[bucket] = (out[bucket] ?? 0) + 1 + } + return out + } + + watermark(): number | null { + return this.wm + } + + async applyBatch(facts: CommitFact[], upTo: number): Promise { + for (const fact of facts) { + this.appliedGenerations.push(fact.generation) + for (const op of fact.ops) { + if (op.kind !== 'noun') continue + if (op.record === null) { + this.latest.set(op.id, null) // tombstone + continue + } + // The stored noun record nests user metadata under `.metadata`. + const stored = op.record.metadata as Record | null + const user = (stored?.metadata ?? stored) as Record | null + const bucket = typeof user?.bucket === 'string' ? user.bucket : null + this.latest.set(op.id, bucket) + } + } + // Durability THEN stamp — the projectionWatermark law. + writeFileSync(join(this.dir, 'data.json'), JSON.stringify([...this.latest])) + writeFileSync(join(this.dir, 'stamp.json'), JSON.stringify(makeProjectionStamp(upTo))) + this.wm = upTo + } + + async discard(): Promise { + rmSync(this.dir, { recursive: true, force: true }) + } +} + +describe('reprojection doors-open — a real brain, a toy secondary projection', () => { + let brainDir: string + let projRoot: string + let brain: Brainy + const ids: string[] = [] + + const openBrain = async (dir: string): Promise => { + const b = new Brainy({ + storage: { type: 'filesystem', path: dir }, + requireSubtype: false, + silent: true, + dimensions: 384 + }) + await b.init() + return b + } + + /** + * The production wiring, callback form: the engine's `from` is an EXCLUSIVE + * lower bound, `scanFacts` bounds are inclusive — hence `from + 1`; the + * first batch is returned and the handle closed (short batches at segment + * boundaries are legal — only EMPTY means caught up). + */ + const sourceFor = (b: Brainy): FactLogSource => + new FactLogSource(async (from, limit) => { + const scan = b.scanFacts({ fromGeneration: from + 1, batchSize: limit }) + if (!scan) throw new Error('this brain hosts no fact log — cannot reproject') + const iterator = scan.batches() + try { + const first = await iterator.next() + return first.done ? [] : first.value.facts + } finally { + if (typeof iterator.return === 'function') await iterator.return(undefined) + } + }) + + beforeAll(async () => { + brainDir = mkdtempSync(join(tmpdir(), 'brainy-reproj-')) + projRoot = mkdtempSync(join(tmpdir(), 'brainy-reproj-artifacts-')) + brain = await openBrain(brainDir) + for (let i = 0; i < ROWS; i++) { + ids.push( + await brain.add({ + data: `record ${i} filed in bucket ${i % BUCKETS}`, + type: 'document', + metadata: { bucket: `b${i % BUCKETS}` } + }) + ) + } + }, 240_000) + + afterAll(async () => { + await brain?.close().catch(() => {}) + rmSync(brainDir, { recursive: true, force: true }) + rmSync(projRoot, { recursive: true, force: true }) + }) + + it('(i) folds to caught-up through the engine and matches ground-truth counts', async () => { + const projection = await BucketCountProjection.open(join(projRoot, 'i')) + const engine = new ReprojectionEngine({ source: sourceFor(brain), batchSize: 8 }) + engine.register(projection) + + const result = await engine.advance(projection.family, { budgetMs: 60_000 }) + + expect(result.status).toBe('caught-up') + expect(result.watermark).toBeGreaterThanOrEqual(ROWS) // one generation per add, at least + expect(result.applied).toBeGreaterThanOrEqual(ROWS) + expect(engine.quarantined(projection.family)).toEqual([]) + expect(projection.counts()).toEqual(GROUND_TRUTH) + // The stamp on disk is the adapter's own — stamped exactly at the fold head. + const reloaded = await BucketCountProjection.open(join(projRoot, 'i')) + expect(reloaded.watermark()).toBe(result.watermark) + expect(reloaded.counts()).toEqual(GROUND_TRUTH) + }) + + it('(ii) doors stay open mid-fold: find() and get() answer, and a bump preempts the advance', async () => { + const projection = await BucketCountProjection.open(join(projRoot, 'ii')) + const engine = new ReprojectionEngine({ source: sourceFor(brain), batchSize: 4 }) + engine.register(projection) + const head = brain.scanFacts()!.headGeneration + + const inFlight = engine.advance(projection.family, { budgetMs: 60_000 }) + // The read hook: foreground door traffic announces itself, then reads — + // both interleave with the running fold on the same event loop. + engine.doorSignal.bump() + const found = await brain.find({ query: 'record filed in bucket', limit: 3 }) + const got = await brain.get(ids[0]) + const result = await inFlight + + // The doors answered mid-fold. + expect(found.length).toBeGreaterThan(0) + expect(got).toBeTruthy() + const gotMeta = got!.metadata as Record | undefined + expect((gotMeta?.bucket ?? (gotMeta?.metadata as Record)?.bucket)).toBe('b0') + + // THE PREEMPTION PIN — mechanism, not wall-clock: the bump landed before + // the first installment boundary, so the advance yielded after exactly + // one batch (≤ batchSize facts), far short of the head. + expect(result.status).toBe('preempted') + expect(result.applied).toBeGreaterThan(0) + expect(result.applied).toBeLessThanOrEqual(4) + expect(projection.appliedGenerations.length).toBe(result.applied) + expect(projection.watermark()).not.toBeNull() + expect(projection.watermark()!).toBeLessThan(head) + + // Resuming folds the remainder; nothing was lost to the preemption. + const resumed = await engine.advance(projection.family, { budgetMs: 60_000 }) + expect(resumed.status).toBe('caught-up') + expect(projection.counts()).toEqual(GROUND_TRUTH) + }) + + it('(iii) crash mid-fold: reopen and re-advance resumes from the stamp, never refolds from zero', async () => { + const projDir = join(projRoot, 'iii') + const before = await BucketCountProjection.open(projDir) + const engine1 = new ReprojectionEngine({ source: sourceFor(brain), batchSize: 4 }) + engine1.register(before) + + // A zero budget folds exactly one guaranteed batch, then stops. + const partial = await engine1.advance(before.family, { budgetMs: 0 }) + expect(partial.status).toBe('budget-exhausted') + const stamped = before.watermark() + expect(stamped).not.toBeNull() + expect(stamped!).toBeGreaterThan(0) + + // CRASH: abandon the engine and adapter mid-fold; reopen the brain cold. + await brain.close() + brain = await openBrain(brainDir) + + const after = await BucketCountProjection.open(projDir) + expect(after.watermark()).toBe(stamped) // the stamp survived the crash + + const engine2 = new ReprojectionEngine({ source: sourceFor(brain), batchSize: 4 }) + engine2.register(after) + const resumed = await engine2.advance(after.family, { budgetMs: 60_000 }) + expect(resumed.status).toBe('caught-up') + + // NEVER REFOLDS FROM ZERO: every generation the resumed instance applied + // sits strictly above the crash stamp. + expect(after.appliedGenerations.length).toBeGreaterThan(0) + expect(Math.min(...after.appliedGenerations)).toBeGreaterThan(stamped!) + // And the combined state — durable prefix plus resumed fold — is exact. + expect(after.counts()).toEqual(GROUND_TRUTH) + }) +}) diff --git a/tests/unit/reprojection/reprojection-engine.test.ts b/tests/unit/reprojection/reprojection-engine.test.ts new file mode 100644 index 00000000..58d10b44 --- /dev/null +++ b/tests/unit/reprojection/reprojection-engine.test.ts @@ -0,0 +1,590 @@ +/** + * @module tests/unit/reprojection/reprojection-engine + * @description Spec-by-example for the pure-TS reprojection engine — the + * frozen contract mirrored from the native twin (a shared conformance suite + * runs against both, so the shapes pinned here are load-bearing): + * + * (a) register + advance folds a scripted source to caught-up with exact + * watermark/applied counts and adapter-owned stamping; + * (b) budget exhaustion answers mid-stream and a second advance RESUMES from + * the watermark — never a refold; + * (c) a door bump mid-advance preempts within one installment — pinned by + * MECHANISM (no further applyBatch after the bumping step), with only a + * generous wall-clock sanity bound; + * (d) advanceAll round-robins families at batch granularity — no starvation; + * (e) swap builds beside (the old adapter serves throughout), flips + * atomically at parity, refuses a concurrent swap with a typed error; + * (f) quarantine: a typed poison fact is skipped + ledgered, narration + * doubles, a NON-typed throw aborts loudly; + * (g) discard() lands on the LOSING adapter after a swap. + */ +import { describe, it, expect, vi, afterEach } from 'vitest' +import { + ReprojectionEngine, + DoorSignal, + ProjectionApplyError, + SwapInFlightError, + MAX_INSTALLMENT_MS, + type ProjectionAdapter, + type FactSource +} from '../../../src/reprojection/reprojectionEngine.js' +import { FactLogSource } from '../../../src/reprojection/factLogSource.js' +import type { CommitFact } from '../../../src/db/factLog.js' +import { prodLog } from '../../../src/utils/logger.js' + +/** Build one committed fact for a generation. */ +function fact(generation: number): CommitFact { + return { + generation, + timestamp: 1_700_000_000_000 + generation, + ops: [ + { + kind: 'noun', + id: `id-${generation}`, + record: { metadata: { n: generation }, vector: null } + } + ] + } +} + +/** A scripted FactSource over a (possibly mutable) list of generations. */ +function scriptedSource(gens: () => number[]): FactSource { + return { + async scan(from: number, limit: number): Promise { + return gens() + .filter((g) => g > from) + .sort((x, y) => x - y) + .slice(0, limit) + .map(fact) + } + } +} + +/** + * A recording in-memory adapter: stamps after data (the watermark advances + * only after a successful apply), applies idempotently (a Map keyed by + * generation), and can be scripted to poison (typed) or hard-fail (untyped) + * specific generations, or to run a hook inside applyBatch. + */ +class RecordingAdapter implements ProjectionAdapter { + readonly family: string + /** Generations per applyBatch call, in call order (empty arrays included). */ + readonly batches: number[][] = [] + /** The upTo passed to each applyBatch call, in call order. */ + readonly upTos: number[] = [] + /** Latest state per generation — idempotent under at-least-once delivery. */ + readonly state = new Map() + /** Generations that throw a typed ProjectionApplyError. */ + readonly poison = new Set() + /** Generations that throw a plain (untyped) Error. */ + readonly hardFail = new Set() + /** Runs inside applyBatch after validation, before the stamp. */ + onApply?: (gens: number[]) => void | Promise + discarded = 0 + private wm: number | null + + constructor(family: string, watermark: number | null = null) { + this.family = family + this.wm = watermark + } + + watermark(): number | null { + return this.wm + } + + async applyBatch(facts: CommitFact[], upTo: number): Promise { + for (const [i, f] of facts.entries()) { + if (this.hardFail.has(f.generation)) { + throw new Error(`disk exploded at generation ${f.generation}`) + } + if (this.poison.has(f.generation)) { + throw new ProjectionApplyError({ + generation: f.generation, + recordIndex: i, + cause: new Error(`unfoldable payload at ${f.generation}`) + }) + } + } + for (const f of facts) this.state.set(f.generation, f.ops) + const gens = facts.map((f) => f.generation) + this.batches.push(gens) + this.upTos.push(upTo) + if (this.onApply) await this.onApply(gens) + this.wm = upTo // stamp-after-data + } + + async discard(): Promise { + this.discarded++ + } +} + +const range = (from: number, to: number): number[] => + Array.from({ length: to - from + 1 }, (_, i) => from + i) + +afterEach(() => { + vi.restoreAllMocks() +}) + +describe('reprojection engine — (a) register + advance to caught-up', () => { + it('folds a scripted source in order, adapter-stamped, with exact counts', async () => { + const source = scriptedSource(() => range(1, 7)) + const engine = new ReprojectionEngine({ source, batchSize: 3 }) + const adapter = new RecordingAdapter('a') + engine.register(adapter) + + const result = await engine.advance('a', { budgetMs: 10_000 }) + + expect(result.status).toBe('caught-up') + expect(result.watermark).toBe(7) + expect(result.applied).toBe(7) + // Batch shape and the upTo handed to the adapter's own stamp. + expect(adapter.batches).toEqual([[1, 2, 3], [4, 5, 6], [7]]) + expect(adapter.upTos).toEqual([3, 6, 7]) + // The watermark is the ADAPTER's stamp — the engine never wrote one. + expect(adapter.watermark()).toBe(7) + expect(engine.getAdapter('a')).toBe(adapter) + }) + + it('honors upTo as an inclusive cap and answers caught-up at the cap', async () => { + const source = scriptedSource(() => range(1, 9)) + const engine = new ReprojectionEngine({ source, batchSize: 3 }) + const adapter = new RecordingAdapter('a') + engine.register(adapter) + + const result = await engine.advance('a', { budgetMs: 10_000, upTo: 5 }) + + expect(result.status).toBe('caught-up') + expect(result.watermark).toBe(5) + expect(result.applied).toBe(5) + expect(adapter.batches.flat()).toEqual([1, 2, 3, 4, 5]) + }) + + it('a caught-up family answers immediately with zero applied', async () => { + const source = scriptedSource(() => range(1, 4)) + const engine = new ReprojectionEngine({ source, batchSize: 10 }) + const adapter = new RecordingAdapter('a', 4) // already stamped to the head + engine.register(adapter) + + const result = await engine.advance('a', { budgetMs: 10_000 }) + + expect(result).toEqual({ status: 'caught-up', watermark: 4, applied: 0 }) + expect(adapter.batches).toEqual([]) + }) + + it('refuses duplicate registration and unregistered families loudly', async () => { + const engine = new ReprojectionEngine({ source: scriptedSource(() => []) }) + engine.register(new RecordingAdapter('a')) + expect(() => engine.register(new RecordingAdapter('a'))).toThrow(/already registered/) + await expect(engine.advance('ghost', { budgetMs: 0 })).rejects.toThrow(/not registered/) + }) +}) + +describe('reprojection engine — (b) budget exhaustion resumes, never refolds', () => { + it('returns budget-exhausted mid-stream; the next advance resumes from the watermark', async () => { + const source = scriptedSource(() => range(1, 10)) + const engine = new ReprojectionEngine({ source, batchSize: 2 }) + const adapter = new RecordingAdapter('b') + engine.register(adapter) + + // Zero budget: exactly ONE step of guaranteed progress, then the answer. + const first = await engine.advance('b', { budgetMs: 0 }) + expect(first.status).toBe('budget-exhausted') + expect(first.watermark).toBe(2) + expect(first.applied).toBe(2) + expect(adapter.batches).toEqual([[1, 2]]) + + // The second advance RESUMES from the stamp — its first batch starts at 3. + const second = await engine.advance('b', { budgetMs: 10_000 }) + expect(second.status).toBe('caught-up') + expect(second.watermark).toBe(10) + expect(second.applied).toBe(8) + expect(adapter.batches[1]).toEqual([3, 4]) + // No refold: every generation delivered exactly once across both calls. + expect(adapter.batches.flat()).toEqual(range(1, 10)) + }) +}) + +describe('reprojection engine — (c) door bump preempts within one installment', () => { + it('a bump during a step yields preempted at that step boundary — no further applyBatch', async () => { + const source = scriptedSource(() => range(1, 12)) + const engine = new ReprojectionEngine({ source, batchSize: 2 }) + const adapter = new RecordingAdapter('c') + adapter.onApply = (gens) => { + if (gens[0] === 3) engine.doorSignal.bump() // door traffic mid-second-batch + } + engine.register(adapter) + + const started = Date.now() + const result = await engine.advance('c', { budgetMs: 60_000 }) + const elapsed = Date.now() - started + + expect(result.status).toBe('preempted') + expect(result.watermark).toBe(4) + expect(result.applied).toBe(4) + // THE MECHANISM PIN: the batch that observed the bump was the LAST batch — + // preemption landed at the very next boundary, not after more work. + expect(adapter.batches).toEqual([[1, 2], [3, 4]]) + // Generous wall-clock sanity only (the pin above carries the contract): + // two tiny batches plus one installment boundary sit far under 5s. + expect(elapsed).toBeLessThan(5_000) + expect(MAX_INSTALLMENT_MS).toBe(50) + + // Resuming folds the rest — preemption lost nothing. + const resumed = await engine.advance('c', { budgetMs: 60_000 }) + expect(resumed.status).toBe('caught-up') + expect(resumed.watermark).toBe(12) + expect(adapter.batches.flat()).toEqual(range(1, 12)) + }) + + it('bumps are edge-triggered per advance: a stale bump never preempts', async () => { + const source = scriptedSource(() => range(1, 4)) + const doorSignal = new DoorSignal() + const engine = new ReprojectionEngine({ source, doorSignal, batchSize: 2 }) + const adapter = new RecordingAdapter('c2') + engine.register(adapter) + + doorSignal.bump() // BEFORE the advance — belongs to earlier traffic + const result = await engine.advance('c2', { budgetMs: 10_000 }) + expect(result.status).toBe('caught-up') + expect(result.watermark).toBe(4) + }) +}) + +describe('reprojection engine — (d) advanceAll round-robin fairness', () => { + it('a one-batch family is served on the first round despite a huge backlog next to it', async () => { + const source = scriptedSource(() => range(1, 40)) + const engine = new ReprojectionEngine({ source, batchSize: 5 }) + const callOrder: string[] = [] + const big = new RecordingAdapter('big') // 8 batches behind + const small = new RecordingAdapter('small', 35) // 1 batch behind + big.onApply = () => { + callOrder.push('big') + } + small.onApply = () => { + callOrder.push('small') + } + engine.register(big) + engine.register(small) + + const results = await engine.advanceAll({ budgetMs: 10_000 }) + + expect(results.big).toEqual({ status: 'caught-up', watermark: 40, applied: 40 }) + expect(results.small).toEqual({ status: 'caught-up', watermark: 40, applied: 5 }) + // Fairness pin: 'small' folded its single batch on round ONE — it never + // waited behind 'big''s backlog. + expect(callOrder[1]).toBe('small') + expect(callOrder.filter((f) => f === 'small')).toHaveLength(1) + }) + + it('two full-backlog families interleave strictly, one batch each per round', async () => { + const source = scriptedSource(() => range(1, 40)) + const engine = new ReprojectionEngine({ source, batchSize: 5 }) + const callOrder: string[] = [] + const first = new RecordingAdapter('first') + const second = new RecordingAdapter('second') + first.onApply = () => { + callOrder.push('first') + } + second.onApply = () => { + callOrder.push('second') + } + engine.register(first) + engine.register(second) + + const results = await engine.advanceAll({ budgetMs: 10_000 }) + + expect(results.first.status).toBe('caught-up') + expect(results.second.status).toBe('caught-up') + // 8 rounds × (first, second): strict alternation — neither ever ran twice + // while the other waited. + expect(callOrder).toHaveLength(16) + for (let i = 0; i < callOrder.length; i += 2) { + expect(callOrder.slice(i, i + 2)).toEqual(['first', 'second']) + } + }) + + it('budget exhaustion mid-round reports every unfinished family at its own watermark', async () => { + const source = scriptedSource(() => range(1, 40)) + const engine = new ReprojectionEngine({ source, batchSize: 5 }) + const a = new RecordingAdapter('a') + const b = new RecordingAdapter('b') + engine.register(a) + engine.register(b) + + const results = await engine.advanceAll({ budgetMs: 0 }) + + // Zero budget: the leading family gets its one guaranteed step, then the + // budget answer lands for everyone still mid-stream. + expect(results.a.status).toBe('budget-exhausted') + expect(results.b.status).toBe('budget-exhausted') + expect(results.a.applied + results.b.applied).toBeGreaterThanOrEqual(5) + // A later advanceAll resumes both to the head. + const finished = await engine.advanceAll({ budgetMs: 10_000 }) + expect(finished.a.status).toBe('caught-up') + expect(finished.b.status).toBe('caught-up') + expect(a.batches.flat()).toEqual(range(1, 40)) + expect(b.batches.flat()).toEqual(range(1, 40)) + }) +}) + +describe('reprojection engine — (e) swap: build-beside, atomic flip, single-flight', () => { + it('the old adapter serves at its own watermark throughout the build; the flip is atomic at parity', async () => { + const log = range(1, 20) + const source = scriptedSource(() => log) + const engine = new ReprojectionEngine({ source, batchSize: 4 }) + const oldAdapter = new RecordingAdapter('e') + engine.register(oldAdapter) + await engine.advance('e', { budgetMs: 10_000 }) + expect(oldAdapter.watermark()).toBe(20) + + // The log grows after the old adapter stamped — the build must reach the + // HEAD (24), not merely the old watermark (20), before the flip. + log.push(21, 22, 23, 24) + + const servingDuringBuild: Array<{ adapter: ProjectionAdapter | undefined; watermark: number | null }> = [] + let replacement!: RecordingAdapter + const result = await engine.swap('e', async () => { + replacement = new RecordingAdapter('e') + replacement.onApply = () => { + servingDuringBuild.push({ + adapter: engine.getAdapter('e'), + watermark: engine.getAdapter('e')!.watermark() + }) + } + return replacement + }) + + // Build-beside pin: EVERY mid-build observation saw the OLD adapter, + // still serving, still at its own stamp. + expect(servingDuringBuild.length).toBeGreaterThan(0) + for (const seen of servingDuringBuild) { + expect(seen.adapter).toBe(oldAdapter) + expect(seen.watermark).toBe(20) + } + // The flip: the registry now serves the replacement, at parity with head. + expect(engine.getAdapter('e')).toBe(replacement) + expect(result.watermark).toBe(24) + expect(result.applied).toBe(24) + expect(replacement.batches.flat()).toEqual(range(1, 24)) + }) + + it('a second concurrent swap on the same family refuses with the typed single-flight error', async () => { + const source = scriptedSource(() => range(1, 8)) + const engine = new ReprojectionEngine({ source, batchSize: 4 }) + engine.register(new RecordingAdapter('e2')) + + let release!: () => void + const gate = new Promise((resolve) => { + release = resolve + }) + const inFlight = engine.swap('e2', async () => { + const building = new RecordingAdapter('e2') + building.onApply = () => gate // the build parks mid-fold + return building + }) + + // While the first swap builds, a second one is refused — typed. + const refusal = await engine.swap('e2', async () => new RecordingAdapter('e2')).catch((e) => e) + expect(refusal).toBeInstanceOf(SwapInFlightError) + expect((refusal as SwapInFlightError).family).toBe('e2') + + release() + const done = await inFlight + expect(done.watermark).toBe(8) + // Single-flight released: a follow-up swap is admitted again. + const again = await engine.swap('e2', async () => new RecordingAdapter('e2')) + expect(again.watermark).toBe(8) + }) + + it('a failed build discards the partial replacement and leaves the old adapter serving', async () => { + const source = scriptedSource(() => range(1, 8)) + const engine = new ReprojectionEngine({ source, batchSize: 4 }) + const oldAdapter = new RecordingAdapter('e3') + engine.register(oldAdapter) + await engine.advance('e3', { budgetMs: 10_000 }) + + let failed!: RecordingAdapter + await expect( + engine.swap('e3', async () => { + failed = new RecordingAdapter('e3') + failed.hardFail.add(5) // an UNTYPED failure mid-build + return failed + }) + ).rejects.toThrow(/disk exploded/) + + expect(failed.discarded).toBe(1) // the partial build was cleaned up + expect(oldAdapter.discarded).toBe(0) + expect(engine.getAdapter('e3')).toBe(oldAdapter) // still serving, untouched + expect(oldAdapter.watermark()).toBe(8) + }) +}) + +describe('reprojection engine — (f) quarantine: the fourth answer class', () => { + it('a typed poison fact is skipped, ledgered, and the rest folds to quarantined', async () => { + const source = scriptedSource(() => range(1, 10)) + const engine = new ReprojectionEngine({ source, batchSize: 4 }) + const adapter = new RecordingAdapter('f') + adapter.poison.add(6) + engine.register(adapter) + + const result = await engine.advance('f', { budgetMs: 10_000 }) + + expect(result.status).toBe('quarantined') + expect(result.watermark).toBe(10) + expect(result.applied).toBe(9) // every generation but the poison + expect(adapter.batches.flat().sort((x, y) => x - y)).toEqual([1, 2, 3, 4, 5, 7, 8, 9, 10]) + expect(adapter.state.has(6)).toBe(false) + + const ledger = engine.quarantined('f') + expect(ledger).toHaveLength(1) + expect(ledger[0].generation).toBe(6) + expect(ledger[0].error).toBeInstanceOf(ProjectionApplyError) + expect(ledger[0].error.recordIndex).toBe(1) // 6 sat at index 1 of [5..8] + expect(typeof ledger[0].at).toBe('number') + }) + + it('narration doubles: warns on the 1st, 2nd, and 4th quarantine — not the 3rd', async () => { + const warnSpy = vi.spyOn(prodLog, 'warn').mockImplementation(() => {}) + const source = scriptedSource(() => range(1, 10)) + const engine = new ReprojectionEngine({ source, batchSize: 10 }) + const adapter = new RecordingAdapter('f2') + for (const g of [2, 4, 6, 8]) adapter.poison.add(g) + engine.register(adapter) + + const result = await engine.advance('f2', { budgetMs: 10_000 }) + + expect(result.status).toBe('quarantined') + expect(result.watermark).toBe(10) + expect(result.applied).toBe(6) + expect(engine.quarantined('f2').map((q) => q.generation)).toEqual([2, 4, 6, 8]) + const quarantineWarns = warnSpy.mock.calls.filter((c) => String(c[0]).includes('quarantined generation')) + // 4 entries, narrated at counts 1, 2, and 4 — the 3rd stayed quiet. + expect(quarantineWarns).toHaveLength(3) + expect(quarantineWarns.map((c) => String(c[0]))).toEqual([ + expect.stringContaining('(1 quarantined total)'), + expect.stringContaining('(2 quarantined total)'), + expect.stringContaining('(4 quarantined total)') + ]) + }) + + it('an all-poison window still advances the stamp via an empty applyBatch', async () => { + const source = scriptedSource(() => range(1, 3)) + const engine = new ReprojectionEngine({ source, batchSize: 3 }) + const adapter = new RecordingAdapter('f3') + for (const g of [1, 2, 3]) adapter.poison.add(g) + engine.register(adapter) + + const result = await engine.advance('f3', { budgetMs: 10_000 }) + + expect(result.status).toBe('quarantined') + expect(result.watermark).toBe(3) + expect(result.applied).toBe(0) + // The final call carried NO facts but a real upTo — the pure watermark + // advance past poison, stamped by the adapter itself. + expect(adapter.batches).toEqual([[]]) + expect(adapter.upTos).toEqual([3]) + expect(engine.quarantined('f3').map((q) => q.generation)).toEqual([1, 2, 3]) + }) + + it('a NON-typed throw aborts the advance loudly — unknown failure is never poison', async () => { + const source = scriptedSource(() => range(1, 8)) + const engine = new ReprojectionEngine({ source, batchSize: 4 }) + const adapter = new RecordingAdapter('f4') + adapter.hardFail.add(5) + engine.register(adapter) + + await expect(engine.advance('f4', { budgetMs: 10_000 })).rejects.toThrow(/disk exploded at generation 5/) + + expect(adapter.watermark()).toBe(4) // the clean first batch landed; nothing after + expect(engine.quarantined('f4')).toEqual([]) // no ledger entry for an unknown failure + }) + + it('an adapter re-condemning an already-quarantined generation is refused loudly', async () => { + const source = scriptedSource(() => range(1, 4)) + const engine = new ReprojectionEngine({ source, batchSize: 4 }) + // A misbehaving adapter: always blames generation 3, even once it is + // filtered out of its batches. + const adapter: ProjectionAdapter = { + family: 'f5', + watermark: () => null, + applyBatch: async () => { + throw new ProjectionApplyError({ generation: 3, cause: new Error('always 3') }) + }, + discard: async () => {} + } + engine.register(adapter) + + await expect(engine.advance('f5', { budgetMs: 10_000 })).rejects.toThrow(/ALREADY quarantined/) + expect(engine.quarantined('f5').map((q) => q.generation)).toEqual([3]) + }) + + it('an adapter that never stamps is refused loudly instead of spinning', async () => { + const source = scriptedSource(() => range(1, 4)) + const engine = new ReprojectionEngine({ source, batchSize: 2 }) + const adapter: ProjectionAdapter = { + family: 'f6', + watermark: () => null, // never advances + applyBatch: async () => {}, + discard: async () => {} + } + engine.register(adapter) + + await expect(engine.advance('f6', { budgetMs: 10_000 })).rejects.toThrow(/not stamping/) + }) +}) + +describe('reprojection engine — (g) discard lands on the losing adapter after a swap', () => { + it('the OLD adapter is discarded exactly once, after the flip; the winner is never discarded', async () => { + const source = scriptedSource(() => range(1, 6)) + const engine = new ReprojectionEngine({ source, batchSize: 3 }) + const losing = new RecordingAdapter('g') + engine.register(losing) + await engine.advance('g', { budgetMs: 10_000 }) + expect(losing.discarded).toBe(0) // serving adapters are never discarded + + let winner!: RecordingAdapter + await engine.swap('g', async () => { + winner = new RecordingAdapter('g') + winner.onApply = () => { + // Mid-build the loser still serves and is still intact. + expect(losing.discarded).toBe(0) + } + return winner + }) + + expect(losing.discarded).toBe(1) + expect(winner.discarded).toBe(0) + expect(engine.getAdapter('g')).toBe(winner) + }) +}) + +describe('FactLogSource — the production source enforces the window contract', () => { + it('delegates to the injected callback and passes clean windows through', async () => { + const calls: Array<[number, number]> = [] + const source = new FactLogSource(async (from, limit) => { + calls.push([from, limit]) + return range(from + 1, Math.min(from + limit, 5)).map(fact) + }) + const facts = await source.scan(2, 2) + expect(facts.map((f) => f.generation)).toEqual([3, 4]) + expect(calls).toEqual([[2, 2]]) + expect(await source.scan(5, 3)).toEqual([]) + }) + + it('refuses out-of-contract callbacks loudly: oversize, non-ascending, at-or-below from', async () => { + const oversize = new FactLogSource(async () => range(1, 5).map(fact)) + await expect(oversize.scan(0, 2)).rejects.toThrow(/contract violation/) + + const unsorted = new FactLogSource(async () => [fact(3), fact(2)]) + await expect(unsorted.scan(0, 10)).rejects.toThrow(/strictly ascending/) + + const stale = new FactLogSource(async () => [fact(2)]) + await expect(stale.scan(2, 10)).rejects.toThrow(/strictly ascending/) + }) + + it('validates its own window arguments', async () => { + const source = new FactLogSource(async () => []) + await expect(source.scan(-1, 5)).rejects.toThrow(/non-negative integer/) + await expect(source.scan(0, 0)).rejects.toThrow(/positive integer/) + }) +})