diff --git a/modules/claude-code/README.md b/modules/claude-code/README.md index 55de6ed..2575ad9 100644 --- a/modules/claude-code/README.md +++ b/modules/claude-code/README.md @@ -30,16 +30,17 @@ this node a subscription token. Nothing else under the home is read or written. ## Over NATS -Everything between this module and the rest of the mesh is NATS, in two kinds: an **event** says that +Everything between this module and the rest of the mesh is NATS, in three kinds: an **event** says that something happened and carries no secret, because a stream keeps it; a **request** carries a token, -because nothing keeps it (hq design 32 §10). +because nothing keeps it (hq design 32 §10); and **state** is the current value of something every node +must see, a node that joins later included — kept, so it carries no secret either (hq ADR 0202). | what | how | |---|---| | the licence manager rotated a licence, or switched this node | its `licence.rotated` / `licence.switched` event; this module then asks `anthropic-licence-manager.current` for its token, sealed to the key it sends | | this node starts | it asks `current` once, so a node that was off catches up | | a person ran `/login` here | the credentials file gains a refresh token this module never writes; it asks `anthropic-licence-manager.adopt` at once with the grant sealed to the manager's key — the one time a refresh token travels, because the login made the manager's stale | -| an MCP server registered for more nodes than this one | an `mcp.registered` / `mcp.unregistered` event every node's claude-code consumes; a node that was off takes it when it is back | +| an MCP server registered through this module | a key in the module's `servers` state — `all.` for every node, `.` for one; every node watches it and renders what applies to it, a node's own entry over the one for every node. A node that joins later, or was off, reads the whole current set at start; unregistering is a delete. An entry with a secret in its `env` or `headers` is refused by the runtime | ## Tools diff --git a/modules/claude-code/module.json b/modules/claude-code/module.json index 3fa6085..841e009 100644 --- a/modules/claude-code/module.json +++ b/modules/claude-code/module.json @@ -11,16 +11,13 @@ "binds": { "mcp-endpoint": "${dir:state}/mcp-endpoint.json" }, - "emits": [ - "mcp.registered", - "mcp.unregistered" - ], "consumes": [ - "claude-code.mcp.registered", - "claude-code.mcp.unregistered", "claude-licence-manager.licence.rotated", "claude-licence-manager.licence.switched" ], + "state": [ + "servers" + ], "tools": [ "claude_code_status", "claude_code_render", diff --git a/modules/claude-code/node.ts b/modules/claude-code/node.ts index 6b1be5b..af22d7f 100644 --- a/modules/claude-code/node.ts +++ b/modules/claude-code/node.ts @@ -9,8 +9,10 @@ // - a login a person made here — a refresh token this module never writes — is offered to the seat at // once, sealed to the seat's key: the one moment a refresh token travels, because the login made the // manager's stale; -// - an MCP server registered for more nodes than this one is an `mcp.registered` event every node's -// claude-code consumes, so a node that was off takes it when it is back. +// - an MCP server registered through this module is **state, not an event** (novox/hq ADR 0202): one +// key per server in the module's `servers` bucket — `all.` for every node, `.` +// for one — which every node watches. A node that joins later, or was off, reads the whole current set +// at start; unregistering is a delete. A secret never goes in an entry: the runtime refuses one. import { chmodSync, existsSync, readFileSync, rmSync, writeFileSync } from "node:fs"; import { join } from "node:path"; @@ -107,7 +109,7 @@ export function apply(p: Paths, handed: Current, write: WriteManaged): Record | null): boolean { - const list = { ...registered(p) } as Record>; - const before = JSON.stringify(list[name] ?? null); - if (entry) list[name] = entry; - else delete list[name]; - if (JSON.stringify(list[name] ?? null) === before) return false; - writeFileSync(registryPath(p), JSON.stringify(list, null, 2) + "\n", { mode: 0o600 }); - return true; +/** The `servers` state, as this module reaches it through the runtime (`state("servers")` in the SDK). */ +export interface ServerState { + put(key: string, value: Record): Promise; + delete(key: string): Promise; + keys(): Promise; } -const targets = (p: Paths, nodes: Registration["nodes"]) => - nodes === "all" ? true : Array.isArray(nodes) ? nodes.includes(p.node) : false; +/** One change to the `servers` state, as a watch hands it over. */ +export interface ServerChange { + key: string; + op: "put" | "delete"; + value?: Record; +} -/** Register (or with no entry, unregister) here, and announce it for the other nodes asked for. */ -export async function registerServer(p: Paths, r: Registration, emit: Emit, write: WriteManaged, - others: () => Promise): Promise> { +/** The key a registration lives at: `all.` for every node, `.` for one. */ +export const keyOf = (scope: string, name: string) => `${scope}.${name}`; + +/** + * What this node takes from the `servers` state: the entries for every node and for this one, by key — + * kept in memory from the watch, and written through to the module's own file whenever what applies here + * changes, so the managed directory can be rendered without the bus. + */ +export class ServerView { + private readonly entries = new Map>(); + constructor(private readonly p: Paths) {} + + /** Take one change; answers whether what applies to this node changed. */ + take(c: ServerChange): boolean { + const dot = c.key.indexOf("."); + const scope = c.key.slice(0, dot), name = c.key.slice(dot + 1); + if (dot <= 0 || (scope !== "all" && scope !== this.p.node)) return false; + if (c.op === "put" && c.value && entryProblem(name, c.value) === null) this.entries.set(c.key, c.value); + else this.entries.delete(c.key); + return this.writeThrough(); + } + + /** What applies here: every node's entries, with this node's own laid over them by server name. */ + effective(): Servers { + const out: Record> = {}; + for (const scope of ["all", this.p.node]) { + for (const [key, entry] of [...this.entries].sort(([a], [b]) => a.localeCompare(b))) { + if (key.startsWith(scope + ".")) out[key.slice(scope.length + 1)] = entry; + } + } + return out; + } + + private writeThrough(): boolean { + const now = JSON.stringify(this.effective(), null, 2) + "\n"; + let before = ""; + try { + before = readFileSync(registryPath(this.p), "utf8"); + } catch { + /* none yet */ + } + if (now === before) return false; + writeFileSync(registryPath(this.p), now, { mode: 0o600 }); + return true; + } +} + +/** A change from the watch: take it, and render when what applies here changed. */ +export function onServerChange(view: ServerView, c: ServerChange, p: Paths, write: WriteManaged): string | null { + if (!view.take(c)) return null; + renderNow(p, write); + return `${c.op === "put" ? "registered" : "unregistered"} ${c.key}`; +} + +const scopesOf = (p: Paths, nodes: Registration["nodes"]): string[] => + nodes === undefined ? [p.node] : nodes === "all" ? ["all"] : nodes; + +/** + * Register (or with no entry, unregister) a server: a put (or delete) per scope in the `servers` state. + * Taken into this node's view at once, so the answer says what it did here; every other node takes it + * from its watch, and a node that joins later from the current state. + */ +export async function registerServer(p: Paths, r: Registration, servers: ServerState, view: ServerView, + write: WriteManaged, others: () => Promise): Promise> { if (r.entry) { const problem = entryProblem(r.name, r.entry); if (problem) return { registered: false, reason: problem }; } - const here = r.nodes === undefined || targets(p, r.nodes); - const changed = here ? setRegistered(p, r.name, r.entry ?? null) : false; - const rendered = here && changed ? renderNow(p, write) : []; - if (r.nodes !== undefined) { - await emit(r.entry ? "mcp.registered" : "mcp.unregistered", { name: r.name, entry: r.entry ?? null, nodes: r.nodes }); + const scopes = scopesOf(p, r.nodes); + // Compared before and after rather than read from take(): this node's own watch may hand the view the + // same change first, and then take() here finds nothing new although this call made it. + const before = JSON.stringify(view.effective()); + for (const scope of scopes) { + const key = keyOf(scope, r.name); + if (r.entry) await servers.put(key, r.entry); + else await servers.delete(key); + view.take({ key, op: r.entry ? "put" : "delete", value: r.entry }); } + const changedHere = JSON.stringify(view.effective()) !== before; + const here = scopes.includes("all") || scopes.includes(p.node); const answer: Record = { [r.entry ? "registered" : "unregistered"]: r.name, on: r.nodes === undefined ? [p.node] : r.nodes, - here: here ? (changed ? "changed" : "already so") : "not this node", - rendered, + here: here ? (changedHere ? "changed" : "already so") : "not this node", + rendered: changedHere ? renderNow(p, write) : [], }; + if (!r.entry && view.effective()[r.name]) { + answer.still = `${r.name} still applies here from another registration (for every node, or for this one); unregister that too`; + } if (r.nodes === undefined) { // The question the operator wanted asked: here only, or more? const elsewhere = (await others().catch(() => [] as string[])).filter((n) => n !== p.node); @@ -195,14 +268,4 @@ export async function registerServer(p: Paths, r: Registration, emit: Emit, writ return answer; } -/** An `mcp.registered`/`mcp.unregistered` event from any node's claude-code: apply it if it names this node. */ -export function onServerEvent(p: Paths, type: string, body: Registration, write: WriteManaged): string | null { - if (!body?.name || !targets(p, body.nodes)) return null; - const entry = type.endsWith("mcp.registered") ? body.entry ?? null : null; - if (entry && entryProblem(body.name, entry)) return null; - if (!setRegistered(p, body.name, entry)) return null; - renderNow(p, write); - return `${entry ? "registered" : "unregistered"} ${body.name} from an event`; -} - export { MANAGED_DIR }; diff --git a/modules/claude-code/package.json b/modules/claude-code/package.json index fca3599..fc05060 100644 --- a/modules/claude-code/package.json +++ b/modules/claude-code/package.json @@ -9,7 +9,7 @@ "test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'" }, "dependencies": { - "@novox/mesh-sdk": "^0.1.0" + "@novox/mesh-sdk": "^0.1.7" }, "devDependencies": { "@types/node": "^22.0.0", diff --git a/modules/claude-code/test/node.test.ts b/modules/claude-code/test/node.test.ts index 50d57dd..d7cd92d 100644 --- a/modules/claude-code/test/node.test.ts +++ b/modules/claude-code/test/node.test.ts @@ -4,7 +4,8 @@ import { existsSync, mkdirSync, mkdtempSync, readFileSync, writeFileSync } from import { tmpdir } from "node:os"; import { join } from "node:path"; import { - apply, concerns, keypair, offerLogin, onServerEvent, pull, registerServer, registered, type Paths, + apply, concerns, keypair, offerLogin, onServerChange, pull, registerServer, registered, ServerView, type Paths, + type ServerChange, type ServerState, } from "../dist/node.js"; import { generateKeyPair, open, seal } from "../dist/seal.js"; @@ -89,31 +90,84 @@ test("no refresh token in the file is no login, and nothing is asked", async () assert.equal(await offerLogin(p, async () => { throw new Error("asked"); }), null); }); -test("registering a server here renders it and asks whether to register it on the other nodes", async () => { - const { p, written } = node(); - const emitted: unknown[] = []; - const r = await registerServer(p, { name: "search", entry: { type: "http", url: "https://s.example/mcp" } }, - async (t, b) => { emitted.push([t, b]); }, writer(written), async () => ["laptop", "server", "desktop"]); +/** The `servers` state as the bus holds it, shared by every node in a test, with each node's watch. */ +function bus() { + const kept = new Map>(); + const watchers: ((c: ServerChange) => void)[] = []; + const state: ServerState = { + put: async (key, value) => { kept.set(key, value); watchers.forEach((w) => w({ key, op: "put", value })); return kept.size; }, + delete: async (key) => { kept.delete(key); watchers.forEach((w) => w({ key, op: "delete" })); }, + keys: async () => [...kept.keys()].sort(), + }; + /** A node joining: its view takes the current state, then every change. */ + const join = (n: { p: Paths; written: Record }) => { + const view = new ServerView(n.p); + for (const [key, value] of kept) onServerChange(view, { key, op: "put", value }, n.p, writer(n.written)); + watchers.push((c) => onServerChange(view, c, n.p, writer(n.written))); + return view; + }; + return { state, join, kept }; +} + +test("registering a server here puts it under this node's key, renders it, and asks about the other nodes", async () => { + const n = node(); + const b = bus(); + const view = b.join(n); + const r = await registerServer(n.p, { name: "search", entry: { type: "http", url: "https://s.example/mcp" } }, + b.state, view, writer(n.written), async () => ["laptop", "server", "desktop"]); assert.equal(r.here, "changed"); assert.match(String(r.also), /server, desktop/); - assert.equal(emitted.length, 0, "a registration for this node alone is announced to nobody"); - assert.ok(JSON.parse(written["managed-mcp.json"]).mcpServers.search); + assert.deepEqual([...b.kept.keys()], ["laptop.search"]); + assert.ok(JSON.parse(n.written["managed-mcp.json"]).mcpServers.search); }); -test("registering for every node emits the event, and another node applies it from the event", async () => { - const a = node("laptop"), b = node("server"); - let event: [string, unknown] | null = null; +test("registering for every node reaches the others through their watch, and a node joining later reads it", async () => { + const a = node("laptop"), s = node("server"); + const b = bus(); + const va = b.join(a); + b.join(s); await registerServer(a.p, { name: "docs", entry: { type: "stdio", command: "docs-mcp" }, nodes: "all" }, - async (t, body) => { event = [t, body]; }, writer(a.written), async () => []); - assert.equal(event![0], "mcp.registered"); - assert.equal(onServerEvent(b.p, "claude-code.mcp.registered", event![1] as never, writer(b.written)), "registered docs from an event"); - assert.deepEqual(registered(b.p).docs, { type: "stdio", command: "docs-mcp" }); - assert.equal(onServerEvent(b.p, "claude-code.mcp.registered", event![1] as never, writer(b.written)), null, "a repeated event changed something"); + b.state, va, writer(a.written), async () => []); + assert.deepEqual([...b.kept.keys()], ["all.docs"]); + assert.deepEqual(registered(s.p).docs, { type: "stdio", command: "docs-mcp" }); + assert.ok(JSON.parse(s.written["managed-mcp.json"]).mcpServers.docs); + // The gap events left: a node assigned after the registration takes the whole current set at start. + const late = node("desktop"); + b.join(late); + assert.deepEqual(registered(late.p).docs, { type: "stdio", command: "docs-mcp" }); + // Unregistering is a delete, and every node's view drops it. + await registerServer(a.p, { name: "docs", nodes: "all" }, b.state, va, writer(a.written), async () => []); + assert.equal(registered(s.p).docs, undefined); + assert.equal(registered(late.p).docs, undefined); }); -test("an event naming other nodes leaves this one alone; a bad entry is refused before anything is written", async () => { - const { p, written } = node(); - assert.equal(onServerEvent(p, "claude-code.mcp.registered", { name: "x", entry: { type: "http", url: "https://x" }, nodes: ["server"] }, writer(written)), null); - const r = await registerServer(p, { name: "mesh", entry: { type: "http", url: "https://x" } }, async () => {}, writer(written), async () => []); +test("a node's own registration overrides the one for every node; other nodes' keys leave this one alone", async () => { + const a = node("laptop"), s = node("server"); + const b = bus(); + const va = b.join(a); + const vs = b.join(s); + await registerServer(a.p, { name: "x", entry: { type: "http", url: "https://all" }, nodes: "all" }, b.state, va, writer(a.written), async () => []); + await registerServer(a.p, { name: "x", entry: { type: "http", url: "https://laptop" } }, b.state, va, writer(a.written), async () => []); + assert.equal(registered(a.p).x.url, "https://laptop"); + assert.equal(registered(s.p).x.url, "https://all"); + await registerServer(a.p, { name: "only", entry: { type: "http", url: "https://o" }, nodes: ["server"] }, b.state, va, writer(a.written), async () => []); + assert.equal(registered(a.p).only, undefined); + assert.equal(registered(s.p).only.url, "https://o"); + // Unregistering here leaves the every-node one applying, and says so. + const r = await registerServer(a.p, { name: "x" }, b.state, va, writer(a.written), async () => []); + assert.match(String(r.still), /still applies here/); + assert.equal(registered(a.p).x.url, "https://all"); + assert.equal(vs.effective().x.url, "https://all"); +}); + +test("a bad entry is refused before anything is put; a repeated change changes nothing", async () => { + const n = node(); + const b = bus(); + const view = b.join(n); + const r = await registerServer(n.p, { name: "mesh", entry: { type: "http", url: "https://x" } }, b.state, view, writer(n.written), async () => []); assert.equal(r.registered, false); + assert.equal(b.kept.size, 0); + assert.equal(onServerChange(view, { key: "all.a", op: "put", value: { type: "http", url: "https://a" } }, n.p, writer(n.written)), "registered all.a"); + assert.equal(onServerChange(view, { key: "all.a", op: "put", value: { type: "http", url: "https://a" } }, n.p, writer(n.written)), null); + assert.equal(onServerChange(view, { key: "server.b", op: "put", value: { type: "http", url: "https://b" } }, n.p, writer(n.written)), null); }); diff --git a/modules/claude-code/tools/index.ts b/modules/claude-code/tools/index.ts index 7a04abe..af4bffb 100644 --- a/modules/claude-code/tools/index.ts +++ b/modules/claude-code/tools/index.ts @@ -4,19 +4,21 @@ // runtime's own words. **stdout is the MCP channel**: everything this module says, it says on stderr. // // At start it renders the agent's managed directory, asks the licence manager for this node's token, -// begins watching the credentials file for a login, and takes the module's events: the manager's -// licence events and every node's MCP server registrations. node.ts holds the logic. +// begins watching the credentials file for a login, takes the manager's licence events, and watches the +// module's `servers` state — every node's MCP server registrations (novox/hq ADR 0202). node.ts holds the +// logic. import { readFileSync, watchFile } from "node:fs"; import { spawnSync } from "node:child_process"; import { join } from "node:path"; import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; import { broker } from "@novox/mesh-sdk/messaging"; -import { emit, on } from "@novox/mesh-sdk/events"; +import { on } from "@novox/mesh-sdk/events"; +import { state } from "@novox/mesh-sdk/state"; import { - MANAGED_DIR, SEAT, concerns, keypair, offerLogin, onServerEvent, pull, readJson, registerServer, - registered, renderNow, type Ask, type Paths, type Registration, type WriteManaged, + MANAGED_DIR, SEAT, ServerView, concerns, keypair, offerLogin, onServerChange, pull, readJson, registerServer, + registered, renderNow, type Ask, type Paths, type Registration, type ServerChange, type ServerState, type WriteManaged, } from "../node.js"; import { grantOf, holdsLogin, readCredentials } from "../grant.js"; import { createHash } from "node:crypto"; @@ -90,6 +92,13 @@ function status(p: Paths): Record { }; } +/** The module's MCP servers on the bus (ADR 0202): its own state, which every node of it watches. */ +const servers = () => state>("servers") as unknown as ServerState; + +/** What this node takes from that state, kept from the watch. One per process. */ +let view: ServerView | null = null; +const viewOf = (p: Paths) => (view ??= new ServerView(p)); + function tools(p: Paths): ToolDefinition[] { const nodesArg = { type: "string", description: 'more nodes: "all" for every node running claude-code, or a comma-separated list; absent is this node only' }; const nodesOf = (v: unknown): Registration["nodes"] => @@ -115,13 +124,13 @@ function tools(p: Paths): ToolDefinition[] { }, { name: "claude_code_mcp_list", - description: "The MCP servers registered on this node through this module, beside the console (`mesh`) and those set in the module's settings.", + description: "The MCP servers registered through this module: those that apply on this node (beside the console, `mesh`, and those set in the module's settings), and every registration on the mesh, by key — `all.` for every node, `.` for one.", input: {}, - run: async () => ({ registered: registered(p) }), + run: async () => ({ here: registered(p), everywhere: await servers().keys() }), }, { name: "claude_code_mcp_register", - description: "Register an MCP server with Claude Code on this node — an http/sse server by url, or a stdio server by command — and say which other nodes run claude-code, so it can be registered there too.", + description: "Register an MCP server with Claude Code on this node, every node, or a list — an http/sse server by url, or a stdio server by command. Kept on the bus, so a node that joins later takes it too. Never put a secret in env or headers: the mesh refuses one.", input: { name: { type: "string", description: "the server's name: letters, digits, - and _" }, type: { type: "string", description: "http, sse or stdio (default stdio when a command is given, http when a url is)" }, @@ -135,14 +144,14 @@ function tools(p: Paths): ToolDefinition[] { run: async (a) => { const entry: Record = { type: a.type ?? (a.url ? "http" : "stdio") }; for (const k of ["url", "command", "args", "env", "headers"]) if (a[k] !== undefined) entry[k] = a[k]; - return registerServer(p, { name: String(a.name ?? ""), entry, nodes: nodesOf(a.nodes) }, emit, writeManaged, nodesRunningMe); + return registerServer(p, { name: String(a.name ?? ""), entry, nodes: nodesOf(a.nodes) }, servers(), viewOf(p), writeManaged, nodesRunningMe); }, }, { name: "claude_code_mcp_unregister", description: "Remove an MCP server registered through this module, on this node or more.", input: { name: { type: "string", description: "the server's name" }, nodes: nodesArg }, - run: async (a) => registerServer(p, { name: String(a.name ?? ""), nodes: nodesOf(a.nodes) }, emit, writeManaged, nodesRunningMe), + run: async (a) => registerServer(p, { name: String(a.name ?? ""), nodes: nodesOf(a.nodes) }, servers(), viewOf(p), writeManaged, nodesRunningMe), }, ]; } @@ -171,10 +180,20 @@ if (p) { say(JSON.stringify(await pull(p, ask, writeManaged).catch((e) => ({ failed: String(e) })))); }).catch(loud("the licence events")); - void on("claude-code.mcp.*", async (event) => { - const done = onServerEvent(p, event.type, event.body, writeManaged); - if (done) say(done); - }).catch(loud("the MCP server events")); + // Every node's MCP servers: the whole current set first, then each change (ADR 0202). Awaited, so the + // managed directory holds every server that applies here before the bundle says what it serves. + try { + await state>("servers").watch((c) => { + try { + const done = onServerChange(viewOf(p), c as ServerChange, p, writeManaged); + if (done) say(done); + } catch (err) { + loud(`taking ${c.op} ${c.key}`)(err); // the view took it; the next render writes it + } + }); + } catch (err) { + loud("watching the MCP servers")(err); + } // Catch up once at start: a node that was off takes its current token now. void pull(p, ask, writeManaged).then((r) => say(`at start: ${JSON.stringify(r)}`), loud("asking for this node's token at start"));