mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
fix(db): fix slow db connection, add redis cache for gmail polling
This commit is contained in:
@@ -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<any>
|
||||
startedAt: number
|
||||
}
|
||||
|
||||
const activePollingTasks = new Map<string, PollingTask>()
|
||||
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(() => {})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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',
|
||||
|
||||
@@ -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<any>[] = []
|
||||
const settledResults: PromiseSettledResult<any>[] = []
|
||||
|
||||
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<string, string> = {}
|
||||
@@ -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)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user