diff --git a/.github/k8s/sam-box-canary-template.yaml b/.github/k8s/sam-box-canary-template.yaml index 0247d16b..476ab62a 100644 --- a/.github/k8s/sam-box-canary-template.yaml +++ b/.github/k8s/sam-box-canary-template.yaml @@ -59,6 +59,7 @@ spec: metadata: labels: app: box-canary-${ENV_NAME} + sam-canary: "true" spec: serviceAccountName: sam-box-sa # nano-init ships as its own image with no shell in it, so the binary is @@ -84,6 +85,11 @@ spec: # the socket's permissions are the credential. - "--bind-addr=" - "--socket-path=/var/run/sam/node.sock" + # The one TCP port, and it carries nothing an API token would gate. + - "--metrics-addr=0.0.0.0:9090" + ports: + - containerPort: 9090 + name: metrics resources: requests: cpu: 50m @@ -107,6 +113,10 @@ spec: - "--sidecar-socket=/var/run/sam/node.sock" - "--egress-allow=example.com" - "--log-level=debug" + - "--metrics-addr=0.0.0.0:9091" + ports: + - containerPort: 9091 + name: box-metrics resources: requests: cpu: 20m diff --git a/.github/k8s/sam-canary-sa-template.yaml b/.github/k8s/sam-canary-sa-template.yaml new file mode 100644 index 00000000..215bebde --- /dev/null +++ b/.github/k8s/sam-canary-sa-template.yaml @@ -0,0 +1,8 @@ +# Identity shared by every canary node in the namespace; deploy.yaml's policy +# seed binds it to the sam-canary role. Applied with the namespaces so it +# exists before any canary references it. +apiVersion: v1 +kind: ServiceAccount +metadata: + name: sam-node-sa + namespace: sam-canary-${ENV_NAME} diff --git a/.github/k8s/sam-monitoring-template.yaml b/.github/k8s/sam-monitoring-template.yaml new file mode 100644 index 00000000..d6dab1cc --- /dev/null +++ b/.github/k8s/sam-monitoring-template.yaml @@ -0,0 +1,53 @@ +# Scrape targets for Google Managed Prometheus, which GKE runs by default. +# A PodMonitoring only sees pods in its own namespace, so the canaries get +# their own below. +apiVersion: monitoring.googleapis.com/v1 +kind: PodMonitoring +metadata: + name: sam-control-plane-${ENV_NAME} + namespace: ${NAMESPACE} +spec: + selector: + matchLabels: + app: sam-control-plane-${ENV_NAME} + endpoints: + - port: http + path: /metrics + interval: 30s +--- +apiVersion: monitoring.googleapis.com/v1 +kind: PodMonitoring +metadata: + name: sam-router-${ENV_NAME} + namespace: ${NAMESPACE} +spec: + selector: + matchLabels: + app: sam-router-${ENV_NAME} + endpoints: + - port: metrics + path: /metrics + interval: 30s +--- +# Every canary pod carries sam-canary=true; a pod without one of these named +# ports is simply not a target for that endpoint. +apiVersion: monitoring.googleapis.com/v1 +kind: PodMonitoring +metadata: + name: sam-canaries-${ENV_NAME} + namespace: sam-canary-${ENV_NAME} +spec: + selector: + matchLabels: + sam-canary: "true" + endpoints: + - port: metrics + path: /metrics + interval: 30s + - port: box-metrics + path: /metrics + interval: 30s + targetLabels: + fromPod: + - from: app + to: canary diff --git a/.github/k8s/sam-node-cop-template.yaml b/.github/k8s/sam-node-cop-template.yaml deleted file mode 100644 index 6a2d7474..00000000 --- a/.github/k8s/sam-node-cop-template.yaml +++ /dev/null @@ -1,98 +0,0 @@ -apiVersion: v1 -kind: ConfigMap -metadata: - name: sam-canary-cop-config-${ENV_NAME} - namespace: sam-canary-${ENV_NAME} -data: - sam-node.yaml: | - version: "v1alpha1" - attenuation: - policies: - - 'allow if service("mcp", "cop");' - - 'allow if service("system", "/sam/catalog");' - checks: [] - rules: [] - services: [] - banana_bot_playground.py: | -${BANANA_BOT_SCRIPT} ---- -apiVersion: apps/v1 -kind: Deployment -metadata: - name: cop-canary-${ENV_NAME} - namespace: sam-canary-${ENV_NAME} -spec: - replicas: 2 - selector: - matchLabels: - app: cop-canary-${ENV_NAME} - template: - metadata: - labels: - app: cop-canary-${ENV_NAME} - spec: - serviceAccountName: sam-node-sa - containers: - - name: cop-agent - image: python:3.11-slim - ports: - - containerPort: 18790 - name: http - command: - - "/bin/bash" - - "-c" - - | - pip install --no-cache-dir httpx 'mcp>=2,<3' && - python -u /app/banana_bot_playground.py - env: - - name: SAM_MCP_URL - value: "http://127.0.0.1:8080/mcp" - - name: SAM_API_TOKEN - value: "secret-token" - - name: GEMINI_API_KEY - valueFrom: - secretKeyRef: - name: openclaw-secret-${ENV_NAME} - key: gemini-api-key - volumeMounts: - - name: config-volume - mountPath: /app/banana_bot_playground.py - subPath: banana_bot_playground.py - - name: sam-node - image: ghcr.io/google/sam-node:${IMAGE_TAG} - env: - - name: SAM_API_TOKEN - value: "secret-token" - args: - - "run" - - "--config=/etc/sam/sam-node.yaml" - - "--control-plane=http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" - - "--insecure-control-plane" - - "--jwt-path=/var/run/secrets/tokens/sam-token" - - "--bind-addr=127.0.0.1:8080" - ports: - - containerPort: 8080 - resources: - requests: - cpu: 50m - memory: 64Mi - limits: - cpu: 200m - memory: 256Mi - volumeMounts: - - name: config-volume - mountPath: /etc/sam - - name: sam-token - mountPath: /var/run/secrets/tokens - readOnly: true - volumes: - - name: config-volume - configMap: - name: sam-canary-cop-config-${ENV_NAME} - - name: sam-token - projected: - sources: - - serviceAccountToken: - path: sam-token - expirationSeconds: 3600 - audience: "sam-control-plane-audience" diff --git a/.github/k8s/sam-node-everything-template.yaml b/.github/k8s/sam-node-everything-template.yaml index b79ac198..b5fa9ae1 100644 --- a/.github/k8s/sam-node-everything-template.yaml +++ b/.github/k8s/sam-node-everything-template.yaml @@ -30,6 +30,7 @@ spec: metadata: labels: app: everything-canary-${ENV_NAME} + sam-canary: "true" spec: serviceAccountName: sam-node-sa containers: @@ -58,8 +59,11 @@ spec: - "--insecure-control-plane" - "--jwt-path=/var/run/secrets/tokens/sam-token" - "--bind-addr=127.0.0.1:8080" + - "--metrics-addr=0.0.0.0:9090" ports: - containerPort: 8080 + - containerPort: 9090 + name: metrics resources: requests: cpu: 50m diff --git a/.github/k8s/sam-node-openclaw-template.yaml b/.github/k8s/sam-node-openclaw-template.yaml index 64c65976..2ce0a7bf 100644 --- a/.github/k8s/sam-node-openclaw-template.yaml +++ b/.github/k8s/sam-node-openclaw-template.yaml @@ -47,6 +47,7 @@ spec: metadata: labels: app: openclaw-canary-${ENV_NAME} + sam-canary: "true" spec: serviceAccountName: sam-node-sa containers: @@ -114,8 +115,11 @@ spec: - "--insecure-control-plane" - "--jwt-path=/var/run/secrets/tokens/sam-token" - "--bind-addr=127.0.0.1:8080" + - "--metrics-addr=0.0.0.0:9090" ports: - containerPort: 8080 + - containerPort: 9090 + name: metrics resources: requests: cpu: 50m diff --git a/.github/k8s/sam-node-openrouter-template.yaml b/.github/k8s/sam-node-openrouter-template.yaml index 578166a0..ae22d884 100644 --- a/.github/k8s/sam-node-openrouter-template.yaml +++ b/.github/k8s/sam-node-openrouter-template.yaml @@ -32,6 +32,7 @@ spec: metadata: labels: app: openrouter-canary-${ENV_NAME} + sam-canary: "true" spec: serviceAccountName: sam-node-sa containers: @@ -67,8 +68,11 @@ spec: - "--insecure-control-plane" - "--jwt-path=/var/run/secrets/tokens/sam-token" - "--bind-addr=127.0.0.1:8080" + - "--metrics-addr=0.0.0.0:9090" ports: - containerPort: 8080 + - containerPort: 9090 + name: metrics resources: requests: cpu: 500m diff --git a/.github/k8s/sam-node-template.yaml b/.github/k8s/sam-node-template.yaml deleted file mode 100644 index 7b3bfdad..00000000 --- a/.github/k8s/sam-node-template.yaml +++ /dev/null @@ -1,89 +0,0 @@ -apiVersion: v1 -kind: ConfigMap -metadata: - name: sam-canary-config-${ENV_NAME} - namespace: sam-canary-${ENV_NAME} -data: - sam-node.yaml: | - version: "v1alpha1" - attenuation: - policies: [] - checks: [] - rules: [] - services: - - type: mcp - name: dummy-http - description: "Canary HTTP tool (k8s agnhost)" - # agnhost netexec responds with request info on the root path - target_url: "http://localhost:9090/" ---- -apiVersion: v1 -kind: ServiceAccount -metadata: - name: sam-node-sa - namespace: sam-canary-${ENV_NAME} ---- -apiVersion: apps/v1 -kind: Deployment -metadata: - name: sam-canary-${ENV_NAME} - namespace: sam-canary-${ENV_NAME} -spec: - replicas: 3 - selector: - matchLabels: - app: sam-canary-${ENV_NAME} - template: - metadata: - labels: - app: sam-canary-${ENV_NAME} - spec: - serviceAccountName: sam-node-sa - containers: - - name: sam-node - image: ghcr.io/google/sam-node:${IMAGE_TAG} - env: - - name: SAM_API_TOKEN - value: "secret-token" - args: - - "run" - - "--config=/etc/sam/sam-node.yaml" - - "--control-plane=http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" - - "--insecure-control-plane" - - "--jwt-path=/var/run/secrets/tokens/sam-token" - ports: - - containerPort: 8080 - resources: - requests: - cpu: 50m - memory: 64Mi - limits: - cpu: 200m - memory: 256Mi - volumeMounts: - - name: config-volume - mountPath: /etc/sam - - name: shared-bin - mountPath: /shared - - name: sam-token - mountPath: /var/run/secrets/tokens - readOnly: true - - name: mock-mcp-server - image: registry.k8s.io/e2e-test-images/agnhost:2.39 - # agnhost's netexec subcommand spins up a robust test HTTP server - args: ["netexec", "--http-port=9090"] - ports: - - containerPort: 9090 - volumes: - - name: config-volume - configMap: - name: sam-canary-config-${ENV_NAME} - - name: shared-bin - emptyDir: {} - - name: sam-token - projected: - sources: - - serviceAccountToken: - path: sam-token - expirationSeconds: 3600 - audience: "sam-control-plane-audience" diff --git a/.github/k8s/sam-node-vllm-template.yaml b/.github/k8s/sam-node-vllm-template.yaml index a956e104..0a9dc86e 100644 --- a/.github/k8s/sam-node-vllm-template.yaml +++ b/.github/k8s/sam-node-vllm-template.yaml @@ -46,6 +46,7 @@ spec: metadata: labels: app: vllm-canary-${ENV_NAME} + sam-canary: "true" spec: serviceAccountName: sam-node-sa nodeSelector: @@ -107,8 +108,11 @@ spec: - "--insecure-control-plane" - "--jwt-path=/var/run/secrets/tokens/sam-token" - "--bind-addr=127.0.0.1:8080" + - "--metrics-addr=0.0.0.0:9090" ports: - containerPort: 8080 + - containerPort: 9090 + name: metrics resources: requests: cpu: 50m diff --git a/.github/k8s/sam-probe-cronjob-template.yaml b/.github/k8s/sam-probe-cronjob-template.yaml new file mode 100644 index 00000000..e4ded877 --- /dev/null +++ b/.github/k8s/sam-probe-cronjob-template.yaml @@ -0,0 +1,176 @@ +# 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. +apiVersion: batch/v1 +kind: CronJob +metadata: + name: sam-probe-${ENV_NAME} + namespace: sam-canary-${ENV_NAME} +spec: + schedule: "*/15 * * * *" + 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. + backoffLimit: 0 + activeDeadlineSeconds: 300 + template: + metadata: + labels: + app: sam-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 + initContainers: + # A native sidecar: it stays up for the probe container's lifetime + # and is stopped when that container exits, so the Job completes. + - name: sam-node + restartPolicy: Always + image: ghcr.io/google/sam-node:${IMAGE_TAG} + args: + - "run" + - "--config=/etc/sam/sam-node.yaml" + - "--control-plane=http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" + - "--insecure-control-plane" + - "--jwt-path=/var/run/secrets/tokens/sam-token" + - "--data-dir=/var/run/sam/data" + # Socket only: no TCP listener, so no API token to hand out. + - "--bind-addr=" + - "--socket-path=/var/run/sam/node.sock" + - "--metrics-addr=127.0.0.1:9090" + resources: + requests: + cpu: 50m + memory: 64Mi + limits: + cpu: 200m + memory: 256Mi + volumeMounts: + - name: config-volume + mountPath: /etc/sam + - name: sam-token + mountPath: /var/run/secrets/tokens + readOnly: true + - name: sam-uds + mountPath: /var/run/sam + containers: + - name: probe + image: alpine/curl:8.12.1 + # The node's socket is 0600 and owned by the uid the node runs as; + # the probe is that uid rather than root with an override. + securityContext: + runAsUser: 65532 + runAsGroup: 65532 + allowPrivilegeEscalation: false + capabilities: + drop: ["ALL"] + env: + - name: PROBE_SERVICE + # The everything canary's MCP server (sam-node-everything-template.yaml): + # a real server, so the node advertises it and initialize succeeds. + value: everything + resources: + requests: + cpu: 10m + memory: 16Mi + limits: + cpu: 100m + memory: 64Mi + volumeMounts: + - name: sam-uds + mountPath: /var/run/sam + command: ["/bin/sh", "-c"] + args: + - | + set -u + SOCK=/var/run/sam/node.sock + READYZ=http://127.0.0.1:9090/readyz + START=$(date +%s) + + # One JSON line per run, for a log-based metric; the exit code is + # what the Job reports. + fail() { + printf '{"probe":"sam-cold-path","ok":false,"stage":"%s","error":"%s","elapsed_s":%d}\n' \ + "$1" "$(printf '%s' "$2" | tr -d '"\n' | cut -c1-300)" "$(( $(date +%s) - START ))" + exit 1 + } + + # 1. Enroll and authenticate to a router: the node's /readyz. + until [ "$(curl -s -o /dev/null -w '%{http_code}' "$READYZ")" = "200" ]; do + [ $(( $(date +%s) - START )) -ge 180 ] && fail enroll "node not ready after 180s" + sleep 2 + done + READY_S=$(( $(date +%s) - START )) + + # 2. The router answers a dial. + CONN=$(curl -sf --unix-socket "$SOCK" http://localhost/debug/connectivity) \ + || fail connectivity "GET /debug/connectivity failed" + printf '%s' "$CONN" | grep -q '"router_error":false' \ + || fail connectivity "router dial failed: $CONN" + ROUTER_MS=$(printf '%s' "$CONN" | grep -oE '"router_latency_ms":[0-9]+' | grep -oE '[0-9]+$') + PEERS=$(printf '%s' "$CONN" | grep -oE '"connected_peers":[0-9]+' | grep -oE '[0-9]+$') + + # 3. A service someone else advertised is discoverable. A node + # that joined seconds ago has a thin routing table, so the lookup + # is retried for a while; how long it takes is itself reported. + PEER="" + DISCOVER_START=$(date +%s) + while [ -z "$PEER" ]; do + if curl -sf --unix-socket "$SOCK" -o /tmp/providers.json \ + "http://localhost/sam/service/discover?type=mcp&name=${PROBE_SERVICE}&timeout=20s"; then + PEER=$(grep -oE '"peer_id":"[^"]+"' /tmp/providers.json | head -1 | cut -d'"' -f4) + fi + [ -n "$PEER" ] && break + [ $(( $(date +%s) - DISCOVER_START )) -ge 90 ] && fail discover "no provider advertises ${PROBE_SERVICE} after 90s" + sleep 3 + done + DISCOVER_S=$(( $(date +%s) - DISCOVER_START )) + + # 4. And answers an MCP initialize through the mesh, end to end. + # No trailing slash: the path maps onto the service's target_url + # as-is. The body is JSON or an SSE frame; either names the server. + set -- $(curl -s --unix-socket "$SOCK" -o /tmp/init.out -w '%{http_code} %{time_total}' \ + -X POST "http://localhost/sam/${PEER}/mcp/${PROBE_SERVICE}" \ + -H 'Content-Type: application/json' -H 'Accept: application/json, text/event-stream' \ + -d '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"sam-probe","version":"0"}}}') + CALL_CODE=$1; CALL_S=$2 + [ "$CALL_CODE" = "200" ] || fail call "initialize on ${PEER} returned HTTP ${CALL_CODE}: $(head -c 200 /tmp/init.out)" + grep -q '"serverInfo"' /tmp/init.out || fail call "initialize on ${PEER} returned no serverInfo: $(head -c 200 /tmp/init.out)" + + printf '{"probe":"sam-cold-path","ok":true,"ready_s":%d,"router_latency_ms":%s,"connected_peers":%s,"discover_s":%d,"call_s":%s,"provider":"%s","elapsed_s":%d}\n' \ + "$READY_S" "${ROUTER_MS:-null}" "${PEERS:-null}" "$DISCOVER_S" "$CALL_S" "$PEER" "$(( $(date +%s) - START ))" + volumes: + - name: config-volume + configMap: + name: sam-probe-config-${ENV_NAME} + - name: sam-uds + emptyDir: {} + - name: sam-token + projected: + sources: + - serviceAccountToken: + path: sam-token + expirationSeconds: 3600 + audience: "sam-control-plane-audience" +--- +apiVersion: v1 +kind: ConfigMap +metadata: + name: sam-probe-config-${ENV_NAME} + namespace: sam-canary-${ENV_NAME} +data: + sam-node.yaml: | + version: "v1alpha1" + attenuation: + policies: [] + checks: [] + rules: [] + services: [] diff --git a/.github/k8s/sam-router-template.yaml b/.github/k8s/sam-router-template.yaml index 681a8a8a..3f916732 100644 --- a/.github/k8s/sam-router-template.yaml +++ b/.github/k8s/sam-router-template.yaml @@ -55,6 +55,10 @@ spec: hostPort: 4501 protocol: UDP name: p2p-udp + # Cluster-internal only: no hostPort, and nothing on it is authenticated. + - containerPort: 9090 + protocol: TCP + name: metrics args: - "--control-plane=http://sam-control-plane-${ENV_NAME}.${NAMESPACE}.svc.cluster.local:8080" - "--insecure-control-plane" @@ -63,6 +67,25 @@ spec: - "--external-addr=/dnsaddr/bootstrap.${ENV_NAME}.sam-mesh.dev" - "--jwt-path=/var/run/secrets/tokens/sam-token" - "--keys-path=/data/router.key" + - "--metrics-addr=0.0.0.0:9090" + # /healthz answers as soon as the process is up; /readyz once the + # router has enrolled and its libp2p host is online. + startupProbe: + httpGet: + path: /healthz + port: metrics + periodSeconds: 5 + failureThreshold: 24 + readinessProbe: + httpGet: + path: /readyz + port: metrics + periodSeconds: 10 + livenessProbe: + httpGet: + path: /healthz + port: metrics + periodSeconds: 20 resources: requests: cpu: 100m diff --git a/.github/workflows/deploy.yaml b/.github/workflows/deploy.yaml index ee95dbee..15ead65a 100644 --- a/.github/workflows/deploy.yaml +++ b/.github/workflows/deploy.yaml @@ -234,6 +234,7 @@ jobs: kubectl create namespace dex --dry-run=client -o yaml | kubectl apply -f - kubectl create namespace ${NAMESPACE} --dry-run=client -o yaml | kubectl apply -f - kubectl create namespace ${CANARY_NAMESPACE} --dry-run=client -o yaml | kubectl apply -f - + envsubst '${ENV_NAME}' < .github/k8s/sam-canary-sa-template.yaml | kubectl apply -f - - name: Provision Dex Secrets env: @@ -374,6 +375,12 @@ jobs: envsubst '${ENV_NAME} ${NAMESPACE} ${GCP_PROJECT_ID} ${CLUSTER_NAME} ${CLUSTER_REGION} ${IMAGE_TAG}' < .github/k8s/sam-router-template.yaml | kubectl apply -f - envsubst '${ENV_NAME} ${NAMESPACE} ${GCP_PROJECT_ID} ${CLUSTER_NAME} ${CLUSTER_REGION} ${IMAGE_TAG}' < .github/k8s/sam-console-template.yaml | kubectl apply -f - envsubst '${ENV_NAME} ${NAMESPACE} ${GCP_PROJECT_ID}' < .github/k8s/dns-sync-cronjob-template.yaml | kubectl apply -f - + # Managed Prometheus ships with GKE; a cluster without it should not block the mesh rollout. + if kubectl get crd podmonitorings.monitoring.googleapis.com >/dev/null 2>&1; then + envsubst '${ENV_NAME} ${NAMESPACE}' < .github/k8s/sam-monitoring-template.yaml | kubectl apply -f - + else + echo "::warning::PodMonitoring CRD not found; managed collection is off on this cluster, metrics will not be scraped" + fi kubectl rollout status deployment/sam-control-plane-${ENV_NAME} -n ${NAMESPACE} --timeout=120s || { echo "Control Plane Deployment failed!" @@ -529,48 +536,6 @@ jobs: --external-addr=/ip4/${EXTERNAL_IP}/tcp/4501 \ --external-addr=/ip4/${EXTERNAL_IP}/udp/4501/quic-v1' - - - name: Deploy COP Canary - env: - VAR_ENV_NAME: ${{ vars.ENV_NAME }} - VAR_IMAGE_TAG: ${{ env.IMAGE_TAG }} - run: | - print_rollout_diagnostics() { - local namespace="$1" - local deployment="$2" - local selector="$3" - - echo "Collecting diagnostics for deployment/${deployment} in namespace ${namespace}..." - - kubectl describe deployment/${deployment} -n ${namespace} || true - kubectl get pods -n ${namespace} -l "${selector}" -o wide || true - kubectl describe pods -n ${namespace} -l "${selector}" || true - kubectl get events -n ${namespace} --sort-by=.lastTimestamp || true - - for pod in $(kubectl get pods -n ${namespace} -l "${selector}" -o name 2>/dev/null); do - echo "==== Describe ${pod} ====" - kubectl describe -n ${namespace} "${pod}" || true - - for container in $(kubectl get -n ${namespace} "${pod}" -o jsonpath='{.spec.containers[*].name}' 2>/dev/null); do - echo "==== Logs for ${pod} container ${container} ====" - kubectl logs -n ${namespace} "${pod#pod/}" -c "${container}" --tail=-1 || true - done - done - } - - export ENV_NAME="${VAR_ENV_NAME}" - export CANARY_NAMESPACE="sam-canary-${ENV_NAME}" - export NAMESPACE="sam-${ENV_NAME}" - export IMAGE_TAG="${VAR_IMAGE_TAG}" - export BANANA_BOT_SCRIPT=$(cat site/content/docs/snippets/banana_bot_playground.py | sed 's/^/ /') - - envsubst '${ENV_NAME} ${NAMESPACE} ${IMAGE_TAG} ${BANANA_BOT_SCRIPT}' < .github/k8s/sam-node-cop-template.yaml | kubectl apply -f - - kubectl rollout status deployment/cop-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} --timeout=120s || { - echo "Cop Canary Deployment failed!" - print_rollout_diagnostics "${CANARY_NAMESPACE}" "cop-canary-${ENV_NAME}" "app=cop-canary-${ENV_NAME}" - exit 1 - } - - name: Deploy SAM Box Canary env: VAR_ENV_NAME: ${{ vars.ENV_NAME }} @@ -611,50 +576,6 @@ jobs: exit 1 } - - name: Deploy SAM Node Canary - env: - VAR_ENV_NAME: ${{ vars.ENV_NAME }} - VAR_IMAGE_TAG: ${{ env.IMAGE_TAG }} - run: | - print_rollout_diagnostics() { - local namespace="$1" - local deployment="$2" - local selector="$3" - - echo "Collecting diagnostics for deployment/${deployment} in namespace ${namespace}..." - - kubectl describe deployment/${deployment} -n ${namespace} || true - kubectl get pods -n ${namespace} -l "${selector}" -o wide || true - kubectl describe pods -n ${namespace} -l "${selector}" || true - kubectl get events -n ${namespace} --sort-by=.lastTimestamp || true - - for pod in $(kubectl get pods -n ${namespace} -l "${selector}" -o name 2>/dev/null); do - echo "==== Describe ${pod} ====" - kubectl describe -n ${namespace} "${pod}" || true - - for container in $(kubectl get -n ${namespace} "${pod}" -o jsonpath='{.spec.containers[*].name}' 2>/dev/null); do - echo "==== Logs for ${pod} container ${container} ====" - kubectl logs -n ${namespace} "${pod#pod/}" -c "${container}" --tail=-1 || true - done - done - } - - 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-node-template.yaml | kubectl apply -f - - kubectl rollout status deployment/sam-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} --timeout=120s || { - echo "Canary Deployment failed!" - print_rollout_diagnostics "${CANARY_NAMESPACE}" "sam-canary-${ENV_NAME}" "app=sam-canary-${ENV_NAME}" - - echo "Rolling back Canary deployment..." - kubectl rollout undo deployment/sam-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} - kubectl rollout status deployment/sam-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} || true - exit 1 - } - - name: Deploy OpenClaw Canary env: VAR_ENV_NAME: ${{ vars.ENV_NAME }} @@ -864,3 +785,50 @@ jobs: kubectl rollout status deployment/openrouter-canary-${ENV_NAME} -n ${CANARY_NAMESPACE} || true exit 1 } + + # A fresh node enrolls, reaches a router and calls the everything + # canary's MCP server through the mesh: what a new user does. The + # CronJob keeps doing it every 15 minutes; running it once here, after + # every canary is up, makes the rollout prove it before the deploy is + # called good. + - name: Deploy and run the 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-probe-rollout-${ENV_NAME}" + + envsubst '${ENV_NAME} ${NAMESPACE} ${IMAGE_TAG}' < .github/k8s/sam-probe-cronjob-template.yaml | kubectl apply -f - + + kubectl delete job "${JOB}" -n "${CANARY_NAMESPACE}" --ignore-not-found + kubectl create job --from="cronjob/sam-probe-${ENV_NAME}" "${JOB}" -n "${CANARY_NAMESPACE}" + + # `kubectl wait` can only wait for one condition; a failed Job would + # otherwise sit out the whole timeout. + for _ in $(seq 1 60); 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 + + echo "==== probe output ====" + kubectl logs "job/${JOB}" -n "${CANARY_NAMESPACE}" -c probe --tail=-1 || true + + case "${status:-}" in + *Complete=True*) echo "Cold-path probe passed." ;; + *) + echo "Cold-path probe did not pass: a fresh node could not enroll, reach a router, or call a mesh service." + echo "==== sam-node sidecar log ====" + kubectl logs "job/${JOB}" -n "${CANARY_NAMESPACE}" -c sam-node --tail=200 || true + kubectl describe job "${JOB}" -n "${CANARY_NAMESPACE}" || true + exit 1 + ;; + esac + diff --git a/.github/workflows/testnet-health.yaml b/.github/workflows/testnet-health.yaml new file mode 100644 index 00000000..10532b3f --- /dev/null +++ b/.github/workflows/testnet-health.yaml @@ -0,0 +1,181 @@ +name: Testnet Health + +# Looks at each public testnet every half hour and says whether a new user +# could join it right now. deploy.yaml only knows the moment of a rollout; +# this is what watches in between. A red badge on the README and an open +# "Testnet is unhealthy" issue are the signals; the issue closes itself +# on the next green run. + +on: + schedule: + - cron: '*/30 * * * *' + workflow_dispatch: + +permissions: + contents: read + +concurrency: + group: testnet-health + cancel-in-progress: false + +jobs: + check: + name: ${{ matrix.environment }} + runs-on: ubuntu-latest + # The same GitHub Environments deploy.yaml uses, for the same cluster + # credentials and per-testnet variables. + environment: ${{ matrix.environment }} + strategy: + fail-fast: false + matrix: + environment: [hub, bananas] + permissions: + contents: read + id-token: write + issues: write + timeout-minutes: 15 + + steps: + - name: Checkout code + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + with: + persist-credentials: false + + - name: Google Auth + uses: google-github-actions/auth@c200f3691d83b41bf9bbd8638997a462592937ed # v2 + with: + workload_identity_provider: ${{ vars.WIF_PROVIDER_NAME }} + service_account: ${{ vars.SERVICE_ACCOUNT_EMAIL }} + + - name: Set up GKE credentials + uses: google-github-actions/get-gke-credentials@3da1e46a907576cefaa90c484278bb5b259dd395 # v3.0.0 + with: + cluster_name: ${{ vars.CLUSTER_NAME }} + location: ${{ vars.CLUSTER_REGION }} + + - name: Check the testnet + id: check + # Every check runs even after one fails, so the report is complete; + # the report step turns the outcome into the issue and the exit code. + continue-on-error: true + env: + ENV_NAME: ${{ vars.ENV_NAME }} + run: | + set -uo pipefail + NS="sam-${ENV_NAME}" + CNS="sam-canary-${ENV_NAME}" + HOST="${ENV_NAME}.sam-mesh.dev" + rows=() + fails=0 + ok() { rows+=("| :white_check_mark: | $1 | $2 |"); } + bad() { rows+=("| :x: | $1 | $2 |"); fails=$((fails + 1)); } + + # 1. The public enrollment surface, as a new node sees it. /info is + # protobuf; router peer ids are ASCII inside it. + code=$(curl -s -m 20 -o info.bin -w '%{http_code}' "https://${HOST}/info" || echo 000) + routers_public=$(grep -aoE '/p2p/12D3Koo[A-Za-z0-9]+' info.bin 2>/dev/null | sort -u | wc -l) + if [[ "${code}" == "200" && "${routers_public}" -ge 1 ]]; then + ok "Public \`/info\`" "HTTP 200, ${routers_public} router(s) advertised" + else + bad "Public \`/info\`" "HTTP ${code}, ${routers_public} router(s) advertised" + fi + + # 2. What the control plane knows, read through the API server so no + # port is opened to the internet for it. + if metrics=$(kubectl get --raw "/api/v1/namespaces/${NS}/services/sam-control-plane-${ENV_NAME}:8080/proxy/metrics" 2>&1); then + routers_active=$(awk '$1=="sam_control_plane_routers_active"{print int($2)}' <<<"${metrics}") + peers=$(awk '$1=="sam_control_plane_mesh_connected_peers"{print int($2)}' <<<"${metrics}") + scrape=$(awk '$1=="sam_control_plane_mesh_state_scrape_success"{print int($2)}' <<<"${metrics}") + if [[ "${scrape:-0}" -eq 1 && "${routers_active:-0}" -ge 2 ]]; then + ok "Routers holding a lease" "${routers_active}" + else + bad "Routers holding a lease" "${routers_active:-?} (store scrape success=${scrape:-?}; want >= 2)" + fi + if [[ "${peers:-0}" -ge 1 ]]; then + ok "Peers attached to the mesh" "${peers}" + else + bad "Peers attached to the mesh" "${peers:-?} (want >= 1)" + fi + else + bad "Control plane metrics" "$(head -c 200 <<<"${metrics}")" + fi + + # 3. Routers and canaries are fully rolled out. + sts=$(kubectl get statefulset "sam-router-${ENV_NAME}" -n "${NS}" -o jsonpath='{.status.readyReplicas}/{.spec.replicas}' 2>&1 || echo "?/?") + if [[ "${sts%%/*}" == "${sts##*/}" && "${sts%%/*}" != "?" && "${sts%%/*}" != "" ]]; then + ok "Router StatefulSet" "${sts} ready" + else + bad "Router StatefulSet" "${sts} ready" + fi + + unhealthy="" + total=0 + while read -r name ready want; do + [[ -z "${name}" ]] && continue + total=$((total + 1)) + [[ "${ready:-0}" == "${want}" ]] || unhealthy="${unhealthy} ${name}(${ready:-0}/${want})" + done < <(kubectl get deployments -n "${CNS}" -o jsonpath='{range .items[*]}{.metadata.name} {.status.availableReplicas} {.spec.replicas}{"\n"}{end}' 2>/dev/null) + if [[ "${total}" -gt 0 && -z "${unhealthy}" ]]; then + ok "Canaries" "${total} deployments fully available" + else + bad "Canaries" "${total} deployments; not available:${unhealthy:- (none found)}" + fi + + # 4. The cold path: the probe CronJob enrolled a fresh node, found a + # service and spoke MCP to it recently. Three schedules of slack. + last=$(kubectl get cronjob "sam-probe-${ENV_NAME}" -n "${CNS}" -o jsonpath='{.status.lastSuccessfulTime}' 2>/dev/null || true) + if [[ -n "${last}" ]]; then + age=$(( $(date +%s) - $(date -d "${last}" +%s) )) + if [[ "${age}" -le 2700 ]]; then + ok "Cold-path probe" "last success ${age}s ago" + else + bad "Cold-path probe" "last success ${age}s ago (want <= 2700s)" + fi + else + bad "Cold-path probe" "no successful run recorded" + fi + + { + echo "## ${ENV_NAME}.sam-mesh.dev" + echo + echo "| | Check | Result |" + echo "|---|---|---|" + printf '%s\n' "${rows[@]}" + } > report.md + cat report.md >> "${GITHUB_STEP_SUMMARY}" + cat report.md + [[ "${fails}" -eq 0 ]] + + - name: Report + if: always() + env: + GH_TOKEN: ${{ github.token }} + ENV_NAME: ${{ vars.ENV_NAME }} + OUTCOME: ${{ steps.check.outcome }} + RUN_URL: ${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }} + run: | + set -euo pipefail + title="Testnet ${ENV_NAME} is unhealthy" + label="testnet-health" + existing=$(gh issue list --state open --label "${label}" --json number,title \ + --jq ".[] | select(.title == \"${title}\") | .number" | head -n1 || true) + + if [[ "${OUTCOME}" == "success" ]]; then + if [[ -n "${existing}" ]]; then + gh issue close "${existing}" --comment "Healthy again: ${RUN_URL}" + fi + exit 0 + fi + + report=$(cat report.md 2>/dev/null || echo "The check step did not produce a report.") + body=$(printf '%s\n\nRun: %s\n_Updated by every failing run; closed automatically by the next healthy one._\n' "${report}" "${RUN_URL}") + + if [[ -n "${existing}" ]]; then + # One open issue per testnet, kept current rather than commented on + # every half hour. + gh issue edit "${existing}" --body "${body}" + else + gh label create "${label}" --description "Opened by the Testnet Health workflow" --color B60205 2>/dev/null || true + gh issue create --title "${title}" --label "${label}" --body "${body}" + fi + exit 1 diff --git a/README.md b/README.md index ec634218..7a581bf2 100644 --- a/README.md +++ b/README.md @@ -40,6 +40,8 @@ SAM provides the open protocols, cryptographic building blocks, and software to > **About the Public Developer Testnets:** > The public endpoints (`bananas.sam-mesh.dev` and `hub.sam-mesh.dev`) are free testbeds created using community resources solely for developer testing, continuous integration, and rapid experimentation. They provide **no guarantees, no uptime commitments, zero SLA, and no sovereign guarantees**. Running on a shared community testnet delegates identity management to the testbed maintainers; true sovereignty requires deploying a dedicated control plane with customer-held keys. > +> [![Testnet Health](https://github.com/google/sam/actions/workflows/testnet-health.yaml/badge.svg)](https://github.com/google/sam/actions/workflows/testnet-health.yaml) Every half hour a fresh node's view of each testnet is checked (enrollment surface, routers, canaries, and a cold-path probe that joins and calls a tool); a red badge means one of them is unhealthy and an issue labelled `testnet-health` says which check failed. +> > 📖 **Deep Dive:** Read our full **[Digital & Data Sovereignty Architecture](site/content/docs/sovereignty.md)** covering the 5 pillars, fail-closed label gates, uncooperative sandbox confinement, and regulatory alignment (GDPR Chapter V, EU Cloud Sovereignty Framework SEAL-3, EU Data Act). --- diff --git a/charts/sam-mesh/README.md b/charts/sam-mesh/README.md index f8c67b84..807a2c79 100644 --- a/charts/sam-mesh/README.md +++ b/charts/sam-mesh/README.md @@ -117,3 +117,26 @@ gateway: There is no bundled Dex. Point `controlPlane.oidcIssuer` at your identity provider and register `https:///auth/callback` as a redirect URI for the OIDC client the control plane reports. + +## Metrics (`monitoring.*`) + +The control plane serves Prometheus metrics on its `http` port and the router +on a dedicated `metrics` containerPort (`router.metricsPort`, default 9090), +which also carries its `/healthz` and `/readyz` probes. Neither endpoint is +authenticated, so the router port is deliberately never a hostPort or a +Service. + +Scraping is off by default because both supported resources are CRDs: + +- `monitoring.podMonitoring.enabled: true` renders a + `monitoring.googleapis.com/v1` `PodMonitoring` per component for Google + Managed Prometheus, which GKE runs out of the box. +- `monitoring.podMonitor.enabled: true` renders a `monitoring.coreos.com/v1` + `PodMonitor` per component for prometheus-operator; use + `monitoring.podMonitor.labels` for the selector your Prometheus matches on. + +Mesh size and state come from the control plane +(`sam_control_plane_mesh_connected_peers`, `sam_control_plane_routers_active`, +`sam_control_plane_enrolled_nodes{role,state}`); the router adds its live view +(`sam_router_authenticated_peers`, `sam_router_auth_handshakes_total{result}`) +alongside libp2p's own relay and connection metrics. diff --git a/charts/sam-mesh/templates/monitoring.yaml b/charts/sam-mesh/templates/monitoring.yaml new file mode 100644 index 00000000..455b0c54 --- /dev/null +++ b/charts/sam-mesh/templates/monitoring.yaml @@ -0,0 +1,48 @@ +{{- /* +One scrape target per component, rendered for whichever collector the +cluster runs. Pod-selecting resources are used for both so the router's +metrics port never has to appear on a Service. +*/ -}} +{{- $targets := list (dict "name" "control-plane" "port" "http") -}} +{{- if .Values.router.enabled -}} +{{- $targets = append $targets (dict "name" "router" "port" "metrics") -}} +{{- end -}} +{{- range $t := $targets }} +{{- if $.Values.monitoring.podMonitoring.enabled }} +--- +apiVersion: monitoring.googleapis.com/v1 +kind: PodMonitoring +metadata: + name: {{ include "sam-mesh.fullname" $ }}-{{ $t.name }} + labels: + {{- include "sam-mesh.labels" $ | nindent 4 }} +spec: + selector: + matchLabels: + app: {{ include "sam-mesh.fullname" $ }}-{{ $t.name }} + endpoints: + - port: {{ $t.port }} + path: /metrics + interval: {{ $.Values.monitoring.interval }} +{{- end }} +{{- if $.Values.monitoring.podMonitor.enabled }} +--- +apiVersion: monitoring.coreos.com/v1 +kind: PodMonitor +metadata: + name: {{ include "sam-mesh.fullname" $ }}-{{ $t.name }} + labels: + {{- include "sam-mesh.labels" $ | nindent 4 }} + {{- with $.Values.monitoring.podMonitor.labels }} + {{- toYaml . | nindent 4 }} + {{- end }} +spec: + selector: + matchLabels: + app: {{ include "sam-mesh.fullname" $ }}-{{ $t.name }} + podMetricsEndpoints: + - port: {{ $t.port }} + path: /metrics + interval: {{ $.Values.monitoring.interval }} +{{- end }} +{{- end }} diff --git a/charts/sam-mesh/templates/router-statefulset.yaml b/charts/sam-mesh/templates/router-statefulset.yaml index ff7c9e88..532fda1d 100644 --- a/charts/sam-mesh/templates/router-statefulset.yaml +++ b/charts/sam-mesh/templates/router-statefulset.yaml @@ -87,13 +87,27 @@ spec: hostPort: {{ .Values.router.hostPort }} {{- end }} protocol: UDP + # Cluster-internal: never given a hostPort, nothing on it is authenticated. + - containerPort: {{ .Values.router.metricsPort }} + name: metrics + protocol: TCP + # /healthz answers as soon as the process is up; /readyz once the + # router has enrolled and its libp2p host is online. + startupProbe: + httpGet: + path: /healthz + port: metrics + periodSeconds: 5 + failureThreshold: 24 readinessProbe: - tcpSocket: - port: p2p-tcp + httpGet: + path: /readyz + port: metrics periodSeconds: 5 livenessProbe: - tcpSocket: - port: p2p-tcp + httpGet: + path: /healthz + port: metrics periodSeconds: 15 resources: {{- toYaml .Values.router.resources | nindent 10 }} @@ -112,6 +126,7 @@ spec: - "--insecure-control-plane" - "--listen=/ip4/0.0.0.0/tcp/4501" - "--listen=/ip4/0.0.0.0/udp/4501/quic-v1" + - "--metrics-addr=0.0.0.0:{{ .Values.router.metricsPort }}" {{- if .Values.router.externalAddrs }} {{- range .Values.router.externalAddrs }} - "--external-addr={{ . }}" diff --git a/charts/sam-mesh/tests/monitoring_test.yaml b/charts/sam-mesh/tests/monitoring_test.yaml new file mode 100644 index 00000000..90b5358e --- /dev/null +++ b/charts/sam-mesh/tests/monitoring_test.yaml @@ -0,0 +1,84 @@ +suite: monitoring +templates: + - templates/monitoring.yaml +release: + # collapses fullname to "sam-mesh" (see _helpers.tpl) + name: sam-mesh +tests: + - it: renders nothing by default, since both kinds are CRDs the cluster may lack + asserts: + - hasDocuments: + count: 0 + + - it: podMonitoring scrapes the control plane and router on their own ports + set: + monitoring.podMonitoring.enabled: true + asserts: + - hasDocuments: + count: 2 + - isKind: + of: PodMonitoring + - isAPIVersion: + of: monitoring.googleapis.com/v1 + - equal: + path: spec.selector.matchLabels.app + value: sam-mesh-control-plane + documentIndex: 0 + - equal: + path: spec.endpoints[0].port + value: http + documentIndex: 0 + - equal: + path: spec.selector.matchLabels.app + value: sam-mesh-router + documentIndex: 1 + - equal: + path: spec.endpoints[0].port + value: metrics + documentIndex: 1 + - equal: + path: spec.endpoints[0].interval + value: 30s + documentIndex: 1 + + - it: podMonitor renders the prometheus-operator shape with extra labels + set: + monitoring.podMonitor.enabled: true + monitoring.podMonitor.labels: + release: kube-prometheus-stack + monitoring.interval: 15s + asserts: + - hasDocuments: + count: 2 + - isKind: + of: PodMonitor + - isAPIVersion: + of: monitoring.coreos.com/v1 + - equal: + path: metadata.labels.release + value: kube-prometheus-stack + - equal: + path: spec.podMetricsEndpoints[0].path + value: /metrics + - equal: + path: spec.podMetricsEndpoints[0].interval + value: 15s + + - it: skips the router target when the router is disabled + set: + monitoring.podMonitoring.enabled: true + router.enabled: false + asserts: + - hasDocuments: + count: 1 + - equal: + path: spec.selector.matchLabels.app + value: sam-mesh-control-plane + + - it: both collectors can be enabled at once + set: + monitoring.podMonitoring.enabled: true + monitoring.podMonitor.enabled: true + asserts: + - hasDocuments: + count: 4 diff --git a/charts/sam-mesh/tests/router-statefulset_test.yaml b/charts/sam-mesh/tests/router-statefulset_test.yaml index 72636ddc..e38f9a74 100644 --- a/charts/sam-mesh/tests/router-statefulset_test.yaml +++ b/charts/sam-mesh/tests/router-statefulset_test.yaml @@ -154,3 +154,48 @@ tests: - equal: path: spec.template.spec.containers[0].image value: sam-router:v9 + + - it: serves metrics on a containerPort that probes use and no hostPort reaches + set: + router.hostPort: 4501 + documentSelector: + path: kind + value: StatefulSet + asserts: + - contains: + path: spec.template.spec.containers[0].args + content: --metrics-addr=0.0.0.0:9090 + - equal: + path: spec.template.spec.containers[0].ports[2].name + value: metrics + - equal: + path: spec.template.spec.containers[0].ports[2].containerPort + value: 9090 + - notExists: + path: spec.template.spec.containers[0].ports[2].hostPort + - equal: + path: spec.template.spec.containers[0].readinessProbe.httpGet.path + value: /readyz + - equal: + path: spec.template.spec.containers[0].readinessProbe.httpGet.port + value: metrics + - equal: + path: spec.template.spec.containers[0].livenessProbe.httpGet.path + value: /healthz + - equal: + path: spec.template.spec.containers[0].startupProbe.httpGet.path + value: /healthz + + - it: metricsPort moves the listener, the port and the probes together + set: + router.metricsPort: 9999 + documentSelector: + path: kind + value: StatefulSet + asserts: + - contains: + path: spec.template.spec.containers[0].args + content: --metrics-addr=0.0.0.0:9999 + - equal: + path: spec.template.spec.containers[0].ports[2].containerPort + value: 9999 diff --git a/charts/sam-mesh/values.yaml b/charts/sam-mesh/values.yaml index c67b9d53..9a013dfa 100644 --- a/charts/sam-mesh/values.yaml +++ b/charts/sam-mesh/values.yaml @@ -143,6 +143,9 @@ router: # Overrides the announced multiaddrs; derived from hostPort when left empty. externalAddrs: [] allowLoopback: false + # Plain HTTP port for /metrics, /healthz and /readyz. Unauthenticated, so it + # is only ever a containerPort, never a hostPort or Service. + metricsPort: 9090 # PVC for /data/router.key, so the libp2p peer ID survives rescheduling. storageSize: 1Gi # PVC StorageClass; cluster default when null. @@ -213,3 +216,18 @@ bootstrap: # labels in its sam-node.yaml. Empty means nodes cannot enroll with any label # (fail closed). nodeLabels: [] + +# Scraping of the control-plane and router /metrics endpoints. Both are +# cluster-internal and unauthenticated; enable the resource for whichever +# collector the cluster runs, since each is a CRD that may not be installed. +monitoring: + interval: 30s + # Google Managed Prometheus (monitoring.googleapis.com/v1 PodMonitoring), + # present by default on GKE. + podMonitoring: + enabled: false + # prometheus-operator (monitoring.coreos.com/v1 PodMonitor). + podMonitor: + enabled: false + # Extra labels, e.g. the release selector a Prometheus instance matches on. + labels: {} diff --git a/charts/sam-node/templates/deployment.yaml b/charts/sam-node/templates/deployment.yaml index c314ae61..1868e9dc 100644 --- a/charts/sam-node/templates/deployment.yaml +++ b/charts/sam-node/templates/deployment.yaml @@ -46,9 +46,37 @@ spec: - "--jwt-path=/var/run/secrets/tokens/sam-token" - "--api-token-path=/var/run/secrets/sam/api-token" - "--bind-addr={{ .Values.bindAddr }}" + {{- if .Values.metricsPort }} + - "--metrics-addr=0.0.0.0:{{ .Values.metricsPort }}" + {{- end }} {{- range .Values.extraArgs }} - {{ . | quote }} {{- end }} + {{- if .Values.metricsPort }} + ports: + # Cluster-internal and unauthenticated; the API token is never needed here. + - containerPort: {{ .Values.metricsPort }} + name: metrics + protocol: TCP + # /healthz answers as soon as the process is up; /readyz once the node + # holds an authenticated connection to a router. + startupProbe: + httpGet: + path: /healthz + port: metrics + periodSeconds: 5 + failureThreshold: 24 + readinessProbe: + httpGet: + path: /readyz + port: metrics + periodSeconds: 10 + livenessProbe: + httpGet: + path: /healthz + port: metrics + periodSeconds: 20 + {{- end }} {{- with .Values.resources }} resources: {{- toYaml . | nindent 12 }} diff --git a/charts/sam-node/templates/monitoring.yaml b/charts/sam-node/templates/monitoring.yaml new file mode 100644 index 00000000..02bbc13c --- /dev/null +++ b/charts/sam-node/templates/monitoring.yaml @@ -0,0 +1,39 @@ +{{- if .Values.metricsPort }} +{{- if .Values.monitoring.podMonitoring.enabled }} +--- +apiVersion: monitoring.googleapis.com/v1 +kind: PodMonitoring +metadata: + name: {{ include "sam-node.fullname" . }} + labels: + {{- include "sam-node.labels" . | nindent 4 }} +spec: + selector: + matchLabels: + {{- include "sam-node.selectorLabels" . | nindent 6 }} + endpoints: + - port: metrics + path: /metrics + interval: {{ .Values.monitoring.interval }} +{{- end }} +{{- if .Values.monitoring.podMonitor.enabled }} +--- +apiVersion: monitoring.coreos.com/v1 +kind: PodMonitor +metadata: + name: {{ include "sam-node.fullname" . }} + labels: + {{- include "sam-node.labels" . | nindent 4 }} + {{- with .Values.monitoring.podMonitor.labels }} + {{- toYaml . | nindent 4 }} + {{- end }} +spec: + selector: + matchLabels: + {{- include "sam-node.selectorLabels" . | nindent 6 }} + podMetricsEndpoints: + - port: metrics + path: /metrics + interval: {{ .Values.monitoring.interval }} +{{- end }} +{{- end }} diff --git a/charts/sam-node/tests/deployment_test.yaml b/charts/sam-node/tests/deployment_test.yaml index 62289eb7..cd8a4ceb 100644 --- a/charts/sam-node/tests/deployment_test.yaml +++ b/charts/sam-node/tests/deployment_test.yaml @@ -125,3 +125,48 @@ tests: - equal: path: spec.template.spec.serviceAccountName value: existing-sa + + - it: serves metrics on a containerPort the probes use, apart from the API + template: templates/deployment.yaml + set: + controlPlaneUrl: http://sam-mesh-control-plane:8080 + asserts: + - contains: + path: spec.template.spec.containers[0].args + content: "--metrics-addr=0.0.0.0:9090" + - equal: + path: spec.template.spec.containers[0].ports[0].name + value: metrics + - equal: + path: spec.template.spec.containers[0].ports[0].containerPort + value: 9090 + - equal: + path: spec.template.spec.containers[0].readinessProbe.httpGet.path + value: /readyz + - equal: + path: spec.template.spec.containers[0].readinessProbe.httpGet.port + value: metrics + - equal: + path: spec.template.spec.containers[0].livenessProbe.httpGet.path + value: /healthz + - equal: + path: spec.template.spec.containers[0].startupProbe.httpGet.path + value: /healthz + + - it: metricsPort 0 drops the listener, the port and the probes together + template: templates/deployment.yaml + set: + controlPlaneUrl: http://sam-mesh-control-plane:8080 + metricsPort: 0 + asserts: + - notContains: + path: spec.template.spec.containers[0].args + content: "--metrics-addr=0.0.0.0:9090" + - notExists: + path: spec.template.spec.containers[0].ports + - notExists: + path: spec.template.spec.containers[0].readinessProbe + - notExists: + path: spec.template.spec.containers[0].livenessProbe + - notExists: + path: spec.template.spec.containers[0].startupProbe diff --git a/charts/sam-node/tests/monitoring_test.yaml b/charts/sam-node/tests/monitoring_test.yaml new file mode 100644 index 00000000..09b116ae --- /dev/null +++ b/charts/sam-node/tests/monitoring_test.yaml @@ -0,0 +1,61 @@ +suite: sam-node monitoring +templates: + - templates/monitoring.yaml +release: + name: sam-node +tests: + - it: renders nothing by default, since both kinds are CRDs the cluster may lack + set: + controlPlaneUrl: http://sam-mesh-control-plane:8080 + asserts: + - hasDocuments: + count: 0 + + - it: podMonitoring selects the node pods on the metrics port + set: + controlPlaneUrl: http://sam-mesh-control-plane:8080 + monitoring.podMonitoring.enabled: true + asserts: + - hasDocuments: + count: 1 + - isKind: + of: PodMonitoring + - isAPIVersion: + of: monitoring.googleapis.com/v1 + - equal: + path: spec.selector.matchLabels["app.kubernetes.io/name"] + value: sam-node + - equal: + path: spec.endpoints[0].port + value: metrics + - equal: + path: spec.endpoints[0].interval + value: 30s + + - it: podMonitor renders the prometheus-operator shape with extra labels + set: + controlPlaneUrl: http://sam-mesh-control-plane:8080 + monitoring.podMonitor.enabled: true + monitoring.podMonitor.labels: + release: kube-prometheus-stack + asserts: + - hasDocuments: + count: 1 + - isKind: + of: PodMonitor + - equal: + path: metadata.labels.release + value: kube-prometheus-stack + - equal: + path: spec.podMetricsEndpoints[0].port + value: metrics + + - it: renders nothing when the metrics listener is disabled + set: + controlPlaneUrl: http://sam-mesh-control-plane:8080 + metricsPort: 0 + monitoring.podMonitoring.enabled: true + monitoring.podMonitor.enabled: true + asserts: + - hasDocuments: + count: 0 diff --git a/charts/sam-node/values.yaml b/charts/sam-node/values.yaml index 9f8510c8..64d3d075 100644 --- a/charts/sam-node/values.yaml +++ b/charts/sam-node/values.yaml @@ -43,6 +43,25 @@ securityContext: # pod-private; set 0.0.0.0:8080 to expose it on the pod IP. bindAddr: "127.0.0.1:8080" +# Plain HTTP port for /metrics, /healthz and /readyz, separate from bindAddr +# so neither the scraper nor the kubelet needs the API token. Unauthenticated, +# so it is only ever a containerPort. Set to 0 to disable it and the probes. +metricsPort: 9090 + +# Scraping of the metrics port. Each is a CRD the cluster may not have, so +# both are off until you pick the collector you run. +monitoring: + interval: 30s + # Google Managed Prometheus (monitoring.googleapis.com/v1 PodMonitoring), + # present by default on GKE. + podMonitoring: + enabled: false + # prometheus-operator (monitoring.coreos.com/v1 PodMonitor). + podMonitor: + enabled: false + # Extra labels, e.g. the release selector a Prometheus instance matches on. + labels: {} + # Extra sam-node args, e.g. ["--discovery-interval=200ms"]. extraArgs: [] diff --git a/cmd/sam-control-plane/main.go b/cmd/sam-control-plane/main.go index 62b281d9..f9c7ae7b 100644 --- a/cmd/sam-control-plane/main.go +++ b/cmd/sam-control-plane/main.go @@ -45,6 +45,7 @@ var ( leaseDuration time.Duration biscuitTTL time.Duration oidcSessionTTL time.Duration + nodeRetention time.Duration adminTokenPath string insecureSkipTLSVerify bool logLevel string @@ -126,6 +127,7 @@ func main() { BiscuitTimeout: 10 * time.Second, BiscuitTTL: biscuitTTL, OIDCSessionTTL: oidcSessionTTL, + NodeRetention: nodeRetention, AdminToken: adminToken, AutoApproveEnrollment: autoApproveEnrollment, } @@ -161,6 +163,7 @@ func main() { rootCmd.Flags().DurationVar(&leaseDuration, "lease-duration", 15*time.Minute, "Router lease registration TTL.") rootCmd.Flags().DurationVar(&biscuitTTL, "biscuit-ttl", api.BiscuitTokenTTL, "Lifespan minted into every issued Biscuit's expiration fact. Capped to the OIDC token's own expiry when shorter.") rootCmd.Flags().DurationVar(&oidcSessionTTL, "oidc-session-ttl", api.OIDCSessionTTL, "How long an OIDC enrollment stays refreshable before the identity must re-authenticate with the OIDC provider. Shorter values keep the provider authoritative for offboarding at the cost of more frequent interactive re-enrollment.") + rootCmd.Flags().DurationVar(&nodeRetention, "node-retention", controlplane.DefaultNodeRetention, "How long an enrolled node's record is kept after its session expired before it is deleted. Banned nodes are always kept. 0 keeps every record forever.") rootCmd.Flags().StringVar(&adminTokenPath, "admin-token-path", "", "Path to file containing the token for authenticating policy REST API requests (or env SAM_ADMIN_TOKEN)") rootCmd.Flags().BoolVar(&insecureSkipTLSVerify, "insecure-skip-tls-verify", false, "Skip TLS verification for OIDC providers") rootCmd.Flags().StringVar(&logLevel, "log-level", "info", "Log level (debug, info, warn, error)") diff --git a/cmd/sam-node/main.go b/cmd/sam-node/main.go index 318e3cdf..9bed6b41 100644 --- a/cmd/sam-node/main.go +++ b/cmd/sam-node/main.go @@ -62,6 +62,7 @@ var ( clientSecretFlag string controlPlanePublicKeyFlag string bindAddrFlag string + metricsAddrFlag string socketPathFlag string meshFlag string discoveryIntervalFlag string @@ -644,6 +645,14 @@ func main() { } }() + if metricsAddrFlag != "" { + metricsSrv, err := node.StartMetricsServer(meshNode, metricsAddrFlag) + if err != nil { + logger.Fatalf("Failed to start metrics server: %v", err) + } + defer func() { _ = metricsSrv.Close() }() + } + fmt.Printf("SAM Node Online.\nPeerID: %s\nListening on: %v\n", meshNode.Host.ID(), meshNode.Host.Addrs()) // Block forever @@ -836,6 +845,7 @@ func main() { runCmd.Flags().StringVar(&controlPlanePublicKeyFlag, "control-plane-public-key", "", "Control plane public key (32-byte Hex)") runCmd.Flags().StringVar(&bindAddrFlag, "bind-addr", "127.0.0.1:8080", "Local TCP address for the HTTP server (MCP and Sidecar API); pass an empty value to serve only on the Unix socket") runCmd.Flags().StringVar(&socketPathFlag, "socket-path", "", "Unix socket serving the same API, where the socket's owner-only permissions replace the API token (defaults to /"+node.DefaultSocketName+"; pass an empty value to disable)") + runCmd.Flags().StringVar(&metricsAddrFlag, "metrics-addr", "", "Serve Prometheus /metrics, /healthz and /readyz on this address without authentication (e.g. 0.0.0.0:9090), for scrapers and probes that hold no API token; off by default") runCmd.Flags().StringVar(&meshFlag, "mesh", node.DefaultMeshName, "Mesh federation name") runCmd.Flags().StringVar(&discoveryIntervalFlag, "discovery-interval", node.DefaultDiscoveryInterval, "Polling interval for DHT discovery") runCmd.Flags().DurationVar(&monitorBootstrapFlag, "monitor-bootstrap", 2*time.Minute, "Initial wait before monitoring router connection") diff --git a/cmd/sam-router/main.go b/cmd/sam-router/main.go index ff8ca0a3..da37ad63 100644 --- a/cmd/sam-router/main.go +++ b/cmd/sam-router/main.go @@ -45,6 +45,7 @@ var ( dhtMaxRecordAge time.Duration lowWaterMark int highWaterMark int + metricsAddr string ) var logger = golog.Logger("sam-router-cli") @@ -84,6 +85,7 @@ func main() { DHTMaxRecordAge: dhtMaxRecordAge, LowWaterMark: lowWaterMark, HighWaterMark: highWaterMark, + MetricsAddr: metricsAddr, } r, err := router.NewRouter(cmd.Context(), opts) @@ -122,6 +124,7 @@ func main() { rootCmd.Flags().DurationVar(&dhtMaxRecordAge, "dht-max-record-age", 0, "Maximum age for DHT records (0s uses library default)") rootCmd.Flags().IntVar(&lowWaterMark, "low-watermark", 1000, "Connection manager low watermark limit") rootCmd.Flags().IntVar(&highWaterMark, "high-watermark", 4000, "Connection manager high watermark limit") + rootCmd.Flags().StringVar(&metricsAddr, "metrics-addr", "", "Serve Prometheus /metrics, /healthz and /readyz on this address (e.g. 0.0.0.0:9090); unauthenticated, off by default") ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer cancel() diff --git a/internal/controlplane/config.go b/internal/controlplane/config.go index b090a47a..20ad86bf 100644 --- a/internal/controlplane/config.go +++ b/internal/controlplane/config.go @@ -36,8 +36,13 @@ type Options struct { BiscuitTimeout time.Duration BiscuitTTL time.Duration // Lifespan minted into every issued Biscuit's expiration() fact; defaults to api.BiscuitTokenTTL OIDCSessionTTL time.Duration // How long an OIDC enrollment stays refreshable before the identity must re-authenticate interactively; defaults to api.OIDCSessionTTL - AdminToken string // Optional: administrative bearer token for protecting policy and enrollment queue REST APIs - AutoApproveEnrollment bool // If true, valid bootstrap token enrollment requests are immediately approved without administrative manual gate + // NodeRetention is how long an enrolled node's row is kept after its + // session expired before being deleted; 0 keeps rows forever. Every + // pod restart without a persistent data dir enrolls a fresh identity, + // so without this the nodes table only ever grows. + NodeRetention time.Duration + AdminToken string // Optional: administrative bearer token for protecting policy and enrollment queue REST APIs + AutoApproveEnrollment bool // If true, valid bootstrap token enrollment requests are immediately approved without administrative manual gate } // Default sets default values for control plane options. diff --git a/internal/controlplane/metrics.go b/internal/controlplane/metrics.go new file mode 100644 index 00000000..cce00825 --- /dev/null +++ b/internal/controlplane/metrics.go @@ -0,0 +1,375 @@ +// 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. + +package controlplane + +import ( + "context" + "net/http" + "strconv" + "strings" + "sync" + "time" + + "github.com/google/sam/api" + "github.com/google/sam/internal/storage" + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +var ( + httpRequestsTotal = promauto.NewCounterVec( + prometheus.CounterOpts{ + Name: "sam_control_plane_http_requests_total", + Help: "Control-plane HTTP requests by registered route and status code", + }, + []string{"route", "code"}, + ) + + httpRequestDurationSeconds = promauto.NewHistogramVec( + prometheus.HistogramOpts{ + Name: "sam_control_plane_http_request_duration_seconds", + Help: "Time a control-plane request occupied its handler", + Buckets: []float64{0.001, 0.0025, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10}, + }, + []string{"route"}, + ) +) + +// observeRoute counts and times requests to one registered mux pattern. The +// pattern, not the request path, is the label: paths carry peer ids and token +// ids chosen off-plane, so the raw path can never become a label. +func observeRoute(pattern string, h http.HandlerFunc) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + start := time.Now() + rec := &statusWriter{ResponseWriter: w} + h(rec, r) + httpRequestsTotal.WithLabelValues(pattern, strconv.Itoa(rec.status())).Inc() + httpRequestDurationSeconds.WithLabelValues(pattern).Observe(time.Since(start).Seconds()) + }) +} + +type statusWriter struct { + http.ResponseWriter + code int +} + +func (s *statusWriter) Unwrap() http.ResponseWriter { return s.ResponseWriter } + +func (s *statusWriter) status() int { + if s.code == 0 { + return http.StatusOK + } + return s.code +} + +func (s *statusWriter) WriteHeader(code int) { + if s.code == 0 { + s.code = code + } + s.ResponseWriter.WriteHeader(code) +} + +const ( + // meshStateTTL bounds how often a scrape may hit the store. Router leases + // renew on the order of minutes, so anything fresher is noise. + meshStateTTL = 30 * time.Second + // meshStateTimeout caps one refresh so a slow store cannot hold a scrape + // past the scraper's own deadline. + meshStateTimeout = 5 * time.Second +) + +// Node admission states as reported by meshStateCollector. "admitted" is +// what CheckAdmission says, not liveness: a node that enrolled once and went +// away stays admitted until its session expires. Liveness is +// sam_control_plane_mesh_connected_peers. +const ( + nodeStateAdmitted = "admitted" + nodeStateExpired = "expired" + nodeStateBanned = "banned" +) + +// Bootstrap token states as reported by meshStateCollector. +const ( + tokenStateActive = "active" + tokenStateRevoked = "revoked" + tokenStateExpired = "expired" + tokenStateExhausted = "exhausted" +) + +// meshSnapshot is what one refresh of the store reduces to. Only counts and +// per-router figures survive: a node's peer id is never a label. +type meshSnapshot struct { + nodesByRoleState map[[2]string]int + users int + enrollmentRequests map[string]int + tokensByState map[string]int + routers []storage.RouterLease + meshConnectedPeers int +} + +// meshStateCollector exports the control plane's view of the mesh: who is +// enrolled, which routers hold a lease and what they report attached to them. +// The store is the source of truth shared by every replica, so each replica +// exports the same figures and a dashboard may take any one of them. +type meshStateCollector struct { + store storage.Store + ttl time.Duration + timeout time.Duration + now func() time.Time + + mu sync.Mutex + fetchedAt time.Time + snapshot meshSnapshot + ok bool + + nodesDesc *prometheus.Desc + usersDesc *prometheus.Desc + requestsDesc *prometheus.Desc + tokensDesc *prometheus.Desc + routersDesc *prometheus.Desc + routerPeersDesc *prometheus.Desc + routerDHTDesc *prometheus.Desc + routerLeaseDesc *prometheus.Desc + meshPeersDesc *prometheus.Desc + scrapeOKDesc *prometheus.Desc + scrapeSampleDesc *prometheus.Desc +} + +func newMeshStateCollector(store storage.Store) *meshStateCollector { + return &meshStateCollector{ + store: store, + ttl: meshStateTTL, + timeout: meshStateTimeout, + now: time.Now, + nodesDesc: prometheus.NewDesc( + "sam_control_plane_enrolled_nodes", + "Enrolled identities by role and admission state", + []string{"role", "state"}, nil), + usersDesc: prometheus.NewDesc( + "sam_control_plane_users", + "Human identities that have enrolled", + nil, nil), + requestsDesc: prometheus.NewDesc( + "sam_control_plane_enrollment_requests", + "Enrollment requests by status", + []string{"status"}, nil), + tokensDesc: prometheus.NewDesc( + "sam_control_plane_bootstrap_tokens", + "Bootstrap tokens by state", + []string{"state"}, nil), + routersDesc: prometheus.NewDesc( + "sam_control_plane_routers_active", + "Routers holding an unexpired lease", + nil, nil), + routerPeersDesc: prometheus.NewDesc( + "sam_control_plane_router_connected_peers", + "Peers a router reported connected on its last lease renewal", + []string{"router"}, nil), + routerDHTDesc: prometheus.NewDesc( + "sam_control_plane_router_dht_size", + "DHT routing table size a router reported on its last lease renewal", + []string{"router"}, nil), + routerLeaseDesc: prometheus.NewDesc( + "sam_control_plane_router_lease_renewed_timestamp_seconds", + "Unix time of a router's last lease renewal", + []string{"router"}, nil), + meshPeersDesc: prometheus.NewDesc( + "sam_control_plane_mesh_connected_peers", + "Distinct non-router peers connected to at least one active router", + nil, nil), + scrapeOKDesc: prometheus.NewDesc( + "sam_control_plane_mesh_state_scrape_success", + "1 if the mesh state was read from the store, 0 if the last read failed", + nil, nil), + scrapeSampleDesc: prometheus.NewDesc( + "sam_control_plane_mesh_state_timestamp_seconds", + "Unix time the exported mesh state was read from the store", + nil, nil), + } +} + +func (c *meshStateCollector) Describe(ch chan<- *prometheus.Desc) { + for _, d := range []*prometheus.Desc{ + c.nodesDesc, c.usersDesc, c.requestsDesc, c.tokensDesc, c.routersDesc, + c.routerPeersDesc, c.routerDHTDesc, c.routerLeaseDesc, c.meshPeersDesc, + c.scrapeOKDesc, c.scrapeSampleDesc, + } { + ch <- d + } +} + +func (c *meshStateCollector) Collect(ch chan<- prometheus.Metric) { + c.mu.Lock() + defer c.mu.Unlock() + + now := c.now() + if now.Sub(c.fetchedAt) >= c.ttl { + c.refresh(now) + } + + ok := 0.0 + if c.ok { + ok = 1 + } + ch <- prometheus.MustNewConstMetric(c.scrapeOKDesc, prometheus.GaugeValue, ok) + if c.fetchedAt.IsZero() { + return + } + ch <- prometheus.MustNewConstMetric(c.scrapeSampleDesc, prometheus.GaugeValue, float64(c.fetchedAt.Unix())) + + s := c.snapshot + for k, n := range s.nodesByRoleState { + ch <- prometheus.MustNewConstMetric(c.nodesDesc, prometheus.GaugeValue, float64(n), k[0], k[1]) + } + ch <- prometheus.MustNewConstMetric(c.usersDesc, prometheus.GaugeValue, float64(s.users)) + for status, n := range s.enrollmentRequests { + ch <- prometheus.MustNewConstMetric(c.requestsDesc, prometheus.GaugeValue, float64(n), status) + } + for state, n := range s.tokensByState { + ch <- prometheus.MustNewConstMetric(c.tokensDesc, prometheus.GaugeValue, float64(n), state) + } + ch <- prometheus.MustNewConstMetric(c.routersDesc, prometheus.GaugeValue, float64(len(s.routers))) + for _, r := range s.routers { + ch <- prometheus.MustNewConstMetric(c.routerPeersDesc, prometheus.GaugeValue, float64(len(r.ConnectedPeers)), r.PeerID) + ch <- prometheus.MustNewConstMetric(c.routerDHTDesc, prometheus.GaugeValue, float64(r.DHTSize), r.PeerID) + ch <- prometheus.MustNewConstMetric(c.routerLeaseDesc, prometheus.GaugeValue, float64(r.LastRenewal.Unix()), r.PeerID) + } + ch <- prometheus.MustNewConstMetric(c.meshPeersDesc, prometheus.GaugeValue, float64(s.meshConnectedPeers)) +} + +// refresh replaces the snapshot from the store. A failed read keeps the last +// good snapshot and flips the success gauge, so a store outage is visible +// without the figures vanishing. +func (c *meshStateCollector) refresh(now time.Time) { + ctx, cancel := context.WithTimeout(context.Background(), c.timeout) + defer cancel() + + snap, err := readMeshSnapshot(ctx, c.store, now) + if err != nil { + logger.Warnf("Mesh state metrics refresh failed: %v", err) + c.ok = false + return + } + c.snapshot = snap + c.fetchedAt = now + c.ok = true +} + +func readMeshSnapshot(ctx context.Context, store storage.Store, now time.Time) (meshSnapshot, error) { + snap := meshSnapshot{ + nodesByRoleState: map[[2]string]int{}, + enrollmentRequests: map[string]int{}, + tokensByState: map[string]int{}, + } + + nodes, err := store.ListNodes(ctx) + if err != nil { + return snap, err + } + for i := range nodes { + snap.nodesByRoleState[[2]string{nodes[i].Role, nodeState(&nodes[i], now)}]++ + } + + users, err := store.ListUsers(ctx) + if err != nil { + return snap, err + } + snap.users = len(users) + + reqs, err := store.ListEnrollmentRequests(ctx) + if err != nil { + return snap, err + } + for _, r := range reqs { + snap.enrollmentRequests[enrollmentStatusLabel(r.Status)]++ + } + + tokens, err := store.ListBootstrapTokens(ctx) + if err != nil { + return snap, err + } + for i := range tokens { + snap.tokensByState[tokenState(&tokens[i], now)]++ + } + + routers, err := store.GetActiveRouters(ctx) + if err != nil { + return snap, err + } + snap.routers = routers + snap.meshConnectedPeers = countMeshPeers(routers) + + return snap, nil +} + +// enrollmentStatusLabel turns ENROLLMENT_STATUS_PENDING into "pending". +func enrollmentStatusLabel(s api.EnrollmentStatus) string { + return strings.ToLower(strings.TrimPrefix(s.String(), "ENROLLMENT_STATUS_")) +} + +func nodeState(n *storage.EnrolledNode, now time.Time) string { + switch n.CheckAdmission(now) { + case nil: + return nodeStateAdmitted + case storage.ErrNodeBanned: + return nodeStateBanned + default: + return nodeStateExpired + } +} + +func tokenState(t *storage.BootstrapToken, now time.Time) string { + switch { + case t.IsRevoked(): + return tokenStateRevoked + case !t.ExpiresAt.IsZero() && now.After(t.ExpiresAt): + return tokenStateExpired + case t.UsagesCount >= t.MaxUsages: + return tokenStateExhausted + default: + return tokenStateActive + } +} + +// countMeshPeers is the headline mesh size: every peer some active router has +// attached, counted once, with the routers themselves left out since they +// hold connections to each other. +func countMeshPeers(routers []storage.RouterLease) int { + routerIDs := make(map[string]bool, len(routers)) + for _, r := range routers { + routerIDs[r.PeerID] = true + } + seen := map[string]bool{} + for _, r := range routers { + for _, p := range r.ConnectedPeers { + if !routerIDs[p] { + seen[p] = true + } + } + } + return len(seen) +} + +// metricsHandler serves the process-wide registry, which carries the request +// counters above and the Go runtime, alongside this server's own mesh state. +// The mesh collector is per server rather than global so two servers in one +// process (tests, sam-one) never fight over a registration. +func (s *Server) metricsHandler() http.Handler { + return promhttp.HandlerFor( + prometheus.Gatherers{prometheus.DefaultGatherer, s.metricsRegistry}, + promhttp.HandlerOpts{}, + ) +} diff --git a/internal/controlplane/metrics_test.go b/internal/controlplane/metrics_test.go new file mode 100644 index 00000000..e37c2da0 --- /dev/null +++ b/internal/controlplane/metrics_test.go @@ -0,0 +1,323 @@ +// 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. + +package controlplane + +import ( + "context" + "io" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/google/sam/api" + "github.com/google/sam/internal/storage" + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" +) + +// gathered indexes one Gather() by metric name and sorted label pairs so a +// test can ask for a sample by "name{k=v,...}". +type gathered map[string]float64 + +func gather(t *testing.T, c prometheus.Collector) gathered { + t.Helper() + reg := prometheus.NewPedanticRegistry() + if err := reg.Register(c); err != nil { + t.Fatalf("register: %v", err) + } + fams, err := reg.Gather() + if err != nil { + t.Fatalf("gather: %v", err) + } + out := gathered{} + for _, f := range fams { + for _, m := range f.GetMetric() { + out[sampleKey(f.GetName(), m)] = m.GetGauge().GetValue() + } + } + return out +} + +func sampleKey(name string, m *dto.Metric) string { + var parts []string + for _, l := range m.GetLabel() { + parts = append(parts, l.GetName()+"="+l.GetValue()) + } + if len(parts) == 0 { + return name + } + return name + "{" + strings.Join(parts, ",") + "}" +} + +func (g gathered) want(t *testing.T, key string, v float64) { + t.Helper() + got, ok := g[key] + if !ok { + t.Fatalf("missing sample %s; have %v", key, g) + } + if got != v { + t.Errorf("%s = %v, want %v", key, got, v) + } +} + +func newMetricsTestStore(t *testing.T) storage.Store { + t.Helper() + store, err := storage.NewSQLStore("sqlite", filepath.Join(t.TempDir(), "cp.db")) + if err != nil { + t.Fatalf("store: %v", err) + } + t.Cleanup(func() { _ = store.Close() }) + return store +} + +func seedMeshState(t *testing.T, store storage.Store, now time.Time) { + t.Helper() + ctx := context.Background() + + nodes := []storage.EnrolledNode{ + {PeerID: "node-a", Role: api.RoleNode, EnrolledAt: now, ExpiresAt: now.Add(time.Hour)}, + {PeerID: "node-b", Role: api.RoleNode, EnrolledAt: now, ExpiresAt: now.Add(time.Hour)}, + {PeerID: "node-old", Role: api.RoleNode, EnrolledAt: now.Add(-2 * time.Hour), ExpiresAt: now.Add(-time.Hour)}, + {PeerID: "node-bad", Role: api.RoleNode, EnrolledAt: now, ExpiresAt: now.Add(time.Hour)}, + {PeerID: "router-1", Role: api.RoleRouter, EnrolledAt: now, ExpiresAt: now.Add(time.Hour)}, + {PeerID: "router-2", Role: api.RoleRouter, EnrolledAt: now, ExpiresAt: now.Add(time.Hour)}, + } + for i := range nodes { + nodes[i].PublicKey = []byte("pk-" + nodes[i].PeerID) + nodes[i].Biscuit = []byte("biscuit") + nodes[i].EnrollmentType = "test" + if err := store.EnrollNode(ctx, &nodes[i]); err != nil { + t.Fatalf("enroll %s: %v", nodes[i].PeerID, err) + } + } + if err := store.SetNodeBanned(ctx, "node-bad", true); err != nil { + t.Fatalf("ban: %v", err) + } + + // Routers peer with each other and share node-a; node-c is attached to + // router-2 without an enrollment row, as a peer mid-handshake would be. + leases := []storage.RouterLease{ + {PeerID: "router-1", Addresses: []string{"/ip4/10.0.0.1/tcp/4501"}, LastRenewal: now, ExpiresAt: now.Add(time.Hour), + ConnectedPeers: []string{"router-2", "node-a", "node-b"}, DHTSize: 5}, + {PeerID: "router-2", Addresses: []string{"/ip4/10.0.0.2/tcp/4501"}, LastRenewal: now.Add(-time.Minute), ExpiresAt: now.Add(time.Hour), + ConnectedPeers: []string{"router-1", "node-a", "node-c"}, DHTSize: 4}, + {PeerID: "router-gone", Addresses: []string{"/ip4/10.0.0.3/tcp/4501"}, LastRenewal: now.Add(-2 * time.Hour), ExpiresAt: now.Add(-time.Hour), + ConnectedPeers: []string{"node-z"}, DHTSize: 1}, + } + for i := range leases { + if err := store.UpsertRouterLease(ctx, &leases[i]); err != nil { + t.Fatalf("lease %s: %v", leases[i].PeerID, err) + } + } + + if err := store.SaveUser(ctx, &storage.User{ID: "alice", Issuer: "https://issuer", Email: "alice@example.com", Role: "user", CreatedAt: now}); err != nil { + t.Fatalf("user: %v", err) + } + + tokens := []storage.BootstrapToken{ + {ID: "tok-active", TokenHash: "h1", Role: api.RoleNode, MaxUsages: 5, UsagesCount: 1, CreatedAt: now, ExpiresAt: now.Add(time.Hour)}, + {ID: "tok-exhausted", TokenHash: "h2", Role: api.RoleNode, MaxUsages: 1, UsagesCount: 1, CreatedAt: now, ExpiresAt: now.Add(time.Hour)}, + {ID: "tok-expired", TokenHash: "h3", Role: api.RoleNode, MaxUsages: 5, CreatedAt: now.Add(-2 * time.Hour), ExpiresAt: now.Add(-time.Hour)}, + {ID: "tok-revoked", TokenHash: "h4", Role: api.RoleNode, MaxUsages: 5, CreatedAt: now, ExpiresAt: now.Add(time.Hour)}, + } + for i := range tokens { + if err := store.SaveBootstrapToken(ctx, &tokens[i]); err != nil { + t.Fatalf("token %s: %v", tokens[i].ID, err) + } + } + if err := store.RevokeBootstrapToken(ctx, "tok-revoked"); err != nil { + t.Fatalf("revoke: %v", err) + } + + reqs := []storage.EnrollmentRequest{ + {ID: "req-1", PeerID: "pending-1", TokenID: "tok-active", Status: api.EnrollmentStatus_ENROLLMENT_STATUS_PENDING, CreatedAt: now}, + {ID: "req-2", PeerID: "pending-2", TokenID: "tok-active", Status: api.EnrollmentStatus_ENROLLMENT_STATUS_PENDING, CreatedAt: now}, + {ID: "req-3", PeerID: "node-a", TokenID: "tok-active", Status: api.EnrollmentStatus_ENROLLMENT_STATUS_APPROVED, CreatedAt: now}, + } + for i := range reqs { + reqs[i].PublicKey = []byte("pk-" + reqs[i].PeerID) + if err := store.CreateEnrollmentRequest(ctx, &reqs[i]); err != nil { + t.Fatalf("request %s: %v", reqs[i].ID, err) + } + } +} + +func TestMeshStateCollectorExportsStoreState(t *testing.T) { + store := newMetricsTestStore(t) + now := time.Now() + seedMeshState(t, store, now) + + c := newMeshStateCollector(store) + c.now = func() time.Time { return now } + g := gather(t, c) + + g.want(t, "sam_control_plane_mesh_state_scrape_success", 1) + g.want(t, "sam_control_plane_mesh_state_timestamp_seconds", float64(now.Unix())) + + g.want(t, "sam_control_plane_enrolled_nodes{role="+api.RoleNode+",state=admitted}", 2) + g.want(t, "sam_control_plane_enrolled_nodes{role="+api.RoleNode+",state=expired}", 1) + g.want(t, "sam_control_plane_enrolled_nodes{role="+api.RoleNode+",state=banned}", 1) + g.want(t, "sam_control_plane_enrolled_nodes{role="+api.RoleRouter+",state=admitted}", 2) + g.want(t, "sam_control_plane_users", 1) + + g.want(t, "sam_control_plane_enrollment_requests{status=pending}", 2) + g.want(t, "sam_control_plane_enrollment_requests{status=approved}", 1) + + g.want(t, "sam_control_plane_bootstrap_tokens{state=active}", 1) + g.want(t, "sam_control_plane_bootstrap_tokens{state=exhausted}", 1) + g.want(t, "sam_control_plane_bootstrap_tokens{state=expired}", 1) + g.want(t, "sam_control_plane_bootstrap_tokens{state=revoked}", 1) + + // The lapsed lease is neither counted nor labelled. + g.want(t, "sam_control_plane_routers_active", 2) + g.want(t, "sam_control_plane_router_connected_peers{router=router-1}", 3) + g.want(t, "sam_control_plane_router_connected_peers{router=router-2}", 3) + g.want(t, "sam_control_plane_router_dht_size{router=router-1}", 5) + g.want(t, "sam_control_plane_router_dht_size{router=router-2}", 4) + g.want(t, "sam_control_plane_router_lease_renewed_timestamp_seconds{router=router-2}", float64(now.Add(-time.Minute).Unix())) + if _, ok := g["sam_control_plane_router_connected_peers{router=router-gone}"]; ok { + t.Error("lapsed router lease was exported") + } + + // node-a counted once, routers excluded, node-z behind a dead lease ignored. + g.want(t, "sam_control_plane_mesh_connected_peers", 3) +} + +func TestMeshStateCollectorCachesWithinTTL(t *testing.T) { + store := newMetricsTestStore(t) + now := time.Now() + seedMeshState(t, store, now) + + clock := now + c := newMeshStateCollector(store) + c.now = func() time.Time { return clock } + + gather(t, c).want(t, "sam_control_plane_routers_active", 2) + + extra := storage.RouterLease{PeerID: "router-3", Addresses: []string{"/ip4/10.0.0.4/tcp/4501"}, LastRenewal: now, ExpiresAt: now.Add(time.Hour)} + if err := store.UpsertRouterLease(context.Background(), &extra); err != nil { + t.Fatal(err) + } + + clock = now.Add(c.ttl / 2) + gather(t, c).want(t, "sam_control_plane_routers_active", 2) + + clock = now.Add(c.ttl) + gather(t, c).want(t, "sam_control_plane_routers_active", 3) +} + +func TestMeshStateCollectorKeepsLastSnapshotWhenStoreFails(t *testing.T) { + store := newMetricsTestStore(t) + now := time.Now() + seedMeshState(t, store, now) + + clock := now + c := newMeshStateCollector(store) + c.now = func() time.Time { return clock } + + g := gather(t, c) + g.want(t, "sam_control_plane_mesh_state_scrape_success", 1) + g.want(t, "sam_control_plane_routers_active", 2) + + if err := store.Close(); err != nil { + t.Fatal(err) + } + clock = now.Add(c.ttl) + + g = gather(t, c) + g.want(t, "sam_control_plane_mesh_state_scrape_success", 0) + g.want(t, "sam_control_plane_routers_active", 2) + g.want(t, "sam_control_plane_mesh_state_timestamp_seconds", float64(now.Unix())) +} + +func TestMeshStateCollectorBeforeFirstReadReportsOnlyFailure(t *testing.T) { + store := newMetricsTestStore(t) + if err := store.Close(); err != nil { + t.Fatal(err) + } + + g := gather(t, newMeshStateCollector(store)) + g.want(t, "sam_control_plane_mesh_state_scrape_success", 0) + if len(g) != 1 { + t.Errorf("expected only the success gauge before any snapshot, got %v", g) + } +} + +func TestObserveRouteLabelsByPatternAndCode(t *testing.T) { + before := counterValue(t, httpRequestsTotal, "/admin/nodes/", "404") + + h := observeRoute("/admin/nodes/", func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "no such node", http.StatusNotFound) + }) + rec := httptest.NewRecorder() + h.ServeHTTP(rec, httptest.NewRequest(http.MethodPost, "/admin/nodes/12D3KooWsecret/unban", nil)) + if rec.Code != http.StatusNotFound { + t.Fatalf("status = %d", rec.Code) + } + + if got := counterValue(t, httpRequestsTotal, "/admin/nodes/", "404"); got != before+1 { + t.Errorf("counter = %v, want %v", got, before+1) + } + + // A handler that never calls WriteHeader is a 200. + before = counterValue(t, httpRequestsTotal, "/info", "200") + observeRoute("/info", func(w http.ResponseWriter, r *http.Request) { + _, _ = io.WriteString(w, "ok") + }).ServeHTTP(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/info", nil)) + if got := counterValue(t, httpRequestsTotal, "/info", "200"); got != before+1 { + t.Errorf("counter = %v, want %v", got, before+1) + } +} + +func counterValue(t *testing.T, vec *prometheus.CounterVec, labels ...string) float64 { + t.Helper() + var m dto.Metric + if err := vec.WithLabelValues(labels...).Write(&m); err != nil { + t.Fatal(err) + } + return m.GetCounter().GetValue() +} + +func TestMetricsEndpointServesMeshStateAndRuntime(t *testing.T) { + issuer, _ := startCustomMockOIDC(t) + srv, store, baseURL := setupTestServer(t, issuer) + defer func() { + _ = srv.Close() + _ = store.Close() + }() + + resp, err := http.Get(baseURL + "/metrics") + if err != nil { + t.Fatal(err) + } + defer func() { _ = resp.Body.Close() }() + body, _ := io.ReadAll(resp.Body) + if resp.StatusCode != http.StatusOK { + t.Fatalf("status %d: %s", resp.StatusCode, body) + } + for _, want := range []string{ + "sam_control_plane_mesh_state_scrape_success 1", + "sam_control_plane_routers_active 0", + "sam_control_plane_mesh_connected_peers 0", + "go_goroutines", + } { + if !strings.Contains(string(body), want) { + t.Errorf("/metrics missing %q", want) + } + } +} diff --git a/internal/controlplane/node_gc.go b/internal/controlplane/node_gc.go new file mode 100644 index 00000000..07ff5844 --- /dev/null +++ b/internal/controlplane/node_gc.go @@ -0,0 +1,83 @@ +// 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. + +package controlplane + +import ( + "context" + "time" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +const ( + // DefaultNodeRetention keeps a lapsed enrollment for a month past its + // session, long enough to be looked up when investigating an incident, + // and is what --node-retention defaults to. + DefaultNodeRetention = 30 * 24 * time.Hour + + // nodeGCInterval is how often the sweep runs. Rows lapse on a scale of + // days, so nothing is gained by sweeping more often than hourly. + nodeGCInterval = time.Hour +) + +var nodesDeletedTotal = promauto.NewCounter(prometheus.CounterOpts{ + Name: "sam_control_plane_nodes_deleted_total", + Help: "Enrolled node records removed after outliving their session by the retention period", +}) + +// runNodeGCLoop sweeps once at start and then hourly. Replicas each sweep +// against the shared database; the delete is idempotent, so the race is +// harmless and not worth a claim. +func (s *Server) runNodeGCLoop() { + defer s.wg.Done() + if s.config.NodeRetention <= 0 { + return + } + + s.gcExpiredNodes(time.Now()) + + ticker := time.NewTicker(nodeGCInterval) + defer ticker.Stop() + for { + select { + case <-ticker.C: + s.gcExpiredNodes(time.Now()) + case <-s.ctx.Done(): + return + } + } +} + +// gcExpiredNodes deletes every unbanned node whose session expired more than +// NodeRetention before now. A retention of zero means keep forever, and that +// holds here too, not only in the loop that decides whether to tick. +func (s *Server) gcExpiredNodes(now time.Time) { + if s.config.NodeRetention <= 0 { + return + } + ctx, cancel := context.WithTimeout(s.ctx, 30*time.Second) + defer cancel() + + deleted, err := s.store.DeleteExpiredNodes(ctx, now.Add(-s.config.NodeRetention)) + if err != nil { + logger.Errorf("Failed to delete expired node records: %v", err) + return + } + if deleted > 0 { + nodesDeletedTotal.Add(float64(deleted)) + logger.Infof("Deleted %d node records whose session expired more than %s ago", deleted, s.config.NodeRetention) + } +} diff --git a/internal/controlplane/node_gc_test.go b/internal/controlplane/node_gc_test.go new file mode 100644 index 00000000..6cb1a961 --- /dev/null +++ b/internal/controlplane/node_gc_test.go @@ -0,0 +1,109 @@ +// 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. + +package controlplane + +import ( + "context" + "testing" + "time" + + "github.com/google/sam/api" + "github.com/google/sam/internal/storage" + dto "github.com/prometheus/client_model/go" +) + +func TestGCExpiredNodesAppliesRetention(t *testing.T) { + store := newMetricsTestStore(t) + ctx := context.Background() + now := time.Now() + + enroll := func(id string, expiresAt time.Time) { + t.Helper() + if err := store.EnrollNode(ctx, &storage.EnrolledNode{ + PeerID: id, PublicKey: []byte("pub"), Biscuit: []byte("b"), Role: api.RoleNode, + EnrollmentType: "OIDC", EnrolledAt: now.Add(-90 * 24 * time.Hour), ExpiresAt: expiresAt, + }); err != nil { + t.Fatalf("enroll %s: %v", id, err) + } + } + // Retention is 7d: only a session that lapsed more than a week ago goes. + enroll("lapsed-8d", now.Add(-8*24*time.Hour)) + enroll("lapsed-6d", now.Add(-6*24*time.Hour)) + enroll("admitted", now.Add(24*time.Hour)) + + srv, err := NewServer(Options{ + DriverName: "sqlite", DataSourceName: "unused", OIDCIssuer: "https://issuer", + NodeRetention: 7 * 24 * time.Hour, + }, store) + if err != nil { + t.Fatal(err) + } + defer func() { _ = srv.Close() }() + + var before dto.Metric + _ = nodesDeletedTotal.Write(&before) + + srv.gcExpiredNodes(now) + + nodes, err := store.ListNodes(ctx) + if err != nil { + t.Fatal(err) + } + if len(nodes) != 2 { + t.Fatalf("got %d nodes after sweep, want 2: %+v", len(nodes), nodes) + } + for _, n := range nodes { + if n.PeerID == "lapsed-8d" { + t.Error("lapsed-8d survived a 7d retention") + } + } + + var after dto.Metric + _ = nodesDeletedTotal.Write(&after) + if got := after.GetCounter().GetValue() - before.GetCounter().GetValue(); got != 1 { + t.Errorf("sam_control_plane_nodes_deleted_total advanced by %v, want 1", got) + } +} + +func TestNodeGCLoopIsOffWithoutRetention(t *testing.T) { + store := newMetricsTestStore(t) + ctx := context.Background() + now := time.Now() + if err := store.EnrollNode(ctx, &storage.EnrolledNode{ + PeerID: "lapsed", PublicKey: []byte("pub"), Biscuit: []byte("b"), Role: api.RoleNode, + EnrollmentType: "OIDC", EnrolledAt: now.Add(-2 * time.Hour), ExpiresAt: now.Add(-time.Hour), + }); err != nil { + t.Fatal(err) + } + + srv, err := NewServer(Options{DriverName: "sqlite", DataSourceName: "unused", OIDCIssuer: "https://issuer"}, store) + if err != nil { + t.Fatal(err) + } + defer func() { _ = srv.Close() }() + + // The loop returns immediately with retention 0, and a direct sweep is a + // no-op too: zero means keep forever, not "expired as of now". + srv.wg.Add(1) + srv.runNodeGCLoop() + srv.gcExpiredNodes(now) + nodes, err := store.ListNodes(ctx) + if err != nil { + t.Fatal(err) + } + if len(nodes) != 1 { + t.Fatalf("retention 0 deleted rows: %d left", len(nodes)) + } +} diff --git a/internal/controlplane/server.go b/internal/controlplane/server.go index 6f7dd9e3..54023db8 100644 --- a/internal/controlplane/server.go +++ b/internal/controlplane/server.go @@ -45,7 +45,7 @@ import ( "github.com/libp2p/go-libp2p/core/crypto" "github.com/libp2p/go-libp2p/core/peer" "github.com/multiformats/go-multiaddr" - "github.com/prometheus/client_golang/prometheus/promhttp" + "github.com/prometheus/client_golang/prometheus" "golang.org/x/time/rate" "google.golang.org/protobuf/encoding/protojson" "google.golang.org/protobuf/proto" @@ -93,6 +93,10 @@ type Server struct { catalogMu sync.RWMutex catalog map[string]nodeCatalogEntry + // metricsRegistry holds this server's store-backed mesh state collector; + // see metricsHandler. + metricsRegistry *prometheus.Registry + ctx context.Context cancel context.CancelFunc wg sync.WaitGroup @@ -108,15 +112,19 @@ func NewServer(config Options, store storage.Store) (*Server, error) { ctx, cancel := context.WithCancel(context.Background()) + reg := prometheus.NewRegistry() + reg.MustRegister(newMeshStateCollector(store)) + return &Server{ - config: config, - store: store, - mesh: NewNopMeshAdapter(), - limiter: rate.NewLimiter(rate.Limit(EnrollRateLimit), EnrollBurst), - providers: make(map[string]*oidc.Provider), - catalog: make(map[string]nodeCatalogEntry), - ctx: ctx, - cancel: cancel, + config: config, + store: store, + mesh: NewNopMeshAdapter(), + limiter: rate.NewLimiter(rate.Limit(EnrollRateLimit), EnrollBurst), + providers: make(map[string]*oidc.Provider), + catalog: make(map[string]nodeCatalogEntry), + metricsRegistry: reg, + ctx: ctx, + cancel: cancel, }, nil } @@ -205,33 +213,41 @@ func (s *Server) Init() error { s.wg.Add(1) go s.runKeyRotationLoop() + s.wg.Add(1) + go s.runNodeGCLoop() + return nil } -// RegisterRoutes registers every control-plane HTTP handler on mux. +// RegisterRoutes registers every control-plane HTTP handler on mux. Probes +// and /metrics are left uncounted so scrapers do not dominate the figures. func (s *Server) RegisterRoutes(mux *http.ServeMux) { mux.HandleFunc("/healthz", s.HandleHealthz) mux.HandleFunc("/readyz", s.HandleReadyz) - mux.Handle("/metrics", promhttp.Handler()) - mux.HandleFunc("/info", s.HandleInfo) - mux.HandleFunc("/register", noStore(s.HandleRegister)) - mux.HandleFunc("/keys", s.HandleKeys) - mux.HandleFunc("/routers/lease", s.HandleRouterLease) - mux.HandleFunc("/policies", s.HandlePolicies) - mux.HandleFunc("/enroll", noStore(s.HandleEnroll)) - mux.HandleFunc("/enroll/status", noStore(s.HandleEnrollStatus)) - mux.HandleFunc("/refresh", noStore(s.HandleRefresh)) - mux.HandleFunc("/nodes/catalog", s.HandleNodeCatalog) - mux.HandleFunc("/admin/bootstrap-tokens", noStore(s.HandleAdminBootstrapTokens)) - mux.HandleFunc("/admin/bootstrap-tokens/", noStore(s.HandleAdminBootstrapTokenAction)) - mux.HandleFunc("/admin/enrollments", noStore(s.HandleAdminEnrollments)) - mux.HandleFunc("/admin/enrollments/", noStore(s.HandleAdminEnrollmentAction)) - mux.HandleFunc("/admin/nodes/", noStore(s.HandleAdminNodeAction)) - mux.HandleFunc("/admin/revoke", noStore(s.HandleAdminRevoke)) - mux.HandleFunc("/admin/status", noStore(s.HandleAdminStatus)) - mux.HandleFunc("/user/status", noStore(s.HandleUserStatus)) - mux.HandleFunc("/user/bootstrap-tokens", noStore(s.HandleUserBootstrapTokens)) - mux.HandleFunc("/user/revoke", noStore(s.HandleUserRevoke)) + mux.Handle("/metrics", s.metricsHandler()) + + handle := func(pattern string, h http.HandlerFunc) { + mux.Handle(pattern, observeRoute(pattern, h)) + } + handle("/info", s.HandleInfo) + handle("/register", noStore(s.HandleRegister)) + handle("/keys", s.HandleKeys) + handle("/routers/lease", s.HandleRouterLease) + handle("/policies", s.HandlePolicies) + handle("/enroll", noStore(s.HandleEnroll)) + handle("/enroll/status", noStore(s.HandleEnrollStatus)) + handle("/refresh", noStore(s.HandleRefresh)) + handle("/nodes/catalog", s.HandleNodeCatalog) + handle("/admin/bootstrap-tokens", noStore(s.HandleAdminBootstrapTokens)) + handle("/admin/bootstrap-tokens/", noStore(s.HandleAdminBootstrapTokenAction)) + handle("/admin/enrollments", noStore(s.HandleAdminEnrollments)) + handle("/admin/enrollments/", noStore(s.HandleAdminEnrollmentAction)) + handle("/admin/nodes/", noStore(s.HandleAdminNodeAction)) + handle("/admin/revoke", noStore(s.HandleAdminRevoke)) + handle("/admin/status", noStore(s.HandleAdminStatus)) + handle("/user/status", noStore(s.HandleUserStatus)) + handle("/user/bootstrap-tokens", noStore(s.HandleUserBootstrapTokens)) + handle("/user/revoke", noStore(s.HandleUserRevoke)) } // noStore marks responses that carry credentials (biscuits, bootstrap tokens, diff --git a/internal/node/node.go b/internal/node/node.go index af2d4ae8..b80e2ae9 100644 --- a/internal/node/node.go +++ b/internal/node/node.go @@ -64,6 +64,7 @@ import ( "github.com/libp2p/go-msgio" "github.com/multiformats/go-multiaddr" madns "github.com/multiformats/go-multiaddr-dns" + "github.com/prometheus/client_golang/prometheus" "google.golang.org/protobuf/proto" ) @@ -178,6 +179,10 @@ type SamNode struct { BiscuitTimeout time.Duration cachedIdentity atomic.Value logger *golog.ZapEventLogger + + // metricsRegistry holds this node's state collector; see metricsHandler. + metricsOnce sync.Once + metricsRegistry *prometheus.Registry } // UpdateRelays updates the current relays used by AutoRelay. diff --git a/internal/node/node_metrics.go b/internal/node/node_metrics.go new file mode 100644 index 00000000..0519d259 --- /dev/null +++ b/internal/node/node_metrics.go @@ -0,0 +1,182 @@ +// 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. + +package node + +import ( + "errors" + "net" + "net/http" + "strings" + "time" + + "github.com/google/sam/api" + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +// nodeStateCollector exports the node's own view of its mesh membership: the +// same figures GET /debug/mesh-info and /debug/token-info answer on demand, +// as gauges a scraper can keep. A node that runs unattended, such as a +// canary, becomes a continuous probe of the mesh this way. +type nodeStateCollector struct { + n *SamNode + + readyDesc *prometheus.Desc + meshConnectedDesc *prometheus.Desc + connectedDesc *prometheus.Desc + authenticatedDesc *prometheus.Desc + dhtDesc *prometheus.Desc + biscuitExpiryDesc *prometheus.Desc + servicesDesc *prometheus.Desc +} + +func newNodeStateCollector(n *SamNode) *nodeStateCollector { + return &nodeStateCollector{ + n: n, + readyDesc: prometheus.NewDesc( + "sam_node_ready", + "1 once the libp2p host, DHT and identity store exist", + nil, nil), + meshConnectedDesc: prometheus.NewDesc( + "sam_node_mesh_connected", + "1 while the node holds an authenticated connection to a router", + nil, nil), + connectedDesc: prometheus.NewDesc( + "sam_node_connected_peers", + "Peers with an open libp2p connection, authenticated or not", + nil, nil), + authenticatedDesc: prometheus.NewDesc( + "sam_node_authenticated_peers", + "Peers that have completed the mesh authentication handshake with this node", + nil, nil), + dhtDesc: prometheus.NewDesc( + "sam_node_dht_routing_table_size", + "Peers in the Kademlia routing table", + nil, nil), + biscuitExpiryDesc: prometheus.NewDesc( + "sam_node_biscuit_expiry_timestamp_seconds", + "Unix time the node's mesh credential expires", + nil, nil), + servicesDesc: prometheus.NewDesc( + "sam_node_services_registered", + "Local services this node offers to the mesh, by type", + []string{"type"}, nil), + } +} + +func (c *nodeStateCollector) Describe(ch chan<- *prometheus.Desc) { + for _, d := range []*prometheus.Desc{ + c.readyDesc, c.meshConnectedDesc, c.connectedDesc, c.authenticatedDesc, + c.dhtDesc, c.biscuitExpiryDesc, c.servicesDesc, + } { + ch <- d + } +} + +func (c *nodeStateCollector) Collect(ch chan<- prometheus.Metric) { + n := c.n + if n.debugReady() != nil { + ch <- prometheus.MustNewConstMetric(c.readyDesc, prometheus.GaugeValue, 0) + return + } + ch <- prometheus.MustNewConstMetric(c.readyDesc, prometheus.GaugeValue, 1) + + connected := 0.0 + if n.IsConnected() { + connected = 1 + } + ch <- prometheus.MustNewConstMetric(c.meshConnectedDesc, prometheus.GaugeValue, connected) + ch <- prometheus.MustNewConstMetric(c.connectedDesc, prometheus.GaugeValue, float64(len(n.Host.Network().Peers()))) + + authenticated := 0 + n.authPeers.Range(func(_, _ any) bool { + authenticated++ + return true + }) + ch <- prometheus.MustNewConstMetric(c.authenticatedDesc, prometheus.GaugeValue, float64(authenticated)) + ch <- prometheus.MustNewConstMetric(c.dhtDesc, prometheus.GaugeValue, float64(n.DHT.RoutingTable().Size())) + + if exp, err := n.Store.LoadIdentityExpiration(); err == nil && exp > 0 { + ch <- prometheus.MustNewConstMetric(c.biscuitExpiryDesc, prometheus.GaugeValue, float64(exp)) + } + + if n.services != nil { + byType := map[string]int{} + for _, svc := range n.services.List(api.ServiceType_SERVICE_TYPE_UNSPECIFIED) { + byType[serviceTypeLabel(svc.Type)]++ + } + for t, count := range byType { + ch <- prometheus.MustNewConstMetric(c.servicesDesc, prometheus.GaugeValue, float64(count), t) + } + } +} + +// serviceTypeLabel turns SERVICE_TYPE_MCP into "mcp". +func serviceTypeLabel(t api.ServiceType) string { + return strings.ToLower(strings.TrimPrefix(t.String(), "SERVICE_TYPE_")) +} + +// metricsHandler serves the process-wide registry, which carries the request +// and inference counters and the Go runtime, alongside this node's own state. +// The state collector is per node so two nodes in one process (tests, +// sam-one) never fight over a registration. +func (n *SamNode) metricsHandler() http.Handler { + n.metricsOnce.Do(func() { + n.metricsRegistry = prometheus.NewRegistry() + n.metricsRegistry.MustRegister(newNodeStateCollector(n)) + }) + return promhttp.HandlerFor( + prometheus.Gatherers{prometheus.DefaultGatherer, n.metricsRegistry}, + promhttp.HandlerOpts{}, + ) +} + +// StartMetricsServer serves /metrics, /healthz and /readyz on addr with no +// authentication, for a scraper and a kubelet that hold no sidecar token. +// The sidecar's own /metrics stays token-gated: this listener exists so a +// socket-only node, which has no TCP port at all, can still be observed. It +// is off unless an operator names the address, since nothing on it is gated. +func StartMetricsServer(node *SamNode, addr string) (*http.Server, error) { + if addr == "" { + return nil, errors.New("no metrics address configured") + } + listener, err := net.Listen("tcp", addr) + if err != nil { + return nil, err + } + + mux := http.NewServeMux() + mux.Handle("/metrics", node.metricsHandler()) + mux.HandleFunc("/healthz", handleHealthz) + mux.HandleFunc("/readyz", func(w http.ResponseWriter, _ *http.Request) { + if node.debugReady() != nil || !node.IsConnected() { + http.Error(w, "not connected to the mesh", http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusOK) + }) + + server := &http.Server{ + // Informational once Serve has the listener; lets callers find the port. + Addr: listener.Addr().String(), + Handler: mux, + ReadHeaderTimeout: 10 * time.Second, + } + go func() { + _ = server.Serve(listener) + }() + logger.Infof("Metrics listening on http://%s", listener.Addr()) + return server, nil +} diff --git a/internal/node/node_metrics_test.go b/internal/node/node_metrics_test.go new file mode 100644 index 00000000..5747f9ed --- /dev/null +++ b/internal/node/node_metrics_test.go @@ -0,0 +1,150 @@ +// 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. + +package node + +import ( + "io" + "net/http" + "strconv" + "strings" + "testing" + "time" + + "github.com/libp2p/go-libp2p" + dht "github.com/libp2p/go-libp2p-kad-dht" + "github.com/libp2p/go-libp2p/core/peer" +) + +func getMetricsBody(t *testing.T, url string) (int, string) { + t.Helper() + resp, err := http.Get(url) + if err != nil { + t.Fatalf("GET %s: %v", url, err) + } + defer func() { _ = resp.Body.Close() }() + body, _ := io.ReadAll(resp.Body) + return resp.StatusCode, string(body) +} + +func TestNodeMetricsServer(t *testing.T) { + store, err := NewStore(t.TempDir()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = store.Close() }() + + // Half-built, as a node is between NewSamNode and Start: alive, not + // ready, and nothing read from a host that does not exist yet. + node := &SamNode{Store: store, services: NewServiceRegistry(&fakeDHT{}, 0)} + srv, err := StartMetricsServer(node, "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer func() { _ = srv.Close() }() + base := "http://" + srv.Addr + if base == "http://" { + t.Fatal("metrics server did not record its address") + } + + if code, _ := getMetricsBody(t, base+"/healthz"); code != http.StatusOK { + t.Errorf("/healthz = %d before ready", code) + } + if code, _ := getMetricsBody(t, base+"/readyz"); code != http.StatusServiceUnavailable { + t.Errorf("/readyz = %d before ready, want 503", code) + } + code, body := getMetricsBody(t, base+"/metrics") + if code != http.StatusOK { + t.Fatalf("/metrics = %d: %s", code, body) + } + if !strings.Contains(body, "sam_node_ready 0") { + t.Error("/metrics missing sam_node_ready 0 before start") + } + if strings.Contains(body, "sam_node_connected_peers") { + t.Error("/metrics exported host state before the host existed") + } + + host, err := libp2p.New(libp2p.ListenAddrStrings("/ip4/127.0.0.1/tcp/0")) + if err != nil { + t.Fatal(err) + } + defer func() { _ = host.Close() }() + kad, err := dht.New(host, dht.Mode(dht.ModeServer)) + if err != nil { + t.Fatal(err) + } + defer func() { _ = kad.Close() }() + + expiry := time.Now().Add(time.Hour).Unix() + if err := store.SaveIdentityExpiration(expiry); err != nil { + t.Fatal(err) + } + node.Host = host + node.DHT = kad + node.authPeers.Store(peer.ID("peer-a"), time.Now().Add(time.Hour)) + node.authPeers.Store(peer.ID("peer-b"), time.Now().Add(time.Hour)) + + // Ready, but no authenticated router yet: the kubelet must not route to it. + if code, _ := getMetricsBody(t, base+"/readyz"); code != http.StatusServiceUnavailable { + t.Errorf("/readyz = %d with no router, want 503", code) + } + _, body = getMetricsBody(t, base+"/metrics") + for _, want := range []string{ + "sam_node_ready 1", + "sam_node_mesh_connected 0", + "sam_node_connected_peers 0", + "sam_node_authenticated_peers 2", + "sam_node_dht_routing_table_size 0", + "sam_node_biscuit_expiry_timestamp_seconds " + strconv.FormatFloat(float64(expiry), 'g', -1, 64), + // The process-wide registry rides along: sidecar counters and libp2p. + "sam_node_requests_in_flight", + "libp2p_", + } { + if !strings.Contains(body, want) { + t.Errorf("/metrics missing %q", want) + } + } + if strings.Contains(body, "sam_node_services_registered") { + t.Error("services gauge exported with an empty registry") + } + + // An authenticated router that is actually connected flips both. + routerHost, err := libp2p.New(libp2p.ListenAddrStrings("/ip4/127.0.0.1/tcp/0")) + if err != nil { + t.Fatal(err) + } + defer func() { _ = routerHost.Close() }() + if err := host.Connect(t.Context(), peer.AddrInfo{ID: routerHost.ID(), Addrs: routerHost.Addrs()}); err != nil { + t.Fatal(err) + } + node.mu.Lock() + node.authenticatedRouters = map[peer.ID]bool{routerHost.ID(): true} + node.mu.Unlock() + + if code, _ := getMetricsBody(t, base+"/readyz"); code != http.StatusOK { + t.Errorf("/readyz = %d once a router is authenticated", code) + } + _, body = getMetricsBody(t, base+"/metrics") + for _, want := range []string{"sam_node_mesh_connected 1", "sam_node_connected_peers 1"} { + if !strings.Contains(body, want) { + t.Errorf("/metrics missing %q", want) + } + } +} + +func TestStartMetricsServerRequiresAddress(t *testing.T) { + if _, err := StartMetricsServer(&SamNode{}, ""); err == nil { + t.Fatal("expected an error for an empty address") + } +} diff --git a/internal/node/sidecar.go b/internal/node/sidecar.go index ff3c5cac..ea6ca1e1 100644 --- a/internal/node/sidecar.go +++ b/internal/node/sidecar.go @@ -36,7 +36,6 @@ import ( libp2phttp "github.com/libp2p/go-libp2p-http" "github.com/libp2p/go-libp2p/core/network" "github.com/libp2p/go-libp2p/core/peer" - "github.com/prometheus/client_golang/prometheus/promhttp" ) // StartSidecarServer serves the node's local API on a TCP address, on a Unix @@ -51,7 +50,7 @@ func StartSidecarServer(node *SamNode, addr, socketPath, token, certFile, keyFil // Gated like the rest: the labels carry peer IDs and per-peer request counts, // and this mux is reachable by any local process over TCP. Socket callers are // unaffected, which is how every scrape in this repo reads it. - mux.Handle("/metrics", withAuth(token, true, promhttp.Handler())) + mux.Handle("/metrics", withAuth(token, true, node.metricsHandler())) // Protected endpoints. allowAuthorizationFallback=true is safe here: none of // these ever forward the inbound Authorization header to another service. diff --git a/internal/router/config.go b/internal/router/config.go index 7debe25f..58cf29cd 100644 --- a/internal/router/config.go +++ b/internal/router/config.go @@ -57,6 +57,10 @@ type Options struct { // to a non-loopback host. Off by default: whoever answers that URL is // the trust root. AllowInsecureControlPlane bool + // MetricsAddr, when set, serves /metrics, /healthz and /readyz on a + // plain HTTP listener separate from the libp2p ports. Off by default: + // the listener is unauthenticated, so the operator names where it binds. + MetricsAddr string } // Default sets default values for options. diff --git a/internal/router/metrics.go b/internal/router/metrics.go new file mode 100644 index 00000000..ee50a481 --- /dev/null +++ b/internal/router/metrics.go @@ -0,0 +1,199 @@ +// 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. + +package router + +import ( + "net" + "net/http" + "sync" + "time" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +// Outcomes of an inbound /sam/auth handshake. +const ( + handshakeOK = "ok" + handshakeBanned = "banned" + handshakeRateLimited = "rate_limited" + handshakeReadFailed = "read_failed" + handshakeInvalidFrame = "invalid_frame" + handshakeUnauthorized = "unauthorized" +) + +// Outcomes of a control-plane lease renewal attempt. +const ( + leaseOK = "ok" + leaseUnreachable = "unreachable" + leaseUnauthorized = "unauthorized" + leaseRejected = "rejected" +) + +var ( + authHandshakesTotal = promauto.NewCounterVec( + prometheus.CounterOpts{ + Name: "sam_router_auth_handshakes_total", + Help: "Inbound mesh authentication handshakes by outcome", + }, + []string{"result"}, + ) + + leaseRenewalsTotal = promauto.NewCounterVec( + prometheus.CounterOpts{ + Name: "sam_router_lease_renewals_total", + Help: "Control-plane lease renewal attempts by outcome", + }, + []string{"result"}, + ) +) + +// Every outcome exists from the first scrape, so a rate() over one that has +// not happened yet reads 0 rather than no data. +func init() { + for _, v := range []string{handshakeOK, handshakeBanned, handshakeRateLimited, handshakeReadFailed, handshakeInvalidFrame, handshakeUnauthorized} { + authHandshakesTotal.WithLabelValues(v) + } + for _, v := range []string{leaseOK, leaseUnreachable, leaseUnauthorized, leaseRejected} { + leaseRenewalsTotal.WithLabelValues(v) + } +} + +// routerStateCollector exports the live view only the router has: who has +// proven mesh membership on this hop, as opposed to who merely holds a +// transport connection. Nothing is emitted before Start has finished, since +// the host and DHT do not exist yet. +type routerStateCollector struct { + r *Router + + readyDesc *prometheus.Desc + authenticatedDesc *prometheus.Desc + connectedDesc *prometheus.Desc + bannedDesc *prometheus.Desc + dhtDesc *prometheus.Desc + biscuitExpiryDesc *prometheus.Desc +} + +func newRouterStateCollector(r *Router) *routerStateCollector { + return &routerStateCollector{ + r: r, + readyDesc: prometheus.NewDesc( + "sam_router_ready", + "1 once the router is enrolled and its libp2p host is online", + nil, nil), + authenticatedDesc: prometheus.NewDesc( + "sam_router_authenticated_peers", + "Connected peers that have completed the mesh authentication handshake", + nil, nil), + connectedDesc: prometheus.NewDesc( + "sam_router_connected_peers", + "Peers with an open libp2p connection, authenticated or not", + nil, nil), + bannedDesc: prometheus.NewDesc( + "sam_router_banned_peers", + "Peers on the ban list synced from the control plane", + nil, nil), + dhtDesc: prometheus.NewDesc( + "sam_router_dht_routing_table_size", + "Peers in the Kademlia routing table", + nil, nil), + biscuitExpiryDesc: prometheus.NewDesc( + "sam_router_biscuit_expiry_timestamp_seconds", + "Unix time the router's own mesh credential expires", + nil, nil), + } +} + +func (c *routerStateCollector) Describe(ch chan<- *prometheus.Desc) { + for _, d := range []*prometheus.Desc{ + c.readyDesc, c.authenticatedDesc, c.connectedDesc, c.bannedDesc, c.dhtDesc, c.biscuitExpiryDesc, + } { + ch <- d + } +} + +func (c *routerStateCollector) Collect(ch chan<- prometheus.Metric) { + r := c.r + // isReady is stored after Host and DHT are assigned, so observing it + // true is what makes reading them from this goroutine safe. + if !r.isReady.Load() { + ch <- prometheus.MustNewConstMetric(c.readyDesc, prometheus.GaugeValue, 0) + return + } + ch <- prometheus.MustNewConstMetric(c.readyDesc, prometheus.GaugeValue, 1) + ch <- prometheus.MustNewConstMetric(c.authenticatedDesc, prometheus.GaugeValue, float64(syncMapLen(&r.authenticatedPeers))) + ch <- prometheus.MustNewConstMetric(c.connectedDesc, prometheus.GaugeValue, float64(len(r.Host.Network().Peers()))) + ch <- prometheus.MustNewConstMetric(c.bannedDesc, prometheus.GaugeValue, float64(syncMapLen(&r.bannedPeers))) + ch <- prometheus.MustNewConstMetric(c.dhtDesc, prometheus.GaugeValue, float64(r.DHT.RoutingTable().Size())) + + r.keysMu.RLock() + expiry := r.biscuitExpiration + r.keysMu.RUnlock() + if !expiry.IsZero() { + ch <- prometheus.MustNewConstMetric(c.biscuitExpiryDesc, prometheus.GaugeValue, float64(expiry.Unix())) + } +} + +func syncMapLen(m *sync.Map) int { + n := 0 + m.Range(func(_, _ any) bool { + n++ + return true + }) + return n +} + +// serveMetrics opens the operator listener: /metrics, /healthz and /readyz. +// It is off unless an address is configured, and never shares a port with the +// libp2p listeners, so it can stay cluster-internal while the mesh port is +// public. libp2p's own metrics land in the default registry, so exposing it +// here is what makes relay reservations and connection churn visible. +func (r *Router) serveMetrics(addr string) error { + listener, err := net.Listen("tcp", addr) + if err != nil { + return err + } + + reg := prometheus.NewRegistry() + reg.MustRegister(newRouterStateCollector(r)) + + mux := http.NewServeMux() + mux.Handle("/metrics", promhttp.HandlerFor( + prometheus.Gatherers{prometheus.DefaultGatherer, reg}, + promhttp.HandlerOpts{}, + )) + mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + }) + mux.HandleFunc("/readyz", func(w http.ResponseWriter, _ *http.Request) { + if !r.isReady.Load() { + http.Error(w, "router is not online yet", http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusOK) + }) + + r.metricsServer = &http.Server{ + Handler: mux, + ReadHeaderTimeout: 10 * time.Second, + } + r.metricsAddr = listener.Addr() + go func() { + _ = r.metricsServer.Serve(listener) + }() + logger.Infof("Router metrics listening on http://%s", r.metricsAddr) + return nil +} diff --git a/internal/router/metrics_test.go b/internal/router/metrics_test.go new file mode 100644 index 00000000..3e17ecfd --- /dev/null +++ b/internal/router/metrics_test.go @@ -0,0 +1,124 @@ +// 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. + +package router + +import ( + "context" + "io" + "net/http" + "strconv" + "strings" + "testing" + "time" + + "github.com/google/sam/api" + "github.com/libp2p/go-libp2p" + dht "github.com/libp2p/go-libp2p-kad-dht" + "github.com/libp2p/go-libp2p/core/peer" +) + +func getBody(t *testing.T, url string) (int, string) { + t.Helper() + resp, err := http.Get(url) + if err != nil { + t.Fatalf("GET %s: %v", url, err) + } + defer func() { _ = resp.Body.Close() }() + body, _ := io.ReadAll(resp.Body) + return resp.StatusCode, string(body) +} + +func TestRouterMetricsListener(t *testing.T) { + ctx := context.Background() + r, err := NewRouter(ctx, Options{BiscuitTimeout: time.Second, RequiredRole: api.RoleRouter}) + if err != nil { + t.Fatal(err) + } + if err := r.serveMetrics("127.0.0.1:0"); err != nil { + t.Fatal(err) + } + base := "http://" + r.metricsAddr.String() + + // Before Start completes: alive, not ready, and nothing read from a host + // that does not exist yet. + if code, _ := getBody(t, base+"/healthz"); code != http.StatusOK { + t.Errorf("/healthz = %d before ready", code) + } + if code, _ := getBody(t, base+"/readyz"); code != http.StatusServiceUnavailable { + t.Errorf("/readyz = %d before ready, want 503", code) + } + code, body := getBody(t, base+"/metrics") + if code != http.StatusOK { + t.Fatalf("/metrics = %d: %s", code, body) + } + if !strings.Contains(body, "sam_router_ready 0") { + t.Error("/metrics missing sam_router_ready 0 before start") + } + if strings.Contains(body, "sam_router_authenticated_peers") { + t.Error("/metrics exported host state before the host existed") + } + + host, err := libp2p.New(libp2p.ListenAddrStrings("/ip4/127.0.0.1/tcp/0")) + if err != nil { + t.Fatal(err) + } + defer func() { _ = host.Close() }() + kad, err := dht.New(host, dht.Mode(dht.ModeServer)) + if err != nil { + t.Fatal(err) + } + defer func() { _ = kad.Close() }() + + r.Host = host + r.DHT = kad + r.authenticatedPeers.Store(peer.ID("peer-a"), true) + r.authenticatedPeers.Store(peer.ID("peer-b"), true) + r.bannedPeers.Store(peer.ID("peer-x"), true) + expiry := time.Now().Add(time.Hour).Truncate(time.Second) + r.keysMu.Lock() + r.biscuitExpiration = expiry + r.keysMu.Unlock() + r.isReady.Store(true) + + if code, _ := getBody(t, base+"/readyz"); code != http.StatusOK { + t.Errorf("/readyz = %d once ready", code) + } + _, body = getBody(t, base+"/metrics") + for _, want := range []string{ + "sam_router_ready 1", + "sam_router_authenticated_peers 2", + "sam_router_banned_peers 1", + "sam_router_connected_peers 0", + "sam_router_dht_routing_table_size 0", + // The text format renders values with strconv 'g', so a unix time + // comes out in scientific notation. + "sam_router_biscuit_expiry_timestamp_seconds " + strconv.FormatFloat(float64(expiry.Unix()), 'g', -1, 64), + // libp2p registers into the default registry, which the listener + // serves alongside the router's own state. + "libp2p_", + "sam_router_auth_handshakes_total", + } { + if !strings.Contains(body, want) { + t.Errorf("/metrics missing %q", want) + } + } + + if err := r.Close(); err != nil { + t.Fatal(err) + } + if _, err := http.Get(base + "/healthz"); err == nil { + t.Error("metrics listener still answering after Close") + } +} diff --git a/internal/router/router.go b/internal/router/router.go index 02899279..a8f4a470 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -136,6 +136,9 @@ type Router struct { wg sync.WaitGroup isReady atomic.Bool shutdown bool + + metricsServer *http.Server + metricsAddr net.Addr } // NewRouter initializes the router. @@ -236,6 +239,14 @@ func perIPConnResourceManager(limit int) (network.ResourceManager, error) { // Start performs enrollment, syncs keys, launches libp2p host, and starts tasks. func (r *Router) Start() error { + // The operator listener comes up first so /healthz answers while + // enrollment is still in flight; /readyz turns 200 at the end. + if r.config.MetricsAddr != "" { + if err := r.serveMetrics(r.config.MetricsAddr); err != nil { + return fmt.Errorf("failed to start metrics listener: %w", err) + } + } + // 1. Load or Generate persistent identity key priv, err := getOrGeneratePeerKey(r.config.KeysDBPath) if err != nil { @@ -913,6 +924,7 @@ func (r *Router) renewLease() { resp, err := client.Post(r.config.ControlPlaneURL+"/routers/lease", "application/x-protobuf", bytes.NewReader(data)) if err != nil { logger.Errorf("Failed to renew lease with control plane: %v", err) + leaseRenewalsTotal.WithLabelValues(leaseUnreachable).Inc() return } @@ -921,6 +933,7 @@ func (r *Router) renewLease() { if resp.StatusCode == http.StatusUnauthorized && attempt == 0 { logger.Warnf("Control plane lease renewal rejected (401 Unauthorized: %s), attempting recovery...", string(body)) + leaseRenewalsTotal.WithLabelValues(leaseUnauthorized).Inc() if err := r.recoverAfterLease401(); err != nil { logger.Errorf("Recovery failed after 401 Unauthorized lease renewal: %v", err) return @@ -931,19 +944,23 @@ func (r *Router) renewLease() { if resp.StatusCode != http.StatusOK { logger.Errorf("Control plane lease renewal rejected, status %s: %s", resp.Status, string(body)) + leaseRenewalsTotal.WithLabelValues(leaseRejected).Inc() return } var leaseResp api.RouterLeaseResponse if err := proto.Unmarshal(body, &leaseResp); err != nil { logger.Errorf("Failed to parse lease response: %v", err) + leaseRenewalsTotal.WithLabelValues(leaseRejected).Inc() return } if !leaseResp.Success { logger.Errorf("Lease renewal failed: %s", leaseResp.Error) + leaseRenewalsTotal.WithLabelValues(leaseRejected).Inc() } else { logger.Debugf("Lease renewed successfully. Expires at: %s", time.Unix(leaseResp.ExpiresAt, 0)) + leaseRenewalsTotal.WithLabelValues(leaseOK).Inc() } return } @@ -1137,12 +1154,14 @@ func (r *Router) HandleAuthHandshake(s network.Stream) { if _, banned := r.bannedPeers.Load(remotePeer); banned { logger.Warnf("[AuthN] Rejecting authentication for banned peer %s", remotePeer) + authHandshakesTotal.WithLabelValues(handshakeBanned).Inc() _ = s.Reset() return } if r.handshakeLimiter != nil && !r.handshakeLimiter.Allow(remotePeer.String()) { logger.Warnf("[AuthN] Handshake rate limit exceeded for %s", remotePeer) + authHandshakesTotal.WithLabelValues(handshakeRateLimited).Inc() _ = s.Reset() return } @@ -1154,6 +1173,7 @@ func (r *Router) HandleAuthHandshake(s network.Stream) { msg, err := reader.ReadMsg() if err != nil { logger.Errorf("[AuthN] Failed to read handshake from %s: %v", remotePeer, err) + authHandshakesTotal.WithLabelValues(handshakeReadFailed).Inc() return } defer reader.ReleaseMsg(msg) @@ -1161,6 +1181,7 @@ func (r *Router) HandleAuthHandshake(s network.Stream) { var exchange api.AuthFrame if err := proto.Unmarshal(msg, &exchange); err != nil { logger.Warnf("[AuthN] Invalid protobuf from %s", remotePeer) + authHandshakesTotal.WithLabelValues(handshakeInvalidFrame).Inc() return } @@ -1168,11 +1189,13 @@ func (r *Router) HandleAuthHandshake(s network.Stream) { _, err = identity.VerifyBiscuit(exchange.Biscuit, remotePeer, r.getTrustedPublicKeys(), r.config.BiscuitTimeout) if err != nil { logger.Warnf("[AuthN] Authorization failed for peer %s: %v", remotePeer, err) + authHandshakesTotal.WithLabelValues(handshakeUnauthorized).Inc() _ = s.Reset() return } r.authenticatedPeers.Store(remotePeer, true) + authHandshakesTotal.WithLabelValues(handshakeOK).Inc() logger.Infof("[AuthN] Successfully authenticated peer %s", remotePeer) // Send mutual response (our biscuit) @@ -1250,6 +1273,11 @@ func (r *Router) Close() error { r.cancel() var errs []error + if r.metricsServer != nil { + if err := r.metricsServer.Close(); err != nil { + errs = append(errs, err) + } + } if r.DHT != nil { if err := r.DHT.Close(); err != nil { errs = append(errs, err) diff --git a/internal/storage/round_trip_test.go b/internal/storage/round_trip_test.go index 83b0d033..3dfdf67b 100644 --- a/internal/storage/round_trip_test.go +++ b/internal/storage/round_trip_test.go @@ -243,6 +243,61 @@ func TestReEnrollmentKeepsAnExistingBan(t *testing.T) { } } +// Expiry-based deletion must never take a ban with it, nor touch a node that +// never had a session bound: the first is how a node is kept out, the second +// has no expiry to have lapsed. +func TestDeleteExpiredNodesSparesBannedAndUnbounded(t *testing.T) { + store := newTestStore(t) + ctx := context.Background() + now := time.Now() + + enroll := func(id string, expiresAt time.Time) { + t.Helper() + if err := store.EnrollNode(ctx, &EnrolledNode{ + PeerID: id, PublicKey: []byte("pub"), Biscuit: []byte("b"), Role: api.RoleNode, + EnrollmentType: "OIDC", EnrolledAt: now.Add(-48 * time.Hour), ExpiresAt: expiresAt, + }); err != nil { + t.Fatalf("enroll %s: %v", id, err) + } + } + enroll("lapsed-long-ago", now.Add(-24*time.Hour)) + enroll("lapsed-just-now", now.Add(-time.Minute)) + enroll("still-admitted", now.Add(time.Hour)) + enroll("lapsed-but-banned", now.Add(-24*time.Hour)) + enroll("never-bounded", time.Time{}) + if err := store.SetNodeBanned(ctx, "lapsed-but-banned", true); err != nil { + t.Fatal(err) + } + + deleted, err := store.DeleteExpiredNodes(ctx, now.Add(-time.Hour)) + if err != nil { + t.Fatalf("DeleteExpiredNodes: %v", err) + } + if deleted != 1 { + t.Errorf("deleted %d rows, want 1", deleted) + } + + remaining, err := store.ListNodes(ctx) + if err != nil { + t.Fatal(err) + } + got := map[string]bool{} + for _, n := range remaining { + got[n.PeerID] = true + } + for _, want := range []string{"lapsed-just-now", "still-admitted", "lapsed-but-banned", "never-bounded"} { + if !got[want] { + t.Errorf("%s was deleted", want) + } + } + if got["lapsed-long-ago"] { + t.Error("lapsed-long-ago survived") + } + if banned, _ := store.IsNodeBanned(ctx, "lapsed-but-banned"); !banned { + t.Error("the ban did not survive the sweep") + } +} + func TestBootstrapTokenRoundTripsEveryField(t *testing.T) { store := newTestStore(t) ctx := context.Background() diff --git a/internal/storage/sql_store.go b/internal/storage/sql_store.go index d13051c2..6c3f01fa 100644 --- a/internal/storage/sql_store.go +++ b/internal/storage/sql_store.go @@ -892,6 +892,17 @@ func (s *SQLStore) IsIdentityBanned(ctx context.Context, identity string) (bool, return true, nil } +// DeleteExpiredNodes implements Store. +func (s *SQLStore) DeleteExpiredNodes(ctx context.Context, before time.Time) (int64, error) { + // A zero ExpiresAt round-trips as a negative UnixMilli, hence > 0. + query := s.rebind(`DELETE FROM nodes WHERE banned = ? AND expires_at > 0 AND expires_at < ?`) + res, err := s.db.ExecContext(ctx, query, false, before.UnixMilli()) + if err != nil { + return 0, err + } + return res.RowsAffected() +} + // ListBannedPeerIDs implements Store. func (s *SQLStore) ListBannedPeerIDs(ctx context.Context) ([]string, error) { query := s.rebind(`SELECT peer_id FROM nodes WHERE banned = ?`) diff --git a/internal/storage/storage.go b/internal/storage/storage.go index 46554fdc..48103a99 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -287,6 +287,14 @@ type Store interface { // ListNodes retrieves all enrolled nodes. ListNodes(ctx context.Context) ([]EnrolledNode, error) + // DeleteExpiredNodes removes enrolled nodes whose session expired before + // the given time and that are not banned, returning how many went. A + // lapsed session already refuses every refresh, so the row only records + // that the node once existed; a ban, by contrast, is the row, and must + // outlive the session it was placed on. Nodes with no session bound + // (zero ExpiresAt) are never removed. + DeleteExpiredNodes(ctx context.Context, before time.Time) (int64, error) + // ListBootstrapTokens retrieves all bootstrap tokens. ListBootstrapTokens(ctx context.Context) ([]BootstrapToken, error) diff --git a/site/content/docs/integrations/vscode-copilot.md b/site/content/docs/integrations/vscode-copilot.md index aad1d6f6..2100df91 100644 --- a/site/content/docs/integrations/vscode-copilot.md +++ b/site/content/docs/integrations/vscode-copilot.md @@ -162,8 +162,6 @@ Calls `discover_remote_services` with `{"type":"mcp"}`: ```json [ - {"peer_id": "12D3KooWQ1hk…veSLS", "srv_name": "dummy-http", - "srv_description": "Canary HTTP tool (k8s agnhost)"}, {"peer_id": "12D3KooWFQrX…9Uwe1V", "srv_name": "everything", "srv_description": "MCP everything test server (tools, resources, prompts)"}, {"peer_id": "12D3KooWAjWy…RcZcs", "srv_name": "everything", @@ -171,10 +169,10 @@ Calls `discover_remote_services` with `{"type":"mcp"}`: ] ``` -Note `everything` appearing twice under different peers, and `dummy-http` three -times. Service names are not unique across the mesh and were never meant to be — -the `peer_id` is the identity. Any step that remembers "the everything service" -without remembering which peer will eventually talk to the wrong one. +Note `everything` appearing twice under different peers. Service names are +not unique across the mesh and were never meant to be — the `peer_id` is the +identity. Any step that remembers "the everything service" without remembering +which peer will eventually talk to the wrong one. ### Find tools, and read the failures @@ -185,7 +183,7 @@ failures in one array: ```json [ - {"peer_id": "12D3KooWQ1hk…veSLS", "tool_name": "mcp://dummy-http", + {"peer_id": "12D3KooWFQrX…9Uwe1V", "tool_name": "mcp://everything", "error": "failed to connect: failed to connect client: calling \"initialize\": EOF"}, {"peer_id": "12D3KooWAjWy…RcZcs", "tool_name": "mcp://everything/get-sum", "description": "Returns the sum of two numbers"}, @@ -194,7 +192,7 @@ failures in one array: ] ``` -Discovery is best-effort per peer. Three peers advertising `dummy-http` were +Discovery is best-effort per peer. One of the two `everything` peers was reachable enough to be listed but failed at `initialize`, and that is reported as an `error` field on the entry rather than failing the whole call. A partly broken mesh returns a partly populated array, so it is worth checking whether diff --git a/site/content/docs/user/control-plane-configuration.md b/site/content/docs/user/control-plane-configuration.md index 7fd625ed..551d9aac 100644 --- a/site/content/docs/user/control-plane-configuration.md +++ b/site/content/docs/user/control-plane-configuration.md @@ -28,6 +28,9 @@ The Control Plane is responsible for bridging user identities from trusted OIDC | `--key-rotation-interval` | *None* | `24h` | Key rotation interval (e.g. `24h`). `0s` disables rotation. | | `--key-grace-period` | *None* | `1h` | How long a rotated-out signing key stays accepted. Once it is retired, Biscuits it signed can no longer be verified or refreshed; see [Signing-Key Retirement and Recovery](#signing-key-retirement-and-recovery). | | `--lease-duration` | *None* | `15m` | Router lease registration TTL. | +| `--node-retention` | *None* | `720h` (30 days) | How long an enrolled node's record is kept after its session expired before it is deleted. Every restart without a persistent data directory enrolls a fresh identity, so without this the node table only grows. Banned nodes are always kept. `0s` keeps every record forever. | + +The Control Plane serves Prometheus metrics on `/metrics` alongside the API: mesh size (`sam_control_plane_mesh_connected_peers`, `sam_control_plane_routers_active`, `sam_control_plane_enrolled_nodes{role,state}`), per-route request counts and latencies, and the Go runtime. --- @@ -49,6 +52,7 @@ The Router is a dedicated GossipSub helper that maintains stable network address | `--keys-sync-interval` | `5m` | Key synchronization polling interval. `GET /keys` returns the valid key set signed by every key in it; the router (and nodes, at start-up) accept the set only if one of those signatures verifies under a key they already trust, so the first key always comes from enrollment and a rotation is learned from the retiring key. | | `--lease-renew-interval` | `300s` | Lease renewal registration interval. | | `--allow-loopback` | `false` | Allow loopback and link-local addresses for discovery (development only). | +| `--metrics-addr` | *None* | Serve Prometheus `/metrics`, `/healthz` and `/readyz` on this plain HTTP address (e.g. `0.0.0.0:9090`). Off by default: nothing on it is authenticated, so keep it off the libp2p ports and inside the cluster. `/readyz` turns `200` once the router is enrolled and its libp2p host is online. The same flag exists on `sam-node run`. | --- diff --git a/tests/e2e/canary_manifests.bats b/tests/e2e/canary_manifests.bats index d9100138..eb4bad0b 100644 --- a/tests/e2e/canary_manifests.bats +++ b/tests/e2e/canary_manifests.bats @@ -61,16 +61,16 @@ render() { # unservable reports whether every error is this cluster's missing CRDs rather # than the manifest's fault. The control plane template carries GKE Gateway and -# HealthCheckPolicy resources, which no kind cluster serves and which nothing -# here can install, so failing on them would only teach people to ignore this -# test. +# HealthCheckPolicy resources, and the monitoring template Managed Prometheus +# PodMonitorings, which no kind cluster serves and which nothing here can +# install, so failing on them would only teach people to ignore this test. # # The kinds are named rather than matched on "no matches for kind", because a # Deployment declared against a retired apiVersion fails with that same wording # -- which made an earlier version of this test pass a manifest whose # apps/v1beta1 would have broken the rollout it exists to protect. unservable() { - local kinds='Gateway|HTTPRoute|HealthCheckPolicy|GCPBackendPolicy' + local kinds='Gateway|HTTPRoute|HealthCheckPolicy|GCPBackendPolicy|PodMonitoring' ! grep -qvE "no matches for kind \"(${kinds})\"|ensure CRDs are installed|^[[:space:]]*$" <<<"$1" }