Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 0 additions & 45 deletions apps/sim/background/table-update.ts

This file was deleted.

2 changes: 1 addition & 1 deletion apps/sim/lib/table/delete-runner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,7 @@ export async function runTableDelete(payload: TableDeletePayload): Promise<void>
: undefined
// A filter that was SUPPLIED but compiles to no clause must never widen into
// "delete every row" — `and()` silently drops an undefined clause downstream.
// Mirrors the guard in update-runner and the inline deleteRowsByFilter path;
// Mirrors the guard in the inline deleteRowsByFilter path;
// an absent filter is still legitimate (delete-all is an explicit caller mode).
if (filter && !filterClause) throw new Error('Filter is required for bulk delete')
const excluded = new Set(excludeRowIds ?? [])
Expand Down
64 changes: 0 additions & 64 deletions apps/sim/lib/table/rows/ordering.ts
Original file line number Diff line number Diff line change
Expand Up @@ -671,67 +671,3 @@ export async function deletePageByIds(
}
return deleted
}

/** The patch one update batch writes, or `null` when it writes nothing. */
export interface PagePatch {
patchJson: string
secretProvenance: TableRowSecretProvenanceWrite
}

/**
* Applies a JSONB-merge patch (`data || patchJson`) to a page of row ids, committed in
* UPDATE_BATCH_SIZE chunks (each its own transaction, 60s timeout) so a large background update
* makes incremental, resumable progress. Each batch takes its patch from `prepare`, called inside
* the batch's transaction with the definition `revalidate` read there and the batch's row ids, so a
* caller can derive it, and check the rows it merges into, against the live schema. Returns the
* number of rows updated.
*/
export async function updatePageByIds(
tableId: string,
workspaceId: string,
rowIds: string[],
prepare: (
trx: DbTransaction,
table: TableDefinition | undefined,
batch: string[]
) => Promise<PagePatch | null>,
/** Proof the caller asserted the update lock (see `mutation-locks.ts`). */
_proof: MutationProof<'update'>,
/** Re-asserts the lock inside each batch transaction. See {@link guardBatch}. */
revalidate?: MutationRevalidator
): Promise<number> {
const now = new Date()
let updated = 0
for (let i = 0; i < rowIds.length; i += TABLE_LIMITS.UPDATE_BATCH_SIZE) {
const batch = rowIds.slice(i, i + TABLE_LIMITS.UPDATE_BATCH_SIZE)
const rows = await db.transaction(async (trx) => {
await setTableTxTimeouts(trx, { statementMs: 60_000 })
const patch = await prepare(trx, await guardBatch(trx, tableId, revalidate), batch)
if (!patch) return []
return mutateTableRowsWithSecretProvenance(trx, {
rows: batch.map((rowId) => ({ rowId, provenance: patch.secretProvenance })),
rowState: 'existing',
mode: 'merge',
mutate: async () => {
const rows = await trx
.update(userTableRows)
.set({
data: sql`${userTableRows.data} || ${patch.patchJson}::jsonb`,
updatedAt: now,
})
.where(
and(
eq(userTableRows.tableId, tableId),
eq(userTableRows.workspaceId, workspaceId),
inArray(userTableRows.id, batch)
)
)
.returning({ id: userTableRows.id })
return { value: rows, affectedRowIds: rows.map((row) => row.id) }
},
})
})
updated += rows.length
}
return updated
}
250 changes: 2 additions & 248 deletions apps/sim/lib/table/rows/row-writes.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,9 @@ vi.mock('@/lib/table/trigger', () => tableTriggerMock)
vi.mock('@/lib/table/workflow-columns', () => tableWorkflowColumnsMock)
vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock)

import {
deleteColumn,
updateColumnConstraints,
updateColumnOptions,
} from '@/lib/table/columns/service'
import { getMaxRowSizeBytes, TABLE_LIMITS } from '@/lib/table/constants'
import { deleteColumn, updateColumnConstraints } from '@/lib/table/columns/service'
import { getMaxRowSizeBytes } from '@/lib/table/constants'
import { bulkInsertImportBatch, importReplaceRows } from '@/lib/table/import-data'
import { markTableJobRunningInWorkspace } from '@/lib/table/jobs/service'
import type { DbTransaction } from '@/lib/table/planner'
import { lockLiveTableSchema } from '@/lib/table/rows/live-schema'
import { acquireRowOrderLock } from '@/lib/table/rows/ordering'
Expand All @@ -46,7 +41,6 @@ import { lockUniqueColumns, lockUniqueValues } from '@/lib/table/rows/unique-loc
import { getTableById } from '@/lib/table/service'
import { getOrCreateTableSnapshot } from '@/lib/table/snapshot-cache'
import type { ColumnDefinition, JsonValue, RowData, TableDefinition } from '@/lib/table/types'
import { runTableUpdate, UpdatePatchRejectedError } from '@/lib/table/update-runner'

const url = readTestDatabaseUrl()
if (process.env.DATABASE_URL !== url) {
Expand Down Expand Up @@ -1388,246 +1382,6 @@ describe('table row writes against real PostgreSQL', () => {
})
})

describe('background updates validate their patch against the live schema', () => {
const columns: ColumnDefinition[] = [
{ id: 'email', name: 'email', type: 'string' },
{ id: 'kind', name: 'kind', type: 'string' },
{
id: 'status',
name: 'status',
type: 'select',
options: [{ id: 'opt_open', name: 'Open' }],
},
{ id: 'note', name: 'note', type: 'string' },
{ id: 'due', name: 'due', type: 'date' },
]
/** Two update batches' worth of rows. */
const ROWS = TABLE_LIMITS.UPDATE_BATCH_SIZE + 50

async function seededJob(data: RowData) {
const table = await createTable(columns)
await control`INSERT INTO user_table_rows (id, table_id, workspace_id, data, position, order_key)
SELECT ${table.id} || '-' || lpad(g::text, 4, '0'), ${table.id}, ${workspaceId},
jsonb_build_object('email', 'e' || g || '@example.test', 'kind', 'seed'), g, 'a' || lpad(g::text, 4, '0')
FROM generate_series(1, ${ROWS}) g`
const jobId = generateId()
const filter = { kind: 'seed' }
expect(
await markTableJobRunningInWorkspace(table.id, workspaceId, jobId, 'update', {
filter,
data,
})
).toBe(true)
const run = () =>
runTableUpdate({
jobId,
tableId: table.id,
workspaceId,
filter,
data: { ...data },
cutoff: new Date(),
})
/** Re-runs the job as a retry after a crash would: the job is still running. */
const retry = async () => {
await control`UPDATE table_jobs SET status = 'running', completed_at = NULL
WHERE id = ${jobId}`
return run()
}
return { table, run, retry }
}

/**
* Commits `change` between the job's first and second batch. A held schema lock stops the
* first batch; `change` then queues behind it, and a waiting lock is granted in queue order, so
* the first batch commits, then `change`, then the second batch.
*/
async function changeBetweenBatches(
tableId: string,
run: () => Promise<void>,
change: () => Promise<unknown>
): Promise<PromiseSettledResult<unknown>[]> {
const holder = await control.reserve()
try {
await holder`BEGIN`
await holder`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${tableId}`}, 0))`
const job = run()
await untilSchemaLockWaiters(tableId, 1)
const changed = change()
await untilSchemaLockWaiters(tableId, 2)
await holder`COMMIT`
return await Promise.allSettled([job, changed])
} finally {
await holder`ROLLBACK`.catch(() => {})
holder.release()
}
}

/** Commits `schema = <change>` for the table under its exclusive schema lock. */
async function changeSchemaUnderLock(tableId: string, change: string): Promise<void> {
const changer = await control.reserve()
try {
await changer`BEGIN`
await changer`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${tableId}`}, 0))`
await changer.unsafe(`UPDATE user_table_definitions SET schema = ${change} WHERE id = $1`, [
tableId,
])
await changer`COMMIT`
} finally {
changer.release()
}
}

async function countEmail(tableId: string, email: string): Promise<number> {
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
FROM user_table_rows WHERE table_id = ${tableId} AND data->>'email' = ${email}`
return count
}

it('refuses a patch to a column made unique before the job started, retry included', async () => {
const { table, run } = await seededJob({ email: 'same@example.test' })
await updateColumnConstraints(
{ tableId: table.id, columnName: 'email', unique: true },
'bulk-update'
)

await expect(run()).rejects.toBeInstanceOf(UpdatePatchRejectedError)
await expect(run()).rejects.toBeInstanceOf(UpdatePatchRejectedError)
expect(await countEmail(table.id, 'same@example.test')).toBe(0)
})

it('refuses the next batch once a column it clears is made required mid-job', async () => {
const { table, run } = await seededJob({ note: null })
await control`UPDATE user_table_rows SET data = data || '{"note":"filled"}' WHERE table_id = ${table.id}`

const [job, change] = await changeBetweenBatches(table.id, run, async () => {
const changer = await control.reserve()
try {
await changer`BEGIN`
await changer`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${table.id}`}, 0))`
await changer`UPDATE user_table_definitions
SET schema = jsonb_set(schema, '{columns,3,required}', 'true') WHERE id = ${table.id}`
await changer`COMMIT`
} finally {
changer.release()
}
})

expect(change.status).toBe('fulfilled')
expect(job.status === 'rejected' && job.reason).toBeInstanceOf(UpdatePatchRejectedError)
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
FROM user_table_rows WHERE table_id = ${table.id} AND data->'note' = 'null'::jsonb`
expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE)
})

it('refuses the next batch once a patched column is made unique mid-job, retry included', async () => {
const { table, run } = await seededJob({ email: 'same@example.test' })

const [job, change] = await changeBetweenBatches(table.id, run, () =>
changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,0,unique}', 'true')`)
)

expect(change.status).toBe('fulfilled')
expect(job.status === 'rejected' && job.reason).toBeInstanceOf(UpdatePatchRejectedError)
await expect(run()).rejects.toBeInstanceOf(UpdatePatchRejectedError)
expect(await countEmail(table.id, 'same@example.test')).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE)
})

it('finishes when a column it writes gains a select option mid-job', async () => {
const { table, run } = await seededJob({ status: 'Open' })

const [job, change] = await changeBetweenBatches(table.id, run, () =>
updateColumnOptions(
{
tableId: table.id,
columnName: 'status',
options: [
{ id: 'opt_open', name: 'Open' },
{ id: 'opt_closed', name: 'Closed' },
],
},
'bulk-update'
)
)

expect(change.status).toBe('fulfilled')
expect(job.status).toBe('fulfilled')
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
FROM user_table_rows WHERE table_id = ${table.id} AND data->>'status' = 'opt_open'`
expect(count).toBe(ROWS)
})

it('re-derives the patch for later batches once a patched column is retyped mid-job', async () => {
const { table, run, retry } = await seededJob({ note: '7' })
const notes = async () =>
control<{ kind: string; count: number }[]>`SELECT jsonb_typeof(data->'note') AS kind,
count(*)::int AS count FROM user_table_rows WHERE table_id = ${table.id}
GROUP BY 1 ORDER BY 1`

const [job, change] = await changeBetweenBatches(table.id, run, () =>
changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"number"')`)
)

expect([job.status, change.status]).toEqual(['fulfilled', 'fulfilled'])
expect(await notes()).toEqual([
{ kind: 'number', count: ROWS - TABLE_LIMITS.UPDATE_BATCH_SIZE },
{ kind: 'string', count: TABLE_LIMITS.UPDATE_BATCH_SIZE },
])
const [{ seven }] = await control<{ seven: number }[]>`SELECT count(*)::int AS seven
FROM user_table_rows WHERE table_id = ${table.id} AND data->'note' IN ('"7"', '7')`
expect(seven).toBe(ROWS)

await expect(retry()).resolves.toBeUndefined()
expect(await notes()).toEqual([{ kind: 'number', count: ROWS }])
})

it('drops a column deleted mid-job from later batches', async () => {
const { table, run } = await seededJob({ note: 'x', email: 'y@example.test' })

const [job, change] = await changeBetweenBatches(table.id, run, () =>
changeSchemaUnderLock(table.id, `schema #- '{columns,3}'`)
)

expect([job.status, change.status]).toEqual(['fulfilled', 'fulfilled'])
expect(await countEmail(table.id, 'y@example.test')).toBe(ROWS)
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
FROM user_table_rows WHERE table_id = ${table.id} AND data->>'note' = 'x'`
expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE)
})

it('refuses a batch whose re-derived patch grows a row past the size limit', async () => {
const { table, run } = await seededJob({ note: 7 })
await changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"number"')`)
const bigRowId = `${table.id}-${String(TABLE_LIMITS.UPDATE_BATCH_SIZE + 1).padStart(4, '0')}`
// Stored jsonb orders keys by length, so a merged row reads kind, email, filler, then note.
const email = `e${TABLE_LIMITS.UPDATE_BATCH_SIZE + 1}@example.test`
const base = Buffer.byteLength(JSON.stringify({ kind: 'seed', email, filler: '', note: 7 }))
await control`UPDATE user_table_rows
SET data = data || jsonb_build_object('filler', repeat('x', ${getMaxRowSizeBytes() - base}))
WHERE id = ${bigRowId}`

const [job, change] = await changeBetweenBatches(table.id, run, () =>
changeSchemaUnderLock(table.id, `jsonb_set(schema, '{columns,3,type}', '"string"')`)
)

expect(change.status).toBe('fulfilled')
expect(job.status === 'rejected' && job.reason).toBeInstanceOf(UpdatePatchRejectedError)
const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
FROM user_table_rows WHERE table_id = ${table.id} AND data ? 'note'`
expect(count).toBe(TABLE_LIMITS.UPDATE_BATCH_SIZE)
})

it('writes a date patch given as an epoch number', async () => {
const { table, run } = await seededJob({ due: 1704067200000 })

await run()

const [{ count }] = await control<{ count: number }[]>`SELECT count(*)::int AS count
FROM user_table_rows WHERE table_id = ${table.id}
AND (data->>'due')::timestamptz = '2024-01-01T00:00:00Z'`
expect(count).toBe(ROWS)
})
})

describe.skipIf(!migrated)('rows_version', () => {
it('advances once for a transaction that edits cells across several statements', async () => {
const table = await createTable(textColumns('name'))
Expand Down
Loading
Loading