diff --git a/apps/sim/background/table-update.ts b/apps/sim/background/table-update.ts index 412c241c535..e26fd7e8329 100644 --- a/apps/sim/background/table-update.ts +++ b/apps/sim/background/table-update.ts @@ -1,8 +1,9 @@ -import { task } from '@trigger.dev/sdk' +import { AbortTaskRunError, task } from '@trigger.dev/sdk' import { markTableUpdateFailed, runTableUpdate, type TableUpdatePayload, + UpdatePatchRejectedError, } from '@/lib/table/update-runner' /** @@ -19,7 +20,8 @@ export interface TableUpdateTaskPayload extends Omit { - await runTableUpdate({ ...payload, cutoff: new Date(payload.cutoff) }) + try { + await runTableUpdate({ ...payload, cutoff: new Date(payload.cutoff) }) + } catch (error) { + if (error instanceof UpdatePatchRejectedError) throw new AbortTaskRunError(error.message) + throw error + } }, onFailure: async ({ payload, error }) => { await markTableUpdateFailed(payload.tableId, payload.jobId, error) diff --git a/apps/sim/lib/table/application/copilot-bulk-rows.ts b/apps/sim/lib/table/application/copilot-bulk-rows.ts index 0c29ed56b06..4720ea898d0 100644 --- a/apps/sim/lib/table/application/copilot-bulk-rows.ts +++ b/apps/sim/lib/table/application/copilot-bulk-rows.ts @@ -33,6 +33,7 @@ import { import { assertRowDelete, assertRowUpdate, patchColumnIds } from '@/lib/table/mutation-locks' import { createExactEmptyTableRowSecretProvenance } from '@/lib/table/rows/secret-provenance' import { markTableUpdateFailed, runTableUpdate } from '@/lib/table/update-runner' +import { uniqueColumnsInPatch } from '@/lib/table/validation' const logger = createLogger('CopilotBulkRowsApplication') @@ -199,9 +200,7 @@ export const copilotUpdateRowsByFilter = defineAuthorizedTableUseCase({ validateLimit(input.limit) const idData = rowDataNameToId(input.data, buildIdByName(context.table.schema)) const filter = tablePredicateNamesToFilter(input.filter, context.table) - const patchTouchesUnique = context.table.schema.columns.some( - (column) => column.unique === true && (column.id ?? column.name) in idData - ) + const patchTouchesUnique = uniqueColumnsInPatch(context.table.schema, idData).length > 0 const inlineEligible = input.limit !== undefined && input.limit <= TABLE_LIMITS.MAX_BULK_OPERATION_SIZE diff --git a/apps/sim/lib/table/bulk-update-concurrency.test.ts b/apps/sim/lib/table/bulk-update-concurrency.test.ts index fddffbdbf40..ecfa636c891 100644 --- a/apps/sim/lib/table/bulk-update-concurrency.test.ts +++ b/apps/sim/lib/table/bulk-update-concurrency.test.ts @@ -1,4 +1,5 @@ import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing' +import { tableRowsLiveSchemaMock } from '@sim/testing/mocks/table-rows-live-schema.mock' import { tableRowsSecretProvenanceMock, tableRowsSecretProvenanceMockFns, @@ -22,21 +23,25 @@ vi.mock('@/lib/table/rows/ordering', () => ({ vi.mock('@/lib/table/rows/secret-provenance', () => tableRowsSecretProvenanceMock) +vi.mock('@/lib/table/rows/live-schema', () => tableRowsLiveSchemaMock) + vi.mock('@/lib/table/sql', () => ({ buildFilterClause: vi.fn(() => sql`true`), buildPredicateClause: vi.fn(() => sql`true`), buildSortClause: vi.fn(() => sql`true`), escapeLikePattern: vi.fn((value: string) => value), - fieldPredicate: vi.fn(() => sql`true`), + uniqueValuePredicate: vi.fn(() => sql`true`), })) vi.mock('@/lib/table/trigger', () => tableTriggerMock) -vi.mock('@/lib/table/validation', () => ({ +vi.mock('@/lib/table/validation', async (importOriginal) => ({ + cellOf: (await importOriginal()).cellOf, validateRowSize: hoisted.validateRowSize, coerceRowToSchema: hoisted.coerceRowToSchema, coerceRowValues: vi.fn(), getUniqueColumns: vi.fn(() => []), + uniqueColumnsInPatch: vi.fn(() => []), checkUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })), checkBatchUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })), })) diff --git a/apps/sim/lib/table/import-data.ts b/apps/sim/lib/table/import-data.ts index 48436a6144c..ee70fb34465 100644 --- a/apps/sim/lib/table/import-data.ts +++ b/apps/sim/lib/table/import-data.ts @@ -15,6 +15,7 @@ import { CSV_MAX_BATCH_SIZE } from '@/lib/table/import' import { assertRowDelete, assertRowInsert, assertSchemaMutable } from '@/lib/table/mutation-locks' import { nKeysBetween } from '@/lib/table/order-key' import type { DbTransaction } from '@/lib/table/planner' +import { lockLiveTableSchema, refitRowToSchema, withLiveSchema } from '@/lib/table/rows/live-schema' import { acquireRowOrderLock, guardBatch, @@ -64,9 +65,10 @@ export interface BulkImportBatch { * `runWorkflowColumn`** (a 1M-row import must not dispatch a workflow run per row). * Append and replace imports run this against the live table, so other writers can * race it: the batch holds the table's unique columns exclusively while it checks and - * inserts. `row_count` is maintained set-based by the statement-level - * trigger. There is no surrounding transaction and no rollback: each batch commits on - * its own, so committed batches persist even if a later batch fails. + * inserts, and checks the rows against the schema it reads under the schema lock, which + * may have changed since the job resolved the table. `row_count` is maintained set-based + * by the statement-level trigger. There is no surrounding transaction and no rollback: + * each batch commits on its own, so committed batches persist even if a later batch fails. * * Throws on row-size/schema/unique violations or if the statement-level trigger rejects * the batch for crossing `max_rows`; the caller marks the import failed. @@ -82,6 +84,7 @@ export async function bulkInsertImportBatch( // the caller's snapshot too would reject a since-cleared lock. if (!revalidate) assertRowInsert(table) + const rawRows = data.rows.map((row) => ({ ...row })) for (let i = 0; i < data.rows.length; i++) { const sizeValidation = validateRowSize(data.rows[i]) if (!sizeValidation.valid) { @@ -119,14 +122,23 @@ export async function bulkInsertImportBatch( })) const inserted = await db.transaction(async (trx) => { - await guardBatch(trx, data.tableId, revalidate) - if (getUniqueColumns(table.schema).length > 0) { + const fresh = await guardBatch(trx, data.tableId, revalidate) + const live = fresh ? withLiveSchema(table, fresh.schema) : await lockLiveTableSchema(trx, table) + if (live !== table) { + for (let i = 0; i < data.rows.length; i++) { + const refit = refitRowToSchema(data.rows[i], rawRows[i], table.schema, live.schema, 'null') + if (!refit.valid) { + throw new OrchestrationError('validation', `Row ${i + 1}: ${refit.errors.join(', ')}`) + } + } + } + if (getUniqueColumns(live.schema).length > 0) { // The whole-table unique lock, not per-value: a batch is far more values than the value-lock cap. - await lockUniqueColumns(trx, table) + await lockUniqueColumns(trx, live) const uniqueResult = await checkBatchUniqueConstraintsDb( data.tableId, data.rows, - table.schema, + live.schema, trx ) if (!uniqueResult.valid) { @@ -314,7 +326,7 @@ export async function importAppendRows( generateId().slice(0, 8), { uniqueColumnsLocked: true } ) - inserted.push(...batchInserted) + inserted.push(...batchInserted.rows) } return { inserted, table: working } }) diff --git a/apps/sim/lib/table/rows/bulk-update-patch-validation.test.ts b/apps/sim/lib/table/rows/bulk-update-patch-validation.test.ts index aa6bf13b21b..61f18fa16a3 100644 --- a/apps/sim/lib/table/rows/bulk-update-patch-validation.test.ts +++ b/apps/sim/lib/table/rows/bulk-update-patch-validation.test.ts @@ -4,6 +4,7 @@ * filter matches rows or none. */ import { resetDbChainMock } from '@sim/testing' +import { tableRowsLiveSchemaMock } from '@sim/testing/mocks/table-rows-live-schema.mock' import { tableRowsSecretProvenanceMock, tableRowsSecretProvenanceMockFns, @@ -24,12 +25,14 @@ vi.mock('@/lib/table/rows/ordering', () => ({ vi.mock('@/lib/table/rows/secret-provenance', () => tableRowsSecretProvenanceMock) +vi.mock('@/lib/table/rows/live-schema', () => tableRowsLiveSchemaMock) + vi.mock('@/lib/table/sql', () => ({ buildFilterClause: vi.fn(() => sql`true`), buildPredicateClause: vi.fn(() => sql`true`), buildSortClause: vi.fn(() => sql`true`), escapeLikePattern: vi.fn((value: string) => value), - fieldPredicate: vi.fn(() => sql`true`), + uniqueValuePredicate: vi.fn(() => sql`true`), })) vi.mock('@/lib/table/trigger', () => tableTriggerMock) diff --git a/apps/sim/lib/table/rows/live-schema.ts b/apps/sim/lib/table/rows/live-schema.ts new file mode 100644 index 00000000000..d27f7fcc183 --- /dev/null +++ b/apps/sim/lib/table/rows/live-schema.ts @@ -0,0 +1,104 @@ +/** + * Keeps row writes on the table's live schema rather than the definition their caller resolved + * before the write transaction opened. + * + * Schema changes hold the table's schema lock exclusively (`withLockedTable`) while they check the + * stored rows against the new schema: a column made unique is scanned for duplicates, one made + * required for empty cells. A write validated against an older definition skips the check the new + * schema adds, and if it commits after that scan, the column ends up holding the duplicate or the + * empty cell. Each write transaction therefore takes the schema lock shared and reads the schema + * under it, then validates against that: a schema change waits for writes already in flight, and a + * write that waited sees the change. + */ + +import { compareStrings } from '@sim/utils/string' +import { sql } from 'drizzle-orm' +import { canonicalJson } from '@/lib/api/cursor-binding' +import { OrchestrationError } from '@/lib/core/orchestration/types' +import { getColumnId } from '@/lib/table/column-keys' +import type { DbTransaction } from '@/lib/table/planner' +import { type TableTxTimeouts, tableTxTimeoutSettings } from '@/lib/table/tx' +import type { RowData, TableDefinition, TableSchema, ValidationResult } from '@/lib/table/types' +import { + coerceRowToSchema, + type PatchedKeys, + type UncoercibleValuePolicy, + validateRowSize, +} from '@/lib/table/validation' + +/** Compares schemas by content, whatever order their columns are listed in. */ +function schemaFingerprint(schema: TableSchema): string { + const columns = [...schema.columns].sort((a, b) => compareStrings(getColumnId(a), getColumnId(b))) + return canonicalJson({ ...schema, columns }) +} + +/** `table` itself when `schema` matches its schema, else a copy carrying `schema`. */ +export function withLiveSchema(table: TableDefinition, schema: TableSchema): TableDefinition { + return schemaFingerprint(schema) === schemaFingerprint(table.schema) + ? table + : { ...table, schema } +} + +/** + * Takes the table's schema lock shared and reads its live schema, and returns the definition to + * validate and write against: `table` itself when its schema is still current, else a copy carrying + * the live schema. Call it first in the transaction, before its other locks, passing the + * transaction's `timeouts` (see `setTableTxTimeouts`) in place of a separate timeouts statement. + * + * One statement: `user_table_schema_for_write` (script migration 0026) takes the lock and then + * reads the schema. It is VOLATILE, so under READ COMMITTED its read takes a fresh snapshot and sees + * a schema change that committed while the lock waited. The timeouts are applied in a subquery the call + * reads from, first. The lock waits as long as the transaction's `statement_timeout` allows, not + * its shorter `lock_timeout`, which the function leaves as it found it for the locks that follow. + */ +export async function lockLiveTableSchema( + trx: DbTransaction, + table: TableDefinition, + timeouts?: TableTxTimeouts +): Promise { + const read = sql`SELECT user_table_schema_for_write(${table.id}) AS schema` + const [live] = await trx.execute<{ schema: TableSchema | null }>( + timeouts ? sql`${read} FROM (SELECT ${tableTxTimeoutSettings(timeouts)}) AS settings` : read + ) + if (!live?.schema) throw new OrchestrationError('not_found', 'Table not found') + return withLiveSchema(table, live.schema) +} + +/** + * Removes, in place, the cells of columns `snapshot` defines and `live` no longer does. A column + * delete reclaims its cells in the background, so a cell written after that pass would stay behind. + */ +export function dropDeletedColumns( + rows: readonly RowData[], + snapshot: TableSchema, + live: TableSchema +): void { + const liveIds = new Set(live.columns.map(getColumnId)) + const deleted = snapshot.columns.map(getColumnId).filter((id) => !liveIds.has(id)) + if (deleted.length === 0) return + for (const row of rows) { + for (const id of deleted) delete row[id] + } +} + +/** + * Rebuilds `row`, in place, for `live` once the schema has moved since `snapshot`: from `raw`, the + * row as the caller wrote it before coercing it against `snapshot`, so a value that schema would + * have reshaped (`"007"` read as a number) reaches the live column as it was sent. Then drops the + * cells of deleted columns, coerces and validates against `live` exactly as the write first did, + * and re-checks the row's size, which a coercion to a wider type can grow. + */ +export function refitRowToSchema( + row: RowData, + raw: RowData, + snapshot: TableSchema, + live: TableSchema, + policy?: UncoercibleValuePolicy, + patchedKeys?: PatchedKeys +): ValidationResult { + for (const key of Object.keys(row)) delete row[key] + Object.assign(row, raw) + dropDeletedColumns([row], snapshot, live) + const result = coerceRowToSchema(row, live, policy, patchedKeys) + return result.valid ? validateRowSize(row) : result +} diff --git a/apps/sim/lib/table/rows/ordering.ts b/apps/sim/lib/table/rows/ordering.ts index 53077b92602..0919efa7582 100644 --- a/apps/sim/lib/table/rows/ordering.ts +++ b/apps/sim/lib/table/rows/ordering.ts @@ -333,10 +333,12 @@ export async function insertOrderedRow(params: { /** Proof the caller asserted the insert lock (see `mutation-locks.ts`). */ proof: MutationProof<'insert'> /** - * Runs first in the transaction, before the row-order lock: the caller's unique-value locks and - * unique check (see `unique-locks.ts`), so the check sees any concurrent insert of the same value. + * Opens the transaction in place of the default timeouts, before the row-order lock: the + * caller's schema guard (see `live-schema.ts`), which applies the timeouts, then its unique-value + * locks and unique check (see `unique-locks.ts`), so the check sees any concurrent insert of the + * same value. */ - assertUnique?: (trx: DbTransaction) => Promise + validate?: (trx: DbTransaction) => Promise }): Promise<{ id: string data: RowData @@ -358,8 +360,8 @@ export async function insertOrderedRow(params: { secretProvenance, } = params const [row] = await db.transaction(async (trx) => { - await setTableTxTimeouts(trx) - await params.assertUnique?.(trx) + if (params.validate) await params.validate(trx) + else await setTableTxTimeouts(trx) await acquireRowOrderLock(trx, tableId) // Resolve the authoritative order key from neighbor ids when given, else from the requested @@ -670,17 +672,29 @@ export async function deletePageByIds( return deleted } +/** The patch one update batch writes, or `null` when it writes nothing. */ +export interface PagePatch { + patchJson: string + secretProvenance: TableRowSecretProvenanceWrite +} + /** * Applies a JSONB-merge patch (`data || patchJson`) to a page of row ids, committed in * UPDATE_BATCH_SIZE chunks (each its own transaction, 60s timeout) so a large background update - * makes incremental, resumable progress. Returns the number of rows updated. + * makes incremental, resumable progress. Each batch takes its patch from `prepare`, called inside + * the batch's transaction with the definition `revalidate` read there and the batch's row ids, so a + * caller can derive it, and check the rows it merges into, against the live schema. Returns the + * number of rows updated. */ export async function updatePageByIds( tableId: string, workspaceId: string, rowIds: string[], - patchJson: string, - secretProvenance: TableRowSecretProvenanceWrite, + prepare: ( + trx: DbTransaction, + table: TableDefinition | undefined, + batch: string[] + ) => Promise, /** Proof the caller asserted the update lock (see `mutation-locks.ts`). */ _proof: MutationProof<'update'>, /** Re-asserts the lock inside each batch transaction. See {@link guardBatch}. */ @@ -692,15 +706,19 @@ export async function updatePageByIds( const batch = rowIds.slice(i, i + TABLE_LIMITS.UPDATE_BATCH_SIZE) const rows = await db.transaction(async (trx) => { await setTableTxTimeouts(trx, { statementMs: 60_000 }) - await guardBatch(trx, tableId, revalidate) + const patch = await prepare(trx, await guardBatch(trx, tableId, revalidate), batch) + if (!patch) return [] return mutateTableRowsWithSecretProvenance(trx, { - rows: batch.map((rowId) => ({ rowId, provenance: secretProvenance })), + rows: batch.map((rowId) => ({ rowId, provenance: patch.secretProvenance })), rowState: 'existing', mode: 'merge', mutate: async () => { const rows = await trx .update(userTableRows) - .set({ data: sql`${userTableRows.data} || ${patchJson}::jsonb`, updatedAt: now }) + .set({ + data: sql`${userTableRows.data} || ${patch.patchJson}::jsonb`, + updatedAt: now, + }) .where( and( eq(userTableRows.tableId, tableId), diff --git a/apps/sim/lib/table/rows/row-writes.integration.ts b/apps/sim/lib/table/rows/row-writes.integration.ts index 74a9a0584b7..e5c3a5de422 100644 --- a/apps/sim/lib/table/rows/row-writes.integration.ts +++ b/apps/sim/lib/table/rows/row-writes.integration.ts @@ -9,7 +9,7 @@ import { userTableDefinitions, userTableRows } from '@sim/db/schema' import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure' import { storageServiceMock, storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock' import { tableBillingMock, tableBillingMockFns } from '@sim/testing/mocks/table-billing.mock' -import { tableTriggerMock } from '@sim/testing/mocks/table-trigger.mock' +import { tableTriggerMock, tableTriggerMockFns } from '@sim/testing/mocks/table-trigger.mock' import { tableWorkflowColumnsMock } from '@sim/testing/mocks/table-workflow-columns.mock' import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' @@ -22,14 +22,31 @@ vi.mock('@/lib/table/trigger', () => tableTriggerMock) vi.mock('@/lib/table/workflow-columns', () => tableWorkflowColumnsMock) vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock) +import { + deleteColumn, + updateColumnConstraints, + updateColumnOptions, +} from '@/lib/table/columns/service' +import { getMaxRowSizeBytes, TABLE_LIMITS } from '@/lib/table/constants' import { bulkInsertImportBatch, importReplaceRows } from '@/lib/table/import-data' +import { markTableJobRunningInWorkspace } from '@/lib/table/jobs/service' import type { DbTransaction } from '@/lib/table/planner' +import { lockLiveTableSchema } from '@/lib/table/rows/live-schema' import { acquireRowOrderLock } from '@/lib/table/rows/ordering' -import { batchInsertRows, batchUpdateRows, insertRow, updateRow } from '@/lib/table/rows/service' +import { + batchInsertRows, + batchUpdateRows, + insertRow, + replaceTableRows, + updateRow, + updateRowsByFilter, + upsertRow, +} from '@/lib/table/rows/service' import { lockUniqueColumns, lockUniqueValues } from '@/lib/table/rows/unique-locks' import { getTableById } from '@/lib/table/service' import { getOrCreateTableSnapshot } from '@/lib/table/snapshot-cache' -import type { ColumnDefinition, RowData, TableDefinition } from '@/lib/table/types' +import type { ColumnDefinition, JsonValue, RowData, TableDefinition } from '@/lib/table/types' +import { runTableUpdate, UpdatePatchRejectedError } from '@/lib/table/update-runner' const url = readTestDatabaseUrl() if (process.env.DATABASE_URL !== url) { @@ -80,6 +97,27 @@ async function rowsVersion(tableId: string): Promise { const textColumns = (...ids: string[]): ColumnDefinition[] => ids.map((id) => ({ id, name: id, type: 'string' })) +/** Sessions waiting on the table's schema lock, matched by the key's hash as `pg_locks` shows it. */ +async function schemaLockWaiters(tableId: string): Promise { + const [{ waiting }] = await control<{ waiting: number }[]>` + WITH lock AS (SELECT hashtextextended(${`user_table_schema:${tableId}`}, 0) AS key) + SELECT count(*)::int AS waiting FROM pg_locks l CROSS JOIN lock + WHERE l.locktype = 'advisory' AND NOT l.granted + AND l.classid = ((lock.key >> 32) & 4294967295)::oid + AND l.objid = (lock.key & 4294967295)::oid` + return waiting +} + +/** Polls until exactly `expected` sessions wait on the table's schema lock, and asserts it. */ +async function untilSchemaLockWaiters(tableId: string, expected: number) { + let waiting = 0 + for (let attempt = 0; attempt < 400 && waiting !== expected; attempt++) { + await sleep(5) + waiting = await schemaLockWaiters(tableId) + } + expect(waiting).toBe(expected) +} + describe('table row writes against real PostgreSQL', () => { beforeAll(async () => { await control`INSERT INTO "user" (id, name, email, email_verified, created_at, updated_at) @@ -121,26 +159,39 @@ describe('table row writes against real PostgreSQL', () => { /** * Polls until exactly `expected.onOrderLock` sessions wait on the table's row-order lock and - * `expected.onValueLock` wait on a unique-value lock, and returns the last count seen. + * `expected.onValueLock` wait on one of its unique locks, and returns the last count seen. Both + * counts are scoped to the table, so concurrent suites cannot skew them: a unique-lock waiter + * waits on the table's unique lock itself, or on a value lock while holding the table's unique + * lock shared. */ async function waitForLockWaiters(tableId: string, expected: LockWaiters) { const orderLockKey = `user_table_rows_pos:${tableId}` + const uniqueLockKey = `user_table_unique:${tableId}` let waiting: LockWaiters = { onOrderLock: 0, onValueLock: 0 } for (let attempt = 0; attempt < 400; attempt++) { await sleep(5) ;[waiting] = await control` - WITH lock AS (SELECT hashtextextended(${orderLockKey}, 0) AS key) + WITH lock AS ( + SELECT hashtextextended(${orderLockKey}, 0) AS order_key, + hashtextextended(${uniqueLockKey}, 0) AS unique_key + ), advisory AS ( + SELECT l.pid, l.granted, + l.classid = ((lock.order_key >> 32) & 4294967295)::oid + AND l.objid = (lock.order_key & 4294967295)::oid AS is_order, + l.classid = ((lock.unique_key >> 32) & 4294967295)::oid + AND l.objid = (lock.unique_key & 4294967295)::oid AS is_unique + FROM pg_locks l CROSS JOIN lock + WHERE l.locktype = 'advisory' + AND l.database = (SELECT oid FROM pg_database WHERE datname = current_database()) + ) SELECT + count(*) FILTER (WHERE NOT granted AND is_order)::int AS "onOrderLock", count(*) FILTER ( - WHERE l.classid = ((lock.key >> 32) & 4294967295)::oid - AND l.objid = (lock.key & 4294967295)::oid - )::int AS "onOrderLock", - count(*) FILTER (WHERE a.query LIKE '%user_table_unique_value%')::int AS "onValueLock" - FROM pg_locks l - JOIN pg_stat_activity a ON a.pid = l.pid - CROSS JOIN lock - WHERE l.locktype = 'advisory' AND NOT l.granted - AND l.database = (SELECT oid FROM pg_database WHERE datname = current_database())` + WHERE NOT granted AND NOT is_order AND ( + is_unique OR pid IN (SELECT pid FROM advisory WHERE granted AND is_unique) + ) + )::int AS "onValueLock" + FROM advisory` if ( waiting.onOrderLock === expected.onOrderLock && waiting.onValueLock === expected.onValueLock @@ -177,10 +228,7 @@ describe('table row writes against real PostgreSQL', () => { } } - async function storedCount( - tableId: string, - match: Record> - ) { + async function storedCount(tableId: string, match: Record) { const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count FROM user_table_rows WHERE table_id = ${tableId} AND data @> ${control.json(match)}` return count @@ -483,6 +531,1101 @@ describe('table row writes against real PostgreSQL', () => { ).rejects.toThrow(/must be unique/) expect(await storedCount(table.id, { meta: { x: 1, y: 2 } })).toBe(0) }) + + describe('json columns compare whole values', () => { + const jsonColumns: ColumnDefinition[] = [ + { id: 'meta', name: 'meta', type: 'json', unique: true }, + ] + + const insertMeta = (table: TableDefinition, meta: JsonValue) => + insertRow( + { + tableId: table.id, + workspaceId, + data: { meta }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'json-unique' + ) + + async function metaValues(tableId: string): Promise { + const rows = await control<{ meta: JsonValue }[]>`SELECT data->'meta' AS meta + FROM user_table_rows WHERE table_id = ${tableId} ORDER BY created_at, id` + return rows.map((row) => row.meta) + } + + it('accepts an object or array that a stored value merely contains', async () => { + const table = await createTable(jsonColumns) + await seedRows(table.id, [ + { id: `${table.id}-a`, data: { meta: { a: 1, b: 2 } }, orderKey: 'a0' }, + { id: `${table.id}-b`, data: { meta: [1, 2, 3] }, orderKey: 'a1' }, + ]) + + await expect(insertMeta(table, { a: 1 })).resolves.toBeDefined() + await expect(insertMeta(table, [1])).resolves.toBeDefined() + expect(await storedCount(table.id, { meta: { a: 1 } })).toBe(2) + }) + + it('rejects the same object in a different key order', async () => { + const table = await createTable(jsonColumns) + await seedRows(table.id, [ + { id: `${table.id}-a`, data: { meta: { a: 1, b: 2 } }, orderKey: 'a0' }, + ]) + const batch = (rows: RowData[]) => ({ + tableId: table.id, + workspaceId, + rows, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }) + + await expect(insertMeta(table, { b: 2, a: 1 })).rejects.toThrow(/must be unique/) + await expect( + batchInsertRows(batch([{ meta: { b: 2, a: 1 } }]), table, 'json-unique') + ).rejects.toThrow(/must be unique/) + await expect( + replaceTableRows( + batch([{ meta: { b: 2, a: 1 } }, { meta: { a: 1, b: 2 } }]), + table, + 'json-unique' + ) + ).rejects.toThrow(/must be unique/) + expect(await metaValues(table.id)).toEqual([{ a: 1, b: 2 }]) + }) + + it('upserts on a JSON conflict target without overwriting a row that contains it', async () => { + const table = await createTable(jsonColumns) + await seedRows(table.id, [ + { id: `${table.id}-a`, data: { meta: { a: 1, b: 2 } }, orderKey: 'a0' }, + ]) + + const result = await upsertRow( + { + tableId: table.id, + workspaceId, + data: { meta: { a: 1 } }, + conflictTarget: 'meta', + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'json-unique' + ) + + expect(result.operation).toBe('insert') + expect(await metaValues(table.id)).toEqual([{ a: 1, b: 2 }, { a: 1 }]) + }) + + it('accepts a batch update of values contained in other rows', async () => { + const table = await createTable(jsonColumns) + await seedRows(table.id, [ + { id: `${table.id}-a`, data: { meta: { n: 1 } }, orderKey: 'a0' }, + { id: `${table.id}-b`, data: { meta: { n: 2 } }, orderKey: 'a1' }, + { id: `${table.id}-c`, data: { meta: { a: 1, b: 2, c: 3 } }, orderKey: 'a2' }, + ]) + + await batchUpdateRows( + { + tableId: table.id, + workspaceId, + updates: [ + { rowId: `${table.id}-a`, data: { meta: { a: 1 } } }, + { rowId: `${table.id}-b`, data: { meta: { a: 1, b: 2 } } }, + ], + capabilityGovernedUserId: null, + }, + table, + 'json-unique' + ) + + expect(await metaValues(table.id)).toEqual([{ a: 1 }, { a: 1, b: 2 }, { a: 1, b: 2, c: 3 }]) + }) + + it('locks JSON values one by one, equating objects whose keys differ only in order', async () => { + const table = await createTable(jsonColumns) + + const different = await raceUnderHeldOrderLock( + table.id, + [() => insertMeta(table, { a: 1 }), () => insertMeta(table, { a: 1, b: 2 })], + { onOrderLock: 2, onValueLock: 0 } + ) + expect(fulfilled(different)).toBe(2) + + const reordered = await raceUnderHeldOrderLock( + table.id, + [ + () => insertMeta(table, { x: 1, y: [1, 2] }), + () => insertMeta(table, { y: [1, 2], x: 1 }), + ], + { onOrderLock: 1, onValueLock: 1 } + ) + expect(fulfilled(reordered)).toBe(1) + expect(await storedCount(table.id, { meta: { x: 1 } })).toBe(1) + }) + }) + }) + + describe('writers validate against the live schema, not their snapshot', () => { + const columns: ColumnDefinition[] = [ + { id: 'key', name: 'key', type: 'string', unique: true }, + { id: 'email', name: 'email', type: 'string' }, + { id: 'note', name: 'note', type: 'string' }, + ] + + /** A table whose `k1`/`k2` rows hold distinct emails and filled notes, and its snapshot. */ + async function seededTable() { + const table = await createTable(columns) + await seedRows(table.id, [ + { + id: `${table.id}-1`, + data: { key: 'k1', email: 'dup@example.test', note: 'n' }, + orderKey: 'a0', + }, + { + id: `${table.id}-2`, + data: { key: 'k2', email: 'b@example.test', note: 'n' }, + orderKey: 'a1', + }, + ]) + return table + } + + type StaleWrite = (table: TableDefinition, row: RowData) => Promise + + const rows = (table: TableDefinition, data: RowData[]) => ({ + tableId: table.id, + workspaceId, + rows: data, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }) + + /** + * One entry per row writer. Each writes `row`, a new row (`key: 'k3'`) or a patch to row `k2`, + * holding the snapshot it was handed. + */ + const writers: Array<[string, StaleWrite]> = [ + [ + 'insertRow', + (table, row) => + insertRow( + { + tableId: table.id, + workspaceId, + data: { key: 'k3', ...row }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ), + ], + [ + 'batchInsertRows', + (table, row) => + batchInsertRows(rows(table, [{ key: 'k3', ...row }]), table, 'stale-schema'), + ], + [ + 'bulkInsertImportBatch', + (table, row) => + bulkInsertImportBatch( + { tableId: table.id, workspaceId, rows: [{ key: 'k3', ...row }], startPosition: 2 }, + table, + 'stale-schema' + ), + ], + [ + 'replaceTableRows', + (table, row) => + replaceTableRows( + rows(table, [ + { key: 'k1', email: 'dup@example.test', note: 'n' }, + { key: 'k3', ...row }, + ]), + table, + 'stale-schema' + ), + ], + [ + 'upsertRow', + (table, row) => + upsertRow( + { + tableId: table.id, + workspaceId, + data: { key: 'k3', ...row }, + conflictTarget: 'key', + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ), + ], + [ + 'updateRow', + (table, row) => + updateRow( + { + tableId: table.id, + rowId: `${table.id}-2`, + workspaceId, + data: row, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ), + ], + [ + 'batchUpdateRows', + (table, row) => + batchUpdateRows( + { + tableId: table.id, + workspaceId, + updates: [{ rowId: `${table.id}-2`, data: row }], + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ), + ], + [ + 'updateRowsByFilter', + (table, row) => + updateRowsByFilter( + table, + { + filter: { key: 'k2' }, + data: row, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + 'stale-schema' + ), + ], + [ + 'updateRowsByFilter with a limit', + (table, row) => + updateRowsByFilter( + table, + { + filter: { key: 'k2' }, + data: row, + limit: 1, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + 'stale-schema' + ), + ], + ] + + async function countRows( + tableId: string, + where: 'email-dup' | 'email-c' | 'note-empty' | 'note-restored' + ) { + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${tableId} AND ${ + where === 'email-dup' + ? control`data->>'email' = 'dup@example.test'` + : where === 'email-c' + ? control`data->>'email' = 'c@example.test'` + : where === 'note-empty' + ? control`data->>'note' IS NULL` + : control`data->>'note' = 'restored'` + }` + return count + } + + it.each(writers)( + '%s rejects a duplicate in a column made unique since its snapshot', + async (_, write) => { + const table = await seededTable() + await updateColumnConstraints( + { tableId: table.id, columnName: 'email', unique: true }, + 'stale-schema' + ) + + await expect(write(table, { email: 'dup@example.test', note: 'n' })).rejects.toThrow( + /must be unique|would violate uniqueness/ + ) + expect(await countRows(table.id, 'email-dup')).toBe(1) + } + ) + + it.each(writers)( + '%s rejects an empty cell in a column made required since its snapshot', + async (_, write) => { + const table = await seededTable() + await updateColumnConstraints( + { tableId: table.id, columnName: 'note', required: true }, + 'stale-schema' + ) + + await expect(write(table, { email: 'c@example.test', note: null })).rejects.toThrow( + /Missing required field/ + ) + expect(await countRows(table.id, 'note-empty')).toBe(0) + } + ) + + it.each(writers)('%s drops a cell of a column deleted since its snapshot', async (_, write) => { + const table = await seededTable() + await deleteColumn({ tableId: table.id, columnName: 'note' }, 'stale-schema') + + await write(table, { email: 'c@example.test', note: 'restored' }) + expect(await countRows(table.id, 'note-restored')).toBe(0) + }) + + it.each([ + [ + 'updateRow', + (table: TableDefinition, rowId: string) => + updateRow( + { + tableId: table.id, + rowId, + workspaceId, + data: { out: null }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ), + ], + [ + 'batchUpdateRows', + (table: TableDefinition, rowId: string) => + batchUpdateRows( + { + tableId: table.id, + workspaceId, + updates: [{ rowId, data: { out: null } }], + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ), + ], + ])('%s clears the run of a workflow output added since its snapshot', async (_, write) => { + const table = await createTable([ + { id: 'name', name: 'name', type: 'string' }, + { id: 'out', name: 'out', type: 'string', workflowGroupId: 'group-1' }, + ]) + await control`UPDATE user_table_definitions SET schema = jsonb_set(schema, '{workflowGroups}', + ${control.json([ + { + id: 'group-1', + workflowId: 'workflow-1', + outputs: [{ blockId: 'b', path: 'p', columnName: 'out' }], + }, + ])}) WHERE id = ${table.id}` + const stale: TableDefinition = { + ...table, + schema: { columns: textColumns('name', 'out') }, + } + const rowId = `${table.id}-a` + await seedRows(table.id, [{ id: rowId, data: { name: 'a', out: 'done' }, orderKey: 'a0' }]) + await control`INSERT INTO table_row_executions (table_id, row_id, group_id, status, workflow_id) + VALUES (${table.id}, ${rowId}, 'group-1', 'completed', 'workflow-1')` + + await write(stale, rowId) + + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM table_row_executions WHERE row_id = ${rowId}` + expect(count).toBe(0) + }) + + /** A table whose `code` column is text, and a snapshot from when it was a number. */ + async function retypedTable(): Promise<{ table: TableDefinition; stale: TableDefinition }> { + const table = await createTable([ + { id: 'key', name: 'key', type: 'string', unique: true }, + { id: 'code', name: 'code', type: 'string' }, + { id: 'filler', name: 'filler', type: 'string' }, + ]) + await seedRows(table.id, [ + { id: `${table.id}-2`, data: { key: 'k2', filler: '' }, orderKey: 'a0' }, + ]) + const stale: TableDefinition = { + ...table, + schema: { + columns: table.schema.columns.map((column) => + column.id === 'code' ? { ...column, type: 'number' as const } : column + ), + }, + } + return { table, stale } + } + + /** Writes `row` (a new row, or a patch to row `k2`) through each writer holding `stale`. */ + const retypeWriters: Array< + [string, (table: TableDefinition, row: RowData) => Promise] + > = [ + [ + 'insertRow', + (table, row) => + insertRow( + { + tableId: table.id, + workspaceId, + data: { key: 'k3', ...row }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'retype' + ), + ], + [ + 'upsertRow', + (table, row) => + upsertRow( + { + tableId: table.id, + workspaceId, + data: { key: 'k3', ...row }, + conflictTarget: 'key', + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'retype' + ), + ], + [ + 'bulkInsertImportBatch', + (table, row) => + bulkInsertImportBatch( + { tableId: table.id, workspaceId, rows: [{ key: 'k3', ...row }], startPosition: 1 }, + table, + 'retype' + ), + ], + [ + 'updateRow', + (table, row) => + updateRow( + { + tableId: table.id, + rowId: `${table.id}-2`, + workspaceId, + data: row, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'retype' + ), + ], + [ + 'batchUpdateRows', + (table, row) => + batchUpdateRows( + { + tableId: table.id, + workspaceId, + updates: [{ rowId: `${table.id}-2`, data: row }], + capabilityGovernedUserId: null, + }, + table, + 'retype' + ), + ], + ] + + it.each(retypeWriters)( + '%s stores the value it was sent in a column retyped since its snapshot', + async (_, write) => { + const { table, stale } = await retypedTable() + + await write(stale, { code: '007' }) + + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id} AND data->'code' = '"007"'` + expect(count).toBe(1) + } + ) + + it.each(retypeWriters)( + '%s refuses a row a retype since its snapshot grows past the size limit', + async (writer, write) => { + const { table, stale } = await retypedTable() + // Exactly at the limit with `code` a number; the live text column stores it with quotes. + const inserting = writer !== 'updateRow' && writer !== 'batchUpdateRows' + const shape = (filler: string) => + inserting ? { key: 'k3', code: 7, filler } : { key: 'k2', filler, code: 7 } + const filler = 'x'.repeat( + getMaxRowSizeBytes() - Buffer.byteLength(JSON.stringify(shape(''))) + ) + + await expect(write(stale, { code: 7, filler })).rejects.toThrow(/Row size exceeds limit/) + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id} AND data->>'code' = '7'` + expect(count).toBe(0) + } + ) + + it('fires the insert trigger of a batch insert with the column names of the live schema', async () => { + const table = await seededTable() + const stale: TableDefinition = { + ...table, + schema: { + columns: table.schema.columns.map((column) => + column.id === 'note' ? { ...column, name: 'old_note' } : column + ), + }, + } + tableTriggerMockFns.mockFireTableTrigger.mockClear() + + await batchInsertRows( + { + tableId: table.id, + workspaceId, + rows: [{ key: 'k3', note: 'n' }], + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + stale, + 'stale-schema' + ) + + const [trigger] = tableTriggerMockFns.mockFireTableTrigger.mock.calls + const schema = trigger?.[6] as { columns: ColumnDefinition[] } + expect(schema.columns.map((column) => column.name)).toContain('note') + }) + + it('does not count a legacy unique column named after an object prototype key as patched', async () => { + const table = await createTable([ + { name: 'constructor', type: 'string', unique: true }, + { id: 'note', name: 'note', type: 'string' }, + ]) + await seedRows(table.id, [ + { id: `${table.id}-a`, data: { constructor: 'a', note: 'n' }, orderKey: 'a0' }, + { id: `${table.id}-b`, data: { constructor: 'b', note: 'n' }, orderKey: 'a1' }, + ]) + + const result = await updateRowsByFilter( + table, + { + filter: { note: 'n' }, + data: { note: 'patched' }, + // The limited path: the paged one skips rows created after its JS-clock cutoff. + limit: 10, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + 'prototype-key' + ) + + expect(result.affectedCount).toBe(2) + }) + + it('does not count a required legacy column named after a prototype key as supplied', async () => { + const table = await createTable([ + { name: 'constructor', type: 'string', required: true }, + { id: 'note', name: 'note', type: 'string' }, + ]) + await seedRows(table.id, [ + { id: `${table.id}-a`, data: { constructor: 'a', note: 'n' }, orderKey: 'a0' }, + ]) + + const result = await updateRowsByFilter( + table, + { + filter: { note: 'n' }, + data: { note: 'patched' }, + limit: 10, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + 'prototype-key' + ) + + expect(result.affectedCount).toBe(1) + }) + + it('replaces rows that leave a unique legacy column named after a prototype key empty', async () => { + const table = await createTable([ + { name: 'constructor', type: 'string', unique: true }, + { id: 'note', name: 'note', type: 'string' }, + ]) + + await replaceTableRows( + { + tableId: table.id, + workspaceId, + rows: [{ note: 'a' }, { note: 'b' }], + secretProvenance: undefined, + }, + table, + 'prototype-key' + ) + + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id}` + expect(count).toBe(2) + }) + + it('refuses an upsert missing its conflict target named after a prototype key', async () => { + const table = await createTable([ + { name: 'constructor', type: 'string', unique: true }, + { id: 'note', name: 'note', type: 'string' }, + ]) + + await expect( + upsertRow( + { + tableId: table.id, + workspaceId, + data: { note: 'a' }, + conflictTarget: 'constructor', + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'prototype-key' + ) + ).rejects.toThrow(/requires a value for the conflict target/) + }) + + it('refuses a bulk update writing one value to rows of a column made unique since its snapshot', async () => { + const table = await seededTable() + await updateColumnConstraints( + { tableId: table.id, columnName: 'email', unique: true }, + 'stale-schema' + ) + + await expect( + updateRowsByFilter( + table, + { + filter: { note: 'n' }, + data: { email: 'same@example.test' }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + 'stale-schema' + ) + ).rejects.toThrow(/would violate uniqueness/) + expect(await countRows(table.id, 'email-dup')).toBe(1) + }) + + it('refuses an upsert on a column no longer unique since its snapshot', async () => { + const table = await seededTable() + await updateColumnConstraints( + { tableId: table.id, columnName: 'key', unique: false }, + 'stale-schema' + ) + + await expect( + upsertRow( + { + tableId: table.id, + workspaceId, + data: { key: 'k2', email: 'c@example.test', note: 'n' }, + conflictTarget: 'key', + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ) + ).rejects.toThrow(/requires at least one unique column/) + const [row] = await control<{ email: string }[]>`SELECT data->>'email' AS email + FROM user_table_rows WHERE id = ${`${table.id}-2`}` + expect(row.email).toBe('b@example.test') + }) + + /** + * Holds the table's schema lock exclusively in a transaction that marks `email` unique, starts + * `write`, and commits only once `write` is seen waiting on the lock, then settles as `write` does. + */ + async function writeDuringUniqueChange( + table: TableDefinition, + write: () => Promise + ): Promise { + const schemaLockKey = `user_table_schema:${table.id}` + const holder = await control.reserve() + try { + await holder`BEGIN` + await holder`SELECT pg_advisory_xact_lock(hashtextextended(${schemaLockKey}, 0))` + await holder`UPDATE user_table_definitions + SET schema = jsonb_set(schema, '{columns,1,unique}', 'true') + WHERE id = ${table.id}` + const pending = write() + const settled = pending.then( + () => 'fulfilled', + () => 'rejected' + ) + + await untilSchemaLockWaiters(table.id, 1) + expect(await Promise.race([settled, sleep(50).then(() => 'pending')])).toBe('pending') + + await holder`COMMIT` + return await pending + } finally { + await holder`ROLLBACK`.catch(() => {}) + holder.release() + } + } + + /** + * Holds the table's schema lock exclusively, starts `write`, and commits once `write` is seen + * waiting on it and `holdMs` more has passed. Settles as `write` does. + */ + async function writeBehindSchemaLock( + table: TableDefinition, + holdMs: number, + write: () => Promise + ): Promise { + const holder = await control.reserve() + try { + await holder`BEGIN` + await holder`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${table.id}`}, 0))` + const pending = write() + pending.catch(() => {}) + await untilSchemaLockWaiters(table.id, 1) + await sleep(holdMs) + await holder`COMMIT` + return await pending + } finally { + await holder`ROLLBACK`.catch(() => {}) + holder.release() + } + } + + it.each([ + [ + 'insertRow, which sets a 3 s lock timeout', + (table: TableDefinition) => + insertRow( + { + tableId: table.id, + workspaceId, + data: { key: 'k3', email: 'c@example.test' }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'long-schema-change' + ), + ], + [ + 'updateRow, which sets none', + (table: TableDefinition) => + updateRow( + { + tableId: table.id, + rowId: `${table.id}-2`, + workspaceId, + data: { email: 'c@example.test' }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'long-schema-change' + ), + ], + ])( + '%s waits out a schema change held past the lock timeout', + async (_, write) => { + const table = await seededTable() + + await expect(writeBehindSchemaLock(table, 3_500, () => write(table))).resolves.toBeDefined() + expect(await countRows(table.id, 'email-c')).toBe(1) + }, + 15_000 + ) + + it('bounds the schema-lock wait by the statement timeout', async () => { + const table = await seededTable() + + await expect( + writeBehindSchemaLock(table, 1_500, () => + db.transaction((trx) => lockLiveTableSchema(trx, table, { statementMs: 300 })) + ) + ).rejects.toMatchObject({ cause: { code: '55P03' } }) + }) + + it('leaves the lock timeout the transaction set for the locks after the guard', async () => { + const table = await seededTable() + + const lockTimeout = await db.transaction(async (trx) => { + await lockLiveTableSchema(trx, table, {}) + const [{ setting }] = await trx.execute<{ setting: string }>( + sql`SELECT current_setting('lock_timeout') AS setting` + ) + return setting + }) + expect(lockTimeout).toBe('3s') + }) + + it('waits for a schema change in flight and then writes against it', async () => { + const table = await seededTable() + + await expect( + writeDuringUniqueChange(table, () => + insertRow( + { + tableId: table.id, + workspaceId, + data: { key: 'k3', email: 'dup@example.test' }, + secretProvenance: undefined, + capabilityGovernedUserId: null, + }, + table, + 'stale-schema' + ) + ) + ).rejects.toThrow(/must be unique/) + expect(await countRows(table.id, 'email-dup')).toBe(1) + }) + }) + + describe('background updates validate their patch against the live schema', () => { + const columns: ColumnDefinition[] = [ + { id: 'email', name: 'email', type: 'string' }, + { id: 'kind', name: 'kind', type: 'string' }, + { + id: 'status', + name: 'status', + type: 'select', + options: [{ id: 'opt_open', name: 'Open' }], + }, + { id: 'note', name: 'note', type: 'string' }, + { id: 'due', name: 'due', type: 'date' }, + ] + /** Two update batches' worth of rows. */ + const ROWS = TABLE_LIMITS.UPDATE_BATCH_SIZE + 50 + + async function seededJob(data: RowData) { + const table = await createTable(columns) + await control`INSERT INTO user_table_rows (id, table_id, workspace_id, data, position, order_key) + SELECT ${table.id} || '-' || lpad(g::text, 4, '0'), ${table.id}, ${workspaceId}, + jsonb_build_object('email', 'e' || g || '@example.test', 'kind', 'seed'), g, 'a' || lpad(g::text, 4, '0') + FROM generate_series(1, ${ROWS}) g` + const jobId = generateId() + const filter = { kind: 'seed' } + expect( + await markTableJobRunningInWorkspace(table.id, workspaceId, jobId, 'update', { + filter, + data, + }) + ).toBe(true) + const run = () => + runTableUpdate({ + jobId, + tableId: table.id, + workspaceId, + filter, + data: { ...data }, + cutoff: new Date(), + }) + /** Re-runs the job as a retry after a crash would: the job is still running. */ + const retry = async () => { + await control`UPDATE table_jobs SET status = 'running', completed_at = NULL + WHERE id = ${jobId}` + return run() + } + return { table, run, retry } + } + + /** + * Commits `change` between the job's first and second batch. A held schema lock stops the + * first batch; `change` then queues behind it, and a waiting lock is granted in queue order, so + * the first batch commits, then `change`, then the second batch. + */ + async function changeBetweenBatches( + tableId: string, + run: () => Promise, + change: () => Promise + ): Promise[]> { + const holder = await control.reserve() + try { + await holder`BEGIN` + await holder`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${tableId}`}, 0))` + const job = run() + await untilSchemaLockWaiters(tableId, 1) + const changed = change() + await untilSchemaLockWaiters(tableId, 2) + await holder`COMMIT` + return await Promise.allSettled([job, changed]) + } finally { + await holder`ROLLBACK`.catch(() => {}) + holder.release() + } + } + + /** Commits `schema = ` for the table under its exclusive schema lock. */ + async function changeSchemaUnderLock(tableId: string, change: string): Promise { + const changer = await control.reserve() + try { + await changer`BEGIN` + await changer`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${tableId}`}, 0))` + await changer.unsafe(`UPDATE user_table_definitions SET schema = ${change} WHERE id = $1`, [ + tableId, + ]) + await changer`COMMIT` + } finally { + changer.release() + } + } + + async function countEmail(tableId: string, email: string): Promise { + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${tableId} AND data->>'email' = ${email}` + return count + } + + it('refuses a patch to a column made unique before the job started, retry included', async () => { + const { table, run } = await seededJob({ email: 'same@example.test' }) + await updateColumnConstraints( + { tableId: table.id, columnName: 'email', unique: true }, + 'bulk-update' + ) + + await expect(run()).rejects.toBeInstanceOf(UpdatePatchRejectedError) + await expect(run()).rejects.toBeInstanceOf(UpdatePatchRejectedError) + expect(await countEmail(table.id, 'same@example.test')).toBe(0) + }) + + it('refuses the next batch once a column it clears is made required mid-job', async () => { + const { table, run } = await seededJob({ note: null }) + await control`UPDATE user_table_rows SET data = data || '{"note":"filled"}' WHERE table_id = ${table.id}` + + const [job, change] = await changeBetweenBatches(table.id, run, async () => { + const changer = await control.reserve() + try { + await changer`BEGIN` + await changer`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${table.id}`}, 0))` + await changer`UPDATE user_table_definitions + SET schema = jsonb_set(schema, '{columns,3,required}', 'true') WHERE id = ${table.id}` + await changer`COMMIT` + } finally { + changer.release() + } + }) + + expect(change.status).toBe('fulfilled') + expect(job.status === 'rejected' && job.reason).toBeInstanceOf(UpdatePatchRejectedError) + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id} AND data->'note' = 'null'::jsonb` + expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE) + }) + + it('refuses the next batch once a patched column is made unique mid-job, retry included', async () => { + const { table, run } = await seededJob({ email: 'same@example.test' }) + + const [job, change] = await changeBetweenBatches(table.id, run, () => + changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,0,unique}', 'true')`) + ) + + expect(change.status).toBe('fulfilled') + expect(job.status === 'rejected' && job.reason).toBeInstanceOf(UpdatePatchRejectedError) + await expect(run()).rejects.toBeInstanceOf(UpdatePatchRejectedError) + expect(await countEmail(table.id, 'same@example.test')).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE) + }) + + it('finishes when a column it writes gains a select option mid-job', async () => { + const { table, run } = await seededJob({ status: 'Open' }) + + const [job, change] = await changeBetweenBatches(table.id, run, () => + updateColumnOptions( + { + tableId: table.id, + columnName: 'status', + options: [ + { id: 'opt_open', name: 'Open' }, + { id: 'opt_closed', name: 'Closed' }, + ], + }, + 'bulk-update' + ) + ) + + expect(change.status).toBe('fulfilled') + expect(job.status).toBe('fulfilled') + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id} AND data->>'status' = 'opt_open'` + expect(count).toBe(ROWS) + }) + + it('re-derives the patch for later batches once a patched column is retyped mid-job', async () => { + const { table, run, retry } = await seededJob({ note: '7' }) + const notes = async () => + control<{ kind: string; count: number }[]>`SELECT jsonb_typeof(data->'note') AS kind, + count(*)::int AS count FROM user_table_rows WHERE table_id = ${table.id} + GROUP BY 1 ORDER BY 1` + + const [job, change] = await changeBetweenBatches(table.id, run, () => + changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"number"')`) + ) + + expect([job.status, change.status]).toEqual(['fulfilled', 'fulfilled']) + expect(await notes()).toEqual([ + { kind: 'number', count: ROWS - TABLE_LIMITS.UPDATE_BATCH_SIZE }, + { kind: 'string', count: TABLE_LIMITS.UPDATE_BATCH_SIZE }, + ]) + const [{ seven }] = await control<{ seven: number }[]>`SELECT count(*)::int AS seven + FROM user_table_rows WHERE table_id = ${table.id} AND data->'note' IN ('"7"', '7')` + expect(seven).toBe(ROWS) + + await expect(retry()).resolves.toBeUndefined() + expect(await notes()).toEqual([{ kind: 'number', count: ROWS }]) + }) + + it('drops a column deleted mid-job from later batches', async () => { + const { table, run } = await seededJob({ note: 'x', email: 'y@example.test' }) + + const [job, change] = await changeBetweenBatches(table.id, run, () => + changeSchemaUnderLock(table.id, `schema #- '{columns,3}'`) + ) + + expect([job.status, change.status]).toEqual(['fulfilled', 'fulfilled']) + expect(await countEmail(table.id, 'y@example.test')).toBe(ROWS) + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id} AND data->>'note' = 'x'` + expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE) + }) + + it('refuses a batch whose re-derived patch grows a row past the size limit', async () => { + const { table, run } = await seededJob({ note: 7 }) + await changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"number"')`) + const bigRowId = `${table.id}-${String(TABLE_LIMITS.UPDATE_BATCH_SIZE + 1).padStart(4, '0')}` + // Stored jsonb orders keys by length, so a merged row reads kind, email, filler, then note. + const email = `e${TABLE_LIMITS.UPDATE_BATCH_SIZE + 1}@example.test` + const base = Buffer.byteLength(JSON.stringify({ kind: 'seed', email, filler: '', note: 7 })) + await control`UPDATE user_table_rows + SET data = data || jsonb_build_object('filler', repeat('x', ${getMaxRowSizeBytes() - base})) + WHERE id = ${bigRowId}` + + const [job, change] = await changeBetweenBatches(table.id, run, () => + changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"string"')`) + ) + + expect(change.status).toBe('fulfilled') + expect(job.status === 'rejected' && job.reason).toBeInstanceOf(UpdatePatchRejectedError) + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id} AND data ? 'note'` + expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE) + }) + + it('writes a date patch given as an epoch number', async () => { + const { table, run } = await seededJob({ due: 1704067200000 }) + + await run() + + const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count + FROM user_table_rows WHERE table_id = ${table.id} + AND (data->>'due')::timestamptz = '2024-01-01T00:00:00Z'` + expect(count).toBe(ROWS) + }) }) describe.skipIf(!migrated)('rows_version', () => { diff --git a/apps/sim/lib/table/rows/secret-provenance.integration.ts b/apps/sim/lib/table/rows/secret-provenance.integration.ts index 1eaae4ba5ee..a818d4aa5d2 100644 --- a/apps/sim/lib/table/rows/secret-provenance.integration.ts +++ b/apps/sim/lib/table/rows/secret-provenance.integration.ts @@ -188,7 +188,7 @@ describe('table provenance in PostgreSQL', () => { await connection`CREATE SCHEMA ${connection(testSchema)}` database.current = drizzle(connection, { schema }) await connection.unsafe(` - CREATE TABLE user_table_definitions (id text PRIMARY KEY, workspace_id text NOT NULL, rows_version integer NOT NULL); + CREATE TABLE user_table_definitions (id text PRIMARY KEY, workspace_id text NOT NULL, rows_version integer NOT NULL, schema jsonb NOT NULL); CREATE TABLE user_table_rows ( id text PRIMARY KEY, table_id text NOT NULL, workspace_id text NOT NULL, data jsonb NOT NULL, updated_at timestamp NOT NULL, secret_provenance_version integer, @@ -207,6 +207,9 @@ describe('table provenance in PostgreSQL', () => { row_id text PRIMARY KEY, content_updated_at timestamp NOT NULL, status text NOT NULL, entries jsonb NOT NULL, updated_at timestamp DEFAULT now() ); + CREATE FUNCTION user_table_schema_for_write(p_table_id text) RETURNS jsonb LANGUAGE sql VOLATILE AS $body$ + SELECT schema FROM user_table_definitions WHERE id = p_table_id + $body$; CREATE FUNCTION demote_changed_row() RETURNS trigger LANGUAGE plpgsql AS $body$ BEGIN IF NEW.data IS DISTINCT FROM OLD.data THEN @@ -225,7 +228,7 @@ describe('table provenance in PostgreSQL', () => { await connection.unsafe( 'TRUNCATE table_row_executions, user_table_rows, user_table_row_secret_provenance, user_table_definitions' ) - await connection`INSERT INTO user_table_definitions VALUES ('table-1', 'workspace-1', 7)` + await connection`INSERT INTO user_table_definitions VALUES ('table-1', 'workspace-1', 7, ${JSON.stringify(table.schema)}::jsonb)` }) afterAll(async () => { diff --git a/apps/sim/lib/table/rows/service.ts b/apps/sim/lib/table/rows/service.ts index b06adec9091..ae42156c51a 100644 --- a/apps/sim/lib/table/rows/service.ts +++ b/apps/sim/lib/table/rows/service.ts @@ -24,7 +24,7 @@ import { wouldExceedRowLimit, } from '@/lib/table/billing' import { getColumnId } from '@/lib/table/column-keys' -import { columnTypeOf, columnValueForEquality } from '@/lib/table/column-types' +import { columnTypeOf } from '@/lib/table/column-types' import { getMaxPageBytes, TABLE_LIMITS, USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants' import { TableQueryValidationError } from '@/lib/table/errors' import { @@ -50,6 +50,11 @@ import { tableMayHaveRunState, writeExecutionsPatch, } from '@/lib/table/rows/executions' +import { + dropDeletedColumns, + lockLiveTableSchema, + refitRowToSchema, +} from '@/lib/table/rows/live-schema' import { acquireRowOrderLock, type DeletedTableRow, @@ -73,7 +78,7 @@ import { buildPredicateClause, buildSortClause, escapeLikePattern, - fieldPredicate, + uniqueValuePredicate, } from '@/lib/table/sql' import { fireTableTrigger } from '@/lib/table/trigger' import { scaledStatementTimeoutMs, setTableTxTimeouts } from '@/lib/table/tx' @@ -99,17 +104,20 @@ import type { TableDefinition, TableRow, TableRowsCursor, + TableSchema, UpdateRowData, UpsertResult, UpsertRowData, } from '@/lib/table/types' import { + cellOf, checkBatchUniqueConstraintsDb, checkUniqueConstraintsDb, coerceRowToSchema, coerceRowValues, getUniqueColumns, type UncoercibleValuePolicy, + uniqueColumnsInPatch, uniqueValueKey, validateRowSize, } from '@/lib/table/validation' @@ -171,6 +179,7 @@ export async function insertRow( } // Validate against schema + const raw = { ...data.data } const schemaValidation = coerceRowToSchema(data.data, table.schema, options.uncoercibleValues) if (!schemaValidation.valid) { throw new OrchestrationError( @@ -188,6 +197,7 @@ export async function insertRow( const rowId = `row_${generateId().replace(/-/g, '')}` const now = new Date() + let written = table const row = await insertOrderedRow({ tableId: data.tableId, @@ -202,22 +212,37 @@ export async function insertRow( secretProvenance: data.secretProvenance, proof: insertProof, readProvenance: options.readProvenance, - assertUnique: - getUniqueColumns(table.schema).length > 0 - ? async (trx) => { - await lockUniqueValues(trx, table, [data.data]) - const uniqueValidation = await checkUniqueConstraintsDb( - data.tableId, - data.data, - table.schema, - undefined, - trx - ) - if (!uniqueValidation.valid) { - throw new OrchestrationError('validation', uniqueValidation.errors.join(', ')) - } - } - : undefined, + validate: async (trx) => { + const live = await lockLiveTableSchema(trx, table, {}) + written = live + if (live !== table) { + const refit = refitRowToSchema( + data.data, + raw, + table.schema, + live.schema, + options.uncoercibleValues + ) + if (!refit.valid) { + throw new OrchestrationError( + 'validation', + `Schema validation failed: ${refit.errors.join(', ')}` + ) + } + } + if (getUniqueColumns(live.schema).length === 0) return + await lockUniqueValues(trx, live, [data.data]) + const uniqueValidation = await checkUniqueConstraintsDb( + data.tableId, + data.data, + live.schema, + undefined, + trx + ) + if (!uniqueValidation.valid) { + throw new OrchestrationError('validation', uniqueValidation.errors.join(', ')) + } + }, }) notifyTableRowUsage({ @@ -246,7 +271,7 @@ export async function insertRow( 'insert', [insertedRow], null, - table.schema, + written.schema, requestId ) void runWorkflowColumn({ @@ -286,20 +311,30 @@ export async function batchInsertRows( addedRows: data.rows.length, }) - const result = await db.transaction((trx) => - batchInsertRowsWithTx(trx, data, table, requestId, options) + const { rows, table: written } = await db.transaction((trx) => + batchInsertRowsWithTx(trx, data, table, requestId, { ...options, lockSchema: true }) ) notifyTableRowUsage({ workspaceId: table.workspaceId, currentRowCount: table.rowCount, - addedRows: result.length, + addedRows: rows.length, limit: rowLimit, }) - dispatchAfterBatchInsert(table, result, requestId, data.userId, data.capabilityGovernedUserId) - return result + dispatchAfterBatchInsert(written, rows, requestId, data.userId, data.capabilityGovernedUserId) + return rows +} + +/** Options for the transaction-bound whole-row writers. */ +export interface TxRowWriteOptions extends RowWriteOptions { + /** + * Opens the transaction with the table's schema guard (see `live-schema.ts`) and validates + * against the live schema it reads. Callers that already read `table` under the schema lock in + * this transaction (imports, `withLockedTable`) leave it off. + */ + lockSchema?: boolean } -export interface BatchInsertOptions extends RowWriteOptions { +export interface BatchInsertOptions extends TxRowWriteOptions { /** * The caller already holds `lockUniqueColumns` for the table in this transaction (bulk imports, * which take it before the row-order lock), so the batch takes no per-value locks of its own. @@ -311,23 +346,28 @@ export interface BatchInsertOptions extends RowWriteOptions { * Transaction-bound variant of `batchInsertRows`. Validates rows and unique * constraints, then performs INSERTs inside the provided transaction. Caller * is responsible for opening the transaction. Use when row inserts must be - * atomic with other writes (e.g., schema mutations) on the same tx. + * atomic with other writes (e.g., schema mutations) on the same tx. Returns the inserted rows and + * the definition they were validated against, for the caller's post-commit dispatch. * * Capacity is NOT checked here (it would mean a billing-pool read inside the tx). * Callers gate it before opening the tx — see `batchInsertRows` and the import paths. * * Takes the rows' unique-value locks (see `unique-locks.ts`) before the unique check and the * row-order lock, so a caller must not already hold the row-order lock unless it passes - * `uniqueColumnsLocked`. + * `uniqueColumnsLocked`. Validates against `table` unless it passes `lockSchema`, so a caller that + * does not reads `table` under the table's schema lock in this transaction. */ export async function batchInsertRowsWithTx( trx: DbTransaction, data: BatchInsertData, - table: TableDefinition, + snapshot: TableDefinition, requestId: string, options: BatchInsertOptions = {} -): Promise { - assertRowInsert(table) +): Promise<{ rows: TableRow[]; table: TableDefinition }> { + assertRowInsert(snapshot) + const timeouts = { statementMs: 60_000 } + const table = options.lockSchema ? await lockLiveTableSchema(trx, snapshot, timeouts) : snapshot + if (table !== snapshot) dropDeletedColumns(data.rows, snapshot.schema, table.schema) for (let i = 0; i < data.rows.length; i++) { const row = data.rows[i] @@ -368,7 +408,7 @@ export async function batchInsertRowsWithTx( const now = new Date() - await setTableTxTimeouts(trx, { statementMs: 60_000 }) + if (!options.lockSchema) await setTableTxTimeouts(trx, timeouts) const buildRow = (rowData: RowData, position: number, orderKey: string) => ({ id: `row_${generateId().replace(/-/g, '')}`, @@ -418,7 +458,7 @@ export async function batchInsertRowsWithTx( })) await options.readProvenance?.capture(trx, result) - return result + return { rows: result, table } } /** @@ -490,7 +530,7 @@ export async function replaceTableRows( addedRows: data.rows.length, }) const result = await db.transaction((trx) => - replaceTableRowsWithTx(trx, data, table, requestId, options) + replaceTableRowsWithTx(trx, data, table, requestId, { ...options, lockSchema: true }) ) notifyTableRowUsage({ workspaceId: table.workspaceId, @@ -509,17 +549,29 @@ export async function replaceTableRows( * Callers gate it before opening the tx — see `replaceTableRows` and `importReplaceRows`. * * Takes the table's unique lock before its row-order lock, so a caller already holding the - * row-order lock must take `lockUniqueColumns` first. + * row-order lock must take `lockUniqueColumns` first. Validates against `table` unless it passes + * `lockSchema`, so a caller that does not reads `table` under the table's schema lock in this + * transaction. */ export async function replaceTableRowsWithTx( trx: DbTransaction, data: ReplaceRowsData, - table: TableDefinition, + snapshot: TableDefinition, requestId: string, - options: RowWriteOptions = {} + options: TxRowWriteOptions = {} ): Promise { - assertRowDelete(table) - assertRowInsert(table) + assertRowDelete(snapshot) + assertRowInsert(snapshot) + + const totalRowWork = Math.max(0, snapshot.rowCount ?? 0) + data.rows.length + const statementMs = scaledStatementTimeoutMs(totalRowWork, { + baseMs: 120_000, + perRowMs: 3, + }) + const table = options.lockSchema + ? await lockLiveTableSchema(trx, snapshot, { statementMs }) + : snapshot + if (table !== snapshot) dropDeletedColumns(data.rows, snapshot.schema, table.schema) if (data.tableId !== table.id) { throw new Error(`Table ID mismatch: ${data.tableId} vs ${table.id}`) @@ -560,11 +612,9 @@ export async function replaceTableRowsWithTx( // Coerced rows are keyed by column id, not display name — reading // `row[col.name]` silently misses renamed columns and lets dupes through. const colId = getColumnId(col) - const value = row[colId] + const value = cellOf(row, colId) if (value === null || value === undefined) continue - // Case-sensitive, consistent with the unique-constraint check leaf. - const comparable = columnValueForEquality(value, col) - const normalized = typeof comparable === 'string' ? comparable : JSON.stringify(comparable) + const normalized = uniqueValueKey(value, col) const map = seen.get(colId)! if (map.has(normalized)) { throw new OrchestrationError( @@ -579,13 +629,7 @@ export async function replaceTableRowsWithTx( const now = new Date() - const totalRowWork = Math.max(0, table.rowCount ?? 0) + data.rows.length - const statementMs = scaledStatementTimeoutMs(totalRowWork, { - baseMs: 120_000, - perRowMs: 3, - }) - - await setTableTxTimeouts(trx, { statementMs }) + if (!options.lockSchema) await setTableTxTimeouts(trx, { statementMs }) // Every current row is about to go, so any concurrent write of a unique value conflicts with the // replacement set: hold the unique columns exclusively, ahead of the row-order lock. @@ -662,28 +706,13 @@ export async function replaceTableRowsWithTx( } /** - * Upserts a row: updates an existing row if a match is found on the conflict target - * column, otherwise inserts a new row. - * - * Uses a single unique column for matching (not OR across all unique columns) to avoid - * ambiguous matches when multiple unique columns exist. Capacity is checked best-effort - * against the current plan limit on the insert path. On the insert path we acquire the - * per-table advisory lock and re-check for an existing match before inserting, so a - * concurrent upsert racing on the same conflict target cannot produce a duplicate row. - * - * @param data - Upsert data including optional conflictTarget - * @param table - Table definition - * @param requestId - Request ID for logging - * @returns The upserted row and whether it was an insert or update - * @throws Error if no unique columns, ambiguous conflict target, or capacity exceeded + * The single unique column an upsert matches on, resolved from `conflictTarget` — an id + * (first-party) or a name (legacy/internal) — or, without one, the table's only unique column. */ -export async function upsertRow( - data: UpsertRowData, - table: TableDefinition, - requestId: string, - options: RowWriteOptions = {} -): Promise { - const schema = table.schema +function resolveUpsertTarget( + schema: TableSchema, + conflictTarget: string | undefined +): ColumnDefinition { const uniqueColumns = getUniqueColumns(schema) if (uniqueColumns.length === 0) { @@ -693,13 +722,9 @@ export async function upsertRow( ) } - // Determine the single conflict target column, resolving to its stable - // storage id (the row-data key). `conflictTarget` may arrive as an id - // (first-party) or a name (legacy/internal) — match either. - let targetColumnKey: string - if (data.conflictTarget) { + if (conflictTarget) { const col = uniqueColumns.find( - (c) => getColumnId(c) === data.conflictTarget || c.name === data.conflictTarget + (c) => getColumnId(c) === conflictTarget || c.name === conflictTarget ) if (!col) { /** @@ -707,25 +732,65 @@ export async function upsertRow( * `conflictTarget` to its storage id before this call, so echoing the * argument verbatim answers a request naming `email` with a `col_…` id * the caller has never seen and cannot map back. Same rule as the - * missing-value branch below. + * missing-value branch in {@link upsertConflictProbe}. */ const requested = - schema.columns.find((c) => getColumnId(c) === data.conflictTarget)?.name ?? - data.conflictTarget + schema.columns.find((c) => getColumnId(c) === conflictTarget)?.name ?? conflictTarget throw new OrchestrationError( 'validation', `Column "${requested}" is not a unique column. Available unique columns: ${uniqueColumns.map((c) => c.name).join(', ')}` ) } - targetColumnKey = getColumnId(col) - } else if (uniqueColumns.length === 1) { - targetColumnKey = getColumnId(uniqueColumns[0]) - } else { + return col + } + if (uniqueColumns.length === 1) return uniqueColumns[0] + throw new OrchestrationError( + 'validation', + `Table has multiple unique columns (${uniqueColumns.map((c) => c.name).join(', ')}). Specify a conflict column to indicate which one to match on.` + ) +} + +/** + * The upsert's conflict probe: the coerced row's `target` value, matched through the same + * {@link uniqueValuePredicate} as the unique check, so "find the row to update" and "is this value + * unique" agree. Read after coercion so the probe matches the persisted type (a coerced `"123"` + * matches a stored `123`). + */ +function upsertConflictProbe(target: ColumnDefinition, row: RowData): SQL { + const targetValue = cellOf(row, getColumnId(target)) + if (targetValue === undefined || targetValue === null) { + // Surface the display name, not the internal id — v1 callers pass a name. throw new OrchestrationError( 'validation', - `Table has multiple unique columns (${uniqueColumns.map((c) => c.name).join(', ')}). Specify a conflict column to indicate which one to match on.` + `Upsert requires a value for the conflict target column "${target.name}"` ) } + return uniqueValuePredicate(USER_TABLE_ROWS_SQL_NAME, getColumnId(target), targetValue, target) +} + +/** + * Upserts a row: updates an existing row if a match is found on the conflict target + * column, otherwise inserts a new row. + * + * Uses a single unique column for matching (not OR across all unique columns) to avoid + * ambiguous matches when multiple unique columns exist. Capacity is checked best-effort + * against the current plan limit on the insert path. On the insert path we acquire the + * per-table advisory lock and re-check for an existing match before inserting, so a + * concurrent upsert racing on the same conflict target cannot produce a duplicate row. + * + * @param data - Upsert data including optional conflictTarget + * @param table - Table definition + * @param requestId - Request ID for logging + * @returns The upserted row and whether it was an insert or update + * @throws Error if no unique columns, ambiguous conflict target, or capacity exceeded + */ +export async function upsertRow( + data: UpsertRowData, + table: TableDefinition, + requestId: string, + options: RowWriteOptions = {} +): Promise { + const target = resolveUpsertTarget(table.schema, data.conflictTarget) // Validate row data const sizeValidation = validateRowSize(data.data) @@ -733,7 +798,8 @@ export async function upsertRow( throw new OrchestrationError('validation', sizeValidation.errors.join(', ')) } - const schemaValidation = coerceRowToSchema(data.data, schema, options.uncoercibleValues) + const raw = { ...data.data } + const schemaValidation = coerceRowToSchema(data.data, table.schema, options.uncoercibleValues) if (!schemaValidation.valid) { throw new OrchestrationError( 'validation', @@ -741,49 +807,44 @@ export async function upsertRow( ) } - // Read the conflict-target value *after* coercion so `matchFilter` branches on - // the persisted type (e.g. a coerced `"123"` → `123` matches existing rows). - const targetValue = data.data[targetColumnKey] - if (targetValue === undefined || targetValue === null) { - // Surface the display name, not the internal id — v1 callers pass a name. - const targetColumnName = - uniqueColumns.find((c) => getColumnId(c) === targetColumnKey)?.name ?? targetColumnKey - throw new OrchestrationError( - 'validation', - `Upsert requires a value for the conflict target column "${targetColumnName}"` - ) - } - - // Build the conflict probe through the SAME leaf as the unique-constraint check - // (`fieldPredicate` → case-sensitive JSONB containment). This is what makes - // "find the row to update" and "is this value unique" agree: a value differing - // only in case can no longer slip past the probe and then trip the guard. - // `eq` always yields a clause for a non-null value (guaranteed above). - const matchFilter = fieldPredicate( - USER_TABLE_ROWS_SQL_NAME, - targetColumnKey, - 'eq', - targetValue, - table.schema.columns.find((c) => getColumnId(c) === targetColumnKey) - ) - if (!matchFilter) { - throw new Error('Failed to build upsert conflict predicate') - } + const snapshotMatchFilter = upsertConflictProbe(target, data.data) // Resolve the plan limit BEFORE the tx (the lookup is a separate pool read; doing // it inside the tx would hold a connection + the row-order lock during it). The // insert branch enforces it; the update path doesn't add a row, so it's exempt. const rowLimit = await getMaxRowsPerTable(table.workspaceId) + let written = table const result = await db.transaction(async (trx) => { - await setTableTxTimeouts(trx) + const live = await lockLiveTableSchema(trx, table, {}) + written = live // The conflict lookups below match on `data->>key` — unestimatable, and an // insert-path upsert (no existing match) can't exit early, so the planner // would seq-scan the whole shared relation. See withSeqscanOff. await trx.execute(sql`SET LOCAL enable_seqscan = off`) + let matchFilter = snapshotMatchFilter + if (live !== table) { + const refit = refitRowToSchema( + data.data, + raw, + table.schema, + live.schema, + options.uncoercibleValues + ) + if (!refit.valid) { + throw new OrchestrationError( + 'validation', + `Schema validation failed: ${refit.errors.join(', ')}` + ) + } + matchFilter = upsertConflictProbe( + resolveUpsertTarget(live.schema, data.conflictTarget), + data.data + ) + } // Holds every unique value this row writes, the conflict target included, before the lookup // and the unique check below: concurrent upserts and inserts of the same value serialize here. - await lockUniqueValues(trx, table, [data.data]) + await lockUniqueValues(trx, live, [data.data]) // Find existing row by single conflict target column const [existingRow] = await trx @@ -802,7 +863,7 @@ export async function upsertRow( const uniqueValidation = await checkUniqueConstraintsDb( data.tableId, data.data, - schema, + live.schema, existingRow?.id, // exclude the matched row on updates trx ) @@ -945,7 +1006,7 @@ export async function upsertRow( 'insert', [result.row], null, - table.schema, + written.schema, requestId ) } else if (result.operation === 'update' && result.previousData) { @@ -957,7 +1018,7 @@ export async function upsertRow( 'update', [result.row], oldRows, - table.schema, + written.schema, requestId ) } @@ -1767,7 +1828,7 @@ export async function updateRow( // Auto-clear exec records for workflow output columns the user just wiped // AND for downstream groups whose deps just changed. Surfaces the in-flight // downstream groups so the caller can cancel + re-run them. - const { executionsPatch: effectiveExecutionsPatch, inFlightDownstreamGroups } = + let { executionsPatch: effectiveExecutionsPatch, inFlightDownstreamGroups } = deriveExecClearsForDataPatch( data.data, table.schema, @@ -1775,7 +1836,7 @@ export async function updateRow( data.executionsPatch, mergedData ) - const mergedExecutions = applyExecutionsPatch(existingRow.executions, effectiveExecutionsPatch) + let mergedExecutions = applyExecutionsPatch(existingRow.executions, effectiveExecutionsPatch) // Validate size const sizeValidation = validateRowSize(mergedData) @@ -1810,12 +1871,8 @@ export async function updateRow( // unrelated cell edit should not fail on data it did not write, and blocking // it was never a repair mechanism. const patchedColumnIds = new Set(Object.keys(data.data)) - const patchedUniqueColumns = getUniqueColumns(table.schema).filter((column) => - patchedColumnIds.has(getColumnId(column)) - ) const now = new Date() - const persistedData = jsonbMergePatch(Object.keys(data.data), mergedData) // Cell-task partial writes pass `cancellationGuard` so the upsert into // `tableRowExecutions` is a no-op when (a) a stop click already wrote @@ -1826,12 +1883,48 @@ export async function updateRow( // and the row out of sync. const guard = data.cancellationGuard let persistedRow: typeof userTableRows.$inferSelect + let written = table try { persistedRow = await db.transaction(async (trx) => { + // An executions-only write stores no cell, so it has no schema to hold still. + const live = patchedColumnIds.size > 0 ? await lockLiveTableSchema(trx, table) : table + written = live + if (live !== table) { + const refit = refitRowToSchema( + mergedData, + { ...(existingRow.data as RowData), ...data.data }, + table.schema, + live.schema, + options.uncoercibleValues, + Object.keys(data.data) + ) + if (!refit.valid) { + throw new OrchestrationError( + 'validation', + `Schema validation failed: ${refit.errors.join(', ')}` + ) + } + ;({ executionsPatch: effectiveExecutionsPatch, inFlightDownstreamGroups } = + deriveExecClearsForDataPatch( + data.data, + live.schema, + existingRow.executions, + data.executionsPatch, + mergedData + )) + mergedExecutions = applyExecutionsPatch(existingRow.executions, effectiveExecutionsPatch) + } + const persistedData = jsonbMergePatch( + Object.keys(data.data).filter((columnId) => Object.hasOwn(mergedData, columnId)), + mergedData + ) + const patchedUniqueColumns = getUniqueColumns(live.schema).filter((column) => + patchedColumnIds.has(getColumnId(column)) + ) if (patchedUniqueColumns.length > 0) { - // The value locks come first in the transaction, so a concurrent write of the same value - // commits before this check runs, or waits for this edit to commit. - await lockUniqueValues(trx, table, [mergedData], patchedColumnIds) + // The value locks come first after the schema lock, so a concurrent write of the same + // value commits before this check runs, or waits for this edit to commit. + await lockUniqueValues(trx, live, [mergedData], patchedColumnIds) const uniqueValidation = await checkUniqueConstraintsDb( data.tableId, mergedData, @@ -1839,7 +1932,7 @@ export async function updateRow( // probe issues one SELECT per unique column it is given, so a table with // several would otherwise re-check all of them to validate a patch that // touched one. `schema` is read only for its unique columns here. - { ...table.schema, columns: patchedUniqueColumns }, + { ...live.schema, columns: patchedUniqueColumns }, data.rowId, // Exclude current row trx ) @@ -1917,7 +2010,7 @@ export async function updateRow( 'update', [updatedRow], oldRows, - table.schema, + written.schema, requestId ) @@ -2038,7 +2131,7 @@ function validateBulkUpdatePatch( policy: UncoercibleValuePolicy | undefined ): void { const suppliedColumns = table.schema.columns.filter((column) => { - const value = patch[getColumnId(column)] + const value = cellOf(patch, getColumnId(column)) return value !== null && value !== undefined }) if (suppliedColumns.length === 0) return @@ -2053,6 +2146,26 @@ function validateBulkUpdatePatch( } } +/** + * The patch a bulk update writes against `live`: `raw` with the cells of columns deleted since + * `previous` dropped, coerced with `policy`, and validated on its own (see + * {@link validateBulkUpdatePatch}). The inline update derives its patch here, again whenever the + * schema moves under it, and the background runner derives each batch's patch here, so the two + * cannot disagree about what one request writes. + */ +export function deriveBulkUpdatePatch( + raw: RowData, + previous: TableSchema, + live: TableDefinition, + policy: UncoercibleValuePolicy | undefined +): RowData { + const patch = { ...raw } + dropDeletedColumns([patch], previous, live.schema) + coerceRowValues(patch, live.schema, policy) + validateBulkUpdatePatch(live, patch, policy) + return patch +} + /** Validates a bounded page of rows against a bulk merge patch. */ function validateBulkUpdateMatches( table: TableDefinition, @@ -2066,10 +2179,55 @@ function validateBulkUpdateMatches( } } -/** Persists one bounded bulk-update page and returns the locked rows actually changed. */ +/** + * Refuses a bulk patch that writes a unique column into more than one row, and otherwise locks and + * checks the one row's unique values. The value locks come first after the schema lock, so a + * concurrent write of the same value is committed and visible to the check. + */ +async function assertBulkUpdateUnique( + trx: DbTransaction, + table: TableDefinition, + rows: BulkUpdateMatch[], + patch: RowData +): Promise { + const uniqueColumnsInUpdate = getUniqueColumns(table.schema).filter((column) => + Object.hasOwn(patch, getColumnId(column)) + ) + if (uniqueColumnsInUpdate.length === 0 || rows.length === 0) return + if (rows.length > 1) { + throw new OrchestrationError( + 'validation', + `Cannot set unique column values when updating multiple rows. ` + + `Columns with unique constraint: ${uniqueColumnsInUpdate.map((column) => column.name).join(', ')}. ` + + `Updating ${rows.length} rows with the same value would violate uniqueness.` + ) + } + await lockUniqueValues(trx, table, [patch]) + const uniqueValidation = await checkUniqueConstraintsDb( + table.id, + { ...rows[0].data, ...patch }, + table.schema, + rows[0].id, + trx + ) + if (!uniqueValidation.valid) { + throw new OrchestrationError( + 'validation', + `Unique constraint violation: ${uniqueValidation.errors.join(', ')}` + ) + } +} + +/** + * Persists one bounded bulk-update page and returns the locked rows actually changed, with the + * patch as written and the definition it was written against: re-coerced, and live, when the + * schema changed since the caller prepared it. + */ async function persistBulkUpdateBatch(params: { table: TableDefinition rows: BulkUpdateMatch[] + /** The request's patch before coercion, re-derived when the schema moved. */ + raw: RowData patch: RowData patchJson: string filterClause: SQL @@ -2077,26 +2235,27 @@ async function persistBulkUpdateBatch(params: { secretProvenance: BulkUpdateData['secretProvenance'] requestId: string uncoercibleValues: UncoercibleValuePolicy | undefined - /** Runs first in the write transaction: the unique-value locks and check of a unique patch. */ - assertUnique?: (trx: DbTransaction) => Promise -}): Promise<{ rows: BulkUpdateMatch[]; affectedRowIds: string[] }> { - const { - table, - rows, - patch, - patchJson, - filterClause, - now, - secretProvenance, - requestId, - uncoercibleValues, - assertUnique, - } = params +}): Promise<{ + rows: BulkUpdateMatch[] + affectedRowIds: string[] + patch: RowData + table: TableDefinition +}> { + const { table: snapshot, rows, raw, filterClause, now, secretProvenance, requestId } = params + const { uncoercibleValues } = params + let { patch, patchJson } = params const ids = rows.map((row) => row.id) const persistedRows: BulkUpdateMatch[] = [] + let written = snapshot const affectedRowIds = await db.transaction(async (trx) => { - await setTableTxTimeouts(trx, { statementMs: 60_000 }) - await assertUnique?.(trx) + const table = await lockLiveTableSchema(trx, snapshot, { statementMs: 60_000 }) + written = table + if (table !== snapshot) { + patch = deriveBulkUpdatePatch(raw, snapshot.schema, table, uncoercibleValues) + validateBulkUpdateMatches(table, rows, patch, uncoercibleValues) + patchJson = JSON.stringify(patch) + } + await assertBulkUpdateUnique(trx, table, rows, patch) return mutateTableRowsWithSecretProvenance(trx, { rows: ids.map((rowId) => ({ rowId, provenance: secretProvenance })), rowState: 'existing', @@ -2152,7 +2311,7 @@ async function persistBulkUpdateBatch(params: { }, }) }) - return { rows: persistedRows, affectedRowIds } + return { rows: persistedRows, affectedRowIds, patch, table: written } } /** Emits trigger and enrichment side effects for one committed bulk-update page. */ @@ -2235,40 +2394,15 @@ export async function updateRowsByFilter( eq(userTableRows.workspaceId, table.workspaceId) ) - coerceRowValues(data.data, table.schema, options.uncoercibleValues) - validateBulkUpdatePatch(table, data.data, options.uncoercibleValues) - const uniqueColumns = getUniqueColumns(table.schema) - const uniqueColumnsInUpdate = uniqueColumns.filter((col) => getColumnId(col) in data.data) - const patchJson = JSON.stringify(data.data) + const patch = deriveBulkUpdatePatch(data.data, table.schema, table, options.uncoercibleValues) + const uniqueColumnsInUpdate = uniqueColumnsInPatch(table.schema, patch) + const patchJson = JSON.stringify(patch) const now = new Date() const limit = data.limit - /** - * A unique patch reaches exactly one row. Its value locks and check run first in the write - * transaction, so a concurrent write of the same value is committed and visible to the check. - */ - const uniqueGuard = - (row: BulkUpdateMatch) => - async (trx: DbTransaction): Promise => { - await lockUniqueValues(trx, table, [data.data]) - const uniqueValidation = await checkUniqueConstraintsDb( - table.id, - { ...row.data, ...data.data }, - table.schema, - row.id, - trx - ) - if (!uniqueValidation.valid) { - throw new OrchestrationError( - 'validation', - `Unique constraint violation: ${uniqueValidation.errors.join(', ')}` - ) - } - } if (limit === undefined) { const cutoff = new Date() let matchingRowCount = 0 - let singleMatchingRow: BulkUpdateMatch | undefined let afterId: string | undefined while (true) { @@ -2282,9 +2416,8 @@ export async function updateRowsByFilter( }) if (page.length === 0) break - validateBulkUpdateMatches(table, page, data.data, options.uncoercibleValues) + validateBulkUpdateMatches(table, page, patch, options.uncoercibleValues) matchingRowCount += page.length - singleMatchingRow ??= page[0] afterId = page[page.length - 1].id if (page.length < TABLE_LIMITS.UPDATE_BATCH_SIZE) break } @@ -2302,14 +2435,7 @@ export async function updateRowsByFilter( `Updating ${matchingRowCount} rows with the same value would violate uniqueness.` ) } - if (!singleMatchingRow) { - throw new Error('Bulk update lost its selected row') - } } - const assertUnique = - uniqueColumnsInUpdate.length > 0 && singleMatchingRow - ? uniqueGuard(singleMatchingRow) - : undefined const affectedRowIds: string[] = [] afterId = undefined @@ -2328,21 +2454,21 @@ export async function updateRowsByFilter( const persisted = await persistBulkUpdateBatch({ table, rows: batchRows, - patch: data.data, + raw: data.data, + patch, patchJson, filterClause, now, secretProvenance: data.secretProvenance, requestId, uncoercibleValues: options.uncoercibleValues, - assertUnique, }) affectedRowIds.push(...persisted.affectedRowIds) dispatchBulkUpdateEffects( - table, + persisted.table, persisted.rows, persisted.affectedRowIds, - data.data, + persisted.patch, now, requestId, data.actorUserId, @@ -2369,7 +2495,7 @@ export async function updateRowsByFilter( return { affectedCount: 0, affectedRowIds: [] } } - validateBulkUpdateMatches(table, matchingRows, data.data, options.uncoercibleValues) + validateBulkUpdateMatches(table, matchingRows, patch, options.uncoercibleValues) if (uniqueColumnsInUpdate.length > 0) { if (matchingRows.length > 1) { throw new OrchestrationError( @@ -2384,23 +2510,23 @@ export async function updateRowsByFilter( const persisted = await persistBulkUpdateBatch({ table, rows: matchingRows, - patch: data.data, + raw: data.data, + patch, patchJson, filterClause, now, secretProvenance: data.secretProvenance, requestId, uncoercibleValues: options.uncoercibleValues, - assertUnique: uniqueColumnsInUpdate.length > 0 ? uniqueGuard(matchingRows[0]) : undefined, }) const { affectedRowIds } = persisted logger.info(`[${requestId}] Updated ${affectedRowIds.length} rows in table ${table.id}`) dispatchBulkUpdateEffects( - table, + persisted.table, persisted.rows, affectedRowIds, - data.data, + persisted.patch, now, requestId, data.actorUserId, @@ -2413,6 +2539,54 @@ export async function updateRowsByFilter( } } +/** One update of a batch, merged over its stored row and coerced to the schema. */ +interface MergedRowUpdate { + rowId: string + changedColumnIds: string[] + mergedData: RowData + mergedExecutions: RowExecutions + executionsPatch?: Record + inFlightDownstreamGroups: string[] +} + +/** + * The unique columns each update of a batch changes, for its value locks and unique check. Like + * `updateRow`, an update checks only the unique columns it changes: a value it leaves alone is the + * one already stored, and updates that change none cost nothing here. The DB check runs before any + * write and the value locks are deduplicated, so two updates writing the same value would both pass + * it; they are rejected here instead, keyed like the locks. + */ +function planBatchUpdateUniqueChecks( + schema: TableSchema, + updates: readonly MergedRowUpdate[] +): Array<{ rowId: string; mergedData: RowData; columns: ColumnDefinition[] }> { + const uniqueColumns = getUniqueColumns(schema) + const uniqueChecks: Array<{ rowId: string; mergedData: RowData; columns: ColumnDefinition[] }> = + [] + if (uniqueColumns.length === 0) return uniqueChecks + const firstWriter = new Map() + for (const { rowId, changedColumnIds, mergedData } of updates) { + const changed = new Set(changedColumnIds) + const columns = uniqueColumns.filter((column) => changed.has(getColumnId(column))) + if (columns.length === 0) continue + uniqueChecks.push({ rowId, mergedData, columns }) + for (const column of columns) { + const value = mergedData[getColumnId(column)] + if (value === null || value === undefined) continue + const key = `${getColumnId(column)}:${uniqueValueKey(value, column)}` + const otherRowId = firstWriter.get(key) + if (otherRowId !== undefined && otherRowId !== rowId) { + throw new OrchestrationError( + 'validation', + `Row ${rowId}: Column "${column.name}" must be unique. Value "${String(value)}" duplicates row ${otherRowId} in batch` + ) + } + firstWriter.set(key, rowId) + } + } + return uniqueChecks +} + export interface BatchUpdateRowsOptions extends RowWriteOptions { /** * Marks the batch as workflow/enrichment output cells (the backfill runner), @@ -2476,14 +2650,7 @@ export async function batchUpdateRows( throw new OrchestrationError('validation', `Rows not found: ${missing.join(', ')}`) } - const mergedUpdates: Array<{ - rowId: string - changedColumnIds: string[] - mergedData: RowData - mergedExecutions: RowExecutions - executionsPatch?: Record - inFlightDownstreamGroups: string[] - }> = [] + const mergedUpdates: MergedRowUpdate[] = [] for (const update of data.updates) { const existing = existingMap.get(update.rowId)! const merged = { ...existing.data, ...update.data } @@ -2532,44 +2699,52 @@ export async function batchUpdateRows( }) } - // Like `updateRow`, each update locks and checks only the unique columns it changes: a value it - // leaves alone is the one already stored. Updates that change none cost nothing extra here. - const uniqueColumns = getUniqueColumns(table.schema) - const uniqueChecks: Array<{ rowId: string; mergedData: RowData; columns: ColumnDefinition[] }> = - [] - if (uniqueColumns.length > 0) { - // The DB check runs before any write and the value locks are deduplicated, so two updates in - // this batch writing the same value would both pass it; catch them here, keyed like the locks. - const firstWriter = new Map() - for (const { rowId, changedColumnIds, mergedData } of mergedUpdates) { - const changed = new Set(changedColumnIds) - const columns = uniqueColumns.filter((column) => changed.has(getColumnId(column))) - if (columns.length === 0) continue - uniqueChecks.push({ rowId, mergedData, columns }) - for (const column of columns) { - const value = mergedData[getColumnId(column)] - if (value === null || value === undefined) continue - const key = `${getColumnId(column)}:${uniqueValueKey(value, column)}` - const otherRowId = firstWriter.get(key) - if (otherRowId !== undefined && otherRowId !== rowId) { + let uniqueChecks = planBatchUpdateUniqueChecks(table.schema, mergedUpdates) + + const now = new Date() + let written = table + + const affectedRowIds = await db.transaction(async (trx) => { + const live = await lockLiveTableSchema(trx, table, { statementMs: 60_000 }) + written = live + if (live !== table) { + for (const [index, update] of mergedUpdates.entries()) { + const request = data.updates[index] + const existing = existingMap.get(update.rowId)! + const refit = refitRowToSchema( + update.mergedData, + { ...existing.data, ...request.data }, + table.schema, + live.schema, + options.uncoercibleValues, + update.changedColumnIds + ) + if (!refit.valid) { throw new OrchestrationError( 'validation', - `Row ${rowId}: Column "${column.name}" must be unique. Value "${String(value)}" duplicates row ${otherRowId} in batch` + `Row ${update.rowId}: ${refit.errors.join(', ')}` ) } - firstWriter.set(key, rowId) + update.changedColumnIds = update.changedColumnIds.filter((columnId) => + Object.hasOwn(update.mergedData, columnId) + ) + const cleared = deriveExecClearsForDataPatch( + request.data, + live.schema, + existing.executions, + request.executionsPatch, + update.mergedData + ) + update.executionsPatch = cleared.executionsPatch + update.inFlightDownstreamGroups = cleared.inFlightDownstreamGroups + update.mergedExecutions = applyExecutionsPatch(existing.executions, cleared.executionsPatch) } + uniqueChecks = planBatchUpdateUniqueChecks(live.schema, mergedUpdates) } - } - - const now = new Date() - - const affectedRowIds = await db.transaction(async (trx) => { - await setTableTxTimeouts(trx, { statementMs: 60_000 }) if (uniqueChecks.length > 0) { await lockUniqueValues( trx, - table, + live, uniqueChecks.map(({ mergedData, columns }) => Object.fromEntries( columns.map((column) => [getColumnId(column), mergedData[getColumnId(column)]]) @@ -2580,7 +2755,7 @@ export async function batchUpdateRows( const uniqueValidation = await checkUniqueConstraintsDb( data.tableId, mergedData, - { ...table.schema, columns }, + { ...live.schema, columns }, rowId, trx ) @@ -2653,7 +2828,7 @@ export async function batchUpdateRows( 'update', updatedRowsForTrigger, oldRowsForTrigger, - table.schema, + written.schema, requestId ) } diff --git a/apps/sim/lib/table/rows/unique-locks.ts b/apps/sim/lib/table/rows/unique-locks.ts index 313c4f79076..e5b2ee8b40e 100644 --- a/apps/sim/lib/table/rows/unique-locks.ts +++ b/apps/sim/lib/table/rows/unique-locks.ts @@ -7,9 +7,9 @@ * values it writes, so the second writer's check runs after the first commits and sees its row. * Writers of different values never wait on each other. * - * Lock order, everywhere: the table's schema lock (when taken), then the table's unique lock and its - * value locks in sorted order, in one statement, then the table's row-order lock, then the - * definition row. + * Lock order, everywhere: the table's schema lock (shared by row writes, exclusive by schema + * changes; see `live-schema.ts`), then the table's unique lock and its value locks in sorted order, + * in one statement, then the table's row-order lock, then the definition row. */ import { compareStrings } from '@sim/utils/string' @@ -17,7 +17,7 @@ import { type AdvisoryXactLockRequest, acquireAdvisoryXactLocks } from '@/lib/db import { getColumnId } from '@/lib/table/column-keys' import type { DbTransaction } from '@/lib/table/planner' import type { RowData, TableDefinition } from '@/lib/table/types' -import { getUniqueColumns, uniqueValueKey } from '@/lib/table/validation' +import { cellOf, getUniqueColumns, uniqueValueKey } from '@/lib/table/validation' const UNIQUE_LOCK_TAG = 'user_table_unique_value' @@ -45,11 +45,11 @@ function tableLockKey(tableId: string): string { * lock only the unique columns a patch changes; values it leaves alone are already stored. * * The table lock is taken shared, then an exclusive lock per value. The value key includes - * `uniqueValueKey`, so two values the check treats as equal always share a key; null cells never - * conflict and take none. The table lock is taken exclusively instead, with no value locks, when a - * value is an object or array (a `json` column's check matches by containment, which no single key - * can express) or when the transaction would exceed {@link MAX_VALUE_LOCKS}. An exclusive holder - * waits for, and blocks, every shared holder, so it still serializes against per-value writers. + * `uniqueValueKey`, so two values the check treats as equal always share a key, objects and arrays + * included, since the check compares those by jsonb equality; null cells never conflict and take + * none. The table lock is taken exclusively instead, with no value locks, when the transaction would + * exceed {@link MAX_VALUE_LOCKS}. An exclusive holder waits for, and blocks, every shared holder, so + * it still serializes against per-value writers. */ export async function lockUniqueValues( trx: DbTransaction, @@ -63,9 +63,8 @@ export async function lockUniqueValues( const columnId = getColumnId(column) if (columnIds && !columnIds.has(columnId)) continue for (const row of rows) { - const value = row[columnId] + const value = cellOf(row, columnId) if (value === null || value === undefined) continue - if (typeof value === 'object') return lockUniqueColumns(trx, table) const key = `${tableKey}:${columnId}:${uniqueValueKey(value, column)}` if (!valueKeys.has(key) && valueKeys.size === MAX_VALUE_LOCKS) { return lockUniqueColumns(trx, table) diff --git a/apps/sim/lib/table/sql.ts b/apps/sim/lib/table/sql.ts index d154a6f44db..0f287ead4cb 100644 --- a/apps/sim/lib/table/sql.ts +++ b/apps/sim/lib/table/sql.ts @@ -532,10 +532,11 @@ function buildFieldCondition( /** * The single leaf primitive: compiles one `field op value` into SQL. Every * matcher routes through here — both filter compilers (`buildFilterClause` for - * the legacy `$`-grammar, `buildPredicateClause` for the v2 grammar), the upsert - * conflict probe, and the unique-constraint checks. Centralizing the leaf means - * equality/case/null/cast semantics are defined exactly once, so "find the row" - * and "is this value unique" can never disagree. + * the legacy `$`-grammar, `buildPredicateClause` for the v2 grammar), and, through + * {@link uniqueValuePredicate}, the upsert conflict probe and the unique-constraint + * checks. Centralizing the leaf means equality/case/null/cast semantics are + * defined exactly once, so "find the row" and "is this value unique" can never + * disagree. * * Returns `undefined` when the predicate is a no-op (empty `in`/`nin` array), * matching the legacy behavior of emitting no clause. @@ -869,6 +870,29 @@ function buildSystemColumnClause( } } +/** + * Whether a row holds `value` in a unique column: the upsert conflict probe and every unique check + * match with it. It is the `eq` leaf of {@link fieldPredicate}, made exact for an object or array. + * Containment is not equality there (`{"a":1,"b":2} @> {"a":1}`), so the leaf ANDs jsonb equality + * on the cell, keeping the containment half for the GIN index. jsonb equality ignores object key + * order and keeps array order, as `uniqueValueKey` does. A scalar leaf's containment is already + * equality. + */ +export function uniqueValuePredicate( + tableName: string, + field: string, + value: JsonValue, + column: ColumnDefinition +): SQL { + const clause = fieldPredicate(tableName, field, 'eq', value, column) + if (!clause) { + throw new Error(`Failed to build unique-constraint predicate for column "${column.name}"`) + } + if (value === null || typeof value !== 'object') return clause + const cell = sql`${sql.raw(`${tableName}.data`)}->${field}::text` + return sql`(${clause} AND ${cell} = ${JSON.stringify(value)}::jsonb)` +} + /** Builds JSONB containment clause: `data @> '{"field": value}'::jsonb` (uses GIN index) */ function buildContainmentClause(tableName: string, field: string, value: JsonValue): SQL { const jsonObj = JSON.stringify({ [field]: value }) diff --git a/apps/sim/lib/table/tx.ts b/apps/sim/lib/table/tx.ts index 43e5f57bf2d..5b5b05322f9 100644 --- a/apps/sim/lib/table/tx.ts +++ b/apps/sim/lib/table/tx.ts @@ -5,7 +5,7 @@ * directly from `@/lib/table/tx`. */ -import { sql } from 'drizzle-orm' +import { type SQL, sql } from 'drizzle-orm' import type { DbTransaction } from '@/lib/table/planner' const TIMEOUT_CAP_MS = 10 * 60_000 @@ -28,19 +28,29 @@ const TIMEOUT_CAP_MS = 10 * 60_000 * Safe under pgBouncer transaction pooling — the settings are transaction-scoped and cleared at * COMMIT/ROLLBACK before the session returns to the pool. */ -export async function setTableTxTimeouts( - trx: DbTransaction, - opts?: { statementMs?: number; lockMs?: number; idleMs?: number } -) { +export async function setTableTxTimeouts(trx: DbTransaction, opts?: TableTxTimeouts) { + await trx.execute(sql`select ${tableTxTimeoutSettings(opts)}`) +} + +/** Per-transaction timeouts, in milliseconds. See {@link setTableTxTimeouts}. */ +export interface TableTxTimeouts { + statementMs?: number + lockMs?: number + idleMs?: number +} + +/** + * The `set_config` calls {@link setTableTxTimeouts} runs, for a statement that applies them + * alongside its own work. + */ +export function tableTxTimeoutSettings(opts?: TableTxTimeouts): SQL { const s = opts?.statementMs ?? 10_000 const l = opts?.lockMs ?? 3_000 const i = opts?.idleMs ?? 5_000 - await trx.execute(sql` - select + return sql` set_config('statement_timeout', ${`${s}ms`}, true), set_config('lock_timeout', ${`${l}ms`}, true), - set_config('idle_in_transaction_session_timeout', ${`${i}ms`}, true) - `) + set_config('idle_in_transaction_session_timeout', ${`${i}ms`}, true)` } /** diff --git a/apps/sim/lib/table/update-row.test.ts b/apps/sim/lib/table/update-row.test.ts index 6df3dde836f..27ca42f3bbc 100644 --- a/apps/sim/lib/table/update-row.test.ts +++ b/apps/sim/lib/table/update-row.test.ts @@ -1,6 +1,7 @@ import { tableRowExecutions, userTableRows } from '@sim/db/schema' import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing' import { tableBillingMock, tableBillingMockFns } from '@sim/testing/mocks/table-billing.mock' +import { tableRowsLiveSchemaMock } from '@sim/testing/mocks/table-rows-live-schema.mock' import { beforeEach, describe, expect, it, vi } from 'vitest' import { batchUpdateRows, @@ -20,18 +21,24 @@ tableBillingMockFns.mockWouldExceedRowLimit.mockReturnValue(false) // suites can use large synthetic row counts without tripping the plan limit. vi.mock('@/lib/table/billing', () => tableBillingMock) -vi.mock('@/lib/table/validation', async (importOriginal) => ({ - uniqueValueKey: (await importOriginal()).uniqueValueKey, - validateRowSize: vi.fn(() => ({ valid: true, errors: [] })), - validateRowAgainstSchema: vi.fn(() => ({ valid: true, errors: [] })), - coerceRowToSchema: vi.fn(() => ({ valid: true, errors: [] })), - coerceRowValues: vi.fn(), - validateTableName: vi.fn(() => ({ valid: true, errors: [] })), - validateTableSchema: vi.fn(() => ({ valid: true, errors: [] })), - getUniqueColumns: vi.fn(() => []), - checkUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })), - checkBatchUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })), -})) +vi.mock('@/lib/table/rows/live-schema', () => tableRowsLiveSchemaMock) + +vi.mock('@/lib/table/validation', async (importOriginal) => { + const actual = await importOriginal() + return { + cellOf: actual.cellOf, + uniqueValueKey: actual.uniqueValueKey, + validateRowSize: vi.fn(() => ({ valid: true, errors: [] })), + validateRowAgainstSchema: vi.fn(() => ({ valid: true, errors: [] })), + coerceRowToSchema: vi.fn(() => ({ valid: true, errors: [] })), + coerceRowValues: vi.fn(), + validateTableName: vi.fn(() => ({ valid: true, errors: [] })), + validateTableSchema: vi.fn(() => ({ valid: true, errors: [] })), + getUniqueColumns: vi.fn(() => []), + checkUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })), + checkBatchUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })), + } +}) /** * Inspects the queued `trx.execute(...)` calls for SQL containing `substring`. @@ -52,31 +59,6 @@ function findExecutedSqlContaining(substring: string): boolean { }) } -/** Every string reachable from a drizzle `sql` fragment — its literal chunks AND its bound values. */ -function collectStrings(node: unknown, out: string[] = []): string[] { - if (typeof node === 'string') out.push(node) - else if (Array.isArray(node)) for (const entry of node) collectStrings(entry, out) - else if (node && typeof node === 'object') - for (const entry of Object.values(node as Record)) collectStrings(entry, out) - return out -} - -/** - * Whether one `set_config(, '', true)` guard was executed. - * - * The guards bind their values as parameters, so the setting name lives in the statement's - * literal chunks while the duration lives in its bound values — asserting on a rendered string - * would only re-check the placeholder. - */ -function executedTxTimeout(setting: string, value: string): boolean { - return dbChainMockFns.execute.mock.calls.some(([arg]) => { - const strings = collectStrings(arg) - return ( - strings.some((entry) => entry.includes(`set_config('${setting}'`)) && strings.includes(value) - ) - }) -} - /** * The `data` payload of the last `.set(...)` row write. `updateRow` always writes a JSONB merge * (`data = data || {changed}::jsonb`), so this is a `sql` fragment exposing `{ strings, values }`. @@ -342,29 +324,6 @@ describe('mutation paths — SET LOCAL timeouts', () => { dbChainMockFns.execute.mockResolvedValue([{ count: 0 }]) }) - it('replaceTableRows scales statement_timeout with (existing + new) row count', async () => { - const bigTable: TableDefinition = { ...TABLE, rowCount: 100_000, maxRows: 1_000_000 } - const payload = Array.from({ length: 50_000 }, (_, i) => ({ name: `row-${i}` })) - - await replaceTableRows( - { tableId: 'tbl-1', workspaceId: 'ws-1', rows: payload }, - bigTable, - 'req-1' - ) - - // (100_000 + 50_000) × 3ms/row = 450_000ms; above 120_000 floor, below 600_000 cap - expect(executedTxTimeout('statement_timeout', '450000ms')).toBe(true) - }) - - it('replaceTableRows caps scaled timeout at 10 minutes for very large tables', async () => { - const hugeTable: TableDefinition = { ...TABLE, rowCount: 10_000_000, maxRows: 20_000_000 } - - await replaceTableRows({ tableId: 'tbl-1', workspaceId: 'ws-1', rows: [] }, hugeTable, 'req-1') - - // 10M × 3ms = 30M ms, capped at 600_000ms (10 min) - expect(executedTxTimeout('statement_timeout', '600000ms')).toBe(true) - }) - it('replaceTableRows acquires the per-table advisory lock to serialize concurrent replaces', async () => { await replaceTableRows( { tableId: 'tbl-1', workspaceId: 'ws-1', rows: [{ name: 'a' }] }, diff --git a/apps/sim/lib/table/update-runner.test.ts b/apps/sim/lib/table/update-runner.test.ts index 3d5930f68ba..0ff5b940b1b 100644 --- a/apps/sim/lib/table/update-runner.test.ts +++ b/apps/sim/lib/table/update-runner.test.ts @@ -31,11 +31,16 @@ vi.mock('@/lib/table/rows/ordering', () => ({ })) vi.mock('@/lib/table/events', () => tableEventsMock) vi.mock('@/lib/table/sql', () => ({ buildFilterClause: mockBuildFilterClause })) -vi.mock('@/lib/table/validation', () => ({ - validateRowSize: mockValidateRowSize, - coerceRowToSchema: mockCoerceRowToSchema, - coerceRowValues: mockCoerceRowValues, -})) +vi.mock('@/lib/table/validation', async (importOriginal) => { + const actual = await importOriginal() + return { + cellOf: actual.cellOf, + uniqueColumnsInPatch: actual.uniqueColumnsInPatch, + validateRowSize: mockValidateRowSize, + coerceRowToSchema: mockCoerceRowToSchema, + coerceRowValues: mockCoerceRowValues, + } +}) vi.mock('@/lib/table/constants', () => ({ ...tableConstantsMock, TABLE_LIMITS: { ...tableConstantsMock.TABLE_LIMITS, DELETE_PAGE_SIZE: 2, UPDATE_BATCH_SIZE: 100 }, @@ -57,7 +62,12 @@ const UNLOCKED = { updateLocked: false, deleteLocked: false, } -const table = { id: 'tbl_1', workspaceId: 'ws_1', schema: { columns: [] }, locks: UNLOCKED } +const table = { + id: 'tbl_1', + workspaceId: 'ws_1', + schema: { columns: [{ id: 'flag', name: 'flag', type: 'boolean' }] }, + locks: UNLOCKED, +} const cutoff = new Date('2026-06-05T00:00:00Z') function basePayload(overrides = {}) { diff --git a/apps/sim/lib/table/update-runner.ts b/apps/sim/lib/table/update-runner.ts index b3ed3380002..e1a75e84b15 100644 --- a/apps/sim/lib/table/update-runner.ts +++ b/apps/sim/lib/table/update-runner.ts @@ -1,8 +1,12 @@ +import { userTableRows } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { getErrorMessage, toError } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' import { truncate } from '@sim/utils/string' -import type { Filter, RowData, TableDefinition } from '@/lib/table' +import { and, eq, inArray } from 'drizzle-orm' +import { OrchestrationError } from '@/lib/core/orchestration/types' +import type { Filter, RowData, TableDefinition, TableSchema } from '@/lib/table' +import { getColumnId } from '@/lib/table/column-keys' import { TABLE_LIMITS, USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants' import { appendTableEvent } from '@/lib/table/events' import { @@ -19,11 +23,13 @@ import { TableLockedError, } from '@/lib/table/mutation-locks' import type { DbTransaction } from '@/lib/table/planner' +import { withLiveSchema } from '@/lib/table/rows/live-schema' import { selectRowDataPage, updatePageByIds } from '@/lib/table/rows/ordering' import { createExactEmptyTableRowSecretProvenance } from '@/lib/table/rows/secret-provenance' +import { deriveBulkUpdatePatch } from '@/lib/table/rows/service' import { getTableById } from '@/lib/table/service' import { buildFilterClause } from '@/lib/table/sql' -import { coerceRowToSchema, coerceRowValues, validateRowSize } from '@/lib/table/validation' +import { coerceRowToSchema, uniqueColumnsInPatch, validateRowSize } from '@/lib/table/validation' const logger = createLogger('TableUpdateRunner') @@ -36,6 +42,88 @@ const PROGRESS_INTERVAL_ROWS = 5000 */ class JobSupersededError extends Error {} +/** + * The patch cannot be applied to the table as it now stands. Retrying cannot help: the schema that + * refused it is the one a retry would read. + */ +export class UpdatePatchRejectedError extends Error { + constructor(message: string) { + super(message) + this.name = 'UpdatePatchRejectedError' + } +} + +/** + * The patch one batch writes against `live`, derived from the job's raw `patch` exactly as the + * inline bulk update derives its own ({@link deriveBulkUpdatePatch}): cells of columns deleted since + * `previous` dropped, the rest coerced and validated. It is refused, with an + * {@link UpdatePatchRejectedError}, when it cannot be written to every matched row: a value the + * columns refuse, a null in a required column, or any unique column, since one value in many rows + * cannot stay unique. Batches derive it under the schema lock, so a column made required or unique + * while the job waited is refused before the batch writes. + */ +function deriveJobPatch(patch: RowData, previous: TableSchema, live: TableDefinition): RowData { + let derived: RowData + try { + derived = deriveBulkUpdatePatch(patch, previous, live, undefined) + } catch (err) { + if (err instanceof OrchestrationError && err.code === 'validation') { + throw new UpdatePatchRejectedError(err.message) + } + throw err + } + const cleared = live.schema.columns.filter((column) => { + const columnId = getColumnId(column) + return ( + column.required && Object.hasOwn(derived, columnId) && (derived[columnId] ?? null) === null + ) + }) + if (cleared.length > 0) { + throw new UpdatePatchRejectedError( + `Missing required field: ${cleared.map((column) => column.name).join(', ')}` + ) + } + const unique = uniqueColumnsInPatch(live.schema, derived) + if (unique.length > 0) { + throw new UpdatePatchRejectedError( + `Cannot set unique column values when updating multiple rows: ${unique.map((column) => column.name).join(', ')}` + ) + } + return derived +} + +/** Refuses a row that `patch`, merged over it, would leave oversized or invalid under `schema`. */ +function assertMergedRowFits( + schema: TableSchema, + row: { id: string; data: RowData }, + patch: RowData +): void { + const merged = { ...row.data, ...patch } + const sizeValidation = validateRowSize(merged) + if (!sizeValidation.valid) { + throw new UpdatePatchRejectedError(`Row ${row.id}: ${sizeValidation.errors.join(', ')}`) + } + const schemaValidation = coerceRowToSchema(merged, schema) + if (!schemaValidation.valid) { + throw new UpdatePatchRejectedError(`Row ${row.id}: ${schemaValidation.errors.join(', ')}`) + } +} + +/** Reads a batch's rows in its transaction and refuses any `patch` would not fit under `live`. */ +async function validateMergedRows( + trx: DbTransaction, + live: TableDefinition, + rowIds: string[], + patch: RowData +): Promise { + const rows = await trx + .select({ id: userTableRows.id, data: userTableRows.data }) + .from(userTableRows) + .where(and(eq(userTableRows.tableId, live.id), inArray(userTableRows.id, rowIds))) + for (const row of rows) + assertMergedRowFits(live.schema, { id: row.id, data: row.data as RowData }, patch) +} + export interface TableUpdatePayload { jobId: string tableId: string @@ -65,7 +153,8 @@ export interface TableUpdatePayload { * the affected columns explicitly afterward if downstream recompute is needed. * * Unexpected errors are rethrown for the caller's retry machinery; the caller marks the job - * failed via `markTableUpdateFailed`. A superseded run returns quietly. + * failed via `markTableUpdateFailed`. An {@link UpdatePatchRejectedError} is rethrown too, and is + * not worth a retry. A superseded run returns quietly. */ export async function runTableUpdate(payload: TableUpdatePayload): Promise { const { jobId, tableId, workspaceId, filter, data, cutoff, maxRows } = payload @@ -109,7 +198,8 @@ export async function runTableUpdate(payload: TableUpdatePayload): Promise if ((await stopIfLocked(table, 0)) === null) return // Runs inside each batch's transaction, under the same advisory lock the - // lock toggle holds, so no page can be written after a lock commits. + // lock toggle and schema changes hold, so no page can be written after a + // lock commits, nor under a schema that would not store the patch. const revalidate = async (trx: DbTransaction) => { const fresh = await getTableById(tableId, { tx: trx, includeArchived: true }) if (fresh) assertRowUpdate(fresh, patchColumnIds(data)) @@ -119,10 +209,10 @@ export async function runTableUpdate(payload: TableUpdatePayload): Promise const filterClause = buildFilterClause(filter, USER_TABLE_ROWS_SQL_NAME, table.schema.columns) if (!filterClause) throw new Error('Filter is required for bulk update') - // Coerce the patch once to the schema's types — the merged validation below and the persisted - // JSONB merge both use this normalized copy. - coerceRowValues(data, table.schema) - const patchJson = JSON.stringify(data) + // `data` stays the raw payload. Each page, and each batch under the schema lock, derives the + // patch it writes from it against the schema of that moment, as the inline update would. Derive + // it once up front too, so a patch the table refuses fails before any page is read. + deriveJobPatch(data, table.schema, table) // Resume the persisted count: a retried attempt's earlier pages are already committed, so // starting at zero would overwrite cumulative progress. Doubles as the initial ownership gate. @@ -145,6 +235,10 @@ export async function runTableUpdate(payload: TableUpdatePayload): Promise if (!current) throw new JobSupersededError() const pageProof = await stopIfLocked(current, processed) if (pageProof === null) return + const pagePatch = deriveJobPatch(data, table.schema, current) + // Every column the update wrote has since been deleted: nothing is left to write. + if (Object.keys(pagePatch).length === 0) break + const patchJson = JSON.stringify(pagePatch) const page = await selectRowDataPage({ tableId, @@ -163,25 +257,25 @@ export async function runTableUpdate(payload: TableUpdatePayload): Promise // Validate each merged result before writing the page — a row that would overflow the size // cap or violate the schema fails the job (earlier pages stay applied; best-effort). - for (const row of page) { - const merged = { ...row.data, ...data } - const sizeValidation = validateRowSize(merged) - if (!sizeValidation.valid) { - throw new Error(`Row ${row.id}: ${sizeValidation.errors.join(', ')}`) - } - const schemaValidation = coerceRowToSchema(merged, table.schema) - if (!schemaValidation.valid) { - throw new Error(`Row ${row.id}: ${schemaValidation.errors.join(', ')}`) - } - } + for (const row of page) assertMergedRowFits(current.schema, row, pagePatch) try { processed += await updatePageByIds( tableId, workspaceId, page.map((r) => r.id), - patchJson, - createExactEmptyTableRowSecretProvenance(data), + async (trx, fresh, batch) => { + const live = fresh ? withLiveSchema(current, fresh.schema) : current + const batchPatch = deriveJobPatch(data, table.schema, live) + if (Object.keys(batchPatch).length === 0) return null + // The page was checked against `current`; a batch under a schema that moved since is + // checked again, against the rows as they stand, before it writes. + if (live !== current) await validateMergedRows(trx, live, batch, batchPatch) + return { + patchJson: JSON.stringify(batchPatch), + secretProvenance: createExactEmptyTableRowSecretProvenance(batchPatch), + } + }, pageProof, revalidate ) diff --git a/apps/sim/lib/table/validation.ts b/apps/sim/lib/table/validation.ts index cbe13d617f8..29bd9157078 100644 --- a/apps/sim/lib/table/validation.ts +++ b/apps/sim/lib/table/validation.ts @@ -27,7 +27,7 @@ import { } from '@/lib/table/constants' import { withSeqscanOff } from '@/lib/table/planner' import { resolveSelectOptionId, splitMultiSelectInput } from '@/lib/table/select-options' -import { fieldPredicate } from '@/lib/table/sql' +import { uniqueValuePredicate } from '@/lib/table/sql' import type { ColumnDefinition, JsonValue, @@ -252,12 +252,20 @@ export function validateTableSchema(schema: TableSchema): ValidationResult { return { valid: errors.length === 0, errors } } +/** + * The cell `row` holds for `columnId`, read as an own property: a legacy column keyed by its name + * can be called `constructor`, which a plain index would find on every row's prototype. + */ +export function cellOf(row: RowData, columnId: string): JsonValue | undefined { + return Object.hasOwn(row, columnId) ? row[columnId] : undefined +} + /** Validates row data matches schema column types and required fields. */ export function validateRowAgainstSchema(data: RowData, schema: TableSchema): ValidationResult { const errors: string[] = [] for (const column of schema.columns) { - const value = data[getColumnId(column)] + const value = cellOf(data, getColumnId(column)) if (column.required && (value === undefined || value === null)) { errors.push(`Missing required field: ${column.name}`) @@ -348,7 +356,7 @@ export function coerceRowValues( const policyFor = policyResolver(policy, patchedKeys) for (const column of schema.columns) { const key = getColumnId(column) - const value = data[key] + const value = cellOf(data, key) if (value === null || value === undefined) continue const coerced = coerceValueToColumnType(value, column) @@ -407,6 +415,14 @@ export function getUniqueColumns(schema: TableSchema): ColumnDefinition[] { return schema.columns.filter((col) => col.unique === true) } +/** + * The unique columns `patch` writes. A patch applied to more than one row cannot write any of them + * without storing a duplicate. + */ +export function uniqueColumnsInPatch(schema: TableSchema, patch: RowData): ColumnDefinition[] { + return getUniqueColumns(schema).filter((column) => Object.hasOwn(patch, getColumnId(column))) +} + /** * The key two unique-column values share when the unique check treats them as equal. Object keys * are sorted, since the check compares JSONB, where key order carries no meaning. In-batch @@ -428,13 +444,13 @@ export function validateUniqueConstraints( for (const column of uniqueColumns) { const key = getColumnId(column) - const value = data[key] + const value = cellOf(data, key) if (value === null || value === undefined) continue const duplicate = existingRows.find((row) => { if (excludeRowId && row.id === excludeRowId) return false // Case-sensitive, matching the DB unique-check leaf (`fieldPredicate` eq). - const existing = row.data[key] + const existing = cellOf(row.data, key) return ( existing !== undefined && columnValueForEquality(value, column) === columnValueForEquality(existing, column) @@ -481,13 +497,14 @@ export async function checkUniqueConstraintsDb( for (const column of uniqueColumns) { const key = getColumnId(column) - const value = data[key] + const value = cellOf(data, key) if (value === null || value === undefined) continue - // Same leaf as the upsert conflict probe → case-sensitive JSONB containment - // (GIN-indexed). `eq` always yields a clause for a non-null value. - const clause = fieldPredicate(USER_TABLE_ROWS_SQL_NAME, key, 'eq', value, column) - if (clause) conditions.push({ column, value, sql: clause }) + conditions.push({ + column, + value, + sql: uniqueValuePredicate(USER_TABLE_ROWS_SQL_NAME, key, value, column), + }) } if (conditions.length === 0) { @@ -495,7 +512,7 @@ export async function checkUniqueConstraintsDb( } // Query for each unique column separately to provide specific error messages. - // The predicate is now case-sensitive JSONB containment (`data @> {...}`), + // The predicate leads with case-sensitive JSONB containment (`data @> {...}`), // which can use the GIN index. We still pin `enable_seqscan = off` (tenant- // bounded) defensively for the small-table / cold-stats case. With an external // transaction the flag is set on it for the check only (see withSeqscanOffOn) @@ -600,7 +617,7 @@ export async function checkBatchUniqueConstraintsDb( for (const column of uniqueColumns) { const key = getColumnId(column) - const value = rowData[key] + const value = cellOf(rowData, key) if (value === null || value === undefined) continue const normalizedValue = uniqueValueKey(value, column) @@ -645,18 +662,7 @@ export async function checkBatchUniqueConstraintsDb( // string — a unique `date` column normalized to a bare `2024-01-01` // and then threw `SyntaxError` trying to parse it back. const originalValue: JsonValue = JSON.parse(normalizedValue) - // Same case-sensitive containment leaf as every other matcher. - const clause = fieldPredicate( - USER_TABLE_ROWS_SQL_NAME, - columnId, - 'eq', - originalValue, - column - ) - if (!clause) { - throw new Error(`Failed to build unique-constraint predicate for column "${column.name}"`) - } - return clause + return uniqueValuePredicate(USER_TABLE_ROWS_SQL_NAME, columnId, originalValue, column) }) const conflictingRows = await ex @@ -672,22 +678,17 @@ export async function checkBatchUniqueConstraintsDb( // Map conflicts back to batch rows for (const conflict of conflictingRows) { const conflictData = conflict.data as RowData - const conflictValue = columnValueForEquality(conflictData[columnId], column) - const normalizedConflictValue = - typeof conflictValue === 'string' ? conflictValue : JSON.stringify(conflictValue) + const conflictValue = cellOf(conflictData, columnId) + if (conflictValue === null || conflictValue === undefined) continue + // Keyed like the batch, since stored jsonb comes back with its keys reordered. + const normalizedConflictValue = uniqueValueKey(conflictValue, column) // Find which batch rows have this conflicting value for (let i = 0; i < rows.length; i++) { - const rowValue = rows[i][columnId] + const rowValue = cellOf(rows[i], columnId) if (rowValue === null || rowValue === undefined) continue - const comparableRowValue = columnValueForEquality(rowValue, column) - const normalizedRowValue = - typeof comparableRowValue === 'string' - ? comparableRowValue - : JSON.stringify(comparableRowValue) - - if (normalizedRowValue === normalizedConflictValue) { + if (uniqueValueKey(rowValue, column) === normalizedConflictValue) { // Check if this row already has errors for this column let rowError = rowErrors.find((e) => e.row === i) if (!rowError) { diff --git a/packages/db/script-migrations-paused-billing-attribution.test.ts b/packages/db/script-migrations-paused-billing-attribution.test.ts index f7901d6c65a..131364665ac 100644 --- a/packages/db/script-migrations-paused-billing-attribution.test.ts +++ b/packages/db/script-migrations-paused-billing-attribution.test.ts @@ -378,6 +378,7 @@ describe('script migration registry', () => { '0023_projection_acl_skip_unfilled', '0024_knowledge_projection_async', '0025_scope_keyword_projections', + '0026_user_table_schema_for_write', ]) }) }) diff --git a/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts b/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts index d9972d14d3d..fc9cc6f741c 100644 --- a/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts +++ b/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts @@ -396,6 +396,7 @@ describe('search projection upgrade in PostgreSQL', () => { { name: '0023_projection_acl_skip_unfilled' }, { name: '0024_knowledge_projection_async' }, { name: '0025_scope_keyword_projections' }, + { name: '0026_user_table_schema_for_write' }, ]) const [{ complete }] = await sql`SELECT count(*)::int AS complete FROM embedding e JOIN embedding_search s ON s.id = e.id JOIN embedding_keyword_search k ON k.id = e.id diff --git a/packages/db/script-migrations/0026_user_table_schema_for_write.ts b/packages/db/script-migrations/0026_user_table_schema_for_write.ts new file mode 100644 index 00000000000..d651e9da730 --- /dev/null +++ b/packages/db/script-migrations/0026_user_table_schema_for_write.ts @@ -0,0 +1,66 @@ +import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url' +import type { ScriptMigration } from '@sim/db/script-migrations/types' +import postgres, { type Sql } from 'postgres' + +/** + * Installs `user_table_schema_for_write(table_id)`, which takes a table's schema lock shared and + * returns the schema a row write must validate against. Every row write in `apps/sim` calls it, so + * it is a script migration: `db:push` builds a database from `schema.ts` and the script migrations + * it lists, never from the SQL migrations. + * + * Schema changes hold the advisory lock `user_table_schema:` exclusively while they + * check stored rows against the new schema (`withLockedTable`), so a row write holds it shared + * for its whole transaction: a change waits for writes in flight, and a write that waited must + * validate against the change it waited for. The key and hash match `acquireAdvisoryXactLock`. + * + * The function must stay VOLATILE. Under READ COMMITTED a volatile function takes a fresh + * snapshot for each query it runs, so the read below sees a schema change that committed while + * the lock waited, even though the calling statement's snapshot predates it. A STABLE function + * would read through the caller's snapshot and return the schema from before the wait. Folding + * the lock and the read into one call lets a writer take both inside a statement it already runs. + * + * The schema lock waits under the caller's `statement_timeout`, not its `lock_timeout`. A schema + * change (a retype, a constraint scan, an import) can hold the lock far longer than the short + * `lock_timeout` a writer sets for its row locks, and failing an edit there turns an ordinary wait + * into an error; a writer that set no timeouts keeps waiting as it did before the guard existed. + * The value is read here rather than left to `statement_timeout` itself because a writer applies + * its timeouts in the same statement that calls this function, and a statement's own timeout is + * armed before it runs. The caller's `lock_timeout` is restored before the function returns, so the + * rest of the transaction's locks keep their bound. + * + * Returns NULL when the table does not exist. `CREATE OR REPLACE` makes it safe to run again. + */ +export async function installUserTableSchemaForWrite(sql: Sql): Promise { + await sql.unsafe(`CREATE OR REPLACE FUNCTION "public"."user_table_schema_for_write"(p_table_id text) +RETURNS jsonb +LANGUAGE plpgsql +VOLATILE +SET search_path = pg_catalog, public +AS $$ +DECLARE + caller_lock_timeout text := current_setting('lock_timeout'); +BEGIN + PERFORM set_config('lock_timeout', current_setting('statement_timeout'), true); + PERFORM pg_advisory_xact_lock_shared(hashtextextended('user_table_schema:' || p_table_id, 0)); + PERFORM set_config('lock_timeout', caller_lock_timeout, true); + RETURN (SELECT "schema" FROM "public"."user_table_definitions" WHERE "id" = p_table_id); +END; +$$`) +} + +export const userTableSchemaForWriteMigration: ScriptMigration = { + name: '0026_user_table_schema_for_write', + up: installUserTableSchemaForWrite, +} + +/** Run directly by `db:push`. */ +if (import.meta.main) { + const url = resolveMigrationDatabaseUrl() + if (!url) throw new Error('DATABASE_URL is required to install user_table_schema_for_write') + const sql = postgres(url, { max: 1, max_lifetime: null, onnotice: () => undefined }) + try { + await installUserTableSchemaForWrite(sql) + } finally { + await sql.end() + } +} diff --git a/packages/db/script-migrations/index.ts b/packages/db/script-migrations/index.ts index b469e505669..a98a4cfcd6e 100644 --- a/packages/db/script-migrations/index.ts +++ b/packages/db/script-migrations/index.ts @@ -9,6 +9,7 @@ import { projectionSourceAclBackfillMigration } from '@sim/db/script-migrations/ import { projectionAclSkipUnfilledMigration } from '@sim/db/script-migrations/0023_projection_acl_skip_unfilled' import { knowledgeProjectionAsyncMigration } from '@sim/db/script-migrations/0024_knowledge_projection_async' import { scopeKeywordProjectionsMigration } from '@sim/db/script-migrations/0025_scope_keyword_projections' +import { userTableSchemaForWriteMigration } from '@sim/db/script-migrations/0026_user_table_schema_for_write' import type { Sql } from 'postgres' import { backfillTableOrderKeys } from './0001_backfill_table_order_keys' import { backfillPausedBillingAttribution } from './0002_backfill_paused_billing_attribution' @@ -55,6 +56,8 @@ export const scriptMigrations: readonly ScriptMigration[] = [ knowledgeProjectionAsyncMigration, /** 0025 keeps the keyword projections for search indexes only. */ scopeKeywordProjectionsMigration, + /** 0026 installs the schema guard every table row write takes before it validates. */ + userTableSchemaForWriteMigration, ] /** diff --git a/packages/db/scripts/push.ts b/packages/db/scripts/push.ts index 90b2288954c..2d10c1c4e0e 100644 --- a/packages/db/scripts/push.ts +++ b/packages/db/scripts/push.ts @@ -9,6 +9,7 @@ const RECONCILIATION_COMMANDS = [ ['bun', '--env-file=.env', 'run', './script-migrations/0021_embedding_search_connector.ts'], ['bun', '--env-file=.env', 'run', './script-migrations/0024_knowledge_projection_async.ts'], ['bun', '--env-file=.env', 'run', './script-migrations/0025_scope_keyword_projections.ts'], + ['bun', '--env-file=.env', 'run', './script-migrations/0026_user_table_schema_for_write.ts'], ] /** diff --git a/packages/testing/src/mocks/index.ts b/packages/testing/src/mocks/index.ts index c118f2c6922..1f408765f71 100644 --- a/packages/testing/src/mocks/index.ts +++ b/packages/testing/src/mocks/index.ts @@ -807,6 +807,10 @@ export { tableRouteUtilsMock, tableRouteUtilsMockFns, } from './table-route-utils.mock' +export { + tableRowsLiveSchemaMock, + tableRowsLiveSchemaMockFns, +} from './table-rows-live-schema.mock' export { MockTableRowProvenanceReader, tableRowsSecretProvenanceMock, diff --git a/packages/testing/src/mocks/table-rows-live-schema.mock.ts b/packages/testing/src/mocks/table-rows-live-schema.mock.ts new file mode 100644 index 00000000000..68cfb34c9df --- /dev/null +++ b/packages/testing/src/mocks/table-rows-live-schema.mock.ts @@ -0,0 +1,37 @@ +import { vi } from 'vitest' + +/** + * Controllable mock functions for `@/lib/table/rows/live-schema`. + * + * Defaults: `lockLiveTableSchema` and `withLiveSchema` return the definition they are handed, as + * when its schema is still current, so no refit runs; `dropDeletedColumns` is a no-op and + * `refitRowToSchema` reports a valid row. + * + * @example + * ```ts + * import { tableRowsLiveSchemaMockFns } from '@sim/testing/mocks/table-rows-live-schema.mock' + * + * tableRowsLiveSchemaMockFns.mockLockLiveTableSchema.mockResolvedValueOnce(liveTable) + * ``` + */ +export const tableRowsLiveSchemaMockFns = { + mockLockLiveTableSchema: vi.fn(async (_trx: unknown, table: T): Promise => table), + mockWithLiveSchema: vi.fn((table: T): T => table), + mockDropDeletedColumns: vi.fn((..._args: unknown[]): void => {}), + mockRefitRowToSchema: vi.fn((..._args: unknown[]) => ({ valid: true, errors: [] as string[] })), +} + +/** + * Static mock module for `@/lib/table/rows/live-schema`. + * + * @example + * ```ts + * vi.mock('@/lib/table/rows/live-schema', () => tableRowsLiveSchemaMock) + * ``` + */ +export const tableRowsLiveSchemaMock = { + lockLiveTableSchema: tableRowsLiveSchemaMockFns.mockLockLiveTableSchema, + withLiveSchema: tableRowsLiveSchemaMockFns.mockWithLiveSchema, + dropDeletedColumns: tableRowsLiveSchemaMockFns.mockDropDeletedColumns, + refitRowToSchema: tableRowsLiveSchemaMockFns.mockRefitRowToSchema, +}