From 6f6656f491717a9df02fee3e0f4951cb4821945c Mon Sep 17 00:00:00 2001 From: andrijapanicsb Date: Tue, 22 Sep 2026 12:47:24 +0200 Subject: [PATCH 1/6] CKS: reconcile node capacity after live resize --- .../KubernetesClusterScaleWorker.java | 99 ++++++++- ...bernetesClusterNodeCapacityReconciler.java | 67 ++++++ ...etesClusterNodeCapacityReconcilerImpl.java | 190 ++++++++++++++++++ .../spring-kubernetes-service-context.xml | 1 + .../KubernetesClusterScaleWorkerTest.java | 33 +++ ...ClusterNodeCapacityReconcilerImplTest.java | 71 +++++++ 6 files changed, 455 insertions(+), 6 deletions(-) create mode 100644 plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconciler.java create mode 100644 plugins/integrations/kubernetes-service/src/main/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImpl.java create mode 100644 plugins/integrations/kubernetes-service/src/test/java/com/cloud/kubernetes/cluster/utils/KubernetesClusterNodeCapacityReconcilerImplTest.java 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 08513dbd4487..a1dbf28b2bde 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.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.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 @@ private void scaleKubernetesClusterOffering(KubernetesClusterNodeType nodeType, 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()); - } 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()); + 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 @@ private void scaleKubernetesClusterOffering(KubernetesClusterNodeType nodeType, kubernetesCluster = updateKubernetesClusterEntryForNodeType(null, nodeType, serviceOffering, updateNodeOffering, updateClusterOffering); } + private NodeAccess resolveNodeAccess(UserVm userVM) { + Pair 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 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 000000000000..f52458774aca --- /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 000000000000..a6785ec3fdfd --- /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 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 053366786292..9fe19200d6e7 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 @@ + 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 c9299bdbaa67..ffd9a28288cc 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.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 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 void testGetWorkerNodesToRemoveForDownsize_noRemoval() { 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 000000000000..729803c7187e --- /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; + } +} From a4e4e3752b3d05a1ea14b4cc7a64d17f1c376b8d Mon Sep 17 00:00:00 2001 From: andrijapanicsb Date: Tue, 22 Sep 2026 13:32:09 +0200 Subject: [PATCH 2/6] CKS: cover node capacity reconciliation ownership --- ...etesClusterNodeCapacityReconcilerImpl.java | 4 +- ...ClusterNodeCapacityReconcilerImplTest.java | 82 +++++++++++++++++++ 2 files changed, 84 insertions(+), 2 deletions(-) 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 index a6785ec3fdfd..68576ce12c9e 100644 --- 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 @@ -176,8 +176,8 @@ private boolean kubernetesMatchesGuest(NodeCapacitySnapshot observed, NodeCapaci && (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); } + protected String executeControl(NodeAccess access, String command) throws Exception { return execute(access.getControlAddress(), access.getControlPort(), access, command); } + protected 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 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()); 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 index 729803c7187e..3707c8014bff 100644 --- 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 @@ -16,9 +16,13 @@ // under the License. package com.cloud.kubernetes.cluster.utils; +import java.util.ArrayList; +import java.util.List; + import com.cloud.kubernetes.cluster.KubernetesCluster; import com.cloud.kubernetes.cluster.KubernetesServiceHelper.KubernetesClusterNodeType; import com.cloud.offering.ServiceOffering; +import com.cloud.uservm.UserVm; import org.junit.Assert; import org.junit.Test; import org.mockito.Mockito; @@ -62,10 +66,88 @@ public void testCurrentKubernetesResourcesDoNotRequireASecondKubeletRestart() { Assert.assertFalse(reconciler.isKubernetesResourcesCurrent(stale, offering(4, 4096))); } + @Test + public void testCaptureBeforeReadsStructuredNodeState() throws Exception { + FakeReconciler fakeReconciler = new FakeReconciler(nodeJson(false)); + KubernetesCluster cluster = cluster("cluster-uuid"); + UserVm vm = vm("worker-1"); + + KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot snapshot = fakeReconciler.captureBefore(cluster, vm, access()); + + Assert.assertFalse(snapshot.isUnschedulable()); + Assert.assertTrue(snapshot.isReady()); + Assert.assertEquals(2L, snapshot.getGuestOnlineCpuCount()); + Assert.assertEquals(2097152L, snapshot.getGuestMemoryKiB()); + Assert.assertEquals(3900L, snapshot.getCapacityCpuMillis()); + Assert.assertEquals(4L * 1024L * 1024L * 1024L, snapshot.getCapacityMemoryBytes()); + } + + @Test + public void testCordonAndRestoreOnlyNodesOwnedByCks() throws Exception { + KubernetesCluster cluster = cluster("cluster-uuid"); + UserVm vm = vm("worker-1"); + KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot previouslyCordon = + new KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot(true, false, 0, 0, 0, 0, 0, 0, true); + FakeReconciler operatorCordon = new FakeReconciler(nodeJson(true)); + + operatorCordon.cordonIfNeeded(cluster, vm, previouslyCordon, access(), System.currentTimeMillis() + 1000L); + operatorCordon.restoreSchedulability(cluster, vm, previouslyCordon, access(), System.currentTimeMillis() + 1000L); + Assert.assertTrue(operatorCordon.controlCommands.isEmpty()); + + KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot cksCordon = + new KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot(true, true, 0, 0, 0, 0, 0, 0, true); + FakeReconciler cleanup = new FakeReconciler(nodeJson(false)); + cleanup.restoreSchedulability(cluster, vm, cksCordon, access(), System.currentTimeMillis() + 1000L); + + Assert.assertEquals("sudo /opt/bin/kubectl uncordon worker-1", cleanup.controlCommands.get(0)); + Assert.assertEquals("sudo /opt/bin/kubectl annotate node worker-1 cloudstack.apache.org/cks-live-resize-", cleanup.controlCommands.get(1)); + } + 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; } + + private KubernetesCluster cluster(String uuid) { + KubernetesCluster cluster = Mockito.mock(KubernetesCluster.class); + Mockito.when(cluster.getUuid()).thenReturn(uuid); + return cluster; + } + + private UserVm vm(String hostname) { + UserVm vm = Mockito.mock(UserVm.class); + Mockito.when(vm.getHostName()).thenReturn(hostname); + Mockito.when(vm.getUuid()).thenReturn("vm-uuid"); + return vm; + } + + private KubernetesClusterNodeCapacityReconciler.NodeAccess access() { + return new KubernetesClusterNodeCapacityReconciler.NodeAccess("control", 22, "node", 22, "root", null); + } + + private String nodeJson(boolean unschedulable) { + return String.format("{\"metadata\":{\"annotations\":{}},\"spec\":{\"unschedulable\":%s},\"status\":{\"capacity\":{\"cpu\":\"3900m\",\"memory\":\"4Gi\"},\"allocatable\":{\"cpu\":\"3700m\",\"memory\":\"3900Mi\"},\"conditions\":[{\"type\":\"Ready\",\"status\":\"True\"}]}}", unschedulable); + } + + private static class FakeReconciler extends KubernetesClusterNodeCapacityReconcilerImpl { + private final String nodeJson; + private final List controlCommands = new ArrayList<>(); + + FakeReconciler(String nodeJson) { + this.nodeJson = nodeJson; + } + + @Override + protected String executeControl(KubernetesClusterNodeCapacityReconciler.NodeAccess access, String command) { + controlCommands.add(command); + return command.contains(" get node ") ? nodeJson : ""; + } + + @Override + protected String executeNode(KubernetesClusterNodeCapacityReconciler.NodeAccess access, String command) { + return "2\n2097152\n"; + } + } } From 6c34333920878a62694b4bac53e4b7d248f59a91 Mon Sep 17 00:00:00 2001 From: andrijapanicsb Date: Tue, 22 Sep 2026 13:37:16 +0200 Subject: [PATCH 3/6] CKS: test live resize reconciliation guards --- .../KubernetesClusterScaleWorkerTest.java | 3 +++ ...etesClusterNodeCapacityReconcilerImplTest.java | 15 +++++++++++++++ 2 files changed, 18 insertions(+) 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 ffd9a28288cc..2c16ea755731 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 @@ -219,5 +219,8 @@ public void testShouldReconcileNodeCapacityOnlyForManagedKvmNodes() { Mockito.when(managedNode.isExternalNode()).thenReturn(false); Mockito.when(kvmNode.getHypervisorType()).thenReturn(Hypervisor.HypervisorType.XenServer); Assert.assertFalse(scaleWorker.shouldReconcileNodeCapacity(managedNode, kvmNode, oldOffering, targetOffering, WORKER)); + + Mockito.when(runningManagedCluster.getClusterType()).thenReturn(KubernetesCluster.ClusterType.ExternalManaged); + 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 index 3707c8014bff..eee5b154888b 100644 --- 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 @@ -103,6 +103,21 @@ public void testCordonAndRestoreOnlyNodesOwnedByCks() throws Exception { Assert.assertEquals("sudo /opt/bin/kubectl annotate node worker-1 cloudstack.apache.org/cks-live-resize-", cleanup.controlCommands.get(1)); } + @Test + public void testCordonAddsOwnershipAnnotationBeforeCordoning() throws Exception { + KubernetesCluster cluster = cluster("cluster-uuid"); + UserVm vm = vm("worker-1"); + KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot schedulable = + new KubernetesClusterNodeCapacityReconciler.NodeCapacitySnapshot(false, false, 0, 0, 0, 0, 0, 0, true); + FakeReconciler fakeReconciler = new FakeReconciler(nodeJson(true)); + + fakeReconciler.cordonIfNeeded(cluster, vm, schedulable, access(), System.currentTimeMillis() + 1000L); + + Assert.assertEquals("sudo /opt/bin/kubectl annotate node worker-1 cloudstack.apache.org/cks-live-resize=cluster-uuid --overwrite", + fakeReconciler.controlCommands.get(0)); + Assert.assertEquals("sudo /opt/bin/kubectl cordon worker-1", fakeReconciler.controlCommands.get(1)); + } + private ServiceOffering offering(int cpu, int memory) { ServiceOffering offering = Mockito.mock(ServiceOffering.class); Mockito.when(offering.getCpu()).thenReturn(cpu); From 859f394e4d954ba9db1455118c83b5ed12a4d3ff Mon Sep 17 00:00:00 2001 From: andrijapanicsb Date: Tue, 22 Sep 2026 13:39:38 +0200 Subject: [PATCH 4/6] CKS: improve live resize failure diagnostics --- .../actionworkers/KubernetesClusterScaleWorker.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) 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 a1dbf28b2bde..0dbb53e19a39 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 @@ -421,7 +421,14 @@ private void scaleKubernetesClusterOffering(KubernetesClusterNodeType nodeType, 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); + String recovery = resizeSucceeded + ? "The node remains cordoned; check kubelet on the VM and retry the CKS scale operation." + : "The VM was not resized; CKS attempted to restore its temporary cordon."; + String message = String.format("Scaling Kubernetes cluster %s (UUID %s) failed for %s node VM %s (UUID %s), " + + "offering %s to %s: %s %s", kubernetesCluster.getName(), kubernetesCluster.getUuid(), nodeType, + userVM.getDisplayName(), userVM.getUuid(), oldOffering == null ? null : oldOffering.getId(), + serviceOffering.getId(), e.getMessage(), recovery); + logTransitStateAndThrow(Level.ERROR, message, kubernetesCluster.getId(), KubernetesCluster.Event.OperationFailed, e); } 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); From d512edb7e00888f53e36de12dea31df5c75ecabe Mon Sep 17 00:00:00 2001 From: andrijapanicsb Date: Tue, 22 Sep 2026 13:42:01 +0200 Subject: [PATCH 5/6] CKS: resolve live resize target public IP --- .../cluster/actionworkers/KubernetesClusterScaleWorker.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) 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 0dbb53e19a39..dd13be645889 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 @@ -456,7 +456,11 @@ private NodeAccess resolveNodeAccess(UserVm userVM) { 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(), + IpAddress targetPublicIp = network.getVpcId() == null ? getNetworkSourceNatIp(network) : getVpcTierKubernetesPublicIp(network); + if (targetPublicIp == null) { + throw new CloudRuntimeException(String.format("Unable to resolve target-node public IP for VM %s", userVM.getUuid())); + } + return new NodeAccess(controlAccess.first(), controlAccess.second(), targetPublicIp.getAddress().addr(), sshRule.getSourcePortStart(), getControlNodeLoginUser(), getManagementServerSshPublicKeyFile()); } From 40b50b012e374a764d7b87c84cb8ab71e0c27b47 Mon Sep 17 00:00:00 2001 From: andrijapanicsb Date: Tue, 22 Sep 2026 17:08:33 +0200 Subject: [PATCH 6/6] CKS: register live resize capacity reconciler --- .../KubernetesClusterNodeCapacityReconcilerImpl.java | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) 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 index 68576ce12c9e..2c8d9efc4caf 100644 --- 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 @@ -32,8 +32,15 @@ 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. */ +import org.springframework.stereotype.Component; + +/** + * SSH-backed implementation for a rolling CKS live-resize reconciliation. + * + * Component scanning makes this implementation available to the scale worker, which obtains it + * through {@code ComponentContext}. + */ +@Component 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;