package controller import ( "context" "fmt" "strconv" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" 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/handler" "sigs.k8s.io/controller-runtime/pkg/log" "sigs.k8s.io/controller-runtime/pkg/reconcile" v1alpha1 "git.unkin.net/unkin/kea-operator/api/v1alpha1" "git.unkin.net/unkin/kea-operator/internal/kea" ) // KeaClusterReconciler reconciles a KeaCluster. type KeaClusterReconciler struct { client.Client Scheme *runtime.Scheme Control *kea.ControlClient } // +kubebuilder:rbac:groups=kea.unkin.net,resources=keaclusters,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups=kea.unkin.net,resources=keaclusters/status,verbs=get;update;patch // +kubebuilder:rbac:groups=kea.unkin.net,resources=keasubnets,verbs=get;list;watch // +kubebuilder:rbac:groups=kea.unkin.net,resources=keaclientclasses,verbs=get;list;watch // +kubebuilder:rbac:groups=apps,resources=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 func (r *KeaClusterReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { l := log.FromContext(ctx) var cluster v1alpha1.KeaCluster if err := r.Get(ctx, req.NamespacedName, &cluster); err != nil { return ctrl.Result{}, client.IgnoreNotFound(err) } if err := r.reconcileConfigMap(ctx, &cluster); err != nil { return r.fail(ctx, &cluster, "ConfigError", err) } if err := r.reconcileServices(ctx, &cluster); err != nil { return r.fail(ctx, &cluster, "ServiceError", err) } sts, err := r.reconcileStatefulSet(ctx, &cluster) if err != nil { return r.fail(ctx, &cluster, "WorkloadError", err) } r.reloadReadyPods(ctx, &cluster) ready := sts.Status.ReadyReplicas desired := int32(1) if cluster.Spec.Replicas != nil { desired = *cluster.Spec.Replicas } cluster.Status.ObservedGeneration = cluster.Generation cluster.Status.Replicas = sts.Status.Replicas cluster.Status.ReadyReplicas = ready cluster.Status.ServiceIP = r.serviceIP(ctx, &cluster) if ready > 0 { cluster.Status.ActivePeer = cluster.Name + "-0" } if ready >= desired && desired > 0 { cluster.Status.Phase = "Ready" setReady(&cluster.Status.Conditions, cluster.Generation, true, "Ready", "all replicas ready") } else { cluster.Status.Phase = "Progressing" setReady(&cluster.Status.Conditions, cluster.Generation, false, "Progressing", fmt.Sprintf("%d/%d replicas ready", ready, desired)) } if err := r.Status().Update(ctx, &cluster); err != nil { l.Error(err, "status update") } if cluster.Status.Phase != "Ready" { return ctrl.Result{RequeueAfter: requeueShort}, nil } return ctrl.Result{RequeueAfter: requeueLong}, nil } func (r *KeaClusterReconciler) fail(ctx context.Context, c *v1alpha1.KeaCluster, reason string, err error) (ctrl.Result, error) { c.Status.Phase = "Error" setReady(&c.Status.Conditions, c.Generation, false, reason, err.Error()) _ = r.Status().Update(ctx, c) return ctrl.Result{}, err } // buildInput gathers matching subnets/classes and the stable HA peer list. func (r *KeaClusterReconciler) buildInput(ctx context.Context, c *v1alpha1.KeaCluster) (kea.RenderInput, error) { var subnetList v1alpha1.KeaSubnetList if err := r.List(ctx, &subnetList, client.InNamespace(c.Namespace)); err != nil { return kea.RenderInput{}, err } var subnets []v1alpha1.KeaSubnet for _, s := range subnetList.Items { if s.Spec.ClusterRef == "" || s.Spec.ClusterRef == c.Name { subnets = append(subnets, s) } } var classList v1alpha1.KeaClientClassList if err := r.List(ctx, &classList, client.InNamespace(c.Namespace)); err != nil { return kea.RenderInput{}, err } var classes []v1alpha1.KeaClientClass for _, cc := range classList.Items { if cc.Spec.ClusterRef == "" || cc.Spec.ClusterRef == c.Name { classes = append(classes, cc) } } return kea.RenderInput{ Cluster: *c, Subnets: subnets, ClientClasses: classes, Peers: r.peers(c), }, nil } // peers returns stable HA peer identities (DNS only, no pod IPs). func (r *KeaClusterReconciler) peers(c *v1alpha1.KeaCluster) []kea.Peer { replicas := int32(1) if c.Spec.Replicas != nil { replicas = *c.Spec.Replicas } mode := c.Spec.HA.Mode if mode == "" { mode = v1alpha1.HAHotStandby } peers := make([]kea.Peer, 0, replicas) for i := int32(0); i < replicas; i++ { role := "backup" switch { case i == 0: role = "primary" case i == 1 && mode == v1alpha1.HAHotStandby: role = "standby" case i == 1: role = "secondary" } peers = append(peers, kea.Peer{ Name: fmt.Sprintf("server%d", i), URL: peerDNS(c.Name, c.Namespace, int(i)), Role: role, }) } return peers } func (r *KeaClusterReconciler) reconcileConfigMap(ctx context.Context, c *v1alpha1.KeaCluster) error { in, err := r.buildInput(ctx, c) if err != nil { return err } dhcp4, err := kea.RenderDHCP4(in) if err != nil { return err } agent, err := kea.RenderCtrlAgent() if err != nil { return err } cm := &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: configMapName(c.Name), Namespace: c.Namespace}} _, err = ctrl.CreateOrUpdate(ctx, r.Client, cm, func() error { cm.Labels = commonLabels(c.Name) cm.Data = map[string]string{ "kea-dhcp4.conf": dhcp4, "kea-ctrl-agent.conf": agent, "init.sh": kea.InitScript(), } return ctrl.SetControllerReference(c, cm, r.Scheme) }) return err } func (r *KeaClusterReconciler) reconcileServices(ctx context.Context, c *v1alpha1.KeaCluster) error { // Headless service for stable per-pod DNS (HA peer URLs, ctrl-agent). headless := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: headlessName(c.Name), Namespace: c.Namespace}} if _, err := ctrl.CreateOrUpdate(ctx, r.Client, headless, func() error { headless.Labels = commonLabels(c.Name) headless.Spec.ClusterIP = corev1.ClusterIPNone headless.Spec.PublishNotReadyAddresses = true headless.Spec.Selector = commonLabels(c.Name) headless.Spec.Ports = []corev1.ServicePort{ {Name: "ctrl", Port: kea.CtrlAgentPort, Protocol: corev1.ProtocolTCP}, } return ctrl.SetControllerReference(c, headless, r.Scheme) }); err != nil { return err } // Anycast DHCP service (LoadBalancer via PureLB by default). svc := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: serviceName(c.Name), Namespace: c.Namespace}} _, err := ctrl.CreateOrUpdate(ctx, r.Client, svc, func() error { svc.Labels = commonLabels(c.Name) svc.Annotations = mergeAnnotations(c.Spec.Service) svc.Spec.Selector = commonLabels(c.Name) svcType := c.Spec.Service.Type if svcType == "" { svcType = corev1.ServiceTypeLoadBalancer } svc.Spec.Type = svcType if c.Spec.Service.LoadBalancerIP != "" { svc.Spec.LoadBalancerIP = c.Spec.Service.LoadBalancerIP } if c.Spec.Service.LoadBalancerClass != nil { svc.Spec.LoadBalancerClass = c.Spec.Service.LoadBalancerClass } if svcType == corev1.ServiceTypeLoadBalancer || svcType == corev1.ServiceTypeNodePort { svc.Spec.ExternalTrafficPolicy = corev1.ServiceExternalTrafficPolicyLocal } svc.Spec.Ports = []corev1.ServicePort{ {Name: "dhcp", Port: kea.DHCP4Port, Protocol: corev1.ProtocolUDP}, } return ctrl.SetControllerReference(c, svc, r.Scheme) }) return err } func mergeAnnotations(s v1alpha1.ClusterServiceSpec) map[string]string { out := map[string]string{} for k, v := range s.Annotations { out[k] = v } if s.IPAddressPool != "" { out["purelb.io/service-group"] = s.IPAddressPool } if len(out) == 0 { return nil } return out } func (r *KeaClusterReconciler) reconcileStatefulSet(ctx context.Context, c *v1alpha1.KeaCluster) (*appsv1.StatefulSet, error) { hash, err := configHash(ctx, r.Client, c.Namespace, configMapName(c.Name)) if err != nil { return nil, err } image := c.Spec.Image if image == "" { image = kea.DefaultImage } replicas := int32(1) if c.Spec.Replicas != nil { replicas = *c.Spec.Replicas } sts := &appsv1.StatefulSet{ObjectMeta: metav1.ObjectMeta{Name: stsName(c.Name), Namespace: c.Namespace}} _, err = ctrl.CreateOrUpdate(ctx, r.Client, sts, func() error { sts.Labels = commonLabels(c.Name) sts.Spec.ServiceName = headlessName(c.Name) sts.Spec.Replicas = int32ptr(replicas) sts.Spec.Selector = &metav1.LabelSelector{MatchLabels: commonLabels(c.Name)} sts.Spec.PodManagementPolicy = appsv1.ParallelPodManagement sts.Spec.Template = r.podTemplate(c, image, hash) return ctrl.SetControllerReference(c, sts, r.Scheme) }) if err != nil { return nil, err } return sts, nil } func (r *KeaClusterReconciler) podTemplate(c *v1alpha1.KeaCluster, image, hash string) corev1.PodTemplateSpec { volProjected := corev1.Volume{ Name: "kea-etc", VolumeSource: corev1.VolumeSource{ ConfigMap: &corev1.ConfigMapVolumeSource{ LocalObjectReference: corev1.LocalObjectReference{Name: configMapName(c.Name)}, DefaultMode: int32ptr(0o755), }, }, } volRun := corev1.Volume{Name: "run", VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{}}} mounts := []corev1.VolumeMount{ {Name: "kea-etc", MountPath: kea.ConfigDir, ReadOnly: true}, {Name: "run", MountPath: kea.RunDir}, } // The initContainer finalises the per-pod config and bounded-waits for HA // peer DNS; the main containers then exec kea directly with no wrapper shell. initC := corev1.Container{ Name: kea.ContainerInit, Image: image, Command: []string{"/bin/sh", kea.InitScriptPath}, Env: initEnv(), Resources: c.Spec.Resources, VolumeMounts: mounts, } dhcp4 := corev1.Container{ Name: kea.ContainerDHCP4, Image: image, Command: []string{kea.DHCP4Bin, "-c", kea.DHCP4ConfPath}, Resources: c.Spec.Resources, Ports: []corev1.ContainerPort{ {Name: "dhcp", ContainerPort: kea.DHCP4Port, Protocol: corev1.ProtocolUDP}, }, VolumeMounts: mounts, } agent := corev1.Container{ Name: kea.ContainerCtrlAgent, Image: image, Command: []string{kea.CtrlAgentBin, "-c", kea.CtrlAgentConfPath}, Resources: c.Spec.Resources, Ports: []corev1.ContainerPort{ {Name: "ctrl", ContainerPort: kea.CtrlAgentPort, Protocol: corev1.ProtocolTCP}, }, VolumeMounts: mounts, ReadinessProbe: &corev1.Probe{ ProbeHandler: corev1.ProbeHandler{TCPSocket: &corev1.TCPSocketAction{Port: intstrFromInt(kea.CtrlAgentPort)}}, InitialDelaySeconds: 5, PeriodSeconds: 10, }, } labels := commonLabels(c.Name) return corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ Labels: labels, Annotations: map[string]string{"kea.unkin.net/config-hash": hash}, }, Spec: corev1.PodSpec{ InitContainers: []corev1.Container{initC}, Containers: []corev1.Container{dhcp4, agent}, Volumes: []corev1.Volume{volProjected, volRun}, NodeSelector: c.Spec.NodeSelector, Tolerations: c.Spec.Tolerations, Affinity: c.Spec.Affinity, }, } } // initEnv is the environment the embedded init.sh reads. Passing paths and // tunables as env vars (rather than interpolating them into the script text) // keeps init.sh static, committed and shellcheck-clean. func initEnv() []corev1.EnvVar { return []corev1.EnvVar{ {Name: kea.EnvPodName, ValueFrom: &corev1.EnvVarSource{ FieldRef: &corev1.ObjectFieldSelector{FieldPath: "metadata.name"}, }}, {Name: kea.EnvRunDir, Value: kea.RunDir}, {Name: kea.EnvConfigDir, Value: kea.ConfigDir}, {Name: kea.EnvThisServerPlaceholder, Value: kea.ThisServerPlaceholder}, {Name: kea.EnvDHCP4Bin, Value: kea.DHCP4Bin}, {Name: kea.EnvDHCP4Conf, Value: kea.DHCP4ConfPath}, {Name: kea.EnvCtrlAgentConf, Value: kea.CtrlAgentConfPath}, {Name: kea.EnvWaitAttempts, Value: strconv.Itoa(kea.WaitAttempts)}, {Name: kea.EnvWaitSleep, Value: strconv.Itoa(kea.WaitSleepSeconds)}, } } // reloadReadyPods best-effort hot-reloads config on ready pods via the // ctrl-agent REST channel, analogous to bind-operator's rndc reconfig. func (r *KeaClusterReconciler) reloadReadyPods(ctx context.Context, c *v1alpha1.KeaCluster) { if r.Control == nil { return } var pods corev1.PodList if err := r.List(ctx, &pods, client.InNamespace(c.Namespace), client.MatchingLabels(commonLabels(c.Name))); err != nil { return } l := log.FromContext(ctx) for i := range pods.Items { p := &pods.Items[i] if p.Status.PodIP == "" || !podReady(p) { continue } url := fmt.Sprintf("http://%s:%d/", p.Status.PodIP, kea.CtrlAgentPort) if err := r.Control.ConfigReload(ctx, url); err != nil { l.V(1).Info("config-reload failed", "pod", p.Name, "err", err.Error()) } } } func (r *KeaClusterReconciler) serviceIP(ctx context.Context, c *v1alpha1.KeaCluster) string { var svc corev1.Service if err := r.Get(ctx, types.NamespacedName{Namespace: c.Namespace, Name: serviceName(c.Name)}, &svc); err != nil { return "" } if len(svc.Status.LoadBalancer.Ingress) > 0 { return svc.Status.LoadBalancer.Ingress[0].IP } return svc.Spec.ClusterIP } func podReady(p *corev1.Pod) bool { for _, cond := range p.Status.Conditions { if cond.Type == corev1.PodReady { return cond.Status == corev1.ConditionTrue } } return false } func (r *KeaClusterReconciler) SetupWithManager(mgr ctrl.Manager) error { mapToClusters := func(ctx context.Context, obj client.Object) []reconcile.Request { var list v1alpha1.KeaClusterList if err := r.List(ctx, &list, client.InNamespace(obj.GetNamespace())); err != nil { return nil } var reqs []reconcile.Request for _, c := range list.Items { reqs = append(reqs, reconcile.Request{NamespacedName: types.NamespacedName{Namespace: c.Namespace, Name: c.Name}}) } return reqs } return ctrl.NewControllerManagedBy(mgr). For(&v1alpha1.KeaCluster{}). Owns(&appsv1.StatefulSet{}). Owns(&corev1.Service{}). Owns(&corev1.ConfigMap{}). Watches(&v1alpha1.KeaSubnet{}, handler.EnqueueRequestsFromMapFunc(mapToClusters)). Watches(&v1alpha1.KeaClientClass{}, handler.EnqueueRequestsFromMapFunc(mapToClusters)). Complete(r) }