Add lavinmq amqp provider and amqp-ping consumer
lavinmq becomes a provider of a user-facing amqp interface: a consumer
that requires a message queue is given its OWN broker — a scoped vhost
and user on a lavinmq provider — not an account on the mesh's own
control-plane broker (ADR 0048). Vhost-per-login is the isolation model,
the exact analog of postgres's database-per-login: the provider names a
vhost after the consumer's login and a user with full rights on that
vhost and none elsewhere, so a login is a broker the consumer alone can
reach.
The provider drives lavinmq through its HTTP management API (client.ts,
the module's one impure seam), with a run-once bootstrap that computes
the RabbitMQ-compatible password hash lavinmq's config wants from the
plain admin secret the mesh mints — the value no ${secret:...}
placeholder can produce and the reason the bootstrap exists (ADR 0052).
serves.amqp carries the port so consumers reference ${bound:amqp:port}.
amqp-ping is a demo consumer: it contributes nothing (the vhost is the
login), reads its grant from an env-file the mesh fills, and uses
${bound:amqp:as} for BOTH its username and its vhost — the db-name
lesson applied to AMQP. It speaks AMQP 0-9-1 over a raw socket with no
npm dependency (the way redis speaks RESP) and round-trips one message.
It carries a slug so its identity fits the 20-char backend bound
(ADR 0049).
Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
@@ -0,0 +1,191 @@
|
||||
// amqp-ping's AMQP client — the demo consumer's own code (novox/hq ADR 0039). It speaks AMQP 0-9-1
|
||||
// directly over a raw TCP socket (node:net), the way redis's client speaks RESP: the module carries
|
||||
// NO npm dependency beyond @novox/mesh-sdk — no amqplib, no CLI in the image. It does exactly one
|
||||
// thing, the round-trip that proves the grant works: connect, authenticate with PLAIN to the vhost
|
||||
// the mesh named, declare a queue, publish one message and get it back.
|
||||
//
|
||||
// This is the consumer half of the `amqp` interface. It connects as the login the mesh derived
|
||||
// (`${bound:amqp:as}`) with the password the mesh minted (`${secret:amqp}`) to a vhost of that SAME
|
||||
// name — the provider named the vhost after the login, so the consumer must too. Nothing here is
|
||||
// hardcoded: user AND vhost are both the bound login, and a wrong vhost is refused by the broker.
|
||||
|
||||
import { createConnection, type Socket } from "node:net";
|
||||
import { readFileSync } from "node:fs";
|
||||
|
||||
const FRAME_END = 0xce;
|
||||
const PROTOCOL_HEADER = Buffer.from([0x41, 0x4d, 0x51, 0x50, 0x00, 0x00, 0x09, 0x01]); // "AMQP" 0-9-1
|
||||
|
||||
export interface AmqpConn {
|
||||
readonly host: string;
|
||||
readonly port: number;
|
||||
readonly user: string;
|
||||
readonly password: string;
|
||||
readonly vhost: string;
|
||||
}
|
||||
|
||||
/** Build the connection facts from the environment the mesh's env-file set (see the module manifest). */
|
||||
export function connFromEnv(env: NodeJS.ProcessEnv = process.env): AmqpConn {
|
||||
const host = env.MESH_AMQP_HOST ?? "";
|
||||
const port = Number(env.MESH_AMQP_PORT ?? "5672") || 5672;
|
||||
const user = env.MESH_AMQP_USER ?? "";
|
||||
const vhost = env.MESH_AMQP_VHOST ?? user; // the provider names the vhost after the login
|
||||
const password = env.MESH_AMQP_PASSWORD ?? readMaybe(env.MESH_AMQP_PASSWORD_FILE);
|
||||
if (!host || !user || !password) {
|
||||
throw new Error(`amqp-ping: connection is not fully set yet (host=${host} user=${user} password=${password ? "set" : "unset"})`);
|
||||
}
|
||||
return { host, port, user, password, vhost };
|
||||
}
|
||||
|
||||
// --- wire helpers ---------------------------------------------------------------------------------
|
||||
|
||||
function shortstr(s: string): Buffer {
|
||||
const b = Buffer.from(s, "utf8");
|
||||
const o = Buffer.alloc(1 + b.length);
|
||||
o.writeUInt8(b.length, 0);
|
||||
b.copy(o, 1);
|
||||
return o;
|
||||
}
|
||||
function longstr(s: Buffer | string): Buffer {
|
||||
const b = Buffer.isBuffer(s) ? s : Buffer.from(s, "utf8");
|
||||
const o = Buffer.alloc(4 + b.length);
|
||||
o.writeUInt32BE(b.length, 0);
|
||||
b.copy(o, 4);
|
||||
return o;
|
||||
}
|
||||
function u16(n: number): Buffer {
|
||||
const o = Buffer.alloc(2);
|
||||
o.writeUInt16BE(n, 0);
|
||||
return o;
|
||||
}
|
||||
function u32(n: number): Buffer {
|
||||
const o = Buffer.alloc(4);
|
||||
o.writeUInt32BE(n, 0);
|
||||
return o;
|
||||
}
|
||||
function frame(type: number, channel: number, payload: Buffer): Buffer {
|
||||
const o = Buffer.alloc(7 + payload.length + 1);
|
||||
o.writeUInt8(type, 0);
|
||||
o.writeUInt16BE(channel, 1);
|
||||
o.writeUInt32BE(payload.length, 3);
|
||||
payload.copy(o, 7);
|
||||
o.writeUInt8(FRAME_END, 7 + payload.length);
|
||||
return o;
|
||||
}
|
||||
function method(channel: number, classId: number, methodId: number, ...parts: Buffer[]): Buffer {
|
||||
return frame(1, channel, Buffer.concat([u16(classId), u16(methodId), ...parts]));
|
||||
}
|
||||
|
||||
interface MethodWaiter {
|
||||
classId: number;
|
||||
methodId: number;
|
||||
resolve: (args: Buffer) => void;
|
||||
reject: (e: Error) => void;
|
||||
}
|
||||
|
||||
/**
|
||||
* Connect, authenticate to the vhost, declare a queue, publish one message and get it back. Returns
|
||||
* the body that came back — the caller checks it equals what went out. Throws on any protocol error,
|
||||
* including the broker's `NOT_ALLOWED` refusal of a vhost the login has no permission on (the
|
||||
* isolation the provider builds, seen from the consumer's side).
|
||||
*/
|
||||
export function roundTrip(conn: AmqpConn, queue = "amqp-ping", payload?: string): Promise<string> {
|
||||
const body = Buffer.from(payload ?? `ping-${Date.now()}`);
|
||||
return new Promise<string>((resolve, reject) => {
|
||||
const sock: Socket = createConnection({ host: conn.host, port: conn.port });
|
||||
let buf = Buffer.alloc(0);
|
||||
const waiters: MethodWaiter[] = [];
|
||||
let lastBody: Buffer | null = null;
|
||||
let done = false;
|
||||
|
||||
const fail = (e: Error): void => {
|
||||
if (done) return;
|
||||
done = true;
|
||||
sock.destroy();
|
||||
reject(e);
|
||||
};
|
||||
const expect = (classId: number, methodId: number): Promise<Buffer> =>
|
||||
new Promise((res, rej) => waiters.push({ classId, methodId, resolve: res, reject: rej }));
|
||||
|
||||
sock.on("error", (e) => fail(e));
|
||||
sock.on("close", () => fail(new Error("amqp connection closed before the round-trip completed")));
|
||||
sock.on("data", (chunk: Buffer) => {
|
||||
buf = Buffer.concat([buf, chunk]);
|
||||
for (;;) {
|
||||
if (buf.length < 7) return;
|
||||
const type = buf.readUInt8(0);
|
||||
const size = buf.readUInt32BE(3);
|
||||
if (buf.length < 7 + size + 1) return;
|
||||
const framePayload = buf.subarray(7, 7 + size);
|
||||
buf = buf.subarray(7 + size + 1);
|
||||
if (type === 1) {
|
||||
const classId = framePayload.readUInt16BE(0);
|
||||
const methodId = framePayload.readUInt16BE(2);
|
||||
const args = framePayload.subarray(4);
|
||||
const w = waiters.shift();
|
||||
if (!w) continue;
|
||||
if (w.classId === classId && w.methodId === methodId) w.resolve(args);
|
||||
else w.reject(new Error(`expected method ${w.classId}/${w.methodId}, got ${classId}/${methodId}: ${args.toString("utf8")}`));
|
||||
} else if (type === 3) {
|
||||
lastBody = framePayload; // a content body frame
|
||||
}
|
||||
// type 2 (content header) and type 8 (heartbeat) need no handling for this round-trip.
|
||||
}
|
||||
});
|
||||
|
||||
sock.on("connect", () => {
|
||||
void (async () => {
|
||||
try {
|
||||
sock.write(PROTOCOL_HEADER);
|
||||
await expect(10, 10); // Connection.Start
|
||||
const response = Buffer.concat([
|
||||
Buffer.from([0]), Buffer.from(conn.user, "utf8"), Buffer.from([0]), Buffer.from(conn.password, "utf8"),
|
||||
]);
|
||||
// Connection.Start-Ok: empty client-properties table, PLAIN, the SASL response, locale.
|
||||
sock.write(method(0, 10, 11, u32(0), shortstr("PLAIN"), longstr(response), shortstr("en_US")));
|
||||
const tune = await expect(10, 30); // Connection.Tune
|
||||
const frameMax = tune.readUInt32BE(2) || 131072;
|
||||
sock.write(method(0, 10, 31, u16(tune.readUInt16BE(0)), u32(frameMax), u16(0))); // Tune-Ok, no heartbeat
|
||||
sock.write(method(0, 10, 40, shortstr(conn.vhost), shortstr(""), Buffer.from([0]))); // Connection.Open
|
||||
await expect(10, 41); // Open-Ok — authenticated and into the vhost
|
||||
|
||||
sock.write(method(1, 20, 10, shortstr(""))); // Channel.Open
|
||||
await expect(20, 11);
|
||||
// Queue.Declare: reserved, queue, bits(auto-delete=1), empty arguments table.
|
||||
sock.write(method(1, 50, 10, u16(0), shortstr(queue), Buffer.from([0b00001000]), u32(0)));
|
||||
await expect(50, 11);
|
||||
|
||||
// Basic.Publish to the default exchange, routing-key = queue; then content header + body.
|
||||
sock.write(method(1, 60, 40, u16(0), shortstr(""), shortstr(queue), Buffer.from([0])));
|
||||
const bodySize = Buffer.alloc(8);
|
||||
bodySize.writeBigUInt64BE(BigInt(body.length), 0);
|
||||
sock.write(frame(2, 1, Buffer.concat([u16(60), u16(0), bodySize, u16(0)]))); // content header, no properties
|
||||
sock.write(frame(3, 1, body)); // content body
|
||||
|
||||
await new Promise((r) => setTimeout(r, 200));
|
||||
lastBody = null;
|
||||
sock.write(method(1, 60, 70, u16(0), shortstr(queue), Buffer.from([1]))); // Basic.Get, no-ack
|
||||
await expect(60, 71); // Get-Ok (a Get-Empty would arrive as 60/72 and reject the expect)
|
||||
await new Promise((r) => setTimeout(r, 200));
|
||||
const received = lastBody ? (lastBody as Buffer).toString("utf8") : "";
|
||||
|
||||
sock.write(method(0, 10, 50, u16(200), shortstr("bye"), u16(0), u16(0))); // Connection.Close
|
||||
await expect(10, 51).catch(() => undefined);
|
||||
done = true;
|
||||
sock.end();
|
||||
resolve(received);
|
||||
} catch (e) {
|
||||
fail(e instanceof Error ? e : new Error(String(e)));
|
||||
}
|
||||
})();
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
function readMaybe(path: string | undefined): string {
|
||||
if (!path) return "";
|
||||
try {
|
||||
return readFileSync(path, "utf8").replace(/\n$/, "");
|
||||
} catch {
|
||||
return "";
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
// amqp-ping — a tiny demo consumer of the mesh `amqp` interface, run as a long-lived container by
|
||||
// `mesh-tools run` (it never returns, so the container stays up). It exists to PROVE the grant end to
|
||||
// end: the mesh gave it a scoped login and a vhost of that name on the lavinmq provider, and this
|
||||
// connects with exactly those and round-trips a message.
|
||||
//
|
||||
// The connection facts arrive the way every consumer's do — the mesh writes them into an env-file the
|
||||
// container reads (novox/hq ADR 0048): MESH_AMQP_HOST/PORT from the binding, MESH_AMQP_USER and
|
||||
// MESH_AMQP_VHOST both from `${bound:amqp:as}` (the provider named the vhost after the login, so the
|
||||
// consumer uses the login for both — the db-name lesson applied to AMQP), and MESH_AMQP_PASSWORD from
|
||||
// `${secret:amqp}`.
|
||||
//
|
||||
// It retries: on first boot the provider may not have provisioned this consumer yet (the reconcile is
|
||||
// asynchronous and cross-container), so a refused or unreachable connection is a "not yet", not a
|
||||
// failure — it waits and tries again until the round-trip succeeds, then holds the connection idle
|
||||
// and re-pings on a slow cadence so the container is a stable, running proof.
|
||||
|
||||
import { connFromEnv, roundTrip } from "./client.js";
|
||||
|
||||
async function sleep(ms: number): Promise<void> {
|
||||
await new Promise((r) => setTimeout(r, ms));
|
||||
}
|
||||
|
||||
async function pingOnce(): Promise<boolean> {
|
||||
try {
|
||||
const conn = connFromEnv();
|
||||
const sent = `ping-${Date.now()}`;
|
||||
const got = await roundTrip(conn, "amqp-ping", sent);
|
||||
if (got === sent) {
|
||||
console.log(`[amqp-ping] round-trip ok as ${conn.user} on vhost ${conn.vhost} (${conn.host}:${conn.port})`);
|
||||
return true;
|
||||
}
|
||||
console.error(`[amqp-ping] round-trip mismatch: sent ${sent}, got ${got}`);
|
||||
return false;
|
||||
} catch (err) {
|
||||
console.error(`[amqp-ping] not ready yet: ${err instanceof Error ? err.message : err}`);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// Wait for the first successful round-trip — the proof this consumer's grant works — then stay up.
|
||||
let first = false;
|
||||
for (let i = 0; !first; i++) {
|
||||
first = await pingOnce();
|
||||
if (!first) await sleep(3000);
|
||||
}
|
||||
console.log("[amqp-ping] connected and round-tripped; holding steady");
|
||||
for (;;) {
|
||||
await sleep(30000);
|
||||
await pingOnce();
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
{
|
||||
"module": "amqp-ping",
|
||||
"slug": "ping",
|
||||
"version": "1",
|
||||
"capabilities": [
|
||||
"container-runtime"
|
||||
],
|
||||
"requires": [
|
||||
"amqp"
|
||||
],
|
||||
"contributes": {},
|
||||
"binds": {
|
||||
"amqp": "/var/lib/amqp-ping/amqp.json"
|
||||
},
|
||||
"secrets": {
|
||||
"amqp": "/var/lib/amqp-ping/amqp.secret"
|
||||
},
|
||||
"own-secrets": {
|
||||
"broker": "/var/lib/mesh/amqp-ping/broker"
|
||||
},
|
||||
"resources": [
|
||||
{
|
||||
"id": "mesh-state",
|
||||
"type": "directory",
|
||||
"path": "/var/lib/mesh/amqp-ping",
|
||||
"mode": "0700"
|
||||
},
|
||||
{
|
||||
"id": "state",
|
||||
"type": "directory",
|
||||
"path": "/var/lib/amqp-ping",
|
||||
"mode": "0700"
|
||||
},
|
||||
{
|
||||
"id": "amqp-env",
|
||||
"type": "file",
|
||||
"path": "/var/lib/amqp-ping/amqp.env",
|
||||
"mode": "0600",
|
||||
"content": "MESH_AMQP_HOST=${bound:amqp:at}\nMESH_AMQP_PORT=${bound:amqp:port}\nMESH_AMQP_USER=${bound:amqp:as}\nMESH_AMQP_VHOST=${bound:amqp:as}\nMESH_AMQP_PASSWORD=${secret:amqp}\n"
|
||||
},
|
||||
{
|
||||
"id": "net",
|
||||
"type": "network",
|
||||
"name": "amqp-ping"
|
||||
},
|
||||
{
|
||||
"id": "runtime",
|
||||
"type": "container",
|
||||
"name": "amqp-ping",
|
||||
"image": "mesh-runtime-amqp-ping@sha256:0000000000000000000000000000000000000000000000000000000000000000",
|
||||
"network": "amqp-ping",
|
||||
"env-file": [
|
||||
"/var/lib/amqp-ping/amqp.env"
|
||||
],
|
||||
"args": [
|
||||
"run",
|
||||
"/app/modules/amqp-ping/dist/index.js"
|
||||
],
|
||||
"restart-on": [
|
||||
"amqp-env"
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"name": "@novox/module-amqp-ping",
|
||||
"version": "0.1.0",
|
||||
"description": "amqp-ping — a demo consumer of the mesh amqp interface. Connects with its scoped grant and round-trips one message to prove the broker the mesh gave it (novox/hq ADR 0039).",
|
||||
"type": "module",
|
||||
"private": true,
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.0.0",
|
||||
"typescript": "^5.6.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ES2022",
|
||||
"module": "NodeNext",
|
||||
"moduleResolution": "NodeNext",
|
||||
"strict": true,
|
||||
"esModuleInterop": true,
|
||||
"skipLibCheck": true,
|
||||
"noEmit": true
|
||||
},
|
||||
"include": ["client.ts", "index.ts"]
|
||||
}
|
||||
Reference in New Issue
Block a user