From 880df9aa68f7e219533f68749a820d83562dddd1 Mon Sep 17 00:00:00 2001 From: Daveamhs Date: Wed, 30 Sep 2026 02:02:00 +0100 Subject: [PATCH] feat: implement webhook event delivery system --- .../migration.sql | 1 + prisma/schema.prisma | 3 +- .../integration/tip-flow.integration.test.ts | 49 ++++---- .../creators/verification.service.test.ts | 19 +++ src/domains/creators/verification.service.ts | 4 + src/domains/payments/payment.service.test.ts | 8 +- src/domains/payments/payment.service.ts | 30 +++-- src/domains/webhooks/events.test.ts | 24 ++++ src/domains/webhooks/events.ts | 2 + src/domains/webhooks/webhook.events.ts | 65 ++++++++++ src/domains/webhooks/webhook.routes.test.ts | 25 ++++ src/domains/webhooks/webhook.routes.ts | 56 ++++++--- src/domains/webhooks/webhook.service.test.ts | 111 ++++++++++++++++++ src/domains/webhooks/webhook.service.ts | 77 ++++++++---- src/lib/queue.ts | 25 +++- .../workers/stellar-confirmation.worker.ts | 8 +- src/lib/workers/webhook-dispatch.worker.ts | 62 ++++++---- 17 files changed, 460 insertions(+), 109 deletions(-) create mode 100644 prisma/migrations/20260924000000_webhook_event_delivery_index/migration.sql create mode 100644 src/domains/creators/verification.service.test.ts create mode 100644 src/domains/webhooks/events.test.ts create mode 100644 src/domains/webhooks/events.ts create mode 100644 src/domains/webhooks/webhook.events.ts create mode 100644 src/domains/webhooks/webhook.routes.test.ts create mode 100644 src/domains/webhooks/webhook.service.test.ts diff --git a/prisma/migrations/20260924000000_webhook_event_delivery_index/migration.sql b/prisma/migrations/20260924000000_webhook_event_delivery_index/migration.sql new file mode 100644 index 0000000..25a848d --- /dev/null +++ b/prisma/migrations/20260924000000_webhook_event_delivery_index/migration.sql @@ -0,0 +1 @@ +CREATE INDEX "WebhookEvent_createdAt_idx" ON "WebhookEvent"("createdAt"); diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 1387fe7..5d304d7 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -100,7 +100,7 @@ model Webhook { creator Creator @relation(fields: [creatorId], references: [id], onDelete: Cascade) url String - events String[] // tip.created, tip.confirmed, tip.failed, payout.completed + events String[] // Versioned event subscriptions, e.g. tip.created secret String // HMAC secret for signature verification active Boolean @default(true) @@ -128,6 +128,7 @@ model WebhookEvent { @@index([webhookId]) @@index([status]) + @@index([createdAt]) } model WalletFlag { diff --git a/src/__tests__/integration/tip-flow.integration.test.ts b/src/__tests__/integration/tip-flow.integration.test.ts index d868683..24528ed 100644 --- a/src/__tests__/integration/tip-flow.integration.test.ts +++ b/src/__tests__/integration/tip-flow.integration.test.ts @@ -1,3 +1,4 @@ +import { randomUUID } from 'node:crypto'; import { describe, it, expect, beforeAll, afterAll, beforeEach } from 'vitest'; import { PrismaClient } from '@prisma/client'; import { PaymentService } from '../../domains/payments/payment.service'; @@ -10,34 +11,29 @@ let userService: UserService; let payoutService: PayoutService; // Test data +const createdUserIds: string[] = []; let testUserId: string; let testCreatorId: string; let testCreatorUserId: string; -let isDbAvailable = false; - const isDbAvailable = Boolean(process.env.DATABASE_URL); describe.skipIf(!isDbAvailable)('Tip Flow Integration Tests', () => { beforeAll(async () => { - try { - prisma = new PrismaClient(); - await prisma.$connect(); - paymentService = new PaymentService(prisma); - userService = new UserService(prisma); - payoutService = new PayoutService(prisma); - } catch { - console.warn('Database connection failed, skipping integration tests'); - } + prisma = new PrismaClient(); + await prisma.$connect(); + paymentService = new PaymentService(prisma); + userService = new UserService(prisma); + payoutService = new PayoutService(prisma); }); afterAll(async () => { if (!prisma) return; try { // Clean up test data - await prisma.tip.deleteMany({}); - await prisma.wallet.deleteMany({}); - await prisma.creator.deleteMany({}); - await prisma.user.deleteMany({}); + await prisma.tip.deleteMany({ where: { fromUserId: { in: createdUserIds } } }); + await prisma.wallet.deleteMany({ where: { userId: { in: createdUserIds } } }); + await prisma.creator.deleteMany({ where: { userId: { in: createdUserIds } } }); + await prisma.user.deleteMany({ where: { id: { in: createdUserIds } } }); await prisma.$disconnect(); } catch { // ignore cleanup errors on disconnected DB @@ -59,6 +55,7 @@ describe.skipIf(!isDbAvailable)('Tip Flow Integration Tests', () => { }, }); testUserId = fanUser.id; + createdUserIds.push(fanUser.id); const creatorUser = await prisma.user.create({ data: { @@ -69,26 +66,28 @@ describe.skipIf(!isDbAvailable)('Tip Flow Integration Tests', () => { }, }); testCreatorUserId = creatorUser.id; + createdUserIds.push(creatorUser.id); - // Create creator profile + // Create verified creator profile const creator = await prisma.creator.create({ data: { userId: testCreatorUserId, username: `creator-${Date.now()}`, displayName: 'Test Creator', isPublic: true, - }, - }); - testCreatorId = creator.id; - - // Link verified wallet to fan - await prisma.wallet.create({ - data: { - userId: testUserId, - publicKey: `GBBD47AB2EB00E041B61C1B7AD184E687E24658D52EDFFDD118F5E6221D60E${Math.random().toString().slice(2, 4)}`, verified: true, }, }); + +testCreatorId = creator.id; + // Link verified wallet to fan + await prisma.wallet.create({ + data: { + userId: testUserId, + publicKey: `TEST-${randomUUID()}`, + verified: true, + }, + }); }); describe('Complete tip flow', () => { diff --git a/src/domains/creators/verification.service.test.ts b/src/domains/creators/verification.service.test.ts new file mode 100644 index 0000000..92bd38e --- /dev/null +++ b/src/domains/creators/verification.service.test.ts @@ -0,0 +1,19 @@ +import type { PrismaClient } from '@prisma/client'; +import { beforeEach, expect, it, vi } from 'vitest'; +import { VerificationService } from './verification.service'; +const { publish } = vi.hoisted(() => ({ publish: vi.fn() })); +vi.mock('../webhooks/webhook.service', () => ({ WebhookService: vi.fn(() => ({ dispatchEvent: publish })) })); +beforeEach(() => { vi.clearAllMocks(); }); +it('publishes creator.verified after verification is saved', async () => { + const prisma = { creator: { findUnique: vi.fn().mockResolvedValue({ verified: false }), update: vi.fn().mockResolvedValue({}) } }; + publish.mockImplementation(async () => expect(prisma.creator.update).toHaveBeenCalled()); + await new VerificationService(prisma as unknown as PrismaClient).verifyCreator('creator'); + expect(publish).toHaveBeenCalledWith('creator', 'creator', 'creator.verified', { creatorId: 'creator' }); +}); +it('does not republish an already verified creator', async () => { + const prisma = { creator: { findUnique: vi.fn().mockResolvedValue({ verified: true }), update: vi.fn() } }; + await new VerificationService(prisma as unknown as PrismaClient).verifyCreator('creator'); + expect(publish).not.toHaveBeenCalled(); +}); + + diff --git a/src/domains/creators/verification.service.ts b/src/domains/creators/verification.service.ts index 4bd1e6f..cbf97e1 100644 --- a/src/domains/creators/verification.service.ts +++ b/src/domains/creators/verification.service.ts @@ -2,6 +2,7 @@ import { PrismaClient } from '@prisma/client'; import { BaseService } from '../../services/base.service'; import { ValidationError } from '../../utils/errors'; import { logger } from '../../utils/logger'; +import { WebhookService } from '../webhooks/webhook.service'; export class VerificationService extends BaseService { constructor(private prisma: PrismaClient) { @@ -24,6 +25,9 @@ export class VerificationService extends BaseService { }); logger.info(`Creator verified: ${creatorId}`); + if (!creator.verified) { + await new WebhookService(this.prisma).dispatchEvent(creatorId, creatorId, 'creator.verified', { creatorId }); + } return { verified: true }; }); diff --git a/src/domains/payments/payment.service.test.ts b/src/domains/payments/payment.service.test.ts index cc8e688..e28c60a 100644 --- a/src/domains/payments/payment.service.test.ts +++ b/src/domains/payments/payment.service.test.ts @@ -1,3 +1,7 @@ +import type { PrismaClient } from '@prisma/client'; +vi.mock('../webhooks/webhook.service', () => ({ WebhookService: vi.fn(() => ({ dispatchEvent: publish })) })); +const { publish } = vi.hoisted(() => ({ publish: vi.fn().mockResolvedValue(undefined) })); +vi.mock('../../lib/queue', () => ({ stellarConfirmationQueue: { add: vi.fn() } })); import { describe, it, expect, beforeEach, vi } from 'vitest'; import { PaymentService } from './payment.service'; import { ValidationError, NotFoundError } from '../../utils/errors'; @@ -27,7 +31,7 @@ describe('PaymentService', () => { let paymentService: PaymentService; beforeEach(() => { - paymentService = new PaymentService(mockPrisma as any); + paymentService = new PaymentService(mockPrisma as unknown as PrismaClient); vi.clearAllMocks(); }); @@ -68,6 +72,7 @@ describe('PaymentService', () => { message: 'Great content!', }); + expect(publish).toHaveBeenCalledWith(creatorId, 'tip-123', 'tip.created', expect.objectContaining({ amount: 100 })); expect(result.id).toBe('tip-123'); expect(result.amount).toBe(100); expect(result.status).toBe('pending'); @@ -319,6 +324,7 @@ describe('PaymentService', () => { expect(result.status).toBe('completed'); expect(mockPrisma.creator.update).toHaveBeenCalled(); + expect(publish).toHaveBeenCalledWith(creatorId, tipId, 'payment.completed', expect.objectContaining({ amount: 100 })); }); it('should throw NotFoundError if tip does not exist', async () => { diff --git a/src/domains/payments/payment.service.ts b/src/domains/payments/payment.service.ts index be5dec7..eb4314c 100644 --- a/src/domains/payments/payment.service.ts +++ b/src/domains/payments/payment.service.ts @@ -1,4 +1,4 @@ -import { PrismaClient } from '@prisma/client'; +import { PrismaClient, Tip } from '@prisma/client'; import { BaseService } from '../../services/base.service'; import { CreateTipRequest, @@ -12,7 +12,6 @@ import { ValidationError, NotFoundError, UnauthorizedError } from '../../utils/e import { buildPaymentTransaction, submitSignedTransaction, - checkTransactionStatus, } from '../../lib/stellar/transactions'; import { logger } from '../../utils/logger'; import { stellarConfirmationQueue } from '../../lib/queue'; @@ -22,6 +21,7 @@ import { parseSortParameters, } from '../../utils/pagination'; import { paginateWithCursor } from '../../db/pagination'; +import { WebhookService } from '../webhooks/webhook.service'; export class PaymentService extends BaseService { constructor(private prisma: PrismaClient) { @@ -97,6 +97,9 @@ export class PaymentService extends BaseService { }); logger.info(`Tip created: ${tip.id} from ${userId} to ${data.creatorId} for ${data.amount}`); + await new WebhookService(this.prisma).dispatchEvent(tip.creatorId, tip.id, 'tip.created', { + tipId: tip.id, amount: tip.amount, message: tip.message, status: tip.status, + }); return this.formatTipResponse(tip); }); } @@ -153,7 +156,7 @@ export class PaymentService extends BaseService { const safePage = sanitizePageNumber(page); const safePageSize = sanitizePageSize(pageSize, 20); - const where: any = { creatorId }; + const where: Record = { creatorId }; if (options.status) { where.status = options.status; } @@ -216,12 +219,12 @@ export class PaymentService extends BaseService { throw new NotFoundError('Creator'); } - const where: any = { creatorId }; + const where: Record = { creatorId }; if (params.status) { where.status = params.status; } - const result = await paginateWithCursor( + const result = await paginateWithCursor( this.prisma.tip, { limit: params.limit, @@ -271,7 +274,7 @@ export class PaymentService extends BaseService { const safePage = sanitizePageNumber(page); const safePageSize = sanitizePageSize(pageSize, 20); - const where: any = { fromUserId: userId }; + const where: Record = { fromUserId: userId }; if (options.status) { where.status = options.status; } @@ -326,12 +329,12 @@ export class PaymentService extends BaseService { } = {} ) { return this.executeWithLogging('payment.getUserTipHistoryCursor', async () => { - const where: any = { fromUserId: userId }; + const where: Record = { fromUserId: userId }; if (params.status) { where.status = params.status; } - const result = await paginateWithCursor( + const result = await paginateWithCursor( this.prisma.tip, { limit: params.limit, @@ -403,6 +406,9 @@ export class PaymentService extends BaseService { }); logger.info(`Tip completed and creator earnings updated: ${tipId}, amount: ${tip.amount}`); + await new WebhookService(this.prisma).dispatchEvent(tip.creatorId, tip.id, 'payment.completed', { + tipId: tip.id, amount: tip.amount, transactionHash: tip.transactionHash, + }); } return this.formatTipResponse(updatedTip); @@ -457,13 +463,13 @@ export class PaymentService extends BaseService { memo: `tip-${tipId}`, }); - const transaction = transactionBuilder.build(); - const transactionEnvelope = transaction.toEnvelope().toXDR(); + const transaction = transactionBuilder; + const transactionEnvelope = transaction.toXDR(); logger.debug(`Payment transaction built for tip: ${tipId}`); return { - transactionEnvelope: transactionEnvelope as any as string, + transactionEnvelope: transactionEnvelope, tipId, fee: 100, // Base fee in stroops }; @@ -564,7 +570,7 @@ export class PaymentService extends BaseService { /** * Format database tip record to response DTO */ - private formatTipResponse(tip: any): TipResponse { + private formatTipResponse(tip: { id: string; fromUserId: string; creatorId: string; amount: number; message: string | null; status: string; transactionHash?: string | null; createdAt: Date; updatedAt: Date }): TipResponse { return { id: tip.id, fromUserId: tip.fromUserId, diff --git a/src/domains/webhooks/events.test.ts b/src/domains/webhooks/events.test.ts new file mode 100644 index 0000000..d8e894b --- /dev/null +++ b/src/domains/webhooks/events.test.ts @@ -0,0 +1,24 @@ +import { describe, expect, it } from 'vitest'; +import { + createWebhookSignature, + isWebhookEventType, + verifyWebhookSignature, +} from './events'; + +describe('webhook event contracts', () => { + it('accepts only supported version 1 event types', () => { + expect(isWebhookEventType('tip.created')).toBe(true); + expect(isWebhookEventType('creator.verified')).toBe(true); + expect(isWebhookEventType('payment.completed')).toBe(true); + expect(isWebhookEventType('tip.confirmed')).toBe(false); + }); + + it('creates and verifies an HMAC SHA-256 signature over the raw body', () => { + const body = JSON.stringify({ id: 'evt_1', type: 'tip.created', version: '1' }); + const signature = createWebhookSignature(body, 'secret'); + expect(signature).toMatch(/^sha256=[a-f0-9]{64}$/); + expect(verifyWebhookSignature(body, 'secret', signature)).toBe(true); + expect(verifyWebhookSignature(`${body} `, 'secret', signature)).toBe(false); + expect(verifyWebhookSignature(body, 'wrong-secret', signature)).toBe(false); + }); +}); diff --git a/src/domains/webhooks/events.ts b/src/domains/webhooks/events.ts new file mode 100644 index 0000000..0060dc8 --- /dev/null +++ b/src/domains/webhooks/events.ts @@ -0,0 +1,2 @@ +// Compatibility export for existing consumers. +export * from './webhook.events'; diff --git a/src/domains/webhooks/webhook.events.ts b/src/domains/webhooks/webhook.events.ts new file mode 100644 index 0000000..f6b88c3 --- /dev/null +++ b/src/domains/webhooks/webhook.events.ts @@ -0,0 +1,65 @@ +import { createHmac, timingSafeEqual } from 'node:crypto'; + +export const WEBHOOK_EVENT_TYPES = [ + 'tip.created', + 'creator.verified', + 'payment.completed', +] as const; + +export type WebhookEventType = (typeof WEBHOOK_EVENT_TYPES)[number]; + +export interface WebhookEventEnvelope< + T = Record, +> { + id: string; + type: WebhookEventType; + version: '1'; + createdAt: string; + data: T; +} + +export const isWebhookEventType = ( + event: string, +): event is WebhookEventType => + (WEBHOOK_EVENT_TYPES as readonly string[]).includes(event); + +/** + * Creates an HMAC SHA-256 signature for a webhook payload. + * + * Format: + * sha256= + */ +export const createWebhookSignature = ( + rawBody: string, + secret: string, +): string => { + const digest = createHmac('sha256', secret) + .update(rawBody, 'utf8') + .digest('hex'); + + return `sha256=${digest}`; +}; + +/** + * Verifies a webhook signature using a timing-safe comparison. + */ +export const verifyWebhookSignature = ( + rawBody: string, + secret: string, + signature: string, +): boolean => { + if (!signature?.startsWith('sha256=')) { + return false; + } + + const expectedSignature = createWebhookSignature(rawBody, secret); + + const expected = Buffer.from(expectedSignature, 'utf8'); + const supplied = Buffer.from(signature, 'utf8'); + + if (expected.length !== supplied.length) { + return false; + } + + return timingSafeEqual(expected, supplied); +}; \ No newline at end of file diff --git a/src/domains/webhooks/webhook.routes.test.ts b/src/domains/webhooks/webhook.routes.test.ts new file mode 100644 index 0000000..03d1fbc --- /dev/null +++ b/src/domains/webhooks/webhook.routes.test.ts @@ -0,0 +1,25 @@ +import type { PrismaClient } from '@prisma/client'; +import Fastify, { FastifyRequest, FastifyReply } from 'fastify'; +import { afterEach, expect, it, vi } from 'vitest'; +import { registerWebhookRoutes } from './webhook.routes'; +const { testWebhook } = vi.hoisted(() => ({ testWebhook: vi.fn() })); +vi.mock('./webhook.service', () => ({ WebhookService: vi.fn(() => ({ testWebhook })) })); +vi.mock('../../middleware/auth', () => ({ authMiddleware: async (request: FastifyRequest, reply: FastifyReply) => { + if (!request.headers.authorization) return reply.code(401).send({ error: 'Unauthorized' }); + request.user = { userId: 'user', email: 'user@example.com', role: 'creator' }; +} })); +const app = Fastify(); +const prisma = { creator: { findUnique: vi.fn().mockResolvedValue({ id: 'creator' }) } }; +registerWebhookRoutes(app, prisma as unknown as PrismaClient); +afterEach(() => { vi.clearAllMocks(); }); +it('requires authentication before testing a webhook', async () => { + const response = await app.inject({ method: 'POST', url: '/api/v1/webhooks/hook/test' }); + expect(response.statusCode).toBe(401); + expect(testWebhook).not.toHaveBeenCalled(); +}); +it('queues an authenticated test using the current creator identity', async () => { + const response = await app.inject({ method: 'POST', url: '/api/v1/webhooks/hook/test', headers: { authorization: 'Bearer token' } }); + expect(response.statusCode).toBe(202); + expect(testWebhook).toHaveBeenCalledWith('hook', 'creator'); +}); + diff --git a/src/domains/webhooks/webhook.routes.ts b/src/domains/webhooks/webhook.routes.ts index 9f29e35..f0ab2c1 100644 --- a/src/domains/webhooks/webhook.routes.ts +++ b/src/domains/webhooks/webhook.routes.ts @@ -1,3 +1,4 @@ +import type {} from '@fastify/swagger'; import { FastifyInstance, FastifyRequest, FastifyReply } from 'fastify'; import { PrismaClient } from '@prisma/client'; import { WebhookService, CreateWebhookRequest } from './webhook.service'; @@ -8,13 +9,34 @@ import { ValidationError, AppError } from '../../utils/errors'; export const registerWebhookRoutes = (app: FastifyInstance, prisma: PrismaClient): void => { const webhookService = new WebhookService(prisma); + app.post<{ Params: { id: string } }>( + '/api/v1/webhooks/:id/test', + { preHandler: authMiddleware }, + async (request, reply) => { + try { + const user = request.user; + if (!user) throw new Error('User not found'); + const creator = await prisma.creator.findUnique({ where: { userId: user.userId } }); + if (!creator) { + reply.code(404).send(formatError('Creator not found', 'CREATOR_NOT_FOUND')); + return; + } + await webhookService.testWebhook(request.params.id, creator.id); + reply.code(202).send(formatSuccess({ message: 'Test event queued' })); + } catch (error) { + if (error instanceof AppError) reply.code(error.statusCode).send(formatError(error.message, error.code)); + else throw error; + } + } + ); + // POST /api/v1/webhooks - Register a webhook app.post<{ Body: CreateWebhookRequest }>( '/api/v1/webhooks', { preHandler: authMiddleware, schema: { - description: 'Register a webhook to receive events when tips are created and confirmed.', + description: 'Register a webhook to receive events for tips, creator verification, and completed payments.', body: { type: 'object', required: ['url', 'events'], @@ -24,16 +46,16 @@ export const registerWebhookRoutes = (app: FastifyInstance, prisma: PrismaClient type: 'array', items: { type: 'string' }, description: - 'Events to subscribe to (tip.created, tip.confirmed, tip.failed, payout.completed)', + 'Events to subscribe to (tip.created, creator.verified, payment.completed)', }, }, }, response: { - 201: { description: 'Webhook registered' }, - 400: { description: 'Validation error' }, - 401: { description: 'Unauthorized' }, + 201: { type: 'object', additionalProperties: true, description: 'Webhook registered' }, + 400: { type: 'object', additionalProperties: true, description: 'Validation error' }, + 401: { type: 'object', additionalProperties: true, description: 'Unauthorized' }, }, - } as any, + }, }, async (request: FastifyRequest, reply: FastifyReply) => { try { @@ -73,10 +95,10 @@ export const registerWebhookRoutes = (app: FastifyInstance, prisma: PrismaClient schema: { description: 'Get all webhooks registered for the creator.', response: { - 200: { description: 'List of webhooks' }, - 401: { description: 'Unauthorized' }, + 200: { type: 'object', additionalProperties: true, description: 'List of webhooks' }, + 401: { type: 'object', additionalProperties: true, description: 'Unauthorized' }, }, - } as any, + }, }, async (request: FastifyRequest, reply: FastifyReply) => { try { @@ -118,11 +140,11 @@ export const registerWebhookRoutes = (app: FastifyInstance, prisma: PrismaClient }, }, response: { - 200: { description: 'Webhook deleted' }, - 401: { description: 'Unauthorized' }, - 404: { description: 'Webhook not found' }, + 200: { type: 'object', additionalProperties: true, description: 'Webhook deleted' }, + 401: { type: 'object', additionalProperties: true, description: 'Unauthorized' }, + 404: { type: 'object', additionalProperties: true, description: 'Webhook not found' }, }, - } as any, + }, }, async (request: FastifyRequest, reply: FastifyReply) => { try { @@ -179,11 +201,11 @@ export const registerWebhookRoutes = (app: FastifyInstance, prisma: PrismaClient }, }, response: { - 200: { description: 'Delivery history with pagination' }, - 401: { description: 'Unauthorized' }, - 404: { description: 'Webhook not found' }, + 200: { type: 'object', additionalProperties: true, description: 'Delivery history with pagination' }, + 401: { type: 'object', additionalProperties: true, description: 'Unauthorized' }, + 404: { type: 'object', additionalProperties: true, description: 'Webhook not found' }, }, - } as any, + }, }, async (request: FastifyRequest, reply: FastifyReply) => { try { diff --git a/src/domains/webhooks/webhook.service.test.ts b/src/domains/webhooks/webhook.service.test.ts new file mode 100644 index 0000000..bdca11a --- /dev/null +++ b/src/domains/webhooks/webhook.service.test.ts @@ -0,0 +1,111 @@ +import type { WorkerOptions, JobsOptions } from 'bullmq'; +import type { PrismaClient } from '@prisma/client'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { WebhookService } from './webhook.service'; +import { verifyWebhookSignature } from './webhook.events'; + +const mocks = vi.hoisted(() => ({ + add: vi.fn(), dead: vi.fn(), post: vi.fn(), + prisma: { webhook: { findUnique: vi.fn(), findMany: vi.fn() }, webhookEvent: { create: vi.fn(), update: vi.fn() } }, + process: undefined as ((job: { id: string; data: unknown; opts: JobsOptions; attemptsMade: number }) => Promise) | undefined, options: undefined as WorkerOptions | undefined, +})); +vi.mock('../../lib/queue', () => ({ webhookDispatchQueue: { add: mocks.add }, webhookDeadLetterQueue: { add: mocks.dead }, webhookConnection: {} })); +vi.mock('@prisma/client', () => ({ PrismaClient: vi.fn(() => mocks.prisma) })); +vi.mock('axios', () => ({ default: { post: mocks.post } })); +vi.mock('bullmq', () => ({ Worker: vi.fn((_name, process, options) => { mocks.process = process; mocks.options = options; return { on: vi.fn() }; }) })); +import '../../lib/workers/webhook-dispatch.worker'; + +const subscriber = { id: 'hook', creatorId: 'creator', active: true, events: ['tip.created'], url: 'https://example.com/hook', secret: 'secret' }; +const service = new WebhookService(mocks.prisma as unknown as PrismaClient); +const publish = (n = 0) => service.dispatchEvent('creator', `tip-${n}`, 'tip.created', { n }); +const job = (index = 0, attemptsMade = 0) => ({ id: mocks.add.mock.calls[index][2].jobId, data: mocks.add.mock.calls[index][1], opts: mocks.add.mock.calls[index][2], attemptsMade }); + +beforeEach(() => { + vi.clearAllMocks(); + mocks.prisma.webhook.findMany.mockImplementation(async ({ where }) => [subscriber].filter(w => w.creatorId === where.creatorId && w.active === where.active && w.events.includes(where.events.has) && (!where.id || where.id === w.id))); + mocks.prisma.webhook.findUnique.mockResolvedValue(subscriber); + let id = 0; + mocks.prisma.webhookEvent.create.mockImplementation(async ({ data }) => ({ ...data, id: `delivery-${++id}` })); + mocks.prisma.webhookEvent.update.mockResolvedValue({}); + mocks.add.mockResolvedValue({}); mocks.dead.mockResolvedValue({}); + mocks.post.mockReset().mockResolvedValue({ status: 200 }); +}); + +describe('webhook delivery lifecycle', () => { + it('persists pending records before enqueueing and passes the database eventId', async () => { + mocks.add.mockImplementation(async (_name, data) => { + expect(mocks.prisma.webhookEvent.create).toHaveBeenCalled(); + expect(data.eventId).toBe('delivery-1'); + }); + await publish(); + const stored = mocks.prisma.webhookEvent.create.mock.calls[0][0].data; + expect(stored.status).toBe('pending'); + expect(JSON.parse(stored.payload)).toEqual(job().data.payload); + expect(job().data.payload).toMatchObject({ type: 'tip.created', version: '1', data: { transactionId: 'tip-0' } }); + expect(new Date(job().data.payload.createdAt).toISOString()).toBe(job().data.payload.createdAt); + }); + it('preserves a pending record and enqueue error when Redis rejects a job', async () => { + mocks.add.mockRejectedValueOnce(new Error('Redis unavailable')); + await expect(publish()).rejects.toThrow('Redis unavailable'); + expect(mocks.prisma.webhookEvent.update).toHaveBeenCalledWith({ where: { id: 'delivery-1' }, data: { lastError: 'Redis unavailable' } }); + }); + it('delivers the signed raw envelope to the subscriber and tracks success', async () => { + await publish(); await mocks.process(job()); + const [url, body, options] = mocks.post.mock.calls[0]; + expect(url).toBe(subscriber.url); + expect(verifyWebhookSignature(body, subscriber.secret, options.headers['X-Dorisio-Signature'])).toBe(true); + expect(mocks.prisma.webhookEvent.update).toHaveBeenCalledWith(expect.objectContaining({ where: { id: 'delivery-1' }, data: expect.objectContaining({ status: 'delivered', attempts: 1 }) })); + }); + it('filters subscriptions, inactive hooks, and other creators before enqueueing', async () => { + await service.dispatchEvent('creator', 'tip', 'payment.completed', {}); + await service.dispatchEvent('other', 'tip', 'tip.created', {}); + expect(mocks.add).not.toHaveBeenCalled(); + expect(mocks.prisma.webhook.findMany).toHaveBeenCalledWith({ where: { creatorId: 'creator', active: true, events: { has: 'payment.completed' } } }); + }); + it('uses five exponential attempts, retaining errors and dead-lettering only the last failure', async () => { + await publish(); + expect(job().opts).toMatchObject({ attempts: 5, backoff: { type: 'exponential', delay: 2000 } }); + mocks.post.mockRejectedValue(new Error('HTTP 503')); + for (let attempt = 0; attempt < 5; attempt++) { + await expect(mocks.process(job(0, attempt))).rejects.toThrow('HTTP 503'); + expect(mocks.prisma.webhookEvent.update).toHaveBeenLastCalledWith(expect.objectContaining({ data: expect.objectContaining({ attempts: attempt + 1, status: attempt === 4 ? 'failed' : 'pending', lastError: 'HTTP 503' }) })); + expect(mocks.dead).toHaveBeenCalledTimes(attempt === 4 ? 1 : 0); + } + expect(mocks.dead).toHaveBeenCalledWith('failed-delivery', expect.objectContaining({ eventId: 'delivery-1', error: 'HTTP 503', attempts: 5 }), { jobId: 'delivery-1' }); + }); + it('tracks success after a retry without dead-lettering', async () => { + await publish(); mocks.post.mockRejectedValueOnce(new Error('timeout')); + await expect(mocks.process(job())).rejects.toThrow(); + await mocks.process(job(0, 1)); + expect(mocks.prisma.webhookEvent.update).toHaveBeenLastCalledWith(expect.objectContaining({ data: expect.objectContaining({ status: 'delivered', attempts: 2 }) })); + expect(mocks.dead).not.toHaveBeenCalled(); + }); + it('keeps concurrent events and delivery IDs distinct', async () => { + await Promise.all(Array.from({ length: 25 }, (_, n) => publish(n))); + expect(new Set(mocks.add.mock.calls.map(c => c[1].eventId)).size).toBe(25); + expect(new Set(mocks.add.mock.calls.map(c => c[1].payload.id)).size).toBe(25); + await Promise.all(mocks.add.mock.calls.map((_, n) => mocks.process(job(n)))); + expect(mocks.prisma.webhookEvent.update).toHaveBeenCalledTimes(25); + }); + it('enqueues sequential publications in order and processes with concurrency one', async () => { + await publish(1); await publish(2); + expect(mocks.add.mock.calls.map(c => c[1].payload.data.n)).toEqual([1, 2]); + expect(mocks.options.concurrency).toBe(1); + await mocks.process(job(0)); await mocks.process(job(1)); + expect(mocks.post.mock.calls.map(c => JSON.parse(c[1]).data.n)).toEqual([1, 2]); + }); + it('allows a newer event to complete while an earlier event awaits retry', async () => { + await publish(1); await publish(2); + mocks.post.mockRejectedValueOnce(new Error('timeout')); + await expect(mocks.process(job(0))).rejects.toThrow(); + await mocks.process(job(1)); await mocks.process(job(0, 1)); + expect(mocks.post.mock.calls.map(c => JSON.parse(c[1]).data.n)).toEqual([1, 2, 1]); + }); + it('tests only the requested owned registered webhook', async () => { + await service.testWebhook('hook', 'creator'); + expect(mocks.prisma.webhook.findMany).toHaveBeenCalledWith({ where: { creatorId: 'creator', id: 'hook', active: true, events: { has: 'tip.created' } } }); + expect(job().data.payload.data.test).toBe(true); + await expect(service.testWebhook('hook', 'other')).rejects.toThrow(); + expect(mocks.add).toHaveBeenCalledTimes(1); + }); +}); diff --git a/src/domains/webhooks/webhook.service.ts b/src/domains/webhooks/webhook.service.ts index ed372a9..06aa038 100644 --- a/src/domains/webhooks/webhook.service.ts +++ b/src/domains/webhooks/webhook.service.ts @@ -4,10 +4,11 @@ import { ValidationError, NotFoundError } from '../../utils/errors'; import { webhookDispatchQueue } from '../../lib/queue'; import { logger } from '../../utils/logger'; import crypto from 'crypto'; +import { isWebhookEventType, WebhookEventEnvelope, WebhookEventType } from './webhook.events'; export interface CreateWebhookRequest { url: string; - events: string[]; + events: WebhookEventType[]; } export interface WebhookResponse { @@ -39,9 +40,8 @@ export class WebhookService extends BaseService { } // Validate events - const validEvents = ['tip.created', 'tip.confirmed', 'tip.failed', 'payout.completed']; for (const event of data.events) { - if (!validEvents.includes(event)) { + if (!isWebhookEventType(event)) { throw new ValidationError(`Invalid event type: ${event}`); } } @@ -109,14 +109,26 @@ export class WebhookService extends BaseService { async dispatchEvent( creatorId: string, transactionId: string, - eventType: string, - payload: Record + eventType: WebhookEventType, + payload: Record, + webhookId?: string ): Promise { return this.executeWithLogging('webhook.dispatch', async () => { - // Find active webhooks for this creator that subscribe to this event + if (!isWebhookEventType(eventType)) throw new ValidationError(`Invalid event type: ${eventType}`); + const eventId = crypto.randomUUID(); + const event: WebhookEventEnvelope = { + id: eventId, + type: eventType, + version: '1', + createdAt: new Date().toISOString(), + data: { ...payload, transactionId }, + }; + + // Filtering happens before work is placed on the delivery queue. const webhooks = await this.prisma.webhook.findMany({ where: { creatorId, + ...(webhookId ? { id: webhookId } : {}), active: true, events: { has: eventType, @@ -126,25 +138,48 @@ export class WebhookService extends BaseService { // Queue dispatch jobs for each webhook for (const webhook of webhooks) { - await webhookDispatchQueue.add( - 'dispatch-event', - { - webhookId: webhook.id, - transactionId, - eventType, - payload, - }, - { - attempts: 5, - backoff: { type: 'exponential', delay: 2000 }, - } - ); + const delivery = await this.prisma.webhookEvent.create({ + data: { webhookId: webhook.id, eventType, payload: JSON.stringify(event), status: 'pending' }, + }); + try { + await webhookDispatchQueue.add( + 'dispatch-event', + { + webhookId: webhook.id, + eventId: delivery.id, + eventType, + payload: event, + }, + { + attempts: 5, + backoff: { type: 'exponential', delay: 2000 }, + jobId: delivery.id, + } + ); + + } catch (error) { + // Keep the durable pending record available for recovery if enqueueing fails. + await this.prisma.webhookEvent.update({ + where: { id: delivery.id }, + data: { lastError: error instanceof Error ? error.message : String(error) }, + }); + throw error; + } logger.info(`Queued webhook dispatch for ${webhook.id} (event: ${eventType})`); } }); } + async testWebhook(webhookId: string, creatorId: string): Promise { + const webhook = await this.prisma.webhook.findUnique({ where: { id: webhookId } }); + if (!webhook || webhook.creatorId !== creatorId) throw new NotFoundError('Webhook'); + if (!webhook.active) throw new ValidationError('Webhook is inactive'); + const eventType = webhook.events.find(isWebhookEventType); + if (!eventType) throw new ValidationError('Webhook has no supported subscriptions'); + await this.dispatchEvent(creatorId, `test-${crypto.randomUUID()}`, eventType, { test: true }, webhook.id); + } + /** * Get webhook delivery history */ @@ -189,7 +224,7 @@ export class WebhookService extends BaseService { const safePageSize = sanitizePageSize(pageSize, 20); const skip = (safePage - 1) * safePageSize; - const where: any = { webhookId }; + const where: { webhookId: string; status?: string } = { webhookId }; if (status) { where.status = status; } @@ -228,7 +263,7 @@ export class WebhookService extends BaseService { }); } - private formatWebhookResponse(webhook: any): WebhookResponse { + private formatWebhookResponse(webhook: Omit & { createdAt: Date; updatedAt: Date }): WebhookResponse { return { id: webhook.id, creatorId: webhook.creatorId, diff --git a/src/lib/queue.ts b/src/lib/queue.ts index 1ed9f27..42fa88d 100644 --- a/src/lib/queue.ts +++ b/src/lib/queue.ts @@ -1,4 +1,4 @@ -import { Queue, Worker, QueueEvents } from 'bullmq'; +import { Queue, QueueEvents } from 'bullmq'; import { createClient } from 'redis'; import { config } from '../config/env'; import { logger } from '../utils/logger'; @@ -21,19 +21,33 @@ redis.connect().catch((err) => { logger.error('Failed to connect to Redis:', err); }); +// Let BullMQ create its own compatible Redis connections. +const webhookRedisUrl = new URL(config.REDIS_URL); +export const webhookConnection = { + host: webhookRedisUrl.hostname, + port: Number(webhookRedisUrl.port || 6379), + username: decodeURIComponent(webhookRedisUrl.username) || undefined, + password: decodeURIComponent(webhookRedisUrl.password) || undefined, + db: Number(webhookRedisUrl.pathname.slice(1) || 0), + ...(webhookRedisUrl.protocol === 'rediss:' ? { tls: {} } : {}), +}; // Job queues export const stellarConfirmationQueue = new Queue('stellar-confirmation', { - connection: redis as any, + connection: webhookConnection, }); -export const webhookDispatchQueue = new Queue('webhook-dispatch', { connection: redis as any }); +export const webhookDispatchQueue = new Queue('webhook-dispatch', { + connection: webhookConnection, + defaultJobOptions: { attempts: 5, backoff: { type: 'exponential', delay: 2000 } }, +}); +export const webhookDeadLetterQueue = new Queue('webhook-dead-letter', { connection: webhookConnection }); // Queue event handlers export const stellarConfirmationEvents = new QueueEvents('stellar-confirmation', { - connection: redis as any, + connection: webhookConnection, }); export const webhookDispatchEvents = new QueueEvents('webhook-dispatch', { - connection: redis as any, + connection: webhookConnection, }); // Initialize queue event listeners @@ -56,6 +70,7 @@ webhookDispatchEvents.on('failed', ({ jobId, failedReason }) => { export async function closeQueues() { await stellarConfirmationQueue.close(); await webhookDispatchQueue.close(); + await webhookDeadLetterQueue.close(); await stellarConfirmationEvents.close(); await webhookDispatchEvents.close(); await redis.quit(); diff --git a/src/lib/workers/stellar-confirmation.worker.ts b/src/lib/workers/stellar-confirmation.worker.ts index 094d7c5..9eb6c3d 100644 --- a/src/lib/workers/stellar-confirmation.worker.ts +++ b/src/lib/workers/stellar-confirmation.worker.ts @@ -1,12 +1,8 @@ import { Worker, Job } from 'bullmq'; -import { createClient } from 'redis'; +import { webhookConnection } from '../queue'; import { PrismaClient } from '@prisma/client'; -import { config } from '../../config/env'; import { logger } from '../../utils/logger'; -const redis = createClient({ - url: config.REDIS_URL, -}); const prisma = new PrismaClient(); @@ -39,7 +35,7 @@ export const stellarConfirmationWorker = new Worker( } }, { - connection: redis as any, + connection: webhookConnection, } ); diff --git a/src/lib/workers/webhook-dispatch.worker.ts b/src/lib/workers/webhook-dispatch.worker.ts index 80e97a2..09bc0e9 100644 --- a/src/lib/workers/webhook-dispatch.worker.ts +++ b/src/lib/workers/webhook-dispatch.worker.ts @@ -1,21 +1,16 @@ import { Worker, Job } from 'bullmq'; -import { createClient } from 'redis'; import axios from 'axios'; import { PrismaClient } from '@prisma/client'; -import { config } from '../../config/env'; import { logger } from '../../utils/logger'; -import crypto from 'crypto'; - -const redis = createClient({ - url: config.REDIS_URL, -}); +import { webhookDeadLetterQueue, webhookConnection } from '../queue'; +import { createWebhookSignature } from '../../domains/webhooks/webhook.events'; const prisma = new PrismaClient(); export const webhookDispatchWorker = new Worker( 'webhook-dispatch', async (job: Job) => { - const { webhookId, eventType, payload } = job.data; + const { webhookId, eventType, payload, eventId, deliveryId } = job.data; logger.info(`Dispatching webhook ${webhookId} for ${eventType} event`); @@ -26,28 +21,34 @@ export const webhookDispatchWorker = new Worker( throw new Error(`Webhook ${webhookId} not found`); } + if (!webhook.active || !webhook.events.includes(eventType)) { + throw new Error('Webhook is inactive or no longer subscribed'); + } + // Create signature for webhook verification - const signature = crypto - .createHmac('sha256', webhook.secret) - .update(JSON.stringify(payload)) - .digest('hex'); + const signature = createWebhookSignature( + typeof payload === 'string' ? payload : JSON.stringify(payload), + webhook.secret + ); - const response = await axios.post(webhook.url, payload, { + const response = await axios.post(webhook.url, typeof payload === 'string' ? payload : JSON.stringify(payload), { headers: { 'Content-Type': 'application/json', - 'X-Dorisio-Signature': `sha256=${signature}`, + 'X-Dorisio-Signature': signature, 'X-Dorisio-Event': eventType, - 'X-Dorisio-Delivery-Id': job.id, + 'X-Dorisio-Delivery-Id': deliveryId ?? eventId, + 'X-Dorisio-Event-Version': '1', }, timeout: 30000, + maxRedirects: 0, }); // Track successful dispatch await prisma.webhookEvent.update({ - where: { id: job.data.eventId }, + where: { id: deliveryId ?? eventId }, data: { status: 'delivered', - attempts: { increment: 1 }, + attempts: job.attemptsMade + 1, updatedAt: new Date(), }, }); @@ -61,20 +62,30 @@ export const webhookDispatchWorker = new Worker( // Track failed dispatch attempt await prisma.webhookEvent.update({ - where: { id: job.data.eventId }, + where: { id: deliveryId ?? eventId }, data: { - status: 'pending', - attempts: { increment: 1 }, + status: job.attemptsMade + 1 >= Number(job.opts.attempts ?? 1) ? 'failed' : 'pending', + attempts: job.attemptsMade + 1, lastError: errorMsg, updatedAt: new Date(), }, }); + if (job.attemptsMade + 1 >= Number(job.opts.attempts ?? 1)) { + await webhookDeadLetterQueue.add('failed-delivery', { + ...job.data, + failedAt: new Date().toISOString(), + error: errorMsg, + attempts: job.attemptsMade + 1, + }, { jobId: deliveryId ?? eventId }); + } + throw error; } }, { - connection: redis as any, + connection: webhookConnection, + concurrency: 1, } ); @@ -85,3 +96,12 @@ webhookDispatchWorker.on('completed', (job) => { webhookDispatchWorker.on('failed', (job, err) => { logger.error(`Webhook dispatch worker failed job ${job?.id}:`, err); }); + +webhookDispatchWorker.on('error', (error) => { + logger.error('Webhook worker error:', error); +}); + +export async function closeWebhookWorker(): Promise { + await webhookDispatchWorker.close(); + await prisma.$disconnect(); +}