Skip to content

Commit 86fcaf6

Browse files
authored
fix(webhooks): check existing path owners before claiming on deploy (#8605)
1 parent f67222e commit 86fcaf6

3 files changed

Lines changed: 158 additions & 2 deletions

File tree

Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,141 @@
1+
/**
2+
* Stable webhook registration against real PostgreSQL: a deploy must not claim
3+
* a path another workflow already serves through an unclaimed legacy row.
4+
*/
5+
import { db } from '@sim/db'
6+
import {
7+
user,
8+
webhook,
9+
webhookPathClaim,
10+
workflow,
11+
workflowDeploymentOperation,
12+
workflowDeploymentVersion,
13+
workspace,
14+
} from '@sim/db/schema'
15+
import { generateId } from '@sim/utils/id'
16+
import { eq } from 'drizzle-orm'
17+
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
18+
import { WebhookPathClaimConflictError } from '@/lib/webhooks/path-claims'
19+
import {
20+
prepareWebhookRegistrationIntents,
21+
type WebhookRegistrationOperationFence,
22+
} from '@/lib/webhooks/registration-store'
23+
24+
const owner = `registration-store-owner-${generateId()}`
25+
const workspaceId = generateId()
26+
const victimWorkflow = generateId()
27+
const deployingWorkflow = generateId()
28+
const victimVersion = generateId()
29+
const deployingVersion = generateId()
30+
const legacyPath = `legacy-${generateId()}`
31+
const freePath = `free-${generateId()}`
32+
33+
const fence: WebhookRegistrationOperationFence = {
34+
workflowId: deployingWorkflow,
35+
operationId: generateId(),
36+
generation: 1,
37+
deploymentVersionId: deployingVersion,
38+
}
39+
40+
function desiredFor(path: string) {
41+
return {
42+
blockId: 'trigger',
43+
provider: 'generic',
44+
path,
45+
routingKey: null,
46+
providerConfig: {},
47+
configFingerprint: `fp-${path}`,
48+
}
49+
}
50+
51+
beforeAll(async () => {
52+
const now = new Date()
53+
await db.insert(user).values({
54+
id: owner,
55+
name: 'Registration Store',
56+
email: `${owner}@registration-store.test`,
57+
emailVerified: true,
58+
createdAt: now,
59+
updatedAt: now,
60+
})
61+
await db.insert(workspace).values({
62+
id: workspaceId,
63+
name: 'Registration Store',
64+
ownerId: owner,
65+
billedAccountUserId: owner,
66+
})
67+
await db.insert(workflow).values(
68+
[victimWorkflow, deployingWorkflow].map((id) => ({
69+
id,
70+
userId: owner,
71+
workspaceId,
72+
name: id,
73+
lastSynced: now,
74+
createdAt: now,
75+
updatedAt: now,
76+
isDeployed: id === victimWorkflow,
77+
}))
78+
)
79+
const emptyState = { blocks: {}, edges: [], loops: {}, parallels: {} }
80+
await db.insert(workflowDeploymentVersion).values([
81+
{
82+
id: victimVersion,
83+
workflowId: victimWorkflow,
84+
version: 1,
85+
isActive: true,
86+
state: emptyState,
87+
},
88+
{ id: deployingVersion, workflowId: deployingWorkflow, version: 1, state: emptyState },
89+
])
90+
await db.insert(webhook).values({
91+
id: generateId(),
92+
workflowId: victimWorkflow,
93+
deploymentVersionId: victimVersion,
94+
blockId: 'trigger',
95+
path: legacyPath,
96+
provider: 'generic',
97+
providerConfig: {},
98+
})
99+
await db.insert(workflowDeploymentOperation).values({
100+
id: fence.operationId,
101+
workflowId: deployingWorkflow,
102+
deploymentVersionId: deployingVersion,
103+
version: 1,
104+
action: 'deploy',
105+
protocolVersion: 2,
106+
generation: fence.generation,
107+
status: 'preparing',
108+
requestHash: generateId(),
109+
actorId: owner,
110+
})
111+
})
112+
113+
afterAll(async () => {
114+
await db.delete(workspace).where(eq(workspace.id, workspaceId))
115+
await db.delete(user).where(eq(user.id, owner))
116+
})
117+
118+
async function claimOwner(path: string) {
119+
const [claim] = await db
120+
.select({ workflowId: webhookPathClaim.workflowId })
121+
.from(webhookPathClaim)
122+
.where(eq(webhookPathClaim.path, path))
123+
return claim?.workflowId ?? null
124+
}
125+
126+
describe('prepareWebhookRegistrationIntents path ownership', () => {
127+
it.each([legacyPath, `/${legacyPath}/`])(
128+
'refuses to claim %s while another workflow serves it unclaimed',
129+
async (path) => {
130+
await expect(
131+
prepareWebhookRegistrationIntents({ fence, desired: [desiredFor(path)] })
132+
).rejects.toBeInstanceOf(WebhookPathClaimConflictError)
133+
expect(await claimOwner(legacyPath)).toBeNull()
134+
}
135+
)
136+
137+
it('claims a path no other workflow serves', async () => {
138+
await prepareWebhookRegistrationIntents({ fence, desired: [desiredFor(freePath)] })
139+
expect(await claimOwner(freePath)).toBe(deployingWorkflow)
140+
})
141+
})

‎apps/sim/lib/webhooks/registration-store.test.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,10 @@ vi.mock('@/lib/webhooks/path-claims', () => ({
1414
claimWebhookPath: mockClaimWebhookPath,
1515
}))
1616

17+
vi.mock('@/lib/webhooks/utils.server', () => ({
18+
findConflictingWebhookPathOwner: vi.fn().mockResolvedValue(null),
19+
}))
20+
1721
vi.mock('@/lib/workflows/persistence/deployment-operations', () => ({
1822
isDeploymentOperationCurrent: mockIsDeploymentOperationCurrent,
1923
setDeploymentTxTimeouts: vi.fn(),

‎apps/sim/lib/webhooks/registration-store.ts‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,13 +9,14 @@ import { generateShortId } from '@sim/utils/id'
99
import { isPlainRecord } from '@sim/utils/object'
1010
import type { DbOrTx } from '@sim/workflow-persistence/types'
1111
import { and, eq, exists, gt, inArray, isNull, lt, lte, notExists, sql } from 'drizzle-orm'
12-
import { claimWebhookPath } from '@/lib/webhooks/path-claims'
12+
import { claimWebhookPath, WebhookPathClaimConflictError } from '@/lib/webhooks/path-claims'
1313
import { projectDesiredWebhookProviderConfig } from '@/lib/webhooks/provider-subscriptions'
1414
import {
1515
fingerprintDesiredWebhookRegistration,
1616
normalizeWebhookRegistrationPath,
1717
} from '@/lib/webhooks/registration-identity'
1818
import { planWebhookRegistrationReconciliation } from '@/lib/webhooks/registration-reconciliation'
19+
import { findConflictingWebhookPathOwner } from '@/lib/webhooks/utils.server'
1920
import type { DeploymentOperationStatus } from '@/lib/workflows/deployment-lifecycle'
2021
import {
2122
isDeploymentOperationCurrent,
@@ -230,8 +231,18 @@ export async function prepareWebhookRegistrationIntents(input: {
230231

231232
for (const desired of input.desired) {
232233
if (desired.path) {
234+
// Unclaimed legacy rows of other workflows still own their path.
235+
const path = normalizeWebhookRegistrationPath(desired.path) ?? desired.path
236+
const conflictingOwner = await findConflictingWebhookPathOwner({
237+
path,
238+
workflowId: input.fence.workflowId,
239+
tx,
240+
})
241+
if (conflictingOwner) {
242+
throw new WebhookPathClaimConflictError(path, conflictingOwner)
243+
}
233244
await claimWebhookPath(tx, {
234-
path: desired.path,
245+
path,
235246
workflowId: input.fence.workflowId,
236247
generation: input.fence.generation,
237248
})

0 commit comments

Comments
 (0)