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/src/primitives/index.ts b/src/primitives/index.ts index 5d6edd2..29fbf65 100644 --- a/src/primitives/index.ts +++ b/src/primitives/index.ts @@ -1,42 +1,10 @@ // Small, stable primitives every module's code may need. No behaviour here changes when a module // changes; that is the whole point of it living in the sdk. - -import { createCipheriv, createDecipheriv, randomBytes, scryptSync } from "node:crypto"; - -// --- sealing --- // -// A secret is sealed to a key so a copy of it at rest is not a working credential. AES-256-GCM; -// the key is derived from a per-node passphrase the host holds. The host unseals on the machine; -// nothing else does (novox/hq ADR 0043's link is the boundary — the sdk only carries the mechanism). - -const MAGIC = "msk1"; // versions the sealed format, so it can change without silent misreads - -/** Seal plaintext to a passphrase. Returns `msk1::::`, base64 parts. */ -export function seal(plaintext: string, passphrase: string): string { - const salt = randomBytes(16); - const iv = randomBytes(12); - const key = scryptSync(passphrase, salt, 32); - const cipher = createCipheriv("aes-256-gcm", key, iv); - const enc = Buffer.concat([cipher.update(plaintext, "utf8"), cipher.final()]); - const tag = cipher.getAuthTag(); - return [MAGIC, b64(salt), b64(iv), b64(tag), b64(enc)].join(":"); -} - -/** Unseal what seal produced. Throws — loudly — on any tamper or wrong key. */ -export function unseal(sealed: string, passphrase: string): string { - const parts = sealed.split(":"); - if (parts.length !== 5 || parts[0] !== MAGIC) { - throw new Error("not a sealed value this version understands"); - } - const [, salt, iv, tag, enc] = parts.map((p, i) => (i === 0 ? Buffer.alloc(0) : ub64(p))); - const key = scryptSync(passphrase, salt, 32); - const decipher = createDecipheriv("aes-256-gcm", key, iv); - decipher.setAuthTag(tag); - return Buffer.concat([decipher.update(enc), decipher.final()]).toString("utf8"); -} - -const b64 = (b: Buffer): string => b.toString("base64url"); -const ub64 = (s: string): Buffer => Buffer.from(s, "base64url"); +// There was a symmetric seal()/unseal() here, for a provider to seal a credential to a key before +// writing it. It is gone: a provider is handed the credential the mesh minted and seals nothing +// (novox/hq ADR 0053), the mesh's own secret delivery is asymmetric and belongs to the host, and +// nothing else called it. The primitive left with the provisioner that was its only caller. // --- semver --- // diff --git a/src/provisioner/index.ts b/src/provisioner/index.ts index 2606bf4..f31a445 100644 --- a/src/provisioner/index.ts +++ b/src/provisioner/index.ts @@ -1,74 +1,109 @@ // The reconcile harness every provider shares. The loop below is identical for postgres, redis, -// minio and umami: read the grants the control plane has written, bring each requested resource to -// existence through the provider's adapter, seal and hand back the credential, and remove what is -// no longer requested. A module writes ONLY the adapter — the per-service half — which is why this -// lives in the sdk and the adapter lives in the module (novox/hq ADR 0044). +// minio and umami: read the contributions the mesh delivered, bring each consumer's resource into +// being through the provider's adapter under the login and password the mesh minted, and withdraw +// what the mesh no longer asks for. A module writes ONLY the adapter — the per-service half — which +// is why this lives in the sdk and the adapter lives in the module (novox/hq ADR 0044). +// +// **A provider creates the credential the mesh minted, and seals nothing (novox/hq ADR 0053).** The +// control plane mints one password per consumer and seals it to this node; the host unseals it into +// the file the contribution names. The provider does not generate a password, does not seal one, and +// does not hand one back — the consumer already receives its copy through the mesh's own channel. +// This is why the old $MESH_SEAL_KEY, the symmetric seal, and the *.grant.json / *.credential files +// are gone: they described a second, contradictory credential path the mesh does not have. -import { readdir, readFile, writeFile, rm } from "node:fs/promises"; -import { join } from "node:path"; -import type { Grant, Credential } from "../contracts/index.js"; -import { seal } from "../primitives/index.js"; +import { readFile } from "node:fs/promises"; -/** What a provider implements — the only per-service code. It is handed a Grant and creates (or - * removes) the resource in its own software, returning the credential a consumer receives. */ +/** One consumer's resource to bring into being — everything the mesh derived and delivered. */ +export interface Provision { + /** The login the mesh derived and gave the consumer to present. The provider creates exactly + * this name — a name the consumer cannot learn is a name it cannot authenticate with. */ + readonly as: string; + /** The password the mesh minted for this consumer, read from the file the mesh sealed to this + * node and the host unsealed. The provider sets it; it never invents one. */ + readonly password: string; + /** What the consumer contributed, per the interface's spec keys (e.g. `{ name: "umami" }`). */ + readonly values: Readonly>; + /** Where the consumer is, for a provider that must reach back to it. Usually unused. */ + readonly at?: string; + /** The consumer node, for logging and lifecycle events. */ + readonly consumer?: string; +} + +/** What a provider implements — the only per-service code. It is handed the login and password the + * mesh made and brings the resource into being under them, or withdraws it. It returns nothing: + * the credential is the mesh's, already delivered to the consumer. */ export interface Adapter { - create(grant: Grant): Promise; - remove(grant: Grant): Promise; + create(p: Provision): Promise; + remove(p: { readonly as: string }): Promise; } export interface ProvisionerOptions { - /** Directory the control plane writes grant requests into and the harness writes credentials to. - * Defaults to $GRANTS. */ - grants?: string; - /** Passphrase a credential is sealed to before it is written. Defaults to $MESH_SEAL_KEY. */ - sealKey?: string; + /** The contributions file the mesh writes for this resource (the module's `receives` path). + * Defaults to $MESH_RECEIVES. */ + receives?: string; /** Reconcile interval in ms. Defaults to 5000. */ everyMs?: number; } -export type { Grant, Credential }; +/** One entry in the mesh's contributions file: a consumer the provider must serve. */ +interface Contribution { + readonly as: string; + readonly secret: string; + readonly node?: string; + readonly at?: string; + readonly values?: Readonly>; +} /** * Run the reconcile loop for one provided resource. Returns a stop function. Never throws for a - * single bad grant — it logs and keeps converging, because one consumer's failure must not stop - * the others' provisioning. + * single bad contribution — it logs and keeps converging, because one consumer's failure must not + * stop the others' provisioning. */ export function runProvisioner(resource: string, adapter: Adapter, opts: ProvisionerOptions = {}): () => void { - const dir = opts.grants ?? envOrThrow("GRANTS"); - const sealKey = opts.sealKey ?? envOrThrow("MESH_SEAL_KEY"); + const receives = opts.receives ?? envOrThrow("MESH_RECEIVES"); const everyMs = opts.everyMs ?? 5000; - const applied = new Map(); // consumer -> hash of the grant last applied + const applied = new Map(); // login (`as`) -> hash of what was last applied let stopped = false; async function reconcile(): Promise { - const requested = await readGrants(dir, resource); - const wantById = new Map(requested.map((g) => [g.consumer, g])); + const given = await readContributions(receives, resource); + const wantByAs = new Map(given.map((g) => [g.as, g])); - // Create or update anything requested whose grant has changed. - for (const grant of requested) { - const h = hash(grant); - if (applied.get(grant.consumer) === h) continue; + // Create or update every consumer whose login, password or values changed. + for (const g of given) { + let password: string; try { - const cred = await adapter.create(grant); - await writeSealedCredential(dir, grant, cred, sealKey); - applied.set(grant.consumer, h); + // The file the mesh sealed to this node, unsealed by the host into plaintext. Trailing + // newline trimmed: a sealed value is exactly the secret, and a host writing a file may add + // one. + password = (await readFile(g.secret, "utf8")).replace(/\n$/, ""); } catch (err) { - console.error(`[provisioner:${resource}] ${grant.consumer}: create failed, will retry: ${err}`); + // The contribution names a secret the host has not written yet — a normal race on the + // first pass. Skipped, retried on the next tick, never fatal. + console.error(`[provisioner:${resource}] ${g.as}: secret not readable yet (${g.secret}): ${err}`); + continue; + } + const h = hash(g.as, password, g.values ?? {}); + if (applied.get(g.as) === h) continue; + try { + await adapter.create({ as: g.as, password, values: g.values ?? {}, at: g.at, consumer: g.node }); + applied.set(g.as, h); + } catch (err) { + console.error(`[provisioner:${resource}] ${g.as}: create failed, will retry: ${err}`); } } - // Remove anything applied that is no longer requested. Only the harness's own outputs are - // touched — the credential file it wrote — never anything it did not create. - for (const consumer of [...applied.keys()]) { - if (wantById.has(consumer)) continue; - const grant = lastGrant(dir, resource, consumer); + // Withdraw every login the provider made that the mesh no longer asks for. The mesh drops a + // consumer from the file when it goes away; that is how a login is withdrawn rather than kept + // working for ever. Only logins this harness created are touched. + for (const as of [...applied.keys()]) { + if (wantByAs.has(as)) continue; try { - await adapter.remove(grant); - await rm(credentialPath(dir, resource, consumer), { force: true }); - applied.delete(consumer); + await adapter.remove({ as }); + applied.delete(as); } catch (err) { - console.error(`[provisioner:${resource}] ${consumer}: remove failed, will retry: ${err}`); + console.error(`[provisioner:${resource}] ${as}: remove failed, will retry: ${err}`); } } } @@ -89,47 +124,35 @@ export function runProvisioner(resource: string, adapter: Adapter, opts: Provisi }; } -// --- grant/credential files --- - -async function readGrants(dir: string, resource: string): Promise { - let names: string[]; +/** + * Read the mesh's contributions file for one resource. Absent or empty means "no consumer asks for + * this" — the file is always written, so a provider can tell that from "the mesh never wrote it". + */ +async function readContributions(path: string, resource: string): Promise { + let raw: string; try { - names = await readdir(dir); + raw = await readFile(path, "utf8"); } catch { + return []; // not written yet, or nobody serves anything here — nothing to converge to + } + let doc: { requirement?: string; given?: Contribution[] }; + try { + doc = JSON.parse(raw) as { requirement?: string; given?: Contribution[] }; + } catch (err) { + console.error(`[provisioner:${resource}] contributions file is not JSON (${path}): ${err}`); return []; } - const out: Grant[] = []; - for (const name of names) { - if (!name.endsWith(".grant.json")) continue; - try { - const raw = await readFile(join(dir, name), "utf8"); - const g = JSON.parse(raw) as Grant; - if (g.resource === resource) out.push(g); - } catch (err) { - console.error(`[provisioner:${resource}] unreadable grant ${name}: ${err}`); - } + if (doc.requirement && doc.requirement !== resource) { + console.error(`[provisioner:${resource}] ${path} is for ${doc.requirement}, not ${resource}`); + return []; } - return out; + // A contribution with no `as` is not a credential grant (a module offering something on its own + // machine) — the provisioner has nothing to create for it. + return (doc.given ?? []).filter((g) => g.as && g.secret); } -function credentialPath(dir: string, resource: string, consumer: string): string { - return join(dir, `${consumer}.${resource}.credential`); -} - -async function writeSealedCredential(dir: string, grant: Grant, cred: Credential, key: string): Promise { - const body = JSON.stringify({ resource: grant.resource, consumer: grant.consumer, fields: cred.fields }); - await writeFile(credentialPath(dir, grant.resource, grant.consumer), seal(body, key), { mode: 0o600 }); -} - -function lastGrant(dir: string, resource: string, consumer: string): Grant { - // For removal the harness needs a Grant to hand the adapter; the identity is enough to act on. - return { resource, consumer, node: "", values: {} }; -} - -// --- helpers --- - -function hash(g: Grant): string { - return JSON.stringify([g.resource, g.consumer, g.node, g.values]); +function hash(as: string, password: string, values: Readonly>): string { + return JSON.stringify([as, password, values]); } function envOrThrow(name: string): string { diff --git a/src/tools/index.ts b/src/tools/index.ts index 0fae74f..132cbbc 100644 --- a/src/tools/index.ts +++ b/src/tools/index.ts @@ -66,20 +66,43 @@ export function listTools(env: NodeJS.ProcessEnv = process.env): { module: strin * other — two tools answering one name is a fault, not a race to resolve. */ export async function serveTools(broker: Broker, env: NodeJS.ProcessEnv = process.env): Promise<() => void> { - const byName = new Map(); + // Each tool is served on its own key, namespaced by its module (novox/hq ADR 0052): a caller + // invokes `.`, only the module that serves it answers, and the module's account is + // scoped to serve..* — so one module cannot answer another's calls. The module is the + // namespace, so a tool name need only be unique within its module, not across the whole mesh. + const stops: Array<() => void> = []; for (const { module, tools } of collectTools(env)) { + const seen = new Set(); for (const t of tools) { - if (byName.has(t.name)) { - throw new Error(`tool name ${t.name} is exposed by two modules (one is ${module}) — refused`); + if (seen.has(t.name)) { + throw new Error(`${module} exposes two tools named ${t.name} — refused`); } - byName.set(t.name, t); + seen.add(t.name); + const stop = await broker.handle>, unknown>( + toolKey(module, t.name), + (args) => t.run(args ?? {}), + ); + stops.push(stop); } } - return broker.handle("tools.invoke", async (inv) => { - const tool = byName.get(inv.tool); - if (!tool) throw new Error(`no such tool: ${inv.tool}`); - return tool.run(inv.args ?? {}); - }); + return () => { + for (const stop of stops) stop(); + }; +} + +/** The broker key a tool is served on and invoked by — the module namespaces the tool (ADR 0052). */ +export function toolKey(module: string, tool: string): string { + return `${module}.${tool}`; +} + +/** Invoke a module's tool over the broker — the caller's side of serving. */ +export async function invokeTool( + broker: Broker, + module: string, + tool: string, + args: Readonly> = {}, +): Promise { + return broker.request(toolKey(module, tool), args); } /** Testing/inspection: drop all registrations. */ diff --git a/test/sdk.test.ts b/test/sdk.test.ts index 52faf0e..05294e0 100644 --- a/test/sdk.test.ts +++ b/test/sdk.test.ts @@ -4,19 +4,12 @@ import { mkdtemp, writeFile, readFile, readdir } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { seal, unseal, compareVersions } from "../dist/primitives/index.js"; -import { registerModuleTools, collectTools, resetTools, serveTools, listTools } from "../dist/tools/index.js"; -import { runProvisioner, type Grant } from "../dist/provisioner/index.js"; +import { compareVersions } from "../dist/primitives/index.js"; +import { registerModuleTools, collectTools, resetTools, serveTools, listTools, invokeTool } from "../dist/tools/index.js"; +import { runProvisioner } 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", () => { - const sealed = seal("hunter2", "node-key"); - assert.equal(unseal(sealed, "node-key"), "hunter2"); - assert.notEqual(sealed, "hunter2"); - assert.throws(() => unseal(sealed, "wrong-key")); -}); - test("semver orders releases", () => { assert.equal(compareVersions("1.2.3", "1.2.10"), -1); assert.equal(compareVersions("2.0.0", "1.9.9"), 1); @@ -38,39 +31,46 @@ test("module tools register and collect, a thrower is skipped not fatal", () => assert.equal(broken?.tools.length, 0); }); -test("provisioner creates a sealed credential for a grant, then removes on withdrawal", async () => { +test("provisioner creates each consumer with the mesh's login and password, removes on withdrawal", async () => { const dir = await mkdtemp(join(tmpdir(), "prov-")); - const created: string[] = []; + const created: { as: string; password: string; values: unknown }[] = []; const removed: string[] = []; - const grant: Grant = { resource: "analytics", consumer: "webapp", node: "anchor", values: { name: "webapp" } }; - await writeFile(join(dir, "webapp.grant.json"), JSON.stringify(grant)); + // What the mesh delivers: the password it minted (as the host leaves it after unsealing) and a + // contributions file naming the consumer's login and where that password is. + await writeFile(join(dir, "webapp.secret"), "minted-pw\n"); + const receives = join(dir, "analytics.json"); + const doc = (given: unknown[]): string => + JSON.stringify({ contributions: 1, requirement: "analytics", given }); + await writeFile( + receives, + doc([{ from: "webapp", node: "anchor", as: "webapp-anchor", secret: join(dir, "webapp.secret"), values: { name: "webapp" } }]), + ); const stop = runProvisioner( "analytics", { - async create(g) { - created.push(g.consumer); - return { fields: { siteId: "abc", snippet: "