diff --git a/.dockerignore b/.dockerignore index ca4acce1..05a6cf7b 100644 --- a/.dockerignore +++ b/.dockerignore @@ -33,6 +33,7 @@ __pycache__/ *.out coverage.* mobile/sam-node-app/build/ +sdk/js/build/ site/public/ site/resources/ rootfs.ext4 diff --git a/.github/k8s/sam-probe-cronjob-template.yaml b/.github/k8s/sam-probe-cronjob-template.yaml index b3053864..2b3943de 100644 --- a/.github/k8s/sam-probe-cronjob-template.yaml +++ b/.github/k8s/sam-probe-cronjob-template.yaml @@ -1,9 +1,10 @@ # The cold path, on a schedule: a node with no identity enrolls with the # control plane, authenticates to a router, finds a service on the mesh and -# speaks MCP to it through the mesh, then exits. The long-lived canaries only -# do this when a pod restarts, so between rollouts nobody re-checks that a new -# user could still join. Each run enrolls a fresh identity; the control -# plane's --node-retention sweep reclaims the rows. +# speaks MCP to it through the mesh, then reaches the services the SDK +# canaries serve, and exits. The long-lived canaries only do this when a pod +# restarts, so between rollouts nobody re-checks that a new user could still +# join. Each run enrolls a fresh identity; the control plane's +# --node-retention sweep reclaims the rows. apiVersion: batch/v1 kind: CronJob metadata: @@ -19,8 +20,9 @@ spec: # One attempt: a retry would hide exactly the flakiness this exists to # measure. A failed Job is the signal. backoffLimit: 0 - # 180s to be ready plus 120s to reach a provider, with headroom. - activeDeadlineSeconds: 360 + # 180s to be ready, 120s to reach a provider, 120s for the SDK + # canaries, with headroom. + activeDeadlineSeconds: 480 template: metadata: labels: @@ -78,6 +80,10 @@ spec: # The everything canary's MCP server (sam-node-everything-template.yaml): # a real server, so the node advertises it and initialize succeeds. value: everything + - name: SDK_SERVICES + # What the SDK canaries publish (sam-sdk-canary-template.yaml), + # each with a `greet` tool and an A2A endpoint. + value: "greeter-js greeter-py" resources: requests: cpu: 10m @@ -157,8 +163,57 @@ spec: done REACH_S=$(( $(date +%s) - REACH_START )) - printf '{"probe":"sam-cold-path","ok":true,"ready_s":%d,"router_latency_ms":%s,"connected_peers":%s,"reach_s":%d,"providers_tried":%d,"call_s":%s,"provider":"%s","elapsed_s":%d}\n' \ - "$READY_S" "${ROUTER_MS:-null}" "${PEERS:-null}" "$REACH_S" "$TRIED" "$CALL_S" "$PEER" "$(( $(date +%s) - START ))" + # 4. A service an SDK member serves (sam-sdk-canary-template.yaml) + # is reachable from a node the way an agent behind a node + # reaches it: its MCP tool through the node's own MCP API + # (call_remote_tool opens /sam/mcp/1.0.0 to the provider), and + # its A2A endpoint through the egress proxy. The node's MCP + # endpoint is stateful, so one session is opened first. + MCP=http://localhost/mcp + SESSION="" + mcp_post() { + curl -s --unix-socket "$SOCK" -o "$2" -D /tmp/mcp.hdr -X POST "$MCP" \ + -H 'Content-Type: application/json' -H 'Accept: application/json, text/event-stream' \ + ${SESSION:+-H "Mcp-Session-Id: $SESSION"} -d "$1" + } + mcp_post '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"sam-probe","version":"0"}}}' /tmp/mcp-init.out \ + || fail sdk "initialize on the node's MCP API failed" + SESSION=$(grep -i '^mcp-session-id:' /tmp/mcp.hdr | tr -d '\r' | cut -d' ' -f2) + [ -n "$SESSION" ] || fail sdk "the node's MCP API returned no session" + mcp_post '{"jsonrpc":"2.0","method":"notifications/initialized"}' /dev/null + + SDK_START=$(date +%s) + SDK_PROVIDERS="" + for svc in ${SDK_SERVICES}; do + SDK_PEER="" + LAST_ERR="" + while [ -z "$SDK_PEER" ]; do + if curl -sf --unix-socket "$SOCK" -o /tmp/sdk-providers.json \ + "http://localhost/sam/service/discover?type=mcp&name=${svc}&timeout=20s"; then + for p in $(grep -oE '"peer_id":"[^"]+"' /tmp/sdk-providers.json | cut -d'"' -f4); do + mcp_post '{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"call_remote_tool","arguments":{"peer_id":"'"$p"'","tool_name":"mcp://'"$svc"'/greet","arguments":{"name":"sam-probe"}}}}' /tmp/sdk-call.out + if grep -q 'hello sam-probe' /tmp/sdk-call.out; then + CODE=$(curl -s --unix-socket "$SOCK" -o /tmp/sdk-card.out -w '%{http_code}' "http://localhost/sam/${p}/a2a/${svc}/card") + if [ "$CODE" = "200" ] && grep -q '"caller"' /tmp/sdk-card.out; then + SDK_PEER=$p + break + fi + LAST_ERR="${p}: a2a card HTTP ${CODE} $(head -c 120 /tmp/sdk-card.out | tr -d '"\n')" + else + LAST_ERR="${p}: $(head -c 160 /tmp/sdk-call.out | tr -d '"\n')" + fi + done + fi + [ -n "$SDK_PEER" ] && break + [ $(( $(date +%s) - SDK_START )) -ge 120 ] && fail sdk "no provider of ${svc} answered greet and its card in 120s (last: ${LAST_ERR:-none discovered})" + sleep 5 + done + SDK_PROVIDERS="${SDK_PROVIDERS}${SDK_PROVIDERS:+,}\"${svc}\":\"${SDK_PEER}\"" + done + SDK_S=$(( $(date +%s) - SDK_START )) + + printf '{"probe":"sam-cold-path","ok":true,"ready_s":%d,"router_latency_ms":%s,"connected_peers":%s,"reach_s":%d,"providers_tried":%d,"call_s":%s,"provider":"%s","sdk_s":%d,"sdk_providers":{%s},"elapsed_s":%d}\n' \ + "$READY_S" "${ROUTER_MS:-null}" "${PEERS:-null}" "$REACH_S" "$TRIED" "$CALL_S" "$PEER" "$SDK_S" "$SDK_PROVIDERS" "$(( $(date +%s) - START ))" volumes: - name: config-volume configMap: diff --git a/.github/k8s/sam-sdk-canary-template.yaml b/.github/k8s/sam-sdk-canary-template.yaml new file mode 100644 index 00000000..16fb75d7 --- /dev/null +++ b/.github/k8s/sam-sdk-canary-template.yaml @@ -0,0 +1,124 @@ +# Two members with no sam-node beside them: the JavaScript and Python SDK +# example servers (sdk/js/examples/serve.ts, sdk/python/examples/serve.py), +# unchanged, each publishing mcp://greeter- and a2a://greeter-. +# They enroll with the pod's projected service account token, the way every +# other canary does, and keep their identity in an emptyDir, so a restart +# resumes and a new pod enrolls afresh. What the mesh sees from them is what a +# user of the packages gets. +apiVersion: apps/v1 +kind: Deployment +metadata: + name: js-canary-${ENV_NAME} + namespace: sam-canary-${ENV_NAME} +spec: + replicas: 1 + selector: + matchLabels: + app: js-canary-${ENV_NAME} + template: + metadata: + labels: + app: js-canary-${ENV_NAME} + sam-canary: "true" + spec: + serviceAccountName: sam-node-sa + containers: + - name: sdk + image: ghcr.io/google/sam-sdk-js:${IMAGE_TAG} + args: ["build/examples/serve.js", "greeter-js"] + env: + - name: SAM_CONTROL_PLANE_URL + value: "http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" + - name: SAM_INSECURE_CONTROL_PLANE + value: "true" + - name: SAM_JWT_PATH + value: /var/run/secrets/tokens/sam-token + - name: SAM_STATE_DIR + value: /var/run/sam/state + securityContext: + allowPrivilegeEscalation: false + capabilities: + drop: ["ALL"] + resources: + requests: + cpu: 50m + memory: 128Mi + limits: + cpu: 500m + memory: 512Mi + volumeMounts: + - name: sam-token + mountPath: /var/run/secrets/tokens + readOnly: true + - name: sam-state + mountPath: /var/run/sam + volumes: + - name: sam-state + emptyDir: {} + - name: sam-token + projected: + sources: + - serviceAccountToken: + path: sam-token + expirationSeconds: 3600 + audience: "sam-control-plane-audience" +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + name: python-canary-${ENV_NAME} + namespace: sam-canary-${ENV_NAME} +spec: + replicas: 1 + selector: + matchLabels: + app: python-canary-${ENV_NAME} + template: + metadata: + labels: + app: python-canary-${ENV_NAME} + sam-canary: "true" + spec: + serviceAccountName: sam-node-sa + containers: + - name: sdk + image: ghcr.io/google/sam-sdk-python:${IMAGE_TAG} + args: ["examples/serve.py", "greeter-py"] + env: + - name: SAM_CONTROL_PLANE_URL + value: "http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" + - name: SAM_INSECURE_CONTROL_PLANE + value: "true" + - name: SAM_JWT_PATH + value: /var/run/secrets/tokens/sam-token + - name: SAM_STATE_DIR + value: /var/run/sam/state + - name: PYTHONUNBUFFERED + value: "1" + securityContext: + allowPrivilegeEscalation: false + capabilities: + drop: ["ALL"] + resources: + requests: + cpu: 50m + memory: 128Mi + limits: + cpu: 500m + memory: 512Mi + volumeMounts: + - name: sam-token + mountPath: /var/run/secrets/tokens + readOnly: true + - name: sam-state + mountPath: /var/run/sam + volumes: + - name: sam-state + emptyDir: {} + - name: sam-token + projected: + sources: + - serviceAccountToken: + path: sam-token + expirationSeconds: 3600 + audience: "sam-control-plane-audience" diff --git a/.github/k8s/sam-sdk-probe-cronjob-template.yaml b/.github/k8s/sam-sdk-probe-cronjob-template.yaml new file mode 100644 index 00000000..a02239b9 --- /dev/null +++ b/.github/k8s/sam-sdk-probe-cronjob-template.yaml @@ -0,0 +1,183 @@ +# The cold path through the SDKs, on a schedule: a fresh identity enrolls with +# the control plane through each SDK, joins through a router, and calls out +# across every implementation boundary a user can cross: an MCP service a +# sam-node serves (the everything canary), and an MCP tool and an A2A endpoint +# the other language's SDK serves (the sdk canaries). The programs are the +# examples the READMEs embed, unchanged; the shell around them only retries +# discovery, which a member that joined seconds ago can miss, and prints one +# verdict line. Each run enrolls a fresh identity per language; the control +# plane's --node-retention sweep reclaims the rows. +apiVersion: batch/v1 +kind: CronJob +metadata: + name: sam-sdk-probe-${ENV_NAME} + namespace: sam-canary-${ENV_NAME} +spec: + schedule: "7,22,37,52 * * * *" + concurrencyPolicy: Forbid + successfulJobsHistoryLimit: 3 + failedJobsHistoryLimit: 5 + jobTemplate: + spec: + # One attempt: a retry would hide exactly the flakiness this exists to + # measure. A failed Job is the signal; both containers must succeed. + backoffLimit: 0 + # Three calls with up to 120s of discovery each, plus enrollment. + activeDeadlineSeconds: 480 + template: + metadata: + labels: + app: sam-sdk-probe-${ENV_NAME} + spec: + # Bound to the sam-canary role by deploy.yaml's policy seed, which + # is what lets the probe call the canaries' services. + serviceAccountName: sam-node-sa + restartPolicy: Never + containers: + - name: js + image: ghcr.io/google/sam-sdk-js:${IMAGE_TAG} + env: + - name: SAM_CONTROL_PLANE_URL + value: "http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" + - name: SAM_INSECURE_CONTROL_PLANE + value: "true" + - name: SAM_JWT_PATH + value: /var/run/secrets/tokens/sam-token + - name: SAM_STATE_DIR + value: /var/run/sam/js + - name: PROBE + value: sam-sdk-js + - name: CALL + value: "node build/examples/call.js" + - name: PEER_SERVICE + # The other SDK's canary (sam-sdk-canary-template.yaml). + value: greeter-py + securityContext: + allowPrivilegeEscalation: false + capabilities: + drop: ["ALL"] + resources: + requests: + cpu: 50m + memory: 128Mi + limits: + cpu: 500m + memory: 512Mi + volumeMounts: + - name: sam-token + mountPath: /var/run/secrets/tokens + readOnly: true + - name: sam-state + mountPath: /var/run/sam + - name: probe-script + mountPath: /etc/sam-probe + readOnly: true + command: ["/bin/sh", "/etc/sam-probe/probe.sh"] + - name: python + image: ghcr.io/google/sam-sdk-python:${IMAGE_TAG} + env: + - name: SAM_CONTROL_PLANE_URL + value: "http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" + - name: SAM_INSECURE_CONTROL_PLANE + value: "true" + - name: SAM_JWT_PATH + value: /var/run/secrets/tokens/sam-token + - name: SAM_STATE_DIR + value: /var/run/sam/python + - name: PYTHONUNBUFFERED + value: "1" + - name: PROBE + value: sam-sdk-python + - name: CALL + value: "python examples/call.py" + - name: PEER_SERVICE + value: greeter-js + securityContext: + allowPrivilegeEscalation: false + capabilities: + drop: ["ALL"] + resources: + requests: + cpu: 50m + memory: 128Mi + limits: + cpu: 500m + memory: 512Mi + volumeMounts: + - name: sam-token + mountPath: /var/run/secrets/tokens + readOnly: true + - name: sam-state + mountPath: /var/run/sam + - name: probe-script + mountPath: /etc/sam-probe + readOnly: true + command: ["/bin/sh", "/etc/sam-probe/probe.sh"] + volumes: + - name: sam-state + emptyDir: {} + - name: probe-script + configMap: + name: sam-sdk-probe-script-${ENV_NAME} + - name: sam-token + projected: + sources: + - serviceAccountToken: + path: sam-token + expirationSeconds: 3600 + audience: "sam-control-plane-audience" +--- +apiVersion: v1 +kind: ConfigMap +metadata: + name: sam-sdk-probe-script-${ENV_NAME} + namespace: sam-canary-${ENV_NAME} +data: + probe.sh: | + #!/bin/sh + # Runs the SDK's call example against three services and prints one JSON + # verdict line. PROBE names the run, CALL is the example's command line, + # PEER_SERVICE is the other SDK's canary. Discovery is retried: a member + # that joined seconds ago has a thin routing table, and a provider record + # can name a pod a rollout just replaced. The first call enrolls with the + # projected token; the later ones resume from SAM_STATE_DIR. + set -u + START=$(date +%s) + + fail() { + printf '{"probe":"%s","ok":false,"stage":"%s","error":"%s","elapsed_s":%d}\n' \ + "$PROBE" "$1" "$(printf '%s' "$2" | tr -d '"\n' | cut -c1-300)" "$(( $(date +%s) - START ))" + exit 1 + } + + # call ; sets PROVIDER. + call() { + stage=$1; want=$2; shift 2 + stage_start=$(date +%s) + attempts=0 + while :; do + attempts=$((attempts + 1)) + out=$($CALL "$@" 2>&1); rc=$? + if [ $rc -eq 0 ] && printf '%s' "$out" | grep -q "$want"; then + PROVIDER=$(printf '%s' "$out" | grep -o 'is served by [^ ]*' | tail -1 | cut -d' ' -f4) + return 0 + fi + [ $(( $(date +%s) - stage_start )) -ge 120 ] && fail "$stage" "no answer with '$want' in 120s ($attempts attempts; last: $(printf '%s' "$out" | tail -c 200))" + sleep 5 + done + } + + # 1. Enroll, join and call an MCP service a sam-node serves. + call node-mcp 'Echo: sam-sdk-probe' mcp://everything echo '{"message": "sam-sdk-probe"}' + NODE_PEER=$PROVIDER + NODE_S=$(( $(date +%s) - START )) + + # 2. The other SDK's MCP tool, then its A2A endpoint over /libp2p-http. + SDK_START=$(date +%s) + call sdk-mcp 'hello sam-sdk-probe' "mcp://${PEER_SERVICE}" greet '{"name": "sam-sdk-probe"}' + SDK_PEER=$PROVIDER + call sdk-a2a '"caller"' "a2a://${PEER_SERVICE}" /card + SDK_S=$(( $(date +%s) - SDK_START )) + + printf '{"probe":"%s","ok":true,"node_s":%d,"node_provider":"%s","sdk_s":%d,"sdk_provider":"%s","elapsed_s":%d}\n' \ + "$PROBE" "$NODE_S" "$NODE_PEER" "$SDK_S" "$SDK_PEER" "$(( $(date +%s) - START ))" diff --git a/.github/workflows/deploy.yaml b/.github/workflows/deploy.yaml index b6b40def..84425cc3 100644 --- a/.github/workflows/deploy.yaml +++ b/.github/workflows/deploy.yaml @@ -208,6 +208,54 @@ jobs: tags: ${{ steps.meta-sam-one.outputs.tags }} labels: ${{ steps.meta-sam-one.outputs.labels }} + # The SDK images carry the example programs the docs embed, as the + # testnet canaries and probes run them. Built from the commit, not from + # npm or PyPI, so bananas exercises the SDK at the same commit as the + # Go components it talks to. + - name: Extract metadata for sam-sdk-js + id: meta-sdk-js + uses: docker/metadata-action@dc802804100637a589fabce1cb79ff13a1411302 # v6.2.0 + with: + images: ${{ env.REGISTRY }}/google/sam-sdk-js + tags: | + type=ref,event=branch + type=ref,event=tag + type=raw,value=${{ github.sha }} + type=raw,value=latest,enable={{is_default_branch}} + type=raw,value=stable,enable=${{ startsWith(github.ref, 'refs/tags/v') }} + + - name: Build and push sam-sdk-js image + uses: docker/build-push-action@c3c9e263c25d99ce0380d002d59b67737d91b0dc # v7.4.0 + with: + context: . + file: Dockerfile.sam-sdk-js + platforms: ${{ github.ref_type == 'tag' && 'linux/amd64,linux/arm64' || 'linux/amd64' }} + push: true + tags: ${{ steps.meta-sdk-js.outputs.tags }} + labels: ${{ steps.meta-sdk-js.outputs.labels }} + + - name: Extract metadata for sam-sdk-python + id: meta-sdk-python + uses: docker/metadata-action@dc802804100637a589fabce1cb79ff13a1411302 # v6.2.0 + with: + images: ${{ env.REGISTRY }}/google/sam-sdk-python + tags: | + type=ref,event=branch + type=ref,event=tag + type=raw,value=${{ github.sha }} + type=raw,value=latest,enable={{is_default_branch}} + type=raw,value=stable,enable=${{ startsWith(github.ref, 'refs/tags/v') }} + + - name: Build and push sam-sdk-python image + uses: docker/build-push-action@c3c9e263c25d99ce0380d002d59b67737d91b0dc # v7.4.0 + with: + context: . + file: Dockerfile.sam-sdk-python + platforms: ${{ github.ref_type == 'tag' && 'linux/amd64,linux/arm64' || 'linux/amd64' }} + push: true + tags: ${{ steps.meta-sdk-python.outputs.tags }} + labels: ${{ steps.meta-sdk-python.outputs.labels }} + deploy: name: Deploy to GKE runs-on: ubuntu-latest @@ -758,6 +806,35 @@ jobs: exit 1 } + # Two members with no sam-node: the SDK example servers, each + # publishing an MCP tool and an A2A endpoint. What the probes below + # reach from a node and from the other SDK. + - name: Deploy SDK Canaries + env: + VAR_ENV_NAME: ${{ vars.ENV_NAME }} + VAR_IMAGE_TAG: ${{ env.IMAGE_TAG }} + run: | + export ENV_NAME="${VAR_ENV_NAME}" + export CANARY_NAMESPACE="sam-canary-${ENV_NAME}" + export NAMESPACE="sam-${ENV_NAME}" + export IMAGE_TAG="${VAR_IMAGE_TAG}" + + envsubst '${ENV_NAME} ${NAMESPACE} ${IMAGE_TAG}' < .github/k8s/sam-sdk-canary-template.yaml | kubectl apply -f - + for lang in js python; do + kubectl rollout status deployment/${lang}-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} --timeout=180s || { + echo "${lang} SDK canary deployment failed!" + kubectl describe deployment/${lang}-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} || true + kubectl get pods -n ${CANARY_NAMESPACE} -l app=${lang}-canary-${ENV_NAME} -o wide || true + kubectl describe pods -n ${CANARY_NAMESPACE} -l app=${lang}-canary-${ENV_NAME} || true + kubectl logs -n ${CANARY_NAMESPACE} -l app=${lang}-canary-${ENV_NAME} --tail=-1 || true + + echo "Rolling back ${lang} SDK canary deployment..." + kubectl rollout undo deployment/${lang}-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} + kubectl rollout status deployment/${lang}-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} || true + exit 1 + } + done + - name: Deploy OpenRouter Node to Testnet env: OPENROUTER_API_KEY: ${{ secrets.OPENROUTER_API_KEY }} @@ -857,3 +934,47 @@ jobs: ;; esac + # The same cold path through each SDK: a fresh identity enrolls with + # the JavaScript and the Python SDK, reaches the everything canary's + # MCP server (a sam-node) and the other SDK's canary, MCP and A2A. With + # the node probe above, every pair of implementations is crossed in + # both directions on every rollout and every 15 minutes after. + - name: Deploy and run the SDK cold-path probe + env: + VAR_ENV_NAME: ${{ vars.ENV_NAME }} + VAR_IMAGE_TAG: ${{ env.IMAGE_TAG }} + run: | + export ENV_NAME="${VAR_ENV_NAME}" + export CANARY_NAMESPACE="sam-canary-${ENV_NAME}" + export NAMESPACE="sam-${ENV_NAME}" + export IMAGE_TAG="${VAR_IMAGE_TAG}" + JOB="sam-sdk-probe-rollout-${ENV_NAME}" + + envsubst '${ENV_NAME} ${NAMESPACE} ${IMAGE_TAG}' < .github/k8s/sam-sdk-probe-cronjob-template.yaml | kubectl apply -f - + + kubectl delete job "${JOB}" -n "${CANARY_NAMESPACE}" --ignore-not-found + kubectl create job --from="cronjob/sam-sdk-probe-${ENV_NAME}" "${JOB}" -n "${CANARY_NAMESPACE}" + + for _ in $(seq 1 96); do + status=$(kubectl get job "${JOB}" -n "${CANARY_NAMESPACE}" -o jsonpath='{range .status.conditions[*]}{.type}={.status}{"\n"}{end}') + case "${status}" in + *Complete=True*) break ;; + *Failed=True*) break ;; + esac + sleep 5 + done + + for lang in js python; do + echo "==== ${lang} SDK probe output ====" + kubectl logs "job/${JOB}" -n "${CANARY_NAMESPACE}" -c "${lang}" --tail=-1 || true + done + + case "${status:-}" in + *Complete=True*) echo "SDK cold-path probe passed." ;; + *) + echo "SDK cold-path probe did not pass: an SDK could not enroll, reach a router, or call a service a node or the other SDK serves." + kubectl describe job "${JOB}" -n "${CANARY_NAMESPACE}" || true + exit 1 + ;; + esac + diff --git a/.github/workflows/sdk.yml b/.github/workflows/sdk.yml index 21044d87..1e83b8bf 100644 --- a/.github/workflows/sdk.yml +++ b/.github/workflows/sdk.yml @@ -30,6 +30,7 @@ on: - 'hack/verify-sdk-generated.sh' - 'site/content/docs/guides/native-sdks.md' - 'tests/integration/sdk_*_test.go' + - 'Dockerfile.sam-sdk-*' - '.github/workflows/sdk.yml' pull_request: branches: [main] @@ -41,6 +42,7 @@ on: - 'hack/verify-sdk-generated.sh' - 'site/content/docs/guides/native-sdks.md' - 'tests/integration/sdk_*_test.go' + - 'Dockerfile.sam-sdk-*' - '.github/workflows/sdk.yml' permissions: @@ -88,3 +90,13 @@ jobs: - name: Both SDKs and the documented examples on a mesh with a real control plane, router and sam-node run: go test ./tests/integration -run TestNativeSDK -count=1 -v + + # The images the testnets run the examples from (deploy.yaml builds and + # pushes them on main); built here so a broken Dockerfile fails the + # pull request, not the deploy. + - name: SDK canary images build + run: | + docker build -f Dockerfile.sam-sdk-js -t sam-sdk-js:ci . + docker build -f Dockerfile.sam-sdk-python -t sam-sdk-python:ci . + docker run --rm sam-sdk-js:ci build/examples/call.js 2>&1 | grep -q 'exactly one of' + docker run --rm sam-sdk-python:ci examples/call.py 2>&1 | grep -q 'exactly one of' diff --git a/Dockerfile.sam-sdk-js b/Dockerfile.sam-sdk-js new file mode 100644 index 00000000..c9b6fd0a --- /dev/null +++ b/Dockerfile.sam-sdk-js @@ -0,0 +1,26 @@ +# The JavaScript SDK with its example programs, as the testnet canaries run +# them: `node build/examples/serve.js ` publishes services on the mesh, +# `node build/examples/call.js [args]` calls one. The +# examples import "@sam-mesh/sdk" by name and resolve it through the package's +# own exports, so the image keeps the package layout: package.json, dist/, +# build/examples/ and the production node_modules/. + +# Stage 1: Build +FROM node:22-bookworm-slim@sha256:43ac6c60b8f89723f746e8a92ce91abd5017e627ce1ddfe4238355d3a30b772c AS builder +WORKDIR /src +COPY sdk/js/package.json sdk/js/package-lock.json sdk/js/.npmrc ./ +RUN npm ci --no-fund --no-audit +COPY sdk/js/ ./ +RUN npm run build && npm run examples && npm prune --omit=dev --no-fund --no-audit + +# Stage 2: Final +FROM node:22-bookworm-slim@sha256:43ac6c60b8f89723f746e8a92ce91abd5017e627ce1ddfe4238355d3a30b772c +WORKDIR /app +COPY --from=builder --chown=node:node /src/package.json ./ +COPY --from=builder --chown=node:node /src/dist ./dist +COPY --from=builder --chown=node:node /src/build/examples ./build/examples +COPY --from=builder --chown=node:node /src/node_modules ./node_modules +USER node +ENV HOME=/home/node +ENTRYPOINT ["node"] +CMD ["build/examples/serve.js"] diff --git a/Dockerfile.sam-sdk-python b/Dockerfile.sam-sdk-python new file mode 100644 index 00000000..d7ba4094 --- /dev/null +++ b/Dockerfile.sam-sdk-python @@ -0,0 +1,23 @@ +# The Python SDK with its example programs, as the testnet canaries run them: +# `python examples/serve.py ` publishes services on the mesh, +# `python examples/call.py [args]` calls one. fastecdsa, +# a py-libp2p dependency, builds from source against GMP, so the build stage +# carries a compiler and the final image only libgmp10. + +# Stage 1: Build +FROM python:3.12-slim-bookworm@sha256:392307d22300de8b5986851a12d9176dfc0fc073e65bf6523ebd7dcbeb23564e AS builder +RUN apt-get update && apt-get install -y --no-install-recommends gcc libgmp-dev python3-dev && rm -rf /var/lib/apt/lists/* +COPY sdk/python /src +RUN python -m venv /opt/venv && /opt/venv/bin/pip install --no-cache-dir /src + +# Stage 2: Final +FROM python:3.12-slim-bookworm@sha256:392307d22300de8b5986851a12d9176dfc0fc073e65bf6523ebd7dcbeb23564e +RUN apt-get update && apt-get install -y --no-install-recommends libgmp10 && rm -rf /var/lib/apt/lists/* \ + && groupadd --gid 65532 nonroot && useradd --uid 65532 --gid nonroot --create-home nonroot +WORKDIR /app +COPY --from=builder --chown=nonroot:nonroot /opt/venv /opt/venv +COPY --chown=nonroot:nonroot sdk/python/examples ./examples +USER nonroot:nonroot +ENV PATH=/opt/venv/bin:$PATH +ENTRYPOINT ["python"] +CMD ["examples/serve.py"] diff --git a/sdk/README.md b/sdk/README.md index b3d1cf7f..98791462 100644 --- a/sdk/README.md +++ b/sdk/README.md @@ -436,6 +436,22 @@ holds against the control plane's records. that verifies only under the new one, and still authenticating with the `sam-node`. +### On the testnets + +Both testnets run the example programs, unchanged, as canaries beside the +`sam-node` ones (`.github/k8s/sam-sdk-canary-template.yaml`): the JavaScript +and the Python `serve` publish `mcp://greeter-js` / `a2a://greeter-js` and +`…-py`, enrolled with the pod's projected service account token through +`SAM_JWT_PATH`. Two CronJobs cross every implementation boundary every 15 +minutes and once per rollout: the `sam-node` cold-path probe now also calls +both greeters (node → SDK, MCP through `call_remote_tool` and A2A through +the egress proxy), and the SDK cold-path probe +(`sam-sdk-probe-cronjob-template.yaml`) runs each SDK's `call` against the +everything canary (SDK → node) and the other SDK's greeter (SDK → SDK). The +images (`Dockerfile.sam-sdk-js`, `Dockerfile.sam-sdk-python`) are built per +commit by `deploy.yaml`, so `bananas` runs the SDKs at the same commit as +the Go components they talk to. + ### Later - The connector interface for platforms (issue #480): `Attach`, `Detach`, diff --git a/sdk/js/README.md b/sdk/js/README.md index 4a526932..857f1197 100644 --- a/sdk/js/README.md +++ b/sdk/js/README.md @@ -38,27 +38,45 @@ Find a service and call it: // node call.js a2a://greeter /card // node call.js inference://ollama /v1/models // -// SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -// holding the token the mesh operator gave you; the first run spends it and -// keeps the identity and credential in SAM_STATE_DIR, later runs resume from -// there without it. +// SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +// SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or +// SAM_JWT_PATH (a workload identity token your platform issues, such as a +// Kubernetes projected service account token), and keeps the identity and +// credential in SAM_STATE_DIR; later runs resume from there without it. import { homedir } from "node:os"; -import { AgentMesh } from "@sam-mesh/sdk"; +import { AgentMesh, type DiscoveredProvider } from "@sam-mesh/sdk"; const [service = "mcp://greeter", toolOrPath = "greet", args = '{"name": "world"}'] = process.argv.slice(2); const mesh = await AgentMesh.enroll({ controlPlaneUrl: process.env.SAM_CONTROL_PLANE_URL ?? "https://mesh.example.com", bootstrapTokenPath: process.env.SAM_BOOTSTRAP_TOKEN_PATH, + jwtPath: process.env.SAM_JWT_PATH, stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/caller`, + // A plaintext http:// control plane is otherwise accepted only on loopback. + allowInsecure: process.env.SAM_INSECURE_CONTROL_PLANE === "true", }); const session = await mesh.join(); console.log(`on the mesh as ${session.peerId}`); -const [provider] = await session.discover(service); -if (provider === undefined) { +const providers = await session.discover(service); +if (providers.length === 0) { throw new Error(`no member of the mesh serves ${service}`); } +// A provider record can outlive its member; the first that answers is used. +let provider: DiscoveredProvider | undefined; +for (const candidate of providers) { + try { + await session.connect(candidate); + provider = candidate; + break; + } catch (err) { + console.error(`${candidate.peerId}: ${(err as Error).message}`); + } +} +if (provider === undefined) { + throw new Error(`no provider of ${service} is reachable`); +} console.log(`${service} is served by ${provider.peerId}`); if (service.startsWith("mcp://")) { @@ -85,31 +103,38 @@ Publish services of your own: // members may call; the SDK turns the others away before anything reaches // this code or Ollama. // -// node serve.js +// node serve.js # publishes mcp://greeter and a2a://greeter +// node serve.js greeter-2 # the same under another name // -// SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -// holding the token the mesh operator gave you; the first run spends it and -// keeps the identity and credential in SAM_STATE_DIR, later runs resume from -// there without it. +// SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +// SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or +// SAM_JWT_PATH (a workload identity token your platform issues, such as a +// Kubernetes projected service account token), and keeps the identity and +// credential in SAM_STATE_DIR; later runs resume from there without it. import { homedir } from "node:os"; import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { AgentMesh } from "@sam-mesh/sdk"; import { z } from "zod"; +const [name = "greeter"] = process.argv.slice(2); + const mesh = await AgentMesh.enroll({ controlPlaneUrl: process.env.SAM_CONTROL_PLANE_URL ?? "https://mesh.example.com", bootstrapTokenPath: process.env.SAM_BOOTSTRAP_TOKEN_PATH, - stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/greeter`, + jwtPath: process.env.SAM_JWT_PATH, + stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/${name}`, + // A plaintext http:// control plane is otherwise accepted only on loopback. + allowInsecure: process.env.SAM_INSECURE_CONTROL_PLANE === "true", }); const session = await mesh.join(); await session.serve({ type: "mcp", - name: "greeter", + name, createServer: () => { - const server = new McpServer({ name: "greeter", version: "1.0.0" }); - server.registerTool("greet", { description: "Greets someone by name", inputSchema: { name: z.string() } }, async ({ name }) => ({ - content: [{ type: "text", text: `hello ${name}` }], + const server = new McpServer({ name, version: "1.0.0" }); + server.registerTool("greet", { description: "Greets someone by name", inputSchema: { name: z.string() } }, async ({ name: who }) => ({ + content: [{ type: "text", text: `hello ${who}` }], })); return server; }, @@ -117,8 +142,8 @@ await session.serve({ await session.serve({ type: "a2a", - name: "greeter", - target: (request, caller) => Response.json({ name: "greeter", path: new URL(request.url).pathname, caller: caller.peerId }), + name, + target: (request, caller) => Response.json({ name, path: new URL(request.url).pathname, caller: caller.peerId }), }); if (process.env.OLLAMA_URL !== undefined) { diff --git a/sdk/js/examples/call.ts b/sdk/js/examples/call.ts index ddccacb6..392ccb32 100644 --- a/sdk/js/examples/call.ts +++ b/sdk/js/examples/call.ts @@ -5,27 +5,45 @@ // node call.js a2a://greeter /card // node call.js inference://ollama /v1/models // -// SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -// holding the token the mesh operator gave you; the first run spends it and -// keeps the identity and credential in SAM_STATE_DIR, later runs resume from -// there without it. +// SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +// SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or +// SAM_JWT_PATH (a workload identity token your platform issues, such as a +// Kubernetes projected service account token), and keeps the identity and +// credential in SAM_STATE_DIR; later runs resume from there without it. import { homedir } from "node:os"; -import { AgentMesh } from "@sam-mesh/sdk"; +import { AgentMesh, type DiscoveredProvider } from "@sam-mesh/sdk"; const [service = "mcp://greeter", toolOrPath = "greet", args = '{"name": "world"}'] = process.argv.slice(2); const mesh = await AgentMesh.enroll({ controlPlaneUrl: process.env.SAM_CONTROL_PLANE_URL ?? "https://mesh.example.com", bootstrapTokenPath: process.env.SAM_BOOTSTRAP_TOKEN_PATH, + jwtPath: process.env.SAM_JWT_PATH, stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/caller`, + // A plaintext http:// control plane is otherwise accepted only on loopback. + allowInsecure: process.env.SAM_INSECURE_CONTROL_PLANE === "true", }); const session = await mesh.join(); console.log(`on the mesh as ${session.peerId}`); -const [provider] = await session.discover(service); -if (provider === undefined) { +const providers = await session.discover(service); +if (providers.length === 0) { throw new Error(`no member of the mesh serves ${service}`); } +// A provider record can outlive its member; the first that answers is used. +let provider: DiscoveredProvider | undefined; +for (const candidate of providers) { + try { + await session.connect(candidate); + provider = candidate; + break; + } catch (err) { + console.error(`${candidate.peerId}: ${(err as Error).message}`); + } +} +if (provider === undefined) { + throw new Error(`no provider of ${service} is reachable`); +} console.log(`${service} is served by ${provider.peerId}`); if (service.startsWith("mcp://")) { diff --git a/sdk/js/examples/serve.ts b/sdk/js/examples/serve.ts index a3c3fca0..18a75a69 100644 --- a/sdk/js/examples/serve.ts +++ b/sdk/js/examples/serve.ts @@ -4,31 +4,38 @@ // members may call; the SDK turns the others away before anything reaches // this code or Ollama. // -// node serve.js +// node serve.js # publishes mcp://greeter and a2a://greeter +// node serve.js greeter-2 # the same under another name // -// SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -// holding the token the mesh operator gave you; the first run spends it and -// keeps the identity and credential in SAM_STATE_DIR, later runs resume from -// there without it. +// SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +// SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or +// SAM_JWT_PATH (a workload identity token your platform issues, such as a +// Kubernetes projected service account token), and keeps the identity and +// credential in SAM_STATE_DIR; later runs resume from there without it. import { homedir } from "node:os"; import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { AgentMesh } from "@sam-mesh/sdk"; import { z } from "zod"; +const [name = "greeter"] = process.argv.slice(2); + const mesh = await AgentMesh.enroll({ controlPlaneUrl: process.env.SAM_CONTROL_PLANE_URL ?? "https://mesh.example.com", bootstrapTokenPath: process.env.SAM_BOOTSTRAP_TOKEN_PATH, - stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/greeter`, + jwtPath: process.env.SAM_JWT_PATH, + stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/${name}`, + // A plaintext http:// control plane is otherwise accepted only on loopback. + allowInsecure: process.env.SAM_INSECURE_CONTROL_PLANE === "true", }); const session = await mesh.join(); await session.serve({ type: "mcp", - name: "greeter", + name, createServer: () => { - const server = new McpServer({ name: "greeter", version: "1.0.0" }); - server.registerTool("greet", { description: "Greets someone by name", inputSchema: { name: z.string() } }, async ({ name }) => ({ - content: [{ type: "text", text: `hello ${name}` }], + const server = new McpServer({ name, version: "1.0.0" }); + server.registerTool("greet", { description: "Greets someone by name", inputSchema: { name: z.string() } }, async ({ name: who }) => ({ + content: [{ type: "text", text: `hello ${who}` }], })); return server; }, @@ -36,8 +43,8 @@ await session.serve({ await session.serve({ type: "a2a", - name: "greeter", - target: (request, caller) => Response.json({ name: "greeter", path: new URL(request.url).pathname, caller: caller.peerId }), + name, + target: (request, caller) => Response.json({ name, path: new URL(request.url).pathname, caller: caller.peerId }), }); if (process.env.OLLAMA_URL !== undefined) { diff --git a/sdk/js/src/identity.test.ts b/sdk/js/src/identity.test.ts index dbd6ed4e..52150d2a 100644 --- a/sdk/js/src/identity.test.ts +++ b/sdk/js/src/identity.test.ts @@ -17,7 +17,7 @@ import { readFileSync } from "node:fs"; import { test } from "node:test"; import { decodeBase58, encodeBase58 } from "./base58.ts"; import { enrollChallenge } from "./challenges.ts"; -import { Identity, libp2pPublicKey, peerIdFromPublicKey, verifyEd25519 } from "./identity.ts"; +import { Identity, canonicalPeerId, libp2pPublicKey, peerIdFromPublicKey, verifyEd25519 } from "./identity.ts"; interface Vector { seed: string; @@ -74,6 +74,16 @@ test("libp2pPublicKey refuses the wrong size", () => { assert.throws(() => libp2pPublicKey(new Uint8Array(33)), /32 bytes/); }); +test("canonicalPeerId is the base58 form for every encoding libp2p accepts", () => { + // The CIDv1 form of a known ed25519 peer, as `peer.ToCid(id).String()` prints it. + const base58 = "12D3KooWA4Xop1JaT3MHxwYMkCepYsv4iPVopMXwCz5iHYdBfeSB"; + const cidv1 = "bafzaajaiaejcaa5ba677htqqxyoxbxiy45f4bglh4tldbg5fbvpr3xegmqjfkmny"; + assert.equal(canonicalPeerId(base58), base58); + assert.equal(canonicalPeerId(cidv1), base58); + assert.throws(() => canonicalPeerId("not-a-peer"), /is not a peer ID/); + assert.throws(() => canonicalPeerId(""), /is not a peer ID/); +}); + test("base58btc round-trips and keeps leading zeros", () => { const cases: Uint8Array[] = [new Uint8Array(0), Uint8Array.of(0), Uint8Array.of(0, 0, 1, 2), hex("00ff"), hex("deadbeef")]; for (const c of cases) { diff --git a/sdk/js/src/identity.ts b/sdk/js/src/identity.ts index b7ea122f..1744d8a2 100644 --- a/sdk/js/src/identity.ts +++ b/sdk/js/src/identity.ts @@ -16,6 +16,7 @@ // derives, so the same key works in the SDK, in sam-node and on the wire. import { createPrivateKey, createPublicKey, sign, verify, type KeyObject } from "node:crypto"; +import { peerIdFromString } from "@libp2p/peer-id"; import { encodeBase58 } from "./base58.ts"; const PUBLIC_KEY_SIZE = 32; @@ -66,6 +67,21 @@ export function peerIdFromPublicKey(publicKeyRaw: Uint8Array): string { return encodeBase58(concat(IDENTITY_MULTIHASH_PREFIX, libp2pPublicKey(publicKeyRaw))); } +/** + * The base58btc form of a peer ID written in any encoding libp2p accepts + * (base58btc multihash, CIDv1). Every key, ban set and comparison in SAM is + * on this form, as peer.ID.String() in Go; a string read off the wire or + * from a caller goes through here before it is used as one. Throws when the + * text is not a peer ID at all. + */ +export function canonicalPeerId(text: string): string { + try { + return peerIdFromString(text).toString(); + } catch (err) { + throw new Error(`${JSON.stringify(text)} is not a peer ID: ${err instanceof Error ? err.message : String(err)}`); + } +} + export class Identity { readonly #privateKey: KeyObject; readonly #publicKey: KeyObject; diff --git a/sdk/js/src/index.ts b/sdk/js/src/index.ts index a5fc1441..2713def9 100644 --- a/sdk/js/src/index.ts +++ b/sdk/js/src/index.ts @@ -13,7 +13,7 @@ // limitations under the License. export { AgentMesh, type AgentMeshOptions, type ControlPlaneSync, type EnrollOptions } from "./mesh.ts"; -export { Identity, peerIdFromPublicKey, libp2pPublicKey, verifyEd25519 } from "./identity.ts"; +export { Identity, canonicalPeerId, peerIdFromPublicKey, libp2pPublicKey, verifyEd25519 } from "./identity.ts"; export { ControlPlaneClient, ControlPlaneError, diff --git a/sdk/js/src/mesh.test.ts b/sdk/js/src/mesh.test.ts index 500234bd..53206dd6 100644 --- a/sdk/js/src/mesh.test.ts +++ b/sdk/js/src/mesh.test.ts @@ -24,6 +24,8 @@ import { AuthFrameSchema, AuthResponseSchema, BootstrapEnrollResponseSchema, + EnrollRequestSchema, + EnrollResponseSchema, EnrollmentStatus, KeysResponseSchema, TokenRefreshResponseSchema, @@ -66,6 +68,22 @@ function fakeControlPlane(keysOk = true): { fetch: typeof fetch; issued: number case "POST /refresh": state.issued++; return proto(toBinary(TokenRefreshResponseSchema, create(TokenRefreshResponseSchema, { biscuitToken: text(`biscuit-${state.issued}`), expireTime: timestampFromMs(Date.now() + 7200_000) }))); + case "POST /register": { + // The biscuit names the JWT that was presented, so a test can see which. + const jwt = fromBinary(EnrollRequestSchema, new Uint8Array(await req.arrayBuffer())).jwt; + state.issued++; + return proto( + toBinary( + EnrollResponseSchema, + create(EnrollResponseSchema, { + biscuitToken: text(`biscuit-for-${jwt}`), + controlPlanePublicKey: cpKey.publicKeyRaw, + routerAddresses: ["/dns4/router.example/tcp/4001/p2p/12D3KooWP8iKhDf3iCMo2H3butNVfdTUtYwYWYQ75jTGnynXPFMp"], + expireTime: timestampFromMs(Date.now() + 3600_000), + }), + ), + ); + } case "GET /keys": return keysOk ? proto(toBinary(KeysResponseSchema, signedKeys())) : new Response("boom", { status: 500 }); default: @@ -161,9 +179,24 @@ test("enroll refuses ambiguous credentials", async () => { const cp = fakeControlPlane(); await assert.rejects(AgentMesh.enroll({ controlPlaneUrl: "http://127.0.0.1:1", fetch: cp.fetch }), /exactly one of/); await assert.rejects(AgentMesh.enroll({ controlPlaneUrl: "http://127.0.0.1:1", bootstrapToken: "a", jwt: "b", fetch: cp.fetch }), /exactly one of/); + await assert.rejects(AgentMesh.enroll({ controlPlaneUrl: "http://127.0.0.1:1", jwt: "a", jwtPath: "/nonexistent", fetch: cp.fetch }), /exactly one of/); assert.equal(cp.issued, 0); }); +test("enroll reads a workload identity token from jwtPath", async () => { + const dir = await mkdtemp(join(tmpdir(), "sam-sdk-")); + try { + const cp = fakeControlPlane(); + const jwtPath = join(dir, "token"); + await writeFile(jwtPath, "eyJ.projected.token\n"); + const mesh = await AgentMesh.enroll({ controlPlaneUrl: "http://127.0.0.1:1", jwtPath, fetch: cp.fetch }); + assert.deepEqual(mesh.credential.biscuit, text("biscuit-for-eyJ.projected.token")); + await assert.rejects(AgentMesh.enroll({ controlPlaneUrl: "http://127.0.0.1:1", jwtPath: join(dir, "missing"), fetch: cp.fetch }), /ENOENT/); + } finally { + await rm(dir, { recursive: true, force: true }); + } +}); + test("authFrame is the AuthFrame protobuf with this member's biscuit", async () => { const cp = fakeControlPlane(); const mesh = await AgentMesh.enroll({ controlPlaneUrl: "http://127.0.0.1:1", bootstrapToken: "sbt_secret", fetch: cp.fetch }); diff --git a/sdk/js/src/mesh.ts b/sdk/js/src/mesh.ts index 85cbafab..d5192207 100644 --- a/sdk/js/src/mesh.ts +++ b/sdk/js/src/mesh.ts @@ -51,6 +51,12 @@ export interface EnrollOptions extends AgentMeshOptions { bootstrapTokenPath?: string | undefined; /** An OIDC ID token, for meshes that enroll identities interactively. */ jwt?: string | undefined; + /** + * Path of a file holding an OIDC ID token or a platform's workload identity + * token, such as a Kubernetes projected service account token. Preferred + * over a value; the file is read at enrollment. + */ + jwtPath?: string | undefined; /** Bounds the wait for an operator to approve a pending enrollment. */ signal?: AbortSignal; /** Overrides the control plane's suggested poll interval while pending. */ @@ -104,8 +110,8 @@ export class AgentMesh { * plane for the saved identity, that member is returned and no token is * needed, so a program can call enroll on every start and read the token * from its environment only on the first. Otherwise exactly one of - * bootstrapToken, bootstrapTokenPath or jwt must be given. Delete the - * state directory to enroll afresh, for instance with other labels. + * bootstrapToken, bootstrapTokenPath, jwt or jwtPath must be given. Delete + * the state directory to enroll afresh, for instance with other labels. */ static async enroll(options: EnrollOptions): Promise { const saved = await loadIdentity(options.stateDir); @@ -117,16 +123,17 @@ export class AgentMesh { return new AgentMesh(identity, controlPlane, credential, options.stateDir); } } - const given = [options.bootstrapToken, options.bootstrapTokenPath, options.jwt].filter((v) => v !== undefined).length; + const given = [options.bootstrapToken, options.bootstrapTokenPath, options.jwt, options.jwtPath].filter((v) => v !== undefined).length; if (given !== 1) { const where = options.stateDir !== undefined ? ` (no credential to resume in ${options.stateDir})` : ""; - throw new Error(`exactly one of bootstrapToken, bootstrapTokenPath or jwt is required${where}`); + throw new Error(`exactly one of bootstrapToken, bootstrapTokenPath, jwt or jwtPath is required${where}`); } const role = options.role ?? ROLE_NODE; let enrollment: Enrollment; - if (options.jwt !== undefined) { - enrollment = await controlPlane.register({ identity, jwt: options.jwt, role, ...labelsOf(options) }); + if (options.jwt !== undefined || options.jwtPath !== undefined) { + const jwt = options.jwtPath !== undefined ? (await readFile(options.jwtPath, "utf8")).trim() : (options.jwt as string); + enrollment = await controlPlane.register({ identity, jwt, role, ...labelsOf(options) }); } else { const bootstrapToken = options.bootstrapTokenPath !== undefined ? (await readFile(options.bootstrapTokenPath, "utf8")).trim() : (options.bootstrapToken as string); enrollment = await controlPlane.enrollBootstrap({ diff --git a/sdk/js/src/session.ts b/sdk/js/src/session.ts index 0d93a4e5..a24bcfb9 100644 --- a/sdk/js/src/session.ts +++ b/sdk/js/src/session.ts @@ -19,6 +19,7 @@ import { peerIdFromString } from "@libp2p/peer-id"; import { isMultiaddr, multiaddr, type Multiaddr } from "@multiformats/multiaddr"; import { AUTH_HANDLER_OPTIONS, AUTH_PROTOCOL, MCP_PROTOCOL, authenticateWithPeer, authStreamHandler } from "./auth.ts"; import { ROLE_ROUTER, requireRole, type VerifiedBiscuit } from "./biscuit.ts"; +import { canonicalPeerId } from "./identity.ts"; import { isServiceType, parseServiceTarget, serviceCID, type ServiceType } from "./discovery.ts"; import { createMeshHost, listenThroughRelay, type MeshHost, type MeshHostOptions } from "./host.ts"; import { openMCPSession, type MCPSession, type MCPSessionOptions } from "./mcp.ts"; @@ -189,17 +190,19 @@ export class MeshSession { /** The addresses connect() dials for a peer, in the order libp2p tries them. */ dialTargets(peer: Peer): { peerId: string | undefined; addrs: Multiaddr[] } { if (typeof peer === "string" && !peer.startsWith("/")) { - return { peerId: peer, addrs: this.relayedAddresses(peer) }; + const peerId = canonicalPeerId(peer); + return { peerId, addrs: this.relayedAddresses(peerId) }; } if (typeof peer === "string" || isMultiaddr(peer)) { const ma = typeof peer === "string" ? multiaddr(peer) : peer; return { peerId: targetPeerOf(ma), addrs: [ma] }; } + const peerId = canonicalPeerId(peer.peerId); const advertised = peer.addrs.map((text) => { const ma = multiaddr(text); - return targetPeerOf(ma) === undefined ? ma.encapsulate(`/p2p/${peer.peerId}`) : ma; + return targetPeerOf(ma) === undefined ? ma.encapsulate(`/p2p/${peerId}`) : ma; }); - return { peerId: peer.peerId, addrs: [...advertised, ...this.relayedAddresses(peer.peerId)] }; + return { peerId, addrs: [...advertised, ...this.relayedAddresses(peerId)] }; } /** `/p2p-circuit/p2p/` through every router that admitted this member. */ @@ -312,7 +315,7 @@ export class MeshSession { async #syncOnce(): Promise { const result = await this.mesh.syncControlPlane(); if (result.bannedPeerIds !== undefined) { - const { banned } = this.banned.reconcile(result.bannedPeerIds, result.fetchedAt); + const { banned } = this.banned.reconcile(canonicalPeerIds(result.bannedPeerIds), result.fetchedAt); await Promise.all(banned.map((peerId) => this.#evict(peerId))); } if (this.#serving) { @@ -574,8 +577,21 @@ export async function joinMesh(mesh: AgentMesh, options: JoinOptions = {}): Prom return new MeshSession(mesh, node, admitted, authenticatedPeers, banned, options); } -/** The peer a multiaddr ends at: its trailing `/p2p/`, or undefined for a relay address with no target yet. */ +/** The peer a multiaddr ends at, in canonical form: its trailing `/p2p/`, or undefined for a relay address with no target yet. */ function targetPeerOf(ma: Multiaddr): string | undefined { const last = ma.getComponents().at(-1); - return last?.name === "p2p" ? last.value : undefined; + return last?.name === "p2p" && last.value !== undefined ? canonicalPeerId(last.value) : undefined; +} + +/** Canonicalizes a list from the control plane, dropping entries that are not peer IDs. */ +function canonicalPeerIds(ids: string[]): string[] { + const out: string[] = []; + for (const id of ids) { + try { + out.push(canonicalPeerId(id)); + } catch { + // Not a peer ID; it can match nothing, so it bans nothing. + } + } + return out; } diff --git a/sdk/js/src/sync.test.ts b/sdk/js/src/sync.test.ts index 9f9ae5fe..8c43e78a 100644 --- a/sdk/js/src/sync.test.ts +++ b/sdk/js/src/sync.test.ts @@ -25,6 +25,7 @@ import { circuitRelayServer, circuitRelayTransport } from "@libp2p/circuit-relay import { privateKeyFromProtobuf } from "@libp2p/crypto/keys"; import { identify } from "@libp2p/identify"; import type { Libp2p } from "@libp2p/interface"; +import { peerIdFromString } from "@libp2p/peer-id"; import { tcp } from "@libp2p/tcp"; import { tls } from "@libp2p/tls"; import { createLibp2p } from "libp2p"; @@ -190,15 +191,25 @@ test("a mesh event verifies only under a trusted key and only when fresh", () => return toBinary(MeshEventSchema, { ...event, signature: key.identity.sign(unsigned) }); }; const now = new Date(); - const banned = sign(create(MeshEventSchema, { type: MeshEvent_Type.BANNED, peerId: "12D3KooWx", eventTime: timestampFromMs(now.getTime()) })); - assert.equal(verifyMeshEvent(banned, [key.pub], now)?.peerId, "12D3KooWx"); + const target = Identity.generate().peerId; + const banned = sign(create(MeshEventSchema, { type: MeshEvent_Type.BANNED, peerId: target, eventTime: timestampFromMs(now.getTime()) })); + assert.equal(verifyMeshEvent(banned, [key.pub], now)?.peerId, target); assert.equal(verifyMeshEvent(banned, [new SigningKey().pub], now), undefined, "untrusted key"); const tampered = new Uint8Array(banned); tampered[tampered.length - 1] = (tampered[tampered.length - 1] ?? 0) ^ 1; assert.equal(verifyMeshEvent(tampered, [key.pub], now), undefined, "tampered signature"); - const stale = sign(create(MeshEventSchema, { type: MeshEvent_Type.BANNED, peerId: "12D3KooWx", eventTime: timestampFromMs(now.getTime() - EVENT_FRESHNESS_MS - 1) })); + const stale = sign(create(MeshEventSchema, { type: MeshEvent_Type.BANNED, peerId: target, eventTime: timestampFromMs(now.getTime() - EVENT_FRESHNESS_MS - 1) })); assert.equal(verifyMeshEvent(stale, [key.pub], now), undefined, "stale"); assert.equal(verifyMeshEvent(new Uint8Array([1, 2, 3]), [key.pub], now), undefined, "garbage"); + + // The ban set is keyed on the base58 form; an event naming the peer in + // its CIDv1 form bans the same peer, and one naming no peer bans nobody. + const cidForm = peerIdFromString(target).toCID().toString(); + assert.notEqual(cidForm, target); + const bannedByCID = sign(create(MeshEventSchema, { type: MeshEvent_Type.BANNED, peerId: cidForm, eventTime: timestampFromMs(now.getTime()) })); + assert.equal(verifyMeshEvent(bannedByCID, [key.pub], now)?.peerId, target); + const bannedNobody = sign(create(MeshEventSchema, { type: MeshEvent_Type.BANNED, peerId: "not-a-peer", eventTime: timestampFromMs(now.getTime()) })); + assert.equal(verifyMeshEvent(bannedNobody, [key.pub], now), undefined, "not a peer id"); }); test("a pull learns a key rotation and refreshes the credential under the new key", async () => { @@ -273,13 +284,23 @@ test("a banned peer is hung up on and refused at the gate and the handshake", as // Its token still verifies, and it is still refused: a new connection at the gate ... await assert.rejects(peer.dial(relayed).then((c) => authenticateWithPeer(c, frame, [cp.current.pub]))); - // ... and outbound, before any dial. + // ... and outbound, before any dial, however the peer is named. + const cidForm = peerIdFromString(peerIdentity.peerId).toCID().toString(); await assert.rejects(session.connect(`${routerAddr}/p2p-circuit/p2p/${peerIdentity.peerId}`), /banned/); + await assert.rejects(session.connect(`${routerAddr}/p2p-circuit/p2p/${cidForm}`), /banned/); + await assert.rejects(session.connect(cidForm), /banned/); + await assert.rejects(session.connect({ peerId: cidForm, addrs: [] }), /banned/); - // Lifted by the control plane: the next pull unbans it. + // Lifted by the control plane: the next pull unbans it. A ban list that + // names the peer in its CIDv1 form bans the same peer. cp.banned = []; await session.sync(); assert.deepEqual(session.banned.peers(), []); + cp.banned = [cidForm, "not-a-peer"]; + await session.sync(); + assert.deepEqual(session.banned.peers(), [peerIdentity.peerId]); + cp.banned = []; + await session.sync(); } finally { await peer.stop(); await session.close(); diff --git a/sdk/js/src/sync.ts b/sdk/js/src/sync.ts index c720ec68..4a8bdf33 100644 --- a/sdk/js/src/sync.ts +++ b/sdk/js/src/sync.ts @@ -18,7 +18,7 @@ import { fromBinary, toBinary } from "@bufbuild/protobuf"; import { timestampMs, type Timestamp } from "@bufbuild/protobuf/wkt"; -import { verifyEd25519 } from "./identity.ts"; +import { canonicalPeerId, verifyEd25519 } from "./identity.ts"; import { MeshEvent_Type, MeshEventSchema, type MeshEvent } from "./gen/sam_pb.ts"; /** The GossipSub topic the control plane publishes mesh events on (api.GossipEvents). */ @@ -83,13 +83,14 @@ export class BanSet { } /** A MeshEvent that verified: signed by a trusted key and carrying a fresh event_time. */ +/** A MeshEvent that verified: signed by a trusted key, carrying a fresh event_time, its peer id in canonical form. */ export type VerifiedMeshEvent = MeshEvent & { eventTime: Timestamp }; /** * Verifies a MeshEvent as sam-node's verifyEvent does: the signature covers * the deterministic encoding of the event with the signature cleared, under * any trusted control plane key. Returns the event, or undefined when it - * does not verify or is not fresh. + * does not verify, is not fresh, or bans something that is not a peer ID. */ export function verifyMeshEvent(data: Uint8Array, trustedKeys: Uint8Array[], now: Date = new Date()): VerifiedMeshEvent | undefined { let event: MeshEvent; @@ -110,6 +111,14 @@ export function verifyMeshEvent(data: Uint8Array, trustedKeys: Uint8Array[], now if (skew > EVENT_FRESHNESS_MS) { return undefined; } + if (event.type === MeshEvent_Type.BANNED) { + // Canonicalized after the signature check, which covers the bytes as sent. + try { + return { ...event, peerId: canonicalPeerId(event.peerId) } as VerifiedMeshEvent; + } catch { + return undefined; + } + } return event as VerifiedMeshEvent; } diff --git a/sdk/python/README.md b/sdk/python/README.md index 7a6288c8..6d72c3ac 100644 --- a/sdk/python/README.md +++ b/sdk/python/README.md @@ -41,10 +41,11 @@ path of an inference or A2A service. python call.py a2a://greeter /card python call.py inference://ollama /v1/models -SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -holding the token the mesh operator gave you; the first run spends it and -keeps the identity and credential in SAM_STATE_DIR, later runs resume from -there without it. +SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or SAM_JWT_PATH +(a workload identity token your platform issues, such as a Kubernetes +projected service account token), and keeps the identity and credential in +SAM_STATE_DIR; later runs resume from there without it. """ import json @@ -62,7 +63,10 @@ args = json.loads(argv[2]) if len(argv) > 2 else {"name": "world"} mesh = AgentMesh.enroll( os.environ.get("SAM_CONTROL_PLANE_URL", "https://mesh.example.com"), bootstrap_token_path=os.environ.get("SAM_BOOTSTRAP_TOKEN_PATH"), + jwt_path=os.environ.get("SAM_JWT_PATH"), state_dir=os.environ.get("SAM_STATE_DIR", "~/.config/sam-mesh/caller"), + # A plaintext http:// control plane is otherwise accepted only on loopback. + allow_insecure=os.environ.get("SAM_INSECURE_CONTROL_PLANE") == "true", ) @@ -73,7 +77,15 @@ async def main() -> None: providers = await session.discover(service) if not providers: raise SystemExit(f"no member of the mesh serves {service}") - provider = providers[0] + # A provider record can outlive its member; the first that answers is used. + for provider in providers: + try: + await session.connect(provider) + break + except (ConnectionError, PermissionError) as err: + print(f"{provider.peer_id}: {err}", file=sys.stderr) + else: + raise SystemExit(f"no provider of {service} is reachable") print(f"{service} is served by {provider.peer_id}") if service.startswith("mcp://"): @@ -100,30 +112,38 @@ beside this program as an inference service. The mesh policy decides which members may call; the SDK turns the others away before anything reaches this code or Ollama. - python serve.py + python serve.py # publishes mcp://greeter and a2a://greeter + python serve.py greeter-2 # the same under another name -SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -holding the token the mesh operator gave you; the first run spends it and -keeps the identity and credential in SAM_STATE_DIR, later runs resume from -there without it. +SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or SAM_JWT_PATH +(a workload identity token your platform issues, such as a Kubernetes +projected service account token), and keeps the identity and credential in +SAM_STATE_DIR; later runs resume from there without it. """ import json import os +import sys import trio from agent_mesh import AgentMesh, HTTPRequest, HTTPResponse, HTTPService, MCPService, VerifiedBiscuit from mcp.server.mcpserver import MCPServer +service_name = sys.argv[1] if len(sys.argv) > 1 else "greeter" + mesh = AgentMesh.enroll( os.environ.get("SAM_CONTROL_PLANE_URL", "https://mesh.example.com"), bootstrap_token_path=os.environ.get("SAM_BOOTSTRAP_TOKEN_PATH"), - state_dir=os.environ.get("SAM_STATE_DIR", "~/.config/sam-mesh/greeter"), + jwt_path=os.environ.get("SAM_JWT_PATH"), + state_dir=os.environ.get("SAM_STATE_DIR", f"~/.config/sam-mesh/{service_name}"), + # A plaintext http:// control plane is otherwise accepted only on loopback. + allow_insecure=os.environ.get("SAM_INSECURE_CONTROL_PLANE") == "true", ) def create_server() -> MCPServer: - server = MCPServer("greeter") + server = MCPServer(service_name) @server.tool(description="Greets someone by name") def greet(name: str) -> str: @@ -133,14 +153,14 @@ def create_server() -> MCPServer: async def card(request: HTTPRequest, caller: VerifiedBiscuit) -> HTTPResponse: - body = json.dumps({"name": "greeter", "path": request.path, "caller": caller.peer_id}) + body = json.dumps({"name": service_name, "path": request.path, "caller": caller.peer_id}) return HTTPResponse(status=200, headers={"content-type": "application/json"}, body=body.encode()) async def main() -> None: async with mesh.join() as session: - await session.serve(MCPService(name="greeter", create_server=create_server)) - await session.serve(HTTPService(type="a2a", name="greeter", target=card)) + await session.serve(MCPService(name=service_name, create_server=create_server)) + await session.serve(HTTPService(type="a2a", name=service_name, target=card)) if "OLLAMA_URL" in os.environ: await session.serve(HTTPService(type="inference", name="ollama", target=os.environ["OLLAMA_URL"])) diff --git a/sdk/python/examples/call.py b/sdk/python/examples/call.py index 721838dc..8a1ca068 100644 --- a/sdk/python/examples/call.py +++ b/sdk/python/examples/call.py @@ -5,10 +5,11 @@ python call.py a2a://greeter /card python call.py inference://ollama /v1/models -SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -holding the token the mesh operator gave you; the first run spends it and -keeps the identity and credential in SAM_STATE_DIR, later runs resume from -there without it. +SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or SAM_JWT_PATH +(a workload identity token your platform issues, such as a Kubernetes +projected service account token), and keeps the identity and credential in +SAM_STATE_DIR; later runs resume from there without it. """ import json @@ -26,7 +27,10 @@ mesh = AgentMesh.enroll( os.environ.get("SAM_CONTROL_PLANE_URL", "https://mesh.example.com"), bootstrap_token_path=os.environ.get("SAM_BOOTSTRAP_TOKEN_PATH"), + jwt_path=os.environ.get("SAM_JWT_PATH"), state_dir=os.environ.get("SAM_STATE_DIR", "~/.config/sam-mesh/caller"), + # A plaintext http:// control plane is otherwise accepted only on loopback. + allow_insecure=os.environ.get("SAM_INSECURE_CONTROL_PLANE") == "true", ) @@ -37,7 +41,15 @@ async def main() -> None: providers = await session.discover(service) if not providers: raise SystemExit(f"no member of the mesh serves {service}") - provider = providers[0] + # A provider record can outlive its member; the first that answers is used. + for provider in providers: + try: + await session.connect(provider) + break + except (ConnectionError, PermissionError) as err: + print(f"{provider.peer_id}: {err}", file=sys.stderr) + else: + raise SystemExit(f"no provider of {service} is reachable") print(f"{service} is served by {provider.peer_id}") if service.startswith("mcp://"): diff --git a/sdk/python/examples/serve.py b/sdk/python/examples/serve.py index 7f5e2935..6821732a 100644 --- a/sdk/python/examples/serve.py +++ b/sdk/python/examples/serve.py @@ -4,30 +4,38 @@ members may call; the SDK turns the others away before anything reaches this code or Ollama. - python serve.py - -SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -holding the token the mesh operator gave you; the first run spends it and -keeps the identity and credential in SAM_STATE_DIR, later runs resume from -there without it. + python serve.py # publishes mcp://greeter and a2a://greeter + python serve.py greeter-2 # the same under another name + +SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or SAM_JWT_PATH +(a workload identity token your platform issues, such as a Kubernetes +projected service account token), and keeps the identity and credential in +SAM_STATE_DIR; later runs resume from there without it. """ import json import os +import sys import trio from agent_mesh import AgentMesh, HTTPRequest, HTTPResponse, HTTPService, MCPService, VerifiedBiscuit from mcp.server.mcpserver import MCPServer +service_name = sys.argv[1] if len(sys.argv) > 1 else "greeter" + mesh = AgentMesh.enroll( os.environ.get("SAM_CONTROL_PLANE_URL", "https://mesh.example.com"), bootstrap_token_path=os.environ.get("SAM_BOOTSTRAP_TOKEN_PATH"), - state_dir=os.environ.get("SAM_STATE_DIR", "~/.config/sam-mesh/greeter"), + jwt_path=os.environ.get("SAM_JWT_PATH"), + state_dir=os.environ.get("SAM_STATE_DIR", f"~/.config/sam-mesh/{service_name}"), + # A plaintext http:// control plane is otherwise accepted only on loopback. + allow_insecure=os.environ.get("SAM_INSECURE_CONTROL_PLANE") == "true", ) def create_server() -> MCPServer: - server = MCPServer("greeter") + server = MCPServer(service_name) @server.tool(description="Greets someone by name") def greet(name: str) -> str: @@ -37,14 +45,14 @@ def greet(name: str) -> str: async def card(request: HTTPRequest, caller: VerifiedBiscuit) -> HTTPResponse: - body = json.dumps({"name": "greeter", "path": request.path, "caller": caller.peer_id}) + body = json.dumps({"name": service_name, "path": request.path, "caller": caller.peer_id}) return HTTPResponse(status=200, headers={"content-type": "application/json"}, body=body.encode()) async def main() -> None: async with mesh.join() as session: - await session.serve(MCPService(name="greeter", create_server=create_server)) - await session.serve(HTTPService(type="a2a", name="greeter", target=card)) + await session.serve(MCPService(name=service_name, create_server=create_server)) + await session.serve(HTTPService(type="a2a", name=service_name, target=card)) if "OLLAMA_URL" in os.environ: await session.serve(HTTPService(type="inference", name="ollama", target=os.environ["OLLAMA_URL"])) diff --git a/sdk/python/src/agent_mesh/__init__.py b/sdk/python/src/agent_mesh/__init__.py index 518a2683..3fdf0c1e 100644 --- a/sdk/python/src/agent_mesh/__init__.py +++ b/sdk/python/src/agent_mesh/__init__.py @@ -31,7 +31,7 @@ ) from .credential import MeshCredential, decode_auth_response, encode_auth_frame from .discovery import DHT_PROTOCOL, DiscoveredProvider, find_providers, parse_service_target, provide, service_key -from .identity import Identity, libp2p_public_key, peer_id_from_public_key, verify_ed25519 +from .identity import Identity, canonical_peer_id, libp2p_public_key, peer_id_from_public_key, verify_ed25519 from .mcp_client import LabelsNotSatisfiedError, ToolCallResult, ToolInfo, open_mcp_session, require_labels from .mesh import AgentMesh, ControlPlaneSync from .relay import dial_through_relay, reserve_relay @@ -96,6 +96,7 @@ "auth_stream_handler", "authenticate_with_peer", "authorize_caller", + "canonical_peer_id", "decode_auth_response", "dial_through_relay", "encode_auth_frame", diff --git a/sdk/python/src/agent_mesh/host.py b/sdk/python/src/agent_mesh/host.py index d03d5ec3..69cdd3aa 100644 --- a/sdk/python/src/agent_mesh/host.py +++ b/sdk/python/src/agent_mesh/host.py @@ -23,16 +23,19 @@ from typing import Sequence import multiaddr +import trio from cryptography import x509 from cryptography.x509.oid import NameOID from libp2p import new_host from libp2p.abc import IHost from libp2p.crypto.ed25519 import create_new_key_pair from libp2p.custom_types import TProtocol +from libp2p.peer.peerinfo import PeerInfo, info_from_p2p_addr from libp2p.security.tls.transport import PROTOCOL_ID as TLS_PROTOCOL_ID from libp2p.security.tls.transport import IdentityConfig, TLSTransport from libp2p.stream_muxer.yamux.yamux import PROTOCOL_ID as YAMUX_PROTOCOL_ID from libp2p.stream_muxer.yamux.yamux import Yamux +from multiaddr.resolvers import DNSResolver from .identity import Identity @@ -71,3 +74,41 @@ def create_mesh_host(identity: Identity, listen_addrs: Sequence[str] = ()) -> tu if str(host.get_id()) != identity.peer_id: raise RuntimeError(f"libp2p derived peer {host.get_id()} for identity {identity.peer_id}") return host, [multiaddr.Multiaddr(a) for a in listen_addrs] + + +_DNS_PROTOCOLS = frozenset({"dnsaddr", "dns", "dns4", "dns6"}) +# py-libp2p dials TCP only; a resolved address on another transport is noise. +_UNDIALABLE_PROTOCOLS = frozenset({"quic", "quic-v1", "ws", "wss", "webtransport", "webrtc", "webrtc-direct"}) +_DNS_TIMEOUT = 10.0 + + +async def dial_addrs(addr: multiaddr.Multiaddr) -> list[multiaddr.Multiaddr]: + """The concrete addresses this host can dial for addr. A control plane + hands out router addresses as `/dnsaddr//p2p/`, resolved here + through the host's `_dnsaddr` TXT records (and `/dns4`, `/dns6`, `/dns` + through A and AAAA records) the way go-libp2p and js-libp2p do before + dialing; py-libp2p does not, and would report no transport for them. + Addresses on transports this host lacks are left out.""" + protocols = [p.name for p in addr.protocols()] + if protocols and protocols[0] in _DNS_PROTOCOLS: + resolved: list[multiaddr.Multiaddr] = [] + with trio.move_on_after(_DNS_TIMEOUT): + resolved = list(await DNSResolver().resolve(addr)) + if not resolved: + raise RuntimeError(f"{addr} resolved to no address") + else: + resolved = [addr] + dialable: list[multiaddr.Multiaddr] = [] + for m in resolved: + names = {p.name for p in m.protocols()} + if "tcp" in names and not names & _UNDIALABLE_PROTOCOLS: + dialable.append(m) + if not dialable: + raise RuntimeError(f"{addr} offers no TCP address; this host dials TCP only") + return dialable + + +async def peer_info(addr: multiaddr.Multiaddr) -> PeerInfo: + """The peer an address names and the concrete addresses to reach it on.""" + addrs = await dial_addrs(addr) + return PeerInfo(info_from_p2p_addr(addrs[0]).peer_id, addrs) diff --git a/sdk/python/src/agent_mesh/identity.py b/sdk/python/src/agent_mesh/identity.py index aec1cb90..c85fdfe0 100644 --- a/sdk/python/src/agent_mesh/identity.py +++ b/sdk/python/src/agent_mesh/identity.py @@ -19,6 +19,7 @@ import os +import multiaddr from cryptography.exceptions import InvalidSignature from cryptography.hazmat.primitives import serialization from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey, Ed25519PublicKey @@ -48,6 +49,19 @@ def peer_id_from_public_key(public_key_raw: bytes) -> str: return base58.encode(_IDENTITY_MULTIHASH_PREFIX + libp2p_public_key(public_key_raw)) +def canonical_peer_id(text: str) -> str: + """The base58btc form of a peer ID written in any encoding libp2p accepts + (base58btc multihash, CIDv1). Every key, ban set and comparison in SAM is + on this form, as peer.ID.String() in Go; a string read off the wire or from + a caller goes through here before it is used as one. Raises ValueError when + the text is not a peer ID at all.""" + # py-multiaddr's p2p codec decodes both encodings and prints base58btc. + try: + return multiaddr.Multiaddr("/p2p/" + text).value_for_protocol("p2p") + except Exception as err: # noqa: BLE001 - the library raises its own parse error types + raise ValueError(f"{text!r} is not a peer ID: {err}") from err + + def verify_ed25519(public_key_raw: bytes, data: bytes, signature: bytes) -> bool: """Verifies an ed25519 signature with a raw 32-byte public key.""" if len(public_key_raw) != PUBLIC_KEY_SIZE: diff --git a/sdk/python/src/agent_mesh/mesh.py b/sdk/python/src/agent_mesh/mesh.py index 748c006e..e066f3ff 100644 --- a/sdk/python/src/agent_mesh/mesh.py +++ b/sdk/python/src/agent_mesh/mesh.py @@ -74,6 +74,7 @@ def enroll( bootstrap_token: Optional[str] = None, bootstrap_token_path: Optional[str | os.PathLike[str]] = None, jwt: Optional[str] = None, + jwt_path: Optional[str | os.PathLike[str]] = None, state_dir: Optional[str | os.PathLike[str]] = None, identity: Optional[Identity] = None, role: str = ROLE_NODE, @@ -89,9 +90,11 @@ def enroll( plane for the saved identity, that member is returned and no token is needed, so a program can call enroll on every start and read the token from its environment only on the first. Otherwise exactly one of - bootstrap_token, bootstrap_token_path or jwt must be given; a token is - better read from a file than passed as a value. Delete the state - directory to enroll afresh, for instance with other labels.""" + bootstrap_token, bootstrap_token_path, jwt or jwt_path must be given; a + token is better read from a file than passed as a value, and jwt_path + also takes a platform's workload identity token, such as a Kubernetes + projected service account token. Delete the state directory to enroll + afresh, for instance with other labels.""" state = Path(state_dir).expanduser() if state_dir is not None else None saved = _load_identity(state) identity = identity or saved or Identity.generate() @@ -100,13 +103,16 @@ def enroll( credential = _load_credential(state) if credential is not None and credential.control_plane_url.rstrip("/") == control_plane.url.rstrip("/") and credential.time_to_live_seconds() > _REUSE_MIN_TTL_SECONDS: return cls(identity, control_plane, credential, state) - given = sum(v is not None for v in (bootstrap_token, bootstrap_token_path, jwt)) + given = sum(v is not None for v in (bootstrap_token, bootstrap_token_path, jwt, jwt_path)) if given != 1: where = f" (no credential to resume in {state})" if state is not None else "" - raise ValueError(f"exactly one of bootstrap_token, bootstrap_token_path or jwt is required{where}") + raise ValueError(f"exactly one of bootstrap_token, bootstrap_token_path, jwt or jwt_path is required{where}") enrollment: Enrollment - if jwt is not None: + if jwt is not None or jwt_path is not None: + if jwt_path is not None: + jwt = Path(jwt_path).expanduser().read_text(encoding="utf-8").strip() + assert jwt is not None enrollment = control_plane.register(identity, jwt, role=role, labels=labels) else: if bootstrap_token_path is not None: diff --git a/sdk/python/src/agent_mesh/relay.py b/sdk/python/src/agent_mesh/relay.py index 709896b8..9896027e 100644 --- a/sdk/python/src/agent_mesh/relay.py +++ b/sdk/python/src/agent_mesh/relay.py @@ -34,6 +34,7 @@ from libp2p.utils.varint import encode_varint_prefixed, read_varint_prefixed_bytes from ._proto import circuit_pb2 as circuit +from .identity import canonical_peer_id logger = logging.getLogger("agent_mesh") @@ -52,7 +53,7 @@ def split_circuit_address(addr: multiaddr.Multiaddr) -> tuple[multiaddr.Multiadd target = tail.removeprefix("/p2p/") if not target: raise ValueError(f"{addr} names no target peer after /p2p-circuit") - return multiaddr.Multiaddr(relay_text), ID.from_base58(target) + return multiaddr.Multiaddr(relay_text), ID.from_base58(canonical_peer_id(target)) async def reserve_relay(host: IHost, relay_peer_id: ID) -> circuit.Reservation: diff --git a/sdk/python/src/agent_mesh/session.py b/sdk/python/src/agent_mesh/session.py index 6d227263..cfd51264 100644 --- a/sdk/python/src/agent_mesh/session.py +++ b/sdk/python/src/agent_mesh/session.py @@ -42,7 +42,8 @@ from .authorizer import ProviderAuthorizerOptions from .biscuit import ROLE_ROUTER, VerifiedBiscuit, require_role from .discovery import DiscoveredProvider, ServiceType, find_providers, parse_service_target, provide, service_key -from .host import create_mesh_host +from .host import create_mesh_host, dial_addrs, peer_info +from .identity import canonical_peer_id from .mcp_client import ToolCallResult, ToolInfo, open_mcp_session, tool_call_result from .relay import STOP_PROTOCOL, dial_through_relay, reserve_relay, split_circuit_address, stop_stream_handler from .serve import ( @@ -74,6 +75,11 @@ FIRST_CONTROL_PLANE_SYNC = 2.0 DEFAULT_CONTROL_PLANE_SYNC_JITTER = 2.0 +# sam-node's swarm dial timeout. py-libp2p has none of its own: a SYN to an +# address nobody answers waits on the kernel, about two minutes, and is then +# retried, and a provider record can name a pod a rollout just replaced. +DIAL_TIMEOUT = 15.0 + # How a caller names the peer it wants to reach: a provider `discover` returned, # a peer id, or a multiaddr. For a provider or a peer id the SDK dials the # addresses the peer advertised and then the relayed path through every router @@ -82,6 +88,32 @@ Peer = Union[DiscoveredProvider, str, multiaddr.Multiaddr] +def parse_peer_id(text: str) -> ID: + """The libp2p peer ID for a string in any encoding libp2p accepts.""" + return ID.from_base58(canonical_peer_id(text)) + + +async def dial(host: IHost, info: PeerInfo) -> None: + """host.connect, bounded by DIAL_TIMEOUT.""" + try: + with trio.fail_after(DIAL_TIMEOUT): + await host.connect(info) + except trio.TooSlowError: + raise ConnectionError(f"no connection to {info.peer_id} within {DIAL_TIMEOUT:g}s") from None + + +def canonical_peer_ids(ids: Sequence[str]) -> list[str]: + """Canonicalizes a list from the control plane, dropping entries that are + not peer IDs: they can match nothing, so they ban nothing.""" + out: list[str] = [] + for text in ids: + try: + out.append(canonical_peer_id(text)) + except ValueError: + continue + return out + + class _MeshsubNoise(logging.Filter): """py-libp2p's pubsub opens a meshsub stream to every new peer and the host logs an error for each one that does not answer. Peers that leave during @@ -156,18 +188,25 @@ async def connect(self, peer: Peer) -> ID: if isinstance(peer, multiaddr.Multiaddr) or (isinstance(peer, str) and peer.startswith("/")): return await self._connect_addr(multiaddr.Multiaddr(str(peer))) if isinstance(peer, str): - target, advertised = ID.from_base58(peer), [] + target, advertised = parse_peer_id(peer), [] else: - target, advertised = ID.from_base58(peer.peer_id), [multiaddr.Multiaddr(a) for a in peer.addrs] + target, advertised = parse_peer_id(peer.peer_id), [multiaddr.Multiaddr(a) for a in peer.addrs] self._refuse_banned(target) if target in self.host.get_connected_peers(): return target failures: list[str] = [] suffix = f"/p2p/{target}" - direct = [multiaddr.Multiaddr(str(a).removesuffix(suffix)) for a in advertised if "/p2p-circuit" not in str(a)] + direct: list[multiaddr.Multiaddr] = [] + for a in advertised: + if "/p2p-circuit" in str(a): + continue + try: + direct.extend(multiaddr.Multiaddr(str(m).removesuffix(suffix)) for m in await dial_addrs(a)) + except Exception as err: # noqa: BLE001 - an address this host cannot use; the others are tried + failures.append(f"{a}: {err}") if direct: try: - await self.host.connect(PeerInfo(target, direct)) + await dial(self.host, PeerInfo(target, direct)) return target except Exception as err: # noqa: BLE001 - the relayed path is tried next failures.append(f"direct {[str(a) for a in direct]}: {err}") @@ -186,13 +225,13 @@ async def _connect_addr(self, ma: multiaddr.Multiaddr) -> ID: self._refuse_banned(target) relay = info_from_p2p_addr(relay_addr) if relay.peer_id not in self.host.get_connected_peers(): - await self.host.connect(relay) + await dial(self.host, await peer_info(relay_addr)) if target not in self.host.get_connected_peers(): await dial_through_relay(self.host, relay.peer_id, target) return target - info = info_from_p2p_addr(ma) + info = await peer_info(ma) self._refuse_banned(info.peer_id) - await self.host.connect(info) + await dial(self.host, info) return info.peer_id def _refuse_banned(self, peer_id: ID) -> None: @@ -273,9 +312,9 @@ async def sync(self) -> "ControlPlaneSync": async with self._sync_lock: result = await trio.to_thread.run_sync(self.mesh.sync_control_plane) if result.banned_peer_ids is not None: - newly_banned, _ = self.banned.reconcile(result.banned_peer_ids, result.fetched_at) - for peer_id in newly_banned: - await self._evict(peer_id) + newly_banned, _ = self.banned.reconcile(canonical_peer_ids(result.banned_peer_ids), result.fetched_at) + for peer in newly_banned: + await self._evict(peer) if self._serving: try: await self.sync_policy() @@ -287,11 +326,11 @@ def trigger_sync(self) -> None: """Asks for a pull soon, after a random delay so a fleet told at once does not pull at once.""" self._sync_trigger.set() - async def _evict(self, peer_id: str) -> None: + async def _evict(self, banned_peer: str) -> None: """Drops a banned peer: its admission and its connections.""" - self.authenticated_peers.pop(peer_id, None) + self.authenticated_peers.pop(banned_peer, None) try: - await self.host.disconnect(ID.from_base58(peer_id)) + await self.host.disconnect(ID.from_base58(banned_peer)) except Exception: # noqa: BLE001 - not connected, or already gone pass @@ -521,8 +560,8 @@ async def _admit(host: IHost, mesh: "AgentMesh", router_addrs: list[multiaddr.Mu failures: list[str] = [] for addr in router_addrs: try: - info = info_from_p2p_addr(addr) - await host.connect(info) + info = await peer_info(addr) + await dial(host, info) credential = await authenticate_with_peer(host, info.peer_id, mesh.auth_frame(), mesh.credential.control_plane_keys) # Enforced under the key that verified the token; a relay that # is not a router must not become our way onto the mesh. diff --git a/sdk/python/src/agent_mesh/sync.py b/sdk/python/src/agent_mesh/sync.py index 70841576..460ca142 100644 --- a/sdk/python/src/agent_mesh/sync.py +++ b/sdk/python/src/agent_mesh/sync.py @@ -22,7 +22,7 @@ from typing import Optional, Sequence from ._proto import sam_pb2 as pb -from .identity import verify_ed25519 +from .identity import canonical_peer_id, verify_ed25519 # The GossipSub topic the control plane publishes mesh events on (api.GossipEvents). GOSSIP_EVENTS_TOPIC = "/sam/mesh/events/v1" @@ -71,8 +71,9 @@ def reconcile(self, banned_peer_ids: Sequence[str], fetched_at: float) -> tuple[ def verify_mesh_event(data: bytes, trusted_keys: Sequence[bytes], now_ms: Optional[int] = None) -> Optional[pb.MeshEvent]: """Verifies a MeshEvent as sam-node's verifyEvent does: the signature covers the deterministic encoding of the event with the signature cleared, under - any trusted control plane key. Returns the event, or None when it does not - verify or is not fresh.""" + any trusted control plane key. Returns the event, with a banned peer's id in + canonical form, or None when it does not verify, is not fresh, or bans + something that is not a peer ID.""" try: event = pb.MeshEvent.FromString(data) except Exception: # noqa: BLE001 - undecodable is unverifiable @@ -87,4 +88,10 @@ def verify_mesh_event(data: bytes, trusted_keys: Sequence[bytes], now_ms: Option now_ms = int(time.time() * 1000) if now_ms is None else now_ms if not event.HasField("event_time") or abs(now_ms - event.event_time.ToMilliseconds()) > EVENT_FRESHNESS_MS: return None + if event.type == pb.MeshEvent.BANNED: + # Canonicalized after the signature check, which covers the bytes as sent. + try: + event.peer_id = canonical_peer_id(event.peer_id) + except ValueError: + return None return event diff --git a/sdk/python/tests/test_host.py b/sdk/python/tests/test_host.py new file mode 100644 index 00000000..9e6996cd --- /dev/null +++ b/sdk/python/tests/test_host.py @@ -0,0 +1,116 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""What a dial must do that py-libp2p does not: resolve a /dnsaddr router +address through its TXT records and keep the transports this host has +(dial_addrs), and give up on an address nobody answers (dial).""" + +import multiaddr +import pytest +import trio +import trio.testing +from libp2p.peer.id import ID +from libp2p.peer.peerinfo import PeerInfo +from multiaddr.resolvers import DNSResolver + +from agent_mesh.host import dial_addrs, peer_info +from agent_mesh.session import DIAL_TIMEOUT, dial + +ROUTER = "12D3KooWG1pA6goegCncqwbZLSr8pnjUZ6JMAAe6SmnHTgUNCk88" +OTHER = "12D3KooWGvdRCJLYATauVWfsieF2j3a2wXZoEQJUS2MsvRdDtgLM" + + +class FakeTXT: + """A dnspython answer: one TXT record per string.""" + + def __init__(self, strings): + self._strings = strings + + def __iter__(self): + for s in self._strings: + yield type("TXT", (), {"strings": [s.encode()]})() + + def __len__(self): + return len(self._strings) + + +@pytest.fixture +def testnet_dns(monkeypatch): + """The _dnsaddr records a testnet publishes: two routers, TCP and QUIC each.""" + records = { + "_dnsaddr.bootstrap.example": [ + f"dnsaddr=/ip4/203.0.113.1/udp/4501/quic-v1/p2p/{ROUTER}", + f"dnsaddr=/ip4/203.0.113.1/tcp/4501/p2p/{ROUTER}", + f"dnsaddr=/ip4/203.0.113.2/tcp/4501/p2p/{OTHER}", + f"dnsaddr=/ip4/203.0.113.2/udp/4501/quic-v1/p2p/{OTHER}", + ], + "_dnsaddr.quic-only.example": [f"dnsaddr=/ip4/203.0.113.3/udp/4501/quic-v1/p2p/{ROUTER}"], + } + + class FakeDNS: + async def resolve(self, name, rdtype): + assert rdtype == "TXT" + return FakeTXT(records.get(str(name).rstrip("."), [])) + + def init(self): + self._resolver = FakeDNS() + + monkeypatch.setattr(DNSResolver, "__init__", init) + + +def test_dnsaddr_router_address_resolves_to_its_tcp_address(testnet_dns): + async def main(): + addrs = await dial_addrs(multiaddr.Multiaddr(f"/dnsaddr/bootstrap.example/p2p/{ROUTER}")) + assert [str(a) for a in addrs] == [f"/ip4/203.0.113.1/tcp/4501/p2p/{ROUTER}"] + info = await peer_info(multiaddr.Multiaddr(f"/dnsaddr/bootstrap.example/p2p/{ROUTER}")) + assert str(info.peer_id) == ROUTER and len(info.addrs) == 1 + + trio.run(main) + + +def test_addresses_without_a_transport_this_host_has_are_an_error(testnet_dns): + async def main(): + with pytest.raises(RuntimeError, match="TCP"): + await dial_addrs(multiaddr.Multiaddr(f"/dnsaddr/quic-only.example/p2p/{ROUTER}")) + with pytest.raises(RuntimeError, match="TCP"): + await dial_addrs(multiaddr.Multiaddr(f"/ip4/203.0.113.9/udp/4501/quic-v1/p2p/{ROUTER}")) + with pytest.raises(RuntimeError, match="no address"): + await dial_addrs(multiaddr.Multiaddr(f"/dnsaddr/nowhere.example/p2p/{ROUTER}")) + + trio.run(main) + + +def test_a_concrete_address_passes_through(): + async def main(): + ma = multiaddr.Multiaddr(f"/ip4/127.0.0.1/tcp/4001/p2p/{ROUTER}") + assert await dial_addrs(ma) == [ma] + + trio.run(main) + + +def test_a_dial_nobody_answers_ends_at_the_timeout(): + """A provider record can name a pod a rollout replaced; its SYNs go + unanswered. The dial ends at DIAL_TIMEOUT, not at the kernel's.""" + + class Host: + async def connect(self, info): + await trio.sleep_forever() + + async def main(): + started = trio.current_time() + with pytest.raises(ConnectionError, match=f"within {DIAL_TIMEOUT:g}s"): + await dial(Host(), PeerInfo(ID.from_base58(ROUTER), [])) + assert trio.current_time() - started == pytest.approx(DIAL_TIMEOUT) + + trio.run(main, clock=trio.testing.MockClock(autojump_threshold=0)) diff --git a/sdk/python/tests/test_identity.py b/sdk/python/tests/test_identity.py index 78edba2b..ba472b8a 100644 --- a/sdk/python/tests/test_identity.py +++ b/sdk/python/tests/test_identity.py @@ -19,7 +19,7 @@ from agent_mesh import base58 from agent_mesh.challenges import enroll_challenge -from agent_mesh.identity import Identity, libp2p_public_key, peer_id_from_public_key, verify_ed25519 +from agent_mesh.identity import Identity, canonical_peer_id, libp2p_public_key, peer_id_from_public_key, verify_ed25519 VECTORS = json.loads((Path(__file__).resolve().parents[2] / "testdata" / "identity_vectors.json").read_text())["vectors"] @@ -64,6 +64,17 @@ def test_libp2p_public_key_refuses_wrong_size(): libp2p_public_key(b"\x00" * 33) +def test_canonical_peer_id_is_the_base58_form_for_every_encoding(): + # The CIDv1 form of a known ed25519 peer, as `peer.ToCid(id).String()` prints it. + b58 = "12D3KooWA4Xop1JaT3MHxwYMkCepYsv4iPVopMXwCz5iHYdBfeSB" + cidv1 = "bafzaajaiaejcaa5ba677htqqxyoxbxiy45f4bglh4tldbg5fbvpr3xegmqjfkmny" + assert canonical_peer_id(b58) == b58 + assert canonical_peer_id(cidv1) == b58 + for bad in ("not-a-peer", ""): + with pytest.raises(ValueError, match="is not a peer ID"): + canonical_peer_id(bad) + + @pytest.mark.parametrize("data", [b"", b"\x00", b"\x00\x00\x01\x02", bytes.fromhex("00ff"), bytes.fromhex("deadbeef")]) def test_base58_round_trips(data): assert base58.decode(base58.encode(data)) == data diff --git a/sdk/python/tests/test_mesh.py b/sdk/python/tests/test_mesh.py index ead4917d..0af392d1 100644 --- a/sdk/python/tests/test_mesh.py +++ b/sdk/python/tests/test_mesh.py @@ -55,6 +55,16 @@ def _signed_keys(self): def transport(self, method, url, headers, body): path = urllib.parse.urlsplit(url).path + if (method, path) == ("POST", "/register"): + # The biscuit names the JWT that was presented, so a test can see which. + jwt = pb.EnrollRequest.FromString(body).jwt + self.issued += 1 + return 200, pb.EnrollResponse( + biscuit_token=f"biscuit-for-{jwt}".encode(), + control_plane_public_key=CP_KEY.public_key_raw, + router_addresses=["/dns4/router.example/tcp/4001/p2p/12D3KooWP8iKhDf3iCMo2H3butNVfdTUtYwYWYQ75jTGnynXPFMp"], + expire_time=_ts_s(int(time.time()) + 3600), + ).SerializeToString() if (method, path) == ("POST", "/enroll"): self.issued += 1 return 200, pb.BootstrapEnrollResponse( @@ -149,9 +159,21 @@ def test_enroll_refuses_ambiguous_credentials(): AgentMesh.enroll("http://127.0.0.1:1", transport=cp.transport) with pytest.raises(ValueError, match="exactly one of"): AgentMesh.enroll("http://127.0.0.1:1", bootstrap_token="a", jwt="b", transport=cp.transport) + with pytest.raises(ValueError, match="exactly one of"): + AgentMesh.enroll("http://127.0.0.1:1", jwt="a", jwt_path="/nonexistent", transport=cp.transport) assert cp.issued == 0 +def test_enroll_reads_a_workload_identity_token_from_jwt_path(tmp_path): + cp = FakeControlPlane() + token = tmp_path / "token" + token.write_text("eyJ.projected.token\n") + mesh = AgentMesh.enroll("http://127.0.0.1:1", jwt_path=token, transport=cp.transport) + assert mesh.credential.biscuit == b"biscuit-for-eyJ.projected.token" + with pytest.raises(FileNotFoundError): + AgentMesh.enroll("http://127.0.0.1:1", jwt_path=tmp_path / "missing", transport=cp.transport) + + def test_load_without_identity_says_enroll_first(tmp_path): with pytest.raises(FileNotFoundError, match="no identity"): AgentMesh.load(tmp_path) diff --git a/sdk/python/tests/test_sync.py b/sdk/python/tests/test_sync.py index bcf46907..545f2f6d 100644 --- a/sdk/python/tests/test_sync.py +++ b/sdk/python/tests/test_sync.py @@ -26,12 +26,15 @@ import multiaddr import pytest import trio +from cid import make_cid from libp2p.peer.peerinfo import info_from_p2p_addr +from agent_mesh import base58 from agent_mesh._proto import sam_pb2 as pb from agent_mesh.auth import AUTH_PROTOCOL, auth_stream_handler, authenticate_with_peer from agent_mesh.biscuit import ROLE_ROUTER, verify_peer_biscuit from agent_mesh.controlplane import ROLE_NODE +from agent_mesh.discovery import DiscoveredProvider from agent_mesh.identity import Identity from agent_mesh.mesh import AgentMesh from agent_mesh.relay import HOP_PROTOCOL as RELAY_HOP_PROTOCOL @@ -164,16 +167,31 @@ def sign(event: pb.MeshEvent) -> bytes: return event.SerializeToString(deterministic=True) now_ms = int(time.time() * 1000) - banned = sign(pb.MeshEvent(type=pb.MeshEvent.BANNED, peer_id="12D3KooWx", event_time=_ts_ms(now_ms))) - assert verify_mesh_event(banned, [key.pub], now_ms).peer_id == "12D3KooWx" + target = Identity.generate().peer_id + banned = sign(pb.MeshEvent(type=pb.MeshEvent.BANNED, peer_id=target, event_time=_ts_ms(now_ms))) + assert verify_mesh_event(banned, [key.pub], now_ms).peer_id == target assert verify_mesh_event(banned, [SigningKey().pub], now_ms) is None tampered = bytearray(banned) tampered[-1] ^= 1 assert verify_mesh_event(bytes(tampered), [key.pub], now_ms) is None - stale = sign(pb.MeshEvent(type=pb.MeshEvent.BANNED, peer_id="12D3KooWx", event_time=_ts_ms(now_ms - EVENT_FRESHNESS_MS - 1))) + stale = sign(pb.MeshEvent(type=pb.MeshEvent.BANNED, peer_id=target, event_time=_ts_ms(now_ms - EVENT_FRESHNESS_MS - 1))) assert verify_mesh_event(stale, [key.pub], now_ms) is None assert verify_mesh_event(b"\x01\x02\x03", [key.pub], now_ms) is None + # The ban set is keyed on the base58 form; an event naming the peer in its + # CIDv1 form bans the same peer, and one naming no peer bans nobody. + banned_by_cid = sign(pb.MeshEvent(type=pb.MeshEvent.BANNED, peer_id=cid_form(target), event_time=_ts_ms(now_ms))) + assert verify_mesh_event(banned_by_cid, [key.pub], now_ms).peer_id == target + banned_nobody = sign(pb.MeshEvent(type=pb.MeshEvent.BANNED, peer_id="not-a-peer", event_time=_ts_ms(now_ms))) + assert verify_mesh_event(banned_nobody, [key.pub], now_ms) is None + + +def cid_form(peer_id: str) -> str: + """The CIDv1 base32 encoding of a peer ID, as `peer.ToCid(id).String()` prints it.""" + encoded = make_cid(1, "libp2p-key", base58.decode(peer_id)).encode("base32").decode() + assert encoded != peer_id + return encoded + def test_pull_learns_a_key_rotation_and_refreshes_the_credential(): async def main(): @@ -245,11 +263,25 @@ async def main(): await peer.connect(info_from_p2p_addr(member_addr)) with pytest.raises(Exception): await authenticate_with_peer(peer, session.host.get_id(), frame, [cp.current.pub]) - # ... and outbound, before any dial. + # ... and outbound, before any dial, however the peer is named. + cid = cid_form(peer_identity.peer_id) with pytest.raises(PermissionError, match="banned"): await session.connect(f"{cp.router_addr}/p2p-circuit/p2p/{peer_identity.peer_id}") + with pytest.raises(PermissionError, match="banned"): + await session.connect(f"{cp.router_addr}/p2p-circuit/p2p/{cid}") + with pytest.raises(PermissionError, match="banned"): + await session.connect(cid) + with pytest.raises(PermissionError, match="banned"): + await session.connect(DiscoveredProvider(peer_id=cid)) - # Lifted by the control plane: the next pull unbans it. + # Lifted by the control plane: the next pull unbans it. A ban + # list that names the peer in its CIDv1 form bans the same peer. + cp.banned = [] + await session.sync() + assert session.banned.peers() == [] + cp.banned = [cid, "not-a-peer"] + await session.sync() + assert session.banned.peers() == [peer_identity.peer_id] cp.banned = [] await session.sync() assert session.banned.peers() == [] diff --git a/site/content/docs/guides/native-sdks.md b/site/content/docs/guides/native-sdks.md index 3be874cc..eedb11d9 100644 --- a/site/content/docs/guides/native-sdks.md +++ b/site/content/docs/guides/native-sdks.md @@ -76,7 +76,18 @@ export SAM_BOOTSTRAP_TOKEN_PATH=~/sam-one/join-token A plain `http://` URL is accepted only for a control plane on the same machine. Anything else needs `https://`, because the control plane is the trust root of every member and the SDK refuses to fetch it over plaintext -from a remote address. +from a remote address. Inside a network you already trust, such as a +Kubernetes cluster where the control plane is a cluster-local service, set +`SAM_INSECURE_CONTROL_PLANE=true` (`allowInsecure` in code), the same choice +as `sam-node --insecure-control-plane`. + +On a platform that issues workload identity tokens, a program needs no +bootstrap token. A mesh whose control plane trusts the platform's issuer +enrolls the program from that token instead; on Kubernetes that is a +projected service account token, as [Headless enrollment](../headless-enrollment/) +shows for `sam-node`. Set `SAM_JWT_PATH` to the token file rather than +`SAM_BOOTSTRAP_TOKEN_PATH`. That is how the public testnets run these same +programs as canaries beside the `sam-node` ones. ## 2. Install the SDK @@ -111,31 +122,38 @@ JavaScript, `serve.js`: // members may call; the SDK turns the others away before anything reaches // this code or Ollama. // -// node serve.js +// node serve.js # publishes mcp://greeter and a2a://greeter +// node serve.js greeter-2 # the same under another name // -// SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -// holding the token the mesh operator gave you; the first run spends it and -// keeps the identity and credential in SAM_STATE_DIR, later runs resume from -// there without it. +// SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +// SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or +// SAM_JWT_PATH (a workload identity token your platform issues, such as a +// Kubernetes projected service account token), and keeps the identity and +// credential in SAM_STATE_DIR; later runs resume from there without it. import { homedir } from "node:os"; import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { AgentMesh } from "@sam-mesh/sdk"; import { z } from "zod"; +const [name = "greeter"] = process.argv.slice(2); + const mesh = await AgentMesh.enroll({ controlPlaneUrl: process.env.SAM_CONTROL_PLANE_URL ?? "https://mesh.example.com", bootstrapTokenPath: process.env.SAM_BOOTSTRAP_TOKEN_PATH, - stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/greeter`, + jwtPath: process.env.SAM_JWT_PATH, + stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/${name}`, + // A plaintext http:// control plane is otherwise accepted only on loopback. + allowInsecure: process.env.SAM_INSECURE_CONTROL_PLANE === "true", }); const session = await mesh.join(); await session.serve({ type: "mcp", - name: "greeter", + name, createServer: () => { - const server = new McpServer({ name: "greeter", version: "1.0.0" }); - server.registerTool("greet", { description: "Greets someone by name", inputSchema: { name: z.string() } }, async ({ name }) => ({ - content: [{ type: "text", text: `hello ${name}` }], + const server = new McpServer({ name, version: "1.0.0" }); + server.registerTool("greet", { description: "Greets someone by name", inputSchema: { name: z.string() } }, async ({ name: who }) => ({ + content: [{ type: "text", text: `hello ${who}` }], })); return server; }, @@ -143,8 +161,8 @@ await session.serve({ await session.serve({ type: "a2a", - name: "greeter", - target: (request, caller) => Response.json({ name: "greeter", path: new URL(request.url).pathname, caller: caller.peerId }), + name, + target: (request, caller) => Response.json({ name, path: new URL(request.url).pathname, caller: caller.peerId }), }); if (process.env.OLLAMA_URL !== undefined) { @@ -169,30 +187,38 @@ beside this program as an inference service. The mesh policy decides which members may call; the SDK turns the others away before anything reaches this code or Ollama. - python serve.py + python serve.py # publishes mcp://greeter and a2a://greeter + python serve.py greeter-2 # the same under another name -SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -holding the token the mesh operator gave you; the first run spends it and -keeps the identity and credential in SAM_STATE_DIR, later runs resume from -there without it. +SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or SAM_JWT_PATH +(a workload identity token your platform issues, such as a Kubernetes +projected service account token), and keeps the identity and credential in +SAM_STATE_DIR; later runs resume from there without it. """ import json import os +import sys import trio from agent_mesh import AgentMesh, HTTPRequest, HTTPResponse, HTTPService, MCPService, VerifiedBiscuit from mcp.server.mcpserver import MCPServer +service_name = sys.argv[1] if len(sys.argv) > 1 else "greeter" + mesh = AgentMesh.enroll( os.environ.get("SAM_CONTROL_PLANE_URL", "https://mesh.example.com"), bootstrap_token_path=os.environ.get("SAM_BOOTSTRAP_TOKEN_PATH"), - state_dir=os.environ.get("SAM_STATE_DIR", "~/.config/sam-mesh/greeter"), + jwt_path=os.environ.get("SAM_JWT_PATH"), + state_dir=os.environ.get("SAM_STATE_DIR", f"~/.config/sam-mesh/{service_name}"), + # A plaintext http:// control plane is otherwise accepted only on loopback. + allow_insecure=os.environ.get("SAM_INSECURE_CONTROL_PLANE") == "true", ) def create_server() -> MCPServer: - server = MCPServer("greeter") + server = MCPServer(service_name) @server.tool(description="Greets someone by name") def greet(name: str) -> str: @@ -202,14 +228,14 @@ def create_server() -> MCPServer: async def card(request: HTTPRequest, caller: VerifiedBiscuit) -> HTTPResponse: - body = json.dumps({"name": "greeter", "path": request.path, "caller": caller.peer_id}) + body = json.dumps({"name": service_name, "path": request.path, "caller": caller.peer_id}) return HTTPResponse(status=200, headers={"content-type": "application/json"}, body=body.encode()) async def main() -> None: async with mesh.join() as session: - await session.serve(MCPService(name="greeter", create_server=create_server)) - await session.serve(HTTPService(type="a2a", name="greeter", target=card)) + await session.serve(MCPService(name=service_name, create_server=create_server)) + await session.serve(HTTPService(type="a2a", name=service_name, target=card)) if "OLLAMA_URL" in os.environ: await session.serve(HTTPService(type="inference", name="ollama", target=os.environ["OLLAMA_URL"])) @@ -257,27 +283,45 @@ JavaScript, `call.js`: // node call.js a2a://greeter /card // node call.js inference://ollama /v1/models // -// SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -// holding the token the mesh operator gave you; the first run spends it and -// keeps the identity and credential in SAM_STATE_DIR, later runs resume from -// there without it. +// SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +// SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or +// SAM_JWT_PATH (a workload identity token your platform issues, such as a +// Kubernetes projected service account token), and keeps the identity and +// credential in SAM_STATE_DIR; later runs resume from there without it. import { homedir } from "node:os"; -import { AgentMesh } from "@sam-mesh/sdk"; +import { AgentMesh, type DiscoveredProvider } from "@sam-mesh/sdk"; const [service = "mcp://greeter", toolOrPath = "greet", args = '{"name": "world"}'] = process.argv.slice(2); const mesh = await AgentMesh.enroll({ controlPlaneUrl: process.env.SAM_CONTROL_PLANE_URL ?? "https://mesh.example.com", bootstrapTokenPath: process.env.SAM_BOOTSTRAP_TOKEN_PATH, + jwtPath: process.env.SAM_JWT_PATH, stateDir: process.env.SAM_STATE_DIR ?? `${homedir()}/.config/sam-mesh/caller`, + // A plaintext http:// control plane is otherwise accepted only on loopback. + allowInsecure: process.env.SAM_INSECURE_CONTROL_PLANE === "true", }); const session = await mesh.join(); console.log(`on the mesh as ${session.peerId}`); -const [provider] = await session.discover(service); -if (provider === undefined) { +const providers = await session.discover(service); +if (providers.length === 0) { throw new Error(`no member of the mesh serves ${service}`); } +// A provider record can outlive its member; the first that answers is used. +let provider: DiscoveredProvider | undefined; +for (const candidate of providers) { + try { + await session.connect(candidate); + provider = candidate; + break; + } catch (err) { + console.error(`${candidate.peerId}: ${(err as Error).message}`); + } +} +if (provider === undefined) { + throw new Error(`no provider of ${service} is reachable`); +} console.log(`${service} is served by ${provider.peerId}`); if (service.startsWith("mcp://")) { @@ -305,10 +349,11 @@ path of an inference or A2A service. python call.py a2a://greeter /card python call.py inference://ollama /v1/models -SAM_CONTROL_PLANE_URL names the mesh. SAM_BOOTSTRAP_TOKEN_PATH is the file -holding the token the mesh operator gave you; the first run spends it and -keeps the identity and credential in SAM_STATE_DIR, later runs resume from -there without it. +SAM_CONTROL_PLANE_URL names the mesh. The first run enrolls with the file +SAM_BOOTSTRAP_TOKEN_PATH (a token the mesh operator gave you) or SAM_JWT_PATH +(a workload identity token your platform issues, such as a Kubernetes +projected service account token), and keeps the identity and credential in +SAM_STATE_DIR; later runs resume from there without it. """ import json @@ -326,7 +371,10 @@ args = json.loads(argv[2]) if len(argv) > 2 else {"name": "world"} mesh = AgentMesh.enroll( os.environ.get("SAM_CONTROL_PLANE_URL", "https://mesh.example.com"), bootstrap_token_path=os.environ.get("SAM_BOOTSTRAP_TOKEN_PATH"), + jwt_path=os.environ.get("SAM_JWT_PATH"), state_dir=os.environ.get("SAM_STATE_DIR", "~/.config/sam-mesh/caller"), + # A plaintext http:// control plane is otherwise accepted only on loopback. + allow_insecure=os.environ.get("SAM_INSECURE_CONTROL_PLANE") == "true", ) @@ -337,7 +385,15 @@ async def main() -> None: providers = await session.discover(service) if not providers: raise SystemExit(f"no member of the mesh serves {service}") - provider = providers[0] + # A provider record can outlive its member; the first that answers is used. + for provider in providers: + try: + await session.connect(provider) + break + except (ConnectionError, PermissionError) as err: + print(f"{provider.peer_id}: {err}", file=sys.stderr) + else: + raise SystemExit(f"no provider of {service} is reachable") print(f"{service} is served by {provider.peer_id}") if service.startswith("mcp://"): diff --git a/tests/integration/sdk_examples_test.go b/tests/integration/sdk_examples_test.go index 202df849..d0b869e5 100644 --- a/tests/integration/sdk_examples_test.go +++ b/tests/integration/sdk_examples_test.go @@ -30,6 +30,8 @@ import ( "sync" "testing" "time" + + "github.com/google/sam/api" ) // sdkExampleLauncher starts one of the example programs the SDK READMEs and @@ -75,22 +77,26 @@ var sdkExampleLaunchers = []sdkExampleLauncher{ // sdkExampleServer is a running serve example. type sdkExampleServer struct { - name string - peerID string + name string + // service is what it publishes, greeter-, as the testnet canaries do. + service string + peerID string } var servingLine = regexp.MustCompile(`^serving (.*) as (\S+)$`) // TestNativeSDKExamples runs the programs the SDK READMEs and the Native -// SDKs guide embed, unchanged, against a real mesh. Each SDK's serve -// example publishes mcp://greeter, a2a://greeter and inference://ollama, -// the last forwarding to a stand-in for Ollama this test runs. Each SDK's -// call example then reaches the sam-node's mcp://calc and a greeter, which -// is whichever serve example the DHT lists first when both toolchains are -// present, naming providers only by what discover returned; the SDK finds -// the path through the router by itself. The first run spends a bootstrap -// token, the later runs resume from the state directory without one. -// Cross-language calls with a fixed pairing are TestNativeSDKsMesh's job. +// SDKs guide embed, unchanged, against a real mesh, configured as the +// testnet canaries are (.github/k8s/sam-sdk-canary-template.yaml). Each +// SDK's serve example enrolls with an OIDC token through SAM_JWT_PATH, the +// way a Kubernetes workload does, and publishes mcp://greeter-, +// a2a://greeter- and inference://ollama, the last forwarding to a +// stand-in for Ollama this test runs. Each SDK's call example enrolls with +// a bootstrap token, reaches the sam-node's mcp://calc and the other +// language's greeter (its own when the other toolchain is missing), naming +// providers only by what discover returned; the SDK finds the path through +// the router by itself. The first run spends the token, the later runs +// resume from the state directory without one. func TestNativeSDKExamples(t *testing.T) { mesh := startSDKMesh(t) @@ -107,7 +113,7 @@ func TestNativeSDKExamples(t *testing.T) { var launchers []sdkExampleLauncher var cmds []*exec.Cmd for _, l := range sdkExampleLaunchers { - cmd, skip := l.cmd(context.Background(), mesh.root, "serve") + cmd, skip := l.cmd(context.Background(), mesh.root, "serve", "greeter-"+l.name) if skip != "" { t.Logf("%s SDK skipped: %s", l.name, skip) continue @@ -122,13 +128,14 @@ func TestNativeSDKExamples(t *testing.T) { errs := make([]error, len(launchers)) var wg sync.WaitGroup for i := range launchers { - tokenPath := filepath.Join(t.TempDir(), "join-token") - if err := os.WriteFile(tokenPath, []byte(mintBootstrapToken(t, mesh.baseURL, mesh.adminToken)+"\n"), 0o600); err != nil { + jwtPath := filepath.Join(t.TempDir(), "sam-token") + jwt := mesh.mintToken(map[string]interface{}{"sub": "mock-user", "roles": []string{api.RoleNode}}) + if err := os.WriteFile(jwtPath, []byte(jwt+"\n"), 0o600); err != nil { t.Fatal(err) } cmds[i].Env = append(os.Environ(), "SAM_CONTROL_PLANE_URL="+mesh.baseURL, - "SAM_BOOTSTRAP_TOKEN_PATH="+tokenPath, + "SAM_JWT_PATH="+jwtPath, "SAM_STATE_DIR="+filepath.Join(t.TempDir(), "state"), "OLLAMA_URL="+ollama.URL, ) @@ -136,7 +143,7 @@ func TestNativeSDKExamples(t *testing.T) { wg.Add(1) go func(i int) { defer wg.Done() - servers[i], errs[i] = startExampleServer(t, launchers[i].name, cmds[i]) + servers[i], errs[i] = startExampleServer(t, launchers[i].name, "greeter-"+launchers[i].name, cmds[i]) }(i) } wg.Wait() @@ -149,19 +156,26 @@ func TestNativeSDKExamples(t *testing.T) { waitForPeerOnRouter(t, mesh.cpPort, mesh.adminToken, s.peerID, 5*time.Second) } - servedBy := func(t *testing.T, out, service string) { - t.Helper() - got := expectLine(t, out, service+" is served by ") + // The greeter a caller of one language targets: the other language's, + // as the testnet probes do, or its own when it is the only one present. + peerServer := func(name string) *sdkExampleServer { for _, s := range servers { - if s.peerID == got { - return + if s.name != name { + return s } } - t.Fatalf("%s is served by %s, which is none of the serve examples", service, got) + return servers[0] + } + servedBy := func(t *testing.T, out string, want *sdkExampleServer, serviceType string) { + t.Helper() + if got := expectLine(t, out, serviceType+"://"+want.service+" is served by "); got != want.peerID { + t.Fatalf("%s://%s is served by %s, want the %s serve example %s", serviceType, want.service, got, want.name, want.peerID) + } } for _, l := range launchers { l := l + target := peerServer(l.name) t.Run(l.name+"-calls", func(t *testing.T) { t.Parallel() stateDir := filepath.Join(t.TempDir(), "state") @@ -187,11 +201,11 @@ func TestNativeSDKExamples(t *testing.T) { if err := os.Remove(tokenPath); err != nil { t.Fatal(err) } - out = runExample(t, mesh, l, withoutToken, "a2a://greeter", "/card") + out = runExample(t, mesh, l, withoutToken, "a2a://"+target.service, "/card") if got := expectLine(t, out, "on the mesh as "); got != caller { t.Fatalf("second run joined as %s, want the identity of the first run %s", got, caller) } - servedBy(t, out, "a2a://greeter") + servedBy(t, out, target, "a2a") card := expectLine(t, out, "200 ") if !strings.Contains(card, `"path": "/card"`) && !strings.Contains(card, `"path":"/card"`) { t.Fatalf("a2a card %q does not echo the path", card) @@ -200,7 +214,9 @@ func TestNativeSDKExamples(t *testing.T) { t.Fatalf("a2a card %q does not name the verified caller %s", card, caller) } out = runExample(t, mesh, l, withoutToken, "inference://ollama", "/v1/models") - servedBy(t, out, "inference://ollama") + if got := expectLine(t, out, "inference://ollama is served by "); got != servers[0].peerID && (len(servers) < 2 || got != servers[1].peerID) { + t.Fatalf("inference://ollama is served by %s, which is none of the serve examples", got) + } models := expectLine(t, out, "200 ") if !strings.Contains(models, `"gemma3"`) || !strings.Contains(models, caller) { t.Fatalf("models %q: want gemma3 owned by the caller %s", models, caller) @@ -214,7 +230,7 @@ func TestNativeSDKExamples(t *testing.T) { if other.name == l.name { continue } - out := runExample(t, mesh, other, withoutToken, "mcp://greeter", "greet", `{"name": "sam"}`) + out := runExample(t, mesh, other, withoutToken, "mcp://"+target.service, "greet", `{"name": "sam"}`) if got := expectLine(t, out, "on the mesh as "); got != caller { t.Fatalf("%s resumed %s's state directory as %s, want %s", other.name, l.name, got, caller) } @@ -227,10 +243,16 @@ func TestNativeSDKExamples(t *testing.T) { importedAPI := importedNode.waitForAPI(t) waitForPeerOnRouter(t, mesh.cpPort, mesh.adminToken, caller, 10*time.Second) answer, err := callMCPAllowError(t, importedAPI, "imported-token", "call_remote_tool", map[string]any{ - "peer_id": servers[0].peerID, "tool_name": "mcp://greeter/greet", "arguments": map[string]any{"name": "node"}, + "peer_id": target.peerID, "tool_name": "mcp://" + target.service + "/greet", "arguments": map[string]any{"name": "node"}, }) if err != nil || !strings.Contains(answer, "hello node") { - t.Fatalf("sam-node running %s's identity could not call greeter: %v\n%s", l.name, err, answer) + t.Fatalf("sam-node running %s's identity could not call %s: %v\n%s", l.name, target.service, err, answer) + } + // The node's A2A egress path reaches the SDK's handler too, as the + // testnet's node probe expects. + status, body := egressGet(t, importedAPI, "imported-token", "/sam/"+target.peerID+"/a2a/"+target.service+"/card") + if status != 200 || !strings.Contains(body, `"caller"`) || !strings.Contains(body, caller) { + t.Fatalf("a2a card through the node: %d %s", status, body) } }) } @@ -264,7 +286,7 @@ func importStateIntoNode(t *testing.T, mesh *sdkMesh, stateDir string) *backgrou // startExampleServer starts a configured serve example and waits for its // "serving ... as " line. Safe to call from several goroutines at // once, so it reports failures instead of ending the test. -func startExampleServer(t *testing.T, name string, cmd *exec.Cmd) (*sdkExampleServer, error) { +func startExampleServer(t *testing.T, name, service string, cmd *exec.Cmd) (*sdkExampleServer, error) { stderr := &bytes.Buffer{} cmd.Stderr = stderr stdout, err := cmd.StdoutPipe() @@ -306,12 +328,12 @@ func startExampleServer(t *testing.T, name string, cmd *exec.Cmd) (*sdkExampleSe return nil, fmt.Errorf("%s serve example exited before serving\nstderr:\n%s", name, stderr.String()) } m := servingLine.FindStringSubmatch(line) - for _, want := range []string{"mcp://greeter", "a2a://greeter", "inference://ollama"} { + for _, want := range []string{"mcp://" + service, "a2a://" + service, "inference://ollama"} { if !strings.Contains(m[1], want) { return nil, fmt.Errorf("%s serve example serves %q, want %s among them", name, m[1], want) } } - return &sdkExampleServer{name: name, peerID: m[2]}, nil + return &sdkExampleServer{name: name, service: service, peerID: m[2]}, nil case <-time.After(30 * time.Second): return nil, fmt.Errorf("%s serve example did not report serving within 30s\nstderr:\n%s", name, stderr.String()) } diff --git a/tests/integration/sdk_mesh_test.go b/tests/integration/sdk_mesh_test.go index 60fed942..fec4eb0a 100644 --- a/tests/integration/sdk_mesh_test.go +++ b/tests/integration/sdk_mesh_test.go @@ -67,6 +67,9 @@ type sdkMesh struct { routerPeer string samNode *backgroundNode nodeAPI string + // mintToken mints an OIDC token the control plane accepts, as a platform's + // workload identity token would be. + mintToken func(map[string]interface{}) string } const sdkMeshAdminToken = "test-admin-token" @@ -172,6 +175,7 @@ bindings: routerPeer: extractPeerID(routerAddr), samNode: samNode, nodeAPI: samNode.waitForAPI(t), + mintToken: mintToken, } }