From 843ab036088a69b6606bbd8d4c360600152c957b Mon Sep 17 00:00:00 2001 From: Emir Karabeg Date: Tue, 11 Mar 2025 04:01:05 -0700 Subject: [PATCH] improvement(webhook): added redis for processing duplicate requests --- app/api/webhooks/trigger/[path]/route.ts | 112 ++++++++++++--- lib/redis.ts | 169 +++++++++++++++++++++++ package-lock.json | 89 +++++++++++- package.json | 1 + 4 files changed, 350 insertions(+), 21 deletions(-) create mode 100644 lib/redis.ts diff --git a/app/api/webhooks/trigger/[path]/route.ts b/app/api/webhooks/trigger/[path]/route.ts index 152aea338b..94f46682a0 100644 --- a/app/api/webhooks/trigger/[path]/route.ts +++ b/app/api/webhooks/trigger/[path]/route.ts @@ -2,6 +2,7 @@ import { NextRequest, NextResponse } from 'next/server' import { and, eq } from 'drizzle-orm' import { v4 as uuidv4 } from 'uuid' import { persistExecutionError, persistExecutionLogs } from '@/lib/logging' +import { closeRedisConnection, hasProcessedMessage, markMessageAsProcessed } from '@/lib/redis' import { decryptSecret } from '@/lib/utils' import { mergeSubblockState, mergeSubblockStateAsync } from '@/stores/workflows/utils' import { db } from '@/db' @@ -9,18 +10,15 @@ import { environment, webhook, workflow } from '@/db/schema' import { Executor } from '@/executor' import { Serializer } from '@/serializer' +// Force dynamic rendering for webhook endpoints export const dynamic = 'force-dynamic' - -// Store for tracking processed message IDs in memory -// This is a simple in-memory solution that works for a single instance -// For multi-instance deployments, consider using Redis or another distributed cache -const processedMessageIds = new Set() +// Increase the response size limit for webhook payloads +export const maxDuration = 300 // 5 minutes max execution time for long-running webhooks /** * Consolidated webhook trigger endpoint for all providers * Handles both WhatsApp verification and other webhook providers */ - export async function GET(request: NextRequest, { params }: { params: Promise<{ path: string }> }) { try { const path = (await params).path @@ -95,6 +93,9 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{ } catch (error: any) { console.error('Error processing webhook verification:', error) return new NextResponse(`Internal Server Error: ${error.message}`, { status: 500 }) + } finally { + // Ensure Redis connection is properly closed in serverless environment + await closeRedisConnection() } } @@ -112,6 +113,18 @@ export async function POST( const body = await request.json().catch(() => ({})) console.log(`Webhook POST request received for path: ${path}`) + // Generate a unique request ID based on the request content + const requestHash = await generateRequestHash(path, body) + + // Check if this exact request has been processed before + if (await hasProcessedMessage(requestHash)) { + console.log( + `Duplicate webhook request detected with hash: ${requestHash}. Skipping processing.` + ) + // Return early for duplicate requests to prevent workflow execution + return new NextResponse('Duplicate request', { status: 200 }) + } + // Find the webhook in the database const webhooks = await db .select({ @@ -130,7 +143,7 @@ export async function POST( const { webhook: foundWebhook, workflow: workflowData } = webhooks[0] foundWorkflow = workflowData - // For WhatsApp, check for duplicate messages before processing + // For WhatsApp, also check for duplicate messages using their message ID if (foundWebhook.provider === 'whatsapp') { const data = body?.entry?.[0]?.changes?.[0]?.value const messages = data?.messages || [] @@ -139,27 +152,24 @@ export async function POST( const message = messages[0] const messageId = message.id - // Check if we've already processed this message - if (messageId && processedMessageIds.has(messageId)) { + // Check if we've already processed this message using Redis + if (messageId && (await hasProcessedMessage(messageId))) { console.log( `Duplicate WhatsApp message detected with ID: ${messageId}. Skipping processing.` ) - return new NextResponse('OK - Duplicate message', { status: 200 }) + // Return early for duplicate messages to prevent workflow execution + return new NextResponse('Duplicate message', { status: 200 }) } - // Store the message ID to prevent duplicate processing in future requests + // Store the message ID in Redis to prevent duplicate processing in future requests if (messageId) { - processedMessageIds.add(messageId) - - // Clean up old message IDs periodically (keep last ~1000 messages) - if (processedMessageIds.size > 1000) { - const idsArray = Array.from(processedMessageIds) - for (let i = 0; i < idsArray.length - 1000; i++) { - processedMessageIds.delete(idsArray[i]) - } - } + await markMessageAsProcessed(messageId) } + // Mark this request as processed to prevent duplicates + // Use a shorter TTL for request hashes (24 hours) to save Redis memory + await markMessageAsProcessed(requestHash, 60 * 60 * 24) + // Process the webhook synchronously - complete the workflow before returning const result = await processWebhook(foundWebhook, foundWorkflow, body, request, executionId) @@ -173,6 +183,9 @@ export async function POST( } } + // Mark this request as processed to prevent duplicates + await markMessageAsProcessed(requestHash, 60 * 60 * 24) + // For other providers, continue with synchronous processing return await processWebhook(foundWebhook, foundWorkflow, body, request, executionId) } catch (error: any) { @@ -184,6 +197,65 @@ export async function POST( } return new NextResponse(`Internal Server Error: ${error.message}`, { status: 500 }) + } finally { + // Ensure Redis connection is properly closed in serverless environment + await closeRedisConnection() + } +} + +/** + * Generate a unique hash for a webhook request based on its path and body + * This is used to deduplicate webhook requests + */ +async function generateRequestHash(path: string, body: any): Promise { + try { + // Create a string representation of the request + // Remove any timestamp or random fields that would make identical requests look different + const normalizedBody = normalizeBody(body) + const requestString = `${path}:${JSON.stringify(normalizedBody)}` + + // Use a simple hash function for the request + let hash = 0 + for (let i = 0; i < requestString.length; i++) { + const char = requestString.charCodeAt(i) + hash = (hash << 5) - hash + char + hash = hash & hash // Convert to 32bit integer + } + + return `request:${path}:${hash}` + } catch (error) { + // If hashing fails, use a UUID as fallback + console.error('Error generating request hash:', error) + return `request:${path}:${uuidv4()}` + } +} + +/** + * Normalize webhook body by removing fields that might change between identical requests + * This helps with more accurate deduplication + */ +function normalizeBody(body: any): any { + if (!body || typeof body !== 'object') return body + + // Create a copy to avoid modifying the original + const result = Array.isArray(body) ? [...body] : { ...body } + + // Fields to remove (common timestamp/random fields) + const fieldsToRemove = ['timestamp', 'random', 'nonce', 'requestId'] + + if (Array.isArray(result)) { + // Handle arrays + return result.map((item) => normalizeBody(item)) + } else { + // Handle objects + for (const key in result) { + if (fieldsToRemove.includes(key.toLowerCase())) { + delete result[key] + } else if (typeof result[key] === 'object' && result[key] !== null) { + result[key] = normalizeBody(result[key]) + } + } + return result } } diff --git a/lib/redis.ts b/lib/redis.ts new file mode 100644 index 0000000000..6da2f496e4 --- /dev/null +++ b/lib/redis.ts @@ -0,0 +1,169 @@ +import Redis from 'ioredis' + +// Default to localhost if REDIS_URL is not provided +const redisUrl = process.env.REDIS_URL || 'redis://localhost:6379' + +// Global Redis client for connection pooling +// This is important for serverless environments like Vercel +let globalRedisClient: Redis | null = null + +// Fallback in-memory cache for when Redis is not available +const inMemoryCache = new Map() +const MAX_CACHE_SIZE = 1000 + +/** + * Get a Redis client instance + * Uses connection pooling to avoid creating a new connection for each request + * This is critical for performance in serverless environments like Vercel + */ +export function getRedisClient(): Redis | null { + // For server-side only + if (typeof window !== 'undefined') return null + + if (globalRedisClient) return globalRedisClient + + try { + // Create a new Redis client with optimized settings for serverless + globalRedisClient = new Redis(redisUrl, { + // Keep alive is critical for serverless to reuse connections + keepAlive: 1000, + // Faster connection timeout for serverless + connectTimeout: 5000, + // Disable reconnection attempts in serverless + maxRetriesPerRequest: 3, + // Retry strategy with exponential backoff + retryStrategy: (times) => { + if (times > 5) { + console.warn('Redis connection failed after 5 attempts, using fallback') + return null // Stop retrying + } + return Math.min(times * 200, 2000) // Exponential backoff + }, + }) + + // Handle connection events + globalRedisClient.on('error', (err: any) => { + console.error('Redis connection error:', err) + if (err.code === 'ECONNREFUSED' || err.code === 'ETIMEDOUT') { + globalRedisClient = null + } + }) + + globalRedisClient.on('connect', () => { + console.log('Connected to Redis') + }) + + return globalRedisClient + } catch (error) { + console.error('Failed to initialize Redis client:', error) + return null + } +} + +// Message ID cache functions +const MESSAGE_ID_PREFIX = 'whatsapp:message:' +const MESSAGE_ID_EXPIRY = 60 * 60 * 24 * 7 // 7 days in seconds + +/** + * Check if a message ID has been processed before + * @param messageId The message ID to check + * @returns True if the message has been processed before, false otherwise + */ +export async function hasProcessedMessage(messageId: string): Promise { + try { + const redis = getRedisClient() + + if (redis) { + // Use Redis if available + const key = `${MESSAGE_ID_PREFIX}${messageId}` + const result = await redis.exists(key) + return result === 1 + } else { + // Fallback to in-memory cache + const cacheEntry = inMemoryCache.get(messageId) + if (!cacheEntry) return false + + // Check if the entry has expired + if (cacheEntry.expiry && cacheEntry.expiry < Date.now()) { + inMemoryCache.delete(messageId) + return false + } + + return true + } + } catch (error) { + console.error('Error checking message ID:', error) + // Fallback to in-memory cache on error + const cacheEntry = inMemoryCache.get(messageId) + return !!cacheEntry && (!cacheEntry.expiry || cacheEntry.expiry > Date.now()) + } +} + +/** + * Mark a message ID as processed + * @param messageId The message ID to mark as processed + * @param expirySeconds Optional expiry time in seconds (defaults to 7 days) + */ +export async function markMessageAsProcessed( + messageId: string, + expirySeconds: number = MESSAGE_ID_EXPIRY +): Promise { + try { + const redis = getRedisClient() + + if (redis) { + // Use Redis if available - use pipelining for efficiency + const key = `${MESSAGE_ID_PREFIX}${messageId}` + await redis.set(key, '1', 'EX', expirySeconds) + } else { + // Fallback to in-memory cache + const expiry = expirySeconds ? Date.now() + expirySeconds * 1000 : null + inMemoryCache.set(messageId, { value: '1', expiry }) + + // Clean up old message IDs if cache gets too large + if (inMemoryCache.size > MAX_CACHE_SIZE) { + const now = Date.now() + + // First try to remove expired entries + for (const [key, entry] of inMemoryCache.entries()) { + if (entry.expiry && entry.expiry < now) { + inMemoryCache.delete(key) + } + } + + // If still too large, remove oldest entries + if (inMemoryCache.size > MAX_CACHE_SIZE) { + const keysToDelete = Array.from(inMemoryCache.keys()).slice( + 0, + inMemoryCache.size - MAX_CACHE_SIZE + ) + + for (const key of keysToDelete) { + inMemoryCache.delete(key) + } + } + } + } + } catch (error) { + console.error('Error marking message as processed:', error) + // Fallback to in-memory cache on error + const expiry = expirySeconds ? Date.now() + expirySeconds * 1000 : null + inMemoryCache.set(messageId, { value: '1', expiry }) + } +} + +/** + * Close the Redis connection + * Important for proper cleanup in serverless environments + */ +export async function closeRedisConnection(): Promise { + if (globalRedisClient) { + try { + await globalRedisClient.quit() + } catch (error) { + console.error('Error closing Redis connection:', error) + } finally { + globalRedisClient = null + } + } +} diff --git a/package-lock.json b/package-lock.json index 368d8eada8..8a998ca172 100644 --- a/package-lock.json +++ b/package-lock.json @@ -43,6 +43,7 @@ "drizzle-orm": "^0.39.3", "groq-sdk": "^0.15.0", "input-otp": "^1.4.2", + "ioredis": "^5.6.0", "jwt-decode": "^4.0.0", "lodash.debounce": "^4.0.8", "lucide-react": "^0.469.0", @@ -820,6 +821,12 @@ "url": "https://opencollective.com/libvips" } }, + "node_modules/@ioredis/commands": { + "version": "1.2.0", + "resolved": "https://registry.npmjs.org/@ioredis/commands/-/commands-1.2.0.tgz", + "integrity": "sha512-Sx1pU8EM64o2BrqNpEO1CNLtKQwyhuXuqyfH7oGKCk+1a33d2r5saW8zNwm3j6BTExtjrv2BxTgzzkMwts6vGg==", + "license": "MIT" + }, "node_modules/@isaacs/cliui": { "version": "8.0.2", "license": "ISC", @@ -4471,6 +4478,15 @@ "node": ">=6" } }, + "node_modules/cluster-key-slot": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/cluster-key-slot/-/cluster-key-slot-1.1.2.tgz", + "integrity": "sha512-RMr0FhtfXemyinomL4hrWcYJxmX6deFdCxpJzhDttxgO1+bcCnkk+9drydLVDmAMG7NE6aN/fl4F7ucU/90gAA==", + "license": "Apache-2.0", + "engines": { + "node": ">=0.10.0" + } + }, "node_modules/cmdk": { "version": "1.0.0", "license": "MIT", @@ -5337,7 +5353,6 @@ }, "node_modules/debug": { "version": "4.4.0", - "dev": true, "license": "MIT", "dependencies": { "ms": "^2.1.3" @@ -5452,6 +5467,15 @@ "node": ">=0.4.0" } }, + "node_modules/denque": { + "version": "2.1.0", + "resolved": "https://registry.npmjs.org/denque/-/denque-2.1.0.tgz", + "integrity": "sha512-HVQE3AAb/pxF8fQAoiqpvg9i3evqug3hoiwakOyZAwJm+6vZehbkYXZ0l4JxS+I3QxM97v5aaRNhj8v5oBhekw==", + "license": "Apache-2.0", + "engines": { + "node": ">=0.10" + } + }, "node_modules/dequal": { "version": "2.0.3", "dev": true, @@ -6834,6 +6858,30 @@ "node": ">=8" } }, + "node_modules/ioredis": { + "version": "5.6.0", + "resolved": "https://registry.npmjs.org/ioredis/-/ioredis-5.6.0.tgz", + "integrity": "sha512-tBZlIIWbndeWBWCXWZiqtOF/yxf6yZX3tAlTJ7nfo5jhd6dctNxF7QnYlZLZ1a0o0pDoen7CgZqO+zjNaFbJAg==", + "license": "MIT", + "dependencies": { + "@ioredis/commands": "^1.1.1", + "cluster-key-slot": "^1.1.0", + "debug": "^4.3.4", + "denque": "^2.1.0", + "lodash.defaults": "^4.2.0", + "lodash.isarguments": "^3.1.0", + "redis-errors": "^1.2.0", + "redis-parser": "^3.0.0", + "standard-as-callback": "^2.1.0" + }, + "engines": { + "node": ">=12.22.0" + }, + "funding": { + "type": "opencollective", + "url": "https://opencollective.com/ioredis" + } + }, "node_modules/is-arrayish": { "version": "0.3.2", "license": "MIT", @@ -8258,6 +8306,18 @@ "version": "4.0.8", "license": "MIT" }, + "node_modules/lodash.defaults": { + "version": "4.2.0", + "resolved": "https://registry.npmjs.org/lodash.defaults/-/lodash.defaults-4.2.0.tgz", + "integrity": "sha512-qjxPLHd3r5DnsdGacqOMU6pb/avJzdh9tFX2ymgoZE27BmjXrNy/y4LoaiTeAb+O3gL8AfpJGtqfX/ae2leYYQ==", + "license": "MIT" + }, + "node_modules/lodash.isarguments": { + "version": "3.1.0", + "resolved": "https://registry.npmjs.org/lodash.isarguments/-/lodash.isarguments-3.1.0.tgz", + "integrity": "sha512-chi4NHZlZqZD18a0imDHnZPrDeBbTtVN7GXMwuGdRH9qotxAjYs3aVLKc7zNOG9eddR5Ksd8rvFEBc9SsggPpg==", + "license": "MIT" + }, "node_modules/lodash.memoize": { "version": "4.1.2", "dev": true, @@ -9825,6 +9885,27 @@ "node": ">=8" } }, + "node_modules/redis-errors": { + "version": "1.2.0", + "resolved": "https://registry.npmjs.org/redis-errors/-/redis-errors-1.2.0.tgz", + "integrity": "sha512-1qny3OExCf0UvUV/5wpYKf2YwPcOqXzkwKKSmKHiE6ZMQs5heeE/c8eXK+PNllPvmjgAbfnsbpkGZWy8cBpn9w==", + "license": "MIT", + "engines": { + "node": ">=4" + } + }, + "node_modules/redis-parser": { + "version": "3.0.0", + "resolved": "https://registry.npmjs.org/redis-parser/-/redis-parser-3.0.0.tgz", + "integrity": "sha512-DJnGAeenTdpMEH6uAJRK/uiyEIH9WVsUmoLwzudwGJUwZPp80PDBWPHXSAGNPwNvIXAbe7MSUB1zQFugFml66A==", + "license": "MIT", + "dependencies": { + "redis-errors": "^1.0.0" + }, + "engines": { + "node": ">=4" + } + }, "node_modules/regenerator-runtime": { "version": "0.14.1", "license": "MIT" @@ -10323,6 +10404,12 @@ "node": ">=10" } }, + "node_modules/standard-as-callback": { + "version": "2.1.0", + "resolved": "https://registry.npmjs.org/standard-as-callback/-/standard-as-callback-2.1.0.tgz", + "integrity": "sha512-qoRRSyROncaz1z0mvYqIE4lCd9p2R90i6GxW3uZv5ucSu8tU7B5HXUP1gG8pVZsYNVaXjk8ClXHPttLyxAL48A==", + "license": "MIT" + }, "node_modules/stdin-discarder": { "version": "0.2.2", "license": "MIT", diff --git a/package.json b/package.json index 919adb5262..8016b163eb 100644 --- a/package.json +++ b/package.json @@ -55,6 +55,7 @@ "drizzle-orm": "^0.39.3", "groq-sdk": "^0.15.0", "input-otp": "^1.4.2", + "ioredis": "^5.6.0", "jwt-decode": "^4.0.0", "lodash.debounce": "^4.0.8", "lucide-react": "^0.469.0",