claude-code keeps its MCP servers in state, not events (novox/hq ADR 0202)
One key per registration in the module's servers bucket — all.<server> or <node>.<server> — watched by every node, so a node assigned after a registration takes it at start, which the mcp.registered event could not do. Also narrows apply()'s refusal by hand: the builder compiles without strict, where the discriminated union does not narrow and the build failed.
This commit is contained in:
@@ -30,16 +30,17 @@ this node a subscription token. Nothing else under the home is read or written.
|
|||||||
|
|
||||||
## Over NATS
|
## 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,
|
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 |
|
| 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 |
|
| 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 |
|
| 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 |
|
| 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.<server>` for every node, `<node>.<server>` 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
|
## Tools
|
||||||
|
|
||||||
|
|||||||
@@ -11,16 +11,13 @@
|
|||||||
"binds": {
|
"binds": {
|
||||||
"mcp-endpoint": "${dir:state}/mcp-endpoint.json"
|
"mcp-endpoint": "${dir:state}/mcp-endpoint.json"
|
||||||
},
|
},
|
||||||
"emits": [
|
|
||||||
"mcp.registered",
|
|
||||||
"mcp.unregistered"
|
|
||||||
],
|
|
||||||
"consumes": [
|
"consumes": [
|
||||||
"claude-code.mcp.registered",
|
|
||||||
"claude-code.mcp.unregistered",
|
|
||||||
"claude-licence-manager.licence.rotated",
|
"claude-licence-manager.licence.rotated",
|
||||||
"claude-licence-manager.licence.switched"
|
"claude-licence-manager.licence.switched"
|
||||||
],
|
],
|
||||||
|
"state": [
|
||||||
|
"servers"
|
||||||
|
],
|
||||||
"tools": [
|
"tools": [
|
||||||
"claude_code_status",
|
"claude_code_status",
|
||||||
"claude_code_render",
|
"claude_code_render",
|
||||||
|
|||||||
+96
-33
@@ -9,8 +9,10 @@
|
|||||||
// - a login a person made here — a refresh token this module never writes — is offered to the seat at
|
// - 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
|
// once, sealed to the seat's key: the one moment a refresh token travels, because the login made the
|
||||||
// manager's stale;
|
// manager's stale;
|
||||||
// - an MCP server registered for more nodes than this one is an `mcp.registered` event every node's
|
// - an MCP server registered through this module is **state, not an event** (novox/hq ADR 0202): one
|
||||||
// claude-code consumes, so a node that was off takes it when it is back.
|
// key per server in the module's `servers` bucket — `all.<server>` for every node, `<node>.<server>`
|
||||||
|
// 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 { chmodSync, existsSync, readFileSync, rmSync, writeFileSync } from "node:fs";
|
||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
@@ -107,7 +109,7 @@ export function apply(p: Paths, handed: Current, write: WriteManaged): Record<st
|
|||||||
const local = readCredentials(credentialsPath(p));
|
const local = readCredentials(credentialsPath(p));
|
||||||
const d = decideApply(grantOf(local), grant, switched ? "switch" : "rotation");
|
const d = decideApply(grantOf(local), grant, switched ? "switch" : "rotation");
|
||||||
if (d.apply) writeCredentials(credentialsPath(p), switched ? replacedBy(local, grant) : withGrant(local, grant));
|
if (d.apply) writeCredentials(credentialsPath(p), switched ? replacedBy(local, grant) : withGrant(local, grant));
|
||||||
else outcome = { applied: false, licence: handed.licence, reason: d.reason };
|
else outcome = { applied: false, licence: handed.licence, reason: "reason" in d ? d.reason : undefined }; // narrowed by hand: the build compiles without strict
|
||||||
// Away from the API key: it goes, with its helper.
|
// Away from the API key: it goes, with its helper.
|
||||||
rmSync(apiKeyPath(p), { force: true });
|
rmSync(apiKeyPath(p), { force: true });
|
||||||
rmSync(helperPath(p), { force: true });
|
rmSync(helperPath(p), { force: true });
|
||||||
@@ -153,38 +155,109 @@ export interface Registration {
|
|||||||
nodes?: "all" | string[];
|
nodes?: "all" | string[];
|
||||||
}
|
}
|
||||||
|
|
||||||
function setRegistered(p: Paths, name: string, entry: Record<string, unknown> | null): boolean {
|
/** The `servers` state, as this module reaches it through the runtime (`state("servers")` in the SDK). */
|
||||||
const list = { ...registered(p) } as Record<string, Record<string, unknown>>;
|
export interface ServerState {
|
||||||
const before = JSON.stringify(list[name] ?? null);
|
put(key: string, value: Record<string, unknown>): Promise<number>;
|
||||||
if (entry) list[name] = entry;
|
delete(key: string): Promise<void>;
|
||||||
else delete list[name];
|
keys(): Promise<string[]>;
|
||||||
if (JSON.stringify(list[name] ?? null) === before) return false;
|
|
||||||
writeFileSync(registryPath(p), JSON.stringify(list, null, 2) + "\n", { mode: 0o600 });
|
|
||||||
return true;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
const targets = (p: Paths, nodes: Registration["nodes"]) =>
|
/** One change to the `servers` state, as a watch hands it over. */
|
||||||
nodes === "all" ? true : Array.isArray(nodes) ? nodes.includes(p.node) : false;
|
export interface ServerChange {
|
||||||
|
key: string;
|
||||||
|
op: "put" | "delete";
|
||||||
|
value?: Record<string, unknown>;
|
||||||
|
}
|
||||||
|
|
||||||
/** Register (or with no entry, unregister) here, and announce it for the other nodes asked for. */
|
/** The key a registration lives at: `all.<server>` for every node, `<node>.<server>` for one. */
|
||||||
export async function registerServer(p: Paths, r: Registration, emit: Emit, write: WriteManaged,
|
export const keyOf = (scope: string, name: string) => `${scope}.${name}`;
|
||||||
others: () => Promise<string[]>): Promise<Record<string, unknown>> {
|
|
||||||
|
/**
|
||||||
|
* 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<string, Record<string, unknown>>();
|
||||||
|
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<string, Record<string, unknown>> = {};
|
||||||
|
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<string[]>): Promise<Record<string, unknown>> {
|
||||||
if (r.entry) {
|
if (r.entry) {
|
||||||
const problem = entryProblem(r.name, r.entry);
|
const problem = entryProblem(r.name, r.entry);
|
||||||
if (problem) return { registered: false, reason: problem };
|
if (problem) return { registered: false, reason: problem };
|
||||||
}
|
}
|
||||||
const here = r.nodes === undefined || targets(p, r.nodes);
|
const scopes = scopesOf(p, r.nodes);
|
||||||
const changed = here ? setRegistered(p, r.name, r.entry ?? null) : false;
|
// Compared before and after rather than read from take(): this node's own watch may hand the view the
|
||||||
const rendered = here && changed ? renderNow(p, write) : [];
|
// same change first, and then take() here finds nothing new although this call made it.
|
||||||
if (r.nodes !== undefined) {
|
const before = JSON.stringify(view.effective());
|
||||||
await emit(r.entry ? "mcp.registered" : "mcp.unregistered", { name: r.name, entry: r.entry ?? null, nodes: r.nodes });
|
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<string, unknown> = {
|
const answer: Record<string, unknown> = {
|
||||||
[r.entry ? "registered" : "unregistered"]: r.name,
|
[r.entry ? "registered" : "unregistered"]: r.name,
|
||||||
on: r.nodes === undefined ? [p.node] : r.nodes,
|
on: r.nodes === undefined ? [p.node] : r.nodes,
|
||||||
here: here ? (changed ? "changed" : "already so") : "not this node",
|
here: here ? (changedHere ? "changed" : "already so") : "not this node",
|
||||||
rendered,
|
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) {
|
if (r.nodes === undefined) {
|
||||||
// The question the operator wanted asked: here only, or more?
|
// The question the operator wanted asked: here only, or more?
|
||||||
const elsewhere = (await others().catch(() => [] as string[])).filter((n) => n !== p.node);
|
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;
|
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 };
|
export { MANAGED_DIR };
|
||||||
|
|||||||
@@ -9,7 +9,7 @@
|
|||||||
"test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'"
|
"test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@novox/mesh-sdk": "^0.1.0"
|
"@novox/mesh-sdk": "^0.1.7"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@types/node": "^22.0.0",
|
"@types/node": "^22.0.0",
|
||||||
|
|||||||
@@ -4,7 +4,8 @@ import { existsSync, mkdirSync, mkdtempSync, readFileSync, writeFileSync } from
|
|||||||
import { tmpdir } from "node:os";
|
import { tmpdir } from "node:os";
|
||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import {
|
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";
|
} from "../dist/node.js";
|
||||||
import { generateKeyPair, open, seal } from "../dist/seal.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);
|
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 () => {
|
/** The `servers` state as the bus holds it, shared by every node in a test, with each node's watch. */
|
||||||
const { p, written } = node();
|
function bus() {
|
||||||
const emitted: unknown[] = [];
|
const kept = new Map<string, Record<string, unknown>>();
|
||||||
const r = await registerServer(p, { name: "search", entry: { type: "http", url: "https://s.example/mcp" } },
|
const watchers: ((c: ServerChange) => void)[] = [];
|
||||||
async (t, b) => { emitted.push([t, b]); }, writer(written), async () => ["laptop", "server", "desktop"]);
|
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<string, string> }) => {
|
||||||
|
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.equal(r.here, "changed");
|
||||||
assert.match(String(r.also), /server, desktop/);
|
assert.match(String(r.also), /server, desktop/);
|
||||||
assert.equal(emitted.length, 0, "a registration for this node alone is announced to nobody");
|
assert.deepEqual([...b.kept.keys()], ["laptop.search"]);
|
||||||
assert.ok(JSON.parse(written["managed-mcp.json"]).mcpServers.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 () => {
|
test("registering for every node reaches the others through their watch, and a node joining later reads it", async () => {
|
||||||
const a = node("laptop"), b = node("server");
|
const a = node("laptop"), s = node("server");
|
||||||
let event: [string, unknown] | null = null;
|
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" },
|
await registerServer(a.p, { name: "docs", entry: { type: "stdio", command: "docs-mcp" }, nodes: "all" },
|
||||||
async (t, body) => { event = [t, body]; }, writer(a.written), async () => []);
|
b.state, va, writer(a.written), async () => []);
|
||||||
assert.equal(event![0], "mcp.registered");
|
assert.deepEqual([...b.kept.keys()], ["all.docs"]);
|
||||||
assert.equal(onServerEvent(b.p, "claude-code.mcp.registered", event![1] as never, writer(b.written)), "registered docs from an event");
|
assert.deepEqual(registered(s.p).docs, { type: "stdio", command: "docs-mcp" });
|
||||||
assert.deepEqual(registered(b.p).docs, { type: "stdio", command: "docs-mcp" });
|
assert.ok(JSON.parse(s.written["managed-mcp.json"]).mcpServers.docs);
|
||||||
assert.equal(onServerEvent(b.p, "claude-code.mcp.registered", event![1] as never, writer(b.written)), null, "a repeated event changed something");
|
// 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 () => {
|
test("a node's own registration overrides the one for every node; other nodes' keys leave this one alone", async () => {
|
||||||
const { p, written } = node();
|
const a = node("laptop"), s = node("server");
|
||||||
assert.equal(onServerEvent(p, "claude-code.mcp.registered", { name: "x", entry: { type: "http", url: "https://x" }, nodes: ["server"] }, writer(written)), null);
|
const b = bus();
|
||||||
const r = await registerServer(p, { name: "mesh", entry: { type: "http", url: "https://x" } }, async () => {}, writer(written), async () => []);
|
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(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);
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -4,19 +4,21 @@
|
|||||||
// runtime's own words. **stdout is the MCP channel**: everything this module says, it says on stderr.
|
// 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,
|
// 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
|
// begins watching the credentials file for a login, takes the manager's licence events, and watches the
|
||||||
// licence events and every node's MCP server registrations. node.ts holds the logic.
|
// 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 { readFileSync, watchFile } from "node:fs";
|
||||||
import { spawnSync } from "node:child_process";
|
import { spawnSync } from "node:child_process";
|
||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||||
import { broker } from "@novox/mesh-sdk/messaging";
|
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 {
|
import {
|
||||||
MANAGED_DIR, SEAT, concerns, keypair, offerLogin, onServerEvent, pull, readJson, registerServer,
|
MANAGED_DIR, SEAT, ServerView, concerns, keypair, offerLogin, onServerChange, pull, readJson, registerServer,
|
||||||
registered, renderNow, type Ask, type Paths, type Registration, type WriteManaged,
|
registered, renderNow, type Ask, type Paths, type Registration, type ServerChange, type ServerState, type WriteManaged,
|
||||||
} from "../node.js";
|
} from "../node.js";
|
||||||
import { grantOf, holdsLogin, readCredentials } from "../grant.js";
|
import { grantOf, holdsLogin, readCredentials } from "../grant.js";
|
||||||
import { createHash } from "node:crypto";
|
import { createHash } from "node:crypto";
|
||||||
@@ -90,6 +92,13 @@ function status(p: Paths): Record<string, unknown> {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** The module's MCP servers on the bus (ADR 0202): its own state, which every node of it watches. */
|
||||||
|
const servers = () => state<Record<string, unknown>>("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[] {
|
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 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"] =>
|
const nodesOf = (v: unknown): Registration["nodes"] =>
|
||||||
@@ -115,13 +124,13 @@ function tools(p: Paths): ToolDefinition[] {
|
|||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "claude_code_mcp_list",
|
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.<server>` for every node, `<node>.<server>` for one.",
|
||||||
input: {},
|
input: {},
|
||||||
run: async () => ({ registered: registered(p) }),
|
run: async () => ({ here: registered(p), everywhere: await servers().keys() }),
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "claude_code_mcp_register",
|
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: {
|
input: {
|
||||||
name: { type: "string", description: "the server's name: letters, digits, - and _" },
|
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)" },
|
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) => {
|
run: async (a) => {
|
||||||
const entry: Record<string, unknown> = { type: a.type ?? (a.url ? "http" : "stdio") };
|
const entry: Record<string, unknown> = { type: a.type ?? (a.url ? "http" : "stdio") };
|
||||||
for (const k of ["url", "command", "args", "env", "headers"]) if (a[k] !== undefined) entry[k] = a[k];
|
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",
|
name: "claude_code_mcp_unregister",
|
||||||
description: "Remove an MCP server registered through this module, on this node or more.",
|
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 },
|
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) }))));
|
say(JSON.stringify(await pull(p, ask, writeManaged).catch((e) => ({ failed: String(e) }))));
|
||||||
}).catch(loud("the licence events"));
|
}).catch(loud("the licence events"));
|
||||||
|
|
||||||
void on<Registration>("claude-code.mcp.*", async (event) => {
|
// Every node's MCP servers: the whole current set first, then each change (ADR 0202). Awaited, so the
|
||||||
const done = onServerEvent(p, event.type, event.body, writeManaged);
|
// managed directory holds every server that applies here before the bundle says what it serves.
|
||||||
if (done) say(done);
|
try {
|
||||||
}).catch(loud("the MCP server events"));
|
await state<Record<string, unknown>>("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.
|
// 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"));
|
void pull(p, ask, writeManaged).then((r) => say(`at start: ${JSON.stringify(r)}`), loud("asking for this node's token at start"));
|
||||||
|
|||||||
Reference in New Issue
Block a user