From 78b8cf6ded6d798c78204d5a90daafb717b16d67 Mon Sep 17 00:00:00 2001 From: Minseong Choi Date: Fri, 26 Jun 2026 23:31:58 +0900 Subject: [PATCH] feat(operator): add MinecraftServer controller and reconcilers The Kubernetes controller that drives MinecraftServer resources through their lifecycle and issues RCON where readiness requires it. --- internal/operator/builders.go | 287 ++++++++++++++++++++++++++ internal/operator/prober.go | 41 ++++ internal/operator/reconciler.go | 296 ++++++++++++++++++++++++++ internal/operator/reconciler_test.go | 298 +++++++++++++++++++++++++++ 4 files changed, 922 insertions(+) create mode 100644 internal/operator/builders.go create mode 100644 internal/operator/prober.go create mode 100644 internal/operator/reconciler.go create mode 100644 internal/operator/reconciler_test.go diff --git a/internal/operator/builders.go b/internal/operator/builders.go new file mode 100644 index 0000000..8340c87 --- /dev/null +++ b/internal/operator/builders.go @@ -0,0 +1,287 @@ +package operator + +import ( + "fmt" + "strconv" + + "felis.lolicon.best/internal/apis/felis/v1alpha1" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/intstr" +) + +// Workload constants shared by the builders. +const ( + // GamePort is the Minecraft TCP port the proxy and readiness probe target. + GamePort int32 = 25565 + + // DefaultRconPort is the RCON port used when a server does not override + // spec.rcon.port (see rconPort). It is exported because the platform package's + // allow-rcon NetworkPolicy opens this port for the control plane — sharing the + // constant keeps the policy port and the container's default RCON port a single + // source of truth, so the operator's prober can always reach a default-port + // server through the fence. + DefaultRconPort int32 = 25575 + + containerName = "minecraft" + dataVolumeName = "world" + dataMountPath = "/data" + + // ManagedByValue / ComponentValue are the values of the LabelManagedBy / + // LabelComponent labels stamped on every per-server pod (see labelsFor). They + // are exported because the platform package's minecraft-namespace + // NetworkPolicies select server pods by exactly these labels — keeping the + // selector and the pod labels a single source of truth, so an isolation policy + // can never silently stop matching the pods it is meant to fence. + ManagedByValue = "felis-operator" + ComponentValue = "server" + + defaultGraceSeconds int64 = 300 + defaultStorageSize string = "8Gi" +) + +// selectorFor returns the immutable selector labels (a StatefulSet selector +// must never change after creation, so it carries only the server identity). +func selectorFor(server *v1alpha1.MinecraftServer) map[string]string { + return map[string]string{v1alpha1.LabelServer: server.Name} +} + +// labelsFor returns the full label set applied to managed objects. +func labelsFor(server *v1alpha1.MinecraftServer) map[string]string { + return map[string]string{ + v1alpha1.LabelServer: server.Name, + v1alpha1.LabelManagedBy: ManagedByValue, + v1alpha1.LabelComponent: ComponentValue, + } +} + +func headlessServiceName(name string) string { return name + "-hl" } + +// rconPort resolves the RCON port, defaulting to the conventional DefaultRconPort. +func rconPort(server *v1alpha1.MinecraftServer) int32 { + if server.Spec.Rcon.Port > 0 { + return server.Spec.Rcon.Port + } + return DefaultRconPort +} + +// graceSeconds resolves the pod termination grace period (spec §7). +func graceSeconds(server *v1alpha1.MinecraftServer) int64 { + if server.Spec.Lifecycle.TerminationGracePeriodSeconds > 0 { + return server.Spec.Lifecycle.TerminationGracePeriodSeconds + } + return defaultGraceSeconds +} + +// rconAddress is the in-cluster RCON endpoint the operator probes for readiness. +func rconAddress(server *v1alpha1.MinecraftServer) string { + return fmt.Sprintf("%s.%s.svc.cluster.local:%d", server.Name, server.Namespace, rconPort(server)) +} + +// gameAddress is the in-cluster game endpoint advertised when Running. +func gameAddress(server *v1alpha1.MinecraftServer) string { + return fmt.Sprintf("%s.%s.svc.cluster.local:%d", server.Name, server.Namespace, GamePort) +} + +// preStopScript is the operator-injected graceful-shutdown sequence (spec §7): +// flush the world, then stop the server, both over RCON. It relies on rcon-cli +// being present in the Felis base image and reading the RCON_* env injected +// alongside it. +func preStopScript(server *v1alpha1.MinecraftServer) string { + port := rconPort(server) + return fmt.Sprintf( + `rcon-cli --port %d --password "$RCON_PASSWORD" save-all flush; rcon-cli --port %d --password "$RCON_PASSWORD" stop`, + port, port, + ) +} + +// buildHeadlessService backs the StatefulSet's stable network identity. +func buildHeadlessService(server *v1alpha1.MinecraftServer) *corev1.Service { + svc := &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: headlessServiceName(server.Name), + Namespace: server.Namespace, + Labels: labelsFor(server), + }, + Spec: corev1.ServiceSpec{ + ClusterIP: corev1.ClusterIPNone, + Selector: selectorFor(server), + Ports: servicePorts(server), + }, + } + return svc +} + +// buildClientService is the stable ClusterIP the proxy and operator dial. +func buildClientService(server *v1alpha1.MinecraftServer) *corev1.Service { + return &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: server.Name, + Namespace: server.Namespace, + Labels: labelsFor(server), + }, + Spec: corev1.ServiceSpec{ + Selector: selectorFor(server), + Ports: servicePorts(server), + }, + } +} + +func servicePorts(server *v1alpha1.MinecraftServer) []corev1.ServicePort { + ports := []corev1.ServicePort{{ + Name: "game", + Port: GamePort, + TargetPort: intstr.FromInt32(GamePort), + Protocol: corev1.ProtocolTCP, + }} + if server.Spec.Rcon.Enabled { + p := rconPort(server) + ports = append(ports, corev1.ServicePort{ + Name: "rcon", + Port: p, + TargetPort: intstr.FromInt32(p), + Protocol: corev1.ProtocolTCP, + }) + } + return ports +} + +// buildStatefulSet renders the workload for replicas in {0,1}. It is where +// graceful shutdown is injected: the pod gets terminationGracePeriodSeconds and +// (when enabled) a preStop RCON save+stop hook. +func buildStatefulSet(server *v1alpha1.MinecraftServer, replicas int32) (*appsv1.StatefulSet, error) { + storageSize := server.Spec.Storage.Size + if storageSize == "" { + storageSize = defaultStorageSize + } + storageQty, err := resource.ParseQuantity(storageSize) + if err != nil { + return nil, fmt.Errorf("invalid storage size %q: %w", storageSize, err) + } + + container := corev1.Container{ + Name: containerName, + Image: server.Spec.Image, + Resources: server.Spec.Resources, + Env: buildEnv(server), + Ports: []corev1.ContainerPort{ + {Name: "game", ContainerPort: GamePort, Protocol: corev1.ProtocolTCP}, + }, + VolumeMounts: []corev1.VolumeMount{ + {Name: dataVolumeName, MountPath: dataMountPath}, + }, + // Readiness here is a plain TCP check (spec §5: readinessProbe is only + // tcpSocket; the RCON gate is enforced by the operator, not the kubelet). + ReadinessProbe: &corev1.Probe{ + ProbeHandler: corev1.ProbeHandler{ + TCPSocket: &corev1.TCPSocketAction{Port: intstr.FromInt32(GamePort)}, + }, + InitialDelaySeconds: 20, + PeriodSeconds: 10, + FailureThreshold: 6, + }, + } + if len(server.Spec.Args) > 0 { + container.Args = append([]string(nil), server.Spec.Args...) + } + if server.Spec.Rcon.Enabled { + container.Ports = append(container.Ports, corev1.ContainerPort{ + Name: "rcon", ContainerPort: rconPort(server), Protocol: corev1.ProtocolTCP, + }) + if server.Spec.Lifecycle.PreStopSaveAndStop { + container.Lifecycle = &corev1.Lifecycle{ + PreStop: &corev1.LifecycleHandler{ + Exec: &corev1.ExecAction{ + Command: []string{"/bin/sh", "-c", preStopScript(server)}, + }, + }, + } + } + } + + grace := graceSeconds(server) + pvc := corev1.PersistentVolumeClaim{ + ObjectMeta: metav1.ObjectMeta{Name: dataVolumeName}, + Spec: corev1.PersistentVolumeClaimSpec{ + AccessModes: []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce}, + Resources: corev1.VolumeResourceRequirements{ + Requests: corev1.ResourceList{corev1.ResourceStorage: storageQty}, + }, + }, + } + if sc := server.Spec.Storage.StorageClassName; sc != "" { + pvc.Spec.StorageClassName = &sc + } + + sts := &appsv1.StatefulSet{ + ObjectMeta: metav1.ObjectMeta{ + Name: server.Name, + Namespace: server.Namespace, + Labels: labelsFor(server), + }, + Spec: appsv1.StatefulSetSpec{ + Replicas: &replicas, + ServiceName: headlessServiceName(server.Name), + Selector: &metav1.LabelSelector{MatchLabels: selectorFor(server)}, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{Labels: labelsFor(server)}, + Spec: corev1.PodSpec{ + TerminationGracePeriodSeconds: &grace, + Containers: []corev1.Container{container}, + // A Minecraft server runs untrusted user worlds and plugins and + // has no business calling the K8s API, so its pod must NOT carry the + // default ServiceAccount token: a compromised plugin could otherwise + // authenticate as the namespace default SA (spec §21: user servers + // default to no SA-token mount). The pod keeps the default SA but + // with automounting explicitly disabled. + AutomountServiceAccountToken: boolPtr(false), + }, + }, + VolumeClaimTemplates: []corev1.PersistentVolumeClaim{pvc}, + }, + } + return sts, nil +} + +// buildEnv assembles the container environment: heap sizing, user-supplied +// vars, and the RCON_* pair (password sourced from the referenced Secret, never +// inlined into the CRD). +func buildEnv(server *v1alpha1.MinecraftServer) []corev1.EnvVar { + var env []corev1.EnvVar + if mem := server.Spec.JavaMemory; mem != "" { + env = append(env, corev1.EnvVar{Name: "JAVA_MEMORY", Value: mem}) + } + if len(server.Spec.JavaFlags) > 0 { + env = append(env, corev1.EnvVar{Name: "JAVA_FLAGS", Value: joinFlags(server.Spec.JavaFlags)}) + } + for _, e := range server.Spec.Env { + env = append(env, corev1.EnvVar{Name: e.Name, Value: e.Value}) + } + if server.Spec.Rcon.Enabled && server.Spec.Rcon.SecretRef.Name != "" { + env = append(env, + corev1.EnvVar{Name: "RCON_PORT", Value: strconv.Itoa(int(rconPort(server)))}, + corev1.EnvVar{Name: "RCON_PASSWORD", ValueFrom: &corev1.EnvVarSource{ + SecretKeyRef: &corev1.SecretKeySelector{ + LocalObjectReference: corev1.LocalObjectReference{Name: server.Spec.Rcon.SecretRef.Name}, + Key: server.Spec.Rcon.SecretRef.Key, + }, + }}, + ) + } + return env +} + +func boolPtr(b bool) *bool { return &b } + +func joinFlags(flags []string) string { + out := "" + for i, f := range flags { + if i > 0 { + out += " " + } + out += f + } + return out +} diff --git a/internal/operator/prober.go b/internal/operator/prober.go new file mode 100644 index 0000000..9240dcd --- /dev/null +++ b/internal/operator/prober.go @@ -0,0 +1,41 @@ +package operator + +import ( + "context" + "time" + + "felis.lolicon.best/internal/rcon" +) + +// Prober reports whether a server's RCON endpoint is reachable and accepts the +// password. A nil error is the loader-agnostic readiness gate (spec §5). It is +// an interface so the reconciler can be tested without a live server. +type Prober interface { + Probe(ctx context.Context, addr, password string) error +} + +// RconProber is the production Prober: a successful Dial (TCP connect + auth) +// is sufficient; the connection is closed immediately. +type RconProber struct { + // Timeout bounds a single probe. Defaults to 5s. + Timeout time.Duration +} + +// Probe dials addr and authenticates with password, honoring the smaller of the +// configured timeout and any deadline already on ctx. +func (p RconProber) Probe(ctx context.Context, addr, password string) error { + timeout := p.Timeout + if timeout <= 0 { + timeout = 5 * time.Second + } + if dl, ok := ctx.Deadline(); ok { + if remaining := time.Until(dl); remaining > 0 && remaining < timeout { + timeout = remaining + } + } + conn, err := rcon.Dial(addr, password, timeout) + if err != nil { + return err + } + return conn.Close() +} diff --git a/internal/operator/reconciler.go b/internal/operator/reconciler.go new file mode 100644 index 0000000..a647fc8 --- /dev/null +++ b/internal/operator/reconciler.go @@ -0,0 +1,296 @@ +// Package operator reconciles MinecraftServer objects (spec §4, §5, §7). The +// CRD is the lifecycle source-of-truth; this controller renders the +// StatefulSet/Service/PVC from it, gates readiness on an RCON probe, and injects +// graceful shutdown. It never reads or writes business-layer (Postgres) fields. +package operator + +import ( + "context" + "fmt" + "time" + + "felis.lolicon.best/internal/apis/felis/v1alpha1" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" +) + +// Requeue cadences for the transient phases. +const ( + requeueStarting = 5 * time.Second + requeueStopping = 5 * time.Second + requeueSecret = 10 * time.Second +) + +// Reconciler reconciles a MinecraftServer with its managed children. +type Reconciler struct { + client.Client + Scheme *runtime.Scheme + // Prober gates readiness on RCON reachability. + Prober Prober + // Now is injectable for deterministic timestamps in tests; defaults to + // metav1.Now. + Now func() metav1.Time +} + +func (r *Reconciler) now() metav1.Time { + if r.Now != nil { + return r.Now() + } + return metav1.Now() +} + +// SetupWithManager wires the controller to watch MinecraftServers and the +// children it owns. +func (r *Reconciler) SetupWithManager(mgr ctrl.Manager) error { + return ctrl.NewControllerManagedBy(mgr). + For(&v1alpha1.MinecraftServer{}). + Owns(&appsv1.StatefulSet{}). + Owns(&corev1.Service{}). + Complete(r) +} + +// Reconcile drives a single MinecraftServer toward spec.desiredState. +func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { + var server v1alpha1.MinecraftServer + if err := r.Get(ctx, req.NamespacedName, &server); err != nil { + // Deletion is handled by owner references on the children. + return ctrl.Result{}, client.IgnoreNotFound(err) + } + + desired := server.Spec.DesiredState + if desired == "" { + desired = v1alpha1.DesiredStopped + } + if desired == v1alpha1.DesiredStopped { + return r.reconcileStopped(ctx, &server) + } + return r.reconcileRunning(ctx, &server) +} + +func (r *Reconciler) reconcileRunning(ctx context.Context, server *v1alpha1.MinecraftServer) (ctrl.Result, error) { + if err := r.ensureServices(ctx, server); err != nil { + return ctrl.Result{}, err + } + + desired, err := buildStatefulSet(server, 1) + if err != nil { + // A malformed spec (e.g. bad storage quantity) is terminal until edited. + r.markFailed(server, "InvalidSpec", err.Error()) + return ctrl.Result{}, r.patchStatus(ctx, server) + } + if err := controllerutil.SetControllerReference(server, desired, r.Scheme); err != nil { + return ctrl.Result{}, err + } + if err := r.applyStatefulSet(ctx, desired); err != nil { + return ctrl.Result{}, err + } + + var current appsv1.StatefulSet + if err := r.Get(ctx, client.ObjectKeyFromObject(desired), ¤t); err != nil { + return ctrl.Result{}, err + } + + // The pod must first pass its tcpSocket readiness (readyReplicas >= 1). + if current.Status.ReadyReplicas < 1 { + r.markStarting(server, "PodNotReady", "waiting for pod TCP readiness") + if err := r.patchStatus(ctx, server); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: requeueStarting}, nil + } + + // Then the operator gates true readiness on an RCON probe (spec §5). + if server.Spec.Rcon.Enabled { + password, err := r.rconPassword(ctx, server) + if err != nil { + r.markStarting(server, "RconSecretUnavailable", err.Error()) + if perr := r.patchStatus(ctx, server); perr != nil { + return ctrl.Result{}, perr + } + return ctrl.Result{RequeueAfter: requeueSecret}, nil + } + if err := r.Prober.Probe(ctx, rconAddress(server), password); err != nil { + r.markStarting(server, "RconNotReachable", err.Error()) + if perr := r.patchStatus(ctx, server); perr != nil { + return ctrl.Result{}, perr + } + return ctrl.Result{RequeueAfter: requeueStarting}, nil + } + } + + r.markRunningReady(server) + return ctrl.Result{}, r.patchStatus(ctx, server) +} + +func (r *Reconciler) reconcileStopped(ctx context.Context, server *v1alpha1.MinecraftServer) (ctrl.Result, error) { + var sts appsv1.StatefulSet + err := r.Get(ctx, types.NamespacedName{Namespace: server.Namespace, Name: server.Name}, &sts) + if apierrors.IsNotFound(err) { + r.markStopped(server) + return ctrl.Result{}, r.patchStatus(ctx, server) + } + if err != nil { + return ctrl.Result{}, err + } + + // Scaling to zero triggers each pod's preStop RCON save+stop (spec §7). + if sts.Spec.Replicas == nil || *sts.Spec.Replicas != 0 { + zero := int32(0) + sts.Spec.Replicas = &zero + if err := r.Update(ctx, &sts); err != nil { + return ctrl.Result{}, err + } + } + + if sts.Status.Replicas > 0 || sts.Status.ReadyReplicas > 0 { + r.markStopping(server) + if err := r.patchStatus(ctx, server); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: requeueStopping}, nil + } + + r.markStopped(server) + return ctrl.Result{}, r.patchStatus(ctx, server) +} + +func (r *Reconciler) ensureServices(ctx context.Context, server *v1alpha1.MinecraftServer) error { + for _, svc := range []*corev1.Service{buildHeadlessService(server), buildClientService(server)} { + if err := controllerutil.SetControllerReference(server, svc, r.Scheme); err != nil { + return err + } + if err := r.applyService(ctx, svc); err != nil { + return err + } + } + return nil +} + +func (r *Reconciler) rconPassword(ctx context.Context, server *v1alpha1.MinecraftServer) (string, error) { + ref := server.Spec.Rcon.SecretRef + if ref.Name == "" || ref.Key == "" { + return "", fmt.Errorf("rcon.secretRef.name and .key are required when rcon is enabled") + } + var secret corev1.Secret + if err := r.Get(ctx, types.NamespacedName{Namespace: server.Namespace, Name: ref.Name}, &secret); err != nil { + return "", err + } + b, ok := secret.Data[ref.Key] + if !ok { + return "", fmt.Errorf("secret %q has no key %q", ref.Name, ref.Key) + } + return string(b), nil +} + +// applyStatefulSet creates the StatefulSet or, if it exists, updates only its +// mutable fields (StatefulSet selector/serviceName/volumeClaimTemplates are +// immutable and must not be re-sent). +func (r *Reconciler) applyStatefulSet(ctx context.Context, desired *appsv1.StatefulSet) error { + var existing appsv1.StatefulSet + err := r.Get(ctx, client.ObjectKeyFromObject(desired), &existing) + if apierrors.IsNotFound(err) { + return r.Create(ctx, desired) + } + if err != nil { + return err + } + existing.Labels = desired.Labels + existing.Spec.Replicas = desired.Spec.Replicas + existing.Spec.Template = desired.Spec.Template + return r.Update(ctx, &existing) +} + +// applyService creates or updates a Service, preserving the cluster-assigned +// ClusterIP so the update does not orphan the address. +func (r *Reconciler) applyService(ctx context.Context, desired *corev1.Service) error { + var existing corev1.Service + err := r.Get(ctx, client.ObjectKeyFromObject(desired), &existing) + if apierrors.IsNotFound(err) { + return r.Create(ctx, desired) + } + if err != nil { + return err + } + desired.ResourceVersion = existing.ResourceVersion + desired.Spec.ClusterIP = existing.Spec.ClusterIP + desired.Spec.ClusterIPs = existing.Spec.ClusterIPs + return r.Update(ctx, desired) +} + +func (r *Reconciler) patchStatus(ctx context.Context, server *v1alpha1.MinecraftServer) error { + return r.Status().Update(ctx, server) +} + +// --- status mutators ------------------------------------------------------- + +func (r *Reconciler) markStarting(server *v1alpha1.MinecraftServer, reason, msg string) { + server.Status.Phase = v1alpha1.PhaseStarting + server.Status.Ready = false + server.Status.ObservedGeneration = server.Generation + server.Status.Endpoint = v1alpha1.EndpointStatus{Mode: v1alpha1.EndpointFallback, Address: server.Spec.FallbackServer} + server.Status.LiveMotd = server.Spec.Motd.Starting + r.setCondition(server, v1alpha1.ConditionReady, metav1.ConditionFalse, reason, msg) + r.setCondition(server, v1alpha1.ConditionRconReached, metav1.ConditionFalse, reason, msg) +} + +func (r *Reconciler) markRunningReady(server *v1alpha1.MinecraftServer) { + server.Status.Phase = v1alpha1.PhaseRunning + server.Status.Ready = true + server.Status.ObservedGeneration = server.Generation + if server.Status.ReadySignalAt == nil { + t := r.now() + server.Status.ReadySignalAt = &t + } + server.Status.Endpoint = v1alpha1.EndpointStatus{Mode: v1alpha1.EndpointDirect, Address: gameAddress(server)} + server.Status.LiveMotd = server.Spec.Motd.Running + r.setCondition(server, v1alpha1.ConditionRconReached, metav1.ConditionTrue, "Probed", "RCON probe succeeded") + r.setCondition(server, v1alpha1.ConditionReady, metav1.ConditionTrue, "RconReached", "server is accepting RCON") +} + +func (r *Reconciler) markStopping(server *v1alpha1.MinecraftServer) { + server.Status.Phase = v1alpha1.PhaseStopping + server.Status.Ready = false + server.Status.ObservedGeneration = server.Generation + server.Status.Endpoint = v1alpha1.EndpointStatus{Mode: v1alpha1.EndpointFallback, Address: server.Spec.FallbackServer} + server.Status.LiveMotd = server.Spec.Motd.Stopped + r.setCondition(server, v1alpha1.ConditionReady, metav1.ConditionFalse, "Stopping", "scaling down") +} + +func (r *Reconciler) markStopped(server *v1alpha1.MinecraftServer) { + server.Status.Phase = v1alpha1.PhaseStopped + server.Status.Ready = false + server.Status.ObservedGeneration = server.Generation + server.Status.ReadySignalAt = nil + server.Status.Players = v1alpha1.PlayersStatus{} + server.Status.Endpoint = v1alpha1.EndpointStatus{Mode: v1alpha1.EndpointFallback, Address: server.Spec.FallbackServer} + server.Status.LiveMotd = server.Spec.Motd.Stopped + r.setCondition(server, v1alpha1.ConditionReady, metav1.ConditionFalse, "Stopped", "desiredState is Stopped") + r.setCondition(server, v1alpha1.ConditionRconReached, metav1.ConditionFalse, "Stopped", "server is stopped") +} + +func (r *Reconciler) markFailed(server *v1alpha1.MinecraftServer, reason, msg string) { + server.Status.Phase = v1alpha1.PhaseFailed + server.Status.Ready = false + server.Status.ObservedGeneration = server.Generation + r.setCondition(server, v1alpha1.ConditionReady, metav1.ConditionFalse, reason, msg) + r.setCondition(server, v1alpha1.ConditionProvisioned, metav1.ConditionFalse, reason, msg) +} + +func (r *Reconciler) setCondition(server *v1alpha1.MinecraftServer, condType string, status metav1.ConditionStatus, reason, msg string) { + meta.SetStatusCondition(&server.Status.Conditions, metav1.Condition{ + Type: condType, + Status: status, + Reason: reason, + Message: msg, + ObservedGeneration: server.Generation, + LastTransitionTime: r.now(), + }) +} diff --git a/internal/operator/reconciler_test.go b/internal/operator/reconciler_test.go new file mode 100644 index 0000000..d6285e9 --- /dev/null +++ b/internal/operator/reconciler_test.go @@ -0,0 +1,298 @@ +package operator_test + +import ( + "context" + "errors" + "strings" + "testing" + "time" + + "felis.lolicon.best/internal/apis/felis/v1alpha1" + "felis.lolicon.best/internal/operator" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + clientgoscheme "k8s.io/client-go/kubernetes/scheme" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +type fakeProber struct{ err error } + +func (f fakeProber) Probe(context.Context, string, string) error { return f.err } + +func newScheme(t *testing.T) *runtime.Scheme { + t.Helper() + scheme := runtime.NewScheme() + if err := clientgoscheme.AddToScheme(scheme); err != nil { + t.Fatalf("clientgo scheme: %v", err) + } + if err := v1alpha1.AddToScheme(scheme); err != nil { + t.Fatalf("v1alpha1 scheme: %v", err) + } + return scheme +} + +// fixedNow yields a deterministic timestamp so ReadySignalAt is assertable. +func fixedNow() metav1.Time { + return metav1.Date(2026, 6, 25, 12, 0, 0, 0, time.UTC) +} + +func runningServer() *v1alpha1.MinecraftServer { + return &v1alpha1.MinecraftServer{ + ObjectMeta: metav1.ObjectMeta{Name: "survival", Namespace: "minecraft", Generation: 1}, + Spec: v1alpha1.MinecraftServerSpec{ + Subdomain: "survival", + DesiredState: v1alpha1.DesiredRunning, + Image: "registry.internal/felis/paper:latest", + JavaMemory: "4G", + Storage: v1alpha1.StorageSpec{Size: "10Gi"}, + FallbackServer: "lobby", + Motd: v1alpha1.MotdSpec{Running: "up", Stopped: "down", Starting: "booting"}, + Lifecycle: v1alpha1.LifecycleSpec{PreStopSaveAndStop: true}, + Rcon: v1alpha1.RconSpec{ + Enabled: true, + Port: 25575, + SecretRef: v1alpha1.SecretKeyRef{Name: "survival-rcon", Key: "password"}, + }, + }, + } +} + +func rconSecret() *corev1.Secret { + return &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: "survival-rcon", Namespace: "minecraft"}, + Data: map[string][]byte{"password": []byte("hunter2")}, + } +} + +func newReconciler(t *testing.T, prober operator.Prober, objs ...client.Object) (*operator.Reconciler, client.Client) { + t.Helper() + scheme := newScheme(t) + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(objs...). + WithStatusSubresource(&v1alpha1.MinecraftServer{}, &appsv1.StatefulSet{}). + Build() + return &operator.Reconciler{Client: c, Scheme: scheme, Prober: prober, Now: fixedNow}, c +} + +func reconcile(t *testing.T, r *operator.Reconciler, name string) ctrl.Result { + t.Helper() + res, err := r.Reconcile(context.Background(), ctrl.Request{ + NamespacedName: types.NamespacedName{Namespace: "minecraft", Name: name}, + }) + if err != nil { + t.Fatalf("Reconcile(%s): %v", name, err) + } + return res +} + +func getServer(t *testing.T, c client.Client, name string) *v1alpha1.MinecraftServer { + t.Helper() + var s v1alpha1.MinecraftServer + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: name}, &s); err != nil { + t.Fatalf("get server: %v", err) + } + return &s +} + +func getSTS(t *testing.T, c client.Client, name string) *appsv1.StatefulSet { + t.Helper() + var sts appsv1.StatefulSet + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: name}, &sts); err != nil { + t.Fatalf("get statefulset: %v", err) + } + return &sts +} + +// markPodReady simulates the kubelet flipping the StatefulSet to ready. +func markPodReady(t *testing.T, c client.Client, name string) { + t.Helper() + sts := getSTS(t, c, name) + sts.Status.Replicas = 1 + sts.Status.ReadyReplicas = 1 + if err := c.Status().Update(context.Background(), sts); err != nil { + t.Fatalf("update sts status: %v", err) + } +} + +func TestReconcileRunning_CreatesWorkloadAndInjectsGracefulShutdown(t *testing.T) { + r, c := newReconciler(t, fakeProber{}, runningServer(), rconSecret()) + + res := reconcile(t, r, "survival") + if res.RequeueAfter == 0 { + t.Errorf("expected a requeue while Starting, got %+v", res) + } + + // Services exist. + var svc corev1.Service + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "survival"}, &svc); err != nil { + t.Fatalf("client service not created: %v", err) + } + var hl corev1.Service + if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "survival-hl"}, &hl); err != nil { + t.Fatalf("headless service not created: %v", err) + } + if hl.Spec.ClusterIP != corev1.ClusterIPNone { + t.Errorf("headless service ClusterIP = %q, want None", hl.Spec.ClusterIP) + } + + // StatefulSet exists with graceful-shutdown injection. + sts := getSTS(t, c, "survival") + if got := sts.Spec.Template.Spec.TerminationGracePeriodSeconds; got == nil || *got != 300 { + t.Errorf("terminationGracePeriodSeconds = %v, want 300", got) + } + container := sts.Spec.Template.Spec.Containers[0] + if container.Lifecycle == nil || container.Lifecycle.PreStop == nil || container.Lifecycle.PreStop.Exec == nil { + t.Fatal("expected preStop exec hook to be injected") + } + preStop := strings.Join(container.Lifecycle.PreStop.Exec.Command, " ") + if !strings.Contains(preStop, "save-all flush") || !strings.Contains(preStop, "stop") { + t.Errorf("preStop hook missing save/stop sequence: %q", preStop) + } + // RCON password is sourced from the Secret, never inlined. + var sawRconPassword bool + for _, e := range container.Env { + if e.Name == "RCON_PASSWORD" { + sawRconPassword = true + if e.ValueFrom == nil || e.ValueFrom.SecretKeyRef == nil { + t.Error("RCON_PASSWORD must come from a SecretKeyRef") + } + if e.Value != "" { + t.Error("RCON_PASSWORD must not be inlined as a literal value") + } + } + } + if !sawRconPassword { + t.Error("RCON_PASSWORD env not injected") + } + + // Owner reference points back at the MinecraftServer. + if len(sts.OwnerReferences) != 1 || sts.OwnerReferences[0].Name != "survival" { + t.Errorf("statefulset owner refs = %+v, want one ref to survival", sts.OwnerReferences) + } + + // The server pod must NOT mount a ServiceAccount token (spec §21): it runs + // untrusted user worlds/plugins and has no business reaching the K8s API. + if amt := sts.Spec.Template.Spec.AutomountServiceAccountToken; amt == nil || *amt { + t.Error("server pod must set AutomountServiceAccountToken=false — untrusted workloads must not reach the API") + } + + // Status is Starting (pod not yet ready). + server := getServer(t, c, "survival") + if server.Status.Phase != v1alpha1.PhaseStarting || server.Status.Ready { + t.Errorf("status = %s ready=%v, want Starting not-ready", server.Status.Phase, server.Status.Ready) + } + if server.Status.Endpoint.Mode != v1alpha1.EndpointFallback || server.Status.Endpoint.Address != "lobby" { + t.Errorf("starting endpoint = %+v, want fallback->lobby", server.Status.Endpoint) + } +} + +func TestReconcileRunning_RconProbeGatesReadiness(t *testing.T) { + r, c := newReconciler(t, fakeProber{}, runningServer(), rconSecret()) + + reconcile(t, r, "survival") // creates workload, Starting + markPodReady(t, c, "survival") + + res := reconcile(t, r, "survival") // pod ready + probe OK -> Running + if res.RequeueAfter != 0 { + t.Errorf("expected no requeue once Running, got %+v", res) + } + + server := getServer(t, c, "survival") + if server.Status.Phase != v1alpha1.PhaseRunning || !server.Status.Ready { + t.Fatalf("status = %s ready=%v, want Running ready", server.Status.Phase, server.Status.Ready) + } + if server.Status.ReadySignalAt == nil || !server.Status.ReadySignalAt.Equal(ptrTime(fixedNow())) { + t.Errorf("readySignalAt = %v, want %v", server.Status.ReadySignalAt, fixedNow()) + } + if server.Status.Endpoint.Mode != v1alpha1.EndpointDirect { + t.Errorf("endpoint mode = %s, want direct", server.Status.Endpoint.Mode) + } + if !isConditionTrue(server, v1alpha1.ConditionReady) { + t.Error("Ready condition should be True") + } + if !isConditionTrue(server, v1alpha1.ConditionRconReached) { + t.Error("RconReached condition should be True") + } +} + +func TestReconcileRunning_RconProbeFailureStaysStarting(t *testing.T) { + r, c := newReconciler(t, fakeProber{err: errors.New("connection refused")}, runningServer(), rconSecret()) + + reconcile(t, r, "survival") + markPodReady(t, c, "survival") + + res := reconcile(t, r, "survival") // pod ready but probe fails + if res.RequeueAfter == 0 { + t.Error("expected requeue while RCON not reachable") + } + server := getServer(t, c, "survival") + if server.Status.Phase != v1alpha1.PhaseStarting || server.Status.Ready { + t.Errorf("status = %s ready=%v, want Starting not-ready", server.Status.Phase, server.Status.Ready) + } + if isConditionTrue(server, v1alpha1.ConditionReady) { + t.Error("Ready condition must not be True when the probe fails") + } +} + +func TestReconcileStopped_NoWorkloadIsStopped(t *testing.T) { + s := runningServer() + s.Spec.DesiredState = v1alpha1.DesiredStopped + r, c := newReconciler(t, fakeProber{}, s) + + reconcile(t, r, "survival") + + if _, err := getSTSErr(c, "survival"); !apierrors.IsNotFound(err) { + t.Errorf("stopped server should not create a StatefulSet, got err=%v", err) + } + server := getServer(t, c, "survival") + if server.Status.Phase != v1alpha1.PhaseStopped || server.Status.Ready { + t.Errorf("status = %s ready=%v, want Stopped not-ready", server.Status.Phase, server.Status.Ready) + } +} + +func TestReconcileStopped_ScalesRunningWorkloadDown(t *testing.T) { + r, c := newReconciler(t, fakeProber{}, runningServer(), rconSecret()) + + reconcile(t, r, "survival") + markPodReady(t, c, "survival") + reconcile(t, r, "survival") // now Running with replicas=1 + + // Flip desiredState to Stopped. + server := getServer(t, c, "survival") + server.Spec.DesiredState = v1alpha1.DesiredStopped + if err := c.Update(context.Background(), server); err != nil { + t.Fatalf("flip desiredState: %v", err) + } + + reconcile(t, r, "survival") + sts := getSTS(t, c, "survival") + if sts.Spec.Replicas == nil || *sts.Spec.Replicas != 0 { + t.Errorf("replicas = %v, want 0 after stop", sts.Spec.Replicas) + } +} + +// --- helpers --------------------------------------------------------------- + +func getSTSErr(c client.Client, name string) (*appsv1.StatefulSet, error) { + var sts appsv1.StatefulSet + err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: name}, &sts) + return &sts, err +} + +func isConditionTrue(server *v1alpha1.MinecraftServer, condType string) bool { + for _, cond := range server.Status.Conditions { + if cond.Type == condType { + return cond.Status == metav1.ConditionTrue + } + } + return false +} + +func ptrTime(t metav1.Time) *metav1.Time { return &t }