egress: add an example credential provider which reads from k8s secret (#1335)

Add an example gRPC credential-provider plugin which reads from k8s
secrets for the egress
credential-injection path. 

 - It resolves `ate-secret://k8s.io/default/<namespace>/<secret>/<key>`
URIs to Kubernetes Secret values, so Substrate never stores secrets — it
only
brokers a read that the provider is authorized to perform. 
 - It is the only
component in the injection path with Kubernetes secerts access; the
egress gateway and
its injector never read Secrets directly.

### What's included

- **`cmd/credential-provider/kubernetes-secrets`** — the gRPC service:
- Parses and validates `ate-secret://` URIs, resolving the requested
Secret
    (with single-key fallback when the URI omits a key).
- Enforces an **atespace→namespace authorization policy**
(default-deny),
    derived from the caller's attested actor SPIFFE ID.
- Serves over **mutual TLS**, requiring the caller's client cert to
chain to
    the trust bundle *and* carry the egress injector's SPIFFE SAN.
- **Manifests** (`manifests/egress-credential-injection/`) — Deployment,
Service, ServiceAccount + RBAC, a sample namespace-policy ConfigMap, and
a
  sample Secret.
- Renames the default provider address/service from `credprovider` to
`k8s-credential-provider` across `ate-setup`, install scripts, and the
egress
  injection overlay.

### atespace → namespace authorization

Beyond the mTLS check that only the egress injector may call the
provider, each
request is authorized against a **default-deny atespace→namespace
policy**. The
provider derives the requesting actor's atespace from its attested
SPIFFE ID and
resolves a Secret only if that atespace is explicitly granted access to
the URI's
namespace.
This commit is contained in:
yufan-su
2026-09-22 10:13:25 +00:00
committed by GitHub
parent cdac9baef8
commit 57234a866a
21 changed files with 1455 additions and 16 deletions
+1
View File
@@ -45,6 +45,7 @@ CONTROL_PLANE_IMAGES := ./cmd/ateapi \
./cmd/atecontroller \
./cmd/atelet \
./cmd/atenet \
./cmd/credential-provider/kubernetes-secrets \
./cmd/podcertcontroller
WORKER_IMAGES := ./cmd/ateom-gvisor \
./cmd/ateom-microvm
+2 -2
View File
@@ -94,8 +94,8 @@ func init() {
f.BoolVar(&opts.ExperimentalUseSDSMint, "experimental-use-sdsmint", false, "Deploy egress gateway with dynamic per-SNI certificate minting")
f.StringVar(&opts.AdditionalEgressExtprocService, "experimental-additional-egress-extproc-service", "", "Run an additional ext_proc authorization filter served by NS/SVC:PORT (requires --experimental-use-sdsmint)")
f.BoolVar(&opts.ExperimentalEgressCredentialInjection, "experimental-egress-credential-injection", false, "Point the egress gateway's MITM-leg handler at a credential provider so a matching EgressPolicy rule injects its credential (requires --experimental-use-sdsmint and --atenet-dataplane=envoy)")
f.StringVar(&opts.CredentialProviderName, "credential-provider-name", "", "Credential provider the injector serves, as a ate-secret:// prefix (default ate-secret://kubernetes.io)")
f.StringVar(&opts.CredentialProviderAddress, "credential-provider-address", "", "Address the egress gateway dials the credential provider at (default credprovider.ate-system.svc:50051)")
f.StringVar(&opts.CredentialProviderName, "credential-provider-name", "", "Credential provider the injector serves, as a ate-secret:// prefix (default ate-secret://k8s.io)")
f.StringVar(&opts.CredentialProviderAddress, "credential-provider-address", "", "Address the egress gateway dials the credential provider at (default k8s-credential-provider.ate-system.svc:50051)")
f.BoolVar(&opts.NoDevEnv, "no-dev-env", false, "Do not source .ate-dev-env.sh")
f.StringVar(&opts.ImageRepo, "image-repo", "",
+1
View File
@@ -46,6 +46,7 @@ var Components = []string{
"cmd/atenet",
"cmd/ateom-gvisor",
"cmd/ateom-microvm",
"cmd/credential-provider/kubernetes-secrets",
"cmd/podcertcontroller",
"demos/counter",
"demos/egress",
+2 -2
View File
@@ -122,11 +122,11 @@ func (e *Env) patchAtenetEgressInject(raw []byte) ([]byte, error) {
name := e.Cfg.CredentialProviderName
if name == "" {
name = "ate-secret://kubernetes.io"
name = "ate-secret://k8s.io"
}
address := e.Cfg.CredentialProviderAddress
if address == "" {
address = "credprovider.ate-system.svc:50051"
address = "k8s-credential-provider.ate-system.svc:50051"
}
serverName := address
if i := strings.LastIndex(address, ":"); i >= 0 {
+3 -3
View File
@@ -53,9 +53,9 @@ func TestPatchAtenetEgressInject(t *testing.T) {
}
}
for _, want := range []string{
"--credential-provider-name=ate-secret://kubernetes.io",
"--credential-provider-address=credprovider.ate-system.svc:50051",
"--credential-provider-server-name=credprovider.ate-system.svc",
"--credential-provider-name=ate-secret://k8s.io",
"--credential-provider-address=k8s-credential-provider.ate-system.svc:50051",
"--credential-provider-server-name=k8s-credential-provider.ate-system.svc",
} {
if !strings.Contains(string(patched), want) {
t.Errorf("patched manifest is missing spliced flag %q", want)
@@ -264,6 +264,10 @@ func validCredentialURI(raw string) bool {
return false
}
escapedPath := u.EscapedPath()
// Reject percent-encoding in the path of secret uri.
if escapedPath != u.Path {
return false
}
if !strings.HasPrefix(escapedPath, "/") || strings.HasSuffix(escapedPath, "/") {
return false
}
@@ -887,6 +887,8 @@ func TestCredentialURIValidation(t *testing.T) {
"ate-secret://kubernetes.io//provider/secret",
"ate-secret://kubernetes.io/provider/secret/",
"ate-secret://kubernetes.io:443/provider/secret",
"ate-secret://kubernetes.io/provider/sec%2Fret", // percent-encoded separator
"ate-secret://kubernetes.io/provider/sec%2Dret", // percent-encoding of any kind
} {
if validCredentialURI(uri) {
t.Errorf("validCredentialURI(%q) = true", uri)
+1 -1
View File
@@ -68,7 +68,7 @@ func NewRouterCmd() *cobra.Command {
// and only when injection is enabled: an empty --credential-provider-address
// leaves injection off, so an EgressPolicy rule requiring it is skipped and
// the request passes through without the credential.
cmd.Flags().StringVar(&cfg.CredentialProvider.Name, "credential-provider-name", "", "Credential provider this egress gateway serves, as a ate-secret:// prefix (e.g. ate-secret://kubernetes.io); a policy credential URI naming any other provider is refused. Empty disables the check (dev only)")
cmd.Flags().StringVar(&cfg.CredentialProvider.Name, "credential-provider-name", "", "Credential provider this egress gateway serves, as a ate-secret:// prefix (e.g. ate-secret://k8s.io); a policy credential URI naming any other provider is refused. Empty disables the check (dev only)")
cmd.Flags().StringVar(&cfg.CredentialProvider.Address, "credential-provider-address", "", "gRPC dial target of the credential provider the MITM-leg injector resolves secrets through. Empty (the default) disables egress credential injection")
cmd.Flags().StringVar(&cfg.CredentialProvider.CAFile, "credential-provider-ca-file", "", "CA the credential provider's serving certificate must chain to; required unless --credential-provider-insecure is set")
cmd.Flags().StringVar(&cfg.CredentialProvider.ClientCert, "credential-provider-client-cert", "", "Credential bundle presented to the credential provider as the client certificate; required unless --credential-provider-insecure is set")
@@ -82,7 +82,7 @@ func DialProvider(ctx context.Context, cfg ProviderDialConfig) (*grpc.ClientConn
}
// ProviderName reduces a configured provider — a ate-secret:// prefix such
// as ate-secret://kubernetes.io — to the provider name (the URI host) the
// as ate-secret://k8s.io — to the provider name (the URI host) the
// handler compares credential URIs against, so a URI naming another provider
// fails closed. An empty input returns an empty name, which disables the check
// (dev only).
@@ -94,9 +94,9 @@ func ProviderName(name string) (string, error) {
}
// providerNameFromURI returns the provider name of a ate-secret:// URI —
// the URI host, e.g. "kubernetes.io" in
// ate-secret://kubernetes.io/<namespace>/<secret>. The gateway uses it to
// confirm a URI targets the provider it is configured to serve.
// the URI host, e.g. "k8s.io" in
// ate-secret://k8s.io/default/<namespace>/<secret>/<key>. The gateway reads only the host,
// to confirm a URI targets the provider it is configured to serve.
func providerNameFromURI(raw string) (string, error) {
u, err := url.Parse(raw)
if err != nil {
@@ -0,0 +1,195 @@
// 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 file implements the CredentialProvider plugin API backed by Kubernetes
// Secrets. It resolves ate-secret:// URIs of the provider "k8s.io" to a Secret
// value read straight from the Kubernetes API — so Substrate never stores the
// secret, it only brokers a read the provider is authorized to perform.
package main
import (
"context"
"fmt"
"log/slog"
"net/url"
"strings"
k8serrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/credproviderpb"
)
// ProviderName is the ate-secret:// URI host this backend serves.
const ProviderName = "k8s.io"
// uriScheme is the only scheme a credential URI may carry.
const uriScheme = "ate-secret"
// LocalLocator is the reserved leading path segment naming the local Kubernetes
// API server the provider runs in — the only cluster served today.
const LocalLocator = "default"
// SecretRef is a parsed ate-secret:// URI for the k8s.io provider.
//
// Only Secrets in the local cluster are addressable today:
//
// ate-secret://k8s.io/default/<namespace>/<secret>/<key>
//
// Future work: the secret uri can grow a "cluster/<cluster>" locator for fetching remote Secrets.
//
// ate-secret://k8s.io/cluster/<cluster>/<namespace>/<secret>/<key>
type SecretRef struct {
Namespace string
Name string
// Key is the data key within the Secret to return.
Key string
}
// ParseURI parses a ate-secret:// URI of the kubernetes.io provider. It
// rejects any other scheme or provider name.
func ParseURI(raw string) (SecretRef, error) {
u, err := url.Parse(raw)
if err != nil {
return SecretRef{}, fmt.Errorf("parsing credential URI %q: %w", raw, err)
}
if u.Scheme != uriScheme {
return SecretRef{}, fmt.Errorf("malformed credential URI %q: scheme is %q, want %q", raw, u.Scheme, uriScheme)
}
if u.Host != ProviderName {
return SecretRef{}, fmt.Errorf("credential URI %q: provider is %q, this provider serves %q", raw, u.Host, ProviderName)
}
// The grammar is scheme/host/path only; a query or fragment means the caller
// assumed a syntax this provider does not honor, so reject it rather than
// silently ignore it.
if u.RawQuery != "" || u.ForceQuery || u.Fragment != "" {
return SecretRef{}, fmt.Errorf("credential URI %q: query and fragment components are not allowed", raw)
}
// Reject percent-encoding in the path of secret uri.
if u.EscapedPath() != u.Path {
return SecretRef{}, fmt.Errorf("credential URI %q: path must not contain percent-encoding", raw)
}
segments := strings.Split(strings.Trim(u.Path, "/"), "/")
for i, s := range segments {
if s == "" {
return SecretRef{}, fmt.Errorf("credential URI %q: empty path segment %d", raw, i)
}
}
// The path must begin with the "default" locator; only local Secrets are
// served. See SecretRef for the planned "cluster/<cluster>" remote form.
if segments[0] != LocalLocator {
return SecretRef{}, fmt.Errorf("credential URI %q: path must begin with %q (only local Secrets are supported), got %q", raw, LocalLocator, segments[0])
}
// tail is <namespace>/<secret>/<key>.
tail := segments[1:]
if len(tail) != 3 {
return SecretRef{}, fmt.Errorf("credential URI %q: want %s/<namespace>/<secret>/<key>, got %d trailing segments", raw, LocalLocator, len(tail))
}
return SecretRef{
Namespace: tail[0],
Name: tail[1],
Key: tail[2],
}, nil
}
// Server implements credproviderpb.CredentialProviderServer over the Kubernetes
// API.
type Server struct {
credproviderpb.UnimplementedCredentialProviderServer
client kubernetes.Interface
// nsAuth restricts which namespaces an atespace may resolve secrets from.
// Nil disables authorization (dev only): every URI namespace is allowed.
nsAuth *NamespaceAuthorizer
}
// NewServer builds a Kubernetes-backed credential provider. nsAuth enforces the
// atespace→namespace policy; pass nil to disable authorization (dev only).
func NewServer(client kubernetes.Interface, nsAuth *NamespaceAuthorizer) *Server {
return &Server{client: client, nsAuth: nsAuth}
}
// FetchSecret resolves one ate-secret:// URI to its Secret value.
func (s *Server) FetchSecret(ctx context.Context, req *credproviderpb.FetchSecretRequest) (*credproviderpb.FetchSecretResponse, error) {
ref, err := ParseURI(req.GetUri())
if err != nil {
return nil, status.Error(codes.InvalidArgument, err.Error())
}
if err := s.authorize(ctx, req.GetActorSpiffeId(), ref.Namespace); err != nil {
return nil, err
}
slog.InfoContext(ctx, "resolving credential",
slog.String("provider", ProviderName),
slog.String("namespace", ref.Namespace),
slog.String("secret", ref.Name),
slog.String("actor", req.GetActorSpiffeId()),
)
secret, err := s.client.CoreV1().Secrets(ref.Namespace).Get(ctx, ref.Name, metav1.GetOptions{})
if err != nil {
if k8serrors.IsNotFound(err) {
return nil, status.Errorf(codes.NotFound, "secret %s/%s not found", ref.Namespace, ref.Name)
}
if k8serrors.IsForbidden(err) {
return nil, status.Errorf(codes.PermissionDenied, "not permitted to read secret %s/%s", ref.Namespace, ref.Name)
}
return nil, status.Errorf(codes.Unavailable, "reading secret %s/%s: %v", ref.Namespace, ref.Name, err)
}
value, err := selectKey(secret.Data, ref.Key)
if err != nil {
return nil, status.Errorf(codes.NotFound, "secret %s/%s: %v", ref.Namespace, ref.Name, err)
}
return &credproviderpb.FetchSecretResponse{OpaqueBytes: value}, nil
}
// authorize enforces the atespace→namespace policy. It derives the atespace from
// the attested actor SPIFFE ID and denies unless the URI's namespace is in that
// atespace's allowed list.
func (s *Server) authorize(ctx context.Context, actorSpiffeID, namespace string) error {
if s.nsAuth == nil {
return nil
}
actor, err := resources.ActorRefFromSPIFFEID(actorSpiffeID)
if err != nil {
slog.WarnContext(ctx, "credential request denied: unusable actor identity", slog.Any("err", err))
return status.Error(codes.PermissionDenied, "actor identity is required and must be a valid actor SPIFFE URI")
}
if !s.nsAuth.Allowed(actor.Atespace, namespace) {
slog.WarnContext(ctx, "credential request denied: atespace not permitted for namespace",
slog.String("atespace", actor.Atespace), slog.String("namespace", namespace))
return status.Errorf(codes.PermissionDenied, "atespace %q is not permitted to resolve secrets in namespace %q", actor.Atespace, namespace)
}
return nil
}
// selectKey returns the named data entry, or an error when the Secret has no
// such key.
func selectKey(data map[string][]byte, key string) ([]byte, error) {
v, ok := data[key]
if !ok {
return nil, fmt.Errorf("key %q not present", key)
}
return v, nil
}
@@ -0,0 +1,253 @@
// 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 (
"context"
"testing"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes/fake"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"github.com/agent-substrate/substrate/pkg/proto/credproviderpb"
)
func TestParseURI(t *testing.T) {
tests := []struct {
name string
uri string
want SecretRef
wantErr bool
}{
{
name: "local with key",
uri: "ate-secret://k8s.io/default/ns1/example-api/token",
want: SecretRef{Namespace: "ns1", Name: "example-api", Key: "token"},
},
{name: "key required", uri: "ate-secret://k8s.io/default/ns1/example-api", wantErr: true},
{name: "wrong scheme", uri: "https://k8s.io/default/ns1/example-api/token", wantErr: true},
{name: "wrong provider", uri: "ate-secret://vault.io/default/ns1/example-api/token", wantErr: true},
{name: "missing locator", uri: "ate-secret://k8s.io/ns1/example-api/token", wantErr: true},
{name: "remote form not yet supported", uri: "ate-secret://k8s.io/cluster/remote-east/ns1/example-api/token", wantErr: true},
{name: "too few segments", uri: "ate-secret://k8s.io/default/ns1", wantErr: true},
{name: "too many segments", uri: "ate-secret://k8s.io/default/ns1/example-api/token/extra", wantErr: true},
{name: "query not allowed", uri: "ate-secret://k8s.io/default/ns1/example-api/token?cluster=remote", wantErr: true},
{name: "fragment not allowed", uri: "ate-secret://k8s.io/default/ns1/example-api/token#x", wantErr: true},
{name: "percent-encoded separator", uri: "ate-secret://k8s.io/default/ns1/example-api/tok%2Fen", wantErr: true},
{name: "percent-encoding of any kind", uri: "ate-secret://k8s.io/default/ns1/example-api/tok%2Den", wantErr: true},
{name: "space in path", uri: "ate-secret://k8s.io/default/ns1/example-api/tok en", wantErr: true},
{name: "unparseable", uri: "://://", wantErr: true},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
got, err := ParseURI(tc.uri)
if tc.wantErr {
if err == nil {
t.Fatalf("ParseURI(%q) = %+v, want error", tc.uri, got)
}
return
}
if err != nil {
t.Fatalf("ParseURI(%q) unexpected error: %v", tc.uri, err)
}
if got != tc.want {
t.Errorf("ParseURI(%q) = %+v, want %+v", tc.uri, got, tc.want)
}
})
}
}
func TestNamespaceAuthorizer(t *testing.T) {
authz, err := newNamespaceAuthorizer(namespacePolicyFile{
Policies: []atespaceNamespacePolicy{
{Atespace: "team-a", AllowedNamespaces: []string{"ns1", "shared"}},
{Atespace: "team-b", AllowedNamespaces: []string{"ns2"}},
},
})
if err != nil {
t.Fatalf("newNamespaceAuthorizer: %v", err)
}
tests := []struct {
atespace, namespace string
want bool
}{
{"team-a", "ns1", true},
{"team-a", "shared", true},
{"team-a", "ns2", false}, // namespace not in team-a's list
{"team-b", "ns2", true}, // team-b's own namespace
{"team-c", "ns1", false}, // atespace absent -> default deny
{"team-a", "", false}, // empty namespace
}
for _, tc := range tests {
if got := authz.Allowed(tc.atespace, tc.namespace); got != tc.want {
t.Errorf("Allowed(%q, %q) = %v, want %v", tc.atespace, tc.namespace, got, tc.want)
}
}
// An empty file denies everything.
empty, err := newNamespaceAuthorizer(namespacePolicyFile{})
if err != nil {
t.Fatalf("newNamespaceAuthorizer(empty): %v", err)
}
if empty.Allowed("team-a", "ns1") {
t.Error("empty authorizer allowed team-a/ns1, want deny")
}
// A policy without an atespace is rejected.
if _, err := newNamespaceAuthorizer(namespacePolicyFile{
Policies: []atespaceNamespacePolicy{{AllowedNamespaces: []string{"ns1"}}},
}); err == nil {
t.Error("newNamespaceAuthorizer accepted a policy with no atespace, want error")
}
}
func TestFetchSecretAuthorization(t *testing.T) {
secret := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{Name: "example-api", Namespace: "ns1"},
Data: map[string][]byte{"token": []byte("s3cr3t")},
}
authz, err := newNamespaceAuthorizer(namespacePolicyFile{
Policies: []atespaceNamespacePolicy{{Atespace: "team-a", AllowedNamespaces: []string{"ns1"}}},
})
if err != nil {
t.Fatalf("newNamespaceAuthorizer: %v", err)
}
const teamAURI = "spiffe://substrate-actor.local/atespace/team-a/actor/my-actor"
const teamBURI = "spiffe://substrate-actor.local/atespace/team-b/actor/my-actor"
tests := []struct {
name string
actorSpiffeID string
uri string
wantCode codes.Code
}{
{
name: "allowed",
actorSpiffeID: teamAURI,
uri: "ate-secret://k8s.io/default/ns1/example-api/token",
},
{
name: "namespace not permitted",
actorSpiffeID: teamAURI,
uri: "ate-secret://k8s.io/default/ns2/example-api/token",
wantCode: codes.PermissionDenied,
},
{
name: "unknown atespace",
actorSpiffeID: teamBURI,
uri: "ate-secret://k8s.io/default/ns1/example-api/token",
wantCode: codes.PermissionDenied,
},
{
name: "garbage identity",
actorSpiffeID: "not-a-spiffe-uri",
uri: "ate-secret://k8s.io/default/ns1/example-api/token",
wantCode: codes.PermissionDenied,
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
srv := NewServer(fake.NewSimpleClientset(secret), authz)
resp, err := srv.FetchSecret(context.Background(), &credproviderpb.FetchSecretRequest{Uri: tc.uri, ActorSpiffeId: tc.actorSpiffeID})
if tc.wantCode != codes.OK {
if status.Code(err) != tc.wantCode {
t.Fatalf("code = %v, want %v (err=%v)", status.Code(err), tc.wantCode, err)
}
return
}
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if got := string(resp.GetOpaqueBytes()); got != "s3cr3t" {
t.Errorf("secret = %q, want s3cr3t", got)
}
})
}
// With no authorizer configured, enforcement is bypassed entirely.
t.Run("nil authorizer bypasses", func(t *testing.T) {
srv := NewServer(fake.NewSimpleClientset(secret), nil)
if _, err := srv.FetchSecret(context.Background(), &credproviderpb.FetchSecretRequest{
Uri: "ate-secret://k8s.io/default/ns1/example-api/token",
ActorSpiffeId: "not-a-spiffe-uri",
}); err != nil {
t.Fatalf("nil authorizer should not enforce, got %v", err)
}
})
}
func TestFetchSecret(t *testing.T) {
secret := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{Name: "example-api", Namespace: "ns1"},
Data: map[string][]byte{
"token": []byte("s3cr3t"),
},
}
tests := []struct {
name string
uri string
want string
wantCode codes.Code
}{
{
name: "explicit key",
uri: "ate-secret://k8s.io/default/ns1/example-api/token",
want: "s3cr3t",
},
{
name: "missing key",
uri: "ate-secret://k8s.io/default/ns1/example-api/nope",
wantCode: codes.NotFound,
},
{
name: "secret not found",
uri: "ate-secret://k8s.io/default/ns1/absent/token",
wantCode: codes.NotFound,
},
{
name: "remote form rejected",
uri: "ate-secret://k8s.io/cluster/remote-east/ns1/example-api/token",
wantCode: codes.InvalidArgument,
},
{
name: "bad uri",
uri: "ate-secret://vault.io/default/ns1/example-api",
wantCode: codes.InvalidArgument,
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
client := fake.NewSimpleClientset(secret)
srv := NewServer(client, nil)
resp, err := srv.FetchSecret(context.Background(), &credproviderpb.FetchSecretRequest{Uri: tc.uri})
if tc.wantCode != codes.OK {
if status.Code(err) != tc.wantCode {
t.Fatalf("FetchSecret(%q) code = %v, want %v (err=%v)", tc.uri, status.Code(err), tc.wantCode, err)
}
return
}
if err != nil {
t.Fatalf("FetchSecret(%q) unexpected error: %v", tc.uri, err)
}
if got := string(resp.GetOpaqueBytes()); got != tc.want {
t.Errorf("FetchSecret(%q) = %q, want %q", tc.uri, got, tc.want)
}
})
}
}
@@ -0,0 +1,228 @@
// 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 kubernetes-secrets is the Kubernetes-Secrets credential-provider
// plugin: a gRPC service that resolves ate-secret:// URIs of the k8s.io
// provider to Kubernetes Secret values. It is the only component in the egress
// credential-injection path with Kubernetes access; the egress gateway and its
// injector never read Secrets directly.
package main
import (
"context"
"crypto/tls"
"fmt"
"log/slog"
"net"
"os/signal"
"syscall"
"time"
"github.com/spf13/pflag"
"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
"google.golang.org/grpc/reflection"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"github.com/agent-substrate/substrate/internal/credbundle"
"github.com/agent-substrate/substrate/internal/serverboot"
"github.com/agent-substrate/substrate/internal/version"
"github.com/agent-substrate/substrate/pkg/proto/credproviderpb"
)
const serviceName = "credprovider"
// injectorSPIFFEID is the identity of the only caller allowed to fetch secrets:
// the egress gateway's credential injector.
const injectorSPIFFEID = "spiffe://cluster.local/ns/ate-system/sa/atenet-egress"
var (
listenAddr = pflag.String("listen-address", ":50051", "gRPC listen address")
metricsAddr = pflag.String("metrics-address", ":9090", "Prometheus/health HTTP listen address")
serverBundle = pflag.String("server-cred-bundle", "", "credential bundle (PEM key+chain) presented for serving TLS (required)")
clientCAFile = pflag.String("client-ca-file", "", "CA bundle that caller (injector) client certificates must chain to (required)")
nsPolicyFile = pflag.String("namespace-policy-file", "", "path to the atespace→namespace authorization YAML (required)")
logLevel = pflag.String("log-level", "info", "one of debug, info, warn, error")
drainGrace = pflag.Duration("drain-grace", 5*time.Second, "how long to wait for in-flight RPCs on shutdown before a hard stop")
kubeAPIQPS = pflag.Float32("kube-api-qps", 50, "Sustained queries per second allowed against the Kubernetes API.")
kubeAPIBurst = pflag.Int("kube-api-burst", 100, "Burst queries allowed against the Kubernetes API.")
)
func main() {
pflag.Parse()
ctx := context.Background()
serverboot.InitLogger()
if err := serverboot.SetLogLevel(*logLevel); err != nil {
serverboot.Fatal(ctx, "invalid --log-level", err)
}
slog.InfoContext(ctx, "starting credprovider", slog.String("version", version.String()))
if err := run(ctx); err != nil {
serverboot.Fatal(ctx, "credprovider exited with error", err)
}
}
func run(ctx context.Context) error {
mp, err := serverboot.InitMetrics(ctx, serviceName)
if err != nil {
return fmt.Errorf("init metrics: %w", err)
}
defer serverboot.ShutdownProvider("MeterProvider", mp.Shutdown)
readiness := &serverboot.Readiness{}
go serverboot.StartMetricsServer(ctx, serverboot.MetricsServerOptions{
Addr: *metricsAddr,
Readiness: readiness,
EnableHealthz: true,
})
client, err := newKubeClient()
if err != nil {
return fmt.Errorf("kubernetes client: %w", err)
}
var nsAuth *NamespaceAuthorizer
if *nsPolicyFile == "" {
return fmt.Errorf("--namespace-policy-file is required")
}
nsAuth, err = LoadNamespaceAuthorizer(*nsPolicyFile)
if err != nil {
return fmt.Errorf("namespace policy: %w", err)
}
slog.InfoContext(ctx, "loaded namespace authorization policy", slog.String("file", *nsPolicyFile))
creds, err := buildServerCreds(ctx)
if err != nil {
return fmt.Errorf("server credentials: %w", err)
}
srv := grpc.NewServer(
grpc.StatsHandler(otelgrpc.NewServerHandler()),
grpc.Creds(creds),
)
reflection.Register(srv)
credproviderpb.RegisterCredentialProviderServer(srv, NewServer(client, nsAuth))
lis, err := (&net.ListenConfig{}).Listen(ctx, "tcp", *listenAddr)
if err != nil {
return fmt.Errorf("listen on %s: %w", *listenAddr, err)
}
shutdownCtx, stop := signal.NotifyContext(ctx, syscall.SIGINT, syscall.SIGTERM)
defer stop()
go func() {
<-shutdownCtx.Done()
slog.Info("shutting down")
readiness.MarkNotReady()
done := make(chan struct{})
go func() {
srv.GracefulStop()
close(done)
}()
select {
case <-done:
case <-time.After(*drainGrace):
slog.Warn("graceful shutdown timed out; forcing stop", slog.Duration("grace", *drainGrace))
srv.Stop()
}
}()
slog.InfoContext(ctx, "credprovider listening", slog.String("address", lis.Addr().String()))
if err := srv.Serve(lis); err != nil && err != grpc.ErrServerStopped {
return fmt.Errorf("serving: %w", err)
}
return nil
}
func newKubeClient() (kubernetes.Interface, error) {
if *kubeAPIQPS <= 0 || *kubeAPIBurst <= 0 {
return nil, fmt.Errorf("--kube-api-qps and --kube-api-burst must be positive")
}
cfg, err := rest.InClusterConfig()
if err != nil {
return nil, fmt.Errorf("in-cluster config: %w", err)
}
cfg.QPS = *kubeAPIQPS
cfg.Burst = *kubeAPIBurst
return kubernetes.NewForConfig(cfg)
}
// buildServerCreds composes the mutual-TLS credentials the provider serves with:
// it presents the credential bundle to callers and requires each caller to
// present a certificate that both chains to --client-ca-file and carries the
// injector's SAN. Both --server-cred-bundle and --client-ca-file are required.
func buildServerCreds(ctx context.Context) (credentials.TransportCredentials, error) {
if *serverBundle == "" {
return nil, fmt.Errorf("--server-cred-bundle is required")
}
if *clientCAFile == "" {
return nil, fmt.Errorf("--client-ca-file is required")
}
// Load the client CA pool once so a missing or empty projection fails the
// pod promptly; GetConfigForClient below reloads it for every connection.
loadPool := credbundle.PoolLoader(*clientCAFile)
if _, err := loadPool(); err != nil {
return nil, err
}
serverCert := credbundle.Loader(*serverBundle)
verifySAN := verifyClientSAN(injectorSPIFFEID)
// GetConfigForClient builds the config anew per connection: a certificate
// signed by a newly published CA verifies without a restart.
cfg := &tls.Config{
GetConfigForClient: func(*tls.ClientHelloInfo) (*tls.Config, error) {
pool, err := loadPool()
if err != nil {
return nil, err
}
return &tls.Config{
MinVersion: tls.VersionTLS13,
GetCertificate: serverCert,
ClientAuth: tls.RequireAndVerifyClientCert,
ClientCAs: pool,
VerifyConnection: verifySAN,
}, nil
},
}
slog.InfoContext(ctx, "verifying caller client certificates",
slog.String("ca", *clientCAFile), slog.String("required_san", injectorSPIFFEID))
return credentials.NewTLS(cfg), nil
}
// verifyClientSAN returns a TLS VerifyConnection callback that accepts a caller
// only when its certificate carries expectedSAN as a URI SAN.
func verifyClientSAN(expectedSAN string) func(tls.ConnectionState) error {
return func(state tls.ConnectionState) error {
if len(state.PeerCertificates) == 0 {
return fmt.Errorf("client certificate is required")
}
leaf := state.PeerCertificates[0]
for _, u := range leaf.URIs {
if u.String() == expectedSAN {
return nil
}
}
return fmt.Errorf("client certificate URI SANs %v do not include the expected injector identity %q", leaf.URIs, expectedSAN)
}
}
@@ -0,0 +1,296 @@
// 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 (
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/tls"
"crypto/x509"
"crypto/x509/pkix"
"encoding/pem"
"math/big"
"net"
"net/url"
"os"
"path/filepath"
"testing"
"time"
"google.golang.org/grpc/credentials"
)
func certWithURIs(t *testing.T, uris ...string) *x509.Certificate {
t.Helper()
cert := &x509.Certificate{}
for _, u := range uris {
parsed, err := url.Parse(u)
if err != nil {
t.Fatalf("parsing SAN %q: %v", u, err)
}
cert.URIs = append(cert.URIs, parsed)
}
return cert
}
func TestVerifyClientSAN(t *testing.T) {
const injector = injectorSPIFFEID
tests := []struct {
name string
state tls.ConnectionState
wantErr bool
}{
{
name: "matching SAN",
state: tls.ConnectionState{PeerCertificates: []*x509.Certificate{certWithURIs(t, injector)}},
},
{
name: "matching SAN among several",
state: tls.ConnectionState{PeerCertificates: []*x509.Certificate{certWithURIs(t, "spiffe://cluster.local/ns/other/sa/x", injector)}},
},
{
name: "wrong SAN",
state: tls.ConnectionState{PeerCertificates: []*x509.Certificate{certWithURIs(t, "spiffe://cluster.local/ns/ate-system/sa/impostor")}},
wantErr: true,
},
{
name: "no URI SANs",
state: tls.ConnectionState{PeerCertificates: []*x509.Certificate{certWithURIs(t)}},
wantErr: true,
},
{
name: "no peer certificate",
state: tls.ConnectionState{},
wantErr: true,
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
err := verifyClientSAN(injector)(tc.state)
if tc.wantErr && err == nil {
t.Fatal("expected an error, got nil")
}
if !tc.wantErr && err != nil {
t.Fatalf("unexpected error: %v", err)
}
})
}
}
// TestBuildServerCredsReloadsClientCA drives real TLS handshakes against the
// credentials buildServerCreds returns, then rewrites the mounted trust bundle
// and confirms a client certificate signed by the newly published CA verifies
// on a new connection without rebuilding the credentials, while both the chain
// check and the injector SAN check still hold.
func TestBuildServerCredsReloadsClientCA(t *testing.T) {
serverCA := newCA(t, "server-ca")
serverBundlePath := writeCredBundle(t, serverCA.issue(t, certOpts{dnsNames: []string{"credprovider.test"}}))
serverRoots := x509.NewCertPool()
serverRoots.AddCert(serverCA.cert)
clientCA1 := newCA(t, "client-ca-1")
clientCA2 := newCA(t, "client-ca-2")
caFile := filepath.Join(t.TempDir(), "client-ca.pem")
writeFileWithMtime(t, caFile, clientCA1.certPEM, time.Now())
// buildServerCreds reads these package-level flags.
*serverBundle = serverBundlePath
*clientCAFile = caFile
creds, err := buildServerCreds(context.Background())
if err != nil {
t.Fatalf("buildServerCreds() error = %v", err)
}
fromCA1 := clientCA1.issue(t, certOpts{uris: []string{injectorSPIFFEID}})
fromCA2 := clientCA2.issue(t, certOpts{uris: []string{injectorSPIFFEID}})
wrongSAN := clientCA1.issue(t, certOpts{uris: []string{"spiffe://cluster.local/ns/ate-system/sa/impostor"}})
// Before rotation only CA1 is trusted.
if err := handshake(t, creds, serverRoots, fromCA1); err != nil {
t.Fatalf("handshake with CA1-signed client cert failed before rotation: %v", err)
}
if err := handshake(t, creds, serverRoots, fromCA2); err == nil {
t.Fatal("handshake with CA2-signed client cert succeeded before rotation, want chain failure")
}
// A valid chain with the wrong SAN is still rejected.
if err := handshake(t, creds, serverRoots, wrongSAN); err == nil {
t.Fatal("handshake with wrong-SAN client cert succeeded, want SAN failure")
}
// Publish CA2 as the mounted trust bundle, bumping the mtime so the change
// is visible even where timestamps are coarse.
writeFileWithMtime(t, caFile, clientCA2.certPEM, time.Now().Add(time.Second))
// The same credentials now accept a CA2-signed certificate on a new
// connection, without a restart, and CA1 is no longer trusted.
if err := handshake(t, creds, serverRoots, fromCA2); err != nil {
t.Fatalf("handshake with CA2-signed client cert failed after rotation: %v", err)
}
if err := handshake(t, creds, serverRoots, fromCA1); err == nil {
t.Fatal("handshake with CA1-signed client cert succeeded after rotation, want chain failure")
}
}
// handshake performs one TLS handshake against creds, presenting clientCert and
// trusting the server with serverRoots. It reports the first end to fail.
func handshake(t *testing.T, creds credentials.TransportCredentials, serverRoots *x509.CertPool, clientCert issued) error {
t.Helper()
lis, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
defer lis.Close()
serverErr := make(chan error, 1)
go func() {
conn, err := lis.Accept()
if err != nil {
serverErr <- err
return
}
defer conn.Close()
_, _, err = creds.ServerHandshake(conn)
serverErr <- err
}()
clientCfg := &tls.Config{
MinVersion: tls.VersionTLS13,
RootCAs: serverRoots,
ServerName: "credprovider.test",
Certificates: []tls.Certificate{{Certificate: [][]byte{clientCert.certDER}, PrivateKey: clientCert.key}},
// The gRPC server credentials enforce ALPN, so offer "h2".
NextProtos: []string{"h2"},
}
conn, clientErr := tls.Dial("tcp", lis.Addr().String(), clientCfg)
if clientErr == nil {
clientErr = conn.Handshake()
conn.Close()
}
if err := <-serverErr; err != nil {
return err
}
return clientErr
}
// ca is a self-signed certificate authority used to issue test certificates.
type ca struct {
cert *x509.Certificate
key *ecdsa.PrivateKey
certPEM []byte
}
func newCA(t *testing.T, cn string) *ca {
t.Helper()
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
t.Fatalf("generate CA key: %v", err)
}
tmpl := &x509.Certificate{
SerialNumber: big.NewInt(1),
Subject: pkix.Name{CommonName: cn},
NotBefore: time.Now().Add(-time.Hour),
NotAfter: time.Now().Add(time.Hour),
KeyUsage: x509.KeyUsageCertSign,
BasicConstraintsValid: true,
IsCA: true,
}
der, err := x509.CreateCertificate(rand.Reader, tmpl, tmpl, &key.PublicKey, key)
if err != nil {
t.Fatalf("create CA certificate: %v", err)
}
cert, err := x509.ParseCertificate(der)
if err != nil {
t.Fatalf("parse CA certificate: %v", err)
}
return &ca{cert: cert, key: key, certPEM: pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der})}
}
type certOpts struct {
dnsNames []string
uris []string
}
// issued is a leaf certificate and its private key.
type issued struct {
certDER []byte
key *ecdsa.PrivateKey
}
func (c *ca) issue(t *testing.T, opts certOpts) issued {
t.Helper()
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
t.Fatalf("generate leaf key: %v", err)
}
var uris []*url.URL
for _, u := range opts.uris {
parsed, err := url.Parse(u)
if err != nil {
t.Fatalf("parse URI SAN %q: %v", u, err)
}
uris = append(uris, parsed)
}
tmpl := &x509.Certificate{
SerialNumber: big.NewInt(time.Now().UnixNano()),
Subject: pkix.Name{CommonName: "leaf"},
NotBefore: time.Now().Add(-time.Hour),
NotAfter: time.Now().Add(time.Hour),
KeyUsage: x509.KeyUsageDigitalSignature,
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth, x509.ExtKeyUsageClientAuth},
DNSNames: opts.dnsNames,
URIs: uris,
}
der, err := x509.CreateCertificate(rand.Reader, tmpl, c.cert, &key.PublicKey, c.key)
if err != nil {
t.Fatalf("create leaf certificate: %v", err)
}
return issued{certDER: der, key: key}
}
// writeCredBundle writes a credential bundle (PKCS8 key + leaf certificate) in
// the format credbundle.Parse expects and returns its path.
func writeCredBundle(t *testing.T, leaf issued) string {
t.Helper()
keyDER, err := x509.MarshalPKCS8PrivateKey(leaf.key)
if err != nil {
t.Fatalf("marshal PKCS8 key: %v", err)
}
bundle := append(
pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: leaf.certDER}),
pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: keyDER})...,
)
path := filepath.Join(t.TempDir(), "server-bundle.pem")
if err := os.WriteFile(path, bundle, 0o600); err != nil {
t.Fatalf("write credential bundle: %v", err)
}
return path
}
func writeFileWithMtime(t *testing.T, path string, data []byte, mtime time.Time) {
t.Helper()
if err := os.WriteFile(path, data, 0o600); err != nil {
t.Fatalf("write %s: %v", path, err)
}
if err := os.Chtimes(path, mtime, mtime); err != nil {
t.Fatalf("chtimes %s: %v", path, err)
}
}
@@ -0,0 +1,87 @@
// 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 (
"fmt"
"os"
"sigs.k8s.io/yaml"
)
// namespacePolicyFile is the YAML the authorizer loads: a list of grants, each
// mapping one atespace to the namespaces whose Secrets it may resolve.
type namespacePolicyFile struct {
Policies []atespaceNamespacePolicy `json:"policies"`
}
type atespaceNamespacePolicy struct {
Atespace string `json:"atespace"`
AllowedNamespaces []string `json:"allowedNamespaces"`
}
// NamespaceAuthorizer decides whether an atespace may resolve secrets in a given
// Kubernetes namespace. It is default-deny: an atespace absent from the mapping
// can resolve nothing.
type NamespaceAuthorizer struct {
// allowed maps atespace -> set of permitted namespaces.
allowed map[string]map[string]struct{}
}
// LoadNamespaceAuthorizer reads the YAML policy file at path and builds an
// authorizer, so a malformed file fails startup rather than the first request.
func LoadNamespaceAuthorizer(path string) (*NamespaceAuthorizer, error) {
data, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("reading namespace policy file %q: %w", path, err)
}
var file namespacePolicyFile
if err := yaml.Unmarshal(data, &file); err != nil {
return nil, fmt.Errorf("parsing namespace policy file %q: %w", path, err)
}
return newNamespaceAuthorizer(file)
}
// newNamespaceAuthorizer builds an authorizer over a parsed policy file,
// validating that each grant names an atespace.
func newNamespaceAuthorizer(file namespacePolicyFile) (*NamespaceAuthorizer, error) {
allowed := make(map[string]map[string]struct{})
for i, p := range file.Policies {
if p.Atespace == "" {
return nil, fmt.Errorf("namespace policy %d: atespace is required", i)
}
set := allowed[p.Atespace]
if set == nil {
set = make(map[string]struct{})
allowed[p.Atespace] = set
}
for _, ns := range p.AllowedNamespaces {
set[ns] = struct{}{}
}
}
return &NamespaceAuthorizer{allowed: allowed}, nil
}
// Allowed reports whether atespace may resolve secrets in namespace. Default
// deny: an atespace absent from the mapping, or a namespace not in its list, is
// refused.
func (a *NamespaceAuthorizer) Allowed(atespace, namespace string) bool {
set, ok := a.allowed[atespace]
if !ok {
return false
}
_, ok = set[namespace]
return ok
}
@@ -69,8 +69,8 @@ patch_atenet_egress_inject() {
return 1
fi
local name="${ATE_CREDENTIAL_PROVIDER_NAME:-ate-secret://kubernetes.io}"
local address="${ATE_CREDENTIAL_PROVIDER_ADDRESS:-credprovider.ate-system.svc:50051}"
local name="${ATE_CREDENTIAL_PROVIDER_NAME:-ate-secret://k8s.io}"
local address="${ATE_CREDENTIAL_PROVIDER_ADDRESS:-k8s-credential-provider.ate-system.svc:50051}"
# Pin the provider's serving-cert SAN to its Service DNS name (the address
# without the port), so a rotated cert for the same Service still validates.
local server_name="${address%:*}"
+2 -2
View File
@@ -99,11 +99,11 @@ function usage() {
echo " itself is deployed separately. Implies --experimental-use-sdsmint; requires"
echo " --atenet-dataplane=envoy. (experimental)"
echo " --credential-provider-name NAME Provider the injector serves, as a ate-secret:// prefix"
echo " (default ate-secret://kubernetes.io). Only meaningful with"
echo " (default ate-secret://k8s.io). Only meaningful with"
echo " --experimental-egress-credential-injection. (experimental)"
echo " --credential-provider-address HOST:PORT"
echo " Address the egress gateway dials the credential provider at"
echo " (default credprovider.ate-system.svc:50051). Only meaningful with"
echo " (default k8s-credential-provider.ate-system.svc:50051). Only meaningful with"
echo " --experimental-egress-credential-injection. (experimental)"
echo ""
echo "Infrastructure components:"
+65
View File
@@ -53,6 +53,21 @@ func ClientLoader(path string) func(*tls.CertificateRequestInfo) (*tls.Certifica
}
}
// PoolLoader reads a set of trust anchors from a PEM trust-bundle file, as
// projected from a Kubernetes ClusterTrustBundle, and returns a function that
// yields the parsed *x509.CertPool.
//
// A tls.Config's ClientCAs (and RootCAs) is frozen once the config is in use,
// so a pool built at startup never sees a CA rotation. Calling the returned
// function per connection — from GetConfigForClient on the server side — keeps
// verification current: the parsed pool is cached and the file re-read only
// when it changes, mirroring Loader, so a rotation is picked up on the next
// handshake without paying the read and parse cost when nothing changed.
func PoolLoader(path string) func() (*x509.CertPool, error) {
c := &poolCache{path: path}
return c.get
}
// certCache holds the parse of a credential bundle file together with the stat
// of the file it was parsed from, so unchanged files are not re-parsed on
// every TLS handshake.
@@ -103,6 +118,41 @@ func (c *certCache) get() (*tls.Certificate, error) {
return cert, nil
}
// poolCache holds the parse of a trust-bundle file together with the stat of
// the file it was parsed from, so unchanged files are not re-parsed on every
// TLS handshake. It mirrors certCache; see get for the change-detection and
// concurrency reasoning.
type poolCache struct {
path string
mu sync.Mutex
fi os.FileInfo
pool *x509.CertPool
}
// get returns the parsed trust pool, re-reading the file only when it has
// changed since the last successful parse. Change detection and error handling
// match certCache.get.
func (c *poolCache) get() (*x509.CertPool, error) {
c.mu.Lock()
defer c.mu.Unlock()
fi, err := os.Stat(c.path)
if err != nil {
return nil, fmt.Errorf("while getting file info for trust bundle %q: %w", c.path, err)
}
if c.pool != nil && os.SameFile(c.fi, fi) && fi.ModTime().Equal(c.fi.ModTime()) && fi.Size() == c.fi.Size() {
return c.pool, nil
}
pool, err := ParsePool(c.path)
if err != nil {
return nil, err
}
c.fi, c.pool = fi, pool
return pool, nil
}
// Parse reads a private key and certificate chain from a credential bundle file as written by the
// Kubernetes Pod Certificates mechanism.
func Parse(bundlePath string) (*tls.Certificate, error) {
@@ -155,3 +205,18 @@ func Parse(bundlePath string) (*tls.Certificate, error) {
PrivateKey: leafKey,
}, nil
}
// ParsePool reads a PEM trust-bundle file into an *x509.CertPool. It returns an
// error if the file holds no certificates: an empty trust pool would silently
// reject every peer, which is never what a workload should receive.
func ParsePool(path string) (*x509.CertPool, error) {
pemBytes, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("while reading trust bundle: %w", err)
}
pool := x509.NewCertPool()
if !pool.AppendCertsFromPEM(pemBytes) {
return nil, fmt.Errorf("trust bundle %q contains no certificates", path)
}
return pool, nil
}
+84
View File
@@ -216,6 +216,83 @@ func TestClientLoaderCachesAndReloads(t *testing.T) {
}
}
func TestPoolLoaderServesCachedParseWhileFileUnchanged(t *testing.T) {
trust := makeTrustBundle(t, 5)
path := writeBundle(t, trust)
getPool := PoolLoader(path)
first, err := getPool()
if err != nil {
t.Fatalf("PoolLoader() first call error = %v", err)
}
fi, err := os.Stat(path)
if err != nil {
t.Fatalf("stat bundle: %v", err)
}
// Overwrite with same-length garbage and restore the mtime so identity,
// size, and mtime all still match: the cached pool must be served, since a
// re-read would fail loudly on the garbage.
if err := os.WriteFile(path, bytes.Repeat([]byte("x"), len(trust)), 0o600); err != nil {
t.Fatalf("overwrite bundle: %v", err)
}
if err := os.Chtimes(path, fi.ModTime(), fi.ModTime()); err != nil {
t.Fatalf("restore mtime: %v", err)
}
second, err := getPool()
if err != nil {
t.Fatalf("PoolLoader() with unchanged stat error = %v", err)
}
if !second.Equal(first) {
t.Fatalf("PoolLoader() returned a re-parsed pool, want the cached one")
}
}
func TestPoolLoaderPicksUpProjectedVolumeRotation(t *testing.T) {
path := writeProjectedBundle(t, makeTrustBundle(t, 1))
getPool := PoolLoader(path)
before, err := getPool()
if err != nil {
t.Fatalf("PoolLoader() first call error = %v", err)
}
rotated := makeTrustBundle(t, 2)
if err := rotateProjectedBundle(path, rotated); err != nil {
t.Fatalf("rotate bundle: %v", err)
}
after, err := getPool()
if err != nil {
t.Fatalf("PoolLoader() after rotation error = %v", err)
}
if after.Equal(before) {
t.Fatalf("PoolLoader() did not pick up the rotated trust bundle")
}
want, err := ParsePool(path)
if err != nil {
t.Fatalf("ParsePool() error = %v", err)
}
if !after.Equal(want) {
t.Fatalf("PoolLoader() pool does not match the rotated trust bundle")
}
}
func TestPoolLoaderErrorWhenBundleMissing(t *testing.T) {
getPool := PoolLoader(t.TempDir() + "/absent.pem")
if _, err := getPool(); err == nil {
t.Fatalf("PoolLoader() error = nil, want missing-file error")
}
}
func TestParsePoolRejectsBundleWithoutCertificates(t *testing.T) {
path := writeBundle(t, []byte("not a certificate"))
if _, err := ParsePool(path); err == nil {
t.Fatalf("ParsePool() error = nil, want no-certificates error")
}
}
func TestLoaderConcurrentHandshakes(t *testing.T) {
bundles := [][]byte{makeBundle(t, 1), makeBundle(t, 2)}
path := writeProjectedBundle(t, bundles[0])
@@ -314,6 +391,13 @@ func makeBundle(t *testing.T, serial int64) []byte {
)
}
// makeTrustBundle returns a PEM trust bundle of a single CERTIFICATE block
// whose serial number lets tests tell one bundle from another.
func makeTrustBundle(t *testing.T, serial int64) []byte {
t.Helper()
return pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: generateCertificate(t, serial)})
}
func leafSerial(t *testing.T, cert *tls.Certificate) int64 {
t.Helper()
if cert == nil || cert.Leaf == nil {
@@ -0,0 +1,158 @@
# 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 credential provider: a gRPC service that resolves ate-secret:// URIs
# of the k8s.io class to Kubernetes Secret values. It is the ONLY
# component in the egress credential-injection path with Kubernetes access; the
# egress gateway and the injector never read Secrets.
apiVersion: v1
kind: ServiceAccount
metadata:
name: k8s-credential-provider
namespace: ate-system
---
# POC scope note: this grants read on Secrets cluster-wide so a policy can name a
# Secret in any namespace. A production deployment would scope this to the
# namespaces a provider instance is allowed to serve.
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
name: k8s-credential-provider-secret-reader
rules:
- apiGroups: [""]
resources: ["secrets"]
verbs: ["get"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
name: k8s-credential-provider-secret-reader
roleRef:
apiGroup: rbac.authorization.k8s.io
kind: ClusterRole
name: k8s-credential-provider-secret-reader
subjects:
- kind: ServiceAccount
name: k8s-credential-provider
namespace: ate-system
---
apiVersion: apps/v1
kind: Deployment
metadata:
name: k8s-credential-provider
namespace: ate-system
labels:
app: k8s-credential-provider
spec:
replicas: 1
selector:
matchLabels:
app: k8s-credential-provider
template:
metadata:
labels:
app: k8s-credential-provider
spec:
serviceAccountName: k8s-credential-provider
securityContext:
runAsUser: 65532
runAsGroup: 65532
runAsNonRoot: true
containers:
- name: k8s-credential-provider
image: ko://github.com/agent-substrate/substrate/cmd/credential-provider/kubernetes-secrets
args:
- "--listen-address=:50051"
- "--metrics-address=:9090"
# Serve with the pod's servicedns identity (SAN k8s-credential-provider.ate-system.svc)
# and require the injector to present a podidentity client cert whose chain
# verifies against the trust bundle. The provider additionally pins the
# caller's SAN to the egress gateway's identity
# (spiffe://cluster.local/ns/ate-system/sa/atenet-egress), so no other
# CA-trusted workload can fetch secrets.
- "--server-cred-bundle=/run/servicedns.podcert.ate.dev/credential-bundle.pem"
- "--client-ca-file=/run/podidentity.podcert.ate.dev/trust-bundle.pem"
# Enforce the atespace→namespace authorization policy (default-deny).
- "--namespace-policy-file=/etc/k8s-credential-provider/namespace-policy.yaml"
- "--log-level=info"
ports:
- name: grpc
containerPort: 50051
- name: metrics
containerPort: 9090
readinessProbe:
httpGet:
path: /readyz
port: metrics
periodSeconds: 10
securityContext:
allowPrivilegeEscalation: false
readOnlyRootFilesystem: true
capabilities:
drop: ["ALL"]
volumeMounts:
- name: namespace-policy
mountPath: /etc/k8s-credential-provider
readOnly: true
- name: servicedns
mountPath: /run/servicedns.podcert.ate.dev
readOnly: true
- name: podidentity
mountPath: /run/podidentity.podcert.ate.dev
readOnly: true
volumes:
- name: namespace-policy
configMap:
name: k8s-credential-provider-namespace-policy
- name: servicedns
projected:
sources:
- podCertificate:
signerName: servicedns.podcert.ate.dev/identity
keyType: ECDSAP256
credentialBundlePath: credential-bundle.pem
- clusterTrustBundle:
signerName: servicedns.podcert.ate.dev/identity
labelSelector:
matchLabels:
podcert.ate.dev/canarying: live
path: trust-bundle.pem
- name: podidentity
projected:
sources:
- podCertificate:
signerName: podidentity.podcert.ate.dev/identity
keyType: ECDSAP256
credentialBundlePath: credential-bundle.pem
- clusterTrustBundle:
signerName: podidentity.podcert.ate.dev/identity
labelSelector:
matchLabels:
podcert.ate.dev/canarying: live
path: trust-bundle.pem
---
apiVersion: v1
kind: Service
metadata:
name: k8s-credential-provider
namespace: ate-system
spec:
type: ClusterIP
selector:
app: k8s-credential-provider
ports:
- name: grpc
port: 50051
targetPort: grpc
protocol: TCP
@@ -0,0 +1,34 @@
# 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 atespace→namespace authorization policy k8s-credential-provider enforces: an actor's
# atespace may only resolve Secrets whose URI namespace is listed for it here.
# Enforcement is default-deny — an atespace absent from this file resolves
# nothing.
#
# TODO: k8s-credential-provider loads this once at startup, so editing this ConfigMap
# requires restarting the k8s-credential-provider Deployment. Make it reload dynamically.
apiVersion: v1
kind: ConfigMap
metadata:
name: k8s-credential-provider-namespace-policy
namespace: ate-system
data:
namespace-policy.yaml: |
# atespace "team-a" may resolve secrets in namespace "ns1" (matches the
# sample policy's ate-secret://k8s.io/default/ns1/example-api/token).
policies:
- atespace: team-a
allowedNamespaces:
- ns1
@@ -0,0 +1,31 @@
# 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 sample credential the policy injects. The URI in the sample policy,
# ate-secret://k8s.io/default/ns1/example-api/token
# resolves to namespace "ns1", Secret "example-api". The URI omits a key, and the
# Secret has exactly one ("token"), so that sole key is returned.
apiVersion: v1
kind: Namespace
metadata:
name: ns1
---
apiVersion: v1
kind: Secret
metadata:
name: example-api
namespace: ns1
type: Opaque
stringData:
token: "test-github-token"