A refused tool subscription is said, never fatal (hq issue 218) #46
@@ -277,17 +277,27 @@ export async function connectNats(
|
|||||||
const sub = conn.subscribe(subject, queue ? { queue } : {});
|
const sub = conn.subscribe(subject, queue ? { queue } : {});
|
||||||
subs.push(sub);
|
subs.push(sub);
|
||||||
void (async () => {
|
void (async () => {
|
||||||
for await (const msg of sub) {
|
// **A refused subscription is said, never fatal** (novox/hq 04-ISSUES/218, as 217 for the
|
||||||
let reply: { result?: Res; error?: string; node?: string };
|
// announcements). The grants are the mesh's word on what this account may answer; a subject
|
||||||
try {
|
// they leave out — a seat claimed here and held elsewhere — costs that subject, never the
|
||||||
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
|
// module's other tools, its handlers and its provisioning. Unhandled, the refusal ended the
|
||||||
} catch (err) {
|
// process and a module's runtime crash-looped on 2026-10-04.
|
||||||
// The caller is told, rather than left to time out: a handler that threw is a
|
try {
|
||||||
// different failure from a tool nobody serves, and only one of them is worth retrying.
|
for await (const msg of sub) {
|
||||||
reply = { error: err instanceof Error ? err.message : String(err) };
|
let reply: { result?: Res; error?: string; node?: string };
|
||||||
|
try {
|
||||||
|
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
|
||||||
|
} catch (err) {
|
||||||
|
// The caller is told, rather than left to time out: a handler that threw is a
|
||||||
|
// different failure from a tool nobody serves, and only one of them is worth retrying.
|
||||||
|
reply = { error: err instanceof Error ? err.message : String(err) };
|
||||||
|
}
|
||||||
|
if (node) reply.node = node;
|
||||||
|
msg.respond(sc.encode(JSON.stringify(reply)));
|
||||||
}
|
}
|
||||||
if (node) reply.node = node;
|
} catch (err) {
|
||||||
msg.respond(sc.encode(JSON.stringify(reply)));
|
console.log(`[mesh-tools] the bus refused ${subject}: ${err instanceof Error ? err.message : String(err)}; ` +
|
||||||
|
"not served here, and the rest serves on");
|
||||||
}
|
}
|
||||||
})();
|
})();
|
||||||
return () => sub.unsubscribe();
|
return () => sub.unsubscribe();
|
||||||
|
|||||||
+1
-1
@@ -3,7 +3,7 @@ authorization {
|
|||||||
users = [
|
users = [
|
||||||
{ user: "runtime", password: "runtime", permissions: {
|
{ user: "runtime", password: "runtime", permissions: {
|
||||||
publish: { allow: [">"] }
|
publish: { allow: [">"] }
|
||||||
subscribe: { allow: [">"], deny: ["$SRV.PING.>"] }
|
subscribe: { allow: [">"], deny: ["$SRV.PING.>", "mesh.seat.held-elsewhere.>"] }
|
||||||
} }
|
} }
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,6 +9,7 @@
|
|||||||
import assert from "node:assert/strict";
|
import assert from "node:assert/strict";
|
||||||
import { test } from "node:test";
|
import { test } from "node:test";
|
||||||
import { connectNats } from "../dist/broker-nats.js";
|
import { connectNats } from "../dist/broker-nats.js";
|
||||||
|
import { callTool } from "../dist/client.js";
|
||||||
|
|
||||||
const url = process.env.MESH_TEST_REFUSING_NATS;
|
const url = process.env.MESH_TEST_REFUSING_NATS;
|
||||||
|
|
||||||
@@ -35,3 +36,33 @@ test("a refused discovery subscription is logged and the runtime serves on", asy
|
|||||||
await bus.close();
|
await bus.close();
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// novox/hq issue 218: a tool subject the grants leave out — a seat claimed here and held elsewhere —
|
||||||
|
// is refused, said, and the module's other tools still answer.
|
||||||
|
test("a refused tool subscription is logged and the module's other tools answer", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_REFUSING_NATS unset");
|
||||||
|
const bus = await connectNats({ url, user: "runtime", password: "runtime", module: "alpha", node: "anchor" });
|
||||||
|
const asker = await connectNats({ url, user: "runtime", password: "runtime", module: "console", node: "workstation" });
|
||||||
|
const said: string[] = [];
|
||||||
|
const log = console.log;
|
||||||
|
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
||||||
|
const crashed: unknown[] = [];
|
||||||
|
const onRejection = (e: unknown) => crashed.push(e);
|
||||||
|
process.on("unhandledRejection", onRejection);
|
||||||
|
try {
|
||||||
|
const refused = await bus.handleSubject!("mesh.seat.held-elsewhere.tool.databases", async () => ({ seat: true }));
|
||||||
|
const stop = await bus.handle("alpha.ping", async () => ({ pong: true }));
|
||||||
|
for (let i = 0; i < 50 && !said.some((s) => s.includes("the bus refused mesh.seat.held-elsewhere")); i++) await new Promise((r) => setTimeout(r, 50));
|
||||||
|
console.log = log;
|
||||||
|
assert.ok(said.some((s) => /the bus refused mesh\.seat\.held-elsewhere\.tool\.databases.*serves on/.test(s)), said.join("\n"));
|
||||||
|
assert.equal(crashed.length, 0, `the refusal escaped: ${String(crashed[0])}`);
|
||||||
|
assert.deepEqual((await callTool(asker, "alpha.ping@anchor", {})).result, { pong: true });
|
||||||
|
stop();
|
||||||
|
refused();
|
||||||
|
} finally {
|
||||||
|
console.log = log;
|
||||||
|
process.off("unhandledRejection", onRejection);
|
||||||
|
await asker.close();
|
||||||
|
await bus.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user