Skip to content

Commit 1f47329

Browse files
committed
Merge project enforcement parent with table ordering updates
2 parents 903c669 + ab1e8cb commit 1f47329

12 files changed

Lines changed: 528 additions & 442 deletions
Lines changed: 106 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,106 @@
1+
/**
2+
* The run dispatcher's row windows against real PostgreSQL: rows that share a `position` (two
3+
* inserts that assigned it concurrently) are each dispatched exactly once, wherever the window
4+
* boundary falls.
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 { tableEventsMock } from '@sim/testing/mocks/table-events.mock'
10+
import {
11+
tableWorkflowColumnsMock,
12+
tableWorkflowColumnsMockFns,
13+
} from '@sim/testing/mocks/table-workflow-columns.mock'
14+
import { generateId } from '@sim/utils/id'
15+
import postgres from 'postgres'
16+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
17+
18+
const { mockBatchEnqueueAndWait } = vi.hoisted(() => ({ mockBatchEnqueueAndWait: vi.fn() }))
19+
20+
vi.mock('@/lib/core/async-jobs/config', () => ({
21+
getJobQueue: async () => ({ batchEnqueueAndWait: mockBatchEnqueueAndWait }),
22+
}))
23+
vi.mock('@/lib/table/events', () => tableEventsMock)
24+
vi.mock('@/lib/table/workflow-columns', () => tableWorkflowColumnsMock)
25+
26+
import { insertDispatch, runDispatcherToCompletion } from '@/lib/table/dispatcher'
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: 2, onnotice: () => {} })
33+
const workspaceId = generateId()
34+
const userId = generateId()
35+
36+
describe('table run dispatcher against real PostgreSQL', () => {
37+
beforeAll(async () => {
38+
await control`INSERT INTO "user" (id, name, email, email_verified, created_at, updated_at)
39+
VALUES (${userId}, 'Dispatcher fixture', ${`${userId}@example.test`}, true, now(), now())`
40+
await control.begin(async (tx) => {
41+
await tx`INSERT INTO workspace (id, name, owner_id, billed_account_user_id)
42+
VALUES (${workspaceId}, 'Dispatcher fixtures', ${userId}, ${userId})`
43+
await tx`INSERT INTO project (id, name, owner_id)
44+
VALUES (${workspaceId}, 'Dispatcher fixture project', ${userId})`
45+
await tx`INSERT INTO project_workspace (project_id, workspace_id)
46+
VALUES (${workspaceId}, ${workspaceId})`
47+
})
48+
})
49+
50+
afterAll(async () => {
51+
await control.begin(async (tx) => {
52+
await tx`DELETE FROM workspace WHERE id = ${workspaceId}`
53+
await tx`DELETE FROM project WHERE id = ${workspaceId}`
54+
await tx`DELETE FROM "user" WHERE id = ${userId}`
55+
})
56+
await control.end()
57+
})
58+
59+
it('dispatches every row once when a window boundary splits rows that share a position', async () => {
60+
const tableId = generateId()
61+
await db.insert(userTableDefinitions).values({
62+
id: tableId,
63+
workspaceId,
64+
name: tableId,
65+
schema: { columns: [], workflowGroups: [{ id: 'group-1' }] },
66+
createdBy: userId,
67+
})
68+
const positions = [0, 1, 1, 1, 2, 3, 3, 4]
69+
await db.insert(userTableRows).values(
70+
positions.map((position, i) => ({
71+
id: `${tableId}-${i}`,
72+
tableId,
73+
workspaceId,
74+
data: {},
75+
position,
76+
orderKey: `a${i}`,
77+
}))
78+
)
79+
tableWorkflowColumnsMockFns.mockBuildPendingRuns.mockImplementation(
80+
(_table: unknown, rows: Array<{ id: string }>) =>
81+
rows.map((row) => ({ rowId: row.id, groupId: 'group-1', tableId, workspaceId }))
82+
)
83+
tableWorkflowColumnsMockFns.mockBuildEnqueueItems.mockImplementation(async (runs: unknown[]) =>
84+
runs.map((payload) => ({ payload }))
85+
)
86+
const dispatched: string[] = []
87+
mockBatchEnqueueAndWait.mockImplementation(
88+
async (_kind: string, items: Array<{ payload: { rowId: string } }>) => {
89+
for (const item of items) dispatched.push(item.payload.rowId)
90+
}
91+
)
92+
93+
const dispatchId = await insertDispatch({
94+
tableId,
95+
workspaceId,
96+
requestId: 'dispatcher-ties',
97+
mode: 'all',
98+
scope: { groupIds: ['group-1'] },
99+
isManualRun: true,
100+
capabilityGovernedUserId: null,
101+
})
102+
await runDispatcherToCompletion(dispatchId, 2)
103+
104+
expect([...dispatched].sort()).toEqual(positions.map((_, i) => `${tableId}-${i}`).sort())
105+
})
106+
})

‎apps/sim/lib/table/dispatcher.ts‎

Lines changed: 25 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -556,13 +556,34 @@ export async function dispatcherStep(
556556
.select()
557557
.from(userTableRows)
558558
.where(and(...filters))
559-
.orderBy(asc(userTableRows.position))
559+
.orderBy(asc(userTableRows.position), asc(userTableRows.id))
560560
.limit(windowSize)
561561
// Filtered scopes carry a jsonb predicate the planner can't estimate — left alone it
562562
// seq-scans the whole shared relation per window; keep it on the tenant's position index.
563-
const chunk = hasJsonbFilter
564-
? await withSeqscanOff(async (trx) => windowQuery(trx))
565-
: await windowQuery(db)
563+
const runQuery = async <T>(query: (executor: DbExecutor) => PromiseLike<T>): Promise<T> =>
564+
hasJsonbFilter ? withSeqscanOff(async (trx) => query(trx)) : query(db)
565+
const chunk = await runQuery(windowQuery)
566+
567+
// Inserts assign `position` without a lock, so concurrent ones can share it. The next window
568+
// starts strictly after this one's last position, so a full window takes the rest of that tie.
569+
if (chunk.length === windowSize) {
570+
const lastPosition = chunk[chunk.length - 1].position
571+
const seen = chunk.filter((r) => r.position === lastPosition).map((r) => r.id)
572+
const ties = await runQuery((executor) =>
573+
executor
574+
.select()
575+
.from(userTableRows)
576+
.where(
577+
and(
578+
...filters,
579+
eq(userTableRows.position, lastPosition),
580+
notInArray(userTableRows.id, seen)
581+
)
582+
)
583+
.orderBy(asc(userTableRows.id))
584+
)
585+
chunk.push(...ties)
586+
}
566587

567588
if (chunk.length === 0) {
568589
// Through the shared, guarded completion like the other two exits: this

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

Lines changed: 4 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -16,11 +16,7 @@ import { assertRowDelete, assertRowInsert, assertSchemaMutable } from '@/lib/tab
1616
import { nKeysBetween } from '@/lib/table/order-key'
1717
import type { DbTransaction } from '@/lib/table/planner'
1818
import { lockLiveTableSchema, refitRowToSchema, withLiveSchema } from '@/lib/table/rows/live-schema'
19-
import {
20-
acquireRowOrderLock,
21-
guardBatch,
22-
type MutationRevalidator,
23-
} from '@/lib/table/rows/ordering'
19+
import { guardBatch, type MutationRevalidator } from '@/lib/table/rows/ordering'
2420
import {
2521
createExactEmptyTableRowSecretProvenance,
2622
mutateTableRowsWithSecretProvenance,
@@ -60,7 +56,7 @@ export interface BulkImportBatch {
6056
* Inserts one batch of rows for an async import in a single committed statement.
6157
*
6258
* Differs from {@link batchInsertRowsWithTx} for the bulk-load case: caller-supplied
63-
* contiguous order keys (no `acquireRowOrderLock` scan; the caller threads each batch's
59+
* contiguous order keys (no max-key scan; the caller threads each batch's
6460
* anchor from the previous one), no `RETURNING`, and **no `fireTableTrigger` /
6561
* `runWorkflowColumn`** (a 1M-row import must not dispatch a workflow run per row).
6662
* Append and replace imports run this against the live table, so other writers can
@@ -251,9 +247,6 @@ export async function setTableSchemaForImport(
251247
* parsed must be visible to the asserts in `addTableColumnsWithTx` /
252248
* `batchInsertRowsWithTx` / `replaceTableRowsWithTx`, which all read the
253249
* definition they are handed.
254-
*
255-
* Taken before `acquireRowOrderLock` so the order stays advisory → rows_pos →
256-
* definitions, matching every other advisory-lock holder.
257250
*/
258251
async function refreshUnderLock(
259252
trx: DbTransaction,
@@ -297,16 +290,10 @@ export async function importAppendRows(
297290
})
298291
const result = await db.transaction(async (trx) => {
299292
let working = await refreshUnderLock(trx, table)
300-
// Lock the unique columns whole, ahead of the row-order lock: per-value locks for every row of
301-
// an import would flood the server's lock table.
293+
// Lock the unique columns whole: per-value locks for every row of an import would flood the
294+
// server's lock table.
302295
await lockUniqueColumns(trx, working)
303296
if (additions.length > 0) {
304-
// Take the row-order lock before creating columns so this path uses the
305-
// same rows_pos → user_table_definitions order as plain inserts. Creating
306-
// columns first would lock the definition row before rows_pos, inverting
307-
// the order and deadlocking concurrent inserts on this table. The lock is
308-
// re-entrant, so the per-batch acquire below is a no-op.
309-
await acquireRowOrderLock(trx, table.id)
310297
working = await addTableColumnsWithTx(trx, working, additions, ctx.requestId)
311298
}
312299
const inserted: TableRow[] = []
@@ -368,7 +355,6 @@ export async function importReplaceRows(
368355
let working = await refreshUnderLock(trx, table)
369356
await lockUniqueColumns(trx, working)
370357
if (additions.length > 0) {
371-
await acquireRowOrderLock(trx, table.id)
372358
working = await addTableColumnsWithTx(trx, working, additions, requestId)
373359
}
374360
return replaceTableRowsWithTx(

‎apps/sim/lib/table/lock-order.test.ts‎

Lines changed: 0 additions & 100 deletions
This file was deleted.

‎apps/sim/lib/table/order-key.ts‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@
1212
*/
1313

1414
import { generateKeyBetween, generateNKeysBetween } from '@sim/utils/fractional-indexing'
15+
import { generateRandomBytes } from '@sim/utils/random'
1516

1617
/**
1718
* Returns a key that sorts strictly between `a` and `b`. Pass `null` for an open
@@ -32,3 +33,27 @@ export function keyBetween(a: string | null, b: string | null): string {
3233
export function nKeysBetween(a: string | null, b: string | null, n: number): string[] {
3334
return generateNKeysBetween(a, b, n)
3435
}
36+
37+
/** Random halvings {@link appendKeys} takes: two concurrent appends collide with odds 2^-32. */
38+
const APPEND_SLOT_BITS = 32
39+
40+
/**
41+
* Returns `n` ordered keys after `last`, in a slot no concurrent append to the same table picks.
42+
*
43+
* Appends take no lock, so two writers can read the same `last`. Each mints inside the integer
44+
* range `keyBetween(last, null)` opens, narrowed by {@link APPEND_SLOT_BITS} random halvings to a
45+
* private sub-range: keys never collide, and a batch stays contiguous instead of interleaving with
46+
* another writer's. The next append reads the new max and moves to the next integer, so keys do
47+
* not grow across appends — each carries a fixed few extra characters.
48+
*/
49+
export function appendKeys(last: string | null, n: number): string[] {
50+
let lo = keyBetween(last, null)
51+
let hi = keyBetween(lo, null)
52+
const bits = generateRandomBytes(APPEND_SLOT_BITS / 8)
53+
for (let i = 0; i < APPEND_SLOT_BITS; i++) {
54+
const mid = keyBetween(lo, hi)
55+
if ((bits[i >> 3] >> (i & 7)) & 1) lo = mid
56+
else hi = mid
57+
}
58+
return nKeysBetween(lo, hi, n)
59+
}

0 commit comments

Comments
 (0)