events: modules log activity to the broker, any module reacts
The lighter sibling of provisioning — 1:many and broadcast, no credential,
just the broker's topic routing. A thin, audit-ready surface over the
broker's publish/subscribe:
- emit(type, body): publishes an Event carrying who emitted it (MESH_MODULE),
on which node (MESH_NODE) and when (ISO timestamp) — so a listener can
build a real audit trail.
- on(pattern, handler): react to events by topic pattern. The audit logger
is just on("#", ...).
Tested: a module emits; a targeted listener (module.umami.#) hears only its
events, the audit sink (#) hears every module's, and the metadata audit
needs is present.
The declared side — a manifest's emits/consumes, so the mesh knows the
event graph — and the audit-logger module are the next pieces.
Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
@@ -9,6 +9,7 @@
|
|||||||
"./provisioner": "./dist/provisioner/index.js",
|
"./provisioner": "./dist/provisioner/index.js",
|
||||||
"./tools": "./dist/tools/index.js",
|
"./tools": "./dist/tools/index.js",
|
||||||
"./messaging": "./dist/messaging/index.js",
|
"./messaging": "./dist/messaging/index.js",
|
||||||
|
"./events": "./dist/events/index.js",
|
||||||
"./primitives": "./dist/primitives/index.js"
|
"./primitives": "./dist/primitives/index.js"
|
||||||
},
|
},
|
||||||
"scripts": {
|
"scripts": {
|
||||||
|
|||||||
@@ -0,0 +1,48 @@
|
|||||||
|
// 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.
|
||||||
|
//
|
||||||
|
// 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.
|
||||||
|
|
||||||
|
import { broker } from "../messaging/index.js";
|
||||||
|
|
||||||
|
/** An event on the mesh: a topic key, its source, and a body — with the metadata audit needs. */
|
||||||
|
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 emitting module. */
|
||||||
|
readonly source: string;
|
||||||
|
/** The node it was emitted from. */
|
||||||
|
readonly node: string;
|
||||||
|
/** ISO-8601 emit time. */
|
||||||
|
readonly at: string;
|
||||||
|
readonly body: T;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 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.
|
||||||
|
*/
|
||||||
|
export async function emit<T>(type: string, body: T): Promise<void> {
|
||||||
|
const event: Event<T> = {
|
||||||
|
type,
|
||||||
|
source: process.env.MESH_MODULE ?? "unknown",
|
||||||
|
node: process.env.MESH_NODE ?? "unknown",
|
||||||
|
at: new Date().toISOString(),
|
||||||
|
body,
|
||||||
|
};
|
||||||
|
await broker().publish<Event<T>>({ key: type, node: event.node, body: event });
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 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.
|
||||||
|
*/
|
||||||
|
export async function on<T>(pattern: string, handler: (event: Event<T>) => Promise<void>): Promise<() => void> {
|
||||||
|
return broker().subscribe<Event<T>>(pattern, async (envelope) => {
|
||||||
|
await handler(envelope.body);
|
||||||
|
});
|
||||||
|
}
|
||||||
@@ -8,4 +8,5 @@ export * from "./contracts/index.js";
|
|||||||
export * as tools from "./tools/index.js";
|
export * as tools from "./tools/index.js";
|
||||||
export * as provisioner from "./provisioner/index.js";
|
export * as provisioner from "./provisioner/index.js";
|
||||||
export * as messaging from "./messaging/index.js";
|
export * as messaging from "./messaging/index.js";
|
||||||
|
export * as events from "./events/index.js";
|
||||||
export * as primitives from "./primitives/index.js";
|
export * as primitives from "./primitives/index.js";
|
||||||
|
|||||||
+63
-5
@@ -7,6 +7,8 @@ import { join } from "node:path";
|
|||||||
import { seal, unseal, compareVersions } from "../dist/primitives/index.js";
|
import { seal, unseal, compareVersions } from "../dist/primitives/index.js";
|
||||||
import { registerModuleTools, collectTools, resetTools, serveTools, listTools } from "../dist/tools/index.js";
|
import { registerModuleTools, collectTools, resetTools, serveTools, listTools } from "../dist/tools/index.js";
|
||||||
import { runProvisioner, type Grant } from "../dist/provisioner/index.js";
|
import { runProvisioner, type Grant } from "../dist/provisioner/index.js";
|
||||||
|
import { emit, on, type Event } from "../dist/events/index.js";
|
||||||
|
import { useBroker } from "../dist/messaging/index.js";
|
||||||
|
|
||||||
test("seal round-trips and rejects the wrong key", () => {
|
test("seal round-trips and rejects the wrong key", () => {
|
||||||
const sealed = seal("hunter2", "node-key");
|
const sealed = seal("hunter2", "node-key");
|
||||||
@@ -126,10 +128,12 @@ test("serving refuses two modules exposing one tool name", async () => {
|
|||||||
await assert.rejects(serveTools(memBroker(), {}), /exposed by two modules/);
|
await assert.rejects(serveTools(memBroker(), {}), /exposed by two modules/);
|
||||||
});
|
});
|
||||||
|
|
||||||
// A minimal in-memory broker: request routes to a registered handle. Enough to serve tools; the
|
// A minimal in-memory broker: request routes to a registered handle, and publish routes to every
|
||||||
// real broker binding is the mesh's, provided by the hosting runtime.
|
// subscriber whose topic pattern matches. Enough to exercise tools and events; the real binding is
|
||||||
|
// the mesh's AMQP one, provided by the hosting runtime.
|
||||||
function memBroker() {
|
function memBroker() {
|
||||||
const handlers = new Map<string, (b: unknown) => Promise<unknown>>();
|
const handlers = new Map<string, (b: unknown) => Promise<unknown>>();
|
||||||
|
const subs: { pattern: string; handler: (env: { key: string; node: string; body: unknown }) => Promise<void> }[] = [];
|
||||||
return {
|
return {
|
||||||
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
||||||
const h = handlers.get(key);
|
const h = handlers.get(key);
|
||||||
@@ -140,14 +144,68 @@ function memBroker() {
|
|||||||
handlers.set(key, handler as (b: unknown) => Promise<unknown>);
|
handlers.set(key, handler as (b: unknown) => Promise<unknown>);
|
||||||
return () => handlers.delete(key);
|
return () => handlers.delete(key);
|
||||||
},
|
},
|
||||||
async publish(): Promise<void> {},
|
async publish<T>(env: { key: string; node: string; body: T }): Promise<void> {
|
||||||
async subscribe(): Promise<() => void> {
|
for (const s of subs) if (topicMatch(s.pattern, env.key)) await s.handler(env);
|
||||||
return () => {};
|
},
|
||||||
|
async subscribe<T>(pattern: string, handler: (env: { key: string; node: string; body: T }) => Promise<void>): Promise<() => void> {
|
||||||
|
const entry = { pattern, handler: handler as (env: { key: string; node: string; body: unknown }) => Promise<void> };
|
||||||
|
subs.push(entry);
|
||||||
|
return () => {
|
||||||
|
const i = subs.indexOf(entry);
|
||||||
|
if (i >= 0) subs.splice(i, 1);
|
||||||
|
};
|
||||||
},
|
},
|
||||||
async close(): Promise<void> {},
|
async close(): Promise<void> {},
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// AMQP topic matching: `*` one segment, `#` any run of segments.
|
||||||
|
function topicMatch(pattern: string, key: string): boolean {
|
||||||
|
if (pattern === "#") return true;
|
||||||
|
const p = pattern.split(".");
|
||||||
|
const k = key.split(".");
|
||||||
|
let pi = 0;
|
||||||
|
let ki = 0;
|
||||||
|
while (pi < p.length) {
|
||||||
|
if (p[pi] === "#") return true; // simplification: # only as a trailing wildcard
|
||||||
|
if (ki >= k.length) return false;
|
||||||
|
if (p[pi] !== "*" && p[pi] !== k[ki]) return false;
|
||||||
|
pi++;
|
||||||
|
ki++;
|
||||||
|
}
|
||||||
|
return ki === k.length;
|
||||||
|
}
|
||||||
|
|
||||||
|
test("events: a module emits, a listener and the audit sink (#) both receive it, with metadata", async () => {
|
||||||
|
const broker = memBroker();
|
||||||
|
useBroker(() => broker);
|
||||||
|
|
||||||
|
process.env.MESH_MODULE = "umami";
|
||||||
|
process.env.MESH_NODE = "anchor";
|
||||||
|
|
||||||
|
const heard: Event[] = [];
|
||||||
|
const audited: Event[] = [];
|
||||||
|
await on("module.umami.#", async (e) => void heard.push(e));
|
||||||
|
await on("#", async (e) => void audited.push(e)); // the audit logger is just this
|
||||||
|
|
||||||
|
await emit("module.umami.site.created", { domain: "my-app" });
|
||||||
|
await emit("module.plex.play.started", { title: "x" }); // a different module's event
|
||||||
|
|
||||||
|
// The umami listener heard only umami's event; the audit sink heard both.
|
||||||
|
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.
|
||||||
|
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/);
|
||||||
|
|
||||||
|
delete process.env.MESH_MODULE;
|
||||||
|
delete process.env.MESH_NODE;
|
||||||
|
});
|
||||||
|
|
||||||
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
||||||
const start = Date.now();
|
const start = Date.now();
|
||||||
while (!cond()) {
|
while (!cond()) {
|
||||||
|
|||||||
Reference in New Issue
Block a user