diff --git a/src/index.ts b/src/index.ts index 16017bd..eabc988 100644 --- a/src/index.ts +++ b/src/index.ts @@ -6,6 +6,7 @@ import helmet from "helmet"; import { logger } from "./logger"; import { openApiSpec } from "./openapi"; import { isSafeWebhookUrl } from "./utils/webhookUrl"; +import { generateWebhookSecret, deliverEventToWebhooks, webhookDlq } from "./webhook/delivery"; import { resolveClientIp } from "./utils/clientIp"; import { getStoreAdapter } from "./persistence"; import { @@ -23,6 +24,7 @@ import { pairKey, defaultMeta, recordEvent, + setWebhookDeliveryHook, trimEventLog, EVENT_LOG_CAP, EVENT_LOG_CAP_MAX, @@ -131,6 +133,10 @@ const EVENT_PAYLOAD_MAX_ARRAY_ITEMS = 32; const EVENT_PAYLOAD_MAX_DEPTH = 3; const app = express(); +// Register webhook delivery hook so every recorded event is delivered to subscribers. +setWebhookDeliveryHook((event) => { + deliverEventToWebhooks(event, webhookStore).catch(() => {}); +}); // --- Persistence Hydration on startup --- export const hydrationPromise = (async () => { @@ -1997,10 +2003,11 @@ app.post( const deduped = validateWebhookEvents(res, req, events); if (deduped === null) return; const id = `wh_${randomUUID().replace(/-/g, "").slice(0, 16)}`; - webhookStore.set(id, { url, events: deduped, createdAt: Date.now() }); + const secret = generateWebhookSecret(); + webhookStore.set(id, { url, events: deduped, createdAt: Date.now(), secret }); // Record id and url only — never any webhook secret material. recordEvent("webhook.created", { id, url }); - res.status(201).json({ id, url, events: deduped }); + res.status(201).json({ id, url, events: deduped, secret }); }, ); @@ -2117,6 +2124,47 @@ app.patch("/api/v1/webhooks/:id", (req: Request, res: Response) => { res.json({ id, ...updated }); }); +/** + * List dead-lettered webhook deliveries. + * + * @route GET /api/v1/webhooks/dlq + */ +app.get("/api/v1/webhooks/dlq", (req: Request, res: Response) => { + const query = (req.query ?? {}) as Record; + const limit = Math.min( + Math.max(1, Number(query["limit"] ?? 20)), + 100, + ); + const entries = Array.from(webhookDlq.values()) + .sort((a, b) => b.deadAt - a.deadAt) + .slice(0, limit); + res.json({ entries, total: webhookDlq.size }); +}); + +/** + * Replay a dead-lettered webhook delivery. + * + * @route POST /api/v1/webhooks/dlq/:id/replay + */ +app.post("/api/v1/webhooks/dlq/:id/replay", (req: Request, res: Response) => { + const id = req.params.id ?? ""; + const entry = webhookDlq.get(id); + if (!entry) { + sendError(res, req, 404, "not_found", `dlq entry ${id} not found`); + return; + } + const record = webhookStore.get(entry.webhookId); + if (!record) { + sendError(res, req, 404, "not_found", `webhook ${entry.webhookId} not found`); + return; + } + webhookDlq.delete(id); + deliverEventToWebhooks(entry.event, new Map([[entry.webhookId, record]])).catch( + () => {}, + ); + res.json({ replayed: true, id }); +}); + /** * Normalize the `:source`/`:destination` route params to their canonical asset * codes via {@link normalizeAsset}. On invalid input a `400 invalid_request` is diff --git a/src/stores.ts b/src/stores.ts index 2357414..0203aad 100644 --- a/src/stores.ts +++ b/src/stores.ts @@ -156,6 +156,7 @@ export const verifyApiKeySecret = ( /** Record stored for each registered webhook. */ export type WebhookRecord = { + secret: string; url: string; events: string[]; createdAt: number; @@ -300,6 +301,14 @@ export const trimEventLog = (cap: number): void => { * enforces this at the call site so stray string literals are caught at compile time. * @param payload - Arbitrary structured data attached to the event. */ +/** Optional hook for delivering recorded events to webhook subscribers. */ +let webhookDeliveryHook: ((event: AppEvent) => void) | undefined; + +/** Register a callback that fires after every event is recorded. */ +export function setWebhookDeliveryHook(hook: (event: AppEvent) => void): void { + webhookDeliveryHook = hook; +} + export const recordEvent = ( type: EventType, payload: Record, @@ -313,6 +322,13 @@ export const recordEvent = ( eventLog.push(event); const cap = effectiveEventLogCap(); if (eventLog.length > cap) eventLog.shift(); + if (webhookDeliveryHook) { + try { + webhookDeliveryHook(event); + } catch { + // Delivery errors must not break event recording. + } + } return event; }; @@ -480,7 +496,12 @@ export const hydrateFromSnapshot = (snapshot: unknown): void => { item.length === 2 && typeof item[0] === "string" ) { - webhookStore.set(item[0], item[1] as WebhookRecord); + const record = item[1] as WebhookRecord; + // Backfill secret for webhooks created before the signing feature. + if (!record.secret) { + record.secret = randomBytes(32).toString("base64url"); + } + webhookStore.set(item[0], record); } } } diff --git a/src/webhook/delivery.test.ts b/src/webhook/delivery.test.ts new file mode 100644 index 0000000..483fa4b --- /dev/null +++ b/src/webhook/delivery.test.ts @@ -0,0 +1,120 @@ +import { describe, it, expect, vi } from "vitest"; +import { + generateWebhookSecret, + signWebhookPayload, + verifyWebhookPayload, + webhookDlq, + deliverEventToWebhooks, + WEBHOOK_MAX_RETRIES, +} from "./delivery"; +import type { AppEvent, WebhookRecord } from "../stores"; + +describe("webhook delivery", () => { + it("generates a URL-safe secret", () => { + const secret = generateWebhookSecret(); + expect(secret).toMatch(/^[A-Za-z0-9_-]+$/); + expect(secret.length).toBeGreaterThanOrEqual(40); + }); + + it("signs and verifies a payload", () => { + const secret = generateWebhookSecret(); + const payload = JSON.stringify({ event: "test" }); + const { signature, timestamp } = signWebhookPayload(secret, payload); + expect(signature).toMatch(/^t=\d+,v0=[0-9a-f]{64}$/); + expect(verifyWebhookPayload(secret, payload, signature)).toBe(true); + }); + + it("rejects a tampered signature", () => { + const secret = generateWebhookSecret(); + const payload = JSON.stringify({ event: "test" }); + const { signature } = signWebhookPayload(secret, payload); + expect(verifyWebhookPayload(secret, payload + "x", signature)).toBe(false); + }); + + it("rejects an expired signature", () => { + const secret = generateWebhookSecret(); + const payload = JSON.stringify({ event: "test" }); + const oldTs = Math.floor(Date.now() / 1000) - 400; + const { signature } = signWebhookPayload(secret, payload, oldTs); + expect(verifyWebhookPayload(secret, payload, signature)).toBe(false); + }); + + it("delivers to matching webhooks", async () => { + webhookDlq.clear(); + const fetchMock = vi.fn().mockResolvedValue({ ok: true, status: 200 }); + const event: AppEvent = { + id: "evt-1", + ts: Date.now(), + type: "pair.registered", + payload: { source: "XLM", destination: "USDC" }, + }; + const webhooks = new Map([ + [ + "wh-1", + { + url: "https://example.com/hook", + events: ["pair.registered"], + createdAt: Date.now(), + secret: generateWebhookSecret(), + }, + ], + [ + "wh-2", + { + url: "https://example.com/other", + events: ["pair.disabled"], + createdAt: Date.now(), + secret: generateWebhookSecret(), + }, + ], + ]); + + await deliverEventToWebhooks(event, webhooks, fetchMock as any); + // Give async delivery a tick + await new Promise((r) => setTimeout(r, 10)); + + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(fetchMock).toHaveBeenCalledWith( + "https://example.com/hook", + expect.objectContaining({ + method: "POST", + headers: expect.objectContaining({ + "Content-Type": "application/json", + "X-Signature": expect.stringMatching(/^t=\d+,v0=[0-9a-f]{64}$/), + }), + }), + ); + }); + + it("dead-letters after max retries on 5xx", async () => { + webhookDlq.clear(); + const fetchMock = vi.fn().mockResolvedValue({ ok: false, status: 503 }); + const event: AppEvent = { + id: "evt-2", + ts: Date.now(), + type: "pair.registered", + payload: {}, + }; + const webhooks = new Map([ + [ + "wh-1", + { + url: "https://example.com/hook", + events: ["pair.registered"], + createdAt: Date.now(), + secret: generateWebhookSecret(), + }, + ], + ]); + + await deliverEventToWebhooks(event, webhooks, fetchMock as any); + // Wait for all retries + backoff + await new Promise((r) => setTimeout(r, 8000)); + + expect(fetchMock).toHaveBeenCalledTimes(WEBHOOK_MAX_RETRIES); + expect(webhookDlq.size).toBe(1); + const entry = Array.from(webhookDlq.values())[0]; + expect(entry.webhookId).toBe("wh-1"); + expect(entry.attempts).toBe(WEBHOOK_MAX_RETRIES); + }); +}); diff --git a/src/webhook/delivery.ts b/src/webhook/delivery.ts new file mode 100644 index 0000000..149be19 --- /dev/null +++ b/src/webhook/delivery.ts @@ -0,0 +1,163 @@ +/** + * @module webhook/delivery + * @description Signed webhook delivery with bounded retries and dead-letter queue. + * + * Each event is delivered to every matching webhook subscriber: + * 1. Payload is signed with HMAC-SHA256 (X-Signature header + timestamp) + * 2. Delivery is attempted with exponential backoff on 5xx/timeout + * 3. After MAX_RETRIES attempts, the event moves to the dead-letter queue + */ + +import { createHmac, randomBytes } from "node:crypto"; +import type { AppEvent, WebhookRecord } from "../stores"; + +/** Maximum delivery attempts before an event is dead-lettered. */ +export const WEBHOOK_MAX_RETRIES = 3; + +/** Base delay in ms for exponential backoff (1s, 2s, 4s). */ +export const WEBHOOK_RETRY_BASE_MS = 1000; + +/** How long a delivered payload signature remains valid (seconds). */ +export const WEBHOOK_SIG_TTL_SECS = 300; + +/** Generate a URL-safe random secret for a new webhook subscription. */ +export function generateWebhookSecret(): string { + return randomBytes(32).toString("base64url"); +} + +/** Sign a webhook payload with HMAC-SHA256. Returns `t=,v0=`. */ +export function signWebhookPayload( + secret: string, + payload: string, + timestamp: number = Math.floor(Date.now() / 1000), +): { signature: string; timestamp: number } { + const sig = createHmac("sha256", secret) + .update(`${timestamp}.${payload}`) + .digest("hex"); + return { signature: `t=${timestamp},v0=${sig}`, timestamp }; +} + +/** Verify a webhook signature. */ +export function verifyWebhookPayload( + secret: string, + payload: string, + signature: string, + nowSec: number = Math.floor(Date.now() / 1000), +): boolean { + const m = signature.match(/^t=(\d+),v0=([0-9a-f]{64})$/); + if (!m) return false; + const ts = Number(m[1]); + if (Math.abs(nowSec - ts) > WEBHOOK_SIG_TTL_SECS) return false; + const expected = signWebhookPayload(secret, payload, ts).signature; + return signature === expected; +} + +/** Dead-letter entry for a failed webhook delivery. */ +export interface WebhookDeadLetter { + id: string; + webhookId: string; + url: string; + event: AppEvent; + attempts: number; + lastError: string; + deadAt: number; +} + +/** In-memory DLQ keyed by entry id. */ +export const webhookDlq = new Map(); + +/** In-flight delivery attempts so we can cap concurrency. */ +const inFlight = new Set(); + +/** Deliver a single event to every matching webhook. */ +export async function deliverEventToWebhooks( + event: AppEvent, + webhooks: Map, + httpClient: typeof fetch = fetch, +): Promise { + for (const [id, record] of webhooks) { + if (!record.events.includes(event.type)) continue; + if (!record.secret) continue; // safety: cannot sign without secret + attemptDelivery(id, record, event, 0, httpClient).catch(() => {}); + } +} + +async function attemptDelivery( + webhookId: string, + record: WebhookRecord, + event: AppEvent, + attempt: number, + httpClient: typeof fetch, +): Promise { + const payload = JSON.stringify({ + eventId: event.id, + type: event.type, + timestamp: event.ts, + payload: event.payload, + }); + + const { signature } = signWebhookPayload(record.secret!, payload); + + try { + const res = await httpClient(record.url, { + method: "POST", + headers: { + "Content-Type": "application/json", + "X-Signature": signature, + "User-Agent": "StableRoute-Webhook/1.0", + }, + body: payload, + }); + + if (res.ok) return; // 2xx = success + + const status = res.status; + if (status >= 500 || status === 408 || status === 429) { + // Retryable + await scheduleRetry(webhookId, record, event, attempt, `HTTP ${status}`, httpClient); + return; + } + // 4xx (except 408/429) = non-retryable, dead-letter immediately + deadLetter(webhookId, record, event, attempt + 1, `HTTP ${status}`); + } catch (err) { + // Network/timeout error = retryable + const msg = err instanceof Error ? err.message : String(err); + await scheduleRetry(webhookId, record, event, attempt, msg, httpClient); + } +} + +async function scheduleRetry( + webhookId: string, + record: WebhookRecord, + event: AppEvent, + attempt: number, + errorMsg: string, + httpClient: typeof fetch, +): Promise { + if (attempt >= WEBHOOK_MAX_RETRIES - 1) { + deadLetter(webhookId, record, event, attempt + 1, errorMsg); + return; + } + const delay = WEBHOOK_RETRY_BASE_MS * Math.pow(2, attempt); + await new Promise((r) => setTimeout(r, delay)); + return attemptDelivery(webhookId, record, event, attempt + 1, httpClient); +} + +function deadLetter( + webhookId: string, + record: WebhookRecord, + event: AppEvent, + attempts: number, + lastError: string, +): void { + const id = `${webhookId}::${event.id}::${Date.now()}`; + webhookDlq.set(id, { + id, + webhookId, + url: record.url, + event, + attempts, + lastError, + deadAt: Date.now(), + }); +}