Add NATS account-balance SoT and statement review for customers, CS, and sales
Some checks are pending
offline / test (push) Waiting to run

Internal billing now uses verae.billing.* request-reply and pubs. zappier-account-balance tracks prepaid, credits, usage, and payments. Portal, admin, CS, and sales all review the same statement. Independent Forgejo repos stay split via push-module-repos.
This commit is contained in:
George Lambert 2026-09-11 15:35:37 -04:00
parent a1a5b957fd
commit ac38676645
135 changed files with 3078 additions and 130 deletions

View file

@ -0,0 +1,3 @@
# zappier-account-balance NATS
See README. Queue `account-balance`. No public bind.

View file

@ -0,0 +1,31 @@
# zappier-account-balance
Source of truth for **customer prepaid balances**. Internal traffic is **NATS request-reply**; HTTP is health plus a fallback.
**Forgejo:** https://git.georgelambert.org/marchon/zappier-account-balance
**Catalog:** https://zapier.georgelambert.org/packages/zappier-account-balance/README.pdf
Zapier cloud never connects here. zappier-edge, customer-service, sales-pricing, and accounting-export do.
## Subjects
| Subject | Kind | Who |
|---------|------|-----|
| `verae.billing.statement.get` | request-reply | portal, CS, sales, admin |
| `verae.billing.balance.get` | request-reply | anyone internal |
| `verae.billing.balance.adjust` | request-reply | CS credits, portal reload |
| `verae.billing.usage.recorded` | pub | zappier-edge meter |
| `verae.billing.payment.recorded` | pub | portal reload, invoice paid |
| `verae.billing.credit.applied` | pub | CS |
Queue group: `account-balance`.
## Run
```bash
NATS_URL=nats://127.0.0.1:4222 PORT=3010 npm start
curl http://127.0.0.1:3010/health
curl http://127.0.0.1:3010/statement/cust_1
```
Statement JSON: `{ customerId, prepaidCents, credits[], usage[], payments[] }`.

View file

@ -0,0 +1,46 @@
{
"name": "zappier-account-balance",
"version": "0.1.0",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "zappier-account-balance",
"version": "0.1.0",
"dependencies": {
"nats": "^2.28.2"
}
},
"node_modules/nats": {
"version": "2.29.3",
"resolved": "https://registry.npmjs.org/nats/-/nats-2.29.3.tgz",
"integrity": "sha512-tOQCRCwC74DgBTk4pWZ9V45sk4d7peoE2njVprMRCBXrhJ5q5cYM7i6W+Uvw2qUrcfOSnuisrX7bEx3b3Wx4QA==",
"deprecated": "Package moved. Use @nats-io/transport-node from https://github.com/nats-io/nats.js",
"license": "Apache-2.0",
"dependencies": {
"nkeys.js": "1.1.0"
},
"engines": {
"node": ">= 14.0.0"
}
},
"node_modules/nkeys.js": {
"version": "1.1.0",
"resolved": "https://registry.npmjs.org/nkeys.js/-/nkeys.js-1.1.0.tgz",
"integrity": "sha512-tB/a0shZL5UZWSwsoeyqfTszONTt4k2YS0tuQioMOD180+MbombYVgzDUYHlx+gejYK6rgf08n/2Df99WY0Sxg==",
"license": "Apache-2.0",
"dependencies": {
"tweetnacl": "1.0.3"
},
"engines": {
"node": ">=10.0.0"
}
},
"node_modules/tweetnacl": {
"version": "1.0.3",
"resolved": "https://registry.npmjs.org/tweetnacl/-/tweetnacl-1.0.3.tgz",
"integrity": "sha512-6rt+RN7aOi1nGMyC4Xa5DdYiukl2UWCbcJft7YhxReBGQD7OAM8Pbxw6YMo4r2diNEA8FEmu32YOn9rhaiE5yw==",
"license": "Unlicense"
}
}
}

View file

@ -0,0 +1,14 @@
{
"name": "zappier-account-balance",
"version": "0.1.0",
"private": true,
"type": "module",
"description": "NATS source of truth for customer prepaid balances, credits, usage, payments",
"scripts": {
"start": "node src/server.js",
"test": "node --test test/*.test.js"
},
"dependencies": {
"nats": "^2.28.2"
}
}

View file

@ -0,0 +1,89 @@
/**
* Source of truth for prepaid balances, CS credits, usage, and payments.
*/
export class AccountBooks {
constructor() {
/** @type {Record<string, number>} */
this.prepaid = {};
this.credits = [];
this.usage = [];
this.payments = [];
}
prepaidCents(customerId) {
return this.prepaid[customerId] ?? 0;
}
adjust({ customerId, cents, reason, agent, kind = 'credit' }) {
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) {
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',
});
}
statement(customerId) {
const match = (rows) => rows.filter((r) => r.customerId === customerId).slice(0, 100);
return {
customerId,
prepaidCents: this.prepaidCents(customerId),
credits: match(this.credits),
usage: match(this.usage),
payments: match(this.payments),
};
}
}
export function handle(subject, payload, books) {
const p = payload || {};
if (subject.endsWith('balance.get') || subject.endsWith('statement.get')) {
return books.statement(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}`);
}

View file

@ -0,0 +1,92 @@
#!/usr/bin/env node
import http from 'node:http';
import { AccountBooks, handle } from './books.js';
import { SUBJECTS } from './subjects.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 = new AccountBooks();
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);
if (m.reply) m.respond(sc.encode(JSON.stringify(out)));
}
};
reply(nc.subscribe(SUBJECTS.STATEMENT_GET, { queue: SUBJECTS.QUEUE }), (p) =>
handle(SUBJECTS.STATEMENT_GET, p, books),
);
reply(nc.subscribe(SUBJECTS.BALANCE_GET, { queue: SUBJECTS.QUEUE }), (p) =>
handle(SUBJECTS.BALANCE_GET, p, books),
);
reply(nc.subscribe(SUBJECTS.BALANCE_ADJUST, { queue: SUBJECTS.QUEUE }), (p) =>
handle(SUBJECTS.BALANCE_ADJUST, p, books),
);
(async () => {
for await (const m of nc.subscribe(SUBJECTS.USAGE_RECORDED)) {
handle(SUBJECTS.USAGE_RECORDED, JSON.parse(sc.decode(m.data) || '{}'), books);
}
})();
// 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.statement(st[1]));
}
if (req.method === 'POST' && url.pathname === '/adjust') {
const body = await readBody(req);
return json(200, books.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`);
});

View file

@ -0,0 +1,10 @@
/** 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',
QUEUE: 'account-balance',
};

View file

@ -0,0 +1,16 @@
import { test } from 'node:test';
import assert from 'node:assert/strict';
import { AccountBooks, handle } from '../src/books.js';
import { SUBJECTS } from '../src/subjects.js';
test('adjust credits prepaid and statement lists credits usage payments', () => {
const books = new AccountBooks();
handle(SUBJECTS.BALANCE_ADJUST, { customerId: 'c1', cents: 500, reason: 'goodwill', agent: 'cs' }, books);
handle(SUBJECTS.USAGE_RECORDED, { customerId: 'c1', endpointId: 'timestamp', cents: 4 }, books);
handle(SUBJECTS.PAYMENT_RECORDED, { customerId: 'c1', cents: 1000, reason: 'reload' }, books);
const st = handle(SUBJECTS.STATEMENT_GET, { customerId: 'c1' }, books);
assert.equal(st.prepaidCents, 1496);
assert.equal(st.credits[0].cents, 500);
assert.equal(st.usage[0].endpointId, 'timestamp');
assert.equal(st.payments[0].kind, 'payment');
});

View file

@ -0,0 +1,30 @@
import { test } from 'node:test';
import assert from 'node:assert/strict';
import { spawn } from 'node:child_process';
import path from 'node:path';
import { fileURLToPath } from 'node:url';
const root = path.join(path.dirname(fileURLToPath(import.meta.url)), '..');
test('account-balance health and statement', async () => {
const port = 18010;
const child = spawn(process.execPath, ['src/server.js'], {
cwd: root,
env: { ...process.env, PORT: String(port) },
stdio: ['ignore', 'pipe', 'pipe'],
});
await new Promise((r) => setTimeout(r, 400));
try {
const h = await (await fetch(`http://127.0.0.1:${port}/health`)).json();
assert.equal(h.role, 'zappier-account-balance');
await fetch(`http://127.0.0.1:${port}/adjust`, {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ customerId: 'c1', cents: 200, reason: 'test' }),
});
const st = await (await fetch(`http://127.0.0.1:${port}/statement/c1`)).json();
assert.equal(st.prepaidCents, 200);
} finally {
child.kill('SIGTERM');
}
});