Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
123 changes: 92 additions & 31 deletions src/jobs/process-webhook-retries.ts
Original file line number Diff line number Diff line change
@@ -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<typeof setInterval> | 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<void> {
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<void> {
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
};
47 changes: 47 additions & 0 deletions src/repositories/webhook-attempt.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
import { prisma } from '../utils/prisma';

export class WebhookAttemptRepository {
static async findDueForRetry(): Promise<any[]> {
return prisma.webhookAttempt.findMany({
where: {
nextRetryAt: { lte: new Date() },
succeededAt: null,
failedAt: null
},
orderBy: { nextRetryAt: 'asc' }
});
}

static async markSucceeded(id: string, succeededAt: Date): Promise<void> {
await prisma.webhookAttempt.update({
where: { id },
data: { succeededAt, retryCount: { increment: 1 } }
});
}

static async markFailed(id: string, failedAt: Date, errorMessage: string | null): Promise<void> {
await prisma.webhookAttempt.update({
where: { id },
data: { failedAt, errorMessage, retryCount: { increment: 1 } }
});
}

static async updateRetry(id: string, retryCount: number, nextRetryAt: Date): Promise<void> {
await prisma.webhookAttempt.update({
where: { id },
data: { retryCount, nextRetryAt }
});
}

static async create(webhookId: string, event: string, payload: unknown): Promise<any> {
return prisma.webhookAttempt.create({
data: {
webhookId,
event,
payload,
retryCount: 0,
nextRetryAt: new Date()
}
});
}
}
Loading