improvement(webhook): added redis for processing duplicate requests

This commit is contained in:
Emir Karabeg
2025-03-11 04:32:57 -07:00
parent 3ac9cf5d8c
commit 843ab03608
4 changed files with 350 additions and 21 deletions
+92 -20
View File
@@ -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<string>()
// 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<string> {
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
}
}
+169
View File
@@ -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<string, { value: string; expiry: number | null }>()
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<boolean> {
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<void> {
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<void> {
if (globalRedisClient) {
try {
await globalRedisClient.quit()
} catch (error) {
console.error('Error closing Redis connection:', error)
} finally {
globalRedisClient = null
}
}
}
+88 -1
View File
@@ -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",
+1
View File
@@ -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",