From 38831c5c56b4123e9eb48eec0b0ee32c61f8ff19 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 28 Sep 2026 16:15:03 +0200 Subject: [PATCH] One consumer, one reader, however many patterns a module registers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- src/broker-nats.ts | 65 +++++++++++++++++++++++++----------- test/patient-connect.test.ts | 20 ++++++++++- 2 files changed, 65 insertions(+), 20 deletions(-) diff --git a/src/broker-nats.ts b/src/broker-nats.ts index 6fc48a0..aae89c8 100644 --- a/src/broker-nats.ts +++ b/src/broker-nats.ts @@ -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) => Promise }[] = []; + let reading: Awaited>["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( pattern: string, handler: (env: Envelope) => Promise, ): 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) => Promise }; + 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( +/** + * 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) => Promise, + listeners: { pattern: string; handler: (env: Envelope) => Promise }[], ): Promise { - let env: Envelope; + let env: Envelope; try { - env = toEnvelope(msg); + env = toEnvelope(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( 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 diff --git a/test/patient-connect.test.ts b/test/patient-connect.test.ts index 5815dad..81f4c5d 100644 --- a/test/patient-connect.test.ts +++ b/test/patient-connect.test.ts @@ -1,7 +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 @@ -60,3 +60,21 @@ test("preparing with nothing named fails rather than passing quietly", async () 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")), ["#"]); +}); -- 2.54.0