Initial import of verae-middleware from zapier monorepo
This commit is contained in:
commit
7fe9f616fa
66 changed files with 6590 additions and 0 deletions
108
src/workers/inProcessJobPoller.js
Normal file
108
src/workers/inProcessJobPoller.js
Normal file
|
|
@ -0,0 +1,108 @@
|
|||
/**
|
||||
* @fileoverview In-process job poller when NATS_ENABLED=false.
|
||||
* @module workers/inProcessJobPoller
|
||||
*/
|
||||
|
||||
import { config } from '../config.js';
|
||||
import { veraeClient } from '../clients/veraeClient.js';
|
||||
import {
|
||||
listPendingJobs,
|
||||
updateJobWatcher,
|
||||
removeJobWatcher,
|
||||
} from '../store/jobWatchers.js';
|
||||
import { getActiveWebhooks } from '../store/webhooks.js';
|
||||
import { deliverWebhook } from '../services/webhookService.js';
|
||||
import { createDebugger } from '../debug/logger.js';
|
||||
|
||||
const log = createDebugger('jobs');
|
||||
|
||||
let timer = null;
|
||||
let running = false;
|
||||
|
||||
/**
|
||||
* @param {import('../store/jobWatchers.js').JobWatcher} job
|
||||
*/
|
||||
async function processJob(job) {
|
||||
const attempts = job.attempts + 1;
|
||||
updateJobWatcher(job.id, { attempts });
|
||||
|
||||
if (attempts > config.jobPollMaxAttempts) {
|
||||
updateJobWatcher(job.id, { status: 'timeout' });
|
||||
removeJobWatcher(job.id);
|
||||
log.warn('job timeout', { jobId: job.jobId });
|
||||
return;
|
||||
}
|
||||
|
||||
let status;
|
||||
try {
|
||||
status = await veraeClient.getStatus(job.veraeToken, job.jobId);
|
||||
} catch (err) {
|
||||
log.debug('poll error', { jobId: job.jobId, error: err.message });
|
||||
return;
|
||||
}
|
||||
|
||||
if (status.status === 'pending') {
|
||||
return;
|
||||
}
|
||||
|
||||
const event = status.status === 'completed' ? 'timestamp.completed' : 'timestamp.failed';
|
||||
const hooks = getActiveWebhooks(job.tenantId, event);
|
||||
|
||||
for (const hook of hooks) {
|
||||
try {
|
||||
await deliverWebhook(hook.targetUrl, {
|
||||
event,
|
||||
jobId: job.jobId,
|
||||
tenantId: job.tenantId,
|
||||
status,
|
||||
});
|
||||
} catch (err) {
|
||||
log.error('webhook deliver failed', { hookId: hook.id, error: err.message });
|
||||
}
|
||||
}
|
||||
|
||||
updateJobWatcher(job.id, { status: status.status });
|
||||
removeJobWatcher(job.id);
|
||||
log.debug('job terminal', { jobId: job.jobId, status: status.status, hooks: hooks.length });
|
||||
}
|
||||
|
||||
async function tick() {
|
||||
if (running) return;
|
||||
running = true;
|
||||
try {
|
||||
const jobs = listPendingJobs();
|
||||
await Promise.all(jobs.map((job) => processJob(job)));
|
||||
} finally {
|
||||
running = false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Start interval poller (no-op if already started or NATS enabled).
|
||||
* @returns {void}
|
||||
*/
|
||||
export function startInProcessJobPoller() {
|
||||
if (config.natsEnabled) {
|
||||
log.info('in-process poller skipped (NATS_ENABLED=true)');
|
||||
return;
|
||||
}
|
||||
if (timer) return;
|
||||
|
||||
const interval = config.jobPollIntervalMs;
|
||||
timer = setInterval(() => {
|
||||
tick().catch((err) => log.error('poller tick failed', { error: err.message }));
|
||||
}, interval);
|
||||
|
||||
log.info('in-process job poller started', { intervalMs: interval });
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop interval poller.
|
||||
* @returns {void}
|
||||
*/
|
||||
export function stopInProcessJobPoller() {
|
||||
if (!timer) return;
|
||||
clearInterval(timer);
|
||||
timer = null;
|
||||
log.info('in-process job poller stopped');
|
||||
}
|
||||
158
src/workers/jobPollerWorker.js
Normal file
158
src/workers/jobPollerWorker.js
Normal file
|
|
@ -0,0 +1,158 @@
|
|||
/**
|
||||
* @fileoverview JetStream consumer that polls Verae job status.
|
||||
* @module workers/jobPollerWorker
|
||||
*/
|
||||
|
||||
import { createDebugger } from '../debug/logger.js';
|
||||
import { config } from '../config.js';
|
||||
import { SUBJECTS, CONSUMERS, STREAMS } from '../nats/subjects.js';
|
||||
import { connectNats, ensureStreams } from '../nats/connection.js';
|
||||
import { publishJobEvent } from '../nats/publishers.js';
|
||||
import { veraeClient } from '../clients/veraeClient.js';
|
||||
import { getTenant } from '../store/tenants.js';
|
||||
import { withTrace } from '../debug/trace.js';
|
||||
|
||||
const log = createDebugger('jobs');
|
||||
|
||||
let running = false;
|
||||
/** @type {AbortController|null} */
|
||||
let abort = null;
|
||||
|
||||
/**
|
||||
* Resolve a Verae token for polling (re-login via tenant if needed).
|
||||
* @param {object} msg
|
||||
* @returns {Promise<string>}
|
||||
*/
|
||||
async function resolveVeraeToken(msg) {
|
||||
if (msg.veraeToken) return msg.veraeToken;
|
||||
|
||||
const tenant = getTenant(msg.tenantId);
|
||||
if (!tenant?.veraeUsername) {
|
||||
throw new Error(`Cannot resolve token for tenant ${msg.tenantId}`);
|
||||
}
|
||||
const login = await veraeClient.login({
|
||||
username: tenant.veraeUsername,
|
||||
password: tenant.veraePassword,
|
||||
});
|
||||
return login.token;
|
||||
}
|
||||
|
||||
/**
|
||||
* Process one watch message.
|
||||
* @param {object} data
|
||||
* @param {{ ack: () => Promise<void>, nak: (delay?: number) => Promise<void> }} ctrl
|
||||
*/
|
||||
async function handleWatch(data, ctrl) {
|
||||
await withTrace({ traceId: data.traceId, span: 'job-poll' }, async () => {
|
||||
const attempt = (data.attempt ?? 0) + 1;
|
||||
const maxAttempts = data.maxAttempts ?? config.jobPollMaxAttempts;
|
||||
|
||||
if (attempt > maxAttempts) {
|
||||
await publishJobEvent({
|
||||
event: 'timestamp.timeout',
|
||||
tenantId: data.tenantId,
|
||||
jobId: data.jobId,
|
||||
status: { id: data.jobId, status: 'timeout' },
|
||||
traceId: data.traceId,
|
||||
});
|
||||
await ctrl.ack();
|
||||
return;
|
||||
}
|
||||
|
||||
const token = await resolveVeraeToken(data);
|
||||
let status;
|
||||
try {
|
||||
status = await veraeClient.getStatus(token, data.jobId);
|
||||
} catch (err) {
|
||||
log.debug('poll error, nak', { jobId: data.jobId, error: err.message });
|
||||
await ctrl.nak(config.jobPollIntervalMs);
|
||||
return;
|
||||
}
|
||||
|
||||
log.debug('poll status', { jobId: data.jobId, status: status.status, attempt });
|
||||
|
||||
if (status.status === 'pending') {
|
||||
await ctrl.nak(config.jobPollIntervalMs);
|
||||
return;
|
||||
}
|
||||
|
||||
const event =
|
||||
status.status === 'completed' ? 'timestamp.completed' : 'timestamp.failed';
|
||||
|
||||
await publishJobEvent({
|
||||
event,
|
||||
tenantId: data.tenantId,
|
||||
jobId: data.jobId,
|
||||
status,
|
||||
traceId: data.traceId,
|
||||
});
|
||||
await ctrl.ack();
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Start the durable job-poller worker loop.
|
||||
* @returns {Promise<{ stop: () => Promise<void> }>}
|
||||
*/
|
||||
export async function startJobPollerWorker() {
|
||||
if (running) {
|
||||
return { stop: async () => stopJobPollerWorker() };
|
||||
}
|
||||
|
||||
const { nc, js, jsm } = await connectNats();
|
||||
await ensureStreams(jsm);
|
||||
|
||||
// Ensure durable consumer (workqueue-style via filter + durable name)
|
||||
try {
|
||||
await jsm.consumers.add(STREAMS.ZAPIER_JOBS, {
|
||||
durable_name: CONSUMERS.JOB_POLLER,
|
||||
ack_policy: 'explicit',
|
||||
filter_subject: SUBJECTS.JOBS_WATCH,
|
||||
max_deliver: config.jobPollMaxAttempts + 5,
|
||||
ack_wait: 30_000_000_000, // 30s ns
|
||||
});
|
||||
} catch (err) {
|
||||
// already exists
|
||||
log.debug('consumer may exist', { error: err.message });
|
||||
}
|
||||
|
||||
const consumer = await js.consumers.get(STREAMS.ZAPIER_JOBS, CONSUMERS.JOB_POLLER);
|
||||
abort = new AbortController();
|
||||
running = true;
|
||||
log.info('job poller worker started', { consumer: CONSUMERS.JOB_POLLER });
|
||||
|
||||
(async () => {
|
||||
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 handleWatch(data, {
|
||||
ack: () => msg.ack(),
|
||||
nak: (delayMs = 1000) => msg.nak(delayMs),
|
||||
});
|
||||
} catch (err) {
|
||||
log.error('job poller handle failed', { error: err.message });
|
||||
try {
|
||||
msg.nak(1000);
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
}
|
||||
}
|
||||
})().catch((err) => log.error('job poller loop failed', { error: err.message }));
|
||||
|
||||
return {
|
||||
stop: async () => stopJobPollerWorker(),
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* @returns {Promise<void>}
|
||||
*/
|
||||
export async function stopJobPollerWorker() {
|
||||
abort?.abort();
|
||||
abort = null;
|
||||
running = false;
|
||||
log.info('job poller worker stopped');
|
||||
}
|
||||
138
src/workers/webhookWorker.js
Normal file
138
src/workers/webhookWorker.js
Normal file
|
|
@ -0,0 +1,138 @@
|
|||
/**
|
||||
* @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<void>, nak: (d?: number) => Promise<void> }} 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<void> }>}
|
||||
*/
|
||||
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<void>}
|
||||
*/
|
||||
export async function stopWebhookWorker() {
|
||||
abort?.abort();
|
||||
abort = null;
|
||||
running = false;
|
||||
log.info('webhook worker stopped');
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue