diff --git a/apps/sim/background/table-update.ts b/apps/sim/background/table-update.ts deleted file mode 100644 index e26fd7e8329..00000000000 --- a/apps/sim/background/table-update.ts +++ /dev/null @@ -1,45 +0,0 @@ -import { AbortTaskRunError, task } from '@trigger.dev/sdk' -import { - markTableUpdateFailed, - runTableUpdate, - type TableUpdatePayload, - UpdatePatchRejectedError, -} from '@/lib/table/update-runner' - -/** - * `TableUpdatePayload` with the cutoff as an ISO string — task payloads cross a JSON boundary, so - * the Date is rehydrated in `run` rather than trusting payload serialization. - */ -export interface TableUpdateTaskPayload extends Omit { - cutoff: string -} - -/** - * Trigger.dev wrapper around `runTableUpdate`. Errors propagate out of `run` so the retry policy - * fires; the job is marked failed only in `onFailure`, after the final attempt. Retry-safe: the - * worker keysets by id with a `created_at <= cutoff` floor and the JSONB-merge patch is idempotent - * (re-applying the same patch to an already-patched row is a no-op), so a retried attempt re-walks - * and re-applies whatever remains. The `table_jobs` ownership gate stops a retried run that lost - * the job within one page. A patch the table's schema refuses aborts without a retry: the retry - * would read the same schema. - */ -export const tableUpdateTask = task({ - id: 'table-update', - machine: 'small-1x', - retry: { maxAttempts: 3 }, - queue: { - name: 'table-update', - concurrencyLimit: 10, - }, - run: async (payload: TableUpdateTaskPayload) => { - try { - await runTableUpdate({ ...payload, cutoff: new Date(payload.cutoff) }) - } catch (error) { - if (error instanceof UpdatePatchRejectedError) throw new AbortTaskRunError(error.message) - throw error - } - }, - onFailure: async ({ payload, error }) => { - await markTableUpdateFailed(payload.tableId, payload.jobId, error) - }, -}) diff --git a/apps/sim/lib/table/delete-runner.ts b/apps/sim/lib/table/delete-runner.ts index 6b010f02518..d6823ae6fac 100644 --- a/apps/sim/lib/table/delete-runner.ts +++ b/apps/sim/lib/table/delete-runner.ts @@ -119,7 +119,7 @@ export async function runTableDelete(payload: TableDeletePayload): Promise : 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 ?? []) diff --git a/apps/sim/lib/table/rows/ordering.ts b/apps/sim/lib/table/rows/ordering.ts index 0919efa7582..d80f943b246 100644 --- a/apps/sim/lib/table/rows/ordering.ts +++ b/apps/sim/lib/table/rows/ordering.ts @@ -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, - /** 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 { - 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 -} diff --git a/apps/sim/lib/table/rows/row-writes.integration.ts b/apps/sim/lib/table/rows/row-writes.integration.ts index e5c3a5de422..512ad3d0bdb 100644 --- a/apps/sim/lib/table/rows/row-writes.integration.ts +++ b/apps/sim/lib/table/rows/row-writes.integration.ts @@ -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' @@ -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) { @@ -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, - change: () => Promise - ): Promise[]> { - 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 = ` for the table under its exclusive schema lock. */ - async function changeSchemaUnderLock(tableId: string, change: string): Promise { - 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 { - 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')) diff --git a/apps/sim/lib/table/sql.test.ts b/apps/sim/lib/table/sql.test.ts index f700ca08043..5c1cde0f381 100644 --- a/apps/sim/lib/table/sql.test.ts +++ b/apps/sim/lib/table/sql.test.ts @@ -788,7 +788,7 @@ describe('buildPredicateClause (v2 grammar)', () => { * compiling to no WHERE clause — which on a bulk delete means every row rather * than none. The guard turns that into a loud, self-describing failure, and it * sits at the one choke point every filter path shares (`queryRows`, - * `update-runner`, `delete-runner`, inline and background). + * `delete-runner`, inline and background). */ describe('legacy compiler rejects a v2 predicate (version-mismatch fail-fast)', () => { it('throws on a top-level all/any group instead of emitting no clause', () => { diff --git a/apps/sim/lib/table/update-runner.test.ts b/apps/sim/lib/table/update-runner.test.ts deleted file mode 100644 index 0ff5b940b1b..00000000000 --- a/apps/sim/lib/table/update-runner.test.ts +++ /dev/null @@ -1,172 +0,0 @@ -import { tableConstantsMock } from '@sim/testing/mocks/table-constants.mock' -import { tableEventsMock, tableEventsMockFns } from '@sim/testing/mocks/table-events.mock' -import { - tableJobsServiceMock, - tableJobsServiceMockFns, -} from '@sim/testing/mocks/table-jobs-service.mock' -import { tableServiceMock, tableServiceMockFns } from '@sim/testing/mocks/table-service.mock' -import { beforeEach, describe, expect, it, vi } from 'vitest' - -const { - mockSelectRowDataPage, - mockUpdatePageByIds, - mockBuildFilterClause, - mockValidateRowSize, - mockCoerceRowToSchema, - mockCoerceRowValues, -} = vi.hoisted(() => ({ - mockSelectRowDataPage: vi.fn(), - mockUpdatePageByIds: vi.fn(), - mockBuildFilterClause: vi.fn(), - mockValidateRowSize: vi.fn(), - mockCoerceRowToSchema: vi.fn(), - mockCoerceRowValues: vi.fn(), -})) - -vi.mock('@/lib/table/service', () => tableServiceMock) -vi.mock('@/lib/table/jobs/service', () => tableJobsServiceMock) -vi.mock('@/lib/table/rows/ordering', () => ({ - selectRowDataPage: mockSelectRowDataPage, - updatePageByIds: mockUpdatePageByIds, -})) -vi.mock('@/lib/table/events', () => tableEventsMock) -vi.mock('@/lib/table/sql', () => ({ buildFilterClause: mockBuildFilterClause })) -vi.mock('@/lib/table/validation', async (importOriginal) => { - const actual = await importOriginal() - return { - cellOf: actual.cellOf, - uniqueColumnsInPatch: actual.uniqueColumnsInPatch, - validateRowSize: mockValidateRowSize, - coerceRowToSchema: mockCoerceRowToSchema, - coerceRowValues: mockCoerceRowValues, - } -}) -vi.mock('@/lib/table/constants', () => ({ - ...tableConstantsMock, - TABLE_LIMITS: { ...tableConstantsMock.TABLE_LIMITS, DELETE_PAGE_SIZE: 2, UPDATE_BATCH_SIZE: 100 }, -})) - -import { runTableUpdate } from '@/lib/table/update-runner' - -const mockGetTableById = tableServiceMockFns.mockGetTableById -const mockGetJobProgress = tableJobsServiceMockFns.mockGetJobProgress -const mockUpdateJobProgress = tableJobsServiceMockFns.mockUpdateJobProgress -const mockMarkJobReady = tableJobsServiceMockFns.mockMarkJobReady -const mockMarkJobFailed = tableJobsServiceMockFns.mockMarkJobFailed -const mockMarkJobCanceled = tableJobsServiceMockFns.mockMarkJobCanceled -const mockAppendTableEvent = tableEventsMockFns.mockAppendTableEvent - -const UNLOCKED = { - schemaLocked: false, - insertLocked: false, - updateLocked: false, - deleteLocked: false, -} -const table = { - id: 'tbl_1', - workspaceId: 'ws_1', - schema: { columns: [{ id: 'flag', name: 'flag', type: 'boolean' }] }, - locks: UNLOCKED, -} -const cutoff = new Date('2026-06-05T00:00:00Z') - -function basePayload(overrides = {}) { - return { - jobId: 'job_1', - tableId: 'tbl_1', - workspaceId: 'ws_1', - filter: { status: 'old' }, - data: { flag: true }, - cutoff, - ...overrides, - } -} -const row = (id: string) => ({ id, data: {} }) - -describe('runTableUpdate', () => { - beforeEach(() => { - mockGetTableById.mockResolvedValue(table) - mockGetJobProgress.mockResolvedValue(0) - mockUpdateJobProgress.mockResolvedValue(true) - mockMarkJobReady.mockResolvedValue(true) - mockMarkJobFailed.mockResolvedValue(undefined) - mockUpdatePageByIds.mockImplementation((_t, _w, ids: string[]) => Promise.resolve(ids.length)) - mockBuildFilterClause.mockReturnValue({}) - mockValidateRowSize.mockReturnValue({ valid: true, errors: [] }) - mockCoerceRowToSchema.mockReturnValue({ valid: true, errors: [] }) - }) - - it('cancels without updating when the table was update-locked before the run started', async () => { - // The lock is asserted at enqueue, but a queued or retried job can start - // after an admin locks the table — nothing is written yet, so honor it. - mockGetTableById.mockResolvedValue({ ...table, locks: { ...UNLOCKED, updateLocked: true } }) - mockSelectRowDataPage.mockResolvedValue([row('a')]) - - await expect(runTableUpdate(basePayload())).resolves.toBeUndefined() - - expect(mockUpdatePageByIds).not.toHaveBeenCalled() - expect(mockMarkJobCanceled).toHaveBeenCalledWith('tbl_1', 'job_1') - expect(mockMarkJobReady).not.toHaveBeenCalled() - expect(mockAppendTableEvent).toHaveBeenCalledWith( - expect.objectContaining({ kind: 'job', type: 'update', status: 'canceled' }) - ) - }) - - it('fails (rethrows) when a merged row is invalid, without writing that page', async () => { - mockSelectRowDataPage.mockResolvedValueOnce([row('a')]) - mockValidateRowSize.mockReturnValueOnce({ valid: false, errors: ['row too large'] }) - - await expect(runTableUpdate(basePayload())).rejects.toThrow(/Row a: row too large/) - expect(mockUpdatePageByIds).not.toHaveBeenCalled() - expect(mockMarkJobFailed).not.toHaveBeenCalled() // caller decides via markTableUpdateFailed - }) - - it('stops without marking ready when the ownership gate is lost', async () => { - mockSelectRowDataPage.mockResolvedValue([row('a'), row('b')]) - mockUpdateJobProgress.mockResolvedValueOnce(true).mockResolvedValueOnce(false) - - await runTableUpdate(basePayload()) - - expect(mockUpdatePageByIds).toHaveBeenCalledTimes(1) - expect(mockMarkJobReady).not.toHaveBeenCalled() - }) - - it('rethrows the root cause so the clean message survives serialization', async () => { - const cause = new Error('canceling statement due to statement timeout') - mockSelectRowDataPage.mockRejectedValue(new Error('Failed query: update ...', { cause })) - - await expect(runTableUpdate(basePayload())).rejects.toThrow( - 'canceling statement due to statement timeout' - ) - expect(mockMarkJobFailed).not.toHaveBeenCalled() - }) - - it('resumes cumulative progress on retry instead of resetting to zero', async () => { - mockGetJobProgress.mockResolvedValue(7) - mockSelectRowDataPage.mockResolvedValueOnce([row('a'), row('b')]).mockResolvedValueOnce([]) - - await runTableUpdate(basePayload()) - - expect(mockUpdateJobProgress).toHaveBeenNthCalledWith(1, 'tbl_1', 7, 'job_1') - expect(mockAppendTableEvent).toHaveBeenCalledWith( - expect.objectContaining({ status: 'ready', progress: 9 }) - ) - }) - - it('stops once maxRows is reached and never over-fetches a page', async () => { - // budget 3 with page size 2: first page fills 2, second page is capped to the remaining 1. - mockSelectRowDataPage - .mockResolvedValueOnce([row('a'), row('b')]) - .mockResolvedValueOnce([row('c')]) - - await runTableUpdate(basePayload({ maxRows: 3 })) - - expect(mockSelectRowDataPage).toHaveBeenCalledTimes(2) - expect(mockSelectRowDataPage.mock.calls[0][0]).toMatchObject({ limit: 2 }) - expect(mockSelectRowDataPage.mock.calls[1][0]).toMatchObject({ limit: 1 }) - expect(mockUpdatePageByIds).toHaveBeenCalledTimes(2) - expect(mockAppendTableEvent).toHaveBeenCalledWith( - expect.objectContaining({ status: 'ready', progress: 3 }) - ) - }) -}) diff --git a/apps/sim/lib/table/update-runner.ts b/apps/sim/lib/table/update-runner.ts deleted file mode 100644 index e1a75e84b15..00000000000 --- a/apps/sim/lib/table/update-runner.ts +++ /dev/null @@ -1,360 +0,0 @@ -import { userTableRows } from '@sim/db/schema' -import { createLogger } from '@sim/logger' -import { getErrorMessage, toError } from '@sim/utils/errors' -import { generateId } from '@sim/utils/id' -import { truncate } from '@sim/utils/string' -import { and, eq, inArray } from 'drizzle-orm' -import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { Filter, RowData, TableDefinition, TableSchema } from '@/lib/table' -import { getColumnId } from '@/lib/table/column-keys' -import { TABLE_LIMITS, USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants' -import { appendTableEvent } from '@/lib/table/events' -import { - getJobProgress, - markJobCanceled, - markJobFailed, - markJobReady, - updateJobProgress, -} from '@/lib/table/jobs/service' -import { - assertRowUpdate, - type MutationProof, - patchColumnIds, - TableLockedError, -} from '@/lib/table/mutation-locks' -import type { DbTransaction } from '@/lib/table/planner' -import { withLiveSchema } from '@/lib/table/rows/live-schema' -import { selectRowDataPage, updatePageByIds } from '@/lib/table/rows/ordering' -import { createExactEmptyTableRowSecretProvenance } from '@/lib/table/rows/secret-provenance' -import { deriveBulkUpdatePatch } from '@/lib/table/rows/service' -import { getTableById } from '@/lib/table/service' -import { buildFilterClause } from '@/lib/table/sql' -import { coerceRowToSchema, uniqueColumnsInPatch, validateRowSize } from '@/lib/table/validation' - -const logger = createLogger('TableUpdateRunner') - -/** Emit a progress event / heartbeat at most every this many rows. */ -const PROGRESS_INTERVAL_ROWS = 5000 - -/** - * Thrown when this worker discovers it no longer owns the table's job (canceled, or the - * stale-job janitor marked it failed and a newer job took over). The worker stops updating. - */ -class JobSupersededError extends Error {} - -/** - * The patch cannot be applied to the table as it now stands. Retrying cannot help: the schema that - * refused it is the one a retry would read. - */ -export class UpdatePatchRejectedError extends Error { - constructor(message: string) { - super(message) - this.name = 'UpdatePatchRejectedError' - } -} - -/** - * The patch one batch writes against `live`, derived from the job's raw `patch` exactly as the - * inline bulk update derives its own ({@link deriveBulkUpdatePatch}): cells of columns deleted since - * `previous` dropped, the rest coerced and validated. It is refused, with an - * {@link UpdatePatchRejectedError}, when it cannot be written to every matched row: a value the - * columns refuse, a null in a required column, or any unique column, since one value in many rows - * cannot stay unique. Batches derive it under the schema lock, so a column made required or unique - * while the job waited is refused before the batch writes. - */ -function deriveJobPatch(patch: RowData, previous: TableSchema, live: TableDefinition): RowData { - let derived: RowData - try { - derived = deriveBulkUpdatePatch(patch, previous, live, undefined) - } catch (err) { - if (err instanceof OrchestrationError && err.code === 'validation') { - throw new UpdatePatchRejectedError(err.message) - } - throw err - } - const cleared = live.schema.columns.filter((column) => { - const columnId = getColumnId(column) - return ( - column.required && Object.hasOwn(derived, columnId) && (derived[columnId] ?? null) === null - ) - }) - if (cleared.length > 0) { - throw new UpdatePatchRejectedError( - `Missing required field: ${cleared.map((column) => column.name).join(', ')}` - ) - } - const unique = uniqueColumnsInPatch(live.schema, derived) - if (unique.length > 0) { - throw new UpdatePatchRejectedError( - `Cannot set unique column values when updating multiple rows: ${unique.map((column) => column.name).join(', ')}` - ) - } - return derived -} - -/** Refuses a row that `patch`, merged over it, would leave oversized or invalid under `schema`. */ -function assertMergedRowFits( - schema: TableSchema, - row: { id: string; data: RowData }, - patch: RowData -): void { - const merged = { ...row.data, ...patch } - const sizeValidation = validateRowSize(merged) - if (!sizeValidation.valid) { - throw new UpdatePatchRejectedError(`Row ${row.id}: ${sizeValidation.errors.join(', ')}`) - } - const schemaValidation = coerceRowToSchema(merged, schema) - if (!schemaValidation.valid) { - throw new UpdatePatchRejectedError(`Row ${row.id}: ${schemaValidation.errors.join(', ')}`) - } -} - -/** Reads a batch's rows in its transaction and refuses any `patch` would not fit under `live`. */ -async function validateMergedRows( - trx: DbTransaction, - live: TableDefinition, - rowIds: string[], - patch: RowData -): Promise { - const rows = await trx - .select({ id: userTableRows.id, data: userTableRows.data }) - .from(userTableRows) - .where(and(eq(userTableRows.tableId, live.id), inArray(userTableRows.id, rowIds))) - for (const row of rows) - assertMergedRowFits(live.schema, { id: row.id, data: row.data as RowData }, patch) -} - -export interface TableUpdatePayload { - jobId: string - tableId: string - workspaceId: string - /** Rows matching this filter get the patch. */ - filter: Filter - /** Column-id-keyed partial patch merged into every matched row. */ - data: RowData - /** Only rows created at/before this instant are patched, so mid-job inserts are spared. */ - cutoff: Date - /** Stop after updating this many rows (an explicit caller-supplied limit). Omitted = every match. */ - maxRows?: number -} - -/** - * Background worker for large filtered row updates (trigger.dev task, or detached on the web - * container when trigger.dev is disabled — see the update dispatch in the user_table tool). - * Applies the same `data` patch (JSONB merge) to every row matching `filter` with - * `created_at <= cutoff`, in keyset-paginated pages. Each page validates the merged result per - * row, then commits in batches — **best-effort, not atomic**: committed pages persist even if a - * later page fails validation (unlike the inline `updateRowsByFilter`, which pre-validates all - * rows in one transaction). Reads are not masked: updated rows still exist, so mid-job reads are - * eventually consistent. Ownership-gated per page so a cancel/supersede stops within one page. - * - * Unlike the inline path, the worker does NOT fire per-row table triggers or auto-recompute - * workflow/enrichment columns — that would be a runaway cascade across thousands of rows. Run - * the affected columns explicitly afterward if downstream recompute is needed. - * - * Unexpected errors are rethrown for the caller's retry machinery; the caller marks the job - * failed via `markTableUpdateFailed`. An {@link UpdatePatchRejectedError} is rethrown too, and is - * not worth a retry. A superseded run returns quietly. - */ -export async function runTableUpdate(payload: TableUpdatePayload): Promise { - const { jobId, tableId, workspaceId, filter, data, cutoff, maxRows } = payload - const requestId = generateId().slice(0, 8) - const budget = maxRows ?? Number.POSITIVE_INFINITY - - try { - const table = await getTableById(tableId, { includeArchived: true }) - if (!table) throw new Error(`Update target table ${tableId} not found`) - - // Gate the run on the update lock, then re-gate it before every page (see - // the loop below), so enabling the lock stops a job that is already - // running. Runs through `assertRowUpdate` rather than reading - // `updateLocked` directly so the enqueue site and the worker apply - // identical rules. This is a user-driven bulk patch, so it deliberately - // does not pass `computedWrite` — the workflow-output carve-out belongs to - // the cell-write path alone. - const cancelForLock = async (processedSoFar: number): Promise => { - logger.info(`[${requestId}] Update job stopped — table is update-locked`, { - tableId, - jobId, - processedSoFar, - }) - await markJobCanceled(tableId, jobId) - void appendTableEvent({ kind: 'job', type: 'update', tableId, jobId, status: 'canceled' }) - } - - const stopIfLocked = async ( - fresh: TableDefinition, - processedSoFar: number - ): Promise | null> => { - try { - return assertRowUpdate(fresh, patchColumnIds(data)) - } catch (err) { - if (!(err instanceof TableLockedError)) throw err - await cancelForLock(processedSoFar) - return null - } - } - - if ((await stopIfLocked(table, 0)) === null) return - - // Runs inside each batch's transaction, under the same advisory lock the - // lock toggle and schema changes hold, so no page can be written after a - // lock commits, nor under a schema that would not store the patch. - const revalidate = async (trx: DbTransaction) => { - const fresh = await getTableById(tableId, { tx: trx, includeArchived: true }) - if (fresh) assertRowUpdate(fresh, patchColumnIds(data)) - return fresh ?? undefined - } - - const filterClause = buildFilterClause(filter, USER_TABLE_ROWS_SQL_NAME, table.schema.columns) - if (!filterClause) throw new Error('Filter is required for bulk update') - - // `data` stays the raw payload. Each page, and each batch under the schema lock, derives the - // patch it writes from it against the schema of that moment, as the inline update would. Derive - // it once up front too, so a patch the table refuses fails before any page is read. - deriveJobPatch(data, table.schema, table) - - // Resume the persisted count: a retried attempt's earlier pages are already committed, so - // starting at zero would overwrite cumulative progress. Doubles as the initial ownership gate. - const resumed = await getJobProgress(tableId, jobId) - if (resumed === null) throw new JobSupersededError() - - let processed = resumed - let lastReported = resumed - let afterId: string | undefined - - while (processed < budget) { - const owns = await updateJobProgress(tableId, processed, jobId) - if (!owns) throw new JobSupersededError() - - // Cheap early-out before selecting a page we may not be allowed to - // write. The authoritative gate is `revalidate` below, which re-asserts - // inside each batch transaction. Pages already applied stay applied, as - // with an explicit cancel. - const current = await getTableById(tableId, { includeArchived: true }) - if (!current) throw new JobSupersededError() - const pageProof = await stopIfLocked(current, processed) - if (pageProof === null) return - const pagePatch = deriveJobPatch(data, table.schema, current) - // Every column the update wrote has since been deleted: nothing is left to write. - if (Object.keys(pagePatch).length === 0) break - const patchJson = JSON.stringify(pagePatch) - - const page = await selectRowDataPage({ - tableId, - workspaceId, - cutoff, - filterClause, - afterId, - limit: Math.min(TABLE_LIMITS.DELETE_PAGE_SIZE, budget - processed), - // Skip rows already carrying the patch so a retried run resumes without re-walking / - // double-counting the rows an earlier attempt updated (updated rows still exist and may - // still match the filter, unlike deletes). - excludeIfPatched: patchJson, - }) - if (page.length === 0) break - afterId = page[page.length - 1].id - - // Validate each merged result before writing the page — a row that would overflow the size - // cap or violate the schema fails the job (earlier pages stay applied; best-effort). - for (const row of page) assertMergedRowFits(current.schema, row, pagePatch) - - try { - processed += await updatePageByIds( - tableId, - workspaceId, - page.map((r) => r.id), - async (trx, fresh, batch) => { - const live = fresh ? withLiveSchema(current, fresh.schema) : current - const batchPatch = deriveJobPatch(data, table.schema, live) - if (Object.keys(batchPatch).length === 0) return null - // The page was checked against `current`; a batch under a schema that moved since is - // checked again, against the rows as they stand, before it writes. - if (live !== current) await validateMergedRows(trx, live, batch, batchPatch) - return { - patchJson: JSON.stringify(batchPatch), - secretProvenance: createExactEmptyTableRowSecretProvenance(batchPatch), - } - }, - pageProof, - revalidate - ) - } catch (err) { - if (!(err instanceof TableLockedError)) throw err - // A lock landed between batches. Batches already committed stay - // applied; `processed` undercounts them, which only affects the final - // progress number on an already-canceled job. - await cancelForLock(processed) - return - } - - if ( - processed - lastReported >= PROGRESS_INTERVAL_ROWS || - (lastReported === 0 && processed > 0) - ) { - lastReported = processed - void appendTableEvent({ - kind: 'job', - type: 'update', - tableId, - jobId, - status: 'running', - progress: processed, - }) - } - } - - await updateJobProgress(tableId, processed, jobId) - const becameReady = await markJobReady(tableId, jobId) - if (becameReady) { - void appendTableEvent({ - kind: 'job', - type: 'update', - tableId, - jobId, - status: 'ready', - progress: processed, - }) - logger.info(`[${requestId}] Update complete`, { tableId, rows: processed }) - } else { - logger.info( - `[${requestId}] Update finished but no longer owns the run (canceled/superseded)`, - { - tableId, - jobId, - } - ) - } - } catch (err) { - if (err instanceof JobSupersededError) { - logger.info(`[${requestId}] Update superseded by a newer run; stopping`, { tableId, jobId }) - return - } - const cause = toError(err).cause - const error = cause ? toError(cause) : toError(err) - logger.error(`[${requestId}] Update failed for table ${tableId}:`, error) - throw error - } -} - -/** - * Marks the update job failed and emits the failed SSE event. Called once the caller gives up on - * the run (trigger.dev `onFailure` after retries, or the detached fallback). Scoped to jobId — a - * no-op if a newer job has taken over. - */ -export async function markTableUpdateFailed( - tableId: string, - jobId: string, - error: unknown -): Promise { - const message = truncate(getErrorMessage(toError(error).cause ?? error, 'Update failed'), 500) - await markJobFailed(tableId, jobId, message).catch(() => {}) - void appendTableEvent({ - kind: 'job', - type: 'update', - tableId, - jobId, - status: 'failed', - error: message, - }) -}