diff --git a/src/contracts/index.ts b/src/contracts/index.ts index cee58d7..9d0105a 100644 --- a/src/contracts/index.ts +++ b/src/contracts/index.ts @@ -44,9 +44,37 @@ export interface ToolDefinition { readonly run: (args: Readonly>) => Promise; } -/** 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 { readonly key: string; readonly node: string; readonly body: T; + readonly headers?: EventHeaders; } diff --git a/src/events/index.ts b/src/events/index.ts index e9c8019..245f2cd 100644 --- a/src/events/index.ts +++ b/src/events/index.ts @@ -4,45 +4,85 @@ // credential — just the broker's topic routing). Declared as emits/consumes on the manifest so the // mesh knows the event graph. // -// This is a thin, audit-ready surface over the broker's publish/subscribe: every event carries who -// emitted it, on which node, and when — so a logger consuming `#` can write a real audit trail. +// Identity and provenance ride as AMQP headers, not in the body (novox/hq ADR 0047): 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 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 { /** The routing key, e.g. "module.umami.site.created". Dotted, so listeners can match by prefix. */ 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; - /** The node it was emitted from. */ + /** The node it was emitted from (x-node). */ readonly node: string; - /** ISO-8601 emit time. */ + /** 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. + * (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(type: string, body: T): Promise { - const event: Event = { - type, - source: process.env.MESH_MODULE ?? "unknown", - node: process.env.MESH_NODE ?? "unknown", - at: new Date().toISOString(), - body, +export async function emit(type: string, body: T, opts: EmitOptions = {}): Promise { + 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>({ key: type, node: event.node, body: event }); + await broker().publish({ 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 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(pattern: string, handler: (event: Event) => Promise): Promise<() => void> { - return broker().subscribe>(pattern, async (envelope) => { - await handler(envelope.body); + return broker().subscribe(pattern, async (envelope) => { + await handler(fromEnvelope(envelope)); }); } + +/** Rebuild the Event a handler sees from a broker envelope's headers (ADR 0047) and body. */ +function fromEnvelope(env: Envelope): Event { + 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, + }; +} diff --git a/src/messaging/index.ts b/src/messaging/index.ts index 9becda1..31a7257 100644 --- a/src/messaging/index.ts +++ b/src/messaging/index.ts @@ -3,9 +3,9 @@ // 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. -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. */ export interface Broker { diff --git a/test/sdk.test.ts b/test/sdk.test.ts index 52faf0e..82c5263 100644 --- a/test/sdk.test.ts +++ b/test/sdk.test.ts @@ -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(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]; assert.equal(e.source, "umami"); assert.equal(e.node, "anchor"); assert.equal((e.body as { domain: string }).domain, "my-app"); 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_NODE;