feat(operator): add MinecraftServer controller and reconcilers

The Kubernetes controller that drives MinecraftServer resources through their lifecycle and issues RCON where readiness requires it.
This commit is contained in:
flyemoji committed 2026-06-26 23:31:58 +09:00
1 parent 43ab92151f
commit 78b8cf6ded
4 files changed
+922

No files matched your search

+287
View File
@@ -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
}
+41
View File
@@ -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()
}
+296
View File
@@ -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), &current); 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(),
})
}
+298
View File
@@ -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 }