mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
fix(triggers): dedup + not surfacing deployment status log
This commit is contained in:
@@ -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 })
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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, string>): string {
|
||||
static createWebhookIdempotencyKey(
|
||||
webhookId: string,
|
||||
headers?: Record<string, string>,
|
||||
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}`
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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, (body: any) => 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
|
||||
}
|
||||
Reference in New Issue
Block a user