/** * GATE 8 — Job poller worker + webhook via NATS events */ import { describe, it, before, after } from 'node:test'; import assert from 'node:assert/strict'; import http from 'node:http'; import { config } from '../../src/config.js'; import { useTempStore, seedProTenant } from '../helpers.js'; import { connectNats, ensureStreams, closeNats } from '../../src/nats/connection.js'; import { enqueueWatch } from '../../src/nats/publishers.js'; import { startJobPollerWorker, stopJobPollerWorker } from '../../src/workers/jobPollerWorker.js'; import { startWebhookWorker, stopWebhookWorker } from '../../src/workers/webhookWorker.js'; import { createWebhook } from '../../src/store/webhooks.js'; import { veraeClient, clearMockJobs } from '../../src/clients/veraeClient.js'; describe('NATS workers', () => { /** @type {ReturnType} */ let ctx; /** @type {object[]} */ let deliveries; /** @type {import('http').Server} */ let hookServer; /** @type {number} */ let hookPort; /** @type {string} */ let tenantId; before(async () => { assert.equal(config.mockVerae, true); config.natsEnabled = true; process.env.NATS_FORCE_CONNECT = '1'; config.natsUrl = process.env.NATS_URL || 'nats://127.0.0.1:4222'; config.jobPollIntervalMs = 50; config.jobPollMaxAttempts = 40; clearMockJobs(); ctx = useTempStore(); const { tenant } = seedProTenant(); tenantId = tenant.id; deliveries = []; await new Promise((resolve) => { hookServer = http.createServer((req, res) => { let body = ''; req.on('data', (c) => { body += c; }); req.on('end', () => { deliveries.push(JSON.parse(body || '{}')); res.writeHead(200); res.end('ok'); }); }); hookServer.listen(0, '127.0.0.1', () => { hookPort = hookServer.address().port; resolve(); }); }); createWebhook({ tenantId, targetUrl: `http://127.0.0.1:${hookPort}/hook`, event: 'timestamp.completed', }); await connectNats(config.natsUrl); await ensureStreams(); await startJobPollerWorker(); await startWebhookWorker(); }); after(async () => { await stopJobPollerWorker(); await stopWebhookWorker(); await closeNats(); await new Promise((r) => hookServer.close(r)); ctx.cleanup(); process.env.NATS_FORCE_CONNECT = ''; }); it('watch → poll → event → webhook delivery', async () => { const login = await veraeClient.login({ username: 'prouser', password: 'propass', }); const { jobId } = await veraeClient.createTimestamp(login.token, { data: 'nats-worker-test', }); await enqueueWatch({ tenantId, jobId, veraeToken: login.token, maxAttempts: 40, intervalMs: 50, }); const deadline = Date.now() + 8000; while (deliveries.length === 0 && Date.now() < deadline) { await new Promise((r) => setTimeout(r, 50)); } assert.ok(deliveries.length >= 1, 'expected webhook from NATS path'); assert.equal(deliveries[0].event, 'timestamp.completed'); assert.equal(deliveries[0].jobId, jobId); }); });