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
192 lines
8.1 KiB
TypeScript
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 "";
|
|
}
|
|
}
|