Initial import of verae-middleware from zapier monorepo
This commit is contained in:
commit
fb6db30e0e
66 changed files with 6590 additions and 0 deletions
108
test/integration/nats-workers.test.js
Normal file
108
test/integration/nats-workers.test.js
Normal file
|
|
@ -0,0 +1,108 @@
|
|||
/**
|
||||
* 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<typeof useTempStore>} */
|
||||
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);
|
||||
});
|
||||
});
|
||||
Loading…
Add table
Add a link
Reference in a new issue