diff --git a/packages/db/tests/query/includes-oracle.property.test.ts b/packages/db/tests/query/includes-oracle.property.test.ts index a69fde2da..e7f7ea5d2 100644 --- a/packages/db/tests/query/includes-oracle.property.test.ts +++ b/packages/db/tests/query/includes-oracle.property.test.ts @@ -13,6 +13,12 @@ import { mockSyncCollectionOptions, withExpectedRejection, } from '../utils.js' +import { runTrace } from '../trace-runner.js' +import type { + TraceCheckpoint, + TraceDriver, + TraceProjection, +} from '../trace-runner.js' type IncludeDepth = 1 | 2 | 3 | 4 @@ -27,6 +33,11 @@ type ChildRow = RootRow & { parentGroup: number } +type SyncChange = { + type: `insert` | `update` | `delete` + value: T +} + type HistoryAction = { type: `put` | `delete` | `optimisticConfirm` | `optimisticRollback` level: 0 | IncludeDepth @@ -246,21 +257,27 @@ let nextHarnessId = 0 function createControlledCollection( name: string, initialData: Array = [], + rowUpdateMode: `partial` | `full` = `partial`, ) { const options = mockSyncCollectionOptions({ id: `${name}-${nextHarnessId++}`, getKey: (row) => row.id, initialData, }) + options.sync.rowUpdateMode = rowUpdateMode const collection = createCollection(options) + const writeBatch = (changes: ReadonlyArray>): void => { + options.utils.begin() + changes.forEach((change) => options.utils.write(change)) + options.utils.commit() + } return { collection, write(type: `insert` | `update` | `delete`, value: T): void { - options.utils.begin() - options.utils.write({ type, value }) - options.utils.commit() + writeBatch([{ type, value }]) }, + writeBatch, resolveSync(): void { options.utils.resolveSync() }, @@ -544,14 +561,18 @@ function createIncrementalQuery(depth: IncludeDepth, sources: Sources) { } } -function createSources() { +function createSources(rowUpdateMode: `partial` | `full` = `partial`) { return { - roots: createControlledCollection(`oracle-roots`), + roots: createControlledCollection( + `oracle-roots`, + [], + rowUpdateMode, + ), levels: [ - createControlledCollection(`oracle-level-1`), - createControlledCollection(`oracle-level-2`), - createControlledCollection(`oracle-level-3`), - createControlledCollection(`oracle-level-4`), + createControlledCollection(`oracle-level-1`, [], rowUpdateMode), + createControlledCollection(`oracle-level-2`, [], rowUpdateMode), + createControlledCollection(`oracle-level-3`, [], rowUpdateMode), + createControlledCollection(`oracle-level-4`, [], rowUpdateMode), ] as const, } } @@ -604,7 +625,7 @@ async function applyAction( sources: Sources, roots: Map, levels: Array>, - assertMatches: () => void, + assertMatches: TraceCheckpoint, ): Promise { if (action.level === 0) { const current = roots.get(action.id) @@ -730,32 +751,162 @@ async function cleanupSources(sources: Sources) { ) } -async function expectScenarioMatches(scenario: Scenario): Promise { - const sources = createSources() - const incremental = createIncrementalQuery(scenario.depth, sources) - const roots = new Map() - const levels = Array.from({ length: 4 }, () => new Map()) - - try { - await incremental.preload() - expect(stripVirtualProperties(incremental.toArray)).toEqual([]) - - const assertMatches = () => { - expect(stripVirtualProperties(incremental.toArray)).toEqual( - recompute(roots, levels, scenario.depth), - ) - } +type StructuralTraceContext = { + depth: IncludeDepth + sources: Sources + incremental: ReturnType + roots: Map + levels: Array> +} - for (const action of scenario.history) { - await applyAction(action, sources, roots, levels, assertMatches) - assertMatches() +function createStructuralTraceContext( + depth: IncludeDepth, + rowUpdateMode: `partial` | `full` = `partial`, +): StructuralTraceContext { + const sources = createSources(rowUpdateMode) + return { + depth, + sources, + incremental: createIncrementalQuery(depth, sources), + roots: new Map(), + levels: Array.from({ length: 4 }, () => new Map()), + } +} + +async function cleanupStructuralTrace({ + incremental, + sources, +}: StructuralTraceContext): Promise { + await incremental.cleanup() + await cleanupSources(sources) +} + +function createStructuralTraceDriver( + depth: IncludeDepth, +): TraceDriver { + return { + setup: () => createStructuralTraceContext(depth), + start: ({ incremental }) => incremental.preload(), + apply: (action, context, checkpoint) => + applyAction( + action, + context.sources, + context.roots, + context.levels, + checkpoint, + ), + cleanup: cleanupStructuralTrace, + } +} + +type FullRowBatchStep = + | { level: 0; changes: Array> } + | { level: 1; changes: Array> } + +function updateModel( + model: Map, + changes: ReadonlyArray>, +): void { + for (const change of changes) { + if (change.type === `delete`) { + model.delete(change.value.id) + } else { + model.set(change.value.id, change.value) } - } finally { - await incremental.cleanup() - await cleanupSources(sources) } } +function createFullRowBatchTraceDriver(): TraceDriver< + FullRowBatchStep, + StructuralTraceContext +> { + return { + setup: () => createStructuralTraceContext(1, `full`), + start: ({ incremental }) => incremental.preload(), + apply: (step, { sources, roots, levels }) => { + if (step.level === 0) { + sources.roots.writeBatch(step.changes) + updateModel(roots, step.changes) + return + } + + sources.levels[0].writeBatch(step.changes) + updateModel(levels[0]!, step.changes) + }, + cleanup: cleanupStructuralTrace, + } +} + +function batchRoot( + id: number, + group: number, + value = id * 10, + position = id - 1, +): RootRow { + return { id, group, value, position } +} + +function batchChild( + id: number, + parentGroup: number, + value = id * 10, + position = id - 1, +): ChildRow { + return { id, parentGroup, group: id * 10, value, position } +} + +const fullRowBatchTrace: Array = [ + { + level: 0, + changes: [ + { type: `insert`, value: batchRoot(1, 1) }, + { type: `insert`, value: batchRoot(2, 2) }, + ], + }, + { + level: 1, + changes: [ + { type: `insert`, value: batchChild(1, 1, 10, 0) }, + { type: `insert`, value: batchChild(2, 1, 20, 1) }, + { type: `insert`, value: batchChild(3, 2, 30, 0) }, + ], + }, + { + level: 1, + changes: [ + { type: `update`, value: batchChild(1, 2, 11, 0) }, + { type: `update`, value: batchChild(2, 1, 22, 1) }, + ], + }, + { + level: 1, + changes: [ + { type: `delete`, value: batchChild(3, 2, 30, 0) }, + { type: `insert`, value: batchChild(4, 2, 40, 1) }, + ], + }, +] + +const structuralProjection: TraceProjection< + StructuralTraceContext, + unknown, + Array +> = { + observe: ({ incremental }) => stripVirtualProperties(incremental.toArray), + recompute: ({ roots, levels, depth }) => recompute(roots, levels, depth), + assertEqual: (observed, expected) => { + expect(observed).toEqual(expected) + }, +} + +async function expectScenarioMatches(scenario: Scenario): Promise { + await runTrace({ + steps: scenario.history, + driver: createStructuralTraceDriver(scenario.depth), + projection: structuralProjection, + }) +} + function createMaterializeSources() { return { roots: createControlledCollection(`materialize-roots`), @@ -768,6 +919,25 @@ function createMaterializeSources() { type MaterializeSources = ReturnType +type MaterializeModels = { + roots: Map + middles: Map + shared: Map + leaves: Map +} + +type MaterializeTraceStep = + | { type: `insert`; insert: MaterializeInsert } + | { type: `incrementLeaf`; id: number } + | { type: `redirectMiddle`; id: number; sharedId: number } + +type MaterializeTraceContext = { + sharedIntermediate: boolean + sources: MaterializeSources + live: ReturnType + models: MaterializeModels +} + function createMaterializeQuery(sources: MaterializeSources) { return createLiveQueryCollection((q) => q @@ -843,12 +1013,7 @@ function insertMaterializeRow( insert: MaterializeInsert, sharedIntermediate: boolean, sources: MaterializeSources, - models: { - roots: Map - middles: Map - shared: Map - leaves: Map - }, + models: MaterializeModels, ): void { switch (insert) { case `root-1`: @@ -891,52 +1056,125 @@ async function cleanupMaterializeSources(sources: MaterializeSources) { ) } -async function expectMaterializeScenarioMatches({ - sharedIntermediate, - insertOrder, -}: MaterializeScenario): Promise { - const sources = createMaterializeSources() - const live = createMaterializeQuery(sources) - const models = { - roots: new Map(), - middles: new Map(), - shared: new Map(), - leaves: new Map(), - } +function createMaterializeTraceSteps( + insertOrder: Array, +): Array { + const steps: Array = insertOrder.map((insert) => ({ + type: `insert`, + insert, + })) - const assertMatches = () => { - expect(stripVirtualProperties(live.toArray)).toEqual( - recomputeMaterialize( - models.roots, - models.middles, - models.shared, - models.leaves, - ), - ) + for (const insert of insertOrder) { + if (insert === `leaf-1` || insert === `leaf-2`) { + steps.push({ type: `incrementLeaf`, id: insert === `leaf-1` ? 1 : 2 }) + } } - try { - await live.preload() - assertMatches() + return steps +} - for (const insert of insertOrder) { - insertMaterializeRow(insert, sharedIntermediate, sources, models) - assertMatches() - } +function createMaterializeTraceDriver( + scenarioSharedIntermediate: boolean, +): TraceDriver { + return { + setup: () => { + const sources = createMaterializeSources() + return { + sharedIntermediate: scenarioSharedIntermediate, + sources, + live: createMaterializeQuery(sources), + models: { + roots: new Map(), + middles: new Map(), + shared: new Map(), + leaves: new Map(), + }, + } + }, + start: ({ live }) => live.preload(), + apply: (step, { models, sources, sharedIntermediate }) => { + if (step.type === `insert`) { + insertMaterializeRow(step.insert, sharedIntermediate, sources, models) + return + } - for (const leaf of models.leaves.values()) { + if (step.type === `redirectMiddle`) { + const middle = models.middles.get(step.id) + if (!middle) throw new Error(`Missing middle ${step.id} in trace model`) + const updated = { ...middle, sharedId: step.sharedId } + sources.middles.write(`update`, updated) + models.middles.set(updated.id, updated) + return + } + + const leaf = models.leaves.get(step.id) + if (!leaf) throw new Error(`Missing leaf ${step.id} in trace model`) const updated = { ...leaf, value: leaf.value + 1 } sources.leaves.write(`update`, updated) models.leaves.set(updated.id, updated) - assertMatches() - } - } finally { - await live.cleanup() - await cleanupMaterializeSources(sources) + }, + cleanup: async ({ live, sources }) => { + await live.cleanup() + await cleanupMaterializeSources(sources) + }, } } +const materializeProjection: TraceProjection< + MaterializeTraceContext, + unknown, + Array +> = { + observe: ({ live }) => stripVirtualProperties(live.toArray), + recompute: ({ models }) => + recomputeMaterialize( + models.roots, + models.middles, + models.shared, + models.leaves, + ), + assertEqual: (observed, expected) => { + expect(observed).toEqual(expected) + }, +} + +async function expectMaterializeScenarioMatches({ + sharedIntermediate, + insertOrder, +}: MaterializeScenario): Promise { + await runTrace({ + steps: createMaterializeTraceSteps(insertOrder), + driver: createMaterializeTraceDriver(sharedIntermediate), + projection: materializeProjection, + }) +} + describe(`includes recompute oracle`, () => { + fcTest( + `discovered seed: nested scalar materialization follows a reference update`, + expectAssertionFailure(async () => { + await runTrace({ + steps: [ + { type: `insert`, insert: `root-1` }, + { type: `insert`, insert: `middle-1` }, + { type: `insert`, insert: `shared-1` }, + { type: `insert`, insert: `leaf-1` }, + { type: `redirectMiddle`, id: 1, sharedId: 2 }, + ], + driver: createMaterializeTraceDriver(false), + projection: materializeProjection, + }) + }), + ) + + fcTest(`matches recomputation for full-row sync batches`, async () => { + await runTrace({ + steps: fullRowBatchTrace, + driver: createFullRowBatchTraceDriver(), + projection: structuralProjection, + }) + }) + fcTest(`supports repeated optimistic rollbacks in one history`, async () => { await expectScenarioMatches({ depth: 1, diff --git a/packages/db/tests/trace-runner.test-d.ts b/packages/db/tests/trace-runner.test-d.ts new file mode 100644 index 000000000..a1e6a0704 --- /dev/null +++ b/packages/db/tests/trace-runner.test-d.ts @@ -0,0 +1,12 @@ +import type { TraceProjection } from './trace-runner.js' + +const asyncAssertionProjection: TraceProjection = { + observe: () => 0, + recompute: () => 0, + // @ts-expect-error trace assertions must finish synchronously + assertEqual: async () => { + await Promise.resolve() + }, +} + +void asyncAssertionProjection diff --git a/packages/db/tests/trace-runner.test.ts b/packages/db/tests/trace-runner.test.ts new file mode 100644 index 000000000..75f0eaab6 --- /dev/null +++ b/packages/db/tests/trace-runner.test.ts @@ -0,0 +1,145 @@ +import { describe, expect, it, vi } from 'vitest' +import { runTrace } from './trace-runner.js' +import type { TraceDriver } from './trace-runner.js' + +type Context = { + observed: number + expected: number +} + +describe(`runTrace`, () => { + it(`checks after startup, explicit checkpoints, and every step`, async () => { + const checkpoints: Array = [] + const cleanup = vi.fn() + const driver: TraceDriver = { + setup: () => ({ observed: 0, expected: 0 }), + start: (context) => { + context.observed = 1 + context.expected = 1 + }, + apply: (step, context, checkpoint) => { + context.observed += step + context.expected += step + if (step === 2) checkpoint() + }, + cleanup, + } + + await runTrace({ + steps: [2, 3], + driver, + projection: { + observe: (context) => context.observed, + recompute: (context) => context.expected, + assertEqual: (observed, expected) => { + expect(observed).toBe(expected) + checkpoints.push(observed) + }, + }, + }) + + expect(checkpoints).toEqual([1, 3, 3, 6]) + expect(cleanup).toHaveBeenCalledOnce() + }) + + it(`cleans up when a checkpoint fails`, async () => { + const cleanup = vi.fn() + + await expect( + runTrace({ + steps: [1], + driver: { + setup: () => ({ observed: 0, expected: 1 }), + apply: () => undefined, + cleanup, + }, + projection: { + observe: (context) => context.observed, + recompute: (context) => context.expected, + assertEqual: (observed, expected) => { + expect(observed).toBe(expected) + }, + }, + }), + ).rejects.toThrow() + + expect(cleanup).toHaveBeenCalledOnce() + }) + + it(`checks synchronous steps before queued microtasks run`, async () => { + await expect( + runTrace({ + steps: [1], + driver: { + setup: () => ({ observed: 0, expected: 0 }), + apply: (step, context) => { + context.expected = step + queueMicrotask(() => { + context.observed = step + }) + }, + cleanup: () => undefined, + }, + projection: { + observe: (context) => context.observed, + recompute: (context) => context.expected, + assertEqual: (observed, expected) => { + expect(observed).toBe(expected) + }, + }, + }), + ).rejects.toThrow() + }) + + it(`preserves the trace failure when cleanup also fails`, async () => { + const traceError = new Error(`trace failed`) + const cleanupError = new Error(`cleanup failed`) + + const run = runTrace({ + steps: [], + driver: { + setup: () => ({ observed: 0, expected: 1 }), + apply: () => undefined, + cleanup: () => { + throw cleanupError + }, + }, + projection: { + observe: (context) => context.observed, + recompute: (context) => context.expected, + assertEqual: () => { + throw traceError + }, + }, + }) + + await expect(run).rejects.toBe(traceError) + expect( + (traceError as Error & { suppressed?: Array }).suppressed, + ).toEqual([cleanupError]) + }) + + it(`throws cleanup failures when the trace succeeds`, async () => { + const cleanupError = new Error(`cleanup failed`) + + await expect( + runTrace({ + steps: [], + driver: { + setup: () => ({ observed: 0, expected: 0 }), + apply: () => undefined, + cleanup: () => { + throw cleanupError + }, + }, + projection: { + observe: (context) => context.observed, + recompute: (context) => context.expected, + assertEqual: (observed, expected) => { + expect(observed).toBe(expected) + }, + }, + }), + ).rejects.toBe(cleanupError) + }) +}) diff --git a/packages/db/tests/trace-runner.ts b/packages/db/tests/trace-runner.ts new file mode 100644 index 000000000..0028bb20f --- /dev/null +++ b/packages/db/tests/trace-runner.ts @@ -0,0 +1,101 @@ +type MaybePromise = T | PromiseLike + +type ErrorWithSuppressed = Error & { + suppressed?: Array +} + +export type TraceCheckpoint = () => undefined + +export type TraceDriver = { + setup: () => MaybePromise + start?: (context: TContext) => MaybePromise + apply: ( + step: TStep, + context: TContext, + checkpoint: TraceCheckpoint, + ) => MaybePromise + cleanup: (context: TContext) => MaybePromise +} + +export type TraceProjection = { + observe: (context: TContext) => TObserved + recompute: (context: TContext) => TExpected + assertEqual: (observed: TObserved, expected: TExpected) => undefined +} + +type RunTraceOptions = { + steps: ReadonlyArray + driver: TraceDriver + projection: TraceProjection +} + +function isPromiseLike(value: MaybePromise): value is PromiseLike { + return ( + value !== null && + (typeof value === `object` || typeof value === `function`) && + `then` in value && + typeof value.then === `function` + ) +} + +function attachSuppressedError(error: unknown, suppressed: unknown): void { + if (!(error instanceof Error)) return + + try { + const errorWithSuppressed = error as ErrorWithSuppressed + errorWithSuppressed.suppressed = [ + ...(errorWithSuppressed.suppressed ?? []), + suppressed, + ] + } catch { + // A frozen or otherwise immutable error must still remain the primary one. + } +} + +/** + * Drives a trace against a system and checks its observable state against an + * independent projection after startup, after each step, and at any explicit + * checkpoint requested by the driver. + */ +export async function runTrace({ + steps, + driver, + projection, +}: RunTraceOptions): Promise { + const setupResult = driver.setup() + const context = isPromiseLike(setupResult) ? await setupResult : setupResult + const checkpoint: TraceCheckpoint = () => { + projection.assertEqual( + projection.observe(context), + projection.recompute(context), + ) + return undefined + } + + let traceFailed = false + let traceError: unknown + try { + const startResult = driver.start?.(context) + if (isPromiseLike(startResult)) await startResult + checkpoint() + + for (const step of steps) { + const applyResult = driver.apply(step, context, checkpoint) + if (isPromiseLike(applyResult)) await applyResult + checkpoint() + } + } catch (error) { + traceFailed = true + traceError = error + } + + try { + const cleanupResult = driver.cleanup(context) + if (isPromiseLike(cleanupResult)) await cleanupResult + } catch (cleanupError) { + if (!traceFailed) throw cleanupError + attachSuppressedError(traceError, cleanupError) + } + + if (traceFailed) throw traceError +}