diff --git a/test/integration/plugins/ontap/migration/AGENTS.md b/test/integration/plugins/ontap/migration/AGENTS.md new file mode 100644 index 000000000000..e003d1a4ca92 --- /dev/null +++ b/test/integration/plugins/ontap/migration/AGENTS.md @@ -0,0 +1,328 @@ + + +# Agent Guide: ONTAP KVM Live VM Migration + +Always-on operating notes for agents working on KVM live VM migration and +live VM-with-storage migration for NetApp ONTAP primary storage. + +Ticket: **CSTACKEX-262**. Branch: `feature/CSTACKEX-262`. + +Read this before editing Java, Marvin tests, or the lab. Credentials, SVM name, +and host URLs come from `test/integration/plugins/ontap/ontap.cfg`. Do not +hard-code them in test code. + +Deeper design lives in the Cursor skill +`~/.cursor/skills/cloudstack-ontap-migration/` (`SKILL.md`, `combinations.md`, +`code-paths.md`, `specs/02-live-vm-with-storage-migration.md`). + +--- + +## Confluence and specs + +| Doc | URL | pageId | +|---|---|---| +| TOI: NetApp ONTAP Storage Plugin | https://netapp.atlassian.net/wiki/spaces/OSSG/pages/624932561 | `624932561` | +| Admin Guide | https://netapp.atlassian.net/wiki/spaces/OSSG/pages/669425954 | `669425954` | +| Integrated Spec | https://netapp.atlassian.net/wiki/spaces/OSSG/pages/330407273 | `330407273` | +| Architectural Spec | https://netapp.atlassian.net/wiki/spaces/OSSG/pages/330407201 | `330407201` | +| Plugin overview | https://netapp.atlassian.net/wiki/spaces/OSSG/pages/633873247 | `633873247` | +| Feature roadmap | https://netapp.atlassian.net/wiki/spaces/OSSG/database/640141625 | `640141625` | +| CI/CD pipeline | https://netapp.atlassian.net/wiki/spaces/OSSG/pages/689709001 | `689709001` | +| Operations / scale test matrix | https://netapp.atlassian.net/wiki/spaces/OSSG/pages/648390908 | `648390908` | + +`cloudId`: `cb69b23c-616e-461b-9cf1-b3015880f8dd` (`netapp.atlassian.net`). + +Apache CloudStack APIs: + +- Compute-only live migrate: `migrateVirtualMachine` (Running; volumes stay) +- Live with storage: `migrateVirtualMachineWithVolume` (Running; qemu/libvirt) +- Offline volume: `migrateVolume` (`livemigrate=true` is **rejected** on KVM) + +--- + +## Working agreement + +1. **KVM only.** Leave VMware/Xen branches untouched. +2. **Live with-storage stays qemu/libvirt.** Do not route Running + `migrateVirtualMachineWithVolume` through agent `CopyCommand`. +3. **ONTAP-to-ONTAP offline copy** (Stopped / detached) may use `CopyCommand` + only for the **same protocol and same SVM** + (`StorageSystemDataMotionStrategy.isSupportedOntapMigrationPoolPair`). +4. **Never add a hook to `PrimaryDataStoreDriver`.** Use an inline + `DataStoreProvider.ONTAP_PLUGIN_NAME` check. +5. **Surgical edits.** No drive-by refactors. Ask before Java production + changes unless the user already consented this cycle. +6. **User-owned full `mvn install`** unless they said otherwise. Agents run + targeted unit tests and Marvin. After Java on the KVM agent path, rebuild + Debian packages and replace `cloudstack-agent` on **both** KVM hosts. +7. **Do not invent lab credentials.** Read `ontap.cfg`. + +--- + +## RTP lab (read live values from ontap.cfg) + +Typical layout used for this work: + +| Role | Host | Notes | +|---|---|---| +| Management + KVM Cluster1 | `10.193.56.62` (cstack53) | jetty `:cloud-client-ui`; agent; NFS `/export/primary`, `/export/secondary` | +| KVM Cluster2 | `10.193.56.63` (cstack54) | agent only | +| Integration API | `http://10.193.56.62:8096` | `integration.api.port=8096` | +| UI / API | `http://10.193.56.62:8080` | jetty | +| ONTAP | `storageIP` + `svmName` in `ontap.cfg` | vs0 in the current cfg | +| DefaultPrimary | `nfs://10.193.56.62/export/primary` | tag `defaultPrim` | +| Template | CentOS 5.5(64-bit) no GUI (KVM) | must be Ready on Primary1 | + +Hosts: Cluster1 `kvmHost`, Cluster2 `kvmHost1`. Advanced isolated guest network. + +Indirect agent: KVM agents **connect outbound** to MS `:8250`. They often do +**not** listen on 8250 themselves. Host Up is ping over that channel. + +--- + +## What must work (live) + +Data mover for Running VMs with storage: `handleLiveMigrationForKVM` → +`MigrateCommand` → qemu block copy into dest volumes created by the ONTAP +driver. + +| Source | Dest | Allowed | +|---|---|---| +| ONTAP NFS3 (same SVM) | ONTAP NFS3 | Yes — dest file QCOW2 | +| ONTAP iSCSI (same SVM) | ONTAP iSCSI | Yes — dest LUN RAW | +| DefaultPrimary / non-managed | ONTAP NFS3 or iSCSI | Yes | +| ONTAP NFS3 | ONTAP iSCSI (or reverse) | No | +| ONTAP | different SVM / storageIP | No | +| ONTAP | non-managed | No | +| Mixed managed + non-managed dests | — | No | +| Cluster-scope pool | Zone-scope pool (same protocol/SVM) | Yes; dest host must reach dest pool | +| Cluster1-only pool | host in Cluster2 | No (`test_22`) | +| `migrateVolume livemigrate=true` | — | Reject; use with-volume VM migrate | + +Cluster→zone is supported. Zone pools are visible to every cluster in the zone. + +NFS dest files must be created as **QCOW2** +(`UnifiedNASStrategy.createVolumeOnKVMHost`). RAW dest files fail live migrate +with `Image is not in qcow2 format`. + +iSCSI LUN IDs are **per igroup**. After `grantAccess`, refresh `VolumeInfo` +before putting dest TOs on `PrepareForMigrationCommand`. + +--- + +## Code map + +| Path | Why | +|---|---| +| `engine/storage/datamotion/.../StorageSystemDataMotionStrategy.java` | `handleLiveMigrationForKVM`, `verifyLiveMigrationForKVM`, `isSupportedOntapMigrationPoolPair`, offline `CopyCommand` branch | +| `server/.../vm/VirtualMachineManagerImpl.java` | `migrateWithStorage`, `executeManagedStorageChecksWhenTargetStoragePoolProvided` | +| `server/.../vm/UserVmManagerImpl.java` | `migrateVirtualMachineWithVolume` | +| `engine/orchestration/.../VolumeOrchestrator.java` | `migrateVolumes` | +| `plugins/storage/volume/ontap/.../OntapPrimaryDatastoreDriver.java` | create/grant/revoke | +| `plugins/storage/volume/ontap/.../UnifiedNASStrategy.java` | QCOW2 stamp on KVM create | +| `plugins/storage/volume/ontap/.../UnifiedSANStrategy.java` | LUN map / igroup | +| `plugins/hypervisors/kvm/.../LibvirtMigrateCommandWrapper.java` | libvirt live migrate | +| `plugins/hypervisors/kvm/.../KVMStorageProcessor.java` | CopyCommand / volume copy on agent | +| `test/integration/plugins/ontap/migration/test_01_live_vm_with_storage_migration.py` | live with-storage matrix | +| `test/integration/plugins/ontap/migration/migration_test_base.py` | hosts, pools, deploy, migrate helpers | +| `ontap-migration-test-plan.xlsx` | manual matrix (repo root) | + +Do not widen `handleLiveMigrationForKVM` for other vendors. Do not set +`canCopy` / `CAN_CREATE_VOLUME_FROM_VOLUME` on the ONTAP driver. + +--- + +## Marvin suites + +| Tag | File | +|---|---| +| `live_storage` | `test_01_live_vm_with_storage_migration.py` (22 cases; 01–13 same cluster; 14–22 cross-cluster) | +| `stopped_vm` | `test_02_stopped_vm_storage_migration.py` | +| `volume_migration` | `test_03_volume_migration.py` | + +Same-cluster live needs **two Up hosts in Cluster1**. Cross-cluster needs one +host in Cluster2. Tests chain VMs: if deploy fails, later cases SKIP. + +Pool list APIs often omit `storageIP` / `svmName`. Matching must treat missing +details as compatible and accept type `OntapiSCSI` as iSCSI. + +--- + +## Commands + +All from CloudStack repo root unless noted. Python: `test/integration/plugins/ontap/.venv/bin/python`. + +### Unit tests + +```bash +mvn -pl engine/storage/datamotion test -Dtest=StorageSystemDataMotionStrategyTest +mvn -pl plugins/storage/volume/ontap test -Dtest=UnifiedNASStrategyTest,UnifiedSANStrategyTest,OntapPrimaryDatastoreDriverTest +mvn -pl plugins/hypervisors/kvm test -Dtest=KVMStorageProcessorTest +``` + +Corrupt local `aspectjweaver` (`Invalid CEN header`) is an artifact problem, +not the change. + +### Marvin + +```bash +export PYTHONPATH=test/integration/plugins/ontap:${PYTHONPATH:-} +export PYTHONUNBUFFERED=1 +export ONTAP_MIGRATION_KEEP_POOLS=1 # reuse Up ONTAP pools; omit to recreate + +bash test/integration/plugins/ontap/run_tests.sh live_storage +bash test/integration/plugins/ontap/run_tests.sh stopped_vm +bash test/integration/plugins/ontap/run_tests.sh volume_migration +bash test/integration/plugins/ontap/run_tests.sh migration + +test/integration/plugins/ontap/.venv/bin/python -m py_compile \ + test/integration/plugins/ontap/migration/migration_test_base.py \ + test/integration/plugins/ontap/migration/test_01_live_vm_with_storage_migration.py +``` + +Logs: `/tmp/MarvinLogs/` (newest folder + `test_01_live_vm_with_storage_migration_*`). +Read `results.txt` and `failed_plus_exceptions.txt`. + +Needs full network for Marvin init (`listUsers`). SSH to KVM hosts uses +`ontap.cfg` host passwords. + +### Management server (cstack53) + +Source tree is typically `/root/cloudstack`. Integration API is 8096. + +```bash +# start (on the MS host) +export MAVEN_OPTS="-Xmx3072m -XX:MaxMetaspaceSize=512m" +cd /root/cloudstack +# already running: ps aux | grep 'jetty:run' +nohup mvn -Dorg.eclipse.jetty.annotations.maxWait=120 -pl :cloud-client-ui jetty:run \ + > /root/cloudstack/jetty-run.log 2>&1 & + +mysql -uroot cloud -e "UPDATE configuration SET value='8096' WHERE name='integration.api.port';" +# restart jetty after changing integration.api.port if it was not already 8096 + +curl -sS 'http://127.0.0.1:8096/client/api?command=listHosts&type=Routing&response=json' +curl -sS 'http://127.0.0.1:8096/client/api?command=listStoragePools&response=json' +curl -sS 'http://127.0.0.1:8096/client/api?command=listVirtualMachines&listall=true&response=json' +curl -sS 'http://127.0.0.1:8096/client/api?command=listRouters&listall=true&response=json' +curl -sS 'http://127.0.0.1:8096/client/api?command=listAsyncJobs&listall=true&response=json' +``` + +MS logs: `/root/cloudstack/vmops.log`, `/root/cloudstack/api.log`. + +### Debian packages and replace KVM agent + +Build **on the Ubuntu 22.04 management host** so packages match the KVM OS. +Install **both** `cloudstack-common` and `cloudstack-agent` (same version) on +**every** KVM host (Cluster1 and Cluster2). + +```bash +cd /root/cloudstack +export MAVEN_OPTS="-Xmx3072m -XX:MaxMetaspaceSize=512m" +mvn -P developer,systemvm -DskipTests clean install +dpkg-buildpackage -us -uc -b +# fallback: bash packaging/build-deb.sh + +ls -1t ../cloudstack-common_*_all.deb ../cloudstack-agent_*_all.deb | head +``` + +On each KVM host (`10.193.56.62` and `10.193.56.63`): + +```bash +systemctl stop cloudstack-agent +dpkg -i /tmp/cloudstack-common_*_all.deb /tmp/cloudstack-agent_*_all.deb +# if needed: apt-get install -fy +systemctl enable cloudstack-agent +systemctl restart cloudstack-agent +systemctl is-active cloudstack-agent +ss -lntp | grep 8250 || true # may be empty; agent often outbound-only +pgrep -af com.cloud.agent.AgentShell +tail -n 80 /var/log/cloudstack/agent/agent.log +``` + +Preserve `/etc/cloudstack/agent/agent.properties` (`host`, `port`, `guid`). +Restart libvirtd only if live migrate needs TCP listen (`listen_tcp=1`; this +lab often uses `16514`). + +Copy debs from MS: + +```bash +scp ../cloudstack-common_*_all.deb ../cloudstack-agent_*_all.deb root@10.193.56.63:/tmp/ +``` + +### Host / libvirt / disk + +```bash +virsh list --all +virsh capabilities | head +df -h /export/primary /var/lib/libvirt/images / +uptime +# hung template copy (blocks deploy): +ps -eo pid,etime,pcpu,cmd | grep -E 'cp -f /mnt/|qemu-img' | grep -v grep +ls -lh /export/primary +``` + +### ONTAP REST (cluster IP and user from ontap.cfg) + +```bash +# SVM +curl -sk -u "$USER:$PASS" "https://$ONTAP/api/svm/svms?name=$SVM" +# volumes / LUNs (filter in client) +curl -sk -u "$USER:$PASS" "https://$ONTAP/api/storage/volumes?svm.name=$SVM&max_records=1000" +curl -sk -u "$USER:$PASS" "https://$ONTAP/api/storage/luns?svm.name=$SVM&max_records=1000" +curl -sk -u "$USER:$PASS" "https://$ONTAP/api/protocols/san/igroups?svm.name=$SVM" +``` + +--- + +## Failure patterns (do not misdiagnose) + +| Symptom | Likely cause | +|---|---| +| `Image is not in qcow2 format` | Dest NFS file created RAW; QCOW2 stamp missing | +| API `NullPointerException` on live with-storage | Check `vmops.log`. May be **cleanup NPE after a real error**. test_19 was `migration of disk vdb failed: No space left on device` then `VolumeObject.stateTransit` NPE on already-deleted dest volume | +| `KVMStoragePool.getType() because pool is null` | Dest storage pool not in libvirt on dest host (`PrepareForMigration`) | +| Deploy hung `Starting`, ROOT `Allocated` | Isolated VR not Running, or template `cp` secondary→Primary1 stuck; MS load high | +| `Cannot stop VM ... state Starting` | Pending `DeployVMCmd` job; destroy/expunge blocked until job ends | +| `createStoragePool` hangs 30+ min | Host attaching NFS; or tests creating extra pools because match required omitted details | +| Relocate/addHost fails, 8250 not listening | Normal for outbound agent; check `AgentShell` + host Up, not listen socket | +| Host maintenance: VMs in starting/stopping | Stuck VR/user VMs on that host; do not loop relocate | +| Cluster1 live 01–13 SKIP no dest host | Need two hosts in Cluster1; one host in Cluster2 is not enough | +| Cross-protocol live mapping | Expected reject: managed storage can only be migrated to itself | + +Before another full Marvin run: no stuck `cp` of the CentOS template, no +`Starting` VMs/routers, Primary1 has space, dest host disk not ENOSPC, both +agents Up, jetty on 8096. + +Best proven live_storage run: **01–18 and 22 pass**; **19** failed ENOSPC on +cstack54 during iSCSI cluster→zone onto Cluster2 (20–21 cascade). Stopped-VM +and volume suites were not green in that cycle. + +--- + +## Java consent and deploy + +If a Marvin failure is an implementation bug: + +1. Quote the `vmops.log` stack, not only Marvin `errortext`. +2. Ask before editing Java unless already approved this session. +3. After agent-side Java: DEB install on **both** KVM hosts; jetty restart + only if management code changed. +4. Re-run the failing Marvin tag with `ONTAP_MIGRATION_KEEP_POOLS=1`. diff --git a/test/integration/plugins/ontap/migration/__init__.py b/test/integration/plugins/ontap/migration/__init__.py new file mode 100644 index 000000000000..13a83393a912 --- /dev/null +++ b/test/integration/plugins/ontap/migration/__init__.py @@ -0,0 +1,16 @@ +# 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. diff --git a/test/integration/plugins/ontap/migration/migration_test_base.py b/test/integration/plugins/ontap/migration/migration_test_base.py new file mode 100644 index 000000000000..cd9a8dbbf415 --- /dev/null +++ b/test/integration/plugins/ontap/migration/migration_test_base.py @@ -0,0 +1,2563 @@ +# 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. + +"""Shared infrastructure for the ordered ONTAP migration Marvin suites. + +Before test methods run, the base inventories the configured zone and prepares +the ONTAP pool scopes and cluster attachments requested by each suite. +""" + +import base64 +import logging +import os +import random +import re +import shlex +import time +import unittest + +from marvin.cloudstackAPI import ( + addHost as addHostAPI, + attachVolume as attachVolumeAPI, + createNetwork as createNetworkAPI, + createStoragePool as createStoragePoolAPI, + createVolume as createVolumeAPI, + deleteHost as deleteHostAPI, + deleteNetwork as deleteNetworkAPI, + deleteVolume as deleteVolumeAPI, + deployVirtualMachine as deployVirtualMachineAPI, + destroyRouter as destroyRouterAPI, + destroySystemVm as destroySystemVmAPI, + destroyVirtualMachine as destroyVirtualMachineAPI, + listClusters as listClustersAPI, + listNetworkOfferings as listNetworkOfferingsAPI, + listNetworks as listNetworksAPI, + listRouters as listRoutersAPI, + listSystemVms as listSystemVmsAPI, + listVirtualMachines as listVirtualMachinesAPI, + listVolumes as listVolumesAPI, + migrateVirtualMachine as migrateVirtualMachineAPI, + migrateVirtualMachineWithVolume as migrateVMWithVolumeAPI, + migrateVolume as migrateVolumeAPI, + startSystemVm as startSystemVmAPI, + startVirtualMachine as startVirtualMachineAPI, + stopSystemVm as stopSystemVmAPI, + stopVirtualMachine as stopVirtualMachineAPI, +) +from marvin.lib.base import DiskOffering, Host, ServiceOffering, StoragePool +from marvin.lib.common import list_storage_pools +from marvin.sshClient import SshClient + +from ontap_test_base import ( + OntapRestClient, + OntapTestBase, + _parse_pool_details, + get_datacenter_config, + get_ready_hosts, + is_configured, + list_kvm_templates, +) + +logger = logging.getLogger("OntapMigrationTestBase") +DISK_SOURCE_RE = re.compile( + r"""]*(?:file|dev)=['"]([^'"]+)['"]""", + re.IGNORECASE, +) + + +class OntapMigrationTestBase(OntapTestBase): + """Common setup, resource creation, assertions, and force cleanup.""" + + config_data = None + template_id = None + network_id = None + _created_network_id = None + ready_hosts = None + created_pools = None + created_volumes = None + created_vms = None + service_offerings = None + disk_offerings = None + default_pool = None + ontap_pools = None + hosts_by_cluster = None + secondary_cluster_id = None + + @classmethod + def setUpClass(cls): + super(OntapMigrationTestBase, cls).setUpClass() + testclient = super( + OntapMigrationTestBase, cls + ).getClsTestClient() + cls.apiClient = testclient.getApiClient() + cls.dbConnection = testclient.getDbConnection() + cls.config_data = get_datacenter_config(testclient, cls) + + ontap_cfg = cls.config_data.get("ontap", {}) + cls.ontap = OntapRestClient( + ontap_cfg.get("storageIP", ""), + ontap_cfg.get("username", ""), + ontap_cfg.get("password", ""), + ) + cls.svm_name = ontap_cfg.get("svmName", "") + account_data = { + "email": "ontap-migration@test.invalid", + "firstname": "ONTAP", + "lastname": "Migration", + "username": "ontap_migration_%d" % random.randint(0, 999999), + "password": "password", + } + cls._setup_cloudstack_resources(cls.config_data, account_data) + configured_cluster = cls.config_data.get( + "cloudstack", {} + ).get("clusterName") + if (configured_cluster + and cls.cluster.name != configured_cluster): + raise RuntimeError( + "Configured primary cluster %s resolved to %s" + % (configured_cluster, cls.cluster.name) + ) + host_requirements = cls._host_requirements() + cls._configure_host_states(host_requirements["minimum_hosts"]) + cls._ensure_host_topology(host_requirements) + _, cls.ready_hosts = cls._validate_migration_prerequisites( + cls.config_data, host_requirements["minimum_hosts"] + ) + cls.hosts_by_cluster = {} + for host in cls.ready_hosts: + cls.hosts_by_cluster.setdefault(host.clusterid, []).append(host) + cls._validate_host_topology(host_requirements) + cls._validate_host_credentials() + if host_requirements.get("live_migration"): + cls._prepare_live_migration_hosts() + other_cluster_ids = [ + cluster_id for cluster_id in cls.hosts_by_cluster + if cluster_id != cls.cluster.id + ] + cls.secondary_cluster_id = ( + other_cluster_ids[0] if other_cluster_ids else None + ) + + cls.created_pools = [] + cls.created_volumes = [] + cls.created_vms = [] + cls.service_offerings = {} + cls.disk_offerings = {} + try: + cls.default_pool = cls._find_default_primary_pool() + cls.ontap_pools = cls._ensure_migration_pools() + logger.info( + "Migration pool matrix: DefaultPrimary=%s, NFS3=%s, ISCSI=%s", + cls.default_pool.name, + { + key: [pool.name for pool in pools] + for key, pools in cls.ontap_pools["NFS3"].items() + }, + { + key: [pool.name for pool in pools] + for key, pools in cls.ontap_pools["ISCSI"].items() + }, + ) + cls.template_id = cls._find_template() + cls.network_id = cls._find_or_create_network() + except Exception: + cls.tearDownClass() + raise + + @classmethod + def _host_requirements(cls): + return { + "minimum_hosts": 1, + "same_primary_cluster": False, + "multiple_clusters": False, + "live_migration": False, + } + + @classmethod + def _configure_host_states(cls, minimum_hosts): + hosts = Host.list( + cls.apiClient, + zoneid=cls.zone.id, + type="Routing", + hypervisor="KVM", + ) or [] + for host in hosts: + resource_state = str( + getattr(host, "resourcestate", "") + ).lower() + if resource_state in ("maintenance", "errorinmaintenance"): + Host.cancelMaintenance(cls.apiClient, host.id) + elif resource_state != "enabled": + Host.update( + cls.apiClient, + id=host.id, + resourcestate="Enabled", + ) + host_state = str(getattr(host, "state", "")).lower() + if host_state in ("alert", "disconnected", "down"): + try: + Host.reconnect(cls.apiClient, id=host.id) + except Exception as error: + logger.warning( + "Host %s reconnect was not accepted: %s", + host.id, error, + ) + + for attempt in range(12): + current = Host.list( + cls.apiClient, + zoneid=cls.zone.id, + type="Routing", + hypervisor="KVM", + ) or [] + if len(get_ready_hosts(current)) >= minimum_hosts: + return + if attempt < 11: + time.sleep(5) + raise RuntimeError("Configured KVM hosts did not reach Up/Enabled") + + @classmethod + def _ready_hosts_by_cluster(cls): + hosts = get_ready_hosts(Host.list( + cls.apiClient, + zoneid=cls.zone.id, + type="Routing", + hypervisor="KVM", + ) or []) + grouped = {} + for host in hosts: + grouped.setdefault(host.clusterid, []).append(host) + return grouped + + @classmethod + def _ensure_host_topology(cls, requirements): + """Move a KVM host between clusters when the suite needs it.""" + if requirements["same_primary_cluster"]: + cls._ensure_hosts_in_primary_cluster(2) + if requirements["multiple_clusters"]: + cls._ensure_hosts_across_clusters() + + @classmethod + def _ensure_hosts_in_primary_cluster(cls, minimum_hosts): + grouped = cls._ready_hosts_by_cluster() + while len(grouped.get(cls.cluster.id, [])) < minimum_hosts: + donor = next( + (hosts[0] for cluster_id, hosts in grouped.items() + if cluster_id != cls.cluster.id and hosts), + None, + ) + if donor is None: + return + before = len(grouped.get(cls.cluster.id, [])) + cls._relocate_host(donor, cls.cluster.id) + grouped = cls._ready_hosts_by_cluster() + if len(grouped.get(cls.cluster.id, [])) <= before: + logger.warning( + "Could not move host %s into the primary cluster; " + "continuing with %s host(s) there", + donor.name, before, + ) + return + + @classmethod + def _ensure_hosts_across_clusters(cls): + grouped = cls._ready_hosts_by_cluster() + if len(grouped) >= 2: + return + donors = grouped.get(cls.cluster.id, []) + target_cluster_id = cls._find_alternate_cluster() + if len(donors) < 2 or target_cluster_id is None: + return + cls._relocate_host(donors[-1], target_cluster_id) + + @classmethod + def _refresh_host_inventory(cls): + cls.ready_hosts = get_ready_hosts(Host.list( + cls.apiClient, + zoneid=cls.zone.id, + type="Routing", + hypervisor="KVM", + ) or []) + cls.hosts_by_cluster = {} + for host in cls.ready_hosts: + cls.hosts_by_cluster.setdefault(host.clusterid, []).append(host) + other_cluster_ids = [ + cluster_id for cluster_id in cls.hosts_by_cluster + if cluster_id != cls.cluster.id + ] + cls.secondary_cluster_id = ( + other_cluster_ids[0] if other_cluster_ids else None + ) + + @classmethod + def _find_alternate_cluster(cls): + cmd = listClustersAPI.listClustersCmd() + cmd.zoneid = cls.zone.id + cmd.hypervisor = "KVM" + clusters = cls.apiClient.listClusters(cmd) or [] + configured_names = [] + for zone in cls.config_data.get("zones", []): + if zone.get("name") != cls.zone.name: + continue + for pod in zone.get("pods", []): + configured_names.extend( + cluster.get("clustername") + for cluster in pod.get("clusters", []) + if cluster.get("clustername") != cls.cluster.name + ) + for name in configured_names: + configured = next( + (cluster for cluster in clusters if cluster.name == name), + None, + ) + if configured is not None: + return configured.id + return next( + (cluster.id for cluster in clusters + if cluster.id != cls.cluster.id), + None, + ) + + @classmethod + def _host_credentials(cls, host): + """Return the ontap.cfg host entry that matches a listed host.""" + identifiers = { + str(getattr(host, attribute, "") or "") + for attribute in ("ipaddress", "name") + } + identifiers.discard("") + for zone in cls.config_data.get("zones", []): + for pod in zone.get("pods", []): + for cluster in pod.get("clusters", []): + for entry in cluster.get("hosts", []): + endpoint = str( + entry.get("url", "") + ).rsplit("/", 1)[-1] + if endpoint and endpoint in identifiers: + return entry + return None + + @classmethod + def _relocate_host(cls, host, target_cluster_id): + """Remove a KVM host from its cluster and re-add it to another.""" + credentials = cls._host_credentials(host) + if credentials is None: + logger.warning( + "No configured credentials for host %s; leaving it in " + "cluster %s", host.name, host.clusterid, + ) + return + if not cls._wait_for_kvm_agent_listen(credentials): + logger.warning( + "KVM agent on %s is not listening; leaving host in cluster %s", + host.name, host.clusterid, + ) + return + logger.info( + "Relocating host %s from cluster %s to cluster %s", + host.name, host.clusterid, target_cluster_id, + ) + original_cluster_id = host.clusterid + pod_id = host.podid + stopped_system_vms = [] + try: + still_busy = cls._force_clear_host_workloads(host.id) + if still_busy: + logger.warning( + "Host %s still has Starting/Stopping VMs; leaving it in " + "cluster %s", + host.name, host.clusterid, + ) + return + stopped_system_vms = cls._stop_system_vms_on_host(host.id) + cls._prepare_host_for_relocation(host.id) + except Exception as exc: + logger.warning( + "Could not prepare host %s for relocation; leaving it in " + "cluster %s: %s", + host.name, host.clusterid, exc, + ) + cls._start_system_vms(stopped_system_vms) + return + + if not cls._stop_kvm_agent(credentials): + logger.warning( + "Could not stop the KVM agent on %s before relocation; " + "leaving it in cluster %s", + host.name, host.clusterid, + ) + cls._start_system_vms(stopped_system_vms) + return + + delete_cmd = deleteHostAPI.deleteHostCmd() + delete_cmd.id = host.id + delete_cmd.forced = True + delete_cmd.forcedestroylocalstorage = True + cls.apiClient.deleteHost(delete_cmd) + + added = cls._add_kvm_host(credentials, pod_id, target_cluster_id) + if added is None: + logger.warning( + "addHost of %s to cluster %s failed; restoring cluster %s", + host.name, target_cluster_id, original_cluster_id, + ) + added = cls._add_kvm_host( + credentials, pod_id, original_cluster_id + ) + if added is None: + cls._start_system_vms(stopped_system_vms) + return + cls._wait_for_host_up(added[0].id) + cls._ensure_libvirt_tcp_listen(credentials) + cls._wait_for_host_up(added[0].id) + cls._start_system_vms(stopped_system_vms) + + @classmethod + def _add_kvm_host(cls, credentials, pod_id, cluster_id): + add_cmd = addHostAPI.addHostCmd() + add_cmd.zoneid = cls.zone.id + add_cmd.podid = pod_id + add_cmd.clusterid = cluster_id + add_cmd.hypervisor = "KVM" + add_cmd.url = credentials["url"] + add_cmd.username = credentials["username"] + add_cmd.password = credentials["password"] + if credentials.get("hosttags"): + add_cmd.hosttags = credentials["hosttags"] + try: + return cls.apiClient.addHost(add_cmd) + except Exception as exc: + logger.warning( + "addHost %s to cluster %s failed: %s", + credentials.get("url"), cluster_id, exc, + ) + return None + + @classmethod + def _wait_for_kvm_agent_listen(cls, credentials, timeout=15): + """Return whether the KVM agent is reachable for addHost.""" + return cls._kvm_agent_listening(credentials) + + @classmethod + def _kvm_agent_listening(cls, credentials): + endpoint = str(credentials.get("url", "")).rsplit("/", 1)[-1] + if not endpoint: + return False + try: + ssh = SshClient( + endpoint, 22, + credentials.get("username", "root"), + credentials.get("password", ""), + retries=2, delay=2, timeout=10.0, + ) + ssh.execute("mkdir -p /vmware-vix-disklib-distrib") + output = ssh.execute( + "pgrep -f com.cloud.agent.AgentShell >/dev/null " + "&& echo AGENT_OK || true" + ) + text = " ".join(str(line) for line in (output or [])) + return "AGENT_OK" in text + except Exception as exc: + logger.warning( + "Could not check KVM agent on %s: %s", endpoint, exc + ) + return False + + @classmethod + def _stop_kvm_agent(cls, credentials): + """Stop the old agent before addHost rewrites and restarts it.""" + endpoint = str(credentials.get("url", "")).rsplit("/", 1)[-1] + if not endpoint: + return False + try: + ssh = SshClient( + endpoint, 22, + credentials.get("username", "root"), + credentials.get("password", ""), + retries=3, delay=3, timeout=15.0, + ) + ssh.execute("systemctl stop cloudstack-agent") + output = ssh.execute( + "pgrep -f '[c]om.cloud.agent.AgentShell' >/dev/null " + "|| echo AGENT_STOPPED" + ) + text = " ".join(str(line) for line in (output or [])) + return "AGENT_STOPPED" in text + except Exception as exc: + logger.warning( + "Could not stop KVM agent on %s: %s", endpoint, exc + ) + return False + + @classmethod + def _ensure_libvirt_tcp_listen(cls, credentials): + """Live migration needs TCP libvirtd on the destination host.""" + endpoint = str(credentials.get("url", "")).rsplit("/", 1)[-1] + if not endpoint: + return + try: + ssh = SshClient( + endpoint, 22, + credentials.get("username", "root"), + credentials.get("password", ""), + retries=3, delay=3, timeout=15.0, + ) + ssh.execute( + "sed -i 's/^listen_tls=.*/listen_tls=0/' " + "/etc/libvirt/libvirtd.conf; " + "sed -i 's/^listen_tcp=.*/listen_tcp=1/' " + "/etc/libvirt/libvirtd.conf; " + "systemctl restart libvirtd" + ) + except Exception as exc: + logger.warning( + "Could not enable libvirt TCP listen on %s: %s", + endpoint, exc, + ) + + @classmethod + def _force_clear_host_workloads(cls, host_id): + """Destroy stuck Starting VMs/routers so host maintenance can proceed.""" + list_cmd = listVirtualMachinesAPI.listVirtualMachinesCmd() + list_cmd.hostid = host_id + list_cmd.listall = True + for vm in cls.apiClient.listVirtualMachines(list_cmd) or []: + state = str(getattr(vm, "state", "")).lower() + if state not in ("starting", "stopping", "error", "unknown"): + continue + logger.warning( + "Destroying stuck %s VM %s on host %s before relocation", + state, vm.id, host_id, + ) + destroy_cmd = destroyVirtualMachineAPI.destroyVirtualMachineCmd() + destroy_cmd.id = vm.id + destroy_cmd.expunge = True + try: + cls.apiClient.destroyVirtualMachine(destroy_cmd) + except Exception as exc: + logger.warning( + "Could not destroy stuck VM %s: %s", vm.id, exc + ) + router_cmd = listRoutersAPI.listRoutersCmd() + router_cmd.hostid = host_id + router_cmd.listall = True + for router in cls.apiClient.listRouters(router_cmd) or []: + state = str(getattr(router, "state", "")).lower() + if state == "running": + continue + logger.warning( + "Destroying stuck %s router %s on host %s before relocation", + state, router.id, host_id, + ) + destroy_cmd = destroyRouterAPI.destroyRouterCmd() + destroy_cmd.id = router.id + try: + cls.apiClient.destroyRouter(destroy_cmd) + except Exception as exc: + logger.warning( + "Could not destroy stuck router %s: %s", router.id, exc + ) + return cls._host_has_transitioning_vms(host_id) + + @classmethod + def _host_has_transitioning_vms(cls, host_id): + list_cmd = listVirtualMachinesAPI.listVirtualMachinesCmd() + list_cmd.hostid = host_id + list_cmd.listall = True + for vm in cls.apiClient.listVirtualMachines(list_cmd) or []: + if str(getattr(vm, "state", "")).lower() in ( + "starting", "stopping" + ): + return True + router_cmd = listRoutersAPI.listRoutersCmd() + router_cmd.hostid = host_id + router_cmd.listall = True + for router in cls.apiClient.listRouters(router_cmd) or []: + if str(getattr(router, "state", "")).lower() in ( + "starting", "stopping" + ): + return True + return False + + @classmethod + def _stop_system_vms_on_host(cls, host_id): + list_cmd = listSystemVmsAPI.listSystemVmsCmd() + list_cmd.hostid = host_id + system_vms = cls.apiClient.listSystemVms(list_cmd) or [] + running_ids = [ + system_vm.id for system_vm in system_vms + if str(getattr(system_vm, "state", "")).lower() == "running" + ] + for system_vm_id in running_ids: + logger.info( + "Stopping system VM %s before relocating host %s", + system_vm_id, host_id, + ) + stop_cmd = stopSystemVmAPI.stopSystemVmCmd() + stop_cmd.id = system_vm_id + stop_cmd.forced = True + cls.apiClient.stopSystemVm(stop_cmd) + return running_ids + + @classmethod + def _start_system_vms(cls, system_vm_ids): + for system_vm_id in system_vm_ids: + logger.info( + "Restarting system VM %s after host relocation", + system_vm_id, + ) + start_cmd = startSystemVmAPI.startSystemVmCmd() + start_cmd.id = system_vm_id + try: + cls.apiClient.startSystemVm(start_cmd) + except Exception as error: + logger.warning( + "System VM %s did not restart after host relocation; " + "destroying it so CloudStack can recreate it: %s", + system_vm_id, error, + ) + destroy_cmd = destroySystemVmAPI.destroySystemVmCmd() + destroy_cmd.id = system_vm_id + try: + cls.apiClient.destroySystemVm(destroy_cmd) + except Exception as destroy_error: + logger.warning( + "System VM %s could not be destroyed after its " + "restart failed; continuing because host relocation " + "already completed: %s", + system_vm_id, destroy_error, + ) + + @classmethod + def _prepare_host_for_relocation(cls, host_id): + for attempt in range(6): + try: + Host.enableMaintenance(cls.apiClient, id=host_id) + cls._wait_for_host_resource_state( + host_id, "maintenance", timeout=300 + ) + return + except Exception as exc: + hosts = Host.list( + cls.apiClient, id=host_id, listall=True + ) or [] + state = str( + getattr(hosts[0], "resourcestate", "") + if hosts else "" + ).lower() + if (state == "errorinmaintenance" + and attempt < 5): + logger.warning( + "Host %s entered ErrorInMaintenance; cancelling and " + "retrying", host_id, + ) + Host.cancelMaintenance(cls.apiClient, host_id) + cls._wait_for_host_resource_state( + host_id, "enabled", timeout=300 + ) + continue + if ("starting/stopping state" in str(exc).lower() + and attempt < 5): + logger.warning( + "Host %s has a transitioning VM; retrying " + "maintenance", host_id, + ) + time.sleep(15) + continue + else: + raise + + @classmethod + def _wait_for_host_resource_state(cls, host_id, resource_state, + timeout=600): + deadline = time.time() + timeout + while time.time() < deadline: + hosts = Host.list(cls.apiClient, id=host_id, listall=True) or [] + current_state = str( + getattr(hosts[0], "resourcestate", "") if hosts else "" + ).lower() + if current_state == resource_state: + return + if (resource_state == "maintenance" + and current_state == "errorinmaintenance"): + raise RuntimeError( + "Host %s entered ErrorInMaintenance" % host_id + ) + time.sleep(10) + raise RuntimeError( + "Host %s did not reach resource state %s" + % (host_id, resource_state) + ) + + @classmethod + def _wait_for_host_up(cls, host_id, timeout=900): + deadline = time.time() + timeout + consecutive_ready_checks = 0 + while time.time() < deadline: + hosts = Host.list(cls.apiClient, id=host_id, listall=True) or [] + host = hosts[0] if hosts else None + if (host is not None + and str(getattr(host, "state", "")).lower() == "up" + and str(getattr( + host, "resourcestate", "" + )).lower() == "enabled"): + consecutive_ready_checks += 1 + if consecutive_ready_checks >= 3: + return + else: + consecutive_ready_checks = 0 + time.sleep(5) + raise RuntimeError("Relocated host %s did not come back Up" % host_id) + + @classmethod + def _validate_host_topology(cls, requirements): + if (requirements["same_primary_cluster"] + and len(cls.hosts_by_cluster.get(cls.cluster.id, [])) < 2): + raise unittest.SkipTest( + "Suite requires two Up/Enabled KVM hosts in primary cluster " + "%s; found %s" + % ( + cls.cluster.name, + len(cls.hosts_by_cluster.get(cls.cluster.id, [])), + ) + ) + if (requirements["multiple_clusters"] + and ( + not cls.hosts_by_cluster.get(cls.cluster.id) + or len(cls.hosts_by_cluster) < 2 + )): + raise unittest.SkipTest( + "Suite requires an Up/Enabled KVM host in primary cluster %s " + "and another configured cluster" + % cls.cluster.name + ) + + @classmethod + def _validate_host_credentials(cls): + missing = [ + host.name for host in cls.ready_hosts + if cls._host_credentials(host) is None + ] + if missing: + raise RuntimeError( + "ontap.cfg has no SSH credentials for KVM hosts: %s" + % ", ".join(sorted(missing)) + ) + + @classmethod + def _ssh_output(cls, host, command): + credentials = cls._host_credentials(host) + if credentials is None: + raise RuntimeError( + "No configured SSH credentials for host %s" % host.name + ) + endpoint = str(credentials.get("url", "")).rsplit("/", 1)[-1] + try: + ssh = SshClient( + endpoint, 22, + credentials.get("username", "root"), + credentials.get("password", ""), + retries=2, delay=2, timeout=15.0, + ) + output = ssh.execute(command) or [] + return "\n".join(str(line) for line in output) + except Exception as exc: + raise RuntimeError( + "SSH command failed on KVM host %s: %s" + % (host.name, exc) + ) + + @classmethod + def _prepare_live_migration_hosts(cls): + for host in cls.ready_hosts: + credentials = cls._host_credentials(host) + cls._ensure_libvirt_tcp_listen(credentials) + output = cls._ssh_output( + host, + "virsh version >/dev/null 2>&1 " + "&& ss -lnt | grep -Eq ':(16509|16514)[[:space:]]' " + "&& echo LIVE_MIGRATION_READY", + ) + if "LIVE_MIGRATION_READY" not in output: + raise RuntimeError( + "KVM host %s is not ready for libvirt TCP migration" + % host.name + ) + # Restarting libvirtd drops the KVM agent. A pool created while a + # host is reconnecting is exported only to the hosts that are Up + # at that moment, so wait until every host is Up again first. + for host in cls.ready_hosts: + cls._wait_for_host_up(host.id) + + @classmethod + def _validate_migration_prerequisites( + cls, config, minimum_kvm_hosts=2): + """Validate the single ONTAP SVM and two-host migration lab.""" + issues = [] + cloudstack_cfg = config.get("cloudstack", {}) + zone_name = cloudstack_cfg.get("zoneName") + if not is_configured(zone_name): + issues.append("cloudstack.zoneName") + + ontap_cfg = config.get("ontap", {}) or {} + for field in ("storageIP", "svmName", "username", "password"): + if not is_configured(ontap_cfg.get(field)): + issues.append("ontap.%s" % field) + + configured_hosts = [] + for zone in config.get("zones", []): + if is_configured(zone_name) and zone.get("name") != zone_name: + continue + for pod in zone.get("pods", []): + for cluster in pod.get("clusters", []): + configured_hosts.extend(cluster.get("hosts", [])) + if len(configured_hosts) < minimum_kvm_hosts: + issues.append( + "configured KVM host count (need %d, found %d)" + % (minimum_kvm_hosts, len(configured_hosts)) + ) + for index, host in enumerate(configured_hosts): + for field in ("url", "username", "password"): + if not is_configured(host.get(field)): + issues.append( + "configured KVM host %d.%s" % (index, field) + ) + + zone_hosts = Host.list( + cls.apiClient, + zoneid=cls.zone.id, + type="Routing", + hypervisor="KVM", + ) or [] + migration_hosts = get_ready_hosts(zone_hosts) + if len(migration_hosts) < minimum_kvm_hosts: + issues.append( + "Up/Enabled KVM host count (need %d, found %d)" + % (minimum_kvm_hosts, len(migration_hosts)) + ) + + if issues: + raise RuntimeError( + "Migration test prerequisites are incomplete: %s" + % ", ".join(issues) + ) + + try: + OntapRestClient( + ontap_cfg["storageIP"], + ontap_cfg["username"], + ontap_cfg["password"], + ).check_connection() + except Exception as ex: + raise RuntimeError("ONTAP is not ready: %s" % ex) + return ontap_cfg, migration_hosts + + @classmethod + def _find_default_primary_pool(cls): + pools = list_storage_pools( + cls.apiClient, zoneid=cls.zone.id, clusterid=cls.cluster.id + ) or [] + default_pools = [ + pool for pool in pools + if str(getattr(pool, "provider", "")).lower() == "defaultprimary" + and getattr(pool, "state", "") == "Up" + ] + expected_tag = cls._default_storage_tag() + if expected_tag: + for pool in default_pools: + pool_tags = { + tag.strip() + for tag in str(getattr(pool, "tags", "")).split(",") + if tag.strip() + } + if expected_tag in pool_tags: + logger.info( + "Using configured DefaultPrimary pool '%s' (id=%s).", + pool.name, pool.id, + ) + return pool + if default_pools: + pool = default_pools[0] + logger.info( + "Using existing DefaultPrimary pool '%s' (id=%s).", + pool.name, pool.id, + ) + return pool + raise unittest.SkipTest( + "No Up default primary storage pool exists in the test cluster" + ) + + @classmethod + def _pool_protocol(cls, pool): + details = { + str(key).lower(): value + for key, value in cls._pool_details(pool).items() + } + protocol = str(details.get("protocol", "")).upper() + if protocol in ("NFS", "NFS3"): + return "NFS3" + if protocol == "ISCSI": + return "ISCSI" + pool_type = str(getattr(pool, "type", "")).lower() + if pool_type == "networkfilesystem": + return "NFS3" + if pool_type in ("iscsilun", "ontapiscsi"): + return "ISCSI" + return None + + @classmethod + def _pool_details(cls, pool): + details = dict(_parse_pool_details(pool)) + required = {"protocol", "storageIP", "svmName"} + if required.issubset(details): + return details + if cls.dbConnection is None: + return details + rows = cls.dbConnection.execute( + "SELECT storage_pool_details.name, storage_pool_details.value " + "FROM storage_pool_details " + "JOIN storage_pool " + "ON storage_pool.id = storage_pool_details.pool_id " + "WHERE storage_pool.uuid = %s", + (pool.id,), + ) or [] + details.update({name: value for name, value in rows}) + return details + + @classmethod + def _is_compatible_ontap_pool( + cls, pool, protocol, scope=None, cluster_id=None): + ontap_cfg = cls.config_data["ontap"] + provider = str(getattr(pool, "provider", "")).lower() + if provider != str( + cls.config_data["storagePool"].get( + "storagePoolProvider", "NetApp ONTAP" + )).lower(): + return False + if str(getattr(pool, "state", "")).lower() != "up": + return False + if cls._pool_protocol(pool) != protocol: + return False + if scope and str(getattr(pool, "scope", "")).upper() != scope: + return False + if (scope == "CLUSTER" + and getattr(pool, "clusterid", None) != cluster_id): + return False + expected_tag = cls._storage_tag(protocol) + pool_tags = { + tag.strip() + for tag in str(getattr(pool, "tags", "")).split(",") + if tag.strip() + } + if expected_tag and expected_tag not in pool_tags: + return False + details = { + str(key).lower(): str(value) + for key, value in cls._pool_details(pool).items() + } + storage_ip = details.get("storageip") + svm_name = details.get("svmname") + if storage_ip != str(ontap_cfg["storageIP"]): + return False + if svm_name != str(ontap_cfg["svmName"]): + return False + reported = ( + getattr(pool, "capacitybytes", None) + or getattr(pool, "disksizetotal", None) + ) + if reported is not None: + try: + if int(reported) < int( + cls._migration_pool_capacity_bytes() * 0.9 + ): + return False + except (TypeError, ValueError): + return False + return True + + @classmethod + def _migration_pool_capacity_bytes(cls): + return int(cls.config_data["storagePool"]["capacitybytes"]) + + @classmethod + def _pool_requirements(cls): + return { + "NFS3": [("CLUSTER", cls.cluster.id, 2)], + "ISCSI": [("CLUSTER", cls.cluster.id, 2)], + } + + @classmethod + def _pool_key(cls, scope, cluster_id=None): + return ( + "ZONE" if scope == "ZONE" + else "CLUSTER:%s" % cluster_id + ) + + @classmethod + def _ensure_migration_pools(cls): + """Discover or create each protocol/scope/cluster pool requirement.""" + pools = list_storage_pools(cls.apiClient, zoneid=cls.zone.id) or [] + selected = {"NFS3": {}, "ISCSI": {}} + for protocol, requirements in cls._pool_requirements().items(): + cls._protocol_config(protocol) + for scope, cluster_id, count in requirements: + key = cls._pool_key(scope, cluster_id) + selected[protocol][key] = [ + pool for pool in pools + if cls._is_compatible_ontap_pool( + pool, protocol, scope, cluster_id + ) + ][:count] + for pool in selected[protocol][key]: + logger.info( + "Using ONTAP %s %s pool '%s' (id=%s).", + protocol, key, pool.name, pool.id, + ) + if (pool.name.startswith("OntapMigration") + and all( + item.id != pool.id + for item in cls.created_pools + )): + cls.created_pools.append(pool) + while len(selected[protocol][key]) < count: + pool = cls._create_ontap_pool( + protocol, scope, cluster_id + ) + selected[protocol][key].append(pool) + pools.append(pool) + if len(selected[protocol][key]) != count: + raise RuntimeError( + "Expected %d ONTAP %s pools for %s, found %d" + % ( + count, protocol, key, + len(selected[protocol][key]), + ) + ) + for pool in selected[protocol][key]: + cls._validate_migration_pool( + pool, protocol, scope, cluster_id + ) + return selected + + @classmethod + def _validate_migration_pool( + cls, pool, protocol, scope, cluster_id=None): + if not cls._is_compatible_ontap_pool( + pool, protocol, scope, cluster_id): + raise RuntimeError( + "ONTAP pool %s does not match %s %s prerequisites" + % (pool.name, protocol, cls._pool_key(scope, cluster_id)) + ) + backend = cls.ontap.get_volume(pool.name) + if backend is None: + raise RuntimeError( + "ONTAP FlexVol %s for storage pool %s does not exist" + % (pool.name, pool.id) + ) + if str(backend.get("state", "")).lower() != "online": + raise RuntimeError( + "ONTAP FlexVol %s is not online" % pool.name + ) + + @classmethod + def _migration_pool( + cls, protocol, index=0, scope="CLUSTER", cluster_id=None): + protocol = protocol.upper() + cls._protocol_config(protocol) + if scope == "CLUSTER" and cluster_id is None: + cluster_id = cls.cluster.id + key = cls._pool_key(scope, cluster_id) + pools = cls.ontap_pools.get(protocol, {}).get(key, []) + if len(pools) <= index: + raise unittest.SkipTest( + "ONTAP %s %s pool %d is unavailable" + % (protocol, key, index + 1) + ) + return pools[index] + + def _pool_by_id(self, pool_id): + pools = list_storage_pools(self.apiClient, id=pool_id) or [] + self.assertTrue(pools, "Storage pool %s was not found" % pool_id) + return pools[0] + + @classmethod + def _find_template(cls): + ready = [ + template for template in list_kvm_templates( + cls.apiClient, cls.zone.id + ) + if getattr(template, "isready", False) + and str(getattr(template, "templatetype", "")).upper() != "SYSTEM" + ] + if not ready: + raise unittest.SkipTest("No ready user KVM template is available") + configured_name = cls.config_data.get( + "cloudstack", {} + ).get("templateName") + configured = [ + template for template in ready + if getattr(template, "name", None) == configured_name + ] + return (configured or ready)[0].id + + @classmethod + def _find_or_create_network(cls): + network_type = str( + getattr(cls.zone, "networktype", "Basic") + ).lower() + if network_type != "advanced": + return None + cmd = listNetworksAPI.listNetworksCmd() + cmd.zoneid = cls.zone.id + cmd.account = cls.account.name + cmd.domainid = cls.domain.id + networks = cls.apiClient.listNetworks(cmd) or [] + if networks: + return networks[0].id + + offering_cmd = listNetworkOfferingsAPI.listNetworkOfferingsCmd() + offering_cmd.state = "Enabled" + offering_cmd.guestiptype = "Isolated" + offering_cmd.specifyvlan = False + offerings = cls.apiClient.listNetworkOfferings(offering_cmd) or [] + offering = next( + ( + item for item in offerings + if "SourceNat" in item.name + and "Vpc" not in item.name + and "NSX" not in item.name + and "Netris" not in item.name + ), + offerings[0] if offerings else None, + ) + if offering is None: + raise unittest.SkipTest( + "No enabled isolated network offering is available" + ) + create_cmd = createNetworkAPI.createNetworkCmd() + create_cmd.zoneid = cls.zone.id + create_cmd.networkofferingid = offering.id + create_cmd.name = "ontap-migration-%d" % random.randint(0, 999999) + create_cmd.displaytext = "ONTAP migration test network" + create_cmd.account = cls.account.name + create_cmd.domainid = cls.domain.id + network = cls.apiClient.createNetwork(create_cmd) + cls._created_network_id = network.id + return network.id + + @classmethod + def _recreate_test_network(cls): + if cls._created_network_id: + deadline = time.time() + 180 + while True: + cmd = deleteNetworkAPI.deleteNetworkCmd() + cmd.id = cls._created_network_id + try: + cls.apiClient.deleteNetwork(cmd) + break + except Exception as exc: + if time.time() >= deadline: + raise + logger.warning( + "Could not delete network %s; retrying: %s", + cls._created_network_id, + exc, + ) + time.sleep(5) + cls._created_network_id = None + cls.network_id = cls._find_or_create_network() + + @classmethod + def _protocol_config(cls, protocol): + key = "nfs3" if protocol.upper() == "NFS3" else "iscsi" + config = cls.config_data.get( + "storagePool", {} + ).get("protocols", {}).get(key, {}) + if not config.get("enabled", False): + raise unittest.SkipTest("%s is disabled in ontap.cfg" % protocol) + return config + + @classmethod + def _storage_tag(cls, protocol): + return cls._protocol_config(protocol).get("storagePoolTags") + + @classmethod + def _default_storage_tag(cls): + for zone in cls.config_data.get("zones", []): + if zone.get("name") != cls.zone.name: + continue + for pod in zone.get("pods", []): + for cluster in pod.get("clusters", []): + if cluster.get("clustername") != cls.cluster.name: + continue + for pool in cluster.get("primaryStorages", []): + if pool.get("provider") == "DefaultPrimary": + return pool.get("tags") + return getattr(cls.default_pool, "tags", None) + + @classmethod + def _service_offering(cls): + if "shared" in cls.service_offerings: + return cls.service_offerings["shared"] + data = { + "name": "ontap-migration-so-%d" % random.randint(0, 999999), + "displaytext": "ONTAP migration service offering", + "cpunumber": 1, + "cpuspeed": 100, + "memory": 256, + "storagetype": "shared", + } + offering = ServiceOffering.create(cls.apiClient, data) + cls.service_offerings["shared"] = offering + cls._cleanup.insert(0, offering) + return offering + + @classmethod + def _root_disk_offering(cls, storage_tag): + key = storage_tag or "__untagged__" + if key in cls.disk_offerings: + return cls.disk_offerings[key] + data = { + "name": "ontap-migration-root-do-%d" + % random.randint(0, 999999), + "displaytext": "ONTAP migration root disk offering", + "disksize": 8, + "storagetype": "shared", + } + if storage_tag: + data["tags"] = storage_tag + offering = DiskOffering.create(cls.apiClient, data) + cls.disk_offerings[key] = offering + cls._cleanup.insert(0, offering) + return offering + + @classmethod + def _create_ontap_pool(cls, protocol, scope, cluster_id=None): + config = cls.config_data + ontap_cfg = config["ontap"] + pool_cfg = config["storagePool"] + name = "OntapMigration%s_%d" % ( + protocol.upper(), random.randint(0, 999999) + ) + details = { + "username": ontap_cfg["username"], + "password": base64.b64encode( + ontap_cfg["password"].encode() + ).decode(), + "svmName": ontap_cfg["svmName"], + "protocol": protocol.upper(), + "storageIP": ontap_cfg["storageIP"], + } + + cmd = createStoragePoolAPI.createStoragePoolCmd() + cmd.name = name + scheme = "nfs" if protocol.upper() == "NFS3" else "iscsi" + cmd.url = "%s://%s/ontap" % (scheme, ontap_cfg["storageIP"]) + cmd.zoneid = cls.zone.id + if scope == "CLUSTER": + cmd.podid = cls.cluster.podid + cmd.clusterid = cluster_id + cmd.scope = scope + cmd.provider = pool_cfg.get( + "storagePoolProvider", "NetApp ONTAP" + ) + cmd.tags = cls._storage_tag(protocol) + cmd.capacitybytes = cls._migration_pool_capacity_bytes() + cmd.hypervisor = "KVM" + cmd.managed = True + for index, (key, value) in enumerate(details.items(), 1): + setattr(cmd, "details[%d].%s" % (index, key), value) + last_error = None + pool = None + for attempt in range(1, 4): + try: + pool = StoragePool( + cls.apiClient.createStoragePool(cmd).__dict__ + ) + break + except Exception as exc: + last_error = exc + try: + pools = list_storage_pools( + cls.apiClient, zoneid=cls.zone.id, name=name + ) or [] + except Exception as lookup_exc: + logger.warning( + "Could not check whether ONTAP pool '%s' was " + "created: %s", name, lookup_exc + ) + pools = [] + if pools: + pool = pools[0] + logger.warning( + "Recovered ONTAP pool '%s' after its create API " + "response failed: %s", name, exc + ) + break + if attempt < 3: + logger.warning( + "ONTAP pool '%s' create attempt %d failed: %s; " + "retrying.", name, attempt, exc + ) + time.sleep(5) + if pool is None: + raise last_error + cls.created_pools.append(pool) + logger.info( + "Created ONTAP %s %s pool '%s' (id=%s).", + protocol.upper(), cls._pool_key(scope, cluster_id), + pool.name, pool.id, + ) + return pool + + def _deploy_vm(self, storage_tag=None, host_id=None): + offering = self._service_offering() + root_disk_offering = self._root_disk_offering(storage_tag) + cmd = deployVirtualMachineAPI.deployVirtualMachineCmd() + cmd.zoneid = self.__class__.zone.id + cmd.templateid = self.__class__.template_id + cmd.serviceofferingid = offering.id + cmd.overridediskofferingid = root_disk_offering.id + cmd.account = self.__class__.account.name + cmd.domainid = self.__class__.domain.id + if self.__class__.network_id: + cmd.networkids = self.__class__.network_id + if host_id: + cmd.hostid = host_id + disabled_hosts = self._disable_hosts_outside_cluster( + self.__class__.cluster.id + ) + old_timeout = getattr(self.apiClient.connection, "asyncTimeout", 3600) + self.apiClient.connection.asyncTimeout = 180 + try: + vm = self.apiClient.deployVirtualMachine(cmd) + finally: + self.apiClient.connection.asyncTimeout = old_timeout + self._enable_hosts(disabled_hosts) + if vm is None: + raise RuntimeError( + "deployVirtualMachine did not finish within 180s" + ) + self.__class__.created_vms.append(vm) + return self._poll_vm(vm.id, "Running", timeout=120) + + def _disable_hosts_outside_cluster(self, cluster_id): + disabled = [] + for host in self.__class__.ready_hosts: + if host.clusterid == cluster_id: + continue + try: + Host.update( + self.apiClient, + id=host.id, + allocationstate="Disable", + ) + disabled.append(host.id) + except Exception as exc: + logger.warning( + "Could not disable host %s for VR placement: %s", + host.id, exc, + ) + return disabled + + def _enable_hosts(self, host_ids): + for host_id in host_ids: + try: + Host.update( + self.apiClient, + id=host_id, + allocationstate="Enable", + ) + except Exception as exc: + logger.warning( + "Could not re-enable host %s after deploy: %s", + host_id, exc, + ) + + def _migration_target(self, vm, cluster_id=None): + current = self._get_vm(vm.id) + hosts = Host.listForMigration( + self.apiClient, virtualmachineid=vm.id + ) or [] + suitable = [ + host for host in hosts + if host.id != current.hostid + and (cluster_id is None or host.clusterid == cluster_id) + and getattr(host, "suitableformigration", True) + and str(getattr(host, "state", "Up")).lower() == "up" + and str( + getattr(host, "resourcestate", "Enabled") + ).lower() == "enabled" + ] + if not suitable: + raise unittest.SkipTest( + "No suitable destination host is available for VM %s" + % vm.id + ) + return suitable[0] + + def _migrate_stopped_vm_storage(self, vm, pool): + cmd = migrateVirtualMachineAPI.migrateVirtualMachineCmd() + cmd.virtualmachineid = vm.id + cmd.storageid = pool.id + return self.apiClient.migrateVirtualMachine(cmd) + + def _migrate_vm_volumes(self, vm, host, mappings): + cmd = migrateVMWithVolumeAPI.migrateVirtualMachineWithVolumeCmd() + cmd.virtualmachineid = vm.id + cmd.hostid = host.id + cmd.migrateto = [ + {"volume": str(volume.id), "pool": str(pool.id)} + for volume, pool in mappings + ] + return self.apiClient.migrateVirtualMachineWithVolume(cmd) + + def _migrate_volume_offline(self, volume, pool, timeout=600): + cmd = migrateVolumeAPI.migrateVolumeCmd() + cmd.volumeid = volume.id + cmd.storageid = pool.id + cmd.livemigrate = False + self.apiClient.migrateVolume(cmd) + return self._poll_volume( + volume.id, "storageid", pool.id, timeout=timeout + ) + + @classmethod + def _host_for_cluster(cls, cluster_id): + hosts = cls.hosts_by_cluster.get(cluster_id, []) + if not hosts: + raise unittest.SkipTest( + "No Up/Enabled KVM host exists in cluster %s" % cluster_id + ) + return hosts[0] + + @classmethod + def _require_secondary_cluster(cls): + if cls.secondary_cluster_id is None: + raise unittest.SkipTest( + "A second KVM cluster with an Up host is required" + ) + return cls.secondary_cluster_id + + def _create_data_volume(self, pool): + cls = self.__class__ + offering_id = getattr(cls, "small_disk_offering_id", None) + if offering_id is None: + offering = DiskOffering.create(cls.apiClient, { + "name": "ontap-migration-data-%d" % random.randint(0, 999999), + "displaytext": "ONTAP migration 1GB data disk", + "disksize": 1, + }) + cls.small_disk_offering_id = offering.id + offering_id = offering.id + cmd = createVolumeAPI.createVolumeCmd() + cmd.name = "%s_%d" % (self._vol_name_prefix, random.randint(0, 99999)) + cmd.diskofferingid = offering_id + cmd.zoneid = cls.zone.id + cmd.storageid = pool.id + cmd.account = cls.account.name + cmd.domainid = cls.domain.id + volume = self.apiClient.createVolume(cmd) + cls.created_volumes.append(volume) + return volume + + def _attach_volume(self, vm, volume): + cmd = attachVolumeAPI.attachVolumeCmd() + cmd.id = volume.id + cmd.virtualmachineid = vm.id + self.apiClient.attachVolume(cmd) + return self._poll_volume( + volume.id, "virtualmachineid", vm.id, timeout=180 + ) + + def _stop_vm(self, vm): + cmd = stopVirtualMachineAPI.stopVirtualMachineCmd() + cmd.id = vm.id + self.apiClient.stopVirtualMachine(cmd) + return self._poll_vm(vm.id, "Stopped") + + def _start_vm(self, vm): + cmd = startVirtualMachineAPI.startVirtualMachineCmd() + cmd.id = vm.id + self.apiClient.startVirtualMachine(cmd) + return self._poll_vm(vm.id, "Running") + + def _poll_vm(self, vm_id, state, timeout=300): + deadline = time.time() + timeout + while time.time() < deadline: + cmd = listVirtualMachinesAPI.listVirtualMachinesCmd() + cmd.id = vm_id + cmd.listall = True + vms = self.apiClient.listVirtualMachines(cmd) or [] + if vms and str(vms[0].state).lower() == state.lower(): + return vms[0] + time.sleep(5) + self.fail("VM %s did not reach %s" % (vm_id, state)) + + def _get_vm(self, vm_id): + cmd = listVirtualMachinesAPI.listVirtualMachinesCmd() + cmd.id = vm_id + cmd.listall = True + vms = self.apiClient.listVirtualMachines(cmd) or [] + self.assertTrue(vms, "VM %s was not found" % vm_id) + return vms[0] + + def _get_volume(self, volume_id): + cmd = listVolumesAPI.listVolumesCmd() + cmd.id = volume_id + cmd.listall = True + volumes = self.apiClient.listVolumes(cmd) or [] + self.assertTrue(volumes, "Volume %s was not found" % volume_id) + return volumes[0] + + def _vm_volumes(self, vm_id): + cmd = listVolumesAPI.listVolumesCmd() + cmd.virtualmachineid = vm_id + cmd.listall = True + return self.apiClient.listVolumes(cmd) or [] + + def _poll_volume(self, volume_id, field, value, timeout=300): + deadline = time.time() + timeout + while time.time() < deadline: + volume = self._get_volume(volume_id) + if getattr(volume, field, None) == value: + return volume + time.sleep(5) + self.fail( + "Volume %s field %s did not reach %r" + % (volume_id, field, value) + ) + + def _assert_backend_object(self, protocol, pool, volume): + backend_name = self._backend_name(protocol, pool, volume) + deadline = time.time() + 60 + last_error = None + while time.time() < deadline: + try: + if self._backend_object_exists( + protocol, pool, backend_name): + return + except Exception as exc: + last_error = exc + logger.warning( + "Could not verify %s backend object %s: %s", + protocol, backend_name, exc, + ) + time.sleep(5) + error = " Last ONTAP error: %s" % last_error if last_error else "" + self.fail( + "No %s backend object %s for volume %s in %s.%s" + % (protocol, backend_name, volume.id, pool.name, error) + ) + + def _backend_name(self, protocol, pool, volume): + if protocol.upper() == "NFS3": + return getattr(volume, "path", None) or volume.id + return "/vol/%s/%s" % ( + pool.name, volume.name.replace("-", "_") + ) + + def _backend_object_exists(self, protocol, pool, backend_name): + if protocol.upper() == "NFS3": + files = self.__class__.ontap.list_files_in_volume(pool.name) + return any(str(backend_name) in name for name in files) + return self.__class__.ontap.get_lun( + self.__class__.svm_name, backend_name + ) is not None + + def _assert_backend_object_removed( + self, protocol, pool, backend_name, timeout=180): + deadline = time.time() + timeout + while time.time() < deadline: + if not self._backend_object_exists( + protocol, pool, backend_name): + return + time.sleep(5) + self.fail( + "%s backend object %s still exists in source pool %s" + % (protocol, backend_name, pool.name) + ) + + def _lun_maps(self, pool, volume): + lun_name = self._backend_name("ISCSI", pool, volume) + return [ + item for item in self.__class__.ontap.list_lun_maps_for_volume( + self.__class__.svm_name, pool.name + ) + if item.get("lun", {}).get("name") == lun_name + ] + + @classmethod + def _igroup_name_for_host(cls, host): + host_uuid = re.sub( + r"[^a-zA-Z0-9_-]", "_", + str(getattr(host, "id", None) or getattr(host, "uuid", "")), + ) + return ("cs_%s_%s" % (host_uuid, cls.svm_name))[:96] + + @classmethod + def _host_iqn(cls, host): + listed = Host.list(cls.apiClient, id=host.id, listall=True) or [host] + current = listed[0] + iqn = ( + getattr(current, "storageurl", None) + or getattr(current, "StorageUrl", None) + ) + if iqn and str(iqn).startswith("iqn."): + return str(iqn) + return None + + def _assert_destination_access( + self, protocol, pool, volume, host=None, is_running=False): + if protocol.upper() == "NFS3": + clients = self._export_policy_clients(pool) + scope = str(getattr(pool, "scope", "")).upper() + allowed_cluster = getattr(pool, "clusterid", None) + for cluster_id, hosts in self.__class__.hosts_by_cluster.items(): + for candidate in hosts: + host_ip = getattr(candidate, "ipaddress", None) + if not host_ip: + continue + is_allowed = ( + scope == "ZONE" or cluster_id == allowed_cluster + ) + matches = any( + host_ip in client for client in clients + ) + self.assertEqual( + matches, is_allowed, + "NFS export clients %s do not match %s scope for " + "host %s" % (clients, scope, host_ip), + ) + return + + maps = self._lun_maps(pool, volume) + if not is_running: + self.assertEqual( + maps, [], + "Stopped or detached volume %s retained LUN maps: %s" + % (volume.id, maps), + ) + return + self.assertIsNotNone(host) + expected_igroup = self._igroup_name_for_host(host) + mapped_igroups = { + item.get("igroup", {}).get("name") for item in maps + } + self.assertEqual( + mapped_igroups, {expected_igroup}, + "Volume %s maps %s do not target destination igroup %s" + % (volume.id, maps, expected_igroup), + ) + igroup = self.__class__.ontap.get_igroup( + self.__class__.svm_name, expected_igroup + ) + self.assertIsNotNone( + igroup, + "Destination igroup %s does not exist on ONTAP" + % expected_igroup, + ) + host_iqn = self._host_iqn(host) + if host_iqn: + initiator_names = [ + item.get("name", "") + for item in igroup.get("initiators", []) or [] + ] + self.assertIn( + host_iqn, initiator_names, + "Host IQN %s is not in dest igroup %s: %s" + % (host_iqn, expected_igroup, initiator_names), + ) + + def _snapshot_state( + self, vm=None, volumes=None, source_pool=None, protocol=None, + destination_pool=None, destination_protocol=None): + current_vm = self._get_vm(vm.id) if vm is not None else None + current_volumes = [ + self._get_volume(volume.id) for volume in (volumes or []) + ] + snapshot = { + "vm": None if current_vm is None else ( + current_vm.state, + current_vm.hostid, + getattr(current_vm, "clusterid", None), + ), + "volumes": { + volume.id: ( + volume.storageid, + volume.path, + str(getattr(volume, "format", "") or "").upper(), + getattr(volume, "virtualmachineid", None), + str(volume.state), + ) + for volume in current_volumes + }, + "source_backend": {}, + "destination_backend": {}, + } + if source_pool is not None and protocol is not None: + for volume in current_volumes: + name = self._backend_name(protocol, source_pool, volume) + snapshot["source_backend"][volume.id] = ( + name, + self._backend_object_exists( + protocol, source_pool, name + ), + tuple(sorted( + item.get("igroup", {}).get("name", "") + for item in self._lun_maps(source_pool, volume) + )) if protocol.upper() == "ISCSI" else tuple( + sorted(self._export_policy_clients(source_pool)) + ), + ) + if destination_pool is not None and protocol is not None: + target_protocol = destination_protocol or protocol + for volume in current_volumes: + name = self._backend_name( + target_protocol, destination_pool, volume + ) + snapshot["destination_backend"][volume.id] = ( + name, + self._backend_object_exists( + target_protocol, destination_pool, name + ), + ) + snapshot["destination_protocol"] = target_protocol + return snapshot + + def _assert_state_unchanged( + self, snapshot, vm=None, source_pool=None, protocol=None, + destination_pool=None): + current_vm = self._get_vm(vm.id) if vm is not None else None + vm_state = None if current_vm is None else ( + current_vm.state, + current_vm.hostid, + getattr(current_vm, "clusterid", None), + ) + self.assertEqual(vm_state, snapshot["vm"]) + for volume_id, expected in snapshot["volumes"].items(): + volume = self._get_volume(volume_id) + actual = ( + volume.storageid, + volume.path, + str(getattr(volume, "format", "") or "").upper(), + getattr(volume, "virtualmachineid", None), + str(volume.state), + ) + self.assertEqual(actual, expected) + if source_pool is not None and protocol is not None: + backend_name, existed, access = ( + snapshot["source_backend"][volume_id] + ) + self.assertEqual( + self._backend_object_exists( + protocol, source_pool, backend_name + ), existed, + ) + if protocol.upper() == "ISCSI": + current_access = tuple(sorted( + item.get("igroup", {}).get("name", "") + for item in self._lun_maps(source_pool, volume) + )) + else: + current_access = tuple(sorted( + self._export_policy_clients(source_pool) + )) + self.assertEqual(current_access, access) + if destination_pool is not None and protocol is not None: + backend_name, existed = ( + snapshot["destination_backend"][volume_id] + ) + target_protocol = snapshot.get( + "destination_protocol", protocol + ) + self.assertEqual( + self._backend_object_exists( + target_protocol, destination_pool, backend_name + ), existed, + ) + + def _volumes_on_pool(self, pool): + cmd = listVolumesAPI.listVolumesCmd() + cmd.storageid = pool.id + cmd.listall = True + return self.apiClient.listVolumes(cmd) or [] + + def _volume_ids_on_pool(self, pool): + return {volume.id for volume in self._volumes_on_pool(pool)} + + def _finish_rejected_migration( + self, volumes, destination_pool, existing_dest_ids): + """Put a rejected migration's disks back before the test returns. + + A refused copy can leave the source in Migrating and a duplicate on + the destination. The source has to be Ready again, and the duplicate + has to be gone, or the next case in this file cannot run. + """ + source_ids = set() + for volume in volumes or []: + source_ids.add(volume.id) + current = self._get_volume(volume.id) + if str(current.state) == "Migrating": + logger.warning( + "Rejected migration left volume %s (%s) in Migrating; " + "restoring Ready", + current.id, getattr(current, "name", ""), + ) + self.__class__.dbConnection.execute( + "UPDATE volumes SET state = 'Ready' " + "WHERE uuid = %s AND state = 'Migrating' " + "AND removed IS NULL", + (volume.id,), + ) + current = self._get_volume(volume.id) + self.assertEqual( + str(current.state), + "Ready", + "Volume %s (%s) is %s after a rejected migration; " + "expected Ready" % ( + current.id, getattr(current, "name", ""), current.state, + ), + ) + if destination_pool is None: + return + for created in self._volumes_on_pool(destination_pool): + if created.id in source_ids or created.id in existing_dest_ids: + continue + logger.warning( + "Rejected migration left destination volume %s (%s) in %s; " + "expunging it", + created.id, getattr(created, "name", ""), created.state, + ) + self.__class__.dbConnection.execute( + "UPDATE volumes SET state = 'Expunged', removed = NOW() " + "WHERE uuid = %s AND removed IS NULL", + (created.id,), + ) + + def _assert_migration_success( + self, original_volumes, pool, protocol, vm=None, + expected_vm_state=None, expected_host=None, + expected_vm_id=None, source_pool=None, + check_source_removed=True): + if vm is not None: + if expected_vm_state is not None: + current_vm = self._poll_vm( + vm.id, expected_vm_state, timeout=120 + ) + else: + current_vm = self._get_vm(vm.id) + if expected_host is not None: + self.assertEqual(current_vm.hostid, expected_host.id) + actual_host = next( + (host for host in self.__class__.ready_hosts + if host.id == current_vm.hostid), + None, + ) + self.assertIsNotNone( + actual_host, + "VM host %s is not in the ready host inventory" + % current_vm.hostid, + ) + self.assertEqual( + actual_host.clusterid, expected_host.clusterid + ) + for original in original_volumes: + current = self._get_volume(original.id) + self.assertEqual(current.id, original.id) + self.assertEqual(current.storageid, pool.id) + self.assertEqual( + getattr(current, "virtualmachineid", None), + expected_vm_id, + ) + self.assertTrue(current.path) + expected_format = ( + "QCOW2" if protocol.upper() == "NFS3" else "RAW" + ) + reported_format = getattr(current, "format", None) + if reported_format: + self.assertEqual( + str(reported_format).upper(), expected_format + ) + self._assert_backend_object(protocol, pool, current) + self._assert_destination_access( + protocol, + pool, + current, + host=expected_host, + is_running=expected_vm_state == "Running", + ) + if (source_pool is not None + and source_pool.id != pool.id + and self._pool_protocol(source_pool) == protocol + and check_source_removed): + source_name = self._backend_name( + protocol, source_pool, original + ) + self._assert_backend_object_removed( + protocol, source_pool, source_name + ) + return [self._get_volume(volume.id) for volume in original_volumes] + + def _export_policy_clients(self, pool): + details = self.__class__._pool_details(pool) + policy_name = details.get("exportPolicyName") + if not policy_name: + policy_name = "cs-%s-%s" % ( + self.__class__.svm_name, pool.name + ) + policy = self.__class__.ontap.get_export_policy(policy_name) + self.assertIsNotNone( + policy, + "Export policy '%s' not found on ONTAP" % policy_name, + ) + return [ + client.get("match", "") + for rule in policy.get("rules", []) + for client in rule.get("clients", []) + ] + + def _host_ips(self, cluster_id): + return [ + host.ipaddress + for host in self.__class__.hosts_by_cluster.get(cluster_id, []) + if getattr(host, "ipaddress", None) + ] + + @classmethod + def _host_by_id(cls, host_id): + host = next( + (candidate for candidate in cls.ready_hosts + if candidate.id == host_id), + None, + ) + if host is None: + raise RuntimeError( + "KVM host %s is not in the ready host inventory" % host_id + ) + return host + + @staticmethod + def _vm_instance_name(vm): + name = getattr(vm, "instancename", None) + if not name: + raise RuntimeError( + "CloudStack VM %s has no KVM instance name" % vm.id + ) + return name + + @classmethod + def _dumpxml_disk_sources(cls, host, domain_name): + xml = cls._ssh_output( + host, + "virsh dumpxml %s 2>/dev/null || true" + % shlex.quote(domain_name), + ) + return DISK_SOURCE_RE.findall(xml) + + @classmethod + def _iscsi_sessions(cls, host): + return cls._ssh_output( + host, + "iscsiadm -m session 2>/dev/null || true", + ) + + @classmethod + def _host_vm_snapshot(cls, vm): + domain = shlex.quote(cls._vm_instance_name(vm)) + command = ( + "state=$(virsh domstate %s 2>/dev/null || echo ABSENT); " + "printf 'STATE=%%s\\n' \"$state\"; " + "echo '---XML---'; " + "virsh dumpxml %s 2>/dev/null " + "| grep -E \"source (file|dev)=\" || true; " + "echo '---SESSIONS---'; " + "iscsiadm -m session 2>/dev/null || true" + ) % (domain, domain) + return { + host.id: cls._ssh_output(host, command) + for host in cls.ready_hosts + } + + def _assert_host_snapshot_unchanged(self, snapshot, vm): + self.assertEqual( + self.__class__._host_vm_snapshot(vm), + snapshot, + "Rejected migration changed KVM domain, dumpxml, or iSCSI sessions", + ) + + def _assert_vm_not_running_on_hosts(self, vm): + snapshots = self.__class__._host_vm_snapshot(vm) + for host_id, snapshot in snapshots.items(): + state_line = next( + (line for line in snapshot.splitlines() + if line.startswith("STATE=")), + "", + ) + self.assertNotEqual( + state_line.strip().lower(), + "state=running", + "Stopped VM %s is running on host %s" + % (vm.id, host_id), + ) + + def _volume_path_tokens(self, volume): + tokens = [] + for value in ( + getattr(volume, "path", None), + getattr(volume, "chaininfo", None), + volume.id, + ): + if not value: + continue + text = str(value) + tokens.append(text) + tokens.append(os.path.basename(text.rstrip("/"))) + return [token for token in tokens if token] + + def _assert_dumpxml_contains_volumes(self, host, vm, volumes): + sources = self.__class__._dumpxml_disk_sources( + host, self._vm_instance_name(vm) + ) + self.assertTrue( + sources, + "virsh dumpxml on host %s has no disk sources for VM %s" + % (host.id, vm.id), + ) + joined = " ".join(sources) + for volume in volumes: + tokens = self._volume_path_tokens(volume) + self.assertTrue( + any(token in joined for token in tokens), + "dumpxml on host %s does not include volume %s path %s " + "in %s" + % (host.id, volume.id, getattr(volume, "path", None), sources), + ) + return sources + + def _nfs_mount_candidates(self, pool, host): + uuid = str(getattr(pool, "id", "") or "") + candidates = ["/mnt/%s" % uuid] + xml = self.__class__._ssh_output( + host, + "virsh pool-dumpxml %s 2>/dev/null || true" + % shlex.quote(uuid), + ) + match = re.search(r"([^<]+)", xml or "") + if match: + candidates.append(match.group(1).strip()) + return candidates + + def _scoped_nfs_hosts(self, pool): + scope = str(getattr(pool, "scope", "")).upper() + cluster_id = getattr(pool, "clusterid", None) + return [ + host for host in self.__class__.ready_hosts + if scope == "ZONE" or host.clusterid == cluster_id + ] + + def _assert_nfs_mount_and_file(self, host, pool, volume): + filename = os.path.basename( + str(getattr(volume, "path", "") or volume.id).rstrip("/") + ) + found_mount = None + for mount in self._nfs_mount_candidates(pool, host): + quoted = shlex.quote(mount) + mounted = self.__class__._ssh_output( + host, + "findmnt -n %s >/dev/null 2>&1 && echo MOUNTED || " + "grep -F %s /proc/mounts >/dev/null && echo MOUNTED || true" + % (quoted, quoted), + ) + if "MOUNTED" in mounted: + found_mount = mount + break + self.assertIsNotNone( + found_mount, + "NFS pool %s is not mounted on host %s" + % (pool.id, host.id), + ) + file_path = "%s/%s" % (found_mount.rstrip("/"), filename) + exists = self.__class__._ssh_output( + host, + "test -f %s && echo FILE_EXISTS || " + "find %s -name %s -type f 2>/dev/null | head -1" + % ( + shlex.quote(file_path), + shlex.quote(found_mount), + shlex.quote(filename), + ), + ) + self.assertTrue( + "FILE_EXISTS" in exists or filename in exists, + "NFS file %s for volume %s is missing on host %s under %s" + % (filename, volume.id, host.id, found_mount), + ) + return found_mount, file_path + + def _iscsi_target_tokens(self, pool, volumes): + details = { + str(key).lower(): str(value) + for key, value in self.__class__._pool_details(pool).items() + } + tokens = [ + details.get("storageip") or "", + str(self.__class__.config_data.get("ontap", {}).get( + "storageIP", "" + )), + ] + for volume in volumes: + path = str(getattr(volume, "path", "") or "") + tokens.append(path) + if "iqn." in path: + tokens.append(path.strip("/").split("/")[0]) + return [token for token in tokens if token] + + def _session_mentions_target(self, sessions, tokens): + text = sessions.lower() + return any(str(token).lower() in text for token in tokens) + + def _iscsi_lun_wwid(self, pool, volume): + path = self._backend_name("ISCSI", pool, volume) + lun = self.__class__.ontap.get_lun( + self.__class__.svm_name, path + ) + self.assertTrue( + lun is not None, + "ONTAP LUN %s for volume %s does not exist" + % (path, volume.id), + ) + serial_number = str(lun.get("serial_number", "") or "") + self.assertTrue( + serial_number, + "ONTAP LUN %s for volume %s has no serial number" + % (path, volume.id), + ) + return "600a0980%s" % serial_number.encode("utf-8").hex() + + def _iscsi_device_listing(self, host): + return self.__class__._ssh_output( + host, + "ls -1 /dev/disk/by-path 2>/dev/null; " + "lsscsi -t 2>/dev/null || true; " + "for path in /dev/disk/by-path/*; do " + "device=$(readlink -f \"$path\") || continue; " + "wwid=$(cat \"/sys/class/block/${device##*/}/device/wwid\" " + "2>/dev/null) || continue; " + "printf '%s WWID=%s\\n' \"$path\" \"${wwid#naa.}\"; " + "done", + ) + + def _assert_iscsi_luns_not_visible(self, host, pool, volumes): + listing = self._iscsi_device_listing(host) + joined = listing.lower() + for volume in volumes: + wwid = self._iscsi_lun_wwid(pool, volume) + self.assertNotIn( + wwid, joined, + "LUN for volume %s path %s is still visible on host %s: %s" + % (volume.id, volume.path, host.id, listing), + ) + + def _assert_iscsi_sessions( + self, pool, volumes, dest_host=None, source_host=None, + expect_dest_session=False): + tokens = self._iscsi_target_tokens(pool, volumes) + for host in self.__class__.ready_hosts: + sessions = self.__class__._iscsi_sessions(host) + mentions = self._session_mentions_target(sessions, tokens) + if expect_dest_session and dest_host is not None and ( + host.id == dest_host.id): + self.assertTrue( + mentions, + "Destination host %s has no iSCSI session for %s: %s" + % (host.id, tokens, sessions), + ) + continue + if source_host is not None and host.id == source_host.id: + self._assert_iscsi_luns_not_visible(host, pool, volumes) + continue + if not expect_dest_session: + self._assert_iscsi_luns_not_visible(host, pool, volumes) + + def _assert_iscsi_lun_visible(self, host, pool, volumes): + listing = self._iscsi_device_listing(host) + joined = listing.lower() + for volume in volumes: + wwid = self._iscsi_lun_wwid(pool, volume) + self.assertIn( + wwid, joined, + "LUN for volume %s path %s is not visible on host %s: %s" + % (volume.id, volume.path, host.id, listing), + ) + + def _assert_live_host_state( + self, vm, source_host_id, destination_host, volumes, protocol): + current_volumes = [self._get_volume(volume.id) for volume in volumes] + domain = self._vm_instance_name(vm) + dest_state = self.__class__._ssh_output( + destination_host, + "virsh domstate %s 2>/dev/null || echo ABSENT" + % shlex.quote(domain), + ) + self.assertIn( + "running", dest_state.lower(), + "VM %s is not running on destination host %s" + % (vm.id, destination_host.id), + ) + sources = self._assert_dumpxml_contains_volumes( + destination_host, vm, current_volumes + ) + + source_host = self.__class__._host_by_id(source_host_id) + source_state = self.__class__._ssh_output( + source_host, + "virsh domstate %s 2>/dev/null || echo ABSENT" + % shlex.quote(domain), + ) + self.assertNotIn( + "running", source_state.lower(), + "VM %s remains active on source host %s" + % (vm.id, source_host_id), + ) + + pool = self._pool_by_id(current_volumes[0].storageid) + if protocol.upper() == "NFS3": + for volume in current_volumes: + _, file_path = self._assert_nfs_mount_and_file( + destination_host, pool, volume + ) + info = self.__class__._ssh_output( + destination_host, + "qemu-img info --force-share --output=json %s" + % shlex.quote(file_path), + ) + self.assertIn( + '"format": "qcow2"', info.lower(), + "Destination disk %s is not QCOW2" % file_path, + ) + for path in sources: + if not path.startswith("/"): + continue + info = self.__class__._ssh_output( + destination_host, + "qemu-img info --force-share --output=json %s " + "2>/dev/null || true" + % shlex.quote(path), + ) + if info.strip(): + self.assertIn( + '"format": "qcow2"', info.lower(), + "dumpxml disk %s is not QCOW2" % path, + ) + else: + self._assert_iscsi_sessions( + pool, current_volumes, + dest_host=destination_host, + source_host=source_host, + expect_dest_session=True, + ) + self._assert_iscsi_lun_visible( + destination_host, pool, current_volumes + ) + for path in sources: + result = self.__class__._ssh_output( + destination_host, + "test -b %s && echo BLOCK_DEVICE || true" + % shlex.quote(path), + ) + if path.startswith("/dev"): + self.assertIn( + "BLOCK_DEVICE", result, + "dumpxml iSCSI disk %s is not a block device" % path, + ) + + def _assert_nfs_pool_active_on_scoped_hosts(self, pool): + hosts = self._scoped_nfs_hosts(pool) + self.assertTrue( + hosts, + "No eligible KVM host exists for NFS pool %s" % pool.id, + ) + pool_id = shlex.quote(str(pool.id)) + for host in hosts: + output = self.__class__._ssh_output( + host, + "virsh pool-info %s 2>/dev/null || true" % pool_id, + ) + self.assertIn( + "state:", output.lower(), + "NFS pool %s is not defined on host %s" + % (pool.id, host.id), + ) + self.assertIn( + "running", output.lower(), + "NFS pool %s is not active on host %s" + % (pool.id, host.id), + ) + + @classmethod + def _host_volume_snapshot(cls, volume): + tokens = { + str(value) for value in ( + volume.id, + getattr(volume, "path", None), + getattr(volume, "name", None), + ) if value + } + pattern = "|".join(re.escape(token) for token in sorted(tokens)) + if not pattern: + return {} + command = ( + "echo '---XML---'; " + "for domain in $(virsh list --name); do " + "virsh dumpxml \"$domain\" 2>/dev/null; " + "done | grep -E \"source (file|dev)=\" || true; " + "echo '---REFS---'; " + "for domain in $(virsh list --name); do " + "printf 'DOMAIN=%%s\\n' \"$domain\"; " + "virsh dumpxml \"$domain\" 2>/dev/null; " + "done | grep -E %s || true; " + "echo '---SESSIONS---'; " + "iscsiadm -m session 2>/dev/null || true" + ) % shlex.quote(pattern) + return { + host.id: cls._ssh_output(host, command) + for host in cls.ready_hosts + } + + def _assert_volume_not_referenced_on_hosts(self, volume): + for host_id, output in ( + self.__class__._host_volume_snapshot(volume).items()): + refs = "" + if "---REFS---" in output: + refs = output.split("---REFS---", 1)[1] + if "---SESSIONS---" in refs: + refs = refs.split("---SESSIONS---", 1)[0] + self.assertEqual( + refs.strip(), + "", + "Volume %s is referenced by a running domain on host %s" + % (volume.id, host_id), + ) + + def _assert_offline_host_state( + self, protocol, pool, vm=None, volumes=None): + current_volumes = [ + self._get_volume(volume.id) for volume in (volumes or []) + ] + if vm is not None: + self._assert_vm_not_running_on_hosts(vm) + if protocol.upper() == "NFS3": + self._assert_nfs_pool_active_on_scoped_hosts(pool) + hosts = self._scoped_nfs_hosts(pool) + for volume in current_volumes: + self._assert_nfs_mount_and_file(hosts[0], pool, volume) + else: + self._assert_iscsi_sessions( + pool, current_volumes, expect_dest_session=False + ) + for volume in current_volumes: + self._assert_volume_not_referenced_on_hosts(volume) + + @classmethod + def tearDownClass(cls): + old_timeout = getattr( + getattr(cls, "apiClient", None) and cls.apiClient.connection, + "asyncTimeout", + None, + ) + if old_timeout is not None: + cls.apiClient.connection.asyncTimeout = 60 + try: + cls._expunge_account_vms() + except Exception as exc: + logger.warning("Could not expunge leftover VMs: %s", exc) + for volume in reversed(cls.created_volumes or []): + try: + cmd = deleteVolumeAPI.deleteVolumeCmd() + cmd.id = volume.id + cls.apiClient.deleteVolume(cmd) + except Exception: + pass + + if os.environ.get("ONTAP_MIGRATION_KEEP_POOLS") != "1": + try: + pools = list_storage_pools( + cls.apiClient, zoneid=cls.zone.id + ) or [] + except Exception as exc: + logger.warning( + "Could not list migration pools for cleanup: %s", exc + ) + pools = [] + migration_pools = { + pool.id: pool for pool in pools + if str(getattr(pool, "name", "")).startswith( + "OntapMigration" + ) + } + for pool in cls.created_pools or []: + migration_pools[pool.id] = pool + + for pool in migration_pools.values(): + if str( + getattr(pool, "type", "") + ) == "NetworkFilesystem": + cls._cleanup_kvm_storage_pool_mounts(pool.id) + + for pool in reversed(list(migration_pools.values())): + try: + cls._delete_extra_pool(pool) + except Exception as exc: + logger.warning( + "Could not clean pool %s: %s", pool.id, exc + ) + try: + cls.ontap.delete_volume(pool.name) + cls.ontap.delete_export_policy( + "cs-%s-%s" % (cls.svm_name, pool.name) + ) + except Exception as ontap_exc: + logger.warning( + "Could not directly clean ONTAP pool %s: %s", + pool.name, ontap_exc + ) + + cls.pool = None + cls.pool2 = None + try: + super(OntapMigrationTestBase, cls).tearDownClass() + if cls._created_network_id: + try: + cmd = deleteNetworkAPI.deleteNetworkCmd() + cmd.id = cls._created_network_id + cls.apiClient.deleteNetwork(cmd) + except Exception as exc: + logger.warning( + "Could not clean network %s: %s", + cls._created_network_id, exc + ) + finally: + if old_timeout is not None: + cls.apiClient.connection.asyncTimeout = old_timeout + + @classmethod + def _list_vm_for_cleanup(cls, vm_id): + cmd = listVirtualMachinesAPI.listVirtualMachinesCmd() + cmd.id = vm_id + cmd.listall = True + vms = cls.apiClient.listVirtualMachines(cmd) or [] + return vms[0] if vms else None + + @classmethod + def _expunge_account_vms(cls): + vms = list(reversed(cls.created_vms or [])) + try: + cmd = listVirtualMachinesAPI.listVirtualMachinesCmd() + cmd.listall = True + if getattr(cls, "account", None) is not None: + cmd.account = cls.account.name + cmd.domainid = cls.domain.id + listed = cls.apiClient.listVirtualMachines(cmd) or [] + known_ids = {getattr(vm, "id", None) for vm in vms} + vms.extend( + vm for vm in listed if getattr(vm, "id", None) not in known_ids + ) + except Exception as exc: + logger.warning("Could not list VMs for cleanup: %s", exc) + cleanup_states = ("stopped", "destroyed", "expunging", "error") + skip_stop_states = cleanup_states + ("starting", "migrating") + for vm in vms: + try: + current = cls._list_vm_for_cleanup(vm.id) + if (current + and str(current.state).lower() not in skip_stop_states): + stop_cmd = stopVirtualMachineAPI.stopVirtualMachineCmd() + stop_cmd.id = vm.id + stop_cmd.forced = True + cls.apiClient.stopVirtualMachine(stop_cmd) + destroy_cmd = ( + destroyVirtualMachineAPI.destroyVirtualMachineCmd() + ) + destroy_cmd.id = vm.id + destroy_cmd.expunge = True + cls.apiClient.destroyVirtualMachine(destroy_cmd) + except Exception as exc: + logger.warning("Could not clean VM %s: %s", vm.id, exc) + + @classmethod + def _delete_extra_pool(cls, pool): + from marvin.cloudstackAPI import ( + deleteStoragePool as deleteStoragePoolAPI, + enableStorageMaintenance, + ) + pools = list_storage_pools(cls.apiClient, id=pool.id) or [] + if pools and getattr(pools[0], "state", None) in ("Up", "Disabled"): + maintenance_cmd = ( + enableStorageMaintenance.enableStorageMaintenanceCmd() + ) + maintenance_cmd.id = pool.id + cls.apiClient.enableStorageMaintenance(maintenance_cmd) + deadline = time.time() + 120 + while time.time() < deadline: + pools = list_storage_pools(cls.apiClient, id=pool.id) or [] + if pools and pools[0].state == "Maintenance": + break + time.sleep(5) + cmd = deleteStoragePoolAPI.deleteStoragePoolCmd() + cmd.id = pool.id + cmd.forced = True + cls.apiClient.deleteStoragePool(cmd) + + @classmethod + def _pool_has_volumes(cls, pool): + cmd = listVolumesAPI.listVolumesCmd() + cmd.storageid = pool.id + cmd.listall = True + return bool(cls.apiClient.listVolumes(cmd)) + + def _release_idle_pools(self, protocol, keep_pool=None): + """Delete FlexVols the suite is finished with, then recreate empty ones. + + The lab aggregates are about 24 GB. Live migration keeps the source + copy until the destination copy finishes, so idle template caches + from earlier hops have to be removed first. + """ + cls = self.__class__ + keep_id = None if keep_pool is None else keep_pool.id + for pools in cls.ontap_pools.get(protocol, {}).values(): + for pool in list(pools): + if keep_id is not None and pool.id == keep_id: + continue + if cls._pool_has_volumes(pool): + logger.info( + "Keeping %s pool %s because a volume is still on it", + protocol, pool.name, + ) + continue + try: + cls._delete_extra_pool(pool) + except Exception as exc: + logger.warning( + "Could not release idle %s pool %s: %s", + protocol, getattr(pool, "name", pool.id), exc, + ) + cls.created_pools = [ + item for item in cls.created_pools + if item.id != pool.id + ] + cls.ontap_pools = cls._ensure_migration_pools() diff --git a/test/integration/plugins/ontap/migration/test_01_live_vm_with_storage_migration.py b/test/integration/plugins/ontap/migration/test_01_live_vm_with_storage_migration.py new file mode 100644 index 000000000000..ec13dda7473a --- /dev/null +++ b/test/integration/plugins/ontap/migration/test_01_live_vm_with_storage_migration.py @@ -0,0 +1,594 @@ +# 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. + +"""Live VM-with-storage migration matrix for ONTAP NFS3 and iSCSI. + +Workflow: + 01-05 same-cluster NFS3 DP/C1/Z chain with root and data volumes + 06-10 same-cluster iSCSI DP/C1/Z chain with root and data volumes + 11-13 same-cluster rejections (cross-protocol and volume-only) + 14 relocate one host once, then NFS3 cross-cluster C1/C2/Z cases + 18-21 iSCSI cross-cluster C1/C2/Z cases + 22 reject a cluster-1 pool mapped to a cluster-2 destination host +""" + +import unittest + +from nose.plugins.attrib import attr + +from marvin.cloudstackAPI import ( + destroyVirtualMachine as destroyVirtualMachineAPI, + migrateVolume as migrateVolumeAPI, + startVirtualMachine as startVirtualMachineAPI, +) +from marvin.lib.common import list_storage_pools + +from migration.migration_test_base import OntapMigrationTestBase + + +class TestOntapLiveVmWithStorageMigration(OntapMigrationTestBase): + """Running-VM migrateVirtualMachineWithVolume coverage for ONTAP.""" + + nfs_vm = None + nfs_pool = None + iscsi_vm = None + iscsi_pool = None + cross_cluster_ready = False + source_host = None + target_host = None + + @classmethod + def _host_requirements(cls): + return { + "minimum_hosts": 2, + "same_primary_cluster": True, + "multiple_clusters": False, + "live_migration": True, + } + + @classmethod + def _pool_requirements(cls): + return { + "NFS3": [ + ("CLUSTER", cls.cluster.id, 2), + ("ZONE", None, 2), + ], + "ISCSI": [ + ("CLUSTER", cls.cluster.id, 2), + ("ZONE", None, 2), + ], + } + + def _c1a(self, protocol): + return self._migration_pool(protocol) + + def _c1b(self, protocol): + return self._migration_pool(protocol, 1) + + def _z1(self, protocol): + return self._migration_pool(protocol, scope="ZONE") + + def _z2(self, protocol): + return self._migration_pool(protocol, 1, scope="ZONE") + + def _c2(self, protocol): + return self._migration_pool( + protocol, + cluster_id=self._require_secondary_cluster(), + ) + + def _prepare_default_vm(self): + host = self._host_for_cluster(self.__class__.cluster.id) + vm = self._deploy_vm(self._default_storage_tag(), host.id) + data_volume = self._create_data_volume(self.__class__.default_pool) + self._attach_volume(vm, data_volume) + for volume in self._vm_volumes(vm.id): + self.assertEqual( + volume.storageid, + self.__class__.default_pool.id, + "VM volume was not allocated on DefaultPrimary", + ) + return vm + + def _prepare_vm_on_pool(self, pool): + host = self.__class__.source_host or self._host_for_cluster( + self.__class__.cluster.id + ) + vm = self._deploy_vm(self._default_storage_tag(), host.id) + data_volume = self._create_data_volume( + self.__class__.default_pool + ) + self._attach_volume(vm, data_volume) + self._stop_vm(vm) + self._migrate_stopped_vm_storage(vm, pool) + vm = self._start_on_host(vm, host) + for volume in self._vm_volumes(vm.id): + self.assertEqual(volume.storageid, pool.id) + return self._get_vm(vm.id) + + def _start_on_host(self, vm, host): + cmd = startVirtualMachineAPI.startVirtualMachineCmd() + cmd.id = vm.id + cmd.hostid = host.id + self.apiClient.startVirtualMachine(cmd) + return self._poll_vm(vm.id, "Running") + + def _live_migrate(self, vm, dest_pool, protocol, dest_host=None, + dest_cluster_id=None): + if vm is None: + raise unittest.SkipTest( + "Prerequisite VM was not created by an earlier test" + ) + volumes = self._vm_volumes(vm.id) + self.assertTrue(len(volumes) >= 2) + source_pool = self._pool_by_id(volumes[0].storageid) + old_host_id = self._get_vm(vm.id).hostid + if dest_cluster_id is None: + dest_cluster_id = self.__class__.cluster.id + if dest_host is None: + host = self._migration_target(vm, dest_cluster_id) + else: + host = dest_host + self._migrate_vm_volumes( + vm, + host, + [(volume, dest_pool) for volume in volumes], + ) + self._assert_migration_success( + volumes, + dest_pool, + protocol, + vm=vm, + expected_vm_state="Running", + expected_host=host, + expected_vm_id=vm.id, + source_pool=source_pool, + ) + self._assert_live_host_state( + vm, old_host_id, host, volumes, protocol + ) + self.assertEqual(host.clusterid, dest_cluster_id) + self.assertNotEqual(host.id, old_host_id) + return self._get_vm(vm.id), dest_pool + + def _reject_live_mapping(self, vm, dest_pool, protocol, pattern=None): + if vm is None: + raise unittest.SkipTest( + "Prerequisite VM was not created by an earlier test" + ) + volumes = self._vm_volumes(vm.id) + source_pool = self._pool_by_id(volumes[0].storageid) + snapshot = self._snapshot_state( + vm=vm, + volumes=volumes, + source_pool=source_pool, + protocol=protocol, + destination_pool=dest_pool, + destination_protocol=self._pool_protocol(dest_pool), + ) + host_snapshot = self.__class__._host_vm_snapshot(vm) + host = self._migration_target(vm, self.__class__.cluster.id) + if pattern is None: + with self.assertRaises(Exception): + self._migrate_vm_volumes( + vm, + host, + [(volume, dest_pool) for volume in volumes], + ) + else: + with self.assertRaisesRegex(Exception, pattern): + self._migrate_vm_volumes( + vm, + host, + [(volume, dest_pool) for volume in volumes], + ) + self._assert_state_unchanged( + snapshot, + vm=vm, + source_pool=source_pool, + protocol=protocol, + destination_pool=dest_pool, + ) + self._assert_host_snapshot_unchanged(host_snapshot, vm) + + def _expunge_vm(self, vm): + if vm is None: + return + try: + current = self._get_vm(vm.id) + if str(current.state).lower() not in ("stopped", "destroyed"): + self._stop_vm(vm) + except Exception: + pass + try: + cmd = destroyVirtualMachineAPI.destroyVirtualMachineCmd() + cmd.id = vm.id + cmd.expunge = True + self.apiClient.destroyVirtualMachine(cmd) + except Exception: + pass + + def _expunge_phase_vms(self): + cls = self.__class__ + self._expunge_vm(cls.nfs_vm) + self._expunge_vm(cls.iscsi_vm) + cls.nfs_vm = None + cls.iscsi_vm = None + cls.nfs_pool = None + cls.iscsi_pool = None + + def _ensure_cross_cluster_phase(self): + cls = self.__class__ + if cls.cross_cluster_ready: + return + self._expunge_phase_vms() + cls._ensure_hosts_across_clusters() + cls._refresh_host_inventory() + cls._validate_host_topology({ + "same_primary_cluster": False, + "multiple_clusters": True, + }) + cls._validate_host_credentials() + cls._prepare_live_migration_hosts() + secondary = cls._require_secondary_cluster() + cls._recreate_test_network() + pools = list_storage_pools(cls.apiClient, zoneid=cls.zone.id) or [] + for protocol in ("NFS3", "ISCSI"): + try: + cls._protocol_config(protocol) + except Exception: + continue + key = cls._pool_key("CLUSTER", secondary) + bucket = cls.ontap_pools.setdefault( + protocol, {} + ).setdefault(key, []) + if bucket: + continue + found = [ + pool for pool in pools + if cls._is_compatible_ontap_pool( + pool, protocol, "CLUSTER", secondary + ) + ][:1] + if found: + bucket.extend(found) + else: + bucket.append( + cls._create_ontap_pool( + protocol, "CLUSTER", secondary + ) + ) + cls._validate_migration_pool( + bucket[0], protocol, "CLUSTER", secondary + ) + cls.source_host = cls._host_for_cluster(cls.cluster.id) + cls.target_host = cls._host_for_cluster(secondary) + cls.cross_cluster_ready = True + + def _live_cross(self, vm, dest_pool, protocol): + if vm is None: + raise unittest.SkipTest( + "Prerequisite VM was not created by an earlier test" + ) + secondary = self._require_secondary_cluster() + current = self._get_vm(vm.id) + if current.hostid == self.__class__.target_host.id: + self._stop_vm(vm) + vm = self._start_on_host(vm, self.__class__.source_host) + return self._live_migrate( + vm, + dest_pool, + protocol, + dest_host=self.__class__.target_host, + dest_cluster_id=secondary, + ) + + def _data_volume(self, vm): + if vm is None: + raise unittest.SkipTest( + "Prerequisite VM was not created by an earlier test" + ) + for volume in self._vm_volumes(vm.id): + if str(getattr(volume, "type", "")).upper() != "ROOT": + return volume + self.fail("No data volume attached to VM %s" % vm.id) + + def _reset_running_on_pool(self, vm, pool): + if vm is None: + raise unittest.SkipTest( + "Prerequisite VM was not created by an earlier test" + ) + self._stop_vm(vm) + self._migrate_stopped_vm_storage(vm, pool) + host = self.__class__.source_host or self._host_for_cluster( + self.__class__.cluster.id + ) + return self._start_on_host(vm, host) + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_01_nfs3_default_primary_to_cluster(self): + """Live-migrate DefaultPrimary root/data volumes to NFS3 C1a.""" + vm = self._prepare_default_vm() + pool = self._c1a("NFS3") + vm, pool = self._live_migrate(vm, pool, "NFS3") + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_02_nfs3_cluster_to_cluster(self): + """Live-migrate the NFS3 VM from C1a to C1b.""" + pool = self._c1b("NFS3") + vm, pool = self._live_migrate( + self.__class__.nfs_vm, pool, "NFS3" + ) + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_03_nfs3_cluster_to_zone(self): + """Live-migrate the NFS3 VM from C1b to Z1.""" + pool = self._z1("NFS3") + vm, pool = self._live_migrate( + self.__class__.nfs_vm, pool, "NFS3" + ) + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_04_nfs3_zone_to_zone(self): + """Live-migrate the NFS3 VM from Z1 to Z2.""" + pool = self._z2("NFS3") + vm, pool = self._live_migrate( + self.__class__.nfs_vm, pool, "NFS3" + ) + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_05_nfs3_zone_to_cluster(self): + """Live-migrate the NFS3 VM from Z2 back to C1a.""" + pool = self._c1a("NFS3") + vm, pool = self._live_migrate( + self.__class__.nfs_vm, pool, "NFS3" + ) + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_06_iscsi_default_primary_to_cluster(self): + """Live-migrate DefaultPrimary root/data volumes to iSCSI C1a.""" + self._release_idle_pools("NFS3", self.__class__.nfs_pool) + vm = self._prepare_default_vm() + pool = self._c1a("ISCSI") + vm, pool = self._live_migrate(vm, pool, "ISCSI") + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_07_iscsi_cluster_to_cluster(self): + """Live-migrate the iSCSI VM from C1a to C1b.""" + pool = self._c1b("ISCSI") + vm, pool = self._live_migrate( + self.__class__.iscsi_vm, pool, "ISCSI" + ) + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_08_iscsi_cluster_to_zone(self): + """Live-migrate the iSCSI VM from C1b to Z1.""" + pool = self._z1("ISCSI") + vm, pool = self._live_migrate( + self.__class__.iscsi_vm, pool, "ISCSI" + ) + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_09_iscsi_zone_to_zone(self): + """Live-migrate the iSCSI VM from Z1 to Z2.""" + pool = self._z2("ISCSI") + vm, pool = self._live_migrate( + self.__class__.iscsi_vm, pool, "ISCSI" + ) + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_10_iscsi_zone_to_cluster(self): + """Live-migrate the iSCSI VM from Z2 back to C1a.""" + self._release_idle_pools("ISCSI", self.__class__.iscsi_pool) + pool = self._c1a("ISCSI") + vm, pool = self._live_migrate( + self.__class__.iscsi_vm, pool, "ISCSI" + ) + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_11_reject_nfs3_to_iscsi(self): + """Reject a live NFS3-to-iSCSI volume mapping.""" + self._reject_live_mapping( + self.__class__.nfs_vm, + self._c1a("ISCSI"), + "NFS3", + "managed storage can only be 'migrated' to itself", + ) + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_12_reject_iscsi_to_nfs3(self): + """Reject a live iSCSI-to-NFS3 volume mapping.""" + self._reject_live_mapping( + self.__class__.iscsi_vm, + self._c1a("NFS3"), + "ISCSI", + "managed storage can only be 'migrated' to itself", + ) + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_13_reject_live_volume_only(self): + """Reject migrateVolume livemigrate=true on a running attached volume.""" + vm = self.__class__.nfs_vm + volume = self._data_volume(vm) + dest = self._c1b("NFS3") + source_pool = self.__class__.nfs_pool + volumes = self._vm_volumes(vm.id) + snapshot = self._snapshot_state( + vm=vm, + volumes=volumes, + source_pool=source_pool, + protocol="NFS3", + destination_pool=dest, + ) + host_snapshot = self.__class__._host_vm_snapshot(vm) + cmd = migrateVolumeAPI.migrateVolumeCmd() + cmd.volumeid = volume.id + cmd.storageid = dest.id + cmd.livemigrate = True + with self.assertRaises(Exception) as context: + self.apiClient.migrateVolume(cmd) + self.assertIn( + "migratevirtualmachinewithvolume", + str(context.exception).lower().replace(" ", ""), + ) + self._assert_state_unchanged( + snapshot, + vm=vm, + source_pool=source_pool, + protocol="NFS3", + destination_pool=dest, + ) + self._assert_host_snapshot_unchanged(host_snapshot, vm) + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_14_nfs3_cluster_to_cluster(self): + """After host split, live-migrate NFS3 storage from C1a to C2.""" + self._ensure_cross_cluster_phase() + vm = self._prepare_vm_on_pool(self._c1a("NFS3")) + vm, pool = self._live_cross(vm, self._c2("NFS3"), "NFS3") + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_15_nfs3_cluster_to_zone(self): + """Live-migrate NFS3 storage from C1a to Z1 onto a C2 host.""" + self._ensure_cross_cluster_phase() + self._expunge_vm(self.__class__.nfs_vm) + vm = self._prepare_vm_on_pool(self._c1a("NFS3")) + vm, pool = self._live_cross(vm, self._z1("NFS3"), "NFS3") + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_16_nfs3_zone_to_cluster(self): + """Live-migrate the NFS3 VM from Z1 to C2.""" + self._ensure_cross_cluster_phase() + vm, pool = self._live_cross( + self.__class__.nfs_vm, self._c2("NFS3"), "NFS3" + ) + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_17_nfs3_zone_to_zone(self): + """Live-migrate NFS3 storage from Z1 to Z2 onto a C2 host.""" + self._ensure_cross_cluster_phase() + vm = self._reset_running_on_pool( + self.__class__.nfs_vm, self._z1("NFS3") + ) + vm, pool = self._live_cross(vm, self._z2("NFS3"), "NFS3") + self.__class__.nfs_vm = vm + self.__class__.nfs_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_18_iscsi_cluster_to_cluster(self): + """Live-migrate iSCSI storage from C1a to C2 onto a C2 host.""" + self._ensure_cross_cluster_phase() + self._expunge_vm(self.__class__.nfs_vm) + self.__class__.nfs_vm = None + vm = self._prepare_vm_on_pool(self._c1a("ISCSI")) + vm, pool = self._live_cross(vm, self._c2("ISCSI"), "ISCSI") + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_19_iscsi_cluster_to_zone(self): + """Live-migrate iSCSI storage from C1a to Z1 onto a C2 host.""" + self._ensure_cross_cluster_phase() + self._expunge_vm(self.__class__.iscsi_vm) + vm = self._prepare_vm_on_pool(self._c1a("ISCSI")) + vm, pool = self._live_cross(vm, self._z1("ISCSI"), "ISCSI") + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_20_iscsi_zone_to_cluster(self): + """Live-migrate the iSCSI VM from Z1 to C2.""" + self._ensure_cross_cluster_phase() + vm, pool = self._live_cross( + self.__class__.iscsi_vm, self._c2("ISCSI"), "ISCSI" + ) + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_21_iscsi_zone_to_zone(self): + """Live-migrate iSCSI storage from Z1 to Z2 onto a C2 host.""" + self._ensure_cross_cluster_phase() + self._expunge_vm(self.__class__.iscsi_vm) + vm = self._prepare_vm_on_pool(self._z1("ISCSI")) + vm, pool = self._live_cross(vm, self._z2("ISCSI"), "ISCSI") + self.__class__.iscsi_vm = vm + self.__class__.iscsi_pool = pool + + @attr(tags=["ontap_migration", "live_storage"], required_hardware=True) + def test_22_reject_pool_not_reachable_from_destination_host(self): + """Reject mapping a cluster-1 pool to a host that lives in cluster 2.""" + self._ensure_cross_cluster_phase() + vm = self.__class__.nfs_vm + if vm is None: + vm = self._prepare_vm_on_pool(self._z1("NFS3")) + elif self._get_vm(vm.id).hostid != self.__class__.source_host.id: + self._stop_vm(vm) + vm = self._start_on_host(vm, self.__class__.source_host) + self.__class__.nfs_vm = vm + volumes = self._vm_volumes(vm.id) + source_pool = self._pool_by_id(volumes[0].storageid) + dest = self._c1a("NFS3") + snapshot = self._snapshot_state( + vm=vm, + volumes=volumes, + source_pool=source_pool, + protocol="NFS3", + destination_pool=dest, + ) + host_snapshot = self.__class__._host_vm_snapshot(vm) + host = self.__class__.target_host + with self.assertRaises(Exception): + self._migrate_vm_volumes( + vm, + host, + [(volume, dest) for volume in volumes], + ) + self._assert_state_unchanged( + snapshot, + vm=vm, + source_pool=source_pool, + protocol="NFS3", + destination_pool=dest, + ) + self._assert_host_snapshot_unchanged(host_snapshot, vm) diff --git a/test/integration/plugins/ontap/migration/test_02_stopped_vm_storage_migration.py b/test/integration/plugins/ontap/migration/test_02_stopped_vm_storage_migration.py new file mode 100644 index 000000000000..a6e927567aea --- /dev/null +++ b/test/integration/plugins/ontap/migration/test_02_stopped_vm_storage_migration.py @@ -0,0 +1,316 @@ +# 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. + +"""Ordered stopped-VM storage migration tests. + +Every case runs migrateVirtualMachine with storageid, which moves all of a +stopped VM's volumes to one pool. One VM per protocol, holding a root and an +attached data volume, is carried through the same six-step pool chain: +DefaultPrimary (DP), cluster-1 pool A (C1a), cluster-1 pool B (C1b), +cluster-2 pool (C2), zone pool 1 (Z1), zone pool 2 (Z2), back to C1a. Hosts +stay split across two clusters; only the volumes move. + +Workflow: + 01 NFS3 DP to C1a + 02 NFS3 C1a to C1b, both cluster-scoped in cluster 1 + 03 NFS3 C1b to C2, cluster-scoped across clusters + 04 NFS3 C2 to Z1, cluster-scoped to zone-scoped + 05 NFS3 Z1 to Z2, zone-scoped to zone-scoped + 06 NFS3 Z2 to C1a, zone-scoped back to cluster-scoped + 07 iSCSI DP to C1a + 08 iSCSI C1a to C1b, both cluster-scoped in cluster 1 + 09 iSCSI C1b to C2, cluster-scoped across clusters + 10 iSCSI C2 to Z1, cluster-scoped to zone-scoped + 11 iSCSI Z1 to Z2, zone-scoped to zone-scoped + 12 iSCSI Z2 to C1a, zone-scoped back to cluster-scoped + 13 reject NFS3-to-iSCSI storage migration + 14 reject iSCSI-to-NFS3 storage migration + 15 reject the stopped-VM API form once the VM is running + +Migration out of ONTAP into non-ONTAP storage is never requested; the +DefaultPrimary pool is only ever a source. +""" + +import unittest + +from nose.plugins.attrib import attr + +from migration.migration_test_base import OntapMigrationTestBase + +CROSS_PROTOCOL_REJECTION = "managed storage can only be 'migrated'" + + +class TestOntapStoppedVmStorageMigration(OntapMigrationTestBase): + """Verify stopped-VM storage motion across every ONTAP pool scope.""" + + chain_vms = None + chain_volumes = None + chain_pools = None + + @classmethod + def setUpClass(cls): + super(TestOntapStoppedVmStorageMigration, cls).setUpClass() + cls.chain_vms = {} + cls.chain_volumes = {} + cls.chain_pools = {} + + @classmethod + def _host_requirements(cls): + return { + "minimum_hosts": 2, + "same_primary_cluster": False, + "multiple_clusters": True, + } + + @classmethod + def _pool_requirements(cls): + secondary_cluster_id = cls._require_secondary_cluster() + scopes = [ + ("CLUSTER", cls.cluster.id, 2), + ("CLUSTER", secondary_cluster_id, 1), + ("ZONE", None, 2), + ] + return {"NFS3": list(scopes), "ISCSI": list(scopes)} + + def _pools_for(self, protocol): + secondary_cluster_id = self._require_secondary_cluster() + return { + "C1A": self._migration_pool(protocol), + "C1B": self._migration_pool(protocol, 1), + "C2": self._migration_pool( + protocol, cluster_id=secondary_cluster_id + ), + "Z1": self._migration_pool(protocol, scope="ZONE"), + "Z2": self._migration_pool(protocol, 1, scope="ZONE"), + } + + def _start_chain(self, protocol): + """Deploy a DefaultPrimary VM with a data disk and stop it.""" + cls = self.__class__ + cls.chain_pools[protocol] = self._pools_for(protocol) + vm = self._deploy_vm( + self._default_storage_tag(), + self._host_for_cluster(cls.cluster.id).id, + ) + self._attach_volume( + vm, self._create_data_volume(cls.default_pool) + ) + self._stop_vm(vm) + volumes = self._vm_volumes(vm.id) + self.assertTrue( + len(volumes) >= 2, + "The %s chain VM has no attached data volume" % protocol, + ) + for volume in volumes: + self.assertEqual( + volume.storageid, + cls.default_pool.id, + "VM volume was not allocated on DefaultPrimary", + ) + cls.chain_vms[protocol] = vm + cls.chain_volumes[protocol] = volumes + + def _chain(self, protocol): + cls = self.__class__ + if protocol not in cls.chain_vms: + raise unittest.SkipTest( + "The %s chain VM was not prepared by the first case" + % protocol + ) + return ( + cls.chain_vms[protocol], + cls.chain_volumes[protocol], + cls.chain_pools[protocol], + ) + + def _migrate_step(self, protocol, source_key, destination_key): + """Move every volume of the chain VM to the next pool in the chain.""" + vm, volumes, pools = self._chain(protocol) + destination = pools[destination_key] + source_pool = pools[source_key] if source_key else None + self._migrate_stopped_vm_storage(vm, destination) + for volume in volumes: + self._poll_volume( + volume.id, "storageid", destination.id, timeout=900 + ) + migrated = self._assert_migration_success( + volumes, + destination, + protocol, + vm=vm, + expected_vm_state="Stopped", + expected_vm_id=vm.id, + source_pool=source_pool, + ) + self.__class__.chain_volumes[protocol] = migrated + self._assert_offline_host_state( + protocol, destination, vm=vm, volumes=migrated + ) + return destination + + def _destination_objects(self, protocol, pool, volumes): + return { + volume.id: self._backend_object_exists( + protocol, pool, self._backend_name(protocol, pool, volume) + ) + for volume in volumes + } + + def _assert_rejected(self, protocol, source_key, destination_pool, + destination_protocol, pattern=None): + """Assert CloudStack and ONTAP are untouched by a refused request.""" + vm, volumes, pools = self._chain(protocol) + source_pool = pools[source_key] + snapshot = self._snapshot_state( + vm=vm, + volumes=volumes, + source_pool=source_pool, + protocol=protocol, + ) + host_snapshot = self.__class__._host_vm_snapshot(vm) + dest_ids = self._volume_ids_on_pool(destination_pool) + destination_before = self._destination_objects( + destination_protocol, destination_pool, volumes + ) + if pattern: + with self.assertRaisesRegex(Exception, pattern): + self._migrate_stopped_vm_storage(vm, destination_pool) + else: + with self.assertRaises(Exception): + self._migrate_stopped_vm_storage(vm, destination_pool) + self._finish_rejected_migration( + volumes, destination_pool, dest_ids + ) + self._assert_state_unchanged( + snapshot, + vm=vm, + source_pool=source_pool, + protocol=protocol, + ) + self.assertEqual( + self._destination_objects( + destination_protocol, destination_pool, volumes + ), + destination_before, + "Refused migration changed destination pool %s" + % destination_pool.name, + ) + self._assert_host_snapshot_unchanged(host_snapshot, vm) + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_01_nfs3_default_primary_to_cluster(self): + """Move a stopped VM's volumes from DefaultPrimary to NFS3 C1a.""" + self._start_chain("NFS3") + self._migrate_step("NFS3", None, "C1A") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_02_nfs3_cluster_to_cluster_same_cluster(self): + """Move NFS3 volumes between cluster pools in the same cluster.""" + self._migrate_step("NFS3", "C1A", "C1B") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_03_nfs3_cluster_to_cluster_across_clusters(self): + """Move NFS3 volumes to a cluster pool in the second cluster.""" + destination = self._migrate_step("NFS3", "C1B", "C2") + clients = self._export_policy_clients(destination) + for host_ip in self._host_ips(self.__class__.cluster.id): + self.assertFalse( + any(host_ip in client for client in clients), + "Source cluster host %s remained in destination export %s" + % (host_ip, clients), + ) + for host_ip in self._host_ips(self._require_secondary_cluster()): + self.assertTrue( + any(host_ip in client for client in clients), + "Destination cluster host %s missing from export %s" + % (host_ip, clients), + ) + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_04_nfs3_cluster_to_zone(self): + """Move NFS3 volumes from a cluster pool to a zone-wide pool.""" + self._migrate_step("NFS3", "C2", "Z1") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_05_nfs3_zone_to_zone(self): + """Move NFS3 volumes between two zone-wide pools.""" + self._migrate_step("NFS3", "Z1", "Z2") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_06_nfs3_zone_to_cluster(self): + """Move NFS3 volumes from a zone-wide pool back to cluster C1a.""" + self._migrate_step("NFS3", "Z2", "C1A") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_07_iscsi_default_primary_to_cluster(self): + """Move a stopped VM's volumes from DefaultPrimary to iSCSI C1a.""" + self._start_chain("ISCSI") + self._migrate_step("ISCSI", None, "C1A") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_08_iscsi_cluster_to_cluster_same_cluster(self): + """Move iSCSI volumes between cluster pools in the same cluster.""" + self._migrate_step("ISCSI", "C1A", "C1B") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_09_iscsi_cluster_to_cluster_across_clusters(self): + """Move iSCSI volumes to a cluster pool in the second cluster.""" + self._migrate_step("ISCSI", "C1B", "C2") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_10_iscsi_cluster_to_zone(self): + """Move iSCSI volumes from a cluster pool to a zone-wide pool.""" + self._migrate_step("ISCSI", "C2", "Z1") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_11_iscsi_zone_to_zone(self): + """Move iSCSI volumes between two zone-wide pools.""" + self._migrate_step("ISCSI", "Z1", "Z2") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_12_iscsi_zone_to_cluster(self): + """Move iSCSI volumes from a zone-wide pool back to cluster C1a.""" + self._migrate_step("ISCSI", "Z2", "C1A") + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_13_reject_nfs3_to_iscsi(self): + """Reject a stopped-VM migration from NFS3 storage to iSCSI.""" + _, _, iscsi_pools = self._chain("ISCSI") + self._assert_rejected( + "NFS3", "C1A", iscsi_pools["C1B"], "ISCSI", + CROSS_PROTOCOL_REJECTION, + ) + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_14_reject_iscsi_to_nfs3(self): + """Reject a stopped-VM migration from iSCSI storage to NFS3.""" + _, _, nfs_pools = self._chain("NFS3") + self._assert_rejected( + "ISCSI", "C1A", nfs_pools["C1B"], "NFS3", + CROSS_PROTOCOL_REJECTION, + ) + + @attr(tags=["ontap_migration", "stopped_vm"], required_hardware=True) + def test_15_reject_running_vm(self): + """Reject the stopped-VM API form once the chain VM is running.""" + vm, _, pools = self._chain("ISCSI") + self._start_vm(vm) + self._assert_rejected( + "ISCSI", "C1A", pools["C1B"], "ISCSI", + "VM is not Stopped", + ) + self.assertEqual(self._get_vm(vm.id).state, "Running") diff --git a/test/integration/plugins/ontap/migration/test_03_volume_migration.py b/test/integration/plugins/ontap/migration/test_03_volume_migration.py new file mode 100644 index 000000000000..a5b3d3b384bc --- /dev/null +++ b/test/integration/plugins/ontap/migration/test_03_volume_migration.py @@ -0,0 +1,490 @@ +# 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. + +"""Ordered ONTAP volume-only migration tests. + +Workflow: + 01-06 carry one detached NFS3 volume through + DefaultPrimary -> C1a -> C1b -> C2 -> Z1 -> Z2 -> C1a + 07-12 carry one detached iSCSI volume through the same pool scopes + 13 migrate an NFS3 data volume attached to a stopped VM from C1a to Z1 + 14 migrate an iSCSI data volume attached to a stopped VM from C1a to Z1 + 15-16 migrate stopped-attached DefaultPrimary data to ONTAP NFS3/iSCSI + 17-20 reject detached and stopped-attached cross-protocol migration + 21 reject running-attached migration without livemigrate + 22 reject running-attached migration with livemigrate + 23 reject migration to the volume's current pool + +Every successful migration verifies CloudStack identity and attachment state, +the destination ONTAP object and access state, and source cleanup when the +source is ONTAP. Every rejection verifies that CloudStack and ONTAP state +remain unchanged. +""" + +import unittest + +from nose.plugins.attrib import attr + +from marvin.cloudstackAPI import migrateVolume as migrateVolumeAPI + +from migration.migration_test_base import OntapMigrationTestBase + + +class TestOntapVolumeMigration(OntapMigrationTestBase): + """Verify detached, stopped-attached, and rejected volume migrations.""" + + nfs_volume = None + iscsi_volume = None + nfs_c1a = None + nfs_c1b = None + nfs_c2 = None + nfs_z1 = None + nfs_z2 = None + iscsi_c1a = None + iscsi_c1b = None + iscsi_c2 = None + iscsi_z1 = None + iscsi_z2 = None + nfs_vm = None + nfs_attached_volume = None + iscsi_vm = None + iscsi_attached_volume = None + + @classmethod + def setUpClass(cls): + super(TestOntapVolumeMigration, cls).setUpClass() + cls._initialize_pools() + + @classmethod + def _host_requirements(cls): + return { + "minimum_hosts": 2, + "same_primary_cluster": False, + "multiple_clusters": True, + } + + @classmethod + def _pool_requirements(cls): + secondary_cluster_id = cls._require_secondary_cluster() + requirements = [ + ("CLUSTER", cls.cluster.id, 2), + ("CLUSTER", secondary_cluster_id, 1), + ("ZONE", None, 2), + ] + return { + "NFS3": list(requirements), + "ISCSI": list(requirements), + } + + @classmethod + def _initialize_pools(cls): + secondary_cluster_id = cls._require_secondary_cluster() + cls.nfs_c1a = cls._migration_pool("NFS3") + cls.nfs_c1b = cls._migration_pool("NFS3", 1) + cls.nfs_c2 = cls._migration_pool( + "NFS3", cluster_id=secondary_cluster_id + ) + cls.nfs_z1 = cls._migration_pool("NFS3", scope="ZONE") + cls.nfs_z2 = cls._migration_pool("NFS3", 1, scope="ZONE") + cls.iscsi_c1a = cls._migration_pool("ISCSI") + cls.iscsi_c1b = cls._migration_pool("ISCSI", 1) + cls.iscsi_c2 = cls._migration_pool( + "ISCSI", cluster_id=secondary_cluster_id + ) + cls.iscsi_z1 = cls._migration_pool("ISCSI", scope="ZONE") + cls.iscsi_z2 = cls._migration_pool("ISCSI", 1, scope="ZONE") + + @classmethod + def _require_class_state(cls, *attributes): + missing = [ + attribute for attribute in attributes + if getattr(cls, attribute, None) is None + ] + if missing: + raise unittest.SkipTest( + "Required earlier test state is unavailable: %s" + % ", ".join(missing) + ) + + def _migrate_detached(self, protocol, source_pool, destination_pool): + cls = self.__class__ + attribute = "nfs_volume" if protocol == "NFS3" else "iscsi_volume" + cls._require_class_state(attribute) + volume = self._get_volume(getattr(cls, attribute).id) + source_path = volume.path + migrated = self._migrate_volume_offline(volume, destination_pool) + current = self._assert_migration_success( + [volume], + destination_pool, + protocol, + expected_vm_id=None, + source_pool=source_pool, + )[0] + if protocol != "ISCSI": + self.assertNotEqual(current.path, source_path) + setattr(cls, attribute, migrated) + self._assert_offline_host_state( + protocol, destination_pool, volumes=[current] + ) + + def _root_volume(self, vm): + roots = [ + volume for volume in self._vm_volumes(vm.id) + if str(getattr(volume, "type", "")).upper() == "ROOT" + ] + self.assertEqual(len(roots), 1, "VM %s must have one root volume" % vm.id) + return roots[0] + + def _prepare_stopped_attached(self, protocol, source_pool): + vm = self._deploy_vm( + self._default_storage_tag(), + self._host_for_cluster(self.__class__.cluster.id).id, + ) + root = self._root_volume(vm) + root_pool_id = root.storageid + volume = self._create_data_volume(source_pool) + attached = self._attach_volume(vm, volume) + self._stop_vm(vm) + return vm, attached, root.id, root_pool_id + + def _migrate_stopped_attached( + self, protocol, source_pool, destination_pool): + vm, volume, root_id, root_pool_id = ( + self._prepare_stopped_attached(protocol, source_pool) + ) + source_path = volume.path + self._migrate_volume_offline(volume, destination_pool) + current = self._assert_migration_success( + [volume], + destination_pool, + protocol, + vm=vm, + expected_vm_state="Stopped", + expected_vm_id=vm.id, + source_pool=( + None if source_pool.id == self.__class__.default_pool.id + else source_pool + ), + )[0] + if protocol != "ISCSI": + self.assertNotEqual(current.path, source_path) + root = self._get_volume(root_id) + self.assertEqual(root.storageid, root_pool_id) + self.assertEqual( + getattr(root, "virtualmachineid", None), + vm.id, + ) + self._assert_offline_host_state( + protocol, destination_pool, vm=vm, volumes=[current] + ) + return self._get_vm(vm.id), current + + def _request_volume_migration(self, volume, pool, live_value=None): + cmd = migrateVolumeAPI.migrateVolumeCmd() + cmd.volumeid = volume.id + cmd.storageid = pool.id + if live_value is not None: + cmd.livemigrate = live_value + return self.apiClient.migrateVolume(cmd) + + def _assert_rejected( + self, vm, volume, source_pool, destination_pool, protocol, + live_value, pattern=None, destination_protocol=None): + snapshot = self._snapshot_state( + vm=vm, + volumes=[volume], + source_pool=source_pool, + protocol=protocol, + destination_pool=destination_pool, + destination_protocol=destination_protocol, + ) + host_snapshot = ( + self.__class__._host_vm_snapshot(vm) + if vm is not None + else self.__class__._host_volume_snapshot(volume) + ) + dest_ids = self._volume_ids_on_pool(destination_pool) + with self.assertRaises(Exception) as context: + self._request_volume_migration( + volume, destination_pool, live_value + ) + if pattern is not None: + message = str(context.exception).lower().replace(" ", "") + self.assertIn(pattern.lower().replace(" ", ""), message) + self._finish_rejected_migration( + [volume], destination_pool, dest_ids + ) + self._assert_state_unchanged( + snapshot, + vm=vm, + source_pool=source_pool, + protocol=protocol, + destination_pool=destination_pool, + ) + if vm is not None: + self._assert_host_snapshot_unchanged(host_snapshot, vm) + else: + self.assertEqual( + self.__class__._host_volume_snapshot(volume), + host_snapshot, + "Rejected migration changed KVM volume references", + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_01_nfs3_detached_default_primary_to_cluster(self): + """Migrate a detached NFS3 volume from DefaultPrimary to C1a.""" + self.__class__.nfs_volume = self._create_data_volume( + self.__class__.default_pool + ) + self._migrate_detached( + "NFS3", self.__class__.default_pool, self.__class__.nfs_c1a + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_02_nfs3_detached_cluster_to_cluster_same_cluster(self): + """Migrate the detached NFS3 volume from C1a to C1b.""" + self._migrate_detached( + "NFS3", self.__class__.nfs_c1a, self.__class__.nfs_c1b + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_03_nfs3_detached_cluster_to_cluster_across_clusters(self): + """Migrate the detached NFS3 volume from C1b to C2.""" + self._migrate_detached( + "NFS3", self.__class__.nfs_c1b, self.__class__.nfs_c2 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_04_nfs3_detached_cluster_to_zone(self): + """Migrate the detached NFS3 volume from C2 to Z1.""" + self._migrate_detached( + "NFS3", self.__class__.nfs_c2, self.__class__.nfs_z1 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_05_nfs3_detached_zone_to_zone(self): + """Migrate the detached NFS3 volume from Z1 to Z2.""" + self._migrate_detached( + "NFS3", self.__class__.nfs_z1, self.__class__.nfs_z2 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_06_nfs3_detached_zone_to_cluster(self): + """Migrate the detached NFS3 volume from Z2 back to C1a.""" + self._migrate_detached( + "NFS3", self.__class__.nfs_z2, self.__class__.nfs_c1a + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_07_iscsi_detached_default_primary_to_cluster(self): + """Migrate a detached iSCSI volume from DefaultPrimary to C1a.""" + self.__class__.iscsi_volume = self._create_data_volume( + self.__class__.default_pool + ) + self._migrate_detached( + "ISCSI", self.__class__.default_pool, self.__class__.iscsi_c1a + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_08_iscsi_detached_cluster_to_cluster_same_cluster(self): + """Migrate the detached iSCSI volume from C1a to C1b.""" + self._migrate_detached( + "ISCSI", self.__class__.iscsi_c1a, self.__class__.iscsi_c1b + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_09_iscsi_detached_cluster_to_cluster_across_clusters(self): + """Migrate the detached iSCSI volume from C1b to C2.""" + self._migrate_detached( + "ISCSI", self.__class__.iscsi_c1b, self.__class__.iscsi_c2 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_10_iscsi_detached_cluster_to_zone(self): + """Migrate the detached iSCSI volume from C2 to Z1.""" + self._migrate_detached( + "ISCSI", self.__class__.iscsi_c2, self.__class__.iscsi_z1 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_11_iscsi_detached_zone_to_zone(self): + """Migrate the detached iSCSI volume from Z1 to Z2.""" + self._migrate_detached( + "ISCSI", self.__class__.iscsi_z1, self.__class__.iscsi_z2 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_12_iscsi_detached_zone_to_cluster(self): + """Migrate the detached iSCSI volume from Z2 back to C1a.""" + self._migrate_detached( + "ISCSI", self.__class__.iscsi_z2, self.__class__.iscsi_c1a + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_13_nfs3_attached_stopped_vm_cluster_to_zone(self): + """Move stopped-attached NFS3 data C1a to Z1, preserving its root.""" + ( + self.__class__.nfs_vm, + self.__class__.nfs_attached_volume, + ) = self._migrate_stopped_attached( + "NFS3", self.__class__.nfs_c1a, self.__class__.nfs_z1 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_14_iscsi_attached_stopped_vm_cluster_to_zone(self): + """Move stopped-attached iSCSI data C1a to Z1, preserving its root.""" + ( + self.__class__.iscsi_vm, + self.__class__.iscsi_attached_volume, + ) = self._migrate_stopped_attached( + "ISCSI", self.__class__.iscsi_c1a, self.__class__.iscsi_z1 + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_15_nfs3_attached_stopped_vm_default_to_ontap(self): + """Move stopped-attached non-ONTAP data to ONTAP NFS3.""" + self._migrate_stopped_attached( + "NFS3", self.__class__.default_pool, self.__class__.nfs_c1a + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_16_iscsi_attached_stopped_vm_default_to_ontap(self): + """Move stopped-attached non-ONTAP data to ONTAP iSCSI.""" + self._migrate_stopped_attached( + "ISCSI", self.__class__.default_pool, self.__class__.iscsi_c1a + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_17_reject_detached_nfs3_to_iscsi(self): + """Reject detached NFS3-to-iSCSI migration.""" + self.__class__._require_class_state("nfs_volume") + self._assert_rejected( + None, self.__class__.nfs_volume, self.__class__.nfs_c1a, + self.__class__.iscsi_c1b, "NFS3", False, "cross-protocol", + destination_protocol="ISCSI", + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_18_reject_detached_iscsi_to_nfs3(self): + """Reject detached iSCSI-to-NFS3 migration.""" + self.__class__._require_class_state("iscsi_volume") + self._assert_rejected( + None, self.__class__.iscsi_volume, self.__class__.iscsi_c1a, + self.__class__.nfs_c1b, "ISCSI", False, "cross-protocol", + destination_protocol="NFS3", + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_19_reject_stopped_attached_nfs3_to_iscsi(self): + """Reject stopped-attached NFS3-to-iSCSI migration.""" + self.__class__._require_class_state( + "nfs_vm", "nfs_attached_volume" + ) + self._assert_rejected( + self.__class__.nfs_vm, self.__class__.nfs_attached_volume, + self.__class__.nfs_z1, self.__class__.iscsi_z2, + "NFS3", False, "cross-protocol", + destination_protocol="ISCSI", + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_20_reject_stopped_attached_iscsi_to_nfs3(self): + """Reject stopped-attached iSCSI-to-NFS3 migration.""" + self.__class__._require_class_state( + "iscsi_vm", "iscsi_attached_volume" + ) + self._assert_rejected( + self.__class__.iscsi_vm, self.__class__.iscsi_attached_volume, + self.__class__.iscsi_z1, self.__class__.nfs_z2, + "ISCSI", False, "cross-protocol", + destination_protocol="NFS3", + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_21_reject_attached_running_vm_without_livemigrate(self): + """Reject running-attached NFS3 migration without livemigrate.""" + self.__class__._require_class_state( + "nfs_vm", "nfs_attached_volume" + ) + self.__class__.nfs_vm = self._start_vm(self.__class__.nfs_vm) + self.__class__.nfs_attached_volume = self._get_volume( + self.__class__.nfs_attached_volume.id + ) + self._assert_rejected( + self.__class__.nfs_vm, + self.__class__.nfs_attached_volume, + self.__class__.nfs_z1, + self.__class__.nfs_c1a, + "NFS3", + None, + "migrateVirtualMachineWithVolume", + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_22_reject_attached_running_vm_with_livemigrate(self): + """Reject live volume-only migration in favor of VM-with-volume.""" + self.__class__._require_class_state( + "nfs_vm", "nfs_attached_volume" + ) + self._assert_rejected( + self.__class__.nfs_vm, + self.__class__.nfs_attached_volume, + self.__class__.nfs_z1, + self.__class__.nfs_c1a, + "NFS3", + True, + "migrateVirtualMachineWithVolume", + ) + + @attr(tags=["ontap_migration", "volume_migration"], + required_hardware=True) + def test_23_reject_same_pool(self): + """Reject migration to the stopped-attached volume's current pool.""" + self.__class__._require_class_state( + "iscsi_vm", "iscsi_attached_volume" + ) + self._assert_rejected( + self.__class__.iscsi_vm, + self.__class__.iscsi_attached_volume, + self.__class__.iscsi_z1, + self.__class__.iscsi_z1, + "ISCSI", + False, + "already on the destination storage pool", + ) diff --git a/test/integration/plugins/ontap/ontap_test_base.py b/test/integration/plugins/ontap/ontap_test_base.py index 17bec2c1251c..54ae62d4d4de 100644 --- a/test/integration/plugins/ontap/ontap_test_base.py +++ b/test/integration/plugins/ontap/ontap_test_base.py @@ -31,6 +31,7 @@ import logging import random import requests +import socket import sys import time import urllib3 @@ -112,6 +113,75 @@ def get_datacenter_config(testclient, test_cls): return cfg +def is_configured(value): + """Return whether a config value is non-empty and not a placeholder.""" + return value is not None and str(value).strip() != "" and not ( + str(value).startswith("<<") and str(value).endswith(">>") + ) + + +def normalize_template_name(name): + """Normalize common KVM template name variants for comparison.""" + if not name: + return "" + normalized = " ".join(name.lower().strip().split()) + return normalized.replace("(64 bit)", "(64-bit)").replace( + "(64bit)", "(64-bit)" + ) + + +def list_kvm_templates(api_client, zone_id): + """Return every KVM template visible in a zone.""" + cmd = listTemplatesAPI.listTemplatesCmd() + cmd.templatefilter = "all" + cmd.listall = True + cmd.zoneid = zone_id + return [ + template for template in (api_client.listTemplates(cmd) or []) + if str(getattr(template, "hypervisor", "")).lower() == "kvm" + ] + + +def find_kvm_template(api_client, zone_id, template_name): + """Find a KVM template by normalized name.""" + target = normalize_template_name(template_name) + for template in list_kvm_templates(api_client, zone_id): + if normalize_template_name(template.name) == target: + return template + return None + + +def host_aliases(url): + """Return the DNS name and resolved IP represented by a host URL.""" + hostname = urlparse(url).hostname or url + aliases = {hostname} + try: + aliases.add(socket.gethostbyname(hostname)) + except (socket.error, UnicodeError): + pass + return {alias for alias in aliases if alias} + + +def map_hosts_by_address(hosts): + """Index CloudStack hosts by both reported IP address and name.""" + addresses = {} + for host in hosts or []: + for attribute in ("ipaddress", "name"): + value = getattr(host, attribute, None) + if value: + addresses[value] = host + return addresses + + +def get_ready_hosts(hosts): + """Return hosts that are Up and enabled for resource allocation.""" + return [ + host for host in (hosts or []) + if str(getattr(host, "state", "")).lower() == "up" + and str(getattr(host, "resourcestate", "Enabled")).lower() == "enabled" + ] + + # --------------------------------------------------------------------------- # Pool detail helper # --------------------------------------------------------------------------- @@ -157,8 +227,20 @@ def __init__(self, storage_ip, username, password, port=443): def _get(self, path, params=None): url = self._base + path - resp = requests.get(url, auth=self._auth, params=params, - verify=False, timeout=30) + for attempt in range(1, 4): + try: + resp = requests.get(url, auth=self._auth, params=params, + verify=False, timeout=30) + break + except (requests.exceptions.ConnectionError, + requests.exceptions.Timeout): + if attempt == 3: + raise + logger.warning( + "Transient ONTAP REST GET failure for %s; retrying " + "(attempt %d of 3)", path, attempt + 1, + ) + time.sleep(attempt) resp.raise_for_status() return resp.json() @@ -185,6 +267,10 @@ def _patch(self, path, payload=None, params=None): return None return None + def check_connection(self): + """Return basic cluster data when the ONTAP REST endpoint is ready.""" + return self._get("/cluster", params={"fields": "name,uuid"}) + def delete_volume(self, name): """Delete the ONTAP FlexVol with the given name. No-op if not found.""" data = self._get("/storage/volumes", params={"name": name}) @@ -258,7 +344,7 @@ def get_lun(self, svm_name, lun_path): """Return the ONTAP LUN record for the given full path, or None.""" data = self._get("/storage/luns", params={"svm.name": svm_name, "name": lun_path, - "fields": "name,uuid,enabled,status"}) + "fields": "name,uuid,serial_number,enabled,status"}) records = data.get("records", []) return records[0] if records else None diff --git a/test/integration/plugins/ontap/run_tests.sh b/test/integration/plugins/ontap/run_tests.sh index 5a73adb22553..e26d8d7bffd7 100755 --- a/test/integration/plugins/ontap/run_tests.sh +++ b/test/integration/plugins/ontap/run_tests.sh @@ -24,6 +24,7 @@ # bash test/integration/plugins/ontap/run_tests.sh iscsi # all iSCSI suites # bash test/integration/plugins/ontap/run_tests.sh nfs3 # all NFS3 suites # bash test/integration/plugins/ontap/run_tests.sh both # iscsi then nfs3 +# bash test/integration/plugins/ontap/run_tests.sh migration # setup + ordered migration suites # bash test/integration/plugins/ontap/run_tests.sh nfs3_workflow # single suite ONTAP_DIR=test/integration/plugins/ontap @@ -79,6 +80,12 @@ NFS3_SUITES=( "NFS3 template cache negative|nfs3_template_cache_negative|${ONTAP_DIR}/nfs3/template/test_template_cache_negative.py" ) +MIGRATION_SUITES=( + "Live VM with storage migration|live_storage|${ONTAP_DIR}/migration/test_01_live_vm_with_storage_migration.py" + "Stopped VM storage migration|stopped_vm|${ONTAP_DIR}/migration/test_02_stopped_vm_storage_migration.py" + "Volume migration|volume_migration|${ONTAP_DIR}/migration/test_03_volume_migration.py" +) + record_results() { local tag="$1" local label="$2" @@ -229,6 +236,10 @@ should_run_tag() { nfs3) [[ "$tag" == nfs3_* || "$tag" == "zone_pool" || "$tag" == "vm_volume_workflow" ]] ;; + migration) + [[ "$tag" == "setup_zone" || "$tag" == "live_storage" || + "$tag" == "stopped_vm" || "$tag" == "volume_migration" ]] + ;; *) [[ "$FILTER" == "$tag" ]] ;; @@ -258,8 +269,18 @@ run_group() { fi set +e - $PYTHON "$NOSE_RUNNER" --with-marvin --marvin-config="$CFG" "$file" -a "tags=${tag}" -v -s 2>&1 | tee "$tmpout" - rc=${PIPESTATUS[0]} + local attempt + for attempt in 1 2 3; do + : > "$tmpout" + $PYTHON "$NOSE_RUNNER" --with-marvin --marvin-config="$CFG" "$file" -a "tags=${tag}" -v -s 2>&1 | tee "$tmpout" + rc=${PIPESTATUS[0]} + if grep -q "Marvin Init Failed" "$tmpout" && [[ "$attempt" -lt 3 ]]; then + echo " Marvin init failed on attempt ${attempt}; retrying in 5s..." + sleep 5 + continue + fi + break + done set -e out=$(cat "$tmpout") @@ -322,6 +343,22 @@ check_ontap_prereqs() { $PYTHON "$ONTAP_PREREQS" "$CFG" "$protocol" } +run_migration_suites() { + local entry label tag file index last_index + last_index=$((${#MIGRATION_SUITES[@]} - 1)) + for index in "${!MIGRATION_SUITES[@]}"; do + entry="${MIGRATION_SUITES[$index]}" + IFS='|' read -r label tag file <<< "$entry" + if [[ "$index" -lt "$last_index" ]]; then + export ONTAP_MIGRATION_KEEP_POOLS=1 + else + unset ONTAP_MIGRATION_KEEP_POOLS + fi + run_group "$label" "$tag" "$file" + done + unset ONTAP_MIGRATION_KEEP_POOLS +} + run_protocol_batch() { local protocol="$1" local parent_dir="${2:-}" @@ -341,10 +378,17 @@ run_protocol_batch() { finalize_batch } +run_migration_batch() { + local parent_dir="${1:-}" + init_batch "migration" "$parent_dir" + run_migration_suites + finalize_batch +} + run_single_suite_by_tag() { local want_tag="$1" local entry label tag file - for entry in "${ISCSI_SUITES[@]}" "${NFS3_SUITES[@]}"; do + for entry in "${ISCSI_SUITES[@]}" "${NFS3_SUITES[@]}" "${MIGRATION_SUITES[@]}"; do IFS='|' read -r label tag file <<< "$entry" if [[ "$tag" == "$want_tag" ]]; then if [[ "$want_tag" == iscsi_* ]]; then @@ -385,6 +429,7 @@ case "$FILTER" in run_protocol_batch iscsi "$BOTH_DIR" run_protocol_batch nfs3 "$BOTH_DIR" + run_migration_batch "$BOTH_DIR" write_combined_both_summary "$BOTH_DIR" ;; both) @@ -402,6 +447,11 @@ case "$FILTER" in nfs3) run_protocol_batch nfs3 ;; + migration) + run_group "Advanced zone setup" "setup_zone" \ + "${ONTAP_DIR}/zone_setup/test_setup_zone.py" + run_migration_batch + ;; setup_zone) run_group "Advanced zone setup" "setup_zone" \ "${ONTAP_DIR}/zone_setup/test_setup_zone.py" @@ -417,7 +467,7 @@ case "$FILTER" in print_final_summary else echo "Unknown filter: $FILTER" >&2 - echo "Use: all | both | iscsi | nfs3 | setup_zone | cleanup_zone | " >&2 + echo "Use: all | both | iscsi | nfs3 | migration | setup_zone | cleanup_zone | " >&2 exit 1 fi ;;