From 2817a7d0e53d6022a7cf2be5d9750a551d06cc8a Mon Sep 17 00:00:00 2001 From: Ramgopal Nagaboina Date: Tue, 6 Oct 2026 16:44:00 -0400 Subject: [PATCH] kvm: dedicated live migration network with TLS, parallel streams and cluster CPU baseline Carry KVM live migration traffic on a dedicated network instead of the management network. A Migration traffic type is added to the zone's physical network with a per hypervisor label; each KVM host resolves its own IP on the labelled NIC and reports it, and the libvirt data stream is routed to it only when both the source and the destination host have one, otherwise migration falls back to the existing behaviour unchanged. listHosts returns each host's migration IP, shown on the host detail view, so an operator can confirm the dedicated network is in use. Migrations can be encrypted with libvirt native TLS and parallelised with multiple file-descriptor streams. Both are opt-in and gated on the host libvirt version, failing closed with a clear error when the host is too old rather than silently downgrading. Add a cluster-scoped CPU baseline model so a mixed cluster can present a uniform guest CPU, keeping a VM's CPU definition stable as it migrates across hosts. The baseline guest is pinned with fallback=forbid so it refuses a host that cannot provide the exact model rather than silently degrading the CPU the guest sees; the forbid is scoped to the baseline and leaves a VM with its own explicit CPU mode or model untouched. The baseline is validated against every host in the cluster before it is accepted, and hosts that cannot run the model are reported back instead of failing the migration later. Setting the baseline to "auto" computes the common-denominator model of the cluster's reachable hosts, by collecting each host's CPU and running cpu-baseline, and persists the computed model; it does not re-compute as hosts change, so the committed baseline stays stable. Unit tests cover the IP resolver (both-ends gating, non-KVM skip, blank values), the libvirt flag and typed-parameter construction, the version gate, the CPU baseline injection, the forbid fallback emission and scoping, the auto-compute and model parsing, and the host compatibility check including the fail-closed path on an unknown model. A Marvin smoke test exercises the migration traffic type API. --- agent/conf/agent.properties | 35 +++ .../agent/properties/AgentProperties.java | 57 +++++ api/src/main/java/com/cloud/host/Host.java | 1 + .../main/java/com/cloud/network/Networks.java | 4 +- .../network/PhysicalNetworkSetupInfo.java | 9 + .../java/com/cloud/vm/VmDetailConstants.java | 3 + .../cloudstack/api/response/HostResponse.java | 12 + .../apache/cloudstack/query/QueryService.java | 2 +- .../cloud/agent/api/BaselineCpuCommand.java | 47 ++++ .../api/CheckCpuCompatibilityCommand.java | 50 ++++ .../cloud/agent/api/CheckNetworkAnswer.java | 10 + .../agent/api/GetHostCpuModelCommand.java | 35 +++ .../com/cloud/agent/api/MigrateCommand.java | 18 ++ .../com/cloud/vm/VirtualMachineManager.java | 30 +++ .../cloud/vm/VirtualMachineManagerImpl.java | 213 +++++++++++++++- .../orchestration/NetworkOrchestrator.java | 31 +++ ...tualMachineManagerImplCpuBaselineTest.java | 117 +++++++++ ...irtualMachineManagerImplCpuCompatTest.java | 126 ++++++++++ ...tualMachineManagerImplMigrationIpTest.java | 128 ++++++++++ .../vm/VirtualMachineManagerImplTest.java | 64 +++++ .../NetworkOrchestratorTest.java | 70 ++++++ .../resource/LibvirtComputingResource.java | 16 ++ .../hypervisor/kvm/resource/LibvirtVMDef.java | 9 +- .../kvm/resource/MigrateKVMAsync.java | 108 +++++++- .../LibvirtBaselineCpuCommandWrapper.java | 73 ++++++ ...rtCheckCpuCompatibilityCommandWrapper.java | 82 ++++++ .../LibvirtCheckNetworkCommandWrapper.java | 8 +- .../LibvirtGetHostCpuModelCommandWrapper.java | 61 +++++ .../wrapper/LibvirtMigrateCommandWrapper.java | 234 ++++++++++++++++-- .../LibvirtComputingResourceTest.java | 37 +++ .../kvm/resource/LibvirtVMDefTest.java | 20 ++ .../kvm/resource/MigrateKVMAsyncTest.java | 142 ++++++++++- .../LibvirtBaselineCpuCommandWrapperTest.java | 69 ++++++ ...eckCpuCompatibilityCommandWrapperTest.java | 184 ++++++++++++++ ...LibvirtCheckNetworkCommandWrapperTest.java | 69 ++++++ ...virtGetHostCpuModelCommandWrapperTest.java | 43 ++++ .../LibvirtMigrateCommandWrapperTest.java | 124 ++++++++++ .../cloud/api/query/dao/HostJoinDaoImpl.java | 3 + .../ConfigurationManagerImpl.java | 43 ++++ .../com/cloud/network/NetworkModelImpl.java | 2 + .../com/cloud/network/NetworkServiceImpl.java | 29 ++- .../cloud/resource/ResourceManagerImpl.java | 31 +++ .../ConfigurationManagerImplTest.java | 32 +++ .../cloud/network/NetworkServiceImplTest.java | 33 +++ .../smoke/test_migration_network.py | 93 +++++++ ui/public/locales/en.json | 4 + ui/src/config/section/infra/clusters.js | 9 + ui/src/config/section/infra/hosts.js | 2 +- ui/src/config/section/infra/phynetworks.js | 2 +- ui/src/views/infra/ConfigureCpuBaseline.vue | 128 ++++++++++ 50 files changed, 2709 insertions(+), 43 deletions(-) create mode 100644 core/src/main/java/com/cloud/agent/api/BaselineCpuCommand.java create mode 100644 core/src/main/java/com/cloud/agent/api/CheckCpuCompatibilityCommand.java create mode 100644 core/src/main/java/com/cloud/agent/api/GetHostCpuModelCommand.java create mode 100644 engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuBaselineTest.java create mode 100644 engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplCpuCompatTest.java create mode 100644 engine/orchestration/src/test/java/com/cloud/vm/VirtualMachineManagerImplMigrationIpTest.java create mode 100644 plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapper.java create mode 100644 plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtCheckCpuCompatibilityCommandWrapper.java create mode 100644 plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtGetHostCpuModelCommandWrapper.java create mode 100644 plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtBaselineCpuCommandWrapperTest.java create mode 100644 plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtCheckCpuCompatibilityCommandWrapperTest.java create mode 100644 plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtCheckNetworkCommandWrapperTest.java create mode 100644 plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/wrapper/LibvirtGetHostCpuModelCommandWrapperTest.java create mode 100644 test/integration/smoke/test_migration_network.py create mode 100644 ui/src/views/infra/ConfigureCpuBaseline.vue 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