Initial import of zappier-account-balance from zapier monorepo

This commit is contained in:
George Lambert 2026-09-11 19:09:44 -04:00
commit 4cc3d777ab
10 changed files with 457 additions and 0 deletions

138
src/books.js Normal file
View file

@ -0,0 +1,138 @@
/**
* Source of truth for prepaid balances, CS credits, usage, and payments.
*/
export class AccountBooks {
constructor() {
/** @type {Record<string, number>} */
this.prepaid = {};
/** @type {Record<string, string>} */
this.veraeUserIds = {};
/** @type {Record<string, string>} */
this.names = {};
this.credits = [];
this.usage = [];
this.payments = [];
}
remember(customerId, veraeUserId, name) {
if (customerId && veraeUserId) this.veraeUserIds[customerId] = veraeUserId;
if (customerId && name) this.names[customerId] = name;
}
prepaidCents(customerId) {
return this.prepaid[customerId] ?? 0;
}
adjust({ customerId, cents, reason, agent, kind = 'credit', veraeUserId, name }) {
this.remember(customerId, veraeUserId, name);
const delta = Math.trunc(Number(cents) || 0);
const next = this.prepaidCents(customerId) + delta;
this.prepaid[customerId] = next;
const row = {
id: `adj_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 6)}`,
customerId,
cents: delta,
reason: reason || kind,
agent: agent || 'system',
kind,
at: new Date().toISOString(),
prepaidCents: next,
};
if (kind === 'payment' || kind === 'reload') this.payments.unshift(row);
else this.credits.unshift(row);
return row;
}
recordUsage(entry) {
this.remember(entry.customerId, entry.veraeUserId, entry.name);
const cents = Number(entry.cents) || 0;
const customerId = entry.customerId;
const next = this.prepaidCents(customerId) - cents;
this.prepaid[customerId] = next;
const row = {
customerId,
endpointId: entry.endpointId || 'unknown',
cents,
at: entry.at || new Date().toISOString(),
prepaidCents: next,
};
this.usage.unshift(row);
return row;
}
recordPayment(entry) {
return this.adjust({
customerId: entry.customerId,
cents: Number(entry.cents) || 0,
reason: entry.reason || 'payment',
agent: entry.agent || 'payments',
kind: 'payment',
});
}
lookup(idOrName) {
if (this.names[idOrName] || this.prepaid[idOrName] != null || this.veraeUserIds[idOrName]) {
return this.statement(idOrName);
}
const want = String(idOrName || '').toLowerCase();
const hit = Object.entries(this.names).find(([, n]) => String(n).toLowerCase() === want);
return this.statement(hit ? hit[0] : idOrName);
}
statement(customerId) {
const match = (rows) => rows.filter((r) => r.customerId === customerId).slice(0, 100);
return {
customerId,
name: this.names[customerId],
veraeUserId: this.veraeUserIds[customerId],
prepaidCents: this.prepaidCents(customerId),
credits: match(this.credits),
usage: match(this.usage),
payments: match(this.payments),
};
}
dump() {
return {
prepaid: this.prepaid,
veraeUserIds: this.veraeUserIds,
names: this.names,
credits: this.credits,
usage: this.usage,
payments: this.payments,
};
}
load(raw) {
if (!raw || typeof raw !== 'object') return this;
this.prepaid = raw.prepaid || {};
this.veraeUserIds = raw.veraeUserIds || {};
this.names = raw.names || {};
this.credits = Array.isArray(raw.credits) ? raw.credits : [];
this.usage = Array.isArray(raw.usage) ? raw.usage : [];
this.payments = Array.isArray(raw.payments) ? raw.payments : [];
return this;
}
}
export function handle(subject, payload, books) {
const p = payload || {};
if (subject.endsWith('customer.put')) {
books.remember(p.customerId, p.veraeUserId, p.name);
return books.lookup(p.customerId);
}
if (subject.endsWith('balance.get') || subject.endsWith('statement.get')) {
books.remember(p.customerId, p.veraeUserId, p.name);
return books.lookup(p.customerId);
}
if (subject.endsWith('balance.adjust') || subject.endsWith('credit.applied')) {
return books.adjust(p);
}
if (subject.endsWith('usage.recorded')) {
return books.recordUsage(p);
}
if (subject.endsWith('payment.recorded')) {
return books.recordPayment(p);
}
throw new Error(`unknown billing subject ${subject}`);
}

28
src/persist.js Normal file
View file

@ -0,0 +1,28 @@
import fs from 'node:fs';
import path from 'node:path';
export function booksPath() {
return (
process.env.BOOKS_PATH ||
path.join(process.env.FLEET_STATE_DIR || process.cwd(), 'data', 'books.json')
);
}
export function loadInto(books) {
const p = booksPath();
if (!fs.existsSync(p)) return books;
try {
books.load(JSON.parse(fs.readFileSync(p, 'utf8')));
} catch {
/* keep empty */
}
return books;
}
export function saveFrom(books) {
const p = booksPath();
fs.mkdirSync(path.dirname(p), { recursive: true });
const tmp = `${p}.tmp`;
fs.writeFileSync(tmp, JSON.stringify(books.dump(), null, 2));
fs.renameSync(tmp, p);
}

107
src/server.js Normal file
View file

@ -0,0 +1,107 @@
#!/usr/bin/env node
import http from 'node:http';
import { AccountBooks, handle } from './books.js';
import { SUBJECTS } from './subjects.js';
import { loadInto, saveFrom } from './persist.js';
const PORT = Number(process.env.PORT || process.env.FLEET_HEALTH_PORT || 3010);
const BIND = process.env.FLEET_HEALTH_BIND || '0.0.0.0';
const books = loadInto(new AccountBooks());
function apply(subject, payload) {
const out = handle(subject, payload, books);
saveFrom(books);
return out;
}
let natsOk = false;
async function startNats() {
const url = process.env.NATS_URL;
if (!url) return;
const { connect, StringCodec } = await import('nats');
const nc = await connect({ servers: url.split(','), name: 'zappier-account-balance' });
const sc = StringCodec();
const reply = async (sub, fn) => {
for await (const m of sub) {
let payload = {};
try {
payload = JSON.parse(sc.decode(m.data) || '{}');
} catch {
payload = {};
}
const out = fn(payload);
saveFrom(books);
if (m.reply) m.respond(sc.encode(JSON.stringify(out)));
}
};
reply(nc.subscribe(SUBJECTS.STATEMENT_GET, { queue: SUBJECTS.QUEUE }), (p) =>
apply(SUBJECTS.STATEMENT_GET, p),
);
reply(nc.subscribe(SUBJECTS.BALANCE_GET, { queue: SUBJECTS.QUEUE }), (p) =>
apply(SUBJECTS.BALANCE_GET, p),
);
reply(nc.subscribe(SUBJECTS.BALANCE_ADJUST, { queue: SUBJECTS.QUEUE }), (p) =>
apply(SUBJECTS.BALANCE_ADJUST, p),
);
reply(nc.subscribe(SUBJECTS.CUSTOMER_PUT, { queue: SUBJECTS.QUEUE }), (p) =>
apply(SUBJECTS.CUSTOMER_PUT, p),
);
(async () => {
for await (const m of nc.subscribe(SUBJECTS.USAGE_RECORDED)) {
apply(SUBJECTS.USAGE_RECORDED, JSON.parse(sc.decode(m.data) || '{}'));
}
})();
// payment.recorded / credit.applied are fan-out events. Prepaid mutations
// go through balance.adjust so publishers can both pub and request-reply.
natsOk = true;
process.stdout.write(`account-balance nats ${url}\n`);
}
function readBody(req) {
return new Promise((resolve) => {
const chunks = [];
req.on('data', (c) => chunks.push(c));
req.on('end', () => {
try {
resolve(JSON.parse(Buffer.concat(chunks).toString('utf8') || '{}'));
} catch {
resolve({});
}
});
});
}
const server = http.createServer(async (req, res) => {
const url = new URL(req.url || '/', `http://127.0.0.1:${PORT}`);
const json = (code, obj) => {
res.writeHead(code, { 'content-type': 'application/json' });
res.end(JSON.stringify(obj));
};
try {
if (req.method === 'GET' && url.pathname === '/health') {
return json(200, { ok: true, role: 'zappier-account-balance', nats: natsOk, subjects: SUBJECTS });
}
const st = url.pathname.match(/^\/statement\/([^/]+)$/);
if (req.method === 'GET' && st) {
return json(200, books.lookup(decodeURIComponent(st[1])));
}
if (req.method === 'POST' && url.pathname === '/customer') {
const body = await readBody(req);
return json(200, apply(SUBJECTS.CUSTOMER_PUT, body));
}
if (req.method === 'POST' && url.pathname === '/adjust') {
const body = await readBody(req);
return json(200, apply(SUBJECTS.BALANCE_ADJUST, body));
}
json(404, { error: 'not found' });
} catch (err) {
json(500, { error: err.message });
}
});
server.listen(PORT, BIND, () => {
process.stdout.write(`zappier-account-balance http://${BIND}:${PORT}/\n`);
});
startNats().catch((err) => {
process.stderr.write(`nats optional: ${err.message}\n`);
});

11
src/subjects.js Normal file
View file

@ -0,0 +1,11 @@
/** Internal billing bus. Zapier cloud never subscribes. */
export const SUBJECTS = {
BALANCE_GET: 'verae.billing.balance.get',
BALANCE_ADJUST: 'verae.billing.balance.adjust',
STATEMENT_GET: 'verae.billing.statement.get',
USAGE_RECORDED: 'verae.billing.usage.recorded',
PAYMENT_RECORDED: 'verae.billing.payment.recorded',
CREDIT_APPLIED: 'verae.billing.credit.applied',
CUSTOMER_PUT: 'verae.billing.customer.put',
QUEUE: 'account-balance',
};