Make substrate's namespace, Service names and ServiceAccount names configurable (#350)

# Description

**Substrate assumes the canonical install layout, and every deviation
fails closed.** The namespace `ate-system`, the Services `api` /
`atenet-router`, and the ServiceAccounts `atelet` / `atenet-router` are
compiled-in constants. An install in a per-developer namespace, or under
a deployment that prefixes resource names, breaks — and no failure
points at the naming.

This series makes the namespace, the Service names, and the
ServiceAccount names configurable. Every option defaults to the
canonical value. A canonical install is byte-for-byte unaffected.

| Hardcoded assumption | Failure | What it looks like instead |
|---|---|---|
| atelet's namespace in ateapi's SPIFFE check | ateapi rejects every
atelet | an mTLS handshake failure, not a naming error |
| atelet's namespace in the worker's broker check | no actor obtains a
certificate | `credential broker is not atelet` |
| NetworkPolicy's ingress namespace | the CNI drops every request to the
pool | a silent network fault |
| ateapi Service name and namespace in the client | the client cannot
authenticate | `services "api" not found`, then `invalid bearer token` |
| Resource names in `hack/install-ate.sh` and the e2e harness |
authorities land in the wrong namespace | suites die in preflight |

**A SPIFFE ID breaks on two axes.** It names a namespace and a
ServiceAccount. Both were constants, and the failure surfaces as a
rejected peer, not a missing object.

## Design notes

- **ate-controller passes the worker-side identities.** ateom already
took `--atunnel-client-identity` as a flag but relied on its default. A
new `--atunnel-broker-identity` flag alone would be inert: correct only
where the constant was already correct. The controller knows the control
plane's namespace, so it supplies both.
- **ateapi takes the whole expected identity, not a namespace.** The
dialer and `ateletauth` used the namespace only to build one string. One
place decides how atelet's identity is spelled.
- **`ateletauth` deduplicates without taking the dependency it was
avoiding.** Its constants are duplicated rather than imported so the
package does not depend on `controlapi` for three strings. `ateletdial`
would make a third copy, so the strings move to
`internal/installdefaults` instead — a leaf package of constants that
imports only the standard library, so consuming it does not reintroduce
the coupling the duplication was there to prevent.
- **The client reads environment variables, not flags.** It runs outside
the cluster. It has no downward API and nothing to discover from.
- **Namespaces come from the downward API, not flags.** Every supported
topology co-locates these components. Only Service and ServiceAccount
names, which a deployment may legitimately rename, get flags.
- **`hack/install-ate.sh` and the e2e harness are in scope.** A
relocated install cannot be bootstrapped or exercised without them. The
script refuses `ATE_NAMESPACE` or `ATE_API_SERVICE_NAME` overrides on
the manifest-applying subcommands, which would half-install:
`manifests/ate-install/` names `ate-system` and `api` literally.

## Why this is one PR

The series is stacked, not parallel. Commit 1 creates
`internal/installdefaults` and seeds seven imports; commit 2 adds
fourteen more; every later commit builds on those constants. Split into
separate PRs, each blocks on the previous merging and none reviews
independently.

The guard test cannot land first. Applied to `main` it fails on eleven
hardcoded literals across `ateletauth`, `controlapi/informer.go`,
`networkpolicy_controller.go`, both `ateom` mains, `ateclient`,
`ateletdial` and three e2e files. It is green only because the preceding
commits removed them.

A partial series still fails closed. Each commit's audit turned up more
places assuming `ate-system`, so landing commits 1-3 without 4-5 leaves
a relocated install broken later and less legibly than before.

`CONTRIBUTING.md` covers this case: when the intermediate steps are not
useful on their own, keep the change as one PR split into commits at
logical break points. Review commit by commit, and preserve the commits
on merge. If you would still rather split, the one defensible cut is by
axis — relocation (commits 1-3) then renaming (commits 4-5) with the
guard test rebased on top.

## Rebase notes (2026-09-15)

The series is rebased onto `main` at `d0d85c39`, ninety-two commits on
from the base the PR first carried:

- `main` moved actor JWT and certificate minting out of
`cmd/ateapi/internal/actoridentity` into `controlapi`, deleting the
package and dropping its inline atelet authorization check in favor of
the unified authorizer. This series previously threaded the expected
atelet identity into that check; the check no longer exists, so that
adaptation is gone. `ateletauth` survives with one caller,
`workerservice.SetWorkerCapacity`, which still takes the configured
identity.
- The same refactor replaced `BrokerConfig.ExpectedActorUID` with an
`ActorAtespace` / `ActorName` / `ActorUID` triple. `AteletSPIFFEID` is
additive and still required by `ateletdial.TLSConfig`.
- `ateletauth` (ateapi's side) and `ateletdial` (the worker's side) each
re-declared the hardcoded SPIFFE ID. Both take the expected identity as
a parameter.
- `main` deleted the atenet DNS subsystem. The series no longer touches
it, and `installdefaults` carries no DNS Service name.
- The controller passes `--atunnel-broker-identity` only when it differs
from the canonical default. An ateom old enough to predate the flag
exits on it, and `docs/upgrade.md` keeps such a pool serving during a
rolling upgrade. A relocated install gets the flag and necessarily runs
an image that accepts it.
- `main` added `internal/e2e/collector_metrics.go` (agentgateway CI
support), which addresses the router through the package-level
`routerNamespace` and `routerService` that this series replaces. Its two
call sites become `SystemNamespace()` and
`ResourceName("atenet-router")`, matching `router_client.go` and
`statusz.go`. A reviewer diffing against the older base sees those call
sites move; nothing else in that file changes.
- `main` added `cmd/credential-provider/kubernetes-secrets` (`#1335`),
whose `injectorSPIFFEID` constant is the mTLS peer check on the only
caller permitted to read Secrets. The comparison is an exact string
match, so a renamed install rejects every fetch at TLS. It becomes
`--injector-identity`, defaulting to
`installdefaults.EgressSPIFFEID(SystemNamespace)`, alongside a new
`EgressServiceAccount` constant and helper. The flag name follows
`--atunnel-client-identity`; the manifest beside it still names
`ate-system`, so a renaming deployment configures both.
- `main` replaced the dialer's worker indexer with
`DialForAteletOnNode`, dropped `WorkerPodInformer` from `controlapi`,
and removed `kataConfig` from ateom-microvm's `NewService`. The series
adapts to each narrower signature and keeps only its own added
parameter.
- The guard test's allowlist entry for the two inert `ate-system`
literals follows the code from `actoridentity.go` to
`controlapi/actor.go`. The refactor re-introduced exactly the class of
constant this series removes, so the guard earns its place.

# Testing

- `go test -race ./...` passes. `make verify` passes boilerplate,
codegen, go-modules, gofmt, golangci-lint, kube-api-linter, licenses,
metrics and postgresql-migrations. `proto-fmt` needs `clang-format`,
which is absent locally; this series touches no `.proto` files.
- Every commit builds and vets individually, via `git rebase -x 'go
build ./... && go vet ./...'`.
- On a fresh kind cluster, all twelve e2e suites pass: `capabilities`,
`combinedvolumes`, `demo`, `egressauthz`, `egressmitm`, `example`,
`identity`, `metrics`, `networking`, `networkpolicy`, `parking`,
`sizing`. The `demo`, `metrics`, `networkpolicy`, `parking` and
`networking` suites need the `--deploy-demo-counter` and
`--deploy-demo-egress` fixtures, and the MITM path needs
`--deploy-demo-egress-mitm`; without them they fail on `actor template
not found`, which reads as a control-plane fault rather than a missing
fixture.
- With the sdsmint egress gateway (`--deploy-atenet
--experimental-use-sdsmint`, `E2E_EGRESS_MITM=1`),
`TestActorEgressMITMTrust` and `TestActorEgressHTTPSByHostnameMITM` both
pass, and the `networking` suite is green at 25 passed. The two modes
are mutually exclusive by design:
`TestActorEgressHTTPSByHostnamePassthrough` covers the plain gateway and
skips under MITM, and the two MITM tests skip without it. Both modes
were run, so every case executed in one of them.
- The new unit tests use a relocated namespace and renamed
ServiceAccounts. Reintroducing each hardcoding makes them fail.

One e2e caveat that predates this series and misleads: the `identity`
suite calls `ReplaceEgressTrustPool`, which takes over the shared
`egress-mitm-ca-pool` Secret, overwrites it and registers no cleanup.
Any MITM test running afterwards fails with `certificate signed by
unknown authority`, which reads as a broken interception path rather
than a mutated fixture. Reinstall the gateway before running
`egressmitm`.

# Additional Notes

**Credential contents stay canonical.** `controlapi/actor.go` still
mints credentials naming `api.ate-system.svc`: the JWT issuer and the
certificate's Issuer CN. Relying parties validate these, so changing
them is a compatibility decision, not a lookup fix. The code carries a
TODO to make the issuer a globally unique, OIDC-compliant name.
This commit is contained in:
Jonathan Jamroga
2026-09-22 22:28:40 +00:00
committed by GitHub
parent 9e503c42b8
commit 6f45aeb2dd
62 changed files with 1042 additions and 252 deletions
+9
View File
@@ -29,6 +29,7 @@ import (
"time"
"github.com/agent-substrate/substrate/cmd/ate-setup/internal/images"
"github.com/agent-substrate/substrate/internal/installdefaults"
)
// Enumerated values for the install-shaping flags.
@@ -74,6 +75,13 @@ type Config struct {
// Kind selects the local Kind install profile (ATE_INSTALL_KIND).
Kind bool
// Namespace is the namespace the control plane is installed into, from
// ATE_NAMESPACE. It defaults to the canonical installdefaults.SystemNamespace,
// so an install that does not set it is unaffected. The checked-in manifests
// under manifests/ate-install/ name that namespace literally, so the
// manifest-applying steps refuse any other value; see Env.RequireCanonicalNamespace.
Namespace string
// Kubeconfig and Context select the target cluster. Empty Context means
// "use the current context" (the KUBECTL_CONTEXT convention).
//
@@ -301,6 +309,7 @@ func Load(opts Options) (*Config, error) {
cfg := &Config{
Root: root,
Kind: kind,
Namespace: firstNonEmpty(env["ATE_NAMESPACE"], installdefaults.SystemNamespace),
Kubeconfig: kubeconfig,
Context: firstNonEmpty(opts.Context, env["KUBECTL_CONTEXT"]),
ProjectID: env["PROJECT_ID"],
+10 -10
View File
@@ -44,7 +44,7 @@ const envHashAnnotation = "ate.dev/env-hash"
// previous installer left in the ConfigMap.
func (e *Env) CreateAPIServerEnvVars(ctx context.Context) error {
log.Step("create_api_server_env_vars")
if err := e.Kube.EnsureNamespace(ctx, NamespaceAteSystem); err != nil {
if err := e.Kube.EnsureNamespace(ctx, e.Namespace()); err != nil {
return err
}
@@ -77,10 +77,10 @@ func (e *Env) CreateAPIServerEnvVars(ctx context.Context) error {
dsn = withPoolMaxConns(dsn, e.Cfg.PostgresPoolMaxConns, dsnFromOperator)
log.Infof("POSTGRES_CONNECTION_STRING: %s", redactDSN(dsn))
if err := e.Kube.ApplyConfigMap(ctx, NamespaceAteSystem, ConfigMapAPIEnvVars, cloudSQLEnvVars(cloudsql)); err != nil {
if err := e.Kube.ApplyConfigMap(ctx, e.Namespace(), ConfigMapAPIEnvVars, cloudSQLEnvVars(cloudsql)); err != nil {
return err
}
if err := e.Kube.ApplySecret(ctx, NamespaceAteSystem, SecretAPIEnvVars,
if err := e.Kube.ApplySecret(ctx, e.Namespace(), SecretAPIEnvVars,
buildAPIServerEnvVars(dsn, e.Cfg.PostgresSchemaName())); err != nil {
return err
}
@@ -108,7 +108,7 @@ func buildAPIServerEnvVars(connString, schema string) map[string]string {
// recordedDSN reads the connection string the cluster currently runs with.
func (e *Env) recordedDSN(ctx context.Context) (string, error) {
secret, err := e.Kube.GetSecret(ctx, NamespaceAteSystem, SecretAPIEnvVars)
secret, err := e.Kube.GetSecret(ctx, e.Namespace(), SecretAPIEnvVars)
if err != nil || secret == nil {
return "", err
}
@@ -129,7 +129,7 @@ func (e *Env) applyPostgresServerCA(ctx context.Context) error {
if err != nil {
return fmt.Errorf("reading ATE_API_POSTGRES_SERVER_CA_FILE: %w", err)
}
return e.Kube.ApplySecret(ctx, NamespaceAteSystem, SecretPostgresServerCA, map[string]string{
return e.Kube.ApplySecret(ctx, e.Namespace(), SecretPostgresServerCA, map[string]string{
"server-ca.pem": string(pem),
})
}
@@ -184,7 +184,7 @@ func redactDSN(dsn string) string {
// running Deployment without a DSN on its next restart. A full deploy is safe
// because it updates the Deployment in the same run.
func (e *Env) EnsureEnvVarsSafeStandalone(ctx context.Context) error {
dep, err := e.Kube.GetDeployment(ctx, NamespaceAteSystem, "ate-api-server")
dep, err := e.Kube.GetDeployment(ctx, e.Namespace(), "ate-api-server")
if err != nil {
return err
}
@@ -209,7 +209,7 @@ func (e *Env) EnsureEnvVarsSafeStandalone(ctx context.Context) error {
// apiserver's environment, so that a changed DSN starts a rollout. Kubernetes
// does not restart pods when an envFrom ConfigMap or Secret changes.
func (e *Env) annotateAPIServerEnvHash(ctx context.Context) error {
dep, err := e.Kube.GetDeployment(ctx, NamespaceAteSystem, "ate-api-server")
dep, err := e.Kube.GetDeployment(ctx, e.Namespace(), "ate-api-server")
if err != nil {
return err
}
@@ -218,11 +218,11 @@ func (e *Env) annotateAPIServerEnvHash(ctx context.Context) error {
return nil
}
cm, err := e.Kube.GetConfigMap(ctx, NamespaceAteSystem, ConfigMapAPIEnvVars)
cm, err := e.Kube.GetConfigMap(ctx, e.Namespace(), ConfigMapAPIEnvVars)
if err != nil {
return err
}
secret, err := e.Kube.GetSecret(ctx, NamespaceAteSystem, SecretAPIEnvVars)
secret, err := e.Kube.GetSecret(ctx, e.Namespace(), SecretAPIEnvVars)
if err != nil {
return err
}
@@ -237,7 +237,7 @@ func (e *Env) annotateAPIServerEnvHash(ctx context.Context) error {
patch := fmt.Sprintf(`{"spec":{"template":{"metadata":{"annotations":{%q:%q}}}}}`,
envHashAnnotation, envHash(cmData, secretData))
return e.Kube.PatchDeployment(ctx, NamespaceAteSystem, "ate-api-server", []byte(patch))
return e.Kube.PatchDeployment(ctx, e.Namespace(), "ate-api-server", []byte(patch))
}
// envHash digests the apiserver's environment sources. Only changes matter, so
+7 -7
View File
@@ -104,7 +104,7 @@ func cloudSQLSettingsFrom(c config.CloudSQLConfig, recorded map[string]string) c
// recordedAPIServerEnvVars returns the ate-api-server-envvars ConfigMap data,
// or nil when the ConfigMap does not exist.
func (e *Env) recordedAPIServerEnvVars(ctx context.Context) (map[string]string, error) {
cm, err := e.Kube.GetConfigMap(ctx, NamespaceAteSystem, ConfigMapAPIEnvVars)
cm, err := e.Kube.GetConfigMap(ctx, e.Namespace(), ConfigMapAPIEnvVars)
if err != nil || cm == nil {
return nil, err
}
@@ -146,7 +146,7 @@ func (e *Env) resolveCloudSQL(ctx context.Context) (cloudSQLSettings, error) {
s.GSA = e.Cfg.CloudSQL.GSA
if s.GSA == "" {
gsa, err := e.Kube.ServiceAccountAnnotation(ctx, NamespaceAteSystem, "ate-api-server", workloadIdentityAnnotation)
gsa, err := e.Kube.ServiceAccountAnnotation(ctx, e.Namespace(), "ate-api-server", workloadIdentityAnnotation)
if err != nil {
return cloudSQLSettings{}, err
}
@@ -220,7 +220,7 @@ func (e *Env) reconcileCloudSQLProxySidecar(ctx context.Context) error {
if s.Instance != "" {
log.Step("reconcile_cloudsql_proxy_sidecar (add)")
if s.GSA != "" {
if err := e.Kube.SetServiceAccountAnnotation(ctx, NamespaceAteSystem, "ate-api-server",
if err := e.Kube.SetServiceAccountAnnotation(ctx, e.Namespace(), "ate-api-server",
workloadIdentityAnnotation, s.GSA); err != nil {
return err
}
@@ -229,7 +229,7 @@ func (e *Env) reconcileCloudSQLProxySidecar(ctx context.Context) error {
if err != nil {
return fmt.Errorf("reading the Cloud SQL proxy sidecar patch: %w", err)
}
return e.Kube.PatchDeployment(ctx, NamespaceAteSystem, "ate-api-server", patch)
return e.Kube.PatchDeployment(ctx, e.Namespace(), "ate-api-server", patch)
}
installed, err := e.cloudSQLProxyInstalled(ctx)
@@ -239,17 +239,17 @@ func (e *Env) reconcileCloudSQLProxySidecar(ctx context.Context) error {
log.Step("reconcile_cloudsql_proxy_sidecar (remove)")
const removePatch = `{"spec":{"template":{"spec":{"initContainers":[{"name":"` +
cloudSQLProxyContainer + `","$patch":"delete"}]}}}}`
if err := e.Kube.PatchDeployment(ctx, NamespaceAteSystem, "ate-api-server", []byte(removePatch)); err != nil {
if err := e.Kube.PatchDeployment(ctx, e.Namespace(), "ate-api-server", []byte(removePatch)); err != nil {
return err
}
return e.Kube.SetServiceAccountAnnotation(ctx, NamespaceAteSystem, "ate-api-server",
return e.Kube.SetServiceAccountAnnotation(ctx, e.Namespace(), "ate-api-server",
workloadIdentityAnnotation, "")
}
// cloudSQLProxyInstalled reports whether ate-api-server currently runs the
// proxy sidecar.
func (e *Env) cloudSQLProxyInstalled(ctx context.Context) (bool, error) {
dep, err := e.Kube.GetDeployment(ctx, NamespaceAteSystem, "ate-api-server")
dep, err := e.Kube.GetDeployment(ctx, e.Namespace(), "ate-api-server")
if err != nil || dep == nil {
return false, err
}
+11 -11
View File
@@ -54,32 +54,32 @@ const caValidity = 365 * 24 * time.Hour
// to kubectl-ate for.
func (e *Env) CreateJWTAuthorityPoolSecret(ctx context.Context) error {
log.Step("create_jwt_authority_pool_secret")
return e.createJWTPool(ctx, NamespaceAteSystem, SecretActorIDJWTPool)
return e.createJWTPool(ctx, e.Namespace(), SecretActorIDJWTPool)
}
// CreateActorIDCAPoolSecret generates the actor-identity CA pool.
func (e *Env) CreateActorIDCAPoolSecret(ctx context.Context) error {
log.Step("create_actor_id_ca_pool_secret")
return e.createCAPool(ctx, NamespaceAteSystem, SecretActorIDCAPool)
return e.createCAPool(ctx, e.Namespace(), SecretActorIDCAPool)
}
// CreateEgressMITMCAPoolSecret generates the egress MITM CA pool.
func (e *Env) CreateEgressMITMCAPoolSecret(ctx context.Context) error {
log.Step("create_egress_mitm_ca_pool_secret")
exists, err := e.Kube.SecretExists(ctx, NamespaceAteSystem, SecretEgressMITMCAPool)
exists, err := e.Kube.SecretExists(ctx, e.Namespace(), SecretEgressMITMCAPool)
if err != nil {
return err
}
if exists {
log.Infof(" CA pool %s/%s already exists; keeping it", NamespaceAteSystem, SecretEgressMITMCAPool)
log.Infof(" CA pool %s/%s already exists; keeping it", e.Namespace(), SecretEgressMITMCAPool)
return nil
}
data, err := newCAPoolSecretData(poolKeyID, localca.KeyTypeECDSAP256)
if err != nil {
return fmt.Errorf("while generating the CA pool for %s/%s: %w", NamespaceAteSystem, SecretEgressMITMCAPool, err)
return fmt.Errorf("while generating the CA pool for %s/%s: %w", e.Namespace(), SecretEgressMITMCAPool, err)
}
return e.createPoolSecret(ctx, NamespaceAteSystem, SecretEgressMITMCAPool, corev1.SecretTypeTLS, data)
return e.createPoolSecret(ctx, e.Namespace(), SecretEgressMITMCAPool, corev1.SecretTypeTLS, data)
}
// EnsureEgressMITMCAPoolSecret creates the egress MITM CA pool secret if
@@ -89,7 +89,7 @@ func (e *Env) EnsureEgressMITMCAPoolSecret(ctx context.Context) error {
if !e.Cfg.ExperimentalUseSDSMint {
return nil
}
return e.ensureSecret(ctx, NamespaceAteSystem, SecretEgressMITMCAPool, e.CreateEgressMITMCAPoolSecret)
return e.ensureSecret(ctx, e.Namespace(), SecretEgressMITMCAPool, e.CreateEgressMITMCAPoolSecret)
}
// CreatePodCertificateControllerCAs generates the two signer pools the
@@ -113,11 +113,11 @@ func (e *Env) CreatePodCertificateControllerCAs(ctx context.Context) error {
// root.
func (e *Env) CreateActorIDCACertsSecret(ctx context.Context) error {
log.Step("create_actor_id_ca_certs_secret")
root, err := e.Kube.CAPoolRootPEM(ctx, NamespaceAteSystem, SecretActorIDCAPool)
root, err := e.Kube.CAPoolRootPEM(ctx, e.Namespace(), SecretActorIDCAPool)
if err != nil {
return fmt.Errorf("while building %s: %w", SecretActorIDCACerts, err)
}
return e.Kube.ApplySecret(ctx, NamespaceAteSystem, SecretActorIDCACerts, map[string]string{
return e.Kube.ApplySecret(ctx, e.Namespace(), SecretActorIDCACerts, map[string]string{
"ca.crt": string(root),
})
}
@@ -126,7 +126,7 @@ func (e *Env) CreateActorIDCACertsSecret(ctx context.Context) error {
// authentication config, pointing it at the cluster's service account issuer.
func (e *Env) CreateAPIAuthenticationConfig(ctx context.Context) error {
log.Step("create_api_authentication_config")
if err := e.Kube.EnsureNamespace(ctx, NamespaceAteSystem); err != nil {
if err := e.Kube.EnsureNamespace(ctx, e.Namespace()); err != nil {
return err
}
@@ -137,7 +137,7 @@ func (e *Env) CreateAPIAuthenticationConfig(ctx context.Context) error {
for _, line := range strings.Split(authnConfig, "\n") {
log.Infof(" | %s", line)
}
return e.Kube.ApplyConfigMap(ctx, NamespaceAteSystem, ConfigMapAPIAuthn, map[string]string{
return e.Kube.ApplyConfigMap(ctx, e.Namespace(), ConfigMapAPIAuthn, map[string]string{
"authentication.yaml": authnConfig,
})
}
+1 -1
View File
@@ -45,7 +45,7 @@ func (e *Env) DeleteAteSystem(ctx context.Context) error {
}
// atelet DaemonSet names carry a version suffix.
if err := e.Kube.Typed.AppsV1().DaemonSets(NamespaceAteSystem).DeleteCollection(ctx,
if err := e.Kube.Typed.AppsV1().DaemonSets(e.Namespace()).DeleteCollection(ctx,
metav1.DeleteOptions{}, metav1.ListOptions{LabelSelector: "app=atelet"}); err != nil {
return fmt.Errorf("while deleting atelet daemonsets: %w", err)
}
+10 -5
View File
@@ -61,6 +61,11 @@ func (o DeployOptions) Validate() error {
func (e *Env) DeployAteSystem(ctx context.Context, opts DeployOptions) error {
log.Step("deploy_ate_system")
// This step applies the checked-in manifests, so it refuses a relocated
// namespace before creating anything.
if err := e.RequireCanonicalNamespace("deploy ate-system"); err != nil {
return err
}
// Fail fast on an unusable build version before touching the cluster.
if _, _, err := e.SubstrateVersion(); err != nil {
return err
@@ -189,7 +194,7 @@ func (e *Env) DeployAteSystem(ctx context.Context, opts DeployOptions) error {
rollout{kube.KindDaemonSet, ateletName},
)
for _, w := range waits {
if err := e.Kube.RolloutStatus(ctx, w.kind, NamespaceAteSystem, w.name, e.Cfg.RolloutTimeout); err != nil {
if err := e.Kube.RolloutStatus(ctx, w.kind, e.Namespace(), w.name, e.Cfg.RolloutTimeout); err != nil {
return err
}
}
@@ -265,7 +270,7 @@ func (e *Env) DeployAteAPIServer(ctx context.Context) error {
if err := e.reconcileCloudSQLProxySidecar(ctx); err != nil {
return err
}
return e.Kube.RolloutStatus(ctx, kube.KindDeployment, NamespaceAteSystem, "ate-api-server", e.Cfg.RolloutTimeout)
return e.Kube.RolloutStatus(ctx, kube.KindDeployment, e.Namespace(), "ate-api-server", e.Cfg.RolloutTimeout)
}
// DeployAteController redeploys only ate-controller.
@@ -284,7 +289,7 @@ func (e *Env) DeployAteController(ctx context.Context) error {
if err := e.ResolveAndApply(ctx, e.Cfg.Manifest("ate-controller.yaml")); err != nil {
return err
}
return e.Kube.RolloutStatus(ctx, kube.KindDeployment, NamespaceAteSystem, "ate-controller", e.Cfg.RolloutTimeout)
return e.Kube.RolloutStatus(ctx, kube.KindDeployment, e.Namespace(), "ate-controller", e.Cfg.RolloutTimeout)
}
// DeployAtelet redeploys only the atelet DaemonSet.
@@ -329,7 +334,7 @@ func (e *Env) DeployAtelet(ctx context.Context) error {
if err != nil {
return err
}
return e.Kube.RolloutStatus(ctx, kube.KindDaemonSet, NamespaceAteSystem, ateletName, e.Cfg.RolloutTimeout)
return e.Kube.RolloutStatus(ctx, kube.KindDaemonSet, e.Namespace(), ateletName, e.Cfg.RolloutTimeout)
}
// DeployAtenet redeploys the atenet dataplane: router and egress.
@@ -364,7 +369,7 @@ func (e *Env) DeployAtenet(ctx context.Context) error {
}
for _, name := range []string{"atenet-router", "atenet-egress"} {
if err := e.Kube.RolloutStatus(ctx, kube.KindDeployment, NamespaceAteSystem, name, e.Cfg.RolloutTimeout); err != nil {
if err := e.Kube.RolloutStatus(ctx, kube.KindDeployment, e.Namespace(), name, e.Cfg.RolloutTimeout); err != nil {
return err
}
}
+39 -5
View File
@@ -46,6 +46,10 @@ const (
// Well-known namespaces.
const (
// NamespaceAteSystem is the canonical control-plane namespace. It is the
// default for Config.Namespace and the only value the checked-in manifests
// under manifests/ate-install/ carry; steps address the installed control
// plane through Env.Namespace rather than this constant.
NamespaceAteSystem = "ate-system"
NamespacePodCert = "podcertificate-controller-system"
)
@@ -75,6 +79,27 @@ type Env struct {
substrateVersionSuffix string
}
// Namespace is the namespace the control plane is installed into. It is
// Config.Namespace, which defaults to NamespaceAteSystem.
func (e *Env) Namespace() string {
if e.Cfg != nil && e.Cfg.Namespace != "" {
return e.Cfg.Namespace
}
return NamespaceAteSystem
}
// RequireCanonicalNamespace refuses a relocated install for the steps that
// apply the checked-in manifests. Those manifests name ate-system literally,
// so proceeding would put the workloads there while this tool created their
// secrets and ConfigMaps somewhere else — an install that comes up far enough
// to look healthy and then fails on a missing envFrom source.
func (e *Env) RequireCanonicalNamespace(step string) error {
if ns := e.Namespace(); ns != NamespaceAteSystem {
return fmt.Errorf("%s cannot be used with ATE_NAMESPACE=%s: manifests/ate-install/ names %s literally; install into another namespace with a deployment that renders them, and use ate-setup only for the create steps", step, ns, NamespaceAteSystem)
}
return nil
}
// NewEnv connects to the cluster described by cfg.
func NewEnv(cfg *config.Config) (*Env, error) {
client, err := kube.New(cfg.Kubeconfig, cfg.Context)
@@ -179,14 +204,23 @@ func (e *Env) KustomizeResolve(ctx context.Context, overlay string) ([]byte, err
return e.ResolveManifestBytes(ctx, built)
}
// EnsureAteSystemNamespace applies the ate-system namespace manifest and waits
// for it to go Active. Every deploy path starts here so that RBAC, ConfigMaps,
// and workloads have somewhere to land.
// EnsureAteSystemNamespace creates the control-plane namespace and waits for it
// to go Active. Every deploy path starts here so that RBAC, ConfigMaps, and
// workloads have somewhere to land.
//
// The canonical namespace comes from the checked-in manifest, which carries
// labels of its own and stays the source of truth for it. Any other namespace
// is created plainly, because that manifest names ate-system literally.
func (e *Env) EnsureAteSystemNamespace(ctx context.Context) error {
if err := e.Kube.ApplyPath(ctx, e.Cfg.Manifest("ate-system-namespace.yaml")); err != nil {
ns := e.Namespace()
if ns == NamespaceAteSystem {
if err := e.Kube.ApplyPath(ctx, e.Cfg.Manifest("ate-system-namespace.yaml")); err != nil {
return err
}
} else if err := e.Kube.EnsureNamespace(ctx, ns); err != nil {
return err
}
return e.Kube.WaitNamespaceActive(ctx, NamespaceAteSystem, NamespaceTimeout)
return e.Kube.WaitNamespaceActive(ctx, ns, NamespaceTimeout)
}
// RequireKind fails a step that only makes sense on a local Kind cluster. The
@@ -0,0 +1,67 @@
// 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 steps
import (
"strings"
"testing"
"github.com/agent-substrate/substrate/cmd/ate-setup/internal/config"
)
func TestEnvNamespace(t *testing.T) {
tests := []struct {
name string
cfg *config.Config
want string
}{
{"falls back to the canonical namespace when unset", &config.Config{}, NamespaceAteSystem},
{"uses the configured namespace", &config.Config{Namespace: "substrate-dev"}, "substrate-dev"},
{"tolerates a nil config", nil, NamespaceAteSystem},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
e := &Env{Cfg: tt.cfg}
if got := e.Namespace(); got != tt.want {
t.Errorf("Namespace() = %q, want %q", got, tt.want)
}
})
}
}
// The checked-in manifests under manifests/ate-install/ name ate-system
// literally, so the steps that apply them must refuse any other namespace
// rather than scatter the install across two.
func TestRequireCanonicalNamespace(t *testing.T) {
t.Run("permits the canonical namespace", func(t *testing.T) {
e := &Env{Cfg: &config.Config{Namespace: NamespaceAteSystem}}
if err := e.RequireCanonicalNamespace("deploy ate-system"); err != nil {
t.Errorf("RequireCanonicalNamespace() = %v, want nil", err)
}
})
t.Run("refuses a relocated namespace and names both the step and the value", func(t *testing.T) {
e := &Env{Cfg: &config.Config{Namespace: "substrate-dev"}}
err := e.RequireCanonicalNamespace("deploy ate-system")
if err == nil {
t.Fatal("RequireCanonicalNamespace() = nil, want an error")
}
for _, want := range []string{"deploy ate-system", "substrate-dev", NamespaceAteSystem} {
if !strings.Contains(err.Error(), want) {
t.Errorf("error %q does not mention %q", err, want)
}
}
})
}
+7 -7
View File
@@ -315,7 +315,7 @@ func (e *Env) applyAtenetEgress(ctx context.Context) error {
return err
}
running, err := e.Kube.DeploymentExists(ctx, NamespaceAteSystem, "atenet-egress")
running, err := e.Kube.DeploymentExists(ctx, e.Namespace(), "atenet-egress")
if err != nil {
return err
}
@@ -325,7 +325,7 @@ func (e *Env) applyAtenetEgress(ctx context.Context) error {
}
if running && (e.Cfg.AdditionalEgressExtprocService != "" || e.Cfg.ExperimentalEgressCredentialInjection) {
if err := e.Kube.RolloutRestartDeployment(ctx, NamespaceAteSystem, "atenet-egress", time.Now()); err != nil {
if err := e.Kube.RolloutRestartDeployment(ctx, e.Namespace(), "atenet-egress", time.Now()); err != nil {
return err
}
}
@@ -382,7 +382,7 @@ func (e *Env) applyOtelEndpointOverride(ctx context.Context) error {
return nil
}
cm, err := e.Kube.GetConfigMap(ctx, NamespaceAteSystem, otelConfigMap)
cm, err := e.Kube.GetConfigMap(ctx, e.Namespace(), otelConfigMap)
if err != nil {
return err
}
@@ -391,25 +391,25 @@ func (e *Env) applyOtelEndpointOverride(ctx context.Context) error {
}
log.Infof("Overriding %s with %s", otelEndpointKey, endpoint)
if err := e.Kube.MergePatchConfigMap(ctx, NamespaceAteSystem, otelConfigMap,
if err := e.Kube.MergePatchConfigMap(ctx, e.Namespace(), otelConfigMap,
map[string]string{otelEndpointKey: endpoint}); err != nil {
return err
}
now := time.Now()
for _, name := range otelOverrideDeployments {
if err := e.Kube.RolloutRestartDeployment(ctx, NamespaceAteSystem, name, now); err != nil {
if err := e.Kube.RolloutRestartDeployment(ctx, e.Namespace(), name, now); err != nil {
return err
}
}
// atelet DaemonSet names carry a version suffix; restart whichever
// versions are installed.
daemonSets, err := e.Kube.DaemonSetNames(ctx, NamespaceAteSystem, "app=atelet")
daemonSets, err := e.Kube.DaemonSetNames(ctx, e.Namespace(), "app=atelet")
if err != nil {
return err
}
for _, name := range daemonSets {
if err := e.Kube.RolloutRestart(ctx, NamespaceAteSystem, name, now); err != nil {
if err := e.Kube.RolloutRestart(ctx, e.Namespace(), name, now); err != nil {
return err
}
}
+1 -1
View File
@@ -107,5 +107,5 @@ func (e *Env) DeployPostgres(ctx context.Context) error {
if err := e.applyPostgresManifest(ctx); err != nil {
return err
}
return e.Kube.RolloutStatus(ctx, kube.KindStatefulSet, NamespaceAteSystem, "postgres", e.Cfg.RolloutTimeout)
return e.Kube.RolloutStatus(ctx, kube.KindStatefulSet, e.Namespace(), "postgres", e.Cfg.RolloutTimeout)
}
+4 -4
View File
@@ -33,14 +33,14 @@ var trustBundleNames = []string{
func (e *Env) EnsureAPIServerPrerequisites(ctx context.Context) error {
log.Step("ensure_apiserver_prerequisites")
if err := e.ensureSecret(ctx, NamespaceAteSystem, SecretActorIDJWTPool, e.CreateJWTAuthorityPoolSecret); err != nil {
if err := e.ensureSecret(ctx, e.Namespace(), SecretActorIDJWTPool, e.CreateJWTAuthorityPoolSecret); err != nil {
return err
}
if err := e.ensureSecret(ctx, NamespaceAteSystem, SecretActorIDCAPool, e.CreateActorIDCAPoolSecret); err != nil {
if err := e.ensureSecret(ctx, e.Namespace(), SecretActorIDCAPool, e.CreateActorIDCAPoolSecret); err != nil {
return err
}
// Derived from actor-id-ca-pool above, so it must come after it.
if err := e.ensureSecret(ctx, NamespaceAteSystem, SecretActorIDCACerts, e.CreateActorIDCACertsSecret); err != nil {
if err := e.ensureSecret(ctx, e.Namespace(), SecretActorIDCACerts, e.CreateActorIDCACertsSecret); err != nil {
return err
}
if err := e.ensureSecret(ctx, NamespacePodCert, SecretServiceDNSCA, e.CreatePodCertificateControllerCAs); err != nil {
@@ -52,7 +52,7 @@ func (e *Env) EnsureAPIServerPrerequisites(ctx context.Context) error {
return err
}
exists, err := e.Kube.ConfigMapExists(ctx, NamespaceAteSystem, ConfigMapAPIAuthn)
exists, err := e.Kube.ConfigMapExists(ctx, e.Namespace(), ConfigMapAPIAuthn)
if err != nil {
return err
}
+1 -1
View File
@@ -124,7 +124,7 @@ func WaitActorTemplateGolden(ctx context.Context, client *ateclient.Client, ref
// atespaces. A cluster without a reachable
// ate-api-server is not an error -- there is nothing to clean up.
func (e *Env) DeleteSubstrateDemo(ctx context.Context, refs []resources.ActorTemplateRef, atespaces []string) error {
present, err := e.Kube.DeploymentExists(ctx, NamespaceAteSystem, "ate-api-server")
present, err := e.Kube.DeploymentExists(ctx, e.Namespace(), "ate-api-server")
if err != nil {
return err
}
+2 -2
View File
@@ -92,14 +92,14 @@ func (e *Env) SubstituteVersion(manifest []byte) ([]byte, error) {
// RestartAteletDaemonSets pod-restarts every atelet DaemonSet by label.
func (e *Env) RestartAteletDaemonSets(ctx context.Context) error {
list, err := e.Kube.Typed.AppsV1().DaemonSets(NamespaceAteSystem).List(ctx, metav1.ListOptions{
list, err := e.Kube.Typed.AppsV1().DaemonSets(e.Namespace()).List(ctx, metav1.ListOptions{
LabelSelector: "app=atelet",
})
if err != nil {
return fmt.Errorf("while listing atelet daemonsets: %w", err)
}
for _, ds := range list.Items {
if err := e.Kube.RolloutRestart(ctx, NamespaceAteSystem, ds.Name, time.Now()); err != nil {
if err := e.Kube.RolloutRestart(ctx, e.Namespace(), ds.Name, time.Now()); err != nil {
return err
}
}
+5 -22
View File
@@ -19,8 +19,6 @@ package ateletauth
import (
"context"
"log/slog"
"net/url"
"path"
"github.com/agent-substrate/substrate/internal/substratex509"
"google.golang.org/grpc/codes"
@@ -29,32 +27,21 @@ import (
"google.golang.org/grpc/status"
)
// The SPIFFE identity that atelet client certs carry, as minted by the
// podidentity signer (cmd/podcertcontroller/internal/podidentitysigner).
//
// These mirror the constants the atelet dialer verifies against in
// cmd/ateapi/internal/controlapi/dialer.go, duplicated rather than imported so
// that this package does not depend on controlapi for three strings.
const (
TrustDomain = "cluster.local"
Namespace = "ate-system"
ServiceAccount = "atelet"
)
// Caller is the verified identity of an atelet.
type Caller struct {
PodName string
NodeName string
}
// Authenticate verifies that the RPC arrived over mTLS from an atelet, and
// returns the identity that atelet's certificate asserts.
// Authenticate verifies that the RPC arrived over mTLS from an atelet
// presenting ateletSPIFFEID, and returns the identity that atelet's
// certificate asserts.
//
// The certificate chain is already verified by the TLS layer against the
// pod-identity CA (see buildServerCreds in cmd/ateapi/main.go), so the
// extensions read here are trustworthy: only the pod-identity signer can mint
// a certificate carrying a given pod's node name.
func Authenticate(ctx context.Context) (*Caller, error) {
func Authenticate(ctx context.Context, ateletSPIFFEID string) (*Caller, error) {
p, ok := peer.FromContext(ctx)
if !ok {
return nil, status.Errorf(codes.Unauthenticated, "no peer transport information found")
@@ -73,11 +60,7 @@ func Authenticate(ctx context.Context) (*Caller, error) {
// Only atelet may call these RPCs. Everything else with a valid
// pod-identity certificate — including the actor workloads themselves — is
// rejected here.
expected := (&url.URL{
Scheme: "spiffe",
Host: TrustDomain,
Path: path.Join("ns", Namespace, "sa", ServiceAccount),
}).String()
expected := ateletSPIFFEID
if len(leaf.URIs) == 0 || leaf.URIs[0].String() != expected {
slog.WarnContext(ctx, "Denied: caller is not atelet",
slog.Any("uris", leaf.URIs), slog.String("expected", expected))
@@ -0,0 +1,53 @@
// 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 ateletauth_test
import (
"testing"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/ateletauth"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/ateletauth/ateletauthtest"
"github.com/agent-substrate/substrate/internal/installdefaults"
)
// TestAuthenticateHonorsConfiguredIdentity checks that the identity
// Authenticate is given is the one it accepts, and that the canonical
// identity is rejected when the install lives elsewhere.
//
// The authorization table test in actoridentity covers a caller with the
// wrong identity, but it always expects the default identity, so it holds
// against a hardcoded "ate-system"/"atelet" too. Only the relocated case
// distinguishes "reads its configuration" from "happens to agree with the
// constant".
func TestAuthenticateHonorsConfiguredIdentity(t *testing.T) {
const relocated = "substrate-test"
const node = "test-node"
relocatedID := installdefaults.AteletSPIFFEID(relocated)
t.Run("accepts atelet with the configured identity", func(t *testing.T) {
ctx := ateletauthtest.ContextWith(ateletauthtest.CertIn(t, relocated, node))
if _, err := ateletauth.Authenticate(ctx, relocatedID); err != nil {
t.Errorf("Authenticate() = %v, want success for an atelet in %q", err, relocated)
}
})
t.Run("rejects atelet with the canonical identity", func(t *testing.T) {
ctx := ateletauthtest.ContextWith(ateletauthtest.CertIn(t, installdefaults.SystemNamespace, node))
if _, err := ateletauth.Authenticate(ctx, relocatedID); err == nil {
t.Errorf("Authenticate() accepted an atelet from %q, want rejection when configured for %q",
installdefaults.SystemNamespace, relocated)
}
})
}
@@ -29,7 +29,7 @@ import (
"testing"
"time"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/ateletauth"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/substratex509"
"google.golang.org/grpc/credentials"
"google.golang.org/grpc/peer"
@@ -53,7 +53,7 @@ func Cert(t *testing.T, spiffePath string, podIdentity *substratex509.PodIdentit
NotAfter: time.Now().Add(time.Hour),
}
if spiffePath != "" {
template.URIs = []*url.URL{{Scheme: "spiffe", Host: ateletauth.TrustDomain, Path: spiffePath}}
template.URIs = []*url.URL{{Scheme: "spiffe", Host: installdefaults.AteletTrustDomain, Path: spiffePath}}
}
if podIdentity != nil {
if err := substratex509.AddPodIdentityToCertificate(podIdentity, template); err != nil {
@@ -72,11 +72,18 @@ func Cert(t *testing.T, spiffePath string, podIdentity *substratex509.PodIdentit
return cert
}
// PodIdentityOn returns a well-formed atelet PodIdentity pinned to nodeName.
// PodIdentityOn returns a well-formed atelet PodIdentity pinned to nodeName,
// for an atelet running in the default install namespace.
func PodIdentityOn(nodeName string) *substratex509.PodIdentity {
return PodIdentityIn(installdefaults.SystemNamespace, nodeName)
}
// PodIdentityIn returns a well-formed atelet PodIdentity pinned to nodeName,
// for an atelet running in namespace.
func PodIdentityIn(namespace, nodeName string) *substratex509.PodIdentity {
return &substratex509.PodIdentity{
Namespace: ateletauth.Namespace,
ServiceAccountName: ateletauth.ServiceAccount,
Namespace: namespace,
ServiceAccountName: installdefaults.AteletServiceAccount,
ServiceAccountUID: "sa-uid",
PodName: "atelet-xyz",
PodUID: "pod-uid",
@@ -85,10 +92,17 @@ func PodIdentityOn(nodeName string) *substratex509.PodIdentity {
}
}
// CertOn returns the certificate of the atelet running on nodeName.
// CertOn returns the certificate of the atelet running on nodeName in the
// default install namespace.
func CertOn(t *testing.T, nodeName string) *x509.Certificate {
t.Helper()
return Cert(t, path.Join("ns", ateletauth.Namespace, "sa", ateletauth.ServiceAccount), PodIdentityOn(nodeName))
return CertIn(t, installdefaults.SystemNamespace, nodeName)
}
// CertIn returns the certificate of the atelet running on nodeName in namespace.
func CertIn(t *testing.T, namespace, nodeName string) *x509.Certificate {
t.Helper()
return Cert(t, path.Join("ns", namespace, "sa", installdefaults.AteletServiceAccount), PodIdentityIn(namespace, nodeName))
}
// ContextWith injects cert as the transport-authenticated peer certificate. A
+11 -16
View File
@@ -25,6 +25,7 @@ import (
"github.com/agent-substrate/substrate/internal/atelet"
"github.com/agent-substrate/substrate/internal/credbundle"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/substratex509"
"github.com/spiffe/go-spiffe/v2/bundle/x509bundle"
"github.com/spiffe/go-spiffe/v2/spiffeid"
@@ -42,14 +43,6 @@ import (
// Retryable.
var ErrNoAteletOnNode = errors.New("no atelet pod found on node")
// The SPIFFE identity that atelet serving certs carry, as minted by the
// podidentity signer (cmd/podcertcontroller/internal/podidentitysigner).
// The namespace part is ateletNamespace, declared in informer.go.
const (
trustDomainName = "cluster.local"
ateletSA = "atelet"
)
// AteletDialer handles gRPC connections to Atelet pods.
type AteletDialer struct {
ateletIndexer cache.Indexer
@@ -71,13 +64,15 @@ func WithDialCredentials(build func(expectedPodUID string) (credentials.Transpor
}
// NewAteletDialer creates a new AteletDialer. clientBundlePath and serverCAPath
// are used to build the per-atelet mTLS credentials used for every atelet connection.
func NewAteletDialer(ateletIndexer cache.Indexer, clientBundlePath, serverCAPath string, opts ...DialerOption) *AteletDialer {
// are used to build the per-atelet mTLS credentials used for every atelet
// connection, and ateletSPIFFEID is the identity those credentials expect on
// the atelet serving cert.
func NewAteletDialer(ateletIndexer cache.Indexer, ateletSPIFFEID, clientBundlePath, serverCAPath string, opts ...DialerOption) *AteletDialer {
d := &AteletDialer{
ateletIndexer: ateletIndexer,
ateletConns: newAteletConnCache(1024),
dialCredentials: func(expectedPodUID string) (credentials.TransportCredentials, error) {
tlsConfig, err := buildTLSConfig(clientBundlePath, serverCAPath, expectedPodUID)
tlsConfig, err := buildTLSConfig(ateletSPIFFEID, clientBundlePath, serverCAPath, expectedPodUID)
if err != nil {
return nil, err
}
@@ -155,18 +150,18 @@ func (d *AteletDialer) DialForAteletOnNode(nodeName string) (*grpc.ClientConn, e
return ateletConn, nil
}
func buildTLSConfig(clientBundlePath, serverCAPath, expectedPodUID string) (*tls.Config, error) {
trustDomain, err := spiffeid.TrustDomainFromString(trustDomainName)
func buildTLSConfig(ateletSPIFFEID, clientBundlePath, serverCAPath, expectedPodUID string) (*tls.Config, error) {
trustDomain, err := spiffeid.TrustDomainFromString(installdefaults.AteletTrustDomain)
if err != nil {
return nil, fmt.Errorf("while parsing trust domain %q: %w", trustDomainName, err)
return nil, fmt.Errorf("while parsing trust domain %q: %w", installdefaults.AteletTrustDomain, err)
}
bundle, err := x509bundle.Load(trustDomain, serverCAPath)
if err != nil {
return nil, fmt.Errorf("while loading CA bundle from %s: %w", serverCAPath, err)
}
expectedID, err := spiffeid.FromSegments(trustDomain, "ns", ateletNamespace, "sa", ateletSA)
expectedID, err := spiffeid.FromString(ateletSPIFFEID)
if err != nil {
return nil, fmt.Errorf("while building expected atelet SPIFFE ID: %w", err)
return nil, fmt.Errorf("while parsing expected atelet SPIFFE ID %q: %w", ateletSPIFFEID, err)
}
verify, err := verifyAteletServerCert(bundle, expectedID, expectedPodUID)
@@ -27,6 +27,7 @@ import (
"testing"
"time"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/substratex509"
"github.com/spiffe/go-spiffe/v2/bundle/x509bundle"
"github.com/spiffe/go-spiffe/v2/spiffeid"
@@ -138,7 +139,7 @@ func makeLeafCert(t *testing.T, ca *x509.Certificate, caKey *ecdsa.PrivateKey, o
// with insecure test credentials.
func dialerWithAtelets(t *testing.T, pods ...*corev1.Pod) *AteletDialer {
t.Helper()
return NewAteletDialer(newTestAteletIndexer(t, pods...), "", "",
return NewAteletDialer(newTestAteletIndexer(t, pods...), installdefaults.SystemNamespace, "", "",
WithDialCredentials(func(string) (credentials.TransportCredentials, error) {
return insecure.NewCredentials(), nil
}))
@@ -170,7 +171,7 @@ func TestDialForAteletOnNodeTarget(t *testing.T) {
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
ateletPod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{Namespace: ateletNamespace, Name: "atelet-abc", UID: "atelet-uid"},
ObjectMeta: metav1.ObjectMeta{Namespace: installdefaults.SystemNamespace, Name: "atelet-abc", UID: "atelet-uid"},
Spec: corev1.PodSpec{NodeName: "node-1"},
Status: corev1.PodStatus{PodIPs: []corev1.PodIP{{IP: tc.ateletIP}}},
}
@@ -191,7 +192,7 @@ func TestDialForAteletOnNodeTarget(t *testing.T) {
func TestDialForAteletOnNodeNoIPs(t *testing.T) {
ateletPod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{Namespace: ateletNamespace, Name: "atelet-abc", UID: "atelet-uid"},
ObjectMeta: metav1.ObjectMeta{Namespace: installdefaults.SystemNamespace, Name: "atelet-abc", UID: "atelet-uid"},
Spec: corev1.PodSpec{NodeName: "node-1"},
}
d := dialerWithAtelets(t, ateletPod)
@@ -316,7 +317,7 @@ func TestDialForAteletOnNode(t *testing.T) {
}
t.Run("no atelet on node", func(t *testing.T) {
d := NewAteletDialer(newTestAteletIndexer(t), "", "")
d := NewAteletDialer(newTestAteletIndexer(t), installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "", "")
if _, err := d.DialForAteletOnNode("node1"); !errors.Is(err, ErrNoAteletOnNode) {
t.Fatalf("DialForAteletOnNode = %v, want ErrNoAteletOnNode", err)
}
@@ -326,7 +327,7 @@ func TestDialForAteletOnNode(t *testing.T) {
d := NewAteletDialer(newTestAteletIndexer(t,
ateletPod("atelet-1", "uid-1", "node1", "10.0.0.1"),
ateletPod("atelet-2", "uid-2", "node1", "10.0.0.2"),
), "", "")
), installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "", "")
_, err := d.DialForAteletOnNode("node1")
if err == nil || errors.Is(err, ErrNoAteletOnNode) {
t.Fatalf("DialForAteletOnNode = %v, want a non-ErrNoAteletOnNode error", err)
@@ -336,7 +337,7 @@ func TestDialForAteletOnNode(t *testing.T) {
t.Run("dials and caches the node's atelet", func(t *testing.T) {
d := NewAteletDialer(newTestAteletIndexer(t,
ateletPod("atelet-1", "uid-1", "node1", "10.0.0.1"),
), "", "")
), installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "", "")
var credsUID string
d.dialCredentials = func(expectedPodUID string) (credentials.TransportCredentials, error) {
credsUID = expectedPodUID
@@ -363,7 +364,7 @@ func TestDialForAteletOnNode(t *testing.T) {
d := NewAteletDialer(newTestAteletIndexer(t,
ateletPod("atelet-1", "uid-1", "node1", "10.0.0.1"),
ateletPod("atelet-2", "uid-2", "node2", "10.0.0.2"),
), "", "", WithDialCredentials(func(string) (credentials.TransportCredentials, error) {
), installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "", "", WithDialCredentials(func(string) (credentials.TransportCredentials, error) {
return insecure.NewCredentials(), nil
}))
d.ateletConns = newAteletConnCache(1)
@@ -27,6 +27,7 @@ import (
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store/storetest"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/workercache"
"github.com/agent-substrate/substrate/internal/ateinterceptors"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/localca"
"github.com/agent-substrate/substrate/internal/localjwtauthority"
"github.com/agent-substrate/substrate/internal/objectstore/objectstoretest"
@@ -131,7 +132,7 @@ func setupTestWithVolumePlugins(t *testing.T, ns string, plugins map[string]volu
}
// 3. Initialize Informers
ateletFactory, ateletInformer := controlapi.AteletInformer(k8sClient)
ateletFactory, ateletInformer := controlapi.AteletInformer(k8sClient, installdefaults.SystemNamespace)
scFactory := informers.NewSharedInformerFactory(k8sClient, 0)
scLister := scFactory.Storage().V1().StorageClasses().Lister()
@@ -160,7 +161,7 @@ func setupTestWithVolumePlugins(t *testing.T, ns string, plugins map[string]volu
// Dial the fake atelet over insecure transport instead of per-atelet mTLS,
// so DialForAteletOnNode's real lookup/dial/cache path is exercised under test.
dialer := controlapi.NewAteletDialer(ateletInformer.GetIndexer(), "", "",
dialer := controlapi.NewAteletDialer(ateletInformer.GetIndexer(), installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "", "",
controlapi.WithDialCredentials(func(_ string) (credentials.TransportCredentials, error) {
return insecure.NewCredentials(), nil
}))
+4 -4
View File
@@ -23,12 +23,12 @@ import (
)
const (
ateletNamespace = "ate-system"
byNode = "by-node"
byNode = "by-node"
)
// AteletInformer creates a SharedInformerFactory and SharedIndexInformer for Atelet pods.
func AteletInformer(kc kubernetes.Interface) (informers.SharedInformerFactory, cache.SharedIndexInformer) {
// AteletInformer creates a SharedInformerFactory and SharedIndexInformer for
// Atelet pods in the given namespace.
func AteletInformer(kc kubernetes.Interface, ateletNamespace string) (informers.SharedInformerFactory, cache.SharedIndexInformer) {
factory := informers.NewSharedInformerFactoryWithOptions(kc, 0,
informers.WithNamespace(ateletNamespace),
informers.WithTweakListOptions(func(options *metav1.ListOptions) {
@@ -27,6 +27,7 @@ import (
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store/storetest"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/workercache"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/proto/ateletpb"
"github.com/agent-substrate/substrate/internal/resources"
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
@@ -1265,10 +1266,10 @@ func newWireCaptureWorkflow(t *testing.T, persistence store.Interface) (*ActorWo
})
ateletPod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{Namespace: ateletNamespace, Name: "atelet-1", UID: "atelet-uid"},
ObjectMeta: metav1.ObjectMeta{Namespace: installdefaults.SystemNamespace, Name: "atelet-1", UID: "atelet-uid"},
Spec: corev1.PodSpec{NodeName: "node-1"},
}
dialer := NewAteletDialer(newTestAteletIndexer(t, ateletPod), "", "")
dialer := NewAteletDialer(newTestAteletIndexer(t, ateletPod), installdefaults.SystemNamespace, "", "")
dialer.ateletConns.Add("atelet-uid", conn)
lister := sandboxConfigListerFor(t, []*atev1alpha1.SandboxConfig{{
@@ -21,6 +21,7 @@ import (
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store/storetest"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"google.golang.org/grpc/codes"
@@ -206,7 +207,7 @@ func newDanglingDialer() *AteletDialer {
empty := cache.NewIndexer(cache.MetaNamespaceKeyFunc, cache.Indexers{
byNode: func(obj any) ([]string, error) { return nil, nil },
})
return NewAteletDialer(empty, "", "")
return NewAteletDialer(empty, installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "", "")
}
func TestEnsureAteletSuspended_DialFailureLeavesActorRetryable(t *testing.T) {
@@ -41,18 +41,21 @@ type Server struct {
// store is where a Worker's reported capacity is recorded, and the
// authoritative state the report is authorized against.
store store.Interface
// ateletSPIFFEID is the identity the calling atelet must present.
ateletSPIFFEID string
}
var _ ateapipb.WorkerServiceServer = (*Server)(nil)
func New(store store.Interface) *Server {
return &Server{store: store}
func New(store store.Interface, ateletSPIFFEID string) *Server {
return &Server{store: store, ateletSPIFFEID: ateletSPIFFEID}
}
// SetWorkerCapacity records a Worker's reported capacity. As with MintCert,
// the caller must be an atelet running on the Worker's node.
func (s *Server) SetWorkerCapacity(ctx context.Context, req *ateapipb.SetWorkerCapacityRequest) (*ateapipb.SetWorkerCapacityResponse, error) {
caller, err := ateletauth.Authenticate(ctx)
caller, err := ateletauth.Authenticate(ctx, s.ateletSPIFFEID)
if err != nil {
return nil, err
}
@@ -24,6 +24,7 @@ import (
"github.com/agent-substrate/substrate/cmd/ateapi/internal/ateletauth/ateletauthtest"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store/storetest"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"google.golang.org/grpc/codes"
@@ -70,7 +71,7 @@ func setRequest(actors int32) *ateapipb.SetWorkerCapacityRequest {
func TestSetWorkerCapacity(t *testing.T) {
st, cleanup := storetest.SetupTestStore(t)
defer cleanup()
s := New(st)
s := New(st, installdefaults.SPIFFEID(installdefaults.SystemNamespace, installdefaults.AteletServiceAccount))
seedReportedWorker(t, st, capNode, &ateapipb.WorkerResources{Actors: 1, Resources: resources.CPUMemory(2000, 0)})
got, err := s.SetWorkerCapacity(ateletauthtest.ContextWith(ateletauthtest.CertOn(t, capNode)), setRequest(4094))
@@ -94,7 +95,7 @@ func TestSetWorkerCapacity(t *testing.T) {
func TestSetWorkerCapacity_OtherNodeIsNotFound(t *testing.T) {
st, cleanup := storetest.SetupTestStore(t)
defer cleanup()
s := New(st)
s := New(st, installdefaults.SPIFFEID(installdefaults.SystemNamespace, installdefaults.AteletServiceAccount))
seedReportedWorker(t, st, capNode, &ateapipb.WorkerResources{Actors: 1})
_, err := s.SetWorkerCapacity(ateletauthtest.ContextWith(ateletauthtest.CertOn(t, "some-other-node")), setRequest(4094))
@@ -118,7 +119,7 @@ func TestSetWorkerCapacity_OtherNodeIsNotFound(t *testing.T) {
func TestSetWorkerCapacity_UnchangedDoesNotWrite(t *testing.T) {
st, cleanup := storetest.SetupTestStore(t)
defer cleanup()
s := New(st)
s := New(st, installdefaults.SPIFFEID(installdefaults.SystemNamespace, installdefaults.AteletServiceAccount))
seeded := seedReportedWorker(t, st, capNode, &ateapipb.WorkerResources{Actors: 4094})
for range 3 {
@@ -138,7 +139,7 @@ func TestSetWorkerCapacity_UnchangedDoesNotWrite(t *testing.T) {
func TestSetWorkerCapacity_Errors(t *testing.T) {
st, cleanup := storetest.SetupTestStore(t)
defer cleanup()
s := New(st)
s := New(st, installdefaults.SPIFFEID(installdefaults.SystemNamespace, installdefaults.AteletServiceAccount))
seedReportedWorker(t, st, capNode, &ateapipb.WorkerResources{Actors: 1})
authed := ateletauthtest.ContextWith(ateletauthtest.CertOn(t, capNode))
@@ -176,7 +177,7 @@ func TestSetWorkerCapacity_Errors(t *testing.T) {
func TestSetWorkerCapacity_RejectsNonsense(t *testing.T) {
st, cleanup := storetest.SetupTestStore(t)
defer cleanup()
s := New(st)
s := New(st, installdefaults.SPIFFEID(installdefaults.SystemNamespace, installdefaults.AteletServiceAccount))
seeded := seedReportedWorker(t, st, capNode, &ateapipb.WorkerResources{Actors: 4094})
authed := ateletauthtest.ContextWith(ateletauthtest.CertOn(t, capNode))
+17 -3
View File
@@ -38,6 +38,7 @@ import (
"github.com/agent-substrate/substrate/internal/ateinterceptors"
"github.com/agent-substrate/substrate/internal/authz"
"github.com/agent-substrate/substrate/internal/credbundle"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/localca"
"github.com/agent-substrate/substrate/internal/localjwtauthority"
"github.com/agent-substrate/substrate/internal/objectstore"
@@ -82,6 +83,7 @@ var (
actorIDCAPoolFile = pflag.String("actor-id-ca-pool", "", "The file that contains the CA pool for signing actor JWTs")
podIdentityCACerts = pflag.String("pod-identity-ca-certs", "", "The file that contains the pod-identity CA bundle, used both for verifying client certificates presented to the gRPC server and for verifying atelet serving certificates when dialing atelet. If empty, client-cert verification is disabled and atelet dials will fail.")
ateletClientCredBundle = pflag.String("atelet-client-cred-bundle", "", "Credential bundle presented as the client certificate when dialing atelet.")
ateletServiceAccount = pflag.String("atelet-service-account", installdefaults.AteletServiceAccount, "ServiceAccount atelet runs as. It is the service-account segment of the SPIFFE ID expected on atelet's certificate, so it has to match what the deployment actually creates; a deployment that prefixes resource names needs it set.")
drainDelay = pflag.Duration("drain-delay", 13*time.Second, "How long to keep accepting new work after SIGTERM, before starting the gRPC drain.")
drainTimeout = pflag.Duration("drain-timeout", 15*time.Second, "Deadline for the graceful gRPC drain on shutdown. In-flight RPCs still running past it are forcefully cancelled.")
@@ -196,7 +198,19 @@ func main() {
sandboxConfigLister := ateFactory.Api().V1alpha1().SandboxConfigs().Lister()
csiDriverConfigLister := ateFactory.Api().V1alpha1().CSIDriverConfigs().Lister()
ateletPodInformerFactory, ateletPodInformer := controlapi.AteletInformer(clientset)
// atelet shares ateapi's namespace in every supported deployment topology,
// so we read it from Kubernetes' downward API rather than expose a flag.
ateletNamespace := installdefaults.NamespaceFromPodEnv()
// An empty ServiceAccount would not fail here: path.Join drops the empty
// segment, yielding an identity that parses but matches nothing, so every
// atelet dial would be rejected with no hint at the cause.
if *ateletServiceAccount == "" {
serverboot.Fatal(ctx, "Invalid flags", fmt.Errorf("--atelet-service-account must not be empty"))
}
ateletSPIFFEID := installdefaults.SPIFFEID(ateletNamespace, *ateletServiceAccount)
slog.InfoContext(ctx, "Resolved atelet namespace", slog.String("atelet-namespace", ateletNamespace), slog.String("atelet-spiffe-id", ateletSPIFFEID))
ateletPodInformerFactory, ateletPodInformer := controlapi.AteletInformer(clientset, ateletNamespace)
scInformerFactory := informers.NewSharedInformerFactory(clientset, 0)
storageClassLister := scInformerFactory.Storage().V1().StorageClasses().Lister()
@@ -228,7 +242,7 @@ func main() {
}
volPlugins := make(map[string]volume.VolumePluginControlPlane)
ateletDialer := controlapi.NewAteletDialer(ateletPodInformer.GetIndexer(), *ateletClientCredBundle, *podIdentityCACerts)
ateletDialer := controlapi.NewAteletDialer(ateletPodInformer.GetIndexer(), ateletSPIFFEID, *ateletClientCredBundle, *podIdentityCACerts)
actorIDCAPool, err := localca.NewRefreshingPool(*actorIDCAPoolFile)
if err != nil {
@@ -292,7 +306,7 @@ func main() {
)
reflection.Register(mux)
ateapipb.RegisterControlServer(mux, controlSrv)
ateapipb.RegisterWorkerServiceServer(mux, workerservice.New(persistence))
ateapipb.RegisterWorkerServiceServer(mux, workerservice.New(persistence, ateletSPIFFEID))
readiness := &serverboot.Readiness{}
go serverboot.StartMetricsServer(ctx, serverboot.MetricsServerOptions{
@@ -40,13 +40,16 @@ import (
// CA pool.
type EgressMITMTrustReconciler struct {
client.Client
// SystemNamespace is the namespace holding the egress MITM CA pool Secret.
SystemNamespace string
}
// EgressMITMCAPoolRef names the Secret holding the CA pool the egress gateway's
// sdsmint sidecar signs per-SNI leaves with.
func EgressMITMCAPoolRef() types.NamespacedName {
func EgressMITMCAPoolRef(systemNamespace string) types.NamespacedName {
const egressMITMCAPoolSecret = "egress-mitm-ca-pool"
return types.NamespacedName{Namespace: ateSystemNamespace, Name: egressMITMCAPoolSecret}
return types.NamespacedName{Namespace: systemNamespace, Name: egressMITMCAPoolSecret}
}
//+kubebuilder:rbac:groups=core,resources=secrets,verbs=get;list;watch
@@ -162,7 +165,7 @@ func (r *EgressMITMTrustReconciler) deleteTrustBundle(ctx context.Context) error
}
func (r *EgressMITMTrustReconciler) SetupWithManager(mgr ctrl.Manager) error {
poolRef := EgressMITMCAPoolRef()
poolRef := EgressMITMCAPoolRef(r.SystemNamespace)
// The pool Secret is the only object reconciled from. The bundle is watched
// as well so that deleting or hand-editing the derived object is reverted
@@ -32,6 +32,7 @@ import (
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/localca"
)
@@ -66,7 +67,7 @@ func secretForPool(t *testing.T, pool *localca.ConcretePool) *corev1.Secret {
if err != nil {
t.Fatalf("marshal CA pool: %v", err)
}
ref := EgressMITMCAPoolRef()
ref := EgressMITMCAPoolRef(installdefaults.SystemNamespace)
return &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{Namespace: ref.Namespace, Name: ref.Name},
Data: map[string][]byte{"pool": wire},
@@ -86,8 +87,8 @@ func rootPEM(t *testing.T, pool *localca.ConcretePool) string {
func reconcilePool(t *testing.T, c client.Client) error {
t.Helper()
r := &EgressMITMTrustReconciler{Client: c}
_, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: EgressMITMCAPoolRef()})
r := &EgressMITMTrustReconciler{Client: c, SystemNamespace: installdefaults.SystemNamespace}
_, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: EgressMITMCAPoolRef(installdefaults.SystemNamespace)})
return err
}
@@ -170,7 +171,7 @@ func TestEgressMITMTrustFollowsPoolRotation(t *testing.T) {
rotated, rotatedPool := caPoolSecret(t, "mitm", "mitm-next")
current := &corev1.Secret{}
if err := c.Get(context.Background(), EgressMITMCAPoolRef(), current); err != nil {
if err := c.Get(context.Background(), EgressMITMCAPoolRef(installdefaults.SystemNamespace), current); err != nil {
t.Fatalf("get pool secret: %v", err)
}
current.Data = rotated.Data
@@ -279,7 +280,7 @@ func TestEgressMITMTrustKeepsLastGoodBundleOnBadPool(t *testing.T) {
}
current := &corev1.Secret{}
if err := c.Get(context.Background(), EgressMITMCAPoolRef(), current); err != nil {
if err := c.Get(context.Background(), EgressMITMCAPoolRef(installdefaults.SystemNamespace), current); err != nil {
t.Fatalf("get pool secret: %v", err)
}
current.Data = tc.data
@@ -33,13 +33,17 @@ import (
const (
networkPolicyFieldOwner = "ate-networkpolicy"
ateSystemNamespace = "ate-system"
atenetRouterAppName = "atenet-router"
)
type NetworkPolicyReconciler struct {
client.Client
Scheme *runtime.Scheme
// SystemNamespace is the namespace atenet-router runs in. The generated
// ingress policy admits only that namespace, so a value that does not
// match the running router blocks every request to the worker pool.
SystemNamespace string
}
//+kubebuilder:rbac:groups=ate.dev,resources=workerpools,verbs=get;list;watch
@@ -74,7 +78,7 @@ func (r *NetworkPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Reques
func (r *NetworkPolicyReconciler) reconcileImpl(ctx context.Context, wp *atev1alpha1.WorkerPool) error {
log := log.FromContext(ctx)
npAC := buildNetworkPolicyApplyConfig(wp)
npAC := r.buildNetworkPolicyApplyConfig(wp)
if err := r.Apply(ctx, npAC, client.FieldOwner(networkPolicyFieldOwner), client.ForceOwnership); err != nil {
return fmt.Errorf("failed to apply NetworkPolicy %s:%s: %w", *npAC.Namespace, *npAC.Name, err)
@@ -86,7 +90,7 @@ func (r *NetworkPolicyReconciler) reconcileImpl(ctx context.Context, wp *atev1al
return nil
}
func buildNetworkPolicyApplyConfig(wp *atev1alpha1.WorkerPool) *networkingv1ac.NetworkPolicyApplyConfiguration {
func (r *NetworkPolicyReconciler) buildNetworkPolicyApplyConfig(wp *atev1alpha1.WorkerPool) *networkingv1ac.NetworkPolicyApplyConfiguration {
np := networkingv1ac.NetworkPolicy(resources.NetworkPolicyName(wp.Name), wp.Namespace).
WithLabels(map[string]string{
"ate.dev/worker-pool": wp.Name,
@@ -110,7 +114,7 @@ func buildNetworkPolicyApplyConfig(wp *atev1alpha1.WorkerPool) *networkingv1ac.N
WithFrom(
networkingv1ac.NetworkPolicyPeer().
WithNamespaceSelector(metav1ac.LabelSelector().
WithMatchLabels(map[string]string{"kubernetes.io/metadata.name": ateSystemNamespace})).
WithMatchLabels(map[string]string{"kubernetes.io/metadata.name": r.SystemNamespace})).
WithPodSelector(metav1ac.LabelSelector().
WithMatchLabels(map[string]string{"app": atenetRouterAppName})),
),
@@ -21,6 +21,7 @@ import (
networkingv1 "k8s.io/api/networking/v1"
"k8s.io/apimachinery/pkg/types"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/resources"
)
@@ -75,7 +76,7 @@ func TestWorkerPoolCreatesNetworkPolicy(t *testing.T) {
return false, nil
}
fromPeer := ingressRule.From[0]
if fromPeer.NamespaceSelector == nil || fromPeer.NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"] != ateSystemNamespace {
if fromPeer.NamespaceSelector == nil || fromPeer.NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"] != installdefaults.SystemNamespace {
return false, nil
}
if fromPeer.PodSelector == nil || fromPeer.PodSelector.MatchLabels["app"] != atenetRouterAppName {
@@ -90,3 +91,24 @@ func TestWorkerPoolCreatesNetworkPolicy(t *testing.T) {
return true, nil
})
}
// TestBuildNetworkPolicyRelocatedNamespace pins the ingress peer to the
// reconciler's SystemNamespace rather than the canonical install namespace.
// The rest of the suite configures the reconciler with the default, so it
// passes just as well against a hardcoded "ate-system"; this is the case that
// catches that. A policy naming the wrong namespace admits nobody, and the CNI
// drops every request to the pool with no error from substrate itself.
func TestBuildNetworkPolicyRelocatedNamespace(t *testing.T) {
const relocated = "substrate-test"
r := &NetworkPolicyReconciler{SystemNamespace: relocated}
np := r.buildNetworkPolicyApplyConfig(testWorkerPoolApplyConfig(nil))
if len(np.Spec.Ingress) != 1 || len(np.Spec.Ingress[0].From) != 1 {
t.Fatalf("expected exactly one ingress rule with one peer, got %+v", np.Spec.Ingress)
}
got := np.Spec.Ingress[0].From[0].NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"]
if got != relocated {
t.Errorf("ingress namespace selector = %q, want %q", got, relocated)
}
}
@@ -28,6 +28,7 @@ import (
"github.com/agent-substrate/substrate/internal/ateomcapacity"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/deviceplugin"
"github.com/agent-substrate/substrate/internal/installdefaults"
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
)
@@ -92,7 +93,7 @@ const (
// Deployment managed by a WorkerPool. Only fields owned by this controller
// are declared here. otel, when it carries an endpoint, is propagated to the
// ateom container so it pushes telemetry to that collector.
func buildDeploymentApplyConfig(wp *atev1alpha1.WorkerPool, otel ateomOTelSettings) *appsv1ac.DeploymentApplyConfiguration {
func buildDeploymentApplyConfig(wp *atev1alpha1.WorkerPool, otel ateomOTelSettings, systemNamespace, ateletServiceAccount, routerServiceAccount string) *appsv1ac.DeploymentApplyConfiguration {
labels := map[string]string{}
annotations := map[string]string{}
if wp.Spec.Template != nil {
@@ -105,18 +106,42 @@ func buildDeploymentApplyConfig(wp *atev1alpha1.WorkerPool, otel ateomOTelSettin
}
labels["ate.dev/worker-pool"] = wp.Name
args := []string{
"--pod-uid=$(POD_UID)",
"--atunnel-listen-address=:443",
"--atunnel-connect-listen-address=:8443",
"--atunnel-credential-bundle=" + atunnelIdentityMountPath + "/credential-bundle.pem",
"--atunnel-trust-bundle=" + atunnelIdentityMountPath + "/trust-bundle.pem",
// The peers atunnel authenticates live in substrate's namespace, not
// the worker's, so the controller passes their identities rather than
// letting ateom assume the default install. --atunnel-client-identity
// has been accepted by every ateom that carries this controller's
// contemporaries, so it is always safe to pass.
"--atunnel-client-identity=" + installdefaults.SPIFFEID(systemNamespace, routerServiceAccount),
}
// --atunnel-broker-identity is newer than the oldest ateom a rolling
// upgrade still has running. docs/upgrade.md keeps the outgoing worker pool
// serving alongside the new one, and that pool's Deployment is reconciled
// by this controller while still pinned to its old image, which exits on an
// unrecognized flag. An ateom without the flag hardcodes the canonical
// identity, and an ateom with it defaults to the same, so omitting the flag
// when it carries that value is equivalent for both and keeps the upgrade
// intact. A relocated or renamed install passes something else and needs an
// image new enough to accept it, which it necessarily has.
if brokerIdentity := installdefaults.SPIFFEID(systemNamespace, ateletServiceAccount); brokerIdentity != installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace) {
args = append(args, "--atunnel-broker-identity="+brokerIdentity)
}
args = append(args,
"--atunnel-egress-listen-address=0.0.0.0:15001",
"--atunnel-egress-trust-bundle="+atunnelEgressTrustMountPath+"/trust-bundle.pem",
)
containerAC := corev1ac.Container().
WithName("ateom").
WithImage(wp.Spec.WorkerImage).
WithArgs(
"--pod-uid=$(POD_UID)",
"--atunnel-listen-address=:443",
"--atunnel-connect-listen-address=:8443",
"--atunnel-credential-bundle="+atunnelIdentityMountPath+"/credential-bundle.pem",
"--atunnel-trust-bundle="+atunnelIdentityMountPath+"/trust-bundle.pem",
"--atunnel-egress-listen-address=0.0.0.0:15001",
"--atunnel-egress-trust-bundle="+atunnelEgressTrustMountPath+"/trust-bundle.pem",
).
WithArgs(args...).
WithPorts(corev1ac.ContainerPort().
WithName("https").
WithContainerPort(443).
@@ -32,6 +32,7 @@ import (
"github.com/agent-substrate/substrate/internal/ateomcapacity"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/deviceplugin"
"github.com/agent-substrate/substrate/internal/installdefaults"
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
)
@@ -207,7 +208,7 @@ func TestBuildDeploymentApplyConfig(t *testing.T) {
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got := buildDeploymentApplyConfig(tt.wp, ateomOTelSettings{})
got := buildDeploymentApplyConfig(tt.wp, ateomOTelSettings{}, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount)
if diff := cmp.Diff(tt.want, got); diff != "" {
t.Fatalf("buildDeploymentApplyConfig() mismatch (-want +got):\n%s", diff)
}
@@ -227,7 +228,7 @@ func TestBuildDeploymentApplyConfigMetadata(t *testing.T) {
},
})
got := buildDeploymentApplyConfig(wp, ateomOTelSettings{})
got := buildDeploymentApplyConfig(wp, ateomOTelSettings{}, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount)
wantLabels := map[string]string{
"project": "agent-substrate",
"team": "compute",
@@ -269,7 +270,7 @@ func TestMicroVMPodShape(t *testing.T) {
t.Run(tt.name, func(t *testing.T) {
wp := testWorkerPoolApplyConfig(nil)
wp.Spec.SandboxClass = tt.class
ps := buildDeploymentApplyConfig(wp, ateomOTelSettings{}).Spec.Template.Spec
ps := buildDeploymentApplyConfig(wp, ateomOTelSettings{}, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).Spec.Template.Spec
// /dev/kvm must come from the device plugin, never a hostPath: a
// hostPath mount carries no cgroup device allow rule, and the
@@ -362,7 +363,7 @@ func TestMicroVMDeviceRequestsPreserveTemplateResources(t *testing.T) {
},
})
wp.Spec.SandboxClass = atev1alpha1.SandboxClassMicroVM
c := buildDeploymentApplyConfig(wp, ateomOTelSettings{}).Spec.Template.Spec.Containers[0]
c := buildDeploymentApplyConfig(wp, ateomOTelSettings{}, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).Spec.Template.Spec.Containers[0]
if got, ok := deviceLimit(c, string(corev1.ResourceMemory)); !ok || got != "2Gi" {
t.Errorf("memory limit = %q (present=%v), want 2Gi", got, ok)
@@ -427,7 +428,7 @@ func TestAteomSecurityContextByClass(t *testing.T) {
// TestTerminationGracePeriodSeconds asserts the pod's grace period is hardcoded to 3600s.
func TestTerminationGracePeriodSeconds(t *testing.T) {
wp := testWorkerPoolApplyConfig(nil)
ps := buildDeploymentApplyConfig(wp, ateomOTelSettings{}).Spec.Template.Spec
ps := buildDeploymentApplyConfig(wp, ateomOTelSettings{}, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).Spec.Template.Spec
if ps.TerminationGracePeriodSeconds == nil {
t.Fatalf("TerminationGracePeriodSeconds not set")
}
@@ -441,7 +442,7 @@ func TestTerminationGracePeriodSeconds(t *testing.T) {
// actors' drain window.
func TestRolloutStrategy(t *testing.T) {
wp := testWorkerPoolApplyConfig(nil)
spec := buildDeploymentApplyConfig(wp, ateomOTelSettings{}).Spec
spec := buildDeploymentApplyConfig(wp, ateomOTelSettings{}, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).Spec
if spec.Strategy == nil || spec.Strategy.RollingUpdate == nil {
t.Fatalf("Strategy.RollingUpdate not set")
}
@@ -475,7 +476,7 @@ func TestBuildDeploymentApplyConfigOTelEndpoint(t *testing.T) {
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), ateomOTelSettings{Endpoint: tt.endpoint}).
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), ateomOTelSettings{Endpoint: tt.endpoint}, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).
Spec.Template.Spec.Containers[0]
env := envByName(c.Env)
@@ -561,7 +562,7 @@ func TestBuildDeploymentApplyConfigMetricExportTuning(t *testing.T) {
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), tt.otel).
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), tt.otel, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).
Spec.Template.Spec.Containers[0]
env := envByName(c.Env)
for _, k := range []string{"OTEL_METRIC_EXPORT_INTERVAL", "OTEL_METRIC_EXPORT_TIMEOUT"} {
@@ -618,7 +619,7 @@ func TestBuildDeploymentApplyConfigTracesSamplerPropagation(t *testing.T) {
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), tt.otel).
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), tt.otel, installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).
Spec.Template.Spec.Containers[0]
env := envByName(c.Env)
for _, k := range []string{"OTEL_TRACES_SAMPLER", "OTEL_TRACES_SAMPLER_ARG"} {
@@ -753,6 +754,7 @@ func expectedDeploymentApplyConfig(mutatePodSpec func(*corev1ac.PodSpecApplyConf
"--atunnel-connect-listen-address=:8443",
"--atunnel-credential-bundle="+atunnelIdentityMountPath+"/credential-bundle.pem",
"--atunnel-trust-bundle="+atunnelIdentityMountPath+"/trust-bundle.pem",
"--atunnel-client-identity="+installdefaults.RouterSPIFFEID(installdefaults.SystemNamespace),
"--atunnel-egress-listen-address=0.0.0.0:15001",
"--atunnel-egress-trust-bundle="+atunnelEgressTrustMountPath+"/trust-bundle.pem",
).
@@ -842,3 +844,98 @@ func expectedDeploymentApplyConfig(mutatePodSpec func(*corev1ac.PodSpecApplyConf
WithLabels(map[string]string{"ate.dev/worker-pool": wp.Name}).
WithSpec(podSpecAC)))
}
// TestBuildDeploymentAtunnelIdentitiesRelocatedNamespace pins the SPIFFE
// identities handed to ateom to the controller's namespace. atunnel runs in
// the actor's pod, so it cannot derive atelet's or the router's namespace
// itself; if these carry the wrong one, the credential broker handshake and
// actor ingress both fail closed. Every other case here passes the canonical
// namespace and so would pass against a hardcoded value too.
func TestBuildDeploymentAtunnelIdentitiesRelocatedNamespace(t *testing.T) {
const relocated = "substrate-test"
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), ateomOTelSettings{}, relocated, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).
Spec.Template.Spec.Containers[0]
want := map[string]string{
"--atunnel-client-identity=": "spiffe://cluster.local/ns/substrate-test/sa/atenet-router",
"--atunnel-broker-identity=": "spiffe://cluster.local/ns/substrate-test/sa/atelet",
}
for flag, wantVal := range want {
var got string
for _, arg := range c.Args {
if strings.HasPrefix(arg, flag) {
got = strings.TrimPrefix(arg, flag)
}
}
if got == "" {
t.Fatalf("no %s argument found in %v", flag, c.Args)
}
if got != wantVal {
t.Errorf("%s%s, want %s%s", flag, got, flag, wantVal)
}
}
}
// TestBuildDeploymentAtunnelIdentitiesPrefixedServiceAccounts covers a
// deployment that renames the ServiceAccounts, which is what a packaging layer
// does when it prefixes every resource name. The SPIFFE ID embeds the
// ServiceAccount name, so identities built from the compiled-in defaults name
// accounts that do not exist and atunnel rejects the peer.
func TestBuildDeploymentAtunnelIdentitiesPrefixedServiceAccounts(t *testing.T) {
const (
namespace = "kagent-system"
atelet = "kagent-atelet"
router = "kagent-atenet-router"
)
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), ateomOTelSettings{}, namespace, atelet, router).
Spec.Template.Spec.Containers[0]
want := map[string]string{
"--atunnel-client-identity=": "spiffe://cluster.local/ns/kagent-system/sa/kagent-atenet-router",
"--atunnel-broker-identity=": "spiffe://cluster.local/ns/kagent-system/sa/kagent-atelet",
}
for flag, wantVal := range want {
var got string
for _, arg := range c.Args {
if strings.HasPrefix(arg, flag) {
got = strings.TrimPrefix(arg, flag)
}
}
if got != wantVal {
t.Errorf("%s%s, want %s%s", flag, got, flag, wantVal)
}
}
}
// TestBuildDeploymentOmitsBrokerIdentityForCanonicalInstall pins the flag's
// absence, which is what keeps a rolling upgrade working. docs/upgrade.md runs
// the outgoing worker pool alongside the new one, and this controller
// reconciles that pool's Deployment while it is still pinned to its old image.
// An ateom from before --atunnel-broker-identity existed exits on the
// unrecognized flag, so passing it would crashloop every old worker the moment
// the control plane rolled out.
func TestBuildDeploymentOmitsBrokerIdentityForCanonicalInstall(t *testing.T) {
c := buildDeploymentApplyConfig(testWorkerPoolApplyConfig(nil), ateomOTelSettings{},
installdefaults.SystemNamespace, installdefaults.AteletServiceAccount, installdefaults.RouterServiceAccount).
Spec.Template.Spec.Containers[0]
for _, arg := range c.Args {
if strings.HasPrefix(arg, "--atunnel-broker-identity=") {
t.Fatalf("canonical install passed %q; an ateom predating the flag exits on it", arg)
}
}
// The value it would have carried is the one such an ateom already assumes,
// so omitting it changes nothing for either binary.
var clientIdentity string
for _, arg := range c.Args {
if strings.HasPrefix(arg, "--atunnel-client-identity=") {
clientIdentity = strings.TrimPrefix(arg, "--atunnel-client-identity=")
}
}
if want := installdefaults.RouterSPIFFEID(installdefaults.SystemNamespace); clientIdentity != want {
t.Errorf("--atunnel-client-identity=%s, want %s", clientIdentity, want)
}
}
@@ -52,6 +52,15 @@ type WorkerPoolReconciler struct {
// OTelTracesSamplerArg is the OTEL_TRACES_SAMPLER_ARG propagated to ateom
// pods. Ignored unless OTelTracesSampler is set.
OTelTracesSamplerArg string
// SystemNamespace is the namespace substrate's control plane runs in, and
// AteletServiceAccount / RouterServiceAccount are the ServiceAccounts those
// components run as. Together they name the SPIFFE identities that atunnel
// authenticates inside each worker, which is why the ServiceAccount names
// are configuration and not constants: a deployment that prefixes
// resource names changes them.
SystemNamespace string
AteletServiceAccount string
RouterServiceAccount string
desiredWorkers metric.Int64ObservableUpDownCounter
readyWorkers metric.Int64ObservableUpDownCounter
@@ -116,7 +125,7 @@ func (r *WorkerPoolReconciler) applyDeployment(ctx context.Context, wp *atev1alp
MetricExportTimeout: r.OTelMetricExportTimeout,
TracesSampler: r.OTelTracesSampler,
TracesSamplerArg: r.OTelTracesSamplerArg,
})
}, r.SystemNamespace, r.AteletServiceAccount, r.RouterServiceAccount)
if err := r.Apply(ctx, depAC, client.FieldOwner(workerPoolFieldOwner), client.ForceOwnership); err != nil {
return fmt.Errorf("failed to apply Deployment: %w", err)
}
@@ -43,6 +43,7 @@ import (
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
"github.com/agent-substrate/substrate/internal/ateattr"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/testenv"
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
)
@@ -86,8 +87,9 @@ func TestMain(m *testing.M) {
}
if err := (&NetworkPolicyReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
SystemNamespace: installdefaults.SystemNamespace,
}).SetupWithManager(mgr); err != nil {
fmt.Fprintf(os.Stderr, "netpolicy controller setup failed: %v\n", err)
os.Exit(1)
+24 -4
View File
@@ -22,6 +22,7 @@ import (
"github.com/agent-substrate/substrate/cmd/atecontroller/internal/controllers"
"github.com/agent-substrate/substrate/cmd/atecontroller/internal/workersync"
"github.com/agent-substrate/substrate/internal/ateapiauth"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/serverboot"
"github.com/agent-substrate/substrate/internal/version"
clientv1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
@@ -73,6 +74,9 @@ var (
otelTracesSamplerArg = pflag.String("otel-traces-sampler-arg", os.Getenv("OTEL_TRACES_SAMPLER_ARG"),
"Trace sampler argument set on ateom worker pods, ignored unless --otel-traces-sampler is set. Defaults to the controller's own OTEL_TRACES_SAMPLER_ARG.")
ateletServiceAccount = pflag.String("atelet-service-account", installdefaults.AteletServiceAccount, "ServiceAccount atelet runs as. It is the service-account segment of the SPIFFE ID each worker's atunnel expects on the credential broker, so it has to match what the deployment actually creates.")
routerServiceAccount = pflag.String("router-service-account", installdefaults.RouterServiceAccount, "ServiceAccount atenet-router runs as. It is the service-account segment of the SPIFFE ID each worker's atunnel accepts on actor ingress, so it has to match what the deployment actually creates.")
ateapiCAFile = pflag.String("ateapi-ca-file", ateapiauth.DefaultServiceAccountCAFile, "PEM file with CAs trusted to verify the ateapi server cert.")
ateapiServerName = pflag.String("ateapi-server-name", "", "SNI / hostname expected on the ateapi server cert. Optional.")
ateapiClientCert = pflag.String("ateapi-client-cert", "", "Credential bundle presented as the client certificate when dialing ateapi. Required.")
@@ -159,8 +163,19 @@ func main() {
ateapiClient := ateapipb.NewControlClient(ateapiConn)
// Empty ServiceAccounts would not fail here: path.Join drops the empty
// segment, yielding identities that parse but match nothing, so every
// worker would reject both of its peers with no hint at the cause.
for flag, value := range map[string]string{"--atelet-service-account": *ateletServiceAccount, "--router-service-account": *routerServiceAccount} {
if value == "" {
setupLog.Error(nil, "invalid flag", "flag", flag, "reason", "must not be empty")
os.Exit(1)
}
}
// EgressMITMTrustReconciler watches the Secret `egress-mitm-ca-pool`.
egressMITMCAPool := controllers.EgressMITMCAPoolRef()
systemNamespace := installdefaults.NamespaceFromPodEnv()
egressMITMCAPool := controllers.EgressMITMCAPoolRef(systemNamespace)
mgr, err := ctrl.NewManager(k8sConfig, ctrl.Options{
Scheme: scheme,
Cache: cache.Options{
@@ -188,21 +203,26 @@ func main() {
OTelMetricExportTimeout: *otelMetricExportTimeout,
OTelTracesSampler: *otelTracesSampler,
OTelTracesSamplerArg: *otelTracesSamplerArg,
SystemNamespace: systemNamespace,
AteletServiceAccount: *ateletServiceAccount,
RouterServiceAccount: *routerServiceAccount,
}).SetupWithManager(mgr); err != nil {
setupLog.Error(err, "unable to create controller", "controller", "WorkerPool")
os.Exit(1)
}
if err = (&controllers.NetworkPolicyReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
SystemNamespace: systemNamespace,
}).SetupWithManager(mgr); err != nil {
setupLog.Error(err, "unable to create controller", "controller", "NetPolicy")
os.Exit(1)
}
if err = (&controllers.EgressMITMTrustReconciler{
Client: mgr.GetClient(),
Client: mgr.GetClient(),
SystemNamespace: systemNamespace,
}).SetupWithManager(mgr); err != nil {
setupLog.Error(err, "unable to create controller", "controller", "EgressMITMTrust")
os.Exit(1)
+2
View File
@@ -23,6 +23,7 @@ import (
"github.com/agent-substrate/substrate/cmd/atenet/internal/router/egress"
"github.com/agent-substrate/substrate/cmd/atenet/internal/router/ingress"
"github.com/agent-substrate/substrate/internal/installdefaults"
)
func NewRouterCmd() *cobra.Command {
@@ -47,6 +48,7 @@ func NewRouterCmd() *cobra.Command {
cmd.Flags().StringVar(&cfg.MetricsAddr, "metrics-listen-addr", ":9090", "Address and port the prometheus metrics server should listen on.")
cmd.Flags().StringVar(&cfg.AtenetRouter, "atenet-dataplane", string(atenetRouterEnvoy), "Atenet ingress and egress dataplane: envoy or agentgateway")
cmd.Flags().StringVar(&cfg.Namespace, "namespace", "default", "Target operations namespace")
cmd.Flags().StringVar(&cfg.RouterServiceName, "router-service-name", installdefaults.RouterServiceName, "Service name of this atenet-router in the operations namespace. Override when the deployment renames the Service.")
cmd.Flags().StringVar(&cfg.Kubeconfig, "kubeconfig", "", "Absolute path to the kubeconfig configuration file")
cmd.Flags().StringVar(&cfg.AteapiAddr, "ateapi-address", "k8s:///api.ate-system.svc:443", "gRPC dial target for the cluster ateapi Control instance.")
cmd.Flags().IntVar(&cfg.HttpPort, "port-http", 8080, "TCP port for workload traffic entering through the Envoy Router")
+16 -12
View File
@@ -97,18 +97,22 @@ type credentialProviderConfig struct {
// routerConfig holds deployment setup and endpoint options for the router node instance.
type routerConfig struct {
// Mode restricts the instance to one traffic direction. Empty means ModeAll.
Mode Mode
AtenetRouter string
Namespace string
Kubeconfig string
AteapiAddr string
HttpPort int
XdsPort int
ExtprocPort int
ExtprocAddr string
StatusPort int
HealthInterval time.Duration
HttpsPort int
Mode Mode
AtenetRouter string
Namespace string
// RouterServiceName is the Service name of this atenet-router in the
// operations namespace, used by /statusz to look up its own ClusterIP.
// Defaults to installdefaults.RouterServiceName.
RouterServiceName string
Kubeconfig string
AteapiAddr string
HttpPort int
XdsPort int
ExtprocPort int
ExtprocAddr string
StatusPort int
HealthInterval time.Duration
HttpsPort int
// ConnectPlainTextPort and ConnectTLSPort are the plaintext and TLS
// listener ports for CONNECT-tunneled traffic. Non-positive disables the
// corresponding listener.
+1 -1
View File
@@ -68,7 +68,7 @@ func (s *RouterServer) getRouterIP(ctx context.Context) string {
return "Offline Mode (No Cluster IP)"
}
svc, err := s.clientset.CoreV1().Services(s.cfg.Namespace).Get(ctx, "atenet-router", metav1.GetOptions{})
svc, err := s.clientset.CoreV1().Services(s.cfg.Namespace).Get(ctx, s.cfg.RouterServiceName, metav1.GetOptions{})
if err != nil {
return fmt.Sprintf("Lookup Failed: %v", err)
}
+12 -3
View File
@@ -45,6 +45,7 @@ import (
"github.com/agent-substrate/substrate/internal/childreap"
"github.com/agent-substrate/substrate/internal/contextlogging"
"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/ocispec"
"github.com/agent-substrate/substrate/internal/otlprelay"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
@@ -74,7 +75,8 @@ var (
atunnelConnectListenAddress = pflag.String("atunnel-connect-listen-address", ":8443", "Address for actor ingress mTLS CONNECT")
workerCredentialBundle = pflag.String("atunnel-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "Worker Pod credential bundle used by atunnel for inbound serving and outbound mTLS")
podIdentityTrustBundle = pflag.String("atunnel-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "Pod identity trust bundle used for router clients and the node-local atelet")
atunnelClientIdentity = pflag.String("atunnel-client-identity", "spiffe://cluster.local/ns/ate-system/sa/atenet-router", "SPIFFE identity allowed to call actor ingress HTTPS")
atunnelClientIdentity = pflag.String("atunnel-client-identity", installdefaults.RouterSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity allowed to call actor ingress HTTPS")
ateletIdentity = pflag.String("atunnel-broker-identity", installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity the node-local atelet must present on the credential broker connection. Override when atelet runs outside the default namespace.")
atunnelEgressListenAddress = pflag.String("atunnel-egress-listen-address", "0.0.0.0:15001", "Address for transparently intercepted actor egress TCP")
egressGatewayTrustBundle = pflag.String("atunnel-egress-trust-bundle", "/run/servicedns.podcert.ate.dev/trust-bundle.pem", "Service DNS trust bundle for the remote egress gateway")
readinessListenAddress = pflag.String("readiness-listen-address", "0.0.0.0:8080", "Address for HTTP readiness checks")
@@ -223,7 +225,7 @@ func do(ctx context.Context) error {
return err
}
ateomService := NewService(interiorNetNS, actorLogger, atunnelIngress, atunnelEgress, atunnelEgressPort, *workerCredentialBundle, *podIdentityTrustBundle, *egressGatewayTrustBundle)
ateomService := NewService(interiorNetNS, actorLogger, atunnelIngress, atunnelEgress, atunnelEgressPort, *workerCredentialBundle, *podIdentityTrustBundle, *egressGatewayTrustBundle, *ateletIdentity)
svr := grpc.NewServer(
grpc.StatsHandler(otelgrpc.NewServerHandler()),
@@ -259,6 +261,7 @@ func do(ctx context.Context) error {
SocketPath: ateompath.AteomSupportSocket,
CredentialBundlePath: *workerCredentialBundle,
TrustBundlePath: *podIdentityTrustBundle,
AteletSPIFFEID: *ateletIdentity,
})
if err != nil && ctx.Err() == nil {
serverboot.Fatal(ctx, "Failed to report worker capacity", err)
@@ -397,6 +400,10 @@ type AteomService struct {
podIdentityTrustBundlePath string
// egressGatewayTrustBundlePath verifies the remote gateway's serving cert.
egressGatewayTrustBundlePath string
// ateletSPIFFEID is the identity the node-local atelet must present on the
// credential broker connection. It names atelet's namespace, not this
// worker's, so it is configured rather than derived from the downward API.
ateletSPIFFEID string
// activeActor is the actor whose workload this ateom is currently running,
// or nil when it is "available". An ateom serves one actor at a time, so a
@@ -448,7 +455,7 @@ type AteomService struct {
var _ ateompb.AteomServer = (*AteomService)(nil)
// NewService creates a new AteomService.
func NewService(interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger, atunnelIngress *atunnel.Server, atunnelEgress *atunnel.Egress, atunnelEgressPort uint16, workerCredentialBundlePath, podIdentityTrustBundlePath, egressGatewayTrustBundlePath string) *AteomService {
func NewService(interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger, atunnelIngress *atunnel.Server, atunnelEgress *atunnel.Egress, atunnelEgressPort uint16, workerCredentialBundlePath, podIdentityTrustBundlePath, egressGatewayTrustBundlePath, ateletSPIFFEID string) *AteomService {
return &AteomService{
lock: newCancelableMutex(),
interiorNetNS: interiorNetNS,
@@ -459,6 +466,7 @@ func NewService(interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger,
workerCredentialBundlePath: workerCredentialBundlePath,
podIdentityTrustBundlePath: podIdentityTrustBundlePath,
egressGatewayTrustBundlePath: egressGatewayTrustBundlePath,
ateletSPIFFEID: ateletSPIFFEID,
cgroupRoot: defaultCgroupRoot,
}
}
@@ -1118,6 +1126,7 @@ func (s *AteomService) prepareActorEgress(ctx context.Context, actorAtespace, ac
ActorAtespace: actorAtespace,
ActorName: actorName,
ActorUID: actorUID,
AteletSPIFFEID: s.ateletSPIFFEID,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor certificate broker: %w", err)
+15 -6
View File
@@ -46,6 +46,7 @@ import (
"github.com/agent-substrate/substrate/internal/ateomnet"
"github.com/agent-substrate/substrate/internal/ateompath"
"github.com/agent-substrate/substrate/internal/atunnel"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/otlprelay"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/agent-substrate/substrate/internal/resources"
@@ -77,7 +78,8 @@ var (
atunnelConnectListenAddress = flag.String("atunnel-connect-listen-address", ":8443", "Address for actor ingress mTLS CONNECT")
workerCredentialBundle = flag.String("atunnel-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "Worker Pod credential bundle used by atunnel for inbound serving and outbound mTLS")
podIdentityTrustBundle = flag.String("atunnel-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "Pod identity trust bundle used for router clients and the node-local atelet")
atunnelClientIdentity = flag.String("atunnel-client-identity", "spiffe://cluster.local/ns/ate-system/sa/atenet-router", "SPIFFE identity allowed to call actor ingress HTTPS")
atunnelClientIdentity = flag.String("atunnel-client-identity", installdefaults.RouterSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity allowed to call actor ingress HTTPS")
ateletIdentity = flag.String("atunnel-broker-identity", installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity the node-local atelet must present on the credential broker connection. Override when atelet runs outside the default namespace.")
atunnelEgressListenAddress = flag.String("atunnel-egress-listen-address", "0.0.0.0:15001", "Address for transparently intercepted actor egress TCP")
egressGatewayTrustBundle = flag.String("atunnel-egress-trust-bundle", "/run/servicedns.podcert.ate.dev/trust-bundle.pem", "Service DNS trust bundle for the remote egress gateway")
readinessListenAddress = flag.String("readiness-listen-address", "0.0.0.0:8080", "Address for HTTP readiness checks")
@@ -265,7 +267,7 @@ func do(ctx context.Context) error {
}()
slog.InfoContext(ctx, "atunnel egress serving", slog.String("address", *atunnelEgressListenAddress))
ateomService := NewService(*podUID, *chBinary, *kataDebug, *vmmMemReserve, interiorNetNS, actorLogger, atunnelIngress, atunnelEgress, atunnelEgressPort, *workerCredentialBundle, *podIdentityTrustBundle, *egressGatewayTrustBundle)
ateomService := NewService(*podUID, *chBinary, *kataDebug, *vmmMemReserve, interiorNetNS, actorLogger, atunnelIngress, atunnelEgress, atunnelEgressPort, *workerCredentialBundle, *podIdentityTrustBundle, *egressGatewayTrustBundle, *ateletIdentity)
svr := grpc.NewServer(
grpc.StatsHandler(otelgrpc.NewServerHandler()),
@@ -305,6 +307,7 @@ func do(ctx context.Context) error {
SocketPath: ateompath.AteomSupportSocket,
CredentialBundlePath: *workerCredentialBundle,
TrustBundlePath: *podIdentityTrustBundle,
AteletSPIFFEID: *ateletIdentity,
})
if err != nil && ctx.Err() == nil {
serverboot.Fatal(ctx, "Failed to report worker capacity", err)
@@ -447,6 +450,10 @@ type AteomService struct {
podIdentityTrustBundlePath string
// egressGatewayTrustBundlePath verifies the remote gateway's serving cert.
egressGatewayTrustBundlePath string
// ateletSPIFFEID is the identity the node-local atelet must present on the
// credential broker connection. It names atelet's namespace, not this
// worker's, so it is configured rather than derived from the downward API.
ateletSPIFFEID string
// running maps actor UID -> the live micro-VM, kept so CheckpointWorkload can
// pause+snapshot+teardown the same sandbox (and RestoreWorkload can track the
@@ -500,7 +507,7 @@ type AteomService struct {
var _ ateompb.AteomServer = (*AteomService)(nil)
// NewService creates a new AteomService.
func NewService(podUID, chBinary string, kataDebug bool, memReserveMiB int, interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger, atunnelIngress *atunnel.Server, atunnelEgress *atunnel.Egress, atunnelEgressPort uint16, workerCredentialBundlePath, podIdentityTrustBundlePath, egressGatewayTrustBundlePath string) *AteomService {
func NewService(podUID, chBinary string, kataDebug bool, memReserveMiB int, interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger, atunnelIngress *atunnel.Server, atunnelEgress *atunnel.Egress, atunnelEgressPort uint16, workerCredentialBundlePath, podIdentityTrustBundlePath, egressGatewayTrustBundlePath, ateletSPIFFEID string) *AteomService {
return &AteomService{
lock: newCancelableMutex(),
podUID: podUID,
@@ -515,6 +522,7 @@ func NewService(podUID, chBinary string, kataDebug bool, memReserveMiB int, inte
workerCredentialBundlePath: workerCredentialBundlePath,
podIdentityTrustBundlePath: podIdentityTrustBundlePath,
egressGatewayTrustBundlePath: egressGatewayTrustBundlePath,
ateletSPIFFEID: ateletSPIFFEID,
running: map[string]*runningActor{},
}
}
@@ -543,9 +551,10 @@ func (s *AteomService) prepareActorEgress(ctx context.Context, actorAtespace, ac
CredentialBundlePath: s.workerCredentialBundlePath,
TrustBundlePath: s.podIdentityTrustBundlePath,
ActorAtespace: actorAtespace,
ActorName: actorName,
ActorUID: actorUID,
ActorAtespace: actorAtespace,
ActorName: actorName,
ActorUID: actorUID,
AteletSPIFFEID: s.ateletSPIFFEID,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor certificate broker: %w", err)
@@ -38,6 +38,7 @@ import (
"k8s.io/client-go/rest"
"github.com/agent-substrate/substrate/internal/credbundle"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/serverboot"
"github.com/agent-substrate/substrate/internal/version"
"github.com/agent-substrate/substrate/pkg/proto/credproviderpb"
@@ -45,20 +46,20 @@ import (
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.")
// The injector is the only caller allowed to fetch secrets. Its identity
// names the namespace and ServiceAccount atenet-egress runs as, so a
// deployment that relocates or renames substrate must set it.
injectorIdentity = pflag.String("injector-identity", installdefaults.EgressSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity of the credential injector allowed to fetch secrets")
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() {
@@ -186,7 +187,7 @@ func buildServerCreds(ctx context.Context) (credentials.TransportCredentials, er
}
serverCert := credbundle.Loader(*serverBundle)
verifySAN := verifyClientSAN(injectorSPIFFEID)
verifySAN := verifyClientSAN(*injectorIdentity)
// GetConfigForClient builds the config anew per connection: a certificate
// signed by a newly published CA verifies without a restart.
@@ -206,7 +207,7 @@ func buildServerCreds(ctx context.Context) (credentials.TransportCredentials, er
},
}
slog.InfoContext(ctx, "verifying caller client certificates",
slog.String("ca", *clientCAFile), slog.String("required_san", injectorSPIFFEID))
slog.String("ca", *clientCAFile), slog.String("required_san", *injectorIdentity))
return credentials.NewTLS(cfg), nil
}
@@ -31,6 +31,7 @@ import (
"testing"
"time"
"github.com/agent-substrate/substrate/internal/installdefaults"
"google.golang.org/grpc/credentials"
)
@@ -48,7 +49,7 @@ func certWithURIs(t *testing.T, uris ...string) *x509.Certificate {
}
func TestVerifyClientSAN(t *testing.T) {
const injector = injectorSPIFFEID
injector := installdefaults.EgressSPIFFEID(installdefaults.SystemNamespace)
tests := []struct {
name string
@@ -118,8 +119,8 @@ func TestBuildServerCredsReloadsClientCA(t *testing.T) {
t.Fatalf("buildServerCreds() error = %v", err)
}
fromCA1 := clientCA1.issue(t, certOpts{uris: []string{injectorSPIFFEID}})
fromCA2 := clientCA2.issue(t, certOpts{uris: []string{injectorSPIFFEID}})
fromCA1 := clientCA1.issue(t, certOpts{uris: []string{installdefaults.EgressSPIFFEID(installdefaults.SystemNamespace)}})
fromCA2 := clientCA2.issue(t, certOpts{uris: []string{installdefaults.EgressSPIFFEID(installdefaults.SystemNamespace)}})
wrongSAN := clientCA1.issue(t, certOpts{uris: []string{"spiffe://cluster.local/ns/ate-system/sa/impostor"}})
// Before rotation only CA1 is trusted.
+52 -8
View File
@@ -24,6 +24,7 @@ import (
"strings"
"sync"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/portforward"
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
@@ -43,9 +44,52 @@ import (
metricsv1beta1 "k8s.io/metrics/pkg/client/clientset/versioned"
)
const (
apiServerName = "api.ate-system.svc"
// NamespaceEnv overrides the namespace the client looks for substrate in. The
// client runs outside the cluster, so it has no downward API to read and no pod
// namespace to fall back on.
const NamespaceEnv = "ATE_NAMESPACE"
// APIServiceEnv and ClientServiceAccountEnv override the ateapi Service and the
// ServiceAccount the client mints its token from. A deployment that prefixes
// resource names changes both, and the client has no way to discover that from
// outside the cluster.
const (
APIServiceEnv = "ATE_API_SERVICE_NAME"
ClientServiceAccountEnv = "ATE_CLIENT_SERVICE_ACCOUNT"
)
// apiServiceName is the Service that fronts ateapi.
func apiServiceName() string {
if n := os.Getenv(APIServiceEnv); n != "" {
return n
}
return installdefaults.APIServiceName
}
// clientServiceAccount is the ServiceAccount the bearer token is minted from.
func clientServiceAccount() string {
if n := os.Getenv(ClientServiceAccountEnv); n != "" {
return n
}
return installdefaults.ClientServiceAccount
}
// systemNamespace is the namespace the client expects ateapi to be running in.
func systemNamespace() string {
if ns := os.Getenv(NamespaceEnv); ns != "" {
return ns
}
return installdefaults.SystemNamespace
}
// apiServerName is the in-cluster DNS name of the ateapi Service. It is both the
// name checked on ateapi's serving cert and the audience of the bearer token
// minted for it, so it has to track the namespace ateapi actually runs in.
func apiServerName() string {
return fmt.Sprintf("%s.%s.svc", apiServiceName(), systemNamespace())
}
const (
// serviceDNSSignerName and liveBundleSelector mirror the
// clusterTrustBundle projected-volume sources that in-cluster clients
// mount to verify ateapi's serving cert.
@@ -82,8 +126,8 @@ func (c *Client) Close() {
}
}
// NewClient creates a new Ate API client. If endpoint is empty, it automatically port-forwards
// to the ate-api-server pod in the ate-system namespace.
// NewClient creates a new Ate API client. If endpoint is empty, it automatically
// port-forwards to the ate-api-server pod in substrate's system namespace.
func NewClient(ctx context.Context, kubeconfigPath, k8sContext, endpoint, tokenFile string, traceEnabled bool) (*Client, error) {
tp, err := initTracing(ctx, traceEnabled)
if err != nil {
@@ -184,7 +228,7 @@ func dialPortForward(ctx context.Context, kubeconfigPath, k8sContext, tokenFile
// TODO: Should we special-case a LoadBalancer "api" Service and dial its
// address directly instead of port-forwarding?
localPort, stopForward, err := portforward.ServicePortForward(ctx, config, clientset, "ate-system", "api", 443)
localPort, stopForward, err := portforward.ServicePortForward(ctx, config, clientset, systemNamespace(), apiServiceName(), 443)
if err != nil {
return nil, err
}
@@ -249,7 +293,7 @@ func serverTLSConfig(ctx context.Context, clientset kubernetes.Interface) (*tls.
return &tls.Config{
MinVersion: tls.VersionTLS13,
RootCAs: pool,
ServerName: apiServerName,
ServerName: apiServerName(),
}, nil
}
@@ -269,11 +313,11 @@ func bearerTokenDialOption(ctx context.Context, clientset *kubernetes.Clientset,
expirationSeconds := int64(3600)
tokenRequest := &authv1.TokenRequest{
Spec: authv1.TokenRequestSpec{
Audiences: []string{apiServerName},
Audiences: []string{apiServerName()},
ExpirationSeconds: &expirationSeconds,
},
}
token, err := clientset.CoreV1().ServiceAccounts("ate-system").CreateToken(ctx, "ate-client", tokenRequest, metav1.CreateOptions{})
token, err := clientset.CoreV1().ServiceAccounts(systemNamespace()).CreateToken(ctx, clientServiceAccount(), tokenRequest, metav1.CreateOptions{})
if err != nil {
return nil, fmt.Errorf("failed to request ateapi bearer token: %w", err)
}
+7 -8
View File
@@ -22,9 +22,7 @@ import (
"crypto/x509"
"fmt"
"net"
"net/url"
"os"
"path"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
@@ -34,10 +32,12 @@ import (
)
// TLSConfig authenticates this worker to atelet with its Pod certificate, and
// accepts only the atelet on this worker's own node.
func TLSConfig(credentialBundlePath, trustBundlePath string) (*tls.Config, error) {
if credentialBundlePath == "" || trustBundlePath == "" {
return nil, fmt.Errorf("worker credentials and trust bundle are required")
// accepts only the atelet on this worker's own node. ateletSPIFFEID is the
// identity atelet presents; it names atelet's namespace, not this worker's, so
// the caller passes it in rather than deriving it from the downward API.
func TLSConfig(credentialBundlePath, trustBundlePath, ateletSPIFFEID string) (*tls.Config, error) {
if credentialBundlePath == "" || trustBundlePath == "" || ateletSPIFFEID == "" {
return nil, fmt.Errorf("worker credentials, trust bundle, and atelet SPIFFE ID are required")
}
localCert, err := credbundle.Parse(credentialBundlePath)
if err != nil {
@@ -55,7 +55,6 @@ func TLSConfig(credentialBundlePath, trustBundlePath string) (*tls.Config, error
if !roots.AppendCertsFromPEM(trustPEM) {
return nil, fmt.Errorf("atelet trust bundle contains no certificates")
}
expectedURI := (&url.URL{Scheme: "spiffe", Host: "cluster.local", Path: path.Join("ns", "ate-system", "sa", "atelet")}).String()
return &tls.Config{
MinVersion: tls.VersionTLS13,
InsecureSkipVerify: true, // Verification below supports SPIFFE Pod certificates without a DNS name.
@@ -75,7 +74,7 @@ func TLSConfig(credentialBundlePath, trustBundlePath string) (*tls.Config, error
return fmt.Errorf("verify atelet certificate: %w", err)
}
leaf := state.PeerCertificates[0]
if len(leaf.URIs) != 1 || leaf.URIs[0].String() != expectedURI {
if len(leaf.URIs) != 1 || leaf.URIs[0].String() != ateletSPIFFEID {
return fmt.Errorf("node-local peer is not atelet")
}
identity, err := substratex509.PodIdentityFromCertificate(leaf)
+5 -1
View File
@@ -99,6 +99,10 @@ type ReportConfig struct {
SocketPath string
CredentialBundlePath string
TrustBundlePath string
// AteletSPIFFEID is the identity the node-local atelet must present. It
// names atelet's namespace, not this worker's, so it is configured rather
// than derived from the downward API.
AteletSPIFFEID string
}
// Report tells the node-local atelet what this ateom can supply, retrying
@@ -109,7 +113,7 @@ type ReportConfig struct {
// an ateom first comes up. Nothing else reports this, so giving up would leave
// the Worker holding no capacity and hosting nothing.
func Report(ctx context.Context, cfg ReportConfig) error {
tlsConfig, err := ateletdial.TLSConfig(cfg.CredentialBundlePath, cfg.TrustBundlePath)
tlsConfig, err := ateletdial.TLSConfig(cfg.CredentialBundlePath, cfg.TrustBundlePath, cfg.AteletSPIFFEID)
if err != nil {
return fmt.Errorf("capacity report: %w", err)
}
+9 -1
View File
@@ -61,6 +61,11 @@ type BrokerConfig struct {
ActorAtespace string
ActorName string
ActorUID string
// AteletSPIFFEID is the identity the node-local atelet must present on the
// credential broker connection. It names atelet's namespace, not this
// worker's, so it is configured rather than derived from the downward API.
AteletSPIFFEID string
}
// NewBrokerCertificateSource creates one actor key for this activation. The key
@@ -75,7 +80,10 @@ func NewBrokerCertificateSource(cfg BrokerConfig) (*BrokerCertificateSource, err
return nil, fmt.Errorf("actor information is required")
}
tlsConfig, err := ateletdial.TLSConfig(cfg.CredentialBundlePath, cfg.TrustBundlePath)
if cfg.AteletSPIFFEID == "" {
return nil, fmt.Errorf("atunnel: expected atelet SPIFFE ID is required")
}
tlsConfig, err := ateletdial.TLSConfig(cfg.CredentialBundlePath, cfg.TrustBundlePath, cfg.AteletSPIFFEID)
if err != nil {
return nil, fmt.Errorf("atunnel: %w", err)
}
+2
View File
@@ -29,6 +29,7 @@ import (
"testing"
"time"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/proto/ateletpb"
"github.com/agent-substrate/substrate/internal/substratex509"
"google.golang.org/grpc"
@@ -187,6 +188,7 @@ func newTestBrokerCertificateSource(t *testing.T, ateletIdentity *substratex509.
ActorAtespace: "actor-atespace",
ActorName: "actor-name",
ActorUID: "actor-uid",
AteletSPIFFEID: installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace),
})
if err != nil {
t.Fatalf("Error creating broker certificate source: %v", err)
+1 -1
View File
@@ -65,7 +65,7 @@ func ScrapeAgentGatewayRouterMetrics(ctx context.Context) (string, error) {
if err != nil {
return "", fmt.Errorf("creating k8s client: %w", err)
}
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, routerNamespace, routerService, agentGatewayRouterStatsPort)
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, SystemNamespace(), ResourceName("atenet-router"), agentGatewayRouterStatsPort)
if err != nil {
return "", err
}
+26
View File
@@ -17,6 +17,8 @@ package e2e
import (
"fmt"
"os"
"github.com/agent-substrate/substrate/internal/installdefaults"
)
// CheckEnv checks the list of env vars exist and returns their value.
@@ -32,3 +34,27 @@ func CheckEnv(keys ...string) (map[string]string, error) {
}
return env, nil
}
// SystemNamespaceEnv names the namespace the substrate control plane under test
// was installed into. It mirrors the ATE_NAMESPACE hack/install-ate.sh
// installed with, and the namespace the deployment was installed into.
const SystemNamespaceEnv = "E2E_SYSTEM_NAMESPACE"
// SystemNamespace returns the namespace the control plane under test runs in,
// falling back to the canonical install namespace.
func SystemNamespace() string {
if ns := os.Getenv(SystemNamespaceEnv); ns != "" {
return ns
}
return installdefaults.SystemNamespace
}
// ResourcePrefixEnv is the prefix the install under test puts on substrate's
// resource names. A packaging layer that renames resources prefixes them, and
// the harness addresses several of those resources by name.
const ResourcePrefixEnv = "E2E_RESOURCE_PREFIX"
// ResourceName returns name as the install under test renders it.
func ResourceName(name string) string {
return os.Getenv(ResourcePrefixEnv) + name
}
@@ -15,7 +15,7 @@
# The egress probe used by internal/e2e/suites/egressauthz: testserver's probe
# subcommand, with the credential volumes that suite mints. It is not a
# ServerPod because it needs those volumes and a namespace the suite populates
# first, so it stays a bespoke manifest. ${NAMESPACE} is substituted by the
# first, so it stays a bespoke manifest. ${NAMESPACE} and ${SYSTEM_NAMESPACE} are substituted by the
# suite with the randomized namespace it created, so the probe is torn down with
# that namespace and leaves nothing behind.
apiVersion: v1
@@ -33,6 +33,10 @@ spec:
args:
- "egressprobe"
- "--listen=:8080"
# The gateway lives in substrate's namespace, which the probe cannot infer
# from inside the sandbox. ${SYSTEM_NAMESPACE} is substituted alongside
# ${NAMESPACE} by the suite.
- "--gateway-address=${EGRESS_GATEWAY_SERVICE}.${SYSTEM_NAMESPACE}.svc:443"
ports:
- name: http
containerPort: 8080
+3 -3
View File
@@ -38,10 +38,10 @@ func PreflightChecks() error {
// Check deployments.
deployments := []string{
"ate-controller",
"ate-api-server",
ResourceName("ate-controller"),
ResourceName("ate-api-server"),
}
namespace := "ate-system"
namespace := SystemNamespace()
for _, depName := range deployments {
dep, err := clients.K8s.AppsV1().Deployments(namespace).Get(ctx, depName, metav1.GetOptions{})
if err != nil {
+2 -4
View File
@@ -37,8 +37,6 @@ import (
)
const (
routerNamespace = "ate-system"
routerService = "atenet-router"
// routerConnectServicePort is atenet-router's Service port for
// CONNECT-tunneled traffic (see manifests/ate-install/atenet-router.yaml).
// It is a distinct listener from the plain HTTP one Get/PostJSON use:
@@ -80,7 +78,7 @@ func NewRouterClient(ctx context.Context) (*RouterClient, error) {
return nil, fmt.Errorf("creating k8s client: %w", err)
}
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, routerNamespace, routerService, 80)
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, SystemNamespace(), ResourceName("atenet-router"), 80)
if err != nil {
return nil, err
}
@@ -190,7 +188,7 @@ func (c *RouterClient) Connect(ctx context.Context, actorRef resources.ActorRef,
// in one test don't each pay for a fresh port-forward.
func (c *RouterClient) ensureConnectPortForward(ctx context.Context) error {
c.connectOnce.Do(func() {
localPort, stop, err := portforward.ServicePortForward(ctx, c.config, c.clientset, routerNamespace, routerService, routerConnectServicePort)
localPort, stop, err := portforward.ServicePortForward(ctx, c.config, c.clientset, SystemNamespace(), ResourceName("atenet-router"), routerConnectServicePort)
if err != nil {
c.connectErr = fmt.Errorf("port-forwarding to the router's CONNECT listener: %w", err)
return
+1 -1
View File
@@ -51,7 +51,7 @@ func NewStatuszClient(ctx context.Context) (*StatuszClient, error) {
return nil, fmt.Errorf("creating k8s client: %w", err)
}
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, routerNamespace, routerService, routerStatusPort)
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, SystemNamespace(), ResourceName("atenet-router"), routerStatusPort)
if err != nil {
return nil, err
}
+2 -2
View File
@@ -1348,13 +1348,13 @@ func callActorPathOnce(t *testing.T, actorRef resources.ActorRef, method, path s
t.Helper()
clients := e2e.GetClients()
svc, err := clients.K8s.CoreV1().Services("ate-system").Get(context.Background(), "atenet-router", metav1.GetOptions{})
svc, err := clients.K8s.CoreV1().Services(e2e.SystemNamespace()).Get(context.Background(), e2e.ResourceName("atenet-router"), metav1.GetOptions{})
if err != nil {
return "", fmt.Errorf("failed to get atenet-router service: %w", err)
}
selector := labels.SelectorFromSet(svc.Spec.Selector).String()
pods, err := clients.K8s.CoreV1().Pods("ate-system").List(context.Background(), metav1.ListOptions{LabelSelector: selector})
pods, err := clients.K8s.CoreV1().Pods(e2e.SystemNamespace()).List(context.Background(), metav1.ListOptions{LabelSelector: selector})
if err != nil {
return "", fmt.Errorf("failed to list atenet-router pods: %w", err)
}
@@ -73,16 +73,16 @@ const (
// the secret ateapi signs with.
func actorIdentityCA(t *testing.T, ctx context.Context) *localca.CA {
t.Helper()
secret, err := e2e.GetClients().K8s.CoreV1().Secrets(egressNamespace).Get(ctx, actorIDCASecret, metav1.GetOptions{})
secret, err := e2e.GetClients().K8s.CoreV1().Secrets(e2e.SystemNamespace()).Get(ctx, actorIDCASecret, metav1.GetOptions{})
if err != nil {
t.Fatalf("reading actor-identity CA pool secret %s/%s: %v", egressNamespace, actorIDCASecret, err)
t.Fatalf("reading actor-identity CA pool secret %s/%s: %v", e2e.SystemNamespace(), actorIDCASecret, err)
}
pool, err := localca.Unmarshal(secret.Data[actorIDCASecretKey])
if err != nil {
t.Fatalf("parsing actor-identity CA pool from %s/%s key %q: %v", egressNamespace, actorIDCASecret, actorIDCASecretKey, err)
t.Fatalf("parsing actor-identity CA pool from %s/%s key %q: %v", e2e.SystemNamespace(), actorIDCASecret, actorIDCASecretKey, err)
}
if len(pool.CAs) == 0 {
t.Fatalf("actor-identity CA pool %s/%s contains no CA", egressNamespace, actorIDCASecret)
t.Fatalf("actor-identity CA pool %s/%s contains no CA", e2e.SystemNamespace(), actorIDCASecret)
}
// CAs[0] is the one that signs: ateapi's MintCert makes the same choice.
return pool.CAs[0]
@@ -54,9 +54,6 @@ import (
)
const (
// Where the gateway's CA lives, fixed by hack/install-ate.sh.
egressNamespace = "ate-system"
probeName = "egressprobe"
)
@@ -163,6 +160,8 @@ func startProbe(t *testing.T, ctx context.Context) *probeClient {
}
manifest := filepath.Join(t.TempDir(), "egressprobe.yaml")
rendered := strings.ReplaceAll(string(tmpl), "${NAMESPACE}", ns)
rendered = strings.ReplaceAll(rendered, "${SYSTEM_NAMESPACE}", e2e.SystemNamespace())
rendered = strings.ReplaceAll(rendered, "${EGRESS_GATEWAY_SERVICE}", e2e.ResourceName("atenet-egress"))
if err := os.WriteFile(manifest, []byte(rendered), 0o644); err != nil {
t.Fatalf("writing rendered egressprobe manifest: %v", err)
}
@@ -311,7 +311,7 @@ func routerAddress(t *testing.T, ctx context.Context) string {
if err != nil {
t.Fatalf("creating k8s client: %v", err)
}
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, "ate-system", "atenet-router", 80)
localPort, stop, err := portforward.ServicePortForward(ctx, config, clientset, e2e.SystemNamespace(), e2e.ResourceName("atenet-router"), 80)
if err != nil {
t.Fatalf("port-forwarding to the router: %v", err)
}
@@ -304,9 +304,9 @@ type gatewayAccessLogLine struct {
// replica, until predicate accepts the lines written since.
func waitForAccessLog(t *testing.T, ctx context.Context, since metav1.Time, want string, predicate func(lines []gatewayAccessLogLine) bool) {
t.Helper()
gatewayNamespace := e2e.SystemNamespace()
const (
gatewayNamespace = "ate-system"
gatewaySelector = "app=atenet-egress"
gatewaySelector = "app=atenet-egress"
)
clients := e2e.GetClients()
@@ -36,7 +36,6 @@ import (
)
const (
ateSystemNamespace = "ate-system"
atenetRouterAppName = "atenet-router"
)
@@ -103,8 +102,8 @@ func TestNetworkPolicyLifecycleAndReconciliation(t *testing.T) {
t.Fatalf("expected exactly 1 ingress from peer, got %d", len(ingressRule.From))
}
fromPeer := ingressRule.From[0]
if fromPeer.NamespaceSelector == nil || fromPeer.NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"] != ateSystemNamespace {
t.Errorf("expected namespace selector for %s, got %v", ateSystemNamespace, fromPeer.NamespaceSelector)
if fromPeer.NamespaceSelector == nil || fromPeer.NamespaceSelector.MatchLabels["kubernetes.io/metadata.name"] != e2e.SystemNamespace() {
t.Errorf("expected namespace selector for %s, got %v", e2e.SystemNamespace(), fromPeer.NamespaceSelector)
}
if fromPeer.PodSelector == nil || fromPeer.PodSelector.MatchLabels["app"] != atenetRouterAppName {
t.Errorf("expected pod selector for %s, got %v", atenetRouterAppName, fromPeer.PodSelector)
+11 -8
View File
@@ -39,9 +39,12 @@ const (
// EgressTrustBundleObjectName is the reconciler-owned ClusterTrustBundle.
EgressTrustBundleObjectName = "egress-mitm.ate.dev:mitm:primary-bundle"
egressCAPoolNamespace = "ate-system"
// egressCAPoolSecretName is not release-prefixed: hack/install-ate.sh
// creates it under this fixed name and the reconciler looks it up the same
// way, so it does not follow the deployment's naming.
egressCAPoolSecretName = "egress-mitm-ca-pool"
egressCAPoolSecretKey = "pool"
egressCAPoolSecretKey = "pool"
)
// EnsureEgressTrustBundle makes sure the egress trust bundle exists, then
@@ -68,12 +71,12 @@ func ReplaceEgressTrustPool(t *testing.T, ctx context.Context, clients *Clients,
if !createEgressTrustPool(t, ctx, clients, secret) {
// Took over an existing pool: overwrite its contents without adopting
// its cleanup, since whoever created it registered one already.
existing, err := clients.K8s.CoreV1().Secrets(egressCAPoolNamespace).Get(ctx, egressCAPoolSecretName, metav1.GetOptions{})
existing, err := clients.K8s.CoreV1().Secrets(SystemNamespace()).Get(ctx, egressCAPoolSecretName, metav1.GetOptions{})
if err != nil {
t.Fatalf("reading existing CA pool secret: %v", err)
}
existing.Data = secret.Data
if _, err := clients.K8s.CoreV1().Secrets(egressCAPoolNamespace).Update(ctx, existing, metav1.UpdateOptions{}); err != nil {
if _, err := clients.K8s.CoreV1().Secrets(SystemNamespace()).Update(ctx, existing, metav1.UpdateOptions{}); err != nil {
t.Fatalf("updating CA pool secret: %v", err)
}
}
@@ -103,7 +106,7 @@ func newEgressTrustPool(t *testing.T) (*corev1.Secret, string) {
t.Fatalf("encoding the egress CA private key: %v", err)
}
return &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{Namespace: egressCAPoolNamespace, Name: egressCAPoolSecretName},
ObjectMeta: metav1.ObjectMeta{Namespace: SystemNamespace(), Name: egressCAPoolSecretName},
Type: corev1.SecretTypeTLS,
Data: map[string][]byte{
egressCAPoolSecretKey: poolBytes,
@@ -120,14 +123,14 @@ func newEgressTrustPool(t *testing.T) (*corev1.Secret, string) {
// behind, and no caller removes a pool it merely found.
func createEgressTrustPool(t *testing.T, ctx context.Context, clients *Clients, secret *corev1.Secret) bool {
t.Helper()
if _, err := clients.K8s.CoreV1().Secrets(egressCAPoolNamespace).Create(ctx, secret, metav1.CreateOptions{}); err != nil {
if _, err := clients.K8s.CoreV1().Secrets(SystemNamespace()).Create(ctx, secret, metav1.CreateOptions{}); err != nil {
if !apierrors.IsAlreadyExists(err) {
t.Fatalf("creating CA pool secret %s/%s: %v", egressCAPoolNamespace, egressCAPoolSecretName, err)
t.Fatalf("creating CA pool secret %s/%s: %v", SystemNamespace(), egressCAPoolSecretName, err)
}
return false
}
t.Cleanup(func() {
_ = clients.K8s.CoreV1().Secrets(egressCAPoolNamespace).Delete(context.Background(), egressCAPoolSecretName, metav1.DeleteOptions{})
_ = clients.K8s.CoreV1().Secrets(SystemNamespace()).Delete(context.Background(), egressCAPoolSecretName, metav1.DeleteOptions{})
})
return true
}
+143
View File
@@ -0,0 +1,143 @@
// 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 installdefaults
import (
"go/ast"
"go/parser"
"go/token"
"os"
"path/filepath"
"runtime"
"strconv"
"strings"
"testing"
)
// This guard exists because hardcoded install-layout assumptions fail closed
// in a relocated or renamed install, in ways that do not point at the naming
// (see the ateletauth, ateletdial and NetworkPolicy histories). New namespace
// or SPIFFE-identity literals belong in this package, in a flag default, or —
// deliberately — in the allowlist below.
// excludedTrees are directories whose tooling targets the canonical install
// layout by design and is not expected to work against a relocated one.
var excludedTrees = []string{
"cmd/ate-setup", // canonical GCP installer
"cmd/benchmarking", // canonical load-testing workers
"internal/benchmarking", // canonical load-testing harness
}
// allowedLiterals maps a repo-relative file to the canonical-layout string
// literals it may carry. Everything here is either a flag/env default that an
// operator overrides (the sanctioned pattern: default to the canonical
// layout, configure the rest), or a string nothing verifies against.
var allowedLiterals = map[string][]string{
// Flag defaults, overridden by the deployment or the e2e manifest templates.
"cmd/atecontroller/main.go": {"k8s:///api.ate-system.svc:443"},
"cmd/atelet/main.go": {"k8s:///api.ate-system.svc:443", "api.ate-system.svc"},
"cmd/atenet/internal/router/cmd.go": {"k8s:///api.ate-system.svc:443", "spiffe://cluster.local/"},
"internal/e2e/fixtures/testserver/egressprobe.go": {"atenet-egress.ate-system.svc:443"},
// Env-var defaults, overridden by ATE_* / E2E_* variables.
"internal/ateclient/builder.go": {"api.ate-system.svc"},
// Inert: the actor JWT issuer is an upstream TODO nothing validates, and
// an x509 template's Issuer field is overwritten by the signer.
"cmd/ateapi/internal/controlapi/actor.go": {"https://api.ate-system.svc", "api.ate-system.svc.cluster.local"},
}
// suspect reports whether a string literal encodes canonical install layout.
func suspect(s string) bool {
return strings.Contains(s, "ate-system") || strings.Contains(s, "spiffe://cluster.local")
}
func repoRoot(t *testing.T) string {
t.Helper()
_, thisFile, _, ok := runtime.Caller(0)
if !ok {
t.Fatal("cannot locate this test file")
}
return filepath.Dir(filepath.Dir(filepath.Dir(thisFile)))
}
func allowed(relPath, lit string) bool {
for _, tree := range excludedTrees {
if strings.HasPrefix(relPath, tree+string(filepath.Separator)) {
return true
}
}
for _, a := range allowedLiterals[filepath.ToSlash(relPath)] {
if a == lit {
return true
}
}
return false
}
// TestNoNewHardcodedInstallLayoutInGo walks every non-test Go file under cmd/
// and internal/ and fails on string literals that name the canonical
// namespace or a full SPIFFE identity outside this package. Comments are not
// scanned, so kubebuilder markers (which must name a literal namespace) pass.
func TestNoNewHardcodedInstallLayoutInGo(t *testing.T) {
root := repoRoot(t)
selfDir := filepath.Join("internal", "installdefaults")
var offenders []string
for _, tree := range []string{"cmd", "internal"} {
err := filepath.WalkDir(filepath.Join(root, tree), func(path string, d os.DirEntry, err error) error {
if err != nil {
return err
}
if d.IsDir() || !strings.HasSuffix(path, ".go") || strings.HasSuffix(path, "_test.go") {
return nil
}
rel, err := filepath.Rel(root, path)
if err != nil {
return err
}
if strings.HasPrefix(rel, selfDir) {
return nil
}
fset := token.NewFileSet()
f, err := parser.ParseFile(fset, path, nil, 0)
if err != nil {
return err
}
ast.Inspect(f, func(n ast.Node) bool {
lit, ok := n.(*ast.BasicLit)
if !ok || lit.Kind != token.STRING {
return true
}
val, err := strconv.Unquote(lit.Value)
if err != nil {
return true
}
if suspect(val) && !allowed(rel, val) {
offenders = append(offenders, fset.Position(lit.Pos()).String()+": "+lit.Value)
}
return true
})
return nil
})
if err != nil {
t.Fatal(err)
}
}
if len(offenders) > 0 {
t.Errorf("hardcoded install-layout strings outside internal/installdefaults; "+
"derive them from installdefaults or a configured value, or allowlist them "+
"in this test with a justification:\n %s", strings.Join(offenders, "\n "))
}
}
@@ -0,0 +1,89 @@
// 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 installdefaults holds the default namespace and Service names
// that match the canonical install layout in manifests/ate-install/.
// Binaries use these as flag defaults; deployments that diverge from
// the canonical layout pass actual values via the corresponding flags.
package installdefaults
import (
"net/url"
"os"
"path"
)
const (
// SystemNamespace is the namespace where substrate's control-plane
// components and the atelet DaemonSet run.
SystemNamespace = "ate-system"
// APIServiceName is the Service name of ate-api-server.
APIServiceName = "api"
// RouterServiceName is the Service name of atenet-router.
RouterServiceName = "atenet-router"
// ClientServiceAccount is the ServiceAccount an out-of-cluster client mints
// its ateapi bearer token from.
ClientServiceAccount = "ate-client"
// AteletTrustDomain and the ServiceAccount constants are the trust-domain
// and service-account segments of the SPIFFE IDs that atelet,
// atenet-router and atenet-egress Pod certificates carry, as minted by the
// podidentity signer (cmd/podcertcontroller/internal/podidentitysigner).
// The namespace segment is the namespace they run in, which callers
// resolve themselves rather than assume.
AteletTrustDomain = "cluster.local"
AteletServiceAccount = "atelet"
RouterServiceAccount = "atenet-router"
EgressServiceAccount = "atenet-egress"
// PodNamespaceEnv is the conventional env var name for the namespace
// a pod is running in, exposed via Kubernetes' downward API.
PodNamespaceEnv = "POD_NAMESPACE"
)
// NamespaceFromPodEnv returns the namespace from the PodNamespaceEnv env
// var when set (typically populated via Kubernetes' downward API), and
// falls back to SystemNamespace for non-k8s invocations (tests, local dev).
func NamespaceFromPodEnv() string {
if ns := os.Getenv(PodNamespaceEnv); ns != "" {
return ns
}
return SystemNamespace
}
// SPIFFEID returns the SPIFFE ID that Pod certificates for serviceAccount in
// namespace carry. Peers authenticate by comparing against this exact string.
func SPIFFEID(namespace, serviceAccount string) string {
return (&url.URL{
Scheme: "spiffe",
Host: AteletTrustDomain,
Path: path.Join("ns", namespace, "sa", serviceAccount),
}).String()
}
// AteletSPIFFEID returns the SPIFFE ID atelet presents when it runs in namespace.
func AteletSPIFFEID(namespace string) string {
return SPIFFEID(namespace, AteletServiceAccount)
}
// RouterSPIFFEID returns the SPIFFE ID atenet-router presents when it runs in namespace.
func RouterSPIFFEID(namespace string) string {
return SPIFFEID(namespace, RouterServiceAccount)
}
// EgressSPIFFEID returns the SPIFFE ID atenet-egress presents when it runs in
// namespace. The credential provider verifies it on the injector connection.
func EgressSPIFFEID(namespace string) string {
return SPIFFEID(namespace, EgressServiceAccount)
}
@@ -0,0 +1,81 @@
// 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 installdefaults
import "testing"
func TestAteletSPIFFEID(t *testing.T) {
tests := []struct {
name string
namespace string
want string
}{
{
// The canonical install. Peers reject any other string, so this
// value is effectively wire format: changing it breaks the atelet
// mTLS handshake for every existing deployment.
name: "default namespace",
namespace: SystemNamespace,
want: "spiffe://cluster.local/ns/ate-system/sa/atelet",
},
{
name: "namespace the install was relocated to",
namespace: "team-a-substrate",
want: "spiffe://cluster.local/ns/team-a-substrate/sa/atelet",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := AteletSPIFFEID(tt.namespace); got != tt.want {
t.Errorf("AteletSPIFFEID(%q) = %q, want %q", tt.namespace, got, tt.want)
}
})
}
}
func TestRouterSPIFFEID(t *testing.T) {
// Matches the --atunnel-client-identity default the ateom binaries ship
// with, which is what actor ingress authenticates the router against.
const want = "spiffe://cluster.local/ns/ate-system/sa/atenet-router"
if got := RouterSPIFFEID(SystemNamespace); got != want {
t.Errorf("RouterSPIFFEID(%q) = %q, want %q", SystemNamespace, got, want)
}
}
func TestEgressSPIFFEID(t *testing.T) {
// Matches the --injector-identity default the credential provider ships
// with, which is what it authenticates the egress injector against.
const want = "spiffe://cluster.local/ns/ate-system/sa/atenet-egress"
if got := EgressSPIFFEID(SystemNamespace); got != want {
t.Errorf("EgressSPIFFEID(%q) = %q, want %q", SystemNamespace, got, want)
}
}
func TestNamespaceFromPodEnv(t *testing.T) {
t.Run("falls back to the install default when unset", func(t *testing.T) {
t.Setenv(PodNamespaceEnv, "")
if got := NamespaceFromPodEnv(); got != SystemNamespace {
t.Errorf("NamespaceFromPodEnv() = %q, want %q", got, SystemNamespace)
}
})
t.Run("prefers the downward API value", func(t *testing.T) {
t.Setenv(PodNamespaceEnv, "team-a-substrate")
if got := NamespaceFromPodEnv(); got != "team-a-substrate" {
t.Errorf("NamespaceFromPodEnv() = %q, want %q", got, "team-a-substrate")
}
})
}