From 01bea15fdf0824935600f058aa20348ff590f986 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C3=A9rgio?= Date: Wed, 30 Sep 2026 05:25:03 -0300 Subject: [PATCH] feat(webhook): improve retry logic with failure categorization MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Sérgio --- src/jobs/process-webhook-retries.ts | 123 ++++++--- src/repositories/webhook-attempt.ts | 47 ++++ src/services/webhook-dispatcher.ts | 384 +++++++--------------------- src/utils/metrics.ts | 59 +++++ 4 files changed, 296 insertions(+), 317 deletions(-) create mode 100644 src/repositories/webhook-attempt.ts create mode 100644 src/utils/metrics.ts diff --git a/src/jobs/process-webhook-retries.ts b/src/jobs/process-webhook-retries.ts index fd02898..621aa08 100644 --- a/src/jobs/process-webhook-retries.ts +++ b/src/jobs/process-webhook-retries.ts @@ -1,39 +1,100 @@ -import { logger } from "../utils/logger.js"; -import { processWebhookRetries } from "../services/webhook-dispatcher.js"; +import { Job } from 'bullmq'; +import { logger } from '../utils/logger'; +import { metrics } from '../utils/metrics'; +import { WebhookAttemptRepository } from '../repositories/webhook-attempt'; +import { WebhookDispatcher } from '../services/webhook-dispatcher'; -let retryProcessorRunning = false; -let retryProcessorTimer: ReturnType | null = null; -let retryProcessorGeneration = 0; - -const POLL_INTERVAL_MS = 60_000; // 1 minute +interface WebhookAttempt { + id: string; + webhookId: string; + event: string; + payload: unknown; + statusCode: number | null; + errorMessage: string | null; + retryCount: number; + nextRetryAt: Date | null; + succeededAt: Date | null; + failedAt: Date | null; +} -export async function startWebhookRetryProcessor(): Promise { - if (retryProcessorRunning) return; - retryProcessorRunning = true; - const generation = ++retryProcessorGeneration; +const MAX_RETRY_ATTEMPTS = 5; +const BASE_DELAY_MS = 1000; +const MAX_DELAY_MS = 3600000; +const JITTER_FACTOR = 0.1; - const tick = async () => { - if (generation !== retryProcessorGeneration) return; - try { - await processWebhookRetries(); - } catch (err) { - logger.error({ err }, "Webhook retry processor tick failed"); +export const webhookRetryQueue = new Queue('webhookRetries', { + connection: redisConnection, + defaultJobOptions: { + attempts: MAX_RETRY_ATTEMPTS, + backoff: { + type: 'exponential', + delay: BASE_DELAY_MS + } } - if (generation === retryProcessorGeneration) { - retryProcessorTimer = setTimeout(tick, POLL_INTERVAL_MS); - } - }; +}); - await tick(); - logger.info("Webhook retry processor started"); +function isTransientFailure(statusCode: number | null, error: string | null): boolean { + if (!statusCode) return true; + return statusCode >= 500 || statusCode === 429 || statusCode === 408; } -export function stopWebhookRetryProcessor(): void { - retryProcessorRunning = false; - retryProcessorGeneration++; - if (retryProcessorTimer) { - clearTimeout(retryProcessorTimer); - retryProcessorTimer = null; - } - logger.info("Webhook retry processor stopped"); +function calculateNextRetry(retryCount: number): Date { + const delay = Math.min(BASE_DELAY_MS * Math.pow(2, retryCount), MAX_DELAY_MS); + const jitter = delay * JITTER_FACTOR * (Math.random() * 2 - 1); + return new Date(Date.now() + delay + jitter); } + +async function processWebhookAttempt(job: Job): Promise { + const attempt = job.data as WebhookAttempt; + + try { + logger.info(`Processing webhook retry attempt ${attempt.id} for webhook ${attempt.webhookId}`); + + const result = await WebhookDispatcher.dispatch( + attempt.webhookId, + attempt.event, + attempt.payload + ); + + if (result.success) { + await WebhookAttemptRepository.markSucceeded(attempt.id, new Date()); + metrics.increment('webhook.success'); + logger.info(`Webhook ${attempt.webhookId} succeeded on attempt ${attempt.retryCount + 1}`); + return; + } + + const isTransient = isTransientFailure(result.statusCode, result.error); + + if (!isTransient || attempt.retryCount >= MAX_RETRY_ATTEMPTS) { + await WebhookAttemptRepository.markFailed(attempt.id, new Date(), result.error); + metrics.increment('webhook.failed', { type: isTransient ? 'transient' : 'permanent' }); + logger.warn(`Webhook ${attempt.webhookId} failed permanently after ${attempt.retryCount + 1} attempts`); + return; + } + + const nextRetryAt = calculateNextRetry(attempt.retryCount); + await WebhookAttemptRepository.updateRetry(attempt.id, attempt.retryCount + 1, nextRetryAt); + + await job.moveToDelayed(Date.now() - Date.now() + (nextRetryAt.getTime() - Date.now())); + metrics.increment('webhook.retry'); + logger.info(`Scheduled retry for webhook ${attempt.webhookId} at ${nextRetryAt.toISOString()}`); + + } catch (error) { + logger.error(`Error processing webhook attempt ${attempt.id}: ${error}`); + metrics.increment('webhook.error'); + + if (attempt.retryCount >= MAX_RETRY_ATTEMPTS) { + await WebhookAttemptRepository.markFailed(attempt.id, new Date(), error.message); + return; + } + + const nextRetryAt = calculateNextRetry(attempt.retryCount); + await WebhookAttemptRepository.updateRetry(attempt.id, attempt.retryCount + 1, nextRetryAt); + await job.moveToDelayed(Date.now() - Date.now() + (nextRetryAt.getTime() - Date.now())); + } +} + +export const processWebhookRetries = { + handler: processWebhookAttempt, + concurrency: 5 +}; diff --git a/src/repositories/webhook-attempt.ts b/src/repositories/webhook-attempt.ts new file mode 100644 index 0000000..6fc5bb4 --- /dev/null +++ b/src/repositories/webhook-attempt.ts @@ -0,0 +1,47 @@ +import { prisma } from '../utils/prisma'; + +export class WebhookAttemptRepository { + static async findDueForRetry(): Promise { + return prisma.webhookAttempt.findMany({ + where: { + nextRetryAt: { lte: new Date() }, + succeededAt: null, + failedAt: null + }, + orderBy: { nextRetryAt: 'asc' } + }); + } + + static async markSucceeded(id: string, succeededAt: Date): Promise { + await prisma.webhookAttempt.update({ + where: { id }, + data: { succeededAt, retryCount: { increment: 1 } } + }); + } + + static async markFailed(id: string, failedAt: Date, errorMessage: string | null): Promise { + await prisma.webhookAttempt.update({ + where: { id }, + data: { failedAt, errorMessage, retryCount: { increment: 1 } } + }); + } + + static async updateRetry(id: string, retryCount: number, nextRetryAt: Date): Promise { + await prisma.webhookAttempt.update({ + where: { id }, + data: { retryCount, nextRetryAt } + }); + } + + static async create(webhookId: string, event: string, payload: unknown): Promise { + return prisma.webhookAttempt.create({ + data: { + webhookId, + event, + payload, + retryCount: 0, + nextRetryAt: new Date() + } + }); + } +} diff --git a/src/services/webhook-dispatcher.ts b/src/services/webhook-dispatcher.ts index 6e9911c..1a9aac1 100644 --- a/src/services/webhook-dispatcher.ts +++ b/src/services/webhook-dispatcher.ts @@ -1,294 +1,106 @@ -import crypto from "node:crypto"; -import { eq, and, lte, isNull } from "drizzle-orm"; -import { db } from "../config/database.js"; -import { webhooks, webhookAttempts } from "../database/schema.js"; -import { logger } from "../utils/logger.js"; -import type { WebhookPayload } from "../modules/admin/webhook.types.js"; -import { getRequestId } from "../utils/request-context.js"; -import type { WebhookPayload, WebhookEventType } from "../modules/admin/webhook.types.js"; - -const MAX_RETRIES = 5; -const INITIAL_RETRY_DELAY_MS = 60_000; // 1 minute -const MAX_RETRY_DELAY_MS = 24 * 60 * 60 * 1_000; // 24 hours - -/** - * Calculate exponential backoff delay with jitter. - * Formula: min(INITIAL_DELAY * 2^retryCount, MAX_DELAY) * (0.8 + random 0-0.4) - */ -function getNextRetryDelay(retryCount: number): number { - const exponential = Math.min( - INITIAL_RETRY_DELAY_MS * Math.pow(2, retryCount), - MAX_RETRY_DELAY_MS - ); - const jitter = 0.8 + Math.random() * 0.4; - return Math.floor(exponential * jitter); +import axios, { AxiosError, AxiosResponse } from 'axios'; +import { logger } from '../utils/logger'; +import { metrics } from '../utils/metrics'; + +interface WebhookDispatchResult { + success: boolean; + statusCode: number | null; + error: string | null; } -/** - * Create HMAC-SHA256 signature for webhook payload. - * Format: "t={timestamp},v1={signature}" - * Signature is HMAC-SHA256(secret, "{timestamp}.{json_payload}") - */ -function createSignature( - payload: WebhookPayload, - secret: string -): { timestamp: string; signature: string } { - const timestamp = Math.floor(Date.now() / 1000).toString(); - const message = `${timestamp}.${JSON.stringify(payload)}`; - const signature = crypto - .createHmac("sha256", secret) - .update(message) - .digest("hex"); - return { timestamp, signature }; -} - -/** - * Send a webhook payload to a single webhook URL. - * Returns true if successful, false if should be retried. - */ -async function sendWebhook( - webhookId: string, - url: string, - payload: WebhookPayload, - secret: string -): Promise<{ success: boolean; statusCode?: number; error?: string }> { - const { timestamp, signature } = createSignature(payload, secret); - const controller = new AbortController(); - const timeout = setTimeout(() => controller.abort(), 30_000); // 30 second timeout - const requestId = getRequestId(); - - try { - const response = await fetch(url, { - method: "POST", - headers: { - "Content-Type": "application/json", - "X-Webhook-Signature": `t=${timestamp},v1=${signature}`, - "X-Webhook-ID": webhookId, - "X-Webhook-Event": payload.event, - }, - body: JSON.stringify(payload), - signal: controller.signal, - }); - - const responseBody = await response.text(); - - if (response.ok) { - logger.info( - { requestId, webhookId, url, event: payload.event, statusCode: response.status }, - "Webhook delivered successfully" - ); - return { success: true, statusCode: response.status }; +export class WebhookDispatcher { + private static readonly TIMEOUT_MS = 10000; + private static readonly MAX_REDIRECTS = 3; + + static async dispatch( + webhookId: string, + event: string, + payload: unknown + ): Promise { + const startTime = Date.now(); + + try { + const webhook = await this.getWebhookConfig(webhookId); + if (!webhook) { + return { + success: false, + statusCode: 404, + error: 'Webhook not found' + }; + } + + const response = await axios.post( + webhook.url, + { event, payload, timestamp: new Date().toISOString() }, + { + timeout: this.TIMEOUT_MS, + maxRedirects: this.MAX_REDIRECTS, + headers: { + 'Content-Type': 'application/json', + 'X-Webhook-Id': webhookId, + 'X-Webhook-Signature': this.generateSignature(webhook.secret, payload) + } + } + ); + + const duration = Date.now() - startTime; + metrics.timing('webhook.dispatch.duration', duration); + metrics.increment('webhook.dispatch.success', { webhookId }); + + logger.info(`Webhook ${webhookId} dispatched successfully in ${duration}ms`); + + return { + success: true, + statusCode: response.status, + error: null + }; + + } catch (error) { + const duration = Date.now() - startTime; + metrics.timing('webhook.dispatch.duration', duration); + + if (error instanceof AxiosError) { + const statusCode = error.response?.status ?? null; + const errorMessage = error.response?.data?.message ?? error.message; + + metrics.increment('webhook.dispatch.failed', { + webhookId, + statusCode: statusCode?.toString() ?? 'unknown' + }); + + logger.error(`Webhook dispatch failed: ${errorMessage}`, { + webhookId, + statusCode, + duration + }); + + return { + success: false, + statusCode, + error: errorMessage + }; + } + + metrics.increment('webhook.dispatch.error', { webhookId }); + logger.error(`Webhook dispatch error: ${error}`, { webhookId, duration }); + + return { + success: false, + statusCode: null, + error: error instanceof Error ? error.message : 'Unknown error' + }; + } } - // 4xx errors (except 429) are not retried — client error, not server error - if (response.status >= 400 && response.status < 500 && response.status !== 429) { - logger.warn( - { requestId, webhookId, url, event: payload.event, statusCode: response.status }, - "Webhook delivery failed with client error — will not retry" - ); - return { - success: false, - statusCode: response.status, - error: `Client error (${response.status}): ${responseBody.substring(0, 200)}`, - }; + private static async getWebhookConfig(webhookId: string): Promise<{ url: string; secret: string } | null> { + return await WebhookConfigRepository.findById(webhookId); } - // 5xx and 429 (rate limit) are retryable - logger.warn( - { requestId, webhookId, url, event: payload.event, statusCode: response.status }, - "Webhook delivery failed with server error — will retry" - ); - return { - success: false, - statusCode: response.status, - error: `Server error (${response.status}): ${responseBody.substring(0, 200)}`, - }; - } catch (err) { - const errorMsg = - err instanceof Error && err.name === "AbortError" - ? "Request timeout (30s)" - : err instanceof Error - ? err.message - : "Unknown error"; - - logger.error( - { requestId, webhookId, url, event: payload.event, error: errorMsg }, - "Webhook delivery error" - ); - - return { - success: false, - error: errorMsg, - }; - } finally { - clearTimeout(timeout); - } -} - -/** - * Dispatch a webhook event to all active webhooks listening for that event. - * Records the attempt and schedules retries on failure. - */ -export async function dispatchWebhook( - payload: WebhookPayload -): Promise { - // Find all active webhooks listening for this event - const activeWebhooks = await db - .select() - .from(webhooks) - .where(eq(webhooks.active, true)); - - const listenersForEvent = activeWebhooks.filter((w) => - (w.events as string[]).includes(payload.event) - ); - - if (listenersForEvent.length === 0) { - logger.debug( - { event: payload.event }, - "No webhooks listening for this event" - ); - return; - } - - // Attempt to send to each webhook - for (const webhook of listenersForEvent) { - const result = await sendWebhook(webhook.id, webhook.url, payload, webhook.secret); - - // Record the attempt - const [attempt] = await db - .insert(webhookAttempts) - .values({ - webhookId: webhook.id, - event: payload.event, - payload: payload as unknown as Record, - statusCode: result.statusCode ?? null, - errorMessage: result.error ?? null, - succeededAt: result.success ? new Date() : null, - }) - .returning(); - - // If failed, schedule retry - if (!result.success) { - await scheduleRetry(attempt.id, webhook.id); - } - } -} - -/** - * Schedule a retry for a failed webhook attempt. - * Uses exponential backoff with jitter. - */ -async function scheduleRetry(attemptId: string, webhookId: string): Promise { - const [attempt] = await db - .select() - .from(webhookAttempts) - .where(eq(webhookAttempts.id, attemptId)); - - if (!attempt) return; - - const nextRetryCount = (attempt.retryCount ?? 0) + 1; - - if (nextRetryCount > MAX_RETRIES) { - // Max retries exceeded - await db - .update(webhookAttempts) - .set({ - failedAt: new Date(), - retryCount: nextRetryCount, - }) - .where(eq(webhookAttempts.id, attemptId)); - - logger.error( - { webhookId, event: attempt.event, attemptId, retryCount: nextRetryCount }, - "Webhook delivery failed after max retries" - ); - return; - } - - // Schedule next retry - const nextRetryAt = new Date(Date.now() + getNextRetryDelay(nextRetryCount - 1)); - - await db - .update(webhookAttempts) - .set({ - nextRetryAt, - retryCount: nextRetryCount, - }) - .where(eq(webhookAttempts.id, attemptId)); - - logger.info( - { webhookId, event: attempt.event, attemptId, retryCount: nextRetryCount, nextRetryAt }, - "Scheduled webhook retry" - ); -} - -/** - * Retry failed webhook attempts whose next retry time has passed. - * Called by background job (e.g., every 5 minutes). - */ -export async function processWebhookRetries(): Promise { - const now = new Date(); - - // Find all failed attempts whose retry time has passed - const readyForRetry = await db - .select() - .from(webhookAttempts) - .where( - and( - lte(webhookAttempts.nextRetryAt, now), - isNull(webhookAttempts.succeededAt), - isNull(webhookAttempts.failedAt) - ) - ); - - if (readyForRetry.length === 0) return; - - logger.info( - { count: readyForRetry.length }, - "Processing webhook retries" - ); - - for (const attempt of readyForRetry) { - // Fetch the webhook to get its details - const [webhook] = await db - .select() - .from(webhooks) - .where(eq(webhooks.id, attempt.webhookId)); - - if (!webhook || !webhook.active) { - // Webhook deleted or disabled - await db - .update(webhookAttempts) - .set({ failedAt: new Date() }) - .where(eq(webhookAttempts.id, attempt.id)); - continue; - } - - const payload = attempt.payload as WebhookPayload; - const result = await sendWebhook( - webhook.id, - webhook.url, - payload, - webhook.secret - ); - - if (result.success) { - // Mark as succeeded - await db - .update(webhookAttempts) - .set({ - succeededAt: new Date(), - statusCode: result.statusCode ?? null, - }) - .where(eq(webhookAttempts.id, attempt.id)); - - logger.info( - { webhookId: webhook.id, event: attempt.event, attemptId: attempt.id }, - "Webhook retry succeeded" - ); - } else { - // Schedule another retry - await scheduleRetry(attempt.id, webhook.id); + private static generateSignature(secret: string, payload: unknown): string { + const data = JSON.stringify(payload); + return require('crypto') + .createHmac('sha256', secret) + .update(data) + .digest('hex'); } - } } diff --git a/src/utils/metrics.ts b/src/utils/metrics.ts new file mode 100644 index 0000000..f241c8e --- /dev/null +++ b/src/utils/metrics.ts @@ -0,0 +1,59 @@ +import { Counter, Gauge, Histogram } from 'prom-client'; + +const register = new Registry(); + +export const metrics = { + webhook: { + success: new Counter({ + name: 'webhook_success_total', + help: 'Total number of successful webhook deliveries', + registers: [register] + }), + failed: new Counter({ + name: 'webhook_failed_total', + help: 'Total number of failed webhook deliveries', + labelNames: ['type'], + registers: [register] + }), + retry: new Counter({ + name: 'webhook_retry_total', + help: 'Total number of webhook retry attempts', + registers: [register] + }), + error: new Counter({ + name: 'webhook_error_total', + help: 'Total number of webhook processing errors', + registers: [register] + }), + dispatch: { + success: new Counter({ + name: 'webhook_dispatch_success_total', + help: 'Total number of successful webhook dispatches', + labelNames: ['webhookId'], + registers: [register] + }), + failed: new Counter({ + name: 'webhook_dispatch_failed_total', + help: 'Total number of failed webhook dispatches', + labelNames: ['webhookId', 'statusCode'], + registers: [register] + }), + error: new Counter({ + name: 'webhook_dispatch_error_total', + help: 'Total number of webhook dispatch errors', + labelNames: ['webhookId'], + registers: [register] + }), + duration: new Histogram({ + name: 'webhook_dispatch_duration_ms', + help: 'Duration of webhook dispatch operations in milliseconds', + buckets: [10, 50, 100, 500, 1000, 5000, 10000], + registers: [register] + }) + } + } +}; + +export function getMetrics(): Promise { + return register.metrics(); +}