// 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 { const body = Buffer.from(payload ?? `ping-${Date.now()}`); return new Promise((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 => 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 ""; } }