From 0164c785ab6957ecafdc4bbcc954725ca4f13750 Mon Sep 17 00:00:00 2001 From: Enrico Candino Date: Tue, 27 Jan 2026 15:56:37 +0100 Subject: [PATCH] Show correct allocatable resources when a Policy is applied (#638) * wip * wip * wip * fix lint and tests * fixed bugs for missing resources * cleanup and refactor * removed coreClient from configureNode * added comments to distribute algorithm --- charts/k3k/templates/rbac.yaml | 2 + k3k-kubelet/kubelet.go | 13 +- k3k-kubelet/main.go | 1 + k3k-kubelet/provider/configure.go | 82 +------ k3k-kubelet/provider/configure_capacity.go | 200 ++++++++++++++++++ .../provider/configure_capacity_test.go | 147 +++++++++++++ main.go | 1 + pkg/controller/cluster/agent/shared.go | 5 + pkg/controller/policy/policy.go | 1 + pkg/log/zap.go | 1 + tests/common_test.go | 2 + 11 files changed, 377 insertions(+), 78 deletions(-) create mode 100644 k3k-kubelet/provider/configure_capacity.go create mode 100644 k3k-kubelet/provider/configure_capacity_test.go diff --git a/charts/k3k/templates/rbac.yaml b/charts/k3k/templates/rbac.yaml index b9923a48..76257dca 100644 --- a/charts/k3k/templates/rbac.yaml +++ b/charts/k3k/templates/rbac.yaml @@ -23,9 +23,11 @@ rules: resources: - "nodes" - "nodes/proxy" + - "namespaces" verbs: - "get" - "list" + - "watch" --- kind: ClusterRoleBinding apiVersion: rbac.authorization.k8s.io/v1 diff --git a/k3k-kubelet/kubelet.go b/k3k-kubelet/kubelet.go index 6aa1bb54..64d21f4c 100644 --- a/k3k-kubelet/kubelet.go +++ b/k3k-kubelet/kubelet.go @@ -270,7 +270,18 @@ func (k *kubelet) newProviderFunc(cfg config) nodeutil.NewProviderFunc { return nil, nil, errors.New("unable to make nodeutil provider: " + err.Error()) } - provider.ConfigureNode(k.logger, pc.Node, cfg.AgentHostname, k.port, k.agentIP, utilProvider.CoreClient, utilProvider.VirtualClient, k.virtualCluster, cfg.Version, cfg.MirrorHostNodes) + provider.ConfigureNode( + k.logger, + pc.Node, + cfg.AgentHostname, + k.port, + k.agentIP, + utilProvider.HostClient, + utilProvider.VirtualClient, + k.virtualCluster, + cfg.Version, + cfg.MirrorHostNodes, + ) return utilProvider, &provider.Node{}, nil } diff --git a/k3k-kubelet/main.go b/k3k-kubelet/main.go index 744cf470..2cf0a894 100644 --- a/k3k-kubelet/main.go +++ b/k3k-kubelet/main.go @@ -38,6 +38,7 @@ func main() { logger = zapr.NewLogger(log.New(debug, logFormat)) ctrlruntimelog.SetLogger(logger) + return nil }, RunE: run, diff --git a/k3k-kubelet/provider/configure.go b/k3k-kubelet/provider/configure.go index 895d72e7..a225e970 100644 --- a/k3k-kubelet/provider/configure.go +++ b/k3k-kubelet/provider/configure.go @@ -2,7 +2,6 @@ package provider import ( "context" - "time" "github.com/go-logr/logr" "k8s.io/apimachinery/pkg/types" @@ -10,16 +9,16 @@ import ( corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - typedv1 "k8s.io/client-go/kubernetes/typed/core/v1" "github.com/rancher/k3k/pkg/apis/k3k.io/v1beta1" ) -func ConfigureNode(logger logr.Logger, node *corev1.Node, hostname string, servicePort int, ip string, coreClient typedv1.CoreV1Interface, virtualClient client.Client, virtualCluster v1beta1.Cluster, version string, mirrorHostNodes bool) { +func ConfigureNode(logger logr.Logger, node *corev1.Node, hostname string, servicePort int, ip string, hostClient client.Client, virtualClient client.Client, virtualCluster v1beta1.Cluster, version string, mirrorHostNodes bool) { ctx := context.Background() + if mirrorHostNodes { - hostNode, err := coreClient.Nodes().Get(ctx, node.Name, metav1.GetOptions{}) - if err != nil { + var hostNode corev1.Node + if err := hostClient.Get(ctx, types.NamespacedName{Name: node.Name}, &hostNode); err != nil { logger.Error(err, "error getting host node for mirroring", err) } @@ -49,16 +48,7 @@ func ConfigureNode(logger logr.Logger, node *corev1.Node, hostname string, servi // configure versions node.Status.NodeInfo.KubeletVersion = version - updateNodeCapacityInterval := 10 * time.Second - ticker := time.NewTicker(updateNodeCapacityInterval) - - go func() { - for range ticker.C { - if err := updateNodeCapacity(ctx, coreClient, virtualClient, node.Name); err != nil { - logger.Error(err, "error updating node capacity") - } - } - }() + startNodeCapacityUpdater(ctx, logger, hostClient, virtualClient, virtualCluster, node.Name) } } @@ -107,65 +97,3 @@ func nodeConditions() []corev1.NodeCondition { }, } } - -// updateNodeCapacity will update the virtual node capacity (and the allocatable field) with the sum of all the resource in the host nodes. -// If the nodeLabels are specified only the matching nodes will be considered. -func updateNodeCapacity(ctx context.Context, coreClient typedv1.CoreV1Interface, virtualClient client.Client, virtualNodeName string) error { - capacity, allocatable, err := getResourcesFromNodes(ctx, coreClient, virtualNodeName) - if err != nil { - return err - } - - var virtualNode corev1.Node - if err := virtualClient.Get(ctx, types.NamespacedName{Name: virtualNodeName}, &virtualNode); err != nil { - return err - } - - virtualNode.Status.Capacity = capacity - virtualNode.Status.Allocatable = allocatable - - return virtualClient.Status().Update(ctx, &virtualNode) -} - -// getResourcesFromNodes will return a sum of all the resource capacity of the host nodes, and the allocatable resources. -// If some node labels are specified only the matching nodes will be considered. -func getResourcesFromNodes(ctx context.Context, coreClient typedv1.CoreV1Interface, virtualNodeName string) (corev1.ResourceList, corev1.ResourceList, error) { - node, err := coreClient.Nodes().Get(ctx, virtualNodeName, metav1.GetOptions{}) - if err != nil { - return nil, nil, err - } - - // sum all - virtualCapacityResources := corev1.ResourceList{} - virtualAvailableResources := corev1.ResourceList{} - - // check if the node is Ready - for _, condition := range node.Status.Conditions { - if condition.Type != corev1.NodeReady { - continue - } - - // if the node is not Ready then we can skip it - if condition.Status != corev1.ConditionTrue { - break - } - } - - // add all the available metrics to the virtual node - // TODO when using Quotas we should use that to actually limits the virtual node - for resourceName, resourceQuantity := range node.Status.Capacity { - virtualResource := virtualCapacityResources[resourceName] - - (&virtualResource).Add(resourceQuantity) - virtualCapacityResources[resourceName] = virtualResource - } - - for resourceName, resourceQuantity := range node.Status.Allocatable { - virtualResource := virtualAvailableResources[resourceName] - - (&virtualResource).Add(resourceQuantity) - virtualAvailableResources[resourceName] = virtualResource - } - - return virtualCapacityResources, virtualAvailableResources, nil -} diff --git a/k3k-kubelet/provider/configure_capacity.go b/k3k-kubelet/provider/configure_capacity.go new file mode 100644 index 00000000..0ecf5089 --- /dev/null +++ b/k3k-kubelet/provider/configure_capacity.go @@ -0,0 +1,200 @@ +package provider + +import ( + "context" + "sort" + "time" + + "github.com/go-logr/logr" + "k8s.io/apimachinery/pkg/api/resource" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + + corev1 "k8s.io/api/core/v1" + + "github.com/rancher/k3k/pkg/apis/k3k.io/v1beta1" +) + +const ( + // UpdateNodeCapacityInterval is the interval at which the node capacity is updated. + UpdateNodeCapacityInterval = 10 * time.Second +) + +// milliScaleResources is a set of resource names that are measured in milli-units (e.g., CPU). +// This is used to determine whether to use MilliValue() for calculations. +var milliScaleResources = map[corev1.ResourceName]struct{}{ + corev1.ResourceCPU: {}, + corev1.ResourceMemory: {}, + corev1.ResourceStorage: {}, + corev1.ResourceEphemeralStorage: {}, + corev1.ResourceRequestsCPU: {}, + corev1.ResourceRequestsMemory: {}, + corev1.ResourceRequestsStorage: {}, + corev1.ResourceRequestsEphemeralStorage: {}, + corev1.ResourceLimitsCPU: {}, + corev1.ResourceLimitsMemory: {}, + corev1.ResourceLimitsEphemeralStorage: {}, +} + +// StartNodeCapacityUpdater starts a goroutine that periodically updates the capacity +// of the virtual node based on host node capacity and any applied ResourceQuotas. +func startNodeCapacityUpdater(ctx context.Context, logger logr.Logger, hostClient client.Client, virtualClient client.Client, virtualCluster v1beta1.Cluster, virtualNodeName string) { + go func() { + ticker := time.NewTicker(UpdateNodeCapacityInterval) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + updateNodeCapacity(ctx, logger, hostClient, virtualClient, virtualCluster, virtualNodeName) + case <-ctx.Done(): + logger.Info("Stopping node capacity updates for node", "node", virtualNodeName) + return + } + } + }() +} + +// updateNodeCapacity will update the virtual node capacity (and the allocatable field) with the sum of all the resource in the host nodes. +// If the nodeLabels are specified only the matching nodes will be considered. +func updateNodeCapacity(ctx context.Context, logger logr.Logger, hostClient client.Client, virtualClient client.Client, virtualCluster v1beta1.Cluster, virtualNodeName string) { + // by default we get the resources of the same Node where the kubelet is running + var node corev1.Node + if err := hostClient.Get(ctx, types.NamespacedName{Name: virtualNodeName}, &node); err != nil { + logger.Error(err, "error getting virtual node for updating node capacity") + return + } + + allocatable := node.Status.Allocatable.DeepCopy() + + // we need to check if the virtual cluster resources are "limited" through ResourceQuotas + // If so we will use the minimum resources + + var quotas corev1.ResourceQuotaList + if err := hostClient.List(ctx, "as, &client.ListOptions{Namespace: virtualCluster.Namespace}); err != nil { + logger.Error(err, "error getting namespace for updating node capacity") + } + + if len(quotas.Items) > 0 { + resourceLists := []corev1.ResourceList{allocatable} + for _, q := range quotas.Items { + resourceLists = append(resourceLists, q.Status.Hard) + } + + mergedResourceLists := mergeResourceLists(resourceLists...) + + m, err := distributeQuotas(ctx, logger, virtualClient, mergedResourceLists) + if err != nil { + logger.Error(err, "error distributing policy quota") + } + + allocatable = m[virtualNodeName] + } + + var virtualNode corev1.Node + if err := virtualClient.Get(ctx, types.NamespacedName{Name: virtualNodeName}, &virtualNode); err != nil { + logger.Error(err, "error getting virtual node for updating node capacity") + return + } + + virtualNode.Status.Capacity = allocatable + virtualNode.Status.Allocatable = allocatable + + if err := virtualClient.Status().Update(ctx, &virtualNode); err != nil { + logger.Error(err, "error updating node capacity") + } +} + +// mergeResourceLists takes multiple resource lists and returns a single list that represents +// the most restrictive set of resources. For each resource name, it selects the minimum +// quantity found across all the provided lists. +func mergeResourceLists(resourceLists ...corev1.ResourceList) corev1.ResourceList { + merged := corev1.ResourceList{} + + for _, resourceList := range resourceLists { + for resName, qty := range resourceList { + existingQty, found := merged[resName] + + // If it's the first time we see it OR the new one is smaller -> Update + if !found || qty.Cmp(existingQty) < 0 { + merged[resName] = qty.DeepCopy() + } + } + } + + return merged +} + +// distributeQuotas divides the total resource quotas evenly among all active virtual nodes. +// This ensures that each virtual node reports a fair share of the available resources, +// preventing the scheduler from overloading a single node. +// +// The algorithm iterates over each resource, divides it as evenly as possible among the +// sorted virtual nodes, and distributes any remainder to the first few nodes to ensure +// all resources are allocated. Sorting the nodes by name guarantees a deterministic +// distribution. +func distributeQuotas(ctx context.Context, logger logr.Logger, virtualClient client.Client, quotas corev1.ResourceList) (map[string]corev1.ResourceList, error) { + // List all virtual nodes to distribute the quota stably. + var virtualNodeList corev1.NodeList + if err := virtualClient.List(ctx, &virtualNodeList); err != nil { + logger.Error(err, "error listing virtual nodes for stable capacity distribution, falling back to full quota") + return nil, err + } + + // If there are no virtual nodes, there's nothing to distribute. + numNodes := int64(len(virtualNodeList.Items)) + if numNodes == 0 { + logger.Info("error listing virtual nodes for stable capacity distribution, falling back to full quota") + return nil, nil + } + + // Sort nodes by name for a deterministic distribution of resources. + sort.Slice(virtualNodeList.Items, func(i, j int) bool { + return virtualNodeList.Items[i].Name < virtualNodeList.Items[j].Name + }) + + // Initialize the resource map for each virtual node. + resourceMap := make(map[string]corev1.ResourceList) + for _, virtualNode := range virtualNodeList.Items { + resourceMap[virtualNode.Name] = corev1.ResourceList{} + } + + // Distribute each resource type from the policy's hard quota + for resourceName, totalQuantity := range quotas { + // Use MilliValue for precise division, especially for resources like CPU, + // which are often expressed in milli-units. Otherwise, use the standard Value(). + var totalValue int64 + if _, found := milliScaleResources[resourceName]; found { + totalValue = totalQuantity.MilliValue() + } else { + totalValue = totalQuantity.Value() + } + + // Calculate the base quantity of the resource to be allocated per node. + // and the remainder that needs to be distributed among the nodes. + // + // For example, if totalValue is 2000 (e.g., 2 CPU) and there are 3 nodes: + // - quantityPerNode would be 666 (2000 / 3) + // - remainder would be 2 (2000 % 3) + // The first two nodes would get 667 (666 + 1), and the last one would get 666. + quantityPerNode := totalValue / numNodes + remainder := totalValue % numNodes + + // Iterate through the sorted virtual nodes to distribute the resource. + for _, virtualNode := range virtualNodeList.Items { + nodeQuantity := quantityPerNode + if remainder > 0 { + nodeQuantity++ + remainder-- + } + + if _, found := milliScaleResources[resourceName]; found { + resourceMap[virtualNode.Name][resourceName] = *resource.NewMilliQuantity(nodeQuantity, totalQuantity.Format) + } else { + resourceMap[virtualNode.Name][resourceName] = *resource.NewQuantity(nodeQuantity, totalQuantity.Format) + } + } + } + + return resourceMap, nil +} diff --git a/k3k-kubelet/provider/configure_capacity_test.go b/k3k-kubelet/provider/configure_capacity_test.go new file mode 100644 index 00000000..3dbda312 --- /dev/null +++ b/k3k-kubelet/provider/configure_capacity_test.go @@ -0,0 +1,147 @@ +package provider + +import ( + "context" + "testing" + + "github.com/go-logr/zapr" + "github.com/stretchr/testify/assert" + "go.uber.org/zap" + "k8s.io/apimachinery/pkg/api/resource" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func Test_distributeQuotas(t *testing.T) { + scheme := runtime.NewScheme() + err := corev1.AddToScheme(scheme) + assert.NoError(t, err) + + node1 := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "node-1"}} + node2 := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "node-2"}} + node3 := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "node-3"}} + + tests := []struct { + name string + virtualNodes []client.Object + quotas corev1.ResourceList + want map[string]corev1.ResourceList + wantErr bool + }{ + { + name: "no virtual nodes", + virtualNodes: []client.Object{}, + quotas: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("2"), + }, + want: map[string]corev1.ResourceList{}, + wantErr: false, + }, + { + name: "no quotas", + virtualNodes: []client.Object{node1, node2}, + quotas: corev1.ResourceList{}, + want: map[string]corev1.ResourceList{ + "node-1": {}, + "node-2": {}, + }, + wantErr: false, + }, + { + name: "even distribution of cpu and memory", + virtualNodes: []client.Object{node1, node2}, + quotas: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("2"), + corev1.ResourceMemory: resource.MustParse("4Gi"), + }, + want: map[string]corev1.ResourceList{ + "node-1": { + corev1.ResourceCPU: resource.MustParse("1"), + corev1.ResourceMemory: resource.MustParse("2Gi"), + }, + "node-2": { + corev1.ResourceCPU: resource.MustParse("1"), + corev1.ResourceMemory: resource.MustParse("2Gi"), + }, + }, + wantErr: false, + }, + { + name: "uneven distribution with remainder", + virtualNodes: []client.Object{node1, node2, node3}, + quotas: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("2"), // 2000m / 3 = 666m with 2m remainder + }, + want: map[string]corev1.ResourceList{ + "node-1": {corev1.ResourceCPU: resource.MustParse("667m")}, + "node-2": {corev1.ResourceCPU: resource.MustParse("667m")}, + "node-3": {corev1.ResourceCPU: resource.MustParse("666m")}, + }, + wantErr: false, + }, + { + name: "distribution of number resources", + virtualNodes: []client.Object{node1, node2, node3}, + quotas: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("2"), + corev1.ResourcePods: resource.MustParse("11"), + corev1.ResourceSecrets: resource.MustParse("9"), + "custom": resource.MustParse("8"), + }, + want: map[string]corev1.ResourceList{ + "node-1": { + corev1.ResourceCPU: resource.MustParse("667m"), + corev1.ResourcePods: resource.MustParse("4"), + corev1.ResourceSecrets: resource.MustParse("3"), + "custom": resource.MustParse("3"), + }, + "node-2": { + corev1.ResourceCPU: resource.MustParse("667m"), + corev1.ResourcePods: resource.MustParse("4"), + corev1.ResourceSecrets: resource.MustParse("3"), + "custom": resource.MustParse("3"), + }, + "node-3": { + corev1.ResourceCPU: resource.MustParse("666m"), + corev1.ResourcePods: resource.MustParse("3"), + corev1.ResourceSecrets: resource.MustParse("3"), + "custom": resource.MustParse("2"), + }, + }, + wantErr: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + fakeClient := fake.NewClientBuilder().WithScheme(scheme).WithObjects(tt.virtualNodes...).Build() + logger := zapr.NewLogger(zap.NewNop()) + + got, gotErr := distributeQuotas(context.Background(), logger, fakeClient, tt.quotas) + if tt.wantErr { + assert.Error(t, gotErr) + } else { + assert.NoError(t, gotErr) + } + + assert.Equal(t, len(tt.want), len(got), "Number of nodes in result should match") + + for nodeName, expectedResources := range tt.want { + actualResources, ok := got[nodeName] + assert.True(t, ok, "Node %s not found in result", nodeName) + + assert.Equal(t, len(expectedResources), len(actualResources), "Number of resources for node %s should match", nodeName) + + for resName, expectedQty := range expectedResources { + actualQty, ok := actualResources[resName] + assert.True(t, ok, "Resource %s not found for node %s", resName, nodeName) + assert.True(t, expectedQty.Equal(actualQty), "Resource %s for node %s did not match. want: %s, got: %s", resName, nodeName, expectedQty.String(), actualQty.String()) + } + } + }) + } +} diff --git a/main.go b/main.go index 48dbfca5..c1ede9ff 100644 --- a/main.go +++ b/main.go @@ -57,6 +57,7 @@ func main() { }, PersistentPreRun: func(cmd *cobra.Command, args []string) { cmds.InitializeConfig(cmd) + logger = zapr.NewLogger(log.New(debug, logFormat)) }, RunE: run, diff --git a/pkg/controller/cluster/agent/shared.go b/pkg/controller/cluster/agent/shared.go index 1d599810..ecd96fe5 100644 --- a/pkg/controller/cluster/agent/shared.go +++ b/pkg/controller/cluster/agent/shared.go @@ -386,6 +386,11 @@ func (s *SharedAgent) role(ctx context.Context) error { Resources: []string{"events"}, Verbs: []string{"create"}, }, + { + APIGroups: []string{""}, + Resources: []string{"resourcequotas"}, + Verbs: []string{"get", "watch", "list"}, + }, { APIGroups: []string{"networking.k8s.io"}, Resources: []string{"ingresses"}, diff --git a/pkg/controller/policy/policy.go b/pkg/controller/policy/policy.go index 2e89ecb7..b64817d0 100644 --- a/pkg/controller/policy/policy.go +++ b/pkg/controller/policy/policy.go @@ -165,6 +165,7 @@ func nodeEventHandler(r *VirtualClusterPolicyReconciler) handler.Funcs { if oldNode.Spec.PodCIDR != newNode.Spec.PodCIDR { podCIDRChanged = true } + if !reflect.DeepEqual(oldNode.Spec.PodCIDRs, newNode.Spec.PodCIDRs) { podCIDRChanged = true } diff --git a/pkg/log/zap.go b/pkg/log/zap.go index 8e3ddcdb..8129cc39 100644 --- a/pkg/log/zap.go +++ b/pkg/log/zap.go @@ -27,6 +27,7 @@ func newEncoder(format string) zapcore.Encoder { encCfg.EncodeTime = zapcore.ISO8601TimeEncoder var encoder zapcore.Encoder + if format == "text" { encCfg.EncodeLevel = zapcore.CapitalColorLevelEncoder encoder = zapcore.NewConsoleEncoder(encCfg) diff --git a/tests/common_test.go b/tests/common_test.go index ea6307cf..cab06d4c 100644 --- a/tests/common_test.go +++ b/tests/common_test.go @@ -233,6 +233,7 @@ func NewVirtualK8sClientAndConfig(cluster *v1beta1.Cluster) (*kubernetes.Clients kubeletAltName := fmt.Sprintf("k3k-%s-kubelet", cluster.Name) vKubeconfig.AltNames = certs.AddSANs([]string{hostIP, kubeletAltName}) config, err = vKubeconfig.Generate(ctx, k8sClient, cluster, hostIP, 0) + return err }). WithTimeout(time.Minute * 2). @@ -266,6 +267,7 @@ func NewVirtualK8sClientAndKubeconfig(cluster *v1beta1.Cluster) (*kubernetes.Cli kubeletAltName := fmt.Sprintf("k3k-%s-kubelet", cluster.Name) vKubeconfig.AltNames = certs.AddSANs([]string{hostIP, kubeletAltName}) config, err = vKubeconfig.Generate(ctx, k8sClient, cluster, hostIP, 0) + return err }). WithTimeout(time.Minute * 2).