/** * @fileoverview JetStream consumer that POSTs Zapier REST Hook payloads. * @module workers/webhookWorker */ import { createDebugger } from '../debug/logger.js'; import { SUBJECTS, CONSUMERS, STREAMS } from '../nats/subjects.js'; import { connectNats, ensureStreams } from '../nats/connection.js'; import { deliverWebhook } from '../services/webhookService.js'; import { getActiveWebhooks } from '../store/webhooks.js'; import { withTrace } from '../debug/trace.js'; const log = createDebugger('webhooks'); let running = false; /** @type {AbortController|null} */ let abort = null; /** * Route job events → per-hook deliver messages (inline or via re-publish). * Also handles direct deliver subjects. * * @param {object} data * @param {{ ack: () => Promise, nak: (d?: number) => Promise }} ctrl */ async function handleDeliver(data, ctrl) { await withTrace({ traceId: data.traceId, span: 'webhook-deliver' }, async () => { // Event router path: expand tenant hooks if (data.event && data.jobId && !data.targetUrl) { const hooks = getActiveWebhooks(data.tenantId, data.event); for (const hook of hooks) { await deliverWebhook(hook.targetUrl, { event: data.event, jobId: data.jobId, tenantId: data.tenantId, status: data.status, }); } await ctrl.ack(); return; } if (!data.targetUrl) { log.warn('deliver missing targetUrl', { dataKeys: Object.keys(data) }); await ctrl.ack(); return; } const result = await deliverWebhook(data.targetUrl, data.payload ?? data); if (result.ok) { await ctrl.ack(); } else { log.debug('deliver non-2xx, nak', { status: result.status }); await ctrl.nak(2000); } }); } /** * Start webhook delivery worker (consumes WEBHOOKS stream + optional events). * @returns {Promise<{ stop: () => Promise }>} */ export async function startWebhookWorker() { if (running) { return { stop: async () => stopWebhookWorker() }; } const { js, jsm } = await connectNats(); await ensureStreams(jsm); // Events consumer → deliver try { await jsm.consumers.add(STREAMS.ZAPIER_EVENTS, { durable_name: CONSUMERS.EVENT_WEBHOOK_ROUTER, ack_policy: 'explicit', filter_subject: SUBJECTS.JOBS_EVENTS, max_deliver: 10, }); } catch (err) { log.debug('events consumer may exist', { error: err.message }); } try { await jsm.consumers.add(STREAMS.ZAPIER_WEBHOOKS, { durable_name: CONSUMERS.WEBHOOK_DELIVER, ack_policy: 'explicit', filter_subject: SUBJECTS.WEBHOOKS_DELIVER, max_deliver: 10, }); } catch (err) { log.debug('webhook consumer may exist', { error: err.message }); } abort = new AbortController(); running = true; log.info('webhook worker started'); const runConsumer = async (stream, durable) => { const consumer = await js.consumers.get(stream, durable); const messages = await consumer.consume({ max_messages: 10 }); for await (const msg of messages) { if (abort?.signal.aborted) break; try { const data = JSON.parse(msg.string()); await handleDeliver(data, { ack: () => msg.ack(), nak: (d = 1000) => msg.nak(d), }); } catch (err) { log.error('webhook handle failed', { error: err.message }); try { msg.nak(1000); } catch { /* ignore */ } } } }; runConsumer(STREAMS.ZAPIER_EVENTS, CONSUMERS.EVENT_WEBHOOK_ROUTER).catch((err) => log.error('events consumer failed', { error: err.message }), ); runConsumer(STREAMS.ZAPIER_WEBHOOKS, CONSUMERS.WEBHOOK_DELIVER).catch((err) => log.error('webhooks consumer failed', { error: err.message }), ); return { stop: async () => stopWebhookWorker() }; } /** * @returns {Promise} */ export async function stopWebhookWorker() { abort?.abort(); abort = null; running = false; log.info('webhook worker stopped'); }