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..12ab333ab6ab 100644 --- a/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsService.java +++ b/api/src/main/java/org/apache/cloudstack/cluster/ClusterDrsService.java @@ -88,6 +88,46 @@ public interface ClusterDrsService extends Manager, Configurable, Scheduler { true, ConfigKey.Scope.Cluster, null, "DRS imbalance skip threshold for Condensed algorithm", null, null, null); + ConfigKey ClusterDrsPowerManagementEnabled = new ConfigKey<>(Boolean.class, "drs.power.management.enable", + ConfigKey.CATEGORY_ADVANCED, "false", + "Enable distributed power management on the cluster. When the cluster is under-utilized, a host that has " + + "been emptied by DRS is powered off through its out-of-band management; when the cluster is over-utilized, " + + "a powered-off host is powered back on. Requires automatic DRS and out-of-band management to be configured.", + true, ConfigKey.Scope.Cluster, null, "Enable DRS power management", null, null, null); + + ConfigKey ClusterDrsPowerManagementLowThreshold = new ConfigKey<>(Float.class, "drs.power.management.low.threshold", + ConfigKey.CATEGORY_ADVANCED, "0.3", + "Cluster utilization (0.0 to 1.0, on the configured DRS metric) below which an idle host becomes a candidate " + + "to power off, provided the remaining hosts can carry the load without exceeding the high threshold.", + true, ConfigKey.Scope.Cluster, null, "DRS power management low threshold", null, null, null); + + ConfigKey ClusterDrsPowerManagementHighThreshold = new ConfigKey<>(Float.class, "drs.power.management.high.threshold", + ConfigKey.CATEGORY_ADVANCED, "0.75", + "Cluster utilization (0.0 to 1.0, on the configured DRS metric) above which a powered-off host is powered " + + "back on. Also the ceiling the remaining hosts must stay under before a host is powered off.", + true, ConfigKey.Scope.Cluster, null, "DRS power management high threshold", null, null, null); + + ConfigKey ClusterDrsPowerManagementEvacuate = new ConfigKey<>(Boolean.class, "drs.power.management.evacuate.enable", + ConfigKey.CATEGORY_ADVANCED, "false", + "When DRS power management decides a host can be released but the host still runs VMs, migrate those VMs " + + "onto the remaining hosts (only when every VM can be placed) so the host can then be powered off. " + + "Off by default: power management otherwise waits for DRS to empty a host on its own.", + true, ConfigKey.Scope.Cluster, null, "Actively evacuate before power-off", null, null, null); + + ConfigKey ClusterDrsPredictiveEnabled = new ConfigKey<>(Boolean.class, "drs.predictive.enable", + ConfigKey.CATEGORY_ADVANCED, "false", + "Use the recent trend of cluster utilization, not only the instantaneous value, when deciding power " + + "management actions. A rising trend powers a host back on earlier and holds off powering one off; a falling " + + "trend is ignored for power-off so a brief dip does not churn hosts. Requires DRS power management. The trend " + + "window is kept in memory by the management server that runs the poll, so after a restart or in a multi-management-server " + + "deployment it rebuilds from the samples seen since that server last took the poll.", + true, ConfigKey.Scope.Cluster, null, "Enable predictive DRS", null, null, null); + + ConfigKey ClusterDrsPredictiveWindow = new ConfigKey<>(Integer.class, "drs.predictive.window", + ConfigKey.CATEGORY_ADVANCED, "5", + "Number of recent DRS poll samples of cluster utilization used to compute the trend for predictive DRS.", + true, ConfigKey.Scope.Cluster, null, "Predictive DRS window", null, null, null); + /** * Generate a DRS plan for a cluster and save it as per the parameters 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..1535a5a5e7b6 100644 --- a/server/src/main/java/org/apache/cloudstack/cluster/ClusterDrsServiceImpl.java +++ b/server/src/main/java/org/apache/cloudstack/cluster/ClusterDrsServiceImpl.java @@ -19,6 +19,7 @@ package org.apache.cloudstack.cluster; +import com.cloud.agent.AgentManager; import com.cloud.api.ApiGsonHelper; import com.cloud.api.query.dao.HostJoinDao; import com.cloud.api.query.vo.HostJoinVO; @@ -32,9 +33,13 @@ import com.cloud.event.EventVO; import com.cloud.event.dao.EventDao; import com.cloud.exception.InvalidParameterValueException; +import com.cloud.host.DetailVO; import com.cloud.host.Host; import com.cloud.host.HostVO; +import com.cloud.host.Status; import com.cloud.host.dao.HostDao; +import com.cloud.host.dao.HostDetailsDao; +import com.cloud.resource.ResourceState; import com.cloud.offering.ServiceOffering; import com.cloud.org.Cluster; import com.cloud.server.ManagementServer; @@ -54,6 +59,8 @@ import com.cloud.vm.VMInstanceDetailVO; import com.cloud.vm.VMInstanceVO; import com.cloud.vm.VirtualMachine; +import org.apache.cloudstack.outofbandmanagement.OutOfBandManagement; +import org.apache.cloudstack.outofbandmanagement.OutOfBandManagementService; import com.cloud.vm.VirtualMachineProfile; import com.cloud.vm.VirtualMachineProfileImpl; import com.cloud.vm.VmDetailConstants; @@ -89,8 +96,12 @@ import java.util.Collections; import java.util.Date; import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.LinkedList; +import java.util.concurrent.ConcurrentHashMap; import java.util.HashSet; import java.util.List; +import java.util.TreeMap; import java.util.Map; import java.util.Set; import java.util.Timer; @@ -117,12 +128,37 @@ public class ClusterDrsServiceImpl extends ManagerBase implements ClusterDrsServ @Inject HostDao hostDao; + @Inject + HostDetailsDao hostDetailsDao; + + @Inject + OutOfBandManagementService outOfBandManagementService; + @Inject EventDao eventDao; + // Host detail marking a host that DRS power management powered off, so only those are powered back on + // (never a host that is down for another reason). + protected static final String DRS_POWER_STATE_DETAIL = "drs.power.state"; + protected static final String DRS_POWER_STATE_OFF = "off"; + // Marks a host DRS is actively draining so it can be powered off once empty; distinct from OFF so a draining + // host is never treated as a wake candidate and is picked up across polls to continue or abandon the drain. + protected static final String DRS_POWER_STATE_DRAINING = "draining"; + + // Recent cluster utilization samples (ratio on the DRS metric) per cluster, newest last, used by predictive DRS. + // Written only from the single-threaded poll under the clusterDRS.poll lock. + protected final Map> clusterUtilizationHistory = new ConcurrentHashMap<>(); + + // VM count on each host at the last drain batch, so a drain that stops making progress (migrations rejected for + // affinity/tags/storage) is abandoned instead of re-submitted forever. Written only under the poll lock. + protected final Map drainingVmCountByHost = new ConcurrentHashMap<>(); + @Inject HostJoinDao hostJoinDao; + @Inject + AgentManager agentManager; + @Inject VMInstanceDao vmInstanceDao; @@ -199,6 +235,7 @@ public void poll(Date timestamp) { processPlans(); generateDrsPlanForAllClusters(); processPlans(); + managePowerForAllClusters(); } finally { lock.unlock(); } @@ -851,11 +888,542 @@ public String getConfigComponentName() { return ClusterDrsService.class.getSimpleName(); } + protected void managePowerForAllClusters() { + for (ClusterVO cluster : clusterDao.listAll()) { + try { + managePowerForCluster(cluster); + } catch (Exception e) { + logger.warn("DRS power management skipped cluster [{}] due to [{}].", cluster.getId(), e.getMessage(), e); + } + } + } + + protected void managePowerForCluster(ClusterVO cluster) { + if (cluster.getAllocationState() == Disabled + || Boolean.FALSE.equals(ClusterDrsEnabled.valueIn(cluster.getId())) + || Boolean.FALSE.equals(ClusterDrsPowerManagementEnabled.valueIn(cluster.getId()))) { + return; + } + // Do not power-manage a cluster with a DRS migration plan still in flight: a host that is the source or + // destination of a pending or running migration must not be disabled or powered off underneath it. + if (hasInFlightDrsPlan(cluster.getId())) { + logger.debug("DRS power management: skipping cluster [{}] while a DRS migration plan is in flight.", cluster.getId()); + return; + } + + final float lowThreshold = ClusterDrsPowerManagementLowThreshold.valueIn(cluster.getId()); + final float highThreshold = ClusterDrsPowerManagementHighThreshold.valueIn(cluster.getId()); + final boolean useCpu = "cpu".equals(getClusterDrsMetric(cluster.getId())); + + List routingHosts = hostDao.findByClusterId(cluster.getId(), Host.Type.Routing); + List upHosts = new ArrayList<>(); + List poweredOffByDrs = new ArrayList<>(); + List drainingHosts = new ArrayList<>(); + for (HostVO host : routingHosts) { + if (isDrainingByDrs(host)) { + if (host.getStatus() == Status.Up) { + // still reachable: continue (or finish) the drain below. + drainingHosts.add(host); + } else { + // went down or away mid-drain: DRS did not power it off, so stop tracking it and release the + // marker rather than leaking it and blocking a later drain. + abandonDrain(host, "host is no longer up"); + } + continue; + } + if (host.getStatus() == Status.Up) { + if (isPoweredOffByDrs(host)) { + // Our host is back in service: either a wake completed, or a power-off we issued never took + // effect (command accepted but the host stayed up). Re-enable it if we had disabled it, drop + // the marker, and return it to the capacity pool. Keying only on Up avoids leaving such a host + // stranded Disabled and marked, in neither list, forever. + if (host.getResourceState() == ResourceState.Disabled) { + hostDao.updateResourceState(ResourceState.Disabled, ResourceState.Event.Enable, ResourceState.Enabled, host); + } + clearStalePowerMarker(host); + upHosts.add(host); + } else if (host.getResourceState() == ResourceState.Enabled) { + upHosts.add(host); + } + // Up but Disabled by someone other than DRS: leave it alone, it is not ours to schedule onto. + } else if (isPoweredOffByDrs(host)) { + // Down and marked by DRS: a wake candidate. The marker is kept across the wake (power-on is + // idempotent) and cleared only once the host is actually Up again, above. + poweredOffByDrs.add(host); + } + } + if (upHosts.isEmpty() && poweredOffByDrs.isEmpty() && drainingHosts.isEmpty()) { + return; + } + + // Capacity reflects every host still carrying load: the schedulable up hosts plus any host being drained + // (its VMs still run until it empties), so utilization is not understated while a drain is in progress. + List loadHosts = new ArrayList<>(upHosts); + loadHosts.addAll(drainingHosts); + Map> capacityMap = getHostCapacityMap(loadHosts, useCpu); + double clusterUsed = 0d; + double clusterTotal = 0d; + for (Ternary capacity : capacityMap.values()) { + // count used + reserved as the committed load, so a host holding HA/allocation reservations is not + // powered off just because its live usage is low. + clusterUsed += capacity.first() + capacity.second(); + clusterTotal += capacity.third(); + } + + // predictive DRS: fold the recent trend in. Use the more conservative of the instantaneous and forecast + // utilization, so a rising trend wakes a host sooner and holds off power-off, while a transient dip never + // triggers an aggressive power-off. + double effectiveUsed = clusterUsed; + double instantaneousRatio = clusterTotal > 0d ? clusterUsed / clusterTotal : 0d; + if (Boolean.TRUE.equals(ClusterDrsPredictiveEnabled.valueIn(cluster.getId()))) { + int window = ClusterDrsPredictiveWindow.valueIn(cluster.getId()); + double forecastRatio = forecastUtilization(recordAndGetUtilizationHistory(cluster.getId(), instantaneousRatio, window)); + effectiveUsed = Math.max(instantaneousRatio, forecastRatio) * clusterTotal; + } + + // prefer restoring capacity over saving power: if the cluster is hot and we have a host we powered + // off earlier that we can still reach over out-of-band management, bring it back before any power-off. + if (clusterNeedsWakeup(effectiveUsed, clusterTotal, highThreshold, poweredOffByDrs.size())) { + for (HostVO poweredOff : poweredOffByDrs) { + if (isPowerManageable(poweredOff)) { + powerOnHost(poweredOff, cluster); + return; + } + } + } + + // power off at most one empty, out-of-band-manageable host per poll, and only if the rest can carry the load. + for (HostVO host : upHosts) { + Ternary capacity = capacityMap.get(host.getId()); + if (capacity == null) { + continue; + } + if (hostHasNoRunningVms(host) && isPowerManageable(host) && !isPoweredOffByDrs(host) + && clusterCanReleaseHost(effectiveUsed, clusterTotal, capacity.third(), upHosts.size(), lowThreshold, highThreshold)) { + powerOffHost(host, cluster); + return; + } + } + + // Continue a drain already in progress even if evacuation was since turned off, so a host is never stranded + // mid-drain; only START a new drain when the operator has opted in and no host is empty to power off. + if (!drainingHosts.isEmpty()) { + drainHostBatch(cluster, drainingHosts.get(0), upHosts, capacityMap, effectiveUsed, clusterTotal, lowThreshold, highThreshold, useCpu); + } else if (Boolean.TRUE.equals(ClusterDrsPowerManagementEvacuate.valueIn(cluster.getId()))) { + HostVO candidate = selectDrainCandidate(upHosts, capacityMap, effectiveUsed, clusterTotal, lowThreshold, highThreshold); + if (candidate != null) { + drainHostBatch(cluster, candidate, upHosts, capacityMap, effectiveUsed, clusterTotal, lowThreshold, highThreshold, useCpu); + } + } + } + + protected boolean isDrainingByDrs(HostVO host) { + DetailVO detail = hostDetailsDao.findDetail(host.getId(), DRS_POWER_STATE_DETAIL); + return detail != null && DRS_POWER_STATE_DRAINING.equals(detail.getValue()); + } + + protected boolean hasMigratingVm(long hostId) { + List vms = vmInstanceDao.listByHostId(hostId); + if (vms == null) { + return false; + } + for (VMInstanceVO vm : vms) { + if (vm.getState() == VirtualMachine.State.Migrating) { + return true; + } + } + return false; + } + + protected void setPowerMarker(long hostId, String value) { + DetailVO existing = hostDetailsDao.findDetail(hostId, DRS_POWER_STATE_DETAIL); + if (existing != null) { + hostDetailsDao.remove(existing.getId()); + } + hostDetailsDao.persist(new DetailVO(hostId, DRS_POWER_STATE_DETAIL, value)); + } + + protected void abandonDrain(HostVO host, String reason) { + logger.info("DRS power management: abandoning drain of host [{}]: {}", host.getId(), reason); + // the host was disabled when the drain started so the scheduler would leave it alone while it emptied; + // put it back in service now that the drain is being given up. + if (host.getResourceState() == ResourceState.Disabled) { + hostDao.updateResourceState(ResourceState.Disabled, ResourceState.Event.Enable, ResourceState.Enabled, host); + } + clearStalePowerMarker(host); + drainingVmCountByHost.remove(host.getId()); + } + + /** + * The least-loaded host the cluster can give up: powerable, holding only running user VMs, and such that the + * remaining hosts can still carry the load. Null if no host qualifies. + */ + protected HostVO selectDrainCandidate(List upHosts, Map> capacityMap, + double effectiveUsed, double clusterTotal, float lowThreshold, float highThreshold) { + HostVO candidate = null; + double candidateUsed = Double.MAX_VALUE; + for (HostVO host : upHosts) { + Ternary cap = capacityMap.get(host.getId()); + if (cap == null || hostHasNoRunningVms(host) || !isPowerManageable(host) || !hostEvacuatable(host) + || !clusterCanReleaseHost(effectiveUsed, clusterTotal, cap.third(), upHosts.size(), lowThreshold, highThreshold)) { + continue; + } + if (cap.first() < candidateUsed) { + candidate = host; + candidateUsed = cap.first(); + } + } + return candidate; + } + + /** + * Submits one batch of migrations off {@code candidate}. Waits while a previous batch is still in flight; if the + * host is quiescent and still holds as many VMs as at the previous batch the drain made no progress and is + * abandoned; otherwise the host is marked draining and up to ClusterDrsMaxMigrations of its remaining VMs are + * placed by capacity and migrated away. The host empties over successive polls and is powered off once empty. + */ + protected void drainHostBatch(ClusterVO cluster, HostVO candidate, List upHosts, + Map> capacityMap, double effectiveUsed, double clusterTotal, + float lowThreshold, float highThreshold, boolean useCpu) { + List vms = vmInstanceDao.listByHostId(candidate.getId()); + if (vms == null || vms.isEmpty()) { + // Fully drained. It was disabled when the drain started, so the empty-host power-off loop (which only + // scans enabled up hosts) does not see it; power it off here, applying the same guards that loop does. + // If it can no longer be powered off (out-of-band management gone, or the cluster is no longer + // under-utilized because load rose during the drain), put it back in service instead of flapping it. + Ternary cap = capacityMap.get(candidate.getId()); + if (isPowerManageable(candidate) && cap != null + && clusterCanReleaseHost(effectiveUsed, clusterTotal, cap.third(), upHosts.size() + 1, lowThreshold, highThreshold)) { + powerOffHost(candidate, cluster); + } else { + abandonDrain(candidate, "host is drained but can no longer be powered off; returning it to service"); + } + return; + } + if (hasMigratingVm(candidate.getId())) { + // a batch is still in flight; wait for it before judging progress or submitting more. + return; + } + int current = vms.size(); + Integer previous = drainingVmCountByHost.get(candidate.getId()); + if (previous != null && current >= previous) { + abandonDrain(candidate, "no progress since the last batch; remaining VMs cannot be migrated off"); + return; + } + + Map vmNeeds = new HashMap<>(); + for (VMInstanceVO vm : vms) { + vmNeeds.put(vm.getId(), vmResourceNeed(vm, useCpu)); + } + Map hostFree = new HashMap<>(); + for (HostVO host : upHosts) { + if (host.getId() == candidate.getId()) { + continue; + } + Ternary cap = capacityMap.get(host.getId()); + if (cap != null) { + // Placeable room = (high threshold of total) - (used + reserved), floored at zero, so draining + // never pushes a destination past the high threshold it is meant to respect. + double placeable = (double) highThreshold * cap.third() - (cap.first() + cap.second()); + hostFree.put(host.getId(), Math.max(0d, placeable)); + } + } + + Map plan = planHostEvacuation(vmNeeds, hostFree); + if (plan.isEmpty()) { + if (previous != null) { + abandonDrain(candidate, "remaining VMs no longer fit on the other hosts"); + } else { + logger.debug("DRS power management: host [{}] in cluster [{}] cannot be evacuated; not all VMs fit on the remaining hosts.", + candidate.getId(), cluster.getId()); + } + return; + } + + // Resolve the concrete VM/destination pairs to migrate this poll (bounded by the DRS migration budget) + // before starting any event, so the start event is only opened when there is real work to do. + int maxMigrations = ClusterDrsMaxMigrations.valueIn(cluster.getId()); + Map batch = new LinkedHashMap<>(); + for (Map.Entry entry : plan.entrySet()) { + if (batch.size() >= maxMigrations) { + break; + } + VirtualMachine vm = vmInstanceDao.findById(entry.getKey()); + HostVO destination = hostDao.findById(entry.getValue()); + if (vm != null && destination != null) { + batch.put(vm, destination); + } + } + if (batch.isEmpty()) { + return; + } + // Starting a new drain: record the durable draining marker BEFORE disabling the host, so a crash or an + // exception in the window between the two never leaves the host disabled with no marker (which would strand + // it out of the scheduling pool with no self-heal). Then disable it so the allocator stops placing new VMs + // on it (it is the least-loaded host, which the allocator would otherwise prefer, and that would fight the + // drain). If the host cannot be disabled, undo the marker and do not start. A drain already in progress + // (previous != null) already has its marker and is already disabled. + if (previous == null) { + setPowerMarker(candidate.getId(), DRS_POWER_STATE_DRAINING); + if (candidate.getResourceState() == ResourceState.Enabled + && !hostDao.updateResourceState(ResourceState.Enabled, ResourceState.Event.Disable, ResourceState.Disabled, candidate)) { + logger.warn("DRS power management: could not disable host [{}] to begin draining it; skipping.", candidate.getId()); + clearStalePowerMarker(candidate); + return; + } + } + logger.info("DRS power management: cluster [{}] is under-utilized; draining host [{}] ({} VMs left, {} this poll) to power it off.", + cluster.getId(), candidate.getId(), plan.size(), batch.size()); + // One start event for the batch. Each migration job carries it as its start event id and completes it, as in + // executeDrsPlan, so the event is never left open; opening it only when the batch is non-empty avoids an orphan. + long eventId = ActionEventUtils.onStartedActionEvent(User.UID_SYSTEM, Account.ACCOUNT_ID_SYSTEM, + EventTypes.EVENT_VM_MIGRATE, + String.format("DRS power management draining host %d in cluster %s", candidate.getId(), cluster.getUuid()), + candidate.getId(), ApiCommandResourceType.Host.toString(), true, 0); + for (Map.Entry migration : batch.entrySet()) { + createMigrateVMAsyncJob(migration.getKey(), migration.getValue(), eventId); + } + drainingVmCountByHost.put(candidate.getId(), current); + } + + protected double vmResourceNeed(VMInstanceVO vm, boolean useCpu) { + ServiceOffering offering = serviceOfferingDao.findByIdIncludingRemoved(vm.getId(), vm.getServiceOfferingId()); + if (offering == null) { + return 0d; + } + if (useCpu) { + return (double) offering.getCpu() * offering.getSpeed(); + } + return (double) offering.getRamSize() * 1024L * 1024L; + } + + /** + * First-fit-decreasing placement of the given VMs (id to resource need) onto hosts (id to free capacity), + * returning a VM-to-host assignment or an empty map if any VM does not fit on capacity alone. This is a + * capacity feasibility pre-check only: it does NOT account for affinity/anti-affinity, host tags, storage + * access or dedication. Each migration is still validated authoritatively by the migration job, which may + * reject a placement this map proposed; in that case the host is drained over later polls rather than at once. + * Pure function, unit-tested. + */ + protected Map planHostEvacuation(Map vmNeeds, Map hostFreeCapacity) { + Map assignment = new HashMap<>(); + // TreeMap so host iteration is by id and placement is deterministic run-to-run; secondary sort by vm id + // breaks ties among equal-sized VMs for the same reason. + Map remaining = new TreeMap<>(hostFreeCapacity); + List> vms = new ArrayList<>(vmNeeds.entrySet()); + vms.sort((a, b) -> b.getValue().equals(a.getValue()) ? Long.compare(a.getKey(), b.getKey()) : Double.compare(b.getValue(), a.getValue())); + for (Map.Entry vm : vms) { + Long target = null; + for (Map.Entry host : remaining.entrySet()) { + if (host.getValue() >= vm.getValue()) { + target = host.getKey(); + break; + } + } + if (target == null) { + return Collections.emptyMap(); + } + assignment.put(vm.getKey(), target); + remaining.put(target, remaining.get(target) - vm.getValue()); + } + return assignment; + } + + /** + * Appends the latest cluster utilization ratio to the cluster's rolling history (trimmed to {@code window} + * samples) and returns a snapshot of it, newest last. + */ + protected List recordAndGetUtilizationHistory(long clusterId, double ratio, int window) { + int cap = Math.max(1, window); + LinkedList history = clusterUtilizationHistory.computeIfAbsent(clusterId, k -> new LinkedList<>()); + synchronized (history) { + history.addLast(ratio); + while (history.size() > cap) { + history.removeFirst(); + } + return new ArrayList<>(history); + } + } + + /** + * Projects one poll interval ahead from the recent utilization samples using the least-squares trend, clamped + * to [0, 1]. With fewer than two samples it returns the latest sample (no trend to project). + */ + protected double forecastUtilization(List history) { + if (history == null || history.isEmpty()) { + return 0d; + } + int n = history.size(); + if (n == 1) { + return history.get(0); + } + double sumX = 0d; + double sumY = 0d; + double sumXY = 0d; + double sumXX = 0d; + for (int i = 0; i < n; i++) { + double y = history.get(i); + sumX += i; + sumY += y; + sumXY += i * y; + sumXX += (double) i * i; + } + double denominator = n * sumXX - sumX * sumX; + double slope = denominator == 0d ? 0d : (n * sumXY - sumX * sumY) / denominator; + double forecast = history.get(n - 1) + slope; + return Math.max(0d, Math.min(1d, forecast)); + } + + protected Map> getHostCapacityMap(List hosts, boolean useCpu) { + List joins = hostJoinDao.searchByIds(hosts.stream().map(HostVO::getId).toArray(Long[]::new)); + Map> map = new HashMap<>(); + for (HostJoinVO join : joins) { + if (useCpu) { + Integer cpus = join.getCpus(); + Long speed = join.getSpeed(); + long cpuTotal = (cpus == null || speed == null) ? 0L : (long) cpus * speed; + map.put(join.getId(), new Ternary<>(join.getCpuUsedCapacity(), join.getCpuReservedCapacity(), cpuTotal)); + } else { + map.put(join.getId(), new Ternary<>(join.getMemUsedCapacity(), join.getMemReservedCapacity(), join.getTotalMemory())); + } + } + return map; + } + + /** + * True when the cluster can give up the candidate host: the cluster is below the low utilization threshold, + * there is more than one host up, and the remaining hosts can carry the current load without crossing the + * high threshold. + */ + protected boolean clusterCanReleaseHost(double clusterUsed, double clusterTotal, double candidateHostTotal, + int upHostCount, float lowThreshold, float highThreshold) { + if (upHostCount <= 1 || clusterTotal <= 0d) { + return false; + } + if ((clusterUsed / clusterTotal) >= lowThreshold) { + return false; + } + double remainingTotal = clusterTotal - candidateHostTotal; + if (remainingTotal <= 0d) { + return false; + } + return (clusterUsed / remainingTotal) <= highThreshold; + } + + /** + * True when the cluster is above the high utilization threshold and there is a host that DRS powered off + * earlier which can be powered back on. + */ + protected boolean clusterNeedsWakeup(double clusterUsed, double clusterTotal, float highThreshold, int poweredOffHostCount) { + if (poweredOffHostCount <= 0 || clusterTotal <= 0d) { + return false; + } + return (clusterUsed / clusterTotal) > highThreshold; + } + + protected boolean hostHasNoRunningVms(HostVO host) { + List vms = vmInstanceDao.listByHostId(host.getId()); + return vms == null || vms.isEmpty(); + } + + protected boolean isPowerManageable(HostVO host) { + try { + return outOfBandManagementService.isOutOfBandManagementEnabled(host); + } catch (Exception e) { + logger.debug("Could not determine out-of-band management status for host [{}]: {}", host.getId(), e.getMessage()); + return false; + } + } + + /** + * A host can be drained by power management only if every VM on it is a running user VM: a system VM cannot be + * moved with a user-VM migration, and a VM in a transitional state must not be forced, so such a host is never + * chosen (it could not be fully emptied). + */ + protected boolean hostEvacuatable(HostVO host) { + List vms = vmInstanceDao.listByHostId(host.getId()); + if (vms == null || vms.isEmpty()) { + return false; + } + for (VMInstanceVO vm : vms) { + if (vm.getType().isUsedBySystem() || vm.getState() != VirtualMachine.State.Running) { + return false; + } + } + return true; + } + + protected boolean isPoweredOffByDrs(HostVO host) { + DetailVO detail = hostDetailsDao.findDetail(host.getId(), DRS_POWER_STATE_DETAIL); + return detail != null && DRS_POWER_STATE_OFF.equals(detail.getValue()); + } + + protected void clearStalePowerMarker(HostVO host) { + DetailVO detail = hostDetailsDao.findDetail(host.getId(), DRS_POWER_STATE_DETAIL); + if (detail != null) { + hostDetailsDao.remove(detail.getId()); + } + } + + protected boolean hasInFlightDrsPlan(long clusterId) { + return !drsPlanDao.listByClusterIdAndStatus(clusterId, ClusterDrsPlan.Status.UNDER_REVIEW).isEmpty() + || !drsPlanDao.listByClusterIdAndStatus(clusterId, ClusterDrsPlan.Status.READY).isEmpty() + || !drsPlanDao.listByClusterIdAndStatus(clusterId, ClusterDrsPlan.Status.IN_PROGRESS).isEmpty(); + } + + protected void powerOffHost(HostVO host, ClusterVO cluster) { + logger.info("DRS power management: cluster [{}] is under-utilized; disabling and powering off empty host [{}].", cluster.getId(), host.getId()); + // Disable first so CloudStack stops scheduling to the host and does not treat the imminent agent disconnect + // as a failure (host monitor / HA). A host that was being drained is already Disabled, so only transition an + // Enabled host, and abort only when that transition does not take effect, so the host is never powered off + // while CloudStack still believes it is schedulable. + if (host.getResourceState() == ResourceState.Enabled + && !hostDao.updateResourceState(ResourceState.Enabled, ResourceState.Event.Disable, ResourceState.Disabled, host)) { + logger.warn("DRS power management: could not disable host [{}]; skipping power-off.", host.getId()); + return; + } + // Record the durable wake intent BEFORE the irreversible power-off. If the management server dies between + // the power-off and here, the host is still recognised on the next poll and powered back on, rather than + // left off forever with no marker. Replace any existing marker (e.g. a draining marker when powering off a + // host that has just finished draining) so only the off marker remains. + DetailVO existing = hostDetailsDao.findDetail(host.getId(), DRS_POWER_STATE_DETAIL); + if (existing != null) { + hostDetailsDao.remove(existing.getId()); + } + DetailVO marker = new DetailVO(host.getId(), DRS_POWER_STATE_DETAIL, DRS_POWER_STATE_OFF); + hostDetailsDao.persist(marker); + try { + // Detach the agent without investigation first so the imminent link drop from cutting power is not + // reported as a host-down failure (the host is intentionally going away, and it is already empty). + agentManager.disconnectWithoutInvestigation(host.getId(), Status.Event.ShutdownRequested); + outOfBandManagementService.executePowerOperation(host, OutOfBandManagement.PowerOperation.OFF, null); + } catch (Exception e) { + // A failure response does not prove the host is still running: a chassis power-off can take effect and + // still report a timeout or error. Rather than roll back (which would strand a host that did power off + // as enabled and unmarked), leave it disabled and marked. The next poll reclaims it (re-enable, clear + // marker) if it is still up, or treats it as a DRS-powered-off host if it is actually down. + logger.warn("DRS power management: power-off command for host [{}] failed; leaving it disabled and marked for reconciliation on the next poll.", + host.getId(), e); + } + drainingVmCountByHost.remove(host.getId()); + } + + protected void powerOnHost(HostVO host, ClusterVO cluster) { + logger.info("DRS power management: cluster [{}] is over-utilized; powering on host [{}].", cluster.getId(), host.getId()); + // Issue the power-on but keep the marker and leave the host Disabled: the host is still down until its + // agent reconnects, and an accepted-but-unconfirmed power-on must not look done. The classification loop + // re-enables the host and clears the marker only once it is actually Up. Power-on is idempotent, so a host + // still booting is simply re-issued the command on a later poll until it connects. + outOfBandManagementService.executePowerOperation(host, OutOfBandManagement.PowerOperation.ON, null); + } + @Override public ConfigKey[] getConfigKeys() { return new ConfigKey[]{ClusterDrsPlanExpireInterval, ClusterDrsEnabled, ClusterDrsInterval, ClusterDrsMaxMigrations, ClusterDrsAlgorithm, ClusterDrsImbalanceThreshold, ClusterDrsMetric, ClusterDrsMetricType, ClusterDrsMetricUseRatio, - ClusterDrsImbalanceSkipThreshold}; + ClusterDrsImbalanceSkipThreshold, ClusterDrsPowerManagementEnabled, ClusterDrsPowerManagementLowThreshold, + ClusterDrsPowerManagementHighThreshold, ClusterDrsPredictiveEnabled, ClusterDrsPredictiveWindow, + ClusterDrsPowerManagementEvacuate}; } @Override diff --git a/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsPowerManagementTest.java b/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsPowerManagementTest.java new file mode 100644 index 000000000000..41a6b26fd400 --- /dev/null +++ b/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsPowerManagementTest.java @@ -0,0 +1,150 @@ +// 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 org.apache.cloudstack.cluster; + +import java.util.Arrays; +import java.util.Collections; + +import org.junit.Assert; +import org.junit.Test; + +public class ClusterDrsPowerManagementTest { + + private final ClusterDrsServiceImpl drs = new ClusterDrsServiceImpl(); + + @Test + public void releasesHostWhenUnderUtilizedAndRemainingHostsCanCarryLoad() { + // 4 hosts of 100 each, 90 used total (22.5%): below low 0.30, and after dropping one host 90/300 = 30% <= high 0.75. + Assert.assertTrue(drs.clusterCanReleaseHost(90d, 400d, 100d, 4, 0.30f, 0.75f)); + } + + @Test + public void keepsHostWhenClusterIsNotUnderUtilized() { + // 250/400 = 62.5% is above the low threshold, so nothing is powered off. + Assert.assertFalse(drs.clusterCanReleaseHost(250d, 400d, 100d, 4, 0.30f, 0.75f)); + } + + @Test + public void keepsHostWhenRemainingHostsWouldExceedHighThreshold() { + // 40/200 = 20% used (below the low threshold), but the candidate is a large host (170 of 200): dropping it + // leaves only 30 of capacity, so 40/30 = 133% would blow past the high threshold. Keep it. + Assert.assertFalse(drs.clusterCanReleaseHost(40d, 200d, 170d, 2, 0.30f, 0.75f)); + } + + @Test + public void neverReleasesTheLastHost() { + Assert.assertFalse(drs.clusterCanReleaseHost(10d, 100d, 100d, 1, 0.30f, 0.75f)); + } + + @Test + public void wakesHostWhenOverUtilizedAndOneWasPoweredOff() { + // 320/400 = 80% is above high 0.75 and one host is available to wake. + Assert.assertTrue(drs.clusterNeedsWakeup(320d, 400d, 0.75f, 1)); + } + + @Test + public void doesNotWakeWhenNoHostWasPoweredOff() { + Assert.assertFalse(drs.clusterNeedsWakeup(320d, 400d, 0.75f, 0)); + } + + @Test + public void doesNotWakeWhenBelowHighThreshold() { + Assert.assertFalse(drs.clusterNeedsWakeup(200d, 400d, 0.75f, 1)); + } + + @Test + public void forecastProjectsARisingTrendAboveTheLastSample() { + double forecast = drs.forecastUtilization(Arrays.asList(0.40d, 0.50d, 0.60d, 0.70d)); + Assert.assertTrue("rising trend should forecast above the last sample", forecast > 0.70d); + } + + @Test + public void forecastOnFlatSeriesEqualsTheLastSample() { + Assert.assertEquals(0.50d, drs.forecastUtilization(Arrays.asList(0.50d, 0.50d, 0.50d)), 0.0001d); + } + + @Test + public void forecastOnFallingTrendIsBelowTheLastSample() { + double forecast = drs.forecastUtilization(Arrays.asList(0.80d, 0.70d, 0.60d, 0.50d)); + Assert.assertTrue("falling trend should forecast below the last sample", forecast < 0.50d); + } + + @Test + public void forecastClampsToOne() { + double forecast = drs.forecastUtilization(Arrays.asList(0.70d, 0.85d, 0.99d)); + Assert.assertTrue("forecast must never exceed 1.0", forecast <= 1.0d); + } + + @Test + public void forecastOfSingleSampleIsThatSample() { + Assert.assertEquals(0.42d, drs.forecastUtilization(Collections.singletonList(0.42d)), 0.0001d); + } + + @Test + public void forecastOfEmptyHistoryIsZero() { + Assert.assertEquals(0d, drs.forecastUtilization(Collections.emptyList()), 0.0001d); + } + + @Test + public void rollingHistoryIsTrimmedToWindowNewestLast() { + drs.recordAndGetUtilizationHistory(99L, 0.10d, 3); + drs.recordAndGetUtilizationHistory(99L, 0.20d, 3); + drs.recordAndGetUtilizationHistory(99L, 0.30d, 3); + java.util.List history = drs.recordAndGetUtilizationHistory(99L, 0.40d, 3); + Assert.assertEquals(Arrays.asList(0.20d, 0.30d, 0.40d), history); + } + + @Test + public void evacuationPlanPlacesEveryVmWhenTheyFit() { + java.util.Map vmNeeds = new java.util.HashMap<>(); + vmNeeds.put(1L, 30d); + vmNeeds.put(2L, 20d); + java.util.Map hostFree = new java.util.HashMap<>(); + hostFree.put(10L, 50d); + hostFree.put(11L, 40d); + + java.util.Map plan = drs.planHostEvacuation(vmNeeds, hostFree); + + Assert.assertEquals(2, plan.size()); + Assert.assertTrue(plan.containsKey(1L) && plan.containsKey(2L)); + } + + @Test + public void evacuationPlanIsEmptyWhenAVmCannotBePlaced() { + java.util.Map vmNeeds = new java.util.HashMap<>(); + vmNeeds.put(1L, 100d); + java.util.Map hostFree = new java.util.HashMap<>(); + hostFree.put(10L, 50d); + + Assert.assertTrue(drs.planHostEvacuation(vmNeeds, hostFree).isEmpty()); + } + + @Test + public void evacuationPlanPacksLargestFirst() { + java.util.Map vmNeeds = new java.util.HashMap<>(); + vmNeeds.put(1L, 60d); + vmNeeds.put(2L, 40d); + vmNeeds.put(3L, 40d); + java.util.Map hostFree = new java.util.HashMap<>(); + hostFree.put(10L, 60d); + hostFree.put(11L, 80d); + + java.util.Map plan = drs.planHostEvacuation(vmNeeds, hostFree); + + Assert.assertEquals("all three VMs placed", 3, plan.size()); + } +} diff --git a/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsPowerOrchestrationTest.java b/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsPowerOrchestrationTest.java new file mode 100644 index 000000000000..12705705bff6 --- /dev/null +++ b/server/src/test/java/org/apache/cloudstack/cluster/ClusterDrsPowerOrchestrationTest.java @@ -0,0 +1,325 @@ +// 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 org.apache.cloudstack.cluster; + +import java.util.Arrays; +import java.util.Collections; + +import org.apache.cloudstack.outofbandmanagement.OutOfBandManagement; +import org.apache.cloudstack.outofbandmanagement.OutOfBandManagementService; +import org.apache.cloudstack.cluster.dao.ClusterDrsPlanDao; +import com.cloud.agent.AgentManager; +import com.cloud.utils.Ternary; +import com.cloud.host.Status; +import org.junit.Assert; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.InjectMocks; +import org.mockito.InOrder; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; + +import com.cloud.host.DetailVO; +import com.cloud.host.HostVO; +import com.cloud.host.dao.HostDao; +import com.cloud.host.dao.HostDetailsDao; +import com.cloud.dc.ClusterVO; +import com.cloud.resource.ResourceState; +import com.cloud.service.ServiceOfferingVO; +import com.cloud.service.dao.ServiceOfferingDao; +import com.cloud.vm.VMInstanceVO; +import com.cloud.vm.VirtualMachine; +import com.cloud.vm.dao.VMInstanceDao; + +@RunWith(MockitoJUnitRunner.class) +public class ClusterDrsPowerOrchestrationTest { + + @Mock + private VMInstanceDao vmInstanceDao; + @Mock + private ServiceOfferingDao serviceOfferingDao; + @Mock + private OutOfBandManagementService outOfBandManagementService; + @Mock + private HostDao hostDao; + @Mock + private HostDetailsDao hostDetailsDao; + @Mock + private ClusterDrsPlanDao drsPlanDao; + @Mock + private AgentManager agentManager; + + @InjectMocks + private ClusterDrsServiceImpl drs = new ClusterDrsServiceImpl(); + + private ClusterVO cluster(long id) { + ClusterVO cluster = Mockito.mock(ClusterVO.class); + Mockito.lenient().when(cluster.getId()).thenReturn(id); + return cluster; + } + + private VMInstanceVO vm(long id, VirtualMachine.Type type, VirtualMachine.State state) { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.lenient().when(vm.getId()).thenReturn(id); + Mockito.lenient().when(vm.getType()).thenReturn(type); + Mockito.lenient().when(vm.getState()).thenReturn(state); + return vm; + } + + private HostVO host(long id) { + HostVO host = Mockito.mock(HostVO.class); + Mockito.lenient().when(host.getId()).thenReturn(id); + Mockito.lenient().when(host.getResourceState()).thenReturn(ResourceState.Enabled); + return host; + } + + @Test + public void hostEvacuatableTrueForRunningUserVmsOnly() { + VMInstanceVO a = vm(1L, VirtualMachine.Type.User, VirtualMachine.State.Running); + VMInstanceVO b = vm(2L, VirtualMachine.Type.User, VirtualMachine.State.Running); + Mockito.when(vmInstanceDao.listByHostId(10L)).thenReturn(Arrays.asList(a, b)); + Assert.assertTrue(drs.hostEvacuatable(host(10L))); + } + + @Test + public void hostEvacuatableFalseWhenASystemVmIsPresent() { + VMInstanceVO user = vm(1L, VirtualMachine.Type.User, VirtualMachine.State.Running); + VMInstanceVO router = vm(3L, VirtualMachine.Type.DomainRouter, VirtualMachine.State.Running); + Mockito.when(vmInstanceDao.listByHostId(11L)).thenReturn(Arrays.asList(user, router)); + Assert.assertFalse(drs.hostEvacuatable(host(11L))); + } + + @Test + public void hostEvacuatableFalseWhenAVmIsNotRunning() { + VMInstanceVO stopping = vm(4L, VirtualMachine.Type.User, VirtualMachine.State.Stopping); + Mockito.when(vmInstanceDao.listByHostId(12L)).thenReturn(Collections.singletonList(stopping)); + Assert.assertFalse(drs.hostEvacuatable(host(12L))); + } + + @Test + public void hostEvacuatableFalseWhenEmpty() { + Mockito.when(vmInstanceDao.listByHostId(13L)).thenReturn(Collections.emptyList()); + Assert.assertFalse(drs.hostEvacuatable(host(13L))); + } + + @Test + public void vmResourceNeedZeroWhenOfferingMissing() { + VMInstanceVO vm = vm(5L, VirtualMachine.Type.User, VirtualMachine.State.Running); + Mockito.when(serviceOfferingDao.findByIdIncludingRemoved(Mockito.eq(5L), Mockito.anyLong())).thenReturn(null); + Assert.assertEquals(0d, drs.vmResourceNeed(vm, true), 0.0001d); + } + + @Test + public void vmResourceNeedComputesCpuAndMemory() { + VMInstanceVO vm = vm(6L, VirtualMachine.Type.User, VirtualMachine.State.Running); + Mockito.when(vm.getServiceOfferingId()).thenReturn(99L); + ServiceOfferingVO offering = Mockito.mock(ServiceOfferingVO.class); + Mockito.when(offering.getCpu()).thenReturn(4); + Mockito.when(offering.getSpeed()).thenReturn(2000); + Mockito.when(offering.getRamSize()).thenReturn(2048); + Mockito.when(serviceOfferingDao.findByIdIncludingRemoved(6L, 99L)).thenReturn(offering); + + Assert.assertEquals(8000d, drs.vmResourceNeed(vm, true), 0.0001d); + Assert.assertEquals(2048d * 1024L * 1024L, drs.vmResourceNeed(vm, false), 0.0001d); + } + + @Test + public void isPowerManageableReflectsOobmEnabled() { + HostVO h = host(20L); + Mockito.when(outOfBandManagementService.isOutOfBandManagementEnabled(h)).thenReturn(true); + Assert.assertTrue(drs.isPowerManageable(h)); + } + + @Test + public void isPowerManageableFalseWhenOobmLookupThrows() { + HostVO h = host(21L); + Mockito.when(outOfBandManagementService.isOutOfBandManagementEnabled(h)).thenThrow(new RuntimeException("boom")); + Assert.assertFalse(drs.isPowerManageable(h)); + } + + @Test + public void powerOffPersistsWakeMarkerBeforeThePowerOff() { + HostVO h = host(30L); + Mockito.when(hostDao.updateResourceState(ResourceState.Enabled, ResourceState.Event.Disable, ResourceState.Disabled, h)).thenReturn(true); + + drs.powerOffHost(h, cluster(1L)); + + // the durable wake marker must be written BEFORE the irreversible power-off, so a crash in between + // still leaves a host that can be recognised and powered back on; and the agent is detached without + // investigation before the power is cut so the link drop is not reported as a host-down failure. + InOrder inOrder = Mockito.inOrder(hostDao, hostDetailsDao, agentManager, outOfBandManagementService); + inOrder.verify(hostDao).updateResourceState(ResourceState.Enabled, ResourceState.Event.Disable, ResourceState.Disabled, h); + inOrder.verify(hostDetailsDao).persist(Mockito.any(DetailVO.class)); + inOrder.verify(agentManager).disconnectWithoutInvestigation(30L, Status.Event.ShutdownRequested); + inOrder.verify(outOfBandManagementService).executePowerOperation(Mockito.eq(h), Mockito.eq(OutOfBandManagement.PowerOperation.OFF), Mockito.any()); + } + + @Test + public void powerOffAbortsWhenTheHostCannotBeDisabled() { + HostVO h = host(31L); + Mockito.when(hostDao.updateResourceState(ResourceState.Enabled, ResourceState.Event.Disable, ResourceState.Disabled, h)).thenReturn(false); + + drs.powerOffHost(h, cluster(1L)); + + // a host that could not be disabled must never be powered off, and must not be marked. + Mockito.verify(outOfBandManagementService, Mockito.never()).executePowerOperation(Mockito.any(), Mockito.any(), Mockito.any()); + Mockito.verify(hostDetailsDao, Mockito.never()).persist(Mockito.any(DetailVO.class)); + } + + @Test + public void powerOffLeavesHostDisabledAndMarkedWhenThePowerOffFails() { + HostVO h = host(32L); + Mockito.when(hostDao.updateResourceState(ResourceState.Enabled, ResourceState.Event.Disable, ResourceState.Disabled, h)).thenReturn(true); + Mockito.doThrow(new RuntimeException("oobm down")).when(outOfBandManagementService) + .executePowerOperation(Mockito.eq(h), Mockito.eq(OutOfBandManagement.PowerOperation.OFF), Mockito.any()); + + drs.powerOffHost(h, cluster(1L)); + + // a failed power-off must not roll back: the host stays disabled and marked so the next poll reconciles it, + // because the failure does not prove the chassis is still on. No marker removal, no re-enable. + Mockito.verify(hostDetailsDao, Mockito.never()).remove(Mockito.anyLong()); + Mockito.verify(hostDao, Mockito.never()).updateResourceState(ResourceState.Disabled, ResourceState.Event.Enable, ResourceState.Enabled, h); + } + + @Test + public void powerOnKeepsTheMarkerAndDoesNotEnableUntilTheHostIsUp() { + HostVO h = host(33L); + + drs.powerOnHost(h, cluster(1L)); + + // the power-on is issued, but the marker is kept and the host is not re-enabled: an unconfirmed power-on + // must not look done, or a host that never actually boots is stranded out of the wake set. + Mockito.verify(outOfBandManagementService).executePowerOperation(Mockito.eq(h), Mockito.eq(OutOfBandManagement.PowerOperation.ON), Mockito.any()); + Mockito.verify(hostDetailsDao, Mockito.never()).remove(Mockito.anyLong()); + Mockito.verify(hostDao, Mockito.never()).updateResourceState(Mockito.any(), Mockito.any(), Mockito.any(), Mockito.any()); + } + + @Test + public void hasInFlightDrsPlanTrueWhenAPlanIsInProgress() { + Mockito.when(drsPlanDao.listByClusterIdAndStatus(5L, ClusterDrsPlan.Status.UNDER_REVIEW)).thenReturn(Collections.emptyList()); + Mockito.when(drsPlanDao.listByClusterIdAndStatus(5L, ClusterDrsPlan.Status.READY)).thenReturn(Collections.emptyList()); + Mockito.when(drsPlanDao.listByClusterIdAndStatus(5L, ClusterDrsPlan.Status.IN_PROGRESS)) + .thenReturn(Collections.singletonList(Mockito.mock(ClusterDrsPlanVO.class))); + Assert.assertTrue(drs.hasInFlightDrsPlan(5L)); + } + + @Test + public void hasInFlightDrsPlanFalseWhenNoPlansArePending() { + Mockito.when(drsPlanDao.listByClusterIdAndStatus(Mockito.eq(5L), Mockito.any())).thenReturn(Collections.emptyList()); + Assert.assertFalse(drs.hasInFlightDrsPlan(5L)); + } + + @Test + public void hasMigratingVmDetectsAnInFlightMigration() { + VMInstanceVO running = vm(1L, VirtualMachine.Type.User, VirtualMachine.State.Running); + VMInstanceVO migrating = vm(2L, VirtualMachine.Type.User, VirtualMachine.State.Migrating); + Mockito.when(vmInstanceDao.listByHostId(40L)).thenReturn(Arrays.asList(running, migrating)); + Assert.assertTrue(drs.hasMigratingVm(40L)); + } + + @Test + public void hasMigratingVmFalseWhenNothingIsMigrating() { + VMInstanceVO running = vm(1L, VirtualMachine.Type.User, VirtualMachine.State.Running); + Mockito.when(vmInstanceDao.listByHostId(40L)).thenReturn(Collections.singletonList(running)); + Assert.assertFalse(drs.hasMigratingVm(40L)); + } + + @Test + public void isDrainingByDrsTrueOnlyForTheDrainingMarker() { + HostVO h = host(41L); + Mockito.when(hostDetailsDao.findDetail(41L, "drs.power.state")) + .thenReturn(new DetailVO(41L, "drs.power.state", "draining")); + Assert.assertTrue(drs.isDrainingByDrs(h)); + } + + @Test + public void drainBatchWaitsWhileAMigrationIsInFlight() { + HostVO candidate = host(42L); + VMInstanceVO migrating = vm(1L, VirtualMachine.Type.User, VirtualMachine.State.Migrating); + Mockito.when(vmInstanceDao.listByHostId(42L)).thenReturn(Collections.singletonList(migrating)); + drs.drainingVmCountByHost.put(42L, 3); + + drs.drainHostBatch(cluster(1L), candidate, Collections.emptyList(), Collections.emptyMap(), 0d, 100d, 0.30f, 0.75f, true); + + // a batch is still running: do not touch the marker or the recorded progress, just wait. + Mockito.verify(hostDetailsDao, Mockito.never()).remove(Mockito.anyLong()); + Assert.assertEquals(Integer.valueOf(3), drs.drainingVmCountByHost.get(42L)); + } + + @Test + public void drainBatchAbandonsWhenNoProgressSinceLastBatch() { + HostVO candidate = host(43L); + VMInstanceVO a = vm(1L, VirtualMachine.Type.User, VirtualMachine.State.Running); + VMInstanceVO b = vm(2L, VirtualMachine.Type.User, VirtualMachine.State.Running); + Mockito.when(vmInstanceDao.listByHostId(43L)).thenReturn(Arrays.asList(a, b)); + Mockito.when(hostDetailsDao.findDetail(43L, "drs.power.state")) + .thenReturn(new DetailVO(43L, "drs.power.state", "draining")); + drs.drainingVmCountByHost.put(43L, 2); // same count as now: the last batch moved nothing + + drs.drainHostBatch(cluster(1L), candidate, Collections.emptyList(), Collections.emptyMap(), 0d, 100d, 0.30f, 0.75f, true); + + // no progress: the drain is abandoned (marker cleared, progress forgotten) instead of re-submitted forever. + Mockito.verify(hostDetailsDao).remove(Mockito.anyLong()); + Assert.assertFalse(drs.drainingVmCountByHost.containsKey(43L)); + } + + @Test + public void drainBatchPowersOffAnEmptiedDrainingHostWhenStillReleasable() { + HostVO candidate = host(44L); + HostVO other = host(45L); + // a drained host is already disabled, so powerOffHost must proceed without re-disabling it. + Mockito.when(candidate.getResourceState()).thenReturn(ResourceState.Disabled); + Mockito.when(vmInstanceDao.listByHostId(44L)).thenReturn(Collections.emptyList()); + Mockito.when(outOfBandManagementService.isOutOfBandManagementEnabled(candidate)).thenReturn(true); + java.util.Map> capacityMap = new java.util.HashMap<>(); + capacityMap.put(44L, new Ternary<>(0L, 0L, 100L)); + capacityMap.put(45L, new Ternary<>(10L, 0L, 100L)); + drs.drainingVmCountByHost.put(44L, 1); + + // cluster at 10/200 = 5% (below low 0.30), remaining after release 10/100 = 10% (below high 0.75): releasable. + drs.drainHostBatch(cluster(1L), candidate, Collections.singletonList(other), capacityMap, 10d, 200d, 0.30f, 0.75f, true); + + // the emptied host is powered off (the empty-host loop never sees it because it is disabled) and its + // progress entry is cleared. + Mockito.verify(outOfBandManagementService).executePowerOperation(Mockito.eq(candidate), Mockito.eq(OutOfBandManagement.PowerOperation.OFF), Mockito.any()); + Assert.assertFalse(drs.drainingVmCountByHost.containsKey(44L)); + } + + @Test + public void drainBatchReturnsAnEmptiedHostToServiceWhenNoLongerPowerManageable() { + HostVO candidate = host(46L); + HostVO other = host(47L); + Mockito.when(candidate.getResourceState()).thenReturn(ResourceState.Disabled); + Mockito.when(vmInstanceDao.listByHostId(46L)).thenReturn(Collections.emptyList()); + // out-of-band management is no longer available for the drained host. + Mockito.when(outOfBandManagementService.isOutOfBandManagementEnabled(candidate)).thenReturn(false); + Mockito.when(hostDetailsDao.findDetail(46L, "drs.power.state")) + .thenReturn(new DetailVO(46L, "drs.power.state", "draining")); + java.util.Map> capacityMap = new java.util.HashMap<>(); + capacityMap.put(46L, new Ternary<>(0L, 0L, 100L)); + capacityMap.put(47L, new Ternary<>(10L, 0L, 100L)); + drs.drainingVmCountByHost.put(46L, 1); + + drs.drainHostBatch(cluster(1L), candidate, Collections.singletonList(other), capacityMap, 10d, 200d, 0.30f, 0.75f, true); + + // it is not powered off; instead it is re-enabled and unmarked so it goes back into the scheduling pool. + Mockito.verify(outOfBandManagementService, Mockito.never()).executePowerOperation(Mockito.any(), Mockito.any(), Mockito.any()); + Mockito.verify(hostDao).updateResourceState(ResourceState.Disabled, ResourceState.Event.Enable, ResourceState.Enabled, candidate); + Assert.assertFalse(drs.drainingVmCountByHost.containsKey(46L)); + } +}