Skip to content

Commit fbc3db1

Browse files
committed
fix(tables): lock a table's unique values under one table key so wide schemas stay within the lock budget
1 parent 97cb317 commit fbc3db1

4 files changed

Lines changed: 121 additions & 49 deletions

File tree

‎apps/sim/lib/table/import-data.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -121,7 +121,7 @@ export async function bulkInsertImportBatch(
121121
const inserted = await db.transaction(async (trx) => {
122122
await guardBatch(trx, data.tableId, revalidate)
123123
if (getUniqueColumns(table.schema).length > 0) {
124-
// Whole-column locks, not per-value: a batch is far more values than the value-lock cap.
124+
// The whole-table unique lock, not per-value: a batch is far more values than the value-lock cap.
125125
await lockUniqueColumns(trx, table)
126126
const uniqueResult = await checkBatchUniqueConstraintsDb(
127127
data.tableId,

‎apps/sim/lib/table/rows/row-writes.integration.ts‎

Lines changed: 81 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import { tableTriggerMock } from '@sim/testing/mocks/table-trigger.mock'
1313
import { tableWorkflowColumnsMock } from '@sim/testing/mocks/table-workflow-columns.mock'
1414
import { sleep } from '@sim/utils/helpers'
1515
import { generateId } from '@sim/utils/id'
16+
import { sql } from 'drizzle-orm'
1617
import postgres from 'postgres'
1718
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
1819

@@ -21,8 +22,11 @@ vi.mock('@/lib/table/trigger', () => tableTriggerMock)
2122
vi.mock('@/lib/table/workflow-columns', () => tableWorkflowColumnsMock)
2223
vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock)
2324

24-
import { bulkInsertImportBatch } from '@/lib/table/import-data'
25+
import { bulkInsertImportBatch, importReplaceRows } from '@/lib/table/import-data'
26+
import type { DbTransaction } from '@/lib/table/planner'
27+
import { acquireRowOrderLock } from '@/lib/table/rows/ordering'
2528
import { batchInsertRows, batchUpdateRows, insertRow, updateRow } from '@/lib/table/rows/service'
29+
import { lockUniqueColumns, lockUniqueValues } from '@/lib/table/rows/unique-locks'
2630
import { getTableById } from '@/lib/table/service'
2731
import { getOrCreateTableSnapshot } from '@/lib/table/snapshot-cache'
2832
import type { ColumnDefinition, TableDefinition } from '@/lib/table/types'
@@ -271,8 +275,13 @@ describe('table row writes against real PostgreSQL', () => {
271275
table.id,
272276
[
273277
() => insertEmail(table, 'dup@example.test'),
274-
() =>
275-
updateRow(
278+
async () => {
279+
// Start the edit only once the insert has passed its check and holds its locks.
280+
expect(await waitForLockWaiters(table.id, { onOrderLock: 1, onValueLock: 0 })).toEqual({
281+
onOrderLock: 1,
282+
onValueLock: 0,
283+
})
284+
return updateRow(
276285
{
277286
tableId: table.id,
278287
rowId: `${table.id}-a`,
@@ -283,7 +292,8 @@ describe('table row writes against real PostgreSQL', () => {
283292
},
284293
table,
285294
'unique-race'
286-
),
295+
)
296+
},
287297
],
288298
{ onOrderLock: 1, onValueLock: 1 }
289299
)
@@ -292,6 +302,69 @@ describe('table row writes against real PostgreSQL', () => {
292302
expect(await storedCount(table.id, { email: 'dup@example.test' })).toBe(1)
293303
})
294304

305+
it('lets a replace that adds the first unique column wait out a writer on an older schema', async () => {
306+
const table = await createTable([{ id: 'name', name: 'name', type: 'string' }])
307+
// The writer resolved the table while it still had a unique column, so it holds the unique
308+
// lock shared and will want the row-order lock next.
309+
const stale: TableDefinition = {
310+
...table,
311+
schema: { columns: [...table.schema.columns, ...uniqueColumns] },
312+
}
313+
let holdsUniqueLock!: () => void
314+
const holding = new Promise<void>((resolve) => {
315+
holdsUniqueLock = resolve
316+
})
317+
let release!: () => void
318+
const released = new Promise<void>((resolve) => {
319+
release = resolve
320+
})
321+
const writer = db.transaction(async (trx) => {
322+
await lockUniqueValues(trx, stale, [{ email: 'stale@example.test' }])
323+
holdsUniqueLock()
324+
await released
325+
await acquireRowOrderLock(trx, table.id)
326+
})
327+
await holding
328+
329+
const replace = importReplaceRows(
330+
table,
331+
[{ id: 'code', name: 'code', type: 'string', unique: true }],
332+
{ rows: [{ name: 'fresh', code: 'c-1' }], workspaceId },
333+
'unique-race'
334+
)
335+
expect(await waitForLockWaiters(table.id, { onOrderLock: 0, onValueLock: 1 })).toEqual({
336+
onOrderLock: 0,
337+
onValueLock: 1,
338+
})
339+
release()
340+
341+
const results = await Promise.allSettled([writer, replace])
342+
expect(results.map((result) => result.status)).toEqual(['fulfilled', 'fulfilled'])
343+
})
344+
345+
it('holds a bounded number of locks however many unique columns a table has', async () => {
346+
const wide: ColumnDefinition[] = Array.from({ length: 100 }, (_, i) => ({
347+
id: `code_${i}`,
348+
name: `code_${i}`,
349+
type: 'string',
350+
unique: true,
351+
}))
352+
const table = await createTable(wide)
353+
const row = Object.fromEntries(wide.map((_, i) => [`code_${i}`, `value-${i}`]))
354+
const advisoryLocksHeld = async (lock: (trx: DbTransaction) => Promise<void>) =>
355+
db.transaction(async (trx) => {
356+
await lock(trx)
357+
const [{ held }] = await trx.execute<{ held: number }>(sql`SELECT count(*)::int AS held
358+
FROM pg_locks WHERE locktype = 'advisory' AND pid = pg_backend_pid()`)
359+
return held
360+
})
361+
362+
expect(await advisoryLocksHeld((trx) => lockUniqueValues(trx, table, [row]))).toBe(1)
363+
expect(await advisoryLocksHeld((trx) => lockUniqueColumns(trx, table))).toBe(1)
364+
const narrow = { code_0: 'value-0', code_1: 'value-1' }
365+
expect(await advisoryLocksHeld((trx) => lockUniqueValues(trx, table, [narrow]))).toBe(3)
366+
})
367+
295368
it('serializes a batch over the value-lock cap against a single insert of the same value', async () => {
296369
const table = await createTable(uniqueColumns)
297370
// More values than the per-transaction value-lock cap, so the batch locks the column instead.
@@ -331,7 +404,10 @@ describe('table row writes against real PostgreSQL', () => {
331404
() => insertEmail(table, 'dup@example.test'),
332405
async () => {
333406
// Start the import only once the insert has passed its check and holds its locks.
334-
await waitForLockWaiters(table.id, { onOrderLock: 1, onValueLock: 0 })
407+
expect(await waitForLockWaiters(table.id, { onOrderLock: 1, onValueLock: 0 })).toEqual({
408+
onOrderLock: 1,
409+
onValueLock: 0,
410+
})
335411
return bulkInsertImportBatch(
336412
{
337413
tableId: table.id,

‎apps/sim/lib/table/rows/service.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -508,7 +508,7 @@ export async function replaceTableRows(
508508
* Capacity is NOT checked here (it would mean a billing-pool read inside the tx).
509509
* Callers gate it before opening the tx — see `replaceTableRows` and `importReplaceRows`.
510510
*
511-
* Takes the table's unique-column locks before its row-order lock, so a caller already holding the
511+
* Takes the table's unique lock before its row-order lock, so a caller already holding the
512512
* row-order lock must take `lockUniqueColumns` first.
513513
*/
514514
export async function replaceTableRowsWithTx(

‎apps/sim/lib/table/rows/unique-locks.ts‎

Lines changed: 38 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,9 @@
77
* values it writes, so the second writer's check runs after the first commits and sees its row.
88
* Writers of different values never wait on each other.
99
*
10-
* Lock order, everywhere: the table's schema lock (when taken), then these locks, taken all at once
11-
* in one sorted order, then the table's row-order lock, then the definition row.
10+
* Lock order, everywhere: the table's schema lock (when taken), then the table's unique lock and its
11+
* value locks in sorted order, in one statement, then the table's row-order lock, then the
12+
* definition row.
1213
*/
1314

1415
import { compareStrings } from '@sim/utils/string'
@@ -21,77 +22,72 @@ import { getUniqueColumns, uniqueValueKey } from '@/lib/table/validation'
2122
const UNIQUE_LOCK_TAG = 'user_table_unique_value'
2223

2324
/**
24-
* Most value locks one transaction takes before it locks whole columns instead. Every advisory lock
25-
* holds a slot in the server's shared lock table, which is sized at `max_locks_per_transaction`
26-
* (64 by default) × connections and shared by every transaction, so concurrent writers each holding
27-
* more than their per-connection share can exhaust it and fail unrelated queries with `out of
28-
* shared memory`. Capping at that per-connection budget keeps ordinary writes on per-value locks
29-
* while larger batches serialize on the column.
25+
* Most value locks one transaction takes before it locks the table's unique values as a whole.
26+
* Every advisory lock holds a slot in the server's shared lock table, which is sized at
27+
* `max_locks_per_transaction` (64 by default) × connections and shared by every transaction, so
28+
* concurrent writers each holding more than their per-connection share can exhaust it and fail
29+
* unrelated queries with `out of shared memory`. Capping at that per-connection budget keeps
30+
* ordinary writes on per-value locks while larger batches serialize on the table.
3031
*/
3132
const MAX_VALUE_LOCKS = 64
3233

33-
function columnLockKey(tableId: string, columnId: string): string {
34-
return `user_table_unique:${tableId}:${columnId}`
34+
/**
35+
* One lock per table over all its unique columns: value writers hold it shared, and writers that
36+
* cannot lock by value hold it exclusively. A single key bounds a transaction at one lock plus its
37+
* value locks, however many unique columns the schema has.
38+
*/
39+
function tableLockKey(tableId: string): string {
40+
return `user_table_unique:${tableId}`
3541
}
3642

3743
/**
3844
* Locks the unique-column values `rows` will write, before their unique check. Pass `columnIds` to
3945
* lock only the unique columns a patch changes; values it leaves alone are already stored.
4046
*
41-
* Each written column takes a shared column lock, then an exclusive lock per value. The value key
42-
* is `uniqueValueKey`, so two values the check treats as equal always share a key; null cells
43-
* never conflict and take none.
44-
* A column falls back to one exclusive column lock when a value is an object or array (a `json`
45-
* column's check matches by containment, which no single key can express), and every column does
46-
* when the transaction would exceed {@link MAX_VALUE_LOCKS}. An exclusive column lock waits for,
47-
* and blocks, every shared holder, so the fallback still serializes against per-value writers.
47+
* The table lock is taken shared, then an exclusive lock per value. The value key includes
48+
* `uniqueValueKey`, so two values the check treats as equal always share a key; null cells never
49+
* conflict and take none. The table lock is taken exclusively instead, with no value locks, when a
50+
* value is an object or array (a `json` column's check matches by containment, which no single key
51+
* can express) or when the transaction would exceed {@link MAX_VALUE_LOCKS}. An exclusive holder
52+
* waits for, and blocks, every shared holder, so it still serializes against per-value writers.
4853
*/
4954
export async function lockUniqueValues(
5055
trx: DbTransaction,
5156
table: TableDefinition,
5257
rows: readonly RowData[],
5358
columnIds?: ReadonlySet<string>
5459
): Promise<void> {
55-
const columnLocks = new Map<string, boolean>()
60+
const tableKey = tableLockKey(table.id)
5661
const valueKeys = new Set<string>()
5762
for (const column of getUniqueColumns(table.schema)) {
5863
const columnId = getColumnId(column)
5964
if (columnIds && !columnIds.has(columnId)) continue
60-
const columnKey = columnLockKey(table.id, columnId)
61-
const keys: string[] = []
62-
let columnByValue = true
6365
for (const row of rows) {
6466
const value = row[columnId]
6567
if (value === null || value === undefined) continue
66-
if (typeof value === 'object') {
67-
columnByValue = false
68-
break
68+
if (typeof value === 'object') return lockUniqueColumns(trx, table)
69+
const key = `${tableKey}:${columnId}:${uniqueValueKey(value, column)}`
70+
if (!valueKeys.has(key) && valueKeys.size === MAX_VALUE_LOCKS) {
71+
return lockUniqueColumns(trx, table)
6972
}
70-
keys.push(`${columnKey}:${uniqueValueKey(value, column)}`)
73+
valueKeys.add(key)
7174
}
72-
if (columnByValue && keys.length === 0) continue
73-
columnLocks.set(columnKey, columnByValue)
74-
if (columnByValue) for (const key of keys) valueKeys.add(key)
7575
}
76+
if (valueKeys.size === 0) return
7677

77-
const byValue = valueKeys.size <= MAX_VALUE_LOCKS
78-
const locks: AdvisoryXactLockRequest[] = [...columnLocks.entries()]
79-
.sort(([a], [b]) => compareStrings(a, b))
80-
.map(([key, columnByValue]) => ({ key, shared: byValue && columnByValue }))
81-
if (byValue) {
82-
for (const key of [...valueKeys].sort(compareStrings)) locks.push({ key, shared: false })
83-
}
78+
const locks: AdvisoryXactLockRequest[] = [{ key: tableKey, shared: true }]
79+
for (const key of [...valueKeys].sort(compareStrings)) locks.push({ key, shared: false })
8480
await acquireAdvisoryXactLocks(trx, UNIQUE_LOCK_TAG, locks)
8581
}
8682

8783
/**
88-
* Locks every unique column of `table` exclusively, for writers that replace the table's rows
89-
* wholesale and so conflict with any concurrent write of a unique value.
84+
* Locks every unique value of `table` exclusively, for writers that replace the table's rows
85+
* wholesale and so conflict with any concurrent write of a unique value. It locks even a table with
86+
* no unique columns yet: a writer holding an older schema may still hold the lock shared, and a
87+
* whole-table writer that adds a unique column must already hold it before the row-order lock.
9088
*/
9189
export async function lockUniqueColumns(trx: DbTransaction, table: TableDefinition): Promise<void> {
92-
const locks = getUniqueColumns(table.schema)
93-
.map((column) => columnLockKey(table.id, getColumnId(column)))
94-
.sort(compareStrings)
95-
.map((key) => ({ key, shared: false }))
96-
await acquireAdvisoryXactLocks(trx, UNIQUE_LOCK_TAG, locks)
90+
await acquireAdvisoryXactLocks(trx, UNIQUE_LOCK_TAG, [
91+
{ key: tableLockKey(table.id), shared: false },
92+
])
9793
}

0 commit comments

Comments
 (0)