rajiv-jain-netapp commented on code in PR #14256:
URL: https://github.com/apache/cloudstack/pull/14256#discussion_r4133992802


##########
server/src/main/java/org/apache/cloudstack/vm/VmwareCbtMigrationManagerImpl.java:
##########
@@ -0,0 +1,2978 @@
+// 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.vm;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Date;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.UUID;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+import javax.inject.Inject;
+
+import org.apache.cloudstack.api.ApiErrorCode;
+import org.apache.cloudstack.api.ApiConstants;
+import org.apache.cloudstack.api.ServerApiException;
+import org.apache.cloudstack.api.command.admin.vm.CancelVmwareCbtMigrationCmd;
+import 
org.apache.cloudstack.api.command.admin.vm.CheckVmwareCbtMigrationPrerequisitesCmd;
+import org.apache.cloudstack.api.command.admin.vm.CutoverVmwareCbtMigrationCmd;
+import org.apache.cloudstack.api.command.admin.vm.DeleteVmwareCbtMigrationCmd;
+import org.apache.cloudstack.api.command.admin.vm.ListVmwareCbtMigrationsCmd;
+import org.apache.cloudstack.api.command.admin.vm.StartVmwareCbtMigrationCmd;
+import org.apache.cloudstack.api.command.admin.vm.SyncVmwareCbtMigrationCmd;
+import org.apache.cloudstack.api.response.ListResponse;
+import 
org.apache.cloudstack.api.response.VmwareCbtMigrationPreflightDiskResponse;
+import 
org.apache.cloudstack.api.response.VmwareCbtMigrationPreflightFindingResponse;
+import org.apache.cloudstack.api.response.VmwareCbtMigrationPreflightResponse;
+import org.apache.cloudstack.api.response.VmwareCbtMigrationCycleResponse;
+import org.apache.cloudstack.api.response.VmwareCbtMigrationDiskResponse;
+import org.apache.cloudstack.api.response.VmwareCbtMigrationResponse;
+import org.apache.cloudstack.context.CallContext;
+import org.apache.cloudstack.framework.config.ConfigKey;
+import org.apache.cloudstack.framework.config.Configurable;
+import org.apache.cloudstack.storage.datastore.db.PrimaryDataStoreDao;
+import org.apache.cloudstack.storage.datastore.db.StoragePoolVO;
+import org.apache.commons.collections4.MapUtils;
+import org.apache.commons.collections4.CollectionUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.springframework.beans.factory.annotation.Autowired;
+
+import com.cloud.agent.AgentManager;
+import com.cloud.agent.api.Answer;
+import com.cloud.agent.api.CheckConvertInstanceAnswer;
+import com.cloud.agent.api.CheckConvertInstanceCommand;
+import com.cloud.agent.api.Command;
+import com.cloud.agent.api.VmwareCbtCleanupCommand;
+import com.cloud.agent.api.VmwareCbtCutoverCommand;
+import com.cloud.agent.api.VmwareCbtMigrationAnswer;
+import com.cloud.agent.api.VmwareCbtPrepareCommand;
+import com.cloud.agent.api.VmwareCbtRbdProbeCommand;
+import com.cloud.agent.api.VmwareCbtSyncCommand;
+import com.cloud.agent.api.to.RemoteInstanceTO;
+import com.cloud.agent.api.to.VmwareCbtChangedBlockRangeTO;
+import com.cloud.agent.api.to.VmwareCbtDiskSyncResultTO;
+import com.cloud.agent.api.to.VmwareCbtDiskTO;
+import com.cloud.agent.api.to.VmwareCbtTargetStorageType;
+import com.cloud.dc.ClusterVO;
+import com.cloud.dc.DataCenterVO;
+import com.cloud.dc.VmwareDatacenterVO;
+import com.cloud.dc.dao.ClusterDao;
+import com.cloud.dc.dao.DataCenterDao;
+import com.cloud.dc.dao.VmwareDatacenterDao;
+import com.cloud.storage.DiskOfferingVO;
+import com.cloud.storage.dao.DiskOfferingDao;
+import com.cloud.exception.AgentUnavailableException;
+import com.cloud.exception.InvalidParameterValueException;
+import com.cloud.exception.OperationTimedoutException;
+import com.cloud.host.Host;
+import com.cloud.host.HostVO;
+import com.cloud.host.Status;
+import com.cloud.host.dao.HostDao;
+import com.cloud.hypervisor.Hypervisor;
+import com.cloud.resource.ResourceState;
+import com.cloud.network.Network;
+import com.cloud.service.ServiceOfferingVO;
+import com.cloud.service.dao.ServiceOfferingDao;
+import com.cloud.storage.Storage;
+import com.cloud.storage.dao.StoragePoolHostDao;
+import com.cloud.user.Account;
+import com.cloud.user.AccountService;
+import com.cloud.user.UserVO;
+import com.cloud.utils.Pair;
+import com.cloud.vm.VmDetailConstants;
+import com.cloud.user.dao.UserDao;
+import com.cloud.uservm.UserVm;
+import com.cloud.vm.UserVmVO;
+import com.cloud.vm.VmwareCbtMigrationCycleVO;
+import com.cloud.vm.VmwareCbtMigrationDiskVO;
+import com.cloud.vm.VmwareCbtMigrationVO;
+import com.cloud.vm.dao.UserVmDao;
+import com.cloud.vm.dao.VmwareCbtMigrationCycleDao;
+import com.cloud.vm.dao.VmwareCbtMigrationDao;
+import com.cloud.vm.dao.VmwareCbtMigrationDiskDao;
+import com.google.gson.Gson;
+import com.google.gson.reflect.TypeToken;
+
+public class VmwareCbtMigrationManagerImpl implements 
VmwareCbtMigrationManager, Configurable {
+
+    /** vCenter refuses snapshot removal while the source is still busy right 
after a cancel;
+     *  retry rather than orphaning the snapshot on the customer's VM. */
+    private static final int SNAPSHOT_REMOVAL_ATTEMPTS = 4;
+    private static final long SNAPSHOT_REMOVAL_RETRY_INTERVAL_MS = 15000L;
+
+    private static final String OBJECT_NAME = "vmwarecbtmigration";
+    private static final String DETAIL_VDDK_TRANSPORTS = "vddk.transports";
+    private static final String DETAIL_VDDK_THUMBPRINT = "vddk.thumbprint";
+    private static final String REDACTED_SECRET = "******";
+    private static final String DEFAULT_CBT_DISK_BASE_PATH = 
"/var/lib/libvirt/images/cloudstack-cbt";
+    private static final String KVM_STORAGE_POOL_MOUNT_BASE_PATH = "/mnt";
+    static final String CBT_FINALIZED_DISK_CONTROLLER = "virtio";
+    private static final String RBD_PROBE_IMAGE_PREFIX = 
"cloudstack-cbt-probe-";
+    private static final List<Storage.StoragePoolType> 
CBT_COMPATIBLE_STORAGE_POOL_TYPES = Arrays.asList(
+            Storage.StoragePoolType.NetworkFilesystem,
+            Storage.StoragePoolType.Filesystem,
+            Storage.StoragePoolType.SharedMountPoint,
+            Storage.StoragePoolType.RBD,
+            Storage.StoragePoolType.Linstor);
+    private static final int CLEANUP_WAIT_SECONDS = 300;
+    private static final int DELETE_CLEANUP_WAIT_SECONDS = 30;
+    private static final Logger LOGGER = 
LogManager.getLogger(VmwareCbtMigrationManagerImpl.class);
+    private static final Gson GSON = new Gson();
+
+    static final ConfigKey<Integer> VmwareCbtMigrationMinCycles = new 
ConfigKey<>(Integer.class,
+            "vmware.cbt.migration.min.cycles",
+            "Advanced",
+            "1",
+            "Minimum number of CBT delta synchronization cycles to run before 
CloudStack can recommend final VMware to KVM cutover",
+            true,
+            ConfigKey.Scope.Global,
+            null);
+
+    static final ConfigKey<Integer> VmwareCbtMigrationMaxCycles = new 
ConfigKey<>(Integer.class,
+            "vmware.cbt.migration.max.cycles",
+            "Advanced",
+            "5",
+            "Maximum number of CBT delta synchronization cycles to run before 
CloudStack recommends final VMware to KVM cutover",
+            true,
+            ConfigKey.Scope.Global,
+            null);
+
+    static final ConfigKey<Integer> VmwareCbtMigrationQuietCycles = new 
ConfigKey<>(Integer.class,
+            "vmware.cbt.migration.quiet.cycles",
+            "Advanced",
+            "2",
+            "Number of consecutive quiet CBT delta synchronization cycles 
required before CloudStack recommends final VMware to KVM cutover",
+            true,
+            ConfigKey.Scope.Global,
+            null);
+
+    static final ConfigKey<Long> VmwareCbtMigrationQuietBytes = new 
ConfigKey<>(Long.class,
+            "vmware.cbt.migration.quiet.bytes",
+            "Advanced",
+            "1073741824",
+            "Maximum changed bytes in a CBT delta synchronization cycle for 
the cycle to be considered quiet",
+            true,
+            ConfigKey.Scope.Global,
+            null);
+
+    static final ConfigKey<Long> VmwareCbtMigrationQuietDirtyRate = new 
ConfigKey<>(Long.class,
+            "vmware.cbt.migration.quiet.dirty.rate",
+            "Advanced",
+            "16777216",
+            "Maximum changed bytes per second in a CBT delta synchronization 
cycle for the cycle to be considered quiet",
+            true,
+            ConfigKey.Scope.Global,
+            null);
+
+    static final ConfigKey<Boolean> VmwareCbtAllowNonInPlaceFinalization = new 
ConfigKey<>(Boolean.class,
+            "vmware.cbt.allow.non.inplace.finalization",
+            "Advanced",
+            "false",
+            "If true, VMware CBT cutover may fall back to regular virt-v2v 
finalization for qcow2 file targets when true in-place finalization is 
unavailable. The fallback stages temporary data on the selected primary storage 
and requires additional free space.",
+            true,
+            ConfigKey.Scope.Global,
+            null);
+
+    static final ConfigKey<Integer> VmwareCbtMigrationAgentCommandTimeout = 
new ConfigKey<>(Integer.class,
+            "vmware.cbt.migration.agent.command.timeout",
+            "Advanced",
+            "86400",
+            "Timeout in seconds for long-running VMware CBT data-plane 
commands dispatched to the KVM agent, including initial full sync, delta sync, 
final delta sync, and cutover finalization.",
+            true,
+            ConfigKey.Scope.Global,
+            null);
+
+    @Inject
+    private VmwareCbtMigrationDao vmwareCbtMigrationDao;
+    @Inject
+    private DataCenterDao dataCenterDao;
+    @Inject
+    private ClusterDao clusterDao;
+    @Inject
+    private HostDao hostDao;
+    @Inject
+    private PrimaryDataStoreDao primaryDataStoreDao;
+    @Inject
+    private VmwareDatacenterDao vmwareDatacenterDao;
+    @Inject
+    private AccountService accountService;
+    @Inject
+    private UserDao userDao;
+    @Inject
+    private ServiceOfferingDao serviceOfferingDao;
+    @Inject
+    private DiskOfferingDao diskOfferingDao;
+    @Inject
+    private UserVmDao userVmDao;
+    @Inject
+    private VmwareCbtMigrationDiskDao vmwareCbtMigrationDiskDao;
+    @Inject
+    private VmwareCbtMigrationCycleDao vmwareCbtMigrationCycleDao;
+    @Inject
+    private StoragePoolHostDao storagePoolHostDao;
+    @Inject
+    private AgentManager agentManager;
+    @Inject
+    private UnmanagedVMsManager unmanagedVMsManager;
+    @Autowired(required = false)
+    private VmwareCbtMigrationService vmwareCbtMigrationService;
+
+    @Override
+    public List<Class<?>> getCommands() {
+        final List<Class<?>> cmdList = new ArrayList<>();
+        cmdList.add(CheckVmwareCbtMigrationPrerequisitesCmd.class);
+        cmdList.add(StartVmwareCbtMigrationCmd.class);
+        cmdList.add(ListVmwareCbtMigrationsCmd.class);
+        cmdList.add(SyncVmwareCbtMigrationCmd.class);
+        cmdList.add(CutoverVmwareCbtMigrationCmd.class);
+        cmdList.add(CancelVmwareCbtMigrationCmd.class);
+        cmdList.add(DeleteVmwareCbtMigrationCmd.class);
+        return cmdList;
+    }
+
+    @Override
+    public VmwareCbtMigrationPreflightResponse 
checkVmwareCbtMigrationPrerequisites(CheckVmwareCbtMigrationPrerequisitesCmd 
cmd) {
+        PreflightFindingCollector findings = new PreflightFindingCollector();
+        VmwareCbtMigrationPreflightResponse response = new 
VmwareCbtMigrationPreflightResponse();
+        response.setObjectName("vmwarecbtmigrationpreflight");
+
+        DataCenterVO zone = getPreflightZone(cmd.getZoneId(), findings);
+        if (zone != null) {
+            response.setZoneId(zone.getUuid());
+            response.setZoneName(zone.getName());
+        }
+
+        ClusterVO destinationCluster = 
getPreflightDestinationCluster(cmd.getClusterId(), zone, findings);
+        if (destinationCluster != null) {
+            response.setClusterId(destinationCluster.getUuid());
+            response.setClusterName(destinationCluster.getName());
+        }
+
+        StoragePoolVO storagePool = 
getPreflightStoragePool(cmd.getStoragePoolId(), zone, destinationCluster, 
findings);
+        VmwareCbtStorageTarget storageTarget = null;
+        if (storagePool != null) {
+            storageTarget = VmwareCbtStorageTarget.forPool(storagePool);
+            populateStorageTargetResponse(response, storageTarget);
+            addStorageTargetFinding(storageTarget, findings);
+        }
+
+        HostVO cbtHost = getPreflightCbtHost(cmd.getConvertInstanceHostId(), 
destinationCluster, storageTarget, findings);
+        if (cbtHost != null) {
+            response.setConvertInstanceHostId(cbtHost.getUuid());
+            response.setConvertInstanceHostName(cbtHost.getName());
+            
response.setConvertInstanceHostInPlaceFinalizationSupported(hostSupportsInPlaceFinalization(cbtHost));
+        }
+        addStorageTargetFinalizationFinding(storageTarget, cbtHost, findings);
+        addRbdStorageAccessFinding(storageTarget, cbtHost, findings);
+
+        String sourceVmName = StringUtils.trimToNull(cmd.getSourceVmName());
+        if (sourceVmName == null) {
+            findings.fail("sourceVm.name.present", "vmware", null, "Source VM 
name is required.");
+        } else {
+            response.setSourceVmName(sourceVmName);
+        }
+
+        VmwareSource source = getPreflightVmwareSource(cmd, findings);
+        if (source != null) {
+            response.setVcenter(source.vcenter);
+            response.setDatacenterName(source.datacenterName);
+            response.setSourceHost(source.sourceHost);
+        }
+
+        if (source != null && sourceVmName != null) {
+            populateSourceVmPreflight(response, source, sourceVmName, cbtHost, 
cmd.getServiceOfferingId(),
+                    cmd.getDetails(), zone, findings);
+        }
+
+        response.setFindings(findings.getFindings());
+        response.setReady(!findings.hasFailures());
+        return response;
+    }
+
+    private VmwareCbtPreflightInfo 
validateSourceVmPreflightForStart(VmwareSource source, String sourceVmName) {
+        VmwareCbtPreflightInfo preflightInfo;
+        try {
+            preflightInfo = 
getVmwareCbtMigrationService().getPreflightInfo(source.vcenter, 
source.datacenterName,
+                    source.username, source.password, source.sourceHost, 
sourceVmName);
+        } catch (RuntimeException e) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Unable to validate VMware CBT prerequisites 
for source VM %s: %s",
+                            sourceVmName, 
sanitizeSensitiveMessage(StringUtils.defaultIfBlank(e.getMessage(), 
e.getClass().getSimpleName()), source)));
+        }
+
+        if (Boolean.FALSE.equals(preflightInfo.getChangeTrackingSupported())) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Source VM %s does not report VMware CBT 
support", sourceVmName));
+        }
+        if (Boolean.TRUE.equals(preflightInfo.getConsolidationNeeded())) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Source VM %s reports pending disk 
consolidation; consolidate VMware disks before starting CBT migration",
+                            sourceVmName));
+        }
+        if (CollectionUtils.isEmpty(preflightInfo.getDisks())) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("No source VMware disks were discovered for 
VM %s", sourceVmName));
+        }
+        for (VmwareCbtPreflightDiskInfo disk : preflightInfo.getDisks()) {
+            validateSourceDiskPreflightForStart(disk);
+        }
+        return preflightInfo;
+    }
+
+    private void ensureSourceVmChangeTrackingEnabledForStart(VmwareSource 
source, String sourceVmName,
+                                                             
VmwareCbtPreflightInfo preflightInfo) {
+        if (Boolean.TRUE.equals(preflightInfo.getChangeTrackingEnabled())) {
+            return;
+        }
+        try {
+            
getVmwareCbtMigrationService().ensureChangeTrackingEnabled(source.vcenter, 
source.datacenterName,
+                    source.username, source.password, source.sourceHost, 
sourceVmName);
+        } catch (RuntimeException e) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Unable to enable VMware CBT on source VM %s 
before initial full sync: %s",
+                            sourceVmName, 
sanitizeSensitiveMessage(StringUtils.defaultIfBlank(e.getMessage(), 
e.getClass().getSimpleName()), source)));
+        }
+    }
+
+    private void 
validateSourceDiskPreflightForStart(VmwareCbtPreflightDiskInfo disk) {
+        String diskId = StringUtils.defaultIfBlank(disk.getSourceDiskId(), 
disk.getSourceDiskPath());
+        if (disk.getSourceDiskDeviceKey() == null) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Source disk %s does not expose a VMware 
device key required for CBT delta queries",
+                            diskId));
+        }
+        if (StringUtils.isBlank(disk.getSourceDiskPath())) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Source disk %s does not expose a VMware 
backing path", diskId));
+        }
+        if (disk.isIndependentDisk()) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Source disk %s uses independent disk mode 
%s, which is not supported for VMware CBT migration",
+                            diskId, disk.getDiskMode()));
+        }
+        if (disk.isPhysicalRdm()) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Source disk %s is a physical-mode RDM, 
which is not supported for VMware CBT migration",
+                            diskId));
+        }
+    }
+
+    private ServiceOfferingVO getServiceOfferingForVmwareCbtMigration(Long 
serviceOfferingId, Account owner, DataCenterVO zone) {
+        if (serviceOfferingId == null) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR, "Service 
offering ID cannot be null");
+        }
+        ServiceOfferingVO serviceOffering = 
serviceOfferingDao.findById(serviceOfferingId);
+        if (serviceOffering == null) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Service offering ID: %d cannot be found", 
serviceOfferingId));
+        }
+        accountService.checkAccess(owner, serviceOffering, zone);
+        serviceOfferingDao.loadDetails(serviceOffering);
+        return serviceOffering;
+    }
+
+    static VmwareCbtOfferingResources 
resolveRequestedOfferingResources(ServiceOfferingVO serviceOffering, 
Map<String, String> details) {
+        if (serviceOffering == null) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR, "Service 
offering cannot be null");
+        }
+        Map<String, String> callerDetails = details == null ? new HashMap<>() 
: details;
+        Map<String, String> offeringDetails = serviceOffering.getDetails();
+        Integer cpu = firstInteger(VmDetailConstants.CPU_NUMBER, callerDetails,
+                serviceOffering.getCpu(), ApiConstants.MIN_CPU_NUMBER, 
offeringDetails);
+        Integer cpuSpeed = firstInteger(VmDetailConstants.CPU_SPEED, 
callerDetails,
+                serviceOffering.getSpeed(), null, offeringDetails);
+        Integer memory = firstInteger(VmDetailConstants.MEMORY, callerDetails,
+                serviceOffering.getRamSize(), ApiConstants.MIN_MEMORY, 
offeringDetails);
+        return new VmwareCbtOfferingResources(cpu, cpuSpeed, memory);
+    }
+
+    private static Integer firstInteger(String callerDetailKey, Map<String, 
String> callerDetails,
+                                        Integer offeringValue, String 
offeringDetailKey,
+                                        Map<String, String> offeringDetails) {
+        if (offeringValue != null) {
+            return offeringValue;
+        }
+        Integer callerValue = parseIntegerDetail(callerDetails, 
callerDetailKey);
+        if (callerValue != null) {
+            return callerValue;
+        }
+        return StringUtils.isBlank(offeringDetailKey) ? null : 
parseIntegerDetail(offeringDetails, offeringDetailKey);
+    }
+
+    private static Integer parseIntegerDetail(Map<String, String> details, 
String key) {
+        if (MapUtils.isEmpty(details) || StringUtils.isBlank(key) || 
StringUtils.isBlank(details.get(key))) {
+            return null;
+        }
+        try {
+            return Integer.valueOf(details.get(key));
+        } catch (NumberFormatException e) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Please provide a valid integer value for 
detail '%s'.", key));
+        }
+    }
+
+    static void 
validateSelectedServiceOfferingResourcesForSourceVm(VmwareCbtPreflightInfo 
preflightInfo,
+                                                                    
ServiceOfferingVO serviceOffering,
+                                                                    
Map<String, String> details) {
+        VmwareCbtOfferingResources requestedResources = 
resolveRequestedOfferingResources(serviceOffering, details);
+        validateRequestedResourceAtLeastSource(preflightInfo.getCpuCores(), 
requestedResources.cpuNumber, "CPU number");
+        validateRequestedResourceAtLeastSource(preflightInfo.getCpuSpeed(), 
requestedResources.cpuSpeed, "CPU speed");
+        validateRequestedResourceAtLeastSource(preflightInfo.getMemoryMb(), 
requestedResources.memoryMb, "Memory");
+    }
+
+    private static void validateRequestedResourceAtLeastSource(Integer 
sourceResource, Integer requestedResource, String resourceName) {
+        if (sourceResource == null || requestedResource == null || 
sourceResource <= 0 || requestedResource <= 0) {
+            return;
+        }
+        if (requestedResource < sourceResource) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("The requested %s (%d) is less than the 
source VM %s (%d)",
+                            resourceName, requestedResource, resourceName, 
sourceResource));
+        }
+    }
+
+    @Override
+    public VmwareCbtMigrationResponse 
startVmwareCbtMigration(StartVmwareCbtMigrationCmd cmd) {
+        Account caller = CallContext.current().getCallingAccount();
+        if (caller == null) {
+            throw new ServerApiException(ApiErrorCode.ACCOUNT_ERROR, "Unable 
to determine calling account");
+        }
+        Account owner = 
accountService.getActiveAccountById(cmd.getEntityOwnerId());
+        if (owner == null) {
+            throw new ServerApiException(ApiErrorCode.ACCOUNT_ERROR, "Unable 
to determine target account");
+        }
+
+        DataCenterVO zone = getZone(cmd.getZoneId());
+        ClusterVO destinationCluster = 
getDestinationCluster(cmd.getClusterId(), zone.getId());
+        StoragePoolVO storagePool = getStoragePool(cmd.getStoragePoolId(), 
zone, destinationCluster);
+        VmwareCbtStorageTarget storageTarget = 
VmwareCbtStorageTarget.forPool(storagePool);
+        if (!storageTarget.isSupported()) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR, 
storageTarget.getSupportMessage());
+        }
+        HostVO convertHost = selectCbtHost(cmd.getConvertInstanceHostId(), 
destinationCluster, storageTarget);
+        validateStorageTargetFinalizationSupport(storageTarget, convertHost);
+        validateRbdStorageAccessForStart(storageTarget, convertHost);
+
+        String sourceVmName = StringUtils.trimToNull(cmd.getSourceVmName());
+        if (sourceVmName == null) {
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR, "Source VM 
name is required");
+        }
+
+        VmwareSource source = resolveVmwareSource(cmd);
+        VmwareCbtPreflightInfo preflightInfo = 
validateSourceVmPreflightForStart(source, sourceVmName);
+        ServiceOfferingVO serviceOffering = 
getServiceOfferingForVmwareCbtMigration(cmd.getServiceOfferingId(), owner, 
zone);
+        validateSelectedServiceOfferingResourcesForSourceVm(preflightInfo, 
serviceOffering, cmd.getDetails());
+        validateWindowsGuestConversionSupportForStart(convertHost, 
sourceVmName, preflightInfo, cmd.getDetails());
+        ensureSourceVmChangeTrackingEnabledForStart(source, sourceVmName, 
preflightInfo);
+        List<VmwareCbtDiskInfo> sourceDisks = discoverSourceDisks(source, 
sourceVmName);
+        validateDataDiskOfferingMappingsForStart(sourceVmName, sourceDisks, 
cmd.getDataDiskToDiskOfferingList(), owner, zone);
+        validateNicMappingsForStart(sourceVmName, preflightInfo.getNics(), 
cmd.getNicNetworkList(), cmd.getNicIpAddressList());
+        String displayName = StringUtils.defaultIfBlank(cmd.getDisplayName(), 
sourceVmName);
+
+        VmwareCbtMigrationVO migration = new 
VmwareCbtMigrationVO(zone.getId(), owner.getId(), getUserIdForOwner(owner),
+                destinationCluster.getId(), displayName, source.vcenter, 
source.datacenterName, cmd.getSourceHost(), cmd.getSourceCluster(), 
sourceVmName);
+        migration.setExistingVcenterId(source.existingVcenterId);
+        storeExternalVmwareSourceCredentials(migration, source);
+        if (convertHost != null) {
+            migration.setConvertHostId(convertHost.getId());
+        }
+        if (storagePool != null) {
+            migration.setStoragePoolId(storagePool.getId());
+        }
+        migration.setHostName(StringUtils.trimToNull(cmd.getHostName()));
+        migration.setTemplateId(cmd.getTemplateId());
+        migration.setServiceOfferingId(cmd.getServiceOfferingId());
+        // Derive the guest OS and hardware details from the source VM's 
configuration when the
+        // caller did not choose them: without this the cutover import 
inherits the dummy import
+        // template's generic OS type and the KVM defaults (BIOS firmware, 
cirrus video), which
+        // leaves a UEFI source unbootable and gives Windows guests Linux 
clock semantics.
+        migration.setGuestOsId(cmd.getGuestOsId() != null ? cmd.getGuestOsId()
+                : 
unmanagedVMsManager.resolveGuestOsIdForVmwareImport(preflightInfo.getOperatingSystem(),
 preflightInfo.getOperatingSystemId()));
+        
migration.setDataDiskOfferingMap(serializeMap(cmd.getDataDiskToDiskOfferingList()));
+        migration.setNicNetworkMap(serializeMap(cmd.getNicNetworkList()));
+        // Preserve static source IPs captured during preflight: NICs mapped 
to a network but
+        // given no explicit IP keep their guest-configured address when it 
fits the target
+        // network. The source is running at this point, so Tools data is 
available.
+        Map<String, List<String>> sourceNicCidrs = new HashMap<>();
+        for (VmwareCbtPreflightNicInfo nic : preflightInfo.getNics()) {
+            if (StringUtils.isNotBlank(nic.getSourceNicId()) && 
CollectionUtils.isNotEmpty(nic.getIpv4Cidrs())) {
+                sourceNicCidrs.put(nic.getSourceNicId(), nic.getIpv4Cidrs());
+            }
+        }
+        
migration.setNicIpAddressMap(serializeNicIpAddressMap(unmanagedVMsManager.autoFillStaticNicIpAddresses(
+                cmd.getNicNetworkList(), cmd.getNicIpAddressList(), 
sourceNicCidrs)));
+        
migration.setImportDetails(serializeMap(unmanagedVMsManager.applyVmwareImportHardwareDetails(cmd.getDetails(),
+                preflightInfo.getBootType(), preflightInfo.getBootMode(), 
preflightInfo.getOperatingSystem())));
+        migration.setForced(cmd.isForced());
+        applyVddkDetails(migration, cmd.getDetails());
+        migration.setState(VmwareCbtMigration.State.InitialSync);
+        migration.setCurrentStep(String.format("Discovered %s source disk(s); 
preparing initial VDDK full sync", sourceDisks.size()));
+        migration.setUpdated(new Date());
+        migration = vmwareCbtMigrationDao.persist(migration);
+        persistSourceDisks(migration, sourceDisks);
+        return runInitialFullSync(source, migration, storageTarget);
+    }
+
+    @Override
+    public ListResponse<VmwareCbtMigrationResponse> 
listVmwareCbtMigrations(ListVmwareCbtMigrationsCmd cmd) {
+        VmwareCbtMigration.State state = parseState(cmd.getState());
+        Pair<List<VmwareCbtMigrationVO>, Integer> result = 
vmwareCbtMigrationDao.listMigrations(cmd.getId(), cmd.getZoneId(), 
cmd.getAccountId(),
+                cmd.getVcenter(), cmd.getSourceVmName(), state, 
cmd.getStartIndex(), cmd.getPageSizeVal());
+
+        List<VmwareCbtMigrationResponse> responses = new ArrayList<>();
+        for (VmwareCbtMigrationVO migration : result.first()) {
+            responses.add(createVmwareCbtMigrationResponse(migration));
+        }
+        ListResponse<VmwareCbtMigrationResponse> listResponse = new 
ListResponse<>();
+        listResponse.setResponses(responses, result.second());
+        return listResponse;
+    }
+
+    @Override
+    public VmwareCbtMigrationResponse 
syncVmwareCbtMigration(SyncVmwareCbtMigrationCmd cmd) {
+        VmwareCbtMigrationVO migration = getMigration(cmd.getId());
+        rejectTerminalMigration(migration, "synchronize");
+        requireMigrationState(migration, "synchronize", 
VmwareCbtMigration.State.Replicating,
+                VmwareCbtMigration.State.ReadyForCutover);
+        VmwareSource source = resolveVmwareSource(migration, 
cmd.getUsername(), cmd.getPassword());
+        validateInitialSyncTargetDisks(migration);
+        HostVO cbtHost = getCbtHostForMigration(migration);
+        int cycleNumber = migration.getCompletedCycles() + 1;
+        VmwareCbtMigrationCycleVO cycle = new 
VmwareCbtMigrationCycleVO(migration.getId(), cycleNumber);
+        cycle.setState(VmwareCbtMigrationCycle.State.CopyingChangedBlocks);
+        cycle.setDescription("Dispatching CBT delta synchronization to KVM 
agent");
+        cycle.setUpdated(new Date());
+        cycle = vmwareCbtMigrationCycleDao.persist(cycle);
+
+        migration.setState(VmwareCbtMigration.State.Replicating);
+        migration.setCurrentStep(String.format("Running CBT delta 
synchronization cycle %s", cycleNumber));
+        migration.setLastError(null);
+        migration.setUpdated(new Date());
+        migration = persistMigrationProgress(migration, String.format("start 
of CBT delta cycle %s", cycleNumber));
+
+        VmwareCbtSnapshotInfo snapshot = null;
+        try {
+            snapshot = createDeltaSnapshot(source, migration, cycleNumber);
+            cycle.setState(VmwareCbtMigrationCycle.State.QueryingChangedAreas);
+            cycle.setSnapshotMor(snapshot.getSnapshotMor());
+            cycle.setDescription("Querying VMware CBT changed disk areas");
+            cycle.setUpdated(new Date());
+            vmwareCbtMigrationCycleDao.update(cycle.getId(), cycle);
+
+            VmwareCbtChangedBlockQueryResult changedBlockQuery = 
queryChangedBlocks(source, migration,
+                    snapshot.getSnapshotMor());
+
+            cycle.setState(VmwareCbtMigrationCycle.State.CopyingChangedBlocks);
+            cycle.setDescription(String.format("Dispatching %s VMware CBT 
changed block range(s) to KVM agent",
+                    changedBlockQuery.changedBlocks.size()));
+            cycle.setUpdated(new Date());
+            vmwareCbtMigrationCycleDao.update(cycle.getId(), cycle);
+
+            VmwareCbtSyncCommand syncCommand = new 
VmwareCbtSyncCommand(migration.getUuid(),
+                    createRemoteInstance(source, migration, 
snapshot.getSourceVmMor()), getDiskTransferObjects(migration),
+                    changedBlockQuery.changedBlocks, cycleNumber, 
snapshot.getSnapshotMor(), false);
+            applyVddkDetails(syncCommand, migration);
+            applyTargetStorageDetails(syncCommand, migration);
+            syncCommand.setWait(getVmwareCbtMigrationAgentCommandTimeout());
+
+            VmwareCbtMigrationAnswer answer = sendVmwareCbtCommand(cbtHost, 
syncCommand, "synchronize",
+                    migration.getUuid());
+            if (!answer.getResult()) {
+                String details = sanitizeSensitiveMessage(answer.getDetails(), 
source);
+                markCycleFailed(cycle, details);
+                markMigrationFailed(migration, "CBT delta synchronization 
failed", details);
+                return 
createVmwareCbtMigrationResponse(vmwareCbtMigrationDao.findById(migration.getId()));
+            }
+
+            applyDiskResults(migration, answer.getDiskResults());
+            updateDiskChangeIds(migration, changedBlockQuery.changedDisks);
+            long changedBytes = answer.getChangedBytes();
+            long durationSeconds = answer.getDurationSeconds();
+            long dirtyRate = getDirtyRateBytesPerSecond(changedBytes, 
durationSeconds, answer.getDirtyRateBytesPerSecond());
+            VmwareCbtMigrationCutoverPolicy cutoverPolicy = getCutoverPolicy();
+            VmwareCbtMigrationCutoverPolicy.Decision cutoverDecision = 
cutoverPolicy.decide(cycleNumber,
+                    migration.getQuietCycles(), changedBytes, durationSeconds);
+            int quietCycles = cutoverPolicy.isQuietCycle(changedBytes, 
durationSeconds) ?
+                    migration.getQuietCycles() + 1 : 0;
+
+            cycle.setState(VmwareCbtMigrationCycle.State.Completed);
+            cycle.setChangedBytes(changedBytes);
+            cycle.setDirtyRate(dirtyRate);
+            cycle.setDuration(durationSeconds * 1000);
+            cycle.setDescription(sanitizeSensitiveMessage(answer.getDetails(), 
source));
+            cycle.setUpdated(new Date());
+            vmwareCbtMigrationCycleDao.update(cycle.getId(), cycle);
+
+            migration.setCompletedCycles(cycleNumber);
+            migration.setQuietCycles(quietCycles);
+            migration.setTotalChangedBytes(migration.getTotalChangedBytes() + 
changedBytes);
+            migration.setLastChangedBytes(changedBytes);
+            migration.setLastDirtyRate(dirtyRate);
+            migration.setState(cutoverDecision == 
VmwareCbtMigrationCutoverPolicy.Decision.CONTINUE ?
+                    VmwareCbtMigration.State.Replicating : 
VmwareCbtMigration.State.ReadyForCutover);
+            migration.setCurrentStep(getCutoverDecisionStep(cutoverDecision));
+            migration.setUpdated(new Date());
+            migration = persistMigrationProgress(migration, "CBT delta 
synchronization progress");
+            return createVmwareCbtMigrationResponse(migration);
+        } catch (RuntimeException e) {
+            String error = 
sanitizeSensitiveMessage(StringUtils.defaultIfBlank(e.getMessage(), 
e.getClass().getSimpleName()), source);
+            markCycleFailed(cycle, error);
+            markMigrationFailed(migration, "CBT delta synchronization failed", 
error);
+            return 
createVmwareCbtMigrationResponse(vmwareCbtMigrationDao.findById(migration.getId()));
+        } finally {
+            removeDeltaSnapshotIfPossible(source, migration, snapshot);
+        }
+    }
+
+    @Override
+    public VmwareCbtMigrationResponse 
cutoverVmwareCbtMigration(CutoverVmwareCbtMigrationCmd cmd) {
+        VmwareCbtMigrationVO migration = getMigration(cmd.getId());
+        rejectTerminalMigration(migration, "cut over");
+        requireMigrationState(migration, "cut over", 
VmwareCbtMigration.State.ReadyForCutover,
+                VmwareCbtMigration.State.ReadyForImport);
+        VmwareSource source = resolveVmwareSource(migration, 
cmd.getUsername(), cmd.getPassword());
+        validateInitialSyncTargetDisks(migration);
+        if (migration.getState() == VmwareCbtMigration.State.ReadyForImport) {
+            return importCutoverMigration(source, migration);
+        }
+        HostVO cbtHost = getCbtHostForMigration(migration);
+        VmwareCbtStorageTarget storageTarget = 
getStorageTargetForMigration(migration);
+        validateStorageTargetFinalizationSupport(storageTarget, cbtHost);
+        requireSourceVmPoweredOff(source, migration);
+
+        migration.setState(VmwareCbtMigration.State.CuttingOver);
+        migration.setCurrentStep("Running final CBT delta synchronization 
before cutover");
+        migration.setLastError(null);
+        migration.setUpdated(new Date());
+        if (!vmwareCbtMigrationDao.updateIfNotTerminal(migration)) {
+            // Cancelled between the state check above and here: refuse rather 
than cut over a
+            // migration whose target disks the cancellation has already 
removed.
+            VmwareCbtMigrationVO current = 
vmwareCbtMigrationDao.findById(migration.getId());
+            throw new ServerApiException(ApiErrorCode.PARAM_ERROR,
+                    String.format("Cannot cut over VMware CBT migration %s: it 
reached state %s while the cutover was starting.",
+                            migration.getUuid(), current == null ? "unknown" : 
current.getState()));
+        }
+
+        int finalCycleNumber = migration.getCompletedCycles() + 1;
+        if (!runFinalDeltaSync(source, migration, cbtHost, finalCycleNumber)) {
+            return 
createVmwareCbtMigrationResponse(vmwareCbtMigrationDao.findById(migration.getId()));
+        }
+
+        migration = vmwareCbtMigrationDao.findById(migration.getId());
+        boolean inPlaceFinalization = hostSupportsInPlaceFinalization(cbtHost);
+        String finalizationDescription = inPlaceFinalization ? "in-place 
conversion" : "virt-v2v fallback conversion";
+        migration.setCurrentStep(String.format("Final CBT delta 
synchronization completed; running %s",
+                finalizationDescription));
+        migration.setUpdated(new Date());
+        migration = persistMigrationProgress(migration, "final delta 
synchronization progress");
+
+        VmwareCbtCutoverCommand cutoverCommand = new 
VmwareCbtCutoverCommand(migration.getUuid(), createRemoteInstance(source, 
migration),
+                getDiskTransferObjects(migration), finalCycleNumber, true);
+        applyVddkDetails(cutoverCommand, migration);
+        applyTargetStorageDetails(cutoverCommand, migration);
+        
cutoverCommand.setAllowNonInPlaceFinalization(isNonInPlaceFinalizationFallbackAllowed(storageTarget));
+        cutoverCommand.setWait(getVmwareCbtMigrationAgentCommandTimeout());
+
+        VmwareCbtMigrationAnswer answer = sendVmwareCbtCommand(cbtHost, 
cutoverCommand, "cut over", migration.getUuid());

Review Comment:
   We should consider wrapping this in a try-catch block, as the method can 
throw an exception. Without proper handling, we may lose the opportunity to 
return a meaningful response that includes the underlying exception details and 
context.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to