Skip to content

Commit ccce7f7

Browse files
authored
fix(webhooks): skip duplicate webhook deliveries already in progress instead of polling (#8712)
* fix(webhooks): skip duplicate in-progress deliveries and bound the polling claim lease * fix(webhooks): drop the polling lease change; detached poll passes are not bounded by the route * refactor(idempotency): return the in-progress wait as a thunk instead of overloading on policy * chore(idempotency): log the skipped duplicate once, at the caller * docs(idempotency): describe dedup as bounded by the claim lease and result TTL * docs(idempotency): reflow executeWithIdempotency TSDoc
1 parent e82689a commit ccce7f7

5 files changed

Lines changed: 207 additions & 32 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: 82 additions & 2 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', () => ({
@@ -22,7 +24,11 @@ import {
2224
webhookIdempotency,
2325
} from '@/lib/core/idempotency/service'
2426

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

2733
const SEVEN_DAYS_SECONDS = 60 * 60 * 24 * 7
2834

@@ -32,6 +38,8 @@ afterEach(() => {
3238

3339
beforeEach(() => {
3440
resetDbChainMock()
41+
redisEvalMock.mockReset()
42+
redisGetMock.mockReset()
3543
})
3644

3745
describe('IdempotencyService.createWebhookIdempotencyKey', () => {
@@ -305,6 +313,78 @@ describe('IdempotencyService in-progress deadlines', () => {
305313
})
306314
})
307315

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

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

Lines changed: 74 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,18 @@ export interface IdempotencyExecutionOptions {
7373
inProgressExpiresAt?: number
7474
}
7575

76+
/**
77+
* Outcome of {@link IdempotencyService.executeOrSkipInProgress}. `resolved` covers both a
78+
* fresh run and a replayed completed result; `in-progress` means another holder owns a live
79+
* claim on the key and this caller did nothing.
80+
*/
81+
export type IdempotentExecution<T> = { outcome: 'resolved'; result: T } | { outcome: 'in-progress' }
82+
83+
/** An in-progress outcome carries the wait on the live holder, so each caller decides whether to take it. */
84+
type ClaimedExecution<T> =
85+
| { outcome: 'resolved'; result: T }
86+
| { outcome: 'in-progress'; wait: () => Promise<T> }
87+
7688
export interface AtomicClaimResult {
7789
claimed: boolean
7890
existingResult?: ProcessingResult
@@ -559,13 +571,60 @@ export class IdempotencyService {
559571
return deleted.length > 0
560572
}
561573

574+
/**
575+
* Runs `operation` once per key while its claim lease and stored result last. A caller that
576+
* finds another holder's live claim polls until that holder finishes and returns (or rethrows)
577+
* its outcome, so use this when the caller's own response depends on the result (e.g. Stripe
578+
* must not get a 2xx before the first attempt settles).
579+
*/
562580
async executeWithIdempotency<T>(
563581
provider: string,
564582
identifier: string,
565583
operation: () => Promise<T>,
566-
additionalContext?: Record<string, any>,
584+
additionalContext?: Record<string, unknown>,
567585
options?: IdempotencyExecutionOptions
568586
): Promise<T> {
587+
const execution = await this.execute(
588+
provider,
589+
identifier,
590+
operation,
591+
additionalContext,
592+
options
593+
)
594+
return execution.outcome === 'resolved' ? execution.result : execution.wait()
595+
}
596+
597+
/**
598+
* Like {@link executeWithIdempotency}, but returns `{ outcome: 'in-progress' }` immediately
599+
* when another holder owns a live claim instead of polling until it finishes. Use it when
600+
* nobody consumes the duplicate's result, so waiting would only pin the caller's worker and
601+
* leases. Completed and failed keys behave exactly as in `executeWithIdempotency`.
602+
*/
603+
async executeOrSkipInProgress<T>(
604+
provider: string,
605+
identifier: string,
606+
operation: () => Promise<T>,
607+
additionalContext?: Record<string, unknown>,
608+
options?: IdempotencyExecutionOptions
609+
): Promise<IdempotentExecution<T>> {
610+
const execution = await this.execute(
611+
provider,
612+
identifier,
613+
operation,
614+
additionalContext,
615+
options
616+
)
617+
if (execution.outcome === 'resolved') return execution
618+
return { outcome: 'in-progress' }
619+
}
620+
621+
private async execute<T>(
622+
provider: string,
623+
identifier: string,
624+
operation: () => Promise<T>,
625+
additionalContext: Record<string, unknown> | undefined,
626+
options: IdempotencyExecutionOptions | undefined
627+
): Promise<ClaimedExecution<T>> {
569628
const claimResult = await this.atomicallyClaim(provider, identifier, additionalContext, options)
570629

571630
if (!claimResult.claimed) {
@@ -576,7 +635,7 @@ export class IdempotencyService {
576635
if (existingResult.success === false) {
577636
throw new Error(existingResult.error || 'Previous operation failed')
578637
}
579-
return existingResult.result as T
638+
return { outcome: 'resolved', result: existingResult.result as T }
580639
}
581640

582641
if (existingResult?.status === 'failed') {
@@ -585,30 +644,28 @@ export class IdempotencyService {
585644
observedResult: existingResult,
586645
observedValue: claimResult.observedValue,
587646
})
588-
return this.executeWithIdempotency(
589-
provider,
590-
identifier,
591-
operation,
592-
additionalContext,
593-
options
594-
)
647+
return this.execute(provider, identifier, operation, additionalContext, options)
595648
}
596649
logger.info(`Previous operation failed for: ${claimResult.normalizedKey}`)
597650
throw new Error(existingResult.error || 'Previous operation failed')
598651
}
599652

600653
if (existingResult?.status === 'in-progress') {
601-
logger.info(`Waiting for in-progress operation: ${claimResult.normalizedKey}`)
602-
return await this.waitForResult<T>(
603-
claimResult.normalizedKey,
604-
claimResult.storageMethod,
654+
const { normalizedKey, storageMethod } = claimResult
655+
const deadline =
605656
existingResult.inProgressExpiresAt ??
606-
(existingResult.startedAt ?? Date.now()) + this.config.inProgressTtlSeconds * 1000
607-
)
657+
(existingResult.startedAt ?? Date.now()) + this.config.inProgressTtlSeconds * 1000
658+
return {
659+
outcome: 'in-progress',
660+
wait: () => {
661+
logger.info(`Waiting for in-progress operation: ${normalizedKey}`)
662+
return this.waitForResult<T>(normalizedKey, storageMethod, deadline)
663+
},
664+
}
608665
}
609666

610667
if (existingResult) {
611-
return existingResult.result as T
668+
return { outcome: 'resolved', result: existingResult.result as T }
612669
}
613670

614671
throw new Error(`Unexpected state: key claimed but no existing result found`)
@@ -630,7 +687,7 @@ export class IdempotencyService {
630687
)
631688

632689
logger.debug(`Successfully completed operation: ${claimResult.normalizedKey}`)
633-
return result
690+
return { outcome: 'resolved', result }
634691
} catch (error) {
635692
const errorMessage = getErrorMessage(error, 'Unknown error')
636693

0 commit comments

Comments
 (0)