diff --git a/.cspell.json b/.cspell.json index e200d6fb..61a5c4bb 100644 --- a/.cspell.json +++ b/.cspell.json @@ -56,7 +56,8 @@ "portforward", "livez", "requestheader", - "subjectaccessreviews" + "subjectaccessreviews", + "cabundle" ], "ignorePaths": [ ".git/**", diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index f0660711..06b69792 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -320,6 +320,36 @@ jobs: kubectl wait --for=condition=Available deploy/coder-k8s -n coder-system --timeout=120s kubectl wait --for=condition=Available apiservice/v1alpha1.aggregation.coder.com --timeout=180s + # The aggregated API server sets the APIService caBundle from its CA Secret and turns + # insecureSkipTLSVerify off shortly after it starts. kube-apiserver's Available check does + # not verify the serving certificate, so the proof is a request proxied with verification on. + - name: Verify the APIService trusts the serving CA + if: env.E2E_FULL == 'true' + run: | + set -euo pipefail + apisvc=apiservice/v1alpha1.aggregation.coder.com + want="" have="" insecure="" + for _ in $(seq 1 60); do + want=$(kubectl -n coder-system get secret coder-k8s-apiserver-tls -o jsonpath='{.data.ca\.crt}' 2>/dev/null || true) + have=$(kubectl get "$apisvc" -o jsonpath='{.spec.caBundle}') + insecure=$(kubectl get "$apisvc" -o jsonpath='{.spec.insecureSkipTLSVerify}') + if [[ -n $want && $have == "$want" && -z $insecure ]]; then + break + fi + sleep 2 + done + if [[ -z $want || $have != "$want" || -n $insecure ]]; then + echo "assertion failed: APIService caBundle does not match the serving CA Secret (insecureSkipTLSVerify='$insecure')" >&2 + kubectl get "$apisvc" -o yaml >&2 + exit 1 + fi + kubectl wait --for=condition=Available "$apisvc" --timeout=180s + for _ in $(seq 1 30); do + kubectl get --raw /apis/aggregation.coder.com/v1alpha1 >/dev/null && break + sleep 2 + done + kubectl get --raw /apis/aggregation.coder.com/v1alpha1 >/dev/null + - name: Install CloudNativePG operator if: env.E2E_FULL == 'true' run: | @@ -430,6 +460,55 @@ jobs: E2E_WORKDIR: ${{ runner.temp }}/workspace-lifecycle run: bash ./hack/e2e-workspace-lifecycle.sh + # Last functional step: with the sync's permission revoked (so it cannot repair the value), + # a wrong caBundle makes requests proxied by kube-apiserver fail, and kube-apiserver logs the + # x509 error (its Available check does not verify certificates). Restoring the permission + # lets the sync put the right CA back. + - name: Verify the APIService rejects a wrong CA + if: env.E2E_FULL == 'true' + run: | + set -euo pipefail + apisvc=apiservice/v1alpha1.aggregation.coder.com + can_patch() { + kubectl auth can-i patch apiservices.apiregistration.k8s.io/v1alpha1.aggregation.coder.com \ + --as=system:serviceaccount:coder-system:coder-k8s 2>/dev/null || true + } + kubectl delete clusterrolebinding coder-k8s-apiservice-cabundle + for _ in $(seq 1 30); do + [[ $(can_patch) == no ]] && break + sleep 1 + done + test "$(can_patch)" = no + wrong=$(openssl req -x509 -newkey ec -pkeyopt ec_paramgen_curve:P-256 -nodes -keyout /dev/null \ + -days 1 -subj /CN=e2e-wrong-ca 2>/dev/null | base64 -w0) + since=$(date -u +%Y-%m-%dT%H:%M:%SZ) + kubectl patch "$apisvc" --type=merge -p "{\"spec\":{\"caBundle\":\"$wrong\"}}" + for _ in $(seq 1 30); do + kubectl get --raw /apis/aggregation.coder.com/v1alpha1 >/dev/null 2>&1 || break + sleep 1 + done + if kubectl get --raw /apis/aggregation.coder.com/v1alpha1 >/dev/null; then + echo "assertion failed: aggregated API still served with a wrong caBundle" >&2 + exit 1 + fi + x509="" + for _ in $(seq 1 30); do + x509=$(kubectl -n kube-system logs -l component=kube-apiserver --since-time="$since" --tail=-1 | + grep 'v1alpha1.aggregation.coder.com' | grep -m1 'x509: certificate signed by unknown authority' || true) + [[ -n $x509 ]] && break + sleep 2 + done + echo "${x509:?assertion failed: kube-apiserver logged no x509 error for the wrong caBundle}" + kubectl apply -f config/rbac/apiservice-cabundle-role.yaml + want=$(kubectl -n coder-system get secret coder-k8s-apiserver-tls -o jsonpath='{.data.ca\.crt}') + for _ in $(seq 1 90); do + [[ $(kubectl get "$apisvc" -o jsonpath='{.spec.caBundle}') == "$want" ]] && + kubectl get --raw /apis/aggregation.coder.com/v1alpha1 >/dev/null 2>&1 && break + sleep 2 + done + test "$(kubectl get "$apisvc" -o jsonpath='{.spec.caBundle}')" = "$want" + kubectl get --raw /apis/aggregation.coder.com/v1alpha1 >/dev/null + # ---- Failure diagnostics ---- - name: Dump cluster state on failure if: failure() && env.E2E_FULL == 'true' diff --git a/config/rbac/apiservice-cabundle-role.yaml b/config/rbac/apiservice-cabundle-role.yaml new file mode 100644 index 00000000..ab61a5f0 --- /dev/null +++ b/config/rbac/apiservice-cabundle-role.yaml @@ -0,0 +1,33 @@ +# Lets the aggregated API server keep its own APIService trusting its serving CA +# (spec.caBundle, spec.insecureSkipTLSVerify). Limited to that one APIService; list and watch are +# authorized only for requests that select it by metadata.name. Not part of dist/install.yaml, +# which installs the controller only. +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: coder-k8s-apiservice-cabundle +rules: + - apiGroups: + - apiregistration.k8s.io + resources: + - apiservices + resourceNames: + - v1alpha1.aggregation.coder.com + verbs: + - get + - list + - watch + - patch +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: coder-k8s-apiservice-cabundle +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: coder-k8s-apiservice-cabundle +subjects: + - kind: ServiceAccount + name: coder-k8s + namespace: coder-system diff --git a/deploy/apiserver-apiservice.yaml b/deploy/apiserver-apiservice.yaml index 85c602ff..47f09d0e 100644 --- a/deploy/apiserver-apiservice.yaml +++ b/deploy/apiserver-apiservice.yaml @@ -10,4 +10,3 @@ spec: namespace: coder-system groupPriorityMinimum: 1000 versionPriority: 100 - insecureSkipTLSVerify: true diff --git a/docs/how-to/deploy-aggregated-apiserver.md b/docs/how-to/deploy-aggregated-apiserver.md index e17cd94f..05d6dab4 100644 --- a/docs/how-to/deploy-aggregated-apiserver.md +++ b/docs/how-to/deploy-aggregated-apiserver.md @@ -118,12 +118,50 @@ In a cluster, the aggregated API server serves a certificate signed by its own C - 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. +- The CA is valid for 10 years. To replace it, delete the Secret and restart **every** replica (`kubectl -n coder-system rollout restart deployment/coder-k8s`). The first new pod generates a CA and updates the APIService `caBundle`; until the old pods are gone, requests routed to them fail certificate verification, so expect `503 ServiceUnavailable` for several seconds (about 10 seconds with two replicas in testing). Other 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). +## How kube-apiserver trusts the server + +kube-apiserver verifies the aggregated API server's certificate against the APIService `spec.caBundle`. The aggregated API server keeps that field set to the CA in `coder-k8s-apiserver-tls`, and keeps `insecureSkipTLSVerify` off: + +- It patches only the APIService `v1alpha1.aggregation.coder.com`, and only when the `caBundle` differs from the Secret's CA or `insecureSkipTLSVerify` is set. It reacts to changes of that APIService and of its own CA; it does not poll. +- It writes only a CA that passed the same checks as at startup. If the Secret is missing or invalid, it leaves the APIService unchanged and logs why. +- It needs `config/rbac/apiservice-cabundle-role.yaml`: `get`, `list`, `watch`, and `patch` on that one APIService (`resourceNames`). `list` and `watch` are allowed only for requests that select it by name. The controller-only install bundle (`dist/install.yaml`) does not include this file, because it registers no APIService. +- Without that permission the aggregated API server keeps serving, logs the missing permission (at most every 5 minutes), and retries with backoff up to 60 seconds. +- On a fresh install, requests through kube-apiserver fail with `503 ServiceUnavailable` for a few seconds, until the first patch. `Available=True` alone does not show that verification works, because kube-apiserver's availability check does not verify the certificate. Check a real request instead: `kubectl get --raw /apis/aggregation.coder.com/v1alpha1`. +- Running `kubectl apply -f deploy/apiserver-apiservice.yaml` again keeps the injected `caBundle`. Replacing or re-creating the APIService clears it; the aggregated API server sets it again within seconds. + +### Upgrade from a version that used `insecureSkipTLSVerify` + +Apply the changes in this order: + +1. `kubectl apply -f config/rbac/` +2. Deploy the new image. It sets `caBundle` and turns `insecureSkipTLSVerify` off in one step. +3. `kubectl apply -f deploy/apiserver-apiservice.yaml` (the file no longer sets `insecureSkipTLSVerify`). + +If the new image starts before step 1, it logs a missing-permission error until the RBAC exists, and the APIService keeps working without verification in the meantime. If you apply step 3 before the new image runs, requests through kube-apiserver fail with `503` until the new image starts. Do not re-apply an old copy of `deploy/apiserver-apiservice.yaml`: once a `caBundle` is set, the API rejects `insecureSkipTLSVerify: true`. + +To roll back to a version without a managed serving certificate, restore the old registration before you change the image. The opt-out annotation stops the running server from setting the `caBundle` again: + +```bash +kubectl annotate apiservice v1alpha1.aggregation.coder.com coder.com/manage-ca-bundle=false +kubectl patch apiservice v1alpha1.aggregation.coder.com --type=merge \ + -p '{"spec":{"caBundle":null,"insecureSkipTLSVerify":true}}' +``` + +Then deploy the old image and re-apply its `deploy/apiserver-apiservice.yaml`. + +### Opt out + +If another tool owns the APIService `caBundle` (for example cert-manager's CA injector, or a GitOps tool that sets it from Git), annotate the APIService so the aggregated API server leaves it alone: + +```bash +kubectl annotate apiservice v1alpha1.aggregation.coder.com coder.com/manage-ca-bundle=false +``` + +That tool must then trust a CA that signed the certificate the server actually serves, which is the one in `coder-k8s-apiserver-tls`. Remove the annotation (or set any other value) to hand the field back. diff --git a/hack/update-manifests.sh b/hack/update-manifests.sh index 35ea576e..feba1651 100755 --- a/hack/update-manifests.sh +++ b/hack/update-manifests.sh @@ -29,7 +29,7 @@ list_manifests() { } # RBAC manifests that only the aggregated API server needs; the controller-only install bundle leaves them out. -BUNDLE_EXCLUDED_RBAC=(auth-delegator-binding.yaml authentication-reader-binding.yaml) +BUNDLE_EXCLUDED_RBAC=(apiservice-cabundle-role.yaml auth-delegator-binding.yaml authentication-reader-binding.yaml) # write_default_kustomization lists every generated CRD and RBAC file individually, so the install bundle # picks up new files automatically without placing kustomization files in the directories above. The bundle diff --git a/internal/aggregated/apiservicetrust/controller.go b/internal/aggregated/apiservicetrust/controller.go new file mode 100644 index 00000000..f66602e9 --- /dev/null +++ b/internal/aggregated/apiservicetrust/controller.go @@ -0,0 +1,237 @@ +// Package apiservicetrust keeps the aggregated API server's APIService trusting its serving CA. +// +// The controller runs inside the aggregated API server process. It sets spec.caBundle of the +// v1alpha1.aggregation.coder.com APIService to the CA in the servingcert Secret and turns +// spec.insecureSkipTLSVerify off, in one merge patch and only when something differs. It reacts +// to a single-object informer on that APIService and to servingcert.Manager changes; it never +// polls, and failures never affect serving or readiness. +package apiservicetrust + +import ( + "context" + "encoding/base64" + "encoding/json" + "fmt" + "sync" + "time" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/fields" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apiserver/pkg/server/dynamiccertificates" + "k8s.io/client-go/dynamic" + "k8s.io/client-go/dynamic/dynamicinformer" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" + ctrl "sigs.k8s.io/controller-runtime" + + "github.com/coder/coder-k8s/internal/aggregated/servingcert" +) + +const ( + // APIServiceName is the APIService registered by deploy/apiserver-apiservice.yaml. + APIServiceName = "v1alpha1.aggregation.coder.com" + // OptOutAnnotation set to "false" on the APIService makes the controller leave it alone, for + // clusters where cert-manager or a GitOps tool owns spec.caBundle. + OptOutAnnotation = "coder.com/manage-ca-bundle" + // FieldManager identifies the controller's patches in managedFields. + FieldManager = "coder-k8s-apiservice-cabundle" + // DocsURL explains the RBAC, upgrade order, and opt-out. + DocsURL = "https://coder.github.io/coder-k8s/how-to/deploy-aggregated-apiserver/#how-kube-apiserver-trusts-the-server" + + retryBaseDelay = 500 * time.Millisecond + // retryMaxDelay bounds how long a repaired permission takes to be noticed. + retryMaxDelay = 60 * time.Second + // warnInterval rate-limits repeated warnings (for example while RBAC is missing). + warnInterval = 5 * time.Minute +) + +// APIServiceGVR is the apiregistration.k8s.io/v1 APIService resource. kube-aggregator's typed +// client is not vendored, so the controller uses the dynamic client. +var APIServiceGVR = schema.GroupVersionResource{Group: "apiregistration.k8s.io", Version: "v1", Resource: "apiservices"} + +var log = ctrl.Log.WithName("apiservicetrust") + +// Controller keeps the APIService caBundle in sync with the servingcert Secret. +type Controller struct { + dyn dynamic.Interface + secrets kubernetes.Interface + namespace string + now func() time.Time + + queue workqueue.TypedRateLimitingInterface[string] + factory dynamicinformer.DynamicSharedInformerFactory + informer cache.SharedIndexInformer + + warnMu sync.Mutex + lastWarn map[string]time.Time +} + +var _ dynamiccertificates.Listener = (*Controller)(nil) + +// New returns a Controller for the servingcert Secret in namespace. +func New(dyn dynamic.Interface, secrets kubernetes.Interface, namespace string) (*Controller, error) { + if dyn == nil || secrets == nil { + return nil, fmt.Errorf("assertion failed: Kubernetes clients must not be nil") + } + if namespace == "" { + return nil, fmt.Errorf("assertion failed: namespace must not be empty") + } + // A metadata.name field selector lets RBAC authorize list/watch with resourceNames. + factory := dynamicinformer.NewFilteredDynamicSharedInformerFactory(dyn, 0, metav1.NamespaceAll, func(o *metav1.ListOptions) { + o.FieldSelector = fields.OneTermEqualSelector("metadata.name", APIServiceName).String() + }) + c := &Controller{ + dyn: dyn, + secrets: secrets, + namespace: namespace, + now: time.Now, + queue: workqueue.NewTypedRateLimitingQueueWithConfig( + workqueue.NewTypedItemExponentialFailureRateLimiter[string](retryBaseDelay, retryMaxDelay), + workqueue.TypedRateLimitingQueueConfig[string]{Name: "apiservice-cabundle"}, + ), + factory: factory, + informer: factory.ForResource(APIServiceGVR).Informer(), + lastWarn: map[string]time.Time{}, + } + // Without list/watch permission the informer never syncs; say so instead of waiting silently. + if err := c.informer.SetWatchErrorHandler(func(_ *cache.Reflector, err error) { + if apierrors.IsForbidden(err) { + c.warn("forbidden-watch", err, "missing permission to list and watch the APIService; apply config/rbac/apiservice-cabundle-role.yaml", + "permission", "list,watch apiservices.apiregistration.k8s.io/"+APIServiceName, "docs", DocsURL) + return + } + c.warn("watch", err, "cannot watch the APIService; retrying") + }); err != nil { + return nil, fmt.Errorf("set APIService watch error handler: %w", err) + } + if _, err := c.informer.AddEventHandler(cache.ResourceEventHandlerFuncs{ + AddFunc: func(any) { c.Enqueue() }, + UpdateFunc: func(any, any) { c.Enqueue() }, + DeleteFunc: func(any) { c.Enqueue() }, + }); err != nil { + return nil, fmt.Errorf("add APIService event handler: %w", err) + } + return c, nil +} + +// Enqueue schedules a sync. It implements dynamiccertificates.Listener, so registering the +// controller with servingcert.Manager.AddListener reacts to CA changes. +func (c *Controller) Enqueue() { + c.queue.Add(APIServiceName) +} + +// Run syncs until ctx is done. +func (c *Controller) Run(ctx context.Context) { + defer c.queue.ShutDown() + c.factory.Start(ctx.Done()) + if !cache.WaitForCacheSync(ctx.Done(), c.informer.HasSynced) { + return + } + c.Enqueue() + go func() { + <-ctx.Done() + c.queue.ShutDown() + }() + for c.processNext(ctx) { + } + c.factory.Shutdown() +} + +func (c *Controller) processNext(ctx context.Context) bool { + key, shutdown := c.queue.Get() + if shutdown { + return false + } + defer c.queue.Done(key) + if err := c.sync(ctx); err != nil { + c.queue.AddRateLimited(key) + return true + } + c.queue.Forget(key) + return true +} + +// sync patches the APIService when its trust settings differ from the Secret's CA. It returns an +// error only for failures worth retrying. +func (c *Controller) sync(ctx context.Context) error { + obj, exists, err := c.informer.GetStore().GetByKey(APIServiceName) + if err != nil { + return fmt.Errorf("read APIService from cache: %w", err) + } + if !exists { + c.warn("absent", nil, "APIService is not registered; its caBundle will be set once it exists", "apiService", APIServiceName, "docs", DocsURL) + return nil + } + apiService, ok := obj.(*unstructured.Unstructured) + if !ok { + return fmt.Errorf("assertion failed: APIService cache object is %T", obj) + } + if apiService.GetAnnotations()[OptOutAnnotation] == "false" { + c.warn("optout", nil, "APIService opted out of caBundle management; leaving it unchanged", "annotation", OptOutAnnotation+"=false") + return nil + } + + secret, err := c.secrets.CoreV1().Secrets(c.namespace).Get(ctx, servingcert.SecretName, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + // Mid-rotation: the serving certificate manager creates the Secret and notifies us. + c.warn("nosecret", nil, "serving CA Secret not found; not changing the APIService caBundle", "secret", c.namespace+"/"+servingcert.SecretName) + return nil + } + if err != nil { + c.warn("getsecret", err, "cannot read the serving CA Secret; retrying", "secret", c.namespace+"/"+servingcert.SecretName) + return err + } + // Only validated trust material is ever written: the API does not validate caBundle. + bundle, err := servingcert.Parse(secret, c.namespace, c.now()) + if err != nil { + c.warn("corrupt", err, "serving CA Secret is invalid; not changing the APIService caBundle") + return nil + } + + want := base64.StdEncoding.EncodeToString(bundle.CACertPEM) + have, _, _ := unstructured.NestedString(apiService.Object, "spec", "caBundle") + insecure, _, _ := unstructured.NestedBool(apiService.Object, "spec", "insecureSkipTLSVerify") + if have == want && !insecure { + return nil + } + + // One patch: the API rejects insecureSkipTLSVerify=true together with a caBundle. + patch, err := json.Marshal(map[string]any{"spec": map[string]any{"caBundle": want, "insecureSkipTLSVerify": false}}) + if err != nil { + return fmt.Errorf("assertion failed: marshal patch: %w", err) + } + if _, err := c.dyn.Resource(APIServiceGVR).Patch(ctx, APIServiceName, types.MergePatchType, patch, metav1.PatchOptions{FieldManager: FieldManager}); err != nil { + if apierrors.IsForbidden(err) { + c.warn("forbidden", err, "missing permission to patch the APIService caBundle; apply config/rbac/apiservice-cabundle-role.yaml", + "permission", "patch apiservices.apiregistration.k8s.io/"+APIServiceName, "docs", DocsURL) + } else { + c.warn("patch", err, "cannot patch the APIService caBundle; retrying", "docs", DocsURL) + } + return err + } + log.Info("Set the APIService caBundle to the serving CA and enabled TLS verification", "apiService", APIServiceName) + return nil +} + +// warn logs at most once per warnInterval for each kind of problem. +func (c *Controller) warn(kind string, err error, msg string, keysAndValues ...any) { + c.warnMu.Lock() + now := c.now() + last, seen := c.lastWarn[kind] + if seen && now.Sub(last) < warnInterval { + c.warnMu.Unlock() + return + } + c.lastWarn[kind] = now + c.warnMu.Unlock() + if err != nil { + log.Error(err, msg, keysAndValues...) + return + } + log.Info(msg, keysAndValues...) +} diff --git a/internal/aggregated/apiservicetrust/controller_test.go b/internal/aggregated/apiservicetrust/controller_test.go new file mode 100644 index 00000000..16e75f68 --- /dev/null +++ b/internal/aggregated/apiservicetrust/controller_test.go @@ -0,0 +1,394 @@ +package apiservicetrust + +import ( + "context" + "encoding/base64" + "encoding/json" + "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/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + dynamicfake "k8s.io/client-go/dynamic/fake" + "k8s.io/client-go/kubernetes/fake" + k8stesting "k8s.io/client-go/testing" + + "github.com/coder/coder-k8s/internal/aggregated/servingcert" +) + +const testNS = "coder-system" + +type harness struct { + dyn *dynamicfake.FakeDynamicClient + kube *fake.Clientset + ctrl *Controller + bundle *servingcert.Bundle + patches atomic.Int32 +} + +func newHarness(t *testing.T, apiService *unstructured.Unstructured, withSecret bool) *harness { + t.Helper() + h := &harness{kube: fake.NewClientset()} + var objs []runtime.Object + if apiService != nil { + objs = append(objs, apiService) + } + h.dyn = dynamicfake.NewSimpleDynamicClientWithCustomListKinds(runtime.NewScheme(), map[schemaGVR]string{APIServiceGVR: "APIServiceList"}, objs...) + h.dyn.PrependReactor("patch", "apiservices", func(k8stesting.Action) (bool, runtime.Object, error) { + h.patches.Add(1) + return false, nil, nil + }) + if withSecret { + b, err := servingcert.Generate(testNS, time.Now()) + if err != nil { + t.Fatal(err) + } + h.bundle = b + h.setSecret(t, b) + } + c, err := New(h.dyn, h.kube, testNS) + if err != nil { + t.Fatal(err) + } + h.ctrl = c + return h +} + +func (h *harness) setSecret(t *testing.T, b *servingcert.Bundle) { + t.Helper() + secret := &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: servingcert.SecretName, Namespace: testNS}, + Type: servingcert.SecretType, + Data: b.Data(), + } + if _, err := h.kube.CoreV1().Secrets(testNS).Get(context.Background(), servingcert.SecretName, metav1.GetOptions{}); err == nil { + if _, err := h.kube.CoreV1().Secrets(testNS).Update(context.Background(), secret, metav1.UpdateOptions{}); err != nil { + t.Fatal(err) + } + return + } + if _, err := h.kube.CoreV1().Secrets(testNS).Create(context.Background(), secret, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } +} + +func (h *harness) run(t *testing.T) { + t.Helper() + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { h.ctrl.Run(ctx); close(done) }() + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("controller did not stop") + } + }) +} + +func (h *harness) get(t *testing.T) *unstructured.Unstructured { + t.Helper() + u, err := h.dyn.Resource(APIServiceGVR).Get(context.Background(), APIServiceName, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + return u +} + +func caBundleOf(u *unstructured.Unstructured) string { + v, _, _ := unstructured.NestedString(u.Object, "spec", "caBundle") + return v +} + +func insecureOf(u *unstructured.Unstructured) bool { + v, _, _ := unstructured.NestedBool(u.Object, "spec", "insecureSkipTLSVerify") + return v +} + +func eventually(t *testing.T, what string, cond func() bool) { + t.Helper() + deadline := time.Now().Add(5 * time.Second) + for !cond() { + if time.Now().After(deadline) { + t.Fatalf("timed out waiting for %s", what) + } + time.Sleep(10 * time.Millisecond) + } +} + +// consistently fails if cond becomes false within d. +func consistently(t *testing.T, what string, d time.Duration, cond func() bool) { + t.Helper() + deadline := time.Now().Add(d) + for time.Now().Before(deadline) { + if !cond() { + t.Fatalf("%s changed", what) + } + time.Sleep(10 * time.Millisecond) + } +} + +func newAPIService(caBundle string, insecure bool, annotations map[string]string) *unstructured.Unstructured { + spec := map[string]any{ + "group": "aggregation.coder.com", "version": "v1alpha1", + "service": map[string]any{"name": servingcert.ServiceName, "namespace": testNS}, + "groupPriorityMinimum": int64(1000), "versionPriority": int64(100), + } + if insecure { + spec["insecureSkipTLSVerify"] = true + } + if caBundle != "" { + spec["caBundle"] = caBundle + } + u := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "apiregistration.k8s.io/v1", "kind": "APIService", + "metadata": map[string]any{"name": APIServiceName}, + "spec": spec, + }} + if annotations != nil { + u.SetAnnotations(annotations) + } + return u +} + +func patchActions(h *harness) []k8stesting.PatchActionImpl { + var out []k8stesting.PatchActionImpl + for _, a := range h.dyn.Actions() { + if p, ok := a.(k8stesting.PatchActionImpl); ok { + out = append(out, p) + } + } + return out +} + +func TestSyncPatchesInsecureAPIServiceOnceWithBothFields(t *testing.T) { + h := newHarness(t, newAPIService("", true, nil), true) + h.run(t) + want := base64.StdEncoding.EncodeToString(h.bundle.CACertPEM) + eventually(t, "caBundle", func() bool { return caBundleOf(h.get(t)) == want }) + if insecureOf(h.get(t)) { + t.Fatal("insecureSkipTLSVerify must be off") + } + consistently(t, "patch count", 300*time.Millisecond, func() bool { return h.patches.Load() == 1 }) + p := patchActions(h)[0] + if p.PatchType != types.MergePatchType || p.PatchOptions.FieldManager != FieldManager { + t.Fatalf("patch type %s / field manager %q", p.PatchType, p.PatchOptions.FieldManager) + } + var body map[string]map[string]any + if err := json.Unmarshal(p.Patch, &body); err != nil { + t.Fatal(err) + } + if len(body) != 1 || len(body["spec"]) != 2 || body["spec"]["caBundle"] != want || body["spec"]["insecureSkipTLSVerify"] != false { + t.Fatalf("patch body = %s", p.Patch) + } +} + +func TestSyncLeavesMatchingAPIServiceAlone(t *testing.T) { + b, err := servingcert.Generate(testNS, time.Now()) + if err != nil { + t.Fatal(err) + } + h := newHarness(t, newAPIService(base64.StdEncoding.EncodeToString(b.CACertPEM), false, nil), false) + h.setSecret(t, b) + h.run(t) + h.ctrl.Enqueue() + consistently(t, "patch count", 500*time.Millisecond, func() bool { return h.patches.Load() == 0 }) +} + +func TestSyncPatchesWhenCABundleDiffersButInsecureOff(t *testing.T) { + h := newHarness(t, newAPIService(base64.StdEncoding.EncodeToString([]byte("stale")), false, nil), true) + h.run(t) + want := base64.StdEncoding.EncodeToString(h.bundle.CACertPEM) + eventually(t, "caBundle", func() bool { return caBundleOf(h.get(t)) == want }) +} + +func TestSyncRespectsOptOutAnnotation(t *testing.T) { + h := newHarness(t, newAPIService("", true, map[string]string{OptOutAnnotation: "false"}), true) + h.run(t) + h.ctrl.Enqueue() + consistently(t, "patch count", 500*time.Millisecond, func() bool { return h.patches.Load() == 0 }) + if !insecureOf(h.get(t)) || caBundleOf(h.get(t)) != "" { + t.Fatal("an opted-out APIService must not change") + } +} + +func TestSyncManagesAgainWhenOptOutIsRemoved(t *testing.T) { + h := newHarness(t, newAPIService("", true, map[string]string{OptOutAnnotation: "false"}), true) + h.run(t) + u := h.get(t) + u.SetAnnotations(map[string]string{OptOutAnnotation: "true"}) + if _, err := h.dyn.Resource(APIServiceGVR).Update(context.Background(), u, metav1.UpdateOptions{}); err != nil { + t.Fatal(err) + } + want := base64.StdEncoding.EncodeToString(h.bundle.CACertPEM) + eventually(t, "caBundle after removing the opt-out", func() bool { return caBundleOf(h.get(t)) == want }) +} + +func TestSyncWaitsForSecretThenPatchesOnEnqueue(t *testing.T) { + h := newHarness(t, newAPIService("", true, nil), false) + h.run(t) + consistently(t, "patch count without Secret", 300*time.Millisecond, func() bool { return h.patches.Load() == 0 }) + b, err := servingcert.Generate(testNS, time.Now()) + if err != nil { + t.Fatal(err) + } + h.setSecret(t, b) + h.ctrl.Enqueue() // what servingcert.Manager does after creating the Secret + want := base64.StdEncoding.EncodeToString(b.CACertPEM) + eventually(t, "caBundle after the Secret appears", func() bool { return caBundleOf(h.get(t)) == want }) +} + +func TestSyncNeverWritesAnInvalidSecret(t *testing.T) { + h := newHarness(t, newAPIService("", true, nil), false) + if _, err := h.kube.CoreV1().Secrets(testNS).Create(context.Background(), &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: servingcert.SecretName, Namespace: testNS}, + Type: servingcert.SecretType, + Data: map[string][]byte{servingcert.CACertKey: []byte("-----BEGIN CERTIFICATE-----\nnot a cert\n-----END CERTIFICATE-----\n")}, + }, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + h.run(t) + h.ctrl.Enqueue() + consistently(t, "patch count with a corrupt Secret", 500*time.Millisecond, func() bool { return h.patches.Load() == 0 }) +} + +func TestSyncWaitsForAPIServiceToBeCreated(t *testing.T) { + h := newHarness(t, nil, true) + h.run(t) + consistently(t, "patch count without APIService", 300*time.Millisecond, func() bool { return h.patches.Load() == 0 }) + if _, err := h.dyn.Resource(APIServiceGVR).Create(context.Background(), newAPIService("", true, nil), metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + want := base64.StdEncoding.EncodeToString(h.bundle.CACertPEM) + eventually(t, "caBundle after creation", func() bool { return caBundleOf(h.get(t)) == want }) +} + +func TestSyncRepairsDrift(t *testing.T) { + h := newHarness(t, newAPIService("", true, nil), true) + h.run(t) + want := base64.StdEncoding.EncodeToString(h.bundle.CACertPEM) + eventually(t, "initial caBundle", func() bool { return caBundleOf(h.get(t)) == want }) + u := h.get(t) + if err := unstructured.SetNestedField(u.Object, base64.StdEncoding.EncodeToString([]byte("wrong")), "spec", "caBundle"); err != nil { + t.Fatal(err) + } + if _, err := h.dyn.Resource(APIServiceGVR).Update(context.Background(), u, metav1.UpdateOptions{}); err != nil { + t.Fatal(err) + } + eventually(t, "drift repair", func() bool { return caBundleOf(h.get(t)) == want }) +} + +func TestSyncFollowsCAChangesSignalledByTheManager(t *testing.T) { + h := newHarness(t, newAPIService("", true, nil), true) + h.run(t) + eventually(t, "initial caBundle", func() bool { return h.patches.Load() == 1 }) + rotated, err := servingcert.Generate(testNS, time.Now()) + if err != nil { + t.Fatal(err) + } + h.setSecret(t, rotated) + h.ctrl.Enqueue() + want := base64.StdEncoding.EncodeToString(rotated.CACertPEM) + eventually(t, "rotated caBundle", func() bool { return caBundleOf(h.get(t)) == want }) +} + +// TestSyncWritesTheSecretCAOnly models an old replica during a CA rotation: whatever a pod holds +// in memory, every controller writes the CA that is in the Secret now, so replicas cannot flap. +func TestSyncWritesTheSecretCAOnly(t *testing.T) { + current, err := servingcert.Generate(testNS, time.Now()) + if err != nil { + t.Fatal(err) + } + want := base64.StdEncoding.EncodeToString(current.CACertPEM) + h := newHarness(t, newAPIService(want, false, nil), false) + h.setSecret(t, current) + h.run(t) + for i := 0; i < 5; i++ { + h.ctrl.Enqueue() + } + consistently(t, "caBundle", 500*time.Millisecond, func() bool { return caBundleOf(h.get(t)) == want && h.patches.Load() == 0 }) +} + +func TestSyncRetriesAfterForbiddenWithoutGivingUp(t *testing.T) { + h := newHarness(t, newAPIService("", true, nil), true) + var forbidden atomic.Bool + forbidden.Store(true) + h.dyn.PrependReactor("patch", "apiservices", func(k8stesting.Action) (bool, runtime.Object, error) { + if forbidden.Load() { + h.patches.Add(1) // this reactor runs before the counting one + return true, nil, apierrors.NewForbidden(APIServiceGVR.GroupResource(), APIServiceName, errFake("rbac")) + } + return false, nil, nil + }) + h.run(t) + eventually(t, "retries", func() bool { return h.patches.Load() >= 3 }) + forbidden.Store(false) + want := base64.StdEncoding.EncodeToString(h.bundle.CACertPEM) + eventually(t, "patch after RBAC is fixed", func() bool { return caBundleOf(h.get(t)) == want }) +} + +func TestMissingListPermissionIsReported(t *testing.T) { + h := newHarness(t, newAPIService("", true, nil), true) + h.dyn.PrependReactor("list", "apiservices", func(k8stesting.Action) (bool, runtime.Object, error) { + return true, nil, apierrors.NewForbidden(APIServiceGVR.GroupResource(), "", errFake("rbac")) + }) + h.run(t) + eventually(t, "a missing-permission warning", func() bool { + h.ctrl.warnMu.Lock() + defer h.ctrl.warnMu.Unlock() + _, ok := h.ctrl.lastWarn["forbidden-watch"] + return ok + }) + if h.patches.Load() != 0 { + t.Fatal("nothing may be patched before the APIService can be read") + } +} + +func TestWarnIsRateLimitedPerKind(t *testing.T) { + c := &Controller{lastWarn: map[string]time.Time{}} + now := time.Unix(1000, 0) + c.now = func() time.Time { return now } + for i := 0; i < 3; i++ { + c.warn("forbidden", nil, "x") + } + if got := c.lastWarn["forbidden"]; !got.Equal(now) { + t.Fatalf("first warning must be recorded at %v, got %v", now, got) + } + now = now.Add(warnInterval - time.Second) + c.warn("forbidden", nil, "x") + if !c.lastWarn["forbidden"].Equal(time.Unix(1000, 0)) { + t.Fatal("a repeated warning inside the interval must be suppressed") + } + now = now.Add(2 * time.Second) + c.warn("forbidden", nil, "x") + if !c.lastWarn["forbidden"].Equal(now) { + t.Fatal("a warning after the interval must be logged again") + } + c.warn("absent", nil, "y") + if !c.lastWarn["absent"].Equal(now) { + t.Fatal("each kind of warning is limited separately") + } +} + +func TestNewAssertions(t *testing.T) { + if _, err := New(nil, fake.NewClientset(), testNS); err == nil { + t.Fatal("expected assertion for nil dynamic client") + } + dyn := dynamicfake.NewSimpleDynamicClientWithCustomListKinds(runtime.NewScheme(), map[schemaGVR]string{APIServiceGVR: "APIServiceList"}) + if _, err := New(dyn, fake.NewClientset(), ""); err == nil { + t.Fatal("expected assertion for empty namespace") + } +} + +type errFake string + +func (e errFake) Error() string { return string(e) } + +type schemaGVR = schema.GroupVersionResource diff --git a/internal/aggregated/apiservicetrust/envtest_test.go b/internal/aggregated/apiservicetrust/envtest_test.go new file mode 100644 index 00000000..41624015 --- /dev/null +++ b/internal/aggregated/apiservicetrust/envtest_test.go @@ -0,0 +1,221 @@ +package apiservicetrust + +import ( + "bytes" + "context" + "encoding/base64" + "errors" + "io" + "os" + "testing" + "time" + + corev1 "k8s.io/api/core/v1" + rbacv1 "k8s.io/api/rbac/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/fields" + "k8s.io/apimachinery/pkg/types" + utilyaml "k8s.io/apimachinery/pkg/util/yaml" + "k8s.io/client-go/dynamic" + "k8s.io/client-go/kubernetes" + "sigs.k8s.io/controller-runtime/pkg/envtest" + "sigs.k8s.io/yaml" + + "github.com/coder/coder-k8s/internal/aggregated/servingcert" +) + +const rbacManifest = "../../../config/rbac/apiservice-cabundle-role.yaml" + +// shippedClusterRole returns the ClusterRole from the manifest users apply, so the test checks +// the real least-privilege rules rather than a copy. +func shippedClusterRole(t *testing.T) *rbacv1.ClusterRole { + t.Helper() + data, err := os.ReadFile(rbacManifest) + if err != nil { + t.Fatal(err) + } + decoder := utilyaml.NewDocumentDecoder(io.NopCloser(bytes.NewReader(data))) + defer func() { _ = decoder.Close() }() + buf := make([]byte, len(data)+1) + for { + n, err := decoder.Read(buf) + if errors.Is(err, io.EOF) { + break + } + if err != nil { + t.Fatal(err) + } + var role rbacv1.ClusterRole + if err := yaml.Unmarshal(buf[:n], &role); err != nil { + t.Fatal(err) + } + if role.Kind == "ClusterRole" { + return &role + } + } + t.Fatalf("no ClusterRole in %s", rbacManifest) + return nil +} + +func TestEnvtestCABundleSyncWithShippedLeastPrivilegeRBAC(t *testing.T) { + if os.Getenv("KUBEBUILDER_ASSETS") == "" { + t.Fatal("KUBEBUILDER_ASSETS is not set; run via `make test`") + } + env := &envtest.Environment{} + adminCfg, err := env.Start() + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = env.Stop() }) + ctx := t.Context() + admin := kubernetes.NewForConfigOrDie(adminCfg) + adminDyn := dynamic.NewForConfigOrDie(adminCfg) + + if _, err := admin.CoreV1().Namespaces().Create(ctx, &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: testNS}}, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + manager, err := servingcert.NewManager(admin, testNS) + if err != nil { + t.Fatal(err) + } + bundle, err := manager.Ensure(ctx) + if err != nil { + t.Fatal(err) + } + want := base64.StdEncoding.EncodeToString(bundle.CACertPEM) + + other := newAPIService("", true, nil) + other.SetName("v1alpha1.other.example.com") + if err := unstructured.SetNestedField(other.Object, "other.example.com", "spec", "group"); err != nil { + t.Fatal(err) + } + for _, obj := range []*unstructured.Unstructured{newAPIService("", true, nil), other} { + if _, err := adminDyn.Resource(APIServiceGVR).Create(ctx, obj, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + } + + // The sync identity gets exactly the shipped ClusterRole, plus Secret read in its namespace + // (which manager-role grants in production). + user, err := env.AddUser(envtest.User{Name: "cabundle-sync"}, nil) + if err != nil { + t.Fatal(err) + } + role := shippedClusterRole(t) + role.ObjectMeta = metav1.ObjectMeta{Name: "test-apiservice-cabundle"} + if _, err := admin.RbacV1().ClusterRoles().Create(ctx, role, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + subject := []rbacv1.Subject{{Kind: rbacv1.UserKind, APIGroup: rbacv1.GroupName, Name: "cabundle-sync"}} + if _, err := admin.RbacV1().ClusterRoleBindings().Create(ctx, &rbacv1.ClusterRoleBinding{ + ObjectMeta: metav1.ObjectMeta{Name: "test-apiservice-cabundle"}, + RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "ClusterRole", Name: role.Name}, + Subjects: subject, + }, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + if _, err := admin.RbacV1().Roles(testNS).Create(ctx, &rbacv1.Role{ + ObjectMeta: metav1.ObjectMeta{Name: "test-secret-read"}, + Rules: []rbacv1.PolicyRule{{APIGroups: []string{""}, Resources: []string{"secrets"}, Verbs: []string{"get"}}}, + }, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + if _, err := admin.RbacV1().RoleBindings(testNS).Create(ctx, &rbacv1.RoleBinding{ + ObjectMeta: metav1.ObjectMeta{Name: "test-secret-read"}, + RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "Role", Name: "test-secret-read"}, + Subjects: subject, + }, metav1.CreateOptions{}); err != nil { + t.Fatal(err) + } + userKube := kubernetes.NewForConfigOrDie(user.Config()) + userDyn := dynamic.NewForConfigOrDie(user.Config()) + + // Least privilege: nothing beyond the one APIService. + byName := fields.OneTermEqualSelector("metadata.name", APIServiceName).String() + eventuallyErr(t, "RBAC to propagate", func() error { + _, err := userDyn.Resource(APIServiceGVR).List(ctx, metav1.ListOptions{FieldSelector: byName}) + return err + }) + if _, err := userDyn.Resource(APIServiceGVR).List(ctx, metav1.ListOptions{}); !apierrors.IsForbidden(err) { + t.Fatalf("list without the metadata.name selector must be forbidden, got %v", err) + } + if _, err := userDyn.Resource(APIServiceGVR).Get(ctx, other.GetName(), metav1.GetOptions{}); !apierrors.IsForbidden(err) { + t.Fatalf("get of another APIService must be forbidden, got %v", err) + } + if _, err := userDyn.Resource(APIServiceGVR).Patch(ctx, other.GetName(), types.MergePatchType, []byte(`{"spec":{"insecureSkipTLSVerify":false}}`), metav1.PatchOptions{}); !apierrors.IsForbidden(err) { + t.Fatalf("patch of another APIService must be forbidden, got %v", err) + } + + controller, err := New(userDyn, userKube, testNS) + if err != nil { + t.Fatal(err) + } + runCtx, cancel := context.WithCancel(ctx) + done := make(chan struct{}) + go func() { controller.Run(runCtx); close(done) }() + t.Cleanup(func() { cancel(); <-done }) + + get := func() *unstructured.Unstructured { + u, err := adminDyn.Resource(APIServiceGVR).Get(ctx, APIServiceName, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + return u + } + eventuallyTrue(t, 20*time.Second, "caBundle set and verification on", func() bool { + u := get() + _, hasInsecure, _ := unstructured.NestedBool(u.Object, "spec", "insecureSkipTLSVerify") + return caBundleOf(u) == want && !hasInsecure + }) + managed := false + for _, mf := range get().GetManagedFields() { + managed = managed || mf.Manager == FieldManager + } + if !managed { + t.Fatalf("managedFields must name %q", FieldManager) + } + + // Drift is repaired through the watch, without polling. + if _, err := adminDyn.Resource(APIServiceGVR).Patch(ctx, APIServiceName, types.MergePatchType, + []byte(`{"spec":{"caBundle":"`+base64.StdEncoding.EncodeToString([]byte("wrong"))+`"}}`), metav1.PatchOptions{}); err != nil { + t.Fatal(err) + } + eventuallyTrue(t, 10*time.Second, "drift repair", func() bool { return caBundleOf(get()) == want }) + + // The other APIService was never touched. + u, err := adminDyn.Resource(APIServiceGVR).Get(ctx, other.GetName(), metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + if !insecureOf(u) || caBundleOf(u) != "" { + t.Fatal("another APIService must stay unchanged") + } +} + +func eventuallyTrue(t *testing.T, timeout time.Duration, what string, cond func() bool) { + t.Helper() + deadline := time.Now().Add(timeout) + for !cond() { + if time.Now().After(deadline) { + t.Fatalf("timed out waiting for %s", what) + } + time.Sleep(100 * time.Millisecond) + } +} + +func eventuallyErr(t *testing.T, what string, fn func() error) { + t.Helper() + deadline := time.Now().Add(20 * time.Second) + for { + err := fn() + if err == nil { + return + } + if time.Now().After(deadline) { + t.Fatalf("timed out waiting for %s: %v", what, err) + } + time.Sleep(200 * time.Millisecond) + } +} diff --git a/internal/app/apiserverapp/apiserverapp.go b/internal/app/apiserverapp/apiserverapp.go index e167a45f..6026f603 100644 --- a/internal/app/apiserverapp/apiserverapp.go +++ b/internal/app/apiserverapp/apiserverapp.go @@ -30,6 +30,7 @@ import ( "k8s.io/kube-openapi/pkg/validation/spec" aggregationv1alpha1 "github.com/coder/coder-k8s/api/aggregation/v1alpha1" + "github.com/coder/coder-k8s/internal/aggregated/apiservicetrust" "github.com/coder/coder-k8s/internal/aggregated/coder" "github.com/coder/coder-k8s/internal/aggregated/servingcert" "github.com/coder/coder-k8s/internal/aggregated/storage" @@ -70,6 +71,9 @@ type Options struct { // 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 + // CABundleSync is a test seam. It keeps the APIService caBundle trusting ServingCert's CA; the + // production path creates it together with the managed serving certificate. + CABundleSync *apiservicetrust.Controller } type errClientProvider struct { @@ -370,7 +374,7 @@ func RunWithOptions(ctx context.Context, opts Options) error { secureServingOptions.ServerCert.CertDirectory = "" secureServingOptions.ServerCert.PairName = "" - servingCert := opts.ServingCert + servingCert, caBundleSync := opts.ServingCert, opts.CABundleSync authenticationOptions, authorizationOptions := opts.Authentication, opts.Authorization switch { case authenticationOptions == nil && authorizationOptions == nil: @@ -381,10 +385,13 @@ func RunWithOptions(ctx context.Context, opts Options) error { } authenticationOptions, authorizationOptions = newDelegatedAuthOptions(kubeconfigPath) if servingCert == nil { - servingCert, err = newServingCertManager(kubeconfigPath, serviceAccountNamespaceFile) + tlsSetup, err := newManagedTLS(kubeconfigPath, serviceAccountNamespaceFile) if err != nil { return fmt.Errorf("configure aggregated API server serving certificate: %w", err) } + if tlsSetup != nil { + servingCert, caBundleSync = tlsSetup.manager, tlsSetup.sync + } } case authenticationOptions == nil || authorizationOptions == nil: return fmt.Errorf("assertion failed: authentication and authorization options must be set together") @@ -396,6 +403,11 @@ func RunWithOptions(ctx context.Context, opts Options) error { } secureServingOptions.ServerCert.GeneratedCert = servingCert go servingCert.Run(ctx, servingcert.DefaultCheckInterval) + if caBundleSync != nil { + // CA changes reach the APIService without polling; failures never stop serving. + servingCert.AddListener(caBundleSync) + go caBundleSync.Run(ctx) + } } else { log.Printf("warning: no managed serving certificate (not running in a Pod); serving a self-signed certificate for localhost") } diff --git a/internal/app/apiserverapp/servingcert.go b/internal/app/apiserverapp/servingcert.go index 2996bd8c..aed40890 100644 --- a/internal/app/apiserverapp/servingcert.go +++ b/internal/app/apiserverapp/servingcert.go @@ -7,21 +7,30 @@ import ( "os" "strings" + "k8s.io/client-go/dynamic" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" "k8s.io/client-go/tools/clientcmd" + "github.com/coder/coder-k8s/internal/aggregated/apiservicetrust" "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) { +// managedTLS is the in-Pod TLS setup: the serving certificate manager and the controller that +// keeps the APIService caBundle trusting its CA. +type managedTLS struct { + manager *servingcert.Manager + sync *apiservicetrust.Controller +} + +// newManagedTLS returns the managed TLS setup 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 newManagedTLS(kubeconfigPath, namespaceFile string) (*managedTLS, 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 @@ -47,5 +56,17 @@ func newServingCertManager(kubeconfigPath, namespaceFile string) (*servingcert.M if err != nil { return nil, fmt.Errorf("build Kubernetes client: %w", err) } - return servingcert.NewManager(client, namespace) + dyn, err := dynamic.NewForConfig(cfg) + if err != nil { + return nil, fmt.Errorf("build Kubernetes dynamic client: %w", err) + } + manager, err := servingcert.NewManager(client, namespace) + if err != nil { + return nil, err + } + sync, err := apiservicetrust.New(dyn, client, namespace) + if err != nil { + return nil, err + } + return &managedTLS{manager: manager, sync: sync}, nil } diff --git a/internal/app/apiserverapp/servingcert_test.go b/internal/app/apiserverapp/servingcert_test.go index 5b7a8058..d45d6889 100644 --- a/internal/app/apiserverapp/servingcert_test.go +++ b/internal/app/apiserverapp/servingcert_test.go @@ -4,6 +4,7 @@ import ( "context" "crypto/tls" "crypto/x509" + "encoding/base64" "errors" "fmt" "io" @@ -16,8 +17,13 @@ import ( corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + dynamicfake "k8s.io/client-go/dynamic/fake" "k8s.io/client-go/kubernetes/fake" + "github.com/coder/coder-k8s/internal/aggregated/apiservicetrust" "github.com/coder/coder-k8s/internal/aggregated/servingcert" ) @@ -127,30 +133,30 @@ func TestRunWithOptionsFailsOnCorruptServingCertSecret(t *testing.T) { } } -func TestNewServingCertManager(t *testing.T) { +func TestNewManagedTLS(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) + if m, err := newManagedTLS("", filepath.Join(dir, "absent")); err != nil || m != nil { + t.Fatalf("outside a Pod: setup=%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") { + if _, err := newManagedTLS("", 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") { + if _, err := newManagedTLS(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) + m, err := newManagedTLS(newFakeKubeAPI(t).kubeconfigPath, nsFile) + if err != nil || m == nil || m.manager == nil || m.sync == nil { + t.Fatalf("in a Pod: setup=%+v err=%v; want both the certificate manager and the caBundle sync", 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 !strings.Contains(m.manager.Name(), servingTestNS+"/"+servingcert.SecretName) { + t.Fatalf("manager must target %s/%s, got %s", servingTestNS, servingcert.SecretName, m.manager.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") { + if _, err := newManagedTLS("", nsFile); err == nil || !strings.Contains(err.Error(), "read pod namespace") { t.Fatalf("expected read error, got %v", err) } } @@ -200,3 +206,72 @@ func waitForHealthz(t *testing.T, addr string, caPEM []byte, serverName string) time.Sleep(100 * time.Millisecond) } } + +// TestRunWithOptionsKeepsAPIServiceCABundleInSync checks the wiring: the caBundle sync starts with +// the managed certificate and follows CA changes signalled by the certificate manager. +func TestRunWithOptionsKeepsAPIServiceCABundleInSync(t *testing.T) { + client := fake.NewClientset() + manager := ensuredManager(t, client) + apiService := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "apiregistration.k8s.io/v1", "kind": "APIService", + "metadata": map[string]any{"name": apiservicetrust.APIServiceName}, + "spec": map[string]any{"insecureSkipTLSVerify": true}, + }} + dyn := dynamicfake.NewSimpleDynamicClientWithCustomListKinds(runtime.NewScheme(), + map[schema.GroupVersionResource]string{apiservicetrust.APIServiceGVR: "APIServiceList"}, apiService) + sync, err := apiservicetrust.New(dyn, client, servingTestNS) + if err != nil { + t.Fatal(err) + } + _, 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, CABundleSync: sync}) + }() + t.Cleanup(func() { + cancel() + select { + case <-errCh: + case <-time.After(10 * time.Second): + t.Error("timed out waiting for server shutdown") + } + }) + caBundle := func() string { + u, err := dyn.Resource(apiservicetrust.APIServiceGVR).Get(context.Background(), apiservicetrust.APIServiceName, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + v, _, _ := unstructured.NestedString(u.Object, "spec", "caBundle") + return v + } + waitFor := func(what, want string) { + t.Helper() + deadline := time.Now().Add(10 * time.Second) + for caBundle() != want { + if time.Now().After(deadline) { + t.Fatalf("timed out waiting for %s", what) + } + time.Sleep(50 * time.Millisecond) + } + } + waitFor("initial caBundle", base64.StdEncoding.EncodeToString(manager.CABundle())) + + rotated, err := servingcert.Generate(servingTestNS, time.Now()) + if err != nil { + t.Fatal(err) + } + secret, err := client.CoreV1().Secrets(servingTestNS).Get(ctx, servingcert.SecretName, metav1.GetOptions{}) + if err != nil { + t.Fatal(err) + } + secret.Data = rotated.Data() + if _, err := client.CoreV1().Secrets(servingTestNS).Update(ctx, secret, metav1.UpdateOptions{}); err != nil { + t.Fatal(err) + } + if _, err := manager.Ensure(ctx); err != nil { // adopts the new CA and notifies listeners + t.Fatal(err) + } + waitFor("rotated caBundle", base64.StdEncoding.EncodeToString(rotated.CACertPEM)) +}