From 92b480b36efef34e5eb05d3d0e99dad9e95e77c8 Mon Sep 17 00:00:00 2001 From: Thomas Kosiewski Date: Fri, 25 Sep 2026 09:29:26 +0000 Subject: [PATCH 1/2] feat(apiserver): serve a CA-signed certificate from a managed Secret MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Refs #137 In a Pod, the aggregated API server now serves a certificate signed by its own CA instead of an in-memory self-signed localhost certificate. This is the first half of #137; the APIService still uses insecureSkipTLSVerify, and registering the CA in caBundle follows. - New internal/aggregated/servingcert: the CA and serving certificate live in Secret coder-k8s-apiserver-tls (type coder.com/aggregated-apiserver-serving-ca, coder-k8s labels) in the pod's namespace. Generation uses client-go certutil/keyutil; the serving cert covers coder-k8s-apiserver, ., ..svc and ..svc.cluster.local. - The Manager implements dynamiccertificates.CertKeyContentProvider and is set as SecureServingOptions.ServerCert.GeneratedCert, so the vendored DynamicServingCertificateController hot-swaps renewals. Listeners are notified on change (hook for the caBundle controller). - Get-or-create adopts a Secret created by another replica; renewal (less than a third of the 1-year lifetime left, or SANs for another namespace) keeps the CA and adopts on update conflicts. Checked at startup and every 12h. - An unusable Secret (wrong type, missing key, unparsable PEM, key and certificate mismatch, serving certificate not signed by the CA, expired or not-yet-valid CA) fails startup with a message naming the field; it is never overwritten. - Outside a Pod the old self-signed localhost certificate is kept. - Docs: serving certificate, rotation, CA-key warning, troubleshooting. --- _Generated with `xum` • Model: `anthropic:claude-opus-5-5` • Thinking: `high`_ --- docs/how-to/deploy-aggregated-apiserver.md | 17 +- docs/how-to/troubleshooting.md | 11 + internal/aggregated/servingcert/manager.go | 191 +++++++ .../aggregated/servingcert/servingcert.go | 283 +++++++++++ .../servingcert/servingcert_test.go | 474 ++++++++++++++++++ internal/app/apiserverapp/apiserverapp.go | 31 +- internal/app/apiserverapp/servingcert.go | 51 ++ internal/app/apiserverapp/servingcert_test.go | 202 ++++++++ 8 files changed, 1256 insertions(+), 4 deletions(-) create mode 100644 internal/aggregated/servingcert/manager.go create mode 100644 internal/aggregated/servingcert/servingcert.go create mode 100644 internal/aggregated/servingcert/servingcert_test.go create mode 100644 internal/app/apiserverapp/servingcert.go create mode 100644 internal/app/apiserverapp/servingcert_test.go diff --git a/docs/how-to/deploy-aggregated-apiserver.md b/docs/how-to/deploy-aggregated-apiserver.md index e37a8daf..e17cd94f 100644 --- a/docs/how-to/deploy-aggregated-apiserver.md +++ b/docs/how-to/deploy-aggregated-apiserver.md @@ -112,5 +112,18 @@ If something fails, check `kubectl logs -n coder-system deploy/coder-k8s` and [T These resources are backed by Coder, not etcd, so some Kubernetes behavior differs. Read [Aggregated API behavior](../reference/aggregated-api-behavior.md) before you write manifests. The most important rule: object names must use Coder's canonical names. -!!! warning "TLS" - `deploy/apiserver-apiservice.yaml` sets `insecureSkipTLSVerify: true` for development. Use CA-backed TLS in any real environment. +## Serving certificate + +In a cluster, the aggregated API server serves a certificate signed by its own CA. Both live in the Secret `coder-k8s-apiserver-tls` in the server's namespace (type `coder.com/aggregated-apiserver-serving-ca`, label `app.kubernetes.io/component: aggregated-apiserver-serving-ca`). The certificate is valid for `coder-k8s-apiserver`, `coder-k8s-apiserver.`, `coder-k8s-apiserver..svc`, and `coder-k8s-apiserver..svc.cluster.local`. + +- The server creates the Secret on first start and reuses it afterwards. With several replicas, they all use the same Secret. +- The serving certificate is valid for 1 year. The server checks it at startup and every 12 hours, and renews it with the same CA when less than a third of its lifetime is left. The new certificate is served without a restart. +- The CA is valid for 10 years. To replace it, delete the Secret and restart the Deployment (`kubectl -n coder-system rollout restart deployment/coder-k8s`). Clients that trusted the old CA must then trust the new one. +- If the Secret exists but is unusable (a missing key, unparsable PEM, a key that does not match its certificate, a serving certificate not signed by the CA, or an expired CA), the server does not start and the log names the field. Fix the Secret or delete it. +- Outside a cluster (for example `go run`), the server serves a self-signed certificate for `localhost` instead. + +!!! warning "The Secret holds the CA private key" + Anyone who can read Secrets in the server's namespace can issue certificates that the aggregated API server's CA vouches for. Restrict Secret read access in `coder-system` accordingly. + +!!! warning "TLS verification is still off" + `deploy/apiserver-apiservice.yaml` still sets `insecureSkipTLSVerify: true`, so kube-apiserver does not check this certificate yet. Registering the CA in the APIService `caBundle` is tracked in [#137](https://github.com/coder/coder-k8s/issues/137). diff --git a/docs/how-to/troubleshooting.md b/docs/how-to/troubleshooting.md index 03f21089..53ca5af1 100644 --- a/docs/how-to/troubleshooting.md +++ b/docs/how-to/troubleshooting.md @@ -72,6 +72,17 @@ The aggregated API server checks every caller with the Kubernetes API and refuse - `no Kubernetes configuration for delegated authentication and authorization` (outside a cluster): set `KUBECONFIG` to one kubeconfig file, or create `~/.kube/config`. - `load kubeconfig ...` or `invalid kubeconfig ...`: the file named by `KUBECONFIG` is missing or incomplete. The server does not fall back to another configuration. +## The pod exits with `configure aggregated API server serving certificate` + +The Secret `coder-k8s-apiserver-tls` exists but cannot be used; the message names the field (for example `data["ca.key"] is missing or empty`). Fix the Secret, or delete it so the server generates a new CA on its next start: + +```bash +kubectl -n coder-system delete secret coder-k8s-apiserver-tls +kubectl -n coder-system rollout restart deployment/coder-k8s +``` + +A read or create error instead of a field name means the ServiceAccount cannot get or create Secrets in its namespace. + ## Aggregated requests fail with `401 Unauthorized` or `403 Forbidden` - **`401`:** the request has no valid credential. Requests sent straight to port `6443` need a Kubernetes bearer token; anonymous requests only reach `/healthz`, `/livez`, and `/readyz`. Use `kubectl`, which goes through kube-apiserver. diff --git a/internal/aggregated/servingcert/manager.go b/internal/aggregated/servingcert/manager.go new file mode 100644 index 00000000..9fe36fbe --- /dev/null +++ b/internal/aggregated/servingcert/manager.go @@ -0,0 +1,191 @@ +package servingcert + +import ( + "context" + "crypto/tls" + "fmt" + "sync" + "time" + + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apiserver/pkg/server/dynamiccertificates" + "k8s.io/client-go/kubernetes" + "k8s.io/klog/v2" +) + +// DefaultCheckInterval is how often Run re-reads the Secret and renews the serving certificate. +const DefaultCheckInterval = 12 * time.Hour + +const maxEnsureAttempts = 5 + +// Manager keeps the serving certificate in sync with the Secret. +// +// Ensure is the adopt path: any valid Secret is used as is (so a later "bring your own Secret" +// mode only has to skip generation and renewal). Listeners registered through +// CertKeyContentProvider().AddListener are notified whenever the served certificate or the CA +// changes, so an APIService caBundle controller can react without polling. +type Manager struct { + client kubernetes.Interface + namespace string + now func() time.Time + + mu sync.RWMutex + current *Bundle + listeners []dynamiccertificates.Listener +} + +var _ dynamiccertificates.CertKeyContentProvider = (*Manager)(nil) + +// NewManager returns a Manager for the Secret in namespace. +func NewManager(client kubernetes.Interface, namespace string) (*Manager, error) { + if client == nil { + return nil, fmt.Errorf("assertion failed: Kubernetes client must not be nil") + } + if namespace == "" { + return nil, fmt.Errorf("assertion failed: namespace must not be empty") + } + return &Manager{client: client, namespace: namespace, now: time.Now}, nil +} + +// Ensure loads, creates, or renews the Secret and makes it the served certificate. Invalid +// trust material returns a *CorruptSecretError and is never overwritten. +func (m *Manager) Ensure(ctx context.Context) (*Bundle, error) { + if ctx == nil { + return nil, fmt.Errorf("assertion failed: context must not be nil") + } + secrets := m.client.CoreV1().Secrets(m.namespace) + for attempt := 1; attempt <= maxEnsureAttempts; attempt++ { + now := m.now() + secret, err := secrets.Get(ctx, SecretName, metav1.GetOptions{}) + switch { + case apierrors.IsNotFound(err): + bundle, genErr := Generate(m.namespace, now) + if genErr != nil { + return nil, genErr + } + _, err = secrets.Create(ctx, &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: SecretName, Namespace: m.namespace, Labels: copyLabels()}, + Type: SecretType, + Data: bundle.Data(), + }, metav1.CreateOptions{}) + if apierrors.IsAlreadyExists(err) { + continue // Another replica created it first; adopt theirs. + } + if err != nil { + return nil, fmt.Errorf("create secret %s/%s: %w", m.namespace, SecretName, err) + } + klog.InfoS("Created aggregated API server CA and serving certificate", "secret", klog.KRef(m.namespace, SecretName)) + return bundle, m.serve(bundle) + case err != nil: + return nil, fmt.Errorf("get secret %s/%s: %w", m.namespace, SecretName, err) + } + + bundle, err := Parse(secret, m.namespace, now) + if err != nil { + return nil, err + } + if !bundle.NeedsRenewal(m.namespace, now) { + return bundle, m.serve(bundle) + } + if err := bundle.issueServingCert(m.namespace, now); err != nil { + return nil, err + } + updated := secret.DeepCopy() + updated.Data = bundle.Data() + _, err = secrets.Update(ctx, updated, metav1.UpdateOptions{}) + if apierrors.IsConflict(err) { + continue // Another replica renewed it; re-read and adopt. + } + if err != nil { + return nil, fmt.Errorf("update secret %s/%s: %w", m.namespace, SecretName, err) + } + klog.InfoS("Renewed aggregated API server serving certificate", "secret", klog.KRef(m.namespace, SecretName), "notAfter", bundle.Cert.NotAfter) + return bundle, m.serve(bundle) + } + return nil, fmt.Errorf("secret %s/%s kept changing; gave up after %d attempts", m.namespace, SecretName, maxEnsureAttempts) +} + +// Run calls Ensure every interval until ctx is done. Failures are logged and the current +// certificate keeps being served. +func (m *Manager) Run(ctx context.Context, interval time.Duration) { + if interval <= 0 { + panic("assertion failed: interval must be positive") + } + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + if _, err := m.Ensure(ctx); err != nil { + klog.ErrorS(err, "Could not refresh the aggregated API server serving certificate; still serving the current one") + } + } + } +} + +// CABundle returns the PEM CA that signs the served certificate, or nil before the first Ensure. +func (m *Manager) CABundle() []byte { + m.mu.RLock() + defer m.mu.RUnlock() + if m.current == nil { + return nil + } + return m.current.CACertPEM +} + +// serve makes bundle the served certificate and notifies listeners if anything changed. +func (m *Manager) serve(bundle *Bundle) error { + // Same check the vendored static provider performs. + if _, err := tls.X509KeyPair(bundle.CertPEM, bundle.KeyPEM); err != nil { + return fmt.Errorf("assertion failed: validated serving certificate is not a usable key pair: %w", err) + } + m.mu.Lock() + changed := m.current == nil || !bundleEqual(m.current, bundle) + m.current = bundle + listeners := append([]dynamiccertificates.Listener(nil), m.listeners...) + m.mu.Unlock() + if changed { + for _, l := range listeners { + l.Enqueue() + } + } + return nil +} + +// Name implements dynamiccertificates.CertKeyContentProvider. +func (m *Manager) Name() string { + return "coder-k8s-managed-serving-cert::" + m.namespace + "/" + SecretName +} + +// CurrentCertKeyContent implements dynamiccertificates.CertKeyContentProvider. +func (m *Manager) CurrentCertKeyContent() ([]byte, []byte) { + m.mu.RLock() + defer m.mu.RUnlock() + if m.current == nil { + return nil, nil + } + return m.current.CertPEM, m.current.KeyPEM +} + +// AddListener implements dynamiccertificates.Notifier. +func (m *Manager) AddListener(listener dynamiccertificates.Listener) { + m.mu.Lock() + defer m.mu.Unlock() + m.listeners = append(m.listeners, listener) +} + +func bundleEqual(a, b *Bundle) bool { + return string(a.CACertPEM) == string(b.CACertPEM) && string(a.CertPEM) == string(b.CertPEM) && string(a.KeyPEM) == string(b.KeyPEM) +} + +func copyLabels() map[string]string { + out := make(map[string]string, len(SecretLabels)) + for k, v := range SecretLabels { + out[k] = v + } + return out +} diff --git a/internal/aggregated/servingcert/servingcert.go b/internal/aggregated/servingcert/servingcert.go new file mode 100644 index 00000000..d50aeb07 --- /dev/null +++ b/internal/aggregated/servingcert/servingcert.go @@ -0,0 +1,283 @@ +// Package servingcert manages the aggregated API server's serving certificate. +// +// The server keeps a private CA and a CA-signed serving certificate in one Secret in its own +// namespace. The serving certificate is served through a dynamiccertificates.CertKeyContentProvider, +// so the generic API server's DynamicServingCertificateController picks up renewals without a +// restart. The CA is what the APIService caBundle will trust (#137). +package servingcert + +import ( + "bytes" + "crypto" + "crypto/ecdsa" + "crypto/elliptic" + "crypto/rand" + "crypto/x509" + "crypto/x509/pkix" + "fmt" + "math" + "math/big" + "time" + + corev1 "k8s.io/api/core/v1" + certutil "k8s.io/client-go/util/cert" + "k8s.io/client-go/util/keyutil" +) + +const ( + // ServiceName is the Service that fronts the aggregated API server (deploy/apiserver-service.yaml). + ServiceName = "coder-k8s-apiserver" + // SecretName is the Secret that holds the CA and the serving certificate. + SecretName = "coder-k8s-apiserver-tls" //nolint:gosec // G101: a resource name, not a credential. + // SecretType marks the Secret as coder-k8s's aggregated API server CA. It is deliberately not + // kubernetes.io/tls: the Secret also holds the CA private key. + SecretType corev1.SecretType = "coder.com/aggregated-apiserver-serving-ca" //nolint:gosec // G101: a type name, not a credential. + + // CACertKey is the Secret data key for the PEM CA certificate. + CACertKey = "ca.crt" + // CAKeyKey is the Secret data key for the PEM CA private key. + CAKeyKey = "ca.key" + // CertKey is the Secret data key for the PEM serving certificate. + CertKey = corev1.TLSCertKey + // KeyKey is the Secret data key for the PEM serving private key. + KeyKey = corev1.TLSPrivateKeyKey + + // ServingCertValidity is the lifetime of a serving certificate. The CA lives 10 years + // (client-go certutil.NewSelfSignedCACert). + ServingCertValidity = 365 * 24 * time.Hour + // clockSkewAllowance backdates serving certificates so small clock differences do not reject them. + clockSkewAllowance = time.Hour +) + +// SecretLabels mark the Secret as managed by coder-k8s. +var SecretLabels = map[string]string{ + "app.kubernetes.io/name": "coder-k8s", + "app.kubernetes.io/component": "aggregated-apiserver-serving-ca", + "app.kubernetes.io/managed-by": "coder-k8s", +} + +// DNSNames returns the serving certificate's subject alternative names for the Service in namespace. +// kube-apiserver verifies the aggregated server with ServerName "..svc". +func DNSNames(namespace string) []string { + return []string{ + ServiceName, + ServiceName + "." + namespace, + ServiceName + "." + namespace + ".svc", + ServiceName + "." + namespace + ".svc.cluster.local", + } +} + +// Bundle is parsed, validated trust material from the Secret. +type Bundle struct { + CACertPEM []byte + CAKeyPEM []byte + CertPEM []byte + KeyPEM []byte + + CACert *x509.Certificate + Cert *x509.Certificate + caKey crypto.Signer +} + +// Data returns the Secret data for the bundle. +func (b *Bundle) Data() map[string][]byte { + return map[string][]byte{CACertKey: b.CACertPEM, CAKeyKey: b.CAKeyPEM, CertKey: b.CertPEM, KeyKey: b.KeyPEM} +} + +// CorruptSecretError reports trust material that cannot be used. The server never replaces such a +// Secret on its own, because clients may already trust its CA. +type CorruptSecretError struct { + Namespace string + Problem string +} + +func (e *CorruptSecretError) Error() string { + return fmt.Sprintf("secret %s/%s: %s; fix it, or delete it so coder-k8s generates a new CA (clients that trust the old CA must then be updated)", + e.Namespace, SecretName, e.Problem) +} + +func corrupt(namespace, format string, args ...any) error { + return &CorruptSecretError{Namespace: namespace, Problem: fmt.Sprintf(format, args...)} +} + +// Generate creates a new CA and a serving certificate for namespace. +func Generate(namespace string, now time.Time) (*Bundle, error) { + if namespace == "" { + return nil, fmt.Errorf("assertion failed: namespace must not be empty") + } + caKeyPEM, err := keyutil.MakeEllipticPrivateKeyPEM() + if err != nil { + return nil, fmt.Errorf("generate CA key: %w", err) + } + caKey, err := parseSigner(caKeyPEM) + if err != nil { + return nil, fmt.Errorf("assertion failed: parse generated CA key: %w", err) + } + caCert, err := certutil.NewSelfSignedCACert(certutil.Config{ + CommonName: fmt.Sprintf("%s-ca@%d", ServiceName, now.Unix()), + Organization: []string{"coder-k8s"}, + }, caKey) + if err != nil { + return nil, fmt.Errorf("generate CA certificate: %w", err) + } + caCertPEM, err := certutil.EncodeCertificates(caCert) + if err != nil { + return nil, fmt.Errorf("encode CA certificate: %w", err) + } + b := &Bundle{CACertPEM: caCertPEM, CAKeyPEM: caKeyPEM, CACert: caCert, caKey: caKey} + if err := b.issueServingCert(namespace, now); err != nil { + return nil, err + } + return b, nil +} + +// issueServingCert signs a new serving certificate with the bundle's CA. +func (b *Bundle) issueServingCert(namespace string, now time.Time) error { + if b.CACert == nil || b.caKey == nil { + return fmt.Errorf("assertion failed: bundle has no CA") + } + key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader) + if err != nil { + return fmt.Errorf("generate serving key: %w", err) + } + serial, err := rand.Int(rand.Reader, new(big.Int).SetInt64(math.MaxInt64-1)) + if err != nil { + return fmt.Errorf("generate serial: %w", err) + } + names := DNSNames(namespace) + tmpl := &x509.Certificate{ + SerialNumber: new(big.Int).Add(serial, big.NewInt(1)), + Subject: pkix.Name{CommonName: names[2]}, + DNSNames: names, + NotBefore: now.Add(-clockSkewAllowance).UTC(), + NotAfter: now.Add(ServingCertValidity).UTC(), + KeyUsage: x509.KeyUsageDigitalSignature, + ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}, + } + if tmpl.NotAfter.After(b.CACert.NotAfter) { + tmpl.NotAfter = b.CACert.NotAfter + } + der, err := x509.CreateCertificate(rand.Reader, tmpl, b.CACert, key.Public(), b.caKey) + if err != nil { + return fmt.Errorf("sign serving certificate: %w", err) + } + cert, err := x509.ParseCertificate(der) + if err != nil { + return fmt.Errorf("assertion failed: parse signed serving certificate: %w", err) + } + certPEM, err := certutil.EncodeCertificates(cert) + if err != nil { + return fmt.Errorf("encode serving certificate: %w", err) + } + keyPEM, err := keyutil.MarshalPrivateKeyToPEM(key) + if err != nil { + return fmt.Errorf("encode serving key: %w", err) + } + b.Cert, b.CertPEM, b.KeyPEM = cert, certPEM, keyPEM + return nil +} + +// Parse validates the Secret's trust material. Each failure names the field that is wrong. +func Parse(secret *corev1.Secret, namespace string, now time.Time) (*Bundle, error) { + if secret == nil { + return nil, fmt.Errorf("assertion failed: secret must not be nil") + } + if secret.Type != SecretType { + return nil, corrupt(namespace, "type is %q, want %q", secret.Type, SecretType) + } + for _, key := range []string{CACertKey, CAKeyKey, CertKey, KeyKey} { + if len(secret.Data[key]) == 0 { + return nil, corrupt(namespace, "data[%q] is missing or empty", key) + } + } + b := &Bundle{ + CACertPEM: secret.Data[CACertKey], CAKeyPEM: secret.Data[CAKeyKey], + CertPEM: secret.Data[CertKey], KeyPEM: secret.Data[KeyKey], + } + + var err error + if b.CACert, err = parseSingleCert(b.CACertPEM); err != nil { + return nil, corrupt(namespace, "data[%q]: %v", CACertKey, err) + } + if !b.CACert.IsCA || b.CACert.KeyUsage&x509.KeyUsageCertSign == 0 { + return nil, corrupt(namespace, "data[%q] is not a CA certificate", CACertKey) + } + if now.Before(b.CACert.NotBefore) { + return nil, corrupt(namespace, "the CA in data[%q] is not valid until %s", CACertKey, b.CACert.NotBefore.UTC().Format(time.RFC3339)) + } + if !now.Before(b.CACert.NotAfter) { + return nil, corrupt(namespace, "the CA in data[%q] expired at %s", CACertKey, b.CACert.NotAfter.UTC().Format(time.RFC3339)) + } + if b.caKey, err = parseSigner(b.CAKeyPEM); err != nil { + return nil, corrupt(namespace, "data[%q]: %v", CAKeyKey, err) + } + if !publicKeysEqual(b.caKey.Public(), b.CACert.PublicKey) { + return nil, corrupt(namespace, "data[%q] does not match the certificate in data[%q]", CAKeyKey, CACertKey) + } + + if b.Cert, err = parseSingleCert(b.CertPEM); err != nil { + return nil, corrupt(namespace, "data[%q]: %v", CertKey, err) + } + if err := b.Cert.CheckSignatureFrom(b.CACert); err != nil { + return nil, corrupt(namespace, "data[%q] is not signed by the CA in data[%q]: %v", CertKey, CACertKey, err) + } + servingKey, err := parseSigner(b.KeyPEM) + if err != nil { + return nil, corrupt(namespace, "data[%q]: %v", KeyKey, err) + } + if !publicKeysEqual(servingKey.Public(), b.Cert.PublicKey) { + return nil, corrupt(namespace, "data[%q] does not match the certificate in data[%q]", KeyKey, CertKey) + } + return b, nil +} + +// NeedsRenewal reports whether the serving certificate must be re-issued: less than a third of +// its lifetime is left, it is not yet valid or expired, or it does not cover the Service names of +// namespace (for example after the Secret was copied from another namespace). +func (b *Bundle) NeedsRenewal(namespace string, now time.Time) bool { + lifetime := b.Cert.NotAfter.Sub(b.Cert.NotBefore) + renewAt := b.Cert.NotBefore.Add(lifetime * 2 / 3) + if now.Before(b.Cert.NotBefore) || !now.Before(renewAt) { + return true + } + roots := x509.NewCertPool() + roots.AddCert(b.CACert) + for _, name := range DNSNames(namespace) { + if _, err := b.Cert.Verify(x509.VerifyOptions{ + DNSName: name, Roots: roots, CurrentTime: now, + KeyUsages: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}, + }); err != nil { + return true + } + } + return false +} + +func parseSingleCert(data []byte) (*x509.Certificate, error) { + certs, err := certutil.ParseCertsPEM(data) + if err != nil { + return nil, fmt.Errorf("unparsable PEM certificate: %w", err) + } + if len(certs) != 1 { + return nil, fmt.Errorf("holds %d certificates, want exactly 1", len(certs)) + } + return certs[0], nil +} + +func parseSigner(data []byte) (crypto.Signer, error) { + key, err := keyutil.ParsePrivateKeyPEM(data) + if err != nil { + return nil, fmt.Errorf("unparsable PEM private key: %w", err) + } + signer, ok := key.(crypto.Signer) + if !ok { + return nil, fmt.Errorf("private key of type %T cannot sign", key) + } + return signer, nil +} + +func publicKeysEqual(a, b crypto.PublicKey) bool { + aDER, errA := x509.MarshalPKIXPublicKey(a) + bDER, errB := x509.MarshalPKIXPublicKey(b) + return errA == nil && errB == nil && bytes.Equal(aDER, bDER) +} diff --git a/internal/aggregated/servingcert/servingcert_test.go b/internal/aggregated/servingcert/servingcert_test.go new file mode 100644 index 00000000..23795c72 --- /dev/null +++ b/internal/aggregated/servingcert/servingcert_test.go @@ -0,0 +1,474 @@ +package servingcert + +import ( + "context" + "crypto/ecdsa" + "crypto/elliptic" + "crypto/rand" + "crypto/x509" + "crypto/x509/pkix" + "errors" + "math/big" + "slices" + "strings" + "sync/atomic" + "testing" + "time" + + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/kubernetes/fake" + k8stesting "k8s.io/client-go/testing" + certutil "k8s.io/client-go/util/cert" + "k8s.io/client-go/util/keyutil" +) + +const testNS = "coder-system" + +func TestDNSNames(t *testing.T) { + want := []string{ + "coder-k8s-apiserver", + "coder-k8s-apiserver.coder-system", + "coder-k8s-apiserver.coder-system.svc", + "coder-k8s-apiserver.coder-system.svc.cluster.local", + } + if got := DNSNames(testNS); !slices.Equal(got, want) { + t.Fatalf("DNSNames = %v, want %v", got, want) + } +} + +func TestGenerateProducesVerifiableServingCert(t *testing.T) { + now := time.Now() + b, err := Generate(testNS, now) + if err != nil { + t.Fatal(err) + } + if !b.CACert.IsCA || b.CACert.KeyUsage&x509.KeyUsageCertSign == 0 { + t.Fatal("CA certificate must be a CA with KeyUsageCertSign") + } + if b.Cert.IsCA { + t.Fatal("serving certificate must not be a CA") + } + if !slices.Equal(b.Cert.ExtKeyUsage, []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}) { + t.Fatalf("serving ExtKeyUsage = %v, want [ServerAuth]", b.Cert.ExtKeyUsage) + } + if !slices.Equal(b.Cert.DNSNames, DNSNames(testNS)) { + t.Fatalf("serving SANs = %v", b.Cert.DNSNames) + } + if got := b.Cert.NotAfter.Sub(now); got < ServingCertValidity-time.Minute || got > ServingCertValidity+time.Minute { + t.Fatalf("serving certificate validity %s, want about %s", got, ServingCertValidity) + } + roots := x509.NewCertPool() + roots.AddCert(b.CACert) + for _, name := range DNSNames(testNS) { + if _, err := b.Cert.Verify(x509.VerifyOptions{DNSName: name, Roots: roots, KeyUsages: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}}); err != nil { + t.Errorf("verify for %s: %v", name, err) + } + } + for _, name := range []string{"coder-k8s-apiserver.other.svc", "localhost", "kubernetes.default.svc"} { + if _, err := b.Cert.Verify(x509.VerifyOptions{DNSName: name, Roots: roots}); err == nil { + t.Errorf("serving certificate must not verify for %s", name) + } + } + if _, err := Parse(secretFor(b), testNS, now); err != nil { + t.Fatalf("generated bundle does not parse: %v", err) + } +} + +func TestEnsureCreatesMarkedSecretAndServesIt(t *testing.T) { + client := fake.NewClientset() + m := newTestManager(t, client, time.Now()) + listener := &countingListener{} + m.AddListener(listener) + + b, err := m.Ensure(t.Context()) + if err != nil { + t.Fatal(err) + } + secret, err := client.CoreV1().Secrets(testNS).Get(t.Context(), SecretName, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + if secret.Type != SecretType { + t.Fatalf("secret type = %q, want %q", secret.Type, SecretType) + } + for k, v := range SecretLabels { + if secret.Labels[k] != v { + t.Fatalf("secret label %s = %q, want %q", k, secret.Labels[k], v) + } + } + for _, key := range []string{CACertKey, CAKeyKey, CertKey, KeyKey} { + if len(secret.Data[key]) == 0 { + t.Fatalf("secret data[%q] missing", key) + } + } + cert, key := m.CurrentCertKeyContent() + if string(cert) != string(secret.Data[CertKey]) || string(key) != string(secret.Data[KeyKey]) { + t.Fatal("served certificate must be the one stored in the Secret") + } + if string(m.CABundle()) != string(b.CACertPEM) { + t.Fatal("CABundle must return the Secret's CA") + } + if listener.count.Load() != 1 { + t.Fatalf("listener notified %d times, want 1", listener.count.Load()) + } +} + +func TestEnsureAdoptsValidSecretUnchanged(t *testing.T) { + now := time.Now() + existing, err := Generate(testNS, now.Add(-time.Hour)) + if err != nil { + t.Fatal(err) + } + client := fake.NewClientset(secretFor(existing)) + m := newTestManager(t, client, now) + listener := &countingListener{} + m.AddListener(listener) + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + cert, _ := m.CurrentCertKeyContent() + if string(cert) != string(existing.CertPEM) { + t.Fatal("a valid Secret must be served as is") + } + assertNoWrites(t, client) + // A second Ensure with nothing changed must not notify listeners again. + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + if listener.count.Load() != 1 { + t.Fatalf("listener notified %d times, want 1", listener.count.Load()) + } +} + +func TestEnsureAdoptsSecretCreatedByAnotherReplica(t *testing.T) { + now := time.Now() + theirs, err := Generate(testNS, now) + if err != nil { + t.Fatal(err) + } + client := fake.NewClientset() + var gets atomic.Int32 + client.PrependReactor("get", "secrets", func(k8stesting.Action) (bool, runtime.Object, error) { + if gets.Add(1) == 1 { + return true, nil, apierrors.NewNotFound(schema.GroupResource{Resource: "secrets"}, SecretName) + } + return true, secretFor(theirs), nil + }) + client.PrependReactor("create", "secrets", func(k8stesting.Action) (bool, runtime.Object, error) { + return true, nil, apierrors.NewAlreadyExists(schema.GroupResource{Resource: "secrets"}, SecretName) + }) + m := newTestManager(t, client, now) + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + if string(m.CABundle()) != string(theirs.CACertPEM) { + t.Fatal("must adopt the other replica's CA after AlreadyExists") + } +} + +func TestEnsureRenewsServingCertWithSameCA(t *testing.T) { + issued := time.Now() + existing, err := Generate(testNS, issued) + if err != nil { + t.Fatal(err) + } + client := fake.NewClientset(secretFor(existing)) + // Two thirds of the lifetime (measured from the backdated NotBefore) have passed. + m := newTestManager(t, client, issued.Add(250*24*time.Hour)) + listener := &countingListener{} + m.AddListener(listener) + b, err := m.Ensure(t.Context()) + if err != nil { + t.Fatal(err) + } + if string(b.CACertPEM) != string(existing.CACertPEM) || string(b.CAKeyPEM) != string(existing.CAKeyPEM) { + t.Fatal("renewal must keep the CA") + } + if b.Cert.SerialNumber.Cmp(existing.Cert.SerialNumber) == 0 { + t.Fatal("renewal must issue a new serving certificate") + } + stored, err := client.CoreV1().Secrets(testNS).Get(t.Context(), SecretName, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + if string(stored.Data[CertKey]) != string(b.CertPEM) { + t.Fatal("renewed certificate must be written to the Secret") + } + if listener.count.Load() != 1 { + t.Fatalf("listener notified %d times, want 1", listener.count.Load()) + } +} + +func TestEnsureDoesNotRenewEarly(t *testing.T) { + issued := time.Now() + existing, err := Generate(testNS, issued) + if err != nil { + t.Fatal(err) + } + client := fake.NewClientset(secretFor(existing)) + m := newTestManager(t, client, issued.Add(200*24*time.Hour)) + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + assertNoWrites(t, client) +} + +func TestEnsureRenewsWhenSANsDoNotMatchNamespace(t *testing.T) { + now := time.Now() + copied, err := Generate("another-namespace", now) + if err != nil { + t.Fatal(err) + } + client := fake.NewClientset(secretFor(copied)) + m := newTestManager(t, client, now) + b, err := m.Ensure(t.Context()) + if err != nil { + t.Fatal(err) + } + if !slices.Equal(b.Cert.DNSNames, DNSNames(testNS)) { + t.Fatalf("renewed SANs = %v, want %v", b.Cert.DNSNames, DNSNames(testNS)) + } +} + +func TestEnsureAdoptsAfterRenewalConflict(t *testing.T) { + issued := time.Now() + existing, err := Generate(testNS, issued) + if err != nil { + t.Fatal(err) + } + now := issued.Add(250 * 24 * time.Hour) + // The other replica already renewed: its certificate is fresh relative to now. + theirs := &Bundle{CACertPEM: existing.CACertPEM, CAKeyPEM: existing.CAKeyPEM, CACert: existing.CACert, caKey: existing.caKey} + if err := theirs.issueServingCert(testNS, now); err != nil { + t.Fatal(err) + } + client := fake.NewClientset(secretFor(existing)) + var updates atomic.Int32 + client.PrependReactor("update", "secrets", func(k8stesting.Action) (bool, runtime.Object, error) { + updates.Add(1) + if err := client.Tracker().Update(schema.GroupVersionResource{Version: "v1", Resource: "secrets"}, secretFor(theirs), testNS); err != nil { + t.Fatal(err) + } + return true, nil, apierrors.NewConflict(schema.GroupResource{Resource: "secrets"}, SecretName, errors.New("modified")) + }) + m := newTestManager(t, client, now) + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + cert, _ := m.CurrentCertKeyContent() + if string(cert) != string(theirs.CertPEM) { + t.Fatal("after a conflict the other replica's renewal must be adopted") + } + if updates.Load() != 1 { + t.Fatalf("updates = %d, want 1", updates.Load()) + } +} + +// Every corruption mode fails with a message that names the field, and the Secret is left alone. +func TestEnsureRejectsCorruptSecret(t *testing.T) { + now := time.Now() + good, err := Generate(testNS, now.Add(-time.Hour)) + if err != nil { + t.Fatal(err) + } + other, err := Generate(testNS, now.Add(-time.Hour)) + if err != nil { + t.Fatal(err) + } + notCAPEM := other.CertPEM // a leaf certificate in the CA slot + expiredCA, expiredCAKey := selfSignedCA(t, now.Add(-20*24*time.Hour), now.Add(-time.Hour)) + futureCA, futureCAKey := selfSignedCA(t, now.Add(time.Hour), now.Add(48*time.Hour)) + twoCerts := append(append([]byte{}, good.CACertPEM...), other.CACertPEM...) + + tests := []struct { + name string + mutate func(s *corev1.Secret) + want string + }{ + {"wrong type", func(s *corev1.Secret) { s.Type = corev1.SecretTypeTLS }, `type is "kubernetes.io/tls"`}, + {"missing ca.crt", func(s *corev1.Secret) { delete(s.Data, CACertKey) }, `data["ca.crt"] is missing or empty`}, + {"missing ca.key", func(s *corev1.Secret) { delete(s.Data, CAKeyKey) }, `data["ca.key"] is missing or empty`}, + {"missing tls.crt", func(s *corev1.Secret) { delete(s.Data, CertKey) }, `data["tls.crt"] is missing or empty`}, + {"empty tls.key", func(s *corev1.Secret) { s.Data[KeyKey] = nil }, `data["tls.key"] is missing or empty`}, + {"unparsable ca.crt", func(s *corev1.Secret) { s.Data[CACertKey] = []byte("not pem") }, `data["ca.crt"]: unparsable PEM certificate`}, + {"two certificates in ca.crt", func(s *corev1.Secret) { s.Data[CACertKey] = twoCerts }, `data["ca.crt"]: holds 2 certificates, want exactly 1`}, + {"ca.crt is not a CA", func(s *corev1.Secret) { s.Data[CACertKey] = notCAPEM }, `data["ca.crt"] is not a CA certificate`}, + {"expired CA", func(s *corev1.Secret) { s.Data[CACertKey], s.Data[CAKeyKey] = expiredCA, expiredCAKey }, `the CA in data["ca.crt"] expired at`}, + {"CA not yet valid", func(s *corev1.Secret) { s.Data[CACertKey], s.Data[CAKeyKey] = futureCA, futureCAKey }, `the CA in data["ca.crt"] is not valid until`}, + {"unparsable ca.key", func(s *corev1.Secret) { s.Data[CAKeyKey] = []byte("not pem") }, `data["ca.key"]: unparsable PEM private key`}, + {"ca.key does not match ca.crt", func(s *corev1.Secret) { s.Data[CAKeyKey] = other.CAKeyPEM }, `data["ca.key"] does not match the certificate in data["ca.crt"]`}, + {"unparsable tls.crt", func(s *corev1.Secret) { s.Data[CertKey] = []byte("not pem") }, `data["tls.crt"]: unparsable PEM certificate`}, + {"tls.crt signed by another CA", func(s *corev1.Secret) { s.Data[CertKey], s.Data[KeyKey] = other.CertPEM, other.KeyPEM }, `data["tls.crt"] is not signed by the CA in data["ca.crt"]`}, + {"unparsable tls.key", func(s *corev1.Secret) { s.Data[KeyKey] = []byte("not pem") }, `data["tls.key"]: unparsable PEM private key`}, + {"tls.key does not match tls.crt", func(s *corev1.Secret) { s.Data[KeyKey] = good.CAKeyPEM }, `data["tls.key"] does not match the certificate in data["tls.crt"]`}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + secret := secretFor(good) + tt.mutate(secret) + client := fake.NewClientset(secret) + m := newTestManager(t, client, now) + _, err := m.Ensure(t.Context()) + var corruptErr *CorruptSecretError + if !errors.As(err, &corruptErr) { + t.Fatalf("expected *CorruptSecretError, got %v", err) + } + if !strings.Contains(err.Error(), tt.want) { + t.Fatalf("error %q does not name the problem %q", err, tt.want) + } + if !strings.Contains(err.Error(), "secret coder-system/coder-k8s-apiserver-tls") || !strings.Contains(err.Error(), "delete it") { + t.Fatalf("error must name the Secret and the remedy: %q", err) + } + assertNoWrites(t, client) + if cert, _ := m.CurrentCertKeyContent(); cert != nil { + t.Fatal("nothing may be served from a corrupt Secret") + } + }) + } +} + +func TestRunKeepsServingWhenRefreshFails(t *testing.T) { + now := time.Now() + existing, err := Generate(testNS, now) + if err != nil { + t.Fatal(err) + } + client := fake.NewClientset(secretFor(existing)) + m := newTestManager(t, client, now) + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + var gets atomic.Int32 + client.PrependReactor("get", "secrets", func(k8stesting.Action) (bool, runtime.Object, error) { + gets.Add(1) + return true, nil, apierrors.NewServiceUnavailable("down") + }) + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + go func() { m.Run(ctx, 10*time.Millisecond); close(done) }() + deadline := time.Now().Add(5 * time.Second) + for gets.Load() < 2 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + cancel() + <-done + if gets.Load() < 2 { + t.Fatal("Run did not refresh periodically") + } + if cert, _ := m.CurrentCertKeyContent(); string(cert) != string(existing.CertPEM) { + t.Fatal("a failed refresh must keep serving the current certificate") + } +} + +func TestRunRenewsAndNotifies(t *testing.T) { + issued := time.Now() + existing, err := Generate(testNS, issued) + if err != nil { + t.Fatal(err) + } + client := fake.NewClientset(secretFor(existing)) + var clock atomic.Int64 + clock.Store(issued.UnixNano()) + m, err := NewManager(client, testNS) + if err != nil { + t.Fatal(err) + } + m.now = func() time.Time { return time.Unix(0, clock.Load()) } + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + listener := &countingListener{} + m.AddListener(listener) + clock.Store(issued.Add(300 * 24 * time.Hour).UnixNano()) + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + go func() { m.Run(ctx, 10*time.Millisecond); close(done) }() + deadline := time.Now().Add(5 * time.Second) + for listener.count.Load() == 0 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + cancel() + <-done + if listener.count.Load() == 0 { + t.Fatal("renewal during Run must notify listeners") + } + if cert, _ := m.CurrentCertKeyContent(); string(cert) == string(existing.CertPEM) { + t.Fatal("renewal during Run must swap the served certificate") + } +} + +func TestNewManagerAssertions(t *testing.T) { + if _, err := NewManager(nil, testNS); err == nil { + t.Fatal("expected assertion for nil client") + } + if _, err := NewManager(fake.NewClientset(), ""); err == nil { + t.Fatal("expected assertion for empty namespace") + } +} + +func newTestManager(t *testing.T, client *fake.Clientset, now time.Time) *Manager { + t.Helper() + m, err := NewManager(client, testNS) + if err != nil { + t.Fatal(err) + } + m.now = func() time.Time { return now } + return m +} + +func secretFor(b *Bundle) *corev1.Secret { + return &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: SecretName, Namespace: testNS, Labels: copyLabels(), ResourceVersion: "1"}, + Type: SecretType, + Data: b.Data(), + } +} + +func assertNoWrites(t *testing.T, client *fake.Clientset) { + t.Helper() + for _, a := range client.Actions() { + switch a.GetVerb() { + case "create", "update", "patch", "delete": + t.Fatalf("unexpected %s on %s", a.GetVerb(), a.GetResource().Resource) + } + } +} + +// selfSignedCA builds a CA with an explicit validity window (certutil fixes it at 10 years). +func selfSignedCA(t *testing.T, notBefore, notAfter time.Time) ([]byte, []byte) { + t.Helper() + key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader) + if err != nil { + t.Fatal(err) + } + tmpl := &x509.Certificate{ + SerialNumber: big.NewInt(1), Subject: pkix.Name{CommonName: "test-ca"}, + NotBefore: notBefore, NotAfter: notAfter, + KeyUsage: x509.KeyUsageCertSign, BasicConstraintsValid: true, IsCA: true, + } + der, err := x509.CreateCertificate(rand.Reader, tmpl, tmpl, key.Public(), key) + if err != nil { + t.Fatal(err) + } + cert, err := x509.ParseCertificate(der) + if err != nil { + t.Fatal(err) + } + certPEM, err := certutil.EncodeCertificates(cert) + if err != nil { + t.Fatal(err) + } + keyPEM, err := keyutil.MarshalPrivateKeyToPEM(key) + if err != nil { + t.Fatal(err) + } + return certPEM, keyPEM +} + +type countingListener struct{ count atomic.Int32 } + +func (l *countingListener) Enqueue() { l.count.Add(1) } diff --git a/internal/app/apiserverapp/apiserverapp.go b/internal/app/apiserverapp/apiserverapp.go index feb92b89..e167a45f 100644 --- a/internal/app/apiserverapp/apiserverapp.go +++ b/internal/app/apiserverapp/apiserverapp.go @@ -31,6 +31,7 @@ import ( aggregationv1alpha1 "github.com/coder/coder-k8s/api/aggregation/v1alpha1" "github.com/coder/coder-k8s/internal/aggregated/coder" + "github.com/coder/coder-k8s/internal/aggregated/servingcert" "github.com/coder/coder-k8s/internal/aggregated/storage" ) @@ -65,6 +66,10 @@ type Options struct { // resolveDelegationKubeconfig. There is no anonymous or allow-all mode. Authentication *genericoptions.DelegatingAuthenticationOptions Authorization *genericoptions.DelegatingAuthorizationOptions + // ServingCert is a test seam. When nil and the production authentication path is used, the + // server manages a CA-signed serving certificate in its own namespace if it runs in a cluster, + // and falls back to a self-signed localhost certificate otherwise. + ServingCert *servingcert.Manager } type errClientProvider struct { @@ -200,8 +205,13 @@ func NewRecommendedConfig( return nil, fmt.Errorf("assertion failed: recommended config is nil after successful construction") } - if err := secureServingOptions.MaybeDefaultWithSelfSignedCerts("localhost", []string{"localhost"}, nil); err != nil { - return nil, fmt.Errorf("configure self-signed serving certs: %w", err) + // A managed serving certificate (see servingcert) is already set as GeneratedCert; only fall + // back to an in-memory self-signed localhost certificate without one, because + // MaybeDefaultWithSelfSignedCerts would replace it. + if secureServingOptions.ServerCert.GeneratedCert == nil { + if err := secureServingOptions.MaybeDefaultWithSelfSignedCerts("localhost", []string{"localhost"}, nil); err != nil { + return nil, fmt.Errorf("configure self-signed serving certs: %w", err) + } } if err := secureServingOptions.WithLoopback().ApplyTo(&recommendedConfig.SecureServing, &recommendedConfig.LoopbackClientConfig); err != nil { return nil, fmt.Errorf("configure secure serving: %w", err) @@ -360,6 +370,7 @@ func RunWithOptions(ctx context.Context, opts Options) error { secureServingOptions.ServerCert.CertDirectory = "" secureServingOptions.ServerCert.PairName = "" + servingCert := opts.ServingCert authenticationOptions, authorizationOptions := opts.Authentication, opts.Authorization switch { case authenticationOptions == nil && authorizationOptions == nil: @@ -369,10 +380,26 @@ func RunWithOptions(ctx context.Context, opts Options) error { return fmt.Errorf("configure aggregated API server: %w", err) } authenticationOptions, authorizationOptions = newDelegatedAuthOptions(kubeconfigPath) + if servingCert == nil { + servingCert, err = newServingCertManager(kubeconfigPath, serviceAccountNamespaceFile) + if err != nil { + return fmt.Errorf("configure aggregated API server serving certificate: %w", err) + } + } case authenticationOptions == nil || authorizationOptions == nil: return fmt.Errorf("assertion failed: authentication and authorization options must be set together") } + if servingCert != nil { + if _, err := servingCert.Ensure(ctx); err != nil { + return fmt.Errorf("configure aggregated API server serving certificate: %w", err) + } + secureServingOptions.ServerCert.GeneratedCert = servingCert + go servingCert.Run(ctx, servingcert.DefaultCheckInterval) + } else { + log.Printf("warning: no managed serving certificate (not running in a Pod); serving a self-signed certificate for localhost") + } + recommendedConfig, err := NewRecommendedConfig(scheme, codecs, secureServingOptions, authenticationOptions, authorizationOptions) if err != nil { return fmt.Errorf("configure aggregated API server: %w", err) diff --git a/internal/app/apiserverapp/servingcert.go b/internal/app/apiserverapp/servingcert.go new file mode 100644 index 00000000..2996bd8c --- /dev/null +++ b/internal/app/apiserverapp/servingcert.go @@ -0,0 +1,51 @@ +package apiserverapp + +import ( + "errors" + "fmt" + "io/fs" + "os" + "strings" + + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" + + "github.com/coder/coder-k8s/internal/aggregated/servingcert" +) + +// serviceAccountNamespaceFile is mounted into every Pod with a ServiceAccount token. +const serviceAccountNamespaceFile = "/var/run/secrets/kubernetes.io/serviceaccount/namespace" + +// newServingCertManager returns the serving certificate manager for the Pod's namespace, or nil +// when the process does not run in a Pod (no namespace file), for example `go run` on a laptop. +// It uses the same Kubernetes API authority as delegated authentication (kubeconfigPath, or the +// in-cluster ServiceAccount when empty). +func newServingCertManager(kubeconfigPath, namespaceFile string) (*servingcert.Manager, error) { + data, err := os.ReadFile(namespaceFile) //nolint:gosec // G304: fixed ServiceAccount mount path (tests pass a temp file). + if errors.Is(err, fs.ErrNotExist) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("read pod namespace from %s: %w", namespaceFile, err) + } + namespace := strings.TrimSpace(string(data)) + if namespace == "" { + return nil, fmt.Errorf("pod namespace file %s is empty", namespaceFile) + } + + var cfg *rest.Config + if kubeconfigPath != "" { + cfg, err = clientcmd.BuildConfigFromFlags("", kubeconfigPath) + } else { + cfg, err = rest.InClusterConfig() + } + if err != nil { + return nil, fmt.Errorf("build Kubernetes client config: %w", err) + } + client, err := kubernetes.NewForConfig(cfg) + if err != nil { + return nil, fmt.Errorf("build Kubernetes client: %w", err) + } + return servingcert.NewManager(client, namespace) +} diff --git a/internal/app/apiserverapp/servingcert_test.go b/internal/app/apiserverapp/servingcert_test.go new file mode 100644 index 00000000..5b7a8058 --- /dev/null +++ b/internal/app/apiserverapp/servingcert_test.go @@ -0,0 +1,202 @@ +package apiserverapp + +import ( + "context" + "crypto/tls" + "crypto/x509" + "errors" + "fmt" + "io" + "net/http" + "os" + "path/filepath" + "strings" + "testing" + "time" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes/fake" + + "github.com/coder/coder-k8s/internal/aggregated/servingcert" +) + +const servingTestNS = "test-ns" + +func TestNewRecommendedConfigKeepsManagedServingCert(t *testing.T) { + manager := ensuredManager(t, fake.NewClientset()) + secureServingOptions, _ := newTestSecureServing(t) + secureServingOptions.ServerCert.GeneratedCert = manager + authn, authz := newFakeKubeAPI(t).options(false) + cfg, err := NewRecommendedConfig(NewScheme(), codecsFor(), secureServingOptions, authn, authz) + if err != nil { + t.Fatal(err) + } + if cfg.SecureServing.Cert != manager { + t.Fatalf("serving certificate provider was replaced: %T", cfg.SecureServing.Cert) + } +} + +// TestManagedServingCertVerifiesAndHotSwaps runs the real server with a managed certificate: +// clients that trust only the Secret's CA connect with the Service name, other CAs are rejected, +// and a changed Secret is served without a restart. +func TestManagedServingCertVerifiesAndHotSwaps(t *testing.T) { + client := fake.NewClientset() + manager := ensuredManager(t, client) + _, listener := newTestSecureServing(t) + authn, authz := newFakeKubeAPI(t).options(false) + + ctx, cancel := context.WithCancel(context.Background()) + errCh := make(chan error, 1) + go func() { + errCh <- RunWithOptions(ctx, Options{Listener: listener, Authentication: authn, Authorization: authz, ServingCert: manager}) + }() + t.Cleanup(func() { + cancel() + select { + case err := <-errCh: + if err != nil && !errors.Is(err, context.Canceled) { + t.Errorf("server exited with error: %v", err) + } + case <-time.After(10 * time.Second): + t.Error("timed out waiting for server shutdown") + } + }) + addr := listener.Addr().String() + serverName := "coder-k8s-apiserver." + servingTestNS + ".svc" + + firstCA := manager.CABundle() + waitForHealthz(t, addr, firstCA, serverName) + + if _, err := getHealthz(addr, firstCA, "coder-k8s-apiserver.other-ns.svc"); err == nil || !strings.Contains(err.Error(), "certificate is valid for") { + t.Fatalf("expected a hostname mismatch for another namespace, got %v", err) + } + otherCA, err := servingcert.Generate(servingTestNS, time.Now()) + if err != nil { + t.Fatal(err) + } + if _, err := getHealthz(addr, otherCA.CACertPEM, serverName); err == nil || !strings.Contains(err.Error(), "certificate signed by unknown authority") { + t.Fatalf("expected unknown authority with a wrong CA, got %v", err) + } + + // Replace the Secret (for example after a CA rotation) and let the manager adopt it. + secret, err := client.CoreV1().Secrets(servingTestNS).Get(ctx, servingcert.SecretName, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + secret.Data = otherCA.Data() + if _, err := client.CoreV1().Secrets(servingTestNS).Update(ctx, secret, metav1.UpdateOptions{}); err != nil { + t.Fatal(err) + } + if _, err := manager.Ensure(ctx); err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(10 * time.Second) + for { + status, err := getHealthz(addr, otherCA.CACertPEM, serverName) + if err == nil && status == http.StatusOK { + break + } + if time.Now().After(deadline) { + t.Fatalf("new certificate not served without restart: status=%d err=%v", status, err) + } + time.Sleep(100 * time.Millisecond) + } + if _, err := getHealthz(addr, firstCA, serverName); err == nil { + t.Fatal("the old CA must no longer verify after the swap") + } +} + +func TestRunWithOptionsFailsOnCorruptServingCertSecret(t *testing.T) { + client := fake.NewClientset(&corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: servingcert.SecretName, Namespace: servingTestNS}, + Type: servingcert.SecretType, + Data: map[string][]byte{"ca.crt": []byte("garbage")}, + }) + manager, err := servingcert.NewManager(client, servingTestNS) + if err != nil { + t.Fatal(err) + } + _, listener := newTestSecureServing(t) + authn, authz := newFakeKubeAPI(t).options(false) + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + err = RunWithOptions(ctx, Options{Listener: listener, Authentication: authn, Authorization: authz, ServingCert: manager}) + if err == nil || !strings.Contains(err.Error(), "serving certificate") || !strings.Contains(err.Error(), `data["ca.key"] is missing`) { + t.Fatalf("expected a diagnosable startup failure, got %v", err) + } +} + +func TestNewServingCertManager(t *testing.T) { + dir := t.TempDir() + if m, err := newServingCertManager("", filepath.Join(dir, "absent")); err != nil || m != nil { + t.Fatalf("outside a Pod: manager=%v err=%v, want nil, nil", m, err) + } + empty := filepath.Join(dir, "empty") + mustWrite(t, empty, " \n") + if _, err := newServingCertManager("", empty); err == nil || !strings.Contains(err.Error(), "is empty") { + t.Fatalf("expected empty namespace file error, got %v", err) + } + nsFile := filepath.Join(dir, "namespace") + mustWrite(t, nsFile, servingTestNS+"\n") + if _, err := newServingCertManager(filepath.Join(dir, "no-kubeconfig"), nsFile); err == nil || !strings.Contains(err.Error(), "client config") { + t.Fatalf("expected client config error, got %v", err) + } + m, err := newServingCertManager(newFakeKubeAPI(t).kubeconfigPath, nsFile) + if err != nil || m == nil { + t.Fatalf("in a Pod: manager=%v err=%v", m, err) + } + if !strings.Contains(m.Name(), servingTestNS+"/"+servingcert.SecretName) { + t.Fatalf("manager must target %s/%s, got %s", servingTestNS, servingcert.SecretName, m.Name()) + } + if err := os.Chmod(nsFile, 0o000); err == nil && os.Geteuid() != 0 { + if _, err := newServingCertManager("", nsFile); err == nil || !strings.Contains(err.Error(), "read pod namespace") { + t.Fatalf("expected read error, got %v", err) + } + } +} + +func ensuredManager(t *testing.T, client *fake.Clientset) *servingcert.Manager { + t.Helper() + m, err := servingcert.NewManager(client, servingTestNS) + if err != nil { + t.Fatal(err) + } + if _, err := m.Ensure(t.Context()); err != nil { + t.Fatal(err) + } + return m +} + +func getHealthz(addr string, caPEM []byte, serverName string) (int, error) { + roots := x509.NewCertPool() + if !roots.AppendCertsFromPEM(caPEM) { + return 0, fmt.Errorf("assertion failed: CA PEM did not parse") + } + client := &http.Client{Timeout: 5 * time.Second, Transport: &http.Transport{ + TLSClientConfig: &tls.Config{RootCAs: roots, ServerName: serverName, MinVersion: tls.VersionTLS12}, + DisableKeepAlives: true, + }} + resp, err := client.Get("https://" + addr + "/healthz") + if err != nil { + return 0, err + } + defer func() { _ = resp.Body.Close() }() + _, _ = io.Copy(io.Discard, resp.Body) + return resp.StatusCode, nil +} + +func waitForHealthz(t *testing.T, addr string, caPEM []byte, serverName string) { + t.Helper() + deadline := time.Now().Add(15 * time.Second) + for { + status, err := getHealthz(addr, caPEM, serverName) + if err == nil && status == http.StatusOK { + return + } + if time.Now().After(deadline) { + t.Fatalf("server did not become healthy with verified TLS: status=%d err=%v", status, err) + } + time.Sleep(100 * time.Millisecond) + } +} From a7f93ae9b04863c6526dee708b2d80f50baf0a65 Mon Sep 17 00:00:00 2001 From: Thomas Kosiewski Date: Fri, 25 Sep 2026 09:35:55 +0000 Subject: [PATCH 2/2] refactor(apiserver): log serving certificate events with the controller-runtime logger MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Refs #137 Keeps k8s.io/klog/v2 an indirect dependency (make verify-vendor). --- _Generated with `xum` • Model: `anthropic:claude-opus-5-5` • Thinking: `high`_ --- internal/aggregated/servingcert/manager.go | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/internal/aggregated/servingcert/manager.go b/internal/aggregated/servingcert/manager.go index 9fe36fbe..ad38e28a 100644 --- a/internal/aggregated/servingcert/manager.go +++ b/internal/aggregated/servingcert/manager.go @@ -12,9 +12,11 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apiserver/pkg/server/dynamiccertificates" "k8s.io/client-go/kubernetes" - "k8s.io/klog/v2" + ctrl "sigs.k8s.io/controller-runtime" ) +var log = ctrl.Log.WithName("servingcert") + // DefaultCheckInterval is how often Run re-reads the Secret and renews the serving certificate. const DefaultCheckInterval = 12 * time.Hour @@ -76,7 +78,7 @@ func (m *Manager) Ensure(ctx context.Context) (*Bundle, error) { if err != nil { return nil, fmt.Errorf("create secret %s/%s: %w", m.namespace, SecretName, err) } - klog.InfoS("Created aggregated API server CA and serving certificate", "secret", klog.KRef(m.namespace, SecretName)) + log.Info("Created aggregated API server CA and serving certificate", "namespace", m.namespace, "secret", SecretName) return bundle, m.serve(bundle) case err != nil: return nil, fmt.Errorf("get secret %s/%s: %w", m.namespace, SecretName, err) @@ -101,7 +103,7 @@ func (m *Manager) Ensure(ctx context.Context) (*Bundle, error) { if err != nil { return nil, fmt.Errorf("update secret %s/%s: %w", m.namespace, SecretName, err) } - klog.InfoS("Renewed aggregated API server serving certificate", "secret", klog.KRef(m.namespace, SecretName), "notAfter", bundle.Cert.NotAfter) + log.Info("Renewed aggregated API server serving certificate", "namespace", m.namespace, "secret", SecretName, "notAfter", bundle.Cert.NotAfter) return bundle, m.serve(bundle) } return nil, fmt.Errorf("secret %s/%s kept changing; gave up after %d attempts", m.namespace, SecretName, maxEnsureAttempts) @@ -121,7 +123,7 @@ func (m *Manager) Run(ctx context.Context, interval time.Duration) { return case <-ticker.C: if _, err := m.Ensure(ctx); err != nil { - klog.ErrorS(err, "Could not refresh the aggregated API server serving certificate; still serving the current one") + log.Error(err, "Could not refresh the aggregated API server serving certificate; still serving the current one") } } }