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
44 changes: 38 additions & 6 deletions coriolis/tests/integration/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,12 @@ def _create_pool(
skip_allocation=True,
wait_for_allocation=False,
platform=constants.PROVIDER_PLATFORM_DESTINATION,
minimum_minions=1,
maximum_minions=1,
minion_max_idle_time=3600,
minion_retention_strategy=(
constants.MINION_POOL_MACHINE_RETENTION_STRATEGY_DELETE
),
):
env_options = (
cls._imp_pool_env
Expand All @@ -184,12 +190,10 @@ def _create_pool(
platform=platform,
os_type=constants.OS_TYPE_LINUX,
environment_options=env_options,
minimum_minions=1,
maximum_minions=1,
minion_max_idle_time=3600,
minion_retention_strategy=(
constants.MINION_POOL_MACHINE_RETENTION_STRATEGY_DELETE
),
minimum_minions=minimum_minions,
maximum_minions=maximum_minions,
minion_max_idle_time=minion_max_idle_time,
minion_retention_strategy=minion_retention_strategy,
skip_allocation=skip_allocation,
)
cls.addClassCleanup(cls._safe_delete_pool, pool.id)
Expand Down Expand Up @@ -305,6 +309,15 @@ class ReplicaIntegrationTestBase(CoriolisIntegrationTestBase):
# source_environment.
_EXTRA_SOURCE_ENVIRONMENT = {}

# Overridable params for the pool(s) created when _CREATE_DST_MINION_POOL /
# _CREATE_SRC_MINION_POOL is set.
_POOL_MINIMUM_MINIONS = 1
_POOL_MAXIMUM_MINIONS = 1
_POOL_MINION_MAX_IDLE_TIME = 3600
_POOL_MINION_RETENTION_STRATEGY = (
constants.MINION_POOL_MACHINE_RETENTION_STRATEGY_DELETE
)

@classmethod
def setUpClass(cls):
super().setUpClass()
Expand All @@ -331,6 +344,10 @@ def setUpClass(cls):
"dst-transfer-pool",
skip_allocation=False,
wait_for_allocation=True,
minimum_minions=cls._POOL_MINIMUM_MINIONS,
maximum_minions=cls._POOL_MAXIMUM_MINIONS,
minion_max_idle_time=cls._POOL_MINION_MAX_IDLE_TIME,
minion_retention_strategy=cls._POOL_MINION_RETENTION_STRATEGY,
)
cls._dst_pool_id = pool.id

Expand All @@ -343,6 +360,10 @@ def setUpClass(cls):
skip_allocation=False,
wait_for_allocation=True,
platform=constants.PROVIDER_PLATFORM_SOURCE,
minimum_minions=cls._POOL_MINIMUM_MINIONS,
maximum_minions=cls._POOL_MAXIMUM_MINIONS,
minion_max_idle_time=cls._POOL_MINION_MAX_IDLE_TIME,
minion_retention_strategy=cls._POOL_MINION_RETENTION_STRATEGY,
)
cls._src_pool_id = pool.id

Expand Down Expand Up @@ -430,6 +451,17 @@ 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 = [
self._client.transfer_executions.create(
transfer_id, shutdown_instances=False
)
for transfer_id in transfer_ids
]
for execution in executions:
self.assertExecutionCompleted(execution.id, timeout=timeout)

def _execute_transfer_and_deployment(self, deployment_kwargs=None):
deployment_kwargs = deployment_kwargs or {}

Expand Down
270 changes: 270 additions & 0 deletions coriolis/tests/integration/test_minion_pools.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

from coriolis import constants
from coriolis.db import api as db_api
from coriolis.minion_manager.rpc import server as minion_manager_rpc_server
from coriolis.tests.integration import base

CONF = cfg.CONF
Expand Down Expand Up @@ -194,3 +195,272 @@ def _create_pool(self, endpoint_id, **kwargs):
return super()._create_pool(
endpoint_id, platform=constants.PROVIDER_PLATFORM_SOURCE, **kwargs
)


class _MinionPoolPowerCycleTestMixin:
"""Transfer that reuses pool machines across a power cycle.

The pool allows up to 2 machines (minimum 1) with a tiny idle time and the
"poweroff" retention strategy. Two separate transfers are executed concurrently;
the second execution finds the pre-existing minimum machine already reserved by the
first and allocates a brand new one instead. Once both go idle, refreshing the pool
powers the excess one off. Re-running both transfers concurrently then reuses both
machines, powering the idled-off one back on before healthchecking and reusing it.

Subclasses select which side's pool gets exercised by overriding the ``_pool_id``
property.
"""

_POOL_MAXIMUM_MINIONS = 2
_POOL_MINION_MAX_IDLE_TIME = 1
_POOL_MINION_RETENTION_STRATEGY = (
constants.MINION_POOL_MACHINE_RETENTION_STRATEGY_POWEROFF
)

@property
def _pool_id(self):
raise NotImplementedError

def setUp(self):
super().setUp()

# A second transfer, independent from self._transfer (created by
# ReplicaIntegrationTestBase.setUp). Running it concurrently with self._transfer
# forces the pool to allocate a second machine, since the first is already
# reserved by self._transfer's execution.
self._pool_transfer_b = self._create_transfer(
self._src_endpoint.id,
self._dst_endpoint.id,
instances=[self._instance_name],
source_environment=self._transfer._info["source_environment"],
destination_minion_pool_id=self._dst_pool_id,
origin_minion_pool_id=self._src_pool_id,
)

def _wait_for_power_status(self, status, timeout=120):
"""Poll until one of the pool's machines reaches *status*."""
ctxt = self._get_db_context()
deadline = time.monotonic() + timeout
machines = []

while time.monotonic() < deadline:
pool = db_api.get_minion_pool(ctxt, self._pool_id, include_machines=True)
machines = pool.minion_machines
if any(m.power_status == status for m in machines):
return machines
time.sleep(1)

self.fail(
"No minion machine of pool '%s' reached power status '%s' within %ds "
"(last statuses: %s)"
% (
self._pool_id,
status,
timeout,
[m.power_status for m in machines],
)
)

def test_transfer_after_pool_machine_power_cycle(self):
transfer_ids = [self._transfer.id, self._pool_transfer_b.id]

# Concurrently executing both transfers forces the second one to allocate a new
# machine, since the pre-existing minimum one is already reserved by the first
# (up to the pool's maximum of 2).
self._execute_concurrently_and_wait(transfer_ids)

pool = db_api.get_minion_pool(
self._get_db_context(), self._pool_id, include_machines=True
)
self.assertEqual(2, len(pool.minion_machines))
provider_properties_before = {
machine.id: machine.provider_properties for machine in pool.minion_machines
}

# Let both machines' idle time expire, then refresh the pool: since their count
# exceeds the pool minimum of 1, the excess one gets powered off.
time.sleep(self._POOL_MINION_MAX_IDLE_TIME + 1)
self._client.minion_pools.refresh_minion_pool(self._pool_id)
self._wait_for_power_status(constants.MINION_MACHINE_POWER_STATUS_POWERED_OFF)

# Re-running both transfers concurrently reuses both machines, powering the
# idled-off one back on before healthchecking and reusing it.
self._execute_concurrently_and_wait(transfer_ids)

# The power-cycled machine must genuinely have been reused, not silently deleted
# and recreated from scratch by the healthcheck-failure fallback (which would
# defeat the whole point of the "poweroff" retention strategy)
pool = db_api.get_minion_pool(
self._get_db_context(), self._pool_id, include_machines=True
)
self.assertEqual(2, len(pool.minion_machines))
for machine in pool.minion_machines:
self.assertEqual(
provider_properties_before[machine.id],
machine.provider_properties,
"Minion machine '%s' provider properties changed across the power "
"cycle; it was likely deleted and recreated instead of reused."
% machine.id,
)


class MinionPoolPowerCycleTransferTest(
_MinionPoolPowerCycleTestMixin, base.MinionPoolReplicaTestBase
):
"""Power-cycle test exercising a destination minion pool."""

@property
def _pool_id(self):
return self._dst_pool_id


class SourceMinionPoolPowerCycleTransferTest(
_MinionPoolPowerCycleTestMixin, base.SourceMinionPoolReplicaTestBase
):
"""Power-cycle test exercising a source minion pool."""

@property
def _pool_id(self):
return self._src_pool_id


class _MinionPoolRefreshDeallocationTestMixin:
"""Excess pool machine gets deleted on refresh.

Mirrors _MinionPoolPowerCycleTestMixin but with the default "delete" retention
strategy: once the pool's excess machine (beyond its minimum of 1) goes idle,
refreshing the pool deletes it instead of powering it off, exercising

Subclasses select which side's pool gets exercised by overriding the ``_pool_id``
property.
"""

_POOL_MAXIMUM_MINIONS = 2
_POOL_MINION_MAX_IDLE_TIME = 1

@property
def _pool_id(self):
raise NotImplementedError

def setUp(self):
super().setUp()

# A second transfer, independent from self._transfer (created by
# ReplicaIntegrationTestBase.setUp). Running it concurrently with self._transfer
# forces the pool to allocate a second machine, since the first is already
# reserved by self._transfer's execution.
self._pool_transfer_b = self._create_transfer(
self._src_endpoint.id,
self._dst_endpoint.id,
instances=[self._instance_name],
source_environment=self._transfer._info["source_environment"],
destination_minion_pool_id=self._dst_pool_id,
origin_minion_pool_id=self._src_pool_id,
)

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

# Concurrently executing both transfers forces the second one to allocate a new
# machine, since the pre-existing minimum one is already reserved by the first
# (up to the pool's maximum of 2).
self._execute_concurrently_and_wait(transfer_ids)

pool = db_api.get_minion_pool(
self._get_db_context(), self._pool_id, include_machines=True
)
self.assertEqual(2, len(pool.minion_machines))

# Let both machines' idle time expire, then refresh the pool: since their count
# exceeds the pool minimum of 1, the excess one gets deleted.
time.sleep(self._POOL_MINION_MAX_IDLE_TIME + 1)
self._client.minion_pools.refresh_minion_pool(self._pool_id)

ctxt = self._get_db_context()
deadline = time.monotonic() + 120
pool = None
while time.monotonic() < deadline:
pool = db_api.get_minion_pool(ctxt, self._pool_id, include_machines=True)
if len(pool.minion_machines) == 1:
break
time.sleep(1)

self.assertEqual(
1,
len(pool.minion_machines),
"Expected the excess minion machine to be deleted from pool '%s'; "
"machines still present: %s"
% (self._pool_id, [m.id for m in pool.minion_machines]),
)


class MinionPoolRefreshDeallocationTransferTest(
_MinionPoolRefreshDeallocationTestMixin, base.MinionPoolReplicaTestBase
):
"""Deletion-on-refresh test exercising a destination minion pool."""

@property
def _pool_id(self):
return self._dst_pool_id


class SourceMinionPoolRefreshDeallocationTransferTest(
_MinionPoolRefreshDeallocationTestMixin, base.SourceMinionPoolReplicaTestBase
):
"""Deletion-on-refresh test exercising a source minion pool."""

@property
def _pool_id(self):
return self._src_pool_id


class MinionPoolRefreshCronStartupTest(base.DestinationMinionPoolTestBase):
"""Cron jobs are re-registered for pre-existing pools on startup.

`_init_pools_refresh_cron_jobs` runs once, when a minion manager service endpoint
is instantiated, and scans the DB for already-ALLOCATED pools to re-register their
periodic refresh jobs (e.g.: after a service restart while pools were still
allocated).
"""

def setUp(self):
super().setUp()

self._endpoint = self._create_endpoint(
name="pool-cron-dst",
endpoint_type=self._imp_platform,
connection_info=self._imp_conn_info,
)

def test_startup_registers_refresh_jobs_for_existing_pools(self):
# The harness disables automatic refreshing by default (period 0) to avoid
# interference with other tests. Re-enable it so the new endpoint being
# constructed below actually registers jobs.
CONF.set_override(
"minion_pool_default_refresh_period_minutes", 1, group="minion_manager"
)
self.addCleanup(
CONF.clear_override,
"minion_pool_default_refresh_period_minutes",
group="minion_manager",
)

pool = self._create_pool(
self._endpoint.id, skip_allocation=False, wait_for_allocation=True
)

new_endpoint = minion_manager_rpc_server.MinionManagerServerEndpoint()
self.addCleanup(new_endpoint._cron.stop)

job_prefix = (
minion_manager_rpc_server.MINION_POOL_REFRESH_JOB_PREFIX_FORMAT % pool.id
)
registered = [
name for name in new_endpoint._cron._jobs if name.startswith(job_prefix)
]
self.assertTrue(
registered,
"Expected refresh cron jobs to be registered on startup for pre-existing "
"allocated pool '%s', got jobs: %s"
% (pool.id, list(new_endpoint._cron._jobs)),
)
3 changes: 3 additions & 0 deletions coriolis/tests/integration/test_provider/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -179,5 +179,8 @@ def healthcheck_minion(
username = minion_connection_info.get("username", "root")
pkey = minion_connection_info.get("pkey")

# A freshly power-cycled minion needs a moment to boot sshd back up.
coriolis_utils.wait_for_port_connectivity(ip, port, max_wait=60)

client = coriolis_utils.connect_ssh(ip, port, username, pkey=pkey)
client.close()
Loading
Loading