Skip to content

Commit f90b277

Browse files
committed
fix(webhooks): skip duplicate in-progress deliveries and bound the polling claim lease
1 parent 8fc4c1a commit f90b277

4 files changed

Lines changed: 260 additions & 40 deletions

File tree

‎apps/sim/background/webhook-execution.test.ts‎

Lines changed: 29 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ const {
3232
mockResolveWebhookRecordProviderConfig,
3333
mockExecuteWorkflowCore,
3434
mockWasExecutionFinalizedByCore,
35-
mockExecuteWithIdempotency,
35+
mockExecuteOrSkipInProgress,
3636
mockGetProviderHandler,
3737
mockSetResolvedSecretTraceRegistry,
3838
mockExecutionSnapshot,
@@ -42,7 +42,7 @@ const {
4242
mockResolveWebhookRecordProviderConfig: vi.fn(),
4343
mockExecuteWorkflowCore: vi.fn(),
4444
mockWasExecutionFinalizedByCore: vi.fn(),
45-
mockExecuteWithIdempotency: vi.fn(),
45+
mockExecuteOrSkipInProgress: vi.fn(),
4646
mockGetProviderHandler: vi.fn(() => ({})),
4747
mockSetResolvedSecretTraceRegistry: vi.fn(),
4848
mockExecutionSnapshot: vi.fn(),
@@ -84,7 +84,7 @@ vi.mock('@/lib/workflows/executor/execution-core', () => ({
8484
vi.mock('@/lib/core/idempotency', () => ({
8585
IdempotencyService: { createWebhookIdempotencyKey: vi.fn(() => 'idempotency-key') },
8686
webhookIdempotency: {
87-
executeWithIdempotency: mockExecuteWithIdempotency,
87+
executeOrSkipInProgress: mockExecuteOrSkipInProgress,
8888
},
8989
}))
9090

@@ -254,8 +254,11 @@ describe('executeWebhookJob fault vs error handling', () => {
254254
.mockResolvedValue(undefined)
255255
mockGetProviderHandler.mockReturnValue({})
256256
mockEnqueue.mockReset().mockResolvedValue('run_retry')
257-
mockExecuteWithIdempotency.mockImplementation(
258-
(_provider: string, _key: string, operation: () => Promise<unknown>) => operation()
257+
mockExecuteOrSkipInProgress.mockImplementation(
258+
async (_provider: string, _key: string, operation: () => Promise<unknown>) => ({
259+
outcome: 'resolved',
260+
result: await operation(),
261+
})
259262
)
260263
executionPreprocessingMockFns.mockPreprocessExecution.mockResolvedValue({
261264
success: true,
@@ -634,14 +637,30 @@ describe('executeWebhookJob fault vs error handling', () => {
634637
workflowId: 'workflow-1',
635638
executionId: 'original-execution',
636639
}
637-
mockExecuteWithIdempotency.mockResolvedValueOnce(cachedResult)
640+
mockExecuteOrSkipInProgress.mockResolvedValueOnce({ outcome: 'resolved', result: cachedResult })
638641

639642
await expect(executeWebhookJob(payload)).resolves.toBe(cachedResult)
640643

641644
expect(executionPreprocessingMockFns.mockPreprocessExecution).not.toHaveBeenCalled()
642645
expect(mockReleaseExecutionSlot).toHaveBeenCalledWith('execution-1')
643646
})
644647

648+
it('acknowledges a duplicate of an in-progress delivery without running it and frees its reservation', async () => {
649+
mockExecuteOrSkipInProgress.mockResolvedValueOnce({ outcome: 'in-progress' })
650+
651+
await expect(executeWebhookJob(payload)).resolves.toMatchObject({
652+
success: true,
653+
duplicate: true,
654+
workflowId: 'workflow-1',
655+
executionId: 'execution-1',
656+
})
657+
658+
expect(executionPreprocessingMockFns.mockPreprocessExecution).not.toHaveBeenCalled()
659+
expect(mockExecuteWorkflowCore).not.toHaveBeenCalled()
660+
expect(mockEnqueue).not.toHaveBeenCalled()
661+
expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1')
662+
})
663+
645664
it('rejects queued webhook work without an immutable attribution snapshot', async () => {
646665
await expect(
647666
executeWebhookJob({
@@ -705,7 +724,7 @@ describe('executeWebhookJob fault vs error handling', () => {
705724

706725
expect(redisGet).toHaveBeenCalledWith('usage:reservation:execution-1')
707726
expect(result).toMatchObject({ success: false, requeued: true })
708-
expect(mockExecuteWithIdempotency).not.toHaveBeenCalled()
727+
expect(mockExecuteOrSkipInProgress).not.toHaveBeenCalled()
709728
expect(mockExecuteWorkflowCore).not.toHaveBeenCalled()
710729
expect(loggingSessionMockFns.mockSafeCompleteWithError).not.toHaveBeenCalled()
711730
expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1')
@@ -750,7 +769,7 @@ describe('executeWebhookJob fault vs error handling', () => {
750769
})
751770

752771
expect(mockEnqueue).not.toHaveBeenCalled()
753-
expect(mockExecuteWithIdempotency).not.toHaveBeenCalled()
772+
expect(mockExecuteOrSkipInProgress).not.toHaveBeenCalled()
754773
expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1')
755774
expect(loggingSessionMockFns.mockSafeStart).toHaveBeenCalledTimes(1)
756775
expect(loggingSessionMockFns.mockSafeCompleteWithError).toHaveBeenCalledExactlyOnceWith(
@@ -784,12 +803,12 @@ describe('executeWebhookJob fault vs error handling', () => {
784803

785804
expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1')
786805
expect(mockEnqueue).not.toHaveBeenCalled()
787-
expect(mockExecuteWithIdempotency).not.toHaveBeenCalled()
806+
expect(mockExecuteOrSkipInProgress).not.toHaveBeenCalled()
788807
})
789808

790809
it('does not treat an ambiguous idempotency claim timeout as a safe setup retry', async () => {
791810
const error = new Error('Command timed out')
792-
mockExecuteWithIdempotency.mockRejectedValueOnce(error)
811+
mockExecuteOrSkipInProgress.mockRejectedValueOnce(error)
793812

794813
await expect(executeWebhookJob(payload)).rejects.toBe(error)
795814

‎apps/sim/background/webhook-execution.ts‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -589,7 +589,7 @@ export async function executeWebhookJob(
589589
)
590590
}
591591

592-
const result = await webhookIdempotency.executeWithIdempotency(
592+
const execution = await webhookIdempotency.executeOrSkipInProgress(
593593
authenticatedPayload.provider,
594594
idempotencyKey,
595595
runOperation,
@@ -604,7 +604,26 @@ export async function executeWebhookJob(
604604
if (!operationStarted) {
605605
await releaseExecutionSlot(executionId)
606606
}
607-
return result
607+
if (execution.outcome === 'in-progress') {
608+
// Ingress already acknowledged this delivery and nothing reads this job's result,
609+
// so waiting on the live holder would only pin the machine and queue slot.
610+
logger.info(`[${requestId}] Skipping duplicate webhook delivery already in progress`, {
611+
webhookId: authenticatedPayload.webhookId,
612+
workflowId: authenticatedPayload.workflowId,
613+
provider: authenticatedPayload.provider,
614+
executionId,
615+
})
616+
return {
617+
success: true,
618+
duplicate: true,
619+
workflowId: authenticatedPayload.workflowId,
620+
executionId,
621+
output: {},
622+
executedAt: new Date().toISOString(),
623+
provider: authenticatedPayload.provider,
624+
}
625+
}
626+
return execution.result
608627
})
609628
} catch (error) {
610629
await releaseExecutionSlot(executionId)

‎apps/sim/lib/core/idempotency/service.test.ts‎

Lines changed: 118 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,12 @@
1+
import { flushMicrotasks } from '@sim/testing/helpers/async'
12
import { dbChainMockFns, resetDbChainMock } from '@sim/testing/mocks/database.mock'
23
import { redisConfigMockFns } from '@sim/testing/mocks/redis-config.mock'
34
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
45

5-
const { redisDelMock, redisEvalMock } = vi.hoisted(() => ({
6+
const { redisDelMock, redisEvalMock, redisGetMock } = vi.hoisted(() => ({
67
redisDelMock: vi.fn(),
78
redisEvalMock: vi.fn(),
9+
redisGetMock: vi.fn(),
810
}))
911

1012
vi.mock('@/lib/core/storage', () => ({
@@ -18,11 +20,16 @@ vi.mock('@/lib/core/storage', () => ({
1820
import { RetryableSetupError } from '@/lib/core/errors/retryable-infrastructure'
1921
import {
2022
IdempotencyService,
23+
pollingIdempotency,
2124
WEBHOOK_IN_PROGRESS_LEASE_SECONDS,
2225
webhookIdempotency,
2326
} from '@/lib/core/idempotency/service'
2427

25-
redisConfigMockFns.mockGetRedisClient.mockReturnValue({ del: redisDelMock, eval: redisEvalMock })
28+
redisConfigMockFns.mockGetRedisClient.mockReturnValue({
29+
del: redisDelMock,
30+
eval: redisEvalMock,
31+
get: redisGetMock,
32+
})
2633

2734
const SEVEN_DAYS_SECONDS = 60 * 60 * 24 * 7
2835

@@ -32,6 +39,8 @@ afterEach(() => {
3239

3340
beforeEach(() => {
3441
resetDbChainMock()
42+
redisEvalMock.mockReset()
43+
redisGetMock.mockReset()
3544
})
3645

3746
describe('IdempotencyService.createWebhookIdempotencyKey', () => {
@@ -223,26 +232,41 @@ describe('IdempotencyService in-progress deadlines', () => {
223232
expect(WEBHOOK_IN_PROGRESS_LEASE_SECONDS).toBeLessThan(SEVEN_DAYS_SECONDS)
224233
})
225234

226-
it('leases an untimed webhook claim for the bounded lease, not the seven-day dedupe window', async () => {
227-
vi.useFakeTimers()
228-
vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z'))
229-
redisEvalMock.mockResolvedValue([1, ''])
230-
231-
await webhookIdempotency.atomicallyClaim('gmail', 'wh_1:untimed-delivery')
232-
233-
const [, , redisKey, serialized, ttlSeconds] = redisEvalMock.mock.calls[0] as [
234-
string,
235-
number,
236-
string,
237-
string,
238-
number,
239-
]
240-
expect(redisKey).toBe('idempotency:webhook:gmail:wh_1:untimed-delivery')
241-
expect(ttlSeconds).toBe(WEBHOOK_IN_PROGRESS_LEASE_SECONDS)
242-
expect(JSON.parse(serialized).inProgressExpiresAt).toBe(
243-
Date.now() + WEBHOOK_IN_PROGRESS_LEASE_SECONDS * 1000
244-
)
245-
})
235+
it.each([
236+
{
237+
name: 'webhook',
238+
service: webhookIdempotency,
239+
leaseSeconds: WEBHOOK_IN_PROGRESS_LEASE_SECONDS,
240+
dedupeSeconds: SEVEN_DAYS_SECONDS,
241+
},
242+
{
243+
name: 'polling',
244+
service: pollingIdempotency,
245+
leaseSeconds: 5 * 60,
246+
dedupeSeconds: 60 * 60 * 24 * 3,
247+
},
248+
])(
249+
'leases an untimed $name claim for its bounded lease so a crashed holder frees the key, not the dedupe window',
250+
async ({ name, service, leaseSeconds, dedupeSeconds }) => {
251+
vi.useFakeTimers()
252+
vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z'))
253+
redisEvalMock.mockResolvedValue([1, ''])
254+
255+
await service.atomicallyClaim('gmail', 'wh_1:untimed-delivery')
256+
257+
const [, , redisKey, serialized, ttlSeconds] = redisEvalMock.mock.calls[0] as [
258+
string,
259+
number,
260+
string,
261+
string,
262+
number,
263+
]
264+
expect(redisKey).toBe(`idempotency:${name}:gmail:wh_1:untimed-delivery`)
265+
expect(ttlSeconds).toBe(leaseSeconds)
266+
expect(ttlSeconds).toBeLessThan(dedupeSeconds)
267+
expect(JSON.parse(serialized).inProgressExpiresAt).toBe(Date.now() + leaseSeconds * 1000)
268+
}
269+
)
246270

247271
it('keeps the seven-day dedupe window on a completed webhook result', async () => {
248272
redisEvalMock.mockResolvedValueOnce([1, '']).mockResolvedValueOnce(1)
@@ -305,6 +329,78 @@ describe('IdempotencyService in-progress deadlines', () => {
305329
})
306330
})
307331

332+
describe('IdempotencyService duplicate of an in-progress operation', () => {
333+
const liveClaim = () =>
334+
JSON.stringify({
335+
success: false,
336+
status: 'in-progress',
337+
startedAt: Date.now(),
338+
inProgressExpiresAt: Date.now() + WEBHOOK_IN_PROGRESS_LEASE_SECONDS * 1000,
339+
claimToken: 'other-holder',
340+
})
341+
342+
it('returns in-progress at once instead of polling the live holder when skipping', async () => {
343+
vi.useFakeTimers()
344+
vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z'))
345+
const holder = liveClaim()
346+
redisEvalMock.mockResolvedValueOnce([0, holder])
347+
redisGetMock.mockResolvedValue(holder)
348+
const operation = vi.fn()
349+
const settled = vi.fn()
350+
351+
void webhookIdempotency
352+
.executeOrSkipInProgress('gmail', 'wh_1:running-delivery', operation)
353+
.then(settled, settled)
354+
await flushMicrotasks(10)
355+
356+
expect(settled).toHaveBeenCalledExactlyOnceWith({ outcome: 'in-progress' })
357+
expect(operation).not.toHaveBeenCalled()
358+
expect(redisGetMock).not.toHaveBeenCalled()
359+
})
360+
361+
it('still replays a completed result and rethrows a failed one when skipping', async () => {
362+
const service = new IdempotencyService({ forceStorage: 'redis' })
363+
redisEvalMock
364+
.mockResolvedValueOnce([
365+
0,
366+
JSON.stringify({ success: true, status: 'completed', result: 'first-run' }),
367+
])
368+
.mockResolvedValueOnce([
369+
0,
370+
JSON.stringify({ success: false, status: 'failed', error: 'first run failed' }),
371+
])
372+
373+
await expect(service.executeOrSkipInProgress('provider', 'done', vi.fn())).resolves.toEqual({
374+
outcome: 'resolved',
375+
result: 'first-run',
376+
})
377+
await expect(service.executeOrSkipInProgress('provider', 'failed', vi.fn())).rejects.toThrow(
378+
'first run failed'
379+
)
380+
})
381+
382+
it('keeps waiting for the live holder by default, so a Stripe-style caller never acknowledges early', async () => {
383+
vi.useFakeTimers()
384+
vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z'))
385+
const service = new IdempotencyService({ forceStorage: 'redis' })
386+
const holder = liveClaim()
387+
redisEvalMock.mockResolvedValueOnce([0, holder])
388+
redisGetMock
389+
.mockResolvedValueOnce(holder)
390+
.mockResolvedValueOnce(
391+
JSON.stringify({ success: true, status: 'completed', result: 'holder-result' })
392+
)
393+
const settled = vi.fn()
394+
395+
void service.executeWithIdempotency('stripe', 'evt_1', vi.fn()).then(settled, settled)
396+
await flushMicrotasks(10)
397+
expect(settled).not.toHaveBeenCalled()
398+
399+
await vi.advanceTimersByTimeAsync(1_000)
400+
expect(settled).toHaveBeenCalledExactlyOnceWith('holder-result')
401+
})
402+
})
403+
308404
describe('IdempotencyService retryable setup failures', () => {
309405
it('releases the claim instead of memoizing when the operation throws a RetryableSetupError', async () => {
310406
redisEvalMock.mockResolvedValueOnce([1, '']).mockResolvedValueOnce(1)

0 commit comments

Comments
 (0)