diff --git a/src/services/webhook.ts b/src/services/webhook.ts index 743f10a6..b5d1eeac 100644 --- a/src/services/webhook.ts +++ b/src/services/webhook.ts @@ -84,6 +84,8 @@ interface WebhookServiceOptions { webhookSecret?: string; maxAttempts?: number; baseDelayMs?: number; + maxDelayMs?: number; + jitterFactor?: number; sleep?: (ms: number) => Promise; now?: () => Date; logger?: WebhookLogger; @@ -110,6 +112,52 @@ function wait(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } +/** + * Calculate exponential backoff delay with jitter and max cap. + * @param baseDelayMs - Base delay in milliseconds + * @param attempt - Current attempt number (1-indexed) + * @param maxDelayMs - Maximum delay cap in milliseconds + * @param jitterFactor - Jitter factor (0-1) to add randomness + * @returns Delay in milliseconds + */ +function calculateBackoffDelay( + baseDelayMs: number, + attempt: number, + maxDelayMs: number, + jitterFactor: number, +): number { + const exponentialDelay = baseDelayMs * Math.pow(2, attempt - 1); + const cappedDelay = Math.min(exponentialDelay, maxDelayMs); + const jitter = cappedDelay * jitterFactor * Math.random(); + return Math.floor(cappedDelay + jitter); +} + +/** + * Determine if an error is retryable. + * Retry on: network errors, timeouts, 5xx server errors + * Don't retry on: 4xx client errors (except 429 Too Many Requests) + * @param error - The error that occurred + * @param statusCode - HTTP status code if available + * @returns true if the request should be retried + */ +function isRetryableError( + error: unknown, + statusCode?: number, +): boolean { + if (statusCode !== undefined) { + // Retry on 429 (rate limited) and 5xx server errors + if (statusCode === 429 || (statusCode >= 500 && statusCode < 600)) { + return true; + } + // Don't retry on other 4xx client errors + if (statusCode >= 400 && statusCode < 500) { + return false; + } + } + // Retry on network errors (no status code) + return true; +} + function getStringValue( transaction: Record, ...keys: string[] @@ -160,6 +208,8 @@ export class WebhookService { private readonly webhookSecret: string; private readonly maxAttempts: number; private readonly baseDelayMs: number; + private readonly maxDelayMs: number; + private readonly jitterFactor: number; private readonly sleepImpl: (ms: number) => Promise; private readonly now: () => Date; private readonly logger: WebhookLogger; @@ -173,6 +223,8 @@ export class WebhookService { options.webhookSecret ?? process.env.WEBHOOK_SECRET ?? ""; this.maxAttempts = options.maxAttempts ?? 3; this.baseDelayMs = options.baseDelayMs ?? 500; + this.maxDelayMs = options.maxDelayMs ?? 30000; // 30 seconds max delay + this.jitterFactor = options.jitterFactor ?? 0.2; // 20% jitter this.sleepImpl = options.sleep ?? wait; this.now = options.now ?? (() => new Date()); this.logger = options.logger ?? console; @@ -273,9 +325,11 @@ export class WebhookService { let lastError: string | null = null; let lastStatusCode: number | undefined; let lastAttemptAt: Date | null = null; + let finalAttempt = 0; for (let attempt = 1; attempt <= this.maxAttempts; attempt++) { lastAttemptAt = this.now(); + finalAttempt = attempt; try { const response = await this.fetchImpl(this.webhookUrl, { method: "POST", @@ -306,8 +360,25 @@ export class WebhookService { this.logger.warn( `[webhook] delivery failed event=${event} transactionId=${payload.data.id} attempt=${attempt}/${this.maxAttempts}: ${lastError}`, ); - if (attempt < this.maxAttempts) - await this.sleepImpl(this.baseDelayMs * 2 ** (attempt - 1)); + // Check if we should retry based on error type + if (attempt < this.maxAttempts && isRetryableError(error, lastStatusCode)) { + const delayMs = calculateBackoffDelay( + this.baseDelayMs, + attempt, + this.maxDelayMs, + this.jitterFactor, + ); + this.logger.log( + `[webhook] retrying in ${delayMs}ms event=${event} transactionId=${payload.data.id} attempt=${attempt + 1}/${this.maxAttempts}`, + ); + await this.sleepImpl(delayMs); + } else if (attempt < this.maxAttempts) { + // Non-retryable error, break early + this.logger.warn( + `[webhook] non-retryable error, stopping retries event=${event} transactionId=${payload.data.id}: ${lastError}`, + ); + break; + } } } @@ -316,7 +387,7 @@ export class WebhookService { ); return { status: "failed", - attempts: this.maxAttempts, + attempts: finalAttempt, statusCode: lastStatusCode, lastAttemptAt, deliveredAt: null, @@ -364,9 +435,11 @@ export class WebhookService { let lastError: string | null = null; let lastStatusCode: number | undefined; let lastAttemptAt: Date | null = null; + let finalAttempt = 0; for (let attempt = 1; attempt <= this.maxAttempts; attempt++) { lastAttemptAt = this.now(); + finalAttempt = attempt; try { const response = await this.fetchImpl(this.webhookUrl, { method: "POST", @@ -397,8 +470,25 @@ export class WebhookService { this.logger.warn( `[webhook] delivery failed flat event=${event} transactionId=${payload.transaction_id} attempt=${attempt}/${this.maxAttempts}: ${lastError}`, ); - if (attempt < this.maxAttempts) - await this.sleepImpl(this.baseDelayMs * 2 ** (attempt - 1)); + // Check if we should retry based on error type + if (attempt < this.maxAttempts && isRetryableError(error, lastStatusCode)) { + const delayMs = calculateBackoffDelay( + this.baseDelayMs, + attempt, + this.maxDelayMs, + this.jitterFactor, + ); + this.logger.log( + `[webhook] retrying in ${delayMs}ms flat event=${event} transactionId=${payload.transaction_id} attempt=${attempt + 1}/${this.maxAttempts}`, + ); + await this.sleepImpl(delayMs); + } else if (attempt < this.maxAttempts) { + // Non-retryable error, break early + this.logger.warn( + `[webhook] non-retryable error, stopping retries flat event=${event} transactionId=${payload.transaction_id}: ${lastError}`, + ); + break; + } } } @@ -407,7 +497,7 @@ export class WebhookService { ); return { status: "failed", - attempts: this.maxAttempts, + attempts: finalAttempt, statusCode: lastStatusCode, lastAttemptAt, deliveredAt: null, @@ -456,6 +546,11 @@ export class WebhookService { const errorMessage = error instanceof Error ? error.message : String(error); const attempts = entry.attempts + 1; + // For outbox, we don't have the status code easily available from the error + // We'll treat all errors as potentially retryable for the outbox processor + // since it runs asynchronously and we want to retry on transient failures + const isRetryable = attempts < entry.maxAttempts; + if (attempts >= entry.maxAttempts) { await outboxModel.update(entry.id, { status: "failed", @@ -463,15 +558,31 @@ export class WebhookService { lastAttemptAt: now, errorMessage: `Exhausted retries: ${errorMessage}`, }); - } else { - const backoffMs = this.baseDelayMs * Math.pow(2, attempts - 1); + } else if (isRetryable) { + const delayMs = calculateBackoffDelay( + this.baseDelayMs, + attempts, + this.maxDelayMs, + this.jitterFactor, + ); await outboxModel.update(entry.id, { status: "pending", attempts, lastAttemptAt: now, - nextAttemptAt: new Date(now.getTime() + backoffMs), + nextAttemptAt: new Date(now.getTime() + delayMs), errorMessage, }); + this.logger.log( + `[webhook-outbox] retrying in ${delayMs}ms entry=${entry.id} attempt=${attempts + 1}/${entry.maxAttempts}`, + ); + } else { + // Non-retryable error, mark as failed + await outboxModel.update(entry.id, { + status: "failed", + attempts, + lastAttemptAt: now, + errorMessage: `Non-retryable error: ${errorMessage}`, + }); } this.logger.warn( `[webhook-outbox] Failed to deliver entry=${entry.id} attempt=${attempts}/${entry.maxAttempts}: ${errorMessage}`, diff --git a/tests/services/webhook.test.ts b/tests/services/webhook.test.ts index 745142a3..ad73ba25 100644 --- a/tests/services/webhook.test.ts +++ b/tests/services/webhook.test.ts @@ -86,6 +86,7 @@ describe("WebhookService", () => { webhookSecret: "retry-secret", sleep: sleepMock, baseDelayMs: 250, + jitterFactor: 0, // Disable jitter for deterministic test logger: { log: jest.fn(), warn: jest.fn(), error: jest.fn() }, }); @@ -118,6 +119,159 @@ describe("WebhookService", () => { expect(result.attempts).toBe(0); expect(logger.warn).toHaveBeenCalled(); }); + + it("does not retry on 4xx client errors (except 429)", async () => { + const fetchMock = jest + .fn() + .mockResolvedValue({ ok: false, status: 400 }); // Always return 400 + const sleepMock = jest.fn(async () => undefined); + const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const service = new WebhookService({ + fetchImpl: fetchMock as unknown as typeof fetch, + webhookUrl: "https://example.com/webhooks", + webhookSecret: "retry-secret", + sleep: sleepMock, + baseDelayMs: 250, + jitterFactor: 0, + logger, + }); + + const result = await service.sendTransactionEvent( + "transaction.completed", + buildTransaction(), + ); + + expect(result.status).toBe("failed"); + expect(result.attempts).toBe(1); // Only 1 attempt, no retries + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(sleepMock).not.toHaveBeenCalled(); + }); + + it("retries on 429 rate limit errors", async () => { + const fetchMock = jest + .fn() + .mockResolvedValueOnce({ ok: false, status: 429 }) + .mockResolvedValueOnce({ ok: true, status: 200 }); + const sleepMock = jest.fn(async () => undefined); + const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const service = new WebhookService({ + fetchImpl: fetchMock as unknown as typeof fetch, + webhookUrl: "https://example.com/webhooks", + webhookSecret: "retry-secret", + sleep: sleepMock, + baseDelayMs: 250, + jitterFactor: 0, + logger, + }); + + const result = await service.sendTransactionEvent( + "transaction.completed", + buildTransaction(), + ); + + expect(result.status).toBe("delivered"); + expect(result.attempts).toBe(2); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(sleepMock).toHaveBeenCalledTimes(1); + expect(sleepMock).toHaveBeenCalledWith(250); + }); + + it("retries on 5xx server errors", async () => { + const fetchMock = jest + .fn() + .mockResolvedValueOnce({ ok: false, status: 503 }) + .mockResolvedValueOnce({ ok: true, status: 200 }); + const sleepMock = jest.fn(async () => undefined); + const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const service = new WebhookService({ + fetchImpl: fetchMock as unknown as typeof fetch, + webhookUrl: "https://example.com/webhooks", + webhookSecret: "retry-secret", + sleep: sleepMock, + baseDelayMs: 250, + jitterFactor: 0, + logger, + }); + + const result = await service.sendTransactionEvent( + "transaction.completed", + buildTransaction(), + ); + + expect(result.status).toBe("delivered"); + expect(result.attempts).toBe(2); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(sleepMock).toHaveBeenCalledTimes(1); + }); + + it("respects max delay cap", async () => { + const fetchMock = jest + .fn() + .mockRejectedValue(new Error("network down")); // Always fails + const sleepMock = jest.fn(async () => undefined); + const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const service = new WebhookService({ + fetchImpl: fetchMock as unknown as typeof fetch, + webhookUrl: "https://example.com/webhooks", + webhookSecret: "retry-secret", + sleep: sleepMock, + baseDelayMs: 1000, + maxDelayMs: 500, // Cap at 500ms + jitterFactor: 0, + maxAttempts: 4, + logger, + }); + + const result = await service.sendTransactionEvent( + "transaction.completed", + buildTransaction(), + ); + + expect(result.status).toBe("failed"); + expect(result.attempts).toBe(4); + // Check that delays don't exceed maxDelayMs (500) + // Attempt 1: min(1000*2^0, 500) = 500 + // Attempt 2: min(1000*2^1, 500) = 500 + // Attempt 3: min(1000*2^2, 500) = 500 + expect(sleepMock).toHaveBeenNthCalledWith(1, 500); + expect(sleepMock).toHaveBeenNthCalledWith(2, 500); + expect(sleepMock).toHaveBeenNthCalledWith(3, 500); + }); + + it("applies jitter when enabled", async () => { + const fetchMock = jest + .fn() + .mockRejectedValueOnce(new Error("network down")) + .mockResolvedValueOnce({ ok: true, status: 200 }); + const sleepMock = jest.fn(async () => undefined); + const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const service = new WebhookService({ + fetchImpl: fetchMock as unknown as typeof fetch, + webhookUrl: "https://example.com/webhooks", + webhookSecret: "retry-secret", + sleep: sleepMock, + baseDelayMs: 1000, + jitterFactor: 0.5, // 50% jitter + logger, + }); + + const result = await service.sendTransactionEvent( + "transaction.completed", + buildTransaction(), + ); + + expect(result.status).toBe("delivered"); + expect(result.attempts).toBe(2); + // With 50% jitter, delay should be between 1000 and 1500 + const delay = sleepMock.mock.calls[0][0]; + expect(delay).toBeGreaterThanOrEqual(1000); + expect(delay).toBeLessThanOrEqual(1500); + }); }); describe("notifyTransactionWebhook", () => {