package controller import ( "context" "encoding/json" "fmt" "net/http" "net/http/httptest" "strconv" "strings" "sync" "time" . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" policyv1 "k8s.io/api/policy/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/intstr" 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" ) // fakeVersionString/errJSONKey are shared by every response fakeTerdutServer // writes -- goconst would otherwise flag "test" and "error" as repeated // literals across its handlers. const ( fakeVersionString = "test" errJSONKey = "error" // deadmanSwitchesPath is the literal path (not Printf'd like the others // below) shared by the exact-match collection route and the dispatcher // that routes into it -- goconst flags three occurrences of the same // string, so this is that string, once. deadmanSwitchesPath = "/deadman/switches" ) // 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 nextTeamID int64 teams map[string]int64 // name -> id teamNames map[int64]string // id -> current name (renames update this) teamOIDC map[int64][2]string teamDelete map[int64]bool // id -> true once DELETEd, for 404-on-redelete // users backs GET /api/users for TerdutEscalationRule's username // resolution (DESIGN.md §4.3) -- a fixed, pre-seeded directory, since // nothing in this controller's own flow ever creates a user. users map[string]int64 // username -> id // escalation backs PUT /api/teams/{id}/escalation -- an upsert // server-side (confirmed against source), so this is just "the last // body PUT for this team", keyed by teamID, with no separate create // step to model. escalation map[int64]tdclient.SetEscalationRequest // switches/nextSwitchID/switchDelete back the dead man's switch // endpoints -- no unique-name constraint server-side (DESIGN.md §4.4), // so switches is keyed by id, not name, same as the real API's own // GET-list-and-match-by-name idempotent-create shape requires. nextSwitchID int64 switches map[int64]map[int64]tdclient.DeadmanSwitch // teamID -> switchID -> switch switchDelete map[int64]bool // switchID -> true once DELETEd, for 404-on-redelete // integrations/nextIntegrationID/integrationDelete back the alert // source endpoints -- no unique-name constraint server-side either // (DESIGN.md §4.5), same shape as switches, keyed by id. nextIntegrationID int64 integrations map[int64]map[int64]tdclient.Integration // teamID -> integrationID -> integration integrationDelete map[int64]bool // integrationID -> true once DELETEd, for 404-on-redelete // invites/nextInviteID/inviteDelete back TerdutTeam's own invite-minting // feature -- no unique constraint on an invite server-side either (every // POST mints a brand new row, confirmed against source), same keyed-by-id // shape as switches/integrations. nextInviteID int64 invites map[int64]map[int64]tdclient.Invite // teamID -> inviteID -> invite inviteDelete map[int64]bool // inviteID -> true once DELETEd, for 404-on-redelete } func newFakeTerdutServer() (*fakeTerdutServer, *httptest.Server) { f := &fakeTerdutServer{ accounts: map[string]int64{}, keyMints: map[int64]int{}, teams: map[string]int64{}, teamNames: map[int64]string{}, teamOIDC: map[int64][2]string{}, teamDelete: map[int64]bool{}, users: map[string]int64{}, escalation: map[int64]tdclient.SetEscalationRequest{}, switches: map[int64]map[int64]tdclient.DeadmanSwitch{}, switchDelete: map[int64]bool{}, integrations: map[int64]map[int64]tdclient.Integration{}, integrationDelete: map[int64]bool{}, invites: map[int64]map[int64]tdclient.Invite{}, inviteDelete: map[int64]bool{}, } return f, httptest.NewServer(f) } // seedUser registers a username the fake GET /api/users will return -- // called from test setup, before the controller under test ever runs. // Callers read the assigned id back from f.users themselves, so this has // nothing left to return. func (f *fakeTerdutServer) seedUser(username string) { f.mu.Lock() defer f.mu.Unlock() f.nextID++ f.users[username] = f.nextID } 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": fakeVersionString}) case r.URL.Path == "/api/bootstrap" && r.Method == http.MethodPost: if f.bootstrapped || f.bootstrap403 { writeJSON(w, http.StatusForbidden, map[string]string{errJSONKey: "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{errJSONKey: "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"}}) case r.URL.Path == "/api/teams" && r.Method == http.MethodPost: var req struct { Name string `json:"name"` } _ = json.NewDecoder(r.Body).Decode(&req) if _, exists := f.teams[req.Name]; exists { writeJSON(w, http.StatusConflict, map[string]string{errJSONKey: "a team with that name already exists"}) return } f.nextTeamID++ id := f.nextTeamID f.teams[req.Name] = id f.teamNames[id] = req.Name writeJSON(w, http.StatusCreated, tdclient.Team{ID: id, Name: req.Name}) case r.URL.Path == "/api/teams" && r.Method == http.MethodGet: // The controller never calls this without ?name= (TEAM-LOOKUP.md's // own lookup shape) -- the fake only needs to answer that form. name := r.URL.Query().Get("name") id, exists := f.teams[name] if !exists { writeJSON(w, http.StatusOK, []tdclient.Team{}) return } writeJSON(w, http.StatusOK, []tdclient.Team{{ID: id, Name: f.teamNames[id]}}) case r.URL.Path == "/api/users" && r.Method == http.MethodGet: // No query filter -- GetUserByUsername fetches the whole list and // matches client-side (confirmed against source: no server-side // filter either), so the fake does the same. users := make([]tdclient.User, 0, len(f.users)) for name, id := range f.users { users = append(users, tdclient.User{ID: id, Username: name}) } writeJSON(w, http.StatusOK, users) 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 } if id, rest, ok := parseTeamSubPath(r.URL.Path); ok { f.handleTeamSubPath(w, r, id, rest) return } w.WriteHeader(http.StatusNotFound) } } // handleTeamSubPath answers everything under /api/teams/{id}: PUT (rename), // DELETE, PUT .../oidc-groups, PUT .../escalation, and the dead man's // switch collection/item endpoints. rest is whatever parseTeamSubPath found // after "/api/teams/{id}" -- "" for the bare resource. func (f *fakeTerdutServer) handleTeamSubPath(w http.ResponseWriter, r *http.Request, id int64, rest string) { switch { case rest == "" && r.Method == http.MethodPut: var req struct { Name string `json:"name"` } _ = json.NewDecoder(r.Body).Decode(&req) oldName, exists := f.teamNames[id] if !exists { w.WriteHeader(http.StatusNotFound) return } delete(f.teams, oldName) f.teamNames[id] = req.Name f.teams[req.Name] = id w.WriteHeader(http.StatusNoContent) case rest == "" && r.Method == http.MethodDelete: name, exists := f.teamNames[id] if !exists { w.WriteHeader(http.StatusNotFound) return } delete(f.teams, name) delete(f.teamNames, id) f.teamDelete[id] = true w.WriteHeader(http.StatusNoContent) case rest == "/oidc-groups" && r.Method == http.MethodPut: var req struct { MemberGroup string `json:"member_group"` OwnerGroup string `json:"owner_group"` } _ = json.NewDecoder(r.Body).Decode(&req) if _, exists := f.teamNames[id]; !exists { w.WriteHeader(http.StatusNotFound) return } f.teamOIDC[id] = [2]string{req.MemberGroup, req.OwnerGroup} w.WriteHeader(http.StatusNoContent) case rest == "/escalation" && r.Method == http.MethodPut: var req tdclient.SetEscalationRequest _ = json.NewDecoder(r.Body).Decode(&req) f.escalation[id] = req w.WriteHeader(http.StatusNoContent) case rest == deadmanSwitchesPath || strings.HasPrefix(rest, deadmanSwitchesPath+"/"): f.handleDeadmanSubPath(w, r, id, rest) case rest == "/integrations" || strings.HasPrefix(rest, "/integrations/"): f.handleIntegrationSubPath(w, r, id, rest) case rest == "/invites" || strings.HasPrefix(rest, "/invites/"): f.handleInviteSubPath(w, r, id, rest) default: w.WriteHeader(http.StatusNotFound) } } // handleDeadmanSubPath answers GET/POST /api/teams/{id}/deadman/switches and // PUT/DELETE .../deadman/switches/{switchID} -- split out of // handleTeamSubPath for the same gocyclo reason as handleIntegrationSubPath. func (f *fakeTerdutServer) handleDeadmanSubPath(w http.ResponseWriter, r *http.Request, id int64, rest string) { switch { case rest == deadmanSwitchesPath && r.Method == http.MethodGet: existing := f.switches[id] out := make([]tdclient.DeadmanSwitch, 0, len(existing)) for _, s := range existing { out = append(out, s) } writeJSON(w, http.StatusOK, out) case rest == deadmanSwitchesPath && r.Method == http.MethodPost: var req deadmanSwitchFakeRequest _ = json.NewDecoder(r.Body).Decode(&req) f.nextSwitchID++ switchID := f.nextSwitchID name := req.Name if name == "" { name = "derived-" + req.Matcher } sw := tdclient.DeadmanSwitch{ ID: switchID, Name: name, Matcher: req.Matcher, TimeoutSeconds: req.TimeoutSeconds, Severity: req.Severity, } if f.switches[id] == nil { f.switches[id] = map[int64]tdclient.DeadmanSwitch{} } f.switches[id][switchID] = sw writeJSON(w, http.StatusCreated, sw) case strings.HasPrefix(rest, "/deadman/switches/") && r.Method == http.MethodPut: switchID, ok := parseTrailingID(rest, "/deadman/switches/") if !ok { w.WriteHeader(http.StatusNotFound) return } if _, exists := f.switches[id][switchID]; !exists { w.WriteHeader(http.StatusNotFound) return } var req deadmanSwitchFakeRequest _ = json.NewDecoder(r.Body).Decode(&req) name := req.Name if name == "" { name = f.switches[id][switchID].Name } f.switches[id][switchID] = tdclient.DeadmanSwitch{ ID: switchID, Name: name, Matcher: req.Matcher, TimeoutSeconds: req.TimeoutSeconds, Severity: req.Severity, } w.WriteHeader(http.StatusNoContent) case strings.HasPrefix(rest, "/deadman/switches/") && r.Method == http.MethodDelete: switchID, ok := parseTrailingID(rest, "/deadman/switches/") if !ok { w.WriteHeader(http.StatusNotFound) return } if _, exists := f.switches[id][switchID]; !exists { w.WriteHeader(http.StatusNotFound) return } delete(f.switches[id], switchID) f.switchDelete[switchID] = true w.WriteHeader(http.StatusNoContent) default: w.WriteHeader(http.StatusNotFound) } } // handleInviteSubPath answers POST /api/teams/{id}/invites and // DELETE .../invites/{inviteID} -- split out for the same gocyclo reason as // handleIntegrationSubPath. func (f *fakeTerdutServer) handleInviteSubPath(w http.ResponseWriter, r *http.Request, id int64, rest string) { switch { case rest == "/invites" && r.Method == http.MethodPost: var req struct { Role string `json:"role"` MaxUses int64 `json:"max_uses"` } _ = json.NewDecoder(r.Body).Decode(&req) f.nextInviteID++ inviteID := f.nextInviteID inv := tdclient.Invite{ ID: inviteID, TeamID: id, Role: req.Role, MaxUses: req.MaxUses, ExpiresAt: time.Now().Add(7 * 24 * time.Hour), URL: fmt.Sprintf("https://terdut.example.invalid/signup?invite=invite-token-%d", inviteID), } if f.invites[id] == nil { f.invites[id] = map[int64]tdclient.Invite{} } f.invites[id][inviteID] = inv writeJSON(w, http.StatusCreated, inv) case strings.HasPrefix(rest, "/invites/") && r.Method == http.MethodDelete: inviteID, ok := parseTrailingID(rest, "/invites/") if !ok { w.WriteHeader(http.StatusNotFound) return } if _, exists := f.invites[id][inviteID]; !exists { w.WriteHeader(http.StatusNotFound) return } delete(f.invites[id], inviteID) f.inviteDelete[inviteID] = true w.WriteHeader(http.StatusNoContent) default: w.WriteHeader(http.StatusNotFound) } } // handleIntegrationSubPath answers POST /api/teams/{id}/integrations, // PATCH .../integrations/{integrationID} and DELETE .../integrations/{integrationID} // -- split out of handleTeamSubPath so that switch's own cyclomatic // complexity stays under golangci-lint's gocyclo threshold. func (f *fakeTerdutServer) handleIntegrationSubPath(w http.ResponseWriter, r *http.Request, id int64, rest string) { switch { case rest == "/integrations" && r.Method == http.MethodPost: var req struct { Name string `json:"name"` Kind string `json:"kind"` } _ = json.NewDecoder(r.Body).Decode(&req) f.nextIntegrationID++ integID := f.nextIntegrationID integ := tdclient.Integration{ ID: integID, TeamID: id, Kind: req.Kind, Name: req.Name, // Key/URL are only ever in *this* response -- never again, // matching terdut-server's own one-time-show semantics // (DESIGN.md §4.5) -- so what's stored for later GET/PATCH // calls in this fake deliberately omits them too. Key: fmt.Sprintf("webhook-key-%d", integID), URL: fmt.Sprintf("https://terdut.example.invalid/api/integrations/webhook-key-%d/%s", integID, req.Kind), } if f.integrations[id] == nil { f.integrations[id] = map[int64]tdclient.Integration{} } f.integrations[id][integID] = tdclient.Integration{ID: integID, TeamID: id, Kind: req.Kind, Name: req.Name} writeJSON(w, http.StatusCreated, integ) case strings.HasPrefix(rest, "/integrations/") && r.Method == http.MethodPatch: integID, ok := parseTrailingID(rest, "/integrations/") if !ok { w.WriteHeader(http.StatusNotFound) return } existing, exists := f.integrations[id][integID] if !exists { w.WriteHeader(http.StatusNotFound) return } var req struct { Name string `json:"name"` } _ = json.NewDecoder(r.Body).Decode(&req) existing.Name = req.Name f.integrations[id][integID] = existing w.WriteHeader(http.StatusNoContent) case strings.HasPrefix(rest, "/integrations/") && r.Method == http.MethodDelete: integID, ok := parseTrailingID(rest, "/integrations/") if !ok { w.WriteHeader(http.StatusNotFound) return } if _, exists := f.integrations[id][integID]; !exists { w.WriteHeader(http.StatusNotFound) return } delete(f.integrations[id], integID) f.integrationDelete[integID] = true w.WriteHeader(http.StatusNoContent) default: w.WriteHeader(http.StatusNotFound) } } // deadmanSwitchFakeRequest mirrors tdclient's own (unexported) // deadmanSwitchRequest -- the fake needs its own copy to decode the same // wire shape without reaching across package boundaries for an internal type. type deadmanSwitchFakeRequest struct { Name string `json:"name,omitempty"` Matcher string `json:"matcher"` TimeoutSeconds int64 `json:"timeout_seconds"` Severity string `json:"severity"` } // parseTeamSubPath splits "/api/teams/{id}" from anything after it -- // "" for an exact match, "/oidc-groups", "/escalation", "/deadman/switches" // or "/deadman/switches/{switchID}" otherwise. Doesn't itself validate the // suffix; handleTeamSubPath's own switch does that. func parseTeamSubPath(path string) (id int64, rest string, ok bool) { const prefix = "/api/teams/" if !strings.HasPrefix(path, prefix) { return 0, "", false } trimmed := path[len(prefix):] parts := strings.SplitN(trimmed, "/", 2) parsedID, err := strconv.ParseInt(parts[0], 10, 64) if err != nil { return 0, "", false } if len(parts) == 1 { return parsedID, "", true } return parsedID, "/" + parts[1], true } // parseTrailingID parses the numeric id after prefix within rest, e.g. // parseTrailingID("/deadman/switches/7", "/deadman/switches/") -> 7, true. func parseTrailingID(rest, prefix string) (id int64, ok bool) { parsedID, err := strconv.ParseInt(strings.TrimPrefix(rest, prefix), 10, 64) if err != nil { return 0, false } return parsedID, true } 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: testImageRepo, Tag: testImageTag}, Networking: terdutv1alpha1.NetworkingSpec{Hostname: "terdut.example.invalid", ServicePort: 8080}, Database: terdutv1alpha1.DatabaseSpec{DSN: testDSN}, } } 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", testDSN)) 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(writeOperatorSecret(ctx, k8sClient, operatorNamespace, 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("spec.pod", func() { It("wires pod-level customization onto the right spot on the Deployment", func(ctx SpecContext) { fake, fakeSrv := newFakeTerdutServer() _ = fake DeferCleanup(fakeSrv.Close) reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) } spec := dsnSpec() qty := resource.MustParse("250m") spec.Pod = terdutv1alpha1.PodSpec{ Resources: corev1.ResourceRequirements{Requests: corev1.ResourceList{corev1.ResourceCPU: qty}}, Tolerations: []corev1.Toleration{{Key: "dedicated", Operator: corev1.TolerationOpEqual, Value: "terdut", Effect: corev1.TaintEffectNoSchedule}}, ExtraEnv: []corev1.EnvVar{{Name: "EXTRA_FLAG", Value: "on"}}, ServiceAccountName: "terdut-server-custom", ExtraVolumes: []corev1.Volume{{Name: "extra-ca", VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{}}}}, ExtraVolumeMounts: []corev1.VolumeMount{{Name: "extra-ca", MountPath: "/etc/extra-ca"}}, } createServer(ctx, spec) reconcileOnce(ctx) // finalizer reconcileOnce(ctx) // Deployment/Service/PDB var deploy appsv1.Deployment Expect(k8sClient.Get(ctx, objKey, &deploy)).To(Succeed()) podSpec := deploy.Spec.Template.Spec Expect(podSpec.Tolerations).To(ConsistOf(spec.Pod.Tolerations)) Expect(podSpec.ServiceAccountName).To(Equal("terdut-server-custom")) Expect(podSpec.Volumes).To(ConsistOf(spec.Pod.ExtraVolumes)) main := podSpec.Containers[0] Expect(main.Name).To(Equal("terdut-server")) Expect(main.Resources).To(Equal(spec.Pod.Resources)) Expect(main.VolumeMounts).To(ConsistOf(spec.Pod.ExtraVolumeMounts)) Expect(main.Env).To(ContainElement(corev1.EnvVar{Name: "EXTRA_FLAG", Value: "on"})) initContainer := podSpec.InitContainers[0] Expect(initContainer.Name).To(Equal("wait-for-postgres")) Expect(initContainer.VolumeMounts).To(BeEmpty(), "extraVolumeMounts must not leak onto wait-for-postgres") }) }) Describe("spec.pod.disruptionBudget", func() { It("creates an owned PodDisruptionBudget when set, and deletes it once cleared", func(ctx SpecContext) { fake, fakeSrv := newFakeTerdutServer() _ = fake DeferCleanup(fakeSrv.Close) reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) } spec := dsnSpec() minAvail := intstr.FromInt32(1) spec.Pod.DisruptionBudget = &terdutv1alpha1.PodDisruptionBudgetSpec{MinAvailable: &minAvail} createServer(ctx, spec) reconcileOnce(ctx) // finalizer reconcileOnce(ctx) // Deployment/Service/PDB var pdb policyv1.PodDisruptionBudget Expect(k8sClient.Get(ctx, objKey, &pdb)).To(Succeed()) Expect(pdb.Spec.Selector.MatchLabels).To(Equal(labelsFor(&terdutv1alpha1.TerdutServer{ObjectMeta: metav1.ObjectMeta{Name: name}}))) Expect(pdb.Spec.MinAvailable).To(Equal(&minAvail)) Expect(pdb.OwnerReferences).To(ContainElement(HaveField("Name", name))) srv := &terdutv1alpha1.TerdutServer{} Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed()) srv.Spec.Pod.DisruptionBudget = nil Expect(k8sClient.Update(ctx, srv)).To(Succeed()) reconcileOnce(ctx) err := k8sClient.Get(ctx, objKey, &policyv1.PodDisruptionBudget{}) Expect(apierrors.IsNotFound(err)).To(BeTrue(), "PodDisruptionBudget should be deleted once spec.pod.disruptionBudget is cleared") }) It("rejects both minAvailable and maxUnavailable set together, and neither set", func(ctx SpecContext) { bothSet := dsnSpec() minAvail, maxUnavail := intstr.FromInt32(1), intstr.FromInt32(1) bothSet.Pod.DisruptionBudget = &terdutv1alpha1.PodDisruptionBudgetSpec{MinAvailable: &minAvail, MaxUnavailable: &maxUnavail} Expect(k8sClient.Create(ctx, &terdutv1alpha1.TerdutServer{ ObjectMeta: metav1.ObjectMeta{Name: name + "-both", Namespace: operatorNamespace}, Spec: bothSet, })).To(HaveOccurred()) neitherSet := dsnSpec() neitherSet.Pod.DisruptionBudget = &terdutv1alpha1.PodDisruptionBudgetSpec{} Expect(k8sClient.Create(ctx, &terdutv1alpha1.TerdutServer{ ObjectMeta: metav1.ObjectMeta{Name: name + "-neither", Namespace: operatorNamespace}, Spec: neitherSet, })).To(HaveOccurred()) }) }) 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) }