Files
jschoubben 97707480d6 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
2026-09-06 15:50:19 +02:00

192 lines
8.1 KiB
TypeScript

// 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 "";
}
}