From 626dafae4925a976eb1dc51f25195eb10470a358 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 12 May 2025 13:13:19 -0700 Subject: [PATCH] fix(db): fix slow db connection, add redis cache for gmail polling --- apps/sim/app/api/webhooks/poll/gmail/route.ts | 79 ++---- apps/sim/drizzle.config.ts | 3 + .../sim/lib/webhooks/gmail-polling-service.ts | 253 +++++++++++------- apps/sim/lib/webhooks/utils.ts | 5 + 4 files changed, 185 insertions(+), 155 deletions(-) diff --git a/apps/sim/app/api/webhooks/poll/gmail/route.ts b/apps/sim/app/api/webhooks/poll/gmail/route.ts index 93ca0f8e02..53d06c8560 100644 --- a/apps/sim/app/api/webhooks/poll/gmail/route.ts +++ b/apps/sim/app/api/webhooks/poll/gmail/route.ts @@ -1,43 +1,23 @@ import { NextRequest, NextResponse } from 'next/server' import { nanoid } from 'nanoid' import { Logger } from '@/lib/logs/console-logger' +import { acquireLock, releaseLock } from '@/lib/redis' import { pollGmailWebhooks } from '@/lib/webhooks/gmail-polling-service' const logger = new Logger('GmailPollingAPI') export const dynamic = 'force-dynamic' -export const maxDuration = 300 // Allow up to 5 minutes for polling to complete +export const maxDuration = 180 // Allow up to 3 minutes for polling to complete -interface PollingTask { - promise: Promise - startedAt: number -} - -const activePollingTasks = new Map() -const STALE_TASK_THRESHOLD_MS = 10 * 60 * 1000 // 10 minutes - -function cleanupStaleTasks() { - const now = Date.now() - let removedCount = 0 - - for (const [requestId, task] of activePollingTasks.entries()) { - if (now - task.startedAt > STALE_TASK_THRESHOLD_MS) { - activePollingTasks.delete(requestId) - removedCount++ - } - } - - if (removedCount > 0) { - logger.info(`Cleaned up ${removedCount} stale polling tasks`) - } - - return removedCount -} +const LOCK_KEY = 'gmail-polling-lock' +const LOCK_TTL_SECONDS = 180 // Same as maxDuration (3 min) export async function GET(request: NextRequest) { const requestId = nanoid() logger.info(`Gmail webhook polling triggered (${requestId})`) + let lockValue: string | undefined + try { const authHeader = request.headers.get('authorization') const webhookSecret = process.env.CRON_SECRET || process.env.WEBHOOK_POLLING_SECRET @@ -51,49 +31,42 @@ export async function GET(request: NextRequest) { return new NextResponse('Unauthorized', { status: 401 }) } - cleanupStaleTasks() + lockValue = requestId // unique value to identify the holder + const locked = await acquireLock(LOCK_KEY, lockValue, LOCK_TTL_SECONDS) - const pollingTask: PollingTask = { - promise: null as any, - startedAt: Date.now(), + if (!locked) { + return NextResponse.json( + { + success: true, + message: 'Polling already in progress – skipped', + requestId, + status: 'skip', + }, + { status: 202 } + ) } - pollingTask.promise = pollGmailWebhooks() - .then((results) => { - logger.info(`Gmail polling completed successfully (${requestId})`, { - userCount: results?.total || 0, - successful: results?.successful || 0, - failed: results?.failed || 0, - }) - activePollingTasks.delete(requestId) - return results - }) - .catch((error) => { - logger.error(`Error in background Gmail polling task (${requestId}):`, error) - activePollingTasks.delete(requestId) - throw error - }) - - activePollingTasks.set(requestId, pollingTask) + const results = await pollGmailWebhooks() return NextResponse.json({ success: true, - message: 'Gmail webhook polling started successfully', + message: 'Gmail polling completed', requestId, - status: 'polling_started', - activeTasksCount: activePollingTasks.size, + status: 'completed', + ...results, }) } catch (error) { - logger.error(`Error initiating Gmail webhook polling (${requestId}):`, error) - + logger.error(`Error during Gmail polling (${requestId}):`, error) return NextResponse.json( { success: false, - message: 'Failed to start Gmail webhook polling', + message: 'Gmail polling failed', error: error instanceof Error ? error.message : 'Unknown error', requestId, }, { status: 500 } ) + } finally { + await releaseLock(LOCK_KEY).catch(() => {}) } } diff --git a/apps/sim/drizzle.config.ts b/apps/sim/drizzle.config.ts index c7ddc326a7..8ca4d5358e 100644 --- a/apps/sim/drizzle.config.ts +++ b/apps/sim/drizzle.config.ts @@ -1,5 +1,8 @@ +import * as dotenv from 'dotenv' import type { Config } from 'drizzle-kit' +dotenv.config({ path: '../../.env' }) + export default { schema: './db/schema.ts', out: './db/migrations', diff --git a/apps/sim/lib/webhooks/gmail-polling-service.ts b/apps/sim/lib/webhooks/gmail-polling-service.ts index b53fce000f..77a88211e1 100644 --- a/apps/sim/lib/webhooks/gmail-polling-service.ts +++ b/apps/sim/lib/webhooks/gmail-polling-service.ts @@ -1,6 +1,7 @@ import { and, eq } from 'drizzle-orm' import { nanoid } from 'nanoid' import { Logger } from '@/lib/logs/console-logger' +import { hasProcessedMessage, markMessageAsProcessed } from '@/lib/redis' import { getBaseUrl } from '@/lib/urls/utils' import { getOAuthToken } from '@/app/api/auth/oauth/utils' import { db } from '@/db' @@ -29,6 +30,21 @@ interface GmailEmail { internalDate?: string } +export interface SimplifiedEmail { + id: string + threadId: string + subject: string + from: string + to: string + cc: string + date: string | null + bodyText: string + bodyHtml: string + labels: string[] + hasAttachments: boolean + attachments: Array<{ filename: string; mimeType: string; size: number }> +} + export async function pollGmailWebhooks() { logger.info('Starting Gmail webhook polling') @@ -46,116 +62,136 @@ export async function pollGmailWebhooks() { logger.info(`Found ${activeWebhooks.length} active Gmail webhooks`) - const results = await Promise.allSettled( - activeWebhooks.map(async (webhookData) => { - const webhookId = webhookData.id - const requestId = nanoid() + // Limit the number of webhooks processed in parallel to avoid + // exhausting Postgres or Gmail API connections when many users exist. + const CONCURRENCY = 10 - try { - // Extract user ID from webhook metadata if available - const metadata = webhookData.providerConfig as any - const userId = metadata?.userId + const running: Promise[] = [] + const settledResults: PromiseSettledResult[] = [] - if (!userId) { - logger.error(`[${requestId}] No user ID found for webhook ${webhookId}`) - return { success: false, webhookId, error: 'No user ID' } - } + const enqueue = async (webhookData: (typeof activeWebhooks)[number]) => { + const webhookId = webhookData.id + const requestId = nanoid() - // Get OAuth token for Gmail API - const accessToken = await getOAuthToken(userId, 'google-email') + try { + // Extract user ID from webhook metadata if available + const metadata = webhookData.providerConfig as any + const userId = metadata?.userId - if (!accessToken) { - logger.error(`[${requestId}] Failed to get Gmail access token for webhook ${webhookId}`) - return { success: false, webhookId, error: 'No access token' } - } + if (!userId) { + logger.error(`[${requestId}] No user ID found for webhook ${webhookId}`) + return { success: false, webhookId, error: 'No user ID' } + } - // Get webhook configuration - const config = webhookData.providerConfig as unknown as GmailWebhookConfig + // Get OAuth token for Gmail API + const accessToken = await getOAuthToken(userId, 'google-email') - const now = new Date() + if (!accessToken) { + logger.error(`[${requestId}] Failed to get Gmail access token for webhook ${webhookId}`) + return { success: false, webhookId, error: 'No access token' } + } - // Fetch new emails - const fetchResult = await fetchNewEmails(accessToken, config, requestId) + // Get webhook configuration + const config = webhookData.providerConfig as unknown as GmailWebhookConfig - const { emails, latestHistoryId } = fetchResult + const now = new Date() - if (!emails || !emails.length) { - // Update last checked timestamp - await updateWebhookLastChecked( - webhookId, - now.toISOString(), - latestHistoryId || config.historyId - ) - logger.info(`[${requestId}] No new emails found for webhook ${webhookId}`) - return { success: true, webhookId, status: 'no_emails' } - } + // Fetch new emails + const fetchResult = await fetchNewEmails(accessToken, config, requestId) - logger.info(`[${requestId}] Found ${emails.length} new emails for webhook ${webhookId}`) + const { emails, latestHistoryId } = fetchResult - // Get processed email IDs (to avoid duplicates) - const processedEmailIds = config.processedEmailIds || [] - - // Filter out emails that have already been processed - const newEmails = emails.filter((email) => !processedEmailIds.includes(email.id)) - - if (newEmails.length === 0) { - logger.info( - `[${requestId}] All emails have already been processed for webhook ${webhookId}` - ) - await updateWebhookLastChecked( - webhookId, - now.toISOString(), - latestHistoryId || config.historyId - ) - return { success: true, webhookId, status: 'already_processed' } - } - - logger.info( - `[${requestId}] Processing ${newEmails.length} new emails for webhook ${webhookId}` - ) - - // Process all emails (process each email as a separate workflow trigger) - const emailsToProcess = newEmails - - // Process emails - const processed = await processEmails( - emailsToProcess, - webhookData, - config, - accessToken, - requestId - ) - - // Record which email IDs have been processed - const newProcessedIds = [ - ...processedEmailIds, - ...emailsToProcess.map((email) => email.id), - ] - // Keep only the most recent 100 IDs to prevent the list from growing too large - const trimmedProcessedIds = newProcessedIds.slice(-100) - - // Update webhook with latest history ID, timestamp, and processed email IDs - await updateWebhookData( + if (!emails || !emails.length) { + // Update last checked timestamp + await updateWebhookLastChecked( webhookId, now.toISOString(), - latestHistoryId || config.historyId, - trimmedProcessedIds + latestHistoryId || config.historyId ) - - return { - success: true, - webhookId, - emailsFound: emails.length, - newEmails: newEmails.length, - emailsProcessed: processed, - } - } catch (error) { - const errorMessage = error instanceof Error ? error.message : 'Unknown error' - logger.error(`[${requestId}] Error processing Gmail webhook ${webhookId}:`, error) - return { success: false, webhookId, error: errorMessage } + logger.info(`[${requestId}] No new emails found for webhook ${webhookId}`) + return { success: true, webhookId, status: 'no_emails' } } - }) - ) + + logger.info(`[${requestId}] Found ${emails.length} new emails for webhook ${webhookId}`) + + // Get processed email IDs (to avoid duplicates) + const processedEmailIds = config.processedEmailIds || [] + + // Filter out emails that have already been processed + const newEmails = emails.filter((email) => !processedEmailIds.includes(email.id)) + + if (newEmails.length === 0) { + logger.info( + `[${requestId}] All emails have already been processed for webhook ${webhookId}` + ) + await updateWebhookLastChecked( + webhookId, + now.toISOString(), + latestHistoryId || config.historyId + ) + return { success: true, webhookId, status: 'already_processed' } + } + + logger.info( + `[${requestId}] Processing ${newEmails.length} new emails for webhook ${webhookId}` + ) + + // Process all emails (process each email as a separate workflow trigger) + const emailsToProcess = newEmails + + // Process emails + const processed = await processEmails( + emailsToProcess, + webhookData, + config, + accessToken, + requestId + ) + + // Record which email IDs have been processed + const newProcessedIds = [...processedEmailIds, ...emailsToProcess.map((email) => email.id)] + // Keep only the most recent 100 IDs to prevent the list from growing too large + const trimmedProcessedIds = newProcessedIds.slice(-100) + + // Update webhook with latest history ID, timestamp, and processed email IDs + await updateWebhookData( + webhookId, + now.toISOString(), + latestHistoryId || config.historyId, + trimmedProcessedIds + ) + + return { + success: true, + webhookId, + emailsFound: emails.length, + newEmails: newEmails.length, + emailsProcessed: processed, + } + } catch (error) { + const errorMessage = error instanceof Error ? error.message : 'Unknown error' + logger.error(`[${requestId}] Error processing Gmail webhook ${webhookId}:`, error) + return { success: false, webhookId, error: errorMessage } + } + } + + for (const webhookData of activeWebhooks) { + running.push(enqueue(webhookData)) + + if (running.length >= CONCURRENCY) { + const result = await Promise.race(running) + running.splice(running.indexOf(result), 1) + settledResults.push(result) + } + } + + while (running.length) { + const result = await Promise.race(running) + running.splice(running.indexOf(result), 1) + settledResults.push(result) + } + + const results = settledResults const summary = { total: results.length, @@ -301,8 +337,8 @@ async function searchEmails(accessToken: string, config: GmailWebhookConfig, req // Less than an hour ago // Calculate buffer in seconds - the greater of: // 1. Twice the configured polling interval (or 2 minutes if not set) - // 2. At least 5 minutes (300 seconds) - const bufferSeconds = Math.max((config.pollingInterval || 2) * 60 * 2, 300) + // 2. At least 3 minutes (180 seconds) + const bufferSeconds = Math.max((config.pollingInterval || 2) * 60 * 2, 180) // Calculate the cutoff time with buffer const cutoffTime = new Date(lastCheckedTime.getTime() - bufferSeconds * 1000) @@ -448,6 +484,20 @@ async function processEmails( for (const email of emails) { try { + // Deduplicate at Redis level (guards against races between cron runs) + const dedupeKey = `gmail:${webhookData.id}:${email.id}` + try { + const alreadyProcessed = await hasProcessedMessage(dedupeKey) + if (alreadyProcessed) { + logger.info( + `[${requestId}] Duplicate email ${email.id} for webhook ${webhookData.id} – skipping` + ) + continue + } + } catch (err) { + logger.warn(`[${requestId}] Redis check failed for ${email.id}, continuing`, err) + } + // Extract useful information from email to create a simplified payload // First, extract headers into a map for easy access const headers: Record = {} @@ -525,20 +575,16 @@ async function processEmails( } // Create simplified email object - const simplifiedEmail = { + const simplifiedEmail: SimplifiedEmail = { id: email.id, threadId: email.threadId, - // Basic info subject: headers.subject || '[No Subject]', from: headers.from || '', to: headers.to || '', cc: headers.cc || '', date: date, - // Content bodyText: textContent, bodyHtml: htmlContent, - snippet: email.snippet || '', - // Metadata labels: email.labelIds || [], hasAttachments: attachments.length > 0, attachments: attachments, @@ -560,6 +606,7 @@ async function processEmails( headers: { 'Content-Type': 'application/json', 'X-Webhook-Secret': webhookData.secret || '', + 'User-Agent': 'SimStudio/1.0', }, body: JSON.stringify(payload), }) @@ -579,6 +626,8 @@ async function processEmails( } processedCount++ + + await markMessageAsProcessed(dedupeKey) } catch (error) { const errorMessage = error instanceof Error ? error.message : 'Unknown error' logger.error(`[${requestId}] Error processing email ${email.id}:`, errorMessage) diff --git a/apps/sim/lib/webhooks/utils.ts b/apps/sim/lib/webhooks/utils.ts index efe0a027a8..15878ddaa6 100644 --- a/apps/sim/lib/webhooks/utils.ts +++ b/apps/sim/lib/webhooks/utils.ts @@ -257,6 +257,11 @@ export function formatWebhookInput( } else { return null } + } else if (foundWebhook.provider === 'gmail') { + if (body && typeof body === 'object' && 'email' in body) { + return body // { email: {...}, timestamp: ... } + } + return body } else { // Generic format for Slack and other providers return {