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 a2a687c..b63e1d8 100644 --- a/test/patient-connect.test.ts +++ b/test/patient-connect.test.ts @@ -1,6 +1,6 @@ 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 +41,21 @@ test("a non-Error value does not crash the classifier", () => { assert.equal(fatalBrokerReason("just a string"), null); assert.equal(fatalBrokerReason(undefined), null); }); + +// **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")), ["#"]); +});