Author SHA1 Message Date
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
3 changed files with 79 additions and 63 deletions
+12 -39
View File
@@ -96,10 +96,6 @@ export async function connectNats(
const js = conn.jetstream(); const js = conn.jetstream();
const subs: Subscription[] = []; 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; let closed = false;
return { return {
@@ -184,38 +180,21 @@ export async function connectNats(
* consumes (design 29 §3) — this binds to it and never creates one. A runtime that created * 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 * its own would be a module deciding its own delivery semantics, and its account cannot
* reach the JetStream API to do it anyway. * 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>( async subscribe<T>(
pattern: string, pattern: string,
handler: (env: Envelope<T>) => Promise<void>, handler: (env: Envelope<T>) => Promise<void>,
): Promise<() => void> { ): Promise<() => void> {
const listener = { pattern, handler: handler as (env: Envelope<unknown>) => Promise<void> };
listeners.push(listener);
if (!reading) {
const durable = `${cred.node ?? "?"}_${self}`; const durable = `${cred.node ?? "?"}_${self}`;
const consumer = await js.consumers.get("EVENTS", durable); const consumer = await js.consumers.get("EVENTS", durable);
const messages = await consumer.consume(); const messages = await consumer.consume();
reading = messages;
void (async () => { void (async () => {
for await (const msg of messages) { for await (const msg of messages) {
await deliver(msg, listeners); await deliver(msg, pattern, handler);
} }
})(); })();
}
return () => { return () => {
const at = listeners.indexOf(listener); void messages.close();
if (at >= 0) listeners.splice(at, 1);
if (listeners.length === 0 && reading) {
void reading.close();
reading = undefined;
}
}; };
}, },
@@ -230,20 +209,15 @@ export async function connectNats(
}; };
} }
/** /** Deliver one event, acknowledging only once a handler has taken it. */
* Deliver one event to every handler it is for, acknowledging only once each has taken it. async function deliver<T>(
*
* 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, msg: JsMsg,
listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[], pattern: string,
handler: (env: Envelope<T>) => Promise<void>,
): Promise<void> { ): Promise<void> {
let env: Envelope<unknown>; let env: Envelope<T>;
try { try {
env = toEnvelope<unknown>(msg); env = toEnvelope<T>(msg);
} catch { } catch {
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is // 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 // an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
@@ -251,16 +225,15 @@ async function deliver(
msg.term(); msg.term();
return; return;
} }
const forThis = listeners.filter((l) => topicMatches(l.pattern, env.key)); if (!topicMatches(pattern, env.key)) {
if (forThis.length === 0) { // The consumer's filters are the controller's, and may be wider than one subscription's
// The consumer's filters are the controller's, derived from what the module declared it // pattern when a module subscribes twice. Acknowledge what this handler is not for, or it
// consumes, and may be wider than anything it registered a handler for. Acknowledge it, or it
// would be redelivered until it expired. // would be redelivered until it expired.
msg.ack(); msg.ack();
return; return;
} }
try { try {
for (const l of forThis) await l.handler(env); await handler(env);
msg.ack(); msg.ack();
} catch { } catch {
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another // 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. // 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, // 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. // 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 // 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 // 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 // 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. // 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_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_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) // MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
import { readFileSync } from "node:fs"; import { readFileSync } from "node:fs";
@@ -177,12 +182,49 @@ async function runEntry(entrypoint: string): Promise<void> {
await import(pathToFileURL(entrypoint).href); 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> { async function main(): Promise<void> {
const [command, ...rest] = process.argv.slice(2); const [command, ...rest] = process.argv.slice(2);
if (command === "run") { if (command === "run") {
await runEntry(rest[0] ?? ""); await runEntry(rest[0] ?? "");
return; return;
} }
if (command === "prepare") {
await prepareState();
return;
}
if (command === "invoke") { if (command === "invoke") {
const [module, tool] = rest; const [module, tool] = rest;
if (!module || !tool) { if (!module || !tool) {
+18 -17
View File
@@ -1,6 +1,7 @@
import { spawn } from "node:child_process";
import { test } from "node:test"; import { test } from "node:test";
import assert from "node:assert/strict"; import assert from "node:assert/strict";
import { fatalBrokerReason, PinMismatchError, topicMatches } from "../src/broker-nats.ts"; import { fatalBrokerReason, PinMismatchError } from "../src/broker-nats.ts";
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must // 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 // give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
@@ -42,20 +43,20 @@ test("a non-Error value does not crash the classifier", () => {
assert.equal(fatalBrokerReason(undefined), null); assert.equal(fatalBrokerReason(undefined), null);
}); });
// **One durable consumer feeds one reader, however many patterns a module registers.** // **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
// A module has exactly one consumer, so two readers of it would each take half the messages — and a // image that names nothing was built wrong, and exiting 0 would let that version serve against a
// reader that received one its own pattern does not match acknowledges it, which is right for a // state nobody shaped.
// filter wider than anything registered and silent loss when it is another handler's. The matching is test("preparing with nothing named fails rather than passing quietly", async () => {
// therefore pure and tested as such: what a message is for is decided by the patterns registered, not const runtime = new URL("../dist/main.js", import.meta.url).pathname;
// by which reader happened to fetch it. const ran = await new Promise<{ code: number | null; said: string }>((resolve) => {
test("a message is for every pattern that matches it, and nothing else", () => { const child = spawn(process.execPath, [runtime, "prepare"], {
const registered = ["mesh-build-machine.built", "mesh-controller.built-before"]; env: { ...process.env, MESH_PREPARE: "" },
const matched = (key: string) => registered.filter((p) => topicMatches(p, key)); });
assert.deepEqual(matched("mesh-build-machine.built"), ["mesh-build-machine.built"]); let said = "";
assert.deepEqual(matched("mesh-controller.built-before"), ["mesh-controller.built-before"]); child.stderr.on("data", (chunk) => (said += String(chunk)));
// Nothing registered for it: the consumer's filter is the controller's and may be wider. child.on("close", (code) => resolve({ code, said }));
assert.deepEqual(matched("mesh-catalog.upgraded"), []); });
// And a handler that asked for everything gets both, which is what the audit logger does. assert.notEqual(ran.code, 0, "a module that prepares nothing exited 0, so its version would serve");
assert.deepEqual(["#"].filter((p) => topicMatches(p, "mesh-controller.built-before")), ["#"]); assert.match(ran.said, /MESH_PREPARE/);
}); });