Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 19 additions & 7 deletions sdk/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -460,13 +460,25 @@ holds against the control plane's records.
- `/sam/mcp/1.0.0` client: `session.openMCP(peer, "mcp://<name>")` sends
the `AuthFrame` naming the service, verifies the provider's credential
and the caller's required labels (`checkPeerLabels`: several pairs are
met by any one of them, as `api.LabelCheck` joins them with `or`; the
conjunction is the operator's egress floor, which only `sam-node` has),
then runs the
official MCP client over the varint-framed stream. JS: a `Transport` for
`@modelcontextprotocol/sdk`; Python: a pair of memory streams pumped to
and from the libp2p stream for `mcp.ClientSession`. `""` as the target is
the provider's own catalog (`list_local_services`, `get_mesh_info`).
met by any one of them, as `api.LabelCheck` joins them with `or`), then
runs the official MCP client over the varint-framed stream. JS: a
`Transport` for `@modelcontextprotocol/sdk`; Python: a pair of memory
streams pumped to and from the libp2p stream for `mcp.ClientSession`.
`""` as the target is the provider's own catalog (`list_local_services`,
`get_mesh_info`).
- Egress floor: `join({ egressRequireLabels })` (`join(egress_require_labels=)`)
is `sam-node`'s `egress.require_labels` for an SDK member. The floor is
the conjunction (`api.LabelFloorCheck` joins the pairs with `,`): every
provider the session calls must attest all of them, on top of a call's
required labels, on every outbound call however the peer was named, MCP
and HTTP alike. Stated once at join and held for the session; a call
cannot waive or widen it. The three implementations agree on it, as they
do on the caller's requirement. The HTTP path (`request`, `fetch`,
`MeshTransport`) verifies the provider with or without a floor, through
the mutual `/sam/auth/1.0.0` handshake, as `sam-node`'s `VerifyPeerLabels`
does before its egress proxy sends anything; a positive verdict is kept
per peer for five minutes (`labelGateTTL`); a refusal is not kept. An
unmet floor is a `LabelsNotSatisfiedError` naming the floor.
- `session.listTools(peer, service)` and `session.callTool(peer, service,
tool, args)` on top of that.
- Tests. Unit: each SDK calls a tool on an in-process provider that serves
Expand Down
4 changes: 3 additions & 1 deletion sdk/js/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,9 @@ process.on("SIGTERM", stop);
`enroll` reuses the identity and credential saved in `stateDir` when they
are still valid for that control plane, and needs exactly one of
`bootstrapTokenPath`, `bootstrapToken` or `jwt` otherwise. Read tokens from
a file or the environment; do not put them on a command line.
a file or the environment; do not put them on a command line. `labels` are
attested at enrollment; `join({ egressRequireLabels })` is the floor every
peer the session calls must attest, all of it, held for the session.

A plaintext `http://` control plane is accepted only on loopback. Pass
`allowInsecure: true` for a network you trust.
Expand Down
33 changes: 25 additions & 8 deletions sdk/js/src/conformance-join.ts
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,25 @@ function emit(obj: unknown): void {
process.stdout.write(JSON.stringify(obj) + "\n");
}

/**
* Labels from an environment variable written "k=v,k2=v2"; a pair without
* "=" is refused (parseRequiredLabels), so a typo never switches a floor off.
*/
function labelsFromEnv(name: string): Record<string, string> {
const out: Record<string, string> = {};
for (const pair of (process.env[name] ?? "").split(",")) {
if (pair.trim() === "") {
continue;
}
const eq = pair.indexOf("=");
if (eq === -1) {
throw new Error(`${name}: invalid label ${JSON.stringify(pair.trim())}: expected key=value`);
}
out[pair.slice(0, eq).trim()] = pair.slice(eq + 1).trim();
}
return out;
}

function failure(cmd: string | undefined, err: unknown): unknown {
return { cmd, ok: false, error: err instanceof Error ? `${err.name}: ${err.message}` : String(err) };
}
Expand Down Expand Up @@ -172,14 +191,11 @@ async function main(): Promise<void> {
const stateDir = requireEnv("SAM_SDK_STATE_DIR");
const allowInsecure = process.env.SAM_INSECURE_CONTROL_PLANE === "1";
const listenAddrs = (process.env.SAM_SDK_LISTEN_ADDRS ?? "").split(",").filter((a) => a !== "");
// Labels this member declares at enrollment, "k=v,k2=v2"; the policy's
// allowed_labels decide whether the control plane attests them.
const labels = Object.fromEntries(
(process.env.SAM_SDK_LABELS ?? "")
.split(",")
.filter((pair) => pair.includes("="))
.map((pair) => pair.split("=", 2) as [string, string]),
);
// Labels this member declares at enrollment; the policy's allowed_labels
// decide whether the control plane attests them. SAM_SDK_EGRESS_REQUIRE_LABELS
// is the floor every provider this member calls must attest.
const labels = labelsFromEnv("SAM_SDK_LABELS");
const egressRequireLabels = labelsFromEnv("SAM_SDK_EGRESS_REQUIRE_LABELS");

const mesh = await AgentMesh.enroll({
controlPlaneUrl,
Expand All @@ -197,6 +213,7 @@ async function main(): Promise<void> {
const session = await mesh.join({
listenAddrs,
...(routerAddresses !== undefined ? { routerAddresses } : {}),
...(Object.keys(egressRequireLabels).length > 0 ? { egressRequireLabels } : {}),
signal: AbortSignal.timeout(20_000),
controlPlaneSyncIntervalMs: 0,
controlPlaneSyncJitterMs: 0,
Expand Down
2 changes: 1 addition & 1 deletion sdk/js/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ export { BiscuitVerificationError, ROLE_ROUTER, requireRole, verifyPeerBiscuit,
export { AUTH_HANDLER_OPTIONS, AUTH_PROTOCOL, MCP_PROTOCOL, AuthRejectedError, authenticateWithPeer, authStreamHandler } from "./auth.ts";
export { createMeshHost, type MeshHost, type MeshHostOptions } from "./host.ts";
export { DHT_PROTOCOL, isServiceType, parseServiceTarget, serviceCID, type ServiceType } from "./discovery.ts";
export { LabelsNotSatisfiedError, StreamTransport, openMCPSession, requireLabels, type MCPSession, type MCPSessionOptions } from "./mcp.ts";
export { LabelsNotSatisfiedError, StreamTransport, openMCPSession, requireEgressLabels, requireLabels, type MCPSession, type MCPSessionOptions } from "./mcp.ts";
export { AuthorizationError, authorizeCaller, type AuthorizeRequest, type ProviderAuthorizerOptions } from "./authorizer.ts";
export {
DEFAULT_A2A_NAME,
Expand Down
38 changes: 36 additions & 2 deletions sdk/js/src/mcp.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,12 +31,12 @@ import { readFileSync } from "node:fs";
import { after, before, test } from "node:test";
import { z } from "zod";
import { AuthRejectedError, MAX_AUTH_FRAME_BYTES, MCP_PROTOCOL } from "./auth.ts";
import { loadBiscuit, verifyPeerBiscuit } from "./biscuit.ts";
import { ROLE_ROUTER, loadBiscuit, verifyPeerBiscuit } from "./biscuit.ts";
import { ROLE_NODE } from "./controlplane.ts";
import { parseServiceTarget, serviceCID } from "./discovery.ts";
import { AuthFrameSchema, AuthResponseSchema } from "./gen/sam_pb.ts";
import { Identity } from "./identity.ts";
import { LabelsNotSatisfiedError, StreamTransport, openMCPSession, requireLabels } from "./mcp.ts";
import { LabelsNotSatisfiedError, StreamTransport, openMCPSession, requireEgressLabels, requireLabels } from "./mcp.ts";

type Wasm = Awaited<ReturnType<typeof loadBiscuit>>;

Expand Down Expand Up @@ -186,6 +186,28 @@ test("a requirement of several labels is met by any one of them, as sam-node's c
requireLabels(attesting({}), {});
});

test("the egress floor is met only by every one of its pairs, as sam-node's api.LabelFloorCheck", () => {
const attesting = (labels: Record<string, string>) => ({ peerId: "p", expiration: new Date(), verifyingKey: cpKey, roles: [], labels });
requireEgressLabels(attesting({ region: "eu", team: "platform" }), { region: "eu" });
requireEgressLabels(attesting({ region: "eu", team: "platform" }), { region: "eu", team: "platform" });
// one pair short is a refusal that names the whole floor
assert.throws(() => requireEgressLabels(attesting({ region: "eu" }), { region: "eu", team: "platform" }), /peer p does not attest the egress floor: region=eu, team=platform/);
assert.throws(() => requireEgressLabels(attesting({ region: "us" }), { region: "eu" }), LabelsNotSatisfiedError);
// no floor is no floor
requireEgressLabels(attesting({}), {});
requireEgressLabels(attesting({}), undefined);
});

test("the egress floor is held on the MCP path beside the caller's requirement", async () => {
const conn = await caller.dial(provider.getMultiaddrs()[0] as Parameters<typeof caller.dial>[0]);
const ok = await openMCPSession(conn, frame("mcp://calc"), [cpKey], {}, { region: "eu" });
await ok.close();
await assert.rejects(openMCPSession(conn, frame("mcp://calc"), [cpKey], {}, { region: "eu", team: "platform" }), LabelsNotSatisfiedError);
// Both apply when both are set: neither one's pairs stand in for the other's.
await assert.rejects(openMCPSession(conn, frame("mcp://calc"), [cpKey], { requiredLabels: { region: "eu" } }, { team: "platform" }), LabelsNotSatisfiedError);
await assert.rejects(openMCPSession(conn, frame("mcp://calc"), [cpKey], { requiredLabels: { team: "platform" } }, { region: "eu" }), LabelsNotSatisfiedError);
});

test("a caller the provider cannot verify gets no session", async () => {
const conn = await caller.dial(provider.getMultiaddrs()[0] as Parameters<typeof caller.dial>[0]);
const forged = new wasm.KeyPair(wasm.SignatureAlgorithm.Ed25519);
Expand All @@ -206,3 +228,15 @@ test("a service the provider does not have ends the session before MCP starts",
const conn = await caller.dial(provider.getMultiaddrs()[0] as Parameters<typeof caller.dial>[0]);
await assert.rejects(openMCPSession(conn, frame("mcp://no-such-service"), [cpKey], { signal: AbortSignal.timeout(3000) }));
});

test("only a node is a provider, as sam-node's checkPeerLabels requires", async () => {
// The provider answers with the credential it holds at the time; a router's, attesting the floor, is not a provider's.
const nodeBiscuit = providerBiscuit;
providerBiscuit = mint(provider.peerId.toString(), ROLE_ROUTER, { region: "eu" });
try {
const conn = await caller.dial(provider.getMultiaddrs()[0] as Parameters<typeof caller.dial>[0]);
await assert.rejects(openMCPSession(conn, frame("mcp://calc"), [cpKey], {}, { region: "eu" }), /lacks expected role "sam:role:node"/);
} finally {
providerBiscuit = nodeBiscuit;
}
});
43 changes: 36 additions & 7 deletions sdk/js/src/mcp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,8 @@ import { Client } from "@modelcontextprotocol/sdk/client/index.js";
import type { Transport } from "@modelcontextprotocol/sdk/shared/transport.js";
import { JSONRPCMessageSchema, type JSONRPCMessage } from "@modelcontextprotocol/sdk/types.js";
import { AUTH_HANDSHAKE_TIMEOUT_MS, AuthRejectedError, MAX_AUTH_FRAME_BYTES, MCP_PROTOCOL } from "./auth.ts";
import { BiscuitVerificationError, verifyPeerBiscuit, type VerifiedBiscuit } from "./biscuit.ts";
import { BiscuitVerificationError, requireRole, verifyPeerBiscuit, type VerifiedBiscuit } from "./biscuit.ts";
import { ROLE_NODE } from "./controlplane.ts";
import { decodeAuthResponse } from "./credential.ts";

/** go-msgio's default message cap, which sam-node's StreamTransport uses. */
Expand Down Expand Up @@ -102,19 +103,22 @@ export interface MCPSession {
close(): Promise<void>;
}

/** The provider's credential carries none of the labels the caller requires (checkPeerLabels). */
/**
* The provider's credential lacks what the caller requires (any one pair) or
* what the session's egress floor requires (every pair), as checkPeerLabels refuses.
*/
export class LabelsNotSatisfiedError extends Error {
constructor(peerId: string, required: string[]) {
super(`peer ${peerId} carries none of the required labels: ${required.join(", ")}`);
constructor(peerId: string, required: string[], what = "carries none of the required labels") {
super(`peer ${peerId} ${what}: ${required.join(", ")}`);
this.name = "LabelsNotSatisfiedError";
}
}

/**
* A caller's requirement is satisfied by any one pair, as sam-node's
* api.LabelCheck (`check if label(k1, v1) or label(k2, v2)`): several pairs
* mean "any of these will do". The operator's egress floor is the
* conjunction, and sam-node's alone.
* mean "any of these will do". The egress floor (requireEgressLabels) is the
* conjunction.
*/
export function requireLabels(provider: VerifiedBiscuit, required: Record<string, string> | undefined): void {
if (!required) {
Expand All @@ -130,16 +134,38 @@ export function requireLabels(provider: VerifiedBiscuit, required: Record<string
);
}

/**
* The egress floor is met only by every one of its pairs, as sam-node's
* api.LabelFloorCheck (`check if label(k1, v1), label(k2, v2)`) for
* egress.require_labels: a floor takes no alternatives. Empty is no floor.
*/
export function requireEgressLabels(provider: VerifiedBiscuit, required: Record<string, string> | undefined): void {
if (!required) {
return;
}
const pairs = Object.entries(required);
if (pairs.every(([k, v]) => provider.labels[k] === v)) {
return;
}
throw new LabelsNotSatisfiedError(
provider.peerId,
pairs.map(([k, v]) => `${k}=${v}`),
"does not attest the egress floor",
);
}

/**
* Opens /sam/mcp/1.0.0 to a connected provider for targetService ("" is the
* provider's own catalog), verifies the provider, and returns a connected
* MCP client. frame is this member's AuthFrame for that service.
* MCP client. frame is this member's AuthFrame for that service; egressRequireLabels
* is the session's, not the caller's (requireEgressLabels).
*/
export async function openMCPSession(
conn: Connection,
frame: Uint8Array,
trustedKeys: Uint8Array[],
options: MCPSessionOptions = {},
egressRequireLabels?: Record<string, string>,
): Promise<MCPSession> {
const signal = options.signal ?? AbortSignal.timeout(AUTH_HANDSHAKE_TIMEOUT_MS);
const stream = await conn.newStream(MCP_PROTOCOL, { signal, runOnLimitedConnection: true });
Expand All @@ -153,7 +179,10 @@ export async function openMCPSession(
throw new AuthRejectedError(conn.remotePeer.toString(), resp.error || "no reason given");
}
provider = await verifyPeerBiscuit(resp.biscuit, conn.remotePeer.toString(), trustedKeys);
// Only nodes host services; a router's or an admin's credential is a member, not a provider.
requireRole(provider, ROLE_NODE);
requireLabels(provider, options.requiredLabels);
requireEgressLabels(provider, egressRequireLabels);
} catch (err) {
await stream.close().catch(() => stream.abort(err instanceof Error ? err : new Error(String(err))));
if (err instanceof BiscuitVerificationError) {
Expand Down
Loading
Loading