Author SHA1 Message Date
jschoubben 38831c5c56 One consumer, one reader, however many patterns a module registers
A module has exactly one durable consumer, and each subscribe() started its own reader of it. Two
readers split the stream between them, and a reader that receives a message its own pattern does not
match acknowledges it — which is the right answer for a filter wider than anything registered, and
silent loss when the message was another handler's. The first module to subscribe twice would have
dropped roughly half of each kind of event with nothing reporting it.

Every registration is now dispatched from one reader, and a message is acknowledged once every handler
it is for has taken it.
2026-09-28 16:25:42 +02:00
mesh-admin fdad2f3268 Merge pull request 'A module answers the word the mesh asks: prepare' (#18) from feat/a-module-answers-prepare into main 2026-09-28 10:45:02 +00:00
jschoubben 81972a4995 A module answers the word the mesh asks: prepare
The runtime gains `prepare`, which brings this module's state to the shape this version needs and
exits (novox/hq ADR 0135). The entrypoints come from MESH_PREPARE, which a module's own image names
beside the entrypoints it already lists there — the module knows which of its files prepares its
state and nothing else could. No broker is connected: preparation runs before the version that would
use it. An empty list fails rather than passing quietly, because the mesh asks this only of a module
whose manifest says it prepares something, and exiting 0 would let that version serve against a
state nobody shaped.
2026-09-28 12:45:00 +02:00
mesh-admin 10e8191717 Merge pull request 'The runtime hears on its own inbox' (#17) from fix/the-runtime-hears-on-its-own-inbox into main 2026-09-28 02:30:24 +00:00
3 changed files with 126 additions and 20 deletions
+46 -19
View File
@@ -96,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 {
@@ -180,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;
}
};
},
@@ -209,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
@@ -225,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
View File
@@ -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) {
+38 -1
View File
@@ -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")), ["#"]);
});