e2e: add test for egress credential injection cujs (#1795)

## e2e: add e2e for egress credential injection

Adds the `egresscredinject` e2e suite: an actor fetches
`https://httpbin.org/headers`
through the sdsmint MITM gateway, and the echoed response proves the
injected
`Authorization` header actually reached the upstream.

### What it asserts

- **Injected**: echoed headers contain `Authorization: Bearer <token>`.
- **Overwritten**: an actor-pre-seeded `Authorization` is replaced by
the injected one.
- **Cleartext skip**: plain-HTTP fetch passes through with no
`Authorization`.
- **Fail closed**: nonexistent secret → 403, unserved provider → 500,
  unauthorized namespace → 403.

### Additional changes

- Probe `/fetch` returns the response body, accepts
`header=<name>:<value>`
params, and no longer follows redirects (a cross-scheme redirect would
hop
between the cleartext and TLS legs). Covered by unit tests; egressmitm
rerun green.
- New `e2e.EgressInjectHeader` policy helper.
- New `e2e.DeployCredentialProvider` + `fixtures/credinject`: deploys
the real
provider manifest with a test Secret and namespace-policy ConfigMap,
restarts
  the provider so it picks the policy up, cleans up on test end.
- CI: new envoy-lane steps after the MITM lanes, gated on
`E2E_EGRESS_CREDINJECT=1`.

- [x] Tests pass
- [ ] Appropriate changes to documentation are included in the PR
This commit is contained in:
Yufan Su
2026-09-29 03:47:06 +00:00
committed by GitHub
parent f7ec6acfc0
commit 4f9c293e48
8 changed files with 797 additions and 9 deletions
+14
View File
@@ -239,6 +239,20 @@ jobs:
E2E_SANDBOX_CLASS: microvm
E2E_JUNIT_FILE: ${{ github.workspace }}/_artifacts/e2e-networking-mitm-microvm.xml
run: hack/run-e2e-kind.sh ./internal/e2e/suites/networking -run '^TestActorEgress' -v -args --no-color
- name: Deploy MITM egress with credential injection
# Envoy lane only: the injection flag requires the envoy dataplane.
# TODO: support agent gateway as dataplane.
if: matrix.dataplane == 'envoy'
run: hack/install-ate-kind.sh --deploy-atenet --atenet-dataplane=envoy --experimental-use-sdsmint --experimental-egress-credential-injection
- name: Run E2E tests (egress credential injection)
# Injection happens entirely inside the gateway, so the sandbox class
# does not change the path — per-class trust delivery is already the
# egressmitm steps' job. One (gVisor) run suffices.
if: matrix.dataplane == 'envoy'
env:
E2E_EGRESS_CREDINJECT: "1"
E2E_JUNIT_FILE: ${{ github.workspace }}/_artifacts/e2e-credinject.xml
run: hack/run-e2e-kind.sh ./internal/e2e/suites/egresscredinject -v -args --no-color
# Every lane must have reported tests; a `-run` filter that stops matching
# otherwise leaves its lane green and empty. The set comes from the manifest
# each lane registers itself in, so adding a lane needs no change here, and
+102
View File
@@ -0,0 +1,102 @@
// 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 (
"path/filepath"
"testing"
"github.com/agent-substrate/substrate/internal/installdefaults"
)
// TODO(yufan-su): Move these helpers into an internal/e2e/credprovider package.
const (
// CredentialSecretsNamespace holds the fixture Secret and is the only
// namespace the credinject fixture's authorization policy allows the probe
// atespaces to resolve.
CredentialSecretsNamespace = "ate-e2e-credinject-secrets"
// CredentialInjectionURI resolves to CredentialInjectionToken through the
// k8s-credential-provider, for the atespaces the fixture policy allows.
CredentialInjectionURI = "ate-secret://k8s.io/default/" + CredentialSecretsNamespace + "/api-token/token"
// CredentialInjectionToken is the fixture Secret's value, what an injected
// header carries after any prefix.
CredentialInjectionToken = "e2e-cred-inject-token"
)
// The fixture manifest (Secret plus provider authorization policy) and the
// provider's real deployment manifest — reused rather than copied, so the
// suite cannot drift from what an install deploys.
const (
credinjectFixtureManifest = "internal/e2e/fixtures/credinject/credinject.yaml"
credentialProviderManifest = "manifests/egress-credential-injection/k8s-credential-provider.yaml"
)
// DeployCredentialProvider installs the k8s-credential-provider and the
// credinject fixture (the Secret behind CredentialInjectionURI and the
// authorization policy that lets the probe atespaces resolve it), and removes
// both when the test passes. The provider Deployment is restarted after the
// policy ConfigMap is applied because it reads the policy once at startup, so
// a provider left running by an earlier install would otherwise keep
// enforcing a stale one.
//
// A failed test keeps both so the provider's logs can be inspected; the next
// run re-applies them.
//
// The egress gateway's side of the connection — the --credential-provider-*
// flags on its ext_proc sidecar — is install-time configuration
// (hack/install-ate.sh --experimental-egress-credential-injection), not
// something this helper can retrofit.
func DeployCredentialProvider(t *testing.T) {
t.Helper()
root, err := FindRepoRoot()
if err != nil {
t.Fatalf("FindRepoRoot: %v", err)
}
kubectl := func(args ...string) {
if KubeContext != "" {
args = append([]string{"--context=" + KubeContext}, args...)
}
RunCmd(t, "kubectl", args...)
}
// The policy ConfigMap must exist before the provider pod starts: the
// Deployment mounts it, and the provider loads it at startup.
fixture := filepath.Join(root, credinjectFixtureManifest)
kubectl("apply", "-f", fixture)
deleteOnPass := func(manifest string) {
t.Cleanup(func() {
if t.Failed() {
return
}
kubectl("delete", "--ignore-not-found", "-f", manifest)
})
}
deleteOnPass(fixture)
provider := filepath.Join(root, credentialProviderManifest)
koApply(t, provider)
deleteOnPass(provider)
// Restart unconditionally: if the Deployment already existed, koApply may
// have changed nothing, leaving a pod that started under a previous
// policy ConfigMap. The provider manifest pins the canonical namespace.
ns := installdefaults.SystemNamespace
kubectl("-n", ns, "rollout", "restart", "deployment/k8s-credential-provider")
kubectl("-n", ns, "rollout", "status", "deployment/k8s-credential-provider", "--timeout=3m")
}
+31
View File
@@ -24,6 +24,8 @@ import (
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
)
// TODO(yufan-su): Move these constructors into an internal/e2e/egresspolicy package.
// EgressAllowAll is what a test that is not about egress policy gives its
// actor, since the gateway denies an actor with no policy at all: every name
// and address, as cleartext HTTP on any port and as intercepted HTTPS on 443.
@@ -49,6 +51,35 @@ func EgressAllowHTTPS(patterns ...string) *ateapipb.EgressRule {
return &ateapipb.EgressRule{Https: &ateapipb.HTTPSRule{Hostnames: patterns}}
}
// EgressInjectHeader is an https rule (see EgressAllowHTTPS) that also carries
// a replace_headers effect: on a match, the gateway resolves credentialURI
// through its credential provider and replaces header with prefix plus the
// credential.
func EgressInjectHeader(header, prefix, credentialURI string, patterns ...string) *ateapipb.EgressRule {
rule := EgressAllowHTTPS(patterns...)
rule.Https.Effects = replaceHeaderEffects(header, prefix, credentialURI)
return rule
}
// EgressInjectHeaderHTTP is an http rule (see EgressAllowHTTP) carrying the
// same effect as EgressInjectHeader. The gateway never puts a credential on
// cleartext, so a test uses it to prove the effect is skipped there.
func EgressInjectHeaderHTTP(header, prefix, credentialURI string, patterns ...string) *ateapipb.EgressRule {
rule := EgressAllowHTTP(patterns...)
rule.Http.Effects = replaceHeaderEffects(header, prefix, credentialURI)
return rule
}
func replaceHeaderEffects(header, prefix, credentialURI string) *ateapipb.HttpRuleEffects {
return &ateapipb.HttpRuleEffects{
ReplaceHeaders: []*ateapipb.CredentialHeader{{
Header: header,
Prefix: prefix,
CredentialUri: credentialURI,
}},
}
}
// EnsureEgressPolicy gives actor an EgressPolicy with exactly rules, replacing
// any it had. Deleting the actor deletes the policy, so there is no cleanup.
func EnsureEgressPolicy(t *testing.T, ctx context.Context, clients *Clients, actor *ateapipb.ObjectRef, rules ...*ateapipb.EgressRule) {
@@ -0,0 +1,65 @@
# 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.
# The egress credential-injection fixture: the Secret the suite's policy
# injects, and the k8s-credential-provider authorization policy that lets the
# probe fixture's atespaces resolve it. Applied by e2e.DeployCredentialProvider
# before the provider Deployment itself, which reads the ConfigMap once at
# startup.
apiVersion: v1
kind: Namespace
metadata:
name: ate-e2e-credinject-secrets
labels:
# Lets hack/cleanup-e2e.sh find and delete this namespace when a failed
# test leaves the fixture behind for debugging.
ate.dev/e2e: "true"
---
# The credential the suite's EgressPolicy injects, resolved by the URI
# ate-secret://k8s.io/default/ate-e2e-credinject-secrets/api-token/token
# (see e2e.CredentialInjectionURI).
apiVersion: v1
kind: Secret
metadata:
name: api-token
namespace: ate-e2e-credinject-secrets
type: Opaque
stringData:
token: "e2e-cred-inject-token"
---
# The atespace→namespace authorization policy the provider enforces
# (default-deny). The name and namespace are dictated by the configMap volume
# in manifests/egress-credential-injection/k8s-credential-provider.yaml. Both
# probe-fixture atespaces are listed — the suffixes are deterministic per
# sandbox class (see e2e.DeployProbe and probe-template.yaml.tmpl) — so one
# policy serves the gVisor and micro-VM variants.
apiVersion: v1
kind: ConfigMap
metadata:
name: k8s-credential-provider-namespace-policy
namespace: ate-system
data:
namespace-policy.yaml: |
policies:
- atespace: ate-e2e-probe-egresscredinject
allowedNamespaces:
- ate-e2e-credinject-secrets
- atespace: ate-e2e-probe-microvm-egresscredinject
allowedNamespaces:
- ate-e2e-credinject-secrets
+69 -9
View File
@@ -293,13 +293,49 @@ func memTotalBytes() (int64, error) {
return 0, os.ErrNotExist
}
// fetch GETs ?url= over the actor's normal egress path and reports the
// outcome, doing TLS with the trust anchors selected by ?roots=: "bundle"
// (the default) loads the projected trust bundle at trustFile, "system" uses
// the image's system roots. TestActorEgressMITMTrust documents why each mode
// passes or fails. TLS failures land in the "error" field rather than the
// HTTP status: a verification failure is a result for the suite to assert
// on, not a broken probe.
// maxFetchBody caps the response body fetch echoes back, so a large origin
// response cannot balloon the probe's reply. 64 KiB comfortably holds the
// header-echo documents the suites assert on.
const maxFetchBody = 64 << 10
// parseFetchHeaders parses repeated ?header=<name>:<value> parameters into
// request headers. The value is taken verbatim after the first colon.
func parseFetchHeaders(params []string) (http.Header, error) {
headers := http.Header{}
for _, p := range params {
name, value, ok := strings.Cut(p, ":")
if !ok || name == "" {
return nil, fmt.Errorf("header parameter %q is not <name>:<value>", p)
}
headers.Add(name, value)
}
return headers, nil
}
// fetch causes probe to issue an HTTP(S) GET to exercise the actor's egress
// path. Parameters to the fetch are passed as URL query parameters:
//
// - url=<URL to fetch>: required.
// - roots=bundle|system: the TLS trust anchors. "bundle" (the default)
// loads the projected trust bundle at trustFile; "system" uses the
// image's system roots. TestActorEgressMITMTrust documents why each mode
// passes or fails.
// - header=<name>:<value>: repeatable; set on the request, so a suite can
// pre-seed a header and observe whether the gateway overwrites it.
//
// The reply is a JSON object with the origin's "status" and the first 64 KiB
// of its "body", so a suite can assert on what the origin received (e.g. an
// injected credential echoed back by a headers-echo endpoint). Failures land
// in "error" rather than the HTTP status: a TLS verification failure is a
// result for the suite to assert on, not a broken probe.
//
// Redirects are not followed — a cross-scheme redirect would silently hop
// between the gateway's cleartext and TLS legs, flipping the very behavior
// (credential injection) some suites assert on — so the first response is
// the result.
//
// TODO: Accept the parameters as a JSON request body as well, which avoids
// the escaping that query-string values need.
func fetch(w http.ResponseWriter, r *http.Request) {
resp := map[string]string{}
url := r.URL.Query().Get("url")
@@ -319,6 +355,12 @@ func fetch(w http.ResponseWriter, r *http.Request) {
writeJSON(w, resp)
return
}
headers, err := parseFetchHeaders(r.URL.Query()["header"])
if err != nil {
resp["error"] = err.Error()
writeJSON(w, resp)
return
}
tlsCfg := &tls.Config{}
if roots != "system" {
b, err := os.ReadFile(trustFile)
@@ -338,16 +380,34 @@ func fetch(w http.ResponseWriter, r *http.Request) {
client := &http.Client{
Timeout: 20 * time.Second,
Transport: &http.Transport{TLSClientConfig: tlsCfg},
CheckRedirect: func(*http.Request, []*http.Request) error {
return http.ErrUseLastResponse
},
}
res, err := client.Get(url)
// The transport is per request, so don't leave its connection idle.
defer client.CloseIdleConnections()
req, err := http.NewRequestWithContext(r.Context(), http.MethodGet, url, nil)
if err != nil {
resp["error"] = err.Error()
writeJSON(w, resp)
return
}
for name, values := range headers {
req.Header[name] = values
}
res, err := client.Do(req)
if err != nil {
resp["error"] = err.Error()
writeJSON(w, resp)
return
}
defer res.Body.Close()
_, _ = io.Copy(io.Discard, res.Body)
resp["status"] = strconv.Itoa(res.StatusCode)
body, err := io.ReadAll(io.LimitReader(res.Body, maxFetchBody))
if err != nil {
resp["error"] = "reading response body: " + err.Error()
}
resp["body"] = string(body)
writeJSON(w, resp)
}
+148
View File
@@ -15,8 +15,15 @@
package main
import (
"encoding/json"
"net/http"
"net/http/httptest"
"net/url"
"slices"
"strings"
"testing"
"github.com/google/go-cmp/cmp"
)
// The capabilities e2e assertions are only as trustworthy as this decoder and
@@ -76,6 +83,147 @@ func TestDecodeCapMask(t *testing.T) {
}
}
func TestParseFetchHeaders(t *testing.T) {
tests := []struct {
name string
params []string
want http.Header
wantErr bool
}{{
name: "no parameters",
params: nil,
want: http.Header{},
}, {
name: "name and value",
params: []string{"Authorization:Bearer x"},
want: http.Header{"Authorization": {"Bearer x"}},
}, {
// Only the first colon separates; the value keeps the rest verbatim.
name: "value containing a colon",
params: []string{"X-Test:a:b"},
want: http.Header{"X-Test": {"a:b"}},
}, {
name: "repeated name keeps every value",
params: []string{"x-test:1", "X-Test:2"},
want: http.Header{"X-Test": {"1", "2"}},
}, {
name: "no colon",
params: []string{"no-colon"},
wantErr: true,
}, {
name: "empty name",
params: []string{":empty-name"},
wantErr: true,
}}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := parseFetchHeaders(tt.params)
if tt.wantErr {
if err == nil {
t.Fatalf("parseFetchHeaders(%q) = %v, want an error", tt.params, got)
}
return
}
if err != nil {
t.Fatalf("parseFetchHeaders(%q) failed: %v", tt.params, err)
}
if diff := cmp.Diff(tt.want, got); diff != "" {
t.Errorf("parseFetchHeaders(%q) mismatch (-want +got):\n%s", tt.params, diff)
}
})
}
}
// doFetch drives the fetch handler at an origin URL and decodes its JSON
// reply. roots=system keeps the handler off the projected trust bundle,
// which does not exist outside a cluster.
func doFetch(t *testing.T, origin string, headerParams ...string) map[string]string {
t.Helper()
query := url.Values{"url": {origin}, "roots": {"system"}}
for _, h := range headerParams {
query.Add("header", h)
}
req := httptest.NewRequest(http.MethodGet, "/fetch?"+query.Encode(), nil)
rec := httptest.NewRecorder()
fetch(rec, req)
resp := map[string]string{}
if err := json.Unmarshal(rec.Body.Bytes(), &resp); err != nil {
t.Fatalf("decoding fetch response %q: %v", rec.Body.String(), err)
}
return resp
}
// The suites' credential-injection assertions live in the response body an
// origin echoes back, so fetch must return it — along with the request
// headers set from ?header= parameters, which is how a suite pre-seeds a
// header the gateway should overwrite.
func TestFetchReturnsBodyAndSetsHeaders(t *testing.T) {
origin := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte("auth=" + r.Header.Get("Authorization")))
}))
defer origin.Close()
resp := doFetch(t, origin.URL, "Authorization:Bearer seeded")
if resp["error"] != "" {
t.Fatalf("fetch failed: %s", resp["error"])
}
if resp["status"] != "200" {
t.Errorf("status = %q, want 200", resp["status"])
}
if resp["body"] != "auth=Bearer seeded" {
t.Errorf("body = %q, want %q", resp["body"], "auth=Bearer seeded")
}
}
// A followed cross-scheme redirect would silently hop between the gateway's
// cleartext and TLS legs, so fetch must report the first response instead.
func TestFetchDoesNotFollowRedirects(t *testing.T) {
origin := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/target" {
t.Error("fetch followed the redirect")
}
http.Redirect(w, r, "/target", http.StatusFound)
}))
defer origin.Close()
resp := doFetch(t, origin.URL)
if resp["error"] != "" {
t.Fatalf("fetch failed: %s", resp["error"])
}
if resp["status"] != "302" {
t.Errorf("status = %q, want 302", resp["status"])
}
}
// A malformed ?header= fails the fetch before anything is sent.
func TestFetchRejectsMalformedHeader(t *testing.T) {
origin := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {
t.Error("fetch sent the request despite a malformed header parameter")
}))
defer origin.Close()
resp := doFetch(t, origin.URL, "no-colon")
if !strings.Contains(resp["error"], "not <name>:<value>") {
t.Errorf("error = %q, want the malformed-header error", resp["error"])
}
if resp["status"] != "" {
t.Errorf("status = %q, want none", resp["status"])
}
}
func TestFetchTruncatesBody(t *testing.T) {
origin := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = w.Write([]byte(strings.Repeat("x", maxFetchBody+1)))
}))
defer origin.Close()
resp := doFetch(t, origin.URL)
if got := len(resp["body"]); got != maxFetchBody {
t.Errorf("len(body) = %d, want %d", got, maxFetchBody)
}
}
func TestDecodeCapMaskInvalid(t *testing.T) {
if _, err := decodeCapMask("nothex"); err == nil {
t.Error("decodeCapMask(\"nothex\") succeeded, want an error")
@@ -0,0 +1,344 @@
// 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 egresscredinject e2e-tests egress credential injection: a matching
// EgressPolicy https rule with a replace_headers effect makes the sdsmint
// egress gateway's MITM leg resolve the credential through the
// k8s-credential-provider and replace the actor's placeholder header with it
// before re-originating upstream. See TestActorEgressCredentialInjection for
// the proof structure and how to run this locally.
package egresscredinject
import (
"context"
"encoding/json"
"io"
"net/http"
"net/url"
"os"
"strings"
"testing"
"time"
"github.com/agent-substrate/substrate/internal/e2e"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)
const probeTemplate = "probe"
var probeNamespace string
// The suite's hostnames, one per injection outcome: each is covered by exactly
// one https rule, so each selects exactly one CredentialHeader. echoHost is the
// only one whose response matters — it echoes the request headers it
// received back as JSON, which is what proves the header was on the wire.
const (
echoHost = "httpbin.org"
unfetchableHost = "example.com"
unservedHost = "example.org"
unauthorizedHost = "example.net"
echoOrigin = "https://" + echoHost + "/headers"
echoOriginPlain = "http://" + echoHost + "/headers"
unfetchableOrigin = "https://" + unfetchableHost + "/"
unservedOrigin = "https://" + unservedHost + "/"
unauthorizedOrigin = "https://" + unauthorizedHost + "/"
)
// placeholder is the Authorization value the actor sends. replace_headers
// replaces a header only when the request carries it, so every fetch that
// should trigger injection sends it.
const placeholder = "Bearer actor-placeholder"
var withPlaceholder = []string{"header=" + url.QueryEscape("Authorization:"+placeholder)}
// TestActorEgressCredentialInjection proves the injected credential reaches
// the upstream, and that every way injection can go wrong lands on the
// documented side of fail-open vs fail-closed (see egress.Handler.applyEffects):
//
// - replaced: a fetch of the echo origin carrying the placeholder echoes the
// injected "Authorization: Bearer <token>" among the headers the origin
// received — the on-the-wire proof, not an inference from a status code —
// and not the placeholder, so an actor cannot choose the value that
// leaves.
// - cleartext skip: the same fetch over plain HTTP, allowed by an http rule
// with the same effect, echoes the placeholder and not the credential —
// the secret never rides a cleartext wire, and the request is passed
// through rather than denied.
// - fail closed: a credential the policy requires but the provider will
// not or cannot produce denies the request — 403 for an unfetchable
// secret and for a namespace outside the atespace's authorization
// (default-deny), 500 for a URI naming a provider this gateway does not
// serve.
//
// A request without the header is not covered: the API forwards it without
// the credential, which the gateway does not implement yet.
//
// The gate: this needs the sdsmint egress gateway with injection enabled
// (which replaces the passthrough gateway cluster-wide) plus the
// k8s-credential-provider, which the suite deploys itself. Locally:
//
// hack/install-ate-kind.sh --deploy-atenet --experimental-use-sdsmint --experimental-egress-credential-injection
// E2E_EGRESS_CREDINJECT=1 hack/run-e2e-kind.sh ./internal/e2e/suites/egresscredinject -v -args --no-color
func TestActorEgressCredentialInjection(t *testing.T) {
if os.Getenv("E2E_EGRESS_CREDINJECT") == "" {
t.Skip("needs the sdsmint (MITM) egress gateway with credential injection: deploy with hack/install-ate-kind.sh --deploy-atenet --experimental-use-sdsmint --experimental-egress-credential-injection, then set E2E_EGRESS_CREDINJECT=1")
}
env, err := e2e.CheckEnv("BUCKET_NAME", "KO_DOCKER_REPO")
if err != nil {
t.Fatalf("CheckEnv failed: %v", err)
}
ctx := context.Background()
clients := e2e.GetClients()
e2e.DeployCredentialProvider(t)
probeNamespace, _ = e2e.DeployProbe(t, env["BUCKET_NAME"], "egresscredinject", e2e.WithTrustBundle())
const id = "probe-credinject"
createAndResumeActor(t, ctx, clients, id)
waitForActorState(t, ctx, clients, id, ateapipb.ActorState_ACTOR_STATE_RUNNING)
rc, err := e2e.NewRouterClient(ctx)
if err != nil {
t.Fatalf("NewRouterClient: %v", err)
}
defer rc.Close()
// The gateway discards the placeholder and forwards the credential in its
// place, so the actor cannot choose the value that leaves.
wantHeader := "Bearer " + e2e.CredentialInjectionToken
replaced := fetchEcho(t, ctx, rc, id, echoOrigin, withPlaceholder)
if got := assertEchoedAuthorization(t, "injection fetch", replaced); got != wantHeader {
t.Errorf("upstream received Authorization %q, want the injected %q", got, wantHeader)
}
// The same origin over plain HTTP: the cleartext leg skips injection and
// passes the request through, so the fetch succeeds and the upstream sees
// the placeholder, not the credential. (The probe does not follow
// redirects, so an origin-side upgrade to HTTPS would surface as a non-200
// here rather than silently re-running the TLS case.)
cleartext := fetchEcho(t, ctx, rc, id, echoOriginPlain, withPlaceholder)
if cleartext.Error != "" {
t.Errorf("cleartext fetch of %s failed at the transport: %s", echoOriginPlain, cleartext.Error)
} else if cleartext.Status != "200" {
t.Errorf("cleartext fetch of %s: status %s, want 200 (the request should pass through without the credential)", echoOriginPlain, cleartext.Status)
} else if got := decodeEchoedHeaders(t, "cleartext fetch", cleartext.Body)["Authorization"]; got != placeholder {
t.Errorf("cleartext request arrived with Authorization %q, want the actor's placeholder %q: injection must be skipped on a cleartext wire", got, placeholder)
}
// Fail closed: each of these rules names a credential that cannot be
// injected, and the denial must be the mapped status, not a request that
// went out without the credential. No retries: the gateway answers these
// itself.
tests := []struct {
name string
origin string
wantStatus string
}{{
name: "unfetchable secret",
origin: unfetchableOrigin,
wantStatus: "403",
}, {
name: "unserved provider",
origin: unservedOrigin,
wantStatus: "500",
}, {
name: "unauthorized namespace",
origin: unauthorizedOrigin,
wantStatus: "403",
}}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got := probeFetch(t, ctx, rc, id, tt.origin, withPlaceholder)
if got.Error != "" {
t.Fatalf("fetch of %s failed at the transport (%s), want an HTTP %s from the gateway", tt.origin, got.Error, tt.wantStatus)
}
if got.Status != tt.wantStatus {
t.Errorf("fetch of %s returned status %s, want %s", tt.origin, got.Status, tt.wantStatus)
}
})
}
}
const echoRetryWindow = 2 * time.Minute
// fetchEcho is probeFetch for echo fetches that should return 200. It retries
// transient failures for up to echoRetryWindow:
//
// - certificate errors, until the gateway's signing pool propagates;
// - 503, while the gateway reconnects to a redeployed provider;
// - 502/503/504 from httpbin.org itself.
func fetchEcho(t *testing.T, ctx context.Context, rc *e2e.RouterClient, id, origin string, extraParams []string) fetchResponse {
t.Helper()
deadline := time.Now().Add(echoRetryWindow)
for {
resp := probeFetch(t, ctx, rc, id, origin, extraParams)
if !transientEchoFailure(resp) || time.Now().After(deadline) {
return resp
}
t.Logf("fetch of %s: transient failure (status %q, error %q), retrying", origin, resp.Status, resp.Error)
time.Sleep(5 * time.Second)
}
}
// transientEchoFailure reports whether fetchEcho should retry resp.
func transientEchoFailure(resp fetchResponse) bool {
if resp.Error != "" {
return strings.Contains(resp.Error, "certificate") || strings.Contains(resp.Error, "x509")
}
switch resp.Status {
case "502", "503", "504":
return true
}
return false
}
// echoedHeaders is the echo origin's response shape: the request headers it
// received, echoed back. (httpbin.org/headers returns {"headers": {...}}.)
type echoedHeaders struct {
Headers map[string]string `json:"headers"`
}
// decodeEchoedHeaders parses the echo origin's body into the headers the
// upstream received.
func decodeEchoedHeaders(t *testing.T, step, body string) map[string]string {
t.Helper()
var echoed echoedHeaders
if err := json.Unmarshal([]byte(body), &echoed); err != nil {
t.Fatalf("%s: decoding the echo origin's body: %v (body %q)", step, err, body)
}
return echoed.Headers
}
// assertEchoedAuthorization fails on any transport- or HTTP-level failure of
// an echo fetch and returns the Authorization value the upstream received.
func assertEchoedAuthorization(t *testing.T, step string, resp fetchResponse) string {
t.Helper()
if resp.Error != "" {
t.Fatalf("%s: TLS through the MITM egress gateway failed: %s", step, resp.Error)
}
if resp.Status != "200" {
t.Fatalf("%s: status %s, want 200 (an injection failure would deny with 403/500/503; is the provider deployed and the gateway installed with --experimental-egress-credential-injection?) body %q", step, resp.Status, resp.Body)
}
return decodeEchoedHeaders(t, step, resp.Body)["Authorization"]
}
type fetchResponse struct {
Status string `json:"status"`
Error string `json:"error"`
Body string `json:"body"`
}
// probeFetch asks the probe to fetch origin with the projected trust bundle,
// passing any extra pre-encoded query parameters through. Router-level
// failures are retried for up to 30s (a resume can return before the route
// reaches the router's xDS snapshot); probe-level TLS failures are results,
// returned for the caller to assert on.
func probeFetch(t *testing.T, ctx context.Context, rc *e2e.RouterClient, id, origin string, extraParams []string) fetchResponse {
t.Helper()
path := "/fetch?roots=bundle&url=" + url.QueryEscape(origin)
for _, p := range extraParams {
path += "&" + p
}
ref := resources.ActorRef{Atespace: probeNamespace, Name: id}
deadline := time.Now().Add(30 * time.Second)
for attempt := 1; ; attempt++ {
resp, err := rc.Get(ctx, ref, path)
if err != nil {
t.Fatalf("GET %s for %q: %v", path, id, err)
}
body, readErr := io.ReadAll(resp.Body)
resp.Body.Close()
if readErr != nil {
t.Fatalf("reading %s response for %q: %v", path, id, readErr)
}
if resp.StatusCode == http.StatusOK {
var out fetchResponse
if err := json.Unmarshal(body, &out); err != nil {
t.Fatalf("decoding %s response for %q: %v (body %q)", path, id, err, body)
}
return out
}
if time.Now().After(deadline) {
t.Fatalf("GET %s for %q: status %d after %d attempts, body %q", path, id, resp.StatusCode, attempt, body)
}
t.Logf("GET %s for %q: attempt %d: status %d, body %q; retrying", path, id, attempt, resp.StatusCode, body)
time.Sleep(2 * time.Second)
}
}
// createAndResumeActor first deletes any actor left by an earlier run, since
// actor records outlive the fixture namespace.
func createAndResumeActor(t *testing.T, ctx context.Context, clients *e2e.Clients, id string) {
t.Helper()
ref := &ateapipb.ObjectRef{Atespace: probeNamespace, Name: id}
// NotFound is the normal case on a fresh run. Other errors don't stop the
// test, but they are logged in case CreateActor then fails.
if _, err := clients.SubstrateAPI.SuspendActor(ctx, &ateapipb.SuspendActorRequest{Actor: ref}); err != nil && status.Code(err) != codes.NotFound {
t.Logf("removing leftover actor %q: SuspendActor: %v", id, err)
}
if _, err := clients.SubstrateAPI.DeleteActor(ctx, &ateapipb.DeleteActorRequest{Actor: ref}); err != nil && status.Code(err) != codes.NotFound {
t.Logf("removing leftover actor %q: DeleteActor: %v", id, err)
}
if _, err := clients.SubstrateAPI.CreateActor(ctx, &ateapipb.CreateActorRequest{Actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{Atespace: probeNamespace, Name: id},
ActorTemplate: &ateapipb.ObjectRef{Atespace: probeNamespace, Name: probeTemplate},
}}); err != nil {
t.Fatalf("CreateActor %q: %v", id, err)
}
t.Cleanup(func() {
if _, err := clients.SubstrateAPI.SuspendActor(ctx, &ateapipb.SuspendActorRequest{Actor: ref}); err != nil {
t.Logf("cleanup: SuspendActor %q: %v", id, err)
}
if _, err := clients.SubstrateAPI.DeleteActor(ctx, &ateapipb.DeleteActorRequest{Actor: ref}); err != nil {
t.Logf("cleanup: DeleteActor %q failed, actor leaked (remove with: kubectl ate delete actor %s -a %s): %v", id, id, probeNamespace, err)
}
})
// One https rule per hostname, each carrying the injection whose outcome
// that host is used to observe, and only these hosts are allowed at all.
// echoHost also gets an http rule with the same injection, which the
// cleartext fetch uses to prove the gateway skips it there.
e2e.EnsureEgressPolicy(t, ctx, clients, ref,
e2e.EgressInjectHeader("Authorization", "Bearer ", e2e.CredentialInjectionURI, echoHost),
e2e.EgressInjectHeaderHTTP("Authorization", "Bearer ", e2e.CredentialInjectionURI, echoHost),
e2e.EgressInjectHeader("Authorization", "Bearer ",
"ate-secret://k8s.io/default/"+e2e.CredentialSecretsNamespace+"/no-such-secret/token", unfetchableHost),
e2e.EgressInjectHeader("Authorization", "Bearer ",
"ate-secret://other.io/default/"+e2e.CredentialSecretsNamespace+"/api-token/token", unservedHost),
e2e.EgressInjectHeader("Authorization", "Bearer ",
"ate-secret://k8s.io/default/kube-system/api-token/token", unauthorizedHost),
)
if _, err := clients.SubstrateAPI.ResumeActor(ctx, &ateapipb.ResumeActorRequest{Actor: ref}); err != nil {
t.Fatalf("ResumeActor %q: %v", id, err)
}
}
func waitForActorState(t *testing.T, ctx context.Context, clients *e2e.Clients, actorName string, want ateapipb.ActorState) {
t.Helper()
deadline := time.Now().Add(60 * time.Second)
for time.Now().Before(deadline) {
resp, err := clients.SubstrateAPI.GetActor(ctx, &ateapipb.GetActorRequest{
Actor: &ateapipb.ObjectRef{Atespace: probeNamespace, Name: actorName},
})
if err == nil && resp.GetStatus().GetState() == want {
return
}
time.Sleep(1 * time.Second)
}
t.Fatalf("timed out waiting for actor %q to reach state %v", actorName, want)
}
@@ -0,0 +1,24 @@
// 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 egresscredinject
import (
"os"
"testing"
"github.com/agent-substrate/substrate/internal/e2e"
)
func TestMain(m *testing.M) { os.Exit(e2e.RunTestMain(m)) }