runtime: serve each tool over mesh.rpc, and an invoke subcommand (ADR 0052) #3
@@ -0,0 +1,48 @@
|
|||||||
|
import amqp from "amqplib";
|
||||||
|
|
||||||
|
const PORT = process.argv[2];
|
||||||
|
const MPORT = process.argv[3];
|
||||||
|
const B = `http://127.0.0.1:${MPORT}`;
|
||||||
|
const AUTH = "Basic " + Buffer.from("guest:guest").toString("base64");
|
||||||
|
|
||||||
|
async function api(method, path, body) {
|
||||||
|
const r = await fetch(B + path, {
|
||||||
|
method,
|
||||||
|
headers: { "content-type": "application/json", authorization: AUTH },
|
||||||
|
body: body ? JSON.stringify(body) : undefined,
|
||||||
|
});
|
||||||
|
if (r.status >= 300 && r.status !== 404) throw new Error(`${method} ${path} -> ${r.status}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
await api("PUT", "/api/exchanges/%2f/mesh.events.dead", { type: "topic", durable: true });
|
||||||
|
await api("PUT", "/api/users/al", { password: "s", tags: "" });
|
||||||
|
|
||||||
|
const Q = "anchor.al.events";
|
||||||
|
const D = "mesh.events.dead";
|
||||||
|
const q = Q.replace(/\./g, "\\.");
|
||||||
|
const d = D.replace(/\./g, "\\.");
|
||||||
|
|
||||||
|
// configure, write, read patterns per grant on the dead exchange
|
||||||
|
const combos = {
|
||||||
|
"none": { configure: `^${q}$`, write: `^${q}$`, read: `^${q}$` },
|
||||||
|
"read-dead": { configure: `^${q}$`, write: `^${q}$`, read: `^(${q}|${d})$` },
|
||||||
|
"write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^${q}$` },
|
||||||
|
"configure-dead": { configure: `^(${q}|${d})$`, write: `^${q}$`, read: `^${q}$` },
|
||||||
|
"read+write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^(${q}|${d})$` },
|
||||||
|
};
|
||||||
|
|
||||||
|
let i = 0;
|
||||||
|
for (const [label, perms] of Object.entries(combos)) {
|
||||||
|
await api("PUT", "/api/permissions/%2f/al", perms);
|
||||||
|
const queue = `${Q}.${i++}`; // fresh each time
|
||||||
|
try {
|
||||||
|
const c = await amqp.connect(`amqp://al:s@127.0.0.1:${PORT}/`);
|
||||||
|
const ch = await c.createChannel();
|
||||||
|
ch.on("error", () => {});
|
||||||
|
await ch.assertQueue(queue, { durable: true, deadLetterExchange: D });
|
||||||
|
console.log(`${label}: declare-with-DLX OK`);
|
||||||
|
await c.close();
|
||||||
|
} catch (e) {
|
||||||
|
console.log(`${label}: FAIL - ${String(e.message).slice(0, 70)}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
+265
-16
@@ -1,34 +1,87 @@
|
|||||||
// A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker
|
// A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker
|
||||||
// (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime
|
// (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime
|
||||||
// provides the binding — so a broker-client change never rebuilds the modules.
|
// provides the binding — so a broker-client change never rebuilds the modules. This is where the
|
||||||
|
// ADR 0047 wire shape lives: the two exchanges, persistent events, per-consumer durable queues,
|
||||||
|
// prefetch, dead-letter — none of which a module ever sees.
|
||||||
|
|
||||||
import amqp from "amqplib";
|
import amqp from "amqplib";
|
||||||
import { randomUUID } from "node:crypto";
|
import * as tls from "node:tls";
|
||||||
import type { Broker, Envelope } from "@novox/mesh-sdk/messaging";
|
import { randomUUID, createHash } from "node:crypto";
|
||||||
|
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
// Two topic exchanges, kept apart on purpose: tool invocations are request/reply and are not
|
// Two topic exchanges, kept apart on purpose (ADR 0047): tool invocations are request/reply and are
|
||||||
// events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node
|
// not events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node
|
||||||
// events — never the RPC traffic.
|
// events — never the RPC traffic.
|
||||||
const RPC_EXCHANGE = "mesh.rpc";
|
const RPC_EXCHANGE = "mesh.rpc";
|
||||||
const EVENTS_EXCHANGE = "mesh.events";
|
const EVENTS_EXCHANGE = "mesh.events";
|
||||||
|
// Where an event rejected past its redelivery limit is set aside for inspection.
|
||||||
|
const DEAD_EXCHANGE = "mesh.events.dead";
|
||||||
|
// Bound in-flight events so one slow consumer can't pull the whole backlog into memory (ADR 0047).
|
||||||
|
const EVENT_PREFETCH = 32;
|
||||||
|
|
||||||
interface Reply {
|
interface Reply {
|
||||||
result?: unknown;
|
result?: unknown;
|
||||||
error?: string;
|
error?: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. */
|
/** A broker credential as the mesh delivers it (novox/hq ADR 0048): an amqps URL, the fingerprint
|
||||||
export async function connectAmqp(url: string): Promise<Broker> {
|
* of the certificate the broker must present, and the node and module the account is scoped to (so
|
||||||
const conn = await amqp.connect(url);
|
* the runtime names its queue as the mesh did). A plain string is a bootstrap URL. */
|
||||||
const ch = await conn.createChannel();
|
export interface Credential {
|
||||||
|
url: string;
|
||||||
|
fingerprint?: string;
|
||||||
|
node?: string;
|
||||||
|
module?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down.
|
||||||
|
*
|
||||||
|
* A scoped module (novox/hq ADR 0048) passes `assumeExchanges: true`: its account may not declare
|
||||||
|
* an exchange, and the substrate already owns them, so it declares only its own queue. A credential
|
||||||
|
* carrying a fingerprint is dialled over amqps, pinned to exactly that certificate.
|
||||||
|
*/
|
||||||
|
export async function connectAmqp(
|
||||||
|
target: string | Credential,
|
||||||
|
opts: { assumeExchanges?: boolean } = {},
|
||||||
|
): Promise<Broker> {
|
||||||
|
const cred: Credential = typeof target === "string" ? { url: target } : target;
|
||||||
|
const conn = cred.fingerprint
|
||||||
|
? await amqp.connect(cred.url, await pinnedOptions(cred.url, cred.fingerprint))
|
||||||
|
: await amqp.connect(cred.url);
|
||||||
|
|
||||||
|
// A confirm channel, so an event publish awaits the broker's ack: a publish the broker never
|
||||||
|
// accepted (it was mid-restart, the connection dropped) fails the emit rather than vanishing —
|
||||||
|
// at-least-once starts at the emitter, not only the consumer (ADR 0047).
|
||||||
|
const ch = await conn.createConfirmChannel();
|
||||||
|
|
||||||
|
// The substrate owns the exchanges (ADR 0048). A bootstrap/admin connection declares them; a
|
||||||
|
// scoped module assumes they exist and never tries — its account could not, and the dead-letter
|
||||||
|
// queue behind the exchange is the substrate's to keep, not a module's.
|
||||||
|
if (!opts.assumeExchanges) {
|
||||||
await ch.assertExchange(RPC_EXCHANGE, "topic", { durable: true });
|
await ch.assertExchange(RPC_EXCHANGE, "topic", { durable: true });
|
||||||
await ch.assertExchange(EVENTS_EXCHANGE, "topic", { durable: true });
|
await ch.assertExchange(EVENTS_EXCHANGE, "topic", { durable: true });
|
||||||
|
await ch.assertExchange(DEAD_EXCHANGE, "topic", { durable: true });
|
||||||
|
await ch.assertQueue(DEAD_EXCHANGE, { durable: true });
|
||||||
|
await ch.bindQueue(DEAD_EXCHANGE, DEAD_EXCHANGE, "#");
|
||||||
|
}
|
||||||
|
|
||||||
// Request/reply: one exclusive reply queue, correlationId → resolver.
|
await ch.prefetch(EVENT_PREFETCH);
|
||||||
const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true });
|
|
||||||
|
// Request/reply is set up lazily: a consumer-only module (the audit logger) never calls a tool,
|
||||||
|
// and its scoped account may not declare the exclusive reply queue this would otherwise need.
|
||||||
const pending = new Map<string, (r: Reply) => void>();
|
const pending = new Map<string, (r: Reply) => void>();
|
||||||
|
let replyQueue: string | undefined;
|
||||||
|
async function ensureReply(): Promise<string> {
|
||||||
|
if (replyQueue) return replyQueue;
|
||||||
|
const { queue } = await ch.assertQueue("", { exclusive: true });
|
||||||
|
// Replies come back through the RPC exchange keyed by this queue's own name, not the default
|
||||||
|
// exchange (novox/hq ADR 0052): a serving module's scoped account may write to mesh.rpc but not
|
||||||
|
// the default exchange, which would let it publish into any queue on the broker.
|
||||||
|
await ch.bindQueue(queue, RPC_EXCHANGE, queue);
|
||||||
|
replyQueue = queue;
|
||||||
await ch.consume(
|
await ch.consume(
|
||||||
replyQueue,
|
queue,
|
||||||
(msg) => {
|
(msg) => {
|
||||||
if (!msg) return;
|
if (!msg) return;
|
||||||
const resolve = pending.get(msg.properties.correlationId);
|
const resolve = pending.get(msg.properties.correlationId);
|
||||||
@@ -39,9 +92,46 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
},
|
},
|
||||||
{ noAck: true },
|
{ noAck: true },
|
||||||
);
|
);
|
||||||
|
return queue;
|
||||||
|
}
|
||||||
|
|
||||||
|
// One durable event queue per consumer (ADR 0047: <node>.<module>.events), with many bindings and
|
||||||
|
// a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message
|
||||||
|
// once however many bindings match, so the local match is what keeps a two-pattern module from
|
||||||
|
// running the wrong handler.
|
||||||
|
type EventSub = { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> };
|
||||||
|
const eventSubs: EventSub[] = [];
|
||||||
|
let eventQueue: string | undefined;
|
||||||
|
let eventConsumerTag: string | undefined;
|
||||||
|
|
||||||
|
async function dispatchEvent(msg: amqp.ConsumeMessage): Promise<void> {
|
||||||
|
const key = msg.fields.routingKey;
|
||||||
|
let env: Envelope<unknown>;
|
||||||
|
try {
|
||||||
|
env = toEnvelope(msg);
|
||||||
|
} catch {
|
||||||
|
// An undecodable body will never decode on redelivery — dead-letter it at once rather than
|
||||||
|
// wedge the queue or loop. Decoding sits before the handler try on purpose: a poison message
|
||||||
|
// is a different failure from a handler that threw, and gets no retry.
|
||||||
|
ch.nack(msg, false, false);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
for (const s of eventSubs) {
|
||||||
|
if (topicMatches(s.pattern, key)) await s.handler(env);
|
||||||
|
}
|
||||||
|
ch.ack(msg);
|
||||||
|
} catch {
|
||||||
|
// First handler failure: requeue once. A second (already redelivered) dead-letters it, so a
|
||||||
|
// poison event is set aside rather than looping forever or vanishing (ADR 0047). Redelivery
|
||||||
|
// re-runs every matching handler, so a consumer must be idempotent — which the ADR requires.
|
||||||
|
ch.nack(msg, false, !msg.fields.redelivered);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
return {
|
return {
|
||||||
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
||||||
|
const reply = await ensureReply();
|
||||||
const id = randomUUID();
|
const id = randomUUID();
|
||||||
const answered = new Promise<Res>((resolve, reject) => {
|
const answered = new Promise<Res>((resolve, reject) => {
|
||||||
const timer = setTimeout(() => {
|
const timer = setTimeout(() => {
|
||||||
@@ -55,12 +145,14 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
});
|
});
|
||||||
ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), {
|
ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), {
|
||||||
correlationId: id,
|
correlationId: id,
|
||||||
replyTo: replyQueue,
|
replyTo: reply,
|
||||||
});
|
});
|
||||||
return answered;
|
return answered;
|
||||||
},
|
},
|
||||||
|
|
||||||
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
||||||
|
// A durable, shared serve queue (ADR 0047): several runtimes serving one tool key compete for
|
||||||
|
// invocations rather than each answering the same call.
|
||||||
const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true });
|
const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true });
|
||||||
await ch.bindQueue(queue, RPC_EXCHANGE, key);
|
await ch.bindQueue(queue, RPC_EXCHANGE, key);
|
||||||
const consumer = await ch.consume(queue, (msg) => {
|
const consumer = await ch.consume(queue, (msg) => {
|
||||||
@@ -73,7 +165,9 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
reply = { error: err instanceof Error ? err.message : String(err) };
|
reply = { error: err instanceof Error ? err.message : String(err) };
|
||||||
}
|
}
|
||||||
if (msg.properties.replyTo) {
|
if (msg.properties.replyTo) {
|
||||||
ch.sendToQueue(msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), {
|
// Reply through the RPC exchange, keyed by the caller's reply-queue name, so a scoped
|
||||||
|
// account answers with write on mesh.rpc alone — never the default exchange (ADR 0052).
|
||||||
|
ch.publish(RPC_EXCHANGE, msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), {
|
||||||
correlationId: msg.properties.correlationId,
|
correlationId: msg.properties.correlationId,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
@@ -84,17 +178,79 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
},
|
},
|
||||||
|
|
||||||
async publish<T>(env: Envelope<T>): Promise<void> {
|
async publish<T>(env: Envelope<T>): Promise<void> {
|
||||||
ch.publish(EVENTS_EXCHANGE, env.key, Buffer.from(JSON.stringify(env.body)));
|
const headers = env.headers ?? ({} as EventHeaders);
|
||||||
|
// Events are persistent (delivery-mode 2): an audit trail that loses events on a broker
|
||||||
|
// restart is not one (ADR 0047). Metadata rides as headers; the body is only the payload.
|
||||||
|
// The publish is awaited to the broker's confirm — an unaccepted publish rejects here.
|
||||||
|
await new Promise<void>((resolve, reject) => {
|
||||||
|
ch.publish(
|
||||||
|
EVENTS_EXCHANGE,
|
||||||
|
env.key,
|
||||||
|
Buffer.from(JSON.stringify(env.body)),
|
||||||
|
{
|
||||||
|
persistent: true,
|
||||||
|
contentType:
|
||||||
|
typeof headers["content-type"] === "string" ? headers["content-type"] : "application/json",
|
||||||
|
messageId: headers["x-event-id"],
|
||||||
|
headers: { ...headers },
|
||||||
|
},
|
||||||
|
(err) => (err ? reject(err instanceof Error ? err : new Error(String(err))) : resolve()),
|
||||||
|
);
|
||||||
|
});
|
||||||
},
|
},
|
||||||
|
|
||||||
async subscribe<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): Promise<() => void> {
|
async subscribe<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): Promise<() => void> {
|
||||||
|
const node = process.env.MESH_NODE;
|
||||||
|
const mod = process.env.MESH_MODULE;
|
||||||
|
|
||||||
|
// A module we can name gets its ADR 0047 durable queue; an anonymous subscriber (a test, an
|
||||||
|
// ad-hoc listener) gets a transient exclusive one that dies with the connection.
|
||||||
|
if (node && mod) {
|
||||||
|
if (!eventQueue) {
|
||||||
|
const name = `${node}.${mod}.events`;
|
||||||
|
if (opts.assumeExchanges) {
|
||||||
|
// The mesh pre-declared this queue with its dead-letter when it issued the account: a
|
||||||
|
// scoped account may not declare a dead-lettered queue itself (the broker refuses that
|
||||||
|
// to a non-administrator). Passively check it is there, then bind and consume.
|
||||||
|
await ch.checkQueue(name);
|
||||||
|
} else {
|
||||||
|
await ch.assertQueue(name, { durable: true, deadLetterExchange: DEAD_EXCHANGE });
|
||||||
|
}
|
||||||
|
eventQueue = name;
|
||||||
|
}
|
||||||
|
await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, pattern);
|
||||||
|
const sub: EventSub = { pattern, handler: handler as EventSub["handler"] };
|
||||||
|
eventSubs.push(sub);
|
||||||
|
if (!eventConsumerTag) {
|
||||||
|
const consumer = await ch.consume(eventQueue, (msg) => {
|
||||||
|
if (msg) void dispatchEvent(msg);
|
||||||
|
});
|
||||||
|
eventConsumerTag = consumer.consumerTag;
|
||||||
|
}
|
||||||
|
return () => {
|
||||||
|
const i = eventSubs.indexOf(sub);
|
||||||
|
if (i >= 0) eventSubs.splice(i, 1);
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
const { queue } = await ch.assertQueue("", { exclusive: true });
|
const { queue } = await ch.assertQueue("", { exclusive: true });
|
||||||
await ch.bindQueue(queue, EVENTS_EXCHANGE, pattern);
|
await ch.bindQueue(queue, EVENTS_EXCHANGE, pattern);
|
||||||
const consumer = await ch.consume(queue, (msg) => {
|
const consumer = await ch.consume(queue, (msg) => {
|
||||||
if (!msg) return;
|
if (!msg) return;
|
||||||
void (async () => {
|
void (async () => {
|
||||||
await handler({ key: msg.fields.routingKey, node: "", body: JSON.parse(msg.content.toString()) as T });
|
let env: Envelope<T>;
|
||||||
|
try {
|
||||||
|
env = toEnvelope<T>(msg);
|
||||||
|
} catch {
|
||||||
|
ch.nack(msg, false, false); // undecodable — drop, never retry
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
await handler(env);
|
||||||
ch.ack(msg);
|
ch.ack(msg);
|
||||||
|
} catch {
|
||||||
|
ch.nack(msg, false, !msg.fields.redelivered);
|
||||||
|
}
|
||||||
})();
|
})();
|
||||||
});
|
});
|
||||||
return () => void ch.cancel(consumer.consumerTag);
|
return () => void ch.cancel(consumer.consumerTag);
|
||||||
@@ -106,3 +262,96 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
},
|
},
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Normalise a certificate fingerprint to bare lower-case hex, dropping an `sha256:` prefix and
|
||||||
|
* any colon grouping, so two spellings of the same fingerprint compare equal. */
|
||||||
|
function normalizeFingerprint(fingerprint: string): string {
|
||||||
|
return fingerprint.replace(/^sha256:/i, "").replace(/:/g, "").toLowerCase();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Socket options that pin the broker to exactly the certificate whose fingerprint the mesh
|
||||||
|
* delivered (novox/hq ADR 0048, as the builder does). Done in two phases so a credential never
|
||||||
|
* reaches an impostor: first a bare TLS connection that sends nothing fetches the certificate and
|
||||||
|
* the fingerprint is checked; only then does the real connection trust *that* certificate as its
|
||||||
|
* own authority, so the AMQP login flows solely to the broker that proved it holds the pinned key.
|
||||||
|
* Node's `checkServerIdentity` does not run under `rejectUnauthorized: false`, so a one-phase
|
||||||
|
* "connect then check" would have already sent the password to whoever answered.
|
||||||
|
*/
|
||||||
|
async function pinnedOptions(rawUrl: string, fingerprint: string): Promise<tls.ConnectionOptions> {
|
||||||
|
const url = new URL(rawUrl);
|
||||||
|
const host = url.hostname;
|
||||||
|
const port = url.port ? Number(url.port) : 5671;
|
||||||
|
|
||||||
|
const certificate = await new Promise<tls.DetailedPeerCertificate>((resolve, reject) => {
|
||||||
|
const socket = tls.connect({ host, port, servername: host, rejectUnauthorized: false }, () => {
|
||||||
|
const peer = socket.getPeerCertificate(true);
|
||||||
|
socket.destroy();
|
||||||
|
if (!peer || !peer.raw) reject(new Error("the broker presented no certificate to pin"));
|
||||||
|
else resolve(peer);
|
||||||
|
});
|
||||||
|
socket.setTimeout(15_000, () => {
|
||||||
|
socket.destroy();
|
||||||
|
reject(new Error("timed out fetching the broker's certificate"));
|
||||||
|
});
|
||||||
|
socket.on("error", reject);
|
||||||
|
});
|
||||||
|
|
||||||
|
const seen = createHash("sha256").update(certificate.raw).digest("hex");
|
||||||
|
if (seen !== normalizeFingerprint(fingerprint)) {
|
||||||
|
throw new Error(
|
||||||
|
`the broker's certificate (sha256:${seen}) does not match the pinned ${fingerprint} — refusing`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
const pem =
|
||||||
|
"-----BEGIN CERTIFICATE-----\n" +
|
||||||
|
(certificate.raw.toString("base64").match(/.{1,64}/g) ?? []).join("\n") +
|
||||||
|
"\n-----END CERTIFICATE-----\n";
|
||||||
|
// Trust that one certificate and nothing else; the mesh's own name is not in any public store,
|
||||||
|
// so identity is the pin, not the hostname — checkServerIdentity is satisfied deliberately.
|
||||||
|
return { ca: [pem], checkServerIdentity: () => undefined };
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Read a broker message back into an Envelope: string headers, contentType folded in, body parsed. */
|
||||||
|
function toEnvelope<T>(msg: amqp.ConsumeMessage): Envelope<T> {
|
||||||
|
const raw = msg.properties.headers ?? {};
|
||||||
|
const headers: Record<string, string> = {};
|
||||||
|
for (const [k, v] of Object.entries(raw)) {
|
||||||
|
if (v == null) continue;
|
||||||
|
headers[k] = typeof v === "string" ? v : String(v);
|
||||||
|
}
|
||||||
|
if (!headers["content-type"] && msg.properties.contentType) {
|
||||||
|
headers["content-type"] = msg.properties.contentType;
|
||||||
|
}
|
||||||
|
return {
|
||||||
|
key: msg.fields.routingKey,
|
||||||
|
node: headers["x-node"] ?? "",
|
||||||
|
body: JSON.parse(msg.content.toString()) as T,
|
||||||
|
headers: headers as EventHeaders,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/** AMQP topic matching: `*` matches one word, `#` zero or more. Used to fan a shared queue's
|
||||||
|
* deliveries out to the handlers whose pattern actually matches the routing key. */
|
||||||
|
function topicMatches(pattern: string, key: string): boolean {
|
||||||
|
return matchFrom(pattern.split("."), 0, key.split("."), 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean {
|
||||||
|
while (pi < p.length) {
|
||||||
|
const tok = p[pi];
|
||||||
|
if (tok === "#") {
|
||||||
|
if (pi === p.length - 1) return true; // trailing # swallows the rest, including nothing
|
||||||
|
for (let skip = ki; skip <= k.length; skip++) {
|
||||||
|
if (matchFrom(p, pi + 1, k, skip)) return true;
|
||||||
|
}
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
if (ki >= k.length) return false;
|
||||||
|
if (tok !== "*" && tok !== k[ki]) return false;
|
||||||
|
pi++;
|
||||||
|
ki++;
|
||||||
|
}
|
||||||
|
return ki === k.length;
|
||||||
|
}
|
||||||
|
|||||||
+109
-12
@@ -1,20 +1,70 @@
|
|||||||
// The runnable entrypoint. Reads its configuration from the environment the host resolved for it,
|
// The runnable entrypoint. Two modes:
|
||||||
// connects the mesh broker, and serves the assigned modules' tools until stopped.
|
|
||||||
//
|
//
|
||||||
// MESH_BROKER_URL amqp://… the mesh broker
|
// mesh-tools serve — bind the broker and serve the assigned modules until
|
||||||
// MESH_TOOL_MODULES /path/a,/path/b,… compiled tool entrypoints of the assigned modules
|
// stopped. A module entrypoint that subscribes to events (on("#"))
|
||||||
|
// starts consuming as it is imported, so this also runs consumers.
|
||||||
|
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
|
||||||
|
// and what an events test uses to put a message on the wire.
|
||||||
|
//
|
||||||
|
// The broker, in order of preference:
|
||||||
|
// MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0048) — an
|
||||||
|
// amqps account scoped to this module. Preferred: a module holds its own.
|
||||||
|
// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account.
|
||||||
|
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
|
||||||
|
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0047)
|
||||||
|
|
||||||
|
import { readFileSync } from "node:fs";
|
||||||
import { connectAmqp } from "./broker-amqp.js";
|
import { connectAmqp } from "./broker-amqp.js";
|
||||||
|
import type { Credential } from "./broker-amqp.js";
|
||||||
import { runTools } from "./runtime.js";
|
import { runTools } from "./runtime.js";
|
||||||
|
import { invokeTool } from "@novox/mesh-sdk/tools";
|
||||||
|
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||||
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
import { emit } from "@novox/mesh-sdk/events";
|
||||||
|
|
||||||
async function main(): Promise<void> {
|
/**
|
||||||
const url = requireEnv("MESH_BROKER_URL");
|
* Connect the way this process is meant to: with its sealed credential if the mesh gave it one, and
|
||||||
|
* over the plain bootstrap URL otherwise. A scoped module assumes the substrate's exchanges exist —
|
||||||
|
* its account may not declare them (ADR 0048).
|
||||||
|
*/
|
||||||
|
async function connectBroker(): Promise<Broker> {
|
||||||
|
const file = process.env.MESH_BROKER_FILE;
|
||||||
|
if (file) {
|
||||||
|
let credential: Credential;
|
||||||
|
try {
|
||||||
|
credential = JSON.parse(readFileSync(file, "utf8")) as Credential;
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`);
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
if (!credential.url) {
|
||||||
|
console.error(`mesh-tools: ${file} carries no url — it is not a broker credential`);
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
// The mesh scoped this account to a node and module; take the runtime's identity from the
|
||||||
|
// credential so its queue and the events it emits match what the mesh authorised, no matter
|
||||||
|
// what the environment says.
|
||||||
|
if (credential.node) process.env.MESH_NODE = credential.node;
|
||||||
|
if (credential.module) process.env.MESH_MODULE = credential.module;
|
||||||
|
return connectAmqp(credential, { assumeExchanges: true });
|
||||||
|
}
|
||||||
|
const url = process.env.MESH_BROKER_URL;
|
||||||
|
if (!url) {
|
||||||
|
console.error(
|
||||||
|
"mesh-tools: set MESH_BROKER_FILE (a sealed credential) or MESH_BROKER_URL — there is no broker to reach",
|
||||||
|
);
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
return connectAmqp(url);
|
||||||
|
}
|
||||||
|
|
||||||
|
async function serve(): Promise<void> {
|
||||||
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
|
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
|
||||||
.split(",")
|
.split(",")
|
||||||
.map((s) => s.trim())
|
.map((s) => s.trim())
|
||||||
.filter(Boolean);
|
.filter(Boolean);
|
||||||
|
|
||||||
const broker = await connectAmqp(url);
|
const broker = await connectBroker();
|
||||||
const stop = await runTools({ broker, moduleEntrypoints });
|
const stop = await runTools({ broker, moduleEntrypoints });
|
||||||
|
|
||||||
const shutdown = async (): Promise<void> => {
|
const shutdown = async (): Promise<void> => {
|
||||||
@@ -26,13 +76,60 @@ async function main(): Promise<void> {
|
|||||||
process.on("SIGINT", () => void shutdown());
|
process.on("SIGINT", () => void shutdown());
|
||||||
}
|
}
|
||||||
|
|
||||||
function requireEnv(name: string): string {
|
async function emitOnce(type: string, bodyJson: string): Promise<void> {
|
||||||
const v = process.env[name];
|
let body: unknown = {};
|
||||||
if (!v) {
|
if (bodyJson) {
|
||||||
console.error(`mesh-tools: ${name} is not set — the runtime cannot serve without it`);
|
try {
|
||||||
|
body = JSON.parse(bodyJson);
|
||||||
|
} catch {
|
||||||
|
console.error(`mesh-tools emit: body is not JSON: ${bodyJson}`);
|
||||||
process.exit(1);
|
process.exit(1);
|
||||||
}
|
}
|
||||||
return v;
|
}
|
||||||
|
const broker = await connectBroker();
|
||||||
|
useBroker(() => broker);
|
||||||
|
// emit awaits the broker's publish confirm (ADR 0047), so the event is accepted before we close.
|
||||||
|
await emit(type, body);
|
||||||
|
await broker.close();
|
||||||
|
}
|
||||||
|
|
||||||
|
async function invokeOnce(module: string, tool: string, argsJson: string): Promise<void> {
|
||||||
|
let args: Record<string, unknown> = {};
|
||||||
|
if (argsJson) {
|
||||||
|
try {
|
||||||
|
args = JSON.parse(argsJson) as Record<string, unknown>;
|
||||||
|
} catch {
|
||||||
|
console.error(`mesh-tools invoke: args are not JSON: ${argsJson}`);
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
const broker = await connectBroker();
|
||||||
|
const result = await invokeTool(broker, module, tool, args);
|
||||||
|
process.stdout.write(JSON.stringify(result) + "\n");
|
||||||
|
await broker.close();
|
||||||
|
}
|
||||||
|
|
||||||
|
async function main(): Promise<void> {
|
||||||
|
const [command, ...rest] = process.argv.slice(2);
|
||||||
|
if (command === "invoke") {
|
||||||
|
const [module, tool] = rest;
|
||||||
|
if (!module || !tool) {
|
||||||
|
console.error("mesh-tools invoke <module> <tool> [json-args] — a module and tool are required");
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
await invokeOnce(module, tool, rest[2] ?? "");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (command === "emit") {
|
||||||
|
const type = rest[0];
|
||||||
|
if (!type) {
|
||||||
|
console.error("mesh-tools emit <type> [json-body] — a routing key is required");
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
await emitOnce(type, rest[1] ?? "");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
await serve();
|
||||||
}
|
}
|
||||||
|
|
||||||
void main();
|
void main();
|
||||||
|
|||||||
+4
-1
@@ -25,8 +25,11 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
|||||||
await import(pathToFileURL(resolve(entry)).href);
|
await import(pathToFileURL(resolve(entry)).href);
|
||||||
}
|
}
|
||||||
|
|
||||||
const stop = await serveTools(opts.broker);
|
// Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the
|
||||||
|
// audit logger) registers none, and its scoped account may not declare the serve queue — so a
|
||||||
|
// runtime that always served would fail for exactly the modules that never needed it.
|
||||||
const tools = listTools();
|
const tools = listTools();
|
||||||
|
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {};
|
||||||
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
|
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
|
||||||
return stop;
|
return stop;
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-9
@@ -1,7 +1,7 @@
|
|||||||
import { test } from "node:test";
|
import { test } from "node:test";
|
||||||
import assert from "node:assert/strict";
|
import assert from "node:assert/strict";
|
||||||
|
|
||||||
import { registerModuleTools, resetTools } from "@novox/mesh-sdk/tools";
|
import { registerModuleTools, resetTools, invokeTool } from "@novox/mesh-sdk/tools";
|
||||||
import { connectAmqp } from "../dist/broker-amqp.js";
|
import { connectAmqp } from "../dist/broker-amqp.js";
|
||||||
import { runTools } from "../dist/runtime.js";
|
import { runTools } from "../dist/runtime.js";
|
||||||
|
|
||||||
@@ -27,17 +27,12 @@ test("a module's tool serves and is invoked over a real AMQP broker", { skip: !u
|
|||||||
const serverBroker = await connectAmqp(url!);
|
const serverBroker = await connectAmqp(url!);
|
||||||
const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] });
|
const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] });
|
||||||
|
|
||||||
// A separate connection — a caller, like mesh-control's command API — invokes over the broker.
|
// A separate connection — a caller, like mesh-control's command API — invokes over the broker,
|
||||||
|
// by module and tool (novox/hq ADR 0052: served on serve.demo.greet, invoked as demo.greet).
|
||||||
const caller = await connectAmqp(url!);
|
const caller = await connectAmqp(url!);
|
||||||
const result = await caller.request<{ tool: string; args: Record<string, unknown> }, { hello: string }>(
|
const result = (await invokeTool(caller, "demo", "greet", { who: "mesh" })) as { hello: string };
|
||||||
"tools.invoke",
|
|
||||||
{ tool: "greet", args: { who: "mesh" } },
|
|
||||||
);
|
|
||||||
assert.equal(result.hello, "mesh");
|
assert.equal(result.hello, "mesh");
|
||||||
|
|
||||||
// An unknown tool is refused over the wire, not silently dropped.
|
|
||||||
await assert.rejects(caller.request("tools.invoke", { tool: "nope", args: {} }), /no such tool/);
|
|
||||||
|
|
||||||
stop();
|
stop();
|
||||||
await caller.close();
|
await caller.close();
|
||||||
await serverBroker.close();
|
await serverBroker.close();
|
||||||
|
|||||||
@@ -0,0 +1,142 @@
|
|||||||
|
import { test } from "node:test";
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import amqp from "amqplib";
|
||||||
|
|
||||||
|
import { connectAmqp } from "../dist/broker-amqp.js";
|
||||||
|
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||||
|
import { emit, on, type Event } from "@novox/mesh-sdk/events";
|
||||||
|
|
||||||
|
// Binding conformance: does the mesh-tools AMQP *adapter* honour the ADR 0047 wire contract —
|
||||||
|
// headers on the wire, a body that is only the payload, persistent messages, a durable per-consumer
|
||||||
|
// queue, dead-letter, poison handling? This tests the adapter in isolation, against a disposable
|
||||||
|
// broker that stands in for the one the mesh hosts (ADR 0001). It is NOT an event test of the mesh:
|
||||||
|
// that lives in a full lab scenario where the mesh raises the broker as substrate. Requires
|
||||||
|
// $MESH_BROKER_URL (a throwaway broker); skipped, never failed, when it is not set.
|
||||||
|
const url = process.env.MESH_BROKER_URL;
|
||||||
|
|
||||||
|
test("an event rides the wire with ADR 0047 headers and a body that is only the payload", { skip: !url }, async () => {
|
||||||
|
const sub = await connectAmqp(url!);
|
||||||
|
const pub = await connectAmqp(url!);
|
||||||
|
|
||||||
|
// The consumer's identity names its durable queue (<node>.<module>.events).
|
||||||
|
process.env.MESH_NODE = "lab";
|
||||||
|
process.env.MESH_MODULE = "audit-logger";
|
||||||
|
useBroker(() => sub);
|
||||||
|
const got: Event[] = [];
|
||||||
|
await on("#", async (e) => void got.push(e));
|
||||||
|
await delay(200); // let the binding settle before publishing
|
||||||
|
|
||||||
|
// The emitter is a different module on the same node.
|
||||||
|
process.env.MESH_MODULE = "umami";
|
||||||
|
useBroker(() => pub);
|
||||||
|
await emit("module.umami.site.created", { domain: "my-app" }, { causationId: "cmd-1" });
|
||||||
|
|
||||||
|
await waitFor(() => got.length > 0, 4000);
|
||||||
|
const e = got[0];
|
||||||
|
assert.equal(e.type, "module.umami.site.created");
|
||||||
|
assert.equal(e.source, "umami"); // x-source — read from a header, not the body
|
||||||
|
assert.equal(e.node, "lab"); // x-node
|
||||||
|
assert.ok(e.id, "x-event-id present"); // the handle a consumer dedups on
|
||||||
|
assert.equal(e.causationId, "cmd-1"); // x-causation-id round-trips
|
||||||
|
assert.match(e.at, /^\d{4}-\d{2}-\d{2}T/); // x-time, RFC-3339
|
||||||
|
assert.deepEqual(e.body, { domain: "my-app" }); // provenance never leaked into the body
|
||||||
|
|
||||||
|
// The durable per-consumer queue exists and is bound — a passive assert throws if it does not.
|
||||||
|
const probe = await amqp.connect(url!);
|
||||||
|
const pch = await probe.createChannel();
|
||||||
|
await pch.checkQueue("lab.audit-logger.events");
|
||||||
|
await probe.close();
|
||||||
|
|
||||||
|
await sub.close();
|
||||||
|
await pub.close();
|
||||||
|
delete process.env.MESH_NODE;
|
||||||
|
delete process.env.MESH_MODULE;
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a handler that keeps failing dead-letters the event past the redelivery limit", { skip: !url }, async () => {
|
||||||
|
const module = `flaky-${Date.now()}`;
|
||||||
|
const c = await connectAmqp(url!);
|
||||||
|
process.env.MESH_NODE = "lab";
|
||||||
|
process.env.MESH_MODULE = module;
|
||||||
|
useBroker(() => c);
|
||||||
|
|
||||||
|
let attempts = 0;
|
||||||
|
await on("module.test.boom", async () => {
|
||||||
|
attempts++;
|
||||||
|
throw new Error("boom");
|
||||||
|
});
|
||||||
|
await delay(200);
|
||||||
|
|
||||||
|
await emit("module.test.boom", { n: 1 });
|
||||||
|
|
||||||
|
// First delivery requeues once; the redelivered copy is dead-lettered — two attempts, then it
|
||||||
|
// leaves the consumer queue for good.
|
||||||
|
await waitFor(() => attempts >= 2, 5000);
|
||||||
|
await delay(300);
|
||||||
|
assert.equal(attempts, 2, "attempted twice, not looping forever");
|
||||||
|
|
||||||
|
// The poison event is retained on the dead-letter queue for inspection, not vanished.
|
||||||
|
const probe = await amqp.connect(url!);
|
||||||
|
const pch = await probe.createChannel();
|
||||||
|
const dead = await pch.get("mesh.events.dead", { noAck: true });
|
||||||
|
assert.ok(dead, "the rejected event is on mesh.events.dead");
|
||||||
|
assert.equal((dead as amqp.GetMessage).fields.routingKey, "module.test.boom");
|
||||||
|
await probe.close();
|
||||||
|
|
||||||
|
await c.close();
|
||||||
|
delete process.env.MESH_NODE;
|
||||||
|
delete process.env.MESH_MODULE;
|
||||||
|
});
|
||||||
|
|
||||||
|
test("an undecodable event body is dead-lettered, not looped and not silently swallowed", { skip: !url }, async () => {
|
||||||
|
const module = `poison-${Date.now()}`;
|
||||||
|
const c = await connectAmqp(url!);
|
||||||
|
process.env.MESH_NODE = "lab";
|
||||||
|
process.env.MESH_MODULE = module;
|
||||||
|
useBroker(() => c);
|
||||||
|
|
||||||
|
// Drain any earlier dead events so the assert below sees only this test's.
|
||||||
|
const drain = await amqp.connect(url!);
|
||||||
|
const dch = await drain.createChannel();
|
||||||
|
await dch.purgeQueue("mesh.events.dead");
|
||||||
|
|
||||||
|
let handlerRuns = 0;
|
||||||
|
await on("module.poison.raw", async () => void handlerRuns++);
|
||||||
|
await delay(200);
|
||||||
|
|
||||||
|
// Publish a body that is not JSON straight onto the events exchange — a malformed emitter. A
|
||||||
|
// confirm channel, waited on, so the broker has the message before the connection closes.
|
||||||
|
const raw = await amqp.connect(url!);
|
||||||
|
const rch = await raw.createConfirmChannel();
|
||||||
|
rch.publish("mesh.events", "module.poison.raw", Buffer.from("this is not json{"), {
|
||||||
|
persistent: true,
|
||||||
|
headers: { "x-event-id": "poison-1", "x-source": "bad", "x-node": "lab" },
|
||||||
|
});
|
||||||
|
await rch.waitForConfirms();
|
||||||
|
await raw.close();
|
||||||
|
|
||||||
|
// The handler never ran (the body never decoded), and the message is on the dead queue — set
|
||||||
|
// aside for inspection, not stuck redelivering forever.
|
||||||
|
await delay(600);
|
||||||
|
assert.equal(handlerRuns, 0, "a body that never decodes never reaches the handler");
|
||||||
|
const dead = await dch.get("mesh.events.dead", { noAck: true });
|
||||||
|
assert.ok(dead, "the poison event is retained on mesh.events.dead");
|
||||||
|
assert.equal((dead as amqp.GetMessage).fields.routingKey, "module.poison.raw");
|
||||||
|
await drain.close();
|
||||||
|
|
||||||
|
await c.close();
|
||||||
|
delete process.env.MESH_NODE;
|
||||||
|
delete process.env.MESH_MODULE;
|
||||||
|
});
|
||||||
|
|
||||||
|
function delay(ms: number): Promise<void> {
|
||||||
|
return new Promise((r) => setTimeout(r, ms));
|
||||||
|
}
|
||||||
|
|
||||||
|
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
||||||
|
const start = Date.now();
|
||||||
|
while (!cond()) {
|
||||||
|
if (Date.now() - start > ms) throw new Error("condition not met in time");
|
||||||
|
await delay(25);
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user