Skip to content

Commit 093602b

Browse files
committed
fix(tables): re-check update batches under a moved schema, own-key reads
- The background update runner re-reads a batch's rows inside the batch transaction and re-validates them merged with the re-derived patch when the schema moved since the page was checked, as the inline bulk update does; an unchanged schema adds no query. updatePageByIds takes an async per-batch hook with the transaction for this. - batchInsertRowsWithTx returns the definition it validated against, and batchInsertRows dispatches its insert triggers with it. - Row cells and patch keys are read as own properties (Object.hasOwn), so a legacy column keyed by a prototype name such as `constructor` is not seen in rows or patches that do not hold it. - The long schema-wait test asserts the write is waiting on the schema lock before the holder commits.
1 parent 873c4a5 commit 093602b

8 files changed

Lines changed: 206 additions & 87 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -326,7 +326,7 @@ export async function importAppendRows(
326326
generateId().slice(0, 8),
327327
{ uniqueColumnsLocked: true }
328328
)
329-
inserted.push(...batchInserted)
329+
inserted.push(...batchInserted.rows)
330330
}
331331
return { inserted, table: working }
332332
})

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

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -681,15 +681,20 @@ export interface PagePatch {
681681
/**
682682
* Applies a JSONB-merge patch (`data || patchJson`) to a page of row ids, committed in
683683
* UPDATE_BATCH_SIZE chunks (each its own transaction, 60s timeout) so a large background update
684-
* makes incremental, resumable progress. Each batch takes its patch from `patchFor`, called inside
685-
* the batch's transaction with the definition `revalidate` read there, so a caller can derive it
686-
* from the live schema. Returns the number of rows updated.
684+
* makes incremental, resumable progress. Each batch takes its patch from `prepare`, called inside
685+
* the batch's transaction with the definition `revalidate` read there and the batch's row ids, so a
686+
* caller can derive it, and check the rows it merges into, against the live schema. Returns the
687+
* number of rows updated.
687688
*/
688689
export async function updatePageByIds(
689690
tableId: string,
690691
workspaceId: string,
691692
rowIds: string[],
692-
patchFor: (table: TableDefinition | undefined) => PagePatch | null,
693+
prepare: (
694+
trx: DbTransaction,
695+
table: TableDefinition | undefined,
696+
batch: string[]
697+
) => Promise<PagePatch | null>,
693698
/** Proof the caller asserted the update lock (see `mutation-locks.ts`). */
694699
_proof: MutationProof<'update'>,
695700
/** Re-asserts the lock inside each batch transaction. See {@link guardBatch}. */
@@ -701,7 +706,7 @@ export async function updatePageByIds(
701706
const batch = rowIds.slice(i, i + TABLE_LIMITS.UPDATE_BATCH_SIZE)
702707
const rows = await db.transaction(async (trx) => {
703708
await setTableTxTimeouts(trx, { statementMs: 60_000 })
704-
const patch = patchFor(await guardBatch(trx, tableId, revalidate))
709+
const patch = await prepare(trx, await guardBatch(trx, tableId, revalidate), batch)
705710
if (!patch) return []
706711
return mutateTableRowsWithSecretProvenance(trx, {
707712
rows: batch.map((rowId) => ({ rowId, provenance: patch.secretProvenance })),

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

Lines changed: 102 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ import { userTableDefinitions, userTableRows } from '@sim/db/schema'
99
import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
1010
import { storageServiceMock, storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
1111
import { tableBillingMock, tableBillingMockFns } from '@sim/testing/mocks/table-billing.mock'
12-
import { tableTriggerMock } from '@sim/testing/mocks/table-trigger.mock'
12+
import { tableTriggerMock, tableTriggerMockFns } 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'
@@ -97,6 +97,27 @@ async function rowsVersion(tableId: string): Promise<number> {
9797
const textColumns = (...ids: string[]): ColumnDefinition[] =>
9898
ids.map((id) => ({ id, name: id, type: 'string' }))
9999

100+
/** Sessions waiting on the table's schema lock, matched by the key's hash as `pg_locks` shows it. */
101+
async function schemaLockWaiters(tableId: string): Promise<number> {
102+
const [{ waiting }] = await control<{ waiting: number }[]>`
103+
WITH lock AS (SELECT hashtextextended(${`user_table_schema:${tableId}`}, 0) AS key)
104+
SELECT count(*)::int AS waiting FROM pg_locks l CROSS JOIN lock
105+
WHERE l.locktype = 'advisory' AND NOT l.granted
106+
AND l.classid = ((lock.key >> 32) & 4294967295)::oid
107+
AND l.objid = (lock.key & 4294967295)::oid`
108+
return waiting
109+
}
110+
111+
/** Polls until exactly `expected` sessions wait on the table's schema lock, and asserts it. */
112+
async function untilSchemaLockWaiters(tableId: string, expected: number) {
113+
let waiting = 0
114+
for (let attempt = 0; attempt < 400 && waiting !== expected; attempt++) {
115+
await sleep(5)
116+
waiting = await schemaLockWaiters(tableId)
117+
}
118+
expect(waiting).toBe(expected)
119+
}
120+
100121
describe('table row writes against real PostgreSQL', () => {
101122
beforeAll(async () => {
102123
await control`INSERT INTO "user" (id, name, email, email_verified, created_at, updated_at)
@@ -1050,6 +1071,61 @@ describe('table row writes against real PostgreSQL', () => {
10501071
}
10511072
)
10521073

1074+
it('fires the insert trigger of a batch insert with the column names of the live schema', async () => {
1075+
const table = await seededTable()
1076+
const stale: TableDefinition = {
1077+
...table,
1078+
schema: {
1079+
columns: table.schema.columns.map((column) =>
1080+
column.id === 'note' ? { ...column, name: 'old_note' } : column
1081+
),
1082+
},
1083+
}
1084+
tableTriggerMockFns.mockFireTableTrigger.mockClear()
1085+
1086+
await batchInsertRows(
1087+
{
1088+
tableId: table.id,
1089+
workspaceId,
1090+
rows: [{ key: 'k3', note: 'n' }],
1091+
secretProvenance: undefined,
1092+
capabilityGovernedUserId: null,
1093+
},
1094+
stale,
1095+
'stale-schema'
1096+
)
1097+
1098+
const [trigger] = tableTriggerMockFns.mockFireTableTrigger.mock.calls
1099+
const schema = trigger?.[6] as { columns: ColumnDefinition[] }
1100+
expect(schema.columns.map((column) => column.name)).toContain('note')
1101+
})
1102+
1103+
it('does not count a legacy unique column named after an object prototype key as patched', async () => {
1104+
const table = await createTable([
1105+
{ name: 'constructor', type: 'string', unique: true },
1106+
{ id: 'note', name: 'note', type: 'string' },
1107+
])
1108+
await seedRows(table.id, [
1109+
{ id: `${table.id}-a`, data: { constructor: 'a', note: 'n' }, orderKey: 'a0' },
1110+
{ id: `${table.id}-b`, data: { constructor: 'b', note: 'n' }, orderKey: 'a1' },
1111+
])
1112+
1113+
const result = await updateRowsByFilter(
1114+
table,
1115+
{
1116+
filter: { note: 'n' },
1117+
data: { note: 'patched' },
1118+
// The limited path: the paged one skips rows created after its JS-clock cutoff.
1119+
limit: 10,
1120+
secretProvenance: undefined,
1121+
capabilityGovernedUserId: null,
1122+
},
1123+
'prototype-key'
1124+
)
1125+
1126+
expect(result.affectedCount).toBe(2)
1127+
})
1128+
10531129
it('refuses a bulk update writing one value to rows of a column made unique since its snapshot', async () => {
10541130
const table = await seededTable()
10551131
await updateColumnConstraints(
@@ -1120,17 +1196,7 @@ describe('table row writes against real PostgreSQL', () => {
11201196
() => 'rejected'
11211197
)
11221198

1123-
let waiting = 0
1124-
for (let attempt = 0; attempt < 400 && waiting === 0; attempt++) {
1125-
await sleep(5)
1126-
;[{ waiting }] = await control<{ waiting: number }[]>`
1127-
WITH lock AS (SELECT hashtextextended(${schemaLockKey}, 0) AS key)
1128-
SELECT count(*)::int AS waiting FROM pg_locks l CROSS JOIN lock
1129-
WHERE l.locktype = 'advisory' AND NOT l.granted
1130-
AND l.classid = ((lock.key >> 32) & 4294967295)::oid
1131-
AND l.objid = (lock.key & 4294967295)::oid`
1132-
}
1133-
expect(waiting).toBe(1)
1199+
await untilSchemaLockWaiters(table.id, 1)
11341200
expect(await Promise.race([settled, sleep(50).then(() => 'pending')])).toBe('pending')
11351201

11361202
await holder`COMMIT`
@@ -1143,7 +1209,7 @@ describe('table row writes against real PostgreSQL', () => {
11431209

11441210
/**
11451211
* Holds the table's schema lock exclusively, starts `write`, and commits once `write` is seen
1146-
* waiting and `holdMs` has passed. Settles as `write` does.
1212+
* waiting on it and `holdMs` more has passed. Settles as `write` does.
11471213
*/
11481214
async function writeBehindSchemaLock(
11491215
table: TableDefinition,
@@ -1156,6 +1222,7 @@ describe('table row writes against real PostgreSQL', () => {
11561222
await holder`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${table.id}`}, 0))`
11571223
const pending = write()
11581224
pending.catch(() => {})
1225+
await untilSchemaLockWaiters(table.id, 1)
11591226
await sleep(holdMs)
11601227
await holder`COMMIT`
11611228
return await pending
@@ -1301,25 +1368,6 @@ describe('table row writes against real PostgreSQL', () => {
13011368
return { table, run, retry }
13021369
}
13031370

1304-
async function schemaLockWaiters(tableId: string): Promise<number> {
1305-
const [{ waiting }] = await control<{ waiting: number }[]>`
1306-
WITH lock AS (SELECT hashtextextended(${`user_table_schema:${tableId}`}, 0) AS key)
1307-
SELECT count(*)::int AS waiting FROM pg_locks l CROSS JOIN lock
1308-
WHERE l.locktype = 'advisory' AND NOT l.granted
1309-
AND l.classid = ((lock.key >> 32) & 4294967295)::oid
1310-
AND l.objid = (lock.key & 4294967295)::oid`
1311-
return waiting
1312-
}
1313-
1314-
async function untilSchemaLockWaiters(tableId: string, expected: number) {
1315-
let waiting = 0
1316-
for (let attempt = 0; attempt < 400 && waiting !== expected; attempt++) {
1317-
await sleep(5)
1318-
waiting = await schemaLockWaiters(tableId)
1319-
}
1320-
expect(waiting).toBe(expected)
1321-
}
1322-
13231371
/**
13241372
* Commits `change` between the job's first and second batch. A held schema lock stops the
13251373
* first batch; `change` then queues behind it, and a waiting lock is granted in queue order, so
@@ -1478,6 +1526,28 @@ describe('table row writes against real PostgreSQL', () => {
14781526
expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE)
14791527
})
14801528

1529+
it('refuses a batch whose re-derived patch grows a row past the size limit', async () => {
1530+
const { table, run } = await seededJob({ note: 7 })
1531+
await changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"number"')`)
1532+
const bigRowId = `${table.id}-${String(TABLE_LIMITS.UPDATE_BATCH_SIZE + 1).padStart(4, '0')}`
1533+
// Stored jsonb orders keys by length, so a merged row reads kind, email, filler, then note.
1534+
const email = `e${TABLE_LIMITS.UPDATE_BATCH_SIZE + 1}@example.test`
1535+
const base = Buffer.byteLength(JSON.stringify({ kind: 'seed', email, filler: '', note: 7 }))
1536+
await control`UPDATE user_table_rows
1537+
SET data = data || jsonb_build_object('filler', repeat('x', ${getMaxRowSizeBytes() - base}))
1538+
WHERE id = ${bigRowId}`
1539+
1540+
const [job, change] = await changeBetweenBatches(table.id, run, () =>
1541+
changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"string"')`)
1542+
)
1543+
1544+
expect(change.status).toBe('fulfilled')
1545+
expect(job.status === 'rejected' && job.reason).toBeInstanceOf(UpdatePatchRejectedError)
1546+
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
1547+
FROM user_table_rows WHERE table_id = ${table.id} AND data ? 'note'`
1548+
expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE)
1549+
})
1550+
14811551
it('writes a date patch given as an epoch number', async () => {
14821552
const { table, run } = await seededJob({ due: 1704067200000 })
14831553

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

Lines changed: 13 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -310,17 +310,17 @@ export async function batchInsertRows(
310310
addedRows: data.rows.length,
311311
})
312312

313-
const result = await db.transaction((trx) =>
313+
const { rows, table: written } = await db.transaction((trx) =>
314314
batchInsertRowsWithTx(trx, data, table, requestId, { ...options, lockSchema: true })
315315
)
316316
notifyTableRowUsage({
317317
workspaceId: table.workspaceId,
318318
currentRowCount: table.rowCount,
319-
addedRows: result.length,
319+
addedRows: rows.length,
320320
limit: rowLimit,
321321
})
322-
dispatchAfterBatchInsert(table, result, requestId, data.userId, data.capabilityGovernedUserId)
323-
return result
322+
dispatchAfterBatchInsert(written, rows, requestId, data.userId, data.capabilityGovernedUserId)
323+
return rows
324324
}
325325

326326
/** Options for the transaction-bound whole-row writers. */
@@ -345,7 +345,8 @@ export interface BatchInsertOptions extends TxRowWriteOptions {
345345
* Transaction-bound variant of `batchInsertRows`. Validates rows and unique
346346
* constraints, then performs INSERTs inside the provided transaction. Caller
347347
* is responsible for opening the transaction. Use when row inserts must be
348-
* atomic with other writes (e.g., schema mutations) on the same tx.
348+
* atomic with other writes (e.g., schema mutations) on the same tx. Returns the inserted rows and
349+
* the definition they were validated against, for the caller's post-commit dispatch.
349350
*
350351
* Capacity is NOT checked here (it would mean a billing-pool read inside the tx).
351352
* Callers gate it before opening the tx — see `batchInsertRows` and the import paths.
@@ -361,7 +362,7 @@ export async function batchInsertRowsWithTx(
361362
snapshot: TableDefinition,
362363
requestId: string,
363364
options: BatchInsertOptions = {}
364-
): Promise<TableRow[]> {
365+
): Promise<{ rows: TableRow[]; table: TableDefinition }> {
365366
assertRowInsert(snapshot)
366367
const timeouts = { statementMs: 60_000 }
367368
const table = options.lockSchema ? await lockLiveTableSchema(trx, snapshot, timeouts) : snapshot
@@ -456,7 +457,7 @@ export async function batchInsertRowsWithTx(
456457
}))
457458

458459
await options.readProvenance?.capture(trx, result)
459-
return result
460+
return { rows: result, table }
460461
}
461462

462463
/**
@@ -1913,7 +1914,7 @@ export async function updateRow(
19131914
mergedExecutions = applyExecutionsPatch(existingRow.executions, effectiveExecutionsPatch)
19141915
}
19151916
const persistedData = jsonbMergePatch(
1916-
Object.keys(data.data).filter((columnId) => columnId in mergedData),
1917+
Object.keys(data.data).filter((columnId) => Object.hasOwn(mergedData, columnId)),
19171918
mergedData
19181919
)
19191920
const patchedUniqueColumns = getUniqueColumns(live.schema).filter((column) =>
@@ -2188,8 +2189,8 @@ async function assertBulkUpdateUnique(
21882189
rows: BulkUpdateMatch[],
21892190
patch: RowData
21902191
): Promise<void> {
2191-
const uniqueColumnsInUpdate = getUniqueColumns(table.schema).filter(
2192-
(column) => getColumnId(column) in patch
2192+
const uniqueColumnsInUpdate = getUniqueColumns(table.schema).filter((column) =>
2193+
Object.hasOwn(patch, getColumnId(column))
21932194
)
21942195
if (uniqueColumnsInUpdate.length === 0 || rows.length === 0) return
21952196
if (rows.length > 1) {
@@ -2723,8 +2724,8 @@ export async function batchUpdateRows(
27232724
`Row ${update.rowId}: ${refit.errors.join(', ')}`
27242725
)
27252726
}
2726-
update.changedColumnIds = update.changedColumnIds.filter(
2727-
(columnId) => columnId in update.mergedData
2727+
update.changedColumnIds = update.changedColumnIds.filter((columnId) =>
2728+
Object.hasOwn(update.mergedData, columnId)
27282729
)
27292730
const cleared = deriveExecClearsForDataPatch(
27302731
request.data,

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ import { type AdvisoryXactLockRequest, acquireAdvisoryXactLocks } from '@/lib/db
1717
import { getColumnId } from '@/lib/table/column-keys'
1818
import type { DbTransaction } from '@/lib/table/planner'
1919
import type { RowData, TableDefinition } from '@/lib/table/types'
20-
import { getUniqueColumns, uniqueValueKey } from '@/lib/table/validation'
20+
import { cellOf, getUniqueColumns, uniqueValueKey } from '@/lib/table/validation'
2121

2222
const UNIQUE_LOCK_TAG = 'user_table_unique_value'
2323

@@ -63,7 +63,7 @@ export async function lockUniqueValues(
6363
const columnId = getColumnId(column)
6464
if (columnIds && !columnIds.has(columnId)) continue
6565
for (const row of rows) {
66-
const value = row[columnId]
66+
const value = cellOf(row, columnId)
6767
if (value === null || value === undefined) continue
6868
const key = `${tableKey}:${columnId}:${uniqueValueKey(value, column)}`
6969
if (!valueKeys.has(key) && valueKeys.size === MAX_VALUE_LOCKS) {

‎apps/sim/lib/table/update-row.test.ts‎

Lines changed: 16 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -23,18 +23,22 @@ vi.mock('@/lib/table/billing', () => tableBillingMock)
2323

2424
vi.mock('@/lib/table/rows/live-schema', () => tableRowsLiveSchemaMock)
2525

26-
vi.mock('@/lib/table/validation', async (importOriginal) => ({
27-
uniqueValueKey: (await importOriginal<typeof import('@/lib/table/validation')>()).uniqueValueKey,
28-
validateRowSize: vi.fn(() => ({ valid: true, errors: [] })),
29-
validateRowAgainstSchema: vi.fn(() => ({ valid: true, errors: [] })),
30-
coerceRowToSchema: vi.fn(() => ({ valid: true, errors: [] })),
31-
coerceRowValues: vi.fn(),
32-
validateTableName: vi.fn(() => ({ valid: true, errors: [] })),
33-
validateTableSchema: vi.fn(() => ({ valid: true, errors: [] })),
34-
getUniqueColumns: vi.fn(() => []),
35-
checkUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })),
36-
checkBatchUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })),
37-
}))
26+
vi.mock('@/lib/table/validation', async (importOriginal) => {
27+
const actual = await importOriginal<typeof import('@/lib/table/validation')>()
28+
return {
29+
cellOf: actual.cellOf,
30+
uniqueValueKey: actual.uniqueValueKey,
31+
validateRowSize: vi.fn(() => ({ valid: true, errors: [] })),
32+
validateRowAgainstSchema: vi.fn(() => ({ valid: true, errors: [] })),
33+
coerceRowToSchema: vi.fn(() => ({ valid: true, errors: [] })),
34+
coerceRowValues: vi.fn(),
35+
validateTableName: vi.fn(() => ({ valid: true, errors: [] })),
36+
validateTableSchema: vi.fn(() => ({ valid: true, errors: [] })),
37+
getUniqueColumns: vi.fn(() => []),
38+
checkUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })),
39+
checkBatchUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })),
40+
}
41+
})
3842

3943
/**
4044
* Inspects the queued `trx.execute(...)` calls for SQL containing `substring`.

0 commit comments

Comments
 (0)