egress: add the egress demo and e2e coverage

demos/egress is a small Actor that fetches a URL it is given and echoes the
upstream status and body back, which makes the egress path observable from
outside the sandbox. hack/install-demo-egress.sh registers it as a
--deploy-demo-egress fixture and hack/verify-egress-demo.sh drives it and
checks the atenet-egress logs for the corresponding authorized CONNECT.

TestActorEgress in the networking suite covers the same path automatically:
it creates an Actor from the demo template, POSTs a fetch request through
atenet-router, and asserts 200. The suite's actor helper is parameterised by
template so the ingress test keeps using the counter fixture.
This commit is contained in:
Lior Lieberman
2026-08-10 23:28:07 -07:00
committed by Bowei Du
parent dd764c1e7c
commit 2f05a92f59
20 changed files with 1017 additions and 107 deletions
+4
View File
@@ -96,6 +96,10 @@ jobs:
run: hack/run-microvm-demo-kind.sh --ateapi-client-auth=${{ matrix.ateapi-client-auth }}
- name: Deploy gVisor counter demo
run: hack/install-ate-kind.sh --deploy-demo-counter
- name: Deploy egress demo
# TestActorEgress in the networking suite builds its Actor from the
# ate-demo-egress/egress ActorTemplate, so the fixture has to exist before.
run: hack/install-ate-kind.sh --deploy-demo-egress
- name: Wait for micro-VM golden snapshot
run: |
kubectl --context kind-kind wait --for=condition=Ready \
+1 -1
View File
@@ -57,7 +57,7 @@ build-atenet:
.PHONY: build-demos
build-demos:
$(KO) build --ldflags="$(LDFLAGS)" ./demos/counter
$(KO) build --ldflags="$(LDFLAGS)" ./demos/counter ./demos/egress
.PHONY: test
test:
-8
View File
@@ -100,14 +100,6 @@ func (s *ExtProcServer) Process(stream extprocv3.ExternalProcessor_ProcessServer
// the accepting listener, not by anything in the request itself
// (see isEgressRequest).
//
// SEE(lior): main grew a ResumeOutcome return on
// handleRequestHeaders while this branch was out. Rather than
// splitting the dispatch into two separately-typed call sites, the
// egress handler was widened to the same signature and returns
// ResumeOutcomeNone — egress requires an already-RUNNING actor and
// never resumes one, so "none" is accurate, and it keeps the
// route-duration metric's resume label a closed set (no new empty
// value) across both directions.
handle := s.handleRequestHeaders
if isEgressRequest(req) {
handle = s.handleEgressRequestHeaders
+36 -26
View File
@@ -36,19 +36,29 @@ import (
)
const (
// EgressListenerName is the Envoy listener that terminates actor egress
// CONNECTs. It must stay in sync with the listener name in
// manifests/ate-install/ateway-egress.yaml.
EgressListenerName = "egress"
// ListenerNameAttribute is the CEL attribute carrying the name of the
// listener that accepted the request. The egress Envoy asks for it via
// EgressFilterChainName is the Envoy filter chain that terminates actor
// egress CONNECTs. It must stay in sync with the filter chain name in
// manifests/ate-install/atenet-egress.yaml.
EgressFilterChainName = "egress"
// FilterChainNameAttribute is the CEL attribute carrying the name of the
// filter chain that accepted the request. The egress Envoy asks for it via
// request_attributes on its ext_proc filter.
ListenerNameAttribute = "xds.listener_name"
//
// SEE(lior): this was xds.listener_name, which reads more naturally but
// which Envoy 1.34 cannot parse: it logs "error parsing cel expression
// xds.listener_name" at trace level, then sends the ProcessingRequest with
// an empty attributes map rather than failing config load. Because an
// absent attribute means "ingress" (the fail-safe direction), every egress
// CONNECT silently took the ingress path and 404'd on the actor DNS name
// parse. xds.filter_chain_name parses on the same Envoy build and is
// equally Envoy-asserted, so the trust model is unchanged.
FilterChainNameAttribute = "xds.filter_chain_name"
// forwardedClientCertHeader is the header Envoy fills in with details of
// the mTLS peer, including the PEM chain it validated. The egress listener
// sets forward_client_cert_details: SANITIZE_SET, so whatever a client
// sends under this name is discarded and replaced by Envoy's own value.
// the mTLS peer, including the PEM chain it validated. The egress filter
// chain sets forward_client_cert_details: SANITIZE_SET, so whatever a
// client sends under this name is discarded and replaced by Envoy's own
// value.
//
// This is the only channel that can carry a whole certificate to ext_proc:
// the CEL request attributes Envoy exposes (subject, SANs, SHA-256 digest)
@@ -61,31 +71,31 @@ const (
)
// isEgressRequest reports whether an ext_proc RequestHeaders callback arrived on
// the egress gateway's listener rather than on an ingress listener. This lets
// one ext_proc server handle both directions off the same stream.
// the egress gateway's filter chain rather than on an ingress one. This lets one
// ext_proc server handle both directions off the same stream.
//
// Dispatch is by listener, not by :method, because the two handlers apply
// Dispatch is by filter chain, not by :method, because the two handlers apply
// opposite trust models: on egress the actor identity comes from a client
// certificate Envoy validated, while on ingress every request header is
// unauthenticated client input. Keying on :method would let any external client
// sending CONNECT select the egress handler and use its denial messages as an
// actor-existence and status oracle. Envoy asserts the listener name; the
// actor-existence and status oracle. Envoy asserts the filter chain name; the
// request cannot influence it.
//
// An unrecognized or absent attribute means ingress, the fail-safe direction: an
// egress request misrouted to the ingress handler fails to parse as an actor DNS
// name and 404s, whereas the reverse leaks control-plane state.
func isEgressRequest(req *extprocv3.ProcessingRequest) bool {
return listenerName(req) == EgressListenerName
return filterChainName(req) == EgressFilterChainName
}
// listenerName returns the xds.listener_name attribute Envoy attached to the
// request, or "" when the listener did not request the attribute. The
// filterChainName returns the xds.filter_chain_name attribute Envoy attached to
// the request, or "" when the listener did not request the attribute. The
// attributes map is keyed by the ext_proc filter's name within the HCM chain,
// which we do not want to hardcode here, so scan every entry.
func listenerName(req *extprocv3.ProcessingRequest) string {
func filterChainName(req *extprocv3.ProcessingRequest) string {
for _, attrs := range req.GetAttributes() {
if v, ok := attrs.GetFields()[ListenerNameAttribute]; ok {
if v, ok := attrs.GetFields()[FilterChainNameAttribute]; ok {
return v.GetStringValue()
}
}
@@ -106,19 +116,19 @@ func listenerName(req *extprocv3.ProcessingRequest) string {
// with a single branch. The (target, tmplNs, tmplName) results are unused for
// egress and returned empty.
//
// SEE(lior): the trailing ResumeOutcome exists only to match
// handleRequestHeaders after main added it. Egress never resumes an actor — it
// requires one already RUNNING — so every path returns ResumeOutcomeNone.
// The trailing ResumeOutcome exists only to match handleRequestHeaders.
// Egress never resumes an actor — it requires one already RUNNING -
// so every path returns ResumeOutcomeNone.
func (s *ExtProcServer) handleEgressRequestHeaders(
ctx context.Context,
reqHeaders *extprocv3.HttpHeaders,
) (*extprocv3.HeadersResponse, *requestMetadata, string, string, string, ResumeOutcome, error) {
metadata := newRequestMetadata(reqHeaders.Headers.GetHeaders())
// Dispatch is by listener, so reaching here means the egress listener
// accepted the request. That listener only routes CONNECT (its sole route
// is a connect_matcher), so anything else is a config drift rather than a
// client the gateway should tunnel for.
// Dispatch is by filter chain, so reaching here means the egress chain
// accepted the request. That chain only routes CONNECT (its sole route is a
// connect_matcher), so anything else is config drift rather than a client
// the gateway should tunnel for.
if !strings.EqualFold(metadata.method, "CONNECT") {
return nil, metadata, "", "", "", ResumeOutcomeNone, newReqError(envoy_type.StatusCode_MethodNotAllowed,
"egress denied: expected CONNECT, got %q", metadata.method)
@@ -45,9 +45,9 @@ import (
)
// connectRequest builds a RequestHeaders ProcessingRequest for an egress
// CONNECT, optionally attributed to a listener. filterKey is the ext_proc
// CONNECT, optionally attributed to a filter chain. filterKey is the ext_proc
// filter name Envoy keys the attributes map by.
func connectRequest(filterKey, listener string) *extprocv3.ProcessingRequest {
func connectRequest(filterKey, chain string) *extprocv3.ProcessingRequest {
req := &extprocv3.ProcessingRequest{
Request: &extprocv3.ProcessingRequest_RequestHeaders{
RequestHeaders: &extprocv3.HttpHeaders{
@@ -60,11 +60,11 @@ func connectRequest(filterKey, listener string) *extprocv3.ProcessingRequest {
},
},
}
if listener != "" {
if chain != "" {
req.Attributes = map[string]*structpb.Struct{
filterKey: {
Fields: map[string]*structpb.Value{
ListenerNameAttribute: structpb.NewStringValue(listener),
FilterChainNameAttribute: structpb.NewStringValue(chain),
},
},
}
@@ -76,19 +76,19 @@ func TestIsEgressRequest(t *testing.T) {
tests := []struct {
name string
filterKey string
listener string
chain string
want bool
}{
{
name: "egress listener",
name: "egress filter chain",
filterKey: "envoy.filters.http.ext_proc",
listener: EgressListenerName,
chain: EgressFilterChainName,
want: true,
},
{
name: "egress listener under a renamed filter",
name: "egress filter chain under a renamed filter",
filterKey: "some.custom.ext_proc.name",
listener: EgressListenerName,
chain: EgressFilterChainName,
want: true,
},
{
@@ -98,33 +98,33 @@ func TestIsEgressRequest(t *testing.T) {
// exists and is running.
name: "CONNECT on the ingress HTTP listener",
filterKey: "envoy.filters.http.ext_proc",
listener: IngressHTTPListener,
chain: IngressHTTPListener,
want: false,
},
{
name: "CONNECT on the ingress HTTPS listener",
filterKey: "envoy.filters.http.ext_proc",
listener: IngressHTTPSListener,
chain: IngressHTTPSListener,
want: false,
},
{
// A listener that never requested the attribute falls back to
// ingress, the fail-safe direction.
name: "no attributes at all",
listener: "",
want: false,
name: "no attributes at all",
chain: "",
want: false,
},
{
name: "unrecognised listener",
name: "unrecognised filter chain",
filterKey: "envoy.filters.http.ext_proc",
listener: "some-other-listener",
chain: "some-other-chain",
want: false,
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
if got := isEgressRequest(connectRequest(tc.filterKey, tc.listener)); got != tc.want {
if got := isEgressRequest(connectRequest(tc.filterKey, tc.chain)); got != tc.want {
t.Errorf("isEgressRequest() = %v, want %v", got, tc.want)
}
})
@@ -132,17 +132,17 @@ func TestIsEgressRequest(t *testing.T) {
}
// A request the client dresses up to look like egress must not be enough: only
// the Envoy-asserted listener name selects the egress handler.
// the Envoy-asserted filter chain name selects the egress handler.
func TestIsEgressRequestIgnoresClientSuppliedAttributeHeader(t *testing.T) {
req := connectRequest("envoy.filters.http.ext_proc", IngressHTTPListener)
rh := req.GetRequestHeaders().GetHeaders()
rh.Headers = append(rh.Headers,
&corev3.HeaderValue{Key: ListenerNameAttribute, RawValue: []byte(EgressListenerName)},
&corev3.HeaderValue{Key: "x-envoy-listener-name", RawValue: []byte(EgressListenerName)},
&corev3.HeaderValue{Key: FilterChainNameAttribute, RawValue: []byte(EgressFilterChainName)},
&corev3.HeaderValue{Key: "x-envoy-filter-chain-name", RawValue: []byte(EgressFilterChainName)},
)
if isEgressRequest(req) {
t.Error("isEgressRequest() = true for a client-forged listener header, want false")
t.Error("isEgressRequest() = true for a client-forged filter chain header, want false")
}
}
@@ -290,7 +290,7 @@ func xfccHeaderDER(chain ...[]byte) string {
for _, der := range chain {
_ = pem.Encode(&buf, &pem.Block{Type: "CERTIFICATE", Bytes: der})
}
return fmt.Sprintf(`By=spiffe://cluster.local/ns/ate-system/sa/ateway-egress;Hash=abc123;Chain="%s"`,
return fmt.Sprintf(`By=spiffe://cluster.local/ns/ate-system/sa/atenet-egress;Hash=abc123;Chain="%s"`,
url.PathEscape(buf.String()))
}
@@ -539,7 +539,7 @@ func TestHandleEgressRequestHeadersRejectsBadCertificates(t *testing.T) {
{
name: "XFCC without a Chain value",
xfcc: func(*testing.T) string {
return `By=spiffe://cluster.local/ns/ate-system/sa/ateway-egress;Hash=abc123`
return `By=spiffe://cluster.local/ns/ate-system/sa/atenet-egress;Hash=abc123`
},
want: envoy_type.StatusCode_Forbidden,
},
+148
View File
@@ -0,0 +1,148 @@
# Egress Demo — Pluggable Egress Networking
This demo shows an Actor's outbound traffic being **transparently tunneled through an
egress gateway** and **authenticated by actor identity**, end to end.
The Actor is a tiny service that accepts `{"url":"..."}`, performs an HTTP `GET`, and returns
the upstream response. The Actor believes it is dialing plain HTTP directly — but its egress is
intercepted and carried over mTLS to a gateway that verifies who is making the request.
## What it demonstrates
```
┌──────────────── ateom worker pod ─────────────────┐
│ Actor (gVisor) │
│ GET http://<dst-ip>:80/ (plain HTTP) │
│ │ │
│ ▼ nftables REDIRECT │
│ atunnel egress ──(mTLS with the actor's own │
│ │ certificate + bare CONNECT) │
└────────┼────────────────────────────────────────────┘
▼
┌──────────── atenet-egress pod ───────────────────┐
│ Envoy egress gateway │
│ • downstream mTLS, trusted_ca = actor-id CA │
│ • terminates HTTP CONNECT │
│ • ext_proc ──(localhost)──► atenet router (ext_proc sidecar)
│ (forwards the peer chain │ verify chain + ActorIdentity extension
│ as x-forwarded-client-cert)│ GetActor → UID must match, must be RUNNING
│ • dynamic_forward_proxy │ allow / deny 403
│ │
└───────────┼───────────────────────────────────────┘
▼
real destination (the CONNECT authority, an IP:port)
```
1. **Guide 1 — gateway accepts CONNECT + mTLS.** `atenet-egress` is an Envoy dynamic-forward-proxy
that terminates the actor's mTLS `CONNECT` and tunnels to the requested destination.
2. **Guide 2 — transparent interception.** `nftables` REDIRECTs actor TCP egress into `atunnel`,
which wraps it in mTLS + `CONNECT`.
3. **Guide 3 — HTTP-only actors, identity carried by the certificate.** The Actor only dials plain
HTTP. atunnel presents the actor's own certificate — minted per actor by ateapi off the
actor-identity CA, carrying an `ActorIdentity` X.509 extension — and sends a bare `CONNECT`
with no identity headers at all.
4. **Identity authentication.** Envoy requires a client certificate signed by the actor-identity
CA, so a non-actor client is refused at the handshake. It then forwards the verified chain to
`ext_proc` as `x-forwarded-client-cert`, and the **atenet router** (co-located in the gateway
pod as an ext_proc sidecar, the same ext_proc code used for ingress) re-verifies the chain,
requires exactly one `ActorIdentity` extension with `purpose: atunnel`, and calls the ate API
(`GetActor`). It returns **403** unless the certified **UID** matches a real, `RUNNING` actor.
This mirrors the ingress gateway's Envoy + ext_proc co-location; a standalone/shared ext_proc
is a future step.
## Components
- **Egress app (`main.go`)** — the Actor: `POST /` with `{"url":"..."}` → fetches it → returns
status + body.
- **Egress gateway** — `manifests/ate-install/atenet-egress.yaml`. One pod, two containers:
an Envoy (`envoy`) and the atenet router ext_proc (`ext-proc`), called over localhost.
- **Egress opt-in** — `ate-api-server --egress-gateway-address=atenet-egress.ate-system.svc:443`
(set in `manifests/ate-install/ate-api-server.yaml`). ateapi stamps the address onto every
atelet `Run`/`Restore`, which turns on tunneled egress cluster-wide.
- **Actor-identity trust** — the gateway mounts the `actor-id-ca-certs` Secret, a cert-only copy of
the actor-identity CA root that `hack/install-ate.sh` derives from `actor-id-ca-pool` (which also
holds the CA signing key and is deliberately *not* mounted here).
## Prerequisites
- A kind cluster with Agent Substrate installed (`hack/create-kind-cluster.sh` then
`hack/install-ate-kind.sh --deploy-ate-system`). Egress is enabled by the ateapi flag above.
- `ko`, `kubectl`, and `kubectl-ate` (`go install ./cmd/kubectl-ate`).
## Deploy the demo fixture
```bash
./hack/install-ate.sh --deploy-demo-egress
kubectl wait --for=condition=Ready actortemplate/egress -n ate-demo-egress --timeout=5m
```
## Run the automated test (easiest)
```bash
./demos/egress/test-egress.sh
```
It deploys an in-cluster HTTP target, creates & resumes an Actor, then asserts:
- **positive** — a real Actor's egress reaches the target (`HTTP 200`) *through the gateway*
(the target sees the gateway's IP as its client), and the gateway logs the CONNECT against the
actor's certificate SAN;
- **negative** — a pod holding a valid *pod* identity but no actor certificate cannot open a
tunnel at all: the gateway's `trusted_ca` is the actor-identity CA, so the mTLS handshake is
refused before any CONNECT is answered.
Add `--cleanup` to remove everything the script created.
## Manual walkthrough
```bash
# 1. An in-cluster target the Actor will fetch (any HTTP server works).
kubectl create namespace egress-target
kubectl -n egress-target create deployment whoami --image=traefik/whoami
kubectl -n egress-target expose deployment whoami --port=80
TARGET_IP=$(kubectl -n egress-target get svc whoami -o jsonpath='{.spec.clusterIP}')
# 2. Create and resume an Actor.
kubectl ate create atespace demo
kubectl ate create actor egress-demo -a demo --template ate-demo-egress/egress
kubectl ate resume actor egress-demo -a demo # wait for STATUS_RUNNING
# 3. Drive the Actor's egress through the ingress gateway.
kubectl -n ate-system port-forward service/atenet-router 8000:80 &
curl -s -X POST http://localhost:8000/ \
-H 'Host: egress-demo.demo.actors.resources.substrate.ate.dev' \
-H 'Content-Type: application/json' \
-d "{\"url\":\"http://${TARGET_IP}:80/\"}"
```
### What to observe
```bash
# The egress gateway logs each tunneled CONNECT against the verified peer certificate:
kubectl -n ate-system logs deploy/atenet-egress | grep '\[egress\]'
# [egress] authority=<TARGET_IP>:80 peer_san=spiffe://substrate-actor.local/atespace/demo/actor/egress-demo … code=200 …
# The co-located ext_proc sidecar logs the identity decision, including the UID it authorized on:
kubectl -n ate-system logs deploy/atenet-egress -c ext-proc | grep -i 'egress identity\|egress denied'
# egress identity authenticated atespace=demo actor=egress-demo actorUid=… status=STATUS_RUNNING
```
The `whoami` body shows `RemoteAddr: <atenet-egress pod IP>` — proof the request egressed
*through* the gateway rather than directly.
## Notes / limitations
- This milestone **authenticates** identity (is this a real, running actor?). **Authorizing**
egress by destination and injecting upstream credentials/tokens is a follow-up, implemented in
the same `ext_proc` (policy API TBD).
- Identity comes entirely from the actor certificate: the atespace, actor name, and UID are read
out of the `ActorIdentity` extension and the UID is matched against the live actor, so a
certificate cannot survive its actor being deleted and recreated under the same name. Nothing
the actor can write into the CONNECT contributes to the decision.
## Cleanup
```bash
./demos/egress/test-egress.sh --cleanup
./hack/install-ate.sh --delete-demo-egress
```
+56
View File
@@ -0,0 +1,56 @@
# 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.
apiVersion: v1
kind: Namespace
metadata:
name: ate-demo-egress
---
apiVersion: ate.dev/v1alpha1
kind: WorkerPool
metadata:
name: egress
namespace: ate-demo-egress
labels:
workload: egress
spec:
replicas: 2
ateomImage: ko://github.com/agent-substrate/substrate/cmd/ateom-gvisor
---
apiVersion: ate.dev/v1alpha1
kind: ActorTemplate
metadata:
name: egress
namespace: ate-demo-egress
spec:
pauseImage: "registry.k8s.io/pause:3.10.2@sha256:f548e0e8e3dc1896ca956272154dde3314e8cc4fde0a57577ee9fa1c63f5baf4"
containers:
- name: egress
image: ko://github.com/agent-substrate/substrate/demos/egress
command: ["/ko-app/egress"]
readyz:
httpGet:
path: /readyz
port: 80
workerSelector:
matchLabels:
workload: egress
snapshotsConfig:
onPause: Full
onCommit: Full
location: gs://${BUCKET_NAME}/ate-demo-egress/
+125
View File
@@ -0,0 +1,125 @@
// 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.
// Command egress is a small HTTP service for demonstrating per-Actor egress
// policy. It accepts a URL, fetches it, and returns the upstream response.
package main
import (
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"net/url"
"os"
"time"
)
const (
listenAddress = ":80"
maxRequestBody = 64 << 10
maxResponseBody = 1 << 20
requestTimeout = 15 * time.Second
)
type fetchRequest struct {
URL string `json:"url"`
}
type fetchResponse struct {
StatusCode int `json:"statusCode,omitempty"`
Body string `json:"body,omitempty"`
Error string `json:"error,omitempty"`
}
func main() {
slog.SetDefault(slog.New(slog.NewJSONHandler(os.Stdout, nil)))
client := &http.Client{Timeout: requestTimeout}
slog.Info("starting egress demo", "address", listenAddress)
if err := http.ListenAndServe(listenAddress, newHandler(client)); err != nil {
slog.Error("egress demo stopped", "error", err)
os.Exit(1)
}
}
func newHandler(client *http.Client) http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("/readyz", func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
_, _ = io.WriteString(w, "ok\n")
})
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
w.Header().Set("Allow", http.MethodPost)
writeJSON(w, http.StatusMethodNotAllowed, fetchResponse{Error: "method must be POST"})
return
}
var input fetchRequest
decoder := json.NewDecoder(http.MaxBytesReader(w, r.Body, maxRequestBody))
if err := decoder.Decode(&input); err != nil {
writeJSON(w, http.StatusBadRequest, fetchResponse{Error: fmt.Sprintf("invalid JSON payload: %v", err)})
return
}
if err := validateURL(input.URL); err != nil {
writeJSON(w, http.StatusBadRequest, fetchResponse{Error: err.Error()})
return
}
outbound, err := http.NewRequestWithContext(r.Context(), http.MethodGet, input.URL, nil)
if err != nil {
writeJSON(w, http.StatusBadRequest, fetchResponse{Error: fmt.Sprintf("invalid URL: %v", err)})
return
}
if traceparent := r.Header.Get("traceparent"); traceparent != "" {
outbound.Header.Set("traceparent", traceparent)
}
response, err := client.Do(outbound)
if err != nil {
writeJSON(w, http.StatusBadGateway, fetchResponse{Error: fmt.Sprintf("request failed: %v", err)})
return
}
defer response.Body.Close()
body, err := io.ReadAll(io.LimitReader(response.Body, maxResponseBody))
if err != nil {
writeJSON(w, http.StatusBadGateway, fetchResponse{Error: fmt.Sprintf("reading response: %v", err)})
return
}
writeJSON(w, response.StatusCode, fetchResponse{StatusCode: response.StatusCode, Body: string(body)})
})
return mux
}
func validateURL(raw string) error {
parsed, err := url.Parse(raw)
if err != nil {
return fmt.Errorf("invalid URL: %w", err)
}
if parsed.Scheme != "http" && parsed.Scheme != "https" {
return fmt.Errorf("URL scheme must be http or https")
}
if parsed.Hostname() == "" {
return fmt.Errorf("URL must include a hostname")
}
return nil
}
func writeJSON(w http.ResponseWriter, status int, response fetchResponse) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(response)
}
+109
View File
@@ -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 main
import (
"encoding/json"
"errors"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
func TestFetch(t *testing.T) {
const traceparent = "00-0123456789abcdef0123456789abcdef-0123456789abcdef-01"
client := &http.Client{Transport: roundTripFunc(func(r *http.Request) (*http.Response, error) {
if r.Method != http.MethodGet {
t.Errorf("upstream method = %s, want GET", r.Method)
}
if got := r.Header.Get("traceparent"); got != traceparent {
t.Errorf("upstream traceparent = %q, want %q", got, traceparent)
}
return &http.Response{
StatusCode: http.StatusTeapot,
Body: io.NopCloser(strings.NewReader("hello from upstream")),
Header: make(http.Header),
}, nil
})}
payload, err := json.Marshal(fetchRequest{URL: "https://allowed.example/"})
if err != nil {
t.Fatal(err)
}
recorder := httptest.NewRecorder()
request := httptest.NewRequest(http.MethodPost, "/", strings.NewReader(string(payload)))
request.Header.Set("traceparent", traceparent)
newHandler(client).ServeHTTP(recorder, request)
if recorder.Code != http.StatusTeapot {
t.Fatalf("status = %d, want %d", recorder.Code, http.StatusTeapot)
}
var got fetchResponse
if err := json.NewDecoder(recorder.Body).Decode(&got); err != nil {
t.Fatalf("decoding response: %v", err)
}
if got.StatusCode != http.StatusTeapot || got.Body != "hello from upstream" {
t.Errorf("response = %+v", got)
}
}
func TestInvalidRequests(t *testing.T) {
tests := []struct {
name string
method string
body string
status int
}{
{name: "method", method: http.MethodGet, body: `{}`, status: http.StatusMethodNotAllowed},
{name: "malformed JSON", method: http.MethodPost, body: `{`, status: http.StatusBadRequest},
{name: "missing hostname", method: http.MethodPost, body: `{"url":"https:///path"}`, status: http.StatusBadRequest},
{name: "unsupported scheme", method: http.MethodPost, body: `{"url":"file:///etc/passwd"}`, status: http.StatusBadRequest},
}
handler := newHandler(http.DefaultClient)
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
recorder := httptest.NewRecorder()
request := httptest.NewRequest(test.method, "/", strings.NewReader(test.body))
handler.ServeHTTP(recorder, request)
if recorder.Code != test.status {
t.Errorf("status = %d, want %d; body = %s", recorder.Code, test.status, recorder.Body.String())
}
})
}
}
func TestOutboundFailure(t *testing.T) {
client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
return nil, errors.New("blocked")
})}
handler := newHandler(client)
recorder := httptest.NewRecorder()
request := httptest.NewRequest(http.MethodPost, "/", strings.NewReader(`{"url":"https://example.com/"}`))
handler.ServeHTTP(recorder, request)
if recorder.Code != http.StatusBadGateway {
t.Errorf("status = %d, want %d; body = %s", recorder.Code, http.StatusBadGateway, recorder.Body.String())
}
}
type roundTripFunc func(*http.Request) (*http.Response, error)
func (f roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) {
return f(request)
}
+196
View File
@@ -0,0 +1,196 @@
#!/usr/bin/env bash
# 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.
# End-to-end test for pluggable actor egress. Reproduces:
# * POSITIVE — a real, running Actor's plain-HTTP egress is transparently
# tunneled (nftables -> atunnel -> mTLS + CONNECT) through the Envoy egress
# gateway to an in-cluster target, and the gateway's ext_proc authenticates
# the actor certificate against the ate API (allowed, HTTP 200).
# * NEGATIVE — a pod holding a perfectly valid *pod* identity, but no actor
# certificate, cannot open a tunnel: the gateway's trusted_ca is the
# actor-identity CA, so the mTLS handshake itself is refused.
#
# Prerequisites: a substrate cluster with `--deploy-demo-egress` applied, plus
# kubectl and kubectl-ate on PATH. See demos/egress/README.md.
#
# Usage:
# demos/egress/test-egress.sh # run the tests
# demos/egress/test-egress.sh --cleanup # remove everything this script created
set -o errexit -o nounset -o pipefail
CTX="${KUBECTL_CONTEXT:-kind-kind}"
ATESPACE="${ATESPACE:-demo}"
ACTOR="${ACTOR:-egress-demo}"
TEMPLATE="${TEMPLATE:-ate-demo-egress/egress}"
TARGET_NS="${TARGET_NS:-egress-target}"
PROBE_POD="egress-identity-probe"
K="kubectl --context ${CTX}"
KATE="kubectl-ate --context ${CTX}"
log() { printf '\n\033[1;36m== %s\033[0m\n' "$*"; }
info() { printf ' %s\n' "$*"; }
pass() { printf '\033[1;32mPASS\033[0m %s\n' "$*"; }
fail() { printf '\033[1;31mFAIL\033[0m %s\n' "$*"; FAILED=1; }
FAILED=0
require() { command -v "$1" >/dev/null 2>&1 || { echo "missing required tool: $1"; exit 1; }; }
cleanup() {
log "cleanup"
${K} -n ate-system delete pod "${PROBE_POD}" --ignore-not-found --wait=false >/dev/null 2>&1 || true
${KATE} suspend actor "${ACTOR}" -a "${ATESPACE}" >/dev/null 2>&1 || true
${KATE} delete actor "${ACTOR}" -a "${ATESPACE}" >/dev/null 2>&1 || true
${K} delete namespace "${TARGET_NS}" --ignore-not-found --wait=false >/dev/null 2>&1 || true
info "done"
}
if [[ "${1:-}" == "--cleanup" ]]; then require kubectl; require kubectl-ate; cleanup; exit 0; fi
require kubectl
require kubectl-ate
trap '[[ "${KEEP:-}" == "1" ]] || cleanup' EXIT
log "preflight: egress gateway (Envoy + co-located ext_proc) is running"
${K} -n ate-system rollout status deployment/atenet-egress --timeout=120s
log "deploy an in-cluster HTTP target (whoami)"
${K} create namespace "${TARGET_NS}" >/dev/null 2>&1 || true
${K} -n "${TARGET_NS}" create deployment whoami --image=traefik/whoami >/dev/null 2>&1 || true
${K} -n "${TARGET_NS}" expose deployment whoami --port=80 >/dev/null 2>&1 || true
${K} -n "${TARGET_NS}" rollout status deployment/whoami --timeout=120s
TARGET_IP=$(${K} -n "${TARGET_NS}" get svc whoami -o jsonpath='{.spec.clusterIP}')
info "target ClusterIP = ${TARGET_IP}"
log "create + resume Actor ${ATESPACE}/${ACTOR}"
${KATE} create atespace "${ATESPACE}" >/dev/null 2>&1 || true
${KATE} create actor "${ACTOR}" -a "${ATESPACE}" --template "${TEMPLATE}" >/dev/null 2>&1 || true
${KATE} resume actor "${ACTOR}" -a "${ATESPACE}" >/dev/null 2>&1 || true
for _ in $(seq 1 30); do
${KATE} get actors -a "${ATESPACE}" 2>/dev/null | grep -q "STATUS_RUNNING" && break
sleep 3
done
${KATE} get actors -a "${ATESPACE}" 2>/dev/null | grep "${ACTOR}" || true
${KATE} get actors -a "${ATESPACE}" 2>/dev/null | grep -q "STATUS_RUNNING" || { echo "actor did not reach RUNNING"; exit 1; }
egress_log_since() { ${K} -n ate-system logs deployment/atenet-egress --tail=-1 2>/dev/null | grep '\[egress\]' | tail -n +"$(( $1 + 1 ))"; }
egress_log_count() { ${K} -n ate-system logs deployment/atenet-egress --tail=-1 2>/dev/null | grep -c '\[egress\]' || true; }
# wait_egress_log <before-count> <grep-pattern>: retry for log-shipping lag.
wait_egress_log() { for _ in $(seq 1 10); do if egress_log_since "$1" | grep -qE "$2"; then egress_log_since "$1" | grep -E "$2" | tail -1; return 0; fi; sleep 1; done; return 1; }
##############################################################################
log "POSITIVE — real Actor egress is tunneled through the gateway (expect 200)"
##############################################################################
BEFORE=$(egress_log_count)
${K} -n ate-system port-forward service/atenet-router 18099:80 >/tmp/egress-pf.log 2>&1 &
PF=$!; sleep 4
CODE=$(curl -s -o /tmp/egress-body.txt -w '%{http_code}' -X POST http://localhost:18099/ \
-H "Host: ${ACTOR}.${ATESPACE}.actors.resources.substrate.ate.dev" \
-H 'Content-Type: application/json' \
-d "{\"url\":\"http://${TARGET_IP}:80/\"}" || true)
kill "${PF}" >/dev/null 2>&1 || true
sleep 1
info "actor round-trip HTTP ${CODE}"
GW_IP=$(${K} -n ate-system get pod -l app=atenet-egress -o jsonpath='{.items[0].status.podIP}')
if [[ "${CODE}" == "200" ]]; then pass "actor fetched the target (HTTP 200)"; else fail "expected HTTP 200, got ${CODE}"; fi
if grep -q "RemoteAddr: ${GW_IP}" /tmp/egress-body.txt 2>/dev/null; then
pass "target saw the egress gateway (${GW_IP}) as its client — traffic went through the gateway"
else
info "target body RemoteAddr: $(grep -o 'RemoteAddr: [0-9.]*' /tmp/egress-body.txt 2>/dev/null || echo '?') (gateway IP ${GW_IP})"
fi
# The access log identifies the peer by its certificate SAN
# (spiffe://substrate-actor.local/atespace/<atespace>/actor/<name>), not by any
# header the actor could have written.
if LINE=$(wait_egress_log "${BEFORE}" "actor/${ACTOR}.*code=200"); then
pass "gateway logged the CONNECT: ${LINE}"
else
fail "gateway did not log an allowed CONNECT for ${ACTOR}"
fi
##############################################################################
log "NEGATIVE — a pod identity is not an actor identity (expect a refused handshake)"
##############################################################################
# SEE(lior): this used to be a 403 assertion. On the header-based design the
# gateway trusted any podidentity holder's mTLS and let ext_proc adjudicate the
# X-Ate-* headers it asserted, so a probe pod could reach the CONNECT and be
# denied there. The gateway now trusts only the actor-identity CA, so the same
# probe never gets past the TLS handshake — the denial moved a layer down and
# there is no CONNECT response to read a status code out of. Assert the
# handshake failure instead; it is the stronger property.
${K} apply -f - >/dev/null <<'YAML'
apiVersion: v1
kind: Pod
metadata:
name: egress-identity-probe
namespace: ate-system
spec:
containers:
- name: curl
image: curlimages/curl:latest
command: ["sleep", "600"]
volumeMounts:
- { name: podidentity, mountPath: /run/podidentity.podcert.ate.dev, readOnly: true }
- { name: servicedns, mountPath: /run/servicedns.podcert.ate.dev, readOnly: true }
volumes:
- name: podidentity
projected:
sources:
- podCertificate:
signerName: podidentity.podcert.ate.dev/identity
keyType: ECDSAP256
credentialBundlePath: credential-bundle.pem
- name: servicedns
projected:
sources:
- clusterTrustBundle:
signerName: servicedns.podcert.ate.dev/identity
labelSelector: { matchLabels: { podcert.ate.dev/canarying: live } }
path: trust-bundle.pem
YAML
${K} -n ate-system wait --for=condition=Ready pod/${PROBE_POD} --timeout=60s >/dev/null
BEFORE=$(egress_log_count)
# %{http_connect} carries the proxy's CONNECT response code, and stays 000 when
# the tunnel never opens. curl exits non-zero on a failed handshake, so capture
# both and require that no CONNECT was ever answered.
PROBE=$(${K} -n ate-system exec ${PROBE_POD} -- sh -c "curl -s -o /dev/null -w '%{http_connect}' \
--proxy-cacert /run/servicedns.podcert.ate.dev/trust-bundle.pem \
--proxy-cert /run/podidentity.podcert.ate.dev/credential-bundle.pem \
--proxy-key /run/podidentity.podcert.ate.dev/credential-bundle.pem \
--proxytunnel -x https://atenet-egress.ate-system.svc:443 http://${TARGET_IP}:80/; echo \" exit=\$?\"" || true)
CODE=${PROBE%% *}
info "pod-identity CONNECT attempt: http_connect=${CODE:-000}${PROBE#"${CODE}"}"
if [[ "${CODE}" != "200" ]]; then
pass "a pod identity cannot open an egress tunnel (no CONNECT succeeded)"
else
fail "expected the gateway to refuse a non-actor client certificate, but CONNECT returned 200"
fi
# A rejected handshake never becomes an HTTP request, so it produces no [egress]
# access-log line — only a connection-level TLS error. Report whatever the
# gateway logged so a real failure is diagnosable, but do not assert on it.
if LINE=$(egress_log_since "${BEFORE}" | tail -1); [[ -n "${LINE:-}" ]]; then
info "gateway egress log since the probe: ${LINE}"
else
info "no new [egress] access-log lines — the handshake was refused before HTTP, as expected"
fi
echo
if [[ "${FAILED}" == "0" ]]; then
printf '\033[1;32mALL CHECKS PASSED\033[0m — pluggable egress + identity authentication working.\n'
else
printf '\033[1;31mSOME CHECKS FAILED\033[0m\n'; exit 1
fi
+5 -9
View File
@@ -40,6 +40,7 @@ ATE_DEMOS=()
# Include demos.
source "${ROOT}"/hack/install-demo-counter.sh
source "${ROOT}"/hack/install-demo-egress.sh
source "${ROOT}"/hack/install-demo-sandbox.sh
source "${ROOT}"/hack/install-demo-claude-code-multiplex.sh
source "${ROOT}"/hack/install-demo-multi-template.sh
@@ -446,7 +447,7 @@ deploy_ate_system() {
run_kubectl rollout status deployment/ate-api-server -n ate-system --timeout=120s
run_kubectl rollout status deployment/ate-controller -n ate-system --timeout=120s
run_kubectl rollout status deployment/atenet-router -n ate-system --timeout=120s
run_kubectl rollout status deployment/ateway-egress -n ate-system --timeout=120s
run_kubectl rollout status deployment/atenet-egress -n ate-system --timeout=120s
run_kubectl rollout status statefulset/valkey-cluster -n ate-system --timeout=120s
run_kubectl rollout status daemonset/atelet -n ate-system --timeout=120s
}
@@ -521,15 +522,10 @@ deploy_atenet() {
router_manifest="$(render_atenet_router_manifest)"
echo "${router_manifest}" | run_kubectl apply -f -
run_ko apply -f manifests/ate-install/ateway-egress.yaml
run_ko apply -f manifests/ate-install/atenet-egress.yaml
run_ko apply -f manifests/ate-install/atenet-dns.yaml
run_kubectl rollout status deployment/atenet-router -n ate-system --timeout=120s
run_kubectl rollout status deployment/ateway-egress -n ate-system --timeout=120s
# SEE(lior): this branch also added a `deployment/atenet-dns` wait here, which
# main has since fixed to `deployment/dns` (the Deployment in atenet-dns.yaml
# is named "dns"; every other resource in that file is "atenet-dns", so the
# old name was NotFound on every successful deploy). Dropped this branch's
# duplicate in favour of main's corrected line.
run_kubectl rollout status deployment/atenet-egress -n ate-system --timeout=120s
run_kubectl rollout status deployment/dns -n ate-system --timeout=120s
}
@@ -648,7 +644,7 @@ delete_atenet() {
run_kubectl delete --ignore-not-found -f manifests/ate-install/atenet-router.yaml
run_kubectl delete --ignore-not-found \
-f manifests/ate-install/components/agentgateway/configmap.yaml
run_kubectl delete --ignore-not-found -f manifests/ate-install/ateway-egress.yaml
run_kubectl delete --ignore-not-found -f manifests/ate-install/atenet-egress.yaml
run_kubectl delete --ignore-not-found -f manifests/ate-install/atenet-dns.yaml
}
+51
View File
@@ -0,0 +1,51 @@
#!/usr/bin/env bash
# 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.
#
# This is sourced as part of install-ate.sh. Do not run directly.
ATE_DEMOS+=(demo-egress) # register demo-egress
demo-egress_cmdline() {
case "${1}" in
--deploy-demo-egress) demo-egress_deploy ;;
--delete-demo-egress) demo-egress_delete ;;
*)
return 1
;;
esac
return 0
}
demo-egress_deploy() {
log_step "demo-egress_deploy"
ensure_crds
sed "s|\${BUCKET_NAME}|${BUCKET_NAME}|g" demos/egress/egress.yaml.tmpl \
| run_ko apply -f -
log_step "Waiting for egress demo to be ready..."
# The WorkerPool controller names the Deployment after the WorkerPool
# ("egress"), the same way demo-counter gets "deployment/counter". The old
# "egress-deployment" name was NotFound on every successful deploy.
run_kubectl rollout status deployment/egress -n ate-demo-egress --timeout=300s
run_kubectl wait --for=condition=Ready actortemplate/egress -n ate-demo-egress --timeout=300s
}
demo-egress_delete() {
log_step "demo-egress_delete"
delete_demo_actors ate-demo-egress egress
sed "s|\${BUCKET_NAME}|${BUCKET_NAME}|g" demos/egress/egress.yaml.tmpl \
| run_kubectl delete --ignore-not-found -f -
}
+66
View File
@@ -0,0 +1,66 @@
#!/usr/bin/env bash
# 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.
#
# Preconditions:
# hack/create-kind-cluster.sh
# hack/install-ate-kind.sh --deploy-demo-egress
#
# The egress demo Actor accepts {"url":"..."} and performs an HTTP GET. With
# egress turned on (ate-api-server --egress-gateway-address, which ateapi stamps
# onto every atelet Run/Restore), the Actor's outbound TCP is nftables-REDIRECTed
# into atunnel, wrapped in mTLS + HTTP CONNECT, and sent to the Envoy egress
# gateway, which terminates CONNECT and tunnels to the real destination. This
# script drives that path and shows the gateway's access log proving the actor's
# client certificate + CONNECT authority were seen.
set -o errexit -o nounset -o pipefail
ROOT="$(git rev-parse --show-toplevel)"; cd "${ROOT}"
CTX="${KUBECTL_CONTEXT:-kind-kind}"
K="kubectl --context ${CTX}"
ATESPACE="${ATESPACE:-demo}"
ACTOR="${ACTOR:-egress-demo}"
TARGET_URL="${TARGET_URL:-http://example.com/}"
echo "== gateway should be running =="
${K} -n ate-system rollout status deployment/atenet-egress --timeout=120s
echo "== create atespace + actor =="
kubectl-ate --context "${CTX}" create atespace "${ATESPACE}" 2>/dev/null || true
kubectl-ate --context "${CTX}" create actor "${ACTOR}" \
--atespace "${ATESPACE}" --template ate-demo-egress/egress 2>/dev/null || true
${K} -n ate-system wait --for=condition=Ready "actor/${ACTOR}" 2>/dev/null || sleep 10
echo "== snapshot gateway log offset =="
BEFORE=$(${K} -n ate-system logs deployment/atenet-egress --tail=-1 2>/dev/null | wc -l | tr -d ' ')
echo "== drive actor egress: GET ${TARGET_URL} via the actor =="
${K} -n ate-system port-forward service/atenet-router 18000:80 >/tmp/pf.log 2>&1 &
PF=$!; trap 'kill ${PF} 2>/dev/null || true' EXIT
sleep 3
RESP=$(curl -s -o /dev/null -w "%{http_code}" -X POST http://localhost:18000/ \
-H "Host: ${ACTOR}.${ATESPACE}.actors.resources.substrate.ate.dev" \
-H 'Content-Type: application/json' \
-d "{\"url\":\"${TARGET_URL}\"}") || true
echo "actor round-trip HTTP ${RESP} (200 = the actor fetched ${TARGET_URL} through egress)"
echo "== NEW egress gateway access log lines (proof of CONNECT+mTLS+identity) =="
${K} -n ate-system logs deployment/atenet-egress --tail=-1 2>/dev/null \
| tail -n +"$((BEFORE + 1))" | grep '\[egress\]' || {
echo "!! no [egress] lines — dumping recent gateway logs:"
${K} -n ate-system logs deployment/atenet-egress --tail=20
exit 1
}
echo "== PASS: actor egress traversed the Envoy egress gateway =="
+18 -3
View File
@@ -15,8 +15,10 @@
package e2e
import (
"bytes"
"context"
"fmt"
"io"
"net/http"
"time"
@@ -31,7 +33,7 @@ const (
routerService = "atenet-router"
)
// RouterClient sends HTTP requests to actors through the atenet router, the
// RouterClient sends HTTP requests to actors through the ingress atenet-router, the
// same way real traffic arrives (so the request is routed and, if needed, the
// actor is resumed). It port-forwards the router Service, mirroring the
// approach in internal/ateclient.
@@ -41,7 +43,7 @@ type RouterClient struct {
stop func()
}
// NewRouterClient establishes a port-forward to the atenet router. Call Close
// NewRouterClient establishes a port-forward to the ingress atenet-router. Call Close
// to tear it down.
func NewRouterClient(ctx context.Context) (*RouterClient, error) {
config, err := ateclient.LoadConfig(KubeConfig, KubeContext)
@@ -73,10 +75,23 @@ func (c *RouterClient) Close() {
// Get issues GET path to actor through the router, setting the actor's mesh Host
// so the router routes (and resumes) it. The caller must close the body.
func (c *RouterClient) Get(ctx context.Context, actorRef resources.ActorRef, path string) (*http.Response, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+path, nil)
return c.request(ctx, http.MethodGet, actorRef, path, nil)
}
// PostJSON issues a POST with a JSON body to an Actor through the router. The
// caller must close the response body.
func (c *RouterClient) PostJSON(ctx context.Context, actorRef resources.ActorRef, path string, body []byte) (*http.Response, error) {
return c.request(ctx, http.MethodPost, actorRef, path, bytes.NewReader(body))
}
func (c *RouterClient) request(ctx context.Context, method string, actorRef resources.ActorRef, path string, body io.Reader) (*http.Response, error) {
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, body)
if err != nil {
return nil, err
}
if method == http.MethodPost {
req.Header.Set("Content-Type", "application/json")
}
// The router routes on the Host/:authority, not a header.
req.Host = actorRef.DNSName()
return c.http.Do(req)
+78
View File
@@ -0,0 +1,78 @@
// 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 e2e
import (
"context"
"io"
"net/http"
"strings"
"testing"
"github.com/agent-substrate/substrate/internal/resources"
)
// SEE(lior): this file used to also hold TestResolveHTTPTargetPort and
// TestIsPodReady. Main moved both into internal/portforward/portforward_test.go
// (upstream 4453b5e7, "Prefactoring: Consolidate logic to port-forward to a
// service Pod") and deleted this file, so git rename-matched this branch's
// version of it against the new portforward test and reported the whole thing as
// a conflict. Took main's move as-is and re-added only the net-new PostJSON test
// here, rather than resurrecting the port-forward tests in a package that no
// longer owns that code.
func TestRouterClientPostJSON(t *testing.T) {
client := &RouterClient{
baseURL: "http://router.test",
http: &http.Client{Transport: testRoundTripper(func(request *http.Request) (*http.Response, error) {
if request.Method != http.MethodPost {
t.Errorf("method = %q, want POST", request.Method)
}
if request.Host != "fetcher.demo.actors.resources.substrate.ate.dev" {
t.Errorf("host = %q", request.Host)
}
if request.URL.Path != "/fetch" {
t.Errorf("path = %q, want /fetch", request.URL.Path)
}
if request.Header.Get("Content-Type") != "application/json" {
t.Errorf("content type = %q, want application/json", request.Header.Get("Content-Type"))
}
body, err := io.ReadAll(request.Body)
if err != nil {
t.Fatalf("reading body: %v", err)
}
if string(body) != `{"url":"https://example.com/"}` {
t.Errorf("body = %q", body)
}
return &http.Response{
StatusCode: http.StatusOK,
Body: io.NopCloser(strings.NewReader("ok")),
Header: make(http.Header),
}, nil
})},
}
actorRef := resources.ActorRef{Atespace: "demo", Name: "fetcher"}
response, err := client.PostJSON(context.Background(), actorRef, "/fetch", []byte(`{"url":"https://example.com/"}`))
if err != nil {
t.Fatalf("PostJSON: %v", err)
}
response.Body.Close()
}
type testRoundTripper func(*http.Request) (*http.Response, error)
func (f testRoundTripper) RoundTrip(request *http.Request) (*http.Response, error) {
return f(request)
}
@@ -29,9 +29,24 @@ import (
const networkingAtespace = "networking-e2e"
// actorTemplate identifies a demo ActorTemplate to build test Actors from,
// along with the hack/install-ate.sh flag that deploys it.
type actorTemplate struct {
namespace string
name string
// deployFlag names the install flag that creates the template, so a
// missing fixture reports how to fix it rather than just failing.
deployFlag string
}
var (
counterTemplate = actorTemplate{namespace: "ate-demo-counter", name: "counter", deployFlag: "--deploy-demo-counter"}
egressTemplate = actorTemplate{namespace: "ate-demo-egress", name: "egress", deployFlag: "--deploy-demo-egress"}
)
func TestActorDirectAccess(t *testing.T) {
ctx := context.Background()
actorName, actor := createAndResumeActor(t, ctx, "direct")
actorName, actor := createAndResumeActor(t, ctx, "direct", counterTemplate)
router := mustRouterClient(t, ctx)
defer router.Close()
@@ -68,7 +83,38 @@ func TestActorDirectAccess(t *testing.T) {
})
}
func createAndResumeActor(t *testing.T, ctx context.Context, prefix string) (string, *ateapipb.Actor) {
// TestActorEgress exercises the full egress path. The Actor's outbound TCP
// connection is transparently redirected by nftables into atunnel, wrapped in
// mTLS with the Actor's own actor-identity certificate plus an HTTP CONNECT to
// atenet-egress, authorized there against that certificate, and only then
// dialed out. A masqueraded (pre-gateway) egress would also return 200, so this
// asserts the gateway is deployed and that it did not reject the Actor.
func TestActorEgress(t *testing.T) {
ctx := context.Background()
actorName, _ := createAndResumeActor(t, ctx, "egress", egressTemplate)
router := mustRouterClient(t, ctx)
defer router.Close()
// The egress demo fetches the URL it is given and echoes the upstream
// status and body back.
payload := []byte(`{"url":"http://example.com/"}`)
actorRef := resources.ActorRef{Atespace: networkingAtespace, Name: actorName}
response, err := router.PostJSON(ctx, actorRef, "/", payload)
if err != nil {
t.Fatalf("POST to egress Actor through ingress: %v", err)
}
defer response.Body.Close()
body, err := io.ReadAll(response.Body)
if err != nil {
t.Fatalf("reading egress response body (HTTP %d): %v", response.StatusCode, err)
}
if response.StatusCode != http.StatusOK {
t.Fatalf("Actor egress fetch returned HTTP %d, want 200; body: %s", response.StatusCode, body)
}
t.Logf("Actor egress fetch succeeded; body: %s", body)
}
func createAndResumeActor(t *testing.T, ctx context.Context, prefix string, template actorTemplate) (string, *ateapipb.Actor) {
t.Helper()
clients := e2e.GetClients()
actorName := fmt.Sprintf("%s-%d", prefix, time.Now().UnixNano())
@@ -80,10 +126,10 @@ func createAndResumeActor(t *testing.T, ctx context.Context, prefix string) (str
})
if _, err := clients.SubstrateAPI.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: networkingAtespace, Name: actorName},
ActorTemplateNamespace: "ate-demo-counter",
ActorTemplateName: "counter",
ActorTemplateNamespace: template.namespace,
ActorTemplateName: template.name,
}}); err != nil {
t.Fatalf("CreateActor: %v (deploy the fixture with --deploy-demo-counter)", err)
t.Fatalf("CreateActor from %s/%s: %v (deploy the fixture with %s)", template.namespace, template.name, err, template.deployFlag)
}
t.Cleanup(func() {
_, _ = clients.SubstrateAPI.SuspendActor(context.Background(), &ateapipb.SuspendActorRequest{Actor: actorRef})
+1 -1
View File
@@ -112,7 +112,7 @@ spec:
#
# SEE(lior): this flag lived on atelet on the pre-rebase branch; PR 708
# moved egress-gateway configuration to ateapi, so it is set here now.
- --egress-gateway-address=ateway-egress.ate-system.svc:443
- --egress-gateway-address=atenet-egress.ate-system.svc:443
# Graceful shutdown knobs. The sum must fit within terminationGracePeriodSeconds.
- --drain-delay=13s
- --drain-timeout=15s
@@ -27,7 +27,7 @@
apiVersion: v1
kind: ServiceAccount
metadata:
name: ateway-egress
name: atenet-egress
namespace: ate-system
---
# RBAC for the co-located ext_proc sidecar (atenet router), which watches
@@ -35,7 +35,7 @@ metadata:
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
name: ateway-egress
name: atenet-egress
rules:
- apiGroups:
- "ate.dev"
@@ -49,20 +49,20 @@ rules:
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
name: ateway-egress
name: atenet-egress
subjects:
- kind: ServiceAccount
name: ateway-egress
name: atenet-egress
namespace: ate-system
roleRef:
kind: ClusterRole
name: ateway-egress
name: atenet-egress
apiGroup: rbac.authorization.k8s.io
---
apiVersion: v1
kind: ConfigMap
metadata:
name: ateway-egress
name: atenet-egress
namespace: ate-system
data:
envoy.yaml: |
@@ -75,7 +75,11 @@ data:
address:
socket_address: { address: 0.0.0.0, port_value: 443 }
filter_chains:
- transport_socket:
# Named so ext_proc can read it back as xds.filter_chain_name. Must
# match EgressFilterChainName in
# cmd/atenet/internal/router/extproc_egress.go.
- name: egress
transport_socket:
name: envoy.transport_sockets.tls
typed_config:
"@type": type.googleapis.com/envoy.extensions.transport_sockets.tls.v3.DownstreamTlsContext
@@ -130,10 +134,13 @@ data:
# certificate, not from request headers. The actor sends no
# identity headers any more, and logging attacker-controlled
# ones would have made this log lie about who egressed.
# peer= is the leaf subject (CN carries the actor); the
# atespace/UID live in the ActorIdentity extension, which
# Envoy cannot format — the ext_proc handler logs those.
inline_string: "[egress] authority=%REQ(:AUTHORITY)% peer=%DOWNSTREAM_PEER_SUBJECT% peer_san=%DOWNSTREAM_PEER_URI_SAN% peer_serial=%DOWNSTREAM_PEER_SERIAL% code=%RESPONSE_CODE% flags=%RESPONSE_FLAGS% up_bytes=%BYTES_RECEIVED% down_bytes=%BYTES_SENT%\n"
# peer_san= is the actor's SPIFFE URI SAN,
# spiffe://substrate-actor.local/atespace/<atespace>/actor/<name>
# — the only actor-identifying field Envoy can format
# (ateapi mints these with an empty Subject, and the UID
# lives in the ActorIdentity extension Envoy cannot read).
# The ext_proc sidecar logs atespace/actor/UID separately.
inline_string: "[egress] authority=%REQ(:AUTHORITY)% peer_san=%DOWNSTREAM_PEER_URI_SAN% peer_serial=%DOWNSTREAM_PEER_SERIAL% code=%RESPONSE_CODE% flags=%RESPONSE_FLAGS% up_bytes=%BYTES_RECEIVED% down_bytes=%BYTES_SENT%\n"
route_config:
name: connect_route
virtual_hosts:
@@ -163,13 +170,14 @@ data:
failure_mode_allow: false
# How the ext_proc server tells egress from ingress. It applies
# opposite trust models to the two directions, and dispatches on
# this Envoy-asserted listener name so that no client can select
# the egress path by crafting a request. The value must match
# EgressListenerName in cmd/atenet/internal/router/extproc_egress.go;
# renaming the listener above without updating it fails closed
# (egress requests take the ingress path and 404).
# this Envoy-asserted filter chain name so that no client can
# select the egress path by crafting a request. The value must
# match EgressFilterChainName in
# cmd/atenet/internal/router/extproc_egress.go; renaming the
# filter chain above without updating it fails closed (egress
# requests take the ingress path and 404).
request_attributes:
- xds.listener_name
- xds.filter_chain_name
processing_mode:
request_header_mode: SEND
response_header_mode: SKIP
@@ -224,21 +232,21 @@ data:
apiVersion: apps/v1
kind: Deployment
metadata:
name: ateway-egress
name: atenet-egress
namespace: ate-system
labels:
app: ateway-egress
app: atenet-egress
spec:
replicas: 1
selector:
matchLabels:
app: ateway-egress
app: atenet-egress
template:
metadata:
labels:
app: ateway-egress
app: atenet-egress
spec:
serviceAccountName: ateway-egress
serviceAccountName: atenet-egress
securityContext:
# Allow the non-root envoy user to bind :443.
sysctls:
@@ -258,9 +266,9 @@ spec:
- -c
- /etc/envoy/envoy.yaml
- --service-node
- ateway-egress
- atenet-egress
- --service-cluster
- ateway-egress
- atenet-egress
ports:
- name: https
containerPort: 443
@@ -302,8 +310,9 @@ spec:
- --namespace=ate-system
- --port-extproc=50051
- --extproc-address=127.0.0.1
- --ateapi-address=api.ate-system.svc:443
- --ateapi-auth=mtls
- --ateapi-address=dns:///api.ate-system.svc:443
- --ateapi-ca-file=/run/servicedns.podcert.ate.dev/trust-bundle.pem
- --ateapi-client-cert=/run/podidentity.podcert.ate.dev/credential-bundle.pem
# Same bundle Envoy validates the handshake against. The handler
# re-verifies the chain in Go and reads the ActorIdentity extension out
# of it; without this flag the router has no actor-identity roots and
@@ -327,13 +336,22 @@ spec:
port: extproc
periodSeconds: 10
volumeMounts:
# Trust bundle used to verify ateapi's servicedns serving cert.
- name: servicedns
mountPath: /run/servicedns.podcert.ate.dev
readOnly: true
# ext-proc's own client identity presented to ateapi.
- name: podidentity
mountPath: /run/podidentity.podcert.ate.dev
readOnly: true
# Actor-identity CA roots, for --actor-identity-ca-file above.
- name: actor-id-ca-certs
mountPath: /run/actor-id-ca-certs
readOnly: true
volumes:
- name: config
configMap:
name: ateway-egress
name: atenet-egress
- name: servicedns
projected:
sources:
@@ -380,12 +398,12 @@ spec:
apiVersion: v1
kind: Service
metadata:
name: ateway-egress
name: atenet-egress
namespace: ate-system
spec:
type: ClusterIP
selector:
app: ateway-egress
app: atenet-egress
ports:
- name: https
port: 443
@@ -23,7 +23,7 @@ resources:
- ../ate-controller.yaml
- ../atelet.yaml
- ../atenet-dns.yaml
- ../ateway-egress.yaml
- ../atenet-egress.yaml
- ../atenet-router.yaml
- ../valkey.yaml
- ../pod-certificate-controller.yaml
@@ -26,7 +26,7 @@ resources:
- ../ate-controller.yaml
- ./atelet
- ../atenet-dns.yaml
- ../ateway-egress.yaml
- ../atenet-egress.yaml
- ../atenet-router.yaml
- ../valkey.yaml
- ../pod-certificate-controller.yaml