This is an automated email from the ASF dual-hosted git repository.

andrijapanicsb pushed a commit to branch cks-live-resize-capacity-reconciliation
in repository https://gitbox.apache.org/repos/asf/cloudstack.git

commit 6f6656f491717a9df02fee3e0f4951cb4821945c
Author: andrijapanicsb <[email protected]>
AuthorDate: Tue Sep 22 12:47:24 2026 +0200

    CKS: reconcile node capacity after live resize
---
 .../KubernetesClusterScaleWorker.java              |  99 ++++++++++-
 .../KubernetesClusterNodeCapacityReconciler.java   |  67 ++++++++
 ...ubernetesClusterNodeCapacityReconcilerImpl.java | 190 +++++++++++++++++++++
 .../spring-kubernetes-service-context.xml          |   1 +
 .../KubernetesClusterScaleWorkerTest.java          |  33 ++++
 ...netesClusterNodeCapacityReconcilerImplTest.java |  71 ++++++++
 6 files changed, 455 insertions(+), 6 deletions(-)

diff --git 
a/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorker.java
 
b/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorker.java
index 08513dbd448..a1dbf28b2bd 100644
--- 
a/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorker.java
+++ 
b/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorker.java
@@ -28,6 +28,8 @@ import java.util.Objects;
 import java.util.Set;
 import java.util.stream.Collectors;
 
+import javax.inject.Inject;
+
 import 
com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType;
 import com.cloud.service.ServiceOfferingVO;
 import com.cloud.storage.VMTemplateVO;
@@ -43,16 +45,21 @@ import com.cloud.exception.InsufficientCapacityException;
 import com.cloud.exception.ManagementServerException;
 import com.cloud.exception.NetworkRuleConflictException;
 import com.cloud.exception.ResourceUnavailableException;
-import com.cloud.exception.VirtualMachineMigrationException;
+import com.cloud.hypervisor.Hypervisor;
 import com.cloud.kubernetes.cluster.KubernetesCluster;
 import com.cloud.kubernetes.cluster.KubernetesClusterManagerImpl;
 import com.cloud.kubernetes.cluster.KubernetesClusterService;
 import com.cloud.kubernetes.cluster.KubernetesClusterVO;
 import com.cloud.kubernetes.cluster.KubernetesClusterVmMapVO;
 import com.cloud.kubernetes.cluster.utils.KubernetesClusterUtil;
+import 
com.cloud.kubernetes.cluster.utils.KubernetesClusterNodeCapacityReconciler;
+import 
com.cloud.kubernetes.cluster.utils.KubernetesClusterNodeCapacityReconciler.NodeAccess;
+import 
com.cloud.kubernetes.cluster.utils.KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot;
+import 
com.cloud.kubernetes.cluster.utils.KubernetesClusterNodeCapacityReconcilerImpl;
 import com.cloud.network.IpAddress;
 import com.cloud.network.Network;
 import com.cloud.network.rules.FirewallRule;
+import com.cloud.network.rules.PortForwardingRuleVO;
 import com.cloud.offering.ServiceOffering;
 import com.cloud.storage.LaunchPermissionVO;
 import com.cloud.uservm.UserVm;
@@ -79,6 +86,9 @@ public class KubernetesClusterScaleWorker extends 
KubernetesClusterResourceModif
     private Boolean isAutoscalingEnabled;
     private long scaleTimeoutTime;
 
+    @Inject
+    protected KubernetesClusterNodeCapacityReconciler 
kubernetesClusterNodeCapacityReconciler;
+
     protected KubernetesClusterScaleWorker(final KubernetesCluster 
kubernetesCluster, final KubernetesClusterManagerImpl clusterManager) {
         super(kubernetesCluster, clusterManager);
     }
@@ -359,15 +369,60 @@ public class KubernetesClusterScaleWorker extends 
KubernetesClusterResourceModif
         for (long i = 0; i < tobeScaledVMCount; i++) {
             KubernetesClusterVmMapVO vmMapVO = vmList.get((int) i);
             UserVmVO userVM = userVmDao.findById(vmMapVO.getVmId());
+            if (userVM == null) {
+                logTransitStateAndThrow(Level.ERROR, String.format("Scaling 
Kubernetes cluster : %s failed, unable to find cluster VM %s",
+                        kubernetesCluster.getName(), vmMapVO.getVmId()), 
kubernetesCluster.getId(), KubernetesCluster.Event.OperationFailed);
+            }
+            ServiceOffering oldOffering = 
serviceOfferingDao.findById(userVM.getServiceOfferingId());
+            boolean capacityChanged = 
KubernetesClusterNodeCapacityReconcilerImpl.capacityChanged(oldOffering, 
serviceOffering);
+            if (capacityChanged && KubernetesCluster.State.Running == 
originalState && ETCD == nodeType
+                    && KubernetesCluster.ClusterType.CloudManaged == 
kubernetesCluster.getClusterType()) {
+                logTransitStateAndThrow(Level.ERROR, String.format("Scaling 
Kubernetes cluster : %s cannot live-resize dedicated etcd VM %s; " +
+                        "etcd health reconciliation is not implemented", 
kubernetesCluster.getName(), userVM.getDisplayName()),
+                        kubernetesCluster.getId(), 
KubernetesCluster.Event.OperationFailed);
+            }
+            boolean reconcileKubernetes = shouldReconcileNodeCapacity(vmMapVO, 
userVM, oldOffering, serviceOffering, nodeType);
+            NodeAccess nodeAccess = null;
+            NodeCapacitySnapshot before = null;
             boolean result = false;
+            boolean cksCordonCompleted = false;
+            boolean resizeSucceeded = serviceOffering.getId() == 
userVM.getServiceOfferingId();
             try {
-                result = userVmManager.upgradeVirtualMachine(userVM.getId(), 
serviceOffering.getId(), new HashMap<String, String>());
-            } catch (RuntimeException | ResourceUnavailableException | 
ManagementServerException | VirtualMachineMigrationException e) {
+                if (reconcileKubernetes) {
+                    nodeAccess = resolveNodeAccess(userVM);
+                    before = 
kubernetesClusterNodeCapacityReconciler.captureBefore(kubernetesCluster, 
userVM, nodeAccess);
+                    
kubernetesClusterNodeCapacityReconciler.cordonIfNeeded(kubernetesCluster, 
userVM, before, nodeAccess, scaleTimeoutTime);
+                    cksCordonCompleted = !before.isUnschedulable() || 
before.isCloudStackResizeCordon();
+                }
+                if (serviceOffering.getId() != userVM.getServiceOfferingId()) {
+                    result = 
userVmManager.upgradeVirtualMachine(userVM.getId(), serviceOffering.getId(), 
new HashMap<String, String>());
+                    if (!result) {
+                        logTransitStateAndThrow(Level.WARN, 
String.format("Scaling Kubernetes cluster : %s failed, unable to scale cluster 
VM : %s", kubernetesCluster.getName(), 
userVM.getDisplayName()),kubernetesCluster.getId(), 
KubernetesCluster.Event.OperationFailed);
+                    }
+                    resizeSucceeded = true;
+                    userVM = userVmDao.findById(userVM.getId());
+                    if (userVM == null || serviceOffering.getId() != 
userVM.getServiceOfferingId()) {
+                        logTransitStateAndThrow(Level.WARN, 
String.format("Scaling Kubernetes cluster : %s failed, VM %s did not reach 
target offering", kubernetesCluster.getName(), 
vmMapVO.getVmId()),kubernetesCluster.getId(), 
KubernetesCluster.Event.OperationFailed);
+                    }
+                }
+                if (reconcileKubernetes) {
+                    
kubernetesClusterNodeCapacityReconciler.verifyGuestResources(userVM, 
serviceOffering, before, nodeAccess, scaleTimeoutTime);
+                    if 
(!kubernetesClusterNodeCapacityReconciler.isKubernetesResourcesCurrent(before, 
serviceOffering)) {
+                        
kubernetesClusterNodeCapacityReconciler.restartKubelet(userVM, nodeAccess, 
scaleTimeoutTime);
+                        
kubernetesClusterNodeCapacityReconciler.waitForKubernetesResources(kubernetesCluster,
 userVM, serviceOffering, before, nodeAccess, scaleTimeoutTime);
+                    }
+                    
kubernetesClusterNodeCapacityReconciler.restoreSchedulability(kubernetesCluster,
 userVM, before, nodeAccess, scaleTimeoutTime);
+                }
+            } catch (Exception e) {
+                if (reconcileKubernetes && cksCordonCompleted && 
!resizeSucceeded) {
+                    try {
+                        
kubernetesClusterNodeCapacityReconciler.restoreSchedulability(kubernetesCluster,
 userVM, before, nodeAccess, scaleTimeoutTime);
+                    } catch (Exception cleanupException) {
+                        logger.warn("Unable to restore schedulability for CKS 
node {} after a pre-resize failure", userVM.getUuid(), cleanupException);
+                    }
+                }
                 logTransitStateAndThrow(Level.ERROR, String.format("Scaling 
Kubernetes cluster : %s failed, unable to scale cluster VM : %s due to %s", 
kubernetesCluster.getName(), userVM.getDisplayName(), e.getMessage()), 
kubernetesCluster.getId(), KubernetesCluster.Event.OperationFailed, e);
             }
-            if (!result) {
-                logTransitStateAndThrow(Level.WARN, String.format("Scaling 
Kubernetes cluster : %s failed, unable to scale cluster VM : %s", 
kubernetesCluster.getName(), 
userVM.getDisplayName()),kubernetesCluster.getId(), 
KubernetesCluster.Event.OperationFailed);
-            }
             if (System.currentTimeMillis() > scaleTimeoutTime) {
                 logTransitStateAndThrow(Level.WARN, String.format("Scaling 
Kubernetes cluster : %s failed, scaling action timed out", 
kubernetesCluster.getName()),kubernetesCluster.getId(), 
KubernetesCluster.Event.OperationFailed);
             }
@@ -375,6 +430,38 @@ public class KubernetesClusterScaleWorker extends 
KubernetesClusterResourceModif
         kubernetesCluster = updateKubernetesClusterEntryForNodeType(null, 
nodeType, serviceOffering, updateNodeOffering, updateClusterOffering);
     }
 
+    private NodeAccess resolveNodeAccess(UserVm userVM) {
+        Pair<String, Integer> controlAccess = 
getKubernetesClusterServerIpSshPort(null);
+        if (StringUtils.isBlank(controlAccess.first())) {
+            throw new CloudRuntimeException(String.format("Unable to resolve 
control-plane SSH for Kubernetes cluster %s", kubernetesCluster.getUuid()));
+        }
+        if (manager.isDirectAccess(network)) {
+            if (StringUtils.isBlank(userVM.getPrivateIpAddress())) {
+                throw new CloudRuntimeException(String.format("Unable to 
resolve private SSH address for VM %s", userVM.getUuid()));
+            }
+            return new NodeAccess(controlAccess.first(), 
controlAccess.second(), userVM.getPrivateIpAddress(), DEFAULT_SSH_PORT,
+                    getControlNodeLoginUser(), 
getManagementServerSshPublicKeyFile());
+        }
+        PortForwardingRuleVO sshRule = 
portForwardingRulesDao.listByVm(userVM.getId()).stream()
+                .filter(rule -> rule.getDestinationPortStart() == 
DEFAULT_SSH_PORT)
+                .filter(rule -> 
!FirewallRule.State.Revoke.equals(rule.getState()))
+                .findFirst().orElse(null);
+        if (sshRule == null) {
+            throw new CloudRuntimeException(String.format("Unable to resolve 
SSH port-forwarding rule for VM %s", userVM.getUuid()));
+        }
+        return new NodeAccess(controlAccess.first(), controlAccess.second(), 
controlAccess.first(), sshRule.getSourcePortStart(),
+                getControlNodeLoginUser(), 
getManagementServerSshPublicKeyFile());
+    }
+
+    protected boolean shouldReconcileNodeCapacity(KubernetesClusterVmMapVO 
vmMapVO, UserVm userVM,
+                                                   ServiceOffering 
oldOffering, ServiceOffering targetOffering,
+                                                   KubernetesClusterNodeType 
nodeType) {
+        return !vmMapVO.isExternalNode()
+                && KubernetesCluster.ClusterType.CloudManaged == 
kubernetesCluster.getClusterType()
+                && Hypervisor.HypervisorType.KVM == userVM.getHypervisorType()
+                && 
kubernetesClusterNodeCapacityReconciler.requiresKubeletRefresh(oldOffering, 
targetOffering, nodeType, originalState);
+    }
+
     private void removeNodesFromCluster(List<KubernetesClusterVmMapVO> vmMaps) 
throws CloudRuntimeException {
         for (KubernetesClusterVmMapVO vmMapVO : vmMaps) {
             UserVmVO userVM = userVmDao.findById(vmMapVO.getVmId());
diff --git 
a/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconciler.java
 
b/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconciler.java
new file mode 100644
index 00000000000..f52458774ac
--- /dev/null
+++ 
b/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconciler.java
@@ -0,0 +1,67 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+package com.cloud.kubernetes.cluster.utils;
+
+import java.io.File;
+
+import com.cloud.kubernetes.cluster.KubernetesCluster;
+import 
com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType;
+import com.cloud.offering.ServiceOffering;
+import com.cloud.uservm.UserVm;
+
+/** Reconciles the guest and Kubernetes resource views after a live VM resize. 
*/
+public interface KubernetesClusterNodeCapacityReconciler {
+    NodeCapacitySnapshot captureBefore(KubernetesCluster cluster, UserVm vm, 
NodeAccess access) throws Exception;
+    boolean requiresKubeletRefresh(ServiceOffering oldOffering, 
ServiceOffering newOffering,
+            KubernetesClusterNodeType nodeType, KubernetesCluster.State 
originalState);
+    void cordonIfNeeded(KubernetesCluster cluster, UserVm vm, 
NodeCapacitySnapshot before, NodeAccess access, long deadline) throws Exception;
+    void verifyGuestResources(UserVm vm, ServiceOffering target, 
NodeCapacitySnapshot before, NodeAccess access, long deadline) throws Exception;
+    boolean isKubernetesResourcesCurrent(NodeCapacitySnapshot snapshot, 
ServiceOffering target);
+    void restartKubelet(UserVm vm, NodeAccess access, long deadline) throws 
Exception;
+    NodeCapacitySnapshot waitForKubernetesResources(KubernetesCluster cluster, 
UserVm vm, ServiceOffering target,
+            NodeCapacitySnapshot before, NodeAccess access, long deadline) 
throws Exception;
+    void restoreSchedulability(KubernetesCluster cluster, UserVm vm, 
NodeCapacitySnapshot before, NodeAccess access, long deadline) throws Exception;
+
+    final class NodeAccess {
+        private final String controlAddress; private final int controlPort; 
private final String nodeAddress; private final int nodePort;
+        private final String user; private final File sshKeyFile;
+        public NodeAccess(String controlAddress, int controlPort, String 
nodeAddress, int nodePort, String user, File sshKeyFile) {
+            this.controlAddress = controlAddress; this.controlPort = 
controlPort; this.nodeAddress = nodeAddress; this.nodePort = nodePort;
+            this.user = user; this.sshKeyFile = sshKeyFile;
+        }
+        public String getControlAddress() { return controlAddress; } public 
int getControlPort() { return controlPort; }
+        public String getNodeAddress() { return nodeAddress; } public int 
getNodePort() { return nodePort; }
+        public String getUser() { return user; } public File getSshKeyFile() { 
return sshKeyFile; }
+    }
+
+    final class NodeCapacitySnapshot {
+        private final boolean unschedulable; private final boolean 
cloudStackResizeCordon; private final long guestOnlineCpuCount;
+        private final long guestMemoryKiB; private final long 
capacityCpuMillis; private final long capacityMemoryBytes;
+        private final long allocatableCpuMillis; private final long 
allocatableMemoryBytes; private final boolean ready;
+        public NodeCapacitySnapshot(boolean unschedulable, boolean 
cloudStackResizeCordon, long guestOnlineCpuCount, long guestMemoryKiB,
+                long capacityCpuMillis, long capacityMemoryBytes, long 
allocatableCpuMillis, long allocatableMemoryBytes, boolean ready) {
+            this.unschedulable = unschedulable; this.cloudStackResizeCordon = 
cloudStackResizeCordon; this.guestOnlineCpuCount = guestOnlineCpuCount;
+            this.guestMemoryKiB = guestMemoryKiB; this.capacityCpuMillis = 
capacityCpuMillis; this.capacityMemoryBytes = capacityMemoryBytes;
+            this.allocatableCpuMillis = allocatableCpuMillis; 
this.allocatableMemoryBytes = allocatableMemoryBytes; this.ready = ready;
+        }
+        public boolean isUnschedulable() { return unschedulable; } public 
boolean isCloudStackResizeCordon() { return cloudStackResizeCordon; }
+        public long getGuestOnlineCpuCount() { return guestOnlineCpuCount; } 
public long getGuestMemoryKiB() { return guestMemoryKiB; }
+        public long getCapacityCpuMillis() { return capacityCpuMillis; } 
public long getCapacityMemoryBytes() { return capacityMemoryBytes; }
+        public long getAllocatableCpuMillis() { return allocatableCpuMillis; } 
public long getAllocatableMemoryBytes() { return allocatableMemoryBytes; }
+        public boolean isReady() { return ready; }
+    }
+}
diff --git 
a/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImpl.java
 
b/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImpl.java
new file mode 100644
index 00000000000..a6785ec3fdf
--- /dev/null
+++ 
b/plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImpl.java
@@ -0,0 +1,190 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+package com.cloud.kubernetes.cluster.utils;
+
+import java.util.Objects;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+
+import com.cloud.kubernetes.cluster.KubernetesCluster;
+import 
com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType;
+import com.cloud.offering.ServiceOffering;
+import com.cloud.uservm.UserVm;
+import com.cloud.utils.Pair;
+import com.cloud.utils.exception.CloudRuntimeException;
+import com.cloud.utils.ssh.SshHelper;
+import com.google.gson.JsonArray;
+import com.google.gson.JsonElement;
+import com.google.gson.JsonObject;
+import com.google.gson.JsonParser;
+import org.apache.commons.lang3.StringUtils;
+
+/** SSH-backed implementation for a rolling CKS live-resize reconciliation. */
+public class KubernetesClusterNodeCapacityReconcilerImpl implements 
KubernetesClusterNodeCapacityReconciler {
+    static final long MINIMUM_MEMORY_OVERHEAD_BYTES = 128L * 1024L * 1024L;
+    static final long KUBERNETES_MEMORY_REPORTING_TOLERANCE_BYTES = 16L * 
1024L * 1024L;
+    private static final int COMMAND_TIMEOUT_MS = 30000;
+    private static final int POLL_INTERVAL_MS = 5000;
+    private static final Pattern NODE_NAME = 
Pattern.compile("[A-Za-z0-9][A-Za-z0-9.-]{0,252}");
+    private static final Pattern MEMORY_QUANTITY = 
Pattern.compile("([0-9]+)(Ki|Mi|Gi|Ti|K|M|G|T)?");
+    private static final String RESIZE_ANNOTATION = 
"cloudstack.apache.org/cks-live-resize";
+
+    @Override
+    public boolean requiresKubeletRefresh(ServiceOffering oldOffering, 
ServiceOffering newOffering,
+            KubernetesClusterNodeType nodeType, KubernetesCluster.State 
originalState) {
+        return KubernetesCluster.State.Running == originalState
+                && (KubernetesClusterNodeType.WORKER == nodeType || 
KubernetesClusterNodeType.CONTROL == nodeType)
+                && capacityChanged(oldOffering, newOffering);
+    }
+
+    @Override
+    public NodeCapacitySnapshot captureBefore(KubernetesCluster cluster, 
UserVm vm, NodeAccess access) throws Exception {
+        return nodeSnapshot(cluster, vm, access, guestResources(access));
+    }
+
+    @Override
+    public void cordonIfNeeded(KubernetesCluster cluster, UserVm vm, 
NodeCapacitySnapshot before, NodeAccess access, long deadline) throws Exception 
{
+        if (before.isUnschedulable() && !before.isCloudStackResizeCordon()) 
return;
+        String node = nodeName(vm);
+        executeControl(access, "sudo /opt/bin/kubectl annotate node " + node + 
" " + RESIZE_ANNOTATION + "=" + cluster.getUuid() + " --overwrite");
+        executeControl(access, "sudo /opt/bin/kubectl cordon " + node);
+        while (System.currentTimeMillis() < deadline) {
+            if (nodeSnapshot(cluster, vm, access, null).isUnschedulable()) 
return;
+            sleep();
+        }
+        throw failure("CORDON", vm, "Kubernetes node did not become 
unschedulable");
+    }
+
+    @Override
+    public void verifyGuestResources(UserVm vm, ServiceOffering target, 
NodeCapacitySnapshot before, NodeAccess access, long deadline) throws Exception 
{
+        while (System.currentTimeMillis() < deadline) {
+            if (guestMatchesTarget(guestResources(access), target)) return;
+            sleep();
+        }
+        throw failure("GUEST_VERIFY", vm, "guest CPU or memory did not reach 
the target offering");
+    }
+
+    @Override
+    public boolean isKubernetesResourcesCurrent(NodeCapacitySnapshot snapshot, 
ServiceOffering target) {
+        GuestResources guest = new 
GuestResources(snapshot.getGuestOnlineCpuCount(), snapshot.getGuestMemoryKiB());
+        return guestMatchesTarget(guest, target) && 
kubernetesMatchesGuest(snapshot, snapshot);
+    }
+
+    @Override
+    public void restartKubelet(UserVm vm, NodeAccess access, long deadline) 
throws Exception {
+        executeNode(access, "sudo systemctl restart kubelet");
+        while (System.currentTimeMillis() < deadline) {
+            if ("active".equals(executeNode(access, "sudo systemctl is-active 
kubelet").trim())) return;
+            sleep();
+        }
+        throw failure("KUBELET_RESTART", vm, "kubelet did not become active");
+    }
+
+    @Override
+    public NodeCapacitySnapshot waitForKubernetesResources(KubernetesCluster 
cluster, UserVm vm, ServiceOffering target,
+            NodeCapacitySnapshot before, NodeAccess access, long deadline) 
throws Exception {
+        while (System.currentTimeMillis() < deadline) {
+            GuestResources guest = guestResources(access);
+            NodeCapacitySnapshot observed = nodeSnapshot(cluster, vm, access, 
guest);
+            if (guestMatchesTarget(guest, target) && 
kubernetesMatchesGuest(observed, before)) return observed;
+            sleep();
+        }
+        throw failure("CAPACITY_VERIFY", vm, "Kubernetes capacity did not 
match the resized guest");
+    }
+
+    @Override
+    public void restoreSchedulability(KubernetesCluster cluster, UserVm vm, 
NodeCapacitySnapshot before, NodeAccess access, long deadline) throws Exception 
{
+        if (before.isUnschedulable() && !before.isCloudStackResizeCordon()) 
return;
+        String node = nodeName(vm);
+        executeControl(access, "sudo /opt/bin/kubectl uncordon " + node);
+        executeControl(access, "sudo /opt/bin/kubectl annotate node " + node + 
" " + RESIZE_ANNOTATION + "-");
+        while (System.currentTimeMillis() < deadline) {
+            if (!nodeSnapshot(cluster, vm, access, null).isUnschedulable()) 
return;
+            sleep();
+        }
+        throw failure("RESTORE_SCHEDULABILITY", vm, "Kubernetes node remained 
unschedulable");
+    }
+
+    public static boolean capacityChanged(ServiceOffering oldOffering, 
ServiceOffering targetOffering) {
+        return oldOffering != null && targetOffering != null
+                && (!Objects.equals(oldOffering.getCpu(), 
targetOffering.getCpu())
+                || !Objects.equals(oldOffering.getRamSize(), 
targetOffering.getRamSize()));
+    }
+
+    static long parseCpuMillis(String value) { return value.endsWith("m") ? 
Long.parseLong(value.substring(0, value.length() - 1)) : Long.parseLong(value) 
* 1000L; }
+    static long parseMemoryBytes(String value) {
+        Matcher matcher = MEMORY_QUANTITY.matcher(value);
+        if (!matcher.matches()) throw new 
IllegalArgumentException("Unsupported Kubernetes memory quantity");
+        long number = Long.parseLong(matcher.group(1)); String unit = 
matcher.group(2);
+        if (unit == null) return number;
+        switch (unit) {
+            case "Ki": return number * 1024L; case "Mi": return number * 1024L 
* 1024L; case "Gi": return number * 1024L * 1024L * 1024L;
+            case "Ti": return number * 1024L * 1024L * 1024L * 1024L; case 
"K": return number * 1000L; case "M": return number * 1000L * 1000L;
+            case "G": return number * 1000L * 1000L * 1000L; case "T": return 
number * 1000L * 1000L * 1000L * 1000L;
+            default: throw new IllegalArgumentException("Unsupported 
Kubernetes memory quantity");
+        }
+    }
+
+    private NodeCapacitySnapshot nodeSnapshot(KubernetesCluster cluster, 
UserVm vm, NodeAccess access, GuestResources guest) throws Exception {
+        JsonObject node = new JsonParser().parse(executeControl(access, "sudo 
/opt/bin/kubectl get node " + nodeName(vm) + " -o json")).getAsJsonObject();
+        JsonObject annotations = 
node.getAsJsonObject("metadata").has("annotations") ? 
node.getAsJsonObject("metadata").getAsJsonObject("annotations") : null;
+        boolean cksCordon = annotations != null && 
annotations.has(RESIZE_ANNOTATION) && 
cluster.getUuid().equals(annotations.get(RESIZE_ANNOTATION).getAsString());
+        JsonObject spec = node.has("spec") ? node.getAsJsonObject("spec") : 
new JsonObject(); JsonObject status = node.getAsJsonObject("status");
+        JsonObject capacity = status.getAsJsonObject("capacity"); JsonObject 
allocatable = status.getAsJsonObject("allocatable");
+        return new NodeCapacitySnapshot(spec.has("unschedulable") && 
spec.get("unschedulable").getAsBoolean(), cksCordon,
+                guest == null ? 0 : guest.cpu, guest == null ? 0 : 
guest.memoryKiB, parseCpuMillis(capacity.get("cpu").getAsString()),
+                parseMemoryBytes(capacity.get("memory").getAsString()), 
parseCpuMillis(allocatable.get("cpu").getAsString()),
+                parseMemoryBytes(allocatable.get("memory").getAsString()), 
isReady(status));
+    }
+
+    private boolean isReady(JsonObject status) {
+        JsonArray conditions = status.getAsJsonArray("conditions");
+        for (JsonElement condition : conditions) { JsonObject item = 
condition.getAsJsonObject(); if ("Ready".equals(item.get("type").getAsString()) 
&& "True".equals(item.get("status").getAsString())) return true; }
+        return false;
+    }
+
+    private GuestResources guestResources(NodeAccess access) throws Exception {
+        String[] values = executeNode(access, "getconf _NPROCESSORS_ONLN; awk 
'/^MemTotal:/ {print $2}' /proc/meminfo").trim().split("\\s+");
+        if (values.length != 2) throw new CloudRuntimeException("Unable to 
read guest CPU and memory");
+        return new GuestResources(Long.parseLong(values[0]), 
Long.parseLong(values[1]));
+    }
+
+    private boolean guestMatchesTarget(GuestResources observed, 
ServiceOffering target) {
+        long targetMemory = target.getRamSize() * 1024L * 1024L; long 
observedMemory = observed.memoryKiB * 1024L;
+        long overhead = Math.max(targetMemory / 50L, 
MINIMUM_MEMORY_OVERHEAD_BYTES);
+        return observed.cpu == target.getCpu() && observedMemory <= 
targetMemory && targetMemory - observedMemory <= overhead;
+    }
+
+    private boolean kubernetesMatchesGuest(NodeCapacitySnapshot observed, 
NodeCapacitySnapshot before) {
+        return observed.isReady() && observed.getCapacityCpuMillis() == 
observed.getGuestOnlineCpuCount() * 1000L
+                && Math.abs(observed.getCapacityMemoryBytes() - 
observed.getGuestMemoryKiB() * 1024L) <= 
KUBERNETES_MEMORY_REPORTING_TOLERANCE_BYTES
+                && (observed.getGuestOnlineCpuCount() <= 
before.getGuestOnlineCpuCount() || observed.getAllocatableCpuMillis() > 
before.getAllocatableCpuMillis())
+                && (observed.getGuestMemoryKiB() <= before.getGuestMemoryKiB() 
|| observed.getAllocatableMemoryBytes() > before.getAllocatableMemoryBytes());
+    }
+
+    private String executeControl(NodeAccess access, String command) throws 
Exception { return execute(access.getControlAddress(), access.getControlPort(), 
access, command); }
+    private String executeNode(NodeAccess access, String command) throws 
Exception { return execute(access.getNodeAddress(), access.getNodePort(), 
access, command); }
+    private String execute(String address, int port, NodeAccess access, String 
command) throws Exception {
+        Pair<Boolean, String> result = SshHelper.sshExecute(address, port, 
access.getUser(), access.getSshKeyFile(), null, command, 10000, 10000, 
COMMAND_TIMEOUT_MS);
+        if (Boolean.TRUE.equals(result.first())) return 
StringUtils.defaultString(result.second());
+        throw new CloudRuntimeException("CKS live-resize command failed");
+    }
+    private String nodeName(UserVm vm) { String name = 
StringUtils.lowerCase(vm.getHostName()); if 
(!NODE_NAME.matcher(name).matches()) throw new CloudRuntimeException("Invalid 
Kubernetes node name for live resize"); return name; }
+    private CloudRuntimeException failure(String phase, UserVm vm, String 
message) { return new CloudRuntimeException("CKS live resize " + phase + " 
failed for VM " + vm.getUuid() + ": " + message); }
+    private void sleep() throws InterruptedException { 
Thread.sleep(POLL_INTERVAL_MS); }
+    private static class GuestResources { private final long cpu; private 
final long memoryKiB; GuestResources(long cpu, long memoryKiB) { this.cpu = 
cpu; this.memoryKiB = memoryKiB; } }
+}
diff --git 
a/plugins/integrations/kubernetes-service/src/main/resources/META-INF/cloudstack/kubernetes-service/spring-kubernetes-service-context.xml
 
b/plugins/integrations/kubernetes-service/src/main/resources/META-INF/cloudstack/kubernetes-service/spring-kubernetes-service-context.xml
index 05336678629..9fe19200d6e 100644
--- 
a/plugins/integrations/kubernetes-service/src/main/resources/META-INF/cloudstack/kubernetes-service/spring-kubernetes-service-context.xml
+++ 
b/plugins/integrations/kubernetes-service/src/main/resources/META-INF/cloudstack/kubernetes-service/spring-kubernetes-service-context.xml
@@ -34,6 +34,7 @@
     <bean id="kubernetesClusterVmMapDaoImpl" 
class="com.cloud.kubernetes.cluster.dao.KubernetesClusterVmMapDaoImpl" />
     <bean id="kubernetesClusterAffinityGroupMapDaoImpl" 
class="com.cloud.kubernetes.cluster.dao.KubernetesClusterAffinityGroupMapDaoImpl"
 />
     <bean id="kubernetesClusterManagerImpl" 
class="com.cloud.kubernetes.cluster.KubernetesClusterManagerImpl" />
+    <bean id="kubernetesClusterNodeCapacityReconciler" 
class="com.cloud.kubernetes.cluster.utils.KubernetesClusterNodeCapacityReconcilerImpl"
 />
 
     <bean id="kubernetesServiceHelper" 
class="com.cloud.kubernetes.cluster.KubernetesServiceHelperImpl" >
         <property name="name" value="KubernetesServiceHelper" />
diff --git 
a/plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorkerTest.java
 
b/plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorkerTest.java
index c9299bdbaa6..ffd9a28288c 100644
--- 
a/plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorkerTest.java
+++ 
b/plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/actionworkers/KubernetesClusterScaleWorkerTest.java
@@ -20,6 +20,8 @@ import com.cloud.kubernetes.cluster.KubernetesCluster;
 import com.cloud.kubernetes.cluster.KubernetesClusterVmMapVO;
 import com.cloud.kubernetes.cluster.KubernetesClusterManagerImpl;
 import com.cloud.kubernetes.cluster.dao.KubernetesClusterVmMapDao;
+import 
com.cloud.kubernetes.cluster.utils.KubernetesClusterNodeCapacityReconciler;
+import com.cloud.hypervisor.Hypervisor;
 import com.cloud.offering.ServiceOffering;
 import com.cloud.service.ServiceOfferingVO;
 import com.cloud.service.dao.ServiceOfferingDao;
@@ -40,6 +42,7 @@ import java.util.List;
 
 import static 
com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType.CONTROL;
 import static 
com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType.DEFAULT;
+import static 
com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType.WORKER;
 
 @RunWith(MockitoJUnitRunner.class)
 public class KubernetesClusterScaleWorkerTest {
@@ -54,6 +57,8 @@ public class KubernetesClusterScaleWorkerTest {
     private KubernetesClusterVmMapDao kubernetesClusterVmMapDao;
     @Mock
     private UserVmDao userVmDao;
+    @Mock
+    private KubernetesClusterNodeCapacityReconciler 
kubernetesClusterNodeCapacityReconciler;
 
     private KubernetesClusterScaleWorker worker;
 
@@ -187,4 +192,32 @@ public class KubernetesClusterScaleWorkerTest {
 
         Assert.assertTrue(toRemove.isEmpty());
     }
+
+    @Test
+    public void testShouldReconcileNodeCapacityOnlyForManagedKvmNodes() {
+        KubernetesCluster runningManagedCluster = 
Mockito.mock(KubernetesCluster.class);
+        
Mockito.when(runningManagedCluster.getState()).thenReturn(KubernetesCluster.State.Running);
+        
Mockito.when(runningManagedCluster.getClusterType()).thenReturn(KubernetesCluster.ClusterType.CloudManaged);
+        KubernetesClusterScaleWorker scaleWorker = new 
KubernetesClusterScaleWorker(runningManagedCluster,
+                new java.util.HashMap<>(), 1L, null, false, null, null, 
clusterManager);
+        scaleWorker.kubernetesClusterNodeCapacityReconciler = 
kubernetesClusterNodeCapacityReconciler;
+
+        KubernetesClusterVmMapVO managedNode = 
Mockito.mock(KubernetesClusterVmMapVO.class);
+        Mockito.when(managedNode.isExternalNode()).thenReturn(false);
+        UserVmVO kvmNode = Mockito.mock(UserVmVO.class);
+        
Mockito.when(kvmNode.getHypervisorType()).thenReturn(Hypervisor.HypervisorType.KVM);
+        ServiceOffering oldOffering = Mockito.mock(ServiceOffering.class);
+        ServiceOffering targetOffering = Mockito.mock(ServiceOffering.class);
+        
Mockito.when(kubernetesClusterNodeCapacityReconciler.requiresKubeletRefresh(oldOffering,
 targetOffering, WORKER,
+                KubernetesCluster.State.Running)).thenReturn(true);
+
+        Assert.assertTrue(scaleWorker.shouldReconcileNodeCapacity(managedNode, 
kvmNode, oldOffering, targetOffering, WORKER));
+
+        Mockito.when(managedNode.isExternalNode()).thenReturn(true);
+        
Assert.assertFalse(scaleWorker.shouldReconcileNodeCapacity(managedNode, 
kvmNode, oldOffering, targetOffering, WORKER));
+
+        Mockito.when(managedNode.isExternalNode()).thenReturn(false);
+        
Mockito.when(kvmNode.getHypervisorType()).thenReturn(Hypervisor.HypervisorType.XenServer);
+        
Assert.assertFalse(scaleWorker.shouldReconcileNodeCapacity(managedNode, 
kvmNode, oldOffering, targetOffering, WORKER));
+    }
 }
diff --git 
a/plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImplTest.java
 
b/plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImplTest.java
new file mode 100644
index 00000000000..729803c7187
--- /dev/null
+++ 
b/plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImplTest.java
@@ -0,0 +1,71 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+package com.cloud.kubernetes.cluster.utils;
+
+import com.cloud.kubernetes.cluster.KubernetesCluster;
+import 
com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType;
+import com.cloud.offering.ServiceOffering;
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+public class KubernetesClusterNodeCapacityReconcilerImplTest {
+    private final KubernetesClusterNodeCapacityReconciler reconciler = new 
KubernetesClusterNodeCapacityReconcilerImpl();
+
+    @Test
+    public void testParseKubernetesQuantities() {
+        Assert.assertEquals(3900L, 
KubernetesClusterNodeCapacityReconcilerImpl.parseCpuMillis("3900m"));
+        Assert.assertEquals(4000L, 
KubernetesClusterNodeCapacityReconcilerImpl.parseCpuMillis("4"));
+        Assert.assertEquals(4037048L * 1024L, 
KubernetesClusterNodeCapacityReconcilerImpl.parseMemoryBytes("4037048Ki"));
+        Assert.assertEquals(4L * 1024L * 1024L * 1024L, 
KubernetesClusterNodeCapacityReconcilerImpl.parseMemoryBytes("4Gi"));
+    }
+
+    @Test
+    public void 
testOnlyRunningKubernetesNodesWithChangedCpuOrRamRequireRefresh() {
+        ServiceOffering oldOffering = offering(2, 2048);
+        ServiceOffering cpuOffering = offering(4, 2048);
+        ServiceOffering ramOffering = offering(2, 4096);
+        ServiceOffering capOnlyOffering = offering(2, 2048);
+
+        Assert.assertTrue(reconciler.requiresKubeletRefresh(oldOffering, 
cpuOffering, KubernetesClusterNodeType.WORKER, 
KubernetesCluster.State.Running));
+        Assert.assertTrue(reconciler.requiresKubeletRefresh(oldOffering, 
ramOffering, KubernetesClusterNodeType.CONTROL, 
KubernetesCluster.State.Running));
+        Assert.assertFalse(reconciler.requiresKubeletRefresh(oldOffering, 
capOnlyOffering, KubernetesClusterNodeType.WORKER, 
KubernetesCluster.State.Running));
+        Assert.assertFalse(reconciler.requiresKubeletRefresh(oldOffering, 
cpuOffering, KubernetesClusterNodeType.ETCD, KubernetesCluster.State.Running));
+        Assert.assertFalse(reconciler.requiresKubeletRefresh(oldOffering, 
cpuOffering, KubernetesClusterNodeType.WORKER, 
KubernetesCluster.State.Stopped));
+    }
+
+    @Test
+    public void 
testCurrentKubernetesResourcesDoNotRequireASecondKubeletRestart() {
+        long memoryKiB = 4L * 1024L * 1024L;
+        KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot current =
+                new 
KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot(false, false, 4, 
memoryKiB,
+                        4000, memoryKiB * 1024L, 3900, memoryKiB * 1024L - 
128L * 1024L * 1024L, true);
+        KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot stale =
+                new 
KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot(false, false, 4, 
memoryKiB,
+                        2000, memoryKiB * 1024L, 1900, memoryKiB * 1024L - 
128L * 1024L * 1024L, true);
+
+        Assert.assertTrue(reconciler.isKubernetesResourcesCurrent(current, 
offering(4, 4096)));
+        Assert.assertFalse(reconciler.isKubernetesResourcesCurrent(stale, 
offering(4, 4096)));
+    }
+
+    private ServiceOffering offering(int cpu, int memory) {
+        ServiceOffering offering = Mockito.mock(ServiceOffering.class);
+        Mockito.when(offering.getCpu()).thenReturn(cpu);
+        Mockito.when(offering.getRamSize()).thenReturn(memory);
+        return offering;
+    }
+}

Reply via email to