Skip to content

Commit ecb119e

Browse files
committed
improvement(tables): bump a table's rows version once at commit
1 parent 6a6a882 commit ecb119e

8 files changed

Lines changed: 29894 additions & 47 deletions

File tree

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -332,6 +332,11 @@ export async function insertOrderedRow(params: {
332332
secretProvenance?: TableRowSecretProvenanceWrite
333333
/** Proof the caller asserted the insert lock (see `mutation-locks.ts`). */
334334
proof: MutationProof<'insert'>
335+
/**
336+
* Runs under the row-order lock, before the INSERT. Unique checks belong here: run before the
337+
* lock, two concurrent inserts of the same value would each see no conflict and both commit.
338+
*/
339+
assertUnique?: (trx: DbTransaction) => Promise<void>
335340
}): Promise<{
336341
id: string
337342
data: RowData
@@ -355,6 +360,7 @@ export async function insertOrderedRow(params: {
355360
const [row] = await db.transaction(async (trx) => {
356361
await setTableTxTimeouts(trx)
357362
await acquireRowOrderLock(trx, tableId)
363+
await params.assertUnique?.(trx)
358364

359365
// Resolve the authoritative order key from neighbor ids when given, else from the requested
360366
// position. `order_key` is authoritative — `position` is a best-effort, no-shift companion.
Lines changed: 374 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,374 @@
1+
/**
2+
* Row-write integration tests against the provisioned, disposable TEST_DATABASE_URL database (the
3+
* integration setup points DATABASE_URL at it too). The row-order advisory lock runs for real
4+
* everywhere; the `rows_version` cases also need the migrated deferred trigger and skip without it.
5+
*/
6+
import { db } from '@sim/db'
7+
import { userTableDefinitions, userTableRows } from '@sim/db/schema'
8+
import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
9+
import { storageServiceMock, storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
10+
import { tableBillingMock, tableBillingMockFns } from '@sim/testing/mocks/table-billing.mock'
11+
import { tableTriggerMock } from '@sim/testing/mocks/table-trigger.mock'
12+
import { tableWorkflowColumnsMock } from '@sim/testing/mocks/table-workflow-columns.mock'
13+
import { sleep } from '@sim/utils/helpers'
14+
import { generateId } from '@sim/utils/id'
15+
import postgres from 'postgres'
16+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
17+
18+
vi.mock('@/lib/table/billing', () => tableBillingMock)
19+
vi.mock('@/lib/table/trigger', () => tableTriggerMock)
20+
vi.mock('@/lib/table/workflow-columns', () => tableWorkflowColumnsMock)
21+
vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock)
22+
23+
import { batchInsertRows, insertRow, updateRow } from '@/lib/table/rows/service'
24+
import { getTableById } from '@/lib/table/service'
25+
import { getOrCreateTableSnapshot } from '@/lib/table/snapshot-cache'
26+
import type { ColumnDefinition, TableDefinition } from '@/lib/table/types'
27+
28+
const url = readTestDatabaseUrl()
29+
if (process.env.DATABASE_URL !== url) {
30+
throw new Error('This suite requires only the disposable local test database')
31+
}
32+
const control = postgres(url, { max: 4, onnotice: () => {} })
33+
const workspaceId = generateId()
34+
const userId = generateId()
35+
36+
/**
37+
* `db:push` never installs the migration-only `rows_version` triggers, and a database migrated
38+
* only through 0240 has a statement-level UPDATE trigger under the same name. Only the deferred
39+
* constraint trigger satisfies the `rows_version` cases, so only it enables them.
40+
*/
41+
const [{ migrated }] = await control<{ migrated: boolean }[]>`SELECT EXISTS (
42+
SELECT 1 FROM pg_trigger t
43+
JOIN pg_proc p ON p.oid = t.tgfoid
44+
WHERE t.tgrelid = 'user_table_rows'::regclass
45+
AND t.tgname = 'user_table_rows_version_update_trigger'
46+
AND t.tgconstraint <> 0
47+
AND t.tgdeferrable
48+
AND t.tginitdeferred
49+
AND p.proname = 'bump_user_table_rows_version_at_commit'
50+
) AS migrated`
51+
52+
async function createTable(columns: ColumnDefinition[]): Promise<TableDefinition> {
53+
const id = generateId()
54+
await db
55+
.insert(userTableDefinitions)
56+
.values({ id, workspaceId, name: id, schema: { columns }, createdBy: userId })
57+
const table = await getTableById(id)
58+
if (!table) throw new Error('Fixture table was not created')
59+
return table
60+
}
61+
62+
async function seedRows(
63+
tableId: string,
64+
rows: Array<{ id: string; data: Record<string, string>; orderKey: string | null }>
65+
) {
66+
await db.insert(userTableRows).values(rows.map((row) => ({ ...row, tableId, workspaceId })))
67+
}
68+
69+
async function rowsVersion(tableId: string): Promise<number> {
70+
const [row] = await control`SELECT rows_version FROM user_table_definitions WHERE id = ${tableId}`
71+
return Number(row.rows_version)
72+
}
73+
74+
const textColumns = (...ids: string[]): ColumnDefinition[] =>
75+
ids.map((id) => ({ id, name: id, type: 'string' }))
76+
77+
describe('table row writes against real PostgreSQL', () => {
78+
beforeAll(async () => {
79+
await control`INSERT INTO "user" (id, name, email, email_verified, created_at, updated_at)
80+
VALUES (${userId}, 'Row write fixture', ${`${userId}@example.test`}, true, now(), now())`
81+
await control`INSERT INTO workspace (id, name, owner_id, billed_account_user_id)
82+
VALUES (${workspaceId}, 'Row write fixtures', ${userId}, ${userId})`
83+
})
84+
85+
beforeEach(() => {
86+
tableBillingMockFns.mockAssertRowCapacity.mockResolvedValue(10_000)
87+
})
88+
89+
afterAll(async () => {
90+
await control`DELETE FROM workspace WHERE id = ${workspaceId}`
91+
await control`DELETE FROM "user" WHERE id = ${userId}`
92+
await control.end()
93+
})
94+
95+
describe.skipIf(!migrated)('rows_version', () => {
96+
it('advances once for a transaction that edits cells across several statements', async () => {
97+
const table = await createTable(textColumns('name'))
98+
await seedRows(table.id, [
99+
{ id: `${table.id}-a`, data: { name: 'a' }, orderKey: 'a0' },
100+
{ id: `${table.id}-b`, data: { name: 'b' }, orderKey: 'a1' },
101+
])
102+
const before = await rowsVersion(table.id)
103+
104+
await control.begin(async (tx) => {
105+
await tx`UPDATE user_table_rows SET data = '{"name":"a2"}' WHERE id = ${`${table.id}-a`}`
106+
await tx`UPDATE user_table_rows SET data = '{"name":"b2"}' WHERE id = ${`${table.id}-b`}`
107+
await tx`UPDATE user_table_rows SET data = '{"name":"a3"}' WHERE id = ${`${table.id}-a`}`
108+
})
109+
110+
expect(await rowsVersion(table.id)).toBe(before + 1)
111+
})
112+
113+
it('advances each table once for a transaction that edits rows in two tables', async () => {
114+
const first = await createTable(textColumns('name'))
115+
const second = await createTable(textColumns('name'))
116+
for (const table of [first, second]) {
117+
await seedRows(table.id, [
118+
{ id: `${table.id}-a`, data: { name: 'a' }, orderKey: 'a0' },
119+
{ id: `${table.id}-b`, data: { name: 'b' }, orderKey: 'a1' },
120+
])
121+
}
122+
const [firstBefore, secondBefore] = [
123+
await rowsVersion(first.id),
124+
await rowsVersion(second.id),
125+
]
126+
127+
await control.begin(async (tx) => {
128+
await tx`UPDATE user_table_rows SET data = '{"name":"a2"}' WHERE id = ${`${first.id}-a`}`
129+
await tx`UPDATE user_table_rows SET data = '{"name":"a2"}' WHERE id = ${`${second.id}-a`}`
130+
await tx`UPDATE user_table_rows SET data = '{"name":"b2"}' WHERE id = ${`${first.id}-b`}`
131+
await tx`UPDATE user_table_rows SET order_key = 'a2' WHERE id = ${`${second.id}-b`}`
132+
})
133+
134+
expect(await rowsVersion(first.id)).toBe(firstBefore + 1)
135+
expect(await rowsVersion(second.id)).toBe(secondBefore + 1)
136+
})
137+
138+
it('advances once for a transaction that reorders rows across several statements', async () => {
139+
const table = await createTable(textColumns('name'))
140+
await seedRows(table.id, [
141+
{ id: `${table.id}-a`, data: { name: 'a' }, orderKey: 'a0' },
142+
{ id: `${table.id}-b`, data: { name: 'b' }, orderKey: 'a1' },
143+
])
144+
const before = await rowsVersion(table.id)
145+
146+
await control.begin(async (tx) => {
147+
await tx`UPDATE user_table_rows SET order_key = 'a2' WHERE id = ${`${table.id}-a`}`
148+
await tx`UPDATE user_table_rows SET order_key = 'Zz' WHERE id = ${`${table.id}-b`}`
149+
})
150+
151+
expect(await rowsVersion(table.id)).toBe(before + 1)
152+
})
153+
154+
it('advances once for a service cell edit that also writes executions and provenance', async () => {
155+
const table = await createTable(textColumns('name'))
156+
await seedRows(table.id, [{ id: `${table.id}-a`, data: { name: 'a' }, orderKey: 'a0' }])
157+
const before = await rowsVersion(table.id)
158+
159+
await updateRow(
160+
{
161+
tableId: table.id,
162+
rowId: `${table.id}-a`,
163+
workspaceId,
164+
data: { name: 'edited' },
165+
secretProvenance: { complete: true, columns: {} },
166+
capabilityGovernedUserId: null,
167+
executionsPatch: {
168+
'group-1': {
169+
status: 'completed',
170+
executionId: generateId(),
171+
jobId: null,
172+
workflowId: 'workflow-1',
173+
error: null,
174+
},
175+
},
176+
},
177+
table,
178+
'rows-version-service-edit'
179+
)
180+
181+
expect(await rowsVersion(table.id)).toBe(before + 1)
182+
})
183+
184+
it('stays put for provenance-only, timestamp-only, and executions-only writes', async () => {
185+
const table = await createTable(textColumns('name'))
186+
await seedRows(table.id, [{ id: `${table.id}-a`, data: { name: 'a' }, orderKey: 'a0' }])
187+
const before = await rowsVersion(table.id)
188+
189+
await control`UPDATE user_table_rows SET secret_provenance_version = 1 WHERE id = ${`${table.id}-a`}`
190+
await control`UPDATE user_table_rows SET updated_at = now() WHERE id = ${`${table.id}-a`}`
191+
await control`UPDATE user_table_rows SET data = data WHERE id = ${`${table.id}-a`}`
192+
await updateRow(
193+
{
194+
tableId: table.id,
195+
rowId: `${table.id}-a`,
196+
workspaceId,
197+
data: {},
198+
secretProvenance: undefined,
199+
capabilityGovernedUserId: null,
200+
executionsPatch: {
201+
'group-1': {
202+
status: 'running',
203+
executionId: generateId(),
204+
jobId: null,
205+
workflowId: 'workflow-1',
206+
error: null,
207+
},
208+
},
209+
},
210+
table,
211+
'rows-version-executions-only'
212+
)
213+
214+
expect(await rowsVersion(table.id)).toBe(before)
215+
})
216+
217+
it('lets a second writer commit while the first holds an uncommitted row edit', async () => {
218+
const table = await createTable(textColumns('name'))
219+
await seedRows(table.id, [
220+
{ id: `${table.id}-a`, data: { name: 'a' }, orderKey: 'a0' },
221+
{ id: `${table.id}-b`, data: { name: 'b' }, orderKey: 'a1' },
222+
])
223+
const before = await rowsVersion(table.id)
224+
const first = await control.reserve()
225+
const second = await control.reserve()
226+
try {
227+
await first`BEGIN`
228+
await first`UPDATE user_table_rows SET data = '{"name":"a2"}' WHERE id = ${`${table.id}-a`}`
229+
230+
await second`BEGIN`
231+
await second`SET LOCAL lock_timeout = '1s'`
232+
await second`UPDATE user_table_rows SET data = '{"name":"b2"}' WHERE id = ${`${table.id}-b`}`
233+
await second`COMMIT`
234+
expect(await rowsVersion(table.id)).toBe(before + 1)
235+
236+
await first`COMMIT`
237+
expect(await rowsVersion(table.id)).toBe(before + 2)
238+
} finally {
239+
await first`ROLLBACK`.catch(() => {})
240+
await second`ROLLBACK`.catch(() => {})
241+
first.release()
242+
second.release()
243+
}
244+
})
245+
246+
it('re-keys a snapshot when a writer commits during materialization', async () => {
247+
const table = await createTable(textColumns('name'))
248+
await seedRows(table.id, [{ id: `${table.id}-a`, data: { name: 'a' }, orderKey: 'a0' }])
249+
const before = await rowsVersion(table.id)
250+
251+
const stored = new Map<string, string>()
252+
let concurrentWriteDone = false
253+
storageServiceMockFns.mockHeadObject.mockImplementation(async (key: string) =>
254+
stored.has(key) ? { size: Buffer.byteLength(stored.get(key) ?? '') } : null
255+
)
256+
storageServiceMockFns.mockDeleteFile.mockImplementation(async ({ key }: { key: string }) => {
257+
stored.delete(key)
258+
})
259+
storageServiceMockFns.mockCreateMultipartUpload.mockImplementation(
260+
async ({ key }: { key: string }) => {
261+
let body = ''
262+
return {
263+
write: async (chunk: string) => {
264+
if (!concurrentWriteDone) {
265+
concurrentWriteDone = true
266+
await control`UPDATE user_table_rows SET data = '{"name":"during"}' WHERE id = ${`${table.id}-a`}`
267+
}
268+
body += chunk
269+
},
270+
complete: async () => {
271+
stored.set(key, body)
272+
return { size: Buffer.byteLength(body) }
273+
},
274+
abort: async () => {},
275+
}
276+
}
277+
)
278+
279+
const snapshot = await getOrCreateTableSnapshot(table, 'rows-version-snapshot')
280+
281+
expect(snapshot.version).toBe(before + 1)
282+
expect(snapshot.key).toContain(`/v${before + 1}-`)
283+
expect(stored.get(snapshot.key)).toContain('during')
284+
})
285+
})
286+
287+
describe('unique columns under concurrent inserts', () => {
288+
const uniqueColumns: ColumnDefinition[] = [
289+
{ id: 'email', name: 'email', type: 'string', unique: true },
290+
]
291+
292+
/**
293+
* Holds the table's row-order lock while both inserts start, so each one's pre-insert work runs
294+
* before either can write. Released only once both are seen waiting on this table's lock;
295+
* otherwise the inserts could run one after the other and prove nothing about the race.
296+
*/
297+
async function raceUnderHeldOrderLock(
298+
tableId: string,
299+
insert: () => Promise<unknown>
300+
): Promise<PromiseSettledResult<unknown>[]> {
301+
const holder = await control.reserve()
302+
try {
303+
await holder`BEGIN`
304+
await holder`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_rows_pos:${tableId}`}, 0))`
305+
const racers = Promise.allSettled([insert(), insert()])
306+
let waiting = 0
307+
for (let attempt = 0; attempt < 400 && waiting < 2; attempt++) {
308+
await sleep(5)
309+
;[{ waiting }] = await control<{ waiting: number }[]>`
310+
WITH lock AS (SELECT hashtextextended(${`user_table_rows_pos:${tableId}`}, 0) AS key)
311+
SELECT count(*)::int AS waiting FROM pg_locks, lock
312+
WHERE locktype = 'advisory' AND NOT granted AND objsubid = 1
313+
AND database = (SELECT oid FROM pg_database WHERE datname = current_database())
314+
AND classid = ((lock.key >> 32) & 4294967295)::oid
315+
AND objid = (lock.key & 4294967295)::oid`
316+
}
317+
expect(waiting).toBe(2)
318+
await holder`COMMIT`
319+
return await racers
320+
} finally {
321+
await holder`ROLLBACK`.catch(() => {})
322+
holder.release()
323+
}
324+
}
325+
326+
async function storedEmails(tableId: string): Promise<number> {
327+
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
328+
FROM user_table_rows WHERE table_id = ${tableId} AND data @> '{"email":"dup@example.test"}'`
329+
return count
330+
}
331+
332+
it('rejects the second of two concurrent single-row inserts of the same value', async () => {
333+
const table = await createTable(uniqueColumns)
334+
335+
const results = await raceUnderHeldOrderLock(table.id, () =>
336+
insertRow(
337+
{
338+
tableId: table.id,
339+
workspaceId,
340+
data: { email: 'dup@example.test' },
341+
secretProvenance: undefined,
342+
capabilityGovernedUserId: null,
343+
},
344+
table,
345+
'unique-single-race'
346+
)
347+
)
348+
349+
expect(results.filter((result) => result.status === 'fulfilled')).toHaveLength(1)
350+
expect(await storedEmails(table.id)).toBe(1)
351+
})
352+
353+
it('rejects the second of two concurrent batch inserts of the same value', async () => {
354+
const table = await createTable(uniqueColumns)
355+
356+
const results = await raceUnderHeldOrderLock(table.id, () =>
357+
batchInsertRows(
358+
{
359+
tableId: table.id,
360+
workspaceId,
361+
rows: [{ email: 'dup@example.test' }],
362+
secretProvenance: undefined,
363+
capabilityGovernedUserId: null,
364+
},
365+
table,
366+
'unique-batch-race'
367+
)
368+
)
369+
370+
expect(results.filter((result) => result.status === 'fulfilled')).toHaveLength(1)
371+
expect(await storedEmails(table.id)).toBe(1)
372+
})
373+
})
374+
})

0 commit comments

Comments
 (0)