From 592e80d83ffbd624d149ee71542033fb88169097 Mon Sep 17 00:00:00 2001 From: Tomas Tormo Date: Wed, 30 Sep 2026 20:45:37 +0000 Subject: [PATCH 1/2] sdk: expect the card of an agent behind a node rewritten for the mesh --- tests/integration/sdk_mesh_test.go | 52 +++++++++++++++++++++++++++++- 1 file changed, 51 insertions(+), 1 deletion(-) diff --git a/tests/integration/sdk_mesh_test.go b/tests/integration/sdk_mesh_test.go index 51c94f05..e00b352f 100644 --- a/tests/integration/sdk_mesh_test.go +++ b/tests/integration/sdk_mesh_test.go @@ -212,6 +212,16 @@ egress: // serving the MCP service "calc" from a backend this test runs. backend := httptest.NewServer(newBoundaryMCPHandler(t)) t.Cleanup(backend.Close) + // A stock A2A agent behind the node; its card names its own address. + agentCard := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/.well-known/agent-card.json" { + http.NotFound(w, r) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = io.WriteString(w, sdkMeshStockAgentCard) + })) + t.Cleanup(agentCard.Close) nodeBin := buildBinary(t, "./cmd/sam-node") nodeHome := filepath.Join(t.TempDir(), "node") if err := os.MkdirAll(nodeHome, 0o755); err != nil { @@ -232,7 +242,9 @@ egress: "--allow-loopback", "--api-token-path", tokenPath(t, "node-token"), "--discovery-interval", "100ms", - "--config", writeNodeConfig(t, nodeHome, sdkMeshLabels, svcDecl{Type: "mcp", Name: "calc", TargetURL: backend.URL}), + "--config", writeNodeConfig(t, nodeHome, sdkMeshLabels, + svcDecl{Type: "mcp", Name: "calc", TargetURL: backend.URL}, + svcDecl{Type: "a2a", Name: sdkMeshAgentName, TargetURL: agentCard.URL}), "--secrets-dir", secrets, "--log-level", "debug", ) @@ -259,6 +271,15 @@ egress: } } +const sdkMeshAgentName = "echo-agent" + +const sdkMeshStockAgentCard = `{"name":"echo-agent","description":"stock agent behind a sam-node","version":"1.0.0",` + + `"capabilities":{"streaming":true},` + + `"supportedInterfaces":[{"url":"http://127.0.0.1:7777/","protocolBinding":"JSONRPC","protocolVersion":"1.0"},` + + `{"url":"127.0.0.1:50051","protocolBinding":"GRPC","protocolVersion":"1.0"}],` + + `"signatures":[{"protected":"eyJhbGciOiJFUzI1NiJ9","signature":"c3RhbGU"}],` + + `"skills":[],"defaultInputModes":["text/plain"],"defaultOutputModes":["text/plain"]}` + // sdkMember is a running SDK conformance-join runner: a mesh member written // in another language that the test drives over stdin/stdout. Both runners // (sdk/js/src/conformance-join.ts, sdk/python/src/agent_mesh/conformance_join.py) @@ -468,6 +489,35 @@ func TestNativeSDKsMesh(t *testing.T) { if res := m.callRaw(t, nodeRelayAddr, "mcp://no-such-service", "add", nil); res.OK { t.Fatalf("%s called a service the node does not serve: %+v", m.name, res) } + + // A stock client bootstraps from the node's agent card, so the SDK + // serves it rewritten for the mesh, as the node's egress proxy does. + res := m.http(t, nodeRelayAddr, "a2a://"+sdkMeshAgentName, "/.well-known/agent-card.json") + var card struct { + SupportedInterfaces []struct { + URL string `json:"url"` + ProtocolBinding string `json:"protocolBinding"` + } `json:"supportedInterfaces"` + Capabilities struct { + Streaming bool `json:"streaming"` + } `json:"capabilities"` + Signatures []json.RawMessage `json:"signatures"` + } + if res.Status != 200 || json.Unmarshal([]byte(res.Body), &card) != nil { + t.Fatalf("%s fetching the node's agent card: %+v", m.name, res) + } + meshBase := "http://mesh/sam/" + samNode.peerID.String() + "/a2a/" + sdkMeshAgentName + if len(card.SupportedInterfaces) != 1 || card.SupportedInterfaces[0].URL != meshBase || card.SupportedInterfaces[0].ProtocolBinding != "JSONRPC" { + t.Errorf("%s got interfaces %+v, want one JSONRPC interface at %s", m.name, card.SupportedInterfaces, meshBase) + } + // Streaming stays as the agent declares it: the SDK's transport streams. + if !card.Capabilities.Streaming || len(card.Signatures) != 0 { + t.Errorf("%s got a card with streaming=%v (want true) and %d signatures (want none)", m.name, card.Capabilities.Streaming, len(card.Signatures)) + } + // The bare service root serves the same card, as the node does for a2a-go. + if root := m.http(t, nodeRelayAddr, "a2a://"+sdkMeshAgentName, "/"); root.Body != res.Body { + t.Errorf("%s got a different card at the service root: %s", m.name, root.Body) + } }) } From 1eb670e9dddb033c5a9246e15a0d4a0311bf89e5 Mon Sep 17 00:00:00 2001 From: Tomas Tormo Date: Wed, 30 Sep 2026 20:45:37 +0000 Subject: [PATCH 2/2] sdk: rewrite the card of an agent behind a node for the mesh A stock A2A client on session.fetch() or MeshTransport follows the URL in the agent card. An agent behind a sam-node serves a card naming its own address, so the client left the mesh and the SDK refused the URL. sam-node's egress proxy impersonates that card endpoint for its own callers. The SDKs now do the same at their /libp2p-http client: a GET of /a2a//.well-known/agent-card.json or of the bare service root is held, the SDK fetches the card itself at the well-known path with identity encoding, and answers with it regenerated. The rewrite points every HTTP interface at the mesh URL, drops gRPC interfaces and signatures, and answers 502 for a card that is not JSON or has no interface left; the agent's own non-200 is relayed as it is. Streaming stays as the agent declares it, unlike the node's rewrite, since this transport streams and SDK-hosted agents advertise it. --- sdk/README.md | 10 +- sdk/js/src/index.ts | 2 + sdk/js/src/libp2p-http.test.ts | 75 +++++++++++- sdk/js/src/libp2p-http.ts | 108 +++++++++++++++++- sdk/python/src/agent_mesh/__init__.py | 6 + sdk/python/src/agent_mesh/httpx_transport.py | 11 +- sdk/python/src/agent_mesh/libp2p_http.py | 113 ++++++++++++++++++- sdk/python/src/agent_mesh/session.py | 6 +- sdk/python/tests/test_libp2p_http.py | 97 +++++++++++++++- site/content/docs/guides/native-sdks.md | 4 + 10 files changed, 410 insertions(+), 22 deletions(-) diff --git a/sdk/README.md b/sdk/README.md index b83d595e..0affbdd5 100644 --- a/sdk/README.md +++ b/sdk/README.md @@ -141,7 +141,7 @@ pinned by a test: ID is case-sensitive. The URL an HTTP client uses for a peer's service therefore carries the peer ID in the path, `http://mesh/sam////`, the shape of `sam-node`'s - egress proxy and of an agent card it rewrote; the host is ignored. + egress proxy and of an agent card rewritten for the mesh; the host is ignored. - A member that publishes nothing has announced no address, so nothing in the DHT or a router's peerstore names one. `sam-node` dials `/p2p//p2p-circuit` for every router it authenticated with when @@ -520,7 +520,13 @@ holds against the control plane's records. `open_http_request` / `fetchOverStream` return once the headers are in and stream the body; `MeshTransport` (httpx) and `session.fetch()` (fetch) carry a client's requests to the peer a mesh URL names. The A2A - SDK's client takes either without changes. + SDK's client takes either without changes. The card of an agent behind a + `sam-node` names the agent's own address; a GET of the well-known card + path or of the service root is answered the way the node's egress proxy + answers it: the SDK fetches the card itself, with identity encoding, and + serves it rewritten (`rewriteAgentCard` / `rewrite_agent_card`): HTTP + interfaces point at the mesh URL, gRPC ones are dropped, signatures go. + Streaming stays as the agent declares it, since the transport streams. - `sam-node`: a peer it knows no address for is dialed through every router it authenticated with, so its egress proxy reaches an agent by peer ID (`preparePeerAddrs`). diff --git a/sdk/js/src/index.ts b/sdk/js/src/index.ts index 45da9677..72290db4 100644 --- a/sdk/js/src/index.ts +++ b/sdk/js/src/index.ts @@ -47,6 +47,7 @@ export { LabelsNotSatisfiedError, StreamTransport, openMCPSession, requireEgress export { AuthorizationError, authorizeCaller, type AuthorizeRequest, type ProviderAuthorizerOptions } from "./authorizer.ts"; export { DEFAULT_A2A_NAME, + AGENT_CARD_PATH, HTTP_HANDLER_OPTIONS, HTTP_PROTOCOL, MESH_PATH_PREFIX, @@ -57,6 +58,7 @@ export { httpRequestOverStream, meshHTTPTarget, meshURL, + rewriteAgentCard, splitMeshURL, type A2AEndpoint, type A2AEndpointSpec, diff --git a/sdk/js/src/libp2p-http.test.ts b/sdk/js/src/libp2p-http.test.ts index 697e2501..b05fd7c2 100644 --- a/sdk/js/src/libp2p-http.test.ts +++ b/sdk/js/src/libp2p-http.test.ts @@ -29,6 +29,7 @@ import { after, before, test } from "node:test"; import { loadBiscuit } from "./biscuit.ts"; import { ROLE_NODE } from "./controlplane.ts"; import { + AGENT_CARD_PATH, HTTP_PROTOCOL, a2aEndpoint, admitIngress, @@ -37,6 +38,7 @@ import { httpRequestOverStream, meshHTTPTarget, meshURL, + rewriteAgentCard, splitMeshURL, type ProviderOptions, } from "./libp2p-http.ts"; @@ -55,7 +57,7 @@ let callerBiscuit: Uint8Array; let guestBiscuit: Uint8Array; let backend: http.Server; let backendURL: string; -const backendSeen: { method: string; url: string; peer: string | undefined; body: string }[] = []; +const backendSeen: { method: string; url: string; peer: string | undefined; body: string; encoding: string | undefined }[] = []; const listenerSeen: { url: string; peer: string | undefined; biscuit: string | undefined }[] = []; const authorized: string[] = []; @@ -98,14 +100,40 @@ function providerOptions(biscuit: Uint8Array): ProviderOptions { }; } +const STOCK_CARD = { + name: "echo-agent", + version: "1.0.0", + capabilities: { streaming: true, pushNotifications: false }, + supportedInterfaces: [ + { url: "http://127.0.0.1:7777/", protocolBinding: "JSONRPC", protocolVersion: "1.0" }, + { url: "127.0.0.1:50051", protocolBinding: "GRPC", protocolVersion: "1.0" }, + ], + signatures: [{ protected: "eyJhbGciOiJFUzI1NiJ9", signature: "c3RhbGU" }], + skills: [], + defaultInputModes: ["text/plain"], + defaultOutputModes: ["text/plain"], +}; + // Stands in for an A2A server beside the agent: echoes the request and, on // /stream, answers with three SSE events as message/stream would. function fakeA2AServer(req: http.IncomingMessage, res: http.ServerResponse): void { let body = ""; req.on("data", (c: Buffer) => (body += c.toString())); req.on("end", () => { - backendSeen.push({ method: req.method ?? "", url: req.url ?? "", peer: req.headers["x-peer-id"] as string | undefined, body }); + backendSeen.push({ method: req.method ?? "", url: req.url ?? "", peer: req.headers["x-peer-id"] as string | undefined, body, encoding: req.headers["accept-encoding"] as string | undefined }); assert.equal(req.headers["x-sam-biscuit"], undefined, "biscuit leaked to the backend"); + if (req.url === `/${AGENT_CARD_PATH}` && req.headers["x-card"] === "missing") { + res.writeHead(404, { "content-type": "text/plain" }); + res.end("no card"); + return; + } + if (req.url === `/${AGENT_CARD_PATH}`) { + const grpcOnly = req.headers["x-card"] === "grpc-only"; + const supportedInterfaces = STOCK_CARD.supportedInterfaces.filter((i) => !grpcOnly || i.protocolBinding === "GRPC"); + res.writeHead(200, { "content-type": "application/json" }); + res.end(JSON.stringify({ ...STOCK_CARD, supportedInterfaces })); + return; + } if (req.url === "/stream") { res.writeHead(200, { "content-type": "text/event-stream" }); let i = 0; @@ -242,6 +270,49 @@ test("a fetch over the stream delivers an SSE body event by event", async () => ); }); +test("the card of an agent behind a node comes back rewritten for the mesh, as sam-node serves it", async () => { + const conn = await dial(); + const base = meshURL(agent.peerId.toString(), "a2a://agent"); + const response = await fetchOverStream(conn, callerBiscuit, new Request(`${base}/${AGENT_CARD_PATH}`)); + assert.equal(response.status, 200); + assert.equal(response.headers.get("content-type"), "application/json"); + const card = (await response.json()) as Record; + assert.deepEqual(card.supportedInterfaces, [{ url: base, protocolBinding: "JSONRPC", protocolVersion: "1.0" }]); + assert.deepEqual(card.capabilities, { streaming: true, pushNotifications: false }); + assert.equal("signatures" in card, false); + assert.equal(card.name, "echo-agent"); + assert.deepEqual(card.skills, []); + assert.equal(backendSeen[backendSeen.length - 1]?.url, `/${AGENT_CARD_PATH}`); + + // The bare service root serves the card too, the way a2a-go resolves a pathful base URL. + assert.deepEqual(await (await fetchOverStream(conn, callerBiscuit, new Request(base))).json(), card); + assert.equal(backendSeen[backendSeen.length - 1]?.url, `/${AGENT_CARD_PATH}`); + await fetchOverStream(conn, callerBiscuit, new Request(`${base}/${AGENT_CARD_PATH}`, { headers: { "accept-encoding": "x-test-only" } })); + assert.notEqual(backendSeen[backendSeen.length - 1]?.encoding, "x-test-only", "the client's accept-encoding reached the agent"); + + const missing = await fetchOverStream(conn, callerBiscuit, new Request(`${base}/${AGENT_CARD_PATH}`, { headers: { "x-card": "missing" } })); + assert.equal(missing.status, 404); + assert.equal(await missing.text(), "no card"); + + const raw = await httpRequestOverStream(conn, callerBiscuit, "a2a://agent", `/${AGENT_CARD_PATH}`); + assert.deepEqual(JSON.parse(raw.text()), card); + + const refused = await fetchOverStream(conn, callerBiscuit, new Request(`${base}/${AGENT_CARD_PATH}`, { headers: { "x-card": "grpc-only" } })); + assert.equal(refused.status, 502); + assert.match(await refused.text(), /no supported interface the mesh can carry/); + + const other = await fetchOverStream(conn, callerBiscuit, new Request(`${base}/card`)); + assert.deepEqual(await other.json(), { path: "/card", echo: "" }); +}); + +test("rewriteAgentCard keeps only what the mesh carries", () => { + assert.throws(() => rewriteAgentCard([], "http://mesh/x"), /not a JSON object/); + assert.throws(() => rewriteAgentCard({ supportedInterfaces: [] }, "http://mesh/x"), /no supported interface/); + assert.deepEqual(rewriteAgentCard({ supportedInterfaces: [{ url: "x", protocolBinding: "http+json" }] }, "http://mesh/x"), { + supportedInterfaces: [{ url: "http://mesh/x", protocolBinding: "http+json" }], + }); +}); + test("a fetch handler sees the caller and the path, and its streaming body goes out as it is written", async () => { const conn = await dial(handlerAgent); const res = await httpRequestOverStream(conn, callerBiscuit, "a2a://agent", "/tasks?x=1", { method: "POST", headers: { "x-sam-biscuit": "spoof", "content-type": "text/plain" }, body: "hello" }); diff --git a/sdk/js/src/libp2p-http.ts b/sdk/js/src/libp2p-http.ts index afed2d8c..df95eee3 100644 --- a/sdk/js/src/libp2p-http.ts +++ b/sdk/js/src/libp2p-http.ts @@ -61,12 +61,17 @@ export const DEFAULT_A2A_NAME = "agent"; /** * The path prefix of a mesh URL, http://mesh/sam////: - * the shape of sam-node's egress proxy and of an agent card it rewrote. The - * host is ignored; the peer ID is in the path because URL parsers lowercase - * the host and a peer ID is case-sensitive. + * the shape of sam-node's egress proxy and of an agent card rewritten for the + * mesh, by sam-node or by this SDK. The host is ignored; the peer ID is in the + * path because URL parsers lowercase the host and a peer ID is case-sensitive. */ export const MESH_PATH_PREFIX = "/sam/"; +/** The well-known agent card location (A2A spec / RFC 8615). */ +export const AGENT_CARD_PATH = ".well-known/agent-card.json"; + +const MAX_AGENT_CARD_BYTES = 1 << 20; + /** Largest request body the ingress reads whole for a handler or url target. */ export const MAX_INGRESS_BODY_BYTES = 8 * 1024 * 1024; const REQUEST_TIMEOUT_MS = 60_000; @@ -387,6 +392,94 @@ export function splitMeshURL(url: URL): { peerId: string; target: string } { return { peerId, target: "/" + rest.join("/") + url.search }; } +/** The bare service root counts too: a2a-go treats a pathful base URL as the card location. */ +function agentCardService(method: string, target: string): string | undefined { + const m = /^\/a2a\/([^/?]+)(?:\/(?:\.well-known\/agent-card\.json)?)?(?:\?.*)?$/.exec(target); + return method === "GET" && m !== null ? `a2a://${m[1]}` : undefined; +} + +/** + * An agent card rebuilt for the mesh, as sam-node's egress proxy serves it: + * HTTP interfaces point at base, gRPC ones go, signatures no longer match. + * Streaming stays as declared, this transport streams. Throws when no interface remains. + */ +export function rewriteAgentCard(card: unknown, base: string): Record { + if (typeof card !== "object" || card === null || Array.isArray(card)) { + throw new Error("agent card is not a JSON object"); + } + const out = { ...(card as Record) }; + delete out.signatures; + const interfaces = Array.isArray(out.supportedInterfaces) ? (out.supportedInterfaces as unknown[]) : []; + const kept = interfaces.filter(carriedOverHTTP).map((iface) => ({ ...(iface as Record), url: base })); + if (kept.length === 0) { + throw new Error("agent card advertises no supported interface the mesh can carry (JSONRPC or HTTP+JSON); is the agent serving a pre-1.0 A2A card?"); + } + return { ...out, supportedInterfaces: kept }; +} + +/** Whether an interface's binding can traverse the mesh; gRPC needs its own connection. */ +function carriedOverHTTP(iface: unknown): boolean { + const binding = typeof iface === "object" && iface !== null ? (iface as { protocolBinding?: unknown }).protocolBinding : undefined; + return typeof binding === "string" && ["JSONRPC", "HTTP+JSON"].includes(binding.toUpperCase()); +} + +/** + * Impersonates the agent's card endpoint as sam-node's egress proxy does: holds + * the client's request, fetches the card itself with identity encoding, and + * answers with it regenerated; the agent's own non-200 is relayed as it is. + */ +async function serveAgentCard(conn: Connection, biscuit: Uint8Array, request: Request, options: HTTPStreamOptions, service: string): Promise { + const base = meshURL(conn.remotePeer.toString(), service); + const headers = new Headers(request.headers); + headers.delete("accept-encoding"); + let response: Response; + try { + response = await sendOverStream(conn, biscuit, new Request(`${base}/${AGENT_CARD_PATH}`, { headers, signal: request.signal }), options); + } catch (err) { + return badGateway(`agent card fetch failed: ${err instanceof Error ? err.message : String(err)}`); + } + if (response.status !== 200) { + return response; + } + let card: unknown; + try { + card = JSON.parse(new TextDecoder().decode(await readLimited(response.body, MAX_AGENT_CARD_BYTES))); + } catch { + return badGateway("agent card is not valid JSON"); + } + try { + return new Response(JSON.stringify(rewriteAgentCard(card, base)), { status: 200, headers: { "content-type": "application/json" } }); + } catch (err) { + return badGateway(err instanceof Error ? err.message : String(err)); + } +} + +function badGateway(reason: string): Response { + return new Response(`Bad Gateway: ${reason}`, { status: 502, headers: { "content-type": "text/plain" } }); +} + +/** The first limit bytes of a body, the rest dropped, as io.LimitReader bounds the node. */ +async function readLimited(body: ReadableStream | null, limit: number): Promise { + const chunks: Uint8Array[] = []; + let n = 0; + if (body !== null) { + const reader = body.getReader(); + for (let next = await reader.read(); !next.done && n < limit; next = await reader.read()) { + const chunk = next.value.subarray(0, limit - n); + chunks.push(chunk); + n += chunk.length; + } + await reader.cancel(); + } + const out = new Uint8Array(n); + let at = 0; + for (const chunk of chunks) { + out.set(chunk, at); + at += chunk.length; + } + return out; +} + export interface HTTPStreamOptions { /** The agent this request is made for. */ agent?: string; @@ -398,9 +491,16 @@ export interface HTTPStreamOptions { * Client side of /libp2p-http, as go-libp2p-http's RoundTripper: one stream * per request, plain HTTP/1.1 with Host set to the peer ID and the biscuit in * X-Sam-Biscuit. Resolves once the response headers are in; the body streams - * after, so an SSE response is consumed as the peer sends it. + * after, so an SSE response is consumed as the peer sends it. An agent card + * is served rewritten for the mesh (rewriteAgentCard), as sam-node serves one. */ export async function fetchOverStream(conn: Connection, biscuit: Uint8Array, request: Request, options: HTTPStreamOptions = {}): Promise { + const { target } = splitMeshURL(new URL(request.url)); + const service = agentCardService(request.method, target); + return service === undefined ? sendOverStream(conn, biscuit, request, options) : serveAgentCard(conn, biscuit, request, options, service); +} + +async function sendOverStream(conn: Connection, biscuit: Uint8Array, request: Request, options: HTTPStreamOptions): Promise { const { target } = splitMeshURL(new URL(request.url)); const peerId = conn.remotePeer.toString(); // One controller ends the exchange: the caller's signal at any time, the diff --git a/sdk/python/src/agent_mesh/__init__.py b/sdk/python/src/agent_mesh/__init__.py index e510d3fd..df1ebe96 100644 --- a/sdk/python/src/agent_mesh/__init__.py +++ b/sdk/python/src/agent_mesh/__init__.py @@ -34,6 +34,7 @@ from .httpx_transport import MESH_PATH_PREFIX, MeshTransport, split_mesh_url from .identity import Identity, canonical_peer_id, libp2p_public_key, peer_id_from_public_key, verify_ed25519 from .libp2p_http import ( + AGENT_CARD_PATH, DEFAULT_A2A_NAME, HTTP_PROTOCOL, A2AEndpoint, @@ -45,7 +46,9 @@ http_ingress_handler, http_request_over_stream, mesh_http_target, + mesh_url, open_http_request, + rewrite_agent_card, ) from .mcp_client import LabelsNotSatisfiedError, ToolCallResult, ToolInfo, open_mcp_session, require_egress_labels, require_labels from .mesh import AgentMesh, ControlPlaneSync @@ -57,6 +60,7 @@ __all__ = [ "A2AEndpoint", + "AGENT_CARD_PATH", "AUTH_PROTOCOL", "AdmittedRouter", "AgentMesh", @@ -111,6 +115,7 @@ "http_request_over_stream", "libp2p_public_key", "mesh_http_target", + "mesh_url", "open_http_request", "open_mcp_session", "parse_service_target", @@ -122,6 +127,7 @@ "require_role", "reserve_relay", "service_key", + "rewrite_agent_card", "split_mesh_url", "validate_control_plane_url", "verify_ed25519", diff --git a/sdk/python/src/agent_mesh/httpx_transport.py b/sdk/python/src/agent_mesh/httpx_transport.py index 888b433b..266560bc 100644 --- a/sdk/python/src/agent_mesh/httpx_transport.py +++ b/sdk/python/src/agent_mesh/httpx_transport.py @@ -20,9 +20,11 @@ await client.post("http://mesh/sam//a2a/agent", json=...) The URL has the shape sam-node's egress proxy takes, /sam/// -/, so an agent card a sam-node rewrote for the mesh works here as -it is. The host is ignored: httpx lowercases it, and a peer ID is not -case-insensitive. Response bodies stream, so `message/stream` works.""" +/. The card of an agent behind a sam-node comes back rewritten to +that shape, as the node's egress proxy serves it, so a stock A2A client +bootstraps from it unchanged. The host is ignored: httpx lowercases it, and a +peer ID is not case-insensitive. Response bodies stream, so `message/stream` +works.""" from __future__ import annotations @@ -30,12 +32,11 @@ import httpx -from .libp2p_http import StreamedResponse, open_http_request +from .libp2p_http import MESH_PATH_PREFIX, StreamedResponse, open_http_request if TYPE_CHECKING: from .session import MeshSession -MESH_PATH_PREFIX = "/sam/" _REQUEST_TIMEOUT = 60.0 diff --git a/sdk/python/src/agent_mesh/libp2p_http.py b/sdk/python/src/agent_mesh/libp2p_http.py index 05a826d8..6166673f 100644 --- a/sdk/python/src/agent_mesh/libp2p_http.py +++ b/sdk/python/src/agent_mesh/libp2p_http.py @@ -53,6 +53,19 @@ # The service name an agent answers under unless it picks another. DEFAULT_A2A_NAME = "agent" +# The path prefix of a mesh URL, http://mesh/sam////: +# the shape of sam-node's egress proxy and of an agent card rewritten for the +# mesh, by sam-node or by this SDK. The peer ID is in the path, not the host. +MESH_PATH_PREFIX = "/sam/" + +# The well-known agent card location (A2A spec / RFC 8615). +AGENT_CARD_PATH = ".well-known/agent-card.json" + +_MAX_AGENT_CARD_BYTES = 1 << 20 +# The bare service root counts too: a2a-go treats a pathful base URL as the card location. +_AGENT_CARD_TARGET = re.compile(r"^/a2a/([^/?]+)(?:/(?:\.well-known/agent-card\.json)?)?(?:\?.*)?$") +_HTTP_BINDINGS = ("JSONRPC", "HTTP+JSON") + _MAX_INGRESS_BODY_BYTES = 8 * 1024 * 1024 _READ_CHUNK = 64 * 1024 _REQUEST_TIMEOUT = 60.0 @@ -306,6 +319,34 @@ def mesh_http_target(target_service: str, path: str = "") -> str: return f"/{scheme}/{name}" + (path if path.startswith("/") else "/" + path) +def mesh_url(peer_id: str, target_service: str, path: str = "") -> str: + """The URL an httpx client on MeshTransport uses for a service on a peer: + http://mesh/sam////.""" + return "http://mesh" + MESH_PATH_PREFIX + peer_id + mesh_http_target(target_service, path) + + +def _agent_card_service(method: str, target: str) -> Optional[str]: + m = _AGENT_CARD_TARGET.match(target) + return f"a2a://{m.group(1)}" if method == "GET" and m else None + + +def rewrite_agent_card(card: object, base: str) -> dict: + """An agent card rebuilt for the mesh, as sam-node's egress proxy serves + it: HTTP interfaces point at base, gRPC ones go, signatures no longer match. + Streaming stays as declared, this transport streams. Raises when no interface remains.""" + if not isinstance(card, dict): + raise ValueError("agent card is not a JSON object") + interfaces = card.get("supportedInterfaces") + kept = [ + {**iface, "url": base} + for iface in (interfaces if isinstance(interfaces, list) else []) + if isinstance(iface, dict) and str(iface.get("protocolBinding", "")).upper() in _HTTP_BINDINGS + ] + if not kept: + raise ValueError("agent card advertises no supported interface the mesh can carry (JSONRPC or HTTP+JSON); is the agent serving a pre-1.0 A2A card?") + return {**{k: v for k, v in card.items() if k != "signatures"}, "supportedInterfaces": kept} + + class StreamedResponse: """A response whose body is still arriving on the stream: status and headers are in; `iter_body` yields the body as the peer sends it, which is @@ -341,6 +382,19 @@ async def aclose(self) -> None: await self._stream.close() +class _BufferedResponse(StreamedResponse): + def __init__(self, status: int, headers: dict[str, str], body: bytes) -> None: + self.status = status + self.headers = headers + self._payload = body + + async def iter_body(self) -> AsyncIterator[bytes]: + yield self._payload + + async def aclose(self) -> None: + return None + + async def open_http_request( host: IHost, peer_id: ID, @@ -356,7 +410,60 @@ async def open_http_request( """Client side of /libp2p-http, as go-libp2p-http's RoundTripper: one stream per request, plain HTTP/1.1 with Host set to the peer ID and the biscuit in X-Sam-Biscuit. Returns once the response headers are in; the - body streams after. The timeout bounds the headers, not the body.""" + body streams after. The timeout bounds the headers, not the body. An + agent card is served rewritten for the mesh (rewrite_agent_card), as + sam-node's egress proxy serves one.""" + service = _agent_card_service(method, target) + if service is None: + return await _open_http_request(host, peer_id, biscuit, method, target, headers=headers, body=body, agent=agent, timeout=timeout) + return await _serve_agent_card(host, peer_id, biscuit, headers, agent, timeout, service) + + +async def _serve_agent_card(host: IHost, peer_id: ID, biscuit: bytes, headers: Optional[Mapping[str, str]], agent: str, timeout: float, service: str) -> StreamedResponse: + """Impersonates the agent's card endpoint as sam-node's egress proxy does: + holds the client's request, fetches the card itself with identity encoding, + and answers with it regenerated; the agent's own non-200 is relayed as it is.""" + base = mesh_url(str(peer_id), service) + identity = {k: v for k, v in (headers or {}).items() if k.lower() != "accept-encoding"} + try: + response = await _open_http_request(host, peer_id, biscuit, "GET", mesh_http_target(service, AGENT_CARD_PATH), headers=identity, body=b"", agent=agent, timeout=timeout) + except Exception as err: # noqa: BLE001 - answered as sam-node's 502 + return _bad_gateway(f"agent card fetch failed: {err}") + if response.status != 200: + return response + try: + with trio.fail_after(timeout): + card = json.loads(await response.read(_MAX_AGENT_CARD_BYTES)) + except (json.JSONDecodeError, UnicodeDecodeError): + return _bad_gateway("agent card is not valid JSON") + except Exception as err: # noqa: BLE001 - answered as sam-node's 502 + return _bad_gateway(f"agent card fetch failed: {err}") + finally: + await response.aclose() + try: + payload = json.dumps(rewrite_agent_card(card, base)).encode() + except ValueError as err: + return _bad_gateway(str(err)) + return _BufferedResponse(200, {"content-type": "application/json", "content-length": str(len(payload))}, payload) + + +def _bad_gateway(reason: str) -> StreamedResponse: + payload = f"Bad Gateway: {reason}".encode() + return _BufferedResponse(502, {"content-type": "text/plain", "content-length": str(len(payload))}, payload) + + +async def _open_http_request( + host: IHost, + peer_id: ID, + biscuit: bytes, + method: str, + target: str, + *, + headers: Optional[Mapping[str, str]], + body: bytes, + agent: str, + timeout: float, +) -> StreamedResponse: out = [(k.lower(), v) for k, v in (headers or {}).items() if k.lower() not in ("host", "content-length", HEADER_SAM_BISCUIT, HEADER_PEER_ID)] out.append(("host", str(peer_id))) out.append((HEADER_SAM_BISCUIT, base64.b64encode(biscuit).decode())) @@ -412,9 +519,11 @@ async def http_request_over_stream( __all__: Sequence[str] = ( + "AGENT_CARD_PATH", "A2AEndpoint", "DEFAULT_A2A_NAME", "HTTP_PROTOCOL", + "MESH_PATH_PREFIX", "HTTPHandler", "HTTPRequest", "HTTPResponse", @@ -423,5 +532,7 @@ async def http_request_over_stream( "http_ingress_handler", "http_request_over_stream", "mesh_http_target", + "mesh_url", "open_http_request", + "rewrite_agent_card", ) diff --git a/sdk/python/src/agent_mesh/session.py b/sdk/python/src/agent_mesh/session.py index 4c8a1197..e658afb2 100644 --- a/sdk/python/src/agent_mesh/session.py +++ b/sdk/python/src/agent_mesh/session.py @@ -43,7 +43,6 @@ from .controlplane import ROLE_NODE from .discovery import DiscoveredProvider, find_peer, find_providers, parse_service_target, service_key from .host import create_mesh_host, dial, dial_addrs, peer_info -from .httpx_transport import MESH_PATH_PREFIX from .identity import canonical_peer_id from .libp2p_http import ( DEFAULT_A2A_NAME, @@ -55,6 +54,7 @@ http_ingress_handler, http_request_over_stream, mesh_http_target, + mesh_url, ) from .mcp_client import ToolCallResult, ToolInfo, open_mcp_session, require_egress_labels, tool_call_result from .relay import STOP_PROTOCOL, dial_through_relay, reserve_relay, split_circuit_address, stop_stream_handler @@ -171,8 +171,8 @@ def peer_id(self) -> str: def mesh_url(peer_id: str, target_service: str, path: str = "") -> str: """The URL an httpx client on `MeshTransport` uses for a service on a peer: http://mesh/sam////, the shape of - sam-node's egress proxy and of an agent card it rewrote.""" - return "http://mesh" + MESH_PATH_PREFIX + canonical_peer_id(peer_id) + mesh_http_target(target_service, path) + sam-node's egress proxy and of an agent card rewritten for the mesh.""" + return mesh_url(canonical_peer_id(peer_id), target_service, path) @property def agent_url(self) -> Optional[str]: diff --git a/sdk/python/tests/test_libp2p_http.py b/sdk/python/tests/test_libp2p_http.py index b231ab3e..651d3d6f 100644 --- a/sdk/python/tests/test_libp2p_http.py +++ b/sdk/python/tests/test_libp2p_http.py @@ -35,6 +35,7 @@ from agent_mesh.httpx_transport import MeshTransport, split_mesh_url from agent_mesh.identity import Identity from agent_mesh.libp2p_http import ( + AGENT_CARD_PATH, HTTP_PROTOCOL, A2AEndpoint, HTTPRequest, @@ -44,6 +45,7 @@ http_request_over_stream, mesh_http_target, open_http_request, + rewrite_agent_card, ) from agent_mesh.session import MeshSession @@ -64,6 +66,21 @@ def mint(peer_id: str, role: str) -> bytes: ).build(CP.private_key).to_bytes() +STOCK_CARD = { + "name": "echo-agent", + "version": "1.0.0", + "capabilities": {"streaming": True, "pushNotifications": False}, + "supportedInterfaces": [ + {"url": "http://127.0.0.1:7777/", "protocolBinding": "JSONRPC", "protocolVersion": "1.0"}, + {"url": "127.0.0.1:50051", "protocolBinding": "GRPC", "protocolVersion": "1.0"}, + ], + "signatures": [{"protected": "eyJhbGciOiJFUzI1NiJ9", "signature": "c3RhbGU"}], + "skills": [], + "defaultInputModes": ["text/plain"], + "defaultOutputModes": ["text/plain"], +} + + class FakeA2AServer(BaseHTTPRequestHandler): """Stands in for an A2A server beside the agent: echoes the request and, on /stream, answers with three SSE events as message/stream would.""" @@ -73,7 +90,23 @@ class FakeA2AServer(BaseHTTPRequestHandler): def _answer(self) -> None: length = int(self.headers.get("content-length") or 0) body = self.rfile.read(length) if length else b"" - FakeA2AServer.seen.append({"method": self.command, "path": self.path, "peer": self.headers.get("x-peer-id"), "body": body.decode(), "biscuit": self.headers.get("x-sam-biscuit")}) + FakeA2AServer.seen.append( + { + "method": self.command, + "path": self.path, + "peer": self.headers.get("x-peer-id"), + "body": body.decode(), + "biscuit": self.headers.get("x-sam-biscuit"), + "encoding": self.headers.get("accept-encoding"), + } + ) + if self.path == f"/{AGENT_CARD_PATH}" and self.headers.get("x-card") == "missing": + self.send_response(404) + self.send_header("content-type", "text/plain") + self.send_header("content-length", "7") + self.end_headers() + self.wfile.write(b"no card") + return if self.path == "/stream": self.send_response(200) self.send_header("content-type", "text/event-stream") @@ -82,7 +115,12 @@ def _answer(self) -> None: self.wfile.write(f"data: {json.dumps({'event': i})}\n\n".encode()) self.wfile.flush() return - payload = json.dumps({"path": self.path, "echo": body.decode()}).encode() + if self.path == f"/{AGENT_CARD_PATH}": + grpc_only = self.headers.get("x-card") == "grpc-only" + interfaces = [i for i in STOCK_CARD["supportedInterfaces"] if not grpc_only or i["protocolBinding"] == "GRPC"] + payload = json.dumps({**STOCK_CARD, "supportedInterfaces": interfaces}).encode() + else: + payload = json.dumps({"path": self.path, "echo": body.decode()}).encode() self.send_response(200) self.send_header("content-type", "application/json") self.send_header("x-backend", "fake") @@ -223,6 +261,55 @@ async def with_timeout(): trio.run(with_timeout) +def test_the_card_of_an_agent_behind_a_node_comes_back_rewritten_for_the_mesh(backend_url): + async def main(): + async with trio.open_nursery() as nursery: + agent, addr = await start_agent(nursery, backend_url, []) + caller_identity = Identity.generate() + caller = libp2p_host(caller_identity) + caller_biscuit = mint(caller_identity.peer_id, ROLE_NODE) + async with caller.run(listen_addrs=[]): + await caller.connect(info_from_p2p_addr(addr)) + pid = agent.get_id() + base = MeshSession.mesh_url(str(pid), "a2a://agent") + async with httpx.AsyncClient(transport=MeshTransport(FakeSession(caller, caller_biscuit))) as client: # type: ignore[arg-type] + answer = await client.get(f"{base}/{AGENT_CARD_PATH}") + assert answer.status_code == 200 and answer.headers["content-type"] == "application/json" + card = answer.json() + assert card["supportedInterfaces"] == [{"url": base, "protocolBinding": "JSONRPC", "protocolVersion": "1.0"}] + assert card["capabilities"] == {"streaming": True, "pushNotifications": False} + assert "signatures" not in card and card["name"] == "echo-agent" and card["skills"] == [] + assert FakeA2AServer.seen[-1]["path"] == f"/{AGENT_CARD_PATH}" + # The bare service root serves the card too, the way a2a-go resolves a pathful base URL, + # and the fetch is the SDK's own: at the well-known path, without the client's accept-encoding. + assert (await client.get(base, headers={"accept-encoding": "x-test-only"})).json() == card + assert FakeA2AServer.seen[-1]["path"] == f"/{AGENT_CARD_PATH}" and FakeA2AServer.seen[-1]["encoding"] != "x-test-only" + missing = await client.get(f"{base}/{AGENT_CARD_PATH}", headers={"x-card": "missing"}) + assert missing.status_code == 404 and missing.text == "no card" + refused = await client.get(f"{base}/{AGENT_CARD_PATH}", headers={"x-card": "grpc-only"}) + assert refused.status_code == 502 and "no supported interface the mesh can carry" in refused.text + assert (await client.get(f"{base}/card")).json() == {"path": "/card", "echo": ""} + raw = await http_request_over_stream(caller, pid, caller_biscuit, "a2a://agent", f"/{AGENT_CARD_PATH}") + assert raw.status == 200 and raw.json() == card + nursery.cancel_scope.cancel() + + async def with_timeout(): + with trio.fail_after(30): + await main() + + trio.run(with_timeout) + + +def test_rewrite_agent_card_keeps_only_what_the_mesh_carries(): + with pytest.raises(ValueError, match="not a JSON object"): + rewrite_agent_card([], "http://mesh/x") + with pytest.raises(ValueError, match="no supported interface"): + rewrite_agent_card({"supportedInterfaces": []}, "http://mesh/x") + assert rewrite_agent_card({"supportedInterfaces": [{"url": "x", "protocolBinding": "http+json"}]}, "http://mesh/x") == { + "supportedInterfaces": [{"url": "http://mesh/x", "protocolBinding": "http+json"}], + } + + def test_an_in_process_handler_sees_the_verified_caller(): async def hello(request: HTTPRequest, caller) -> HTTPResponse: return HTTPResponse(status=201, body=f"hello {caller.peer_id} {request.method} {request.path}".encode()) @@ -278,7 +365,7 @@ async def closing(stream) -> None: async def main(): async with trio.open_nursery() as nursery: server = libp2p_host(Identity.generate()) - handlers = {"/a2a/cut": truncating, "/a2a/end": closing} + handlers = {"/inference/cut": truncating, "/inference/end": closing} async def route(stream) -> None: request = b"" @@ -304,12 +391,12 @@ async def run(): await caller.connect(info_from_p2p_addr(addr_box[0])) pid = server.get_id() with trio.fail_after(10): - ended = await open_http_request(caller, pid, b"x", "GET", "/a2a/end") + ended = await open_http_request(caller, pid, b"x", "GET", "/inference/end") assert ended.status == 200 assert (await ended.read()).decode().count("data:") == 2 await ended.aclose() - cut = await open_http_request(caller, pid, b"x", "GET", "/a2a/cut") + cut = await open_http_request(caller, pid, b"x", "GET", "/inference/cut") assert cut.status == 200 with pytest.raises(ConnectionError, match="mid-response"): await cut.read() diff --git a/site/content/docs/guides/native-sdks.md b/site/content/docs/guides/native-sdks.md index a6ef47c6..584c40b0 100644 --- a/site/content/docs/guides/native-sdks.md +++ b/site/content/docs/guides/native-sdks.md @@ -465,6 +465,10 @@ plug the mesh in as its transport. agent card, and what an A2A client on the mesh is given. The peer ID is in the path because URL parsers lowercase the host and a peer ID is case-sensitive. +- An agent behind a `sam-node` serves a card naming its own address. The + SDK hands the client that card rewritten for the mesh, as `sam-node`'s + egress proxy does, so a stock A2A client given the mesh URL bootstraps + from it unchanged. - The A2A JavaScript SDK takes a `fetch` for its client and mounts its server as Express handlers. `session.fetch()` is the fetch; `acceptA2A({ listener: app })` runs the Express app on the mesh, with