fix for portallocator initialization (#423)

This commit is contained in:
Enrico Candino
2025-07-21 17:03:39 +02:00
committed by GitHub
parent c480bc339e
commit 1048e3f82d
3 changed files with 54 additions and 59 deletions
+50 -55
View File
@@ -8,7 +8,6 @@ import (
"gopkg.in/yaml.v2"
v1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/kubernetes/pkg/apis/core"
"k8s.io/kubernetes/pkg/registry/core/service/portallocator"
@@ -20,9 +19,15 @@ import (
const (
kubeletPortRangeConfigMapName = "k3k-kubelet-port-range"
webhookPortRangeConfigMapName = "k3k-webhook-port-range"
rangeKey = "range"
allocatedPortsKey = "allocatedPorts"
snapshotDataKey = "snapshotData"
)
type PortAllocator struct {
ctrlruntimeclient.Client
KubeletCM *v1.ConfigMap
WebhookCM *v1.ConfigMap
}
@@ -31,13 +36,13 @@ func NewPortAllocator(ctx context.Context, client ctrlruntimeclient.Client) (*Po
log := ctrl.LoggerFrom(ctx)
log.Info("starting port allocator")
var kubeletPortRangeCM, webhookPortRangeCM v1.ConfigMap
portRangeConfigMapNamespace := os.Getenv("CONTROLLER_NAMESPACE")
if portRangeConfigMapNamespace == "" {
return nil, fmt.Errorf("failed to find k3k controller namespace")
}
var kubeletPortRangeCM, webhookPortRangeCM v1.ConfigMap
kubeletPortRangeCM.Name = kubeletPortRangeConfigMapName
kubeletPortRangeCM.Namespace = portRangeConfigMapNamespace
@@ -45,6 +50,7 @@ func NewPortAllocator(ctx context.Context, client ctrlruntimeclient.Client) (*Po
webhookPortRangeCM.Namespace = portRangeConfigMapNamespace
return &PortAllocator{
Client: client,
KubeletCM: &kubeletPortRangeCM,
WebhookCM: &webhookPortRangeCM,
}, nil
@@ -52,11 +58,11 @@ func NewPortAllocator(ctx context.Context, client ctrlruntimeclient.Client) (*Po
func (a *PortAllocator) InitPortAllocatorConfig(ctx context.Context, client ctrlruntimeclient.Client, kubeletPortRange, webhookPortRange string) manager.Runnable {
return manager.RunnableFunc(func(ctx context.Context) error {
if err := a.getOrCreate(ctx, client, a.KubeletCM, kubeletPortRange); err != nil {
if err := a.getOrCreate(ctx, a.KubeletCM, kubeletPortRange); err != nil {
return err
}
if err := a.getOrCreate(ctx, client, a.WebhookCM, webhookPortRange); err != nil {
if err := a.getOrCreate(ctx, a.WebhookCM, webhookPortRange); err != nil {
return err
}
@@ -64,39 +70,27 @@ func (a *PortAllocator) InitPortAllocatorConfig(ctx context.Context, client ctrl
})
}
func (a *PortAllocator) cm(name, namespace, portRange string) *v1.ConfigMap {
return &v1.ConfigMap{
TypeMeta: metav1.TypeMeta{
APIVersion: "v1",
Kind: "ConfigMap",
},
ObjectMeta: metav1.ObjectMeta{
Name: name,
Namespace: namespace,
},
Data: map[string]string{
"range": portRange,
"allocatedPorts": "",
},
BinaryData: map[string][]byte{
"snapshotData": []byte(""),
},
}
}
func (a *PortAllocator) getOrCreate(ctx context.Context, client ctrlruntimeclient.Client, configmap *v1.ConfigMap, portRange string) error {
func (a *PortAllocator) getOrCreate(ctx context.Context, configmap *v1.ConfigMap, portRange string) error {
nn := types.NamespacedName{
Name: configmap.Name,
Namespace: configmap.Namespace,
}
if err := client.Get(ctx, nn, configmap); err != nil {
if err := a.Client.Get(ctx, nn, configmap); err != nil {
if !apierrors.IsNotFound(err) {
return err
}
// creating the configMap for the first time
configmap = a.cm(configmap.Name, configmap.Namespace, portRange)
if err := client.Create(ctx, configmap); err != nil {
configmap.Data = map[string]string{
rangeKey: portRange,
allocatedPortsKey: "",
}
configmap.BinaryData = map[string][]byte{
snapshotDataKey: []byte(""),
}
if err := a.Client.Create(ctx, configmap); err != nil {
return fmt.Errorf("failed to create port range configmap: %w", err)
}
}
@@ -104,36 +98,37 @@ func (a *PortAllocator) getOrCreate(ctx context.Context, client ctrlruntimeclien
return nil
}
func (a *PortAllocator) AllocateWebhookPort(ctx context.Context, cfg *Config) (int, error) {
return a.allocatePort(ctx, cfg, a.WebhookCM)
func (a *PortAllocator) AllocateWebhookPort(ctx context.Context, clusterName, clusterNamespace string) (int, error) {
return a.allocatePort(ctx, clusterName, clusterNamespace, a.WebhookCM)
}
func (a *PortAllocator) DeallocateWebhookPort(ctx context.Context, client ctrlruntimeclient.Client, clusterName, clusterNamespace string, webhookPort int) error {
return a.deallocatePort(ctx, client, clusterName, clusterNamespace, a.WebhookCM, webhookPort)
func (a *PortAllocator) DeallocateWebhookPort(ctx context.Context, clusterName, clusterNamespace string, webhookPort int) error {
return a.deallocatePort(ctx, clusterName, clusterNamespace, a.WebhookCM, webhookPort)
}
func (a *PortAllocator) AllocateKubeletPort(ctx context.Context, cfg *Config) (int, error) {
return a.allocatePort(ctx, cfg, a.KubeletCM)
func (a *PortAllocator) AllocateKubeletPort(ctx context.Context, clusterName, clusterNamespace string) (int, error) {
return a.allocatePort(ctx, clusterName, clusterNamespace, a.KubeletCM)
}
func (a *PortAllocator) DeallocateKubeletPort(ctx context.Context, client ctrlruntimeclient.Client, clusterName, clusterNamespace string, kubeletPort int) error {
return a.deallocatePort(ctx, client, clusterName, clusterNamespace, a.KubeletCM, kubeletPort)
func (a *PortAllocator) DeallocateKubeletPort(ctx context.Context, clusterName, clusterNamespace string, kubeletPort int) error {
return a.deallocatePort(ctx, clusterName, clusterNamespace, a.KubeletCM, kubeletPort)
}
// allocatePort will assign port to the cluster from a port Range configured for k3k
func (a *PortAllocator) allocatePort(ctx context.Context, cfg *Config, configMap *v1.ConfigMap) (int, error) {
portRange, ok := configMap.Data["range"]
func (a *PortAllocator) allocatePort(ctx context.Context, clusterName, clusterNamespace string, configMap *v1.ConfigMap) (int, error) {
portRange, ok := configMap.Data[rangeKey]
if !ok {
return 0, fmt.Errorf("port range is not initialized")
}
// get configMap first to avoid conflicts
if err := a.getOrCreate(ctx, cfg.client, configMap, portRange); err != nil {
if err := a.getOrCreate(ctx, configMap, portRange); err != nil {
return 0, err
}
clusterNamespaceName := cfg.cluster.Namespace + "/" + cfg.cluster.Name
clusterNamespaceName := clusterNamespace + "/" + clusterName
portsMap, err := parsePortMap(configMap.Data["allocatedPorts"])
portsMap, err := parsePortMap(configMap.Data[allocatedPortsKey])
if err != nil {
return 0, err
}
@@ -143,8 +138,8 @@ func (a *PortAllocator) allocatePort(ctx context.Context, cfg *Config, configMap
}
// allocate a new port and save the snapshot
snapshot := core.RangeAllocation{
Range: configMap.Data["range"],
Data: configMap.BinaryData["snapshotData"],
Range: configMap.Data[rangeKey],
Data: configMap.BinaryData[snapshotDataKey],
}
pa, err := portallocator.NewFromSnapshot(&snapshot)
@@ -163,7 +158,7 @@ func (a *PortAllocator) allocatePort(ctx context.Context, cfg *Config, configMap
return 0, err
}
if err := cfg.client.Update(ctx, configMap); err != nil {
if err := a.Client.Update(ctx, configMap); err != nil {
return 0, err
}
@@ -171,19 +166,19 @@ func (a *PortAllocator) allocatePort(ctx context.Context, cfg *Config, configMap
}
// deallocatePort will remove the port used by the cluster from the port range
func (a *PortAllocator) deallocatePort(ctx context.Context, client ctrlruntimeclient.Client, clusterName, clusterNamespace string, configMap *v1.ConfigMap, port int) error {
portRange, ok := configMap.Data["range"]
func (a *PortAllocator) deallocatePort(ctx context.Context, clusterName, clusterNamespace string, configMap *v1.ConfigMap, port int) error {
portRange, ok := configMap.Data[rangeKey]
if !ok {
return fmt.Errorf("port range is not initialized")
}
if err := a.getOrCreate(ctx, client, configMap, portRange); err != nil {
if err := a.getOrCreate(ctx, configMap, portRange); err != nil {
return err
}
clusterNamespaceName := clusterNamespace + "/" + clusterName
portsMap, err := parsePortMap(configMap.Data["allocatedPorts"])
portsMap, err := parsePortMap(configMap.Data[allocatedPortsKey])
if err != nil {
return err
}
@@ -194,8 +189,8 @@ func (a *PortAllocator) deallocatePort(ctx context.Context, client ctrlruntimecl
}
snapshot := core.RangeAllocation{
Range: configMap.Data["range"],
Data: configMap.BinaryData["snapshotData"],
Range: configMap.Data[rangeKey],
Data: configMap.BinaryData[snapshotDataKey],
}
pa, err := portallocator.NewFromSnapshot(&snapshot)
@@ -214,7 +209,7 @@ func (a *PortAllocator) deallocatePort(ctx context.Context, client ctrlruntimecl
}
}
return client.Update(ctx, configMap)
return a.Client.Update(ctx, configMap)
}
// parsePortMap will convert ConfigMap Data to a portMap of string keys and values of ints
@@ -243,15 +238,15 @@ func saveSnapshot(portAllocator *portallocator.PortAllocator, snapshot *core.Ran
return err
}
// update the configmap with the new portsMap and the new snapshot
configMap.BinaryData["snapshotData"] = snapshot.Data
configMap.Data["range"] = snapshot.Range
configMap.BinaryData[snapshotDataKey] = snapshot.Data
configMap.Data[rangeKey] = snapshot.Range
allocatedPortsData, err := serializePortMap(portsMap)
if err != nil {
return err
}
configMap.Data["allocatedPorts"] = allocatedPortsData
configMap.Data[allocatedPortsKey] = allocatedPortsData
return nil
}
+2 -2
View File
@@ -682,14 +682,14 @@ func (c *ClusterReconciler) ensureAgent(ctx context.Context, cluster *v1alpha1.C
if cluster.Spec.MirrorHostNodes {
var err error
kubeletPort, err = c.PortAllocator.AllocateKubeletPort(ctx, config)
kubeletPort, err = c.PortAllocator.AllocateKubeletPort(ctx, cluster.Name, cluster.Namespace)
if err != nil {
return err
}
cluster.Status.KubeletPort = kubeletPort
webhookPort, err = c.PortAllocator.AllocateWebhookPort(ctx, config)
webhookPort, err = c.PortAllocator.AllocateWebhookPort(ctx, cluster.Name, cluster.Namespace)
if err != nil {
return err
}
+2 -2
View File
@@ -41,11 +41,11 @@ func (c *ClusterReconciler) finalizeCluster(ctx context.Context, cluster *v1alph
if cluster.Spec.Mode == v1alpha1.SharedClusterMode && cluster.Spec.MirrorHostNodes {
log.Info("dellocating ports for kubelet and webhook")
if err := c.PortAllocator.DeallocateKubeletPort(ctx, c.Client, cluster.Name, cluster.Namespace, cluster.Status.KubeletPort); err != nil {
if err := c.PortAllocator.DeallocateKubeletPort(ctx, cluster.Name, cluster.Namespace, cluster.Status.KubeletPort); err != nil {
return reconcile.Result{}, err
}
if err := c.PortAllocator.DeallocateWebhookPort(ctx, c.Client, cluster.Name, cluster.Namespace, cluster.Status.WebhookPort); err != nil {
if err := c.PortAllocator.DeallocateWebhookPort(ctx, cluster.Name, cluster.Namespace, cluster.Status.WebhookPort); err != nil {
return reconcile.Result{}, err
}
}