Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
38831c5c56 | ||
|
|
fdad2f3268 | ||
|
|
81972a4995 | ||
|
|
10e8191717 | ||
|
|
d703cebff4 | ||
|
|
46b56d53a6 |
+50
-19
@@ -85,6 +85,10 @@ export async function connectNats(
|
||||
pass: cred.password,
|
||||
name: `${cred.node ?? "?"}.${self}`,
|
||||
tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined,
|
||||
// Its own inbox, not a random one: every user's inbox is private to it (design 25 §4), and the
|
||||
// grant names `_INBOX.<user>.>` — a reply space the client invented would be refused, and with
|
||||
// it every pull for the next message and every answer to a tool call.
|
||||
inboxPrefix: cred.user ? `_INBOX.${cred.user}` : undefined,
|
||||
// Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on
|
||||
// the mesh to exit. `close()` stays the only thing that ends the connection.
|
||||
maxReconnectAttempts: -1,
|
||||
@@ -92,6 +96,10 @@ export async function connectNats(
|
||||
const js = conn.jetstream();
|
||||
|
||||
const subs: Subscription[] = [];
|
||||
// Every registration, and the one reader that dispatches to them. A module has one durable
|
||||
// consumer; the loop belongs to the connection rather than to a subscription.
|
||||
const listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[] = [];
|
||||
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined;
|
||||
let closed = false;
|
||||
|
||||
return {
|
||||
@@ -176,21 +184,38 @@ export async function connectNats(
|
||||
* consumes (design 29 §3) — this binds to it and never creates one. A runtime that created
|
||||
* its own would be a module deciding its own delivery semantics, and its account cannot
|
||||
* reach the JetStream API to do it anyway.
|
||||
*
|
||||
* **One consumer, one loop, however many patterns a module registers.** A module has exactly one
|
||||
* durable consumer, so two loops reading it would each take half the messages — and a loop that
|
||||
* received one its own pattern does not match acknowledges it, which is the right answer for a
|
||||
* filter wider than anything registered and silent loss when it is another handler's. Every
|
||||
* registration is therefore dispatched from one reader, and a message is acknowledged once every
|
||||
* handler it is for has taken it.
|
||||
*/
|
||||
async subscribe<T>(
|
||||
pattern: string,
|
||||
handler: (env: Envelope<T>) => Promise<void>,
|
||||
): Promise<() => void> {
|
||||
const durable = `${cred.node ?? "?"}_${self}`;
|
||||
const consumer = await js.consumers.get("EVENTS", durable);
|
||||
const messages = await consumer.consume();
|
||||
void (async () => {
|
||||
for await (const msg of messages) {
|
||||
await deliver(msg, pattern, handler);
|
||||
}
|
||||
})();
|
||||
const listener = { pattern, handler: handler as (env: Envelope<unknown>) => Promise<void> };
|
||||
listeners.push(listener);
|
||||
if (!reading) {
|
||||
const durable = `${cred.node ?? "?"}_${self}`;
|
||||
const consumer = await js.consumers.get("EVENTS", durable);
|
||||
const messages = await consumer.consume();
|
||||
reading = messages;
|
||||
void (async () => {
|
||||
for await (const msg of messages) {
|
||||
await deliver(msg, listeners);
|
||||
}
|
||||
})();
|
||||
}
|
||||
return () => {
|
||||
void messages.close();
|
||||
const at = listeners.indexOf(listener);
|
||||
if (at >= 0) listeners.splice(at, 1);
|
||||
if (listeners.length === 0 && reading) {
|
||||
void reading.close();
|
||||
reading = undefined;
|
||||
}
|
||||
};
|
||||
},
|
||||
|
||||
@@ -205,15 +230,20 @@ export async function connectNats(
|
||||
};
|
||||
}
|
||||
|
||||
/** Deliver one event, acknowledging only once a handler has taken it. */
|
||||
async function deliver<T>(
|
||||
/**
|
||||
* Deliver one event to every handler it is for, acknowledging only once each has taken it.
|
||||
*
|
||||
* Several registrations share one durable consumer, so matching happens here rather than by having
|
||||
* each registration read the stream: two readers of one consumer would split it between them, and a
|
||||
* message that reached the wrong one would be acknowledged as not-for-me and lost.
|
||||
*/
|
||||
async function deliver(
|
||||
msg: JsMsg,
|
||||
pattern: string,
|
||||
handler: (env: Envelope<T>) => Promise<void>,
|
||||
listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[],
|
||||
): Promise<void> {
|
||||
let env: Envelope<T>;
|
||||
let env: Envelope<unknown>;
|
||||
try {
|
||||
env = toEnvelope<T>(msg);
|
||||
env = toEnvelope<unknown>(msg);
|
||||
} catch {
|
||||
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is
|
||||
// an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
|
||||
@@ -221,15 +251,16 @@ async function deliver<T>(
|
||||
msg.term();
|
||||
return;
|
||||
}
|
||||
if (!topicMatches(pattern, env.key)) {
|
||||
// The consumer's filters are the controller's, and may be wider than one subscription's
|
||||
// pattern when a module subscribes twice. Acknowledge what this handler is not for, or it
|
||||
const forThis = listeners.filter((l) => topicMatches(l.pattern, env.key));
|
||||
if (forThis.length === 0) {
|
||||
// The consumer's filters are the controller's, derived from what the module declared it
|
||||
// consumes, and may be wider than anything it registered a handler for. Acknowledge it, or it
|
||||
// would be redelivered until it expired.
|
||||
msg.ack();
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await handler(env);
|
||||
for (const l of forThis) await l.handler(env);
|
||||
msg.ack();
|
||||
} catch {
|
||||
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
|
||||
|
||||
+42
@@ -5,6 +5,10 @@
|
||||
// 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.
|
||||
// mesh-tools prepare bring this module's state to the shape this version needs and exit
|
||||
// — the runtime's answer to the word the mesh asks every module
|
||||
// (novox/hq ADR 0135). The entrypoints come from MESH_PREPARE, which
|
||||
// the module's own image names beside MESH_TOOL_MODULES.
|
||||
// mesh-tools run ENTRYPOINT run one compiled module entrypoint to completion and exit — the
|
||||
// runtime side of a run-once step (novox/hq ADR 0052). It imports
|
||||
// the given entrypoint, whose top-level code does its work — seed a
|
||||
@@ -18,6 +22,7 @@
|
||||
// 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_PREPARE /path/a,/path/b,… compiled entrypoints that prepare this module's state
|
||||
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
|
||||
|
||||
import { readFileSync } from "node:fs";
|
||||
@@ -177,12 +182,49 @@ async function runEntry(entrypoint: string): Promise<void> {
|
||||
await import(pathToFileURL(entrypoint).href);
|
||||
}
|
||||
|
||||
/**
|
||||
* Bring this module's state to the shape this version needs, and exit — the runtime's answer to the
|
||||
* one word the mesh asks every module (novox/hq ADR 0135).
|
||||
*
|
||||
* The entrypoints come from `MESH_PREPARE`, which a module's own image sets beside the entrypoints it
|
||||
* already lists there: the module knows which of its files prepares its state, and nothing else could.
|
||||
* Each is imported in the order given, to completion, with no broker — preparation runs before the
|
||||
* version that would use it, so there is nothing yet to talk to.
|
||||
*
|
||||
* **An empty list is a failure, not a no-op.** The mesh only asks this of a module whose manifest says
|
||||
* it prepares something; a module that says so and names nothing has been built wrong, and exiting 0
|
||||
* would let that version serve against a state nobody shaped.
|
||||
*/
|
||||
async function prepareState(): Promise<void> {
|
||||
const named = (process.env.MESH_PREPARE ?? "")
|
||||
.split(",")
|
||||
.map((entry) => entry.trim())
|
||||
.filter((entry) => entry !== "");
|
||||
if (named.length === 0) {
|
||||
console.error(
|
||||
"mesh-tools prepare: this module was asked to prepare its state and its image names nothing " +
|
||||
"to do it with — set MESH_PREPARE to the compiled entrypoint(s) that prepare it, the way " +
|
||||
"MESH_TOOL_MODULES names the ones it serves",
|
||||
);
|
||||
process.exit(1);
|
||||
}
|
||||
for (const entrypoint of named) {
|
||||
console.log(`[mesh-tools] preparing with ${entrypoint}`);
|
||||
await runEntry(entrypoint);
|
||||
}
|
||||
console.log(`[mesh-tools] prepared: ${named.length} entrypoint(s) ran to completion`);
|
||||
}
|
||||
|
||||
async function main(): Promise<void> {
|
||||
const [command, ...rest] = process.argv.slice(2);
|
||||
if (command === "run") {
|
||||
await runEntry(rest[0] ?? "");
|
||||
return;
|
||||
}
|
||||
if (command === "prepare") {
|
||||
await prepareState();
|
||||
return;
|
||||
}
|
||||
if (command === "invoke") {
|
||||
const [module, tool] = rest;
|
||||
if (!module || !tool) {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { spawn } from "node:child_process";
|
||||
import { test } from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import { fatalBrokerReason, PinMismatchError } from "../src/broker-nats.ts";
|
||||
import { fatalBrokerReason, PinMismatchError, topicMatches } from "../src/broker-nats.ts";
|
||||
|
||||
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must
|
||||
// give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
|
||||
@@ -41,3 +42,39 @@ test("a non-Error value does not crash the classifier", () => {
|
||||
assert.equal(fatalBrokerReason("just a string"), null);
|
||||
assert.equal(fatalBrokerReason(undefined), null);
|
||||
});
|
||||
|
||||
// **A module asked to prepare its state and naming nothing is a failure, not a no-op** (novox/hq
|
||||
// ADR 0135). The mesh asks this only of a module whose manifest says it prepares something, so an
|
||||
// image that names nothing was built wrong, and exiting 0 would let that version serve against a
|
||||
// state nobody shaped.
|
||||
test("preparing with nothing named fails rather than passing quietly", async () => {
|
||||
const runtime = new URL("../dist/main.js", import.meta.url).pathname;
|
||||
const ran = await new Promise<{ code: number | null; said: string }>((resolve) => {
|
||||
const child = spawn(process.execPath, [runtime, "prepare"], {
|
||||
env: { ...process.env, MESH_PREPARE: "" },
|
||||
});
|
||||
let said = "";
|
||||
child.stderr.on("data", (chunk) => (said += String(chunk)));
|
||||
child.on("close", (code) => resolve({ code, said }));
|
||||
});
|
||||
assert.notEqual(ran.code, 0, "a module that prepares nothing exited 0, so its version would serve");
|
||||
assert.match(ran.said, /MESH_PREPARE/);
|
||||
});
|
||||
|
||||
// **One durable consumer feeds one reader, however many patterns a module registers.**
|
||||
//
|
||||
// A module has exactly one consumer, so two readers of it would each take half the messages — and a
|
||||
// reader that received one its own pattern does not match acknowledges it, which is right for a
|
||||
// filter wider than anything registered and silent loss when it is another handler's. The matching is
|
||||
// therefore pure and tested as such: what a message is for is decided by the patterns registered, not
|
||||
// by which reader happened to fetch it.
|
||||
test("a message is for every pattern that matches it, and nothing else", () => {
|
||||
const registered = ["mesh-build-machine.built", "mesh-controller.built-before"];
|
||||
const matched = (key: string) => registered.filter((p) => topicMatches(p, key));
|
||||
assert.deepEqual(matched("mesh-build-machine.built"), ["mesh-build-machine.built"]);
|
||||
assert.deepEqual(matched("mesh-controller.built-before"), ["mesh-controller.built-before"]);
|
||||
// Nothing registered for it: the consumer's filter is the controller's and may be wider.
|
||||
assert.deepEqual(matched("mesh-catalog.upgraded"), []);
|
||||
// And a handler that asked for everything gets both, which is what the audit logger does.
|
||||
assert.deepEqual(["#"].filter((p) => topicMatches(p, "mesh-controller.built-before")), ["#"]);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user