diff --git a/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithm.java b/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithm.java index 368487c2b9b0..b7fd4ca1b8d1 100644 --- a/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithm.java +++ b/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithm.java @@ -33,6 +33,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.function.DoubleBinaryOperator; import static org.apache.cloudstack.cluster.ClusterDrsService.ClusterDrsMetric; import static org.apache.cloudstack.cluster.ClusterDrsService.ClusterDrsMetricType; @@ -84,6 +85,22 @@ Ternary getMetrics(Cluster cluster, VirtualMachine vm, S Boolean requiresStorageMotion, Double preImbalance, double[] baseMetricsArray, Map hostIdToIndexMap) throws ConfigurationException; + /** + * Combines the per-resource (cpu and memory) imbalance values into a single cluster + * imbalance for the "both" metric. Balanced takes the worse (higher) of the two so a + * cluster counts as imbalanced while either resource is imbalanced. Condensed overrides + * this to take the lower of the two: a higher imbalance there means more packed, so using + * the max would let the cluster count as packed as soon as a single resource is packed, + * while the other is still spread out. Using the min keeps packing until both are packed. + * + * @param cpuImbalance the cpu imbalance + * @param memoryImbalance the memory imbalance + * @return the combined imbalance for the "both" metric + */ + default double combineBothMetrics(double cpuImbalance, double memoryImbalance) { + return Math.max(cpuImbalance, memoryImbalance); + } + /** * Calculates the cluster imbalance after migrating a VM to a destination host. * @@ -96,9 +113,20 @@ Ternary getMetrics(Cluster cluster, VirtualMachine vm, S * @return the cluster imbalance after migration */ default Double getImbalancePostMigration(VirtualMachine vm, - Host destHost, Long clusterId, long vmMetric, double[] baseMetricsArray, + Host destHost, Long clusterId, ServiceOffering serviceOffering, double[] baseMetricsArray, Map hostIdToIndexMap, Map> hostCpuMap, - Map> hostMemoryMap) { + Map> hostMemoryMap) throws ConfigurationException { + // Metric "both": combine the cpu and memory post-migration imbalance the same way the + // algorithm combines them pre-migration. The baseMetricsArray fast path is single-metric, + // so evaluate both maps directly instead. + if ("both".equals(getClusterDrsMetric(clusterId))) { + long vmCpuMetric = (long) serviceOffering.getCpu() * serviceOffering.getSpeed(); + long vmMemMetric = serviceOffering.getRamSize() * 1024L * 1024L; + return combineBothMetrics( + imbalancePostMigrationForMap(vm, destHost, clusterId, vmCpuMetric, hostCpuMap), + imbalancePostMigrationForMap(vm, destHost, clusterId, vmMemMetric, hostMemoryMap)); + } + long vmMetric = getVmMetric(serviceOffering, clusterId); // Create a copy of the base array and adjust only the two affected hosts double[] adjustedMetrics = new double[baseMetricsArray.length]; System.arraycopy(baseMetricsArray, 0, adjustedMetrics, 0, baseMetricsArray.length); @@ -129,6 +157,23 @@ default Double getImbalancePostMigration(VirtualMachine vm, return calculateImbalance(adjustedMetrics); } + /** + * Cluster imbalance for a single resource map after hypothetically migrating vm to destHost. + */ + private Double imbalancePostMigrationForMap(VirtualMachine vm, Host destHost, Long clusterId, + long vmMetric, Map> metricMap) { + long destHostId = destHost.getId(); + long vmHostId = vm.getHostId(); + List list = new ArrayList<>(); + for (Map.Entry> entry : metricMap.entrySet()) { + Double value = getMetricValuePostMigration(clusterId, entry.getValue(), vmMetric, entry.getKey(), destHostId, vmHostId); + if (value != null) { + list.add(value); + } + } + return getImbalance(list); + } + /** * Calculate imbalance from an array of metric values. * Imbalance is defined as standard deviation divided by mean. @@ -263,6 +308,14 @@ static String getDrsMetricType(long clusterId) { */ static Double getClusterImbalance(Long clusterId, List> cpuList, List> memoryList, Float skipThreshold) throws ConfigurationException { + // Default "both" aggregation is the worse (higher) of the two resources; an algorithm + // that reads imbalance in the opposite direction passes its own combiner below. + return getClusterImbalance(clusterId, cpuList, memoryList, skipThreshold, Math::max); + } + + static Double getClusterImbalance(Long clusterId, List> cpuList, + List> memoryList, Float skipThreshold, + DoubleBinaryOperator bothCombiner) throws ConfigurationException { String metric = getClusterDrsMetric(clusterId); List list; switch (metric) { @@ -272,6 +325,10 @@ static Double getClusterImbalance(Long clusterId, List case "memory": list = getMetricList(clusterId, memoryList, skipThreshold); break; + case "both": + return bothCombiner.applyAsDouble( + getImbalance(getMetricList(clusterId, cpuList, skipThreshold)), + getImbalance(getMetricList(clusterId, memoryList, skipThreshold))); default: throw new ConfigurationException( String.format("Invalid metric: %s for cluster: %d", metric, clusterId)); diff --git a/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsService.java b/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsService.java index ba6a6464fc20..4f851c41dce1 100644 --- a/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsService.java +++ b/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsService.java @@ -45,6 +45,20 @@ public interface ClusterDrsService extends Manager, Configurable, Scheduler { "The interval in minutes after which a periodic background thread will schedule DRS for a cluster.", true, ConfigKey.Scope.Cluster, null, "Interval for Automatic DRS ", null, null, null); + ConfigKey ClusterDrsEventDrivenEnabled = new ConfigKey<>(Boolean.class, "drs.event.driven.enable", + ConfigKey.CATEGORY_ADVANCED, "false", + "In addition to the periodic drs.automatic.interval timer, trigger DRS for a cluster on VM " + + "power-state events (deploy, start, stop, migrate) so it reacts to imbalance in near-real-time " + + "rather than waiting a whole interval. Requires drs.automatic.enable; rate-limited per cluster " + + "by drs.event.driven.interval.", true, + ConfigKey.Scope.Cluster, null, "Enable event-driven DRS", null, null, null); + + ConfigKey ClusterDrsEventDrivenInterval = new ConfigKey<>(Integer.class, "drs.event.driven.interval", + ConfigKey.CATEGORY_ADVANCED, "5", + "Minimum minutes between event-triggered DRS runs for a cluster (debounce). " + + "Only applies when drs.event.driven.enable is true.", true, + ConfigKey.Scope.Cluster, null, "Event-driven DRS min interval", null, null, null); + ConfigKey ClusterDrsMaxMigrations = new ConfigKey<>(Integer.class, "drs.max.migrations", ConfigKey.CATEGORY_ADVANCED, "50", "Maximum number of live migrations in a DRS execution.", @@ -62,9 +76,10 @@ public interface ClusterDrsService extends Manager, Configurable, Scheduler { ConfigKey ClusterDrsMetric = new ConfigKey<>(String.class, "drs.metric", ConfigKey.CATEGORY_ADVANCED, "memory", - "The allocated resource metric used to measure imbalance in a cluster. Possible values are memory, cpu.", + "The allocated resource metric used to measure imbalance in a cluster. Possible values are memory, cpu, both. " + + "'both' balances on the worse of the cpu and memory imbalance so neither resource is left contended.", true, ConfigKey.Scope.Cluster, null, "DRS metric", null, null, null, ConfigKey.Kind.Select, - "memory,cpu"); + "memory,cpu,both"); ConfigKey ClusterDrsMetricType = new ConfigKey<>(String.class, "drs.metric.type", ConfigKey.CATEGORY_ADVANCED, "used", diff --git a/api/src/test/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithmTest.java b/api/src/test/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithmTest.java index 883798109443..95ee4a2f5951 100644 --- a/api/src/test/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithmTest.java +++ b/api/src/test/java/org/apache/cloudstack/cluster/ClusterDrsAlgorithmTest.java @@ -19,15 +19,23 @@ package org.apache.cloudstack.cluster; +import com.cloud.host.Host; +import com.cloud.offering.ServiceOffering; +import com.cloud.org.Cluster; import com.cloud.utils.Ternary; +import com.cloud.utils.component.AdapterBase; +import com.cloud.vm.VirtualMachine; import junit.framework.TestCase; +import org.apache.cloudstack.framework.config.ConfigKey; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.MockedStatic; import org.mockito.Mockito; import org.mockito.junit.MockitoJUnitRunner; +import java.lang.reflect.Field; import java.util.List; +import java.util.Map; import static org.apache.cloudstack.cluster.ClusterDrsAlgorithm.getMetricValue; import static org.mockito.ArgumentMatchers.any; @@ -94,4 +102,92 @@ public void testGetMetricValueWithSkipThreshold() { } } } + + @Test + public void testGetClusterImbalanceAppliesBothCombiner() throws Exception { + Field defaultValueField = ConfigKey.class.getDeclaredField("_defaultValue"); + defaultValueField.setAccessible(true); + Object originalMetric = defaultValueField.get(ClusterDrsService.ClusterDrsMetric); + try { + List> cpuList = List.of( + new Ternary<>(80L, 0L, 100L), new Ternary<>(20L, 0L, 100L)); + List> memoryList = List.of( + new Ternary<>(50L, 0L, 100L), new Ternary<>(50L, 0L, 100L)); + + defaultValueField.set(ClusterDrsService.ClusterDrsMetric, "cpu"); + double cpu = ClusterDrsAlgorithm.getClusterImbalance(1L, cpuList, memoryList, null); + defaultValueField.set(ClusterDrsService.ClusterDrsMetric, "memory"); + double memory = ClusterDrsAlgorithm.getClusterImbalance(1L, cpuList, memoryList, null); + defaultValueField.set(ClusterDrsService.ClusterDrsMetric, "both"); + + // the default aggregation (and balanced) takes the worse, higher of the per-resource imbalances + double defaultBoth = ClusterDrsAlgorithm.getClusterImbalance(1L, cpuList, memoryList, null); + assertEquals(Math.max(cpu, memory), defaultBoth, 0.0001); + + // the combiner overload controls how the two resources are reduced: condensed passes min, balanced max + double maxBoth = ClusterDrsAlgorithm.getClusterImbalance(1L, cpuList, memoryList, null, Math::max); + double minBoth = ClusterDrsAlgorithm.getClusterImbalance(1L, cpuList, memoryList, null, Math::min); + assertEquals(Math.max(cpu, memory), maxBoth, 0.0001); + assertEquals(Math.min(cpu, memory), minBoth, 0.0001); + } finally { + defaultValueField.set(ClusterDrsService.ClusterDrsMetric, originalMetric); + } + } + + @Test + public void testGetImbalancePostMigrationRoutesBothThroughCombineBothMetrics() throws Exception { + Field defaultValueField = ConfigKey.class.getDeclaredField("_defaultValue"); + defaultValueField.setAccessible(true); + Object originalMetric = defaultValueField.get(ClusterDrsService.ClusterDrsMetric); + try { + defaultValueField.set(ClusterDrsService.ClusterDrsMetric, "both"); + + VirtualMachine vm = Mockito.mock(VirtualMachine.class); + Mockito.when(vm.getHostId()).thenReturn(1L); + Host destHost = Mockito.mock(Host.class); + Mockito.when(destHost.getId()).thenReturn(2L); + ServiceOffering serviceOffering = Mockito.mock(ServiceOffering.class); + Mockito.when(serviceOffering.getCpu()).thenReturn(2); + Mockito.when(serviceOffering.getSpeed()).thenReturn(1000); + Mockito.when(serviceOffering.getRamSize()).thenReturn(2048); + + Map> hostCpuMap = Map.of( + 1L, new Ternary<>(80L, 0L, 100L), 2L, new Ternary<>(20L, 0L, 100L)); + Map> hostMemoryMap = Map.of( + 1L, new Ternary<>(50L, 0L, 100L), 2L, new Ternary<>(50L, 0L, 100L)); + + // the "both" branch must defer to combineBothMetrics; a sentinel override proves it is the reducer used + ClusterDrsAlgorithm algorithm = new TestAlgorithm() { + @Override + public double combineBothMetrics(double cpuImbalance, double memoryImbalance) { + return SENTINEL; + } + }; + // the "both" branch evaluates the cpu and memory maps directly, so the base array and index map are unused + Double imbalance = algorithm.getImbalancePostMigration(vm, destHost, 1L, serviceOffering, null, null, + hostCpuMap, hostMemoryMap); + + assertEquals(SENTINEL, imbalance, 0.0); + } finally { + defaultValueField.set(ClusterDrsService.ClusterDrsMetric, originalMetric); + } + } + + private static final double SENTINEL = 0.4242; + + private static class TestAlgorithm extends AdapterBase implements ClusterDrsAlgorithm { + @Override + public boolean needsDrs(Cluster cluster, List> cpuList, + List> memoryList) { + return false; + } + + @Override + public Ternary getMetrics(Cluster cluster, VirtualMachine vm, ServiceOffering serviceOffering, + Host destHost, Map> hostCpuMap, + Map> hostMemoryMap, Boolean requiresStorageMotion, Double preImbalance, + double[] baseMetricsArray, Map hostIdToIndexMap) { + return null; + } + } } diff --git a/plugins/drs/cluster/balanced/src/main/java/org/apache/cloudstack/cluster/Balanced.java b/plugins/drs/cluster/balanced/src/main/java/org/apache/cloudstack/cluster/Balanced.java index 902ab0900bd7..8770c81a3479 100644 --- a/plugins/drs/cluster/balanced/src/main/java/org/apache/cloudstack/cluster/Balanced.java +++ b/plugins/drs/cluster/balanced/src/main/java/org/apache/cloudstack/cluster/Balanced.java @@ -82,7 +82,7 @@ public Ternary getMetrics(Cluster cluster, VirtualMachin // Use optimized post-imbalance calculation that adjusts only affected hosts Double postImbalance = getImbalancePostMigration(vm, destHost, - cluster.getId(), ClusterDrsAlgorithm.getVmMetric(serviceOffering, cluster.getId()), + cluster.getId(), serviceOffering, baseMetricsArray, hostIdToIndexMap, hostCpuMap, hostMemoryMap); logger.trace("Cluster {} pre-imbalance: {} post-imbalance: {} Algorithm: {} VM: {} srcHost ID: {} destHost: {}", diff --git a/plugins/drs/cluster/balanced/src/test/java/org/apache/cloudstack/cluster/BalancedTest.java b/plugins/drs/cluster/balanced/src/test/java/org/apache/cloudstack/cluster/BalancedTest.java index e955a10f2a56..5ca87fc7935e 100644 --- a/plugins/drs/cluster/balanced/src/test/java/org/apache/cloudstack/cluster/BalancedTest.java +++ b/plugins/drs/cluster/balanced/src/test/java/org/apache/cloudstack/cluster/BalancedTest.java @@ -110,6 +110,13 @@ public void setUp() throws NoSuchFieldException, IllegalAccessException { hostMemoryFreeMap.put(2L, new Ternary<>(2048L * 1024L * 1024L, 0L, 8192L * 1024L * 1024L)); } + @Test + public void combineBothMetricsTakesTheHigher() { + // balanced counts a cluster as imbalanced while either resource is imbalanced, so "both" takes the higher + assertEquals(0.7, balanced.combineBothMetrics(0.2, 0.7), 0.0); + assertEquals(0.9, balanced.combineBothMetrics(0.9, 0.3), 0.0); + } + private void overrideDefaultConfigValue(final ConfigKey configKey, final String name, final Object o) throws IllegalAccessException, NoSuchFieldException { Field f = ConfigKey.class.getDeclaredField(name); diff --git a/plugins/drs/cluster/condensed/src/main/java/org/apache/cloudstack/cluster/Condensed.java b/plugins/drs/cluster/condensed/src/main/java/org/apache/cloudstack/cluster/Condensed.java index d672ddfda615..fc55395723d9 100644 --- a/plugins/drs/cluster/condensed/src/main/java/org/apache/cloudstack/cluster/Condensed.java +++ b/plugins/drs/cluster/condensed/src/main/java/org/apache/cloudstack/cluster/Condensed.java @@ -46,7 +46,7 @@ public boolean needsDrs(Cluster cluster, List> cpuList long clusterId = cluster.getId(); double threshold = getThreshold(clusterId); Float skipThreshold = ClusterDrsImbalanceSkipThreshold.valueIn(clusterId); - Double imbalance = ClusterDrsAlgorithm.getClusterImbalance(clusterId, cpuList, memoryList, skipThreshold); + Double imbalance = ClusterDrsAlgorithm.getClusterImbalance(clusterId, cpuList, memoryList, skipThreshold, this::combineBothMetrics); String drsMetric = ClusterDrsAlgorithm.getClusterDrsMetric(clusterId); String metricType = ClusterDrsAlgorithm.getDrsMetricType(clusterId); Boolean useRatio = ClusterDrsAlgorithm.getDrsMetricUseRatio(clusterId); @@ -66,6 +66,13 @@ private double getThreshold(long clusterId) { return ClusterDrsImbalanceThreshold.valueIn(clusterId); } + @Override + public double combineBothMetrics(double cpuImbalance, double memoryImbalance) { + // Condensed reads a higher imbalance as more packed, so the cluster should only count as + // packed once both resources are packed. Take the lower of the two, not the higher. + return Math.min(cpuImbalance, memoryImbalance); + } + @Override public String getName() { return "condensed"; @@ -80,12 +87,12 @@ public Ternary getMetrics(Cluster cluster, VirtualMachin // Use provided pre-imbalance if available, otherwise calculate it if (preImbalance == null) { preImbalance = ClusterDrsAlgorithm.getClusterImbalance(cluster.getId(), new ArrayList<>(hostCpuMap.values()), - new ArrayList<>(hostMemoryMap.values()), null); + new ArrayList<>(hostMemoryMap.values()), null, this::combineBothMetrics); } // Use optimized post-imbalance calculation that adjusts only affected hosts Double postImbalance = getImbalancePostMigration(vm, destHost, - cluster.getId(), ClusterDrsAlgorithm.getVmMetric(serviceOffering, cluster.getId()), + cluster.getId(), serviceOffering, baseMetricsArray, hostIdToIndexMap, hostCpuMap, hostMemoryMap); logger.trace("Cluster {} pre-imbalance: {} post-imbalance: {} Algorithm: {} VM: {} srcHost ID: {} destHost: {}", diff --git a/plugins/drs/cluster/condensed/src/test/java/org/apache/cloudstack/cluster/CondensedTest.java b/plugins/drs/cluster/condensed/src/test/java/org/apache/cloudstack/cluster/CondensedTest.java index b2a5e6bf84f1..dff6d0ce525f 100644 --- a/plugins/drs/cluster/condensed/src/test/java/org/apache/cloudstack/cluster/CondensedTest.java +++ b/plugins/drs/cluster/condensed/src/test/java/org/apache/cloudstack/cluster/CondensedTest.java @@ -111,6 +111,13 @@ public void setUp() throws NoSuchFieldException, IllegalAccessException { hostMemoryFreeMap.put(2L, new Ternary<>(2048L * 1024L * 1024L, 0L, 8192L * 1024L * 1024L)); } + @Test + public void combineBothMetricsTakesTheLower() { + // condensed reads a higher imbalance as more packed, so "both" must keep the lower of cpu and memory + assertEquals(0.2, condensed.combineBothMetrics(0.2, 0.7), 0.0); + assertEquals(0.3, condensed.combineBothMetrics(0.9, 0.3), 0.0); + } + private void overrideDefaultConfigValue(final ConfigKey configKey, final String name, final Object o) throws IllegalAccessException, NoSuchFieldException { diff --git a/server/src/main/java/org/apache/cloudstack/cluster/ClusterDrsServiceImpl.java b/server/src/main/java/org/apache/cloudstack/cluster/ClusterDrsServiceImpl.java index 62075aae596e..9f6111332a39 100644 --- a/server/src/main/java/org/apache/cloudstack/cluster/ClusterDrsServiceImpl.java +++ b/server/src/main/java/org/apache/cloudstack/cluster/ClusterDrsServiceImpl.java @@ -53,6 +53,14 @@ import com.cloud.utils.exception.CloudRuntimeException; import com.cloud.vm.VMInstanceDetailVO; import com.cloud.vm.VMInstanceVO; +import com.cloud.vm.VirtualMachineManager; +import com.cloud.utils.concurrency.NamedThreadFactory; +import org.apache.cloudstack.framework.messagebus.MessageBus; +import org.apache.cloudstack.framework.messagebus.MessageDispatcher; +import org.apache.cloudstack.framework.messagebus.MessageHandler; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import com.cloud.vm.VirtualMachine; import com.cloud.vm.VirtualMachineProfile; import com.cloud.vm.VirtualMachineProfileImpl; @@ -148,6 +156,13 @@ public class ClusterDrsServiceImpl extends ManagerBase implements ClusterDrsServ Map drsAlgorithmMap = new HashMap<>(); + @Inject + MessageBus messageBus; + // Epoch-ms of the last event-triggered DRS run, per cluster; drives the cooldown. + private final Map lastEventDrsTriggerByCluster = new ConcurrentHashMap<>(); + // Runs event-triggered DRS off the message-bus thread. + private ExecutorService eventDrsExecutor; + public AsyncJobDispatcher getAsyncJobDispatcher() { return asyncJobDispatcher; } @@ -179,6 +194,11 @@ protected void runInContext() { }; Timer vmSchedulerTimer = new Timer("VMSchedulerPollTask"); vmSchedulerTimer.schedule(schedulerPollTask, 5000L, 60 * 1000L); + + // Subscribe to VM power-state events for event-driven DRS (gated per cluster by drs.event.driven.enable). + eventDrsExecutor = Executors.newSingleThreadExecutor(new NamedThreadFactory("Event-Driven-DRS")); + messageBus.subscribe(VirtualMachineManager.Topics.VM_POWER_STATE, MessageDispatcher.getDispatcher(this)); + return true; } @@ -278,48 +298,127 @@ void generateDrsPlanForAllClusters() { List clusterList = clusterDao.listAll(); for (ClusterVO cluster : clusterList) { - if (cluster.getAllocationState() == Disabled || ClusterDrsEnabled.valueIn( - cluster.getId()).equals(Boolean.FALSE)) { - continue; - } + generateDrsPlanForCluster(cluster, ClusterDrsInterval.valueIn(cluster.getId())); + } + } - ClusterDrsPlanVO lastPlan = drsPlanDao.listLatestPlanForClusterId(cluster.getId()); - - // If the last plan is ready or in progress or was executed within the last interval, skip this cluster. - // This is to avoid generating plans for clusters which are already being processed and to avoid - // generating plans for clusters which have been processed recently.This doesn't consider the type - // (manual or automated) of the last plan. - if (lastPlan != null && (lastPlan.getStatus() == ClusterDrsPlan.Status.READY || - lastPlan.getStatus() == ClusterDrsPlan.Status.IN_PROGRESS || - (lastPlan.getStatus() == ClusterDrsPlan.Status.COMPLETED && - lastPlan.getCreated().compareTo(DateUtils.addMinutes(new Date(), -1 * ClusterDrsInterval.valueIn(cluster.getId()))) > 0) - )) { - continue; + /** + * Generates a DRS plan for a single cluster, skipping if DRS is disabled, a plan is already + * pending/running, or one completed within {@code debounceMinutes}. + */ + void generateDrsPlanForCluster(ClusterVO cluster, int debounceMinutes) { + if (cluster.getAllocationState() == Disabled || ClusterDrsEnabled.valueIn(cluster.getId()).equals(Boolean.FALSE)) { + return; + } + + ClusterDrsPlanVO lastPlan = drsPlanDao.listLatestPlanForClusterId(cluster.getId()); + + // Skip if the last plan is ready, in progress, or completed within the debounce window. + if (lastPlan != null && (lastPlan.getStatus() == ClusterDrsPlan.Status.READY || + lastPlan.getStatus() == ClusterDrsPlan.Status.IN_PROGRESS || + (lastPlan.getStatus() == ClusterDrsPlan.Status.COMPLETED && + lastPlan.getCreated().compareTo(DateUtils.addMinutes(new Date(), -1 * debounceMinutes)) > 0) + )) { + return; + } + + long eventId = ActionEventUtils.onStartedActionEvent(User.UID_SYSTEM, Account.ACCOUNT_ID_SYSTEM, + EventTypes.EVENT_CLUSTER_DRS, + String.format("Generating DRS plan for cluster %s", cluster.getUuid()), cluster.getId(), + ApiCommandResourceType.Cluster.toString(), true, 0); + GlobalLock clusterLock = GlobalLock.getInternLock(String.format(CLUSTER_LOCK_STR, cluster.getId())); + try { + if (clusterLock.lock(30)) { + try { + List> plan = getDrsPlan(cluster, + ClusterDrsMaxMigrations.valueIn(cluster.getId())); + savePlan(cluster.getId(), plan, eventId, ClusterDrsPlan.Type.AUTOMATED, + ClusterDrsPlan.Status.READY); + logger.info("Generated DRS plan for cluster {}", cluster); + } catch (Exception e) { + logger.error("Unable to generate DRS plans for cluster {}", cluster, e); + } finally { + clusterLock.unlock(); + } } + } finally { + clusterLock.releaseRef(); + } + } - long eventId = ActionEventUtils.onStartedActionEvent(User.UID_SYSTEM, Account.ACCOUNT_ID_SYSTEM, - EventTypes.EVENT_CLUSTER_DRS, - String.format("Generating DRS plan for cluster %s", cluster.getUuid()), cluster.getId(), - ApiCommandResourceType.Cluster.toString(), true, 0); - GlobalLock clusterLock = GlobalLock.getInternLock(String.format(CLUSTER_LOCK_STR, cluster.getId())); + /** + * Message-bus handler for VM power-state events; triggers event-driven DRS for the VM's cluster. + */ + @MessageHandler(topic = VirtualMachineManager.Topics.VM_POWER_STATE) + protected void handleVmPowerStateEvent(String subject, String senderAddress, Object args) { + if (!(args instanceof Long)) { + return; + } + try { + triggerEventDrivenDrsForVm((Long) args); + } catch (Exception e) { + logger.debug("Event-driven DRS: error handling VM power-state event for {}", args, e); + } + } + + /** + * Resolves the VM's cluster and, subject to the enable flag and cooldown, schedules DRS plan generation. + */ + void triggerEventDrivenDrsForVm(Long vmId) { + if (vmId == null) { + return; + } + VMInstanceVO vm = vmInstanceDao.findById(vmId); + if (vm == null || vm.getHostId() == null) { + return; + } + HostVO host = hostDao.findById(vm.getHostId()); + if (host == null || host.getClusterId() == null) { + return; + } + Long clusterId = host.getClusterId(); + if (!shouldTriggerEventDrivenDrs(clusterId)) { + return; + } + final ClusterVO cluster = clusterDao.findById(clusterId); + if (cluster == null) { + return; + } + final int debounceMinutes = ClusterDrsEventDrivenInterval.valueIn(clusterId); + logger.debug("Event-driven DRS: scheduling plan generation for cluster {} (triggered by VM {})", clusterId, vmId); + submitEventDrivenDrs(cluster, debounceMinutes); + } + + /** + * Runs generateDrsPlanForCluster off the message-bus thread so event publishers are not blocked. + */ + protected void submitEventDrivenDrs(final ClusterVO cluster, final int debounceMinutes) { + eventDrsExecutor.submit(() -> { try { - if (clusterLock.lock(30)) { - try { - List> plan = getDrsPlan(cluster, - ClusterDrsMaxMigrations.valueIn(cluster.getId())); - savePlan(cluster.getId(), plan, eventId, ClusterDrsPlan.Type.AUTOMATED, - ClusterDrsPlan.Status.READY); - logger.info("Generated DRS plan for cluster {}", cluster); - } catch (Exception e) { - logger.error("Unable to generate DRS plans for cluster {}", cluster, e); - } finally { - clusterLock.unlock(); - } - } - } finally { - clusterLock.releaseRef(); + generateDrsPlanForCluster(cluster, debounceMinutes); + } catch (Exception e) { + logger.warn("Event-driven DRS: plan generation failed for cluster {}", cluster, e); } + }); + } + + /** + * Returns true if automatic and event-driven DRS are enabled and the per-cluster cooldown has + * elapsed, recording the trigger time when it does. + */ + boolean shouldTriggerEventDrivenDrs(Long clusterId) { + if (ClusterDrsEnabled.valueIn(clusterId).equals(Boolean.FALSE) + || ClusterDrsEventDrivenEnabled.valueIn(clusterId).equals(Boolean.FALSE)) { + return false; } + long now = System.currentTimeMillis(); + long cooldownMs = ClusterDrsEventDrivenInterval.valueIn(clusterId) * 60L * 1000L; + Long last = lastEventDrsTriggerByCluster.get(clusterId); + if (last != null && (now - last) < cooldownMs) { + return false; + } + lastEventDrsTriggerByCluster.put(clusterId, now); + return true; } /** @@ -593,11 +692,12 @@ Pair getBestMigration(Cluster cluster, ClusterDrsAlgorithm Map> vmToCompatibleHostsCache, Map> vmToStorageMotionCache, Map vmToExcludesMap) throws ConfigurationException { - // Pre-calculate cluster imbalance once per iteration (same for all VM-host combinations) + // Pre-calculate cluster imbalance once per iteration (same for all VM-host combinations). + // Use the algorithm's own "both" aggregation so the pre-imbalance matches its post-imbalance. Double preImbalance = getClusterImbalance(cluster.getId(), new ArrayList<>(hostCpuCapacityMap.values()), new ArrayList<>(hostMemoryCapacityMap.values()), - null); + null, algorithm::combineBothMetrics); // Pre-calculate base metrics array once per iteration for optimized imbalance calculation String metricType = getClusterDrsMetric(cluster.getId()); @@ -855,7 +955,7 @@ public String getConfigComponentName() { public ConfigKey[] getConfigKeys() { return new ConfigKey[]{ClusterDrsPlanExpireInterval, ClusterDrsEnabled, ClusterDrsInterval, ClusterDrsMaxMigrations, ClusterDrsAlgorithm, ClusterDrsImbalanceThreshold, ClusterDrsMetric, ClusterDrsMetricType, ClusterDrsMetricUseRatio, - ClusterDrsImbalanceSkipThreshold}; + ClusterDrsImbalanceSkipThreshold, ClusterDrsEventDrivenEnabled, ClusterDrsEventDrivenInterval}; } @Override diff --git a/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsServiceImplTest.java b/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsServiceImplTest.java index 6390b29097b5..41beb2e50f94 100644 --- a/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsServiceImplTest.java +++ b/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsServiceImplTest.java @@ -74,10 +74,12 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.ExecutorService; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; @RunWith(MockitoJUnitRunner.class) public class ClusterDrsServiceImplTest { @@ -951,4 +953,380 @@ public void testProcessPlans() { Mockito.verify(clusterDrsService, Mockito.times(2)).executeDrsPlan(Mockito.any(ClusterDrsPlanVO.class)); } + + // ---- event-driven DRS ---- + // The ConfigKeys are shared interface constants, so each test that overrides a default restores it + // in a finally block to avoid leaking into other tests. + + private static String getConfigDefault(ConfigKey key) throws Exception { + Field f = ConfigKey.class.getDeclaredField("_defaultValue"); + f.setAccessible(true); + return (String) f.get(key); + } + + private static void setConfigDefault(ConfigKey key, String value) throws Exception { + Field f = ConfigKey.class.getDeclaredField("_defaultValue"); + f.setAccessible(true); + f.set(key, value); + } + + @Test + public void testShouldTriggerEventDrivenDrsDisabledByDefault() throws Exception { + // Automatic DRS enabled but event-driven off -> must not trigger. + String origDrs = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + String origEvt = getConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, "false"); + assertFalse(clusterDrsService.shouldTriggerEventDrivenDrs(1L)); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, origDrs); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, origEvt); + } + } + + @Test + public void testShouldTriggerEventDrivenDrsRequiresAutomaticDrs() throws Exception { + // Event-driven on but automatic DRS off -> must not trigger (event-driven depends on drs.automatic.enable). + String origDrs = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + String origEvt = getConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "false"); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, "true"); + assertFalse(clusterDrsService.shouldTriggerEventDrivenDrs(1L)); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, origDrs); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, origEvt); + } + } + + @Test + public void testShouldTriggerEventDrivenDrsEnabledThenDebounced() throws Exception { + String origDrs = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + String origEvt = getConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled); + String origInt = getConfigDefault(clusterDrsService.ClusterDrsEventDrivenInterval); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, "true"); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenInterval, "5"); + // First event fires; the per-cluster cooldown then suppresses an immediate second event. + assertTrue(clusterDrsService.shouldTriggerEventDrivenDrs(1L)); + assertFalse(clusterDrsService.shouldTriggerEventDrivenDrs(1L)); + // A different cluster has an independent cooldown and still fires. + assertTrue(clusterDrsService.shouldTriggerEventDrivenDrs(2L)); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, origDrs); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, origEvt); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenInterval, origInt); + } + } + + @Test + public void testTriggerEventDrivenDrsForVmSchedulesWhenEnabled() throws Exception { + String origDrs = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + String origEvt = getConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, "true"); + + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getHostId()).thenReturn(10L); + HostVO host = Mockito.mock(HostVO.class); + Mockito.when(host.getClusterId()).thenReturn(1L); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(vmInstanceDao.findById(100L)).thenReturn(vm); + Mockito.when(hostDao.findById(10L)).thenReturn(host); + Mockito.when(clusterDao.findById(1L)).thenReturn(cluster); + // Don't actually run DRS on a background thread in the test. + Mockito.doNothing().when(clusterDrsService).submitEventDrivenDrs(Mockito.any(ClusterVO.class), Mockito.anyInt()); + + clusterDrsService.triggerEventDrivenDrsForVm(100L); + + Mockito.verify(clusterDrsService, Mockito.times(1)).submitEventDrivenDrs(Mockito.eq(cluster), Mockito.anyInt()); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, origDrs); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, origEvt); + } + } + + @Test + public void testTriggerEventDrivenDrsForVmDoesNotScheduleWhenDisabled() { + // Defaults: both flags false -> must not schedule, even though the VM resolves to a cluster. + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getHostId()).thenReturn(10L); + HostVO host = Mockito.mock(HostVO.class); + Mockito.when(host.getClusterId()).thenReturn(1L); + Mockito.when(vmInstanceDao.findById(100L)).thenReturn(vm); + Mockito.when(hostDao.findById(10L)).thenReturn(host); + + clusterDrsService.triggerEventDrivenDrsForVm(100L); + + Mockito.verify(clusterDrsService, Mockito.never()).submitEventDrivenDrs(Mockito.any(ClusterVO.class), Mockito.anyInt()); + } + + @Test + public void testHandleVmPowerStateEventIgnoresNonLongArg() { + clusterDrsService.handleVmPowerStateEvent("subject", "sender", "not-a-long"); + Mockito.verify(clusterDrsService, Mockito.never()).triggerEventDrivenDrsForVm(Mockito.anyLong()); + } + + @Test + public void testHandleVmPowerStateEventForwardsTheVmId() { + Mockito.doNothing().when(clusterDrsService).triggerEventDrivenDrsForVm(100L); + clusterDrsService.handleVmPowerStateEvent("subject", "sender", 100L); + Mockito.verify(clusterDrsService, Mockito.times(1)).triggerEventDrivenDrsForVm(100L); + } + + @Test + public void testHandleVmPowerStateEventSwallowsErrors() { + Mockito.doThrow(new RuntimeException("boom")).when(clusterDrsService).triggerEventDrivenDrsForVm(100L); + // an error must not propagate out of the message-bus handler + clusterDrsService.handleVmPowerStateEvent("subject", "sender", 100L); + Mockito.verify(clusterDrsService, Mockito.times(1)).triggerEventDrivenDrsForVm(100L); + } + + @Test + public void testTriggerEventDrivenDrsForVmNullVmId() { + clusterDrsService.triggerEventDrivenDrsForVm(null); + Mockito.verify(clusterDrsService, Mockito.never()).submitEventDrivenDrs(Mockito.any(ClusterVO.class), Mockito.anyInt()); + } + + @Test + public void testTriggerEventDrivenDrsForVmUnknownVm() { + Mockito.when(vmInstanceDao.findById(100L)).thenReturn(null); + clusterDrsService.triggerEventDrivenDrsForVm(100L); + Mockito.verify(clusterDrsService, Mockito.never()).submitEventDrivenDrs(Mockito.any(ClusterVO.class), Mockito.anyInt()); + } + + @Test + public void testTriggerEventDrivenDrsForVmVmHasNoHost() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getHostId()).thenReturn(null); + Mockito.when(vmInstanceDao.findById(100L)).thenReturn(vm); + clusterDrsService.triggerEventDrivenDrsForVm(100L); + Mockito.verify(clusterDrsService, Mockito.never()).submitEventDrivenDrs(Mockito.any(ClusterVO.class), Mockito.anyInt()); + } + + @Test + public void testTriggerEventDrivenDrsForVmUnknownHost() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getHostId()).thenReturn(10L); + Mockito.when(vmInstanceDao.findById(100L)).thenReturn(vm); + Mockito.when(hostDao.findById(10L)).thenReturn(null); + clusterDrsService.triggerEventDrivenDrsForVm(100L); + Mockito.verify(clusterDrsService, Mockito.never()).submitEventDrivenDrs(Mockito.any(ClusterVO.class), Mockito.anyInt()); + } + + @Test + public void testTriggerEventDrivenDrsForVmClusterNotFound() throws Exception { + String origDrs = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + String origEvt = getConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, "true"); + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getHostId()).thenReturn(10L); + HostVO host = Mockito.mock(HostVO.class); + Mockito.when(host.getClusterId()).thenReturn(1L); + Mockito.when(vmInstanceDao.findById(100L)).thenReturn(vm); + Mockito.when(hostDao.findById(10L)).thenReturn(host); + Mockito.when(clusterDao.findById(1L)).thenReturn(null); + + clusterDrsService.triggerEventDrivenDrsForVm(100L); + + Mockito.verify(clusterDrsService, Mockito.never()).submitEventDrivenDrs(Mockito.any(ClusterVO.class), Mockito.anyInt()); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, origDrs); + setConfigDefault(clusterDrsService.ClusterDrsEventDrivenEnabled, origEvt); + } + } + + @Test + public void testSubmitEventDrivenDrsRunsPlanGeneration() { + ExecutorService direct = Mockito.mock(ExecutorService.class); + Mockito.when(direct.submit(Mockito.any(Runnable.class))).thenAnswer(invocation -> { + ((Runnable) invocation.getArgument(0)).run(); + return null; + }); + ReflectionTestUtils.setField(clusterDrsService, "eventDrsExecutor", direct); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.doNothing().when(clusterDrsService).generateDrsPlanForCluster(cluster, 5); + + clusterDrsService.submitEventDrivenDrs(cluster, 5); + + Mockito.verify(clusterDrsService, Mockito.times(1)).generateDrsPlanForCluster(cluster, 5); + } + + @Test + public void testSubmitEventDrivenDrsSwallowsPlanErrors() { + ExecutorService direct = Mockito.mock(ExecutorService.class); + Mockito.when(direct.submit(Mockito.any(Runnable.class))).thenAnswer(invocation -> { + ((Runnable) invocation.getArgument(0)).run(); + return null; + }); + ReflectionTestUtils.setField(clusterDrsService, "eventDrsExecutor", direct); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.doThrow(new RuntimeException("plan failed")).when(clusterDrsService).generateDrsPlanForCluster(cluster, 5); + + // the background task must swallow the failure, not propagate it + clusterDrsService.submitEventDrivenDrs(cluster, 5); + + Mockito.verify(clusterDrsService, Mockito.times(1)).generateDrsPlanForCluster(cluster, 5); + } + + @Test + public void testGenerateDrsPlanForClusterSkipsWhenClusterDisabled() { + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(cluster.getAllocationState()).thenReturn(Grouping.AllocationState.Disabled); + + clusterDrsService.generateDrsPlanForCluster(cluster, 5); + + Mockito.verify(drsPlanDao, Mockito.never()).listLatestPlanForClusterId(Mockito.anyLong()); + } + + @Test + public void testGenerateDrsPlanForClusterSkipsWhenPlanIsReady() throws Exception { + String orig = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(cluster.getId()).thenReturn(1L); + Mockito.when(cluster.getAllocationState()).thenReturn(Grouping.AllocationState.Enabled); + ClusterDrsPlanVO lastPlan = Mockito.mock(ClusterDrsPlanVO.class); + Mockito.when(lastPlan.getStatus()).thenReturn(ClusterDrsPlan.Status.READY); + Mockito.when(drsPlanDao.listLatestPlanForClusterId(1L)).thenReturn(lastPlan); + + clusterDrsService.generateDrsPlanForCluster(cluster, 5); + + Mockito.verify(clusterDrsService, Mockito.never()).getDrsPlan(Mockito.any(), Mockito.anyInt()); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, orig); + } + } + + @Test + public void testGenerateDrsPlanForClusterGeneratesAndSavesPlan() throws Exception { + String orig = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(cluster.getId()).thenReturn(2L); + Mockito.when(cluster.getUuid()).thenReturn("cluster-uuid"); + Mockito.when(cluster.getAllocationState()).thenReturn(Grouping.AllocationState.Enabled); + Mockito.when(drsPlanDao.listLatestPlanForClusterId(2L)).thenReturn(null); + + GlobalLock lock = Mockito.mock(GlobalLock.class); + Mockito.when(lock.lock(30)).thenReturn(true); + Mockito.when(GlobalLock.getInternLock("drs.plan.cluster.2")).thenReturn(lock); + + Mockito.doReturn(Collections.emptyList()).when(clusterDrsService).getDrsPlan(Mockito.eq(cluster), Mockito.anyInt()); + Mockito.doReturn(null).when(clusterDrsService).savePlan(Mockito.anyLong(), Mockito.anyList(), Mockito.anyLong(), + Mockito.any(), Mockito.any()); + + try (MockedStatic actionEvents = Mockito.mockStatic(ActionEventUtils.class)) { + actionEvents.when(() -> ActionEventUtils.onStartedActionEvent(Mockito.anyLong(), Mockito.anyLong(), + Mockito.anyString(), Mockito.anyString(), Mockito.anyLong(), Mockito.anyString(), + Mockito.anyBoolean(), Mockito.anyLong())).thenReturn(1L); + + clusterDrsService.generateDrsPlanForCluster(cluster, 5); + } + + Mockito.verify(clusterDrsService, Mockito.times(1)).getDrsPlan(Mockito.eq(cluster), Mockito.anyInt()); + Mockito.verify(clusterDrsService, Mockito.times(1)).savePlan(Mockito.eq(2L), Mockito.anyList(), Mockito.anyLong(), + Mockito.any(), Mockito.any()); + Mockito.verify(lock, Mockito.times(1)).unlock(); + Mockito.verify(lock, Mockito.times(1)).releaseRef(); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, orig); + } + } + + @Test + public void testGenerateDrsPlanForClusterSkipsWhenPlanInProgress() throws Exception { + String orig = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(cluster.getId()).thenReturn(1L); + Mockito.when(cluster.getAllocationState()).thenReturn(Grouping.AllocationState.Enabled); + ClusterDrsPlanVO lastPlan = Mockito.mock(ClusterDrsPlanVO.class); + Mockito.when(lastPlan.getStatus()).thenReturn(ClusterDrsPlan.Status.IN_PROGRESS); + Mockito.when(drsPlanDao.listLatestPlanForClusterId(1L)).thenReturn(lastPlan); + + clusterDrsService.generateDrsPlanForCluster(cluster, 5); + + Mockito.verify(clusterDrsService, Mockito.never()).getDrsPlan(Mockito.any(), Mockito.anyInt()); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, orig); + } + } + + @Test + public void testGenerateDrsPlanForClusterSkipsWhenRecentlyCompleted() throws Exception { + String orig = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(cluster.getId()).thenReturn(1L); + Mockito.when(cluster.getAllocationState()).thenReturn(Grouping.AllocationState.Enabled); + ClusterDrsPlanVO lastPlan = Mockito.mock(ClusterDrsPlanVO.class); + Mockito.when(lastPlan.getStatus()).thenReturn(ClusterDrsPlan.Status.COMPLETED); + // created just now, so it is inside the debounce window and must be skipped + Mockito.when(lastPlan.getCreated()).thenReturn(new Date()); + Mockito.when(drsPlanDao.listLatestPlanForClusterId(1L)).thenReturn(lastPlan); + + clusterDrsService.generateDrsPlanForCluster(cluster, 5); + + Mockito.verify(clusterDrsService, Mockito.never()).getDrsPlan(Mockito.any(), Mockito.anyInt()); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, orig); + } + } + + @Test + public void testGenerateDrsPlanForClusterHandlesPlanGenerationFailure() throws Exception { + String orig = getConfigDefault(clusterDrsService.ClusterDrsEnabled); + try { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, "true"); + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(cluster.getId()).thenReturn(2L); + Mockito.when(cluster.getUuid()).thenReturn("cluster-uuid"); + Mockito.when(cluster.getAllocationState()).thenReturn(Grouping.AllocationState.Enabled); + Mockito.when(drsPlanDao.listLatestPlanForClusterId(2L)).thenReturn(null); + + GlobalLock lock = Mockito.mock(GlobalLock.class); + Mockito.when(lock.lock(30)).thenReturn(true); + Mockito.when(GlobalLock.getInternLock("drs.plan.cluster.2")).thenReturn(lock); + + Mockito.doThrow(new RuntimeException("planning failed")).when(clusterDrsService).getDrsPlan(Mockito.eq(cluster), Mockito.anyInt()); + + try (MockedStatic actionEvents = Mockito.mockStatic(ActionEventUtils.class)) { + actionEvents.when(() -> ActionEventUtils.onStartedActionEvent(Mockito.anyLong(), Mockito.anyLong(), + Mockito.anyString(), Mockito.anyString(), Mockito.anyLong(), Mockito.anyString(), + Mockito.anyBoolean(), Mockito.anyLong())).thenReturn(1L); + + // a planning failure must be caught, and the lock still released + clusterDrsService.generateDrsPlanForCluster(cluster, 5); + } + + Mockito.verify(clusterDrsService, Mockito.never()).savePlan(Mockito.anyLong(), Mockito.anyList(), Mockito.anyLong(), + Mockito.any(), Mockito.any()); + Mockito.verify(lock, Mockito.times(1)).unlock(); + Mockito.verify(lock, Mockito.times(1)).releaseRef(); + } finally { + setConfigDefault(clusterDrsService.ClusterDrsEnabled, orig); + } + } + + @Test + public void testGenerateDrsPlanForAllClustersIteratesEachCluster() { + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.when(cluster.getId()).thenReturn(1L); + Mockito.when(clusterDao.listAll()).thenReturn(List.of(cluster)); + Mockito.doNothing().when(clusterDrsService).generateDrsPlanForCluster(Mockito.eq(cluster), Mockito.anyInt()); + + clusterDrsService.generateDrsPlanForAllClusters(); + + Mockito.verify(clusterDrsService, Mockito.times(1)).generateDrsPlanForCluster(Mockito.eq(cluster), Mockito.anyInt()); + } }