diff --git a/plugins/storage/volume/linstor/CHANGELOG.md b/plugins/storage/volume/linstor/CHANGELOG.md index a6ab050b090e..2192c7d96af5 100644 --- a/plugins/storage/volume/linstor/CHANGELOG.md +++ b/plugins/storage/volume/linstor/CHANGELOG.md @@ -24,6 +24,17 @@ All notable changes to Linstor CloudStack plugin will be documented in this file The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [2026-10-07] + +### Fixed + +- Live migration with storage into Linstor failed with "cannot precreate storage for disk type 'block'" + if the auto-placement didn't put a copy of the new resource on the target host; the resource is + now made available on the target host before the migration starts. +- Live migration with storage into Linstor failed on cgroup v2 hosts with + "shares '' must be in range [1, 10000]" for VMs with more than 10000 cpus * MHz, and smaller + VMs ended up with an unscaled CPU weight; the CPU shares calculated by the target host are now used. + ## [2026-06-24] ### Fixed diff --git a/plugins/storage/volume/linstor/src/main/java/org/apache/cloudstack/storage/motion/LinstorDataMotionStrategy.java b/plugins/storage/volume/linstor/src/main/java/org/apache/cloudstack/storage/motion/LinstorDataMotionStrategy.java index f64837e4832e..2f0c7b39e680 100644 --- a/plugins/storage/volume/linstor/src/main/java/org/apache/cloudstack/storage/motion/LinstorDataMotionStrategy.java +++ b/plugins/storage/volume/linstor/src/main/java/org/apache/cloudstack/storage/motion/LinstorDataMotionStrategy.java @@ -24,6 +24,7 @@ import com.linbit.linstor.api.model.ApiCallRcList; import com.linbit.linstor.api.model.ResourceDefinition; import com.linbit.linstor.api.model.ResourceDefinitionModify; +import com.linbit.linstor.api.model.ResourceMakeAvailable; import javax.inject.Inject; @@ -35,8 +36,11 @@ import com.cloud.agent.AgentManager; import com.cloud.agent.api.Answer; +import com.cloud.agent.api.CheckVirtualMachineAnswer; +import com.cloud.agent.api.CheckVirtualMachineCommand; import com.cloud.agent.api.MigrateAnswer; import com.cloud.agent.api.MigrateCommand; +import com.cloud.agent.api.PrepareForMigrationAnswer; import com.cloud.agent.api.PrepareForMigrationCommand; import com.cloud.agent.api.to.DataObjectType; import com.cloud.agent.api.to.VirtualMachineTO; @@ -54,6 +58,7 @@ import com.cloud.storage.dao.VolumeDao; import com.cloud.utils.exception.CloudRuntimeException; import com.cloud.vm.VMInstanceVO; +import com.cloud.vm.VirtualMachine; import com.cloud.vm.dao.VMInstanceDao; import org.apache.cloudstack.engine.subsystem.api.storage.CopyCommandResult; import org.apache.cloudstack.engine.subsystem.api.storage.DataMotionStrategy; @@ -74,7 +79,6 @@ import org.apache.cloudstack.storage.datastore.util.LinstorUtil; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.collections.MapUtils; -import org.apache.commons.lang3.ObjectUtils; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.springframework.stereotype.Component; @@ -169,11 +173,51 @@ private VolumeVO createNewVolumeVO(Volume volume, StoragePoolVO storagePoolVO) { return _volumeDao.persist(newVol); } + private DevelopersApi getLinstorAPI(StoragePoolVO storagePool) { + return LinstorUtil.getLinstorAPI(storagePool.getHostAddress(), + LinstorConfigurationManager.ApiToken.valueIn(storagePool.getId()), + Boolean.TRUE.equals(LinstorConfigurationManager.InsecureSsl.valueIn(storagePool.getId()))); + } + + /** + * Makes the new resource available on the migration target host and returns its device path there. + * + *

The resource is spawned by LINSTOR's auto-placement, which does not know about the target host. + * libvirt can only precreate file disks on the target, so if the block device is missing there the + * migration fails with "cannot precreate storage for disk type 'block'". The source disk is not on + * LINSTOR, so a plain (diskless for DRBD) make-available is enough, no dual-primary needed.

+ */ + private String makeAvailableOnHost(StoragePoolVO storagePool, String rscName, Host host) { + DevelopersApi api = getLinstorAPI(storagePool); + try { + logger.info("Linstor: make resource {} available on migration target {}", rscName, host.getName()); + ApiCallRcList answers = api.resourceMakeAvailableOnNode(rscName, host.getName(), new ResourceMakeAvailable()); + LinstorUtil.checkLinstorAnswersThrow(answers); + return LinstorUtil.getDevicePath(api, rscName); + } catch (ApiException apiEx) { + logger.error("Linstor: ApiEx - {}", apiEx.getMessage()); + throw new CloudRuntimeException(apiEx.getBestMessage(), apiEx); + } + } + + /** + * A paused domain (incoming migration still in progress) is reported as PowerUnknown, so this + * is only true once the VM really runs on the host. + */ + private boolean isVmRunningOnHost(VirtualMachineTO vmTO, Host host) { + try { + Answer answer = _agentManager.send(host.getId(), new CheckVirtualMachineCommand(vmTO.getName())); + return answer instanceof CheckVirtualMachineAnswer && answer.getResult() && + VirtualMachine.PowerState.PowerOn.equals(((CheckVirtualMachineAnswer) answer).getState()); + } catch (AgentUnavailableException | OperationTimedoutException e) { + logger.warn("Unable to check if VM [{}] is running on host [{}]", vmTO, host, e); + return false; + } + } + private void removeExactSizeProperty(VolumeInfo volumeInfo) { StoragePoolVO destStoragePool = _storagePool.findById(volumeInfo.getDataStore().getId()); - DevelopersApi api = LinstorUtil.getLinstorAPI(destStoragePool.getHostAddress(), - LinstorConfigurationManager.ApiToken.valueIn(destStoragePool.getId()), - Boolean.TRUE.equals(LinstorConfigurationManager.InsecureSsl.valueIn(destStoragePool.getId()))); + DevelopersApi api = getLinstorAPI(destStoragePool); ResourceDefinitionModify rdm = new ResourceDefinitionModify(); rdm.setDeleteProps(Collections.singletonList(LinstorUtil.LIN_PROP_DRBDOPT_EXACT_SIZE)); @@ -266,10 +310,10 @@ private void handlePostMigration(boolean success, Map sr _volumeService.expungeVolumeAsync(destVolumeInfo); if (destroyFuture.get().isFailed()) { - logger.debug("Failed to clean up dest volume on storage"); + logger.warn("Failed to clean up dest volume {} on storage", destVolumeInfo); } } catch (Exception e) { - logger.debug("Failed to clean up dest volume on storage", e); + logger.warn("Failed to clean up dest volume {} on storage", destVolumeInfo, e); } } } @@ -293,9 +337,7 @@ private void handlePostMigration(boolean success, Map sr private boolean needsExactSizeProp(VolumeInfo srcVolumeInfo) { StoragePoolVO srcStoragePool = _storagePool.findById(srcVolumeInfo.getDataStore().getId()); if (srcStoragePool.getPoolType() == Storage.StoragePoolType.Linstor) { - DevelopersApi api = LinstorUtil.getLinstorAPI(srcStoragePool.getHostAddress(), - LinstorConfigurationManager.ApiToken.valueIn(srcStoragePool.getId()), - Boolean.TRUE.equals(LinstorConfigurationManager.InsecureSsl.valueIn(srcStoragePool.getId()))); + DevelopersApi api = getLinstorAPI(srcStoragePool); String rscName = LinstorUtil.RSC_PREFIX + srcVolumeInfo.getPath(); try { @@ -335,6 +377,7 @@ public void copyAsync(Map volumeDataStoreMap, VirtualMach Map migrateStorage = new HashMap<>(); Map srcVolumeInfoToDestVolumeInfo = new HashMap<>(); + boolean migrateCommandSent = false; try { for (Map.Entry entry : volumeDataStoreMap.entrySet()) { @@ -351,6 +394,8 @@ public void copyAsync(Map volumeDataStoreMap, VirtualMach VolumeVO destVolume = createNewVolumeVO(srcVolume, destStoragePool); VolumeInfo destVolumeInfo = _volumeDataFactory.getVolume(destVolume.getId(), destDataStore); + // registered right away, so a failure from here on cleans the new volume up again + srcVolumeInfoToDestVolumeInfo.put(srcVolumeInfo, destVolumeInfo); destVolumeInfo.processEvent(ObjectInDataStoreStateMachine.Event.MigrationCopyRequested); destVolumeInfo.processEvent(ObjectInDataStoreStateMachine.Event.MigrationCopySucceeded); @@ -358,13 +403,16 @@ public void copyAsync(Map volumeDataStoreMap, VirtualMach boolean exactSize = needsExactSizeProp(srcVolumeInfo); - String devPath = LinstorUtil.createResource( - destVolumeInfo, destStoragePool, _storagePoolDao, exactSize); + LinstorUtil.createResource(destVolumeInfo, destStoragePool, _storagePoolDao, exactSize); _volumeDao.update(destVolume.getId(), destVolume); destVolume = _volumeDao.findById(destVolume.getId()); destVolumeInfo = _volumeDataFactory.getVolume(destVolume.getId(), destDataStore); + srcVolumeInfoToDestVolumeInfo.put(srcVolumeInfo, destVolumeInfo); + + String devPath = makeAvailableOnHost( + destStoragePool, LinstorUtil.RSC_PREFIX + destVolumeInfo.getUuid(), destHost); MigrateCommand.MigrateDiskInfo migrateDiskInfo = new MigrateCommand.MigrateDiskInfo( srcVolumeInfo.getPath(), @@ -375,13 +423,12 @@ public void copyAsync(Map volumeDataStoreMap, VirtualMach migrateDiskInfoList.add(migrateDiskInfo); migrateStorage.put(srcVolumeInfo.getPath(), migrateDiskInfo); - - srcVolumeInfoToDestVolumeInfo.put(srcVolumeInfo, destVolumeInfo); } PrepareForMigrationCommand pfmc = new PrepareForMigrationCommand(vmTO); + Answer pfma; try { - Answer pfma = _agentManager.send(destHost.getId(), pfmc); + pfma = _agentManager.send(destHost.getId(), pfmc); if (pfma == null || !pfma.getResult()) { String details = pfma != null ? pfma.getDetails() : "null answer returned"; @@ -403,23 +450,41 @@ public void copyAsync(Map volumeDataStoreMap, VirtualMach migrateCommand.setWait(StorageManager.KvmStorageOnlineMigrationWait.value()); migrateCommand.setMigrateStorage(migrateStorage); migrateCommand.setMigrateStorageManaged(true); - migrateCommand.setNewVmCpuShares( - vmTO.getCpus() * ObjectUtils.defaultIfNull(vmTO.getMinSpeed(), vmTO.getSpeed())); + // the target host scales the shares to its cgroup version (v2 only accepts 1-10000), + // the raw cpus * speed value is only valid on cgroup v1 + Integer newVmCpuShares = ((PrepareForMigrationAnswer) pfma).getNewVmCpuShares(); + if (newVmCpuShares != null) { + migrateCommand.setNewVmCpuShares(newVmCpuShares); + } migrateCommand.setMigrateDiskInfoList(migrateDiskInfoList); boolean kvmAutoConvergence = StorageManager.KvmAutoConvergence.value(); migrateCommand.setAutoConvergence(kvmAutoConvergence); - MigrateAnswer migrateAnswer = (MigrateAnswer) _agentManager.send(srcHost.getId(), migrateCommand); - boolean success = migrateAnswer != null && migrateAnswer.getResult(); + // once the migrate command is sent, the VM may end up running on the new volumes, + // so they are only removed if the source host reports the migration as failed + migrateCommandSent = true; + MigrateAnswer migrateAnswer = null; + boolean success; + try { + migrateAnswer = (MigrateAnswer) _agentManager.send(srcHost.getId(), migrateCommand); + success = migrateAnswer != null && migrateAnswer.getResult(); + } catch (OperationTimedoutException ex) { + // no answer from the source host, but the migration may still have finished + if (!isVmRunningOnHost(vmTO, destHost)) { + throw ex; + } + logger.info("VM [{}] is running on the destination host [{}], migration was successful", vmTO, destHost); + success = true; + } handlePostMigration(success, srcVolumeInfoToDestVolumeInfo, vmTO, destHost); - if (migrateAnswer == null) { - throw new CloudRuntimeException("Unable to get an answer to the migrate command"); - } + if (!success) { + if (migrateAnswer == null) { + throw new CloudRuntimeException("Unable to get an answer to the migrate command"); + } - if (!migrateAnswer.getResult()) { errMsg = migrateAnswer.getDetails(); throw new CloudRuntimeException(errMsg); @@ -427,9 +492,18 @@ public void copyAsync(Map volumeDataStoreMap, VirtualMach } catch (AgentUnavailableException | OperationTimedoutException | CloudRuntimeException ex) { errMsg = String.format( "Copy volume(s) of VM [%s] to storage(s) [%s] and VM to host [%s] failed in LinstorDataMotionStrategy.copyAsync. Error message: [%s].", - vmTO, srcHost, destHost, ex.getMessage()); + vmTO, volumeDataStoreMap.values(), destHost, ex.getMessage()); logger.error(errMsg, ex); + if (!migrateCommandSent && !srcVolumeInfoToDestVolumeInfo.isEmpty()) { + // failed before the migration was started: remove the already created destination volumes + try { + handlePostMigration(false, srcVolumeInfoToDestVolumeInfo, vmTO, destHost); + } catch (Exception e) { + logger.warn("Failed to clean up the destination volume(s) of VM [{}]", vmTO, e); + } + } + throw new CloudRuntimeException(errMsg); } finally { CopyCmdAnswer copyCmdAnswer = new CopyCmdAnswer(errMsg); diff --git a/plugins/storage/volume/linstor/src/test/java/org/apache/cloudstack/storage/motion/LinstorDataMotionStrategyTest.java b/plugins/storage/volume/linstor/src/test/java/org/apache/cloudstack/storage/motion/LinstorDataMotionStrategyTest.java new file mode 100644 index 000000000000..9f7aeb8ea465 --- /dev/null +++ b/plugins/storage/volume/linstor/src/test/java/org/apache/cloudstack/storage/motion/LinstorDataMotionStrategyTest.java @@ -0,0 +1,356 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +package org.apache.cloudstack.storage.motion; + +import com.linbit.linstor.api.ApiException; +import com.linbit.linstor.api.DevelopersApi; +import com.linbit.linstor.api.model.ResourceMakeAvailable; + +import java.util.Collections; +import java.util.Map; + +import com.cloud.agent.AgentManager; +import com.cloud.agent.api.CheckVirtualMachineAnswer; +import com.cloud.agent.api.CheckVirtualMachineCommand; +import com.cloud.agent.api.Command; +import com.cloud.agent.api.MigrateAnswer; +import com.cloud.agent.api.MigrateCommand; +import com.cloud.agent.api.PrepareForMigrationAnswer; +import com.cloud.agent.api.PrepareForMigrationCommand; +import com.cloud.agent.api.to.VirtualMachineTO; +import com.cloud.exception.OperationTimedoutException; +import com.cloud.host.Host; +import com.cloud.hypervisor.Hypervisor; +import com.cloud.storage.GuestOSCategoryVO; +import com.cloud.storage.GuestOSVO; +import com.cloud.storage.Storage; +import com.cloud.storage.Volume; +import com.cloud.storage.VolumeVO; +import com.cloud.storage.dao.GuestOSCategoryDao; +import com.cloud.storage.dao.GuestOSDao; +import com.cloud.storage.dao.SnapshotDao; +import com.cloud.storage.dao.VolumeDao; +import com.cloud.utils.exception.CloudRuntimeException; +import com.cloud.vm.VMInstanceVO; +import com.cloud.vm.VirtualMachine; +import com.cloud.vm.dao.VMInstanceDao; +import org.apache.cloudstack.engine.subsystem.api.storage.CopyCommandResult; +import org.apache.cloudstack.engine.subsystem.api.storage.DataStore; +import org.apache.cloudstack.engine.subsystem.api.storage.VolumeDataFactory; +import org.apache.cloudstack.engine.subsystem.api.storage.VolumeInfo; +import org.apache.cloudstack.engine.subsystem.api.storage.VolumeService; +import org.apache.cloudstack.framework.async.AsyncCompletionCallback; +import org.apache.cloudstack.storage.datastore.db.PrimaryDataStoreDao; +import org.apache.cloudstack.storage.datastore.db.StoragePoolVO; +import org.apache.cloudstack.storage.datastore.util.LinstorUtil; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.InOrder; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.MockedStatic; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@RunWith(MockitoJUnitRunner.Silent.class) +public class LinstorDataMotionStrategyTest { + private static final long SRC_HOST_ID = 1L; + private static final long DEST_HOST_ID = 3L; + private static final String DEST_HOST_NAME = "kvm-dest"; + private static final long SRC_VOLUME_ID = 10L; + private static final long DEST_VOLUME_ID = 20L; + private static final long SRC_POOL_ID = 7L; + private static final long DEST_POOL_ID = 8L; + private static final String DEST_VOLUME_UUID = "8ed5dd1d-16f8-4169-8ed9-948178dbd44b"; + private static final String DEST_RSC_NAME = "cs-" + DEST_VOLUME_UUID; + private static final String DEST_DEVICE_PATH = "/dev/drbd/by-res/" + DEST_RSC_NAME + "/0"; + + @Mock + private PrimaryDataStoreDao _storagePool; + @Mock + private PrimaryDataStoreDao _storagePoolDao; + @Mock + private VolumeDao _volumeDao; + @Mock + private VolumeDataFactory _volumeDataFactory; + @Mock + private VMInstanceDao _vmDao; + @Mock + private GuestOSDao _guestOsDao; + @Mock + private VolumeService _volumeService; + @Mock + private GuestOSCategoryDao _guestOsCategoryDao; + @Mock + private SnapshotDao _snapshotDao; + @Mock + private AgentManager _agentManager; + + @InjectMocks + private LinstorDataMotionStrategy strategy; + + private MockedStatic linstorUtil; + private DevelopersApi api; + private Host srcHost; + private Host destHost; + private VirtualMachineTO vmTO; + private VolumeInfo srcVolumeInfo; + private VolumeInfo destVolumeInfo; + private DataStore destDataStore; + private AsyncCompletionCallback callback; + + @Before + @SuppressWarnings("unchecked") + public void setUp() throws Exception { + api = mock(DevelopersApi.class); + linstorUtil = Mockito.mockStatic(LinstorUtil.class); + linstorUtil.when(() -> LinstorUtil.getLinstorAPI(any(), any(), anyBoolean())).thenReturn(api); + linstorUtil.when(() -> LinstorUtil.getDevicePath(api, DEST_RSC_NAME)).thenReturn(DEST_DEVICE_PATH); + + srcHost = mock(Host.class); + when(srcHost.getId()).thenReturn(SRC_HOST_ID); + when(srcHost.getHypervisorType()).thenReturn(Hypervisor.HypervisorType.KVM); + destHost = mock(Host.class); + when(destHost.getId()).thenReturn(DEST_HOST_ID); + when(destHost.getName()).thenReturn(DEST_HOST_NAME); + when(destHost.getPrivateIpAddress()).thenReturn("192.168.0.3"); + + // 2 vCPUs * 5001 MHz: raw shares 10002, more than cgroup v2 accepts + vmTO = mock(VirtualMachineTO.class); + when(vmTO.getId()).thenReturn(100L); + when(vmTO.getName()).thenReturn("i-2-100-VM"); + when(vmTO.getCpus()).thenReturn(2); + when(vmTO.getSpeed()).thenReturn(5001); + when(vmTO.getMinSpeed()).thenReturn(5001); + VMInstanceVO vm = mock(VMInstanceVO.class); + when(vm.getState()).thenReturn(VirtualMachine.State.Migrating); + when(vm.getGuestOSId()).thenReturn(5L); + when(_vmDao.findById(100L)).thenReturn(vm); + GuestOSVO guestOs = mock(GuestOSVO.class); + when(guestOs.getCategoryId()).thenReturn(6L); + when(_guestOsDao.findById(5L)).thenReturn(guestOs); + GuestOSCategoryVO guestOsCategory = mock(GuestOSCategoryVO.class); + when(guestOsCategory.getName()).thenReturn("Other"); + when(_guestOsCategoryDao.findById(6L)).thenReturn(guestOsCategory); + + // source volume on an NFS pool + DataStore srcDataStore = mock(DataStore.class); + when(srcDataStore.getId()).thenReturn(SRC_POOL_ID); + srcVolumeInfo = mock(VolumeInfo.class); + when(srcVolumeInfo.getId()).thenReturn(SRC_VOLUME_ID); + when(srcVolumeInfo.getPath()).thenReturn("216fb792-c5ed-4ad7-98e6-ebdd25816e55"); + when(srcVolumeInfo.getDataStore()).thenReturn(srcDataStore); + when(srcVolumeInfo.getPassphraseId()).thenReturn(null); + StoragePoolVO srcPool = mock(StoragePoolVO.class); + when(srcPool.getPoolType()).thenReturn(Storage.StoragePoolType.NetworkFilesystem); + when(_storagePool.findById(SRC_POOL_ID)).thenReturn(srcPool); + when(_volumeDao.findById(SRC_VOLUME_ID)).thenReturn(new VolumeVO(Volume.Type.ROOT, "ROOT-100", 1L, 1L, 2L, 3L, + Storage.ProvisioningType.THIN, 1024L * 1024 * 1024, null, null, null)); + + // destination volume on the Linstor pool + destDataStore = mock(DataStore.class); + when(destDataStore.getId()).thenReturn(DEST_POOL_ID); + StoragePoolVO destPool = mock(StoragePoolVO.class); + when(destPool.getId()).thenReturn(DEST_POOL_ID); + when(destPool.getHostAddress()).thenReturn("http://linstor-controller"); + when(_storagePool.findById(DEST_POOL_ID)).thenReturn(destPool); + VolumeVO destVolume = mock(VolumeVO.class); + when(destVolume.getId()).thenReturn(DEST_VOLUME_ID); + when(_volumeDao.persist(any(VolumeVO.class))).thenReturn(destVolume); + when(_volumeDao.findById(DEST_VOLUME_ID)).thenReturn(destVolume); + destVolumeInfo = mock(VolumeInfo.class); + when(destVolumeInfo.getId()).thenReturn(DEST_VOLUME_ID); + when(destVolumeInfo.getUuid()).thenReturn(DEST_VOLUME_UUID); + when(destVolumeInfo.getDataStore()).thenReturn(destDataStore); + when(_volumeDataFactory.getVolume(eq(DEST_VOLUME_ID), any(DataStore.class))).thenReturn(destVolumeInfo); + when(_volumeDataFactory.getVolume(DEST_VOLUME_ID)).thenReturn(destVolumeInfo); + when(_volumeDataFactory.getVolume(SRC_VOLUME_ID)).thenReturn(srcVolumeInfo); + + callback = mock(AsyncCompletionCallback.class); + } + + @After + public void tearDown() { + linstorUtil.close(); + } + + private Map volumeMap() { + return Collections.singletonMap(srcVolumeInfo, destDataStore); + } + + private void prepareForMigrationAnswers(Integer cpuShares) throws Exception { + when(_agentManager.send(eq(DEST_HOST_ID), any(PrepareForMigrationCommand.class))).thenAnswer(inv -> { + PrepareForMigrationAnswer answer = new PrepareForMigrationAnswer(inv.getArgument(1)); + answer.setNewVmCpuShares(cpuShares); + return answer; + }); + } + + private void migrateSucceeds() throws Exception { + when(_agentManager.send(eq(SRC_HOST_ID), any(MigrateCommand.class))).thenAnswer( + inv -> new MigrateAnswer(inv.getArgument(1), true, null, null)); + } + + private MigrateCommand sentMigrateCommand() throws Exception { + ArgumentCaptor captor = ArgumentCaptor.forClass(Command.class); + verify(_agentManager, Mockito.atLeastOnce()).send(anyLong(), captor.capture()); + return captor.getAllValues().stream() + .filter(MigrateCommand.class::isInstance) + .map(MigrateCommand.class::cast) + .findFirst() + .orElseThrow(() -> new AssertionError("no MigrateCommand sent")); + } + + private CopyCommandResult completedResult() { + ArgumentCaptor captor = ArgumentCaptor.forClass(CopyCommandResult.class); + verify(callback).complete(captor.capture()); + return captor.getValue(); + } + + @Test + public void copyAsyncMakesNewResourceAvailableOnTargetHostBeforeMigrating() throws Exception { + prepareForMigrationAnswers(8335); + migrateSucceeds(); + + strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback); + + InOrder order = inOrder(api, _agentManager); + order.verify(api).resourceMakeAvailableOnNode(eq(DEST_RSC_NAME), eq(DEST_HOST_NAME), any(ResourceMakeAvailable.class)); + order.verify(_agentManager).send(eq(DEST_HOST_ID), any(PrepareForMigrationCommand.class)); + order.verify(_agentManager).send(eq(SRC_HOST_ID), any(MigrateCommand.class)); + + MigrateCommand migrateCommand = sentMigrateCommand(); + Assert.assertEquals(1, migrateCommand.getMigrateDiskInfoList().size()); + Assert.assertEquals(DEST_DEVICE_PATH, migrateCommand.getMigrateDiskInfoList().get(0).getSourceText()); + Assert.assertTrue(completedResult().isSuccess()); + } + + @Test + public void copyAsyncUsesCpuSharesCalculatedByTargetHost() throws Exception { + prepareForMigrationAnswers(8335); + migrateSucceeds(); + + strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback); + + // not the raw cpus * speed (10002), which cgroup v2 hosts reject + Assert.assertEquals(8335, sentMigrateCommand().getNewVmCpuShares()); + } + + @Test + public void copyAsyncRemovesDestinationVolumeWhenCreateResourceFails() throws Exception { + linstorUtil.when(() -> LinstorUtil.createResource(any(), any(), any(), anyBoolean())) + .thenThrow(new CloudRuntimeException("Not enough available nodes")); + + Assert.assertThrows(CloudRuntimeException.class, + () -> strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback)); + + verify(api, never()).resourceMakeAvailableOnNode(any(), any(), any()); + verify(_agentManager, never()).send(anyLong(), any(MigrateCommand.class)); + verify(_volumeService).destroyVolume(DEST_VOLUME_ID); + verify(_volumeService).expungeVolumeAsync(destVolumeInfo); + Assert.assertTrue(completedResult().isFailed()); + } + + @Test + public void copyAsyncRemovesDestinationVolumeWhenMakeAvailableFails() throws Exception { + when(api.resourceMakeAvailableOnNode(eq(DEST_RSC_NAME), eq(DEST_HOST_NAME), any(ResourceMakeAvailable.class))) + .thenThrow(new ApiException("Autoplacer could not find diskless stor pool")); + + Assert.assertThrows(CloudRuntimeException.class, + () -> strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback)); + + verify(_agentManager, never()).send(anyLong(), any(MigrateCommand.class)); + verify(_volumeService).destroyVolume(DEST_VOLUME_ID); + verify(_volumeService).expungeVolumeAsync(destVolumeInfo); + Assert.assertTrue(completedResult().isFailed()); + } + + @Test + public void copyAsyncRemovesDestinationVolumeWhenPrepareForMigrationFails() throws Exception { + when(_agentManager.send(eq(DEST_HOST_ID), any(PrepareForMigrationCommand.class))).thenAnswer( + inv -> new PrepareForMigrationAnswer(inv.getArgument(1), "failed to prepare")); + + Assert.assertThrows(CloudRuntimeException.class, + () -> strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback)); + + verify(_agentManager, never()).send(anyLong(), any(MigrateCommand.class)); + verify(_volumeService).destroyVolume(DEST_VOLUME_ID); + verify(_volumeService).expungeVolumeAsync(destVolumeInfo); + Assert.assertTrue(completedResult().isFailed()); + } + + @Test + public void copyAsyncRemovesDestinationVolumeWhenMigrationFails() throws Exception { + prepareForMigrationAnswers(8335); + when(_agentManager.send(eq(SRC_HOST_ID), any(MigrateCommand.class))).thenAnswer( + inv -> new MigrateAnswer(inv.getArgument(1), false, "blockdev-add failed", null)); + + Assert.assertThrows(CloudRuntimeException.class, + () -> strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback)); + + verify(_volumeService).destroyVolume(DEST_VOLUME_ID); + verify(_volumeService).expungeVolumeAsync(destVolumeInfo); + verify(_volumeService, never()).destroyVolume(SRC_VOLUME_ID); + Assert.assertTrue(completedResult().isFailed()); + } + + private void migrateTimesOutWithVmOnTarget(VirtualMachine.PowerState powerState) throws Exception { + prepareForMigrationAnswers(8335); + when(_agentManager.send(eq(SRC_HOST_ID), any(MigrateCommand.class))) + .thenThrow(new OperationTimedoutException(null, SRC_HOST_ID, 1L, 86400, false)); + when(_agentManager.send(eq(DEST_HOST_ID), any(CheckVirtualMachineCommand.class))).thenAnswer( + inv -> new CheckVirtualMachineAnswer(inv.getArgument(1), powerState, 5900)); + } + + @Test + public void copyAsyncSucceedsWhenMigrateTimesOutButVmRunsOnTargetHost() throws Exception { + migrateTimesOutWithVmOnTarget(VirtualMachine.PowerState.PowerOn); + + strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback); + + verify(_volumeService, never()).destroyVolume(DEST_VOLUME_ID); + verify(_volumeService).destroyVolume(SRC_VOLUME_ID); + Assert.assertTrue(completedResult().isSuccess()); + } + + @Test + public void copyAsyncKeepsDestinationVolumeWhenMigrateTimesOut() throws Exception { + // a paused domain on the target: the incoming migration may still be running + migrateTimesOutWithVmOnTarget(VirtualMachine.PowerState.PowerUnknown); + + Assert.assertThrows(CloudRuntimeException.class, + () -> strategy.copyAsync(volumeMap(), vmTO, srcHost, destHost, callback)); + + verify(_volumeService, never()).destroyVolume(DEST_VOLUME_ID); + verify(_volumeService, never()).destroyVolume(SRC_VOLUME_ID); + Assert.assertTrue(completedResult().isFailed()); + } +}