package controller import ( "context" "encoding/json" "fmt" "net/http" "net/http/httptest" "sync" . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/types" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/reconcile" terdutv1alpha1 "git.ryuvia.com/niklas/terdut-operator/api/v1alpha1" "git.ryuvia.com/niklas/terdut-operator/internal/tdclient" ) // fakeTerdutServer reproduces the exact stateful semantics of // /api/bootstrap, /api/service-accounts and /api/version that the // bootstrap flow depends on (DESIGN.md §6, §11: "terdut-server's REST API // is faked with a small httptest.Server per controller test ... matching // the real handlers' request/response shapes"), including the 403-after- // first-success and 409-on-name-conflict behavior the adopt-on-conflict // recovery path exists for. type fakeTerdutServer struct { mu sync.Mutex bootstrapped bool bootstrap403 bool // force every /api/bootstrap call to 403, even the first nextID int64 accounts map[string]int64 // name -> id keyMints map[int64]int // id -> number of keys minted so far } func newFakeTerdutServer() (*fakeTerdutServer, *httptest.Server) { f := &fakeTerdutServer{accounts: map[string]int64{}, keyMints: map[int64]int{}} return f, httptest.NewServer(f) } func (f *fakeTerdutServer) ServeHTTP(w http.ResponseWriter, r *http.Request) { f.mu.Lock() defer f.mu.Unlock() switch { case r.URL.Path == "/api/version": writeJSON(w, http.StatusOK, map[string]string{"version": "test"}) case r.URL.Path == "/api/bootstrap" && r.Method == http.MethodPost: if f.bootstrapped || f.bootstrap403 { writeJSON(w, http.StatusForbidden, map[string]string{"error": "bootstrap already completed"}) return } f.bootstrapped = true writeJSON(w, http.StatusCreated, tdclient.BootstrapResult{ APIKey: tdclient.APIKey{ID: 1, Name: "bootstrap", Key: "admin-key-raw"}, }) case r.URL.Path == "/api/service-accounts" && r.Method == http.MethodPost: var req struct { Name string `json:"name"` Scope string `json:"scope"` } _ = json.NewDecoder(r.Body).Decode(&req) if _, exists := f.accounts[req.Name]; exists { writeJSON(w, http.StatusConflict, map[string]string{"error": "a service account with that name already exists"}) return } f.nextID++ id := f.nextID f.accounts[req.Name] = id f.keyMints[id] = 1 writeJSON(w, http.StatusCreated, tdclient.CreateServiceAccountResult{ ServiceAccount: tdclient.ServiceAccount{ID: id, Name: req.Name, Scope: req.Scope}, Key: tdclient.APIKey{ID: 1, Name: "initial", Key: fmt.Sprintf("instance-key-%d-initial", id)}, }) case r.URL.Path == "/api/service-accounts" && r.Method == http.MethodGet: name := r.URL.Query().Get("name") id, exists := f.accounts[name] if !exists { writeJSON(w, http.StatusOK, []tdclient.ServiceAccount{}) return } writeJSON(w, http.StatusOK, []tdclient.ServiceAccount{{ID: id, Name: name, Scope: "instance"}}) default: if id, name, ok := parseKeysPath(r.URL.Path); ok && r.Method == http.MethodPost { f.keyMints[id]++ writeJSON(w, http.StatusCreated, tdclient.APIKey{ ID: int64(f.keyMints[id]), Name: name, Key: fmt.Sprintf("instance-key-%d-mint%d", id, f.keyMints[id]), }) return } w.WriteHeader(http.StatusNotFound) } } func parseKeysPath(path string) (id int64, mintName string, ok bool) { var parsedID int64 n, err := fmt.Sscanf(path, "/api/service-accounts/%d/keys", &parsedID) if err != nil || n != 1 { return 0, "", false } return parsedID, "minted", true } func writeJSON(w http.ResponseWriter, status int, v any) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(status) _ = json.NewEncoder(w).Encode(v) } var _ = Describe("TerdutServer Controller", func() { const operatorNamespace = "default" var ( reconciler *TerdutServerReconciler name string objKey types.NamespacedName ) BeforeEach(func() { reconciler = &TerdutServerReconciler{ Client: k8sClient, Scheme: k8sClient.Scheme(), OperatorNamespace: operatorNamespace, } name = fmt.Sprintf("test-server-%d-%d", GinkgoRandomSeed(), GinkgoParallelProcess()) objKey = types.NamespacedName{Name: name, Namespace: operatorNamespace} }) AfterEach(func(ctx SpecContext) { srv := &terdutv1alpha1.TerdutServer{} if err := k8sClient.Get(ctx, objKey, srv); err == nil { srv.Finalizers = nil _ = k8sClient.Update(ctx, srv) _ = k8sClient.Delete(ctx, srv) } for _, n := range []string{checkpointSecretNameFor(name), credentialsSecretNameFor(name)} { _ = k8sClient.Delete(ctx, &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: n, Namespace: operatorNamespace}}) } }) dsnSpec := func() terdutv1alpha1.TerdutServerSpec { return terdutv1alpha1.TerdutServerSpec{ Image: terdutv1alpha1.ImageSpec{Repository: "example.invalid/terdut-server", Tag: "test"}, Networking: terdutv1alpha1.NetworkingSpec{Hostname: "terdut.example.invalid", ServicePort: 8080}, Database: terdutv1alpha1.DatabaseSpec{DSN: "postgres://terdut@test-postgres:5432/terdut?sslmode=require"}, } } createServer := func(ctx context.Context, spec terdutv1alpha1.TerdutServerSpec) { srv := &terdutv1alpha1.TerdutServer{ ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: operatorNamespace}, Spec: spec, } Expect(k8sClient.Create(ctx, srv)).To(Succeed()) } reconcileOnce := func(ctx context.Context) ctrl.Result { res, err := reconciler.Reconcile(ctx, reconcile.Request{NamespacedName: objKey}) Expect(err).NotTo(HaveOccurred()) return res } markDeploymentReady := func(ctx context.Context) { var deploy appsv1.Deployment Expect(k8sClient.Get(ctx, objKey, &deploy)).To(Succeed()) deploy.Status.ReadyReplicas = 1 deploy.Status.Replicas = 1 Expect(k8sClient.Status().Update(ctx, &deploy)).To(Succeed()) } readyCondition := func(ctx context.Context) metav1.Condition { srv := &terdutv1alpha1.TerdutServer{} Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed()) c := meta.FindStatusCondition(srv.Status.Conditions, terdutv1alpha1.ConditionReady) Expect(c).NotTo(BeNil(), "Ready condition should always be set after a reconcile past the finalizer-add pass") return *c } Describe("bring-your-own DSN path", func() { It("reaches Ready through the full lifecycle: finalizer, Deployment/Service, wait, bootstrap", func(ctx SpecContext) { fake, fakeSrv := newFakeTerdutServer() DeferCleanup(fakeSrv.Close) reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) } createServer(ctx, dsnSpec()) // Pass 1: adds the finalizer and returns early. reconcileOnce(ctx) srv := &terdutv1alpha1.TerdutServer{} Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed()) Expect(srv.Finalizers).To(ContainElement(finalizerName)) // Pass 2: creates Deployment + Service, waits for readiness. reconcileOnce(ctx) Expect(readyCondition(ctx).Reason).To(Equal(terdutv1alpha1.ReasonWaitingForDeployment)) var deploy appsv1.Deployment Expect(k8sClient.Get(ctx, objKey, &deploy)).To(Succeed()) Expect(deploy.Spec.Template.Spec.Containers[0].Image).To(Equal("example.invalid/terdut-server:test")) envNames := map[string]string{} for _, e := range deploy.Spec.Template.Spec.Containers[0].Env { envNames[e.Name] = e.Value } Expect(envNames).To(HaveKeyWithValue("TERDUT_DB_DSN", "postgres://terdut@test-postgres:5432/terdut?sslmode=require")) Expect(envNames).To(HaveKeyWithValue("TERDUT_OPERATOR_MODE", "true")) var svc corev1.Service Expect(k8sClient.Get(ctx, objKey, &svc)).To(Succeed()) Expect(svc.Spec.Ports[0].Port).To(Equal(int32(8080))) // Simulate the Deployment becoming ready (envtest has no // kubelet/deployment-controller to do this for real). markDeploymentReady(ctx) // Pass 3: bootstraps for real against the fake server. reconcileOnce(ctx) Expect(fake.bootstrapped).To(BeTrue()) Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed()) ready := meta.FindStatusCondition(srv.Status.Conditions, terdutv1alpha1.ConditionReady) Expect(ready.Status).To(Equal(metav1.ConditionTrue)) Expect(ready.Reason).To(Equal(terdutv1alpha1.ReasonAdopted)) Expect(srv.Status.CredentialsSecretRef).NotTo(BeNil()) Expect(srv.Status.CredentialsSecretRef.Key).To(Equal(credentialsSecretDataKey)) var credsSecret corev1.Secret Expect(k8sClient.Get(ctx, types.NamespacedName{ Name: srv.Status.CredentialsSecretRef.Name, Namespace: operatorNamespace, }, &credsSecret)).To(Succeed()) Expect(string(credsSecret.Data[credentialsSecretDataKey])).To(Equal("instance-key-1-initial")) // The checkpoint is cleaned up once the lasting credential is // written (DESIGN.md §6 point 2). var checkpoint corev1.Secret err := k8sClient.Get(ctx, types.NamespacedName{Name: checkpointSecretName(srv), Namespace: operatorNamespace}, &checkpoint) Expect(err).To(HaveOccurred()) }) }) Describe("the adopt-on-409 recovery path", func() { It("mints a fresh key instead of erroring when the service account already exists", func(ctx SpecContext) { fake, fakeSrv := newFakeTerdutServer() DeferCleanup(fakeSrv.Close) fake.bootstrapped = true // an earlier attempt already bootstrapped... fake.nextID = 1 fake.accounts[serviceAccountName] = 1 // ...and already created the service account. fake.keyMints[1] = 1 reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) } createServer(ctx, dsnSpec()) reconcileOnce(ctx) // finalizer reconcileOnce(ctx) // Deployment/Service, WaitingForDeployment markDeploymentReady(ctx) // The earlier attempt's checkpoint survived (that's how this // reconcile can authenticate at all to recover). Expect(reconciler.writeSecret(ctx, checkpointSecretName(&terdutv1alpha1.TerdutServer{ ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: operatorNamespace}, }), "admin-key-raw")).To(Succeed()) reconcileOnce(ctx) srv := &terdutv1alpha1.TerdutServer{} Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed()) Expect(meta.FindStatusCondition(srv.Status.Conditions, terdutv1alpha1.ConditionReady).Status).To(Equal(metav1.ConditionTrue)) var credsSecret corev1.Secret Expect(k8sClient.Get(ctx, types.NamespacedName{ Name: srv.Status.CredentialsSecretRef.Name, Namespace: operatorNamespace, }, &credsSecret)).To(Succeed()) // Minted fresh, not the (never-seen-by-this-reconcile) "initial" // key from the account's original creation. Expect(string(credsSecret.Data[credentialsSecretDataKey])).To(Equal("instance-key-1-mint2")) }) }) Describe("the BootstrapStateLost path", func() { It("fails closed when the server reports already-bootstrapped with no checkpoint to recover from", func(ctx SpecContext) { fake, fakeSrv := newFakeTerdutServer() _ = fake DeferCleanup(fakeSrv.Close) fake.bootstrap403 = true reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) } createServer(ctx, dsnSpec()) reconcileOnce(ctx) // finalizer reconcileOnce(ctx) // Deployment/Service markDeploymentReady(ctx) reconcileOnce(ctx) cond := readyCondition(ctx) Expect(cond.Status).To(Equal(metav1.ConditionFalse)) Expect(cond.Reason).To(Equal(terdutv1alpha1.ReasonBootstrapStateLost)) }) }) Describe("the Zalando postgresClusterRef path", func() { zalandoSpec := func(clusterName string) terdutv1alpha1.TerdutServerSpec { s := dsnSpec() s.Database = terdutv1alpha1.DatabaseSpec{ PostgresClusterRef: &terdutv1alpha1.PostgresClusterRef{Name: clusterName}, } return s } It("waits with reason PostgresClusterNotFound when the postgresql CR doesn't exist yet", func(ctx SpecContext) { createServer(ctx, zalandoSpec("missing-cluster")) reconcileOnce(ctx) // finalizer reconcileOnce(ctx) Expect(readyCondition(ctx).Reason).To(Equal(terdutv1alpha1.ReasonPostgresClusterNotFound)) }) It("resolves a real DSN and reaches Ready once the CR and its generated Secret both exist", func(ctx SpecContext) { clusterName := "pg-" + name cluster := &unstructured.Unstructured{} cluster.SetGroupVersionKind(postgresqlGVK) cluster.SetName(clusterName) cluster.SetNamespace(operatorNamespace) Expect(unstructured.SetNestedField(cluster.Object, map[string]any{}, "spec")).To(Succeed()) Expect(k8sClient.Create(ctx, cluster)).To(Succeed()) DeferCleanup(func() { _ = k8sClient.Delete(ctx, cluster) }) credsSecretName := fmt.Sprintf("terdut.%s.credentials.postgresql.acid.zalan.do", clusterName) zalandoSecret := &corev1.Secret{ ObjectMeta: metav1.ObjectMeta{Name: credsSecretName, Namespace: operatorNamespace}, Data: map[string][]byte{"password": []byte("whatever")}, } Expect(k8sClient.Create(ctx, zalandoSecret)).To(Succeed()) DeferCleanup(func() { _ = k8sClient.Delete(ctx, zalandoSecret) }) fake, fakeSrv := newFakeTerdutServer() _ = fake DeferCleanup(fakeSrv.Close) reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) } createServer(ctx, zalandoSpec(clusterName)) reconcileOnce(ctx) // finalizer reconcileOnce(ctx) // Deployment/Service created with resolved DB env var deploy appsv1.Deployment Expect(k8sClient.Get(ctx, objKey, &deploy)).To(Succeed()) var dsn string for _, e := range deploy.Spec.Template.Spec.Containers[0].Env { if e.Name == "TERDUT_DB_DSN" { dsn = e.Value } } Expect(dsn).To(Equal(fmt.Sprintf("postgres://terdut@%s.%s.svc:5432/terdut?sslmode=require", clusterName, operatorNamespace))) markDeploymentReady(ctx) reconcileOnce(ctx) Expect(readyCondition(ctx).Status).To(Equal(metav1.ConditionTrue)) }) }) Describe("deletion", func() { It("removes the credentials and checkpoint Secrets and the finalizer", func(ctx SpecContext) { fake, fakeSrv := newFakeTerdutServer() _ = fake DeferCleanup(fakeSrv.Close) reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) } createServer(ctx, dsnSpec()) reconcileOnce(ctx) reconcileOnce(ctx) markDeploymentReady(ctx) reconcileOnce(ctx) srv := &terdutv1alpha1.TerdutServer{} Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed()) credsName := srv.Status.CredentialsSecretRef.Name Expect(k8sClient.Delete(ctx, srv)).To(Succeed()) reconcileOnce(ctx) // runs the finalizer err := k8sClient.Get(ctx, objKey, srv) Expect(err).To(HaveOccurred(), "the TerdutServer itself should be gone once the finalizer clears") var leftover corev1.Secret err = k8sClient.Get(ctx, types.NamespacedName{Name: credsName, Namespace: operatorNamespace}, &leftover) Expect(err).To(HaveOccurred(), "the credentials Secret should have been cleaned up by the finalizer") }) }) When("the TerdutServer object no longer exists", func() { It("returns no error (deleted between enqueue and reconcile)", func(ctx SpecContext) { _, err := reconciler.Reconcile(ctx, reconcile.Request{ NamespacedName: types.NamespacedName{Name: "never-created", Namespace: operatorNamespace}, }) Expect(err).NotTo(HaveOccurred()) }) }) }) // checkpointSecretNameFor/credentialsSecretNameFor let AfterEach clean up // without needing a live TerdutServer object (it may already be gone by // then in the deletion test). func checkpointSecretNameFor(name string) string { return fmt.Sprintf("default.%s-bootstrap-admin", name) } func credentialsSecretNameFor(name string) string { return fmt.Sprintf("default.%s-instance-credentials", name) }