Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 10 additions & 3 deletions apps/sim/background/table-update.ts
Original file line number Diff line number Diff line change
@@ -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'

/**
Expand All @@ -19,7 +20,8 @@ export interface TableUpdateTaskPayload extends Omit<TableUpdatePayload, 'cutoff
* worker keysets by id with a `created_at <= cutoff` floor and the JSONB-merge patch is idempotent
* (re-applying the same patch to an already-patched row is a no-op), so a retried attempt re-walks
* and re-applies whatever remains. The `table_jobs` ownership gate stops a retried run that lost
* the job within one page.
* the job within one page. A patch the table's schema refuses aborts without a retry: the retry
* would read the same schema.
*/
export const tableUpdateTask = task({
id: 'table-update',
Expand All @@ -30,7 +32,12 @@ export const tableUpdateTask = task({
concurrencyLimit: 10,
},
run: async (payload: TableUpdateTaskPayload) => {
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)
Expand Down
5 changes: 2 additions & 3 deletions apps/sim/lib/table/application/copilot-bulk-rows.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')

Expand Down Expand Up @@ -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

Expand Down
9 changes: 7 additions & 2 deletions apps/sim/lib/table/bulk-update-concurrency.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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<typeof import('@/lib/table/validation')>()).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: [] })),
}))
Expand Down
28 changes: 20 additions & 8 deletions apps/sim/lib/table/import-data.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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.
Expand All @@ -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) {
Expand Down Expand Up @@ -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)
Comment thread
waleedlatif1 marked this conversation as resolved.
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) {
Expand Down Expand Up @@ -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 }
})
Expand Down
5 changes: 4 additions & 1 deletion apps/sim/lib/table/rows/bulk-update-patch-validation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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)
Expand Down
104 changes: 104 additions & 0 deletions apps/sim/lib/table/rows/live-schema.ts
Original file line number Diff line number Diff line change
@@ -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<TableDefinition> {
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
}
40 changes: 29 additions & 11 deletions apps/sim/lib/table/rows/ordering.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>
validate?: (trx: DbTransaction) => Promise<void>
}): Promise<{
id: string
data: RowData
Expand All @@ -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
Expand Down Expand Up @@ -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<PagePatch | null>,
/** 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}. */
Expand All @@ -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),
Expand Down
Loading
Loading