Files
terdut-operator/internal/controller/terdutserver_controller_test.go
Niklas Ye 62664c93ff Keep the instance credential across a TerdutServer delete, and adopt it on recreate
Deleting a TerdutServer removed the credential Secrets but never touched the
database, so a recreated one found a server that was already bootstrapped and
no key for it: /api/bootstrap answered 403 and the operator stopped at
BootstrapStateLost, whose message and DESIGN.md both said "delete and
recreate". That is how the terdut-demo install on the cluster got stuck on
2026-10-03: Helm's cleanupOnFail deleted its TerdutServer after a failed
upgrade, the recreate found the bootstrapped database, and it sat at Ready:
False for five days until the database was reset by hand. Recreating cannot
fix it, because the finalizer clears Secrets and the database is not its to
reset, so "a fresh create starts clean" was only ever true when the database
went with it.

spec.credentials.deletionPolicy is Retain by default: the finalizer keeps the
instance credential Secret (Delete removes it, as before). The bootstrap
checkpoint is always removed. Before calling /api/bootstrap, reconcile now
looks for the retained Secret and asks the server for the operator's own
service account with its token. Accepted: adopt it and skip bootstrap.
Rejected with 401/403: the Secret outlived a database reset, so ignore it and
bootstrap like a first install, which replaces it. Any other error retries.
terdut-server's own tests already call that endpoint with an instance-scoped
key, so the permission is not new.

BootstrapStateLost is still the answer when the server is bootstrapped and no
credential it accepts survives, but its message now names the Secret to
restore and says that recreating does not clear the database. DESIGN.md §6
says the same, and the chart passes the setting through as
terdutServer.credentials.deletionPolicy.

A retained Secret of a TerdutServer that is gone for good is an orphan to
delete by hand. It is inert: nothing adopts it unless the server accepts the
token.

Checked on the kind demo with a locally built image against the real
terdut-server v0.43.0: deleting the TerdutServer kept the Secret, recreating it
reached Ready with the same credential (identical hash) and both TerdutTeams
came back Ready with their original ids. The controller specs cover adoption,
a rejected token after a reset, the bootstrapped-and-rejected failure, and
both deletion policies.

Co-authored-by: Claude <noreply@anthropic.com>
2026-10-08 21:09:56 +02:00

1031 lines
39 KiB
Go

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
// rejectedTokens: bearer tokens the fake answers 401 on GET
// /api/service-accounts, the way a server that never issued the key
// (a database reset since) would.
rejectedTokens map[string]bool
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:
if f.rejectedTokens[strings.TrimPrefix(r.Header.Get("Authorization"), "Bearer ")] {
writeJSON(w, http.StatusUnauthorized, map[string]string{errJSONKey: "invalid or expired API key"})
return
}
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() {
// runToReady brings a server to Ready against a fresh fake and
// returns the name of the instance credential Secret it minted.
runToReady := func(ctx SpecContext, spec terdutv1alpha1.TerdutServerSpec) string {
_, fakeSrv := newFakeTerdutServer()
DeferCleanup(fakeSrv.Close)
reconciler.NewClient = func(string) *tdclient.Client { return tdclient.New(fakeSrv.URL) }
createServer(ctx, spec)
reconcileOnce(ctx)
reconcileOnce(ctx)
markDeploymentReady(ctx)
reconcileOnce(ctx)
srv := &terdutv1alpha1.TerdutServer{}
Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed())
return srv.Status.CredentialsSecretRef.Name
}
deleteAndFinalize := func(ctx SpecContext) {
srv := &terdutv1alpha1.TerdutServer{}
Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed())
Expect(k8sClient.Delete(ctx, srv)).To(Succeed())
reconcileOnce(ctx) // runs the finalizer
Expect(k8sClient.Get(ctx, objKey, srv)).NotTo(Succeed(),
"the TerdutServer itself should be gone once the finalizer clears")
}
secretExists := func(ctx SpecContext, secretName string) bool {
var s corev1.Secret
err := k8sClient.Get(ctx, types.NamespacedName{Name: secretName, Namespace: operatorNamespace}, &s)
return err == nil
}
It("keeps the instance credential by default and removes the checkpoint", func(ctx SpecContext) {
credsName := runToReady(ctx, dsnSpec())
// A checkpoint left behind, as if the best-effort delete had failed.
Expect(k8sClient.Create(ctx, &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{Name: checkpointSecretNameFor(name), Namespace: operatorNamespace},
Data: map[string][]byte{credentialsSecretDataKey: []byte("admin-key-raw")},
})).To(Succeed())
deleteAndFinalize(ctx)
Expect(secretExists(ctx, credsName)).To(BeTrue(),
"the credential is kept so a recreated TerdutServer can adopt it")
Expect(secretExists(ctx, checkpointSecretNameFor(name))).To(BeFalse(),
"the checkpoint is a short-lived admin key and is always removed")
})
It("removes the credentials and checkpoint Secrets and the finalizer when the policy is Delete", func(ctx SpecContext) {
spec := dsnSpec()
spec.Credentials.DeletionPolicy = terdutv1alpha1.CredentialsDelete
credsName := runToReady(ctx, spec)
deleteAndFinalize(ctx)
Expect(secretExists(ctx, credsName)).To(BeFalse(),
"the credentials Secret should have been cleaned up by the finalizer")
})
})
Describe("recreating a TerdutServer against an already-bootstrapped server", func() {
// retainedSecret stands in for what the previous TerdutServer left.
retainedSecret := func(ctx SpecContext, token string) {
Expect(k8sClient.Create(ctx, &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{Name: credentialsSecretNameFor(name), Namespace: operatorNamespace},
Data: map[string][]byte{credentialsSecretDataKey: []byte(token)},
})).To(Succeed())
}
bringUp := func(ctx SpecContext, fakeSrv *httptest.Server) {
DeferCleanup(fakeSrv.Close)
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)
}
credsToken := func(ctx SpecContext) string {
var s corev1.Secret
Expect(k8sClient.Get(ctx, types.NamespacedName{Name: credentialsSecretNameFor(name), Namespace: operatorNamespace}, &s)).To(Succeed())
return string(s.Data[credentialsSecretDataKey])
}
It("adopts the retained credential when the server still accepts it, without bootstrapping", func(ctx SpecContext) {
fake, fakeSrv := newFakeTerdutServer()
fake.bootstrapped = true // every /api/bootstrap call 403s
fake.accounts[serviceAccountName] = 7
retainedSecret(ctx, "retained-key")
bringUp(ctx, fakeSrv)
Expect(readyCondition(ctx).Status).To(Equal(metav1.ConditionTrue))
srv := &terdutv1alpha1.TerdutServer{}
Expect(k8sClient.Get(ctx, objKey, srv)).To(Succeed())
Expect(srv.Status.CredentialsSecretRef).NotTo(BeNil())
Expect(srv.Status.CredentialsSecretRef.Name).To(Equal(credentialsSecretNameFor(name)))
Expect(credsToken(ctx)).To(Equal("retained-key"), "the retained key is reused, not replaced")
})
It("ignores a retained credential the server rejects and bootstraps afresh after a database reset", func(ctx SpecContext) {
fake, fakeSrv := newFakeTerdutServer() // a fresh database: bootstrap succeeds
fake.rejectedTokens = map[string]bool{"stale-key": true}
retainedSecret(ctx, "stale-key")
bringUp(ctx, fakeSrv)
Expect(readyCondition(ctx).Status).To(Equal(metav1.ConditionTrue))
Expect(credsToken(ctx)).NotTo(Equal("stale-key"), "the stale credential is replaced by a freshly minted one")
})
It("fails closed, naming the Secret to restore, when the server is bootstrapped and the retained credential is rejected", func(ctx SpecContext) {
fake, fakeSrv := newFakeTerdutServer()
fake.bootstrapped = true
fake.rejectedTokens = map[string]bool{"stale-key": true}
retainedSecret(ctx, "stale-key")
bringUp(ctx, fakeSrv)
cond := readyCondition(ctx)
Expect(cond.Status).To(Equal(metav1.ConditionFalse))
Expect(cond.Reason).To(Equal(terdutv1alpha1.ReasonBootstrapStateLost))
Expect(cond.Message).To(ContainSubstring(credentialsSecretNameFor(name)),
"the message names the Secret that would restore it")
Expect(cond.Message).To(ContainSubstring("does not clear the database"))
})
})
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)
}