From b6444e0164aef8e0dd1134ea86e68042f95bc261 Mon Sep 17 00:00:00 2001 From: Konstantin Wohlwend Date: Fri, 3 Jul 2026 10:28:53 -0700 Subject: [PATCH] More and better Bulldozer logs --- .../scripts/bulldozer-payments-init.ts | 6 +- .../scripts/profile-bulldozer-performance.ts | 8 +- .../src/databases/bulldozer/example-schema.ts | 2 +- .../src/databases/bulldozer/index.test.ts | 12 +- .../src/databases/bulldozer/index.ts | 395 +++++++++++++++--- .../databases/bulldozer/performance.test.ts | 4 +- .../src/payments/schema/index.test.ts | 24 +- .../payments/schema/item-quantities.test.ts | 4 +- .../payments/schema/schema-test-helpers.ts | 2 +- 9 files changed, 359 insertions(+), 98 deletions(-) diff --git a/apps/backend/scripts/bulldozer-payments-init.ts b/apps/backend/scripts/bulldozer-payments-init.ts index e5f3d6f76..56027471c 100644 --- a/apps/backend/scripts/bulldozer-payments-init.ts +++ b/apps/backend/scripts/bulldozer-payments-init.ts @@ -215,6 +215,7 @@ async function backfillTable( const batch = await fetchBatch(cursor); const fetchMs = performance.now() - fetchStartedAt; if (batch.length === 0) break; + const fetchDoneAt = performance.now(); // Time only the bulldozer write (the HTTP batch request[s]), isolated from // the Prisma read above. Under --continue-on-error the per-row retry writes @@ -242,6 +243,7 @@ async function backfillTable( const reqMs = performance.now() - reqStartedAt; totalReqMs += reqMs; total += batch.length; + const writeDoneAt = performance.now(); const last = batch[batch.length - 1]; const next: Cursor = { tenancyId: last.tenancyId, id: last.id }; @@ -255,8 +257,8 @@ async function backfillTable( cursor = next; batchNumber++; - // req = bulldozer request time (the number to watch), fetch = Prisma read. - log(`[${label}] batch=${batchNumber} req=${formatDuration(reqMs)} fetch=${formatDuration(fetchMs)} rows=${batch.length} total=${total}${failed > 0 ? ` failed=${failed}` : ""} cursor=${cursor.tenancyId},${cursor.id}`); + const allDoneAt = performance.now(); + log(`[${label}] batch=${batchNumber} duration=(r:${formatDuration(fetchMs)} w:${formatDuration(writeDoneAt - fetchDoneAt)} t:${formatDuration(allDoneAt - fetchStartedAt)} rows=${batch.length} total=${total}${failed > 0 ? ` failed=${failed}` : ""} cursor=${cursor.tenancyId},${cursor.id}`); // A short page means we've hit the end; skip the extra empty fetch. if (batch.length < ctx.batchSize) break; diff --git a/apps/bulldozer-js/scripts/profile-bulldozer-performance.ts b/apps/bulldozer-js/scripts/profile-bulldozer-performance.ts index 2668cad50..72f7a38e3 100644 --- a/apps/bulldozer-js/scripts/profile-bulldozer-performance.ts +++ b/apps/bulldozer-js/scripts/profile-bulldozer-performance.ts @@ -225,10 +225,10 @@ async function main() { await db.applyRemainingMigrations(); let snapshot = (await db.getSnapshot()).snapshot; for (let i = 0; i < 8; i++) { - snapshot = await snapshot.setOrDeleteRow({ tableId: "prices", rowIdentifier: `asset-${i}`, newRowData: { asset: `asset-${i}`, usd: 100 + i } }); + snapshot = (await snapshot.setOrDeleteRow({ tableId: "prices", rowIdentifier: `asset-${i}`, newRowData: { asset: `asset-${i}`, usd: 100 + i } })).newSnapshot; } for (let i = 0; i < rowCount; i++) { - snapshot = await snapshot.setOrDeleteRow({ tableId: "events", rowIdentifier: `event-${i}`, newRowData: eventRow(i) }); + snapshot = (await snapshot.setOrDeleteRow({ tableId: "events", rowIdentifier: `event-${i}`, newRowData: eventRow(i) })).newSnapshot; } resetMetrics(); @@ -236,11 +236,11 @@ async function main() { const targetIndex = rowCount - 1; for (let iteration = 0; iteration < measuredIterations; iteration++) { const start = performance.now(); - snapshot = await snapshot.setOrDeleteRow({ + snapshot = (await snapshot.setOrDeleteRow({ tableId: "events", rowIdentifier: `event-${targetIndex}`, newRowData: eventRow(targetIndex, rowCount * 10 + iteration), - }); + })).newSnapshot; durations.push(performance.now() - start); } diff --git a/apps/bulldozer-js/src/databases/bulldozer/example-schema.ts b/apps/bulldozer-js/src/databases/bulldozer/example-schema.ts index b25fbd181..4f18128b6 100644 --- a/apps/bulldozer-js/src/databases/bulldozer/example-schema.ts +++ b/apps/bulldozer-js/src/databases/bulldozer/example-schema.ts @@ -348,7 +348,7 @@ export async function createExampleFungibleLedgerDatabase() { await db.applyRemainingMigrations(); await db.withSnapshotReplicated(async snapshot => { for (const [rowIdentifier, rowData] of Object.entries(exampleLedgerRows)) { - snapshot = await snapshot.setOrDeleteRow({ tableId: storedTableId, rowIdentifier, newRowData: rowData as unknown as PiledriverObject }); + snapshot = (await snapshot.setOrDeleteRow({ tableId: storedTableId, rowIdentifier, newRowData: rowData as unknown as PiledriverObject })).newSnapshot; } return snapshot; }); diff --git a/apps/bulldozer-js/src/databases/bulldozer/index.test.ts b/apps/bulldozer-js/src/databases/bulldozer/index.test.ts index a2a51d28c..5cdd3802b 100644 --- a/apps/bulldozer-js/src/databases/bulldozer/index.test.ts +++ b/apps/bulldozer-js/src/databases/bulldozer/index.test.ts @@ -45,8 +45,8 @@ const initializedSnapshot = async (migrations: Parameters>, tableId: string, range: Record = {}, groupKey: PiledriverObject = null) => asRows(snapshot.listRowsInGroup({ tableId, groupKey, range })); -const set = (snapshot: Awaited>, tableId: string, rowIdentifier: string, newRowData: PiledriverObject | undefined) => - snapshot.setOrDeleteRow({ tableId, rowIdentifier, newRowData }); +const set = async (snapshot: Awaited>, tableId: string, rowIdentifier: string, newRowData: PiledriverObject | undefined) => + (await snapshot.setOrDeleteRow({ tableId, rowIdentifier, newRowData })).newSnapshot; describe("Bulldozer", () => { it("persists an empty snapshot for zero migrations", async () => { @@ -410,13 +410,13 @@ describe("Bulldozer", () => { ]); expect(reducerCalls).toEqual([{ rowIdentifier: "a", trigger: null, state: 0 }]); - snapshot = await snapshot.tick(new Date(firstTrigger)); + snapshot = (await snapshot.tick(new Date(firstTrigger))).newSnapshot; expect(await rows(snapshot, "time")).toEqual([ { groupKey: null, rowIdentifier: JSON.stringify(["a", 0]), rowSortKey: null, rowData: "A:initial" }, { groupKey: null, rowIdentifier: JSON.stringify(["a", 1]), rowSortKey: null, rowData: "A:2" }, ]); - snapshot = await snapshot.tick(new Date(secondTrigger)); + snapshot = (await snapshot.tick(new Date(secondTrigger))).newSnapshot; expect(await rows(snapshot, "time")).toEqual([ { groupKey: null, rowIdentifier: JSON.stringify(["a", 0]), rowSortKey: null, rowData: "A:initial" }, { groupKey: null, rowIdentifier: JSON.stringify(["a", 1]), rowSortKey: null, rowData: "A:2" }, @@ -460,7 +460,7 @@ describe("Bulldozer", () => { ]]); snapshot = await set(snapshot, "store", "a", "A"); - snapshot = await snapshot.tick(new Date(firstTrigger)); + snapshot = (await snapshot.tick(new Date(firstTrigger))).newSnapshot; expect(await rows(snapshot, "time")).toEqual([ { groupKey: null, rowIdentifier: JSON.stringify(["a", 0]), rowSortKey: null, rowData: "A:initial" }, { groupKey: null, rowIdentifier: JSON.stringify(["a", 1]), rowSortKey: null, rowData: "A:2" }, @@ -477,7 +477,7 @@ describe("Bulldozer", () => { // The hook's returned state and trigger drive subsequent timed steps: the next tick appends // (using the carried-over state) instead of replaying history. - snapshot = await snapshot.tick(new Date(secondTrigger)); + snapshot = (await snapshot.tick(new Date(secondTrigger))).newSnapshot; expect(await rows(snapshot, "time")).toEqual([ { groupKey: null, rowIdentifier: JSON.stringify(["a", 0]), rowSortKey: null, rowData: "A:initial" }, { groupKey: null, rowIdentifier: JSON.stringify(["a", 1]), rowSortKey: null, rowData: "A:2" }, diff --git a/apps/bulldozer-js/src/databases/bulldozer/index.ts b/apps/bulldozer-js/src/databases/bulldozer/index.ts index c286ef9e3..022e75152 100644 --- a/apps/bulldozer-js/src/databases/bulldozer/index.ts +++ b/apps/bulldozer-js/src/databases/bulldozer/index.ts @@ -1,4 +1,5 @@ import { isShallowEqual } from "@hexclave/shared/dist/utils/arrays"; +import { inspect } from "node:util"; import { traceSpan } from "../../otel.js"; import { DatabaseSeq } from "../index.js"; import type { LowLevelDatabaseDebugSnapshot } from "../low-level/index.js"; @@ -74,11 +75,157 @@ type TableChanges = { }[], }; type GroupChanges = Omit; +export type TableChangesDebugInfo = { + addedRows: number, + modifiedRows: number, + deletedRows: number, + addedGroups: number, + deletedGroups: number, + rowChanges: number, + groupChanges: number, +}; +export type BulldozerTableMutationDebugInfo = { + tableId: string, + phase: "source" | "downstream", + durationMs: number, + inputChangeCountsByInputTable: Record, + outputChangeCounts: TableChangesDebugInfo, +}; +export type BulldozerAffectedTableDebugInfo = { + tableId: string, + operationCount: number, + sourceOperationCount: number, + emitInputChangesOperationCount: number, + totalDurationMs: number, + sourceDurationMs: number, + emitInputChangesDurationMs: number, + inputChangeCountsByInputTable: Record, + totalInputChangeCounts: TableChangesDebugInfo, + outputChangeCounts: TableChangesDebugInfo, +}; +export type BulldozerSnapshotMutationDebugInfo = { + operation: "setOrDeleteRow" | "setOrDeleteRows" | "tick" | "applyTableMutation", + sourceTableId?: string, + rowsSetOrDeleted: number, + durationMs: number, + tableOperations: BulldozerTableMutationDebugInfo[], + affectedTableIds: string[], + affectedTables: Record, + totalOutputChangeCounts: TableChangesDebugInfo, +}; +export type BulldozerSnapshotMutationResult = { + newSnapshot: BulldozerDatabaseSnapshot, + debugInfo: BulldozerSnapshotMutationDebugInfo, +}; function appendAll(target: T[], values: Iterable) { for (const value of values) target.push(value); } +function tableChangesDebugInfo(changes: TableChanges): TableChangesDebugInfo { + return { + addedRows: changes.addedRows.length, + modifiedRows: changes.modifiedRows.length, + deletedRows: changes.deletedRows.length, + addedGroups: changes.addedGroups.length, + deletedGroups: changes.deletedGroups.length, + rowChanges: changes.addedRows.length + changes.modifiedRows.length + changes.deletedRows.length, + groupChanges: changes.addedGroups.length + changes.deletedGroups.length, + }; +} + +function emptyTableChangesDebugInfo(): TableChangesDebugInfo { + return { + addedRows: 0, + modifiedRows: 0, + deletedRows: 0, + addedGroups: 0, + deletedGroups: 0, + rowChanges: 0, + groupChanges: 0, + }; +} + +function mergeTableChangesDebugInfo(target: TableChangesDebugInfo, value: TableChangesDebugInfo) { + target.addedRows += value.addedRows; + target.modifiedRows += value.modifiedRows; + target.deletedRows += value.deletedRows; + target.addedGroups += value.addedGroups; + target.deletedGroups += value.deletedGroups; + target.rowChanges += value.rowChanges; + target.groupChanges += value.groupChanges; +} + +function inputChangeCountsByInputTable(changes: Record): Record { + return Object.fromEntries(Object.entries(changes).map(([inputTableKey, tableChanges]) => [inputTableKey, tableChangesDebugInfo(tableChanges)])); +} + +function emptyAffectedTableDebugInfo(tableId: string): BulldozerAffectedTableDebugInfo { + return { + tableId, + operationCount: 0, + sourceOperationCount: 0, + emitInputChangesOperationCount: 0, + totalDurationMs: 0, + sourceDurationMs: 0, + emitInputChangesDurationMs: 0, + inputChangeCountsByInputTable: {}, + totalInputChangeCounts: emptyTableChangesDebugInfo(), + outputChangeCounts: emptyTableChangesDebugInfo(), + }; +} + +function mergeInputChangeCounts( + target: Record, + totalTarget: TableChangesDebugInfo, + value: Record, +) { + for (const [inputTableKey, counts] of Object.entries(value)) { + target[inputTableKey] ??= emptyTableChangesDebugInfo(); + mergeTableChangesDebugInfo(target[inputTableKey], counts); + mergeTableChangesDebugInfo(totalTarget, counts); + } +} + +function affectedTablesDebugInfo(tableOperations: BulldozerTableMutationDebugInfo[]): Record { + const result = new Map(); + for (const operation of tableOperations) { + let table = result.get(operation.tableId); + if (table === undefined) { + table = emptyAffectedTableDebugInfo(operation.tableId); + result.set(operation.tableId, table); + } + + table.operationCount++; + table.totalDurationMs += operation.durationMs; + if (operation.phase === "source") { + table.sourceOperationCount++; + table.sourceDurationMs += operation.durationMs; + } else { + table.emitInputChangesOperationCount++; + table.emitInputChangesDurationMs += operation.durationMs; + } + mergeInputChangeCounts(table.inputChangeCountsByInputTable, table.totalInputChangeCounts, operation.inputChangeCountsByInputTable); + mergeTableChangesDebugInfo(table.outputChangeCounts, operation.outputChangeCounts); + } + return Object.fromEntries(result); +} + +function logSnapshotMutationDebugInfo(value: { + operation: BulldozerSnapshotMutationDebugInfo["operation"], + tableId: string | null, + rowsSetOrDeleted: number, + debugInfo: BulldozerSnapshotMutationDebugInfo, +}) { + if (value.rowsSetOrDeleted <= 0) return; + console.debug("bulldozer-js snapshot mutation", inspect(value, { + depth: null, + colors: false, + maxArrayLength: null, + breakLength: 160, + })); +} + function validateTableChanges(changes: TableChanges, context: string) { const deletedGroupKeys = new Set(changes.deletedGroups.map(group => canonicalGroupKeyString(group.groupKey))); const readdedGroup = changes.addedGroups.find(group => deletedGroupKeys.has(canonicalGroupKeyString(group.groupKey))); @@ -381,17 +528,28 @@ class BulldozerDatabaseSnapshot { tableId: string, rowIdentifier: string, newRowData: PiledriverObject | undefined, - }): Promise { + }): Promise { if (!(options.tableId in this.tablesState.tables)) throw new Error(`Table ${options.tableId} does not exist`); const setOrDeleteRow = this.tablesState.tables[options.tableId].table.setOrDeleteRow; if (!setOrDeleteRow) throw new Error("Table is not mutable"); - return await this._applyTableMutation(options.tableId, ({ serializedTable, inputTables }) => setOrDeleteRow({ - serializedTable, - inputTables, - rowIdentifier: options.rowIdentifier, - newRowData: options.newRowData, - })); + const result = await this._applyTableMutation({ + operation: "setOrDeleteRow", + tableId: options.tableId, + mutate: ({ serializedTable, inputTables }) => setOrDeleteRow({ + serializedTable, + inputTables, + rowIdentifier: options.rowIdentifier, + newRowData: options.newRowData, + }), + }); + logSnapshotMutationDebugInfo({ + operation: "setOrDeleteRow", + tableId: options.tableId, + rowsSetOrDeleted: result.debugInfo.rowsSetOrDeleted, + debugInfo: result.debugInfo, + }); + return result; } /** @@ -414,11 +572,30 @@ class BulldozerDatabaseSnapshot { async setOrDeleteRows(options: { tableId: string, rows: { rowIdentifier: string, newRowData: PiledriverObject | undefined }[], - }): Promise { + }): Promise { if (!(options.tableId in this.tablesState.tables)) throw new Error(`Table ${options.tableId} does not exist`); const setOrDeleteRow = this.tablesState.tables[options.tableId].table.setOrDeleteRow; if (!setOrDeleteRow) throw new Error("Table is not mutable"); - if (options.rows.length === 0) return this; + if (options.rows.length === 0) { + const debugInfo: BulldozerSnapshotMutationDebugInfo = { + operation: "setOrDeleteRows", + sourceTableId: options.tableId, + rowsSetOrDeleted: 0, + durationMs: 0, + tableOperations: [], + affectedTableIds: [], + affectedTables: {}, + totalOutputChangeCounts: emptyTableChangesDebugInfo(), + }; + const result = { newSnapshot: this, debugInfo }; + logSnapshotMutationDebugInfo({ + operation: "setOrDeleteRows", + tableId: options.tableId, + rowsSetOrDeleted: 0, + debugInfo, + }); + return result; + } const seenIdentifiers = new Set(); for (const row of options.rows) { @@ -426,93 +603,174 @@ class BulldozerDatabaseSnapshot { seenIdentifiers.add(row.rowIdentifier); } - return await this._applyTableMutation(options.tableId, async ({ serializedTable, inputTables }) => { - let currentSerializedTable = serializedTable; - const combined: TableChanges = { addedRows: [], modifiedRows: [], deletedRows: [], addedGroups: [], deletedGroups: [] }; - for (const row of options.rows) { - const result = await setOrDeleteRow({ - serializedTable: currentSerializedTable, - inputTables, - rowIdentifier: row.rowIdentifier, - newRowData: row.newRowData, - }); - currentSerializedTable = result.newSerializedTable; - mergeTableChanges(combined, result.outputChanges); - } - normalizeGroupLifecycle(combined); - return { newSerializedTable: currentSerializedTable, outputChanges: combined }; + const result = await this._applyTableMutation({ + operation: "setOrDeleteRows", + tableId: options.tableId, + mutate: async ({ serializedTable, inputTables }) => { + let currentSerializedTable = serializedTable; + const combined: TableChanges = { addedRows: [], modifiedRows: [], deletedRows: [], addedGroups: [], deletedGroups: [] }; + for (const row of options.rows) { + const result = await setOrDeleteRow({ + serializedTable: currentSerializedTable, + inputTables, + rowIdentifier: row.rowIdentifier, + newRowData: row.newRowData, + }); + currentSerializedTable = result.newSerializedTable; + mergeTableChanges(combined, result.outputChanges); + } + normalizeGroupLifecycle(combined); + return { newSerializedTable: currentSerializedTable, outputChanges: combined }; + }, }); + logSnapshotMutationDebugInfo({ + operation: "setOrDeleteRows", + tableId: options.tableId, + rowsSetOrDeleted: result.debugInfo.rowsSetOrDeleted, + debugInfo: result.debugInfo, + }); + return result; } - async tick(now: Date): Promise { + async tick(now: Date): Promise { + const startedAt = performance.now(); let snapshot: BulldozerDatabaseSnapshot = this; + const tableOperations: BulldozerTableMutationDebugInfo[] = []; + const affectedTableIds = new Set(); + const totalOutputChangeCounts = emptyTableChangesDebugInfo(); + let rowsSetOrDeleted = 0; for (const [tableId, tableState] of Object.entries(this.tablesState.tables)) { const tick = tableState.table.tick; if (!tick) continue; - snapshot = await snapshot._applyTableMutation(tableId, ({ serializedTable, inputTables }) => tick({ serializedTable, inputTables, now })); + const result = await snapshot._applyTableMutation({ + operation: "tick", + tableId, + mutate: ({ serializedTable, inputTables }) => tick({ serializedTable, inputTables, now }), + }); + snapshot = result.newSnapshot; + tableOperations.push(...result.debugInfo.tableOperations); + for (const affectedTableId of result.debugInfo.affectedTableIds) affectedTableIds.add(affectedTableId); + mergeTableChangesDebugInfo(totalOutputChangeCounts, result.debugInfo.totalOutputChangeCounts); + rowsSetOrDeleted += result.debugInfo.rowsSetOrDeleted; } - return snapshot; + const debugInfo: BulldozerSnapshotMutationDebugInfo = { + operation: "tick", + rowsSetOrDeleted, + durationMs: performance.now() - startedAt, + tableOperations, + affectedTableIds: [...affectedTableIds], + affectedTables: affectedTablesDebugInfo(tableOperations), + totalOutputChangeCounts, + }; + logSnapshotMutationDebugInfo({ + operation: "tick", + tableId: null, + rowsSetOrDeleted, + debugInfo, + }); + return { newSnapshot: snapshot, debugInfo }; } - private async _applyTableMutation( + private async _applyTableMutation(options: { + operation: BulldozerSnapshotMutationDebugInfo["operation"], tableId: string, mutate: (options: { serializedTable: PiledriverObject, inputTables: Record, }) => Promise<{ newSerializedTable: PiledriverObject, outputChanges: TableChanges }>, - ): Promise { - return await traceSpan({ description: "bulldozer-js.bulldozer.applyTableMutation", attributes: { "bulldozer.table_id": tableId } }, async () => { + }): Promise { + const startedAt = performance.now(); + return await traceSpan({ description: "bulldozer-js.bulldozer.applyTableMutation", attributes: { "bulldozer.table_id": options.tableId } }, async () => { const tablesState = this.tablesState; const serializedTables = { ...this.serialized.serializedTables }; const pending = new Map>(); const remainingInputs = new Map(); + const tableOperations: BulldozerTableMutationDebugInfo[] = []; + const affectedTableIds = new Set(); + const totalOutputChangeCounts = emptyTableChangesDebugInfo(); const inputTables = (id: string) => createInputTables(tablesState.tables, inputId => serializedTables[inputId], id); const emptyChanges = (): TableChanges => ({ addedRows: [], modifiedRows: [], deletedRows: [], addedGroups: [], deletedGroups: [] }); const hasChanges = (changes: TableChanges) => changes.addedRows.length || changes.modifiedRows.length || changes.deletedRows.length || changes.addedGroups.length || changes.deletedGroups.length; const addPending = (tableId: string, inputTableKey: string, changes: TableChanges) => { if (!hasChanges(changes)) return; - pending.set(tableId, { ...pending.get(tableId), [inputTableKey]: changes }); + pending.set(tableId, { ...pending.get(tableId), [inputTableKey]: changes }); }; - for (const queue = [tableId], seen = new Set([tableId]); queue.length;) { + for (const queue = [options.tableId], seen = new Set([options.tableId]); queue.length;) { for (const outputTable of tablesState.tables[queue.shift()!].outputTables) { - remainingInputs.set(outputTable.tableId, (remainingInputs.get(outputTable.tableId) ?? 0) + 1); - if (!seen.has(outputTable.tableId)) { - seen.add(outputTable.tableId); + remainingInputs.set(outputTable.tableId, (remainingInputs.get(outputTable.tableId) ?? 0) + 1); + if (!seen.has(outputTable.tableId)) { + seen.add(outputTable.tableId); + queue.push(outputTable.tableId); + } + } + } + + const sourceStartedAt = performance.now(); + const first = await traceSpan({ description: "bulldozer-js.bulldozer.mutateSourceTable", attributes: { "bulldozer.table_id": options.tableId } }, async () => await options.mutate({ serializedTable: serializedTables[options.tableId], inputTables: inputTables(options.tableId) })); + validateTableChanges(first.outputChanges, `Table ${options.tableId} output`); + const sourceOutputChangeCounts = tableChangesDebugInfo(first.outputChanges); + tableOperations.push({ + tableId: options.tableId, + phase: "source", + durationMs: performance.now() - sourceStartedAt, + inputChangeCountsByInputTable: {}, + outputChangeCounts: sourceOutputChangeCounts, + }); + affectedTableIds.add(options.tableId); + mergeTableChangesDebugInfo(totalOutputChangeCounts, sourceOutputChangeCounts); + serializedTables[options.tableId] = first.newSerializedTable; + for (const outputTable of tablesState.tables[options.tableId].outputTables) addPending(outputTable.tableId, outputTable.inputTableKey, first.outputChanges); + + for (const queue = tablesState.tables[options.tableId].outputTables.map(outputTable => outputTable.tableId); queue.length;) { + const downstreamTableId = queue.shift()!; + const left = (remainingInputs.get(downstreamTableId) ?? 1) - 1; + remainingInputs.set(downstreamTableId, left); + if (left > 0) continue; + const table = tablesState.tables[downstreamTableId]; + const changes = pending.get(downstreamTableId); + if (changes) { + const normalizedChanges = Object.fromEntries(Object.keys(table.inputTableIds).map(inputTableKey => [inputTableKey, changes[inputTableKey] ?? emptyChanges()])); + const downstreamStartedAt = performance.now(); + const result = await traceSpan({ description: "bulldozer-js.bulldozer.emitInputChanges", attributes: { "bulldozer.table_id": downstreamTableId } }, async () => await table.table.emitInputChanges({ + serializedTable: serializedTables[downstreamTableId], + inputTables: inputTables(downstreamTableId), + changes: normalizedChanges, + })); + validateTableChanges(result.outputChanges, `Table ${downstreamTableId} output`); + const outputChangeCounts = tableChangesDebugInfo(result.outputChanges); + tableOperations.push({ + tableId: downstreamTableId, + phase: "downstream", + durationMs: performance.now() - downstreamStartedAt, + inputChangeCountsByInputTable: inputChangeCountsByInputTable(normalizedChanges), + outputChangeCounts, + }); + affectedTableIds.add(downstreamTableId); + mergeTableChangesDebugInfo(totalOutputChangeCounts, outputChangeCounts); + serializedTables[downstreamTableId] = result.newSerializedTable; + for (const outputTable of table.outputTables) addPending(outputTable.tableId, outputTable.inputTableKey, result.outputChanges); + } + for (const outputTable of table.outputTables) { queue.push(outputTable.tableId); } - } } - const first = await traceSpan({ description: "bulldozer-js.bulldozer.mutateSourceTable", attributes: { "bulldozer.table_id": tableId } }, async () => await mutate({ serializedTable: serializedTables[tableId], inputTables: inputTables(tableId) })); - validateTableChanges(first.outputChanges, `Table ${tableId} output`); - serializedTables[tableId] = first.newSerializedTable; - for (const outputTable of tablesState.tables[tableId].outputTables) addPending(outputTable.tableId, outputTable.inputTableKey, first.outputChanges); - - for (const queue = tablesState.tables[tableId].outputTables.map(outputTable => outputTable.tableId); queue.length;) { - const downstreamTableId = queue.shift()!; - const left = (remainingInputs.get(downstreamTableId) ?? 1) - 1; - remainingInputs.set(downstreamTableId, left); - if (left > 0) continue; - const table = tablesState.tables[downstreamTableId]; - const changes = pending.get(downstreamTableId); - if (changes) { - const result = await traceSpan({ description: "bulldozer-js.bulldozer.emitInputChanges", attributes: { "bulldozer.table_id": downstreamTableId } }, async () => await table.table.emitInputChanges({ - serializedTable: serializedTables[downstreamTableId], - inputTables: inputTables(downstreamTableId), - changes: Object.fromEntries(Object.keys(table.inputTableIds).map(inputTableKey => [inputTableKey, changes[inputTableKey] ?? emptyChanges()])), - })); - validateTableChanges(result.outputChanges, `Table ${downstreamTableId} output`); - serializedTables[downstreamTableId] = result.newSerializedTable; - for (const outputTable of table.outputTables) addPending(outputTable.tableId, outputTable.inputTableKey, result.outputChanges); - } - for (const outputTable of table.outputTables) { - queue.push(outputTable.tableId); - } - } - - return new BulldozerDatabaseSnapshot({ ...this.serialized, serializedTables, uniqueSnapshotIdentifier: crypto.randomUUID() }, this.tablesState); + const debugInfo: BulldozerSnapshotMutationDebugInfo = { + operation: options.operation, + sourceTableId: options.tableId, + rowsSetOrDeleted: totalOutputChangeCounts.rowChanges, + durationMs: performance.now() - startedAt, + tableOperations, + affectedTableIds: [...affectedTableIds], + affectedTables: affectedTablesDebugInfo(tableOperations), + totalOutputChangeCounts, + }; + return { + newSnapshot: new BulldozerDatabaseSnapshot({ ...this.serialized, serializedTables, uniqueSnapshotIdentifier: crypto.randomUUID() }, this.tablesState), + debugInfo, + }; }); } } @@ -584,8 +842,8 @@ export type BulldozerDatabase = { debugPiledriverSnapshot?(): Promise, debugLowLevelSnapshot?(): Promise, getSnapshot(): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, - withSnapshot(updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, - withSnapshotReplicated(updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, + withSnapshot(updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, + withSnapshotReplicated(updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, applyRemainingMigrations(): Promise<{ seq: DatabaseSeq }>, }; @@ -625,13 +883,14 @@ export function declareBulldozerDatabase(piledriverDatabase: PiledriverDatabase, }; }); const withSnapshot = async ( - updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise, + updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise, options: { replicated: boolean }, ) => { return await traceSpan({ description: "bulldozer-js.bulldozer.withSnapshot", attributes: { "bulldozer.replicated": options.replicated } }, async () => { const result = await withWriteLock(async () => { const { snapshot } = await getSnapshot(); - const newSnapshot = await updateSnapshot(snapshot); + const updateResult = await updateSnapshot(snapshot); + const newSnapshot = updateResult instanceof BulldozerDatabaseSnapshot ? updateResult : updateResult.newSnapshot; const { seq } = await setRoot({ snapshot: newSnapshot.toPiledriverObject() }); await piledriverDatabase.waitUntilAvailable(seq); return { snapshot: newSnapshot, seq }; diff --git a/apps/bulldozer-js/src/databases/bulldozer/performance.test.ts b/apps/bulldozer-js/src/databases/bulldozer/performance.test.ts index db12f5ecc..03373df8a 100644 --- a/apps/bulldozer-js/src/databases/bulldozer/performance.test.ts +++ b/apps/bulldozer-js/src/databases/bulldozer/performance.test.ts @@ -94,8 +94,8 @@ async function initializedSnapshot(migrations: Migration) { await db.applyRemainingMigrations(); return (await db.getSnapshot()).snapshot; } -const set = (snapshot: Snapshot, tableId: string, rowIdentifier: string, newRowData: PiledriverObject | undefined) => - snapshot.setOrDeleteRow({ tableId, rowIdentifier, newRowData }); +const set = async (snapshot: Snapshot, tableId: string, rowIdentifier: string, newRowData: PiledriverObject | undefined) => + (await snapshot.setOrDeleteRow({ tableId, rowIdentifier, newRowData })).newSnapshot; const seedRows = async ( snapshot: Snapshot, tableId: string, diff --git a/apps/bulldozer-js/src/payments/schema/index.test.ts b/apps/bulldozer-js/src/payments/schema/index.test.ts index 788fea0a5..e8d0d4efd 100644 --- a/apps/bulldozer-js/src/payments/schema/index.test.ts +++ b/apps/bulldozer-js/src/payments/schema/index.test.ts @@ -209,7 +209,7 @@ describe("payments schema", () => { // With calendar-anchored repeats, the first monthly boundary off the epoch anchor is // 1970-02-01 (31 days), not the 30-day MONTH_MS approximation. const firstRepeatMillis = Date.UTC(1970, 1, 1); - snapshot = await snapshot.tick(new Date(firstRepeatMillis)); + snapshot = (await snapshot.tick(new Date(firstRepeatMillis))).newSnapshot; expect(await balanceAt(snapshot, group, "credits", 0)).toBe(10); expect(await balanceAt(snapshot, group, "credits", firstRepeatMillis)).toBe(10); @@ -232,8 +232,8 @@ describe("payments schema", () => { endedAtMillis: subEndMillis, }) as unknown as PiledriverObject); - snapshot = await snapshot.tick(new Date(firstRepeatMillis)); - snapshot = await snapshot.tick(new Date(subEndMillis)); + snapshot = (await snapshot.tick(new Date(firstRepeatMillis))).newSnapshot; + snapshot = (await snapshot.tick(new Date(subEndMillis))).newSnapshot; const group = customerGroup("u-repeat"); const txns = ((await rowDatas(snapshot, schema.transactions, group)) as unknown as TransactionRow[]) @@ -270,7 +270,7 @@ describe("payments schema", () => { ...overrides, }) as unknown as PiledriverObject; snapshot = await set(snapshot, schema.subscriptions, "sub-rewrite", subRow()); - snapshot = await snapshot.tick(new Date(firstRepeatMillis)); + snapshot = (await snapshot.tick(new Date(firstRepeatMillis))).newSnapshot; const group = customerGroup("u-rewrite"); const txnIds = async () => ((await rowDatas(snapshot, schema.transactions, group)) as unknown as TransactionRow[]).map(txn => txn.txnId).sort(stringCompare); @@ -287,7 +287,7 @@ describe("payments schema", () => { // Quantity upgrade: history keeps the originally granted quantities; only future repeats scale. snapshot = await set(snapshot, schema.subscriptions, "sub-rewrite", subRow({ quantity: 2, currentPeriodStartMillis: firstRepeatMillis })); expect(await balanceAt(snapshot, group, "credits", firstRepeatMillis)).toBe(20); - snapshot = await snapshot.tick(new Date(secondRepeatMillis)); + snapshot = (await snapshot.tick(new Date(secondRepeatMillis))).newSnapshot; expect(await balanceAt(snapshot, group, "credits", secondRepeatMillis)).toBe(40); const secondGrant = ((await rowDatas(snapshot, schema.transactions, group)) as unknown as TransactionRow[]).find(txn => txn.txnId === `igr:sub-rewrite:${secondRepeatMillis}`); expect(secondGrant?.entries).toMatchObject([{ type: "item-quantity-change", itemId: "credits", quantity: 20 }]); @@ -314,7 +314,7 @@ describe("payments schema", () => { ...overrides, }); snapshot = await set(snapshot, schema.oneTimePurchases, "otp-rewrite", otpRow()); - snapshot = await snapshot.tick(new Date(firstRepeatMillis)); + snapshot = (await snapshot.tick(new Date(firstRepeatMillis))).newSnapshot; const group = customerGroup("u-otp-rewrite"); expect(await balanceAt(snapshot, group, "credits", firstRepeatMillis)).toBe(10); @@ -323,7 +323,7 @@ describe("payments schema", () => { // survive the write, and future repeats stop at the revocation. snapshot = await set(snapshot, schema.oneTimePurchases, "otp-rewrite", otpRow({ revokedAtMillis: Date.UTC(1970, 1, 10) })); expect(await balanceAt(snapshot, group, "credits", firstRepeatMillis)).toBe(10); - snapshot = await snapshot.tick(new Date(Date.UTC(1970, 2, 1))); + snapshot = (await snapshot.tick(new Date(Date.UTC(1970, 2, 1)))).newSnapshot; const txnIds = ((await rowDatas(snapshot, schema.transactions, group)) as unknown as TransactionRow[]).map(txn => txn.txnId).sort(stringCompare); expect(txnIds).toEqual([`igr:otp-rewrite:${firstRepeatMillis}`, "otp:otp-rewrite"]); }); @@ -341,12 +341,12 @@ describe("payments schema", () => { ...overrides, }) as unknown as PiledriverObject; snapshot = await set(snapshot, schema.subscriptions, "sub-seal", subRow()); - snapshot = await snapshot.tick(new Date(firstRepeatMillis)); + snapshot = (await snapshot.tick(new Date(firstRepeatMillis))).newSnapshot; // Cancellation arrives as a rewrite of the live row (that's how the Stripe sync works); the // end event must fire off the *existing* fold state, expiring the actually-emitted grants. snapshot = await set(snapshot, schema.subscriptions, "sub-seal", subRow({ status: "canceled", endedAtMillis: subEndMillis, canceledAtMillis: subEndMillis })); - snapshot = await snapshot.tick(new Date(subEndMillis)); + snapshot = (await snapshot.tick(new Date(subEndMillis))).newSnapshot; const group = customerGroup("u-seal"); const txns = ((await rowDatas(snapshot, schema.transactions, group)) as unknown as TransactionRow[]) @@ -362,7 +362,7 @@ describe("payments schema", () => { // Once ended, further webhook rewrites must not re-arm the fold: no duplicate subscription-end, // no resumed repeats past the end. snapshot = await set(snapshot, schema.subscriptions, "sub-seal", subRow({ status: "canceled", endedAtMillis: subEndMillis, canceledAtMillis: subEndMillis, currentPeriodStartMillis: firstRepeatMillis })); - snapshot = await snapshot.tick(new Date(Date.UTC(1970, 3, 1))); + snapshot = (await snapshot.tick(new Date(Date.UTC(1970, 3, 1)))).newSnapshot; const txnIdsAfter = ((await rowDatas(snapshot, schema.transactions, group)) as unknown as TransactionRow[]).map(txn => txn.txnId).sort(stringCompare); expect(txnIdsAfter).toEqual([`igr:sub-seal:${firstRepeatMillis}`, "sub-end:sub-seal", "sub-start:sub-seal"]); }); @@ -399,7 +399,7 @@ describe("payments schema", () => { // The fold is sealed: later rewrites and ticks add nothing (no duplicate end). snapshot = await set(snapshot, schema.subscriptions, "sub-switch", subRow({ status: "canceled", endedAtMillis: subEndMillis, canceledAtMillis: subEndMillis, currentPeriodStartMillis: subEndMillis })); - snapshot = await snapshot.tick(new Date(Date.UTC(1970, 2, 1))); + snapshot = (await snapshot.tick(new Date(Date.UTC(1970, 2, 1)))).newSnapshot; const txnIdsAfter = ((await rowDatas(snapshot, schema.transactions, group)) as unknown as TransactionRow[]).map(txn => txn.txnId).sort(stringCompare); expect(txnIdsAfter).toEqual(["sub-end:sub-switch", "sub-start:sub-switch"]); }); @@ -661,7 +661,7 @@ describe("transactions-by-tenancy date index", () => { product: product({ credits: { quantity: 10, repeat: [1, "month"], expires: "when-repeated" } }), currentPeriodEndMillis: 2 * MONTH_MS, }) as unknown as PiledriverObject); - snapshot = await snapshot.tick(new Date(firstRepeatMillis)); + snapshot = (await snapshot.tick(new Date(firstRepeatMillis))).newSnapshot; snapshot = await setRefund(snapshot, { txnId: "refund:sub-start:sub-grant:uuid1", customerId: "u-grant", createdAtMillis: 5_000 }); const txnIds = ((await rowDatas(snapshot, schema.transactions, customerGroup("u-grant"))) as unknown as TransactionRow[]) diff --git a/apps/bulldozer-js/src/payments/schema/item-quantities.test.ts b/apps/bulldozer-js/src/payments/schema/item-quantities.test.ts index 7da1ae2fe..d6d761794 100644 --- a/apps/bulldozer-js/src/payments/schema/item-quantities.test.ts +++ b/apps/bulldozer-js/src/payments/schema/item-quantities.test.ts @@ -630,7 +630,7 @@ describe("item quantities: full-pipeline integration", () => { currentPeriodEndMillis: 11 * DAY_MS + MONTH_MS, createdAtMillis: 11 * DAY_MS, }) as unknown as PiledriverObject); - snapshot = await snapshot.tick(new Date(11 * DAY_MS)); + snapshot = (await snapshot.tick(new Date(11 * DAY_MS))).newSnapshot; expect(await balanceAt(snapshot, customerGroup("u-upgrade"), "emails", 11 * DAY_MS)).toBe(500); }); @@ -663,7 +663,7 @@ describe("item quantities: full-pipeline integration", () => { // The first monthly reset off the epoch anchor is 1970-02-01 (calendar-anchored, not 30 days), // still well before the manual grant's 3-month absolute expiry, so the reset ranks soonest. const firstResetMillis = Date.UTC(1970, 1, 1); - snapshot = await snapshot.tick(new Date(firstResetMillis)); + snapshot = (await snapshot.tick(new Date(firstResetMillis))).newSnapshot; expect(await balanceAt(snapshot, g, "emails", firstResetMillis)).toBe(200); }); }); diff --git a/apps/bulldozer-js/src/payments/schema/schema-test-helpers.ts b/apps/bulldozer-js/src/payments/schema/schema-test-helpers.ts index 3292d1d5c..db6b55e67 100644 --- a/apps/bulldozer-js/src/payments/schema/schema-test-helpers.ts +++ b/apps/bulldozer-js/src/payments/schema/schema-test-helpers.ts @@ -42,7 +42,7 @@ export const rowsBySortKey = async (snapshot: Snapshot, tableId: string, groupKe await collect(snapshot.listRowsInGroup({ tableId, groupKey, range: {} })); export const set = async (snapshot: Snapshot, tableId: string, rowIdentifier: string, newRowData: PiledriverObject | undefined) => - await snapshot.setOrDeleteRow({ tableId, rowIdentifier, newRowData }); + (await snapshot.setOrDeleteRow({ tableId, rowIdentifier, newRowData })).newSnapshot; export const customerGroup = (customerId: string, customerType: CustomerType = "user"): PiledriverObject => ({ tenancyId: "t1", customerType, customerId });