diff --git a/apps/sim/background/webhook-execution.test.ts b/apps/sim/background/webhook-execution.test.ts index 6d58bb37ef9..f9f32ca541d 100644 --- a/apps/sim/background/webhook-execution.test.ts +++ b/apps/sim/background/webhook-execution.test.ts @@ -32,7 +32,7 @@ const { mockResolveWebhookRecordProviderConfig, mockExecuteWorkflowCore, mockWasExecutionFinalizedByCore, - mockExecuteWithIdempotency, + mockExecuteOrSkipInProgress, mockGetProviderHandler, mockSetResolvedSecretTraceRegistry, mockExecutionSnapshot, @@ -42,7 +42,7 @@ const { mockResolveWebhookRecordProviderConfig: vi.fn(), mockExecuteWorkflowCore: vi.fn(), mockWasExecutionFinalizedByCore: vi.fn(), - mockExecuteWithIdempotency: vi.fn(), + mockExecuteOrSkipInProgress: vi.fn(), mockGetProviderHandler: vi.fn(() => ({})), mockSetResolvedSecretTraceRegistry: vi.fn(), mockExecutionSnapshot: vi.fn(), @@ -84,7 +84,7 @@ vi.mock('@/lib/workflows/executor/execution-core', () => ({ vi.mock('@/lib/core/idempotency', () => ({ IdempotencyService: { createWebhookIdempotencyKey: vi.fn(() => 'idempotency-key') }, webhookIdempotency: { - executeWithIdempotency: mockExecuteWithIdempotency, + executeOrSkipInProgress: mockExecuteOrSkipInProgress, }, })) @@ -254,8 +254,11 @@ describe('executeWebhookJob fault vs error handling', () => { .mockResolvedValue(undefined) mockGetProviderHandler.mockReturnValue({}) mockEnqueue.mockReset().mockResolvedValue('run_retry') - mockExecuteWithIdempotency.mockImplementation( - (_provider: string, _key: string, operation: () => Promise) => operation() + mockExecuteOrSkipInProgress.mockImplementation( + async (_provider: string, _key: string, operation: () => Promise) => ({ + outcome: 'resolved', + result: await operation(), + }) ) executionPreprocessingMockFns.mockPreprocessExecution.mockResolvedValue({ success: true, @@ -634,7 +637,7 @@ describe('executeWebhookJob fault vs error handling', () => { workflowId: 'workflow-1', executionId: 'original-execution', } - mockExecuteWithIdempotency.mockResolvedValueOnce(cachedResult) + mockExecuteOrSkipInProgress.mockResolvedValueOnce({ outcome: 'resolved', result: cachedResult }) await expect(executeWebhookJob(payload)).resolves.toBe(cachedResult) @@ -642,6 +645,22 @@ describe('executeWebhookJob fault vs error handling', () => { expect(mockReleaseExecutionSlot).toHaveBeenCalledWith('execution-1') }) + it('acknowledges a duplicate of an in-progress delivery without running it and frees its reservation', async () => { + mockExecuteOrSkipInProgress.mockResolvedValueOnce({ outcome: 'in-progress' }) + + await expect(executeWebhookJob(payload)).resolves.toMatchObject({ + success: true, + duplicate: true, + workflowId: 'workflow-1', + executionId: 'execution-1', + }) + + expect(executionPreprocessingMockFns.mockPreprocessExecution).not.toHaveBeenCalled() + expect(mockExecuteWorkflowCore).not.toHaveBeenCalled() + expect(mockEnqueue).not.toHaveBeenCalled() + expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1') + }) + it('rejects queued webhook work without an immutable attribution snapshot', async () => { await expect( executeWebhookJob({ @@ -705,7 +724,7 @@ describe('executeWebhookJob fault vs error handling', () => { expect(redisGet).toHaveBeenCalledWith('usage:reservation:execution-1') expect(result).toMatchObject({ success: false, requeued: true }) - expect(mockExecuteWithIdempotency).not.toHaveBeenCalled() + expect(mockExecuteOrSkipInProgress).not.toHaveBeenCalled() expect(mockExecuteWorkflowCore).not.toHaveBeenCalled() expect(loggingSessionMockFns.mockSafeCompleteWithError).not.toHaveBeenCalled() expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1') @@ -750,7 +769,7 @@ describe('executeWebhookJob fault vs error handling', () => { }) expect(mockEnqueue).not.toHaveBeenCalled() - expect(mockExecuteWithIdempotency).not.toHaveBeenCalled() + expect(mockExecuteOrSkipInProgress).not.toHaveBeenCalled() expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1') expect(loggingSessionMockFns.mockSafeStart).toHaveBeenCalledTimes(1) expect(loggingSessionMockFns.mockSafeCompleteWithError).toHaveBeenCalledExactlyOnceWith( @@ -784,12 +803,12 @@ describe('executeWebhookJob fault vs error handling', () => { expect(mockReleaseExecutionSlot).toHaveBeenCalledExactlyOnceWith('execution-1') expect(mockEnqueue).not.toHaveBeenCalled() - expect(mockExecuteWithIdempotency).not.toHaveBeenCalled() + expect(mockExecuteOrSkipInProgress).not.toHaveBeenCalled() }) it('does not treat an ambiguous idempotency claim timeout as a safe setup retry', async () => { const error = new Error('Command timed out') - mockExecuteWithIdempotency.mockRejectedValueOnce(error) + mockExecuteOrSkipInProgress.mockRejectedValueOnce(error) await expect(executeWebhookJob(payload)).rejects.toBe(error) diff --git a/apps/sim/background/webhook-execution.ts b/apps/sim/background/webhook-execution.ts index 5a1080b5a89..f66b631b425 100644 --- a/apps/sim/background/webhook-execution.ts +++ b/apps/sim/background/webhook-execution.ts @@ -589,7 +589,7 @@ export async function executeWebhookJob( ) } - const result = await webhookIdempotency.executeWithIdempotency( + const execution = await webhookIdempotency.executeOrSkipInProgress( authenticatedPayload.provider, idempotencyKey, runOperation, @@ -604,7 +604,26 @@ export async function executeWebhookJob( if (!operationStarted) { await releaseExecutionSlot(executionId) } - return result + if (execution.outcome === 'in-progress') { + // Ingress already acknowledged this delivery and nothing reads this job's result, + // so waiting on the live holder would only pin the machine and queue slot. + logger.info(`[${requestId}] Skipping duplicate webhook delivery already in progress`, { + webhookId: authenticatedPayload.webhookId, + workflowId: authenticatedPayload.workflowId, + provider: authenticatedPayload.provider, + executionId, + }) + return { + success: true, + duplicate: true, + workflowId: authenticatedPayload.workflowId, + executionId, + output: {}, + executedAt: new Date().toISOString(), + provider: authenticatedPayload.provider, + } + } + return execution.result }) } catch (error) { await releaseExecutionSlot(executionId) diff --git a/apps/sim/lib/core/idempotency/service.test.ts b/apps/sim/lib/core/idempotency/service.test.ts index 52b678d3218..bc4d5ff18c4 100644 --- a/apps/sim/lib/core/idempotency/service.test.ts +++ b/apps/sim/lib/core/idempotency/service.test.ts @@ -1,10 +1,12 @@ +import { flushMicrotasks } from '@sim/testing/helpers/async' import { dbChainMockFns, resetDbChainMock } from '@sim/testing/mocks/database.mock' import { redisConfigMockFns } from '@sim/testing/mocks/redis-config.mock' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' -const { redisDelMock, redisEvalMock } = vi.hoisted(() => ({ +const { redisDelMock, redisEvalMock, redisGetMock } = vi.hoisted(() => ({ redisDelMock: vi.fn(), redisEvalMock: vi.fn(), + redisGetMock: vi.fn(), })) vi.mock('@/lib/core/storage', () => ({ @@ -22,7 +24,11 @@ import { webhookIdempotency, } from '@/lib/core/idempotency/service' -redisConfigMockFns.mockGetRedisClient.mockReturnValue({ del: redisDelMock, eval: redisEvalMock }) +redisConfigMockFns.mockGetRedisClient.mockReturnValue({ + del: redisDelMock, + eval: redisEvalMock, + get: redisGetMock, +}) const SEVEN_DAYS_SECONDS = 60 * 60 * 24 * 7 @@ -32,6 +38,8 @@ afterEach(() => { beforeEach(() => { resetDbChainMock() + redisEvalMock.mockReset() + redisGetMock.mockReset() }) describe('IdempotencyService.createWebhookIdempotencyKey', () => { @@ -305,6 +313,78 @@ describe('IdempotencyService in-progress deadlines', () => { }) }) +describe('IdempotencyService duplicate of an in-progress operation', () => { + const liveClaim = () => + JSON.stringify({ + success: false, + status: 'in-progress', + startedAt: Date.now(), + inProgressExpiresAt: Date.now() + WEBHOOK_IN_PROGRESS_LEASE_SECONDS * 1000, + claimToken: 'other-holder', + }) + + it('returns in-progress at once instead of polling the live holder when skipping', async () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z')) + const holder = liveClaim() + redisEvalMock.mockResolvedValueOnce([0, holder]) + redisGetMock.mockResolvedValue(holder) + const operation = vi.fn() + const settled = vi.fn() + + void webhookIdempotency + .executeOrSkipInProgress('gmail', 'wh_1:running-delivery', operation) + .then(settled, settled) + await flushMicrotasks(10) + + expect(settled).toHaveBeenCalledExactlyOnceWith({ outcome: 'in-progress' }) + expect(operation).not.toHaveBeenCalled() + expect(redisGetMock).not.toHaveBeenCalled() + }) + + it('still replays a completed result and rethrows a failed one when skipping', async () => { + const service = new IdempotencyService({ forceStorage: 'redis' }) + redisEvalMock + .mockResolvedValueOnce([ + 0, + JSON.stringify({ success: true, status: 'completed', result: 'first-run' }), + ]) + .mockResolvedValueOnce([ + 0, + JSON.stringify({ success: false, status: 'failed', error: 'first run failed' }), + ]) + + await expect(service.executeOrSkipInProgress('provider', 'done', vi.fn())).resolves.toEqual({ + outcome: 'resolved', + result: 'first-run', + }) + await expect(service.executeOrSkipInProgress('provider', 'failed', vi.fn())).rejects.toThrow( + 'first run failed' + ) + }) + + it('keeps waiting for the live holder by default, so a Stripe-style caller never acknowledges early', async () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z')) + const service = new IdempotencyService({ forceStorage: 'redis' }) + const holder = liveClaim() + redisEvalMock.mockResolvedValueOnce([0, holder]) + redisGetMock + .mockResolvedValueOnce(holder) + .mockResolvedValueOnce( + JSON.stringify({ success: true, status: 'completed', result: 'holder-result' }) + ) + const settled = vi.fn() + + void service.executeWithIdempotency('stripe', 'evt_1', vi.fn()).then(settled, settled) + await flushMicrotasks(10) + expect(settled).not.toHaveBeenCalled() + + await vi.advanceTimersByTimeAsync(1_000) + expect(settled).toHaveBeenCalledExactlyOnceWith('holder-result') + }) +}) + describe('IdempotencyService retryable setup failures', () => { it('releases the claim instead of memoizing when the operation throws a RetryableSetupError', async () => { redisEvalMock.mockResolvedValueOnce([1, '']).mockResolvedValueOnce(1) diff --git a/apps/sim/lib/core/idempotency/service.ts b/apps/sim/lib/core/idempotency/service.ts index aa885aa8e55..86f020da6a7 100644 --- a/apps/sim/lib/core/idempotency/service.ts +++ b/apps/sim/lib/core/idempotency/service.ts @@ -73,6 +73,18 @@ export interface IdempotencyExecutionOptions { inProgressExpiresAt?: number } +/** + * Outcome of {@link IdempotencyService.executeOrSkipInProgress}. `resolved` covers both a + * fresh run and a replayed completed result; `in-progress` means another holder owns a live + * claim on the key and this caller did nothing. + */ +export type IdempotentExecution = { outcome: 'resolved'; result: T } | { outcome: 'in-progress' } + +/** An in-progress outcome carries the wait on the live holder, so each caller decides whether to take it. */ +type ClaimedExecution = + | { outcome: 'resolved'; result: T } + | { outcome: 'in-progress'; wait: () => Promise } + export interface AtomicClaimResult { claimed: boolean existingResult?: ProcessingResult @@ -559,13 +571,60 @@ export class IdempotencyService { return deleted.length > 0 } + /** + * Runs `operation` once per key while its claim lease and stored result last. A caller that + * finds another holder's live claim polls until that holder finishes and returns (or rethrows) + * its outcome, so use this when the caller's own response depends on the result (e.g. Stripe + * must not get a 2xx before the first attempt settles). + */ async executeWithIdempotency( provider: string, identifier: string, operation: () => Promise, - additionalContext?: Record, + additionalContext?: Record, options?: IdempotencyExecutionOptions ): Promise { + const execution = await this.execute( + provider, + identifier, + operation, + additionalContext, + options + ) + return execution.outcome === 'resolved' ? execution.result : execution.wait() + } + + /** + * Like {@link executeWithIdempotency}, but returns `{ outcome: 'in-progress' }` immediately + * when another holder owns a live claim instead of polling until it finishes. Use it when + * nobody consumes the duplicate's result, so waiting would only pin the caller's worker and + * leases. Completed and failed keys behave exactly as in `executeWithIdempotency`. + */ + async executeOrSkipInProgress( + provider: string, + identifier: string, + operation: () => Promise, + additionalContext?: Record, + options?: IdempotencyExecutionOptions + ): Promise> { + const execution = await this.execute( + provider, + identifier, + operation, + additionalContext, + options + ) + if (execution.outcome === 'resolved') return execution + return { outcome: 'in-progress' } + } + + private async execute( + provider: string, + identifier: string, + operation: () => Promise, + additionalContext: Record | undefined, + options: IdempotencyExecutionOptions | undefined + ): Promise> { const claimResult = await this.atomicallyClaim(provider, identifier, additionalContext, options) if (!claimResult.claimed) { @@ -576,7 +635,7 @@ export class IdempotencyService { if (existingResult.success === false) { throw new Error(existingResult.error || 'Previous operation failed') } - return existingResult.result as T + return { outcome: 'resolved', result: existingResult.result as T } } if (existingResult?.status === 'failed') { @@ -585,30 +644,28 @@ export class IdempotencyService { observedResult: existingResult, observedValue: claimResult.observedValue, }) - return this.executeWithIdempotency( - provider, - identifier, - operation, - additionalContext, - options - ) + return this.execute(provider, identifier, operation, additionalContext, options) } logger.info(`Previous operation failed for: ${claimResult.normalizedKey}`) throw new Error(existingResult.error || 'Previous operation failed') } if (existingResult?.status === 'in-progress') { - logger.info(`Waiting for in-progress operation: ${claimResult.normalizedKey}`) - return await this.waitForResult( - claimResult.normalizedKey, - claimResult.storageMethod, + const { normalizedKey, storageMethod } = claimResult + const deadline = existingResult.inProgressExpiresAt ?? - (existingResult.startedAt ?? Date.now()) + this.config.inProgressTtlSeconds * 1000 - ) + (existingResult.startedAt ?? Date.now()) + this.config.inProgressTtlSeconds * 1000 + return { + outcome: 'in-progress', + wait: () => { + logger.info(`Waiting for in-progress operation: ${normalizedKey}`) + return this.waitForResult(normalizedKey, storageMethod, deadline) + }, + } } if (existingResult) { - return existingResult.result as T + return { outcome: 'resolved', result: existingResult.result as T } } throw new Error(`Unexpected state: key claimed but no existing result found`) @@ -630,7 +687,7 @@ export class IdempotencyService { ) logger.debug(`Successfully completed operation: ${claimResult.normalizedKey}`) - return result + return { outcome: 'resolved', result } } catch (error) { const errorMessage = getErrorMessage(error, 'Unknown error') diff --git a/scripts/check-explicit-any.baseline.json b/scripts/check-explicit-any.baseline.json index af7872b706f..f8dd08ae862 100644 --- a/scripts/check-explicit-any.baseline.json +++ b/scripts/check-explicit-any.baseline.json @@ -210,7 +210,7 @@ "apps/sim/lib/content/mdx.tsx": 15, "apps/sim/lib/content/registry-factory.ts": 3, "apps/sim/lib/core/config/redis.ts": 1, - "apps/sim/lib/core/idempotency/service.ts": 7, + "apps/sim/lib/core/idempotency/service.ts": 6, "apps/sim/lib/core/security/redaction.ts": 6, "apps/sim/lib/core/utils/display-filters.ts": 11, "apps/sim/lib/core/utils/response-format.ts": 10,