Files
kea-operator/internal/controller/keacluster_controller.go
T
unkinben 66ae5f5f3c
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
Point HA peer URLs at per-pod ClusterIP Services
## Why
kea-dhcp4 crash-loops at HA hook load: kea 2.6's HA hook parses each peer url host as an IP literal and never resolves DNS, so the StatefulSet headless hostnames are rejected ("Failed to convert string to address ..."). Verified in-cluster that only an IP works (short name, FQDN both fail; `kea-dhcp4 -t` does not exercise this, which is why the v0.1.3 wait did not catch it). Pod IPs cannot be baked into the config because they change on restart and would roll-loop the StatefulSet via the config hash.

## How
- create one ClusterIP Service per HA peer, selecting the pod by its statefulset.kubernetes.io/pod-name label, with publishNotReadyAddresses so peers are routable during bootstrap
- render each HA peer url as its peer Service ClusterIP (a stable IP literal, safe in the config hash); reconcile Services before the ConfigMap and requeue until the ClusterIPs are allocated
2026-08-09 18:59:20 +10:00

473 lines
16 KiB
Go

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)
}
// Services first: the per-pod ClusterIP Services must exist (and have their
// ClusterIPs allocated) before the ConfigMap is rendered, because the HA
// peer URLs baked into kea-dhcp4.conf are those stable ClusterIPs.
if err := r.reconcileServices(ctx, &cluster); err != nil {
return r.fail(ctx, &cluster, "ServiceError", err)
}
if err := r.reconcileConfigMap(ctx, &cluster); err != nil {
return r.fail(ctx, &cluster, "ConfigError", 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)
}
}
peers, err := r.peers(ctx, c)
if err != nil {
return kea.RenderInput{}, err
}
return kea.RenderInput{
Cluster: *c,
Subnets: subnets,
ClientClasses: classes,
Peers: peers,
}, nil
}
// peers builds the HA peer list. Each URL points at the peer's per-pod ClusterIP
// Service address (a stable IP literal): kea 2.6's HA hook parses the peer URL
// host as an IP and never resolves DNS, so hostnames are rejected. Using the
// stable ClusterIP (not the pod IP) also keeps the config hash stable across
// pod restarts, so the StatefulSet does not roll-loop.
func (r *KeaClusterReconciler) peers(ctx context.Context, c *v1alpha1.KeaCluster) ([]kea.Peer, error) {
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"
}
svcName := peerServiceName(c.Name, int(i))
var svc corev1.Service
if err := r.Get(ctx, types.NamespacedName{Namespace: c.Namespace, Name: svcName}, &svc); err != nil {
return nil, fmt.Errorf("peer service %s: %w", svcName, err)
}
ip := svc.Spec.ClusterIP
if ip == "" || ip == corev1.ClusterIPNone {
return nil, fmt.Errorf("peer service %s has no ClusterIP allocated yet", svcName)
}
peers = append(peers, kea.Peer{
Name: fmt.Sprintf("server%d", i),
URL: peerURL(ip),
Role: role,
})
}
return peers, nil
}
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 (ctrl-agent discovery).
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
}
// Per-pod ClusterIP Services: one stable IP per HA peer. Kea's HA hook needs
// an IP literal for each peer URL, and a ClusterIP survives pod restarts, so
// it is safe to bake into the (roll-triggering) config hash. Not-ready
// addresses are published so peers are routable during HA bootstrap.
replicas := int32(1)
if c.Spec.Replicas != nil {
replicas = *c.Spec.Replicas
}
for i := int32(0); i < replicas; i++ {
ord := int(i)
peerSvc := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: peerServiceName(c.Name, ord), Namespace: c.Namespace}}
if _, err := ctrl.CreateOrUpdate(ctx, r.Client, peerSvc, func() error {
peerSvc.Labels = commonLabels(c.Name)
peerSvc.Spec.Type = corev1.ServiceTypeClusterIP
peerSvc.Spec.PublishNotReadyAddresses = true
peerSvc.Spec.Selector = map[string]string{statefulSetPodNameLabel: podName(c.Name, ord)}
peerSvc.Spec.Ports = []corev1.ServicePort{
{Name: "ctrl", Port: kea.CtrlAgentPort, TargetPort: intstrFromInt(kea.CtrlAgentPort), Protocol: corev1.ProtocolTCP},
}
return ctrl.SetControllerReference(c, peerSvc, 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)
}