From 82f82faba4169126716c5d7801561fbe8728a00c Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Mon, 17 Nov 2025 15:34:51 -0800 Subject: [PATCH] fix(triggers): dedup + not surfacing deployment status log --- .../app/api/webhooks/trigger/[path]/route.ts | 30 +------ apps/sim/background/webhook-execution.ts | 4 +- apps/sim/lib/idempotency/service.ts | 32 ++++++- apps/sim/lib/webhooks/provider-utils.ts | 85 +++++++++++++++++++ 4 files changed, 120 insertions(+), 31 deletions(-) create mode 100644 apps/sim/lib/webhooks/provider-utils.ts diff --git a/apps/sim/app/api/webhooks/trigger/[path]/route.ts b/apps/sim/app/api/webhooks/trigger/[path]/route.ts index 4f8704489d..0de3c1a632 100644 --- a/apps/sim/app/api/webhooks/trigger/[path]/route.ts +++ b/apps/sim/app/api/webhooks/trigger/[path]/route.ts @@ -1,7 +1,5 @@ import { type NextRequest, NextResponse } from 'next/server' -import { v4 as uuidv4 } from 'uuid' import { createLogger } from '@/lib/logs/console/logger' -import { LoggingSession } from '@/lib/logs/execution/logging-session' import { generateRequestId } from '@/lib/utils' import { checkRateLimits, @@ -139,34 +137,10 @@ export async function POST( if (foundWebhook.blockId) { const blockExists = await blockExistsInDeployment(foundWorkflow.id, foundWebhook.blockId) if (!blockExists) { - logger.warn( + logger.info( `[${requestId}] Trigger block ${foundWebhook.blockId} not found in deployment for workflow ${foundWorkflow.id}` ) - - const executionId = uuidv4() - const loggingSession = new LoggingSession(foundWorkflow.id, executionId, 'webhook', requestId) - - const actorUserId = foundWorkflow.workspaceId - ? (await import('@/lib/workspaces/utils')).getWorkspaceBilledAccountUserId( - foundWorkflow.workspaceId - ) || foundWorkflow.userId - : foundWorkflow.userId - - await loggingSession.safeStart({ - userId: actorUserId, - workspaceId: foundWorkflow.workspaceId || '', - variables: {}, - }) - - await loggingSession.safeCompleteWithError({ - error: { - message: `Trigger block not deployed. The webhook trigger (block ${foundWebhook.blockId}) is not present in the deployed workflow. Please redeploy the workflow.`, - stackTrace: undefined, - }, - traceSpans: [], - }) - - return new NextResponse('Trigger block not deployed', { status: 404 }) + return new NextResponse('Trigger block not found in deployment', { status: 404 }) } } diff --git a/apps/sim/background/webhook-execution.ts b/apps/sim/background/webhook-execution.ts index 14e8ca40e2..7698616b59 100644 --- a/apps/sim/background/webhook-execution.ts +++ b/apps/sim/background/webhook-execution.ts @@ -112,7 +112,9 @@ export async function executeWebhookJob(payload: WebhookExecutionPayload) { const idempotencyKey = IdempotencyService.createWebhookIdempotencyKey( payload.webhookId, - payload.headers + payload.headers, + payload.body, + payload.provider ) const runOperation = async () => { diff --git a/apps/sim/lib/idempotency/service.ts b/apps/sim/lib/idempotency/service.ts index e20dc6dbfe..31495c2a99 100644 --- a/apps/sim/lib/idempotency/service.ts +++ b/apps/sim/lib/idempotency/service.ts @@ -4,6 +4,7 @@ import { idempotencyKey } from '@sim/db/schema' import { and, eq } from 'drizzle-orm' import { createLogger } from '@/lib/logs/console/logger' import { getRedisClient } from '@/lib/redis' +import { extractProviderIdentifierFromBody } from '@/lib/webhooks/provider-utils' const logger = createLogger('IdempotencyService') @@ -451,13 +452,25 @@ export class IdempotencyService { /** * Create an idempotency key from a webhook payload following RFC best practices - * Standard webhook headers (webhook-id, x-webhook-id, etc.) + * Checks both headers and body for unique identifiers to prevent duplicate executions + * + * @param webhookId - The webhook database ID + * @param headers - HTTP headers from the webhook request + * @param body - Parsed webhook body (optional, used for provider-specific identifiers) + * @param provider - Provider name for body extraction (optional) + * @returns A unique idempotency key for this webhook event */ - static createWebhookIdempotencyKey(webhookId: string, headers?: Record): string { + static createWebhookIdempotencyKey( + webhookId: string, + headers?: Record, + body?: any, + provider?: string + ): string { const normalizedHeaders = headers ? Object.fromEntries(Object.entries(headers).map(([k, v]) => [k.toLowerCase(), v])) : undefined + // Check standard webhook headers first const webhookIdHeader = normalizedHeaders?.['webhook-id'] || normalizedHeaders?.['x-webhook-id'] || @@ -470,7 +483,22 @@ export class IdempotencyService { return `${webhookId}:${webhookIdHeader}` } + // Check body for provider-specific unique identifiers + if (body && provider) { + const bodyIdentifier = extractProviderIdentifierFromBody(provider, body) + + if (bodyIdentifier) { + return `${webhookId}:${bodyIdentifier}` + } + } + + // No unique identifier found - generate random UUID + // This means duplicate detection will not work for this webhook const uniqueId = randomUUID() + logger.warn('No unique identifier found, duplicate executions may occur', { + webhookId, + provider, + }) return `${webhookId}:${uniqueId}` } } diff --git a/apps/sim/lib/webhooks/provider-utils.ts b/apps/sim/lib/webhooks/provider-utils.ts new file mode 100644 index 0000000000..4c94335cb9 --- /dev/null +++ b/apps/sim/lib/webhooks/provider-utils.ts @@ -0,0 +1,85 @@ +/** + * Provider-specific unique identifier extractors for webhook idempotency + */ + +function extractSlackIdentifier(body: any): string | null { + if (body.event_id) { + return body.event_id + } + + if (body.event?.ts && body.team_id) { + return `${body.team_id}:${body.event.ts}` + } + + return null +} + +function extractTwilioIdentifier(body: any): string | null { + return body.MessageSid || body.CallSid || null +} + +function extractStripeIdentifier(body: any): string | null { + if (body.id && body.object === 'event') { + return body.id + } + return null +} + +function extractHubSpotIdentifier(body: any): string | null { + if (Array.isArray(body) && body.length > 0 && body[0]?.eventId) { + return String(body[0].eventId) + } + return null +} + +function extractLinearIdentifier(body: any): string | null { + if (body.action && body.data?.id) { + return `${body.action}:${body.data.id}` + } + return null +} + +function extractJiraIdentifier(body: any): string | null { + if (body.webhookEvent && (body.issue?.id || body.project?.id)) { + return `${body.webhookEvent}:${body.issue?.id || body.project?.id}` + } + return null +} + +function extractMicrosoftTeamsIdentifier(body: any): string | null { + if (body.value && Array.isArray(body.value) && body.value.length > 0) { + const notification = body.value[0] + if (notification.subscriptionId && notification.resourceData?.id) { + return `${notification.subscriptionId}:${notification.resourceData.id}` + } + } + return null +} + +function extractAirtableIdentifier(body: any): string | null { + if (body.cursor && typeof body.cursor === 'string') { + return body.cursor + } + return null +} + +const PROVIDER_EXTRACTORS: Record string | null> = { + slack: extractSlackIdentifier, + twilio: extractTwilioIdentifier, + twilio_voice: extractTwilioIdentifier, + stripe: extractStripeIdentifier, + hubspot: extractHubSpotIdentifier, + linear: extractLinearIdentifier, + jira: extractJiraIdentifier, + microsoftteams: extractMicrosoftTeamsIdentifier, + airtable: extractAirtableIdentifier, +} + +export function extractProviderIdentifierFromBody(provider: string, body: any): string | null { + if (!body || typeof body !== 'object') { + return null + } + + const extractor = PROVIDER_EXTRACTORS[provider] + return extractor ? extractor(body) : null +}