diff --git a/agent/conf/agent.properties b/agent/conf/agent.properties index 4e36eff4d75f..71549ecddb6f 100644 --- a/agent/conf/agent.properties +++ b/agent/conf/agent.properties @@ -66,6 +66,41 @@ zone=default # If this is commented, the value of the private NIC device will be used. #guest.network.device= +# Policy for encrypting the live-migration data stream using QEMU-native TLS. +# "Disabled" (default, plaintext) or "Required" (use TLS, and fail the migration if +# the host libvirt is older than 3.2.0, rather than silently sending plaintext). +# "Required" needs a migration x509 trust set up on every participating host first: +# a CA plus server and client certs under migrate_tls_x509_cert_dir in qemu.conf +# (default /etc/pki/qemu), with the cert SAN covering the host's migration address, +# then restart libvirtd. Without that, libvirt rejects the TLS migration. +#migrate.encryption.policy=Disabled + +# Use multiple parallel TCP streams (multifd) for live migration to fill fast NICs. +#migrate.parallel.enabled=false + +# Number of parallel connections (multifd channels) for live migration, used only +# when migrate.parallel.enabled is true. 0 lets libvirt/QEMU pick its default. +# To use more than one physical NIC for migration bandwidth, bond the NICs under the +# migration bridge; libvirt sends all channels to a single destination address and +# cannot bind them to separate NICs itself. For the channels to spread across bond +# members the bond must hash on the L4 ports (802.3ad or balance-xor with +# xmit_hash_policy=layer3+4); the default layer2 hash pins all channels to one member +# because they share a MAC and IP pair, giving no spread. The spread is statistical, +# not guaranteed. +#migrate.parallel.connections=0 + +# Allow migrating a VM whose disk cache mode libvirt considers unsafe (e.g. writeback). +# Only enable on coherent shared storage such as Ceph RBD. +#migrate.allow.unsafe=false + +# Run a virsh cpu-compare precheck against the destination before migrating, and abort +# early with a clear error if the destination CPU cannot run the guest. +#migrate.cpu.precheck.enabled=false + +# Memory-compression method for non-parallel migration, xbzrle or mt. Blank leaves libvirt's +# default method. Ignored when migrate.parallel.enabled is true. +#migrate.compression.method= + # Local storage path. Multiple values can be entered and separated by commas. #local.storage.path=/var/lib/libvirt/images/ diff --git a/agent/src/main/java/com/cloud/agent/properties/AgentProperties.java b/agent/src/main/java/com/cloud/agent/properties/AgentProperties.java index e4775188d0ca..e77f3dbe2820 100644 --- a/agent/src/main/java/com/cloud/agent/properties/AgentProperties.java +++ b/agent/src/main/java/com/cloud/agent/properties/AgentProperties.java @@ -146,6 +146,63 @@ public class AgentProperties{ */ public static final Property QEMU_SOCKETS_PATH = new Property<>("qemu.sockets.path", "/var/lib/libvirt/qemu"); + /** + * Policy for encrypting the live-migration data stream (guest RAM, and disk contents + * during storage migration) using QEMU-native TLS (VIR_MIGRATE_TLS). Values: "Disabled" + * (default, plaintext TCP) or "Required" (encrypt, and fail the migration if the host has no + * TLS migration environment). Required needs migrate_tls_x509_cert_dir / default_tls_x509_cert_dir + * configured in qemu.conf on every host. + * Data type: String.
+ * Default value: "Disabled". + */ + public static final Property MIGRATE_ENCRYPTION_POLICY = new Property<>("migrate.encryption.policy", "Disabled"); + + /** + * Use multiple parallel TCP streams (multifd) for live migration to saturate fast NICs + * (25/40/100GbE). Single-stream migration cannot fill such links. Requires libvirt 5.2.0+. + * Data type: Boolean.
+ * Default value: false. + */ + public static final Property MIGRATE_PARALLEL_ENABLED = new Property<>("migrate.parallel.enabled", false); + + /** + * Number of parallel connections (multifd channels) to use for live migration when migrate.parallel.enabled + * is set. 0 lets libvirt/QEMU choose its default. To spread migration across more than one physical NIC, + * bond the NICs (for example LACP) under the migration bridge; libvirt sends all channels to a single + * destination address, so it cannot bind them to separate NICs itself.
+ * Data type: Integer.
+ * Default value: 0. + */ + public static final Property MIGRATE_PARALLEL_CONNECTIONS = new Property<>("migrate.parallel.connections", 0); + + /** + * Allow live migration of a VM whose disk cache mode libvirt considers unsafe (anything + * other than none/directsync), by setting VIR_MIGRATE_UNSAFE. Safe only on coherent shared + * storage such as Ceph RBD, where writeback caching does not risk data loss on migration. + * Data type: Boolean.
+ * Default value: false. + */ + public static final Property MIGRATE_ALLOW_UNSAFE = new Property<>("migrate.allow.unsafe", false); + + /** + * Run a CPU-compatibility precheck (virsh cpu-compare against the destination host) before + * a live migration, so an incompatible destination fails fast with a clear message instead of a + * cryptic mid-migration libvirt error. Best-effort and fail-open: if the check cannot run it does + * not block the migration. Requires the source host's virsh to reach the destination libvirt. + * Data type: Boolean.
+ * Default value: false. + */ + public static final Property MIGRATE_CPU_PRECHECK_ENABLED = new Property<>("migrate.cpu.precheck.enabled", false); + + /** + * Migration memory-compression method, "xbzrle" or "mt". Empty (default) leaves libvirt's + * own default method in effect. Only takes effect when compression is enabled (libvirt >= 1.0.3, + * which sets VIR_MIGRATE_COMPRESSED) and is skipped under multifd. + * Data type: String.
+ * Default value: "" (empty). + */ + public static final Property MIGRATE_COMPRESSION_METHOD = new Property<>("migrate.compression.method", ""); + /** * MANDATORY: The UUID for the local storage pool.
* This property allows multiple values to be entered in a single String. The different values must be separated by commas.
diff --git a/api/src/main/java/com/cloud/host/Host.java b/api/src/main/java/com/cloud/host/Host.java index 79bec7a5871a..c794db436d73 100644 --- a/api/src/main/java/com/cloud/host/Host.java +++ b/api/src/main/java/com/cloud/host/Host.java @@ -66,6 +66,7 @@ public static String[] toStrings(Host.Type... types) { String HOST_VIRTV2V_VERSION = "host.virtv2v.version"; String HOST_SSH_PORT = "host.ssh.port"; String HOST_CDROM_MAX_COUNT = "host.cdrom.max.count"; + String HOST_MIGRATION_IP = "host.migration.ip"; String GUEST_OS_CATEGORY_ID = "guest.os.category.id"; String GUEST_OS_RULE = "guest.os.rule"; diff --git a/api/src/main/java/com/cloud/network/Networks.java b/api/src/main/java/com/cloud/network/Networks.java index 61a1c820723f..fbf2dccc769a 100644 --- a/api/src/main/java/com/cloud/network/Networks.java +++ b/api/src/main/java/com/cloud/network/Networks.java @@ -306,7 +306,7 @@ public static URI encodeStringIntoBroadcastUri(String candidate, BroadcastDomain * Different types of network traffic in the data center. */ public enum TrafficType { - None, Public, Guest, Storage, Management, Control, Vpn; + None, Public, Guest, Storage, Management, Control, Vpn, Migration; public static boolean isSystemNetwork(TrafficType trafficType) { if (Storage.equals(trafficType) || Management.equals(trafficType) || Control.equals(trafficType)) { @@ -322,6 +322,8 @@ public static TrafficType getTrafficType(String type) { return Guest; } else if ("Storage".equals(type)) { return Storage; + } else if ("Migration".equals(type)) { + return Migration; } else if ("Management".equals(type)) { return Management; } else if ("Control".equals(type)) { diff --git a/api/src/main/java/com/cloud/network/PhysicalNetworkSetupInfo.java b/api/src/main/java/com/cloud/network/PhysicalNetworkSetupInfo.java index fea3d9274fdf..858a2ec992b0 100644 --- a/api/src/main/java/com/cloud/network/PhysicalNetworkSetupInfo.java +++ b/api/src/main/java/com/cloud/network/PhysicalNetworkSetupInfo.java @@ -28,6 +28,7 @@ public class PhysicalNetworkSetupInfo { String publicNetworkName; String guestNetworkName; String storageNetworkName; + String migrationNetworkName; String mgmtVlan; public PhysicalNetworkSetupInfo() { @@ -49,6 +50,14 @@ public String getStorageNetworkName() { return storageNetworkName; } + public String getMigrationNetworkName() { + return migrationNetworkName; + } + + public void setMigrationNetworkName(String migrationNetworkName) { + this.migrationNetworkName = migrationNetworkName; + } + public void setPrivateNetworkName(String privateNetworkName) { this.privateNetworkName = privateNetworkName; } diff --git a/api/src/main/java/com/cloud/vm/VmDetailConstants.java b/api/src/main/java/com/cloud/vm/VmDetailConstants.java index 877df55c6d67..51c745117909 100644 --- a/api/src/main/java/com/cloud/vm/VmDetailConstants.java +++ b/api/src/main/java/com/cloud/vm/VmDetailConstants.java @@ -118,6 +118,9 @@ public interface VmDetailConstants { // CPU mode and model, ADMIN only String GUEST_CPU_MODE = "guest.cpu.mode"; String GUEST_CPU_MODEL = "guest.cpu.model"; + // Fallback policy for a custom CPU model ("allow" or "forbid"). "forbid" refuses a host that + // cannot provide the exact model instead of silently degrading it; set by the cluster CPU baseline. + String GUEST_CPU_MODEL_FALLBACK = "guest.cpu.model.fallback"; // Lease related String INSTANCE_LEASE_EXPIRY_DATE = "leaseexpirydate"; diff --git a/api/src/main/java/org/apache/cloudstack/api/response/HostResponse.java b/api/src/main/java/org/apache/cloudstack/api/response/HostResponse.java index 10bd62804fb2..d556a3ba6541 100644 --- a/api/src/main/java/org/apache/cloudstack/api/response/HostResponse.java +++ b/api/src/main/java/org/apache/cloudstack/api/response/HostResponse.java @@ -307,6 +307,10 @@ public class HostResponse extends BaseResponseWithAnnotations { @Param(description = "True if the host has capability to support UEFI boot") private Boolean uefiCapability; + @SerializedName("migrationip") + @Param(description = "the IP address the host uses for live migration traffic, set when a dedicated migration network is configured on this host", since = "24.0.0") + private String migrationIp; + @SerializedName(ApiConstants.ENCRYPTION_SUPPORTED) @Param(description = "True if the host supports encryption", since = "4.18") private Boolean encryptionSupported; @@ -896,6 +900,14 @@ public void setUefiCapability(Boolean hostCapability) { this.uefiCapability = hostCapability; } + public void setMigrationIp(String migrationIp) { + this.migrationIp = migrationIp; + } + + public String getMigrationIp() { + return migrationIp; + } + public void setEncryptionSupported(Boolean encryptionSupported) { this.encryptionSupported = encryptionSupported; } diff --git a/api/src/main/java/org/apache/cloudstack/query/QueryService.java b/api/src/main/java/org/apache/cloudstack/query/QueryService.java index b6362e9a9c9b..d1a6a0f85057 100644 --- a/api/src/main/java/org/apache/cloudstack/query/QueryService.java +++ b/api/src/main/java/org/apache/cloudstack/query/QueryService.java @@ -110,7 +110,7 @@ */ public interface QueryService { - List RootAdminOnlyVmSettings = Arrays.asList(VmDetailConstants.GUEST_CPU_MODE, VmDetailConstants.GUEST_CPU_MODEL); + List RootAdminOnlyVmSettings = Arrays.asList(VmDetailConstants.GUEST_CPU_MODE, VmDetailConstants.GUEST_CPU_MODEL, VmDetailConstants.GUEST_CPU_MODEL_FALLBACK); // Config keys ConfigKey AllowUserViewDestroyedVM = new ConfigKey<>("Advanced", Boolean.class, "allow.user.view.destroyed.vm", "false", diff --git a/core/src/main/java/com/cloud/agent/api/BaselineCpuCommand.java b/core/src/main/java/com/cloud/agent/api/BaselineCpuCommand.java new file mode 100644 index 000000000000..6eef21f46bcb --- /dev/null +++ b/core/src/main/java/com/cloud/agent/api/BaselineCpuCommand.java @@ -0,0 +1,47 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package com.cloud.agent.api; + +import java.util.List; + +/** + * Runs {@code virsh cpu-baseline} over the given per-host {@code } elements and returns the + * most feature-rich CPU compatible with all of them: the common-denominator cluster baseline. + */ +public class BaselineCpuCommand extends Command { + + private List hostCpuXmls; + + protected BaselineCpuCommand() { + } + + public BaselineCpuCommand(List hostCpuXmls) { + this.hostCpuXmls = hostCpuXmls; + } + + public List getHostCpuXmls() { + return hostCpuXmls; + } + + @Override + public boolean executeInSequence() { + return false; + } +} diff --git a/core/src/main/java/com/cloud/agent/api/CheckCpuCompatibilityCommand.java b/core/src/main/java/com/cloud/agent/api/CheckCpuCompatibilityCommand.java new file mode 100644 index 000000000000..fb0e1872f5df --- /dev/null +++ b/core/src/main/java/com/cloud/agent/api/CheckCpuCompatibilityCommand.java @@ -0,0 +1,50 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package com.cloud.agent.api; + +/** + * Checks on the agent, via {@code virsh cpu-compare}, that the host CPU can run the given {@code } model. + */ +public class CheckCpuCompatibilityCommand extends Command { + + private String cpuXml; + private String vmName; + + protected CheckCpuCompatibilityCommand() { + } + + public CheckCpuCompatibilityCommand(String vmName, String cpuXml) { + this.vmName = vmName; + this.cpuXml = cpuXml; + } + + public String getCpuXml() { + return cpuXml; + } + + public String getVmName() { + return vmName; + } + + @Override + public boolean executeInSequence() { + return false; + } +} diff --git a/core/src/main/java/com/cloud/agent/api/CheckNetworkAnswer.java b/core/src/main/java/com/cloud/agent/api/CheckNetworkAnswer.java index 392ad35e7cb6..b5ce7d2611aa 100644 --- a/core/src/main/java/com/cloud/agent/api/CheckNetworkAnswer.java +++ b/core/src/main/java/com/cloud/agent/api/CheckNetworkAnswer.java @@ -22,6 +22,8 @@ public class CheckNetworkAnswer extends Answer { // indicate if agent reconnect is needed after setupNetworkNames command private boolean _reconnect; + // the local IP the host resolved on its dedicated migration network, if one is configured + private String migrationIp; public CheckNetworkAnswer() { } @@ -39,4 +41,12 @@ public boolean needReconnect() { return _reconnect; } + public String getMigrationIp() { + return migrationIp; + } + + public void setMigrationIp(String migrationIp) { + this.migrationIp = migrationIp; + } + } diff --git a/core/src/main/java/com/cloud/agent/api/GetHostCpuModelCommand.java b/core/src/main/java/com/cloud/agent/api/GetHostCpuModelCommand.java new file mode 100644 index 000000000000..e11ebb088a36 --- /dev/null +++ b/core/src/main/java/com/cloud/agent/api/GetHostCpuModelCommand.java @@ -0,0 +1,35 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package com.cloud.agent.api; + +/** + * Asks the agent for this host's {@code } capabilities element, so the management server can + * compute a cluster CPU baseline that every host supports. + */ +public class GetHostCpuModelCommand extends Command { + + public GetHostCpuModelCommand() { + } + + @Override + public boolean executeInSequence() { + return false; + } +} diff --git a/core/src/main/java/com/cloud/agent/api/MigrateCommand.java b/core/src/main/java/com/cloud/agent/api/MigrateCommand.java index 7196247ffc23..62167dcab8f7 100644 --- a/core/src/main/java/com/cloud/agent/api/MigrateCommand.java +++ b/core/src/main/java/com/cloud/agent/api/MigrateCommand.java @@ -31,10 +31,12 @@ public class MigrateCommand extends Command { private String vmName; private String destinationIp; + private String migrateIp; private Map migrateStorage; private boolean migrateStorageManaged; private boolean migrateNonSharedInc; private boolean autoConvergence; + private String migrationEncryptionPolicy; private String hostGuid; private boolean windows; private VirtualMachineTO virtualMachine; @@ -98,6 +100,14 @@ public void setMigrateNonSharedInc(boolean migrateNonSharedInc) { this.migrateNonSharedInc = migrateNonSharedInc; } + public String getMigrationEncryptionPolicy() { + return migrationEncryptionPolicy; + } + + public void setMigrationEncryptionPolicy(String migrationEncryptionPolicy) { + this.migrationEncryptionPolicy = migrationEncryptionPolicy; + } + public void setAutoConvergence(boolean autoConvergence) { this.autoConvergence = autoConvergence; } @@ -118,6 +128,14 @@ public String getDestinationIp() { return destinationIp; } + public void setMigrateIp(String migrateIp) { + this.migrateIp = migrateIp; + } + + public String getMigrateIp() { + return migrateIp; + } + public String getVmName() { return vmName; } diff --git a/engine/api/src/main/java/com/cloud/vm/VirtualMachineManager.java b/engine/api/src/main/java/com/cloud/vm/VirtualMachineManager.java index c7238a513693..315d1ffdde74 100644 --- a/engine/api/src/main/java/com/cloud/vm/VirtualMachineManager.java +++ b/engine/api/src/main/java/com/cloud/vm/VirtualMachineManager.java @@ -67,6 +67,36 @@ public interface VirtualMachineManager extends Manager { ConfigKey VmConfigDriveLabel = new ConfigKey<>("Hidden", String.class, "vm.configdrive.label", "config-2", "The default label name for the config drive", false); + ConfigKey VmMigrationEncryptionPolicy = new ConfigKey<>("Advanced", String.class, "vm.migrate.encryption.policy", "Disabled", + "Policy for encrypting the KVM live-migration data stream. Disabled: plaintext. Required: use TLS " + + "(the host must have a libvirt/QEMU migration TLS environment configured, otherwise the migration fails). " + + "Overrides the per-host migrate.encryption.policy agent property when set to Required.", + true, ConfigKey.Scope.Zone); + + // Config key name for the per-cluster CPU baseline, defined here so other managers can reference it + // without depending on the implementation. + String CLUSTER_CPU_BASELINE_MODEL_KEY = "cluster.cpu.baseline.model"; + + /** + * Returns the names of the hosts in the cluster that are NOT compatible with the given CPU model (empty when + * all are, or when the model is blank). Only reachable KVM hosts are checked; a host that cannot be reached or + * verified is logged and skipped rather than reported incompatible. + */ + List findHostsIncompatibleWithCpuModel(long clusterId, String cpuModel); + + /** + * The configured per-cluster CPU baseline model (empty when none), used to pin instances and to validate hosts. + */ + String getClusterCpuBaselineModel(long clusterId); + + /** + * Computes a CPU baseline model for the cluster as the common denominator of its reachable KVM hosts, by + * collecting each host's CPU and running cpu-baseline. Returns the computed model name, or null when no + * reachable host returned a CPU definition or the cpu-baseline computation failed. Used to resolve a + * baseline of "auto" to a concrete, committed model. + */ + String computeClusterCpuBaseline(long clusterId); + ConfigKey VmConfigDriveOnPrimaryPool = new ConfigKey<>("Advanced", Boolean.class, "vm.configdrive.primarypool.enabled", "false", "If config drive need to be created and hosted on primary storage pool. Currently only supported for KVM.", true, ConfigKey.Scope.Zone); diff --git a/engine/orchestration/src/main/java/com/cloud/vm/VirtualMachineManagerImpl.java b/engine/orchestration/src/main/java/com/cloud/vm/VirtualMachineManagerImpl.java index 8c9a1653a260..a79a5ec1579e 100755 --- a/engine/orchestration/src/main/java/com/cloud/vm/VirtualMachineManagerImpl.java +++ b/engine/orchestration/src/main/java/com/cloud/vm/VirtualMachineManagerImpl.java @@ -44,6 +44,8 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.regex.Matcher; +import java.util.regex.Pattern; import java.util.stream.Collectors; import javax.inject.Inject; @@ -116,6 +118,10 @@ import com.cloud.agent.api.AgentControlAnswer; import com.cloud.agent.api.AgentControlCommand; import com.cloud.agent.api.Answer; +import com.cloud.agent.api.BaselineCpuCommand; +import com.cloud.agent.api.CheckCpuCompatibilityCommand; +import com.cloud.agent.api.GetHostCpuModelCommand; +import com.cloud.agent.api.UnsupportedAnswer; import com.cloud.agent.api.AttachOrDettachConfigDriveCommand; import com.cloud.agent.api.CheckVirtualMachineAnswer; import com.cloud.agent.api.CheckVirtualMachineCommand; @@ -217,6 +223,7 @@ import com.cloud.exception.StorageUnavailableException; import com.cloud.ha.HighAvailabilityManager; import com.cloud.ha.HighAvailabilityManager.WorkType; +import com.cloud.host.DetailVO; import com.cloud.host.Host; import com.cloud.host.HostVO; import com.cloud.host.Status; @@ -516,6 +523,15 @@ public class VirtualMachineManagerImpl extends ManagerBase implements VirtualMac Long.class, "systemvm.root.disk.size", "-1", "Size of root volume (in GB) of system VMs and virtual routers", true); + public static final ConfigKey ClusterCpuBaselineModel = new ConfigKey("Advanced", + String.class, CLUSTER_CPU_BASELINE_MODEL_KEY, "", + "Named CPU model (e.g. Haswell-noTSX, Skylake-Server) that user KVM instances in this cluster are pinned to when started, " + + "so they stay live-migratable across hosts with different CPU generations. Empty keeps the per-host guest.cpu.mode. " + + "Set to 'auto' to compute the common-denominator model of the cluster's current hosts and store that; the computed " + + "model is persisted (it does not re-compute as hosts change). Every host in the cluster must support the model; a " + + "host that does not is rejected by the CPU compatibility check.", + true, ConfigKey.Scope.Cluster); + private boolean syncTransitioningVmPowerState; ScheduledExecutorService _executor = null; @@ -1554,6 +1570,7 @@ public void orchestrateStart(final String vmUuid, final Map sshAccessDetails = _networkMgr.getSystemVMAccessDetails(vm); @@ -3397,8 +3414,19 @@ protected MigrateCommand buildMigrateCommand(VMInstanceVO vmInstance, VirtualMac logger.debug("Setting auto convergence to: {}", StorageManager.KvmAutoConvergence.value()); migrateCommand.setAutoConvergence(StorageManager.KvmAutoConvergence.value()); + // Thread the MS-central migration-encryption policy to the agent; it overrides the per-host property + // only when set to Required, while Disabled or blank defers to the per-host migrate.encryption.policy. + migrateCommand.setMigrationEncryptionPolicy(VmMigrationEncryptionPolicy.valueIn(vmInstance.getDataCenterId())); migrateCommand.setHostGuid(destination.getHost().getGuid()); + String migrationIp = resolveMigrationIp(vmInstance, destination.getHost()); + if (org.apache.commons.lang3.StringUtils.isNotBlank(migrationIp)) { + logger.debug("Live migration of VM [{}] will use the dedicated migration IP [{}] on destination host [{}] for the data stream.", + vmInstance, migrationIp, destination.getHost()); + migrateCommand.setMigrateIp(migrationIp); + } + preflightMigrationNetwork(vmInstance, destination.getHost(), migrationIp); + PrepareForMigrationAnswer prepareForMigrationAnswer = (PrepareForMigrationAnswer) answer; Map answerDpdkInterfaceMapping = prepareForMigrationAnswer.getDpdkInterfaceMapping(); @@ -3418,6 +3446,189 @@ protected MigrateCommand buildMigrateCommand(VMInstanceVO vmInstance, VirtualMac return migrateCommand; } + /** + * Resolves the IP that should carry the live-migration data stream to the destination host, or null to keep + * using the host management IP (the pre-existing behaviour). When a Migration traffic type is designated on the + * zone's physical network, the KVM agent resolves the local IP of that labelled NIC and the management server + * records it as the {@link Host#HOST_MIGRATION_IP} host detail; when that is present the data stream targets it + * instead of the management IP. Only KVM is handled here; VMware and XenServer select the migration network + * inside their own platform. + */ + protected String resolveMigrationIp(VMInstanceVO vmInstance, Host destinationHost) { + if (!HypervisorType.KVM.equals(vmInstance.getHypervisorType()) || vmInstance.getHostId() == null) { + return null; + } + // Use the dedicated migration network only when BOTH hosts have one. The source qemu dials the + // destination's migration IP, so if the source has no interface on that network the migration would fail; + // when either side lacks it, fall back to the management IP (the pre-existing behaviour). + if (org.apache.commons.lang3.StringUtils.isBlank(hostMigrationIp(vmInstance.getHostId()))) { + return null; + } + return org.apache.commons.lang3.StringUtils.trimToNull(hostMigrationIp(destinationHost.getId())); + } + + private String hostMigrationIp(long hostId) { + DetailVO detail = hostDetailsDao.findDetail(hostId, Host.HOST_MIGRATION_IP); + return detail == null ? null : detail.getValue(); + } + + /** + * Pre-flight check that warns when the source host of a KVM live migration uses a dedicated migration network + * but the destination host has none configured. The migration still proceeds over the management network (see + * {@link #resolveMigrationIp}), but without this the fallback is silent; this makes it visible so the operator + * can finish wiring the migration network on the destination. + */ + protected void preflightMigrationNetwork(VMInstanceVO vmInstance, Host destinationHost, String resolvedMigrationIp) { + // The dedicated migration network is already in use when a migration IP was resolved; only the fallback + // case (resolved IP blank) is worth a warning, so reuse the caller's result instead of re-resolving. + if (!HypervisorType.KVM.equals(vmInstance.getHypervisorType()) || vmInstance.getHostId() == null + || org.apache.commons.lang3.StringUtils.isNotBlank(resolvedMigrationIp)) { + return; + } + DetailVO sourceDetail = hostDetailsDao.findDetail(vmInstance.getHostId(), Host.HOST_MIGRATION_IP); + boolean sourceUsesMigrationNetwork = sourceDetail != null && org.apache.commons.lang3.StringUtils.isNotBlank(sourceDetail.getValue()); + if (sourceUsesMigrationNetwork) { + logger.warn("Live migration of VM [{}] has a dedicated migration network on source host [{}] but destination host [{}] has none configured; " + + "the data stream will fall back to the management network. Designate the Migration traffic type on the destination's zone physical network to keep it off the management NIC.", + vmInstance, vmInstance.getHostId(), destinationHost); + } + } + + /** + * Pins a KVM instance to the destination cluster's CPU baseline model ({@link #ClusterCpuBaselineModel}) at start, + * so it presents the same CPU to the guest on every host in the cluster and stays live-migratable across mixed CPU + * generations. The model is injected as the VM's {@code guest.cpu.mode=custom} / {@code guest.cpu.model} details, + * which the KVM agent already honours over its per-host default. Only user instances are pinned; system VMs and + * virtual routers are left on their per-host CPU. An explicit per-VM CPU model is left untouched, a blank cluster + * baseline is a no-op (pre-existing behaviour), and only KVM is affected. + */ + protected void applyClusterCpuBaseline(VirtualMachineTO vmTO, VMInstanceVO vm, DeployDestination dest) { + if (!VirtualMachine.Type.User.equals(vm.getType())) { + return; + } + if (!HypervisorType.KVM.equals(vm.getHypervisorType()) || dest == null || dest.getHost() == null || dest.getHost().getClusterId() == null) { + return; + } + String baselineModel = getClusterCpuBaselineModel(dest.getHost().getClusterId()); + if (org.apache.commons.lang3.StringUtils.isBlank(baselineModel)) { + return; + } + Map details = vmTO.getDetails(); + if (details != null && (org.apache.commons.lang3.StringUtils.isNotBlank(details.get(VmDetailConstants.GUEST_CPU_MODEL)) + || org.apache.commons.lang3.StringUtils.isNotBlank(details.get(VmDetailConstants.GUEST_CPU_MODE)))) { + logger.debug("VM [{}] already has an explicit CPU mode/model; leaving the cluster CPU baseline [{}] unapplied.", vm, baselineModel); + return; + } + Map newDetails = details == null ? new HashMap<>() : new HashMap<>(details); + newDetails.put(VmDetailConstants.GUEST_CPU_MODE, "custom"); + newDetails.put(VmDetailConstants.GUEST_CPU_MODEL, baselineModel); + newDetails.put(VmDetailConstants.GUEST_CPU_MODEL_FALLBACK, "forbid"); + vmTO.setDetails(newDetails); + logger.debug("Pinning VM [{}] to cluster [{}] CPU baseline model [{}] for live-migration compatibility.", + vm, dest.getHost().getClusterId(), baselineModel); + } + + @Override + public String getClusterCpuBaselineModel(long clusterId) { + return ClusterCpuBaselineModel.valueIn(clusterId); + } + + protected String buildCpuModelXml(String cpuModel) { + return String.format("%s", cpuModel); + } + + @Override + public List findHostsIncompatibleWithCpuModel(long clusterId, String cpuModel) { + List incompatible = new ArrayList<>(); + if (org.apache.commons.lang3.StringUtils.isBlank(cpuModel)) { + return incompatible; + } + String cpuXml = buildCpuModelXml(cpuModel); + for (HostVO host : _hostDao.findByClusterId(clusterId, Host.Type.Routing)) { + if (host.getStatus() != Status.Up || !HypervisorType.KVM.equals(host.getHypervisorType())) { + continue; + } + try { + Answer answer = _agentMgr.send(host.getId(), new CheckCpuCompatibilityCommand("cpu-baseline-check", cpuXml)); + // During a rolling upgrade a pre-feature agent cannot deserialize this new command, so the send + // times out and the host is skipped by the OperationTimedoutException catch below. If an agent + // does return an UnsupportedAnswer, treat it the same way: "cannot verify", not "incompatible". + if (answer instanceof UnsupportedAnswer) { + logger.warn("Host [{}] could not run the CPU compatibility check; skipping it for baseline validation.", host); + continue; + } + // Report a host only on an explicit incompatible verdict or an explicitly unknown model (virsh + // prints "Unknown CPU model ..."), so a non-existent model is rejected. A reachable host whose + // check could not run (e.g. libvirtd momentarily down) stays unverified and is skipped, not + // reported, matching this method's contract; the forbid pin and migrate-time compare still guard it. + String details = answer == null ? null : org.apache.commons.lang3.StringUtils.lowerCase(answer.getDetails()); + boolean unknownModel = details != null && details.contains("unknown"); + if (answer != null && (!answer.getResult() || unknownModel)) { + incompatible.add(host.getName()); + } + } catch (AgentUnavailableException | OperationTimedoutException e) { + logger.warn("Could not verify host [{}] against the cluster CPU baseline [{}]; skipping it: {}", host, cpuModel, e.getMessage()); + } + } + return incompatible; + } + + @Override + public String computeClusterCpuBaseline(long clusterId) { + List hostCpuXmls = new ArrayList<>(); + Long computeHostId = null; + for (HostVO host : _hostDao.findByClusterId(clusterId, Host.Type.Routing)) { + if (host.getStatus() != Status.Up || !HypervisorType.KVM.equals(host.getHypervisorType())) { + continue; + } + try { + // A pre-feature agent cannot deserialize this new command, so the send times out and the host is + // skipped by the catch below; an explicit UnsupportedAnswer is handled the same way. + Answer answer = _agentMgr.send(host.getId(), new GetHostCpuModelCommand()); + if (answer instanceof UnsupportedAnswer) { + logger.warn("Host [{}] could not report its CPU model; skipping it for baseline computation.", host); + continue; + } + if (answer != null && answer.getResult() && org.apache.commons.lang3.StringUtils.isNotBlank(answer.getDetails())) { + hostCpuXmls.add(answer.getDetails()); + if (computeHostId == null) { + computeHostId = host.getId(); + } + } + } catch (AgentUnavailableException | OperationTimedoutException e) { + logger.warn("Could not read the CPU model of host [{}] for baseline computation; skipping it: {}", host, e.getMessage()); + } + } + if (hostCpuXmls.isEmpty() || computeHostId == null) { + logger.warn("Could not compute a CPU baseline for cluster [{}]: no reachable KVM host returned a CPU definition.", clusterId); + return null; + } + try { + Answer answer = _agentMgr.send(computeHostId, new BaselineCpuCommand(hostCpuXmls)); + if (answer == null || !answer.getResult()) { + logger.warn("cpu-baseline computation failed for cluster [{}]: {}", clusterId, answer == null ? "no answer" : answer.getDetails()); + return null; + } + String model = parseModelFromCpuXml(answer.getDetails()); + logger.debug("Computed CPU baseline model [{}] for cluster [{}] from {} host(s).", model, clusterId, hostCpuXmls.size()); + return model; + } catch (AgentUnavailableException | OperationTimedoutException e) { + logger.warn("Could not compute the CPU baseline for cluster [{}]: {}", clusterId, e.getMessage()); + return null; + } + } + + protected String parseModelFromCpuXml(String cpuXml) { + if (org.apache.commons.lang3.StringUtils.isBlank(cpuXml)) { + return null; + } + Matcher matcher = Pattern.compile("]*>([^<]+)").matcher(cpuXml); + if (matcher.find()) { + return matcher.group(1).trim(); + } + return null; + } + private void updateVmPod(VMInstanceVO vm, long dstHostId) { // update the VMs pod HostVO host = _hostDao.findById(dstHostId); @@ -5400,7 +5611,7 @@ public ConfigKey[] getConfigKeys() { VmConfigDriveLabel, VmConfigDriveOnPrimaryPool, VmConfigDriveForceHostCacheUse, VmConfigDriveUseHostCacheOnUnsupportedPool, HaVmRestartHostUp, ResourceCountRunningVMsonly, AllowExposeHypervisorHostname, AllowExposeHypervisorHostnameAccountLevel, SystemVmRootDiskSize, AllowExposeDomainInMetadata, MetadataCustomCloudName, VmMetadataManufacturer, VmMetadataProductName, - VmSyncPowerStateTransitioning, SystemVmEnableUserData + VmSyncPowerStateTransitioning, SystemVmEnableUserData, VmMigrationEncryptionPolicy, ClusterCpuBaselineModel }; } diff --git a/engine/orchestration/src/main/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestrator.java b/engine/orchestration/src/main/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestrator.java index 8af75562b31c..6343e0778ae3 100644 --- a/engine/orchestration/src/main/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestrator.java +++ b/engine/orchestration/src/main/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestrator.java @@ -124,10 +124,12 @@ import com.cloud.exception.ResourceAllocationException; import com.cloud.exception.ResourceUnavailableException; import com.cloud.exception.UnsupportedServiceException; +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.hypervisor.Hypervisor.HypervisorType; import com.cloud.network.IpAddress; import com.cloud.network.IpAddressManager; @@ -445,6 +447,8 @@ public void setDhcpProviders(final List dhcpProviders) { @Inject HostDao _hostDao; @Inject + HostDetailsDao _hostDetailsDao; + @Inject NetworkServiceMapDao _ntwkSrvcDao; @Inject VpcManager _vpcMgr; @@ -4469,6 +4473,7 @@ public void processConnect(final Host host, final StartupCommand cmd, final bool final String privateName = _pNTrafficTypeDao.getNetworkTag(pNtwk.getId(), TrafficType.Management, hypervisorType); final String guestName = _pNTrafficTypeDao.getNetworkTag(pNtwk.getId(), TrafficType.Guest, hypervisorType); final String storageName = _pNTrafficTypeDao.getNetworkTag(pNtwk.getId(), TrafficType.Storage, hypervisorType); + final String migrationName = _pNTrafficTypeDao.getNetworkTag(pNtwk.getId(), TrafficType.Migration, hypervisorType); // String controlName = _pNTrafficTypeDao._networkModel.getNetworkTag(pNtwk.getId(), TrafficType.Control, hypervisorType); final PhysicalNetworkSetupInfo info = new PhysicalNetworkSetupInfo(); info.setPhysicalNetworkId(pNtwk.getId()); @@ -4476,6 +4481,7 @@ public void processConnect(final Host host, final StartupCommand cmd, final bool info.setPrivateNetworkName(privateName); info.setPublicNetworkName(publicName); info.setStorageNetworkName(storageName); + info.setMigrationNetworkName(migrationName); final PhysicalNetworkTrafficTypeVO mgmtTraffic = _pNTrafficTypeDao.findBy(pNtwk.getId(), TrafficType.Management); if (mgmtTraffic != null) { final String vlan = mgmtTraffic.getVlan(); @@ -4504,11 +4510,36 @@ public void processConnect(final Host host, final StartupCommand cmd, final bool if (answer.needReconnect()) { throw new ConnectionException(false, "Reinitialize agent after network setup."); } + persistMigrationIp(host, answer.getMigrationIp()); logger.debug("Network setup is correct on Agent"); return; } } + /** + * Records the IP the host resolved on its dedicated migration network as the {@link Host#HOST_MIGRATION_IP} + * host detail, or clears it when the host no longer resolves one, so that a migration network added or removed + * at the zone level takes effect on the next host connect without any per host configuration. + */ + protected void persistMigrationIp(final Host host, final String migrationIp) { + // The migration IP is an informational host detail, so a failure to record it must never fail the + // host connect (processConnect runs on the connect critical path); log it and carry on. + try { + final DetailVO existing = _hostDetailsDao.findDetail(host.getId(), Host.HOST_MIGRATION_IP); + if (StringUtils.isNotBlank(migrationIp)) { + if (existing == null || !migrationIp.equals(existing.getValue())) { + _hostDetailsDao.persist(host.getId(), Collections.singletonMap(Host.HOST_MIGRATION_IP, migrationIp)); + logger.debug("Host {} will use {} for live migration traffic", host, migrationIp); + } + } else if (existing != null) { + _hostDetailsDao.remove(existing.getId()); + logger.debug("Host {} no longer has a dedicated migration network; live migration will use the management address", host); + } + } catch (final Exception e) { + logger.warn("Could not record the migration IP for host {}; live migration will use the management address until the next host connect", host, e); + } + } + @Override public boolean processDisconnect(final long agentId, final Status state) { return false; diff --git a/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuBaselineTest.java b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuBaselineTest.java new file mode 100644 index 000000000000..6d5c9944cf8b --- /dev/null +++ b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuBaselineTest.java @@ -0,0 +1,117 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +package com.cloud.vm; + +import java.util.HashMap; +import java.util.Map; + +import org.junit.Assert; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; + +import com.cloud.agent.api.to.VirtualMachineTO; +import com.cloud.deploy.DeployDestination; +import com.cloud.host.Host; +import com.cloud.hypervisor.Hypervisor.HypervisorType; + +@RunWith(MockitoJUnitRunner.class) +public class VirtualMachineManagerImplCpuBaselineTest { + + private final VirtualMachineManagerImpl vmm = Mockito.spy(new VirtualMachineManagerImpl()); + + private DeployDestination destInCluster(long clusterId) { + Host host = Mockito.mock(Host.class); + Mockito.when(host.getClusterId()).thenReturn(clusterId); + DeployDestination dest = Mockito.mock(DeployDestination.class); + Mockito.when(dest.getHost()).thenReturn(host); + return dest; + } + + private VMInstanceVO kvmVm() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getType()).thenReturn(VirtualMachine.Type.User); + Mockito.when(vm.getHypervisorType()).thenReturn(HypervisorType.KVM); + return vm; + } + + @Test + public void pinsBaselineModelWhenClusterHasOne() { + Mockito.doReturn("Haswell-noTSX").when(vmm).getClusterCpuBaselineModel(5L); + VirtualMachineTO vmTO = Mockito.mock(VirtualMachineTO.class); + Mockito.when(vmTO.getDetails()).thenReturn(null); + + vmm.applyClusterCpuBaseline(vmTO, kvmVm(), destInCluster(5L)); + + ArgumentCaptor> captor = ArgumentCaptor.forClass(Map.class); + Mockito.verify(vmTO).setDetails(captor.capture()); + Assert.assertEquals("custom", captor.getValue().get(VmDetailConstants.GUEST_CPU_MODE)); + Assert.assertEquals("Haswell-noTSX", captor.getValue().get(VmDetailConstants.GUEST_CPU_MODEL)); + Assert.assertEquals("forbid", captor.getValue().get(VmDetailConstants.GUEST_CPU_MODEL_FALLBACK)); + } + + @Test + public void doesNothingWhenClusterBaselineIsBlank() { + Mockito.doReturn("").when(vmm).getClusterCpuBaselineModel(5L); + VirtualMachineTO vmTO = Mockito.mock(VirtualMachineTO.class); + + vmm.applyClusterCpuBaseline(vmTO, kvmVm(), destInCluster(5L)); + + Mockito.verify(vmTO, Mockito.never()).setDetails(Mockito.anyMap()); + } + + @Test + public void doesNotOverrideAnExplicitVmCpuModel() { + Mockito.doReturn("Haswell-noTSX").when(vmm).getClusterCpuBaselineModel(5L); + VirtualMachineTO vmTO = Mockito.mock(VirtualMachineTO.class); + Map explicit = new HashMap<>(); + explicit.put(VmDetailConstants.GUEST_CPU_MODE, "host-passthrough"); + explicit.put(VmDetailConstants.GUEST_CPU_MODEL, "EPYC"); + Mockito.when(vmTO.getDetails()).thenReturn(explicit); + + vmm.applyClusterCpuBaseline(vmTO, kvmVm(), destInCluster(5L)); + + Mockito.verify(vmTO, Mockito.never()).setDetails(Mockito.anyMap()); + } + + @Test + public void ignoresNonKvmWithoutReadingClusterBaseline() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getType()).thenReturn(VirtualMachine.Type.User); + Mockito.when(vm.getHypervisorType()).thenReturn(HypervisorType.VMware); + VirtualMachineTO vmTO = Mockito.mock(VirtualMachineTO.class); + + vmm.applyClusterCpuBaseline(vmTO, vm, Mockito.mock(DeployDestination.class)); + + Mockito.verify(vmm, Mockito.never()).getClusterCpuBaselineModel(Mockito.anyLong()); + Mockito.verify(vmTO, Mockito.never()).setDetails(Mockito.anyMap()); + } + + @Test + public void ignoresSystemVmWithoutReadingClusterBaseline() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + Mockito.when(vm.getType()).thenReturn(VirtualMachine.Type.DomainRouter); + VirtualMachineTO vmTO = Mockito.mock(VirtualMachineTO.class); + + vmm.applyClusterCpuBaseline(vmTO, vm, Mockito.mock(DeployDestination.class)); + + Mockito.verify(vmm, Mockito.never()).getClusterCpuBaselineModel(Mockito.anyLong()); + Mockito.verify(vmTO, Mockito.never()).setDetails(Mockito.anyMap()); + } +} diff --git a/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuCompatTest.java b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuCompatTest.java new file mode 100644 index 000000000000..6fc023496f42 --- /dev/null +++ b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuCompatTest.java @@ -0,0 +1,126 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +package com.cloud.vm; + +import java.util.Arrays; +import java.util.List; + +import org.junit.Assert; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; + +import com.cloud.agent.AgentManager; +import com.cloud.agent.api.Answer; +import com.cloud.agent.api.CheckCpuCompatibilityCommand; +import com.cloud.host.HostVO; +import com.cloud.host.Status; +import com.cloud.host.dao.HostDao; +import com.cloud.hypervisor.Hypervisor.HypervisorType; + +@RunWith(MockitoJUnitRunner.class) +public class VirtualMachineManagerImplCpuCompatTest { + + @Mock + private HostDao _hostDao; + + @Mock + private AgentManager _agentMgr; + + @InjectMocks + private VirtualMachineManagerImpl vmm = new VirtualMachineManagerImpl(); + + private HostVO kvmHostUp(long id, String name) { + HostVO host = Mockito.mock(HostVO.class); + Mockito.when(host.getId()).thenReturn(id); + Mockito.when(host.getName()).thenReturn(name); + Mockito.when(host.getStatus()).thenReturn(Status.Up); + Mockito.when(host.getHypervisorType()).thenReturn(HypervisorType.KVM); + return host; + } + + @Test + public void blankModelReturnsEmptyWithoutQueryingHosts() { + List result = vmm.findHostsIncompatibleWithCpuModel(1L, " "); + Assert.assertTrue(result.isEmpty()); + Mockito.verifyNoInteractions(_hostDao); + } + + @Test + public void reportsHostsThatFailTheCompatibilityCheck() throws Exception { + HostVO good = kvmHostUp(10L, "host-good"); + HostVO bad = kvmHostUp(11L, "host-bad"); + Mockito.when(_hostDao.findByClusterId(2L, com.cloud.host.Host.Type.Routing)).thenReturn(Arrays.asList(good, bad)); + Mockito.when(_agentMgr.send(Mockito.eq(10L), Mockito.any(CheckCpuCompatibilityCommand.class))) + .thenReturn(new Answer(null, true, "compatible")); + Mockito.when(_agentMgr.send(Mockito.eq(11L), Mockito.any(CheckCpuCompatibilityCommand.class))) + .thenReturn(new Answer(null, false, "incompatible")); + + List result = vmm.findHostsIncompatibleWithCpuModel(2L, "Haswell-noTSX"); + + Assert.assertEquals(Arrays.asList("host-bad"), result); + } + + @Test + public void allCompatibleReturnsEmpty() throws Exception { + HostVO good = kvmHostUp(10L, "host-good"); + Mockito.when(_hostDao.findByClusterId(2L, com.cloud.host.Host.Type.Routing)).thenReturn(Arrays.asList(good)); + Mockito.when(_agentMgr.send(Mockito.eq(10L), Mockito.any(CheckCpuCompatibilityCommand.class))) + .thenReturn(new Answer(null, true, "compatible")); + + Assert.assertTrue(vmm.findHostsIncompatibleWithCpuModel(2L, "Haswell-noTSX").isEmpty()); + } + + @Test + public void unknownModelFailsClosed() throws Exception { + // virsh reports an unknown model as an error with result=true (fail-open); setting a baseline must still + // reject it, so a non-existent model name is never accepted as a cluster baseline. + HostVO host = kvmHostUp(10L, "host-a"); + Mockito.when(_hostDao.findByClusterId(2L, com.cloud.host.Host.Type.Routing)).thenReturn(Arrays.asList(host)); + Mockito.when(_agentMgr.send(Mockito.eq(10L), Mockito.any(CheckCpuCompatibilityCommand.class))) + .thenReturn(new Answer(null, true, "error: internal error: Unknown CPU model NotARealModel")); + + Assert.assertEquals(Arrays.asList("host-a"), vmm.findHostsIncompatibleWithCpuModel(2L, "NotARealModel")); + } + + @Test + public void preFeatureAgentUnsupportedAnswerIsSkipped() throws Exception { + // an old agent that cannot handle the command returns an UnsupportedAnswer; that is "cannot verify", + // not "incompatible", so the host is skipped rather than blocking the baseline during a rolling upgrade. + HostVO host = kvmHostUp(10L, "host-a"); + Mockito.when(_hostDao.findByClusterId(2L, com.cloud.host.Host.Type.Routing)).thenReturn(Arrays.asList(host)); + Mockito.when(_agentMgr.send(Mockito.eq(10L), Mockito.any(CheckCpuCompatibilityCommand.class))) + .thenReturn(new com.cloud.agent.api.UnsupportedAnswer(null, "unsupported command")); + + Assert.assertTrue(vmm.findHostsIncompatibleWithCpuModel(2L, "Haswell-noTSX").isEmpty()); + } + + @Test + public void reachableHostWithConnectionErrorIsSkipped() throws Exception { + // a reachable host whose libvirtd was momentarily down (virsh "failed to connect") cannot be verified + // and must be skipped, not reported incompatible, so one transient failure does not block the baseline. + HostVO host = kvmHostUp(10L, "host-a"); + Mockito.when(_hostDao.findByClusterId(2L, com.cloud.host.Host.Type.Routing)).thenReturn(Arrays.asList(host)); + Mockito.when(_agentMgr.send(Mockito.eq(10L), Mockito.any(CheckCpuCompatibilityCommand.class))) + .thenReturn(new Answer(null, true, "error: failed to connect to the hypervisor")); + + Assert.assertTrue(vmm.findHostsIncompatibleWithCpuModel(2L, "Haswell-noTSX").isEmpty()); + } +} diff --git a/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplMigrationIpTest.java b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplMigrationIpTest.java new file mode 100644 index 000000000000..3b24c808fcc7 --- /dev/null +++ b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplMigrationIpTest.java @@ -0,0 +1,128 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +package com.cloud.vm; + +import static org.mockito.Mockito.when; + +import org.junit.Assert; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; + +import com.cloud.host.DetailVO; +import com.cloud.host.Host; +import com.cloud.host.dao.HostDetailsDao; +import com.cloud.hypervisor.Hypervisor.HypervisorType; + +@RunWith(MockitoJUnitRunner.class) +public class VirtualMachineManagerImplMigrationIpTest { + + private static final long SRC = 7L; + private static final long DST = 5L; + + @Mock + private HostDetailsDao hostDetailsDao; + + @InjectMocks + private VirtualMachineManagerImpl virtualMachineManagerImpl = new VirtualMachineManagerImpl(); + + private VMInstanceVO kvmVm() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + when(vm.getHypervisorType()).thenReturn(HypervisorType.KVM); + when(vm.getHostId()).thenReturn(SRC); + return vm; + } + + private Host hostWithId(long id) { + Host host = Mockito.mock(Host.class); + when(host.getId()).thenReturn(id); + return host; + } + + private void stubMigrationIp(long hostId, String value) { + DetailVO detail = Mockito.mock(DetailVO.class); + when(detail.getValue()).thenReturn(value); + when(hostDetailsDao.findDetail(hostId, Host.HOST_MIGRATION_IP)).thenReturn(detail); + } + + @Test + public void usesDedicatedIpWhenBothHostsHaveMigrationNetwork() { + stubMigrationIp(SRC, "10.0.7.7"); + stubMigrationIp(DST, "192.0.2.17"); + Assert.assertEquals("192.0.2.17", virtualMachineManagerImpl.resolveMigrationIp(kvmVm(), hostWithId(DST))); + } + + @Test + public void returnsNullWhenSourceHasNoMigrationNetwork() { + when(hostDetailsDao.findDetail(SRC, Host.HOST_MIGRATION_IP)).thenReturn(null); + Assert.assertNull(virtualMachineManagerImpl.resolveMigrationIp(kvmVm(), hostWithId(DST))); + } + + @Test + public void returnsNullWhenDestinationHasNoMigrationNetwork() { + stubMigrationIp(SRC, "10.0.7.7"); + when(hostDetailsDao.findDetail(DST, Host.HOST_MIGRATION_IP)).thenReturn(null); + Assert.assertNull(virtualMachineManagerImpl.resolveMigrationIp(kvmVm(), hostWithId(DST))); + } + + @Test + public void returnsNullWhenDestinationMigrationIpBlank() { + stubMigrationIp(SRC, "10.0.7.7"); + stubMigrationIp(DST, " "); + Assert.assertNull(virtualMachineManagerImpl.resolveMigrationIp(kvmVm(), hostWithId(DST))); + } + + @Test + public void returnsNullForNonKvmWithoutQueryingHostDetails() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + when(vm.getHypervisorType()).thenReturn(HypervisorType.VMware); + + Assert.assertNull(virtualMachineManagerImpl.resolveMigrationIp(vm, Mockito.mock(Host.class))); + Mockito.verify(hostDetailsDao, Mockito.never()).findDetail(Mockito.anyLong(), Mockito.anyString()); + } + + @Test + public void preflightChecksSourceWhenResolvedIpBlankAndSourceUsesMigrationNetwork() { + stubMigrationIp(SRC, "10.0.7.7"); + + // resolved IP is blank (destination lacks a migration network): preflight checks only the source, + // reusing the caller's resolved IP rather than re-resolving (no destination lookup). + virtualMachineManagerImpl.preflightMigrationNetwork(kvmVm(), hostWithId(DST), ""); + + Mockito.verify(hostDetailsDao, Mockito.times(1)).findDetail(SRC, Host.HOST_MIGRATION_IP); + Mockito.verify(hostDetailsDao, Mockito.never()).findDetail(DST, Host.HOST_MIGRATION_IP); + } + + @Test + public void preflightSkipsWhenResolvedIpPresent() { + virtualMachineManagerImpl.preflightMigrationNetwork(kvmVm(), hostWithId(DST), "192.0.2.17"); + Mockito.verify(hostDetailsDao, Mockito.never()).findDetail(Mockito.anyLong(), Mockito.anyString()); + } + + @Test + public void preflightSkipsNonKvmWithoutQueryingHostDetails() { + VMInstanceVO vm = Mockito.mock(VMInstanceVO.class); + when(vm.getHypervisorType()).thenReturn(HypervisorType.VMware); + + virtualMachineManagerImpl.preflightMigrationNetwork(vm, Mockito.mock(Host.class), ""); + + Mockito.verify(hostDetailsDao, Mockito.never()).findDetail(Mockito.anyLong(), Mockito.anyString()); + } +} diff --git a/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplTest.java b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplTest.java index 3b680f0bc74a..5e4b6294d37c 100644 --- a/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplTest.java +++ b/engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplTest.java @@ -68,6 +68,7 @@ import org.apache.cloudstack.engine.subsystem.api.storage.VolumeDataFactory; import org.apache.cloudstack.engine.subsystem.api.storage.VolumeInfo; import org.apache.cloudstack.framework.config.ConfigKey; +import com.cloud.agent.api.MigrateCommand; import org.apache.cloudstack.framework.config.impl.ConfigDepotImpl; import org.apache.cloudstack.framework.extensions.dao.ExtensionDetailsDao; import org.apache.cloudstack.framework.extensions.manager.ExtensionsManager; @@ -2241,4 +2242,67 @@ public void testUpdateClvmLockHostForVmVolumes_MultipleClvmPools() throws Except verify(clvmPoolManagerMock, times(1)).setClvmLockHostId(3L, destHostId); } + @Test + public void vmMigrationEncryptionPolicyConfigKeyDefaultAndCommandRoundTrip() { + // MS-central migration-encryption policy defaults to Disabled and round-trips on MigrateCommand. + Assert.assertEquals("Disabled", VirtualMachineManager.VmMigrationEncryptionPolicy.defaultValue()); + MigrateCommand mc = new MigrateCommand("vm", "1.1.1.1", false, null, false); + mc.setMigrationEncryptionPolicy("Required"); + Assert.assertEquals("Required", mc.getMigrationEncryptionPolicy()); + } + + @Test + public void computeClusterCpuBaselineReturnsComputedModel() throws Exception { + HostVO host = mock(HostVO.class); + when(host.getStatus()).thenReturn(com.cloud.host.Status.Up); + when(host.getHypervisorType()).thenReturn(HypervisorType.KVM); + when(host.getId()).thenReturn(11L); + doReturn(Arrays.asList(host)).when(hostDaoMock).findByClusterId(5L, Host.Type.Routing); + + when(agentManagerMock.send(eq(11L), any(com.cloud.agent.api.GetHostCpuModelCommand.class))) + .thenReturn(new com.cloud.agent.api.Answer(null, true, "Skylake-Server")); + when(agentManagerMock.send(eq(11L), any(com.cloud.agent.api.BaselineCpuCommand.class))) + .thenReturn(new com.cloud.agent.api.Answer(null, true, "Haswell-noTSX")); + + Assert.assertEquals("Haswell-noTSX", virtualMachineManagerImpl.computeClusterCpuBaseline(5L)); + } + + @Test + public void computeClusterCpuBaselineReturnsNullWhenNoHostAnswers() throws Exception { + HostVO host = mock(HostVO.class); + when(host.getStatus()).thenReturn(com.cloud.host.Status.Up); + when(host.getHypervisorType()).thenReturn(HypervisorType.KVM); + when(host.getId()).thenReturn(11L); + doReturn(Arrays.asList(host)).when(hostDaoMock).findByClusterId(5L, Host.Type.Routing); + + when(agentManagerMock.send(eq(11L), any(com.cloud.agent.api.GetHostCpuModelCommand.class))) + .thenReturn(new com.cloud.agent.api.Answer(null, false, "error")); + + Assert.assertNull(virtualMachineManagerImpl.computeClusterCpuBaseline(5L)); + } + + @Test + public void computeClusterCpuBaselineSkipsPreFeatureAgent() throws Exception { + // an old agent returns an UnsupportedAnswer for the new command; it is skipped, so with it as the only + // host no CPU is collected and the computation returns null rather than failing. + HostVO host = mock(HostVO.class); + when(host.getStatus()).thenReturn(com.cloud.host.Status.Up); + when(host.getHypervisorType()).thenReturn(HypervisorType.KVM); + when(host.getId()).thenReturn(11L); + doReturn(Arrays.asList(host)).when(hostDaoMock).findByClusterId(5L, Host.Type.Routing); + + when(agentManagerMock.send(eq(11L), any(com.cloud.agent.api.GetHostCpuModelCommand.class))) + .thenReturn(new com.cloud.agent.api.UnsupportedAnswer(null, "unsupported command")); + + Assert.assertNull(virtualMachineManagerImpl.computeClusterCpuBaseline(5L)); + } + + @Test + public void parseModelFromCpuXmlExtractsModelIgnoringAttributes() { + Assert.assertEquals("Haswell-noTSX", + virtualMachineManagerImpl.parseModelFromCpuXml("Haswell-noTSX")); + Assert.assertNull(virtualMachineManagerImpl.parseModelFromCpuXml("error: command not found")); + Assert.assertNull(virtualMachineManagerImpl.parseModelFromCpuXml(null)); + } + } diff --git a/engine/orchestration/src/test/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestratorTest.java b/engine/orchestration/src/test/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestratorTest.java index 66f5b699cc46..cbe027f7a2cf 100644 --- a/engine/orchestration/src/test/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestratorTest.java +++ b/engine/orchestration/src/test/java/org/apache/cloudstack/engine/orchestration/NetworkOrchestratorTest.java @@ -54,6 +54,9 @@ import com.cloud.dc.VlanVO; import com.cloud.dc.dao.VlanDao; import com.cloud.deploy.DeployDestination; +import com.cloud.host.DetailVO; +import com.cloud.host.Host; +import com.cloud.host.dao.HostDetailsDao; import com.cloud.exception.InsufficientAddressCapacityException; import com.cloud.exception.InsufficientCapacityException; import com.cloud.exception.InvalidParameterValueException; @@ -1073,4 +1076,71 @@ public void getNetworkElementsIncludingExtensionsReturnsBaseListWhenExtensionHel assertNotNull(result); assertEquals(1, result.size()); } + + @Test + public void testPersistMigrationIpStoresResolvedIp() { + testOrchestrator._hostDetailsDao = mock(HostDetailsDao.class); + Host host = mock(Host.class); + when(host.getId()).thenReturn(42L); + when(testOrchestrator._hostDetailsDao.findDetail(42L, Host.HOST_MIGRATION_IP)).thenReturn(null); + + testOrchestrator.persistMigrationIp(host, "10.9.9.9"); + + verify(testOrchestrator._hostDetailsDao, times(1)).persist(42L, Collections.singletonMap(Host.HOST_MIGRATION_IP, "10.9.9.9")); + } + + @Test + public void testPersistMigrationIpSkipsWhenUnchanged() { + testOrchestrator._hostDetailsDao = mock(HostDetailsDao.class); + Host host = mock(Host.class); + when(host.getId()).thenReturn(42L); + DetailVO existing = mock(DetailVO.class); + when(existing.getValue()).thenReturn("10.9.9.9"); + when(testOrchestrator._hostDetailsDao.findDetail(42L, Host.HOST_MIGRATION_IP)).thenReturn(existing); + + testOrchestrator.persistMigrationIp(host, "10.9.9.9"); + + verify(testOrchestrator._hostDetailsDao, never()).persist(ArgumentMatchers.anyLong(), any()); + } + + @Test + public void testPersistMigrationIpClearsWhenBlank() { + testOrchestrator._hostDetailsDao = mock(HostDetailsDao.class); + Host host = mock(Host.class); + when(host.getId()).thenReturn(42L); + DetailVO existing = mock(DetailVO.class); + when(existing.getId()).thenReturn(7L); + when(testOrchestrator._hostDetailsDao.findDetail(42L, Host.HOST_MIGRATION_IP)).thenReturn(existing); + + testOrchestrator.persistMigrationIp(host, null); + + verify(testOrchestrator._hostDetailsDao, times(1)).remove(7L); + } + + @Test + public void testPersistMigrationIpUpdatesWhenChanged() { + testOrchestrator._hostDetailsDao = mock(HostDetailsDao.class); + Host host = mock(Host.class); + when(host.getId()).thenReturn(42L); + DetailVO existing = mock(DetailVO.class); + when(existing.getValue()).thenReturn("10.1.1.1"); + when(testOrchestrator._hostDetailsDao.findDetail(42L, Host.HOST_MIGRATION_IP)).thenReturn(existing); + + testOrchestrator.persistMigrationIp(host, "10.2.2.2"); + + verify(testOrchestrator._hostDetailsDao, times(1)).persist(42L, Collections.singletonMap(Host.HOST_MIGRATION_IP, "10.2.2.2")); + } + + @Test + public void testPersistMigrationIpSwallowsPersistErrors() { + testOrchestrator._hostDetailsDao = mock(HostDetailsDao.class); + Host host = mock(Host.class); + when(host.getId()).thenReturn(42L); + when(testOrchestrator._hostDetailsDao.findDetail(42L, Host.HOST_MIGRATION_IP)).thenReturn(null); + Mockito.doThrow(new RuntimeException("db down")).when(testOrchestrator._hostDetailsDao) + .persist(ArgumentMatchers.anyLong(), ArgumentMatchers.any()); + + // a failure recording the informational migration IP must not propagate out and fail the host connect + testOrchestrator.persistMigrationIp(host, "10.9.9.9"); + } } diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResource.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResource.java index 9009ec629ca3..07fc94900ee3 100644 --- a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResource.java +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResource.java @@ -1068,6 +1068,19 @@ public NetworkInterface getPublicNic() { return publicNic; } + public String resolveMigrationNetworkIp(String label) { + if (StringUtils.isBlank(label)) { + return null; + } + NetworkInterface migrateNic = NetUtils.getNetworkInterface(label); + String[] migrateNicParams = migrateNic == null ? null : NetUtils.getNetworkParams(migrateNic); + if (migrateNicParams != null && StringUtils.isNotBlank(migrateNicParams[0])) { + return migrateNicParams[0]; + } + LOGGER.warn("Migration network label [{}] could not be resolved to an IPv4 address on this host; live migration will use the management address.", label); + return null; + } + protected String getDefaultTungstenScriptsDir() { return TUNGSTEN_PATH; } @@ -3447,6 +3460,9 @@ private CpuModeDef createCpuModeDef(VirtualMachineTO vmTO, int vcpus) { String cpuModel = MapUtils.isNotEmpty(details) && details.get(VmDetailConstants.GUEST_CPU_MODEL) != null ? details.get(VmDetailConstants.GUEST_CPU_MODEL) : guestCpuModel; cmd.setMode(cpuMode); cmd.setModel(cpuModel); + if (MapUtils.isNotEmpty(details) && details.get(VmDetailConstants.GUEST_CPU_MODEL_FALLBACK) != null) { + cmd.setModelFallback(details.get(VmDetailConstants.GUEST_CPU_MODEL_FALLBACK)); + } cmd.setFeatures(cpuFeatures); int vCpusInDef = vmTO.getVcpuMaxLimit() == null ? vcpus : vmTO.getVcpuMaxLimit(); setCpuTopology(cmd, vCpusInDef, vmTO.getDetails()); diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDef.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDef.java index 439e4f663416..24f307502559 100644 --- a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDef.java +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDef.java @@ -1996,6 +1996,7 @@ public String toString() { public static class CpuModeDef { private String _mode; private String _model; + private String _modelFallback = "allow"; private List _features; private int _coresPerSocket = -1; private int _threadsPerCore = -1; @@ -2005,6 +2006,12 @@ public void setMode(String mode) { _mode = mode; } + public void setModelFallback(String modelFallback) { + if ("allow".equals(modelFallback) || "forbid".equals(modelFallback)) { + _modelFallback = modelFallback; + } + } + public void setFeatures(List features) { if (features != null) { _features = features; @@ -2027,7 +2034,7 @@ public String toString() { // start cpu def, adding mode, model if ("custom".equalsIgnoreCase(_mode) && _model != null){ - modeBuilder.append("" + _model + ""); + modeBuilder.append("" + _model + ""); } else if ("host-model".equals(_mode)) { modeBuilder.append(""); } else if ("host-passthrough".equals(_mode)) { diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsync.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsync.java index 8f027e01ca4f..b49482706931 100644 --- a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsync.java +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsync.java @@ -22,11 +22,16 @@ import java.util.Set; import java.util.concurrent.Callable; +import org.apache.commons.lang3.StringUtils; + +import com.cloud.utils.exception.CloudRuntimeException; + import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.libvirt.Connect; import org.libvirt.Domain; import org.libvirt.LibvirtException; +import org.libvirt.TypedIntParameter; import org.libvirt.TypedParameter; import org.libvirt.TypedStringParameter; import org.libvirt.TypedUlongParameter; @@ -44,6 +49,12 @@ public class MigrateKVMAsync implements Callable { private boolean migrateStorage; private boolean migrateNonSharedInc; private boolean autoConvergence; + private boolean encryptMigration; + private boolean parallelMigration; + private int parallelConnections; + private boolean allowUnsafeMigration; + private String compressionMethod; + private String migrateListenAddress = null; protected Set migrateDiskLabels; @@ -96,8 +107,30 @@ public class MigrateKVMAsync implements Callable { // Libvirt 1.2.3 supports auto converge. private static final int LIBVIRT_VERSION_SUPPORTS_AUTO_CONVERGE = 1002003; + // Encrypt the migration connection using the TLS environment configured in qemu.conf + // (migrate_tls_x509_cert_dir / default_tls_x509_cert_dir). Without this, guest RAM, and full + // disk contents during storage migration, cross the network in plaintext TCP. + private static final long VIR_MIGRATE_TLS = 65536L; // 1 << 16 + + // Libvirt 3.2.0 supports VIR_MIGRATE_TLS. + private static final int LIBVIRT_VERSION_SUPPORTS_MIGRATE_TLS = 3002000; + + // Use multiple parallel network connections (multifd) to transfer memory. Without this, + // migration uses a single TCP stream and cannot fill a fast (25/40/100GbE) link. + private static final long VIR_MIGRATE_PARALLEL = 131072L; // 1 << 17 + + // Libvirt 5.2.0 supports VIR_MIGRATE_PARALLEL. + private static final int LIBVIRT_VERSION_SUPPORTS_PARALLEL = 5002000; + + // Migrate even if libvirt considers the migration unsafe (e.g. a disk cache mode other + // than none/directsync). On coherent shared storage such as Ceph RBD this is safe, but libvirt + // refuses such a migration unless this flag is set. + private static final long VIR_MIGRATE_UNSAFE = 512L; // 1 << 9 + public MigrateKVMAsync(final LibvirtComputingResource libvirtComputingResource, final Domain dm, final Connect dconn, final String dxml, - final boolean migrateStorage, final boolean migrateNonSharedInc, final boolean autoConvergence, final String vmName, final String destIp, Set migrateDiskLabels) { + final boolean migrateStorage, final boolean migrateNonSharedInc, final boolean autoConvergence, final boolean encryptMigration, + final boolean parallelMigration, final boolean allowUnsafeMigration, final String compressionMethod, + final int parallelConnections, final String vmName, final String destIp, final String migrateListenAddress, Set migrateDiskLabels) { this.libvirtComputingResource = libvirtComputingResource; this.dm = dm; @@ -106,16 +139,39 @@ public MigrateKVMAsync(final LibvirtComputingResource libvirtComputingResource, this.migrateStorage = migrateStorage; this.migrateNonSharedInc = migrateNonSharedInc; this.autoConvergence = autoConvergence; + this.encryptMigration = encryptMigration; + this.parallelMigration = parallelMigration; + this.parallelConnections = parallelConnections; + this.allowUnsafeMigration = allowUnsafeMigration; + this.compressionMethod = compressionMethod; this.vmName = vmName; this.destIp = destIp; + this.migrateListenAddress = migrateListenAddress; this.migrateDiskLabels = migrateDiskLabels; } @Override public Domain call() throws LibvirtException { + long flags = buildMigrateFlags(dconn.getLibVirVersion()); + + TypedParameter [] parameters = createTypedParameterList(dconn.getLibVirVersion()); + + logger.debug(String.format("Migrating [%s] with flags [%s], destination [%s] and speed [%s]. The disks with the following labels will be migrated [%s].", vmName, flags, + destIp, libvirtComputingResource.getMigrateSpeed(), migrateDiskLabels)); + + return dm.migrate(dconn, parameters, flags); + + } + + // extracted from call() so the flag computation (including the new VIR_MIGRATE_TLS) + // is unit-testable without a live libvirt connection. + protected long buildMigrateFlags(final long libvirtVersion) { long flags = VIR_MIGRATE_LIVE; - if (dconn.getLibVirVersion() >= LIBVIRT_VERSION_SUPPORTS_MIGRATE_COMPRESSED) { + // legacy compression (VIR_MIGRATE_COMPRESSED, which QEMU maps to xbzrle) is INCOMPATIBLE + // with multifd (VIR_MIGRATE_PARALLEL), QEMU refuses to combine them, so setting both would fail + // every parallel migration. Skip legacy compression whenever multifd is enabled. + if (libvirtVersion >= LIBVIRT_VERSION_SUPPORTS_MIGRATE_COMPRESSED && !parallelMigration) { flags |= VIR_MIGRATE_COMPRESSED; } @@ -130,30 +186,62 @@ public Domain call() throws LibvirtException { } } - if (autoConvergence && dconn.getLibVirVersion() >= LIBVIRT_VERSION_SUPPORTS_AUTO_CONVERGE) { + if (autoConvergence && libvirtVersion >= LIBVIRT_VERSION_SUPPORTS_AUTO_CONVERGE) { flags |= VIR_MIGRATE_AUTO_CONVERGE; } - TypedParameter [] parameters = createTypedParameterList(); + if (encryptMigration) { + if (libvirtVersion < LIBVIRT_VERSION_SUPPORTS_MIGRATE_TLS) { + throw new CloudRuntimeException(String.format( + "Live migration of %s requires encryption but libvirt %d does not support TLS migration (needs >= %d); failing instead of sending the memory stream in plaintext.", + vmName, libvirtVersion, LIBVIRT_VERSION_SUPPORTS_MIGRATE_TLS)); + } + flags |= VIR_MIGRATE_TLS; + } - logger.debug(String.format("Migrating [%s] with flags [%s], destination [%s] and speed [%s]. The disks with the following labels will be migrated [%s].", vmName, flags, - destIp, libvirtComputingResource.getMigrateSpeed(), migrateDiskLabels)); + if (parallelMigration && libvirtVersion >= LIBVIRT_VERSION_SUPPORTS_PARALLEL) { + flags |= VIR_MIGRATE_PARALLEL; + } - return dm.migrate(dconn, parameters, flags); + // permit migration of VMs libvirt deems unsafe (e.g. writeback disk cache) when the + // operator asserts the storage is coherent (Ceph RBD). Opt-in; off by default. + if (allowUnsafeMigration) { + flags |= VIR_MIGRATE_UNSAFE; + } + return flags; } - protected TypedParameter[] createTypedParameterList() { + protected TypedParameter[] createTypedParameterList(final long libvirtVersion) { int sizeOfMigrateDiskLabels = 0; if (migrateDiskLabels != null) { sizeOfMigrateDiskLabels = migrateDiskLabels.size(); } - TypedParameter[] parameters = new TypedParameter[4 + sizeOfMigrateDiskLabels]; + // Each tuning parameter must be gated on the same condition as the flag that activates it, or libvirt + // rejects the migration (e.g. the parallel-connections param without VIR_MIGRATE_PARALLEL). + final boolean hasCompressionMethod = StringUtils.isNotBlank(compressionMethod) && !parallelMigration + && libvirtVersion >= LIBVIRT_VERSION_SUPPORTS_MIGRATE_COMPRESSED; + final boolean bindMigrateListenAddress = StringUtils.isNotBlank(migrateListenAddress); + final boolean setParallelConnections = parallelMigration && parallelConnections > 0 + && libvirtVersion >= LIBVIRT_VERSION_SUPPORTS_PARALLEL; + final int fixedParams = 4 + (hasCompressionMethod ? 1 : 0) + (bindMigrateListenAddress ? 1 : 0) + (setParallelConnections ? 1 : 0); + + TypedParameter[] parameters = new TypedParameter[fixedParams + sizeOfMigrateDiskLabels]; parameters[0] = new TypedStringParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_DEST_NAME, vmName); parameters[1] = new TypedStringParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_DEST_XML, dxml); parameters[2] = new TypedStringParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_URI, "tcp:" + destIp); parameters[3] = new TypedUlongParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_BANDWIDTH, libvirtComputingResource.getMigrateSpeed()); + int nextParam = 4; + if (hasCompressionMethod) { + parameters[nextParam++] = new TypedStringParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_COMPRESSION, compressionMethod); + } + if (bindMigrateListenAddress) { + parameters[nextParam++] = new TypedStringParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_LISTEN_ADDRESS, migrateListenAddress); + } + if (setParallelConnections) { + parameters[nextParam++] = new TypedIntParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_PARALLEL_CONNECTIONS, parallelConnections); + } if (sizeOfMigrateDiskLabels == 0) { return parameters; @@ -161,7 +249,7 @@ protected TypedParameter[] createTypedParameterList() { Iterator iterator = migrateDiskLabels.iterator(); for (int i = 0; i < sizeOfMigrateDiskLabels; i++) { - parameters[4 + i] = new TypedStringParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_MIGRATE_DISKS, iterator.next()); + parameters[fixedParams + i] = new TypedStringParameter(Domain.DomainMigrateParameters.VIR_MIGRATE_PARAM_MIGRATE_DISKS, iterator.next()); } return parameters; diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapper.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapper.java new file mode 100644 index 000000000000..e23e6246c157 --- /dev/null +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapper.java @@ -0,0 +1,73 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package com.cloud.hypervisor.kvm.resource.wrapper; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; + +import org.apache.commons.collections.CollectionUtils; + +import com.cloud.agent.api.Answer; +import com.cloud.agent.api.BaselineCpuCommand; +import com.cloud.hypervisor.kvm.resource.LibvirtComputingResource; +import com.cloud.resource.CommandWrapper; +import com.cloud.resource.ResourceWrapper; +import com.cloud.utils.script.Script; + +@ResourceWrapper(handles = BaselineCpuCommand.class) +public final class LibvirtBaselineCpuCommandWrapper extends CommandWrapper { + + private static final int VIRSH_TIMEOUT_SECONDS = 120; + + @Override + public Answer execute(final BaselineCpuCommand command, final LibvirtComputingResource libvirtComputingResource) { + final List hostCpuXmls = command.getHostCpuXmls(); + if (CollectionUtils.isEmpty(hostCpuXmls)) { + return new Answer(command, false, "no host CPU definitions supplied"); + } + Path tempFile = null; + try { + tempFile = Files.createTempFile("cloudstack-cpu-baseline-", ".xml"); + Files.write(tempFile, String.join("\n", hostCpuXmls).getBytes(StandardCharsets.UTF_8)); + + // full result, not runSimpleBashScript: that keeps only the first line (the opening tag), + // so the would never be seen. + final String output = Script.runSimpleBashScriptWithFullResult("virsh cpu-baseline " + tempFile.toString() + " 2>&1 || true", VIRSH_TIMEOUT_SECONDS); + final boolean computed = output != null && output.contains(" { + + private static final int VIRSH_TIMEOUT_SECONDS = 120; + + @Override + public Answer execute(final CheckCpuCompatibilityCommand command, final LibvirtComputingResource libvirtComputingResource) { + Path tempFile = null; + try { + tempFile = Files.createTempFile("cloudstack-cpu-compare-", ".xml"); + Files.write(tempFile, command.getCpuXml() == null ? new byte[0] : command.getCpuXml().getBytes(StandardCharsets.UTF_8)); + + // full result (not runSimpleBashScript, which keeps only the first line) with a forced zero exit, + // so the full verdict is parsed: an incompatible host, and the "Unknown CPU model" detail that + // virsh prints on a later line for a bogus model, both reach the caller instead of being dropped. + final String output = Script.runSimpleBashScriptWithFullResult("virsh cpu-compare " + tempFile.toString() + " 2>&1 || true", VIRSH_TIMEOUT_SECONDS); + final boolean compatible = isCompatible(output); + logger.debug("CPU compatibility check for VM [{}] on this host returned compatible={}, output=[{}]", + command.getVmName(), compatible, output); + return new Answer(command, compatible, output); + } catch (IOException e) { + logger.warn("Failed to run CPU compatibility check for VM [{}]", command.getVmName(), e); + return new Answer(command, false, e.getMessage()); + } finally { + if (tempFile != null) { + try { + Files.deleteIfExists(tempFile); + } catch (IOException ignored) { + logger.trace("Failed to delete temp CPU XML file: {}", tempFile); + } + } + } + } + + // Incompatible when virsh reports "incompatible" (the verdict current libvirt emits); "not a superset" + // is also treated as incompatible defensively. Blank/unknown output fails open (compatible=true) so a + // missing or odd virsh does not block migration. + protected static boolean isCompatible(final String output) { + if (output == null || output.trim().isEmpty()) { + return true; + } + final String lower = output.toLowerCase(); + if (lower.contains("incompatible") || lower.contains("not a superset")) { + return false; + } + return true; + } +} diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtCheckNetworkCommandWrapper.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtCheckNetworkCommandWrapper.java index abf61935a044..0ff6a2a9a2f4 100644 --- a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtCheckNetworkCommandWrapper.java +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtCheckNetworkCommandWrapper.java @@ -37,6 +37,7 @@ public final class LibvirtCheckNetworkCommandWrapper extends CommandWrapper phyNics = command.getPhysicalNetworkInfoList(); String errMsg = null; + String migrationIp = null; for (final PhysicalNetworkSetupInfo nic : phyNics) { if (!libvirtComputingResource.checkNetwork(Networks.TrafficType.Guest, nic.getGuestNetworkName())) { @@ -49,12 +50,17 @@ public Answer execute(final CheckNetworkCommand command, final LibvirtComputingR errMsg = "Can not find network: " + nic.getPublicNetworkName(); break; } + if (migrationIp == null) { + migrationIp = libvirtComputingResource.resolveMigrationNetworkIp(nic.getMigrationNetworkName()); + } } if (errMsg != null) { return new CheckNetworkAnswer(command, false, errMsg); } else { - return new CheckNetworkAnswer(command, true, null); + final CheckNetworkAnswer answer = new CheckNetworkAnswer(command, true, null); + answer.setMigrationIp(migrationIp); + return answer; } } } diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtGetHostCpuModelCommandWrapper.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtGetHostCpuModelCommandWrapper.java new file mode 100644 index 000000000000..3b1db505470f --- /dev/null +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtGetHostCpuModelCommandWrapper.java @@ -0,0 +1,61 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package com.cloud.hypervisor.kvm.resource.wrapper; + +import org.apache.commons.lang3.StringUtils; + +import com.cloud.agent.api.Answer; +import com.cloud.agent.api.GetHostCpuModelCommand; +import com.cloud.hypervisor.kvm.resource.LibvirtComputingResource; +import com.cloud.resource.CommandWrapper; +import com.cloud.resource.ResourceWrapper; +import com.cloud.utils.script.Script; + +@ResourceWrapper(handles = GetHostCpuModelCommand.class) +public final class LibvirtGetHostCpuModelCommandWrapper extends CommandWrapper { + + private static final int VIRSH_TIMEOUT_SECONDS = 120; + + @Override + public Answer execute(final GetHostCpuModelCommand command, final LibvirtComputingResource libvirtComputingResource) { + // full result, not runSimpleBashScript: that keeps only the first line, which would be "". + final String caps = Script.runSimpleBashScriptWithFullResult("virsh capabilities 2>&1 || true", VIRSH_TIMEOUT_SECONDS); + final String hostCpu = extractHostCpu(caps); + if (StringUtils.isBlank(hostCpu)) { + logger.warn("Could not read the host element from virsh capabilities; output was [{}]", caps); + return new Answer(command, false, caps); + } + return new Answer(command, true, hostCpu); + } + + // The host element from 'virsh capabilities', or null if absent. + protected static String extractHostCpu(final String capabilitiesXml) { + if (StringUtils.isBlank(capabilitiesXml)) { + return null; + } + final int start = capabilitiesXml.indexOf(""); + final int end = capabilitiesXml.indexOf(""); + if (start < 0 || end < 0 || end < start) { + return null; + } + return capabilitiesXml.substring(start, end + "".length()); + } +} diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtMigrateCommandWrapper.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtMigrateCommandWrapper.java index ed02ae6da38d..e5b5ccd5599b 100644 --- a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtMigrateCommandWrapper.java +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtMigrateCommandWrapper.java @@ -21,6 +21,11 @@ import java.io.InputStream; import java.net.URISyntaxException; import java.nio.charset.StandardCharsets; +import java.io.File; +import java.nio.file.Files; +import java.util.regex.Matcher; +import java.util.regex.Pattern; +import com.cloud.utils.script.Script; import java.util.ArrayList; import java.util.HashSet; import java.util.List; @@ -112,6 +117,17 @@ public Answer execute(final MigrateCommand command, final LibvirtComputingResour logger.debug(String.format("Trying to migrate VM [%s] to destination host: [%s].", vmName, destinationUri)); } + // opt-in CPU-compatibility precheck (virsh cpu-compare) run BEFORE any migration setup, + // so an incompatible destination fails fast with a clear message and the source domain is never + // touched. Fail-open: if the check cannot run it does not block the migration. + if (Boolean.TRUE.equals(AgentPropertiesFileHandler.getPropertyValue(AgentProperties.MIGRATE_CPU_PRECHECK_ENABLED))) { + final String cpuError = precheckDestinationCpu(vmName, destinationUri, libvirtComputingResource); + if (cpuError != null) { + logger.warn(cpuError); + return new MigrateAnswer(command, false, cpuError, null); + } + } + String result = null; Command.State commandState = null; @@ -250,9 +266,32 @@ Use VIR_DOMAIN_XML_SECURE (value = 1) prior to v1.0.0. libvirtComputingResource.createOrUpdateLogFileForCommand(command, Command.State.PROCESSING); + // Encrypt the migration data stream when the effective policy resolves to "Required". A blank or + // "Disabled" MS policy defers to the per-host migrate.encryption.policy (the MS ConfigKey default + // is "Disabled", not blank), so the per-host setting is never dead code; requires + // migrate_tls_x509_* configured in qemu.conf. + final String hostEncryptionPolicy = AgentPropertiesFileHandler.getPropertyValue(AgentProperties.MIGRATE_ENCRYPTION_POLICY); + final boolean encryptMigration = resolveEncryptMigration(command.getMigrationEncryptionPolicy(), hostEncryptionPolicy); + final boolean parallelMigration = Boolean.TRUE.equals(AgentPropertiesFileHandler.getPropertyValue(AgentProperties.MIGRATE_PARALLEL_ENABLED)); + // allow libvirt-"unsafe" migrations (e.g. writeback cache on coherent Ceph storage). + final boolean allowUnsafeMigration = Boolean.TRUE.equals(AgentPropertiesFileHandler.getPropertyValue(AgentProperties.MIGRATE_ALLOW_UNSAFE)); + // optional migration compression method (xbzrle or mt); blank leaves libvirt's default. + final String compressionMethod = AgentPropertiesFileHandler.getPropertyValue(AgentProperties.MIGRATE_COMPRESSION_METHOD); + final int parallelConnections = AgentPropertiesFileHandler.getPropertyValue(AgentProperties.MIGRATE_PARALLEL_CONNECTIONS); + // when a dedicated migration network is configured, the management server resolves the destination + // host's migration-NIC IP and sets it on the command; route the data stream (URI + listen address) + // there instead of the management IP. The libvirt control connection stays on the management IP. + final String configuredMigrateIp = command.getMigrateIp(); + final boolean dedicatedMigrationNetwork = StringUtils.isNotBlank(configuredMigrateIp); + final String migrateDataIp = dedicatedMigrationNetwork ? configuredMigrateIp : command.getDestinationIp(); + if (dedicatedMigrationNetwork) { + logger.info("Live migration of VM {} will use dedicated migration address {} for the data stream instead of the management address {}.", + vmName, migrateDataIp, command.getDestinationIp()); + } final Callable worker = new MigrateKVMAsync(libvirtComputingResource, dm, dconn, xmlDesc, migrateStorage, migrateNonSharedInc, - command.isAutoConvergence(), vmName, command.getDestinationIp(), migrateDiskLabels); + command.isAutoConvergence(), encryptMigration, parallelMigration, allowUnsafeMigration, compressionMethod, parallelConnections, vmName, migrateDataIp, + dedicatedMigrationNetwork ? configuredMigrateIp : null, migrateDiskLabels); final Future migrateThread = executor.submit(worker); executor.shutdown(); long sleeptime = 0; @@ -280,7 +319,16 @@ Use VIR_DOMAIN_XML_SECURE (value = 1) prior to v1.0.0. } } if (sleeptime % 1000 == 0) { - logger.info("Waiting for migration of " + vmName + " to complete, waited " + sleeptime + "ms"); + // surface migration progress (percent of migration data transferred) in the periodic log. + int progressPercent = -1; + try { + final DomainJobInfo job = dm.getJobInfo(); + progressPercent = computeMigrationProgressPercent(job.getDataProcessed(), job.getDataRemaining()); + } catch (final LibvirtException e) { + logger.trace("Could not read migration job info for progress reporting: {}", e.getMessage()); + } + logger.info("Waiting for migration of {} to complete, waited {}ms, progress: {}", vmName, sleeptime, + progressPercent < 0 ? "unknown" : progressPercent + "%"); } // abort the vm migration if the job is executed more than vm.migrate.wait @@ -304,7 +352,12 @@ Use VIR_DOMAIN_XML_SECURE (value = 1) prior to v1.0.0. logger.debug(result); break; } catch (final LibvirtException e) { - logger.error(String.format("Failed to abort the VM migration job of VM [%s] due to: [%s].", vmName, e.getMessage()), e); + // Do NOT mark the migration failed here: abortJob throws both when the migration just + // completed (no active job left to abort) and on a transient error while it is still + // running. Log and let the loop continue - a completed migration ends the loop with + // destDomain != null and is reported successful (no split brain), while a still-running + // one is retried on the next pass and ultimately bounded by migrateThread.get below. + logger.warn(String.format("Could not abort the migration job of VM [%s] after the vm.migrate.wait timeout: %s", vmName, e.getMessage()), e); } } } @@ -337,7 +390,15 @@ Use VIR_DOMAIN_XML_SECURE (value = 1) prior to v1.0.0. if (logger.isDebugEnabled()) { logger.debug(String.format("Cleaning the disks of VM [%s] in the source pool after VM migration finished.", vmName)); } - resumeDomainIfPaused(destDomain, vmName); + // The guest is now on the destination and the source domain is about to be undefined, so the + // migration (the move) has succeeded. If it could not be brought out of PAUSED we must NOT report + // failure - that would make the management server believe the VM is still on the source and could + // trigger HA against a host that no longer runs it (split brain). Surface it loudly instead; the + // power-state sync and the operator can resume the paused guest on the destination. + final String resumeFailure = resumeDomainIfPaused(destDomain, vmName); + if (resumeFailure != null) { + logger.warn(String.format("VM [%s] migrated to the destination but is still PAUSED there: [%s]. The migration is reported as successful because the VM now lives on the destination; resume it on the destination host.", vmName, resumeFailure)); + } // For cross-pool CLVM migration, skip deactivation so the source LV stays // active (in shared mode) and deletion can route directly to the source host @@ -403,16 +464,7 @@ Use VIR_DOMAIN_XML_SECURE (value = 1) prior to v1.0.0. if (result == null) { logger.info("Post-migration cleanup for VM {}: ", vmName); - libvirtComputingResource.destroyNetworkRulesForVM(conn, vmName); - for (final InterfaceDef iface : ifaces) { - String vlanId = libvirtComputingResource.getVlanIdFromBridgeName(iface.getBrName()); - // We don't know which "traffic type" is associated with - // each interface at this point, so inform all vif drivers - final List allVifDrivers = libvirtComputingResource.getAllVifDrivers(); - for (final VifDriver vifDriver : allVifDrivers) { - vifDriver.unplug(iface, libvirtComputingResource.shouldDeleteBridge(vlanToPersistenceMap, vlanId)); - } - } + cleanupSourceNetworkingAfterMigration(conn, vmName, ifaces, vlanToPersistenceMap, libvirtComputingResource); commandState = Command.State.COMPLETED; libvirtComputingResource.createOrUpdateLogFileForCommand(command, commandState); } else if (commandState == null) { @@ -424,6 +476,110 @@ Use VIR_DOMAIN_XML_SECURE (value = 1) prior to v1.0.0. return new MigrateAnswer(command, result == null, result, null); } + // cap the CPU precheck so a firewalled/unreachable destination fails fast instead of + // hanging virsh forever (the opposite of the "fail fast" the precheck is meant to provide). + private static final int CPU_PRECHECK_TIMEOUT_SECONDS = 30; + private static final Pattern SAFE_MIGRATION_URI = Pattern.compile("qemu\\+(tcp|tls)://[A-Za-z0-9._:\\[\\]-]+/system"); + + /** + * effective migration-encryption decision. A non-Disabled MS-central policy wins; a blank or + * "Disabled" MS policy defers to the per-host agent property (so the per-host setting is never dead code). + */ + protected boolean resolveEncryptMigration(final String commandPolicy, final String hostPolicy) { + String effective = commandPolicy; + if (StringUtils.isBlank(effective) || "Disabled".equalsIgnoreCase(effective.trim())) { + effective = hostPolicy; + } + // Only the exact "Required" policy enables TLS; anything else (Disabled, blank or an unrecognised value) + // stays plaintext, so a typo never silently turns encryption on. libvirt has no opportunistic/fallback + // TLS mode, so Required means TLS or the migration fails - there is no partial state to model here. + return effective != null && "Required".equalsIgnoreCase(effective.trim()); + } + + // matches the domain's top-level element (paired or self-closing). \b after "cpu" + // excludes ; non-greedy .*? is safe because does not nest another . + private static final Pattern CPU_ELEMENT_PATTERN = Pattern.compile("]*/>|", Pattern.DOTALL); + + /** + * pull the VM's <cpu> definition out of its domain XML so it can be checked against a + * destination host with {@code virsh cpu-compare}. Returns null when the VM has no explicit CPU + * model (any host can then run it). + */ + protected String extractCpuElement(final String domainXml) { + if (domainXml == null) { + return null; + } + final Matcher matcher = CPU_ELEMENT_PATTERN.matcher(domainXml); + return matcher.find() ? matcher.group() : null; + } + + /** + * interpret {@code virsh cpu-compare} output. libvirt prints "incompatible" / "not a + * superset" when the host cannot run the given CPU. Fail-open: blank/unknown output is treated as + * compatible, so a check that could not run never blocks a migration. + */ + protected boolean isCpuCompareOutputCompatible(final String virshOutput) { + return LibvirtCheckCpuCompatibilityCommandWrapper.isCompatible(virshOutput); + } + + protected String runDestinationCpuCompare(final String cpuXml, final String destinationUri) throws IOException { + if (destinationUri == null || !SAFE_MIGRATION_URI.matcher(destinationUri).matches()) { + throw new IOException("Unexpected destination URI for the CPU precheck: " + destinationUri); + } + final File tmp = File.createTempFile("cloudstack-cpucheck-", ".xml"); + try { + Files.write(tmp.toPath(), cpuXml.getBytes(StandardCharsets.UTF_8)); + // full result with a forced zero exit so the complete verdict is parsed, not just virsh's first line. + return Script.runSimpleBashScriptWithFullResult(String.format("timeout %d virsh -c %s cpu-compare %s 2>&1 || true", + CPU_PRECHECK_TIMEOUT_SECONDS, destinationUri, tmp.getAbsolutePath()), CPU_PRECHECK_TIMEOUT_SECONDS + 10); + } finally { + if (!tmp.delete()) { + tmp.deleteOnExit(); + } + } + } + + /** + * opt-in CPU-compatibility precheck, run BEFORE any migration setup so a rejection fails + * fast and never touches the source domain. Returns an error message if the destination host CPU + * is incompatible, or null if it is compatible / the check could not run (best-effort, fail-open). + */ + protected String precheckDestinationCpu(final String vmName, final String destinationUri, final LibvirtComputingResource libvirtComputingResource) { + try { + final LibvirtUtilitiesHelper helper = libvirtComputingResource.getLibvirtUtilitiesHelper(); + final Connect conn = helper.getConnectionByVmName(vmName); + final Domain dm = conn.domainLookupByName(vmName); + final int xmlFlag = conn.getLibVirVersion() >= 1000000 ? 8 : 1; + final String cpuXml = extractCpuElement(dm.getXMLDesc(xmlFlag)); + if (cpuXml == null) { + return null; + } + final String output = runDestinationCpuCompare(cpuXml, destinationUri); + if (!isCpuCompareOutputCompatible(output)) { + return String.format("Cannot migrate VM [%s]: the destination host CPU is not compatible with the VM's CPU " + + "(virsh cpu-compare: %s). Choose a destination host with a compatible or superset CPU.", + vmName, output == null ? "no result" : output.trim()); + } + return null; + } catch (final Exception e) { + logger.warn(String.format("CPU-compatibility precheck for VM [%s] could not run; proceeding with migration: %s", vmName, e.getMessage())); + return null; + } + } + + /** + * migration progress as a 0-100 percent of migration data transferred, from libvirt job stats. + * Returns -1 when the total is not yet known, so the caller logs "unknown" rather than a + * misleading 0%. + */ + protected int computeMigrationProgressPercent(final long dataProcessed, final long dataRemaining) { + final long total = dataProcessed + dataRemaining; + if (total <= 0) { + return -1; + } + return (int) Math.min(100, (dataProcessed * 100) / total); + } + private DomainState getDestDomainState(Domain destDomain, String vmName) { DomainState dmState = null; try { @@ -434,15 +590,51 @@ private DomainState getDestDomainState(Domain destDomain, String vmName) { return dmState; } - private void resumeDomainIfPaused(Domain destDomain, String vmName) { + /** + * Resume a destination domain that libvirt left paused after migration. The migration itself has already + * succeeded (the guest now lives on the destination), so the caller keeps the result successful; this returns a + * message only so the caller can surface the still-paused state as a warning for the operator to resume. + * + * @return {@code null} if the domain is running (or was never paused); otherwise a message describing the + * still-not-running state on the destination. + */ + protected String resumeDomainIfPaused(Domain destDomain, String vmName) { DomainState dmState = getDestDomainState(destDomain, vmName); - if (dmState == DomainState.VIR_DOMAIN_PAUSED) { - logger.info("Resuming VM " + vmName + " on destination after migration"); - try { - destDomain.resume(); - } catch (final Exception e) { - logger.error("Failed to resume vm " + vmName + " on destination after migration due to : " + e.getMessage()); + if (dmState != DomainState.VIR_DOMAIN_PAUSED) { + return null; + } + logger.info("Resuming VM " + vmName + " on destination after migration"); + try { + destDomain.resume(); + } catch (final Exception e) { + logger.error("Failed to resume vm " + vmName + " on destination after migration due to : " + e.getMessage()); + } + DomainState afterState = getDestDomainState(destDomain, vmName); + if (afterState == DomainState.VIR_DOMAIN_RUNNING) { + return null; + } + return String.format("the guest is on the destination but not running (state: %s); resume it on the destination host", afterState); + } + + /** + * Tears down source-side networking after a successful migration. The guest is already live on the + * destination, so a failure here must not flip a successful migration to FAILED (which would make + * orchestration roll back or mark a running VM inconsistent); the error is logged and swallowed. + */ + protected void cleanupSourceNetworkingAfterMigration(Connect conn, String vmName, List ifaces, + Map vlanToPersistenceMap, LibvirtComputingResource libvirtComputingResource) { + try { + libvirtComputingResource.destroyNetworkRulesForVM(conn, vmName); + for (final InterfaceDef iface : ifaces) { + String vlanId = libvirtComputingResource.getVlanIdFromBridgeName(iface.getBrName()); + // the traffic type of each interface is unknown here, so inform all vif drivers + for (final VifDriver vifDriver : libvirtComputingResource.getAllVifDrivers()) { + vifDriver.unplug(iface, libvirtComputingResource.shouldDeleteBridge(vlanToPersistenceMap, vlanId)); + } } + } catch (final Exception e) { + logger.warn("Migration of VM [{}] succeeded, but source-side network cleanup failed: {}. " + + "Keeping the migration result as successful.", vmName, e.getMessage(), e); } } diff --git a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResourceTest.java b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResourceTest.java index 008d9444fe5f..066f03f8ebb0 100644 --- a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResourceTest.java +++ b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResourceTest.java @@ -5928,6 +5928,43 @@ public void testConfigureLocalStorageWithInvalidUUID() throws ConfigurationExcep } } + @Test + public void resolveMigrationNetworkIpTestReturnsNullForBlankLabel() { + Assert.assertNull(libvirtComputingResourceSpy.resolveMigrationNetworkIp(null)); + Assert.assertNull(libvirtComputingResourceSpy.resolveMigrationNetworkIp(" ")); + } + + @Test + public void resolveMigrationNetworkIpTestResolvesTheLabelToItsIp() { + NetworkInterface migrateNic = Mockito.mock(NetworkInterface.class); + try (MockedStatic netUtilsMockedStatic = Mockito.mockStatic(NetUtils.class)) { + netUtilsMockedStatic.when(() -> NetUtils.getNetworkInterface("cloudbr5")).thenReturn(migrateNic); + netUtilsMockedStatic.when(() -> NetUtils.getNetworkParams(migrateNic)).thenReturn(new String[] {"10.2.3.4", "aa:bb:cc:dd:ee:ff", "255.255.255.0"}); + + Assert.assertEquals("10.2.3.4", libvirtComputingResourceSpy.resolveMigrationNetworkIp("cloudbr5")); + } + } + + @Test + public void resolveMigrationNetworkIpTestReturnsNullWhenLabelHasNoIp() { + NetworkInterface migrateNic = Mockito.mock(NetworkInterface.class); + try (MockedStatic netUtilsMockedStatic = Mockito.mockStatic(NetUtils.class)) { + netUtilsMockedStatic.when(() -> NetUtils.getNetworkInterface("cloudbr5")).thenReturn(migrateNic); + netUtilsMockedStatic.when(() -> NetUtils.getNetworkParams(migrateNic)).thenReturn(new String[] {"", "aa:bb:cc:dd:ee:ff", ""}); + + Assert.assertNull(libvirtComputingResourceSpy.resolveMigrationNetworkIp("cloudbr5")); + } + } + + @Test + public void resolveMigrationNetworkIpTestReturnsNullWhenLabelNotFound() { + try (MockedStatic netUtilsMockedStatic = Mockito.mockStatic(NetUtils.class)) { + netUtilsMockedStatic.when(() -> NetUtils.getNetworkInterface("cloudbr5")).thenReturn(null); + + Assert.assertNull(libvirtComputingResourceSpy.resolveMigrationNetworkIp("cloudbr5")); + } + } + @Test public void defineResourceNetworkInterfacesTestUseProperties() { NetworkInterface networkInterfaceMock1 = Mockito.mock(NetworkInterface.class); diff --git a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDefTest.java b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDefTest.java index 56ad267eac7e..9b32c9cb9e26 100644 --- a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDefTest.java +++ b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtVMDefTest.java @@ -213,6 +213,26 @@ public void testCpuModeDef() { } + @Test + public void testCpuModeDefModelFallbackForbid() { + LibvirtVMDef.CpuModeDef cpuModeDef = new LibvirtVMDef.CpuModeDef(); + cpuModeDef.setMode("custom"); + cpuModeDef.setModel("Haswell-noTSX"); + cpuModeDef.setModelFallback("forbid"); + + assertEquals("Haswell-noTSX", cpuModeDef.toString()); + } + + @Test + public void testCpuModeDefModelFallbackDefaultsToAllowAndIgnoresGarbage() { + LibvirtVMDef.CpuModeDef cpuModeDef = new LibvirtVMDef.CpuModeDef(); + cpuModeDef.setMode("custom"); + cpuModeDef.setModel("Nehalem"); + cpuModeDef.setModelFallback("' onerror"); + + assertEquals("Nehalem", cpuModeDef.toString()); + } + @Test public void testCpuModeDefCpuFeatures() { LibvirtVMDef.CpuModeDef cpuModeDef = new LibvirtVMDef.CpuModeDef(); diff --git a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsyncTest.java b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsyncTest.java index 28633b925b21..7d904223d20c 100644 --- a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsyncTest.java +++ b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/MigrateKVMAsyncTest.java @@ -45,11 +45,11 @@ public class MigrateKVMAsyncTest { @Test public void createTypedParameterListTestNoMigrateDiskLabels() { MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "testxml", - false, false, false, "tst", "1.1.1.1", null); + false, false, false, false, false, false, "", 0, "tst", "1.1.1.1", null, null); Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); - TypedParameter[] result = migrateKVMAsync.createTypedParameterList(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); Assert.assertEquals(4, result.length); @@ -64,11 +64,11 @@ public void createTypedParameterListTestNoMigrateDiskLabels() { public void createTypedParameterListTestWithMigrateDiskLabels() { Set labels = Set.of("vda", "vdb"); MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "testxml", - false, false, false, "tst", "1.1.1.1", labels); + false, false, false, false, false, false, "", 0, "tst", "1.1.1.1", null, labels); Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); - TypedParameter[] result = migrateKVMAsync.createTypedParameterList(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); Assert.assertEquals(6, result.length); @@ -80,4 +80,138 @@ public void createTypedParameterListTestWithMigrateDiskLabels() { Assert.assertEquals(labels, Set.of(result[4].getValueAsString(), result[5].getValueAsString())); } + @Test + public void buildMigrateFlagsSetsTlsWhenEncryptionEnabled() { + // with migrate encryption enabled and a TLS-capable libvirt, VIR_MIGRATE_TLS (1<<16) is set. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, true, false, false, "", 0, "tst", "1.1.1.1", null, null); + long flags = migrateKVMAsync.buildMigrateFlags(9000000L); + Assert.assertTrue("VIR_MIGRATE_TLS must be set when encryption is enabled", (flags & 65536L) != 0L); + } + + @Test(expected = com.cloud.utils.exception.CloudRuntimeException.class) + public void buildMigrateFlagsFailsClosedWhenEncryptionRequiredButLibvirtTooOld() { + // encryption Required on libvirt < 3.2.0 must FAIL, not silently send the memory stream in plaintext. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, true, false, false, "", 0, "tst", "1.1.1.1", null, null); + migrateKVMAsync.buildMigrateFlags(3000000L); + } + + @Test + public void buildMigrateFlagsOmitsTlsWhenEncryptionDisabled() { + // with encryption disabled, VIR_MIGRATE_TLS must not be set (plaintext, prior behavior). + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, false, "", 0, "tst", "1.1.1.1", null, null); + long flags = migrateKVMAsync.buildMigrateFlags(9000000L); + Assert.assertEquals("VIR_MIGRATE_TLS must NOT be set when encryption is disabled", 0L, (flags & 65536L)); + } + + @Test + public void buildMigrateFlagsSetsParallelWhenEnabled() { + // with parallel migration enabled and a capable libvirt, VIR_MIGRATE_PARALLEL (1<<17) is set. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, true, false, "", 0, "tst", "1.1.1.1", null, null); + long flags = migrateKVMAsync.buildMigrateFlags(9000000L); + Assert.assertTrue("VIR_MIGRATE_PARALLEL must be set when parallel migration is enabled", (flags & 131072L) != 0L); + // legacy xbzrle compression must NOT be combined with multifd (QEMU refuses it). + Assert.assertEquals("legacy VIR_MIGRATE_COMPRESSED must be OFF with multifd", 0L, (flags & 2048L)); + } + + @Test + public void buildMigrateFlagsLegacyCompressionOnlyWhenNotParallel() { + // VIR_MIGRATE_COMPRESSED (2048) is set for a normal migration but omitted with multifd. + MigrateKVMAsync notParallel = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, false, "", 0, "tst", "1.1.1.1", null, null); + Assert.assertTrue("compression on for non-parallel migration", (notParallel.buildMigrateFlags(9000000L) & 2048L) != 0L); + MigrateKVMAsync parallel = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, true, false, "", 0, "tst", "1.1.1.1", null, null); + Assert.assertEquals("compression off for multifd migration", 0L, (parallel.buildMigrateFlags(9000000L) & 2048L)); + } + + @Test + public void buildMigrateFlagsSetsUnsafeWhenAllowed() { + // with migrate.allow.unsafe, VIR_MIGRATE_UNSAFE (1<<9 = 512) is set so writeback-cache + // VMs on coherent storage (Ceph) can live-migrate instead of being refused by libvirt. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, true, "", 0, "tst", "1.1.1.1", null, null); + long flags = migrateKVMAsync.buildMigrateFlags(9000000L); + Assert.assertTrue("VIR_MIGRATE_UNSAFE must be set when migrate.allow.unsafe is on", (flags & 512L) != 0L); + } + + @Test + public void createTypedParameterListIncludesCompressionWhenMethodSet() { + // a configured compression method adds the VIR_MIGRATE_PARAM_COMPRESSION ("compression") param. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, false, "xbzrle", 0, "tst", "1.1.1.1", null, null); + Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); + Assert.assertEquals("5 fixed params when a compression method is set", 5, result.length); + Assert.assertEquals("xbzrle", result[4].getValueAsString()); + } + + @Test + public void createTypedParameterListOmitsCompressionWhenMethodBlank() { + // a blank compression method preserves the prior 4-param behaviour. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, false, "", 0, "tst", "1.1.1.1", null, null); + Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); + Assert.assertEquals("no compression param when method blank", 4, result.length); + } + + @Test + public void createTypedParameterListBindsListenAddressWhenDedicatedMigrationNetwork() { + // a dedicated migration address adds the VIR_MIGRATE_PARAM_LISTEN_ADDRESS ("listen_address") param + // and the data URI targets that address, so the stream lands on the migration NIC. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, false, "", 0, "tst", "192.0.2.17", "192.0.2.17", null); + Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); + Assert.assertEquals("5 fixed params when a migration listen address is set", 5, result.length); + Assert.assertEquals("tcp:192.0.2.17", result[2].getValueAsString()); + Assert.assertEquals("192.0.2.17", result[4].getValueAsString()); + } + + @Test + public void createTypedParameterListOmitsListenAddressWhenBlank() { + // with no dedicated migration address, the prior 4-param behaviour is preserved. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, false, "", 0, "tst", "1.1.1.1", null, null); + Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); + Assert.assertEquals("no listen-address param when blank", 4, result.length); + } + + @Test + public void createTypedParameterListIncludesParallelConnectionsWhenParallelEnabled() { + // parallel migration enabled with an explicit channel count adds the parallel.connections param. + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, true, false, "", 4, "tst", "1.1.1.1", null, null); + Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); + Assert.assertEquals("extra param for parallel connections", 5, result.length); + Assert.assertEquals(4, ((org.libvirt.TypedIntParameter) result[4]).value); + } + + @Test + public void createTypedParameterListOmitsParallelConnectionsWhenParallelDisabled() { + // a channel count is ignored unless parallel migration is enabled (libvirt ignores it anyway). + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, false, false, "", 4, "tst", "1.1.1.1", null, null); + Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(6000000L); + Assert.assertEquals("no parallel-connections param when parallel disabled", 4, result.length); + } + + @Test + public void createTypedParameterListOmitsParallelConnectionsOnOldLibvirt() { + // libvirt < 5.2.0 does not set VIR_MIGRATE_PARALLEL, so the connections param must be omitted too, + // otherwise libvirt rejects the migration with "Turn parallel migration on to tune it". + MigrateKVMAsync migrateKVMAsync = new MigrateKVMAsync(libvirtComputingResource, domain, connect, "xml", + false, false, false, false, true, false, "", 4, "tst", "1.1.1.1", null, null); + Mockito.doReturn(10).when(libvirtComputingResource).getMigrateSpeed(); + TypedParameter[] result = migrateKVMAsync.createTypedParameterList(5001000L); + Assert.assertEquals("no parallel-connections param on old libvirt", 4, result.length); + } + } diff --git a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapperTest.java b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapperTest.java new file mode 100644 index 000000000000..23d9cb3e9217 --- /dev/null +++ b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapperTest.java @@ -0,0 +1,69 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package com.cloud.hypervisor.kvm.resource.wrapper; + +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mockStatic; + +import java.util.Arrays; +import java.util.Collections; + +import org.junit.Assert; +import org.junit.Test; +import org.mockito.MockedStatic; + +import com.cloud.agent.api.Answer; +import com.cloud.agent.api.BaselineCpuCommand; +import com.cloud.hypervisor.kvm.resource.LibvirtComputingResource; +import com.cloud.utils.script.Script; + +public class LibvirtBaselineCpuCommandWrapperTest { + + private final LibvirtBaselineCpuCommandWrapper wrapper = new LibvirtBaselineCpuCommandWrapper(); + private final LibvirtComputingResource resource = new LibvirtComputingResource(); + + @Test + public void emptyInputReturnsFailure() { + Answer answer = wrapper.execute(new BaselineCpuCommand(Collections.emptyList()), resource); + Assert.assertFalse(answer.getResult()); + } + + @Test + public void computesBaselineWhenOutputContainsModel() { + try (MockedStatic