"use strict"; var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) { if (k2 === undefined) k2 = k; var desc = Object.getOwnPropertyDescriptor(m, k); if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) { desc = { enumerable: true, get: function() { return m[k]; } }; } Object.defineProperty(o, k2, desc); }) : (function(o, m, k, k2) { if (k2 === undefined) k2 = k; o[k2] = m[k]; })); var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) { Object.defineProperty(o, "default", { enumerable: true, value: v }); }) : function(o, v) { o["default"] = v; }); var __decorate = (this && this.__decorate) || function (decorators, target, key, desc) { var c = arguments.length, r = c < 3 ? target : desc === null ? desc = Object.getOwnPropertyDescriptor(target, key) : desc, d; if (typeof Reflect === "object" && typeof Reflect.decorate === "function") r = Reflect.decorate(decorators, target, key, desc); else for (var i = decorators.length - 1; i >= 0; i--) if (d = decorators[i]) r = (c < 3 ? d(r) : c > 3 ? d(target, key, r) : d(target, key)) || r; return c > 3 && r && Object.defineProperty(target, key, r), r; }; var __importStar = (this && this.__importStar) || (function () { var ownKeys = function(o) { ownKeys = Object.getOwnPropertyNames || function (o) { var ar = []; for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k; return ar; }; return ownKeys(o); }; return function (mod) { if (mod && mod.__esModule) return mod; var result = {}; if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]); __setModuleDefault(result, mod); return result; }; })(); var __metadata = (this && this.__metadata) || function (k, v) { if (typeof Reflect === "object" && typeof Reflect.metadata === "function") return Reflect.metadata(k, v); }; var __param = (this && this.__param) || function (paramIndex, decorator) { return function (target, key) { decorator(target, key, paramIndex); } }; var WebhookProcessor_1; Object.defineProperty(exports, "__esModule", { value: true }); exports.WebhookProcessor = void 0; const bullmq_1 = require("@nestjs/bullmq"); const common_1 = require("@nestjs/common"); const typeorm_1 = require("@nestjs/typeorm"); const typeorm_2 = require("typeorm"); const bullmq_2 = require("bullmq"); const crypto = __importStar(require("crypto")); const entities_1 = require("../entities"); const RETRY_DELAYS = [ 60 * 1000, 5 * 60 * 1000, 30 * 60 * 1000, 2 * 60 * 60 * 1000, 6 * 60 * 60 * 1000, ]; let WebhookProcessor = WebhookProcessor_1 = class WebhookProcessor extends bullmq_1.WorkerHost { constructor(deliveryRepo) { super(); this.deliveryRepo = deliveryRepo; this.logger = new common_1.Logger(WebhookProcessor_1.name); } async process(job) { const { deliveryId, url, secret, headers, eventType, payload } = job.data; this.logger.log(`Processing webhook delivery: ${deliveryId}`); const delivery = await this.deliveryRepo.findOne({ where: { id: deliveryId }, }); if (!delivery) { this.logger.warn(`Delivery not found: ${deliveryId}`); return; } try { const timestamp = Date.now(); const body = JSON.stringify(payload); const signature = crypto .createHmac('sha256', secret) .update(`${timestamp}.${body}`) .digest('hex'); const signatureHeader = `t=${timestamp},v1=${signature}`; const controller = new AbortController(); const timeout = setTimeout(() => controller.abort(), 30000); const response = await fetch(url, { method: 'POST', headers: { 'Content-Type': 'application/json', 'X-Webhook-Signature': signatureHeader, 'X-Webhook-Id': delivery.webhookId, 'X-Webhook-Event': eventType, 'X-Webhook-Timestamp': timestamp.toString(), 'X-Webhook-Delivery': deliveryId, ...headers, }, body, signal: controller.signal, redirect: 'error', }); clearTimeout(timeout); const responseBody = await response.text().catch(() => ''); const responseHeaders = {}; response.headers.forEach((value, key) => { responseHeaders[key] = value; }); delivery.responseStatus = response.status; delivery.responseBody = responseBody.substring(0, 10000); delivery.responseHeaders = responseHeaders; if (response.ok) { delivery.status = entities_1.DeliveryStatus.DELIVERED; delivery.deliveredAt = new Date(); delivery.completedAt = new Date(); this.logger.log(`Webhook delivered successfully: ${deliveryId}`); } else { await this.handleFailure(delivery, `HTTP ${response.status}: ${responseBody.substring(0, 500)}`); } await this.deliveryRepo.save(delivery); } catch (error) { delivery.responseStatus = undefined; delivery.responseBody = undefined; await this.handleFailure(delivery, error.message || 'Unknown error'); await this.deliveryRepo.save(delivery); } } async handleFailure(delivery, errorMessage) { delivery.lastError = errorMessage; if (delivery.attempt >= delivery.maxAttempts) { delivery.status = entities_1.DeliveryStatus.FAILED; delivery.completedAt = new Date(); this.logger.warn(`Webhook delivery failed permanently: ${delivery.id} after ${delivery.attempt} attempts`); } else { delivery.status = entities_1.DeliveryStatus.RETRYING; delivery.attempt += 1; delivery.nextRetryAt = new Date(Date.now() + RETRY_DELAYS[delivery.attempt - 1]); this.logger.log(`Webhook delivery scheduled for retry: ${delivery.id} (attempt ${delivery.attempt})`); } } onCompleted(job) { this.logger.debug(`Job completed: ${job.id}`); } onFailed(job, error) { this.logger.error(`Job failed: ${job.id} - ${error.message}`); } }; exports.WebhookProcessor = WebhookProcessor; __decorate([ (0, bullmq_1.OnWorkerEvent)('completed'), __metadata("design:type", Function), __metadata("design:paramtypes", [bullmq_2.Job]), __metadata("design:returntype", void 0) ], WebhookProcessor.prototype, "onCompleted", null); __decorate([ (0, bullmq_1.OnWorkerEvent)('failed'), __metadata("design:type", Function), __metadata("design:paramtypes", [bullmq_2.Job, Error]), __metadata("design:returntype", void 0) ], WebhookProcessor.prototype, "onFailed", null); exports.WebhookProcessor = WebhookProcessor = WebhookProcessor_1 = __decorate([ (0, bullmq_1.Processor)('webhooks'), __param(0, (0, typeorm_1.InjectRepository)(entities_1.WebhookDeliveryEntity)), __metadata("design:paramtypes", [typeorm_2.Repository]) ], WebhookProcessor); //# sourceMappingURL=webhook.processor.js.map