From 63c59c9cf3e50fd4990fd0a71e9780e19466e3e5 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 10 Oct 2026 12:49:29 -0700 Subject: [PATCH 1/2] fix(realtime): scope batch-add upserts to the caller's workflow Block, edge, and subflow ids are global primary keys, so the ON CONFLICT (id) DO UPDATE in batch-add-blocks and batch-add-edges could overwrite a row owned by another workflow. Route every upsert through upsertWorkflowRows, which confines the update to the socket's workflow with setWhere and aborts the transaction when RETURNING shows a skipped conflict. --- apps/realtime/src/database/operations.ts | 203 +++++++++++------------ 1 file changed, 95 insertions(+), 108 deletions(-) diff --git a/apps/realtime/src/database/operations.ts b/apps/realtime/src/database/operations.ts index 7d6f444eb85..1ce03d2492c 100644 --- a/apps/realtime/src/database/operations.ts +++ b/apps/realtime/src/database/operations.ts @@ -38,7 +38,7 @@ import { normalizeWorkflowEdgeSourceHandle, normalizeWorkflowEdgeTargetHandle, } from '@sim/workflow-types/workflow' -import { and, eq, inArray, isNull, or, sql } from 'drizzle-orm' +import { and, eq, inArray, isNull, or, type SQL, sql } from 'drizzle-orm' import { drizzle } from 'drizzle-orm/postgres-js' import postgres from 'postgres' import { env } from '@/env' @@ -77,6 +77,45 @@ function canonicalizeEdgeAddCandidate(edge: EdgeAddCandidate): EdgeAddCandidate } } +type WorkflowScopedTable = typeof workflowBlocks | typeof workflowEdges | typeof workflowSubflows + +/** + * Inserts rows into a workflow-scoped table, updating rows that already exist + * in `workflowId`. Row ids are unique across every workflow, so an upsert keyed + * on `id` alone would overwrite another workflow's row: `setWhere` confines + * the update to this workflow, and a conflict it skips leaves a row unreturned, + * which aborts the transaction instead of dropping that row silently. + */ +async function upsertWorkflowRows( + tx: any, + table: T, + workflowId: string, + values: Array<{ id: string; workflowId: string } & Record>, + set: Partial> +) { + const written: Array<{ id: string }> = await tx + .insert(table) + .values(values) + .onConflictDoUpdate({ target: table.id, set, setWhere: eq(table.workflowId, workflowId) }) + .returning({ id: table.id }) + + if (written.length !== values.length) { + throw new Error('One or more ids already belong to another workflow') + } +} + +const SUBFLOW_UPSERT_SET = { + config: sql`excluded.config`, + updatedAt: sql`now()`, +} + +const EDGE_UPSERT_SET = { + sourceBlockId: sql`excluded.source_block_id`, + targetBlockId: sql`excluded.target_block_id`, + sourceHandle: sql`excluded.source_handle`, + targetHandle: sql`excluded.target_handle`, +} + interface FilterEdgesForPersistResult { safeEdges: T[] droppedCounts: Record @@ -965,79 +1004,61 @@ async function handleBlocksOperationTx( } }) - await tx - .insert(workflowBlocks) - .values(blockValues) - .onConflictDoUpdate({ - target: workflowBlocks.id, - set: { - type: sql`excluded.type`, - name: sql`excluded.name`, - positionX: sql`excluded.position_x`, - positionY: sql`excluded.position_y`, - enabled: sql`excluded.enabled`, - horizontalHandles: sql`excluded.horizontal_handles`, - advancedMode: sql`excluded.advanced_mode`, - triggerMode: sql`excluded.trigger_mode`, - errorEnabled: sql`excluded.error_enabled`, - retry: sql`excluded.retry`, - locked: sql`excluded.locked`, - height: sql`excluded.height`, - subBlocks: sql`excluded.sub_blocks`, - outputs: sql`excluded.outputs`, - data: sql`excluded.data`, - updatedAt: sql`now()`, - }, - }) + await upsertWorkflowRows(tx, workflowBlocks, workflowId, blockValues, { + type: sql`excluded.type`, + name: sql`excluded.name`, + positionX: sql`excluded.position_x`, + positionY: sql`excluded.position_y`, + enabled: sql`excluded.enabled`, + horizontalHandles: sql`excluded.horizontal_handles`, + advancedMode: sql`excluded.advanced_mode`, + triggerMode: sql`excluded.trigger_mode`, + errorEnabled: sql`excluded.error_enabled`, + retry: sql`excluded.retry`, + locked: sql`excluded.locked`, + height: sql`excluded.height`, + subBlocks: sql`excluded.sub_blocks`, + outputs: sql`excluded.outputs`, + data: sql`excluded.data`, + updatedAt: sql`now()`, + }) // Create subflow entries for loop/parallel blocks (skip if already in payload) const loopIds = new Set(loops ? Object.keys(loops) : []) const parallelIds = new Set(parallels ? Object.keys(parallels) : []) + const defaultSubflows = [] for (const block of allowedBlocks) { const blockId = block.id as string if (block.type === 'loop' && !loopIds.has(blockId)) { - await tx - .insert(workflowSubflows) - .values({ - id: blockId, - workflowId, - type: 'loop', - config: { - loopType: 'for', - iterations: DEFAULT_LOOP_ITERATIONS, - nodes: [], - }, - }) - .onConflictDoUpdate({ - target: workflowSubflows.id, - set: { - config: sql`excluded.config`, - updatedAt: sql`now()`, - }, - }) + defaultSubflows.push({ + id: blockId, + workflowId, + type: 'loop', + config: { loopType: 'for', iterations: DEFAULT_LOOP_ITERATIONS, nodes: [] }, + }) } else if (block.type === 'parallel' && !parallelIds.has(blockId)) { - await tx - .insert(workflowSubflows) - .values({ - id: blockId, - workflowId, - type: 'parallel', - config: { - parallelType: 'count', - count: DEFAULT_PARALLEL_COUNT, - batchSize: DEFAULT_PARALLEL_BATCH_SIZE, - nodes: [], - }, - }) - .onConflictDoUpdate({ - target: workflowSubflows.id, - set: { - config: sql`excluded.config`, - updatedAt: sql`now()`, - }, - }) + defaultSubflows.push({ + id: blockId, + workflowId, + type: 'parallel', + config: { + parallelType: 'count', + count: DEFAULT_PARALLEL_COUNT, + batchSize: DEFAULT_PARALLEL_BATCH_SIZE, + nodes: [], + }, + }) } } + if (defaultSubflows.length > 0) { + await upsertWorkflowRows( + tx, + workflowSubflows, + workflowId, + defaultSubflows, + SUBFLOW_UPSERT_SET + ) + } const parentIds = new Set() for (const block of allowedBlocks) { @@ -1087,18 +1108,7 @@ async function handleBlocksOperationTx( targetHandle: normalizeWorkflowEdgeTargetHandle(edge.targetHandle), })) - await tx - .insert(workflowEdges) - .values(edgeValues) - .onConflictDoUpdate({ - target: workflowEdges.id, - set: { - sourceBlockId: sql`excluded.source_block_id`, - targetBlockId: sql`excluded.target_block_id`, - sourceHandle: sql`excluded.source_handle`, - targetHandle: sql`excluded.target_handle`, - }, - }) + await upsertWorkflowRows(tx, workflowEdges, workflowId, edgeValues, EDGE_UPSERT_SET) } } @@ -1110,16 +1120,7 @@ async function handleBlocksOperationTx( config: loop as Record, })) - await tx - .insert(workflowSubflows) - .values(loopValues) - .onConflictDoUpdate({ - target: workflowSubflows.id, - set: { - config: sql`excluded.config`, - updatedAt: sql`now()`, - }, - }) + await upsertWorkflowRows(tx, workflowSubflows, workflowId, loopValues, SUBFLOW_UPSERT_SET) } if (parallels && Object.keys(parallels).length > 0) { @@ -1130,16 +1131,13 @@ async function handleBlocksOperationTx( config: parallel as Record, })) - await tx - .insert(workflowSubflows) - .values(parallelValues) - .onConflictDoUpdate({ - target: workflowSubflows.id, - set: { - config: sql`excluded.config`, - updatedAt: sql`now()`, - }, - }) + await upsertWorkflowRows( + tx, + workflowSubflows, + workflowId, + parallelValues, + SUBFLOW_UPSERT_SET + ) } logger.info(`Successfully batch added blocks to workflow ${workflowId}`) @@ -1791,18 +1789,7 @@ async function handleEdgesOperationTx( targetHandle: normalizeWorkflowEdgeTargetHandle(edge.targetHandle), })) - await tx - .insert(workflowEdges) - .values(edgeValues) - .onConflictDoUpdate({ - target: workflowEdges.id, - set: { - sourceBlockId: sql`excluded.source_block_id`, - targetBlockId: sql`excluded.target_block_id`, - sourceHandle: sql`excluded.source_handle`, - targetHandle: sql`excluded.target_handle`, - }, - }) + await upsertWorkflowRows(tx, workflowEdges, workflowId, edgeValues, EDGE_UPSERT_SET) logger.debug(`Batch added ${safeEdges.length} edges to workflow ${workflowId}`) break From ec84c553b353de58d734aff07042782ee647ae6b Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 10 Oct 2026 12:55:58 -0700 Subject: [PATCH 2/2] fix(realtime): type the scoped upsert helper against its table --- apps/realtime/src/database/operations.ts | 17 +++++++++-------- 1 file changed, 9 insertions(+), 8 deletions(-) diff --git a/apps/realtime/src/database/operations.ts b/apps/realtime/src/database/operations.ts index 1ce03d2492c..7dff7dffe72 100644 --- a/apps/realtime/src/database/operations.ts +++ b/apps/realtime/src/database/operations.ts @@ -39,6 +39,7 @@ import { normalizeWorkflowEdgeTargetHandle, } from '@sim/workflow-types/workflow' import { and, eq, inArray, isNull, or, type SQL, sql } from 'drizzle-orm' +import type { PgInsertValue } from 'drizzle-orm/pg-core' import { drizzle } from 'drizzle-orm/postgres-js' import postgres from 'postgres' import { env } from '@/env' @@ -87,13 +88,13 @@ type WorkflowScopedTable = typeof workflowBlocks | typeof workflowEdges | typeof * which aborts the transaction instead of dropping that row silently. */ async function upsertWorkflowRows( - tx: any, + tx: DbOrTx, table: T, workflowId: string, - values: Array<{ id: string; workflowId: string } & Record>, + values: PgInsertValue[], set: Partial> ) { - const written: Array<{ id: string }> = await tx + const written = await tx .insert(table) .values(values) .onConflictDoUpdate({ target: table.id, set, setWhere: eq(table.workflowId, workflowId) }) @@ -107,14 +108,14 @@ async function upsertWorkflowRows( const SUBFLOW_UPSERT_SET = { config: sql`excluded.config`, updatedAt: sql`now()`, -} +} as const const EDGE_UPSERT_SET = { sourceBlockId: sql`excluded.source_block_id`, targetBlockId: sql`excluded.target_block_id`, sourceHandle: sql`excluded.source_handle`, targetHandle: sql`excluded.target_handle`, -} +} as const interface FilterEdgesForPersistResult { safeEdges: T[] @@ -988,8 +989,8 @@ async function handleBlocksOperationTx( workflowId, type: block.type as string, name: block.name as string, - positionX: (block.position as { x: number; y: number }).x, - positionY: (block.position as { x: number; y: number }).y, + positionX: String((block.position as { x: number; y: number }).x), + positionY: String((block.position as { x: number; y: number }).y), data: (block.data as Record | undefined) || {}, subBlocks: mergedSubBlocks, outputs: (block.outputs as Record) || {}, @@ -999,7 +1000,7 @@ async function handleBlocksOperationTx( triggerMode: (block.triggerMode as boolean) ?? false, errorEnabled: (block.errorEnabled as boolean) ?? false, retry: (block.retry as Record | undefined) ?? null, - height: (block.height as number) || 0, + height: String((block.height as number) || 0), locked: (block.locked as boolean) ?? false, } })