Task 3.7 of novox/hq ADR 0116, and the whole of the sdk's diff for the bus change. Three comments said AMQP where they meant 'message headers' and 'a broker client'; the code never spoke it, which is why no module is rebuilt for any of this (ADR 0039).
89 lines
3.8 KiB
TypeScript
89 lines
3.8 KiB
TypeScript
// Events — a module logs its activity onto the mesh broker, and any module reacts. The lighter
|
|
// sibling of provisioning: provisioning is 1:1 and credentialed (a provider creates a resource for
|
|
// one consumer); an event is 1:many and broadcast (a module emits, any number listen, no
|
|
// credential — just the broker's topic routing). Declared as emits/consumes on the manifest so the
|
|
// mesh knows the event graph.
|
|
//
|
|
// Identity and provenance ride as message headers, not in the body (novox/hq ADR 0042): a consumer —
|
|
// or the broker, or an audit tool — reads who/when/what without parsing the payload, and the body
|
|
// is only the domain payload. This surface hides the header/exchange/queue mechanics; a module
|
|
// names a type and a body and never sees the wire.
|
|
|
|
import { randomUUID } from "node:crypto";
|
|
|
|
import { broker } from "../messaging/index.js";
|
|
import type { Envelope, EventHeaders } from "../contracts/index.js";
|
|
|
|
/** An event on the mesh: a topic key, its provenance, and a body — the shape a handler receives. */
|
|
export interface Event<T = unknown> {
|
|
/** The routing key, e.g. "module.umami.site.created". Dotted, so listeners can match by prefix. */
|
|
readonly type: string;
|
|
/** The unique event id (x-event-id) — at-least-once delivery means a handler must dedup on it. */
|
|
readonly id: string;
|
|
/** The emitting module, context or node (x-source). */
|
|
readonly source: string;
|
|
/** The node it was emitted from (x-node). */
|
|
readonly node: string;
|
|
/** RFC-3339 emit time (x-time). */
|
|
readonly at: string;
|
|
/** The event or command that caused this one, if any (x-causation-id) — tracing. */
|
|
readonly causationId?: string;
|
|
/** A version tag for the body's shape, if the emitter set one (x-schema). */
|
|
readonly schema?: string;
|
|
readonly body: T;
|
|
}
|
|
|
|
/** Extra provenance a caller may attach when emitting. */
|
|
export interface EmitOptions {
|
|
/** The event or command that caused this one (x-causation-id). */
|
|
readonly causationId?: string;
|
|
/** A version tag for the body's shape (x-schema). */
|
|
readonly schema?: string;
|
|
}
|
|
|
|
/**
|
|
* Emit an event. Source and node come from the environment the runtime set for the module
|
|
* (MESH_MODULE, MESH_NODE), so a module names only the type and the body; the sdk stamps the
|
|
* ADR 0042 headers (id, source, node, time) and the runtime rides them on the broker.
|
|
*/
|
|
export async function emit<T>(type: string, body: T, opts: EmitOptions = {}): Promise<void> {
|
|
const source = process.env.MESH_MODULE ?? "unknown";
|
|
const node = process.env.MESH_NODE ?? "unknown";
|
|
const headers: EventHeaders = {
|
|
"x-event-id": randomUUID(),
|
|
"x-source": source,
|
|
"x-node": node,
|
|
"x-time": new Date().toISOString(),
|
|
"content-type": "application/json",
|
|
...(opts.causationId ? { "x-causation-id": opts.causationId } : {}),
|
|
...(opts.schema ? { "x-schema": opts.schema } : {}),
|
|
};
|
|
await broker().publish<T>({ key: type, node, body, headers });
|
|
}
|
|
|
|
/**
|
|
* React to events whose type matches a topic pattern (`*` one segment, `#` any). The audit logger
|
|
* is just `on("#", …)`. The handler receives the reconstructed event — its metadata read back from
|
|
* the headers, its body the domain payload.
|
|
*/
|
|
export async function on<T>(pattern: string, handler: (event: Event<T>) => Promise<void>): Promise<() => void> {
|
|
return broker().subscribe<T>(pattern, async (envelope) => {
|
|
await handler(fromEnvelope<T>(envelope));
|
|
});
|
|
}
|
|
|
|
/** Rebuild the Event a handler sees from a broker envelope's headers (ADR 0042) and body. */
|
|
function fromEnvelope<T>(env: Envelope<T>): Event<T> {
|
|
const h = env.headers ?? ({} as EventHeaders);
|
|
return {
|
|
type: env.key,
|
|
id: h["x-event-id"] ?? "",
|
|
source: h["x-source"] ?? "unknown",
|
|
node: h["x-node"] ?? env.node ?? "unknown",
|
|
at: h["x-time"] ?? "",
|
|
causationId: h["x-causation-id"],
|
|
schema: h["x-schema"],
|
|
body: env.body,
|
|
};
|
|
}
|