diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index eee7b03..62c5e51 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -81,6 +81,10 @@ export async function connectAmqp( opts: { assumeExchanges?: boolean } = {}, ): Promise { const cred: Credential = typeof target === "string" ? { url: target } : target; + // This module's own name, for turning a local event name into this bus's routing key. From the + // credential where the mesh issued one, and from the environment for an ad-hoc client that has no + // credential of its own — the same two places subscribe already looks. + const self = cred.module ?? process.env.MESH_MODULE ?? ""; const conn = cred.fingerprint ? await amqp.connect(cred.url, await pinnedOptions(cred.url, cred.fingerprint)) : await amqp.connect(cred.url); @@ -140,7 +144,7 @@ export async function connectAmqp( let eventConsumerTag: string | undefined; async function dispatchEvent(msg: amqp.ConsumeMessage): Promise { - const key = msg.fields.routingKey; + const key = localKeyFor(msg.fields.routingKey); let env: Envelope; try { env = toEnvelope(msg); @@ -220,7 +224,7 @@ export async function connectAmqp( await new Promise((resolve, reject) => { ch.publish( EVENTS_EXCHANGE, - env.key, + routingKeyFor(env.key, self), Buffer.from(JSON.stringify(env.body)), { persistent: true, @@ -253,7 +257,7 @@ export async function connectAmqp( } eventQueue = name; } - await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, pattern); + await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, bindingFor(pattern)); const sub: EventSub = { pattern, handler: handler as EventSub["handler"] }; eventSubs.push(sub); if (!eventConsumerTag) { @@ -269,7 +273,7 @@ export async function connectAmqp( } const { queue } = await ch.assertQueue("", { exclusive: true }); - await ch.bindQueue(queue, EVENTS_EXCHANGE, pattern); + await ch.bindQueue(queue, EVENTS_EXCHANGE, bindingFor(pattern)); const consumer = await ch.consume(queue, (msg) => { if (!msg) return; void (async () => { @@ -376,7 +380,9 @@ function topicMatches(pattern: string, key: string): boolean { function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { while (pi < p.length) { const tok = p[pi]; - if (tok === "#") { + // `**` is the mesh's wildcard for the rest of a name; `#` is this bus's, accepted so a pattern + // written either way behaves the same while both buses ship (novox/hq design 29 §1). + if (tok === "#" || tok === "**") { if (pi === p.length - 1) return true; // trailing # swallows the rest, including nothing for (let skip = ki; skip <= k.length; skip++) { if (matchFrom(p, pi + 1, k, skip)) return true; @@ -390,3 +396,27 @@ function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { } return ki === k.length; } + +/** This bus spells an event as a routing key that repeats the emitter's name: `module..`. + * A module names its events locally and the mesh derives where they land (novox/hq design 29 §1), so + * the mapping lives here rather than in every module. + * + * **Why it exists at all.** Until 04-ISSUES/127 every module passed the routing key itself, which + * worked on this bus and derived into a namespace nobody owns on the one being built. Converting the + * modules to local names without this would have broken the mesh that is actually running. */ +function routingKeyFor(key: string, self: string): string { + return key.startsWith("module.") ? key : `module.${self}.${key}`; +} + +/** A local pattern as this bus's binding. `**` is the mesh's wildcard for the rest of a name; here + * that is `#`, and on the bus being built it is `>`. Neither spelling appears in a manifest. */ +function bindingFor(pattern: string): string { + const here = pattern.split(".").map((part) => (part === "**" ? "#" : part)).join("."); + if (here === "#") return "#"; + return here.startsWith("module.") ? here : `module.${here}`; +} + +/** A routing key as the local name a handler and a manifest both use: the emitter and the event. */ +function localKeyFor(routingKey: string): string { + return routingKey.startsWith("module.") ? routingKey.slice("module.".length) : routingKey; +} diff --git a/src/broker-nats.ts b/src/broker-nats.ts index bd82bbb..f6f4ff7 100644 --- a/src/broker-nats.ts +++ b/src/broker-nats.ts @@ -255,11 +255,23 @@ function toEnvelope(msg: JsMsg): Envelope { }; } -/** The event key inside a module's event subject. */ +/** The event key a module sees: the emitter and the event, which is exactly how its manifest names + * what it consumes (novox/hq design 29 §1, 04-ISSUES/127). + * + * **One vocabulary for the declaration and the handler.** This returned the event name alone, so a + * manifest declaring `consumes: builder.built` produced a handler pattern that could never match + * the key it was compared against — and a module consuming the same event from two emitters could + * not tell them apart except by reading a header. The subject already carries the emitter; naming it + * here makes a mismatch between manifest and code a typo rather than a category error. */ function keyFromSubject(subject: string): string { const marker = ".event."; const at = subject.indexOf(marker); - return at < 0 ? subject : subject.slice(at + marker.length); + if (at < 0) return subject; + const event = subject.slice(at + marker.length); + // `mesh.mod..event.…` — the emitter is the token before the marker. + const before = subject.slice(0, at).split("."); + const emitter = before[before.length - 1]; + return emitter ? `${emitter}.${event}` : event; } /** A module's own event subject. Derived, never taken from the caller: the module names its @@ -325,7 +337,9 @@ export function topicMatches(pattern: string, key: string): boolean { function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { if (pi === p.length) return ki === k.length; - if (p[pi] === "#") { + // `**` is the mesh's wildcard for the rest of a name; `#` is the old bus's, accepted so a pattern + // written either way behaves the same while both buses ship (novox/hq design 29 §1). + if (p[pi] === "#" || p[pi] === "**") { for (let skip = ki; skip <= k.length; skip++) { if (matchFrom(p, pi + 1, k, skip)) return true; }