From e1f58221983f2c7e07075ee02c7d44173081d8fb Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 3 Oct 2026 12:29:38 -0700 Subject: [PATCH] fix(webhooks): check existing path owners before claiming on deploy --- .../registration-store.integration.ts | 141 ++++++++++++++++++ .../lib/webhooks/registration-store.test.ts | 4 + apps/sim/lib/webhooks/registration-store.ts | 15 +- 3 files changed, 158 insertions(+), 2 deletions(-) create mode 100644 apps/sim/lib/webhooks/registration-store.integration.ts diff --git a/apps/sim/lib/webhooks/registration-store.integration.ts b/apps/sim/lib/webhooks/registration-store.integration.ts new file mode 100644 index 00000000000..d46fa413659 --- /dev/null +++ b/apps/sim/lib/webhooks/registration-store.integration.ts @@ -0,0 +1,141 @@ +/** + * Stable webhook registration against real PostgreSQL: a deploy must not claim + * a path another workflow already serves through an unclaimed legacy row. + */ +import { db } from '@sim/db' +import { + user, + webhook, + webhookPathClaim, + workflow, + workflowDeploymentOperation, + workflowDeploymentVersion, + workspace, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { WebhookPathClaimConflictError } from '@/lib/webhooks/path-claims' +import { + prepareWebhookRegistrationIntents, + type WebhookRegistrationOperationFence, +} from '@/lib/webhooks/registration-store' + +const owner = `registration-store-owner-${generateId()}` +const workspaceId = generateId() +const victimWorkflow = generateId() +const deployingWorkflow = generateId() +const victimVersion = generateId() +const deployingVersion = generateId() +const legacyPath = `legacy-${generateId()}` +const freePath = `free-${generateId()}` + +const fence: WebhookRegistrationOperationFence = { + workflowId: deployingWorkflow, + operationId: generateId(), + generation: 1, + deploymentVersionId: deployingVersion, +} + +function desiredFor(path: string) { + return { + blockId: 'trigger', + provider: 'generic', + path, + routingKey: null, + providerConfig: {}, + configFingerprint: `fp-${path}`, + } +} + +beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: owner, + name: 'Registration Store', + email: `${owner}@registration-store.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Registration Store', + ownerId: owner, + billedAccountUserId: owner, + }) + await db.insert(workflow).values( + [victimWorkflow, deployingWorkflow].map((id) => ({ + id, + userId: owner, + workspaceId, + name: id, + lastSynced: now, + createdAt: now, + updatedAt: now, + isDeployed: id === victimWorkflow, + })) + ) + const emptyState = { blocks: {}, edges: [], loops: {}, parallels: {} } + await db.insert(workflowDeploymentVersion).values([ + { + id: victimVersion, + workflowId: victimWorkflow, + version: 1, + isActive: true, + state: emptyState, + }, + { id: deployingVersion, workflowId: deployingWorkflow, version: 1, state: emptyState }, + ]) + await db.insert(webhook).values({ + id: generateId(), + workflowId: victimWorkflow, + deploymentVersionId: victimVersion, + blockId: 'trigger', + path: legacyPath, + provider: 'generic', + providerConfig: {}, + }) + await db.insert(workflowDeploymentOperation).values({ + id: fence.operationId, + workflowId: deployingWorkflow, + deploymentVersionId: deployingVersion, + version: 1, + action: 'deploy', + protocolVersion: 2, + generation: fence.generation, + status: 'preparing', + requestHash: generateId(), + actorId: owner, + }) +}) + +afterAll(async () => { + await db.delete(workspace).where(eq(workspace.id, workspaceId)) + await db.delete(user).where(eq(user.id, owner)) +}) + +async function claimOwner(path: string) { + const [claim] = await db + .select({ workflowId: webhookPathClaim.workflowId }) + .from(webhookPathClaim) + .where(eq(webhookPathClaim.path, path)) + return claim?.workflowId ?? null +} + +describe('prepareWebhookRegistrationIntents path ownership', () => { + it.each([legacyPath, `/${legacyPath}/`])( + 'refuses to claim %s while another workflow serves it unclaimed', + async (path) => { + await expect( + prepareWebhookRegistrationIntents({ fence, desired: [desiredFor(path)] }) + ).rejects.toBeInstanceOf(WebhookPathClaimConflictError) + expect(await claimOwner(legacyPath)).toBeNull() + } + ) + + it('claims a path no other workflow serves', async () => { + await prepareWebhookRegistrationIntents({ fence, desired: [desiredFor(freePath)] }) + expect(await claimOwner(freePath)).toBe(deployingWorkflow) + }) +}) diff --git a/apps/sim/lib/webhooks/registration-store.test.ts b/apps/sim/lib/webhooks/registration-store.test.ts index be5f1d84f31..fcd3cbadbc5 100644 --- a/apps/sim/lib/webhooks/registration-store.test.ts +++ b/apps/sim/lib/webhooks/registration-store.test.ts @@ -14,6 +14,10 @@ vi.mock('@/lib/webhooks/path-claims', () => ({ claimWebhookPath: mockClaimWebhookPath, })) +vi.mock('@/lib/webhooks/utils.server', () => ({ + findConflictingWebhookPathOwner: vi.fn().mockResolvedValue(null), +})) + vi.mock('@/lib/workflows/persistence/deployment-operations', () => ({ isDeploymentOperationCurrent: mockIsDeploymentOperationCurrent, setDeploymentTxTimeouts: vi.fn(), diff --git a/apps/sim/lib/webhooks/registration-store.ts b/apps/sim/lib/webhooks/registration-store.ts index eeeb1e350f6..78617e1f7b4 100644 --- a/apps/sim/lib/webhooks/registration-store.ts +++ b/apps/sim/lib/webhooks/registration-store.ts @@ -9,13 +9,14 @@ import { generateShortId } from '@sim/utils/id' import { isPlainRecord } from '@sim/utils/object' import type { DbOrTx } from '@sim/workflow-persistence/types' import { and, eq, exists, gt, inArray, isNull, lt, lte, notExists, sql } from 'drizzle-orm' -import { claimWebhookPath } from '@/lib/webhooks/path-claims' +import { claimWebhookPath, WebhookPathClaimConflictError } from '@/lib/webhooks/path-claims' import { projectDesiredWebhookProviderConfig } from '@/lib/webhooks/provider-subscriptions' import { fingerprintDesiredWebhookRegistration, normalizeWebhookRegistrationPath, } from '@/lib/webhooks/registration-identity' import { planWebhookRegistrationReconciliation } from '@/lib/webhooks/registration-reconciliation' +import { findConflictingWebhookPathOwner } from '@/lib/webhooks/utils.server' import type { DeploymentOperationStatus } from '@/lib/workflows/deployment-lifecycle' import { isDeploymentOperationCurrent, @@ -230,8 +231,18 @@ export async function prepareWebhookRegistrationIntents(input: { for (const desired of input.desired) { if (desired.path) { + // Unclaimed legacy rows of other workflows still own their path. + const path = normalizeWebhookRegistrationPath(desired.path) ?? desired.path + const conflictingOwner = await findConflictingWebhookPathOwner({ + path, + workflowId: input.fence.workflowId, + tx, + }) + if (conflictingOwner) { + throw new WebhookPathClaimConflictError(path, conflictingOwner) + } await claimWebhookPath(tx, { - path: desired.path, + path, workflowId: input.fence.workflowId, generation: input.fence.generation, })