events: metadata rides as headers, not in the body (ADR 0047)
emit stamps the ADR 0047 headers — x-event-id, x-source, x-node, x-time, content-type, and optional x-causation-id / x-schema — and publishes the body as only the domain payload. on() reconstructs the Event from those headers. Event gains id (the x-event-id a consumer dedups on) plus the optional causation/schema. EventHeaders joins the contracts spine. Supersedes the first cut that carried source/node/time in the body. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
+29
-1
@@ -44,9 +44,37 @@ export interface ToolDefinition {
|
|||||||
readonly run: (args: Readonly<Record<string, unknown>>) => Promise<unknown>;
|
readonly run: (args: Readonly<Record<string, unknown>>) => Promise<unknown>;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** A message crossing the broker: a routing key and a JSON body, per-node addressed. */
|
/**
|
||||||
|
* The metadata that rides an event as AMQP headers (novox/hq ADR 0047). An event's identity and
|
||||||
|
* provenance live here, not in the body, so a consumer — or the broker, or an audit tool — reads
|
||||||
|
* who/when/what without parsing the payload. An unknown `x-` header is ignored, not refused: an
|
||||||
|
* event is observed by parties that need not all understand every header.
|
||||||
|
*/
|
||||||
|
export interface EventHeaders {
|
||||||
|
/** A unique id — for dedup and audit (delivery is at-least-once). */
|
||||||
|
readonly "x-event-id": string;
|
||||||
|
/** The emitter: the module, context or node name. */
|
||||||
|
readonly "x-source": string;
|
||||||
|
/** The node it was emitted from. */
|
||||||
|
readonly "x-node": string;
|
||||||
|
/** Emit time, RFC-3339. */
|
||||||
|
readonly "x-time": string;
|
||||||
|
/** Always `application/json`. */
|
||||||
|
readonly "content-type": string;
|
||||||
|
/** The event or command that caused this one — tracing. */
|
||||||
|
readonly "x-causation-id"?: string;
|
||||||
|
/** A version of the body's shape, so a body evolves without silent misreads. */
|
||||||
|
readonly "x-schema"?: string;
|
||||||
|
readonly [header: string]: string | undefined;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A message crossing the broker: a routing key and a JSON body, per-node addressed. For an event,
|
||||||
|
* `headers` carries the ADR 0047 metadata; plain request/reply transport leaves it absent.
|
||||||
|
*/
|
||||||
export interface Envelope<T = unknown> {
|
export interface Envelope<T = unknown> {
|
||||||
readonly key: string;
|
readonly key: string;
|
||||||
readonly node: string;
|
readonly node: string;
|
||||||
readonly body: T;
|
readonly body: T;
|
||||||
|
readonly headers?: EventHeaders;
|
||||||
}
|
}
|
||||||
|
|||||||
+58
-18
@@ -4,45 +4,85 @@
|
|||||||
// credential — just the broker's topic routing). Declared as emits/consumes on the manifest so the
|
// credential — just the broker's topic routing). Declared as emits/consumes on the manifest so the
|
||||||
// mesh knows the event graph.
|
// mesh knows the event graph.
|
||||||
//
|
//
|
||||||
// This is a thin, audit-ready surface over the broker's publish/subscribe: every event carries who
|
// Identity and provenance ride as AMQP headers, not in the body (novox/hq ADR 0047): a consumer —
|
||||||
// emitted it, on which node, and when — so a logger consuming `#` can write a real audit trail.
|
// 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 { broker } from "../messaging/index.js";
|
||||||
|
import type { Envelope, EventHeaders } from "../contracts/index.js";
|
||||||
|
|
||||||
/** An event on the mesh: a topic key, its source, and a body — with the metadata audit needs. */
|
/** An event on the mesh: a topic key, its provenance, and a body — the shape a handler receives. */
|
||||||
export interface Event<T = unknown> {
|
export interface Event<T = unknown> {
|
||||||
/** The routing key, e.g. "module.umami.site.created". Dotted, so listeners can match by prefix. */
|
/** The routing key, e.g. "module.umami.site.created". Dotted, so listeners can match by prefix. */
|
||||||
readonly type: string;
|
readonly type: string;
|
||||||
/** The emitting module. */
|
/** 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;
|
readonly source: string;
|
||||||
/** The node it was emitted from. */
|
/** The node it was emitted from (x-node). */
|
||||||
readonly node: string;
|
readonly node: string;
|
||||||
/** ISO-8601 emit time. */
|
/** RFC-3339 emit time (x-time). */
|
||||||
readonly at: string;
|
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;
|
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
|
* 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.
|
* (MESH_MODULE, MESH_NODE), so a module names only the type and the body; the sdk stamps the
|
||||||
|
* ADR 0047 headers (id, source, node, time) and the runtime rides them on the broker.
|
||||||
*/
|
*/
|
||||||
export async function emit<T>(type: string, body: T): Promise<void> {
|
export async function emit<T>(type: string, body: T, opts: EmitOptions = {}): Promise<void> {
|
||||||
const event: Event<T> = {
|
const source = process.env.MESH_MODULE ?? "unknown";
|
||||||
type,
|
const node = process.env.MESH_NODE ?? "unknown";
|
||||||
source: process.env.MESH_MODULE ?? "unknown",
|
const headers: EventHeaders = {
|
||||||
node: process.env.MESH_NODE ?? "unknown",
|
"x-event-id": randomUUID(),
|
||||||
at: new Date().toISOString(),
|
"x-source": source,
|
||||||
body,
|
"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<Event<T>>({ key: type, node: event.node, body: event });
|
await broker().publish<T>({ key: type, node, body, headers });
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* React to events whose type matches a topic pattern (`*` one segment, `#` any). The audit logger
|
* React to events whose type matches a topic pattern (`*` one segment, `#` any). The audit logger
|
||||||
* is just `on("#", …)`. The handler receives the whole event, metadata included.
|
* 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> {
|
export async function on<T>(pattern: string, handler: (event: Event<T>) => Promise<void>): Promise<() => void> {
|
||||||
return broker().subscribe<Event<T>>(pattern, async (envelope) => {
|
return broker().subscribe<T>(pattern, async (envelope) => {
|
||||||
await handler(envelope.body);
|
await handler(fromEnvelope<T>(envelope));
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Rebuild the Event a handler sees from a broker envelope's headers (ADR 0047) 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,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|||||||
@@ -3,9 +3,9 @@
|
|||||||
// binding is provided by the runtime that hosts a module's code — the sdk defines the contract so
|
// binding is provided by the runtime that hosts a module's code — the sdk defines the contract so
|
||||||
// module code, the tool runtime and provisioners all speak it the same way.
|
// module code, the tool runtime and provisioners all speak it the same way.
|
||||||
|
|
||||||
import type { Envelope } from "../contracts/index.js";
|
import type { Envelope, EventHeaders } from "../contracts/index.js";
|
||||||
|
|
||||||
export type { Envelope };
|
export type { Envelope, EventHeaders };
|
||||||
|
|
||||||
/** A request/reply call and a publish/subscribe surface over the mesh broker. */
|
/** A request/reply call and a publish/subscribe surface over the mesh broker. */
|
||||||
export interface Broker {
|
export interface Broker {
|
||||||
|
|||||||
+6
-1
@@ -195,12 +195,17 @@ test("events: a module emits, a listener and the audit sink (#) both receive it,
|
|||||||
assert.deepEqual(heard.map((e) => e.type), ["module.umami.site.created"]);
|
assert.deepEqual(heard.map((e) => e.type), ["module.umami.site.created"]);
|
||||||
assert.deepEqual(audited.map((e) => e.type), ["module.umami.site.created", "module.plex.play.started"]);
|
assert.deepEqual(audited.map((e) => e.type), ["module.umami.site.created", "module.plex.play.started"]);
|
||||||
|
|
||||||
// The metadata an audit trail needs is present.
|
// The metadata an audit trail needs is present — read back from the ADR 0047 headers, not the body.
|
||||||
const e = heard[0];
|
const e = heard[0];
|
||||||
assert.equal(e.source, "umami");
|
assert.equal(e.source, "umami");
|
||||||
assert.equal(e.node, "anchor");
|
assert.equal(e.node, "anchor");
|
||||||
assert.equal((e.body as { domain: string }).domain, "my-app");
|
assert.equal((e.body as { domain: string }).domain, "my-app");
|
||||||
assert.match(e.at, /^\d{4}-\d{2}-\d{2}T/);
|
assert.match(e.at, /^\d{4}-\d{2}-\d{2}T/);
|
||||||
|
// Every event carries a unique id (x-event-id) — the handle a consumer dedups on.
|
||||||
|
assert.ok(e.id, "event has an x-event-id");
|
||||||
|
assert.notEqual(audited[0].id, audited[1].id);
|
||||||
|
// The body is exactly the domain payload — provenance never leaks into it.
|
||||||
|
assert.deepEqual(Object.keys(e.body as object), ["domain"]);
|
||||||
|
|
||||||
delete process.env.MESH_MODULE;
|
delete process.env.MESH_MODULE;
|
||||||
delete process.env.MESH_NODE;
|
delete process.env.MESH_NODE;
|
||||||
|
|||||||
Reference in New Issue
Block a user