diff --git a/.gitignore b/.gitignore index 6b11521a9..09a2a95ae 100644 --- a/.gitignore +++ b/.gitignore @@ -7,4 +7,5 @@ build/ coverage/ .temp/ .turbo/ -plan.md \ No newline at end of file +plan.md +/.kiroo \ No newline at end of file diff --git a/apps/dashboard-api/src/__tests__/webhook.controller.test.js b/apps/dashboard-api/src/__tests__/webhook.controller.test.js new file mode 100644 index 000000000..46484b493 --- /dev/null +++ b/apps/dashboard-api/src/__tests__/webhook.controller.test.js @@ -0,0 +1,382 @@ +'use strict'; + +const mongoose = require('mongoose'); + +// Mock @urbackend/common +jest.mock('@urbackend/common', () => ({ + Webhook: { + create: jest.fn(), + find: jest.fn(), + findOne: jest.fn(), + findOneAndUpdate: jest.fn(), + findOneAndDelete: jest.fn(), + }, + WebhookDelivery: { + find: jest.fn(), + countDocuments: jest.fn(), + }, + Project: { + findOne: jest.fn(), + }, + encrypt: jest.fn((val) => ({ encrypted: 'enc', iv: 'iv', tag: 'tag' })), + decrypt: jest.fn(() => 'decrypted-secret'), + createWebhookSchema: { + safeParse: jest.fn(), + }, + updateWebhookSchema: { + safeParse: jest.fn(), + }, + generateSignature: jest.fn(() => 'sha256=test-signature'), +})); + +const { + createWebhook, + getWebhooks, + getWebhook, + updateWebhook, + deleteWebhook, + getDeliveries, + testWebhook, +} = require('../controllers/webhook.controller'); + +const { + Webhook, + WebhookDelivery, + Project, + createWebhookSchema, + updateWebhookSchema, +} = require('@urbackend/common'); + +describe('webhook.controller', () => { + let req, res; + // Use valid MongoDB ObjectId format + const validProjectId = new mongoose.Types.ObjectId().toString(); + const validWebhookId = new mongoose.Types.ObjectId().toString(); + + beforeEach(() => { + jest.clearAllMocks(); + req = { + params: { projectId: validProjectId }, + user: { _id: 'user123', email: 'test@example.com' }, + body: {}, + query: {}, + }; + res = { + status: jest.fn().mockReturnThis(), + json: jest.fn().mockReturnThis(), + }; + }); + + describe('createWebhook', () => { + test('creates webhook with valid input', async () => { + Project.findOne.mockResolvedValue({ _id: validProjectId }); + createWebhookSchema.safeParse.mockReturnValue({ + success: true, + data: { + name: 'Test Webhook', + url: 'https://example.com/hook', + secret: 'whsec_testsecret12345678', + events: { posts: { insert: true } }, + enabled: true, + }, + }); + + const mockWebhook = { + _id: validWebhookId, + projectId: validProjectId, + name: 'Test Webhook', + url: 'https://example.com/hook', + events: new Map([['posts', { insert: true, update: false, delete: false }]]), + enabled: true, + createdAt: new Date(), + }; + Webhook.create.mockResolvedValue(mockWebhook); + + req.body = { + name: 'Test Webhook', + url: 'https://example.com/hook', + secret: 'whsec_testsecret12345678', + events: { posts: { insert: true } }, + }; + + await createWebhook(req, res); + + expect(Project.findOne).toHaveBeenCalledWith({ + _id: validProjectId, + owner: 'user123', + }); + expect(Webhook.create).toHaveBeenCalled(); + expect(res.status).toHaveBeenCalledWith(201); + expect(res.json).toHaveBeenCalledWith( + expect.objectContaining({ + message: 'Webhook created', + data: expect.objectContaining({ + name: 'Test Webhook', + }), + }) + ); + }); + + test('returns 404 if project not found', async () => { + Project.findOne.mockResolvedValue(null); + + await createWebhook(req, res); + + expect(res.status).toHaveBeenCalledWith(404); + expect(res.json).toHaveBeenCalledWith({ error: 'Project not found' }); + }); + + test('returns 400 on validation failure', async () => { + Project.findOne.mockResolvedValue({ _id: validProjectId }); + createWebhookSchema.safeParse.mockReturnValue({ + success: false, + error: { errors: [{ message: 'Invalid URL' }] }, + }); + + await createWebhook(req, res); + + expect(res.status).toHaveBeenCalledWith(400); + expect(res.json).toHaveBeenCalledWith( + expect.objectContaining({ error: 'Validation failed' }) + ); + }); + }); + + describe('getWebhooks', () => { + test('returns all webhooks for a project', async () => { + Project.findOne.mockResolvedValue({ _id: validProjectId }); + Webhook.find.mockReturnValue({ + lean: jest.fn().mockResolvedValue([ + { _id: 'wh1', name: 'Hook 1', url: 'https://a.com', enabled: true }, + { _id: 'wh2', name: 'Hook 2', url: 'https://b.com', enabled: false }, + ]), + }); + + await getWebhooks(req, res); + + expect(Webhook.find).toHaveBeenCalledWith({ projectId: validProjectId }); + expect(res.json).toHaveBeenCalledWith({ + data: expect.arrayContaining([ + expect.objectContaining({ name: 'Hook 1' }), + expect.objectContaining({ name: 'Hook 2' }), + ]), + }); + }); + }); + + describe('getWebhook', () => { + test('returns single webhook', async () => { + req.params.webhookId = validWebhookId; + Project.findOne.mockResolvedValue({ _id: validProjectId }); + Webhook.findOne.mockReturnValue({ + lean: jest.fn().mockResolvedValue({ + _id: validWebhookId, + name: 'Test Hook', + url: 'https://example.com', + enabled: true, + }), + }); + + await getWebhook(req, res); + + expect(res.json).toHaveBeenCalledWith({ + data: expect.objectContaining({ name: 'Test Hook' }), + }); + }); + + test('returns 404 if webhook not found', async () => { + req.params.webhookId = validWebhookId; + Project.findOne.mockResolvedValue({ _id: validProjectId }); + Webhook.findOne.mockReturnValue({ + lean: jest.fn().mockResolvedValue(null), + }); + + await getWebhook(req, res); + + expect(res.status).toHaveBeenCalledWith(404); + expect(res.json).toHaveBeenCalledWith({ error: 'Webhook not found' }); + }); + }); + + describe('updateWebhook', () => { + test('updates webhook successfully', async () => { + req.params.webhookId = validWebhookId; + req.body = { name: 'Updated Name', enabled: false }; + + Project.findOne.mockResolvedValue({ _id: validProjectId }); + updateWebhookSchema.safeParse.mockReturnValue({ + success: true, + data: { name: 'Updated Name', enabled: false }, + }); + Webhook.findOneAndUpdate.mockReturnValue({ + lean: jest.fn().mockResolvedValue({ + _id: validWebhookId, + name: 'Updated Name', + url: 'https://example.com', + enabled: false, + }), + }); + + await updateWebhook(req, res); + + expect(Webhook.findOneAndUpdate).toHaveBeenCalled(); + expect(res.json).toHaveBeenCalledWith( + expect.objectContaining({ + message: 'Webhook updated', + data: expect.objectContaining({ name: 'Updated Name' }), + }) + ); + }); + }); + + describe('deleteWebhook', () => { + test('deletes webhook successfully', async () => { + req.params.webhookId = validWebhookId; + Project.findOne.mockResolvedValue({ _id: validProjectId }); + Webhook.findOneAndDelete.mockResolvedValue({ _id: validWebhookId }); + + await deleteWebhook(req, res); + + expect(Webhook.findOneAndDelete).toHaveBeenCalledWith({ + _id: validWebhookId, + projectId: validProjectId, + }); + expect(res.json).toHaveBeenCalledWith({ message: 'Webhook deleted' }); + }); + + test('returns 404 if webhook not found', async () => { + req.params.webhookId = validWebhookId; + Project.findOne.mockResolvedValue({ _id: validProjectId }); + Webhook.findOneAndDelete.mockResolvedValue(null); + + await deleteWebhook(req, res); + + expect(res.status).toHaveBeenCalledWith(404); + expect(res.json).toHaveBeenCalledWith({ error: 'Webhook not found' }); + }); + }); + + describe('getDeliveries', () => { + test('returns paginated delivery history', async () => { + req.params.webhookId = validWebhookId; + req.query = { limit: '10', page: '1' }; + + Project.findOne.mockResolvedValue({ _id: validProjectId }); + Webhook.findOne.mockResolvedValue({ _id: validWebhookId }); + + const mockDeliveries = [ + { _id: 'd1', event: 'posts.insert', finalStatus: 'delivered' }, + { _id: 'd2', event: 'posts.update', finalStatus: 'failed' }, + ]; + + WebhookDelivery.find.mockReturnValue({ + sort: jest.fn().mockReturnThis(), + skip: jest.fn().mockReturnThis(), + limit: jest.fn().mockReturnThis(), + lean: jest.fn().mockResolvedValue(mockDeliveries), + }); + WebhookDelivery.countDocuments.mockResolvedValue(2); + + await getDeliveries(req, res); + + expect(res.json).toHaveBeenCalledWith({ + data: mockDeliveries, + pagination: expect.objectContaining({ + page: 1, + limit: 10, + total: 2, + }), + }); + }); + }); + + describe('testWebhook', () => { + beforeEach(() => { + global.fetch = jest.fn(); + }); + + afterEach(() => { + delete global.fetch; + }); + + test('sends test payload and returns success', async () => { + req.params.webhookId = validWebhookId; + + Project.findOne.mockResolvedValue({ _id: validProjectId }); + + const mockWebhook = { + _id: validWebhookId, + name: 'Test Webhook', + url: 'https://example.com/webhook', + secret: { encrypted: 'enc', iv: 'iv', tag: 'tag' }, + enabled: true, + }; + Webhook.findOne.mockReturnValue({ + select: jest.fn().mockResolvedValue(mockWebhook), + }); + + global.fetch.mockResolvedValue({ + status: 200, + text: jest.fn().mockResolvedValue('{"received": true}'), + }); + + await testWebhook(req, res); + + expect(global.fetch).toHaveBeenCalledWith( + 'https://example.com/webhook', + expect.objectContaining({ + method: 'POST', + headers: expect.objectContaining({ + 'Content-Type': 'application/json', + 'X-urBackend-Signature': expect.any(String), + 'X-urBackend-Event': 'test.ping', + }), + }) + ); + expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ + success: true, + statusCode: 200, + })); + }); + + test('returns 404 when webhook not found', async () => { + req.params.webhookId = validWebhookId; + + Project.findOne.mockResolvedValue({ _id: validProjectId }); + Webhook.findOne.mockReturnValue({ + select: jest.fn().mockResolvedValue(null), + }); + + await testWebhook(req, res); + + expect(res.status).toHaveBeenCalledWith(404); + expect(res.json).toHaveBeenCalledWith({ error: 'Webhook not found' }); + }); + + test('handles fetch failure gracefully', async () => { + req.params.webhookId = validWebhookId; + + Project.findOne.mockResolvedValue({ _id: validProjectId }); + + const mockWebhook = { + _id: validWebhookId, + name: 'Test Webhook', + url: 'https://example.com/webhook', + secret: { encrypted: 'enc', iv: 'iv', tag: 'tag' }, + enabled: true, + }; + Webhook.findOne.mockReturnValue({ + select: jest.fn().mockResolvedValue(mockWebhook), + }); + + global.fetch.mockRejectedValue(new Error('Network error')); + + await testWebhook(req, res); + + expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ + success: false, + error: 'Network error', + })); + }); + }); +}); diff --git a/apps/dashboard-api/src/app.js b/apps/dashboard-api/src/app.js index 8c1da3584..7d26e86a9 100644 --- a/apps/dashboard-api/src/app.js +++ b/apps/dashboard-api/src/app.js @@ -80,7 +80,7 @@ app.use(capture({ supabaseUrl: process.env.SUPABASE_URL, supabaseKey: process.env.SUPABASE_KEY, bucket: process.env.SUPABASE_BUCKET, - sampleRate: 0.1 + sampleRate: 0 })); @@ -88,9 +88,11 @@ app.use(capture({ const authRoute = require('./routes/auth'); const projectRoute = require('./routes/projects'); const releaseRoute = require('./routes/releases'); +const webhookRoute = require('./routes/webhooks'); app.use('/api/auth', authRoute); app.use('/api/projects', dashboardLimiter, projectRoute); +app.use('/api/projects', dashboardLimiter, webhookRoute); app.use('/api/releases', releaseRoute); diff --git a/apps/dashboard-api/src/controllers/webhook.controller.js b/apps/dashboard-api/src/controllers/webhook.controller.js new file mode 100644 index 000000000..2727a744c --- /dev/null +++ b/apps/dashboard-api/src/controllers/webhook.controller.js @@ -0,0 +1,437 @@ +const mongoose = require("mongoose"); +const { + Webhook, + WebhookDelivery, + Project, + encrypt, + decrypt, + createWebhookSchema, + updateWebhookSchema, + generateSignature, +} = require("@urbackend/common"); +const crypto = require("crypto"); + +// Validate MongoDB ObjectId +const isValidId = (id) => mongoose.Types.ObjectId.isValid(id); + +/** + * Create a new webhook for a project + */ +module.exports.createWebhook = async (req, res) => { + try { + const { projectId } = req.params; + + if (!isValidId(projectId)) { + return res.status(400).json({ error: "Invalid project ID" }); + } + + // Verify project ownership + const project = await Project.findOne({ + _id: projectId, + owner: req.user._id, + }); + if (!project) { + return res.status(404).json({ error: "Project not found" }); + } + + // Validate input + const validation = createWebhookSchema.safeParse(req.body); + if (!validation.success) { + return res.status(400).json({ + error: "Validation failed", + details: validation.error.errors, + }); + } + + const { name, url, secret, events, enabled } = validation.data; + + // Encrypt the secret + const encryptedSecret = encrypt(secret); + + const webhook = await Webhook.create({ + projectId, + name, + url, + secret: encryptedSecret, + events: events || {}, + enabled: enabled !== false, + }); + + // Return without secret + res.status(201).json({ + message: "Webhook created", + data: { + _id: webhook._id, + projectId: webhook.projectId, + name: webhook.name, + url: webhook.url, + events: Object.fromEntries(webhook.events || new Map()), + enabled: webhook.enabled, + createdAt: webhook.createdAt, + }, + }); + } catch (err) { + console.error("[Webhook] Create error:", err); + res.status(500).json({ error: err.message }); + } +}; + +/** + * Get all webhooks for a project + */ +module.exports.getWebhooks = async (req, res) => { + try { + const { projectId } = req.params; + + if (!isValidId(projectId)) { + return res.status(400).json({ error: "Invalid project ID" }); + } + + // Verify project ownership + const project = await Project.findOne({ + _id: projectId, + owner: req.user._id, + }); + if (!project) { + return res.status(404).json({ error: "Project not found" }); + } + + const webhooks = await Webhook.find({ projectId }).lean(); + + // Transform events Map to plain object and exclude secret + const data = webhooks.map((wh) => ({ + _id: wh._id, + projectId: wh.projectId, + name: wh.name, + url: wh.url, + events: wh.events || {}, + enabled: wh.enabled, + createdAt: wh.createdAt, + updatedAt: wh.updatedAt, + })); + + res.json({ data }); + } catch (err) { + console.error("[Webhook] List error:", err); + res.status(500).json({ error: err.message }); + } +}; + +/** + * Get a single webhook + */ +module.exports.getWebhook = async (req, res) => { + try { + const { projectId, webhookId } = req.params; + + if (!isValidId(projectId) || !isValidId(webhookId)) { + return res.status(400).json({ error: "Invalid ID format" }); + } + + // Verify project ownership + const project = await Project.findOne({ + _id: projectId, + owner: req.user._id, + }); + if (!project) { + return res.status(404).json({ error: "Project not found" }); + } + + const webhook = await Webhook.findOne({ + _id: webhookId, + projectId, + }).lean(); + + if (!webhook) { + return res.status(404).json({ error: "Webhook not found" }); + } + + res.json({ + data: { + _id: webhook._id, + projectId: webhook.projectId, + name: webhook.name, + url: webhook.url, + events: webhook.events || {}, + enabled: webhook.enabled, + createdAt: webhook.createdAt, + updatedAt: webhook.updatedAt, + }, + }); + } catch (err) { + console.error("[Webhook] Get error:", err); + res.status(500).json({ error: err.message }); + } +}; + +/** + * Update a webhook + */ +module.exports.updateWebhook = async (req, res) => { + try { + const { projectId, webhookId } = req.params; + + if (!isValidId(projectId) || !isValidId(webhookId)) { + return res.status(400).json({ error: "Invalid ID format" }); + } + + // Verify project ownership + const project = await Project.findOne({ + _id: projectId, + owner: req.user._id, + }); + if (!project) { + return res.status(404).json({ error: "Project not found" }); + } + + // Validate input + const validation = updateWebhookSchema.safeParse(req.body); + if (!validation.success) { + return res.status(400).json({ + error: "Validation failed", + details: validation.error.errors, + }); + } + + const { name, url, secret, events, enabled } = validation.data; + + const updateData = {}; + if (name !== undefined) updateData.name = name; + if (url !== undefined) updateData.url = url; + if (events !== undefined) updateData.events = events; + if (enabled !== undefined) updateData.enabled = enabled; + + // Re-encrypt if secret is being updated + if (secret !== undefined) { + updateData.secret = encrypt(secret); + } + + const webhook = await Webhook.findOneAndUpdate( + { _id: webhookId, projectId }, + { $set: updateData }, + { new: true } + ).lean(); + + if (!webhook) { + return res.status(404).json({ error: "Webhook not found" }); + } + + res.json({ + message: "Webhook updated", + data: { + _id: webhook._id, + projectId: webhook.projectId, + name: webhook.name, + url: webhook.url, + events: webhook.events || {}, + enabled: webhook.enabled, + createdAt: webhook.createdAt, + updatedAt: webhook.updatedAt, + }, + }); + } catch (err) { + console.error("[Webhook] Update error:", err); + res.status(500).json({ error: err.message }); + } +}; + +/** + * Delete a webhook + */ +module.exports.deleteWebhook = async (req, res) => { + try { + const { projectId, webhookId } = req.params; + + if (!isValidId(projectId) || !isValidId(webhookId)) { + return res.status(400).json({ error: "Invalid ID format" }); + } + + // Verify project ownership + const project = await Project.findOne({ + _id: projectId, + owner: req.user._id, + }); + if (!project) { + return res.status(404).json({ error: "Project not found" }); + } + + const webhook = await Webhook.findOneAndDelete({ + _id: webhookId, + projectId, + }); + + if (!webhook) { + return res.status(404).json({ error: "Webhook not found" }); + } + + // Optionally clean up delivery logs (or keep for audit) + // await WebhookDelivery.deleteMany({ webhookId }); + + res.json({ message: "Webhook deleted" }); + } catch (err) { + console.error("[Webhook] Delete error:", err); + res.status(500).json({ error: err.message }); + } +}; + +/** + * Get delivery history for a webhook + */ +module.exports.getDeliveries = async (req, res) => { + try { + const { projectId, webhookId } = req.params; + const { limit = 50, page = 1 } = req.query; + + if (!isValidId(projectId) || !isValidId(webhookId)) { + return res.status(400).json({ error: "Invalid ID format" }); + } + + // Verify project ownership + const project = await Project.findOne({ + _id: projectId, + owner: req.user._id, + }); + if (!project) { + return res.status(404).json({ error: "Project not found" }); + } + + // Verify webhook belongs to project + const webhook = await Webhook.findOne({ _id: webhookId, projectId }); + if (!webhook) { + return res.status(404).json({ error: "Webhook not found" }); + } + + const pageNum = Math.max(1, parseInt(page) || 1); + const limitNum = Math.min(100, Math.max(1, parseInt(limit) || 50)); + const skip = (pageNum - 1) * limitNum; + + const [deliveries, total] = await Promise.all([ + WebhookDelivery.find({ webhookId }) + .sort({ createdAt: -1 }) + .skip(skip) + .limit(limitNum) + .lean(), + WebhookDelivery.countDocuments({ webhookId }), + ]); + + res.json({ + data: deliveries, + pagination: { + page: pageNum, + limit: limitNum, + total, + totalPages: Math.ceil(total / limitNum), + }, + }); + } catch (err) { + console.error("[Webhook] Get deliveries error:", err); + res.status(500).json({ error: err.message }); + } +}; + +/** + * Send a test webhook + */ +module.exports.testWebhook = async (req, res) => { + try { + const { projectId, webhookId } = req.params; + + if (!isValidId(projectId) || !isValidId(webhookId)) { + return res.status(400).json({ error: "Invalid ID format" }); + } + + // Verify project ownership + const project = await Project.findOne({ + _id: projectId, + owner: req.user._id, + }); + if (!project) { + return res.status(404).json({ error: "Project not found" }); + } + + // Load webhook with secret + const webhook = await Webhook.findOne({ _id: webhookId, projectId }).select( + "+secret.encrypted +secret.iv +secret.tag" + ); + + if (!webhook) { + return res.status(404).json({ error: "Webhook not found" }); + } + + // Decrypt secret + let secret; + try { + secret = decrypt(webhook.secret); + if (!secret) throw new Error("Decryption failed"); + } catch (err) { + return res.status(500).json({ error: "Failed to decrypt webhook secret" }); + } + + // Create test payload + const testPayload = { + event: "test.ping", + timestamp: new Date().toISOString(), + projectId: projectId.toString(), + collection: "test", + action: "ping", + documentId: "test-" + crypto.randomUUID(), + data: { + message: "This is a test webhook from urBackend", + triggeredBy: "dashboard", + }, + }; + + const signature = generateSignature(testPayload, secret); + const startTime = Date.now(); + + let statusCode = null; + let responseBody = null; + let error = null; + + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), 10000); // 10s timeout for test + + try { + const response = await fetch(webhook.url, { + method: "POST", + headers: { + "Content-Type": "application/json", + "X-urBackend-Signature": signature, + "X-urBackend-Event": "test.ping", + "X-urBackend-Delivery-Id": "test-" + crypto.randomUUID(), + }, + body: JSON.stringify(testPayload), + signal: controller.signal, + }); + + statusCode = response.status; + + try { + responseBody = await response.text(); + if (responseBody.length > 1024) { + responseBody = responseBody.substring(0, 1024) + "..."; + } + } catch { + responseBody = "[Could not read response body]"; + } + } catch (err) { + error = err.name === "AbortError" ? "Request timeout (10s)" : err.message; + } finally { + clearTimeout(timeout); + } + + const durationMs = Date.now() - startTime; + const success = statusCode >= 200 && statusCode < 300; + + res.json({ + success, + statusCode, + responseBody, + error, + durationMs, + }); + } catch (err) { + console.error("[Webhook] Test error:", err); + res.status(500).json({ error: err.message }); + } +}; diff --git a/apps/dashboard-api/src/routes/webhooks.js b/apps/dashboard-api/src/routes/webhooks.js new file mode 100644 index 000000000..7372bb0be --- /dev/null +++ b/apps/dashboard-api/src/routes/webhooks.js @@ -0,0 +1,37 @@ +const express = require("express"); +const router = express.Router(); +const authMiddleware = require("../middlewares/authMiddleware"); +const { verifyEmail } = require("@urbackend/common"); + +const { + createWebhook, + getWebhooks, + getWebhook, + updateWebhook, + deleteWebhook, + getDeliveries, + testWebhook, +} = require("../controllers/webhook.controller"); + +// Create webhook +router.post("/:projectId/webhooks", authMiddleware, verifyEmail, createWebhook); + +// List all webhooks for a project +router.get("/:projectId/webhooks", authMiddleware, getWebhooks); + +// Get single webhook +router.get("/:projectId/webhooks/:webhookId", authMiddleware, getWebhook); + +// Update webhook +router.patch("/:projectId/webhooks/:webhookId", authMiddleware, verifyEmail, updateWebhook); + +// Delete webhook +router.delete("/:projectId/webhooks/:webhookId", authMiddleware, verifyEmail, deleteWebhook); + +// Get delivery history +router.get("/:projectId/webhooks/:webhookId/deliveries", authMiddleware, getDeliveries); + +// Test webhook +router.post("/:projectId/webhooks/:webhookId/test", authMiddleware, verifyEmail, testWebhook); + +module.exports = router; diff --git a/apps/public-api/src/app.js b/apps/public-api/src/app.js index 36c4702ed..1dfb8dc8b 100644 --- a/apps/public-api/src/app.js +++ b/apps/public-api/src/app.js @@ -21,6 +21,12 @@ const { capture } = require('@kiroo/sdk'); // Initialize Queue Workers const {emailQueue} = require('@urbackend/common'); const {authEmailQueue} = require('@urbackend/common'); +const {initWebhookWorker} = require('@urbackend/common'); + +// Initialize webhook worker +if (process.env.NODE_ENV !== 'test') { + initWebhookWorker(); +} app.use(express.json()); app.use(express.urlencoded({ extended: true })); @@ -38,7 +44,7 @@ app.use(capture({ supabaseUrl: process.env.SUPABASE_URL, supabaseKey: process.env.SUPABASE_KEY, bucket: process.env.SUPABASE_BUCKET, - sampleRate: 0.2 + sampleRate: 0 })); diff --git a/apps/public-api/src/controllers/data.controller.js b/apps/public-api/src/controllers/data.controller.js index 4bdb70adc..25d2d0171 100644 --- a/apps/public-api/src/controllers/data.controller.js +++ b/apps/public-api/src/controllers/data.controller.js @@ -6,6 +6,7 @@ const { getCompiledModel } = require("@urbackend/common"); const { QueryEngine } = require("@urbackend/common"); const { validateData, validateUpdateData } = require("@urbackend/common"); const { performance } = require('perf_hooks'); +const { dispatchWebhooks } = require('../utils/webhookDispatcher'); const isDebug = process.env.DEBUG === 'true'; @@ -64,6 +65,15 @@ module.exports.insertData = async (req, res) => { ); } + // Fire-and-forget webhook dispatch + dispatchWebhooks({ + projectId: project._id, + collection: collectionName, + action: 'insert', + document: result.toObject ? result.toObject() : result, + documentId: result._id, + }); + if (isDebug) console.log(`[DEBUG] insert data took ${(performance.now() - start).toFixed(2)}ms`); res.status(201).json(result); } catch (err) { @@ -201,6 +211,15 @@ module.exports.updateSingleData = async (req, res) => { if (!result) return res.status(404).json({ error: "Document not found." }); + // Fire-and-forget webhook dispatch + dispatchWebhooks({ + projectId: project._id, + collection: collectionName, + action: 'update', + document: result, + documentId: result._id, + }); + res.json({ message: "Updated", data: result }); } catch (err) { console.error(err); @@ -243,6 +262,9 @@ module.exports.deleteSingleDoc = async (req, res) => { if (!docToDelete) return res.status(404).json({ error: "Document not found." }); + // Capture document data before deletion for webhook + const deletedDoc = docToDelete.toObject ? docToDelete.toObject() : { ...docToDelete._doc }; + let docSize = 0; if (!project.resources.db.isExternal) { docSize = Buffer.byteLength(JSON.stringify(docToDelete)); @@ -255,6 +277,15 @@ module.exports.deleteSingleDoc = async (req, res) => { await Project.updateOne({ _id: project._id }, { $set: { databaseUsed } }); } + // Fire-and-forget webhook dispatch + dispatchWebhooks({ + projectId: project._id, + collection: collectionName, + action: 'delete', + document: deletedDoc, + documentId: id, + }); + res.json({ message: "Document deleted", id }); } catch (err) { console.error(err); diff --git a/apps/public-api/src/utils/webhookDispatcher.js b/apps/public-api/src/utils/webhookDispatcher.js new file mode 100644 index 000000000..6fa61cf63 --- /dev/null +++ b/apps/public-api/src/utils/webhookDispatcher.js @@ -0,0 +1,60 @@ +const { Webhook, enqueueWebhookDelivery } = require("@urbackend/common"); + +/** + * Dispatch webhooks for a data operation + * Fire-and-forget: does not block the API response + * + * @param {Object} options + * @param {string} options.projectId - The project ID + * @param {string} options.collection - The collection name + * @param {string} options.action - The action: 'insert', 'update', or 'delete' + * @param {Object} options.document - The document data (after insert/update, or before delete) + * @param {string} options.documentId - The document _id + */ +async function dispatchWebhooks({ projectId, collection, action, document, documentId }) { + try { + // Find all enabled webhooks for this project that listen to this event + const webhooks = await Webhook.find({ + projectId, + enabled: true, + }); + + if (!webhooks.length) return; + + const event = `${collection}.${action}`; + const timestamp = new Date().toISOString(); + + for (const webhook of webhooks) { + // Check if this webhook listens to this collection+action + const collectionEvents = webhook.events?.get(collection); + if (!collectionEvents || !collectionEvents[action]) { + continue; + } + + const payload = { + event, + timestamp, + projectId: projectId.toString(), + collection, + action, + documentId: documentId?.toString() || document?._id?.toString(), + data: document, + }; + + // Enqueue delivery (fire-and-forget) + enqueueWebhookDelivery({ + webhookId: webhook._id, + projectId, + event, + payload, + }).catch((err) => { + console.error(`[Webhook Dispatch] Failed to enqueue: ${err.message}`); + }); + } + } catch (err) { + // Log but don't throw - webhooks should never block the main operation + console.error(`[Webhook Dispatch] Error: ${err.message}`); + } +} + +module.exports = { dispatchWebhooks }; diff --git a/apps/web-dashboard/src/App.jsx b/apps/web-dashboard/src/App.jsx index c39aba150..8cbf25e94 100644 --- a/apps/web-dashboard/src/App.jsx +++ b/apps/web-dashboard/src/App.jsx @@ -23,6 +23,7 @@ import OtpVerification from './pages/OtpVerification'; import ForgotPassword from './pages/ForgotPassword'; import Settings from './pages/Settings'; import ProjectSettings from './pages/ProjectSettings'; +import Webhooks from './pages/Webhooks'; @@ -100,6 +101,8 @@ function App() { } /> + } /> + } /> } /> diff --git a/apps/web-dashboard/src/components/Layout/ProjectNavbar.jsx b/apps/web-dashboard/src/components/Layout/ProjectNavbar.jsx index 0d14be388..80b509b9a 100644 --- a/apps/web-dashboard/src/components/Layout/ProjectNavbar.jsx +++ b/apps/web-dashboard/src/components/Layout/ProjectNavbar.jsx @@ -1,7 +1,7 @@ import { NavLink, useParams, Link } from 'react-router-dom'; import { LayoutDashboard, Database, Shield, HardDrive, Settings, BarChart2, - ArrowLeft + ArrowLeft, Webhook } from 'lucide-react'; function ProjectNavbar() { @@ -35,6 +35,11 @@ function ProjectNavbar() { Auth + `nav-link ${isActive ? 'active' : ''}`}> + + Webhooks + + `nav-link ${isActive ? 'active' : ''}`}> Storage diff --git a/apps/web-dashboard/src/components/Layout/Sidebar.jsx b/apps/web-dashboard/src/components/Layout/Sidebar.jsx index f9e9cc8f1..d8d8ce082 100644 --- a/apps/web-dashboard/src/components/Layout/Sidebar.jsx +++ b/apps/web-dashboard/src/components/Layout/Sidebar.jsx @@ -2,7 +2,7 @@ import { Link, useLocation, useParams } from 'react-router-dom'; import { useAuth } from '../../context/AuthContext'; import { LayoutDashboard, Database, Shield, HardDrive, Settings, BarChart2, - ArrowLeft, FileText, UserCog, LogOut, X, Rocket // Import Rocket + ArrowLeft, FileText, UserCog, LogOut, X, Rocket, Webhook // Import Webhook } from 'lucide-react'; function Sidebar({ logo, isOpen, onClose }) { // Props received @@ -56,6 +56,10 @@ function Sidebar({ logo, isOpen, onClose }) { // Props received Authentication + + Webhooks + + Storage diff --git a/apps/web-dashboard/src/pages/Webhooks.jsx b/apps/web-dashboard/src/pages/Webhooks.jsx new file mode 100644 index 000000000..3a474a9f5 --- /dev/null +++ b/apps/web-dashboard/src/pages/Webhooks.jsx @@ -0,0 +1,626 @@ +import { useState, useEffect } from 'react'; +import { useParams } from 'react-router-dom'; +import api from '../utils/api'; +import toast from 'react-hot-toast'; +import { + Webhook, Plus, Trash2, Edit2, X, Play, CheckCircle, + XCircle, Clock, RefreshCw, Eye, ChevronDown, ChevronUp, Copy +} from 'lucide-react'; + +export default function Webhooks() { + const { projectId } = useParams(); + + const [webhooks, setWebhooks] = useState([]); + const [project, setProject] = useState(null); + const [loading, setLoading] = useState(true); + + // Modal state + const [isModalOpen, setIsModalOpen] = useState(false); + const [editingWebhook, setEditingWebhook] = useState(null); + const [formData, setFormData] = useState({ + name: '', + url: '', + secret: '', + enabled: true, + events: {} + }); + const [isSaving, setIsSaving] = useState(false); + + // Delivery history state + const [deliveriesWebhookId, setDeliveriesWebhookId] = useState(null); + const [deliveries, setDeliveries] = useState([]); + const [loadingDeliveries, setLoadingDeliveries] = useState(false); + const [expandedDelivery, setExpandedDelivery] = useState(null); + + // Test state + const [testingWebhookId, setTestingWebhookId] = useState(null); + const [testResult, setTestResult] = useState(null); + + // Delete confirmation + const [deleteTarget, setDeleteTarget] = useState(null); + + const collections = project?.collections || []; + + useEffect(() => { + fetchData(); + // eslint-disable-next-line react-hooks/exhaustive-deps + }, [projectId]); + + const fetchData = async () => { + try { + const [projRes, webhooksRes] = await Promise.all([ + api.get(`/api/projects/${projectId}`), + api.get(`/api/projects/${projectId}/webhooks`) + ]); + setProject(projRes.data); + setWebhooks(webhooksRes.data.data || []); + } catch (err) { + console.error(err); + toast.error('Failed to load webhooks'); + } finally { + setLoading(false); + } + }; + + const openCreateModal = () => { + setEditingWebhook(null); + setFormData({ + name: '', + url: '', + secret: generateSecret(), + enabled: true, + events: {} + }); + setIsModalOpen(true); + }; + + const openEditModal = (webhook) => { + setEditingWebhook(webhook); + setFormData({ + name: webhook.name, + url: webhook.url, + secret: '', // Don't show existing secret + enabled: webhook.enabled, + events: webhook.events || {} + }); + setIsModalOpen(true); + }; + + const closeModal = () => { + setIsModalOpen(false); + setEditingWebhook(null); + setFormData({ name: '', url: '', secret: '', enabled: true, events: {} }); + }; + + const generateSecret = () => { + const chars = 'ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789'; + let result = 'whsec_'; + for (let i = 0; i < 32; i++) { + result += chars.charAt(Math.floor(Math.random() * chars.length)); + } + return result; + }; + + const handleEventToggle = (collection, action) => { + setFormData(prev => { + const events = { ...prev.events }; + if (!events[collection]) { + events[collection] = { insert: false, update: false, delete: false }; + } + events[collection] = { ...events[collection], [action]: !events[collection][action] }; + return { ...prev, events }; + }); + }; + + const handleSubmit = async (e) => { + e.preventDefault(); + setIsSaving(true); + + try { + const payload = { + name: formData.name, + url: formData.url, + events: formData.events, + enabled: formData.enabled + }; + + // Only include secret if provided (required for create, optional for update) + if (formData.secret) { + payload.secret = formData.secret; + } + + if (editingWebhook) { + await api.patch(`/api/projects/${projectId}/webhooks/${editingWebhook._id}`, payload); + toast.success('Webhook updated'); + } else { + await api.post(`/api/projects/${projectId}/webhooks`, payload); + toast.success('Webhook created'); + } + + closeModal(); + fetchData(); + } catch (err) { + const msg = err.response?.data?.error || err.response?.data?.details?.[0]?.message || 'Failed to save webhook'; + toast.error(msg); + } finally { + setIsSaving(false); + } + }; + + const handleDelete = async () => { + if (!deleteTarget) return; + try { + await api.delete(`/api/projects/${projectId}/webhooks/${deleteTarget._id}`); + toast.success('Webhook deleted'); + setDeleteTarget(null); + fetchData(); + } catch { + toast.error('Failed to delete webhook'); + } + }; + + const handleTest = async (webhook) => { + setTestingWebhookId(webhook._id); + setTestResult(null); + try { + const res = await api.post(`/api/projects/${projectId}/webhooks/${webhook._id}/test`); + setTestResult({ webhookId: webhook._id, ...res.data }); + } catch (err) { + setTestResult({ + webhookId: webhook._id, + success: false, + error: err.response?.data?.error || 'Test failed' + }); + } finally { + setTestingWebhookId(null); + } + }; + + const openDeliveries = async (webhook) => { + setDeliveriesWebhookId(webhook._id); + setLoadingDeliveries(true); + setDeliveries([]); + try { + const res = await api.get(`/api/projects/${projectId}/webhooks/${webhook._id}/deliveries?limit=50`); + setDeliveries(res.data.data || []); + } catch { + toast.error('Failed to load delivery history'); + } finally { + setLoadingDeliveries(false); + } + }; + + const closeDeliveries = () => { + setDeliveriesWebhookId(null); + setDeliveries([]); + setExpandedDelivery(null); + }; + + const copyToClipboard = (text) => { + navigator.clipboard.writeText(text); + toast.success('Copied to clipboard'); + }; + + const getStatusIcon = (status) => { + switch (status) { + case 'delivered': + return ; + case 'failed': + return ; + default: + return ; + } + }; + + const formatDate = (date) => { + return new Date(date).toLocaleString(); + }; + + if (loading) return
Loading...
; + + return ( +
+ {/* Header */} +
+
+

+ Webhooks +

+

+ Send HTTP callbacks when data changes in your collections +

+
+ +
+ + {/* Webhooks List */} + {webhooks.length === 0 ? ( +
+ +

No webhooks configured

+

+ Create a webhook to receive notifications when data changes +

+ +
+ ) : ( +
+ {webhooks.map((webhook) => ( +
+
+
+
+
+ +
+
+

+ {webhook.name} + + {webhook.enabled ? 'Active' : 'Disabled'} + +

+
+
+ +
+ POST + + {webhook.url} + +
+ +
+ Triggers on: + {Object.entries(webhook.events || {}).map(([coll, events]) => + Object.entries(events).filter(([, v]) => v).map(([action]) => ( + + {coll} {action} + + )) + )} +
+
+ +
+ + +
+ + +
+
+
+ {/* Test Result */} + {testResult && testResult.webhookId === webhook._id && ( +
+
+ {testResult.success ? : } + {testResult.success ? 'Test successful' : 'Test failed'} + {testResult.statusCode && ({testResult.statusCode})} + {testResult.durationMs && {testResult.durationMs}ms} +
+ {testResult.error &&

{testResult.error}

} + {testResult.responseBody && ( +
+                      {testResult.responseBody}
+                    
+ )} +
+ )} +
+ ))} +
+ )} + + {/* Create/Edit Modal */} + {isModalOpen && ( +
+
e.stopPropagation()} style={{ maxWidth: '600px', maxHeight: '90vh', overflow: 'auto' }}> +
+

{editingWebhook ? 'Edit Webhook' : 'Create Webhook'}

+ +
+
+
+ {/* Name */} +
+ + setFormData({ ...formData, name: e.target.value })} + required + /> +
+ + {/* URL */} +
+ + setFormData({ ...formData, url: e.target.value })} + required + /> + Must use HTTPS (or http://localhost for development) +
+ + {/* Secret */} +
+ +
+ setFormData({ ...formData, secret: e.target.value })} + required={!editingWebhook} + style={{ flex: 1 }} + /> + + {formData.secret && ( + + )} +
+ Used for HMAC-SHA256 signature verification +
+ + {/* Enabled */} +
+ setFormData({ ...formData, enabled: e.target.checked })} + /> + +
+ + {/* Events */} +
+ + {collections.length === 0 ? ( +

No collections available. Create a collection first.

+ ) : ( +
+ {collections.map((coll) => ( +
+ {coll.name} +
+ {['insert', 'update', 'delete'].map((action) => ( + + ))} +
+
+ ))} +
+ )} +
+
+ +
+ + +
+
+
+
+ )} + + {/* Delivery History Modal */} + {deliveriesWebhookId && ( +
+
e.stopPropagation()} style={{ maxWidth: '800px', maxHeight: '90vh', overflow: 'auto' }}> +
+

Delivery History

+ +
+
+ {loadingDeliveries ? ( +

Loading...

+ ) : deliveries.length === 0 ? ( +

No deliveries yet

+ ) : ( +
+ {deliveries.map((delivery) => ( +
+
setExpandedDelivery(expandedDelivery === delivery._id ? null : delivery._id)} + > + {getStatusIcon(delivery.finalStatus)} + {delivery.event} + + {delivery.attempts?.length || 0} attempt{(delivery.attempts?.length || 0) !== 1 ? 's' : ''} + + + {formatDate(delivery.createdAt)} + + {expandedDelivery === delivery._id ? : } +
+ {expandedDelivery === delivery._id && ( +
+
+ Payload: +
+                              {JSON.stringify(delivery.payload, null, 2)}
+                            
+
+
+ Attempts: +
+ {(delivery.attempts || []).map((attempt, idx) => ( +
+
+ #{attempt.attemptNumber} + {getStatusIcon(attempt.status)} + {attempt.statusCode && Status: {attempt.statusCode}} + {attempt.durationMs && {attempt.durationMs}ms} + {formatDate(attempt.attemptedAt)} +
+ {attempt.error &&

{attempt.error}

} + {attempt.responseBody && ( +
+                                      {attempt.responseBody}
+                                    
+ )} +
+ ))} +
+
+
+ )} +
+ ))} +
+ )} +
+
+
+ )} + + {/* Delete Confirmation Modal */} + {deleteTarget && ( +
setDeleteTarget(null)}> +
e.stopPropagation()} style={{ maxWidth: '400px' }}> +
+

Delete Webhook

+ +
+
+

Are you sure you want to delete {deleteTarget.name}?

+

This action cannot be undone.

+
+
+ + +
+
+
+ )} + + +
+ ); +} diff --git a/packages/common/src/index.js b/packages/common/src/index.js index d6b8d732b..0b2c1cc91 100644 --- a/packages/common/src/index.js +++ b/packages/common/src/index.js @@ -18,10 +18,18 @@ const Project = require("./models/Project"); const Release = require("./models/Release"); const Log = require("./models/Log"); const Otp = require("./models/otp"); +const Webhook = require("./models/Webhook"); +const WebhookDelivery = require("./models/WebhookDelivery"); // Queues const { authEmailQueue } = require("./queues/authEmailQueue"); const { emailQueue } = require("./queues/emailQueue"); +const { + webhookQueue, + enqueueWebhookDelivery, + initWebhookWorker, + generateSignature, +} = require("./queues/webhookQueue"); // Middleware const checkAuthEnabled = require('./middleware/checkAuthEnabled') @@ -50,6 +58,8 @@ const { userSignupSchema, updateExternalConfigSchema, updateAuthProvidersSchema, + createWebhookSchema, + updateWebhookSchema, } = require("./utils/input.validation"); const { garbageCollect, storageGarbageCollect } = require("./utils/GC"); const { generateApiKey, hashApiKey } = require("./utils/api"); @@ -81,8 +91,14 @@ module.exports = { Release, Log, Otp, + Webhook, + WebhookDelivery, authEmailQueue, emailQueue, + webhookQueue, + enqueueWebhookDelivery, + initWebhookWorker, + generateSignature, sendOtp, sendReleaseEmail, sendAuthOtpEmail, @@ -99,6 +115,8 @@ module.exports = { sanitize, updateExternalConfigSchema, updateAuthProvidersSchema, + createWebhookSchema, + updateWebhookSchema, garbageCollect, storageGarbageCollect, generateApiKey, diff --git a/packages/common/src/models/Webhook.js b/packages/common/src/models/Webhook.js new file mode 100644 index 000000000..5dfee9494 --- /dev/null +++ b/packages/common/src/models/Webhook.js @@ -0,0 +1,60 @@ +const mongoose = require("mongoose"); + +const resourceConfigSchema = new mongoose.Schema( + { + encrypted: { type: String, select: false }, + iv: { type: String, select: false }, + tag: { type: String, select: false }, + }, + { _id: false } +); + +const eventConfigSchema = new mongoose.Schema( + { + insert: { type: Boolean, default: false }, + update: { type: Boolean, default: false }, + delete: { type: Boolean, default: false }, + }, + { _id: false } +); + +const webhookSchema = new mongoose.Schema( + { + projectId: { + type: mongoose.Schema.Types.ObjectId, + ref: "Project", + required: true, + index: true, + }, + name: { + type: String, + required: true, + trim: true, + maxlength: 100, + }, + url: { + type: String, + required: true, + trim: true, + maxlength: 2048, + }, + secret: { + type: resourceConfigSchema, + required: true, + }, + events: { + type: Map, + of: eventConfigSchema, + default: {}, + }, + enabled: { + type: Boolean, + default: true, + }, + }, + { timestamps: true } +); + +webhookSchema.index({ projectId: 1, enabled: 1 }); + +module.exports = mongoose.model("Webhook", webhookSchema); diff --git a/packages/common/src/models/WebhookDelivery.js b/packages/common/src/models/WebhookDelivery.js new file mode 100644 index 000000000..3cc23924b --- /dev/null +++ b/packages/common/src/models/WebhookDelivery.js @@ -0,0 +1,62 @@ +const mongoose = require("mongoose"); + +const attemptSchema = new mongoose.Schema( + { + attemptNumber: { type: Number, required: true }, + status: { + type: String, + enum: ["pending", "success", "failed"], + default: "pending", + }, + statusCode: { type: Number }, + responseBody: { type: String, maxlength: 1024 }, + error: { type: String, maxlength: 500 }, + attemptedAt: { type: Date, default: Date.now }, + durationMs: { type: Number }, + }, + { _id: false } +); + +const webhookDeliverySchema = new mongoose.Schema( + { + webhookId: { + type: mongoose.Schema.Types.ObjectId, + ref: "Webhook", + required: true, + index: true, + }, + projectId: { + type: mongoose.Schema.Types.ObjectId, + ref: "Project", + required: true, + index: true, + }, + event: { + type: String, + required: true, + }, + payload: { + type: mongoose.Schema.Types.Mixed, + }, + attempts: { + type: [attemptSchema], + default: [], + }, + nextRetryAt: { + type: Date, + default: null, + }, + finalStatus: { + type: String, + enum: ["pending", "delivered", "failed"], + default: "pending", + }, + }, + { timestamps: true } +); + +webhookDeliverySchema.index({ projectId: 1, webhookId: 1 }); +webhookDeliverySchema.index({ finalStatus: 1, nextRetryAt: 1 }); +webhookDeliverySchema.index({ createdAt: -1 }); + +module.exports = mongoose.model("WebhookDelivery", webhookDeliverySchema); diff --git a/packages/common/src/queues/webhookQueue.js b/packages/common/src/queues/webhookQueue.js new file mode 100644 index 000000000..b7663c2c2 --- /dev/null +++ b/packages/common/src/queues/webhookQueue.js @@ -0,0 +1,292 @@ +const { Queue, Worker } = require("bullmq"); +const connection = require("../config/redis"); +const crypto = require("crypto"); +const WebhookDelivery = require("../models/WebhookDelivery"); +const Webhook = require("../models/Webhook"); +const { decrypt } = require("../utils/encryption"); + +// Exponential backoff delays in milliseconds: 1min, 5min, 15min, 1hr, 4hr +const RETRY_DELAYS = [ + 60 * 1000, + 5 * 60 * 1000, + 15 * 60 * 1000, + 60 * 60 * 1000, + 4 * 60 * 60 * 1000, +]; +const MAX_ATTEMPTS = 5; + +const webhookQueue = new Queue("webhook-delivery-queue", { connection }); + +/** + * Generate HMAC-SHA256 signature for webhook payload + */ +function generateSignature(payload, secret) { + const hmac = crypto.createHmac("sha256", secret); + hmac.update(JSON.stringify(payload)); + return `sha256=${hmac.digest("hex")}`; +} + +/** + * Truncate string to maxLength + */ +function truncate(str, maxLength = 1024) { + if (!str || typeof str !== "string") return str; + if (str.length <= maxLength) return str; + const ellipsis = "..."; + if (maxLength <= ellipsis.length) return ellipsis.substring(0, maxLength); + return str.substring(0, maxLength - ellipsis.length) + ellipsis; +} + +/** + * Create initial webhook delivery job + */ +async function enqueueWebhookDelivery({ + webhookId, + projectId, + event, + payload, +}) { + // Create delivery record + const delivery = await WebhookDelivery.create({ + webhookId, + projectId, + event, + payload, + finalStatus: "pending", + attempts: [], + }); + + // Add job to queue + await webhookQueue.add( + "deliver", + { + deliveryId: delivery._id.toString(), + webhookId: webhookId.toString(), + attemptNumber: 1, + }, + { + attempts: 1, // BullMQ handles single attempt; we manage retries ourselves + removeOnComplete: true, + removeOnFail: { count: 100 }, // Keep last 100 failed jobs for debugging + } + ); + + return delivery; +} + +/** + * Initialize the webhook worker + * Call this once during app startup + */ +function initWebhookWorker() { + const worker = new Worker( + "webhook-delivery-queue", + async (job) => { + const { deliveryId, webhookId, attemptNumber } = job.data; + + try { + + const delivery = await WebhookDelivery.findById(deliveryId); + if (!delivery) { + console.error(`[Webhook] Delivery ${deliveryId} not found`); + return; + } + + if (delivery.finalStatus !== "pending") { + console.log( + `[Webhook] Delivery ${deliveryId} already ${delivery.finalStatus}, skipping` + ); + return; + } + + // Load webhook with secret + const webhook = await Webhook.findById(webhookId).select( + "+secret.encrypted +secret.iv +secret.tag" + ); + if (!webhook || !webhook.enabled) { + console.log(`[Webhook] Webhook ${webhookId} disabled or not found`); + await WebhookDelivery.findByIdAndUpdate(deliveryId, { + finalStatus: "failed", + $push: { + attempts: { + attemptNumber, + status: "failed", + error: "Webhook disabled or not found", + attemptedAt: new Date(), + }, + }, + }); + return; + } + + // Decrypt secret + let secret; + try { + secret = decrypt(webhook.secret); + if (!secret) throw new Error("Decryption returned null"); + } catch (err) { + console.error(`[Webhook] Failed to decrypt secret: ${err.message}`); + await WebhookDelivery.findByIdAndUpdate(deliveryId, { + finalStatus: "failed", + $push: { + attempts: { + attemptNumber, + status: "failed", + error: "Secret decryption failed", + attemptedAt: new Date(), + }, + }, + }); + return; + } + + // Prepare payload and signature + const signature = generateSignature(delivery.payload, secret); + const startTime = Date.now(); + + let statusCode = null; + let responseBody = null; + let error = null; + let success = false; + + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), 30000); // 30s timeout + + try { + const response = await fetch(webhook.url, { + method: "POST", + headers: { + "Content-Type": "application/json", + "X-urBackend-Signature": signature, + "X-urBackend-Event": delivery.event, + "X-urBackend-Delivery-Id": deliveryId, + }, + body: JSON.stringify(delivery.payload), + signal: controller.signal, + }); + + statusCode = response.status; + + try { + responseBody = await response.text(); + responseBody = truncate(responseBody, 1024); + } catch { + responseBody = "[Could not read response body]"; + } + + success = statusCode >= 200 && statusCode < 300; + } catch (err) { + error = err.name === "AbortError" ? "Request timeout (30s)" : err.message; + } finally { + clearTimeout(timeout); + } + + const durationMs = Date.now() - startTime; + + // Update delivery with attempt result + const attemptRecord = { + attemptNumber, + status: success ? "success" : "failed", + statusCode, + responseBody, + error, + attemptedAt: new Date(), + durationMs, + }; + + // Determine if we should retry + const is4xx = statusCode >= 400 && statusCode < 500; + const shouldRetry = !success && !is4xx && attemptNumber <= MAX_ATTEMPTS; + + if (success) { + await WebhookDelivery.findByIdAndUpdate(deliveryId, { + finalStatus: "delivered", + nextRetryAt: null, + $push: { attempts: attemptRecord }, + }); + console.log( + `[Webhook] Delivery ${deliveryId} succeeded on attempt ${attemptNumber}` + ); + } else if (shouldRetry) { + const nextDelay = RETRY_DELAYS[attemptNumber - 1] || RETRY_DELAYS[RETRY_DELAYS.length - 1]; + const nextRetryAt = new Date(Date.now() + nextDelay); + + await WebhookDelivery.findByIdAndUpdate(deliveryId, { + nextRetryAt, + $push: { attempts: attemptRecord }, + }); + + // Schedule retry job with delay + await webhookQueue.add( + "deliver", + { + deliveryId, + webhookId, + attemptNumber: attemptNumber + 1, + }, + { + delay: nextDelay, + attempts: 1, + removeOnComplete: true, + removeOnFail: { count: 100 }, // Keep last 100 failed jobs for debugging + } + ); + + console.log( + `[Webhook] Delivery ${deliveryId} failed attempt ${attemptNumber}, retrying in ${nextDelay / 1000}s` + ); + } else { + // Final failure + await WebhookDelivery.findByIdAndUpdate(deliveryId, { + finalStatus: "failed", + nextRetryAt: null, + $push: { attempts: attemptRecord }, + }); + console.log( + `[Webhook] Delivery ${deliveryId} permanently failed after ${attemptNumber} attempts` + + (is4xx ? ` (4xx response: ${statusCode})` : "") + ); + } + } catch (err) { + console.error(`[Webhook] Unexpected worker error for delivery ${deliveryId}:`, err.message); + try { + await WebhookDelivery.findByIdAndUpdate(deliveryId, { + finalStatus: "failed", + $push: { + attempts: { + attemptNumber, + status: "failed", + error: `Worker error: ${err.message}`.substring(0, 500), + attemptedAt: new Date(), + }, + }, + }); + } catch (updateErr) { + console.error(`[Webhook] Failed to update delivery ${deliveryId} after error:`, updateErr.message); + } + } + }, + { + connection, + concurrency: 5, + } + ); + + worker.on("completed", (job) => { + console.log(`[Webhook Worker] Job ${job.id} processed`); + }); + + worker.on("failed", (job, err) => { + console.error(`[Webhook Worker] Job ${job?.id} error:`, err.message); + }); + + console.log("[Webhook] Worker initialized"); + return worker; +} + +module.exports = { + webhookQueue, + enqueueWebhookDelivery, + initWebhookWorker, + generateSignature, +}; diff --git a/packages/common/src/utils/input.validation.js b/packages/common/src/utils/input.validation.js index 192da89fe..a45efbd3f 100644 --- a/packages/common/src/utils/input.validation.js +++ b/packages/common/src/utils/input.validation.js @@ -397,3 +397,64 @@ module.exports.userSignupSchema = z.object({ .min(6, { message: "Password must be at least 6 characters." }) .max(100, { message: "Password is too long." }), }); + +// Webhook event config schema for per-collection events +const webhookEventConfigSchema = z.object({ + insert: z.boolean().optional(), + update: z.boolean().optional(), + delete: z.boolean().optional(), +}); + +// URL validation: HTTPS required (or http://localhost for dev) +const webhookUrlSchema = z + .string() + .min(1, "Webhook URL is required") + .max(2048, "URL is too long") + .url("Invalid URL format") + .refine( + (url) => { + try { + const parsed = new URL(url); + return ( + parsed.protocol === "https:" || + (parsed.protocol === "http:" && parsed.hostname === "localhost") + ); + } catch { + return false; + } + }, + "Webhook URL must use HTTPS (or http://localhost for local development)" + ); + +module.exports.createWebhookSchema = z.object({ + name: z + .string() + .min(1, "Webhook name is required") + .max(100, "Webhook name is too long"), + url: webhookUrlSchema, + secret: z + .string() + .min(16, "Secret must be at least 16 characters") + .max(256, "Secret is too long"), + events: z.record(z.string(), webhookEventConfigSchema).optional(), + enabled: z.boolean().optional(), +}); + +module.exports.updateWebhookSchema = z.object({ + name: z + .string() + .min(1, "Webhook name is required") + .max(100, "Webhook name is too long") + .optional(), + url: webhookUrlSchema.optional(), + secret: z + .string() + .min(16, "Secret must be at least 16 characters") + .max(256, "Secret is too long") + .optional(), + events: z.record(z.string(), webhookEventConfigSchema).optional(), + enabled: z.boolean().optional(), +}).refine( + (data) => Object.keys(data).length > 0, + { message: "At least one field must be provided for update." } +);