Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 29 additions & 10 deletions coriolis/tests/integration/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -390,7 +390,7 @@ def setUp(self):
# Safety-net cleanup for destination devices allocated by the provider.
# Must be registered after the transfer, so it runs (LIFO) before the
# transfer delete, while the volumes_info is still in the DB.
self.addCleanup(self._cleanup_provider_dst_devices)
self.addCleanup(self._cleanup_provider_dst_devices, self._transfer.id)

# mock a few commands that are going to be ran through ssh; they won't
# pass anyway.
Expand Down Expand Up @@ -421,25 +421,25 @@ def _dst_device(self):

return None

def _cleanup_provider_dst_devices(self):
"""Remove any devices the provider allocated for this test."""
def _cleanup_provider_dst_devices(self, transfer_id):
"""Remove any devices the provider allocated for this transfer."""
ctxt = self._get_db_context()

try:
executions = db_api.get_transfer_tasks_executions(
ctxt,
self._transfer.id,
transfer_id,
sort_keys=["number"],
sort_dirs=["desc"],
limit=1,
)
if executions:
self._cleanup_execution(self._transfer.id, executions[0].id)
self._cleanup_execution(transfer_id, executions[0].id)
except Exception as ex:
LOG.warn("Could not cancel execution during cleanup. Ex: %s", ex)

try:
deletion = self._client.transfers.delete_disks(self._transfer.id)
deletion = self._client.transfers.delete_disks(transfer_id)
self.wait_for_execution(deletion.id, timeout=60)
except Exception as ex:
LOG.warn("Could not clean up provider dest devices. Ex: %s", ex)
Expand All @@ -451,15 +451,18 @@ def _execute_and_wait(self, transfer_id, timeout=600):
)
self.assertExecutionCompleted(execution.id, timeout=timeout)

def _execute_concurrently_and_wait(self, transfer_ids, timeout=600):
"""Start one execution per transfer id before waiting on any."""
executions = [
def _start_executions(self, transfer_ids):
"""Start one execution per transfer id without waiting on any."""
return [
self._client.transfer_executions.create(
transfer_id, shutdown_instances=False
)
for transfer_id in transfer_ids
]
for execution in executions:

def _execute_concurrently_and_wait(self, transfer_ids, timeout=600):
"""Start one execution per transfer id before waiting on any."""
for execution in self._start_executions(transfer_ids):
self.assertExecutionCompleted(execution.id, timeout=timeout)

def _execute_transfer_and_deployment(self, deployment_kwargs=None):
Expand Down Expand Up @@ -745,6 +748,14 @@ def assertExecutionCompleted(self, execution_id, timeout=600):
self.assertPoolAllocated(self._dst_pool_id)
self.assertMachinesAvailable(self._dst_pool_id)

def _execute_concurrently_and_wait(self, transfer_ids, timeout=600):
# Machines allocated for a still-running execution are not AVAILABLE yet,
# so only check the pool once every execution is done.
for execution in self._start_executions(transfer_ids):
super().assertExecutionCompleted(execution.id, timeout=timeout)
self.assertPoolAllocated(self._dst_pool_id)
self.assertMachinesAvailable(self._dst_pool_id)

def assertDeploymentCompleted(self, deployment_id, timeout=600):
super().assertDeploymentCompleted(deployment_id, timeout=timeout)
self.assertPoolAllocated(self._dst_pool_id)
Expand All @@ -771,3 +782,11 @@ def assertExecutionCompleted(self, execution_id, timeout=600):
super().assertExecutionCompleted(execution_id, timeout=timeout)
self.assertPoolAllocated(self._src_pool_id)
self.assertMachinesAvailable(self._src_pool_id)

def _execute_concurrently_and_wait(self, transfer_ids, timeout=600):
# Machines allocated for a still-running execution are not AVAILABLE yet,
# so only check the pool once every execution is done.
for execution in self._start_executions(transfer_ids):
super().assertExecutionCompleted(execution.id, timeout=timeout)
self.assertPoolAllocated(self._src_pool_id)
self.assertMachinesAvailable(self._src_pool_id)
4 changes: 4 additions & 0 deletions coriolis/tests/integration/test_minion_pools.py
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,8 @@ def setUp(self):
destination_minion_pool_id=self._dst_pool_id,
origin_minion_pool_id=self._src_pool_id,
)
# Registered after the transfer, so it runs (LIFO) before the transfer delete.
self.addCleanup(self._cleanup_provider_dst_devices, self._pool_transfer_b.id)

def _wait_for_power_status(self, status, timeout=120):
"""Poll until one of the pool's machines reaches *status*."""
Expand Down Expand Up @@ -357,6 +359,8 @@ def setUp(self):
destination_minion_pool_id=self._dst_pool_id,
origin_minion_pool_id=self._src_pool_id,
)
# Registered after the transfer, so it runs (LIFO) before the transfer delete.
self.addCleanup(self._cleanup_provider_dst_devices, self._pool_transfer_b.id)

def test_excess_pool_machine_deleted_on_refresh(self):
transfer_ids = [self._transfer.id, self._pool_transfer_b.id]
Expand Down
Loading