Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions src/db/repositories/webhookEndpointRepository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,17 @@ export class WebhookEndpointRepository {
return result.rows[0] ? result.rows[0].cnt : 0;
}

async findCompletedDeliveryByEventId(endpointId: string, eventId: string): Promise<WebhookDelivery | null> {
const result: QueryResult<WebhookDelivery> = await this.db.query(
`SELECT * FROM webhook_deliveries
WHERE endpoint_id = $1 AND status = 'completed' AND payload->>'id' = $2
ORDER BY updated_at DESC
LIMIT 1`,
[endpointId, eventId]
);
return result.rows[0] ? this.mapDelivery(result.rows[0]) : null;
}

private map(row: WebhookEndpoint): WebhookEndpoint {
return {
id: row.id,
Expand Down
5 changes: 5 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ import { EmailDeliverabilityRepository } from "./db/repositories/emailDeliverabi
import { createEmailWebhooksRouter } from "./routes/emailWebhooks";
import { createAdminRouter } from "./routes/admin";
import { createAdminLedgerExportRouter } from "./routes/adminLedgerExport";
import { createAdminWebhookRouter } from "./routes/adminWebhooks";
import { AccountingLedgerService } from "./services/accountingLedgerService";
import { DistributionRepository } from "./db/repositories/distributionRepository";
import { createAdminKycRiskTierRouter } from "./routes/adminKycRiskTier";
Expand Down Expand Up @@ -704,6 +705,10 @@ export function createApp(dependencies: AppDependencies = {}): express.Express {
apiRouter.use("/admin", createAdminRouter(auditLogRepo, retentionLabelService));
apiRouter.use("/admin", createAdminKycRiskTierRouter(pool, amlAuditRepo));

// Mount admin webhook dead-letter routes
const webhookEndpointRepo = new WebhookEndpointRepository(pool);
apiRouter.use("/admin/webhooks", createAdminWebhookRouter({ webhookEndpointRepo }));

// Mount admin ledger double-entry export (RBAC + audited)
apiRouter.use(
"/admin/ledger",
Expand Down
122 changes: 64 additions & 58 deletions src/routes/adminWebhooks.ts
Original file line number Diff line number Diff line change
@@ -1,93 +1,99 @@
import { Router, Request, Response, RequestHandler } from 'express';
import { Router, Request, Response, NextFunction } from 'express';
import { WebhookEndpointRepository } from '../db/repositories/webhookEndpointRepository';
import { WebhookQueue } from '../index';
import { Errors } from '../lib/errors';
import { requireAdmin, AuthenticatedRequest } from '../middleware/auth';

export function createAdminWebhooksRouter(opts: {
repo: WebhookEndpointRepository;
requireAuth: RequestHandler;
}) {
const { repo, requireAuth } = opts;
interface AdminWebhookRouterDependencies {
webhookEndpointRepo: WebhookEndpointRepository;
}

export function createAdminWebhookRouter(deps: AdminWebhookRouterDependencies): Router {
const { webhookEndpointRepo } = deps;
const router = Router();

// Ensure admin role
const requireAdmin: RequestHandler = (req, _res, next) => {
const user = (req as any).user;
if (!user || user.role !== 'admin') {
next({ status: 403, message: 'Forbidden: admin role required' });
return;
}
next();
};
// Apply admin auth to all routes in this router
router.use(requireAdmin);

// List recent dead-letter deliveries for an endpoint
router.get('/:endpointId/dead-letters', requireAuth, requireAdmin, async (req: Request, res: Response, next) => {
router.get('/:endpointId/dead-letters', async (req: Request, res: Response, next) => {
try {
const endpointId = req.params.endpointId;
const limitRaw = parseInt(String(req.query.limit || '50'), 10) || 50;
const pageRaw = parseInt(String(req.query.page || '0'), 10) || 0;
const limit = Math.min(parseInt(req.query.limit as string, 10) || 50, 100);
const offset = Math.max(parseInt(req.query.offset as string, 10) || 0, 0);

const limit = Math.min(Math.max(1, limitRaw), 100); // guard: 1..100
const offset = Math.max(0, pageRaw) * limit;
const endpoint = await webhookEndpointRepo.findById(endpointId);
if (!endpoint) {
return next(Errors.notFound('Webhook endpoint not found'));
}

const [items, total] = await Promise.all([
repo.listDeadLettersByEndpoint(endpointId, limit, offset),
repo.countDeadLettersByEndpoint(endpointId),
]);
const deadLetters = await webhookEndpointRepo.listDeadLettersByEndpoint(endpointId, limit, offset);
const total = await webhookEndpointRepo.countDeadLettersByEndpoint(endpointId);

res.json({ total, limit, page: pageRaw, items });
} catch (err) {
next(err);
res.status(200).json({
endpointId,
deadLetters,
pagination: { limit, offset, total },
});
} catch (error) {
next(error);
}
});

// Replay a dead-letter delivery idempotently by resetting its status to pending
router.post('/dead-letters/:id/replay', requireAuth, requireAdmin, async (req: Request, res: Response, next) => {
router.post('/dead-letters/:id/replay', async (req: Request, res: Response, next) => {
try {
const id = req.params.id;
const delivery = await repo.findDeliveryById(id);
const delivery = await webhookEndpointRepo.findDeliveryById(id);

if (!delivery) {
res.status(404).json({ error: 'Not found' });
return;
return next(Errors.notFound('Delivery not found'));
}

if (delivery.status !== 'dead_letter') {
res.status(400).json({ error: 'Delivery is not dead-lettered' });
return;
return next(Errors.badRequest('Only dead-letter deliveries can be replayed'));
}

const endpoint = await webhookEndpointRepo.findById(delivery.endpoint_id);
if (!endpoint) {
return next(Errors.notFound('Associated webhook endpoint not found'));
}

// Check for existing delivery with same event_id to ensure idempotency
const eventId = (delivery.payload as any)?.id;
if (eventId) {
// Simple idempotency: check if a completed delivery with the same event_id exists
// In production, you might want a more robust approach with a dedicated idempotency table
const existingCompleted = await webhookEndpointRepo.findCompletedDeliveryByEventId(delivery.endpoint_id, eventId);
if (existingCompleted) {
return res.status(200).json({
message: 'Delivery already completed for this event_id',
deliveryId: existingCompleted.id,
status: 'already_completed',
});
}
}

// Idempotent replay: update existing delivery back to pending and clear errors
const updated = await repo.updateDelivery(id, {
// Reset to pending and re-enqueue
await webhookEndpointRepo.updateDelivery(id, {
status: 'pending',
attempts: 0,
last_error: null,
next_retry_at: null,
});

// Trigger immediate re-processing if runtime queue is available
(async () => {
try {
// dynamic import to avoid top-level circular deps
// eslint-disable-next-line @typescript-eslint/no-var-requires
const idx = require('../index');
if (idx && idx.WebhookQueue && typeof idx.WebhookQueue.processDelivery === 'function') {
const endpoint = await repo.findById(delivery.endpoint_id);
if (endpoint) {
void idx.WebhookQueue.processDelivery(endpoint.url, delivery.payload, delivery.id);
}
}
} catch (err) {
// non-fatal: queue may live in a different process
// eslint-disable-next-line no-console
console.warn('[adminWebhooks] Could not trigger immediate replay:', err);
}
})();
// Re-enqueue via WebhookQueue
void WebhookQueue.processDelivery(endpoint.url, delivery.payload, id);

res.json({ success: true, id: updated.id });
} catch (err) {
next(err);
res.status(200).json({
message: 'Dead-letter delivery replayed successfully',
deliveryId: id,
status: 'replayed',
});
} catch (error) {
next(error);
}
});

return router;
}

export default createAdminWebhooksRouter;
Loading