diff --git a/apps/api/src/__tests__/WebhookDelivery.spec.ts b/apps/api/src/__tests__/WebhookDelivery.spec.ts new file mode 100644 index 0000000..7fc553f --- /dev/null +++ b/apps/api/src/__tests__/WebhookDelivery.spec.ts @@ -0,0 +1,831 @@ +import { + Event as EventType, + WebhookDelivery, + WebhookDeliveryAttemptResult, + WebhookEndpointRecord, +} from '@zoneless/shared-types'; +import { Database } from '../modules/Database'; +import { EventService } from '../modules/EventService'; +import { + WebhookDeliveryModule, + WEBHOOK_DELIVERY_LOCK_SECONDS, + WEBHOOK_RETRY_BACKOFF_SECONDS, +} from '../modules/WebhookDelivery'; +import { WebhookDeliveryWorker } from '../modules/WebhookDeliveryWorker'; +import { + WebhookDispatcher, + WebhookResponse, +} from '../modules/WebhookDispatcher'; +import { WebhookEndpointModule } from '../modules/WebhookEndpoint'; +import { DeterministicId, GetFixedTimestamp, ResetIdCounter } from './Setup'; + +jest.mock('../modules/WebhookDispatcher'); +jest.mock('../modules/WebhookEndpoint'); +jest.mock('../utils/IdGenerator', () => ({ + GenerateId: jest.fn((prefix: string) => DeterministicId(prefix)), +})); + +let now = GetFixedTimestamp(); +jest.mock('../utils/Timestamp', () => ({ Now: jest.fn(() => now) })); + +jest.mock('../modules/AppConfig', () => ({ + GetAppConfig: jest.fn(() => ({ + dashboardUrl: 'http://localhost:4200', + livemode: false, + appSecret: 'test-secret', + })), +})); + +const PLATFORM = 'acct_z_platform'; +const EVENT_ID = 'evt_z_test001'; +const ENDPOINT_ID = 'we_z_test001'; +const OTHER_PLATFORM = 'acct_z_other'; +const COLLECTION = 'WebhookDeliveries'; +const RETRY_BACKOFF_SECONDS = WEBHOOK_RETRY_BACKOFF_SECONDS; + +type Row = Record; + +function MatchesFilter(row: Row, filter: Row): boolean { + return Object.entries(filter).every(([field, condition]) => { + if (field === '$or') { + return (condition as Row[]).some((branch) => MatchesFilter(row, branch)); + } + + const value = row[field] ?? null; + + if (condition !== null && typeof condition === 'object') { + const operators = condition as Record; + if ('$in' in operators) + return (operators.$in as unknown[]).includes(value); + if ('$lte' in operators) { + return ( + value !== null && (value as number) <= (operators.$lte as number) + ); + } + if ('$gt' in operators) { + return value !== null && (value as number) > (operators.$gt as number); + } + if ('$exists' in operators) return (value !== null) === operators.$exists; + } + + return value === condition; + }); +} + +function ApplyUpdate(row: Row, data: Row): Row { + const { $set, $inc, $push, ...plain } = data; + const next = { ...row, ...plain, ...($set as Row | undefined) }; + const increments = ($inc as Record | undefined) ?? {}; + const pushes = ($push as Record | undefined) ?? {}; + + for (const [field, amount] of Object.entries(increments)) { + next[field] = ((next[field] as number) ?? 0) + amount; + } + + for (const [field, value] of Object.entries(pushes)) { + next[field] = [...((next[field] as unknown[]) ?? []), value]; + } + + return next; +} + +class DeliveryStore extends Database { + readonly rows = new Map(); + + override async Set( + collection: string, + documentId: string, + data: Partial + ): Promise { + const key = `${collection}:${documentId}`; + const next = { ...(this.rows.get(key) ?? {}), ...(data as Row) }; + this.rows.set(key, next); + return next as T; + } + + override async Get( + collection: string, + documentId: string + ): Promise { + return ( + (this.rows.get(`${collection}:${documentId}`) as T | undefined) ?? null + ); + } + + override async Update( + collection: string, + documentId: string, + data: Partial + ): Promise { + const key = `${collection}:${documentId}`; + const row = this.rows.get(key); + if (!row) return null; + + const next = ApplyUpdate(row, data as Row); + this.rows.set(key, next); + return next as T; + } + + override async FindOneAndUpdateByFilter( + collection: string, + filter: Record, + data: Record + ): Promise { + for (const [key, row] of this.rows) { + if (!key.startsWith(`${collection}:`) || !MatchesFilter(row, filter)) { + continue; + } + + const next = ApplyUpdate(row, data); + this.rows.set(key, next); + return next as T; + } + + return null; + } + + override async Find( + collection: string, + field: string, + value: unknown + ): Promise { + return [...this.rows] + .filter( + ([key, row]) => + key.startsWith(`${collection}:`) && (row[field] ?? null) === value + ) + .map(([, row]) => row as T); + } + + override async Increment( + collection: string, + documentId: string, + field: string, + amount = 1 + ): Promise { + const key = `${collection}:${documentId}`; + const row = this.rows.get(key); + if (!row) return; + + this.rows.set(key, { + ...row, + [field]: ((row[field] as number) ?? 0) + amount, + }); + } + + override async Query(): Promise { + return []; + } +} + +function MakeEvent(overrides: Partial = {}): EventType { + return { + id: EVENT_ID, + object: 'event', + account: PLATFORM, + api_version: null, + context: null, + created: now, + livemode: false, + data: { + object: { id: 'prod_z_test001', object: 'product' }, + previous_attributes: null, + }, + pending_webhooks: 1, + request: null, + type: 'product.created', + platform_account: PLATFORM, + ...overrides, + }; +} + +function MakeEndpoint( + overrides: Partial = {} +): WebhookEndpointRecord { + return { + id: ENDPOINT_ID, + object: 'webhook_endpoint', + account: PLATFORM, + platform_account: PLATFORM, + api_version: null, + application: null, + created: now, + description: null, + enabled_events: ['*'], + livemode: false, + metadata: {}, + secret: 'whsec_z_testsecret', + status: 'enabled', + url: 'https://example.com/webhooks', + ...overrides, + }; +} + +function MakeResult( + result: WebhookDeliveryAttemptResult, + statusCode: number | null, + error: string | null = null +): Awaited> { + return { result, statusCode, error, durationMs: 3206 }; +} + +const Flush = (): Promise => + new Promise((resolve) => setTimeout(resolve, 0)); + +describe('WebhookDelivery', () => { + let store: DeliveryStore; + let send: jest.MockedFunction; + let endpointModule: jest.Mocked; + let deliveries: WebhookDeliveryModule; + let worker: WebhookDeliveryWorker; + let event: EventType; + let endpoint: WebhookEndpointRecord; + + const ReadDelivery = (id: string): WebhookDelivery => + store.rows.get(`${COLLECTION}:${id}`) as unknown as WebhookDelivery; + + const ReadDeliveries = (): WebhookDelivery[] => + [...store.rows] + .filter(([key]) => key.startsWith(`${COLLECTION}:`)) + .map(([, row]) => row as unknown as WebhookDelivery); + + const ReadEvent = (id: string = EVENT_ID): EventType => + store.rows.get(`Events:${id}`) as unknown as EventType; + + const SeedDelivery = async (): Promise => { + const [delivery] = await deliveries.CreateDeliveriesForEvent(event, [ + endpoint, + ]); + return delivery; + }; + + beforeEach(async () => { + jest.clearAllMocks(); + ResetIdCounter(); + now = GetFixedTimestamp(); + + store = new DeliveryStore(); + deliveries = new WebhookDeliveryModule(store); + worker = new WebhookDeliveryWorker(store); + + send = jest.mocked(WebhookDispatcher.prototype.Send); + send.mockResolvedValue(MakeResult('succeeded', 200)); + + event = MakeEvent(); + endpoint = MakeEndpoint(); + await store.Set('Events', EVENT_ID, event); + + endpointModule = jest.mocked(WebhookEndpointModule.prototype); + endpointModule.GetWebhookEndpoint.mockResolvedValue(endpoint); + endpointModule.GetWebhookEndpointsForEvent.mockResolvedValue([endpoint]); + }); + + describe('EventService', () => { + it('should persist a delivery for every subscribed endpoint before the first request', async () => { + const second = MakeEndpoint({ + id: 'we_z_test002', + url: 'https://example.com/second', + }); + endpointModule.GetWebhookEndpointsForEvent.mockResolvedValue([ + endpoint, + second, + ]); + + let storedAtFirstRequest: WebhookDelivery[] = []; + send.mockImplementation(async () => { + storedAtFirstRequest = ReadDeliveries(); + return MakeResult('succeeded', 200); + }); + + await new EventService(store).Emit( + 'product.created', + PLATFORM, + event.data.object + ); + await Flush(); + + expect(send).toHaveBeenCalledTimes(2); + expect(storedAtFirstRequest).toHaveLength(2); + expect( + storedAtFirstRequest + .map((delivery) => delivery.webhook_endpoint_id) + .sort() + ).toEqual([ENDPOINT_ID, second.id].sort()); + + for (const delivery of storedAtFirstRequest) { + expect(delivery.id).toMatch(/^whd_z/); + expect(delivery.event_id).toBe(EVENT_ID); + expect(delivery.status).toBe('pending'); + expect(delivery.next_attempt_at).toBe(now); + expect(delivery.attempts).toEqual([]); + expect(delivery.claim_until).toBeGreaterThan(now); + expect(delivery.claim_token).toMatch(/^claim_z/); + } + + expect(ReadEvent().pending_webhooks).toBe(0); + }); + }); + + describe('WebhookDeliveryWorker', () => { + it('should mark the delivery succeeded on a successful first attempt', async () => { + const delivery = await SeedDelivery(); + + expect(await worker.ProcessBatch()).toEqual({ + processed: 1, + succeeded: 1, + retrying: 0, + failed: 0, + }); + + const stored = ReadDelivery(delivery.id); + expect(stored.status).toBe('succeeded'); + expect(stored.delivered_at).toBe(now); + expect(stored.next_attempt_at).toBeNull(); + expect(stored.claim_until).toBeNull(); + expect(stored.claim_token).toBeNull(); + expect(stored.attempts).toEqual([ + { + attempt_number: 1, + attempted_at: now, + completed_at: now, + result: 'succeeded', + http_status: 200, + duration_ms: 3206, + error: null, + url: endpoint.url, + }, + ]); + expect(ReadEvent().pending_webhooks).toBe(0); + }); + + it('should record a failed attempt and schedule the next retry', async () => { + send.mockResolvedValue(MakeResult('http_error', 500, 'HTTP 500')); + const delivery = await SeedDelivery(); + + const batch = await worker.ProcessBatch(); + + expect(batch).toEqual({ + processed: 1, + succeeded: 0, + retrying: 1, + failed: 0, + }); + expect(send).toHaveBeenCalledWith( + expect.objectContaining({ id: EVENT_ID }), + endpoint.url, + endpoint.secret + ); + + const stored = ReadDelivery(delivery.id); + expect(stored.status).toBe('retrying'); + expect(stored.delivered_at).toBeNull(); + expect(stored.next_attempt_at).toBe(now + RETRY_BACKOFF_SECONDS[0]); + expect(stored.claim_until).toBeNull(); + expect(stored.claim_token).toBeNull(); + expect(stored.attempts).toEqual([ + { + attempt_number: 1, + attempted_at: now, + completed_at: now, + result: 'http_error', + http_status: 500, + duration_ms: 3206, + error: 'HTTP 500', + url: endpoint.url, + }, + ]); + expect(ReadEvent().pending_webhooks).toBe(1); + }); + + it('should follow the retry backoff and give up after the last attempt', async () => { + send.mockResolvedValue(MakeResult('network_error', null, 'fetch failed')); + const delivery = await SeedDelivery(); + + for (const [index, backoff] of RETRY_BACKOFF_SECONDS.entries()) { + if (index > 0) now += RETRY_BACKOFF_SECONDS[index - 1]; + + const batch = await worker.ProcessBatch(); + + expect(batch).toEqual({ + processed: 1, + succeeded: 0, + retrying: 1, + failed: 0, + }); + expect(ReadDelivery(delivery.id).next_attempt_at).toBe(now + backoff); + } + + now += RETRY_BACKOFF_SECONDS[RETRY_BACKOFF_SECONDS.length - 1]; + const finalBatch = await worker.ProcessBatch(); + + expect(finalBatch).toEqual({ + processed: 1, + succeeded: 0, + retrying: 0, + failed: 1, + }); + + const stored = ReadDelivery(delivery.id); + expect(stored.status).toBe('failed'); + expect(stored.next_attempt_at).toBeNull(); + expect(stored.delivered_at).toBeNull(); + expect(stored.attempts.map((attempt) => attempt.attempt_number)).toEqual([ + 1, 2, 3, 4, + ]); + expect(ReadEvent().pending_webhooks).toBe(1); + + now += 86400; + expect(await worker.ProcessBatch()).toEqual({ + processed: 0, + succeeded: 0, + retrying: 0, + failed: 0, + }); + expect(send).toHaveBeenCalledTimes(4); + }); + + it('should stop retrying once a later attempt succeeds', async () => { + send.mockResolvedValue( + MakeResult( + 'timed_out', + null, + 'The operation was aborted due to timeout' + ) + ); + const delivery = await SeedDelivery(); + + await worker.ProcessBatch(); + now += RETRY_BACKOFF_SECONDS[0]; + send.mockResolvedValue(MakeResult('succeeded', 200)); + + const batch = await worker.ProcessBatch(); + + expect(batch).toEqual({ + processed: 1, + succeeded: 1, + retrying: 0, + failed: 0, + }); + + const stored = ReadDelivery(delivery.id); + expect(stored.status).toBe('succeeded'); + expect(stored.delivered_at).toBe(now); + expect(stored.attempts.map((attempt) => attempt.result)).toEqual([ + 'timed_out', + 'succeeded', + ]); + expect(send.mock.calls.map((call) => call[0].id)).toEqual([ + EVENT_ID, + EVENT_ID, + ]); + + now += 86400; + expect(await worker.ProcessBatch()).toEqual({ + processed: 0, + succeeded: 0, + retrying: 0, + failed: 0, + }); + expect(send).toHaveBeenCalledTimes(2); + }); + + it('should track multiple endpoints for the same event independently', async () => { + const failing = MakeEndpoint({ + id: 'we_z_test002', + url: 'https://example.com/failing', + }); + let failingAttempts = 0; + send.mockImplementation(async (_event, url) => { + if (url !== failing.url) return MakeResult('succeeded', 200); + + failingAttempts += 1; + return failingAttempts === 1 + ? MakeResult('http_error', 500, 'HTTP 500') + : MakeResult('succeeded', 200); + }); + endpointModule.GetWebhookEndpoint.mockImplementation(async (id) => + id === failing.id ? failing : endpoint + ); + await store.Set('Events', EVENT_ID, { + ...ReadEvent(), + pending_webhooks: 2, + }); + + const [succeeding, retrying] = await deliveries.CreateDeliveriesForEvent( + event, + [endpoint, failing] + ); + + expect(await worker.ProcessBatch()).toEqual({ + processed: 2, + succeeded: 1, + retrying: 1, + failed: 0, + }); + expect(ReadDelivery(succeeding.id).status).toBe('succeeded'); + expect(ReadDelivery(retrying.id).status).toBe('retrying'); + expect(ReadEvent().pending_webhooks).toBe(1); + + now += RETRY_BACKOFF_SECONDS[0]; + expect(await worker.ProcessBatch()).toEqual({ + processed: 1, + succeeded: 1, + retrying: 0, + failed: 0, + }); + expect(ReadDelivery(retrying.id).status).toBe('succeeded'); + expect(ReadEvent().pending_webhooks).toBe(0); + }); + + it('should only process deliveries for the requested platform', async () => { + const otherEndpoint = MakeEndpoint({ + id: 'we_z_test002', + account: OTHER_PLATFORM, + platform_account: OTHER_PLATFORM, + url: 'https://example.com/other', + }); + const otherEvent = MakeEvent({ + id: 'evt_z_test002', + account: OTHER_PLATFORM, + platform_account: OTHER_PLATFORM, + }); + await store.Set('Events', otherEvent.id, otherEvent); + endpointModule.GetWebhookEndpoint.mockImplementation(async (id) => + id === otherEndpoint.id ? otherEndpoint : endpoint + ); + + const [mine] = await deliveries.CreateDeliveriesForEvent(event, [ + endpoint, + ]); + const [other] = await deliveries.CreateDeliveriesForEvent(otherEvent, [ + otherEndpoint, + ]); + + expect( + await worker.ProcessBatch({ platformAccountId: PLATFORM }) + ).toEqual({ + processed: 1, + succeeded: 1, + retrying: 0, + failed: 0, + }); + expect(ReadDelivery(mine.id).status).toBe('succeeded'); + expect(ReadDelivery(other.id).status).toBe('pending'); + expect(ReadEvent(otherEvent.id).pending_webhooks).toBe(1); + + expect(await worker.ProcessBatch()).toEqual({ + processed: 1, + succeeded: 1, + retrying: 0, + failed: 0, + }); + expect(ReadDelivery(other.id).status).toBe('succeeded'); + expect(ReadEvent(otherEvent.id).pending_webhooks).toBe(0); + }); + + it('should claim one delivery at a time so a claim never expires mid-batch', async () => { + const endpoints = [ + endpoint, + MakeEndpoint({ id: 'we_z_test002', url: 'https://example.com/second' }), + MakeEndpoint({ id: 'we_z_test003', url: 'https://example.com/third' }), + ]; + endpointModule.GetWebhookEndpoint.mockImplementation( + async (id) => endpoints.find((candidate) => candidate.id === id) ?? null + ); + await store.Set('Events', EVENT_ID, { + ...ReadEvent(), + pending_webhooks: endpoints.length, + }); + await deliveries.CreateDeliveriesForEvent(event, endpoints); + + const claimedAtSend: number[] = []; + send.mockImplementation(async () => { + claimedAtSend.push( + ReadDeliveries().filter((delivery) => delivery.claim_until !== null) + .length + ); + return MakeResult('succeeded', 200); + }); + + expect(await worker.ProcessBatch()).toEqual({ + processed: 3, + succeeded: 3, + retrying: 0, + failed: 0, + }); + + expect(claimedAtSend).toEqual([1, 1, 1]); + }); + + it('should not let a stale worker resettle a delivery another worker reclaimed', async () => { + const delivery = await SeedDelivery(); + const stale = await deliveries.ClaimById(delivery.id); + if (!stale) throw new Error('delivery was not claimed'); + const staleAt = now; + const attempt = { + attempt_number: 1, + attempted_at: now, + completed_at: now, + result: 'http_error' as const, + http_status: 500, + duration_ms: 10, + error: 'HTTP 500', + url: endpoint.url, + }; + + now += WEBHOOK_DELIVERY_LOCK_SECONDS + 1; + const owner = await deliveries.ClaimById(delivery.id); + if (!owner) throw new Error('delivery was not reclaimed'); + + expect(owner.claim_token).not.toBe(stale.claim_token); + + expect(await deliveries.RecordFailure(stale, attempt)).toBeNull(); + expect(await deliveries.RecordSuccess(stale, attempt)).toBe(false); + expect(await deliveries.MarkFailed(delivery.id, stale.claim_token)).toBe( + false + ); + + const stored = ReadDelivery(delivery.id); + expect(stored.status).toBe('pending'); + expect(stored.next_attempt_at).toBe(staleAt); + expect(stored.claim_token).toBe(owner.claim_token); + expect(stored.claim_until).toBe(now + WEBHOOK_DELIVERY_LOCK_SECONDS); + expect(stored.attempts).toEqual([]); + expect(ReadEvent().pending_webhooks).toBe(1); + + await store.Update(COLLECTION, delivery.id, { + status: 'succeeded', + delivered_at: now, + next_attempt_at: null, + claim_until: null, + claim_token: null, + }); + + expect(await deliveries.RecordFailure(stale, attempt)).toBeNull(); + expect(await deliveries.RecordSuccess(stale, attempt)).toBe(false); + + expect(ReadDelivery(delivery.id).status).toBe('succeeded'); + expect(ReadDelivery(delivery.id).attempts).toEqual([]); + expect(ReadEvent().pending_webhooks).toBe(1); + }); + + it('should not decrement pending_webhooks twice for one delivery', async () => { + const delivery = await SeedDelivery(); + + await worker.ProcessBatch(); + expect(ReadEvent().pending_webhooks).toBe(0); + + await store.Update(COLLECTION, delivery.id, { + status: 'retrying', + next_attempt_at: now, + claim_until: null, + claim_token: null, + }); + + await worker.ProcessBatch(); + + expect(ReadDelivery(delivery.id).status).toBe('succeeded'); + expect(ReadEvent().pending_webhooks).toBe(0); + }); + + it('should let only one of two workers claim a delivery', async () => { + await SeedDelivery(); + const otherWorker = new WebhookDeliveryWorker(store); + + const [first, second] = await Promise.all([ + worker.ProcessBatch(), + otherWorker.ProcessBatch(), + ]); + + expect([first.processed, second.processed].sort()).toEqual([0, 1]); + expect(send).toHaveBeenCalledTimes(1); + }); + + it('should not attempt a delivery whose next attempt is not due', async () => { + const delivery = await SeedDelivery(); + await store.Update(COLLECTION, delivery.id, { + next_attempt_at: now + 30, + }); + + expect(await worker.ProcessBatch()).toEqual({ + processed: 0, + succeeded: 0, + retrying: 0, + failed: 0, + }); + expect(send).not.toHaveBeenCalled(); + }); + + const endpointStates: Array<[string, WebhookEndpointRecord | null]> = [ + ['disabled', MakeEndpoint({ status: 'disabled' })], + ['deleted', null], + ]; + + it.each(endpointStates)( + 'should mark the delivery failed when the endpoint is %s', + async (_state, endpointRecord) => { + endpointModule.GetWebhookEndpoint.mockResolvedValue(endpointRecord); + const delivery = await SeedDelivery(); + + expect(await worker.ProcessBatch()).toEqual({ + processed: 1, + succeeded: 0, + retrying: 0, + failed: 1, + }); + expect(send).not.toHaveBeenCalled(); + + const stored = ReadDelivery(delivery.id); + expect(stored.status).toBe('failed'); + expect(stored.next_attempt_at).toBeNull(); + expect(stored.attempts).toEqual([]); + } + ); + }); + + describe('WebhookDispatcher', () => { + const RealWebhookDispatcher = jest.requireActual< + typeof import('../modules/WebhookDispatcher') + >('../modules/WebhookDispatcher').WebhookDispatcher; + const originalFetch = globalThis.fetch; + let fetchMock: jest.Mock; + + const outcomes: Array<[string, unknown, Partial]> = [ + [ + 'a 2xx response', + { ok: true, status: 200 }, + { result: 'succeeded', statusCode: 200, error: null }, + ], + [ + 'a non-2xx response', + { ok: false, status: 500 }, + { result: 'http_error', statusCode: 500, error: 'HTTP 500' }, + ], + [ + 'a timeout', + Object.assign(new Error('The operation was aborted due to timeout'), { + name: 'TimeoutError', + }), + { + result: 'timed_out', + statusCode: null, + error: 'The operation was aborted due to timeout', + }, + ], + [ + 'a request that never reaches the endpoint', + new TypeError('fetch failed'), + { result: 'network_error', statusCode: null, error: 'fetch failed' }, + ], + ]; + + beforeEach(() => { + fetchMock = jest.fn(); + globalThis.fetch = fetchMock as unknown as typeof fetch; + }); + + afterAll(() => { + globalThis.fetch = originalFetch; + }); + + it.each(outcomes)( + 'should report the attempt result for %s', + async (_case, response, expected) => { + if (response instanceof Error) { + fetchMock.mockRejectedValue(response); + } else { + fetchMock.mockResolvedValue(response); + } + + const result = await new RealWebhookDispatcher().Send( + event, + endpoint.url, + endpoint.secret + ); + + expect(result).toMatchObject(expected); + expect(result.durationMs).toBeGreaterThanOrEqual(0); + } + ); + + it('should sign every attempt with a fresh timestamp', async () => { + fetchMock.mockResolvedValue({ ok: true, status: 200 }); + const dispatcher = new RealWebhookDispatcher(); + + await dispatcher.Send(event, endpoint.url, endpoint.secret); + now += 5; + await dispatcher.Send(event, endpoint.url, endpoint.secret); + + const signatures = fetchMock.mock.calls.map( + (call) => + (call[1].headers as Record)['Zoneless-Signature'] + ); + expect(signatures[0]).not.toEqual(signatures[1]); + expect( + fetchMock.mock.calls.map( + (call) => JSON.parse(call[1].body as string).id as string + ) + ).toEqual([EVENT_ID, EVENT_ID]); + }); + }); +}); diff --git a/apps/api/src/modules/Database.ts b/apps/api/src/modules/Database.ts index a179124..7880426 100644 --- a/apps/api/src/modules/Database.ts +++ b/apps/api/src/modules/Database.ts @@ -457,6 +457,7 @@ export class Database { 'UsageCounters', 'TelemetryConfigs', 'VerificationSessions', + 'WebhookDeliveries', 'WebhookEndpoints', ]; @@ -488,6 +489,10 @@ export class Database { flexibleSchema.index({ id: 1 }, { unique: true, sparse: true }); + if (collectionName === 'WebhookDeliveries') { + flexibleSchema.index({ status: 1, next_attempt_at: 1 }); + } + return mongoose.model(collectionName, flexibleSchema); } diff --git a/apps/api/src/modules/EventService.ts b/apps/api/src/modules/EventService.ts index 2d1e9d3..6f32099 100644 --- a/apps/api/src/modules/EventService.ts +++ b/apps/api/src/modules/EventService.ts @@ -13,12 +13,19 @@ * @module EventService */ -import { Event, EventDataObject, EventType } from '@zoneless/shared-types'; +import { + Event, + EventDataObject, + EventType, + WebhookDelivery, + WebhookEndpointRecord, +} from '@zoneless/shared-types'; import { Database } from './Database'; import { EventModule } from './Event'; import { AccountModule } from './Account'; import { WebhookEndpointModule } from './WebhookEndpoint'; -import { WebhookDispatcher } from './WebhookDispatcher'; +import { WebhookDeliveryModule } from './WebhookDelivery'; +import { WebhookDeliveryWorker } from './WebhookDeliveryWorker'; import { GetPlatformAccountId } from './PlatformAccess'; import { GetRequestContext } from '../middleware/RequestContext'; import { Logger } from '../utils/Logger'; @@ -31,18 +38,18 @@ interface EventOptions { } export class EventService { - private readonly db: Database; private readonly eventModule: EventModule; private readonly accountModule: AccountModule; private readonly webhookEndpointModule: WebhookEndpointModule; - private readonly webhookDispatcher: WebhookDispatcher; + private readonly webhookDeliveryModule: WebhookDeliveryModule; + private readonly webhookDeliveryWorker: WebhookDeliveryWorker; constructor(db: Database) { - this.db = db; this.eventModule = new EventModule(db); this.accountModule = new AccountModule(db); this.webhookEndpointModule = new WebhookEndpointModule(db); - this.webhookDispatcher = new WebhookDispatcher(); + this.webhookDeliveryModule = new WebhookDeliveryModule(db); + this.webhookDeliveryWorker = new WebhookDeliveryWorker(db); } /** @@ -111,8 +118,10 @@ export class EventService { pendingWebhooks: pendingWebhooksCount, }); + const deliveries = await this.PersistDeliveries(event, endpoints); + // Dispatch webhooks asynchronously (don't await - fire and forget) - this.DispatchWebhooks(event, platformAccountId).catch((error) => { + this.DeliverFirstAttempts(event, deliveries).catch((error) => { Logger.error('Failed to dispatch webhooks', error, { eventId: event.id, eventType: type, @@ -152,85 +161,74 @@ export class EventService { } /** - * Dispatches webhook for an event to all subscribed webhook endpoints. + * Persists a webhook delivery for the event to each subscribed endpoint. * * @param event - The event to dispatch - * @param platformAccountId - The platform to send webhooks to + * @param endpoints - The endpoints subscribed to the event type + * @returns The created deliveries, empty when there is nothing to deliver */ - private async DispatchWebhooks( + private async PersistDeliveries( event: Event, - platformAccountId: string - ): Promise { - // Get all webhook endpoints that subscribe to this event type - const endpoints = - await this.webhookEndpointModule.GetWebhookEndpointsForEvent( - platformAccountId, - event.type - ); - + endpoints: WebhookEndpointRecord[] + ): Promise { if (endpoints.length === 0) { Logger.debug('No webhook endpoints configured for event type', { eventId: event.id, eventType: event.type, - platformAccountId, + platformAccountId: event.platform_account, }); - return; + return []; } Logger.debug('Dispatching webhooks', { eventId: event.id, eventType: event.type, - platformAccountId, + platformAccountId: event.platform_account, endpointCount: endpoints.length, }); + try { + return await this.webhookDeliveryModule.CreateDeliveriesForEvent( + event, + endpoints + ); + } catch (error) { + Logger.error('Failed to persist webhook deliveries', error, { + eventId: event.id, + eventType: event.type, + }); + return []; + } + } + + /** + * Sends the first delivery attempt for each persisted delivery. + * + * @param event - The event being dispatched + * @param deliveries - The deliveries to attempt + */ + private async DeliverFirstAttempts( + event: Event, + deliveries: WebhookDelivery[] + ): Promise { // Dispatch to all endpoints in parallel const results = await Promise.allSettled( - endpoints.map(async (endpoint) => { - try { - const result = await this.webhookDispatcher.Send( - event, - endpoint.url, - endpoint.secret - ); - - if (!result.success) { - Logger.warn('Webhook delivery failed', { - eventId: event.id, - eventType: event.type, - webhookEndpointId: endpoint.id, - url: endpoint.url, - error: result.error, - statusCode: result.statusCode, - }); - } - - return result; - } catch (error) { - Logger.error('Webhook dispatch error', error, { - eventId: event.id, - eventType: event.type, - webhookEndpointId: endpoint.id, - url: endpoint.url, - }); - throw error; - } - }) + deliveries.map((delivery) => + this.webhookDeliveryWorker.ProcessDelivery(delivery.id) + ) ); // Log summary const successful = results.filter( - (r) => - r.status === 'fulfilled' && (r.value as { success: boolean }).success + (result) => result.status === 'fulfilled' && result.value === 'succeeded' ).length; - const failed = results.length - successful; Logger.info('Webhook dispatch completed', { eventId: event.id, eventType: event.type, - platformAccountId, + platformAccountId: event.platform_account, successful, - failed, + failed: results.length - successful, total: results.length, }); } diff --git a/apps/api/src/modules/WebhookDelivery.ts b/apps/api/src/modules/WebhookDelivery.ts new file mode 100644 index 0000000..b43b6df --- /dev/null +++ b/apps/api/src/modules/WebhookDelivery.ts @@ -0,0 +1,210 @@ +import { + Event, + WebhookDelivery, + WebhookDeliveryAttempt, + WebhookDeliveryStatus, + WebhookEndpointRecord, +} from '@zoneless/shared-types'; +import { Database } from './Database'; +import { GenerateId } from '../utils/IdGenerator'; +import { Now } from '../utils/Timestamp'; +import { WEBHOOK_REQUEST_TIMEOUT_SECONDS } from './WebhookDispatcher'; + +/** Retry delays after the first attempt, in seconds. */ +export const WEBHOOK_RETRY_BACKOFF_SECONDS = [5 * 60, 60 * 60, 6 * 60 * 60]; + +export const WEBHOOK_DELIVERY_LOCK_SECONDS = + WEBHOOK_REQUEST_TIMEOUT_SECONDS + 30; + +export class WebhookDeliveryModule { + private readonly db: Database; + + constructor(db: Database) { + this.db = db; + } + + async CreateDeliveriesForEvent( + event: Event, + endpoints: WebhookEndpointRecord[] + ): Promise { + const now = Now(); + const deliveries = endpoints.map((endpoint) => + this.DeliveryObject(event, endpoint, now) + ); + + await Promise.all( + deliveries.map((delivery) => + this.db.Set('WebhookDeliveries', delivery.id, delivery) + ) + ); + + return deliveries; + } + + async ClaimNext(platformAccountId?: string): Promise { + return this.Claim(this.ClaimFilter({ platformAccountId })); + } + + async ClaimById(deliveryId: string): Promise { + return this.Claim(this.ClaimFilter({ deliveryId })); + } + + /** False when another worker holds the claim, so the attempt was not recorded. */ + async RecordSuccess( + delivery: WebhookDelivery, + attempt: WebhookDeliveryAttempt + ): Promise { + const settled = await this.WriteAttempt( + delivery.id, + delivery.claim_token, + attempt, + { + status: 'succeeded', + delivered_at: attempt.completed_at, + next_attempt_at: null, + claim_until: null, + claim_token: null, + } + ); + + if (settled) { + await this.DecrementPendingWebhooks(delivery.event_id); + } + + return settled; + } + + /** Null when another worker holds the claim, so the attempt was not recorded. */ + async RecordFailure( + delivery: WebhookDelivery, + attempt: WebhookDeliveryAttempt + ): Promise { + const backoff = WEBHOOK_RETRY_BACKOFF_SECONDS[attempt.attempt_number - 1]; + const nextAttemptAt = + backoff === undefined ? null : attempt.completed_at + backoff; + const status: WebhookDeliveryStatus = + nextAttemptAt === null ? 'failed' : 'retrying'; + + const settled = await this.WriteAttempt( + delivery.id, + delivery.claim_token, + attempt, + { + status, + next_attempt_at: nextAttemptAt, + claim_until: null, + claim_token: null, + } + ); + + return settled ? status : null; + } + + async MarkFailed( + deliveryId: string, + claimToken: string | null + ): Promise { + const failed = await this.db.FindOneAndUpdateByFilter( + 'WebhookDeliveries', + { + id: deliveryId, + status: { $in: ['pending', 'retrying'] }, + claim_token: claimToken, + }, + { + $set: { + status: 'failed', + next_attempt_at: null, + claim_until: null, + claim_token: null, + }, + } + ); + + return failed !== null; + } + + private DeliveryObject( + event: Event, + endpoint: WebhookEndpointRecord, + now: number + ): WebhookDelivery { + return { + id: GenerateId('whd_z'), + object: 'webhook_delivery', + event_id: event.id, + webhook_endpoint_id: endpoint.id, + platform_account: event.platform_account, + status: 'pending', + next_attempt_at: now, + delivered_at: null, + claim_until: null, + claim_token: null, + attempts: [], + }; + } + + private async Claim( + filter: Record + ): Promise { + return this.db.FindOneAndUpdateByFilter( + 'WebhookDeliveries', + filter, + { + $set: { + claim_until: Now() + WEBHOOK_DELIVERY_LOCK_SECONDS, + claim_token: GenerateId('claim_z'), + }, + } + ); + } + + private ClaimFilter(scope: { + deliveryId?: string; + platformAccountId?: string; + }): Record { + const now = Now(); + + return { + ...(scope.deliveryId ? { id: scope.deliveryId } : {}), + ...(scope.platformAccountId + ? { platform_account: scope.platformAccountId } + : {}), + status: { $in: ['pending', 'retrying'] }, + next_attempt_at: { $lte: now }, + $or: [ + { claim_until: null }, + { claim_until: { $exists: false } }, + { claim_until: { $lte: now } }, + ], + }; + } + + /** Appends the attempt and moves the delivery in one write guarded by the claim token. */ + private async WriteAttempt( + deliveryId: string, + claimToken: string | null, + attempt: WebhookDeliveryAttempt, + state: Record + ): Promise { + const settled = await this.db.FindOneAndUpdateByFilter( + 'WebhookDeliveries', + { + id: deliveryId, + status: { $in: ['pending', 'retrying'] }, + claim_token: claimToken, + }, + { $set: state, $push: { attempts: attempt } } + ); + + return settled !== null; + } + + private async DecrementPendingWebhooks(eventId: string): Promise { + await this.db.FindOneAndUpdateByFilter( + 'Events', + { id: eventId, pending_webhooks: { $gt: 0 } }, + { $inc: { pending_webhooks: -1 } } + ); + } +} diff --git a/apps/api/src/modules/WebhookDeliveryWorker.ts b/apps/api/src/modules/WebhookDeliveryWorker.ts new file mode 100644 index 0000000..9cead9c --- /dev/null +++ b/apps/api/src/modules/WebhookDeliveryWorker.ts @@ -0,0 +1,132 @@ +import { + WebhookDelivery, + WebhookDeliveryAttempt, + WebhookDeliveryBatch, +} from '@zoneless/shared-types'; +import { Database } from './Database'; +import { EventModule } from './Event'; +import { WebhookDeliveryModule } from './WebhookDelivery'; +import { WebhookEndpointModule } from './WebhookEndpoint'; +import { WebhookDispatcher } from './WebhookDispatcher'; +import { Now } from '../utils/Timestamp'; +import { Logger } from '../utils/Logger'; + +const DEFAULT_BATCH_SIZE = 20; +const MAX_BATCH_SIZE = 100; + +export type WebhookDeliveryOutcome = 'succeeded' | 'retrying' | 'failed'; + +export type WebhookDeliveryBatchResult = Omit; + +export class WebhookDeliveryWorker { + private readonly deliveryModule: WebhookDeliveryModule; + private readonly eventModule: EventModule; + private readonly webhookEndpointModule: WebhookEndpointModule; + private readonly webhookDispatcher: WebhookDispatcher; + + constructor(db: Database) { + this.deliveryModule = new WebhookDeliveryModule(db); + this.eventModule = new EventModule(db); + this.webhookEndpointModule = new WebhookEndpointModule(db); + this.webhookDispatcher = new WebhookDispatcher(); + } + + async ProcessBatch( + options: { limit?: number; platformAccountId?: string } = {} + ): Promise { + const limit = Math.min( + Math.max(options.limit ?? DEFAULT_BATCH_SIZE, 1), + MAX_BATCH_SIZE + ); + const result: WebhookDeliveryBatchResult = { + processed: 0, + succeeded: 0, + retrying: 0, + failed: 0, + }; + + for (let index = 0; index < limit; index++) { + const delivery = await this.deliveryModule.ClaimNext( + options.platformAccountId + ); + if (!delivery) break; + + const outcome = await this.Attempt(delivery); + result.processed += 1; + + if (outcome === 'succeeded') { + result.succeeded += 1; + } else if (outcome === 'retrying') { + result.retrying += 1; + } else if (outcome === 'failed') { + result.failed += 1; + } + } + + return result; + } + + async ProcessDelivery( + deliveryId: string + ): Promise { + const delivery = await this.deliveryModule.ClaimById(deliveryId); + if (!delivery) return null; + + return this.Attempt(delivery); + } + + /** Null when another worker holds the claim, so this attempt settled nothing. */ + private async Attempt( + delivery: WebhookDelivery + ): Promise { + const event = await this.eventModule.GetEvent(delivery.event_id); + const endpoint = await this.webhookEndpointModule.GetWebhookEndpoint( + delivery.webhook_endpoint_id + ); + + if (!event || !endpoint || endpoint.status !== 'enabled') { + Logger.warn('Webhook delivery abandoned: endpoint is gone or disabled', { + deliveryId: delivery.id, + eventId: delivery.event_id, + webhookEndpointId: delivery.webhook_endpoint_id, + }); + + const settled = await this.deliveryModule.MarkFailed( + delivery.id, + delivery.claim_token + ); + return settled ? 'failed' : null; + } + + const attemptedAt = Now(); + const response = await this.webhookDispatcher.Send( + event, + endpoint.url, + endpoint.secret + ); + + const attempt: WebhookDeliveryAttempt = { + attempt_number: delivery.attempts.length + 1, + attempted_at: attemptedAt, + completed_at: Now(), + result: response.result, + http_status: response.statusCode, + duration_ms: response.durationMs, + error: response.error, + url: endpoint.url, + }; + + if (response.result === 'succeeded') { + const settled = await this.deliveryModule.RecordSuccess( + delivery, + attempt + ); + return settled ? 'succeeded' : null; + } + + const status = await this.deliveryModule.RecordFailure(delivery, attempt); + if (status === null) return null; + + return status === 'failed' ? 'failed' : 'retrying'; + } +} diff --git a/apps/api/src/modules/WebhookDispatcher.ts b/apps/api/src/modules/WebhookDispatcher.ts index 43cb29a..a000a30 100644 --- a/apps/api/src/modules/WebhookDispatcher.ts +++ b/apps/api/src/modules/WebhookDispatcher.ts @@ -5,19 +5,25 @@ * @module WebhookDispatcher */ -import { Event as EventType } from '@zoneless/shared-types'; +import { + Event as EventType, + WebhookDeliveryAttemptResult, +} from '@zoneless/shared-types'; import { ComputeSignature } from '../utils/Signature'; import { Now } from '../utils/Timestamp'; import { Logger } from '../utils/Logger'; -interface WebhookResponse { - success: boolean; - statusCode?: number; - error?: string; +export interface WebhookResponse { + result: WebhookDeliveryAttemptResult; + statusCode: number | null; + error: string | null; + durationMs: number; } +export const WEBHOOK_REQUEST_TIMEOUT_SECONDS = 30; + export class WebhookDispatcher { - private readonly defaultTimeout = 30000; // 30 seconds + private readonly defaultTimeout = WEBHOOK_REQUEST_TIMEOUT_SECONDS * 1000; /** * Sends an event to a webhook URL. @@ -34,6 +40,7 @@ export class WebhookDispatcher { ): Promise { const timestamp = Now(); const payload = JSON.stringify(event); + const startedAt = Date.now(); const headers: Record = { 'Content-Type': 'application/json', @@ -52,6 +59,8 @@ export class WebhookDispatcher { signal: AbortSignal.timeout(this.defaultTimeout), }); + const durationMs = Date.now() - startedAt; + if (!response.ok) { Logger.warn('Webhook delivery failed', { eventId: event.id, @@ -60,9 +69,10 @@ export class WebhookDispatcher { }); return { - success: false, + result: 'http_error', statusCode: response.status, error: `HTTP ${response.status}`, + durationMs, }; } @@ -73,12 +83,20 @@ export class WebhookDispatcher { }); return { - success: true, + result: 'succeeded', statusCode: response.status, + error: null, + durationMs, }; } catch (error) { + const durationMs = Date.now() - startedAt; const errorMessage = error instanceof Error ? error.message : 'Unknown error'; + // AbortSignal.timeout rejects with a TimeoutError. + const result = + error instanceof Error && error.name === 'TimeoutError' + ? 'timed_out' + : 'network_error'; Logger.error('Webhook delivery error', error, { eventId: event.id, @@ -86,8 +104,10 @@ export class WebhookDispatcher { }); return { - success: false, + result, + statusCode: null, error: errorMessage, + durationMs, }; } } diff --git a/apps/api/src/routes/index.ts b/apps/api/src/routes/index.ts index 7e0dc56..466d536 100644 --- a/apps/api/src/routes/index.ts +++ b/apps/api/src/routes/index.ts @@ -34,6 +34,7 @@ import invoiceItemsRouter from './invoiceItems.routes'; import invoicesRouter from './invoices.routes'; import reportingRouter from './reporting.routes'; import billingRouter from './billing.routes'; +import webhookDeliveriesRouter from './webhookDeliveries.routes'; import telemetryRouter from './telemetry.routes'; import identityVerificationSessionsRouter from './identityVerificationSessions.routes'; import identityWebhooksRouter from './identityWebhooks.routes'; @@ -57,6 +58,9 @@ router.use('/operator', operatorRouter); // Billing run: operator key (Cloud Scheduler) or platform API key endpoints. router.use('/billing', billingRouter); +// Webhook delivery retries: operator key (Cloud Scheduler) or platform API key endpoints. +router.use('/webhook_deliveries', webhookDeliveriesRouter); + // --- Authenticated Routes --- // All routes below this line require an API Key router.use(ValidateApiKey); diff --git a/apps/api/src/routes/webhookDeliveries.routes.ts b/apps/api/src/routes/webhookDeliveries.routes.ts new file mode 100644 index 0000000..2d2e893 --- /dev/null +++ b/apps/api/src/routes/webhookDeliveries.routes.ts @@ -0,0 +1,54 @@ +import * as express from 'express'; +import { AsyncHandler } from '../utils/AsyncHandler'; +import { Logger } from '../utils/Logger'; +import { db } from '../modules/Database'; +import { WebhookDeliveryWorker } from '../modules/WebhookDeliveryWorker'; +import { ValidateOperatorKey } from '../middleware/OperatorMiddleware'; +import { ValidateApiKey } from '../middleware/AuthMiddleware'; +import { RequirePlatform } from '../middleware/Authorization'; + +const router = express.Router(); +const worker = new WebhookDeliveryWorker(db); + +function ParseBatchSize(body: unknown): number | undefined { + return typeof (body as { batch_size?: unknown })?.batch_size === 'number' + ? (body as { batch_size: number }).batch_size + : undefined; +} + +async function RunRetryBatch( + req: express.Request, + res: express.Response, + platformAccountId?: string +): Promise { + const result = await worker.ProcessBatch({ + limit: ParseBatchSize(req.body), + platformAccountId, + }); + + Logger.info('Webhook delivery retry run completed', { + scope: platformAccountId ?? 'operator', + ...result, + }); + + res.json({ object: 'webhook_delivery.batch', ...result }); +} + +router.post( + '/process', + ValidateOperatorKey, + AsyncHandler(async (req: express.Request, res: express.Response) => + RunRetryBatch(req, res) + ) +); + +router.post( + '/process_for_platform', + ValidateApiKey, + RequirePlatform(), + AsyncHandler(async (req: express.Request, res: express.Response) => + RunRetryBatch(req, res, req.user.account) + ) +); + +export default router; diff --git a/libs/shared-types/src/lib/WebhookDelivery.ts b/libs/shared-types/src/lib/WebhookDelivery.ts new file mode 100644 index 0000000..4e4f6d8 --- /dev/null +++ b/libs/shared-types/src/lib/WebhookDelivery.ts @@ -0,0 +1,67 @@ +/** Status of a webhook delivery. Only `pending` and `retrying` can be attempted. */ +export type WebhookDeliveryStatus = + | 'pending' + | 'retrying' + | 'succeeded' + | 'failed'; + +export type WebhookDeliveryAttemptResult = + | 'succeeded' + | 'http_error' + | 'timed_out' + | 'network_error'; + +/** One attempt to deliver an Event to a webhook endpoint. @internal */ +export interface WebhookDeliveryAttempt { + /** 1-based position in the delivery's attempt history */ + attempt_number: number; + /** Time the request was sent, in seconds since the Unix epoch */ + attempted_at: number; + /** Time the request finished, in seconds since the Unix epoch */ + completed_at: number; + result: WebhookDeliveryAttemptResult; + /** HTTP status, or null when no response was received */ + http_status: number | null; + duration_ms: number; + error: string | null; + url: string; +} + +/** + * Delivery of one Event to one webhook endpoint, persisted for every + * subscribed endpoint before the first attempt. @internal + */ +export interface WebhookDelivery { + id: string; + /** String representing the object's type. Objects of the same type share the same value. */ + object: 'webhook_delivery'; + event_id: string; + webhook_endpoint_id: string; + /** Retries stop once the delivery is `succeeded` or `failed` */ + status: WebhookDeliveryStatus; + /** Time the next attempt is due, or null when the delivery is no longer retried */ + next_attempt_at: number | null; + /** Time the first successful attempt was made, or null */ + delivered_at: number | null; + /** Time the current claim expires, or null when no worker holds it */ + claim_until: number | null; + /** Token of the worker holding the claim. A worker whose token no longer matches has lost it and cannot settle the delivery */ + claim_token: string | null; + attempts: WebhookDeliveryAttempt[]; + + /** + * The platform account that owns the Event this delivery belongs to. + * @zoneless_extension + */ + platform_account: string; +} + +/** Result of one bounded retry run. @internal */ +export interface WebhookDeliveryBatch { + object: 'webhook_delivery.batch'; + /** A claim lost to another worker counts here and in none of the outcomes below */ + processed: number; + succeeded: number; + retrying: number; + failed: number; +} diff --git a/libs/shared-types/src/lib/index.ts b/libs/shared-types/src/lib/index.ts index 2382798..53c402a 100644 --- a/libs/shared-types/src/lib/index.ts +++ b/libs/shared-types/src/lib/index.ts @@ -36,4 +36,5 @@ export * from './SubscriptionItem'; export * from './Telemetry'; export * from './TopUp'; export * from './Transfer'; +export * from './WebhookDelivery'; export * from './WebhookEndpoint';