Skip to content
27 changes: 27 additions & 0 deletions apps/sim/lib/db/advisory-locks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,3 +44,30 @@ export async function tryAcquireAdvisoryXactLock(
)
return Boolean(lock?.acquired)
}

/** One lock of an {@link acquireAdvisoryXactLocks} set. */
export interface AdvisoryXactLockRequest {
key: string
/** Takes the lock in shared mode, which conflicts only with exclusive holders of `key`. */
shared: boolean
}

/**
* Blocks until every transaction-scoped advisory lock in `locks` is held, in one round trip.
* `jsonb_array_elements` emits elements in array order and each row's lock call runs before the
* next row is read, so the locks are taken exactly in the order given: callers that must not
* deadlock pass their locks in one global order. The locks release on commit or rollback.
*/
export async function acquireAdvisoryXactLocks(
tx: DbTransaction,
tag: string,
locks: readonly AdvisoryXactLockRequest[]
): Promise<void> {
if (locks.length === 0) return
await tx.execute(sql`
SELECT CASE WHEN (lock ->> 'shared')::boolean
THEN pg_advisory_xact_lock_shared(hashtextextended(lock ->> 'key', 0))
ELSE pg_advisory_xact_lock(hashtextextended(lock ->> 'key', 0))
END
FROM jsonb_array_elements(${JSON.stringify(locks)}::jsonb) AS lock ${lockTag(tag)}`)
}
54 changes: 31 additions & 23 deletions apps/sim/lib/table/import-data.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import {
mutateTableRowsWithSecretProvenance,
} from '@/lib/table/rows/secret-provenance'
import { batchInsertRowsWithTx, replaceTableRowsWithTx } from '@/lib/table/rows/service'
import { lockUniqueColumns } from '@/lib/table/rows/unique-locks'
import { addTableColumnsWithTx, auditTableColumnsAdded, getTableById } from '@/lib/table/service'
import type {
ReplaceRowsResult,
Expand Down Expand Up @@ -58,10 +59,12 @@ export interface BulkImportBatch {
* Inserts one batch of rows for an async import in a single committed statement.
*
* Differs from {@link batchInsertRowsWithTx} for the bulk-load case: caller-supplied
* contiguous order keys (no `acquireRowOrderLock` scan — an
* import owns its hidden table as the sole writer), no `RETURNING`, and **no
* `fireTableTrigger` / `runWorkflowColumn`** (a 1M-row import must not dispatch a
* workflow run per row). `row_count` is maintained set-based by the statement-level
* contiguous order keys (no `acquireRowOrderLock` scan; the caller threads each batch's
* anchor from the previous one), no `RETURNING`, and **no `fireTableTrigger` /
* `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.
*
Expand Down Expand Up @@ -99,25 +102,9 @@ export async function bulkInsertImportBatch(
}
}

const uniqueColumns = getUniqueColumns(table.schema)
if (uniqueColumns.length > 0) {
const uniqueResult = await checkBatchUniqueConstraintsDb(
data.tableId,
data.rows,
table.schema,
db
)
if (!uniqueResult.valid) {
throw new OrchestrationError(
'validation',
uniqueResult.errors.map((e) => `Row ${e.row + 1}: ${e.errors.join(', ')}`).join('; ')
)
}
}

const now = new Date()
// Import worker is the table's sole writer; append keys after the anchor the caller threads
// from the previous batch's last key — no per-batch max(order_key) scan over a growing table.
// Append keys after the anchor the caller threads from the previous batch's last key — no
// per-batch max(order_key) scan over a growing table.
const orderKeys = nKeysBetween(data.afterOrderKey ?? null, null, data.rows.length)
const rowsToInsert = data.rows.map((rowData, i) => ({
id: `row_${generateId().replace(/-/g, '')}`,
Expand All @@ -133,6 +120,22 @@ export async function bulkInsertImportBatch(

const inserted = await db.transaction(async (trx) => {
await guardBatch(trx, data.tableId, revalidate)
if (getUniqueColumns(table.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)
const uniqueResult = await checkBatchUniqueConstraintsDb(
data.tableId,
data.rows,
table.schema,
trx
)
if (!uniqueResult.valid) {
throw new OrchestrationError(
'validation',
uniqueResult.errors.map((e) => `Row ${e.row + 1}: ${e.errors.join(', ')}`).join('; ')
)
}
}
return mutateTableRowsWithSecretProvenance(trx, {
rows: rowsToInsert.map((row) => ({
rowId: row.id,
Expand Down Expand Up @@ -282,6 +285,9 @@ export async function importAppendRows(
})
const result = await db.transaction(async (trx) => {
let working = await refreshUnderLock(trx, table)
// Lock the unique columns whole, ahead of the row-order lock: per-value locks for every row of
// an import would flood the server's lock table.
await lockUniqueColumns(trx, working)
if (additions.length > 0) {
// Take the row-order lock before creating columns so this path uses the
// same rows_pos → user_table_definitions order as plain inserts. Creating
Expand All @@ -305,7 +311,8 @@ export async function importAppendRows(
secretProvenance: batch.map(createExactEmptyTableRowSecretProvenance),
},
working,
generateId().slice(0, 8)
generateId().slice(0, 8),
{ uniqueColumnsLocked: true }
)
inserted.push(...batchInserted)
}
Expand Down Expand Up @@ -347,6 +354,7 @@ export async function importReplaceRows(
})
const result = await db.transaction(async (trx) => {
let working = await refreshUnderLock(trx, table)
await lockUniqueColumns(trx, working)
if (additions.length > 0) {
await acquireRowOrderLock(trx, table.id)
working = await addTableColumnsWithTx(trx, working, additions, requestId)
Expand Down
6 changes: 6 additions & 0 deletions apps/sim/lib/table/rows/ordering.ts
Original file line number Diff line number Diff line change
Expand Up @@ -332,6 +332,11 @@ export async function insertOrderedRow(params: {
secretProvenance?: TableRowSecretProvenanceWrite
/** 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.
*/
assertUnique?: (trx: DbTransaction) => Promise<void>
}): Promise<{
id: string
data: RowData
Expand All @@ -354,6 +359,7 @@ export async function insertOrderedRow(params: {
} = params
const [row] = await db.transaction(async (trx) => {
await setTableTxTimeouts(trx)
await params.assertUnique?.(trx)
await acquireRowOrderLock(trx, tableId)

// Resolve the authoritative order key from neighbor ids when given, else from the requested
Expand Down
Loading
Loading