Skip to content

Commit cd340d6

Browse files
committed
fix(plane): preserve webhook delivery during replacement
1 parent 3510a76 commit cd340d6

10 files changed

Lines changed: 193 additions & 40 deletions

File tree

‎apps/sim/app/api/webhooks/route.test.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -589,7 +589,7 @@ describe('POST /api/webhooks subscription replacement recovery', () => {
589589
}
590590
queueCurrentRows()
591591
expect((await POST(replacementRequest())).status).toBe(500)
592-
expect(external.size).toBe(priorCleanupFails ? 2 : 1)
592+
expect(external.size).toBe(2)
593593
const retainedConfig = toRecord(persisted.providerConfig)
594594
if (priorCleanupFails) {
595595
expect(retainedConfig.projectId).toBe('next')
@@ -600,11 +600,11 @@ describe('POST /api/webhooks subscription replacement recovery', () => {
600600
webhookSecret: 'old-secret',
601601
})
602602
expect(external.get('external-old')?.active).toBe(true)
603-
expect(external.get(String(retainedConfig.externalId))?.active).toBe(false)
603+
expect(external.get(String(retainedConfig.externalId))?.active).toBe(true)
604604
} else {
605605
expect(retainedConfig.projectId).toBe('next')
606606
expect(retainedConfig.webhookSecret).toBe('new-secret')
607-
expect(external.has('external-old')).toBe(false)
607+
expect(external.get('external-old')?.active).toBe(true)
608608
expect(external.get(String(retainedConfig.externalId))?.active).toBe(false)
609609
}
610610
cleanupUnavailable = false

‎apps/sim/app/api/webhooks/route.ts‎

Lines changed: 33 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -512,6 +512,20 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
512512
if (existingWebhook && shouldRecreateSubscription) {
513513
const pendingPrevious = toRecord(existingWebhook.providerConfig?.previousSubscription)
514514
if (typeof pendingPrevious.provider === 'string') {
515+
if (existingWebhook.providerConfig?.subscriptionActivationPending === true) {
516+
await activateExternalWebhookSubscription(
517+
request,
518+
existingWebhook,
519+
workflowRecord,
520+
userId,
521+
requestId
522+
)
523+
existingWebhook.providerConfig.subscriptionActivationPending = false
524+
await db
525+
.update(webhook)
526+
.set({ providerConfig: existingWebhook.providerConfig })
527+
.where(eq(webhook.id, existingWebhook.id))
528+
}
515529
await cleanupExternalWebhook(
516530
{
517531
...existingWebhook,
@@ -643,25 +657,6 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
643657
}
644658

645659
if (savedWebhook) {
646-
const previousSubscription = toRecord(configToSave.previousSubscription)
647-
if (typeof previousSubscription.provider === 'string') {
648-
await cleanupExternalWebhook(
649-
{
650-
...savedWebhook,
651-
provider: previousSubscription.provider,
652-
providerConfig: previousSubscription.providerConfig,
653-
},
654-
workflowRecord,
655-
requestId,
656-
{ throwOnError: true }
657-
)
658-
configToSave.previousSubscription = undefined
659-
await db
660-
.update(webhook)
661-
.set({ providerConfig: configToSave })
662-
.where(eq(webhook.id, savedWebhook.id))
663-
savedWebhook.providerConfig = configToSave
664-
}
665660
if (
666661
!existingWebhook ||
667662
shouldRecreateSubscription ||
@@ -683,6 +678,25 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
683678
savedWebhook.providerConfig = configToSave
684679
}
685680
}
681+
const previousSubscription = toRecord(configToSave.previousSubscription)
682+
if (typeof previousSubscription.provider === 'string') {
683+
await cleanupExternalWebhook(
684+
{
685+
...savedWebhook,
686+
provider: previousSubscription.provider,
687+
providerConfig: previousSubscription.providerConfig,
688+
},
689+
workflowRecord,
690+
requestId,
691+
{ throwOnError: true }
692+
)
693+
configToSave.previousSubscription = undefined
694+
await db
695+
.update(webhook)
696+
.set({ providerConfig: configToSave })
697+
.where(eq(webhook.id, savedWebhook.id))
698+
savedWebhook.providerConfig = configToSave
699+
}
686700
}
687701

688702
if (existingWebhook && shouldRecreateSubscription && !deferPreviousCleanup) {

‎apps/sim/lib/webhooks/provider-subscriptions.test.ts‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,34 @@ describe('cleanupExternalWebhook', () => {
146146
}
147147
)
148148

149+
it('cleans the tracked prior subscription when the replacement has no deletion handler', async () => {
150+
const external = new Set(['previous-registration'])
151+
mockGetProviderHandler.mockImplementation((provider: string) =>
152+
provider === 'plane'
153+
? {
154+
deleteSubscription: async () => {
155+
external.delete('previous-registration')
156+
},
157+
}
158+
: {}
159+
)
160+
await cleanupExternalWebhook(
161+
{
162+
provider: 'webhook',
163+
providerConfig: {
164+
previousSubscription: {
165+
provider: 'plane',
166+
providerConfig: { externalId: 'previous-registration' },
167+
},
168+
},
169+
},
170+
{ userId: 'user-1', workspaceId: 'workspace-1' },
171+
'request-1',
172+
{ throwOnError: true }
173+
)
174+
expect(external.size).toBe(0)
175+
})
176+
149177
it('resolves {{ENV_VAR}} references before deleting the provider subscription', async () => {
150178
const deleteSubscription = vi.fn().mockResolvedValue(undefined)
151179
mockGetProviderHandler.mockReturnValue({ deleteSubscription })

‎apps/sim/lib/webhooks/provider-subscriptions.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -282,8 +282,9 @@ export async function cleanupExternalWebhook(
282282
): Promise<void> {
283283
const provider = webhook.provider as string
284284
const handler = getProviderHandler(provider)
285+
const previousSubscription = toRecord(toRecord(webhook.providerConfig).previousSubscription)
285286

286-
if (!handler.deleteSubscription) {
287+
if (!handler.deleteSubscription && typeof previousSubscription.provider !== 'string') {
287288
return
288289
}
289290

@@ -303,7 +304,6 @@ export async function cleanupExternalWebhook(
303304
{ envVars, onResolved: (name, value) => secrets.set(name, value) }
304305
)
305306
resolvedProviderConfig = resolvedWebhook.providerConfig
306-
const previousSubscription = toRecord(resolvedProviderConfig.previousSubscription)
307307
if (typeof previousSubscription.provider === 'string') {
308308
await cleanupExternalWebhook(
309309
{
@@ -317,6 +317,8 @@ export async function cleanupExternalWebhook(
317317
)
318318
}
319319

320+
if (!handler.deleteSubscription) return
321+
320322
/** Workspace archival precedes provider cleanup; routing still uses its canonical owner. */
321323
await withResourceOutboundScope(
322324
{ workspaceId },

‎apps/sim/lib/webhooks/providers/plane.test.ts‎

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,48 @@ describe('Plane signed webhook delivery', () => {
6969
).toBe(401)
7070
})
7171

72+
it.each(['webhook-a', 'untracked-webhook'])(
73+
'accepts the tracked previous secret only for its subscription: %s',
74+
async (webhookId) => {
75+
const raw = JSON.stringify({ ...V2, webhook_id: webhookId })
76+
const signature = createHmac('sha256', SECRET).update(raw).digest('hex')
77+
const ctx = authContext(raw, signature, 'replacement-secret')
78+
ctx.providerConfig.previousSubscription = {
79+
provider: 'plane',
80+
providerConfig: { externalId: 'webhook-a', webhookSecret: SECRET },
81+
}
82+
if (!planeHandler.verifyAuth) throw new Error('Plane authentication is missing')
83+
const response = await planeHandler.verifyAuth(ctx)
84+
expect(response?.status ?? 200).toBe(webhookId === 'webhook-a' ? 200 : 401)
85+
}
86+
)
87+
88+
it('keeps the previous event scope until replacement activation succeeds', async () => {
89+
if (!planeHandler.matchEvent) throw new Error('Plane event filtering is missing')
90+
const config = {
91+
triggerId: 'plane_workitem_created',
92+
projectId: 'project-b',
93+
subscriptionActivationPending: true,
94+
previousSubscription: {
95+
provider: 'plane',
96+
providerConfig: {
97+
externalId: 'webhook-a',
98+
triggerId: 'plane_workitem_updated',
99+
projectId: 'project-a',
100+
},
101+
},
102+
}
103+
expect(await planeHandler.matchEvent(matchContext(V1, config))).toBe(true)
104+
expect(
105+
await planeHandler.matchEvent(
106+
matchContext(V1, { ...config, subscriptionActivationPending: false })
107+
)
108+
).toBe(false)
109+
expect(
110+
await planeHandler.matchEvent(matchContext({ ...V1, webhook_id: 'replacement' }, config))
111+
).toBe(false)
112+
})
113+
72114
it('deduplicates v2 retries by event ID while separating distinct events', () => {
73115
if (!planeHandler.extractIdempotencyId) throw new Error('Plane deduplication is missing')
74116
const original = planeHandler.extractIdempotencyId(V2)

‎apps/sim/lib/webhooks/providers/plane.ts‎

Lines changed: 42 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -75,39 +75,67 @@ async function deletePlaneWebhook(
7575
throw new Error(`Plane webhook deletion failed (HTTP ${response.status})`)
7676
}
7777

78+
function previousPlaneDeliveryConfig(
79+
providerConfig: Record<string, unknown>,
80+
payload: unknown
81+
): Record<string, unknown> | undefined {
82+
const previous = toRecord(providerConfig.previousSubscription)
83+
const config = toRecord(previous.providerConfig)
84+
return previous.provider === 'plane' &&
85+
typeof config.externalId === 'string' &&
86+
toRecord(payload).webhook_id === config.externalId
87+
? config
88+
: undefined
89+
}
90+
91+
const verifyPlaneSignature = createHmacVerifier({
92+
configKey: 'webhookSecret',
93+
headerName: 'X-Plane-Signature',
94+
providerLabel: 'Plane',
95+
requireSecret: true,
96+
validateFn: (secret, signature, rawBody) =>
97+
typeof secret === 'string' &&
98+
secret.length > 0 &&
99+
/^[a-f0-9]{64}$/.test(signature) &&
100+
safeCompare(hmacSha256Hex(rawBody, secret), signature),
101+
})
102+
78103
export const planeHandler: WebhookProviderHandler = {
79-
verifyAuth: createHmacVerifier({
80-
configKey: 'webhookSecret',
81-
headerName: 'X-Plane-Signature',
82-
providerLabel: 'Plane',
83-
requireSecret: true,
84-
validateFn: (secret, signature, rawBody) =>
85-
typeof secret === 'string' &&
86-
secret.length > 0 &&
87-
/^[a-f0-9]{64}$/.test(signature) &&
88-
safeCompare(hmacSha256Hex(rawBody, secret), signature),
89-
}),
104+
verifyAuth(ctx) {
105+
let providerConfig = ctx.providerConfig
106+
if (providerConfig.previousSubscription) {
107+
try {
108+
const payload: unknown = JSON.parse(ctx.rawBody)
109+
providerConfig = previousPlaneDeliveryConfig(providerConfig, payload) ?? providerConfig
110+
} catch {}
111+
}
112+
return verifyPlaneSignature({ ...ctx, providerConfig })
113+
},
90114

91115
async matchEvent({ body, providerConfig }) {
92116
const { planeEventName, PLANE_TRIGGER_EVENTS } = await import('@/triggers/plane/utils')
93117
const payload = toRecord(body)
118+
const matchConfig =
119+
(providerConfig.subscriptionActivationPending === true
120+
? previousPlaneDeliveryConfig(providerConfig, payload)
121+
: undefined) ?? providerConfig
94122
const eventName = planeEventName(body)
95123
if (!eventName) return false
96-
const triggerId = providerConfig.triggerId
124+
const triggerId = matchConfig.triggerId
97125
if (
98126
typeof triggerId === 'string' &&
99127
triggerId !== 'plane_webhook' &&
100128
PLANE_TRIGGER_EVENTS[triggerId] !== eventName
101129
)
102130
return false
103-
const workspaceId = providerConfig.workspaceId
131+
const workspaceId = matchConfig.workspaceId
104132
if (
105133
typeof workspaceId === 'string' &&
106134
workspaceId.trim() &&
107135
workspaceId.trim() !== payload.workspace_id
108136
)
109137
return false
110-
const projectId = providerConfig.projectId
138+
const projectId = matchConfig.projectId
111139
if (typeof projectId === 'string' && projectId.trim()) {
112140
const records = Array.isArray(payload.data)
113141
? payload.data

‎apps/sim/lib/workflows/deployment-outbox.test.ts‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -633,6 +633,21 @@ describe('versioned deployment preparation outbox', () => {
633633
)
634634
})
635635

636+
it('retires obsolete subscriptions even when new activation fails', async () => {
637+
mockIsDeploymentOperationCurrent.mockResolvedValue(true)
638+
mockGetDeploymentOperation.mockResolvedValue(operation({ status: 'active', completedAt: NOW }))
639+
queueTableRows(schemaMock.workflow, [
640+
{ id: 'workflow-1', name: 'Workflow', workspaceId: 'workspace-1' },
641+
])
642+
const external = new Set(['retired-subscription'])
643+
mockActivatePendingWebhookSubscriptions.mockRejectedValue(new Error('activation unavailable'))
644+
mockCleanupRetiredWebhookRegistrations.mockImplementation(async () => {
645+
external.delete('retired-subscription')
646+
})
647+
await expect(handler()(payload(), context())).rejects.toThrow('activation unavailable')
648+
expect(external.size).toBe(0)
649+
})
650+
636651
it('continues through the outbox while stale webhooks remain, then checkpoints the cleanup', async () => {
637652
mockIsDeploymentOperationCurrent.mockResolvedValue(true)
638653
mockGetDeploymentOperation.mockResolvedValue(operation({ status: 'active', completedAt: NOW }))

‎apps/sim/lib/workflows/deployment-outbox.ts‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -596,6 +596,7 @@ async function runPostActivationWork(params: {
596596
context: OutboxEventContext
597597
}): Promise<DeferredOutboxHandlerResult | undefined> {
598598
await emitPostActivationSideEffects(params)
599+
let activationFailure: Error | undefined
599600
const activationHasMore = await activatePendingWebhookSubscriptionsAfterActivation({
600601
request: new NextRequest(new URL('/api/webhooks', getBaseUrl())),
601602
fence: {
@@ -608,8 +609,10 @@ async function runPostActivationWork(params: {
608609
userId: params.payload.userId,
609610
requestId: params.payload.requestId,
610611
signal: params.context.signal,
612+
}).catch((error: unknown) => {
613+
activationFailure = toError(error)
614+
return false
611615
})
612-
if (activationHasMore) return continueOutboxHandler('webhook_activation_pending')
613616
await cleanupRetiredWebhooksForOperation({
614617
payload: params.payload,
615618
workflow: params.workflow,
@@ -622,7 +625,9 @@ async function runPostActivationWork(params: {
622625
checkpoint: params.checkpoint,
623626
context: params.context,
624627
})
625-
return cleanupComplete ? undefined : continueOutboxHandler(INACTIVE_CLEANUP_CONTINUATION_REASON)
628+
if (!cleanupComplete) return continueOutboxHandler(INACTIVE_CLEANUP_CONTINUATION_REASON)
629+
if (activationFailure) throw activationFailure
630+
return activationHasMore ? continueOutboxHandler('webhook_activation_pending') : undefined
626631
}
627632

628633
async function prepareReadinessComponent(params: {

‎apps/sim/scripts/test-plane-e2e.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -208,6 +208,11 @@ async function receiveDelivery(request: Request): Promise<Response> {
208208
}
209209
const server = publicCallback
210210
? createServer(async (incoming, outgoing) => {
211+
if (incoming.method !== 'POST') {
212+
outgoing.writeHead(405, { Allow: 'POST' })
213+
outgoing.end()
214+
return
215+
}
211216
try {
212217
const chunks: Buffer[] = []
213218
let size = 0
@@ -243,6 +248,19 @@ const server = publicCallback
243248
})
244249
: undefined
245250
server?.listen(Number(process.env.PLANE_E2E_LISTENER_PORT || 49187), '127.0.0.1')
251+
if (server) {
252+
await check('webhook ingress rejects method probes safely', async () => {
253+
for (const method of ['GET', 'HEAD']) {
254+
// boundary-raw-fetch: Exercise the local webhook ingress protocol over real HTTP.
255+
const response = await fetch(
256+
`http://127.0.0.1:${process.env.PLANE_E2E_LISTENER_PORT || 49187}/${callbackPath}`,
257+
{ method }
258+
)
259+
assert.equal(response.status, 405)
260+
}
261+
assert.equal(deliveryErrors.length, 0)
262+
})
263+
}
246264

247265
try {
248266
for (const version of ['v2', 'v1'] as const) {

‎apps/sim/tools/plane/list_labels.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -195,6 +195,7 @@ export const planeListLabelsTool: ToolConfig<PlaneListLabelsParams, PlaneListLab
195195
count: { key: 'count', type: 'boolean', required: false },
196196
fields: { key: 'fields', type: 'string', required: false },
197197
cursor: { key: 'cursor', type: 'string', required: false },
198+
parent_id__isnull: { key: 'parent_id__isnull', type: 'boolean', required: false },
198199
})
199200
)
200201
},

0 commit comments

Comments
 (0)