refactor(table): split the 5.3k-line service.ts god-file into per-concern modules (#5069)

* refactor(table): extract row ordering, executions, and tx helpers from service.ts

Move the row position/fractional-ordering internals to rows/ordering.ts,
the row-execution (workflow-group result) internals to rows/executions.ts,
and the shared tx-timeout helpers to tx.ts. Pure code-motion — verbatim
bodies, identical behavior. service.ts: 5324 -> 4442 lines.

* refactor(table): extract row CRUD and query into rows/service.ts

Move insert/update/upsert/delete/replace/batch row writes, queryRows/
getRowById reads, and findRowMatches into rows/service.ts. Verbatim
bodies; consumers repointed; @/lib/table barrel re-exports the new module
so callers are unchanged. service.ts: 4442 -> 2788 lines.

* refactor(table): extract column/schema management into columns/service.ts

Move add/rename/delete column ops and column-type/constraint updates into
columns/service.ts. addTableColumnsWithTx stays in service.ts (table-creation
primitive) to avoid a cycle. Verbatim bodies; barrel re-exports the module.
service.ts: 2788 -> 2149 lines.

* refactor(table): extract job state machine and export jobs into jobs/service.ts

Move tableJobs reads/mapping, the job lifecycle state machine, and export-job
queries into jobs/service.ts. service.ts imports latestJob* one-way for table
metadata enrichment (no cycle). Import-data orchestration helpers stay in
service.ts for now. Verbatim bodies. service.ts: 2149 -> 1791 lines.

* refactor(table): extract workflow-group management into workflow-groups/service.ts

Move add/update/delete workflow groups + outputs and pruneStale into
workflow-groups/service.ts, preserving the dynamic backfill-runner import
(cycle-breaker). Verbatim bodies; no cycle. service.ts: 1791 -> 851 lines.

* refactor(table): extract import-job data ops into import-data.ts

Move bulk insert, schema setup, and append/replace import operations into
import-data.ts (consumed by import-runner + import route). service.ts is now
the pure table-entity module (root CRUD + shared lock/column primitives).
Verbatim bodies; no cycle. service.ts: 851 -> 664 lines (5324 at start).

* refactor(table): restore private visibility and drop dead helpers post-split

Un-export DerivedJobFields/JOB_PROJECTION/mapJobRow (file-local in jobs/service.ts,
were private pre-split) so they no longer leak into the @/lib/table barrel. Remove
dead code carried into the new modules: countTables (createTable does its own inline
count check) and the unused buildOrderedRowValues/OrderedRowValue pair.

* refactor(table): use absolute imports throughout and fix an orphaned TSDoc

Normalize all relative imports under lib/table to absolute @/lib/table/...
per the project import rule (the split modules and a few pre-existing
holdouts), and relocate the getTableById TSDoc that had drifted above
applyColumnOrderToSchema. Comment/import-only; zero behavior change.
This commit is contained in:
Waleed
2026-06-15 15:16:41 -07:00
committed by GitHub
parent 02022e9ab2
commit 39d0b56e91
70 changed files with 5027 additions and 4823 deletions
@@ -28,7 +28,7 @@ vi.mock('@sim/utils/id', () => ({
generateId: vi.fn().mockReturnValue('job-id-xyz'),
generateShortId: vi.fn().mockReturnValue('short-id'),
}))
vi.mock('@/lib/table/service', () => ({
vi.mock('@/lib/table/jobs/service', () => ({
markTableJobRunning: mockMarkTableJobRunning,
releaseJobClaim: mockReleaseJobClaim,
}))
@@ -9,7 +9,7 @@ import { runDetached } from '@/lib/core/utils/background'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { markTableDeleteFailed, runTableDelete } from '@/lib/table/delete-runner'
import { markTableJobRunning, releaseJobClaim } from '@/lib/table/service'
import { markTableJobRunning, releaseJobClaim } from '@/lib/table/jobs/service'
import type { TableDeleteJobPayload } from '@/lib/table/types'
import { accessError, checkAccess, tableFilterError } from '@/app/api/table/utils'
@@ -16,7 +16,7 @@ vi.mock('@sim/utils/id', () => ({
generateId: vi.fn().mockReturnValue('job-id-xyz'),
generateShortId: vi.fn().mockReturnValue('short-id'),
}))
vi.mock('@/lib/table/service', () => ({ markTableJobRunning: mockMarkTableJobRunning }))
vi.mock('@/lib/table/jobs/service', () => ({ markTableJobRunning: mockMarkTableJobRunning }))
vi.mock('@/lib/table/export-runner', () => ({ runTableExport: mockRunTableExport }))
vi.mock('@/lib/core/utils/background', () => ({
runDetached: (_label: string, work: () => Promise<unknown>) => {
@@ -9,7 +9,7 @@ import { runDetached } from '@/lib/core/utils/background'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { runTableExport, type TableExportPayload } from '@/lib/table/export-runner'
import { markTableJobRunning, releaseJobClaim } from '@/lib/table/service'
import { markTableJobRunning, releaseJobClaim } from '@/lib/table/jobs/service'
import type { TableExportJobPayload } from '@/lib/table/types'
import { accessError, checkAccess } from '@/app/api/table/utils'
@@ -12,7 +12,7 @@ const { mockCheckAccess, mockGetTableJob, mockGeneratePresignedDownloadUrl } = v
mockGeneratePresignedDownloadUrl: vi.fn(),
}))
vi.mock('@/lib/table/service', () => ({ getTableJob: mockGetTableJob }))
vi.mock('@/lib/table/jobs/service', () => ({ getTableJob: mockGetTableJob }))
vi.mock('@/lib/uploads/core/storage-service', () => ({
generatePresignedDownloadUrl: mockGeneratePresignedDownloadUrl,
}))
@@ -5,7 +5,7 @@ import { parseRequest } from '@/lib/api/server'
import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { getTableJob } from '@/lib/table/service'
import { getTableJob } from '@/lib/table/jobs/service'
import type { TableExportJobPayload } from '@/lib/table/types'
import { generatePresignedDownloadUrl } from '@/lib/uploads/core/storage-service'
import { accessError, checkAccess } from '@/app/api/table/utils'
@@ -20,7 +20,7 @@ vi.mock('@/app/api/table/utils', async () => {
}
})
vi.mock('@/lib/table/service', () => ({
vi.mock('@/lib/table/rows/service', () => ({
queryRows: mockQueryRows,
}))
@@ -7,7 +7,7 @@ import { neutralizeCsvFormula } from '@/lib/core/utils/csv'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { buildNameById, getColumnId, rowDataIdToName } from '@/lib/table/column-keys'
import { queryRows } from '@/lib/table/service'
import { queryRows } from '@/lib/table/rows/service'
import { accessError, checkAccess } from '@/app/api/table/utils'
const logger = createLogger('TableExport')
@@ -9,7 +9,11 @@ import { parseRequest } from '@/lib/api/server'
import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { addWorkflowGroup, deleteWorkflowGroup, updateWorkflowGroup } from '@/lib/table/service'
import {
addWorkflowGroup,
deleteWorkflowGroup,
updateWorkflowGroup,
} from '@/lib/table/workflow-groups/service'
import { accessError, checkAccess, normalizeColumn } from '@/app/api/table/utils'
const logger = createLogger('TableWorkflowGroupsAPI')
@@ -16,7 +16,7 @@ vi.mock('@sim/utils/id', () => ({
generateId: vi.fn().mockReturnValue('import-id-xyz'),
generateShortId: vi.fn().mockReturnValue('short-id'),
}))
vi.mock('@/lib/table/service', () => ({ markTableJobRunning: mockMarkTableImporting }))
vi.mock('@/lib/table/jobs/service', () => ({ markTableJobRunning: mockMarkTableImporting }))
vi.mock('@/lib/table/import-runner', () => ({ runTableImport: mockRunTableImport }))
vi.mock('@/lib/core/utils/background', () => ({
runDetached: (_label: string, work: () => Promise<unknown>) => {
@@ -9,7 +9,7 @@ import { runDetached } from '@/lib/core/utils/background'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { runTableImport, type TableImportPayload } from '@/lib/table/import-runner'
import { markTableJobRunning, releaseJobClaim } from '@/lib/table/service'
import { markTableJobRunning, releaseJobClaim } from '@/lib/table/jobs/service'
import { accessError, checkAccess } from '@/app/api/table/utils'
const logger = createLogger('TableImportIntoAsync')
@@ -45,20 +45,26 @@ vi.mock('@/app/api/table/utils', async () => {
})
/**
* The route imports `importAppendRows` / `importReplaceRows` from the barrel,
* which forwards them from `./service`. These functions own the import
* transaction (column adds + row writes); mocking the service module replaces
* them without touching the other real helpers (`coerceRowsForTable`,
* `createCsvParser`, etc.) exported through the barrel.
* The route imports `importAppendRows` / `importReplaceRows` from
* `@/lib/table/import-data`. These functions own the import transaction (column
* adds + row writes); mocking that module replaces them without touching the
* other real helpers (`coerceRowsForTable`, `createCsvParser`, etc.) exported
* through the barrel.
*/
vi.mock('@/lib/table/service', () => ({
vi.mock('@/lib/table/import-data', () => ({
importAppendRows: mockImportAppendRows,
importReplaceRows: mockImportReplaceRows,
dispatchAfterBatchInsert: mockDispatchAfterBatchInsert,
}))
vi.mock('@/lib/table/jobs/service', () => ({
markTableJobRunning: mockMarkTableImporting,
releaseJobClaim: mockReleaseImportClaim,
}))
vi.mock('@/lib/table/rows/service', () => ({
dispatchAfterBatchInsert: mockDispatchAfterBatchInsert,
}))
import { POST } from '@/app/api/table/[tableId]/import/route'
function createCsvFile(contents: string, name = 'data.csv', type = 'text/csv'): File {
@@ -25,8 +25,6 @@ import {
createCsvParser,
dispatchAfterBatchInsert,
generateColumnId,
importAppendRows,
importReplaceRows,
inferColumnType,
markTableJobRunning,
releaseJobClaim,
@@ -35,6 +33,7 @@ import {
type TableSchema,
validateMapping,
} from '@/lib/table'
import { importAppendRows, importReplaceRows } from '@/lib/table/import-data'
import {
accessError,
checkAccess,
@@ -15,7 +15,7 @@ const { mockCheckAccess, mockMarkJobCanceled, mockGetTableJob, mockAppendTableEv
})
)
vi.mock('@/lib/table/service', () => ({
vi.mock('@/lib/table/jobs/service', () => ({
markJobCanceled: mockMarkJobCanceled,
getTableJob: mockGetTableJob,
}))
@@ -6,7 +6,7 @@ import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { appendTableEvent } from '@/lib/table/events'
import { getTableJob, markJobCanceled } from '@/lib/table/service'
import { getTableJob, markJobCanceled } from '@/lib/table/jobs/service'
import type { TableJobType } from '@/lib/table/types'
import { accessError, checkAccess } from '@/app/api/table/utils'
@@ -23,7 +23,7 @@ vi.mock('@/app/api/table/utils', async () => {
}
})
vi.mock('@/lib/table/service', () => ({
vi.mock('@/lib/table/rows/service', () => ({
findRowMatches: mockFindRowMatches,
}))
@@ -6,7 +6,7 @@ import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import type { Sort } from '@/lib/table'
import { findRowMatches } from '@/lib/table/service'
import { findRowMatches } from '@/lib/table/rows/service'
import { TableQueryValidationError } from '@/lib/table/sql'
import { accessError, checkAccess } from '@/app/api/table/utils'
@@ -40,7 +40,7 @@ vi.mock('@/lib/table', async () => {
}
})
vi.mock('@/lib/table/service', () => ({
vi.mock('@/lib/table/rows/service', () => ({
queryRows: mockQueryRows,
}))
@@ -25,7 +25,7 @@ import {
validateRowData,
validateRowSize,
} from '@/lib/table'
import { queryRows } from '@/lib/table/service'
import { queryRows } from '@/lib/table/rows/service'
import { TableQueryValidationError } from '@/lib/table/sql'
import { rowWireTranslators } from '@/app/api/table/row-wire'
import { accessError, checkAccess, rowWriteErrorResponse } from '@/app/api/table/utils'
@@ -22,9 +22,12 @@ vi.mock('@sim/utils/id', () => ({
// streaming multipart + CSV pipeline is exercised end-to-end.
vi.mock('@/lib/table/service', () => ({
createTable: mockCreateTable,
batchInsertRows: mockBatchInsertRows,
deleteTable: mockDeleteTable,
}))
vi.mock('@/lib/table/rows/service', () => ({
batchInsertRows: mockBatchInsertRows,
}))
vi.mock('@/lib/table/billing', () => ({ getWorkspaceTableLimits: mockGetLimits }))
vi.mock('@/app/api/table/utils', async () => {
const { NextResponse } = await import('next/server')
+1 -1
View File
@@ -5,7 +5,7 @@ import { parseRequest } from '@/lib/api/server'
import { checkSessionOrInternalAuth } from '@/lib/auth/hybrid'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { listWorkspaceExportJobs } from '@/lib/table/service'
import { listWorkspaceExportJobs } from '@/lib/table/jobs/service'
import { checkWorkspaceAccess } from '@/lib/workspaces/permissions/utils'
const logger = createLogger('TableJobsAPI')
@@ -31,7 +31,7 @@ import {
validateRowData,
validateRowSize,
} from '@/lib/table'
import { queryRows } from '@/lib/table/service'
import { queryRows } from '@/lib/table/rows/service'
import { TableQueryValidationError } from '@/lib/table/sql'
import { accessError, checkAccess, rowWriteErrorResponse } from '@/app/api/table/utils'
import {
+3 -2
View File
@@ -50,7 +50,7 @@ export async function executeResumeJob(payload: ResumeExecutionPayload) {
// philosophy). Aborting here also stops the wasted compute the guard alone
// can't prevent. Read the cell's current exec and bail if cancelled.
if (cellContext) {
const { getRowById } = await import('@/lib/table/service')
const { getRowById } = await import('@/lib/table/rows/service')
const cellRow = await getRowById(
cellContext.tableId,
cellContext.rowId,
@@ -323,7 +323,8 @@ async function continueCascadeAfterResume(cellContext: {
workspaceId: string
groupId: string
}): Promise<void> {
const { getTableById, getRowById } = await import('@/lib/table/service')
const { getTableById } = await import('@/lib/table/service')
const { getRowById } = await import('@/lib/table/rows/service')
const { pickNextEligibleGroupForRow } = await import('@/lib/table/workflow-columns')
const { runRowCascadeLoop } = await import('@/background/workflow-column-execution')
@@ -51,7 +51,8 @@ export async function executeWorkflowGroupCellJob(
signal?: AbortSignal
) {
const { tableId, rowId, workspaceId } = payload
const { getTableById, getRowById } = await import('@/lib/table/service')
const { getTableById } = await import('@/lib/table/service')
const { getRowById } = await import('@/lib/table/rows/service')
const { pickNextEligibleGroupForRow } = await import('@/lib/table/workflow-columns')
let currentPayload = payload
@@ -105,7 +106,8 @@ export async function runRowCascadeLoop(
signal?: AbortSignal
): Promise<'blocked' | undefined> {
const { tableId, rowId, workspaceId } = payload
const { getTableById, getRowById } = await import('@/lib/table/service')
const { getTableById } = await import('@/lib/table/service')
const { getRowById } = await import('@/lib/table/rows/service')
const { pickNextEligibleGroupForRow } = await import('@/lib/table/workflow-columns')
let currentGroupId = payload.groupId
@@ -175,7 +177,7 @@ async function runWorkflowAndWriteTerminal(
const requestId = `wfgrp-${executionId}`
return runWithRequestContext({ requestId }, async () => {
const { getRowById } = await import('@/lib/table/service')
const { getRowById } = await import('@/lib/table/rows/service')
const { executeWorkflow } = await import('@/lib/workflows/executor/execute-workflow')
const { loadWorkflowFromNormalizedTables, loadDeployedWorkflowState } = await import(
'@/lib/workflows/persistence/utils'
@@ -258,7 +260,7 @@ async function runWorkflowAndWriteTerminal(
logger.warn(
`Usage limit reached — halting enrichment (table=${tableId} row=${rowId} group=${groupId})`
)
const { updateRow } = await import('@/lib/table/service')
const { updateRow } = await import('@/lib/table/rows/service')
await updateRow(
{ tableId, rowId, data: {}, workspaceId, executionsPatch: { [groupId]: null } },
table,
@@ -525,7 +527,7 @@ async function runWorkflowAndWriteTerminal(
// cell's exec so it reverts to un-run (no error/cancelled badge —
// matching "don't mark"; re-runnable after upgrade). Each blocked
// cell clears its own.
const { updateRow } = await import('@/lib/table/service')
const { updateRow } = await import('@/lib/table/rows/service')
await updateRow(
{ tableId, rowId, data: {}, workspaceId, executionsPatch: { [groupId]: null } },
table,
@@ -12,6 +12,9 @@ const { mockGetTableById, mockReplaceTableRows } = vi.hoisted(() => ({
vi.mock('@/lib/table/service', () => ({
getTableById: mockGetTableById,
}))
vi.mock('@/lib/table/rows/service', () => ({
replaceTableRows: mockReplaceTableRows,
}))
+2 -1
View File
@@ -11,7 +11,8 @@ import { withCopilotSpan } from '@/lib/copilot/request/otel'
import type { ExecutionContext, ToolCallResult } from '@/lib/copilot/request/types'
import type { RowData, TableDefinition } from '@/lib/table'
import { buildIdByName, rowDataNameToId } from '@/lib/table/column-keys'
import { getTableById, replaceTableRows } from '@/lib/table/service'
import { replaceTableRows } from '@/lib/table/rows/service'
import { getTableById } from '@/lib/table/service'
const logger = createLogger('CopilotToolResultTables')
@@ -3,7 +3,8 @@ import { decodeVfsPathSegments, encodeVfsPathSegments } from '@/lib/copilot/vfs/
import { resolveWorkflowAliasForWorkspace } from '@/lib/copilot/vfs/workflow-alias-resolver'
import { isPlanAliasPath, workflowAliasSandboxPath } from '@/lib/copilot/vfs/workflow-aliases'
import { isMothershipBetaFeaturesEnabled } from '@/lib/core/config/feature-flags'
import { getTableById, listTables, queryRows } from '@/lib/table/service'
import { queryRows } from '@/lib/table/rows/service'
import { getTableById, listTables } from '@/lib/table/service'
import { listWorkspaceFileFolders } from '@/lib/uploads/contexts/workspace/workspace-file-folder-manager'
import {
fetchWorkspaceFileBuffer,
@@ -56,26 +56,39 @@ vi.mock('@/enrichments/registry', () => ({
}))
vi.mock('@/lib/table/service', () => ({
addTableColumn: vi.fn(),
addWorkflowGroup: mockAddWorkflowGroup,
batchInsertRows: mockBatchInsertRows,
batchUpdateRows: vi.fn(),
createTable: mockCreateTable,
deleteTable: mockDeleteTable,
getTableById: mockGetTableById,
renameTable: vi.fn(),
}))
vi.mock('@/lib/table/workflow-groups/service', () => ({
addWorkflowGroup: mockAddWorkflowGroup,
addWorkflowGroupOutput: vi.fn(),
deleteWorkflowGroup: vi.fn(),
deleteWorkflowGroupOutput: vi.fn(),
updateWorkflowGroup: vi.fn(),
}))
vi.mock('@/lib/table/columns/service', () => ({
addTableColumn: vi.fn(),
deleteColumn: vi.fn(),
deleteColumns: vi.fn(),
renameColumn: vi.fn(),
updateColumnConstraints: vi.fn(),
updateColumnType: vi.fn(),
}))
vi.mock('@/lib/table/rows/service', () => ({
batchInsertRows: mockBatchInsertRows,
batchUpdateRows: vi.fn(),
deleteRow: vi.fn(),
deleteRowsByFilter: vi.fn(),
deleteRowsByIds: vi.fn(),
deleteTable: mockDeleteTable,
getRowById: vi.fn(),
getTableById: mockGetTableById,
insertRow: vi.fn(),
queryRows: vi.fn(),
renameColumn: vi.fn(),
renameTable: vi.fn(),
replaceTableRows: mockReplaceTableRows,
updateColumnConstraints: vi.fn(),
updateColumnType: vi.fn(),
updateRow: vi.fn(),
updateRowsByFilter: vi.fn(),
}))
@@ -30,32 +30,26 @@ import {
import { columnTypeForLeaf, deriveOutputColumnName } from '@/lib/table/column-naming'
import {
addTableColumn,
addWorkflowGroup,
addWorkflowGroupOutput,
batchInsertRows,
batchUpdateRows,
createTable,
deleteColumn,
deleteColumns,
renameColumn,
updateColumnConstraints,
updateColumnType,
} from '@/lib/table/columns/service'
import {
batchInsertRows,
batchUpdateRows,
deleteRow,
deleteRowsByFilter,
deleteRowsByIds,
deleteTable,
deleteWorkflowGroup,
deleteWorkflowGroupOutput,
getRowById,
getTableById,
insertRow,
queryRows,
renameColumn,
renameTable,
replaceTableRows,
updateColumnConstraints,
updateColumnType,
updateRow,
updateRowsByFilter,
updateWorkflowGroup,
} from '@/lib/table/service'
} from '@/lib/table/rows/service'
import { createTable, deleteTable, getTableById, renameTable } from '@/lib/table/service'
import type {
ColumnDefinition,
RowData,
@@ -67,6 +61,13 @@ import type {
WorkflowGroupOutput,
} from '@/lib/table/types'
import { cancelWorkflowGroupRuns, runWorkflowColumn } from '@/lib/table/workflow-columns'
import {
addWorkflowGroup,
addWorkflowGroupOutput,
deleteWorkflowGroup,
deleteWorkflowGroupOutput,
updateWorkflowGroup,
} from '@/lib/table/workflow-groups/service'
import {
fetchWorkspaceFileBuffer,
resolveWorkspaceFileReference,
@@ -38,7 +38,7 @@ vi.mock('@/lib/table/validation', () => ({
checkBatchUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })),
}))
import { findRowMatches } from '@/lib/table/service'
import { findRowMatches } from '@/lib/table/rows/service'
import { buildFilterClause, buildSortClause } from '@/lib/table/sql'
const COLUMNS: ColumnDefinition[] = [
@@ -10,7 +10,7 @@
import { userTableDefinitions } from '@sim/db/schema'
import { dbChainMock, dbChainMockFns, resetDbChainMock } from '@sim/testing'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import { importAppendRows } from '@/lib/table/service'
import { importAppendRows } from '@/lib/table/import-data'
import type { TableDefinition } from '@/lib/table/types'
vi.mock('@sim/db', () => dbChainMock)
@@ -43,7 +43,7 @@ vi.mock('@/lib/table/validation', () => ({
checkBatchUniqueConstraintsDb: vi.fn(async () => ({ valid: true, errors: [] })),
}))
import { deleteRowsByFilter, queryRows, updateRowsByFilter } from '@/lib/table/service'
import { deleteRowsByFilter, queryRows, updateRowsByFilter } from '@/lib/table/rows/service'
const COLUMNS: ColumnDefinition[] = [
{ name: 'name', type: 'string' },
@@ -3,15 +3,14 @@
*/
import { dbChainMock, dbChainMockFns, resetDbChainMock } from '@sim/testing'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import { deleteColumn, renameColumn } from '@/lib/table/columns/service'
import {
batchInsertRows,
deleteColumn,
insertRow,
renameColumn,
replaceTableRows,
updateRow,
upsertRow,
} from '@/lib/table/service'
} from '@/lib/table/rows/service'
import type { TableDefinition } from '@/lib/table/types'
import { getUniqueColumns } from '@/lib/table/validation'
@@ -2,7 +2,7 @@
* @vitest-environment node
*/
import { describe, expect, it } from 'vitest'
import { TABLE_LIMITS } from '../constants'
import { TABLE_LIMITS } from '@/lib/table/constants'
import {
type ColumnDefinition,
coerceRowToSchema,
@@ -15,7 +15,7 @@ import {
validateTableName,
validateTableSchema,
validateUniqueConstraints,
} from '../validation'
} from '@/lib/table/validation'
describe('Validation', () => {
describe('validateTableName', () => {
+5 -5
View File
@@ -10,20 +10,20 @@ import { MATERIALIZE_CONCURRENCY, mapWithConcurrency } from '@/lib/core/utils/co
import { materializeExecutionData } from '@/lib/logs/execution/trace-store'
import { appendTableEvent } from '@/lib/table/events'
import {
batchUpdateRows,
getTableById,
markJobFailed,
markJobReady,
markTableJobRunning,
updateJobProgress,
} from '@/lib/table/service'
} from '@/lib/table/jobs/service'
import { pluckByPath } from '@/lib/table/pluck'
import { batchUpdateRows } from '@/lib/table/rows/service'
import { getTableById } from '@/lib/table/service'
import type {
RowData,
TableBackfillJobPayload,
TableDefinition,
WorkflowGroupOutput,
} from '@/lib/table/types'
import { pluckByPath } from './pluck'
const logger = createLogger('TableBackfillRunner')
@@ -330,7 +330,7 @@ export async function maybeBackfillGroupOutputs(opts: {
// Release the claim so a ghost `running` job doesn't block imports/deletes.
// Swallowed (warn only): a failed backfill never fails the schema change —
// the data stays backfillable.
const { releaseJobClaim } = await import('./service')
const { releaseJobClaim } = await import('@/lib/table/jobs/service')
await releaseJobClaim(table.id, jobId).catch(() => {})
logger.warn(
`[${requestId}] Backfill dispatch failed for table ${table.id} group ${groupId}; skipping`,
+1 -1
View File
@@ -7,8 +7,8 @@
import { createLogger } from '@sim/logger'
import { getHighestPrioritySubscription } from '@/lib/billing/core/subscription'
import { getPlanTypeForLimits } from '@/lib/billing/plan-helpers'
import { getTablePlanLimits, type PlanName, type TablePlanLimits } from '@/lib/table/constants'
import { getWorkspaceBilledAccountUserId } from '@/lib/workspaces/utils'
import { getTablePlanLimits, type PlanName, type TablePlanLimits } from './constants'
const logger = createLogger('TableBilling')
+2 -1
View File
@@ -44,7 +44,8 @@ export async function writeWorkflowGroupState(
): Promise<'wrote' | 'skipped'> {
const { tableId, rowId, workspaceId, groupId, executionId } = ctx
const requestId = ctx.requestId ?? `wfgrp-${executionId}`
const { getTableById, getRowById, updateRow } = await import('@/lib/table/service')
const { getTableById } = await import('@/lib/table/service')
const { getRowById, updateRow } = await import('@/lib/table/rows/service')
const table = await getTableById(tableId)
if (!table) {
+8 -1
View File
@@ -9,7 +9,14 @@
*/
import { generateId } from '@sim/utils/id'
import type { ColumnDefinition, Filter, RowData, Sort, TableSchema, WorkflowGroup } from './types'
import type {
ColumnDefinition,
Filter,
RowData,
Sort,
TableSchema,
WorkflowGroup,
} from '@/lib/table/types'
/**
* Resolves a column's stable storage key. Falls back to `name` for legacy
+1 -1
View File
@@ -6,7 +6,7 @@
* get from the sidebar.
*/
import type { ColumnDefinition } from './types'
import type { ColumnDefinition } from '@/lib/table/types'
/**
* Slugifies a string into a `NAME_PATTERN`-safe column name. Lowercase,
+668
View File
@@ -0,0 +1,668 @@
/**
* Column and schema-management service for user tables.
*
* Standalone column-mutation operations (add, rename, delete, type change,
* constraint change) extracted from the table service. Each acquires the
* table's advisory lock via {@link withLockedTable} from `@/lib/table/service`.
*
* Use this for: workflow executor, background jobs, testing business logic.
* Use API routes for: HTTP requests, frontend clients.
*/
import { db } from '@sim/db'
import { userTableDefinitions, userTableRows } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { and, count, eq, sql } from 'drizzle-orm'
import { columnMatchesRef, generateColumnId, getColumnId } from '@/lib/table/column-keys'
import { COLUMN_TYPES, NAME_PATTERN, TABLE_LIMITS } from '@/lib/table/constants'
import { stripGroupExecutions } from '@/lib/table/rows/executions'
import { withLockedTable } from '@/lib/table/service'
import { scaledStatementTimeoutMs, setTableTxTimeouts } from '@/lib/table/tx'
import type {
DeleteColumnData,
RenameColumnData,
RowData,
TableDefinition,
TableMetadata,
TableSchema,
UpdateColumnConstraintsData,
UpdateColumnTypeData,
} from '@/lib/table/types'
import { assertValidSchema, stripGroupDeps } from '@/lib/table/workflow-columns'
const logger = createLogger('TableColumnService')
/**
* Adds a column to an existing table's schema.
*
* @param tableId - Table ID to update
* @param column - Column definition to add
* @param requestId - Request ID for logging
* @returns Updated table definition
* @throws Error if table not found or column name already exists
*/
export async function addTableColumn(
tableId: string,
column: {
id?: string
name: string
type: string
required?: boolean
unique?: boolean
position?: number
},
requestId: string
): Promise<TableDefinition> {
return withLockedTable(tableId, async (table, trx) => {
if (!NAME_PATTERN.test(column.name)) {
throw new Error(
`Invalid column name "${column.name}". Must start with a letter or underscore and contain only alphanumeric characters and underscores.`
)
}
if (column.name.length > TABLE_LIMITS.MAX_COLUMN_NAME_LENGTH) {
throw new Error(
`Column name exceeds maximum length (${TABLE_LIMITS.MAX_COLUMN_NAME_LENGTH} characters)`
)
}
if (!COLUMN_TYPES.includes(column.type as (typeof COLUMN_TYPES)[number])) {
throw new Error(
`Invalid column type "${column.type}". Must be one of: ${COLUMN_TYPES.join(', ')}`
)
}
const schema = table.schema
if (schema.columns.some((c) => c.name.toLowerCase() === column.name.toLowerCase())) {
throw new Error(`Column "${column.name}" already exists`)
}
if (schema.columns.length >= TABLE_LIMITS.MAX_COLUMNS_PER_TABLE) {
throw new Error(
`Table has reached maximum column limit (${TABLE_LIMITS.MAX_COLUMNS_PER_TABLE})`
)
}
const newColumn: TableSchema['columns'][number] = {
// Honor a caller-provided id (undo of a delete reuses the original id);
// otherwise mint a fresh one.
id: column.id ?? generateColumnId(),
name: column.name,
type: column.type as TableSchema['columns'][number]['type'],
required: column.required ?? false,
unique: column.unique ?? false,
}
const newColumnId = getColumnId(newColumn)
const columns = [...schema.columns]
if (column.position !== undefined && column.position >= 0 && column.position < columns.length) {
columns.splice(column.position, 0, newColumn)
} else {
columns.push(newColumn)
}
const updatedSchema: TableSchema = { ...schema, columns }
// Keep `metadata.columnOrder` (a list of column ids) in sync: splicing the
// new column's id at the same index we used in `columns` keeps display
// ordering aligned with the user's intent for `position`-based inserts.
const existingOrder = table.metadata?.columnOrder
let updatedMetadata = table.metadata
if (existingOrder && existingOrder.length > 0 && !existingOrder.includes(newColumnId)) {
let insertIdx = existingOrder.length
if (column.position !== undefined && column.position >= 0) {
// Anchor on the column previously at `position` — that column shifted
// right by one in `columns`, so the new id slots in at its old spot.
const anchor = schema.columns[column.position]
if (anchor) {
const anchorIdx = existingOrder.indexOf(getColumnId(anchor))
if (anchorIdx !== -1) insertIdx = anchorIdx
}
}
const nextOrder = [...existingOrder]
nextOrder.splice(insertIdx, 0, newColumnId)
updatedMetadata = { ...table.metadata, columnOrder: nextOrder }
}
assertValidSchema(updatedSchema, updatedMetadata?.columnOrder)
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, tableId))
logger.info(`[${requestId}] Added column "${column.name}" to table ${tableId}`)
return {
...table,
schema: updatedSchema,
metadata: updatedMetadata,
updatedAt: now,
}
})
}
/**
* Renames a column in a table's schema and updates all row data keys.
*
* @param data - Rename column data
* @param requestId - Request ID for logging
* @returns Updated table definition
* @throws Error if table not found, column not found, or new name conflicts
*/
export async function renameColumn(
data: RenameColumnData,
requestId: string
): Promise<TableDefinition> {
return withLockedTable(data.tableId, async (table, trx) => {
if (!NAME_PATTERN.test(data.newName)) {
throw new Error(
`Invalid column name "${data.newName}". Column names must start with a letter or underscore, followed by alphanumeric characters or underscores.`
)
}
if (data.newName.length > TABLE_LIMITS.MAX_COLUMN_NAME_LENGTH) {
throw new Error(
`Column name exceeds maximum length (${TABLE_LIMITS.MAX_COLUMN_NAME_LENGTH} characters)`
)
}
const schema = table.schema
const columnIndex = schema.columns.findIndex((c) => columnMatchesRef(c, data.oldName))
if (columnIndex === -1) {
throw new Error(`Column "${data.oldName}" not found`)
}
if (
schema.columns.some(
(c, i) => i !== columnIndex && c.name.toLowerCase() === data.newName.toLowerCase()
)
) {
throw new Error(`Column "${data.newName}" already exists`)
}
const targetColumn = schema.columns[columnIndex]
const actualOldName = targetColumn.name
// Rename is metadata-only: stored rows, metadata, and workflow-group refs all
// key on the column's stable id, which a rename never changes — so this is a
// pure schema write, no per-row JSONB rewrite or group/metadata cascade.
// Stamp the current storage key as the id (for any not-yet-backfilled column)
// so existing rows stay reachable as the display name changes.
const columnId = targetColumn.id ?? actualOldName
const updatedColumns = schema.columns.map((c, i) =>
i === columnIndex ? { ...c, id: columnId, name: data.newName } : c
)
const updatedSchema: TableSchema = { ...schema, columns: updatedColumns }
assertValidSchema(updatedSchema, table.metadata?.columnOrder)
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
logger.info(
`[${requestId}] Renamed column "${actualOldName}" to "${data.newName}" in table ${data.tableId}`
)
return { ...table, schema: updatedSchema, updatedAt: now }
})
}
/** Removes the given column-id keys from a metadata blob (widths/order/pinned). */
function stripColumnIdsFromMetadata(
metadata: TableMetadata | null,
ids: ReadonlySet<string>
): TableMetadata | null {
if (!metadata) return metadata
let next = metadata
if (metadata.columnWidths) {
const widths = { ...metadata.columnWidths }
let changed = false
for (const id of ids)
if (id in widths) {
delete widths[id]
changed = true
}
if (changed) next = { ...next, columnWidths: widths }
}
if (metadata.columnOrder?.some((id) => ids.has(id))) {
next = { ...next, columnOrder: metadata.columnOrder.filter((id) => !ids.has(id)) }
}
if (metadata.pinnedColumns?.some((id) => ids.has(id))) {
next = { ...next, pinnedColumns: metadata.pinnedColumns.filter((id) => !ids.has(id)) }
}
return next
}
/**
* Fire-and-forget reclamation of a deleted column's row storage. The column is
* already gone from the schema, so reads never surface the orphaned id —
* dropping the JSONB key just frees space. Runs in its own transaction with a
* row-count-scaled timeout; failures are logged, not propagated.
*/
function stripColumnDataInBackground(
tableId: string,
columnIds: string[],
rowCount: number,
requestId: string
): void {
if (columnIds.length === 0) return
void (async () => {
try {
await db.transaction(async (trx) => {
const statementMs = scaledStatementTimeoutMs(rowCount, {
baseMs: 60_000,
perRowMs: 2 * columnIds.length,
})
await setTableTxTimeouts(trx, { statementMs })
for (const id of columnIds) {
await trx.execute(
sql`UPDATE user_table_rows SET data = data - ${id}::text WHERE table_id = ${tableId} AND data ? ${id}::text`
)
}
})
logger.info(
`[${requestId}] Background-stripped deleted column data [${columnIds.join(', ')}] from table ${tableId}`
)
} catch (err) {
logger.error(
`[${requestId}] Background column-data strip failed for table ${tableId} [${columnIds.join(', ')}]:`,
err
)
}
})()
}
/**
* Deletes a column from a table's schema. When id-keyed, returns once the schema
* is updated and reclaims the column's row-data storage in the background
* (fire-and-forget); the legacy path strips the row key synchronously.
*
* @param data - Delete column data
* @param requestId - Request ID for logging
* @returns Updated table definition
* @throws Error if table not found, column not found, or it's the last column
*/
export async function deleteColumn(
data: DeleteColumnData,
requestId: string
): Promise<TableDefinition> {
const { def, stripKey } = await withLockedTable(data.tableId, async (table, trx) => {
const schema = table.schema
const columnIndex = schema.columns.findIndex((c) => columnMatchesRef(c, data.columnName))
if (columnIndex === -1) {
throw new Error(`Column "${data.columnName}" not found`)
}
if (schema.columns.length <= 1) {
throw new Error('Cannot delete the last column in a table')
}
const targetColumn = schema.columns[columnIndex]
const actualName = targetColumn.name
const columnId = getColumnId(targetColumn)
const ownerGroupId = targetColumn.workflowGroupId
// Drop this column's reference (by id) from every group's outputs and
// `columns` dependency. If the column is the last output of its parent
// group, the group itself is also removed (a group with zero outputs is
// invalid).
let groupRemovedId: string | null = null
const updatedGroups = (schema.workflowGroups ?? [])
.map((group) => {
let next = group
if (ownerGroupId && group.id === ownerGroupId) {
const remaining = group.outputs.filter((o) => o.columnName !== columnId)
if (remaining.length === 0) {
groupRemovedId = group.id
}
next = { ...next, outputs: remaining }
}
return stripGroupDeps(next, new Set([columnId]))
})
.filter((g) => g.id !== groupRemovedId)
const updatedSchema: TableSchema = {
...schema,
columns: schema.columns.filter((_, i) => i !== columnIndex),
...(updatedGroups.length > 0 ? { workflowGroups: updatedGroups } : {}),
}
const updatedMetadata = stripColumnIdsFromMetadata(
table.metadata as TableMetadata | null,
new Set([columnId])
)
assertValidSchema(updatedSchema, updatedMetadata?.columnOrder)
const now = new Date()
// Schema/metadata update commits now; the column's row-data storage is
// reclaimed in the background (fire-and-forget) — reads never surface the
// orphaned id since the column is already gone from the schema.
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
if (groupRemovedId) await stripGroupExecutions(trx, data.tableId, [groupRemovedId])
logger.info(`[${requestId}] Deleted column "${actualName}" from table ${data.tableId}`)
return {
def: { ...table, schema: updatedSchema, metadata: updatedMetadata, updatedAt: now },
stripKey: columnId,
}
})
stripColumnDataInBackground(data.tableId, [stripKey], def.rowCount ?? 0, requestId)
return def
}
/**
* Deletes multiple columns from a table in a single transaction.
* Avoids the race condition of calling deleteColumn multiple times in parallel.
*/
export async function deleteColumns(
data: { tableId: string; columnNames: string[] },
requestId: string
): Promise<TableDefinition> {
const { def, stripKeys } = await withLockedTable(data.tableId, async (table, trx) => {
const schema = table.schema
const namesToDelete = new Set<string>()
const idsToDelete = new Set<string>()
const notFound: string[] = []
for (const name of data.columnNames) {
const col = schema.columns.find((c) => columnMatchesRef(c, name))
if (!col) {
notFound.push(name)
} else {
namesToDelete.add(col.name)
idsToDelete.add(getColumnId(col))
}
}
if (notFound.length > 0) {
throw new Error(`Columns not found: ${notFound.join(', ')}`)
}
const remaining = schema.columns.filter((c) => !namesToDelete.has(c.name))
if (remaining.length === 0) {
throw new Error('Cannot delete all columns from a table')
}
// For each group, drop outputs whose column (by id) is being deleted. Groups
// that end up with zero outputs are removed entirely (they'd be invalid).
// Then any remaining group's dependencies referencing a removed column are
// cleaned up.
const removedGroupIds = new Set<string>()
let updatedGroups = (schema.workflowGroups ?? []).map((group) => {
const remainingOutputs = group.outputs.filter((o) => !idsToDelete.has(o.columnName))
if (remainingOutputs.length === 0) {
removedGroupIds.add(group.id)
}
return remainingOutputs.length === group.outputs.length
? group
: { ...group, outputs: remainingOutputs }
})
updatedGroups = updatedGroups
.filter((g) => !removedGroupIds.has(g.id))
.map((group) => stripGroupDeps(group, idsToDelete))
const updatedSchema: TableSchema = {
...schema,
columns: remaining,
...(updatedGroups.length > 0 ? { workflowGroups: updatedGroups } : {}),
}
const updatedMetadata = stripColumnIdsFromMetadata(
table.metadata as TableMetadata | null,
idsToDelete
)
assertValidSchema(updatedSchema, updatedMetadata?.columnOrder)
const now = new Date()
// Schema/metadata commit now; row storage for the deleted columns is
// reclaimed in the background (fire-and-forget).
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
await stripGroupExecutions(trx, data.tableId, removedGroupIds)
logger.info(
`[${requestId}] Deleted columns [${[...namesToDelete].join(', ')}] from table ${data.tableId}`
)
return {
def: { ...table, schema: updatedSchema, metadata: updatedMetadata, updatedAt: now },
stripKeys: Array.from(idsToDelete),
}
})
if (stripKeys.length > 0) {
stripColumnDataInBackground(data.tableId, stripKeys, def.rowCount ?? 0, requestId)
}
return def
}
/**
* Changes the type of a column. Validates that existing data is compatible.
*
* @param data - Update column type data
* @param requestId - Request ID for logging
* @returns Updated table definition
* @throws Error if table not found, column not found, or existing data is incompatible
*/
export async function updateColumnType(
data: UpdateColumnTypeData,
requestId: string
): Promise<TableDefinition> {
return withLockedTable(data.tableId, async (table, trx) => {
// Scale both statement and idle timeouts to row count: the compatibility
// check below iterates every row in Node between the row SELECT and the
// schema UPDATE, leaving the transaction idle for that gap. The default 5s
// `idle_in_transaction_session_timeout` would abort a valid type change on
// a large table.
const timeoutMs = scaledStatementTimeoutMs(table.rowCount ?? 0, {
baseMs: 60_000,
perRowMs: 2,
})
await setTableTxTimeouts(trx, { statementMs: timeoutMs, idleMs: timeoutMs })
if (!(COLUMN_TYPES as readonly string[]).includes(data.newType)) {
throw new Error(
`Invalid column type "${data.newType}". Valid types: ${COLUMN_TYPES.join(', ')}`
)
}
const schema = table.schema
const columnIndex = schema.columns.findIndex((c) => columnMatchesRef(c, data.columnName))
if (columnIndex === -1) {
throw new Error(`Column "${data.columnName}" not found`)
}
const column = schema.columns[columnIndex]
if (column.type === data.newType) {
return table
}
const columnKey = getColumnId(column)
// Validate existing data is compatible with the new type
const rows = await trx
.select({ id: userTableRows.id, data: userTableRows.data })
.from(userTableRows)
.where(
and(
eq(userTableRows.tableId, data.tableId),
sql`${userTableRows.data} ? ${columnKey}`,
sql`${userTableRows.data}->>${columnKey}::text IS NOT NULL`
)
)
let incompatibleCount = 0
for (const row of rows) {
const rowData = row.data as RowData
const value = rowData[columnKey]
if (value === null || value === undefined) continue
if (!isValueCompatibleWithType(value, data.newType)) {
incompatibleCount++
}
}
if (incompatibleCount > 0) {
throw new Error(
`Cannot change column "${column.name}" to type "${data.newType}": ${incompatibleCount} row(s) have incompatible values. Fix or remove the incompatible values first.`
)
}
const updatedColumns = schema.columns.map((c, i) =>
i === columnIndex ? { ...c, type: data.newType } : c
)
const updatedSchema: TableSchema = { ...schema, columns: updatedColumns }
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
logger.info(
`[${requestId}] Changed column "${column.name}" type from "${column.type}" to "${data.newType}" in table ${data.tableId}`
)
return { ...table, schema: updatedSchema, updatedAt: now }
})
}
/**
* Updates constraints (required, unique) on a column.
*
* @param data - Update column constraints data
* @param requestId - Request ID for logging
* @returns Updated table definition
* @throws Error if table not found, column not found, or existing data violates the constraint
*/
export async function updateColumnConstraints(
data: UpdateColumnConstraintsData,
requestId: string
): Promise<TableDefinition> {
return withLockedTable(data.tableId, async (table, trx) => {
// Scale both statement and idle timeouts to row count: the required/unique
// validation runs between separate queries inside this transaction, leaving
// it briefly idle. Match `updateColumnType` so the default 5s
// `idle_in_transaction_session_timeout` can't abort a valid change on a
// large table.
const timeoutMs = scaledStatementTimeoutMs(table.rowCount ?? 0, {
baseMs: 60_000,
perRowMs: 2,
})
await setTableTxTimeouts(trx, { statementMs: timeoutMs, idleMs: timeoutMs })
const schema = table.schema
const columnIndex = schema.columns.findIndex((c) => columnMatchesRef(c, data.columnName))
if (columnIndex === -1) {
throw new Error(`Column "${data.columnName}" not found`)
}
const column = schema.columns[columnIndex]
const columnKey = getColumnId(column)
if (column.workflowGroupId) {
throw new Error(
`Cannot change constraints on workflow-output column "${column.name}". Constraints aren't applicable to columns whose values come from workflow execution.`
)
}
if (data.required === true && !column.required) {
const [result] = await trx
.select({ count: count() })
.from(userTableRows)
.where(
and(
eq(userTableRows.tableId, data.tableId),
sql`(NOT (${userTableRows.data} ? ${columnKey}) OR ${userTableRows.data}->>${columnKey}::text IS NULL)`
)
)
if (result.count > 0) {
throw new Error(
`Cannot set column "${column.name}" as required: ${result.count} row(s) have null or missing values`
)
}
}
if (data.unique === true && !column.unique) {
const duplicates = (await trx.execute(
sql`SELECT ${userTableRows.data}->>${columnKey}::text AS val, count(*) AS cnt FROM ${userTableRows} WHERE table_id = ${data.tableId} AND ${userTableRows.data} ? ${columnKey} AND ${userTableRows.data}->>${columnKey}::text IS NOT NULL GROUP BY val HAVING count(*) > 1 LIMIT 1`
)) as { val: string; cnt: number }[]
if (duplicates.length > 0) {
throw new Error(`Cannot set column "${column.name}" as unique: duplicate values exist`)
}
}
const updatedColumns = schema.columns.map((c, i) =>
i === columnIndex
? {
...c,
...(data.required !== undefined ? { required: data.required } : {}),
...(data.unique !== undefined ? { unique: data.unique } : {}),
}
: c
)
const updatedSchema: TableSchema = { ...schema, columns: updatedColumns }
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
logger.info(
`[${requestId}] Updated constraints for column "${column.name}" in table ${data.tableId}`
)
return { ...table, schema: updatedSchema, updatedAt: now }
})
}
/**
* Checks if a value is compatible with a target column type.
*/
function isValueCompatibleWithType(
value: unknown,
targetType: (typeof COLUMN_TYPES)[number]
): boolean {
if (value === null || value === undefined) return true
switch (targetType) {
case 'string':
return true
case 'number': {
if (typeof value === 'number') return Number.isFinite(value)
if (typeof value === 'string') {
const num = Number(value)
return Number.isFinite(num) && value.trim() !== ''
}
return false
}
case 'boolean': {
if (typeof value === 'boolean') return true
if (typeof value === 'string')
return ['true', 'false', '1', '0'].includes(value.toLowerCase())
if (typeof value === 'number') return value === 0 || value === 1
return false
}
case 'date': {
if (value instanceof Date) return !Number.isNaN(value.getTime())
if (typeof value === 'string') return !Number.isNaN(Date.parse(value))
return false
}
case 'json':
return true
default:
return false
}
}
+6 -2
View File
@@ -27,13 +27,17 @@ const {
vi.mock('@/lib/table/service', () => ({
getTableById: mockGetTableById,
}))
vi.mock('@/lib/table/jobs/service', () => ({
getJobProgress: mockGetJobProgress,
selectRowIdPage: mockSelectRowIdPage,
deletePageByIds: mockDeletePageByIds,
updateJobProgress: mockUpdateJobProgress,
markJobReady: mockMarkJobReady,
markJobFailed: mockMarkJobFailed,
}))
vi.mock('@/lib/table/rows/ordering', () => ({
selectRowIdPage: mockSelectRowIdPage,
deletePageByIds: mockDeletePageByIds,
}))
vi.mock('@/lib/table/events', () => ({ appendTableEvent: mockAppendTableEvent }))
vi.mock('@/lib/table/sql', () => ({ buildFilterClause: mockBuildFilterClause }))
vi.mock('@/lib/table/constants', () => ({
+3 -4
View File
@@ -6,14 +6,13 @@ import type { Filter } from '@/lib/table'
import { TABLE_LIMITS, USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants'
import { appendTableEvent } from '@/lib/table/events'
import {
deletePageByIds,
getJobProgress,
getTableById,
markJobFailed,
markJobReady,
selectRowIdPage,
updateJobProgress,
} from '@/lib/table/service'
} from '@/lib/table/jobs/service'
import { deletePageByIds, selectRowIdPage } from '@/lib/table/rows/ordering'
import { getTableById } from '@/lib/table/service'
import { buildFilterClause } from '@/lib/table/sql'
const logger = createLogger('TableDeleteRunner')
+7 -1
View File
@@ -5,7 +5,13 @@
*/
import { createLogger } from '@sim/logger'
import type { RowData, RowExecutionMetadata, RowExecutions, TableRow, WorkflowGroup } from './types'
import type {
RowData,
RowExecutionMetadata,
RowExecutions,
TableRow,
WorkflowGroup,
} from '@/lib/table/types'
const logger = createLogger('OptimisticCascade')
+2 -2
View File
@@ -30,7 +30,7 @@ import {
TABLE_CONCURRENCY_LIMIT,
toTableRow,
type WorkflowGroupCellPayload,
} from './workflow-columns'
} from '@/lib/table/workflow-columns'
const logger = createLogger('TableRunDispatcher')
@@ -375,7 +375,7 @@ export async function dispatcherStep(dispatchId: string): Promise<DispatcherStep
}
if (dispatch.status === 'cancelled' || dispatch.status === 'complete') return 'done'
const { getTableById } = await import('./service')
const { getTableById } = await import('@/lib/table/service')
const table = await getTableById(dispatch.tableId)
if (!table) {
logger.warn(`[${dispatchId}] table ${dispatch.tableId} missing — completing dispatch`)
+2
View File
@@ -27,6 +27,8 @@ const {
vi.mock('@/lib/table/service', () => ({
getTableById: mockGetTableById,
}))
vi.mock('@/lib/table/jobs/service', () => ({
selectExportRowPage: mockSelectExportRowPage,
updateJobProgress: mockUpdateJobProgress,
markJobReady: mockMarkJobReady,
+2 -2
View File
@@ -10,13 +10,13 @@ import {
toCsvRow,
} from '@/lib/table/export-format'
import {
getTableById,
markJobFailed,
markJobReady,
selectExportRowPage,
setJobResultKey,
updateJobProgress,
} from '@/lib/table/service'
} from '@/lib/table/jobs/service'
import { getTableById } from '@/lib/table/service'
import { deleteFile, uploadFile } from '@/lib/uploads/core/storage-service'
const logger = createLogger('TableExportRunner')
+1 -1
View File
@@ -1 +1 @@
export * from './use-table-columns'
export * from '@/lib/table/hooks/use-table-columns'
@@ -1,7 +1,7 @@
import { useMemo } from 'react'
import type { ColumnOption } from '@/lib/table/types'
import { useTable } from '@/hooks/queries/tables'
import { useWorkflowRegistry } from '@/stores/workflows/registry/store'
import type { ColumnOption } from '../types'
interface UseTableColumnsOptions {
tableId: string | null | undefined
+200
View File
@@ -0,0 +1,200 @@
/**
* Import-job table-data write operations — bulk insert, schema setup, and
* append/replace used by `import-runner.ts` and the import route. Distinct from
* `import.ts` (CSV parsing) and `import-runner.ts` (the job runner).
*/
import { db } from '@sim/db'
import { userTableDefinitions, userTableRows } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { generateId } from '@sim/utils/id'
import { eq } from 'drizzle-orm'
import { CSV_MAX_BATCH_SIZE } from '@/lib/table/import'
import { nKeysBetween } from '@/lib/table/order-key'
import { acquireRowOrderLock } from '@/lib/table/rows/ordering'
import { batchInsertRowsWithTx, replaceTableRowsWithTx } from '@/lib/table/rows/service'
import { addTableColumnsWithTx } from '@/lib/table/service'
import type {
ReplaceRowsResult,
RowData,
TableDefinition,
TableRow,
TableSchema,
} from '@/lib/table/types'
import {
checkBatchUniqueConstraintsDb,
coerceRowToSchema,
getUniqueColumns,
validateRowSize,
} from '@/lib/table/validation'
const logger = createLogger('TableImportData')
/** One batch of rows for a background import (see {@link bulkInsertImportBatch}). */
export interface BulkImportBatch {
tableId: string
workspaceId: string
userId?: string
rows: RowData[]
/** Position of the first row in this batch; rows get contiguous positions from here. */
startPosition: number
/** Previous batch's last `order_key` (the append anchor); null for the first batch / empty table. */
afterOrderKey?: string | null
}
/**
* Inserts one batch of rows for an async import in a single committed statement.
*
* Differs from {@link batchInsertRowsWithTx} for the bulk-load case: caller-supplied
* contiguous positions (no `acquireTablePositionLock` / `nextAutoPosition` scan — an
* import owns its hidden table as the sole writer), no `RETURNING`, and **no
* `fireTableTrigger` / `runWorkflowColumn`** (a 1M-row import must not dispatch a
* workflow run per row). `row_count` is maintained set-based by the statement-level
* trigger. There is no surrounding transaction and no rollback: each batch commits on
* its own, so committed batches persist even if a later batch fails.
*
* Throws on row-size/schema/unique violations or if the statement-level trigger rejects
* the batch for crossing `max_rows`; the caller marks the import failed.
*/
export async function bulkInsertImportBatch(
data: BulkImportBatch,
table: TableDefinition,
requestId: string
): Promise<{ inserted: number; lastOrderKey: string | null }> {
for (let i = 0; i < data.rows.length; i++) {
const sizeValidation = validateRowSize(data.rows[i])
if (!sizeValidation.valid) {
throw new Error(`Row ${i + 1}: ${sizeValidation.errors.join(', ')}`)
}
const schemaValidation = coerceRowToSchema(data.rows[i], table.schema)
if (!schemaValidation.valid) {
throw new Error(`Row ${i + 1}: ${schemaValidation.errors.join(', ')}`)
}
}
const uniqueColumns = getUniqueColumns(table.schema)
if (uniqueColumns.length > 0) {
const uniqueResult = await checkBatchUniqueConstraintsDb(
data.tableId,
data.rows,
table.schema,
db
)
if (!uniqueResult.valid) {
throw new Error(
uniqueResult.errors.map((e) => `Row ${e.row + 1}: ${e.errors.join(', ')}`).join('; ')
)
}
}
const now = new Date()
// Import worker is the table's sole writer; append keys after the anchor the caller threads
// from the previous batch's last key — no per-batch max(order_key) scan over a growing table.
const orderKeys = nKeysBetween(data.afterOrderKey ?? null, null, data.rows.length)
const rowsToInsert = data.rows.map((rowData, i) => ({
id: `row_${generateId().replace(/-/g, '')}`,
tableId: data.tableId,
workspaceId: data.workspaceId,
data: rowData,
position: data.startPosition + i,
orderKey: orderKeys[i],
createdAt: now,
updatedAt: now,
...(data.userId ? { createdBy: data.userId } : {}),
}))
await db.insert(userTableRows).values(rowsToInsert)
logger.info(`[${requestId}] Bulk-imported ${rowsToInsert.length} rows into table ${data.tableId}`)
return {
inserted: rowsToInsert.length,
lastOrderKey: orderKeys[orderKeys.length - 1] ?? data.afterOrderKey ?? null,
}
}
/** Deletes every row of a table (set-based; the statement-level trigger zeroes `row_count`). */
export async function deleteAllTableRows(tableId: string): Promise<void> {
await db.delete(userTableRows).where(eq(userTableRows.tableId, tableId))
}
/**
* Adds columns to a table during an import (the `createColumns` flow), wrapping the
* tx-bound {@link addTableColumnsWithTx} in its own transaction. Returns the updated table.
*/
export async function addImportColumns(
table: TableDefinition,
additions: { name: string; type: string }[],
requestId: string
): Promise<TableDefinition> {
return db.transaction((trx) => addTableColumnsWithTx(trx, table, additions, requestId))
}
/** Overwrites a table's schema during an import (used when inferring columns from the file). */
export async function setTableSchemaForImport(tableId: string, schema: TableSchema): Promise<void> {
await db
.update(userTableDefinitions)
.set({ schema, updatedAt: new Date() })
.where(eq(userTableDefinitions.id, tableId))
}
/**
* Owns the append-import transaction so the API route never holds a `trx`:
* optionally creates the new columns, then inserts every row in CSV-sized
* batches — all atomic. Caller fires {@link dispatchAfterBatchInsert} after this
* resolves (post-commit), mirroring the other batch-insert sites.
*/
export async function importAppendRows(
table: TableDefinition,
additions: { id?: string; name: string; type: string; required?: boolean; unique?: boolean }[],
rows: RowData[],
ctx: { workspaceId: string; userId?: string; requestId: string }
): Promise<{ inserted: TableRow[]; table: TableDefinition }> {
return db.transaction(async (trx) => {
let working = table
if (additions.length > 0) {
// Take the row-order lock before creating columns so this path uses the
// same rows_pos → user_table_definitions order as plain inserts. Creating
// columns first would lock the definition row before rows_pos, inverting
// the order and deadlocking concurrent inserts on this table. The lock is
// re-entrant, so the per-batch acquire below is a no-op.
await acquireRowOrderLock(trx, table.id)
working = await addTableColumnsWithTx(trx, table, additions, ctx.requestId)
}
const inserted: TableRow[] = []
for (let i = 0; i < rows.length; i += CSV_MAX_BATCH_SIZE) {
const batch = rows.slice(i, i + CSV_MAX_BATCH_SIZE)
const batchInserted = await batchInsertRowsWithTx(
trx,
{ tableId: working.id, rows: batch, workspaceId: ctx.workspaceId, userId: ctx.userId },
working,
generateId().slice(0, 8)
)
inserted.push(...batchInserted)
}
return { inserted, table: working }
})
}
/**
* Owns the replace-import transaction: optionally creates the new columns, then
* replaces all rows — atomically. Keeps `trx` out of the API route.
*/
export async function importReplaceRows(
table: TableDefinition,
additions: { id?: string; name: string; type: string; required?: boolean; unique?: boolean }[],
data: { rows: RowData[]; workspaceId: string; userId?: string },
requestId: string
): Promise<ReplaceRowsResult> {
return db.transaction(async (trx) => {
let working = table
if (additions.length > 0) {
await acquireRowOrderLock(trx, table.id)
working = await addTableColumnsWithTx(trx, table, additions, requestId)
}
return replaceTableRowsWithTx(
trx,
{ tableId: working.id, rows: data.rows, workspaceId: data.workspaceId, userId: data.userId },
working,
requestId
)
})
}
+4 -7
View File
@@ -23,14 +23,11 @@ import {
addImportColumns,
bulkInsertImportBatch,
deleteAllTableRows,
getTableById,
markJobFailed,
markJobReady,
nextImportStartOrderKey,
nextImportStartPosition,
setTableSchemaForImport,
updateJobProgress,
} from '@/lib/table/service'
} from '@/lib/table/import-data'
import { markJobFailed, markJobReady, updateJobProgress } from '@/lib/table/jobs/service'
import { nextImportStartOrderKey, nextImportStartPosition } from '@/lib/table/rows/ordering'
import { getTableById } from '@/lib/table/service'
import { deleteFile, downloadFileStream, headObject } from '@/lib/uploads/core/storage-service'
import { normalizeColumn } from '@/app/api/table/utils'
+15 -10
View File
@@ -5,13 +5,18 @@
* Import hooks directly from '@/lib/table/hooks' in client components.
*/
export * from './billing'
export * from './column-keys'
export * from './constants'
export * from './import'
export * from './llm'
export * from './query-builder'
export * from './service'
export * from './sql'
export * from './types'
export * from './validation'
export * from '@/lib/table/billing'
export * from '@/lib/table/column-keys'
export * from '@/lib/table/columns/service'
export * from '@/lib/table/constants'
export * from '@/lib/table/import'
export * from '@/lib/table/import-data'
export * from '@/lib/table/jobs/service'
export * from '@/lib/table/llm'
export * from '@/lib/table/query-builder'
export * from '@/lib/table/rows/service'
export * from '@/lib/table/service'
export * from '@/lib/table/sql'
export * from '@/lib/table/types'
export * from '@/lib/table/validation'
export * from '@/lib/table/workflow-groups/service'
+383
View File
@@ -0,0 +1,383 @@
/**
* Table background-job service for user tables.
*
* The `table_jobs` state machine (claim / progress / terminal transitions), the
* latest-job reads that enrich a {@link TableDefinition}, and the export-job read
* paths — extracted from the table service. Operates purely on the `table_jobs`
* table (plus `selectExportRowPage`, which pages rows through the shared
* `pendingDeleteMask`), so it never imports the table-root service.
*
* Use this for: workflow executor, background jobs, testing business logic.
* Use API routes for: HTTP requests, frontend clients.
*/
import { db } from '@sim/db'
import { tableJobs, userTableDefinitions, userTableRows } from '@sim/db/schema'
import { and, asc, desc, eq, gt, inArray, ne, or, sql } from 'drizzle-orm'
import type { DbOrTx } from '@/lib/db/types'
import { pendingDeleteMask } from '@/lib/table/rows/service'
import type {
RowData,
TableDefinition,
TableDeleteJobPayload,
TableExportJobPayload,
TableJobType,
} from '@/lib/table/types'
/** Job fields projected onto a {@link TableDefinition}, derived from its latest `table_jobs` row. */
interface DerivedJobFields {
jobStatus: TableDefinition['jobStatus']
jobId: string | null
jobType: TableDefinition['jobType']
jobError: string | null
jobRowsProcessed: number
/**
* Rows a running delete job still has to remove (its doomed estimate minus
* deletions so far). Internal to count adjustment — callers subtract it from
* the raw `row_count` so list/detail counts match the read path's delete
* mask (a mid-delete refresh must not resurrect the count). Not on the wire.
*/
pendingDeleteRemaining: number
}
export const EMPTY_JOB_FIELDS: DerivedJobFields = {
jobStatus: null,
jobId: null,
jobType: null,
jobError: null,
jobRowsProcessed: 0,
pendingDeleteRemaining: 0,
}
function mapJobRow(
row:
| {
id: string
type: string
status: string
rowsProcessed: number
error: string | null
payload: unknown
}
| undefined
): DerivedJobFields {
if (!row) return EMPTY_JOB_FIELDS
const doomedCount =
row.type === 'delete' && row.status === 'running'
? ((row.payload as TableDeleteJobPayload | null)?.doomedCount ?? 0)
: 0
return {
jobStatus: row.status as TableDefinition['jobStatus'],
jobId: row.id,
jobType: row.type as TableDefinition['jobType'],
jobError: row.error,
jobRowsProcessed: row.rowsProcessed,
pendingDeleteRemaining: Math.max(0, doomedCount - row.rowsProcessed),
}
}
const JOB_PROJECTION = {
id: tableJobs.id,
type: tableJobs.type,
status: tableJobs.status,
rowsProcessed: tableJobs.rowsProcessed,
error: tableJobs.error,
payload: tableJobs.payload,
} as const
/**
* The latest job for one table (the running one if present, else the most recent terminal).
* Exports are excluded: they're read-only, run concurrently with other jobs, and have their own
* client surface — surfacing one here would clobber the import/delete/backfill status the tray
* and SSE consumer derive from these fields.
*/
export async function latestJobForTable(
tableId: string,
executor: DbOrTx = db
): Promise<DerivedJobFields> {
const [row] = await executor
.select(JOB_PROJECTION)
.from(tableJobs)
.where(and(eq(tableJobs.tableId, tableId), ne(tableJobs.type, 'export')))
.orderBy(desc(tableJobs.startedAt))
.limit(1)
return mapJobRow(row)
}
/** Latest non-export job per table for a batch of ids, via `DISTINCT ON (table_id)`. */
export async function latestJobsForTables(
tableIds: string[]
): Promise<Map<string, DerivedJobFields>> {
const map = new Map<string, DerivedJobFields>()
if (tableIds.length === 0) return map
const rows = await db
.selectDistinctOn([tableJobs.tableId], { tableId: tableJobs.tableId, ...JOB_PROJECTION })
.from(tableJobs)
.where(and(inArray(tableJobs.tableId, tableIds), ne(tableJobs.type, 'export')))
.orderBy(tableJobs.tableId, desc(tableJobs.startedAt))
for (const row of rows) map.set(row.tableId, mapJobRow(row))
return map
}
/**
* Atomically claims a table's single background-job slot by inserting a `running` row into
* `table_jobs`. The partial-unique index on `table_id WHERE status = 'running'` is the
* concurrency gate: a second insert while a job runs hits `ON CONFLICT DO NOTHING` and returns no
* row, so import and delete (and two imports) are mutually exclusive for free. Returns whether it
* claimed the slot; the caller returns 409 when it didn't.
*/
export async function markTableJobRunning(
tableId: string,
jobId: string,
type: TableJobType,
/** Type-specific scope persisted to `table_jobs.payload` (e.g. {@link TableDeleteJobPayload})
* so read paths can mask the job's effect while it runs. */
payload?: unknown
): Promise<boolean> {
// workspace_id is immutable; the atomic gate is the INSERT's conflict, not this read.
const [def] = await db
.select({ workspaceId: userTableDefinitions.workspaceId })
.from(userTableDefinitions)
.where(eq(userTableDefinitions.id, tableId))
.limit(1)
if (!def) return false
const inserted = await db
.insert(tableJobs)
.values({
id: jobId,
tableId,
workspaceId: def.workspaceId,
type,
status: 'running',
payload: payload ?? null,
})
.onConflictDoNothing()
.returning({ id: tableJobs.id })
return inserted.length > 0
}
/**
* Releases a claim taken by {@link markTableJobRunning} for a synchronous job — deletes the
* transient claim row. Scoped to `jobId` + still-running so it only clears its own claim, never a
* newer run. A sync route claims, writes, then releases here in a `finally`.
*/
export async function releaseJobClaim(tableId: string, jobId: string): Promise<void> {
await db
.delete(tableJobs)
.where(
and(eq(tableJobs.id, jobId), eq(tableJobs.tableId, tableId), eq(tableJobs.status, 'running'))
)
}
/**
* Records job progress (rows processed so far) and bumps `updated_at` so the stale-job janitor
* (`cleanup-stale-executions`) sees a live heartbeat.
*
* Scoped to `jobId` AND `status = 'running'`: a stale/superseded worker no longer matches (its
* write is a no-op), and once the job is terminal (e.g. canceled) the match fails too — so this
* returning `false` is the worker's signal to stop. Returns whether this worker still owns an
* in-flight job.
*/
export async function updateJobProgress(
tableId: string,
rowsProcessed: number,
jobId: string
): Promise<boolean> {
const updated = await db
.update(tableJobs)
.set({ rowsProcessed, updatedAt: new Date() })
.where(ownsActiveJob(tableId, jobId))
.returning({ id: tableJobs.id })
return updated.length > 0
}
/**
* Reads the persisted progress of an in-flight job this worker still owns (`null` when the job
* was canceled/superseded). A retried run seeds its counter from this so progress stays
* cumulative — earlier attempts' batches are already committed, and restarting from zero would
* clobber `rows_processed` (and every count derived from it) with the retry's smaller number.
*/
export async function getJobProgress(tableId: string, jobId: string): Promise<number | null> {
const [job] = await db
.select({ rowsProcessed: tableJobs.rowsProcessed })
.from(tableJobs)
.where(ownsActiveJob(tableId, jobId))
.limit(1)
return job ? job.rowsProcessed : null
}
/**
* One keyset page of rows for the export worker, ordered by `(position, id)`. Keyset (not
* OFFSET) keeps each page O(page) — offset paging re-scans every prior row per page, which is
* O(N²) across a large export. `(position, id)` is total (position exists on every row; id breaks
* ties) and served by the `(table_id, position)` index; under fractional ordering a manually
* reordered table may export in near-grid rather than exact grid order — the right trade for a
* bulk dump. The delete-job visibility mask applies, like every user-facing read.
*/
export async function selectExportRowPage(
table: TableDefinition,
after: { position: number; id: string } | null,
limit: number
): Promise<Array<{ id: string; data: RowData; position: number }>> {
const deleteMask = await pendingDeleteMask(table)
const rows = await db
.select({ id: userTableRows.id, data: userTableRows.data, position: userTableRows.position })
.from(userTableRows)
.where(
and(
eq(userTableRows.tableId, table.id),
eq(userTableRows.workspaceId, table.workspaceId),
deleteMask,
after
? sql`(${userTableRows.position}, ${userTableRows.id}) > (${after.position}, ${after.id})`
: undefined
)
)
.orderBy(asc(userTableRows.position), asc(userTableRows.id))
.limit(limit)
return rows as Array<{ id: string; data: RowData; position: number }>
}
/** How long a terminal export stays listable (and re-downloadable from the tray). */
const EXPORT_JOB_VISIBILITY_MS = 10 * 60 * 1000
export interface WorkspaceExportJob {
jobId: string
tableId: string
tableName: string
status: string
rowsProcessed: number
format: 'csv' | 'json'
hasResult: boolean
error: string | null
}
/**
* Export jobs the tray surfaces for a workspace: everything running, plus terminals from the last
* {@link EXPORT_JOB_VISIBILITY_MS} so a just-finished export stays re-downloadable. Exports live
* outside the table-level job derivation (which excludes them), so this is their read path.
*/
export async function listWorkspaceExportJobs(workspaceId: string): Promise<WorkspaceExportJob[]> {
const visibilityCutoff = new Date(Date.now() - EXPORT_JOB_VISIBILITY_MS)
const rows = await db
.select({
jobId: tableJobs.id,
tableId: tableJobs.tableId,
tableName: userTableDefinitions.name,
status: tableJobs.status,
rowsProcessed: tableJobs.rowsProcessed,
payload: tableJobs.payload,
error: tableJobs.error,
})
.from(tableJobs)
.innerJoin(userTableDefinitions, eq(userTableDefinitions.id, tableJobs.tableId))
.where(
and(
eq(tableJobs.workspaceId, workspaceId),
eq(tableJobs.type, 'export'),
or(eq(tableJobs.status, 'running'), gt(tableJobs.updatedAt, visibilityCutoff))
)
)
.orderBy(desc(tableJobs.startedAt))
return rows.map((r) => {
const payload = r.payload as TableExportJobPayload | null
return {
jobId: r.jobId,
tableId: r.tableId,
tableName: r.tableName,
status: r.status,
rowsProcessed: r.rowsProcessed,
format: payload?.format ?? 'csv',
hasResult: Boolean(payload?.resultKey),
error: r.error,
}
})
}
/** Reads one job row (type/status/payload) scoped to its table. Null when absent. */
export async function getTableJob(
tableId: string,
jobId: string
): Promise<{ id: string; type: string; status: string; payload: unknown } | null> {
const [job] = await db
.select({
id: tableJobs.id,
type: tableJobs.type,
status: tableJobs.status,
payload: tableJobs.payload,
})
.from(tableJobs)
.where(and(eq(tableJobs.id, jobId), eq(tableJobs.tableId, tableId)))
.limit(1)
return job ?? null
}
/**
* Stamps an export job's generated-file storage key onto its payload (`{ resultKey }` merge).
* Scoped to the still-running job so a superseded attempt can't clobber a newer run's result.
* The download route reads it; the janitor deletes the file when the terminal job is pruned.
*/
export async function setJobResultKey(
tableId: string,
jobId: string,
resultKey: string
): Promise<void> {
await db
.update(tableJobs)
.set({
payload: sql`coalesce(${tableJobs.payload}, '{}'::jsonb) || jsonb_build_object('resultKey', ${resultKey}::text)`,
updatedAt: new Date(),
})
.where(ownsActiveJob(tableId, jobId))
}
/** Shared WHERE for terminal transitions: this job run, and still in-flight (write-once). */
function ownsActiveJob(tableId: string, jobId: string) {
return and(
eq(tableJobs.id, jobId),
eq(tableJobs.tableId, tableId),
eq(tableJobs.status, 'running')
)
}
/**
* Marks a job complete. No-op unless it's still this in-flight run. Returns whether it
* transitioned, so the worker only emits the `ready` event when it actually won (and not after a
* cancel / supersede).
*/
export async function markJobReady(tableId: string, jobId: string): Promise<boolean> {
const now = new Date()
const updated = await db
.update(tableJobs)
.set({ status: 'ready', error: null, completedAt: now, updatedAt: now })
.where(ownsActiveJob(tableId, jobId))
.returning({ id: tableJobs.id })
return updated.length > 0
}
/**
* Marks a job failed, leaving any already-committed work in place. No-op unless it's still this
* in-flight run (so a stale worker can't clobber a newer job or a cancel).
*/
export async function markJobFailed(tableId: string, jobId: string, error: string): Promise<void> {
const now = new Date()
await db
.update(tableJobs)
.set({ status: 'failed', error: error.slice(0, 2000), completedAt: now, updatedAt: now })
.where(ownsActiveJob(tableId, jobId))
}
/**
* Marks an in-flight job canceled (user-initiated). No-op unless it's still running. The
* worker's next ownership check then returns `false` and it stops; committed work is left in
* place (no rollback). Returns whether a running job was actually canceled.
*/
export async function markJobCanceled(tableId: string, jobId: string): Promise<boolean> {
const now = new Date()
const updated = await db
.update(tableJobs)
.set({ status: 'canceled', completedAt: now, updatedAt: now })
.where(ownsActiveJob(tableId, jobId))
.returning({ id: tableJobs.id })
return updated.length > 0
}
+1 -1
View File
@@ -5,7 +5,7 @@
* with table-specific information so LLMs can construct proper queries.
*/
import type { TableSummary } from '../types'
import type { TableSummary } from '@/lib/table/types'
/**
* Operations that use filters and need filter-specific enrichment.
+1 -1
View File
@@ -1 +1 @@
export * from './enrichment'
export * from '@/lib/table/llm/enrichment'
+1 -1
View File
@@ -6,7 +6,7 @@ import { db } from '@sim/db'
import { userTableDefinitions } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { and, eq, isNull } from 'drizzle-orm'
import type { TableSchema } from '../types'
import type { TableSchema } from '@/lib/table/types'
const logger = createLogger('TableWandEnricher')
@@ -2,7 +2,7 @@
* Constants for table query builder UI (filtering and sorting).
*/
export type { FilterRule, SortRule } from '../types'
export type { FilterRule, SortRule } from '@/lib/table/types'
export const COMPARISON_OPERATORS = [
{ value: 'eq', label: 'equals' },
@@ -4,7 +4,14 @@
import { generateShortId } from '@sim/utils/id'
import { isRecordLike } from '@sim/utils/object'
import type { Filter, FilterRule, JsonValue, Sort, SortDirection, SortRule } from '../types'
import type {
Filter,
FilterRule,
JsonValue,
Sort,
SortDirection,
SortRule,
} from '@/lib/table/types'
/** Converts UI filter rules to a Filter object for API queries. */
export function filterRulesToFilter(rules: FilterRule[]): Filter | null {
+3 -3
View File
@@ -2,6 +2,6 @@
* Query builder UI utilities for filtering and sorting tables.
*/
export * from './constants'
export * from './converters'
export * from './use-query-builder'
export * from '@/lib/table/query-builder/constants'
export * from '@/lib/table/query-builder/converters'
export * from '@/lib/table/query-builder/use-query-builder'
@@ -4,14 +4,14 @@
import { useCallback } from 'react'
import { generateShortId } from '@sim/utils/id'
import type { ColumnOption } from '../types'
import {
COMPARISON_OPERATORS,
type FilterRule,
LOGICAL_OPERATORS,
SORT_DIRECTIONS,
type SortRule,
} from './constants'
} from '@/lib/table/query-builder/constants'
import type { ColumnOption } from '@/lib/table/types'
const comparisonOptions: ColumnOption[] = COMPARISON_OPERATORS.map((op) => ({
value: op.value,
+298
View File
@@ -0,0 +1,298 @@
/**
* Row-executions (workflow-group results) internals for the table service layer.
*
* Internal module: not exposed via the `@/lib/table` barrel. Consumers import
* directly from `@/lib/table/rows/executions`.
*/
import { tableRowExecutions } from '@sim/db/schema'
import { and, eq, inArray, type SQL, sql } from 'drizzle-orm'
import type { DbOrTx } from '@/lib/db/types'
import { getColumnId } from '@/lib/table/column-keys'
import { areGroupDepsSatisfied } from '@/lib/table/deps'
import type {
RowData,
RowExecutionMetadata,
RowExecutions,
TableRow,
TableSchema,
} from '@/lib/table/types'
/**
* Loads `tableRowExecutions` rows for the given row ids and groups them into a
* `Map<rowId, RowExecutions>` suitable for plugging into `TableRow.executions`.
*/
export async function loadExecutionsByRow(
trx: DbOrTx,
rowIds: Iterable<string>
): Promise<Map<string, RowExecutions>> {
const ids = Array.from(new Set(rowIds))
const result = new Map<string, RowExecutions>()
if (ids.length === 0) return result
const rows = await trx
.select()
.from(tableRowExecutions)
.where(inArray(tableRowExecutions.rowId, ids))
for (const r of rows) {
const existing = result.get(r.rowId) ?? {}
const meta: RowExecutionMetadata = {
status: r.status as RowExecutionMetadata['status'],
executionId: r.executionId ?? null,
jobId: r.jobId ?? null,
workflowId: r.workflowId,
error: r.error ?? null,
...(r.runningBlockIds && r.runningBlockIds.length > 0
? { runningBlockIds: r.runningBlockIds }
: {}),
...(r.blockErrors && Object.keys(r.blockErrors as Record<string, string>).length > 0
? { blockErrors: r.blockErrors as Record<string, string> }
: {}),
...(r.cancelledAt ? { cancelledAt: r.cancelledAt.toISOString() } : {}),
}
existing[r.groupId] = meta
result.set(r.rowId, existing)
}
return result
}
/** Convenience: load executions for one row, returning `{}` when missing. */
export async function loadExecutionsForRow(trx: DbOrTx, rowId: string): Promise<RowExecutions> {
const byRow = await loadExecutionsByRow(trx, [rowId])
return byRow.get(rowId) ?? {}
}
/**
* Derive automatic clears + cancellation candidates from a row's data patch.
*
* Walks `schema.workflowGroups` left-to-right with a propagating `dirtied`
* column set. For each group whose deps overlap the dirty set, decide to
* clear (terminal exec) or cancel+rerun (in-flight exec), then add the
* group's outputs to the dirty set so later groups in the chain see them
* as dirty too. This models transitive dep chains as a single forward pass —
* editing column A propagates through group 1 (deps on A) to group 2 (deps
* on group 1's output) without explicit DAG traversal.
*
* Returns:
* - `executionsPatch`: caller's patch + nulls for cleared groups (or
* undefined if nothing applied).
* - `inFlightDownstreamGroups`: groups whose dep was dirtied and that are
* currently in-flight. Cancel-and-restart is the caller's job.
*
* Assumption: `workflowGroups[]` is in topological order — a group's deps
* may only reference columns to its left (enforced by `workflow-sidebar`'s
* "Run after" picker + the reorder scrub via `stripGroupDeps`). Violating
* this would silently miss the propagation.
*/
export function deriveExecClearsForDataPatch(
dataPatch: RowData,
schema: TableSchema,
existingExecutions: RowExecutions,
callerPatch: Record<string, RowExecutionMetadata | null> | undefined,
mergedData: RowData
): {
executionsPatch: Record<string, RowExecutionMetadata | null> | undefined
inFlightDownstreamGroups: string[]
} {
const dirtied = new Set(Object.keys(dataPatch))
const groupsToClear = new Set<string>()
const inFlightDownstreamGroups: string[] = []
// Own-output clears: when the user wipes a workflow output column, drop
// that group's exec entry so the auto-fire reactor re-arms the cell.
// Also flags the cleared output column as dirty so transitive downstream
// groups see it.
for (const [columnId, value] of Object.entries(dataPatch)) {
const cleared = value === null || value === undefined || value === ''
if (!cleared) continue
const col = schema.columns.find((c) => getColumnId(c) === columnId)
if (col?.workflowGroupId) groupsToClear.add(col.workflowGroupId)
}
// Left-to-right walk, propagating dirty columns forward.
const groups = schema.workflowGroups ?? []
const afterRow = { data: mergedData } as TableRow
for (const group of groups) {
const deps = group.dependencies?.columns ?? []
const depMatched = deps.some((d) => dirtied.has(d))
if (!depMatched) continue
// A dep column changed, but if the group's deps are no longer satisfied
// after the patch — a checkbox was unchecked or a text dep cleared — there's
// nothing to recompute. Leave the prior result alone instead of re-arming or
// cancelling it; only checking a box / filling a dep drives downstream work.
if (!areGroupDepsSatisfied(group, afterRow)) continue
const exec = existingExecutions[group.id]
if (exec) {
const status = exec.status
if (status === 'completed' || status === 'error' || status === 'cancelled') {
groupsToClear.add(group.id)
} else if (status === 'queued' || status === 'running' || status === 'pending') {
inFlightDownstreamGroups.push(group.id)
}
} else {
// No exec entry yet — `mode: 'new'` already covers this group. We
// still propagate the dirty signal forward so later groups in the
// chain see this group's outputs as dirty too.
groupsToClear.add(group.id)
}
// Propagate: this group is about to be re-computed, so groups whose
// deps reference its output columns are also dirty.
for (const out of group.outputs) dirtied.add(out.columnName)
}
if (groupsToClear.size === 0) {
return { executionsPatch: callerPatch, inFlightDownstreamGroups }
}
const merged: Record<string, RowExecutionMetadata | null> = { ...(callerPatch ?? {}) }
for (const gid of groupsToClear) {
if (!(gid in merged)) merged[gid] = null
}
return { executionsPatch: merged, inFlightDownstreamGroups }
}
/** Merges an `executionsPatch` into the row's existing executions blob. */
export function applyExecutionsPatch(
existing: RowExecutions,
patch: Record<string, RowExecutionMetadata | null> | undefined
): RowExecutions {
if (!patch) return existing
const next: RowExecutions = { ...existing }
for (const [gid, value] of Object.entries(patch)) {
if (value === null) {
delete next[gid]
} else {
next[gid] = value
}
}
return next
}
/**
* Writes a per-group execution patch for one row against the `tableRowExecutions`
* sidecar. Non-null values upsert into the table; nulls delete the entry. When
* `guard` is set, the upsert is gated to:
* - reject if a `cancelled` row for the same execution already exists, and
* - reject if the row exists but is owned by a different executionId
* (with carve-outs for missing rows and null executionIds — the dispatcher's
* pre-batch `pending` stamp leaves executionId unset so the first cell-task
* can claim).
*
* Returns `'guard-rejected'` when the guarded group's upsert affected 0 rows
* (callers signal failure to the cell-task path). Returns `'wrote'` otherwise.
*/
export async function writeExecutionsPatch(
trx: DbOrTx,
tableId: string,
rowId: string,
patch: Record<string, RowExecutionMetadata | null> | undefined,
guard?: { groupId: string; executionId: string }
): Promise<'wrote' | 'guard-rejected'> {
if (!patch) return 'wrote'
const entries = Object.entries(patch)
if (entries.length === 0) return 'wrote'
for (const [gid, value] of entries) {
if (value === null) {
await trx
.delete(tableRowExecutions)
.where(and(eq(tableRowExecutions.rowId, rowId), eq(tableRowExecutions.groupId, gid)) as SQL)
continue
}
const insertValues = {
tableId,
rowId,
groupId: gid,
status: value.status,
executionId: value.executionId,
jobId: value.jobId,
workflowId: value.workflowId,
error: value.error,
runningBlockIds: value.runningBlockIds ?? [],
blockErrors: value.blockErrors ?? {},
cancelledAt: value.cancelledAt ? new Date(value.cancelledAt) : null,
updatedAt: new Date(),
} as const
const isGuarded = guard && guard.groupId === gid
if (isGuarded) {
// Gate by guard semantics. The original JSONB guard had two AND'd
// clauses; we collapse them onto the upsert's WHERE so a non-matching
// existing row leaves the table untouched and we observe 0 affected.
const guardExecutionId = guard.executionId
const updated = await trx
.insert(tableRowExecutions)
.values(insertValues)
.onConflictDoUpdate({
target: [tableRowExecutions.rowId, tableRowExecutions.groupId],
set: {
status: insertValues.status,
executionId: insertValues.executionId,
jobId: insertValues.jobId,
workflowId: insertValues.workflowId,
error: insertValues.error,
runningBlockIds: insertValues.runningBlockIds,
blockErrors: insertValues.blockErrors,
cancelledAt: insertValues.cancelledAt,
updatedAt: insertValues.updatedAt,
},
where: and(
// Reject any guarded worker write when the cell is `cancelled` — a
// stop click wrote it authoritatively. SQL mirror of `isExecCancelled`
// (deps.ts). Status-only (not executionId-scoped): the cancel can
// only carry the pre-stamp's executionId (often null), so matching on
// id would let the worker's real-id claim resurrect a killed cell.
sql`${tableRowExecutions.status} <> 'cancelled'`,
// Stale-worker: the cell's active run has moved on. Carve-outs
// permit a fresh worker to take over when the row's executionId
// is unset (dispatcher's pre-batch `pending` stamp).
sql`(${tableRowExecutions.executionId} IS NULL OR ${tableRowExecutions.executionId} = ${guardExecutionId})`
) as SQL,
})
.returning({ rowId: tableRowExecutions.rowId })
if (updated.length === 0) return 'guard-rejected'
continue
}
await trx
.insert(tableRowExecutions)
.values(insertValues)
.onConflictDoUpdate({
target: [tableRowExecutions.rowId, tableRowExecutions.groupId],
set: {
status: insertValues.status,
executionId: insertValues.executionId,
jobId: insertValues.jobId,
workflowId: insertValues.workflowId,
error: insertValues.error,
runningBlockIds: insertValues.runningBlockIds,
blockErrors: insertValues.blockErrors,
cancelledAt: insertValues.cancelledAt,
updatedAt: insertValues.updatedAt,
},
})
}
return 'wrote'
}
/**
* Strips the given workflow group ids from every row's executions on a table —
* used by the column / group delete paths so stale running/queued exec records
* don't linger and inflate counters after the group is gone. The caller wraps
* in their own transaction.
*/
export async function stripGroupExecutions(
trx: DbOrTx,
tableId: string,
groupIds: Iterable<string>
): Promise<void> {
const ids = Array.from(new Set(groupIds))
if (ids.length === 0) return
await trx
.delete(tableRowExecutions)
.where(
and(eq(tableRowExecutions.tableId, tableId), inArray(tableRowExecutions.groupId, ids)) as SQL
)
}
+557
View File
@@ -0,0 +1,557 @@
/**
* Row position / fractional-ordering internals for the table service layer.
*
* Internal module: only the import/delete-runner entry points are exposed via
* the `@/lib/table/rows/ordering` path. Not re-exported through the
* `@/lib/table` barrel.
*/
import { db } from '@sim/db'
import { userTableRows } from '@sim/db/schema'
import { and, asc, desc, eq, gt, gte, inArray, lt, lte, type SQL, sql } from 'drizzle-orm'
import { isTablesFractionalOrderingEnabled } from '@/lib/core/config/feature-flags'
import type { DbOrTx } from '@/lib/db/types'
import { TABLE_LIMITS } from '@/lib/table/constants'
import { keyBetween, nKeysBetween } from '@/lib/table/order-key'
import { type DbExecutor, type DbTransaction, withSeqscanOff } from '@/lib/table/planner'
import { setTableTxTimeouts } from '@/lib/table/tx'
import type { RowData } from '@/lib/table/types'
/**
* Starting `position` for an append import — `max(position) + 1`, or 0 when empty. Read once,
* unlocked, before streaming: the import worker is the table's sole writer, so it can assign
* contiguous positions from this offset without per-batch position scans.
*/
export async function nextImportStartPosition(tableId: string): Promise<number> {
const [{ maxPos }] = await db
.select({
maxPos: sql<number>`coalesce(max(${userTableRows.position}), -1)`.mapWith(Number),
})
.from(userTableRows)
.where(eq(userTableRows.tableId, tableId))
return maxPos + 1
}
/**
* Append anchor `order_key` for an import — `max(order_key)`, or null when empty. Read once,
* unlocked, before streaming (the import worker is the table's sole writer); each batch threads
* the previous batch's last key forward so no per-batch max scan is needed.
*/
export async function nextImportStartOrderKey(tableId: string): Promise<string | null> {
return maxOrderKey(db, tableId)
}
/**
* Serializes writers that assign `position` for the same table. The row-count
* trigger (migration 0198) serializes capacity via a row lock on
* `user_table_definitions`, but it fires AFTER INSERT, so two concurrent
* auto-positioned inserts could read the same snapshot and assign the same
* position (the `(table_id, position)` index is non-unique). This advisory lock
* restores per-table serialization. Released at COMMIT/ROLLBACK.
*/
export async function acquireRowOrderLock(trx: DbTransaction, tableId: string) {
await trx.execute(
sql`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_rows_pos:${tableId}`}, 0))`
)
}
/** Next append position for a table (max(position) + 1, or 0 if empty). */
export async function nextRowPosition(trx: DbTransaction, tableId: string): Promise<number> {
const [{ maxPos }] = await trx
.select({
maxPos: sql<number>`coalesce(max(${userTableRows.position}), -1)`.mapWith(Number),
})
.from(userTableRows)
.where(eq(userTableRows.tableId, tableId))
return maxPos + 1
}
/** Largest `order_key` for a table, or `null` when empty — the append anchor for new keys. */
export async function maxOrderKey(executor: DbOrTx, tableId: string): Promise<string | null> {
const [{ maxKey }] = await executor
.select({ maxKey: sql<string | null>`max(${userTableRows.orderKey})` })
.from(userTableRows)
.where(eq(userTableRows.tableId, tableId))
return maxKey ?? null
}
/** Shifts every row at or after `position` up by one (`position + 1`). */
export async function shiftRowsUpFrom(trx: DbTransaction, tableId: string, position: number) {
await trx
.update(userTableRows)
.set({ position: sql`position + 1` })
.where(and(eq(userTableRows.tableId, tableId), gte(userTableRows.position, position)))
}
/** Shifts every row after `position` down by one (`position - 1`). */
export async function shiftRowsDownAfter(trx: DbTransaction, tableId: string, position: number) {
await trx
.update(userTableRows)
.set({ position: sql`position - 1` })
.where(and(eq(userTableRows.tableId, tableId), gt(userTableRows.position, position)))
}
/**
* Reserves the `position` for a single inserted row and returns where to INSERT.
* Acquires the row-order lock, then opens a slot at `requestedPosition` (shifting
* the occupant + tail up) or computes the append position. Caller runs inside a
* transaction.
*/
export async function reserveInsertPosition(
trx: DbTransaction,
tableId: string,
requestedPosition?: number
): Promise<number> {
await acquireRowOrderLock(trx, tableId)
if (requestedPosition === undefined) {
return nextRowPosition(trx, tableId)
}
const [existing] = await trx
.select({ id: userTableRows.id })
.from(userTableRows)
.where(and(eq(userTableRows.tableId, tableId), eq(userTableRows.position, requestedPosition)))
.limit(1)
if (existing) {
await shiftRowsUpFrom(trx, tableId, requestedPosition)
}
return requestedPosition
}
/**
* Reserves positions for a batch of `count` rows. Opens each requested slot
* (ascending, preserving prior gaps) and returns the requested positions in
* original order; otherwise returns a contiguous append range.
*/
export async function reserveBatchPositions(
trx: DbTransaction,
tableId: string,
count: number,
requestedPositions?: number[]
): Promise<number[]> {
await acquireRowOrderLock(trx, tableId)
if (requestedPositions && requestedPositions.length > 0) {
for (const pos of [...requestedPositions].sort((a, b) => a - b)) {
await shiftRowsUpFrom(trx, tableId, pos)
}
return requestedPositions
}
const start = await nextRowPosition(trx, tableId)
return Array.from({ length: count }, (_, i) => start + i)
}
/**
* Recompacts row positions to be contiguous after a bulk delete. With
* `minDeletedPos`, only rows at/after it are re-numbered; single-row deletes use
* the cheaper {@link shiftRowsDownAfter}.
*/
export async function compactPositions(
trx: DbTransaction,
tableId: string,
minDeletedPos?: number
) {
if (minDeletedPos === undefined) {
await trx.execute(sql`
UPDATE user_table_rows t
SET position = r.new_pos
FROM (
SELECT id, ROW_NUMBER() OVER (ORDER BY position) - 1 AS new_pos
FROM user_table_rows
WHERE table_id = ${tableId}
) r
WHERE t.id = r.id AND t.table_id = ${tableId} AND t.position != r.new_pos
`)
return
}
await trx.execute(sql`
UPDATE user_table_rows t
SET position = r.new_pos
FROM (
SELECT id, ${minDeletedPos}::int + ROW_NUMBER() OVER (ORDER BY position) - 1 AS new_pos
FROM user_table_rows
WHERE table_id = ${tableId} AND position >= ${minDeletedPos}
) r
WHERE t.id = r.id AND t.table_id = ${tableId} AND t.position != r.new_pos
`)
}
/**
* Computes the fractional `order_key` for a row inserted at the integer
* `requestedPosition` (or appended when omitted). Used by position-based callers
* (mothership tool, v1 API, undo position-fallback, transient old clients).
*
* The neighbor at slot `s` is resolved differently per flag state:
* - **off**: `WHERE position = s` (positions are contiguous, so the row at
* position `s` is the `s`-th row — an indexed O(1) lookup).
* - **on**: the `s`-th row in `order_key, id` order (`OFFSET s`) — positions are
* gappy and non-authoritative, so `position = s` would miss; the visual
* ordinal is the key's ordinal. O(s), acceptable for these low-volume callers.
*
* Caller holds the row-order lock.
*/
export async function resolveInsertOrderKey(
trx: DbTransaction,
tableId: string,
requestedPosition?: number
): Promise<string> {
const orderKeyAtSlot = async (slot: number): Promise<string | null> => {
if (slot < 0) return null
if (isTablesFractionalOrderingEnabled) {
const [r] = await trx
.select({ orderKey: userTableRows.orderKey })
.from(userTableRows)
.where(eq(userTableRows.tableId, tableId))
.orderBy(asc(userTableRows.orderKey), asc(userTableRows.id))
.limit(1)
.offset(slot)
return r?.orderKey ?? null
}
const [r] = await trx
.select({ orderKey: userTableRows.orderKey })
.from(userTableRows)
.where(and(eq(userTableRows.tableId, tableId), eq(userTableRows.position, slot)))
.limit(1)
return r?.orderKey ?? null
}
if (requestedPosition === undefined) {
return keyBetween(await maxOrderKey(trx, tableId), null)
}
const lo = await orderKeyAtSlot(requestedPosition - 1)
const hi = await orderKeyAtSlot(requestedPosition)
return keyBetween(lo, hi)
}
/**
* Resolves the `order_key` for an insert expressed by an anchor row id —
* `afterRowId` (place directly after) or `beforeRowId` (directly before). Finds
* the anchor and its adjacent key via the `(table_id, order_key, id)` index
* (O(1)) and mints a key between them. Also returns a legacy integer `position`
* (anchor's position ±) so the flag-off shift path still works. Caller holds the
* row-order lock.
*/
export async function resolveInsertByNeighbor(
trx: DbTransaction,
tableId: string,
afterRowId?: string,
beforeRowId?: string
): Promise<{ orderKey: string; position: number }> {
const anchorId = afterRowId ?? beforeRowId!
const [anchor] = await trx
.select({ orderKey: userTableRows.orderKey, position: userTableRows.position })
.from(userTableRows)
.where(and(eq(userTableRows.tableId, tableId), eq(userTableRows.id, anchorId)))
.limit(1)
// The client targets a specific neighbor; a missing one (concurrent delete /
// stale view) is an error, not a silent insert at the front.
if (!anchor) throw new Error(`Row not found: ${anchorId}`)
const anchorKey = anchor.orderKey ?? null
// A null key on the anchor means the table isn't backfilled. With the flag on
// (key is authoritative) the adjacent-key lookup below can't work — fail
// loudly rather than mint a wrong key. Flag off keeps `position` authoritative,
// so a best-effort key here is fine (the backfill re-keys before the flip).
if (anchorKey === null && isTablesFractionalOrderingEnabled) {
throw new Error(`Row ${anchorId} has no order_key yet (table not backfilled)`)
}
if (afterRowId) {
// hi = the smallest key strictly GREATER than the anchor key. Comparing keys
// (not the `(order_key, id)` row tuple) skips past any sibling that shares the
// anchor's key, so `keyBetween` always gets strictly-ordered bounds and can't
// throw on a stray duplicate. Identical to the row tuple when keys are distinct.
// A null anchorKey (flag off, un-backfilled) has no key to compare — leave the
// upper bound open, matching the prior best-effort behavior.
let nextKey: string | null = null
if (anchorKey !== null) {
const [next] = await trx
.select({ orderKey: userTableRows.orderKey })
.from(userTableRows)
.where(and(eq(userTableRows.tableId, tableId), gt(userTableRows.orderKey, anchorKey)))
.orderBy(asc(userTableRows.orderKey))
.limit(1)
nextKey = next?.orderKey ?? null
}
return {
orderKey: keyBetween(anchorKey, nextKey),
position: anchor.position + 1,
}
}
// beforeRowId: lo = the largest key strictly LESS than the anchor key (distinct,
// same rationale as the afterRowId branch above).
let prevKey: string | null = null
if (anchorKey !== null) {
const [prev] = await trx
.select({ orderKey: userTableRows.orderKey })
.from(userTableRows)
.where(and(eq(userTableRows.tableId, tableId), lt(userTableRows.orderKey, anchorKey)))
.orderBy(desc(userTableRows.orderKey))
.limit(1)
prevKey = prev?.orderKey ?? null
}
return {
orderKey: keyBetween(prevKey, anchorKey),
position: anchor.position,
}
}
/**
* Computes fractional `order_key`s for a batch insert. With no `positions`,
* appends a contiguous run after the current max key. With explicit `positions`
* (undo restore), keys each row between its pre-shift position neighbors —
* correct because requested positions are distinct. Caller holds the lock.
*
* The explicit-`positions` path is meaningful only when `position` is
* authoritative (flag off): with the flag on, a saved `position` is a gappy
* column value, not a visual rank, so feeding it to {@link resolveInsertOrderKey}
* (which reads `position` as an `OFFSET` rank under the flag) would mint keys at
* the wrong ranks. Callers needing exact placement under the flag pass
* `orderKeys` (handled before this function); here we just append a run.
*/
export async function resolveBatchInsertOrderKeys(
trx: DbTransaction,
tableId: string,
count: number,
positions?: number[]
): Promise<string[]> {
if (!positions || positions.length === 0 || isTablesFractionalOrderingEnabled) {
return nKeysBetween(await maxOrderKey(trx, tableId), null, count)
}
const keys: string[] = []
for (const pos of positions) {
keys.push(await resolveInsertOrderKey(trx, tableId, pos))
}
return keys
}
/**
* Inserts a single row in its own transaction. Always assigns a fractional
* `order_key`. When the fractional-ordering flag is on, `order_key` is
* authoritative and `position` is a best-effort append (no O(N) shift); when
* off, `position` is reserved as before (shifting to open the slot). Validation
* and side-effect dispatch stay with the caller; capacity is enforced by the
* `increment_user_table_row_count` trigger.
*/
export async function insertOrderedRow(params: {
tableId: string
workspaceId: string
data: RowData
rowId: string
position?: number
afterRowId?: string
beforeRowId?: string
createdBy?: string
now: Date
}): Promise<{
id: string
data: RowData
position: number
orderKey: string | null
createdAt: Date
updatedAt: Date
}> {
const { tableId, workspaceId, data, rowId, position, afterRowId, beforeRowId, createdBy, now } =
params
const [row] = await db.transaction(async (trx) => {
await setTableTxTimeouts(trx)
await acquireRowOrderLock(trx, tableId)
// Resolve the order key (and a legacy slot position for the flag-off shift
// path) from neighbor ids when given, else from the requested position.
let orderKey: string
let slotPosition = position
if (afterRowId || beforeRowId) {
const resolved = await resolveInsertByNeighbor(trx, tableId, afterRowId, beforeRowId)
orderKey = resolved.orderKey
slotPosition = resolved.position
} else {
orderKey = await resolveInsertOrderKey(trx, tableId, position)
}
let targetPosition: number
if (isTablesFractionalOrderingEnabled) {
// order_key is authoritative — keep a best-effort, no-shift position.
targetPosition = await nextRowPosition(trx, tableId)
} else if (slotPosition !== undefined) {
const [existing] = await trx
.select({ id: userTableRows.id })
.from(userTableRows)
.where(and(eq(userTableRows.tableId, tableId), eq(userTableRows.position, slotPosition)))
.limit(1)
if (existing) await shiftRowsUpFrom(trx, tableId, slotPosition)
targetPosition = slotPosition
} else {
targetPosition = await nextRowPosition(trx, tableId)
}
return trx
.insert(userTableRows)
.values({
id: rowId,
tableId,
workspaceId,
data,
position: targetPosition,
orderKey,
createdAt: now,
updatedAt: now,
...(createdBy ? { createdBy } : {}),
})
.returning()
})
return {
id: row.id,
data: row.data as RowData,
position: row.position,
orderKey: row.orderKey,
createdAt: row.createdAt,
updatedAt: row.updatedAt,
}
}
/**
* Deletes a single row by id in its own transaction, then closes the positional
* gap. Returns `false` when no row matched.
*/
export async function deleteOrderedRow(params: {
tableId: string
rowId: string
workspaceId: string
}): Promise<boolean> {
const { tableId, rowId, workspaceId } = params
return db.transaction(async (trx) => {
await setTableTxTimeouts(trx)
const [deleted] = await trx
.delete(userTableRows)
.where(
and(
eq(userTableRows.id, rowId),
eq(userTableRows.tableId, tableId),
eq(userTableRows.workspaceId, workspaceId)
)
)
.returning({ position: userTableRows.position })
if (!deleted) return false
// Fractional ordering: deleting a row never changes another row's order_key,
// so the O(N) position reshift is skipped entirely.
if (!isTablesFractionalOrderingEnabled) {
await shiftRowsDownAfter(trx, tableId, deleted.position)
}
return true
})
}
/**
* Deletes the given row ids in batches within one transaction, then recompacts
* positions from the earliest deleted slot. Returns the deleted rows (id + prior
* position). The caller resolves which ids to delete (used by both delete-by-ids
* and delete-by-filter).
*/
export async function deleteOrderedRowsByIds(params: {
tableId: string
workspaceId: string
rowIds: string[]
}): Promise<{ id: string; position: number }[]> {
const { tableId, workspaceId, rowIds } = params
if (rowIds.length === 0) return []
return db.transaction(async (trx) => {
await setTableTxTimeouts(trx, { statementMs: 60_000 })
const deleted: { id: string; position: number }[] = []
for (let i = 0; i < rowIds.length; i += TABLE_LIMITS.DELETE_BATCH_SIZE) {
const batch = rowIds.slice(i, i + TABLE_LIMITS.DELETE_BATCH_SIZE)
const rows = await trx
.delete(userTableRows)
.where(
and(
eq(userTableRows.tableId, tableId),
eq(userTableRows.workspaceId, workspaceId),
inArray(userTableRows.id, batch)
)
)
.returning({ id: userTableRows.id, position: userTableRows.position })
deleted.push(...rows)
}
// Fractional ordering: deletes leave order_key untouched, so no recompaction.
if (!isTablesFractionalOrderingEnabled && deleted.length > 0) {
const minDeletedPos = deleted.reduce(
(min, r) => (r.position < min ? r.position : min),
deleted[0].position
)
await compactPositions(trx, tableId, minDeletedPos)
}
return deleted
})
}
/**
* Selects one page of row ids to delete for the async delete-job worker: base scope plus a
* `created_at <= cutoff` floor (so rows inserted after the job started are never selected) and
* the caller's optional filter clause. Keyset paginated on `id` via `afterId` so excluded rows
* (which are skipped, not deleted) still advance the cursor — no OFFSET, no risk of looping on a
* fully-excluded page.
*/
export async function selectRowIdPage(params: {
tableId: string
workspaceId: string
cutoff: Date
filterClause?: SQL
afterId?: string
limit: number
}): Promise<string[]> {
const { tableId, workspaceId, cutoff, filterClause, afterId, limit } = params
const selectPage = (executor: DbExecutor) =>
executor
.select({ id: userTableRows.id })
.from(userTableRows)
.where(
and(
eq(userTableRows.tableId, tableId),
eq(userTableRows.workspaceId, workspaceId),
lte(userTableRows.createdAt, cutoff),
afterId ? gt(userTableRows.id, afterId) : undefined,
filterClause
)
)
.orderBy(asc(userTableRows.id))
.limit(limit)
// A jsonb filter is unestimatable, so the planner would seq-scan the whole shared relation
// per page (12.6s measured) — keep it on the tenant's (table_id, id) index.
const rows = filterClause
? await withSeqscanOff(async (trx) => selectPage(trx))
: await selectPage(db)
return rows.map((r) => r.id)
}
/**
* Deletes one page of rows for the async delete-job worker, committing each `DELETE_BATCH_SIZE`
* chunk in its own short transaction. One statement per transaction bounds how long the
* statement-level row_count trigger's lock on the definition row is held (a page-wide transaction
* held it for the entire page, starving concurrent inserts and overrunning `statement_timeout`),
* and a mid-page failure loses at most one uncommitted batch — the keyset walker (or a task
* retry) re-walks whatever remains. Skips legacy position compaction: under fractional ordering
* it's unnecessary, and in the legacy path `position` gaps are harmless — rows still order by
* position. Returns the count deleted.
*/
export async function deletePageByIds(
tableId: string,
workspaceId: string,
rowIds: string[]
): Promise<number> {
let deleted = 0
for (let i = 0; i < rowIds.length; i += TABLE_LIMITS.DELETE_BATCH_SIZE) {
const batch = rowIds.slice(i, i + TABLE_LIMITS.DELETE_BATCH_SIZE)
const rows = await db.transaction(async (trx) => {
await setTableTxTimeouts(trx, { statementMs: 60_000 })
return trx
.delete(userTableRows)
.where(
and(
eq(userTableRows.tableId, tableId),
eq(userTableRows.workspaceId, workspaceId),
inArray(userTableRows.id, batch)
)
)
.returning({ id: userTableRows.id })
})
deleted += rows.length
}
return deleted
}
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
+9 -3
View File
@@ -7,9 +7,15 @@
import type { SQL } from 'drizzle-orm'
import { sql } from 'drizzle-orm'
import { getColumnId } from './column-keys'
import { NAME_PATTERN } from './constants'
import type { ColumnDefinition, ConditionOperators, Filter, JsonValue, Sort } from './types'
import { getColumnId } from '@/lib/table/column-keys'
import { NAME_PATTERN } from '@/lib/table/constants'
import type {
ColumnDefinition,
ConditionOperators,
Filter,
JsonValue,
Sort,
} from '@/lib/table/types'
/**
* Error thrown when caller-supplied filter or sort input is malformed.
+50
View File
@@ -0,0 +1,50 @@
/**
* Shared transaction / locking helpers for the table service layer.
*
* Internal module: not exposed via the `@/lib/table` barrel. Consumers import
* directly from `@/lib/table/tx`.
*/
import { sql } from 'drizzle-orm'
import type { DbTransaction } from '@/lib/table/planner'
const TIMEOUT_CAP_MS = 10 * 60_000
/**
* Sets per-transaction Postgres timeouts via `SET LOCAL`.
*
* `lock_timeout` is the critical one: without it, a waiter inherits the full
* `statement_timeout` clock, so one stuck writer can drain the pool.
*
* Safe under pgBouncer transaction pooling — `SET LOCAL` is transaction-scoped
* and cleared at COMMIT/ROLLBACK before the session returns to the pool.
*/
export async function setTableTxTimeouts(
trx: DbTransaction,
opts?: { statementMs?: number; lockMs?: number; idleMs?: number }
) {
const s = opts?.statementMs ?? 10_000
const l = opts?.lockMs ?? 3_000
const i = opts?.idleMs ?? 5_000
await trx.execute(sql.raw(`SET LOCAL statement_timeout = '${s}ms'`))
await trx.execute(sql.raw(`SET LOCAL lock_timeout = '${l}ms'`))
await trx.execute(sql.raw(`SET LOCAL idle_in_transaction_session_timeout = '${i}ms'`))
}
/**
* Scales `statement_timeout` to the expected row-count work.
*
* Bulk operations that rewrite JSONB or cascade row triggers (e.g.
* `replaceTableRows`, `deleteColumn`, `renameColumn`) scale roughly linearly
* with row count. A fixed cap would regress large-table users who never saw a
* timeout before `SET LOCAL` was introduced. This helper picks
* `max(baseMs, rowCount * perRowMs)`, capped at 10 minutes so a single
* runaway transaction cannot indefinitely pin a pool connection.
*/
export function scaledStatementTimeoutMs(
rowCount: number,
opts: { baseMs: number; perRowMs: number }
): number {
const safeRowCount = Math.max(0, rowCount)
return Math.min(TIMEOUT_CAP_MS, Math.max(opts.baseMs, safeRowCount * opts.perRowMs))
}
+1 -1
View File
@@ -2,7 +2,7 @@
* Type definitions for user-defined tables.
*/
import type { COLUMN_TYPES } from './constants'
import type { COLUMN_TYPES } from '@/lib/table/constants'
export type ColumnValue = string | number | boolean | null | Date
export type JsonValue = ColumnValue | JsonValue[] | { [key: string]: JsonValue }
+10 -4
View File
@@ -6,10 +6,16 @@ import { db } from '@sim/db'
import { userTableRows } from '@sim/db/schema'
import { and, eq, or, type SQL, sql } from 'drizzle-orm'
import { NextResponse } from 'next/server'
import { getColumnId } from './column-keys'
import { COLUMN_TYPES, NAME_PATTERN, TABLE_LIMITS } from './constants'
import { withSeqscanOff } from './planner'
import type { ColumnDefinition, JsonValue, RowData, TableSchema, ValidationResult } from './types'
import { getColumnId } from '@/lib/table/column-keys'
import { COLUMN_TYPES, NAME_PATTERN, TABLE_LIMITS } from '@/lib/table/constants'
import { withSeqscanOff } from '@/lib/table/planner'
import type {
ColumnDefinition,
JsonValue,
RowData,
TableSchema,
ValidationResult,
} from '@/lib/table/types'
export type { ColumnDefinition, TableSchema, ValidationResult }
+12 -9
View File
@@ -31,16 +31,16 @@ import type {
const logger = createLogger('WorkflowGroupScheduler')
import { getColumnId } from './column-keys'
import { USER_TABLE_ROWS_SQL_NAME } from './constants'
import { areGroupDepsSatisfied, areOutputsFilled, isExecInFlight } from './deps'
import type { DispatchLimit, DispatchMode } from './dispatcher'
import { buildFilterClause } from './sql'
import { getColumnId } from '@/lib/table/column-keys'
import { USER_TABLE_ROWS_SQL_NAME } from '@/lib/table/constants'
import { areGroupDepsSatisfied, areOutputsFilled, isExecInFlight } from '@/lib/table/deps'
import type { DispatchLimit, DispatchMode } from '@/lib/table/dispatcher'
import { buildFilterClause } from '@/lib/table/sql'
export {
getUnmetGroupDeps,
optimisticallyScheduleNewlyEligibleGroups,
} from './deps'
} from '@/lib/table/deps'
/**
* Per-(row, group) eligibility for both the auto-fire reactor and manual
@@ -367,9 +367,12 @@ export async function cancelWorkflowGroupRuns(
rowId?: string,
options?: { groupIds?: string[]; filter?: Filter; excludeRowIds?: string[] }
): Promise<number> {
const { getTableById, updateRow } = await import('@/lib/table/service')
const { getTableById } = await import('@/lib/table/service')
const { updateRow } = await import('@/lib/table/rows/service')
const { getJobQueue } = await import('@/lib/core/async-jobs/config')
const { listActiveDispatches, markActiveDispatchesCancelled } = await import('./dispatcher')
const { listActiveDispatches, markActiveDispatchesCancelled } = await import(
'@/lib/table/dispatcher'
)
const table = await getTableById(tableId)
if (!table) {
@@ -661,7 +664,7 @@ export async function runWorkflowColumn(opts: {
if (rowIds && rowIds.length === 0) return { dispatchId: null }
// Lazy imports: `./service` and `./dispatcher` both close cycles back to
// this module; `@trigger.dev/sdk` is heavy and only needed on this op.
const { getTableById } = await import('./service')
const { getTableById } = await import('@/lib/table/service')
const table = await getTableById(tableId)
if (!table) throw new Error('Table not found')
if (table.workspaceId !== workspaceId) throw new Error('Invalid workspace ID')
@@ -0,0 +1,964 @@
/**
* Workflow-group operations on user tables.
*
* Extracted from the table service: add/update/delete workflow groups and their
* output columns, plus stale-output pruning after a workflow deploy. These ops
* mutate `schema.workflowGroups` (and the bound output columns + row data) under
* the per-table advisory lock from `withLockedTable`.
*/
import { db } from '@sim/db'
import { userTableDefinitions } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { and, eq, isNull, sql } from 'drizzle-orm'
import type { DbOrTx } from '@/lib/db/types'
import {
columnMatchesRef,
generateColumnId,
getColumnId,
remapGroupColumnRefs,
} from '@/lib/table/column-keys'
import { NAME_PATTERN, TABLE_LIMITS } from '@/lib/table/constants'
import { stripGroupExecutions } from '@/lib/table/rows/executions'
import { getTableById, withLockedTable } from '@/lib/table/service'
import { setTableTxTimeouts } from '@/lib/table/tx'
import type {
AddWorkflowGroupData,
ColumnDefinition,
DeleteWorkflowGroupData,
TableDefinition,
TableMetadata,
TableSchema,
UpdateWorkflowGroupData,
WorkflowGroup,
WorkflowGroupOutput,
} from '@/lib/table/types'
import { assertValidSchema, runWorkflowColumn, stripGroupDeps } from '@/lib/table/workflow-columns'
const logger = createLogger('TableWorkflowGroupsService')
/**
* Drops references to deleted blocks from every workflow group on every table
* that targets the just-deployed workflow. Called from the workflow deploy
* orchestrator after the new deployment commits, so the table UI never holds
* stale `{blockId, path}` entries for blocks the user removed.
*
* - Filters `outputs[]` per group. If every output would be filtered out, the
* group is left untouched and a warning is logged — the user must
* reconfigure it manually.
* - Scoped to the workflow's workspace.
* - Idempotent: running twice with the same `validBlockIds` is a no-op on the
* second pass. Existing row data is left alone.
*/
export async function pruneStaleWorkflowGroupOutputs({
workflowId,
workspaceId,
validBlockIds,
requestId,
tx,
}: {
workflowId: string
workspaceId: string
validBlockIds: Set<string>
requestId: string
tx?: DbOrTx
}): Promise<void> {
const executor = tx ?? db
const tables = await executor
.select({
id: userTableDefinitions.id,
schema: userTableDefinitions.schema,
})
.from(userTableDefinitions)
.where(
and(
eq(userTableDefinitions.workspaceId, workspaceId),
isNull(userTableDefinitions.archivedAt)
)
)
for (const t of tables) {
const schema = t.schema as TableSchema
const groups = schema.workflowGroups ?? []
if (groups.length === 0) continue
let mutated = false
const nextGroups = groups.map((group) => {
if (group.workflowId !== workflowId) return group
const filtered = group.outputs.filter((o) => validBlockIds.has(o.blockId))
if (filtered.length === group.outputs.length) return group
if (filtered.length === 0) {
logger.warn(
`[${requestId}] All outputs for workflow group "${group.name ?? group.id}" in table ${t.id} reference deleted blocks; leaving group intact for user reconfiguration.`
)
return group
}
mutated = true
return { ...group, outputs: filtered }
})
if (!mutated) continue
await executor
.update(userTableDefinitions)
.set({
schema: { ...schema, workflowGroups: nextGroups },
updatedAt: new Date(),
})
.where(eq(userTableDefinitions.id, t.id))
logger.info(`[${requestId}] Pruned stale workflow=${workflowId} block refs from table ${t.id}`)
}
}
/**
* Atomically inserts a workflow group plus its output columns into a table's
* schema. Both arrays update in one DB write so the schema is never observed
* mid-mutation (e.g. columns referencing a group that doesn't yet exist).
*/
export async function addWorkflowGroup(
data: AddWorkflowGroupData,
requestId: string
): Promise<TableDefinition> {
const updatedTable = await withLockedTable(data.tableId, async (table, trx) => {
const schema = table.schema
const groups = schema.workflowGroups ?? []
if (groups.some((g) => g.id === data.group.id)) {
throw new Error(`Workflow group "${data.group.id}" already exists`)
}
const existingNames = new Set(schema.columns.map((c) => c.name.toLowerCase()))
for (const col of data.outputColumns) {
if (!NAME_PATTERN.test(col.name)) {
throw new Error(
`Invalid output column name "${col.name}". Must satisfy ${NAME_PATTERN.source}.`
)
}
if (existingNames.has(col.name.toLowerCase())) {
throw new Error(`Column "${col.name}" already exists`)
}
}
if (schema.columns.length + data.outputColumns.length > TABLE_LIMITS.MAX_COLUMNS_PER_TABLE) {
throw new Error(
`Adding ${data.outputColumns.length} columns would exceed the maximum (${TABLE_LIMITS.MAX_COLUMNS_PER_TABLE}).`
)
}
// Assign stable ids to the new output columns, then rewrite the group's
// column refs from name → id so outputs/deps/inputMappings key on ids —
// matching the row-data storage key and surviving future renames.
const outputColumns = data.outputColumns.map((col) =>
col.id ? col : { ...col, id: generateColumnId() }
)
const updatedColumns = [...schema.columns, ...outputColumns]
const idByName = new Map(updatedColumns.map((c) => [c.name, getColumnId(c)]))
const group = remapGroupColumnRefs(data.group, idByName)
const updatedSchema: TableSchema = {
...schema,
columns: updatedColumns,
workflowGroups: [...groups, group],
}
// Keep `metadata.columnOrder` (column ids) in sync — see `addTableColumn`.
// New output columns get appended in the order the caller supplied.
const existingOrder = table.metadata?.columnOrder
let updatedMetadata = table.metadata
if (existingOrder && existingOrder.length > 0) {
const known = new Set(existingOrder)
const append = outputColumns.map(getColumnId).filter((id) => !known.has(id))
if (append.length > 0) {
updatedMetadata = { ...table.metadata, columnOrder: [...existingOrder, ...append] }
}
}
assertValidSchema(updatedSchema, updatedMetadata?.columnOrder)
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
logger.info(
`[${requestId}] Added workflow group "${data.group.id}" with ${data.outputColumns.length} output column(s) to table ${data.tableId}`
)
return {
...table,
schema: updatedSchema,
metadata: updatedMetadata,
updatedAt: now,
}
})
// Auto-fire existing rows whose deps are already met for the new group.
// Fire-and-forget — the dispatcher bounds queue depth (window of 20) and
// walks the table in the background. HTTP returns instantly; cells fill
// in over the next minutes as the dispatcher walks. Mothership opts out
// by setting `autoRun: false`.
if (data.autoRun !== false) {
void runWorkflowColumn({
tableId: updatedTable.id,
workspaceId: updatedTable.workspaceId,
mode: 'new',
isManualRun: false,
groupIds: [data.group.id],
requestId,
triggeredByUserId: data.actorUserId,
}).catch((err) => logger.error(`[${requestId}] auto-dispatch (addWorkflowGroup) failed:`, err))
}
return updatedTable
}
/**
* Updates a workflow group: any combination of workflowId, name, dependencies,
* outputs[]. Computes added/removed outputs vs current state and inserts /
* removes columns transactionally. Removed outputs also clear their key from
* every row's `data`.
*/
export async function updateWorkflowGroup(
data: UpdateWorkflowGroupData,
requestId: string
): Promise<TableDefinition> {
const mappingUpdates = data.mappingUpdates ?? []
// Phase 1 (no lock): when there are mapping updates, load the workflow once to
// resolve each remap's new leaf type. Kept OFF the advisory-lock critical
// section so concurrent group edits on the same table don't time out waiting
// on this DB load. Best-effort — a resolution failure leaves column types
// unchanged (workflow deleted, block removed). The result is applied against
// the fresh schema under the lock in phase 2.
const remapLeafTypeByColumn = new Map<string, ColumnDefinition['type']>()
// The workflow id the leaf types above were resolved against. Phase 2 only
// applies the resolved types if the group still points at this workflow under
// the lock — a concurrent `workflowId` change would make them stale.
let resolvedForWorkflowId: string | undefined
if (mappingUpdates.length > 0) {
try {
const preTable = await getTableById(data.tableId)
const preGroup = preTable?.schema.workflowGroups?.find((g) => g.id === data.groupId)
const targetWorkflowId = data.workflowId ?? preGroup?.workflowId
if (targetWorkflowId) {
resolvedForWorkflowId = targetWorkflowId
const [
{ loadWorkflowFromNormalizedTables },
{ flattenWorkflowOutputs },
{ columnTypeForLeaf },
] = await Promise.all([
import('@/lib/workflows/persistence/utils'),
import('@/lib/workflows/blocks/flatten-outputs'),
import('@/lib/table/column-naming'),
])
const normalized = await loadWorkflowFromNormalizedTables(targetWorkflowId)
if (normalized) {
const blocks = Object.values(normalized.blocks ?? {}).map((b) => ({
id: b.id,
type: b.type,
name: b.name,
triggerMode: (b as { triggerMode?: boolean }).triggerMode,
subBlocks: b.subBlocks as Record<string, unknown> | undefined,
}))
const flattened = flattenWorkflowOutputs(blocks, normalized.edges ?? [])
const flatByKey = new Map(flattened.map((f) => [`${f.blockId}::${f.path}`, f]))
for (const u of mappingUpdates) {
const match = flatByKey.get(`${u.blockId}::${u.path}`)
if (!match) continue
const newType = columnTypeForLeaf(match.leafType)
if (newType) remapLeafTypeByColumn.set(u.columnName, newType)
}
}
}
} catch (err) {
logger.warn(
`[${requestId}] Could not resolve new leaf types for remap on group ${data.groupId}; leaving column types unchanged:`,
err
)
}
}
const { updatedTable, added, remappedColumnIds, newOutputs, previousAutoRun } =
await withLockedTable(data.tableId, async (table, trx) => {
await setTableTxTimeouts(trx, { statementMs: 60_000 })
const schema = table.schema
const groups = schema.workflowGroups ?? []
const groupIndex = groups.findIndex((g) => g.id === data.groupId)
if (groupIndex === -1) {
throw new Error(`Workflow group "${data.groupId}" not found`)
}
const group = groups[groupIndex]
// Normalize every caller-supplied column reference to its stable id, so
// the diff/splice/clear logic below operates uniformly in id-space (the
// row-data storage key). New output columns get ids first; then output
// `columnName`, deps, input mappings, and mapping-update targets are
// remapped name → id. Callers that already pass ids are unaffected.
const newColDefs = (data.newOutputColumns ?? []).map((col) =>
col.id ? col : { ...col, id: generateColumnId() }
)
const idByName = new Map(
[...schema.columns, ...newColDefs].map((c) => [c.name, getColumnId(c)])
)
const remapRef = (ref: string) => idByName.get(ref) ?? ref
const outputsInput = data.outputs?.map((o) => ({ ...o, columnName: remapRef(o.columnName) }))
const dependenciesInput = data.dependencies
? { columns: data.dependencies.columns?.map(remapRef) }
: undefined
const inputMappingsInput = data.inputMappings?.map((m) => ({
...m,
columnName: remapRef(m.columnName),
}))
const mappingUpdatesNorm = mappingUpdates.map((u) => ({
...u,
columnName: remapRef(u.columnName),
}))
// Re-key the out-of-lock leaf-type resolution to ids to match.
const remapLeafTypeById = new Map<string, ColumnDefinition['type']>()
for (const [name, type] of remapLeafTypeByColumn) remapLeafTypeById.set(remapRef(name), type)
// Apply `mappingUpdates` first: each entry repoints an existing output's
// `(blockId, path)` while preserving the column. We patch the **old** view
// of outputs so the downstream `(blockId, path)`-keyed diff doesn't see the
// swap as a remove+add. The corresponding row data is cleared after the
// schema write so stale values from the old source don't linger.
const remappedColumnIds = new Set<string>()
// Per-column type override (keyed by id) resolved (out-of-lock) from the
// new mapping's leaf type. Only populated when a remap actually changes
// the column's type against the fresh schema.
const remappedColumnTypes = new Map<string, ColumnDefinition['type']>()
let oldOutputs = group.outputs
if (mappingUpdatesNorm.length > 0) {
const updateById = new Map(mappingUpdatesNorm.map((u) => [u.columnName, u]))
for (const u of mappingUpdatesNorm) {
const exists = oldOutputs.some((o) => o.columnName === u.columnName)
if (!exists) {
throw new Error(
`Mapping update for unknown column "${u.columnName}" (group ${data.groupId}).`
)
}
}
oldOutputs = oldOutputs.map((o) => {
const u = updateById.get(o.columnName)
if (!u) return o
remappedColumnIds.add(o.columnName)
return { ...o, blockId: u.blockId, path: u.path }
})
// Only apply the out-of-lock leaf-type resolution if the group still
// points at the workflow we resolved against. If a concurrent writer
// changed `workflowId` between phase 1 and now, those types are stale —
// leave column types unchanged (best-effort, same as a resolution
// failure) rather than stamping types from the old workflow.
const finalWorkflowId = data.workflowId ?? group.workflowId
if (remapLeafTypeById.size > 0 && resolvedForWorkflowId !== finalWorkflowId) {
logger.warn(
`[${requestId}] Workflow group "${data.groupId}" workflowId changed between leaf-type resolution and apply; leaving remapped column types unchanged.`
)
} else {
const colById = new Map(schema.columns.map((c) => [getColumnId(c), c]))
for (const u of mappingUpdatesNorm) {
const newType = remapLeafTypeById.get(u.columnName)
if (!newType) continue
const oldType = colById.get(u.columnName)?.type
if (newType !== oldType) {
remappedColumnTypes.set(u.columnName, newType)
}
}
}
}
// If the caller passed `outputs`, that's the new full set. If only
// `mappingUpdates` was sent, the new set is the remapped old set.
const newOutputs = outputsInput ?? oldOutputs
// Enrichment outputs all share empty `blockId`/`path`, so keying on those
// alone collapses every sibling to one entry (dropping columns on diff). Key
// on the registry `outputId` when present; fall back to `blockId::path` for
// workflow outputs.
const oldKey = (o: WorkflowGroupOutput) =>
o.outputId ? `out::${o.outputId}` : `${o.blockId}::${o.path}`
const oldByKey = new Map(oldOutputs.map((o) => [oldKey(o), o]))
const newByKey = new Map(newOutputs.map((o) => [oldKey(o), o]))
const removed = oldOutputs.filter((o) => !newByKey.has(oldKey(o)))
const added = newOutputs.filter((o) => !oldByKey.has(oldKey(o)))
const newColById = new Map(newColDefs.map((c) => [getColumnId(c), c]))
for (const out of added) {
if (!newColById.has(out.columnName)) {
throw new Error(
`Missing column definition for new output "${out.columnName}" (group ${data.groupId}).`
)
}
}
const removedColumnIds = new Set(removed.map((o) => o.columnName))
let nextColumns = schema.columns
.filter((c) => !removedColumnIds.has(getColumnId(c)))
.map((c) => {
const newType = remappedColumnTypes.get(getColumnId(c))
return newType ? { ...c, type: newType } : c
})
if (newColDefs.length > 0) {
// Splice the new column defs into the group's contiguous run rather than
// appending at the end. The desired in-group order is `newOutputs` (the
// sidebar's BFS-of-the-workflow ordering); we walk it, anchor at the first
// surviving sibling's index in `nextColumns`, and emit each output's
// column def in turn.
const groupColIds = new Set(newOutputs.map((o) => o.columnName))
const firstGroupIdx = nextColumns.findIndex((c) => groupColIds.has(getColumnId(c)))
const anchorIdx = firstGroupIdx === -1 ? nextColumns.length : firstGroupIdx
const orderedGroupCols: ColumnDefinition[] = []
for (const out of newOutputs) {
const fresh = newColById.get(out.columnName)
if (fresh) {
orderedGroupCols.push(fresh)
} else {
const existing = nextColumns.find((c) => getColumnId(c) === out.columnName)
if (existing) orderedGroupCols.push(existing)
}
}
const remaining = nextColumns.filter((c) => !groupColIds.has(getColumnId(c)))
nextColumns = [
...remaining.slice(0, anchorIdx),
...orderedGroupCols,
...remaining.slice(anchorIdx),
]
}
const updatedGroup: WorkflowGroup = {
...group,
workflowId: data.workflowId ?? group.workflowId,
name: data.name ?? group.name,
dependencies: dependenciesInput ?? group.dependencies,
outputs: newOutputs,
...(inputMappingsInput !== undefined ? { inputMappings: inputMappingsInput } : {}),
...(data.deploymentMode !== undefined ? { deploymentMode: data.deploymentMode } : {}),
...(data.type !== undefined ? { type: data.type } : {}),
...(data.autoRun !== undefined ? { autoRun: data.autoRun } : {}),
}
// Removed outputs may be referenced as deps by sibling groups; strip those
// refs so we don't leave dangling-column deps that fail schema validation.
const nextGroups = groups
.map((g, i) => (i === groupIndex ? updatedGroup : g))
.map((g) => (g.id === updatedGroup.id ? g : stripGroupDeps(g, removedColumnIds)))
const updatedSchema: TableSchema = {
...schema,
columns: nextColumns,
workflowGroups: nextGroups,
}
// `columnOrder` (column ids) mirrors the schema layout. Drop removed
// columns, then splice the new ones in at the same anchor as `nextColumns`
// so the table renders them inside the group's contiguous run.
let updatedColumnOrder = table.metadata?.columnOrder?.filter(
(id) => !removedColumnIds.has(id)
)
if (updatedColumnOrder && newColDefs.length > 0) {
const newColIds = new Set(newColDefs.map(getColumnId))
const orderWithoutNew = updatedColumnOrder.filter((id) => !newColIds.has(id))
const groupColIds = new Set(newOutputs.map((o) => o.columnName))
const orderedGroupIds = newOutputs.map((o) => o.columnName)
const firstGroupOrderIdx = orderWithoutNew.findIndex((id) => groupColIds.has(id))
const anchorOrderIdx =
firstGroupOrderIdx === -1 ? orderWithoutNew.length : firstGroupOrderIdx
const remainingOrder = orderWithoutNew.filter((id) => !groupColIds.has(id))
updatedColumnOrder = [
...remainingOrder.slice(0, anchorOrderIdx),
...orderedGroupIds,
...remainingOrder.slice(anchorOrderIdx),
]
}
assertValidSchema(updatedSchema, updatedColumnOrder)
const updatedMetadata: TableMetadata | null =
updatedColumnOrder && table.metadata
? { ...table.metadata, columnOrder: updatedColumnOrder }
: table.metadata
? { ...table.metadata }
: null
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
for (const id of removedColumnIds) {
await trx.execute(
sql`UPDATE user_table_rows SET data = data - ${id}::text WHERE table_id = ${data.tableId} AND data ? ${id}::text`
)
}
// Remapped columns: clear stale values in-tx so rows the backfill can't
// repopulate (no log, no matching span output) end up empty rather than
// retaining the previous mapping's value. The backfill below then writes
// the new mapping's value into rows where it can find one.
for (const id of remappedColumnIds) {
if (removedColumnIds.has(id)) continue
await trx.execute(
sql`UPDATE user_table_rows SET data = data - ${id}::text WHERE table_id = ${data.tableId} AND data ? ${id}::text`
)
}
logger.info(
`[${requestId}] Updated workflow group "${data.groupId}" in table ${data.tableId} (added=${added.length}, removed=${removed.length}, remapped=${remappedColumnIds.size})`
)
const updatedTable: TableDefinition = {
...table,
schema: updatedSchema,
metadata: updatedMetadata,
updatedAt: now,
}
return {
updatedTable,
added,
remappedColumnIds,
newOutputs,
previousAutoRun: group.autoRun,
}
})
// Backfill from saved execution logs so already-completed group runs surface
// the schema changes without re-running the workflow. Two passes:
// - added outputs (new columns): never overwrite hand-edited values.
// - remapped outputs (existing column re-pointed): overwrite, since the
// new mapping is the source of truth and the user expects the cell to
// refresh to the new output's value.
// Small tables backfill inline-awaited (response returns with consistent
// data); large ones run as a background job. A failed backfill is logged
// but doesn't fail the request — the schema change has already committed.
// Lazy import: backfill-runner closes a cycle back to this module.
const { maybeBackfillGroupOutputs } = await import('@/lib/table/backfill-runner')
if (added.length > 0) {
try {
await maybeBackfillGroupOutputs({
table: updatedTable,
groupId: data.groupId,
outputs: added,
overwrite: false,
requestId,
actorUserId: data.actorUserId,
})
} catch (err) {
logger.warn(
`[${requestId}] Backfill from execution logs failed for ${data.tableId} group ${data.groupId}:`,
err
)
}
}
if (remappedColumnIds.size > 0) {
const remappedOutputs = newOutputs.filter((o) => remappedColumnIds.has(o.columnName))
try {
await maybeBackfillGroupOutputs({
table: updatedTable,
groupId: data.groupId,
outputs: remappedOutputs,
overwrite: true,
requestId,
actorUserId: data.actorUserId,
})
} catch (err) {
logger.warn(
`[${requestId}] Remap backfill from execution logs failed for ${data.tableId} group ${data.groupId}:`,
err
)
}
}
// autoRun toggled false → true: fire deps-satisfied rows now via the
// dispatcher. Mirrors the post-add path so re-enabling auto-fire doesn't
// require manual run clicks for rows that are already eligible.
if (previousAutoRun === false && data.autoRun === true) {
void runWorkflowColumn({
tableId: updatedTable.id,
workspaceId: updatedTable.workspaceId,
mode: 'new',
isManualRun: false,
groupIds: [data.groupId],
requestId,
triggeredByUserId: data.actorUserId,
}).catch((err) =>
logger.error(`[${requestId}] auto-dispatch (updateWorkflowGroup autoRun=true) failed:`, err)
)
}
return updatedTable
}
/**
* Adds a single output to an existing workflow group. Mirrors `addTableColumn`
* for plain columns: one canonical op, one column created, type inferred from
* the workflow's flattened outputs (`leafType` for `(blockId, path)`). The
* column is spliced into the group's contiguous run so the table renders the
* new output next to its siblings.
*/
export async function addWorkflowGroupOutput(
data: {
tableId: string
groupId: string
blockId: string
path: string
/** Optional override; defaults to a slug derived from `path`. */
columnName?: string
/** The member adding the output — billed/gated for any backfill-triggered re-run. */
actorUserId?: string | null
},
requestId: string
): Promise<TableDefinition> {
// Phase 1 (no lock): load the workflow and resolve the pickable output plus
// its execution-order index. This depends only on the workflow graph (which
// is stable), so it runs OFF the advisory-lock critical section — holding the
// lock during this DB load would make concurrent adders on the same table
// time out waiting (the Mothership fan-out this fix targets). Phase 2
// re-validates that the group still maps to the same workflow under the lock.
const preTable = await getTableById(data.tableId)
if (!preTable) throw new Error('Table not found')
const preGroup = (preTable.schema.workflowGroups ?? []).find((g) => g.id === data.groupId)
if (!preGroup) {
throw new Error(`Workflow group "${data.groupId}" not found`)
}
const workflowId = preGroup.workflowId
const [
{ loadWorkflowFromNormalizedTables },
{ flattenWorkflowOutputs, getBlockExecutionOrder },
{ columnTypeForLeaf, deriveOutputColumnName },
] = await Promise.all([
import('@/lib/workflows/persistence/utils'),
import('@/lib/workflows/blocks/flatten-outputs'),
import('@/lib/table/column-naming'),
])
const normalized = await loadWorkflowFromNormalizedTables(workflowId)
if (!normalized) {
throw new Error(`Workflow ${workflowId} not found`)
}
const blocks = Object.values(normalized.blocks ?? {}).map((b) => ({
id: b.id,
type: b.type,
name: b.name,
triggerMode: (b as { triggerMode?: boolean }).triggerMode,
subBlocks: b.subBlocks as Record<string, unknown> | undefined,
}))
const flattened = flattenWorkflowOutputs(blocks, normalized.edges ?? [])
const match = flattened.find((f) => f.blockId === data.blockId && f.path === data.path)
if (!match) {
throw new Error(
`Output ${data.blockId}::${data.path} is not a valid pickable output on workflow ${workflowId}`
)
}
const newColumnType = columnTypeForLeaf(match.leafType)
const distances = getBlockExecutionOrder(blocks, normalized.edges ?? [])
const flatIndex = new Map(flattened.map((f, i) => [`${f.blockId}::${f.path}`, i]))
// Phase 2 (locked): re-read fresh, validate against the current schema, and
// write. The critical section holds no I/O — just the in-memory splice + the
// schema UPDATE — so concurrent adders queue behind it quickly.
const { updatedTable, newOutput } = await withLockedTable(data.tableId, async (table, trx) => {
const schema = table.schema
const groups = schema.workflowGroups ?? []
const groupIndex = groups.findIndex((g) => g.id === data.groupId)
if (groupIndex === -1) {
throw new Error(`Workflow group "${data.groupId}" not found`)
}
const group = groups[groupIndex]
if (group.workflowId !== workflowId) {
throw new Error(
`Workflow group "${data.groupId}" was remapped to a different workflow concurrently; retry the add.`
)
}
if (group.outputs.some((o) => o.blockId === data.blockId && o.path === data.path)) {
throw new Error(
`Workflow group "${data.groupId}" already has an output at ${data.blockId}::${data.path}`
)
}
const taken = new Set(schema.columns.map((c) => c.name))
const columnName = data.columnName ?? deriveOutputColumnName(data.path, taken)
if (!NAME_PATTERN.test(columnName)) {
throw new Error(`Invalid column name "${columnName}". Must satisfy ${NAME_PATTERN.source}.`)
}
if (taken.has(columnName)) {
throw new Error(`Column "${columnName}" already exists`)
}
if (schema.columns.length + 1 > TABLE_LIMITS.MAX_COLUMNS_PER_TABLE) {
throw new Error(
`Adding a column would exceed the maximum (${TABLE_LIMITS.MAX_COLUMNS_PER_TABLE}).`
)
}
const newColDef: ColumnDefinition = {
id: generateColumnId(),
name: columnName,
type: newColumnType,
required: false,
unique: false,
workflowGroupId: data.groupId,
}
const newColumnId = getColumnId(newColDef)
const newOutput: WorkflowGroupOutput = {
blockId: data.blockId,
path: data.path,
columnName: newColumnId,
}
// Sort all of the group's outputs (existing + new) in workflow execution
// order: BFS distance from the start block ASC, with discovery order as
// tiebreak. This matches what the column-sidebar does at create time, so
// columns from the same workflow always read in the order their blocks run
// — regardless of whether they were added at create time or one-by-one.
const groupColIdsBefore = new Set(group.outputs.map((o) => o.columnName))
const orderKey = (o: { blockId: string; path: string }) => {
const d = distances[o.blockId]
const dist = d === undefined || d < 0 ? Number.POSITIVE_INFINITY : d
const idx = flatIndex.get(`${o.blockId}::${o.path}`) ?? Number.POSITIVE_INFINITY
return [dist, idx] as const
}
const allGroupOutputs = [...group.outputs, newOutput].sort((a, b) => {
const [da, ia] = orderKey(a)
const [db, ib] = orderKey(b)
return da !== db ? da - db : ia - ib
})
const orderedGroupColIds = allGroupOutputs.map((o) => o.columnName)
const updatedGroup: WorkflowGroup = {
...group,
outputs: allGroupOutputs,
}
const nextGroups = groups.map((g, i) => (i === groupIndex ? updatedGroup : g))
// Splice the new column run into nextColumns: keep the columns outside the
// group where they were, replace the group's contiguous run with the
// BFS-ordered list. Anchor at the position of the first existing sibling
// (or append if the group was empty).
const colById = new Map(schema.columns.map((c) => [getColumnId(c), c]))
const orderedGroupCols: ColumnDefinition[] = orderedGroupColIds.map((id) => {
if (id === newColumnId) return newColDef
const existing = colById.get(id)
if (!existing) {
throw new Error(`Internal: column "${id}" missing while splicing group outputs`)
}
return existing
})
const remainingCols = schema.columns.filter((c) => !groupColIdsBefore.has(getColumnId(c)))
const firstGroupIdx = schema.columns.findIndex((c) => groupColIdsBefore.has(getColumnId(c)))
const colAnchor = firstGroupIdx === -1 ? remainingCols.length : firstGroupIdx
const nextColumns = [
...remainingCols.slice(0, colAnchor),
...orderedGroupCols,
...remainingCols.slice(colAnchor),
]
const updatedSchema: TableSchema = {
...schema,
columns: nextColumns,
workflowGroups: nextGroups,
}
const updatedColumnOrder = table.metadata?.columnOrder
? (() => {
const orderWithoutGroup = table.metadata!.columnOrder!.filter(
(id) => !groupColIdsBefore.has(id)
)
const firstGroupOrderIdx = table.metadata!.columnOrder!.findIndex((id) =>
groupColIdsBefore.has(id)
)
const orderAnchor =
firstGroupOrderIdx === -1 ? orderWithoutGroup.length : firstGroupOrderIdx
return [
...orderWithoutGroup.slice(0, orderAnchor),
...orderedGroupColIds,
...orderWithoutGroup.slice(orderAnchor),
]
})()
: undefined
assertValidSchema(updatedSchema, updatedColumnOrder)
const updatedMetadata: TableMetadata | null =
updatedColumnOrder && table.metadata
? { ...table.metadata, columnOrder: updatedColumnOrder }
: table.metadata
? { ...table.metadata }
: null
const now = new Date()
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
logger.info(
`[${requestId}] Added output "${columnName}" (${newColDef.type}) to workflow group "${data.groupId}" in table ${data.tableId}`
)
const updatedTable: TableDefinition = {
...table,
schema: updatedSchema,
metadata: updatedMetadata,
updatedAt: now,
}
return { updatedTable, newOutput }
})
// Backfill from saved execution logs — same flow `updateWorkflowGroup`
// uses for added outputs. Reads each row's saved trace spans for the
// group's executionId and writes the new output's value back. Existing
// rows that have hand-edited values are left alone (overwrite: false).
// Cheap compared to re-running the workflow on every row, which is what
// an earlier version of this code did — that mistakenly fanned out N
// workflow-group-cell jobs and burned compute the user didn't ask for.
// Small tables backfill inline; large ones run as a background job.
// Lazy import: backfill-runner closes a cycle back to this module.
try {
const { maybeBackfillGroupOutputs } = await import('@/lib/table/backfill-runner')
await maybeBackfillGroupOutputs({
table: updatedTable,
groupId: data.groupId,
outputs: [newOutput],
overwrite: false,
requestId,
actorUserId: data.actorUserId,
})
} catch (err) {
logger.warn(
`[${requestId}] Backfill from execution logs failed for ${data.tableId} group ${data.groupId} after adding output "${newOutput.columnName}":`,
err
)
}
return updatedTable
}
/**
* Removes a single output from a workflow group. Drops the bound column and
* strips the value from every row's `data` JSONB. If the output is the
* group's last, the empty group is left in place — drop it explicitly with
* `deleteWorkflowGroup` if needed.
*/
export async function deleteWorkflowGroupOutput(
data: { tableId: string; groupId: string; columnName: string },
requestId: string
): Promise<TableDefinition> {
return withLockedTable(data.tableId, async (table, trx) => {
const schema = table.schema
const groups = schema.workflowGroups ?? []
const groupIndex = groups.findIndex((g) => g.id === data.groupId)
if (groupIndex === -1) {
throw new Error(`Workflow group "${data.groupId}" not found`)
}
const group = groups[groupIndex]
// `data.columnName` may be a column id (first-party) or display name
// (mothership/legacy); resolve to the stable id used everywhere below.
const targetColumn = schema.columns.find((c) => columnMatchesRef(c, data.columnName))
const columnId = targetColumn ? getColumnId(targetColumn) : data.columnName
if (!group.outputs.some((o) => o.columnName === columnId)) {
throw new Error(
`Workflow group "${data.groupId}" has no output bound to column "${data.columnName}"`
)
}
const updatedGroup: WorkflowGroup = {
...group,
outputs: group.outputs.filter((o) => o.columnName !== columnId),
}
const nextGroups = groups.map((g, i) => (i === groupIndex ? updatedGroup : g))
const nextColumns = schema.columns.filter((c) => getColumnId(c) !== columnId)
const updatedSchema: TableSchema = {
...schema,
columns: nextColumns,
workflowGroups: nextGroups,
}
const updatedColumnOrder = table.metadata?.columnOrder?.filter((id) => id !== columnId)
assertValidSchema(updatedSchema, updatedColumnOrder)
const updatedMetadata: TableMetadata | null =
updatedColumnOrder && table.metadata
? { ...table.metadata, columnOrder: updatedColumnOrder }
: table.metadata
? { ...table.metadata }
: null
const now = new Date()
await setTableTxTimeouts(trx, { statementMs: 60_000 })
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
await trx.execute(
sql`UPDATE user_table_rows SET data = data - ${columnId}::text WHERE table_id = ${data.tableId} AND data ? ${columnId}::text`
)
logger.info(
`[${requestId}] Removed output "${data.columnName}" from workflow group "${data.groupId}" in table ${data.tableId}`
)
return { ...table, schema: updatedSchema, metadata: updatedMetadata, updatedAt: now }
})
}
/**
* Removes a workflow group plus all its output columns. Also strips the
* group's `executions[groupId]` entry from every row.
*/
export async function deleteWorkflowGroup(
data: DeleteWorkflowGroupData,
requestId: string
): Promise<TableDefinition> {
return withLockedTable(data.tableId, async (table, trx) => {
const schema = table.schema
const groups = schema.workflowGroups ?? []
const group = groups.find((g) => g.id === data.groupId)
if (!group) {
throw new Error(`Workflow group "${data.groupId}" not found`)
}
const removedColumnIds = new Set(group.outputs.map((o) => o.columnName))
// Removed group's output columns may be referenced as deps by sibling groups.
// Strip those refs so we don't leave dangling-column deps behind.
const nextGroups = groups
.filter((g) => g.id !== data.groupId)
.map((g) => stripGroupDeps(g, removedColumnIds))
const updatedSchema: TableSchema = {
...schema,
columns: schema.columns.filter((c) => !removedColumnIds.has(getColumnId(c))),
workflowGroups: nextGroups,
}
const updatedColumnOrder = table.metadata?.columnOrder?.filter(
(id) => !removedColumnIds.has(id)
)
assertValidSchema(updatedSchema, updatedColumnOrder)
const updatedMetadata: TableMetadata | null =
updatedColumnOrder && table.metadata
? { ...table.metadata, columnOrder: updatedColumnOrder }
: table.metadata
? { ...table.metadata }
: null
const now = new Date()
await setTableTxTimeouts(trx, { statementMs: 60_000 })
await trx
.update(userTableDefinitions)
.set({ schema: updatedSchema, metadata: updatedMetadata, updatedAt: now })
.where(eq(userTableDefinitions.id, data.tableId))
for (const id of removedColumnIds) {
await trx.execute(
sql`UPDATE user_table_rows SET data = data - ${id}::text WHERE table_id = ${data.tableId} AND data ? ${id}::text`
)
}
await stripGroupExecutions(trx, data.tableId, [data.groupId])
logger.info(
`[${requestId}] Deleted workflow group "${data.groupId}" from table ${data.tableId}`
)
return {
...table,
schema: updatedSchema,
metadata: updatedMetadata,
updatedAt: now,
}
})
}
+1 -1
View File
@@ -615,7 +615,7 @@ async function pruneWorkflowGroupOutputsIfStillActive(params: {
if (!versionRow) return
const { pruneStaleWorkflowGroupOutputs } = await import('@/lib/table/service')
const { pruneStaleWorkflowGroupOutputs } = await import('@/lib/table/workflow-groups/service')
await pruneStaleWorkflowGroupOutputs({
workflowId: params.workflowId,
workspaceId: params.workspaceId,