feat: add Go Kubernetes operator

This commit is contained in:
2026-08-13 02:19:14 +07:00
parent 31a2b8ef28
commit da6eafcae6
30 changed files with 3066 additions and 0 deletions
+32
View File
@@ -0,0 +1,32 @@
package controller
import (
"context"
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/apiutil"
v1alpha1 "github.com/YuzuZensai/Minikura/operator/api/v1alpha1"
)
var applyOpts = []client.PatchOption{
client.FieldOwner(v1alpha1.ManagerName),
client.ForceOwnership,
}
func apply(ctx context.Context, c client.Client, owner client.Object, obj client.Object, scheme *runtime.Scheme) error {
if err := setOwner(owner, obj, scheme); err != nil {
return err
}
gvk, err := apiutil.GVKForObject(obj, scheme)
if err != nil {
return err
}
obj.GetObjectKind().SetGroupVersionKind(gvk)
obj.SetManagedFields(nil)
obj.SetResourceVersion("")
return c.Patch(ctx, obj, client.Apply, applyOpts...)
}
@@ -0,0 +1,91 @@
package controller
import (
"encoding/json"
"testing"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
"sigs.k8s.io/controller-runtime/pkg/client/apiutil"
v1alpha1 "github.com/YuzuZensai/Minikura/operator/api/v1alpha1"
)
func testScheme(t *testing.T) *runtime.Scheme {
t.Helper()
s := runtime.NewScheme()
if err := clientgoscheme.AddToScheme(s); err != nil {
t.Fatalf("add client-go scheme: %v", err)
}
if err := v1alpha1.AddToScheme(s); err != nil {
t.Fatalf("add v1alpha1 scheme: %v", err)
}
return s
}
func TestTypedObjectMarshalsWithoutTypeMeta(t *testing.T) {
cm := &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "a", Namespace: "b"}}
var m map[string]any
raw, err := json.Marshal(cm)
if err != nil {
t.Fatalf("marshal: %v", err)
}
if err := json.Unmarshal(raw, &m); err != nil {
t.Fatalf("unmarshal: %v", err)
}
if _, ok := m["apiVersion"]; ok {
t.Fatal("precondition changed: typed objects now carry apiVersion")
}
if _, ok := m["kind"]; ok {
t.Fatal("precondition changed: typed objects now carry kind")
}
}
func TestApplySetsTypeMeta(t *testing.T) {
scheme := testScheme(t)
owner := &v1alpha1.MinecraftServer{
ObjectMeta: metav1.ObjectMeta{Name: "smp", Namespace: "minikura", UID: "abc"},
}
cm := &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "a", Namespace: "minikura"}}
if err := setOwner(owner, cm, scheme); err != nil {
t.Fatalf("setOwner: %v", err)
}
gvk, err := apiutil.GVKForObject(cm, scheme)
if err != nil {
t.Fatalf("GVKForObject: %v", err)
}
cm.GetObjectKind().SetGroupVersionKind(gvk)
raw, err := json.Marshal(cm)
if err != nil {
t.Fatalf("marshal: %v", err)
}
var m map[string]any
if err := json.Unmarshal(raw, &m); err != nil {
t.Fatalf("unmarshal: %v", err)
}
if m["apiVersion"] != "v1" {
t.Errorf("apiVersion = %v, want v1", m["apiVersion"])
}
if m["kind"] != "ConfigMap" {
t.Errorf("kind = %v, want ConfigMap", m["kind"])
}
refs := cm.GetOwnerReferences()
if len(refs) != 1 {
t.Fatalf("ownerReferences = %d, want 1", len(refs))
}
if refs[0].Kind != "MinecraftServer" || refs[0].Name != "smp" {
t.Errorf("owner = %s/%s, want MinecraftServer/smp", refs[0].Kind, refs[0].Name)
}
if refs[0].Controller == nil || !*refs[0].Controller {
t.Error("owner reference should be a controller reference")
}
}
@@ -0,0 +1,184 @@
package controller
import (
"context"
"fmt"
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"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/log"
v1alpha1 "github.com/YuzuZensai/Minikura/operator/api/v1alpha1"
"github.com/YuzuZensai/Minikura/operator/internal/resources"
)
type MinecraftServerReconciler struct {
client.Client
Scheme *runtime.Scheme
}
// +kubebuilder:rbac:groups=minikura.kirameki.cafe,resources=minecraftservers,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=minikura.kirameki.cafe,resources=minecraftservers/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=minikura.kirameki.cafe,resources=minecraftservers/finalizers,verbs=update
// +kubebuilder:rbac:groups=apps,resources=deployments;statefulsets,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups="",resources=services;configmaps,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups="",resources=pods,verbs=get;list;watch
// +kubebuilder:rbac:groups=coordination.k8s.io,resources=leases,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups="",resources=events,verbs=create;patch
func (r *MinecraftServerReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
logger := log.FromContext(ctx)
var mc v1alpha1.MinecraftServer
if err := r.Get(ctx, req.NamespacedName, &mc); err != nil {
return ctrl.Result{}, client.IgnoreNotFound(err)
}
if !mc.DeletionTimestamp.IsZero() {
return ctrl.Result{}, nil
}
if err := apply(ctx, r.Client, &mc, resources.MinecraftConfigMap(&mc), r.Scheme); err != nil {
return r.fail(ctx, &mc, "ConfigMapFailed", err)
}
if err := apply(ctx, r.Client, &mc, resources.MinecraftService(&mc), r.Scheme); err != nil {
return r.fail(ctx, &mc, "ServiceFailed", err)
}
stateful := mc.Spec.Type == v1alpha1.ServerStateful
if err := r.reconcileWorkload(ctx, &mc, stateful); err != nil {
return r.fail(ctx, &mc, "WorkloadFailed", err)
}
if err := r.pruneOppositeWorkload(ctx, &mc, stateful); err != nil {
logger.Error(err, "failed to prune previous workload")
}
return ctrl.Result{}, r.updateStatus(ctx, &mc, stateful)
}
func (r *MinecraftServerReconciler) reconcileWorkload(ctx context.Context, mc *v1alpha1.MinecraftServer, stateful bool) error {
if !stateful {
return apply(ctx, r.Client, mc, resources.MinecraftDeployment(mc), r.Scheme)
}
sts, err := resources.MinecraftStatefulSet(mc)
if err != nil {
return err
}
return apply(ctx, r.Client, mc, sts, r.Scheme)
}
func (r *MinecraftServerReconciler) pruneOppositeWorkload(ctx context.Context, mc *v1alpha1.MinecraftServer, stateful bool) error {
name := resources.ServerName(mc.Name)
key := client.ObjectKey{Name: name, Namespace: mc.Namespace}
var stale client.Object
if stateful {
stale = &appsv1.Deployment{}
} else {
stale = &appsv1.StatefulSet{}
}
if err := r.Get(ctx, key, stale); err != nil {
return client.IgnoreNotFound(err)
}
if !metav1.IsControlledBy(stale, mc) {
return nil
}
return client.IgnoreNotFound(r.Delete(ctx, stale))
}
func (r *MinecraftServerReconciler) updateStatus(ctx context.Context, mc *v1alpha1.MinecraftServer, stateful bool) error {
name := resources.ServerName(mc.Name)
key := client.ObjectKey{Name: name, Namespace: mc.Namespace}
var ready, replicas int32
if stateful {
var sts appsv1.StatefulSet
if err := r.Get(ctx, key, &sts); err == nil {
ready, replicas = sts.Status.ReadyReplicas, sts.Status.Replicas
} else if !apierrors.IsNotFound(err) {
return err
}
} else {
var dep appsv1.Deployment
if err := r.Get(ctx, key, &dep); err == nil {
ready, replicas = dep.Status.ReadyReplicas, dep.Status.Replicas
} else if !apierrors.IsNotFound(err) {
return err
}
}
phase := v1alpha1.PhasePending
if ready > 0 {
phase = v1alpha1.PhaseRunning
}
endpoint, err := r.endpoint(ctx, mc)
if err != nil {
return err
}
patch := client.MergeFrom(mc.DeepCopy())
mc.Status.Phase = phase
mc.Status.ReadyReplicas = ready
mc.Status.Replicas = replicas
mc.Status.Endpoint = endpoint
mc.Status.ObservedGeneration = mc.Generation
mc.Status.Message = ""
setCondition(&mc.Status.Conditions, metav1.Condition{
Type: v1alpha1.ConditionReady,
Status: conditionStatus(ready > 0),
Reason: phase,
ObservedGeneration: mc.Generation,
})
return r.Status().Patch(ctx, mc, patch)
}
func (r *MinecraftServerReconciler) endpoint(ctx context.Context, mc *v1alpha1.MinecraftServer) (string, error) {
var svc corev1.Service
key := client.ObjectKey{Name: resources.ServerName(mc.Name), Namespace: mc.Namespace}
if err := r.Get(ctx, key, &svc); err != nil {
if apierrors.IsNotFound(err) {
return "", nil
}
return "", err
}
return serviceEndpoint(&svc, mc.Namespace), nil
}
func (r *MinecraftServerReconciler) fail(ctx context.Context, mc *v1alpha1.MinecraftServer, reason string, cause error) (ctrl.Result, error) {
patch := client.MergeFrom(mc.DeepCopy())
mc.Status.Phase = v1alpha1.PhaseFailed
mc.Status.Message = cause.Error()
setCondition(&mc.Status.Conditions, metav1.Condition{
Type: v1alpha1.ConditionReady,
Status: metav1.ConditionFalse,
Reason: reason,
Message: cause.Error(),
ObservedGeneration: mc.Generation,
})
if err := r.Status().Patch(ctx, mc, patch); err != nil {
return ctrl.Result{}, fmt.Errorf("%w (status patch failed: %v)", cause, err)
}
return ctrl.Result{}, cause
}
func (r *MinecraftServerReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
For(&v1alpha1.MinecraftServer{}).
Owns(&appsv1.Deployment{}).
Owns(&appsv1.StatefulSet{}).
Owns(&corev1.Service{}).
Owns(&corev1.ConfigMap{}).
Complete(r)
}
+11
View File
@@ -0,0 +1,11 @@
package controller
import (
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
)
func setOwner(owner, obj client.Object, scheme *runtime.Scheme) error {
return controllerutil.SetControllerReference(owner, obj, scheme)
}
@@ -0,0 +1,172 @@
package controller
import (
"context"
"fmt"
"sort"
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/labels"
"k8s.io/apimachinery/pkg/runtime"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/handler"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
v1alpha1 "github.com/YuzuZensai/Minikura/operator/api/v1alpha1"
"github.com/YuzuZensai/Minikura/operator/internal/resources"
)
type ReverseProxyServerReconciler struct {
client.Client
Scheme *runtime.Scheme
}
// +kubebuilder:rbac:groups=minikura.kirameki.cafe,resources=reverseproxyservers,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=minikura.kirameki.cafe,resources=reverseproxyservers/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=minikura.kirameki.cafe,resources=reverseproxyservers/finalizers,verbs=update
func (r *ReverseProxyServerReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
var rp v1alpha1.ReverseProxyServer
if err := r.Get(ctx, req.NamespacedName, &rp); err != nil {
return ctrl.Result{}, client.IgnoreNotFound(err)
}
if !rp.DeletionTimestamp.IsZero() {
return ctrl.Result{}, nil
}
if err := apply(ctx, r.Client, &rp, resources.ProxyConfigMap(&rp), r.Scheme); err != nil {
return r.fail(ctx, &rp, "ConfigMapFailed", err)
}
if err := apply(ctx, r.Client, &rp, resources.ProxyService(&rp), r.Scheme); err != nil {
return r.fail(ctx, &rp, "ServiceFailed", err)
}
if err := apply(ctx, r.Client, &rp, resources.ProxyDeployment(&rp), r.Scheme); err != nil {
return r.fail(ctx, &rp, "DeploymentFailed", err)
}
return ctrl.Result{}, r.updateStatus(ctx, &rp)
}
func (r *ReverseProxyServerReconciler) backends(ctx context.Context, rp *v1alpha1.ReverseProxyServer) ([]string, error) {
selector := labels.Everything()
if rp.Spec.BackendSelector != nil {
s, err := metav1.LabelSelectorAsSelector(rp.Spec.BackendSelector)
if err != nil {
return nil, fmt.Errorf("invalid backendSelector: %w", err)
}
selector = s
}
var list v1alpha1.MinecraftServerList
if err := r.List(ctx, &list,
client.InNamespace(rp.Namespace),
client.MatchingLabelsSelector{Selector: selector},
); err != nil {
return nil, err
}
names := make([]string, 0, len(list.Items))
for _, mc := range list.Items {
names = append(names, mc.Name)
}
sort.Strings(names)
return names, nil
}
func (r *ReverseProxyServerReconciler) updateStatus(ctx context.Context, rp *v1alpha1.ReverseProxyServer) error {
name := resources.ProxyName(rp.Spec.Type, rp.Name)
key := client.ObjectKey{Name: name, Namespace: rp.Namespace}
var ready, replicas int32
var dep appsv1.Deployment
if err := r.Get(ctx, key, &dep); err == nil {
ready, replicas = dep.Status.ReadyReplicas, dep.Status.Replicas
} else if !apierrors.IsNotFound(err) {
return err
}
backends, err := r.backends(ctx, rp)
if err != nil {
return err
}
endpoint := ""
var svc corev1.Service
if err := r.Get(ctx, key, &svc); err == nil {
endpoint = serviceEndpoint(&svc, rp.Namespace)
} else if !apierrors.IsNotFound(err) {
return err
}
phase := v1alpha1.PhasePending
if ready > 0 {
phase = v1alpha1.PhaseRunning
}
patch := client.MergeFrom(rp.DeepCopy())
rp.Status.Phase = phase
rp.Status.ReadyReplicas = ready
rp.Status.Replicas = replicas
rp.Status.Endpoint = endpoint
rp.Status.Backends = backends
rp.Status.ObservedGeneration = rp.Generation
rp.Status.Message = ""
setCondition(&rp.Status.Conditions, metav1.Condition{
Type: v1alpha1.ConditionReady,
Status: conditionStatus(ready > 0),
Reason: phase,
ObservedGeneration: rp.Generation,
})
return r.Status().Patch(ctx, rp, patch)
}
func (r *ReverseProxyServerReconciler) fail(ctx context.Context, rp *v1alpha1.ReverseProxyServer, reason string, cause error) (ctrl.Result, error) {
patch := client.MergeFrom(rp.DeepCopy())
rp.Status.Phase = v1alpha1.PhaseFailed
rp.Status.Message = cause.Error()
setCondition(&rp.Status.Conditions, metav1.Condition{
Type: v1alpha1.ConditionReady,
Status: metav1.ConditionFalse,
Reason: reason,
Message: cause.Error(),
ObservedGeneration: rp.Generation,
})
if err := r.Status().Patch(ctx, rp, patch); err != nil {
return ctrl.Result{}, fmt.Errorf("%w (status patch failed: %v)", cause, err)
}
return ctrl.Result{}, cause
}
func (r *ReverseProxyServerReconciler) proxiesForServer(ctx context.Context, obj client.Object) []reconcile.Request {
var list v1alpha1.ReverseProxyServerList
if err := r.List(ctx, &list, client.InNamespace(obj.GetNamespace())); err != nil {
return nil
}
reqs := make([]reconcile.Request, 0, len(list.Items))
for _, rp := range list.Items {
reqs = append(reqs, reconcile.Request{
NamespacedName: client.ObjectKey{Name: rp.Name, Namespace: rp.Namespace},
})
}
return reqs
}
func (r *ReverseProxyServerReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
For(&v1alpha1.ReverseProxyServer{}).
Owns(&appsv1.Deployment{}).
Owns(&corev1.Service{}).
Owns(&corev1.ConfigMap{}).
Watches(&v1alpha1.MinecraftServer{}, handler.EnqueueRequestsFromMapFunc(r.proxiesForServer)).
Complete(r)
}
+53
View File
@@ -0,0 +1,53 @@
package controller
import (
"fmt"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
func conditionStatus(ok bool) metav1.ConditionStatus {
if ok {
return metav1.ConditionTrue
}
return metav1.ConditionFalse
}
func setCondition(conditions *[]metav1.Condition, c metav1.Condition) {
if c.Reason == "" {
c.Reason = "Unknown"
}
meta.SetStatusCondition(conditions, c)
}
func serviceEndpoint(svc *corev1.Service, namespace string) string {
internal := fmt.Sprintf("%s.%s.svc.cluster.local", svc.Name, namespace)
port := int32(0)
if len(svc.Spec.Ports) > 0 {
port = svc.Spec.Ports[0].Port
}
switch svc.Spec.Type {
case corev1.ServiceTypeLoadBalancer:
for _, ing := range svc.Status.LoadBalancer.Ingress {
if host := ing.Hostname; host != "" {
return fmt.Sprintf("%s:%d", host, port)
}
if ip := ing.IP; ip != "" {
return fmt.Sprintf("%s:%d", ip, port)
}
}
return ""
case corev1.ServiceTypeNodePort:
if len(svc.Spec.Ports) > 0 && svc.Spec.Ports[0].NodePort != 0 {
return fmt.Sprintf("<node-ip>:%d", svc.Spec.Ports[0].NodePort)
}
return internal
default:
return fmt.Sprintf("%s:%d", internal, port)
}
}