Milestone 0: import zappier billing, Verae middleware, and Zapier research
Compose-ready workspace: packages/zappier (rate card, portal, Stripe), packages/verae-zapier-middleware (timestamp + NATS), packages/verae-zapier (CLI app), vendor/zapier-platform, and research/zapier vendor corpus. Gate 0 structure checks pass. Product code and research are not yet wired.
This commit is contained in:
commit
b4150c8250
1364 changed files with 6814366 additions and 0 deletions
|
|
@ -0,0 +1,111 @@
|
|||
/**
|
||||
* @fileoverview NATS + JetStream connection lifecycle.
|
||||
* @module nats/connection
|
||||
*/
|
||||
|
||||
import { createDebugger } from '../debug/logger.js';
|
||||
import { config } from '../config.js';
|
||||
import { SUBJECTS, STREAMS } from './subjects.js';
|
||||
|
||||
const log = createDebugger('nats');
|
||||
|
||||
/** @type {import('nats').NatsConnection|null} */
|
||||
let nc = null;
|
||||
/** @type {import('nats').JetStreamClient|null} */
|
||||
let js = null;
|
||||
/** @type {import('nats').JetStreamManager|null} */
|
||||
let jsm = null;
|
||||
|
||||
/**
|
||||
* Connect to NATS and return JetStream handles.
|
||||
*
|
||||
* @param {string} [url=config.natsUrl]
|
||||
* @returns {Promise<{ nc: import('nats').NatsConnection, js: import('nats').JetStreamClient, jsm: import('nats').JetStreamManager }>}
|
||||
*/
|
||||
export async function connectNats(url = config.natsUrl) {
|
||||
if (!config.natsEnabled && process.env.NATS_FORCE_CONNECT !== '1') {
|
||||
log.debug('connect skipped — NATS_ENABLED=false');
|
||||
throw new Error('NATS is disabled (NATS_ENABLED=false)');
|
||||
}
|
||||
|
||||
if (nc && js && jsm) {
|
||||
return { nc, js, jsm };
|
||||
}
|
||||
|
||||
log.info('connecting to NATS', { url });
|
||||
|
||||
const { connect } = await import('nats');
|
||||
nc = await connect({ servers: url, name: 'verae-zapier-middleware' });
|
||||
js = nc.jetstream();
|
||||
jsm = await nc.jetstreamManager();
|
||||
|
||||
log.info('NATS connected', { url });
|
||||
return { nc, js, jsm };
|
||||
}
|
||||
|
||||
/**
|
||||
* Idempotently create JetStream streams required by this middleware.
|
||||
*
|
||||
* @param {import('nats').JetStreamManager} [manager]
|
||||
* @returns {Promise<void>}
|
||||
*/
|
||||
export async function ensureStreams(manager) {
|
||||
const m = manager ?? jsm;
|
||||
if (!m) {
|
||||
throw new Error('JetStream manager not available — call connectNats first');
|
||||
}
|
||||
|
||||
/** @type {Array<{ name: string, subjects: string[] }>} */
|
||||
const defs = [
|
||||
{ name: STREAMS.ZAPIER_JOBS, subjects: [SUBJECTS.JOBS_WATCH] },
|
||||
{ name: STREAMS.ZAPIER_EVENTS, subjects: [SUBJECTS.JOBS_EVENTS] },
|
||||
{ name: STREAMS.ZAPIER_WEBHOOKS, subjects: [SUBJECTS.WEBHOOKS_DELIVER] },
|
||||
];
|
||||
|
||||
for (const def of defs) {
|
||||
try {
|
||||
await m.streams.info(def.name);
|
||||
log.debug('stream exists', { stream: def.name });
|
||||
} catch {
|
||||
await m.streams.add({
|
||||
name: def.name,
|
||||
subjects: def.subjects,
|
||||
retention: 'limits',
|
||||
storage: 'file',
|
||||
max_age: 24 * 60 * 60 * 1e9, // 24h in ns
|
||||
num_replicas: 1,
|
||||
});
|
||||
log.info('stream created', { stream: def.name, subjects: def.subjects });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Close the shared NATS connection if open.
|
||||
* @returns {Promise<void>}
|
||||
*/
|
||||
export async function closeNats() {
|
||||
if (!nc) {
|
||||
log.debug('closeNats: no active connection');
|
||||
return;
|
||||
}
|
||||
log.info('closing NATS connection');
|
||||
await nc.drain();
|
||||
nc = null;
|
||||
js = null;
|
||||
jsm = null;
|
||||
}
|
||||
|
||||
/**
|
||||
* @returns {import('nats').JetStreamClient|null}
|
||||
*/
|
||||
export function getJetStream() {
|
||||
return js;
|
||||
}
|
||||
|
||||
/**
|
||||
* @returns {boolean}
|
||||
*/
|
||||
export function isNatsConnected() {
|
||||
return Boolean(nc && !nc.isClosed());
|
||||
}
|
||||
|
|
@ -0,0 +1,95 @@
|
|||
/**
|
||||
* @fileoverview JetStream publishers for jobs, events, and webhooks.
|
||||
* @module nats/publishers
|
||||
*/
|
||||
|
||||
import { createDebugger } from '../debug/logger.js';
|
||||
import { getTraceId } from '../debug/trace-context.js';
|
||||
import { SUBJECTS } from './subjects.js';
|
||||
import { getJetStream, connectNats } from './connection.js';
|
||||
|
||||
const log = createDebugger('nats');
|
||||
|
||||
/**
|
||||
* @returns {Promise<import('nats').JetStreamClient>}
|
||||
*/
|
||||
async function requireJs() {
|
||||
let js = getJetStream();
|
||||
if (!js) {
|
||||
const handles = await connectNats();
|
||||
js = handles.js;
|
||||
}
|
||||
return js;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {object} partial
|
||||
* @returns {Promise<{ seq: number }>}
|
||||
*/
|
||||
export async function enqueueWatch(partial) {
|
||||
const msg = {
|
||||
attempt: 0,
|
||||
enqueuedAt: new Date().toISOString(),
|
||||
traceId: getTraceId() || 'no-trace',
|
||||
...partial,
|
||||
};
|
||||
|
||||
log.debug('enqueueWatch', {
|
||||
subject: SUBJECTS.JOBS_WATCH,
|
||||
tenantId: msg.tenantId,
|
||||
jobId: msg.jobId,
|
||||
attempt: msg.attempt,
|
||||
traceId: msg.traceId,
|
||||
});
|
||||
|
||||
const js = await requireJs();
|
||||
const ack = await js.publish(SUBJECTS.JOBS_WATCH, JSON.stringify(msg));
|
||||
return { seq: Number(ack.seq) };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {object} partial
|
||||
* @returns {Promise<{ seq: number }>}
|
||||
*/
|
||||
export async function publishJobEvent(partial) {
|
||||
const msg = {
|
||||
emittedAt: new Date().toISOString(),
|
||||
traceId: getTraceId() || 'no-trace',
|
||||
...partial,
|
||||
};
|
||||
|
||||
log.debug('publishJobEvent', {
|
||||
subject: SUBJECTS.JOBS_EVENTS,
|
||||
event: msg.event,
|
||||
jobId: msg.jobId,
|
||||
tenantId: msg.tenantId,
|
||||
});
|
||||
|
||||
const js = await requireJs();
|
||||
const ack = await js.publish(SUBJECTS.JOBS_EVENTS, JSON.stringify(msg));
|
||||
return { seq: Number(ack.seq) };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {object} partial
|
||||
* @returns {Promise<{ seq: number }>}
|
||||
*/
|
||||
export async function enqueueWebhook(partial) {
|
||||
const msg = {
|
||||
attempt: 1,
|
||||
traceId: getTraceId() || 'no-trace',
|
||||
...partial,
|
||||
};
|
||||
|
||||
log.debug('enqueueWebhook', {
|
||||
subject: SUBJECTS.WEBHOOKS_DELIVER,
|
||||
hookId: msg.hookId,
|
||||
tenantId: msg.tenantId,
|
||||
event: msg.event,
|
||||
targetUrl: msg.targetUrl,
|
||||
});
|
||||
|
||||
const js = await requireJs();
|
||||
const ack = await js.publish(SUBJECTS.WEBHOOKS_DELIVER, JSON.stringify(msg));
|
||||
return { seq: Number(ack.seq) };
|
||||
}
|
||||
|
|
@ -0,0 +1,43 @@
|
|||
/**
|
||||
* @fileoverview NATS subject and stream name constants.
|
||||
* @module nats/subjects
|
||||
*
|
||||
* See docs/architecture/nats-subjects.md for payload schemas.
|
||||
*/
|
||||
|
||||
/**
|
||||
* Subject strings used by publishers and consumers.
|
||||
* @readonly
|
||||
*/
|
||||
export const SUBJECTS = Object.freeze({
|
||||
/** Work queue: poll Verae job status */
|
||||
JOBS_WATCH: 'verae.zapier.jobs.watch',
|
||||
/** Terminal job outcomes */
|
||||
JOBS_EVENTS: 'verae.zapier.jobs.events',
|
||||
/** Work queue: HTTP POST to Zapier REST Hooks */
|
||||
WEBHOOKS_DELIVER: 'verae.zapier.webhooks.deliver',
|
||||
/** Optional metering stream */
|
||||
USAGE: 'verae.zapier.usage',
|
||||
});
|
||||
|
||||
/**
|
||||
* JetStream stream names.
|
||||
* @readonly
|
||||
*/
|
||||
export const STREAMS = Object.freeze({
|
||||
ZAPIER_JOBS: 'ZAPIER_JOBS',
|
||||
ZAPIER_EVENTS: 'ZAPIER_EVENTS',
|
||||
ZAPIER_WEBHOOKS: 'ZAPIER_WEBHOOKS',
|
||||
ZAPIER_USAGE: 'ZAPIER_USAGE',
|
||||
});
|
||||
|
||||
/**
|
||||
* Durable consumer names (queue groups).
|
||||
* @readonly
|
||||
*/
|
||||
export const CONSUMERS = Object.freeze({
|
||||
JOB_POLLER: 'job-poller',
|
||||
EVENT_WEBHOOK_ROUTER: 'event-webhook-router',
|
||||
WEBHOOK_DELIVER: 'webhook-deliver',
|
||||
USAGE_WRITER: 'usage-writer',
|
||||
});
|
||||
Loading…
Add table
Add a link
Reference in a new issue