From 42b223c1e006f57fdb3609168c29a220063bc6ef Mon Sep 17 00:00:00 2001 From: Claudiu Belu Date: Fri, 11 Sep 2026 08:09:53 +0000 Subject: [PATCH] integration: Use minion pools by default instead of temp workers Change ReplicaIntegrationTestBase to use minion pools by default if the providers supports them, instead of only in Minion-specific test classes. If a provider doesn't support minion pools, temporary workers will be used instead. These minion pools will be created once per test run, instead of per class (setUpClass), meaning that they will be reused across multiple tests. Minion pool-specific tests that disrupt the minion pool will create their own instead. Some tests will still use temporary workers, as they still need to be tested. --- coriolis/tests/integration/base.py | 367 +++++++++++++----- .../deployments/test_deployment.py | 6 - .../deployments/test_luks_osmorphing.py | 13 + .../deployments/test_osmorphing.py | 58 ++- .../integration/management/test_region.py | 27 ++ .../integration/test_failure_recovery.py | 19 +- .../tests/integration/test_minion_pools.py | 28 +- .../tests/integration/test_provider/imp.py | 12 +- .../integration/transfers/test_executions.py | 187 +++++---- .../integration/transfers/test_transfer.py | 35 +- 10 files changed, 513 insertions(+), 239 deletions(-) diff --git a/coriolis/tests/integration/base.py b/coriolis/tests/integration/base.py index 58407efe9..9ecf5bd2f 100644 --- a/coriolis/tests/integration/base.py +++ b/coriolis/tests/integration/base.py @@ -12,6 +12,7 @@ Subclasses must be run as root. """ +import atexit import os import time import unittest @@ -46,6 +47,22 @@ constants.MINION_POOL_STATUS_ERROR, } +# We'll be creating a minion pool per provider (for providers that support them), and +# we'll be using them for most of the integration tests. Some transfer / OS morphing +# tests will run without minion pools, as those scenarios still have to be tested. +_DEFAULT_POOL_MINIMUM_MINIONS = 1 +_DEFAULT_POOL_MAXIMUM_MINIONS = 1 +_DEFAULT_POOL_MINION_MAX_IDLE_TIME = 3600 +_DEFAULT_POOL_MINION_RETENTION_STRATEGY = ( + constants.MINION_POOL_MACHINE_RETENTION_STRATEGY_DELETE +) + +# Default endpoints / pools, created lazily once per run as needed +# (_get_shared_endpoint / _get_shared_pool), keyed by provider platform +# (constants.PROVIDER_PLATFORM_SOURCE / _DESTINATION). +_SHARED_ENDPOINTS = {} +_SHARED_POOL_IDS = {} + class CoriolisIntegrationTestBase(test_base.CoriolisBaseTestCase): """Base class for integration tests.""" @@ -172,13 +189,13 @@ 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 - ), + minimum_minions=_DEFAULT_POOL_MINIMUM_MINIONS, + maximum_minions=_DEFAULT_POOL_MAXIMUM_MINIONS, + minion_max_idle_time=_DEFAULT_POOL_MINION_MAX_IDLE_TIME, + minion_retention_strategy=_DEFAULT_POOL_MINION_RETENTION_STRATEGY, + class_cleanup=True, ): + """Create a pool, deleted at class teardown if *class_cleanup*.""" env_options = ( cls._imp_pool_env if platform == constants.PROVIDER_PLATFORM_DESTINATION @@ -196,7 +213,8 @@ def _create_pool( minion_retention_strategy=minion_retention_strategy, skip_allocation=skip_allocation, ) - cls.addClassCleanup(cls._safe_delete_pool, pool.id) + if class_cleanup: + cls.addClassCleanup(cls._safe_delete_pool, pool.id) if wait_for_allocation: pool_obj = cls._wait_for_pool(pool.id, MINION_ALLOCATED_TERMINAL) @@ -265,10 +283,27 @@ def assertPoolAllocated(self, pool_id): "Pool %s is not ALLOCATED (got %s)" % (pool_id, pool.status), ) - def assertMachinesAvailable(self, pool_id): - """Assert all machines in the pool are AVAILABLE and have been used.""" + def assertMachinesAvailable(self, pool_id, timeout=30): + """Assert all machines in the pool are AVAILABLE and have been used. + + Machines are healthchecked asynchronously after being released, so wait + up to *timeout* seconds for them to settle before asserting. + """ ctxt = self._get_db_context() - pool = db_api.get_minion_pool(ctxt, pool_id, include_machines=True) + deadline = time.monotonic() + timeout + while True: + pool = db_api.get_minion_pool(ctxt, pool_id, include_machines=True) + if ( + pool is None + or all( + m.allocation_status == constants.MINION_MACHINE_STATUS_AVAILABLE + for m in pool.minion_machines + ) + or time.monotonic() >= deadline + ): + break + time.sleep(1) + self.assertIsNotNone(pool, "Pool %s not found" % pool_id) self.assertTrue( pool.minion_machines, @@ -299,73 +334,138 @@ def f(*args, **kwargs): return f + @classmethod + def _get_shared_endpoint(cls, platform): + """Lazily create and cache (once per run) the default endpoint. + + *platform* is constants.PROVIDER_PLATFORM_SOURCE or _DESTINATION. + Created once for the whole run and reused by test classes. Torn down via + atexit, not addClassCleanup. + + Tests that mutate endpoint state (e.g.: mapped_regions) must create their own + endpoint instead of using this one. + """ + endpoint = _SHARED_ENDPOINTS.get(platform) + if endpoint is not None: + return endpoint + + if platform == constants.PROVIDER_PLATFORM_SOURCE: + endpoint = cls._client.endpoints.create( + name="shared-test-src", + endpoint_type=cls._exp_platform, + description="shared integration source endpoint", + connection_info=cls._exp_conn_info, + regions=[], + ) + else: + endpoint = cls._client.endpoints.create( + name="shared-test-dest", + endpoint_type=cls._imp_platform, + description="shared integration destination endpoint", + connection_info=cls._imp_conn_info, + regions=[], + ) + + atexit.register( + cls._ignoreExc(lambda: cls._client.endpoints.delete(endpoint.id)) + ) + _SHARED_ENDPOINTS[platform] = endpoint + + return endpoint + + @classmethod + def _get_shared_pool(cls, platform): + """Lazily create and cache (once per run) the default pool. + + *platform* is constants.PROVIDER_PLATFORM_SOURCE or _DESTINATION. Bound to the + shared endpoint for that platform. Torn down via atexit, not addClassCleanup. + + Tests that need different pool sizing, or that exercise the pool mechanism + disruptively, must create their own dedicated pool instead. + """ + pool_id = _SHARED_POOL_IDS.get(platform) + if pool_id is not None: + return pool_id + + is_dst = platform == constants.PROVIDER_PLATFORM_DESTINATION + endpoint = cls._get_shared_endpoint(platform) + pool = cls._create_pool( + endpoint.id, + "shared-dst-transfer-pool" if is_dst else "shared-src-transfer-pool", + skip_allocation=False, + wait_for_allocation=True, + platform=platform, + class_cleanup=False, + ) + atexit.register(cls._ignoreExc(lambda: cls._safe_delete_pool(pool.id))) + _SHARED_POOL_IDS[platform] = pool.id + + return pool.id + + @classmethod + def _src_minion_pool_supported(cls): + """Whether the export provider advertises source pool support.""" + available = providers_factory.get_available_providers() + exp_types = available.get(cls._exp_platform, {}).get("types", []) + return constants.PROVIDER_TYPE_SOURCE_MINION_POOL in exp_types + + @classmethod + def _dst_minion_pool_supported(cls): + """Whether the import provider advertises destination pool support.""" + available = providers_factory.get_available_providers() + imp_types = available.get(cls._imp_platform, {}).get("types", []) + return constants.PROVIDER_TYPE_DESTINATION_MINION_POOL in imp_types + + @classmethod + def _get_src_endpoint(cls): + return cls._get_shared_endpoint(constants.PROVIDER_PLATFORM_SOURCE) + + @classmethod + def _get_dst_endpoint(cls): + return cls._get_shared_endpoint(constants.PROVIDER_PLATFORM_DESTINATION) + + @classmethod + def _get_src_pool_id(cls): + return cls._get_shared_pool(constants.PROVIDER_PLATFORM_SOURCE) + + @classmethod + def _get_dst_pool_id(cls): + return cls._get_shared_pool(constants.PROVIDER_PLATFORM_DESTINATION) + class ReplicaIntegrationTestBase(CoriolisIntegrationTestBase): - _CREATE_DST_MINION_POOL = False - _CREATE_SRC_MINION_POOL = False + # Minion pools are used by default if the provider supports them, falling back to + # temporary workers if not. Tests that specifically exercise the temporary workers + # (e.g.: transfer-specific tests) must set these to False. + _CREATE_DST_MINION_POOL = True + _CREATE_SRC_MINION_POOL = True + + # Whether the OS Morphing phase of a deployment uses the destination pool + # (when there is one) instead of a temporary OS Morphing minion. Tests without + # a destination pool always use a temporary minion. + _USE_MINION_POOL_FOR_OSMORPHING = True _SRC_DEVICE_SIZE_MB = 16 # Extra source_environment entries merged into the default transfer's # 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() - cls._src_endpoint = cls._create_endpoint( - name="test-src", - endpoint_type=cls._exp_platform, - description="integration source endpoint", - connection_info=cls._exp_conn_info, - ) + cls._src_endpoint = cls._get_src_endpoint() + cls._dst_endpoint = cls._get_dst_endpoint() - cls._dst_endpoint = cls._create_endpoint( - name="test-dest", - endpoint_type=cls._imp_platform, - description="integration destination endpoint", - connection_info=cls._imp_conn_info, - ) - - # Create minion pool if needed. + # Create minion pools whenever the provider supports them and the subclass did + # not opt out, otherwise transfers will use temporary workers. cls._dst_pool_id = None - if cls._CREATE_DST_MINION_POOL: - pool = cls._create_pool( - cls._dst_endpoint.id, - "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 + if cls._CREATE_DST_MINION_POOL and cls._dst_minion_pool_supported(): + cls._dst_pool_id = cls._get_dst_pool_id() - # Create source minion pool if needed. cls._src_pool_id = None - if cls._CREATE_SRC_MINION_POOL: - pool = cls._create_pool( - cls._src_endpoint.id, - "src-transfer-pool", - 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 + if cls._CREATE_SRC_MINION_POOL and cls._src_minion_pool_supported(): + cls._src_pool_id = cls._get_src_pool_id() def setUp(self): super().setUp() @@ -460,12 +560,29 @@ def _execute_concurrently_and_wait(self, transfer_ids, timeout=600): for transfer_id in transfer_ids ] for execution in executions: - self.assertExecutionCompleted(execution.id, timeout=timeout) + self.assertExecutionCompleted( + execution.id, timeout=timeout, check_pools=False + ) + + self._assertMinionPoolsHealthy() def _execute_transfer_and_deployment(self, deployment_kwargs=None): - deployment_kwargs = deployment_kwargs or {} + deployment_kwargs = dict(deployment_kwargs or {}) self._execute_and_wait(self._transfer.id) + + # Reuse the transfer's destination pool for the OS Morphing phase too, unless + # the test class opts out. Callers that need a specific OS Morphing pool can + # pass their own mapping to override this. + if ( + self._USE_MINION_POOL_FOR_OSMORPHING + and self._dst_pool_id + and "instance_osmorphing_minion_pool_mappings" not in deployment_kwargs + ): + deployment_kwargs["instance_osmorphing_minion_pool_mappings"] = { + self._instance_name: self._dst_pool_id, + } + deployment = self._client.deployments.create_from_transfer( self._transfer.id, skip_os_morphing=False, @@ -523,8 +640,28 @@ def wait_for_execution(self, execution_id, timeout=600, desired_statuses=None): % (execution_id, desired_statuses, timeout, execution.status) ) - def assertExecutionCompleted(self, execution_id, timeout=600): - """Assert that *execution_id* completes successfully.""" + def _assertMinionPoolsHealthy(self): + """Assert that any pool(s) used by this test are still usable. + + No-op for pools that weren't created (e.g.: provider doesn't support them, or + the subclass opted out via _CREATE_DST_MINION_POOL / _CREATE_SRC_MINION_POOL). + """ + if self._dst_pool_id: + self.assertPoolAllocated(self._dst_pool_id) + self.assertMachinesAvailable(self._dst_pool_id) + + if self._src_pool_id: + self.assertPoolAllocated(self._src_pool_id) + self.assertMachinesAvailable(self._src_pool_id) + + def assertExecutionCompleted(self, execution_id, timeout=600, check_pools=True): + """Assert that *execution_id* completes successfully. + + :param check_pools: also assert that the minion pools are healthy. Callers + waiting on several concurrent executions should disable this and call + _assertMinionPoolsHealthy() once all of them completed, since machines + of running executions are still in use. + """ execution = self.wait_for_execution(execution_id, timeout=timeout) self.assertEqual( constants.EXECUTION_STATUS_COMPLETED, @@ -545,6 +682,9 @@ def assertExecutionCompleted(self, execution_id, timeout=600): ), ) + if check_pools: + self._assertMinionPoolsHealthy() + def assertExecutionErrored(self, execution_id, timeout=600): """Assert that *execution_id* ends in an error state.""" execution = self.wait_for_execution(execution_id, timeout=timeout) @@ -645,6 +785,7 @@ def assertDeploymentCompleted(self, deployment_id, timeout=600): "Deployment %s ended with status %s" % (deployment_id, deployment.last_execution_status), ) + self._assertMinionPoolsHealthy() def assertDeploymentErrored(self, deployment_id, timeout=600): """Assert that *deployment_id* ends in an error state.""" @@ -683,11 +824,16 @@ class SourceMinionPoolTestBase(CoriolisIntegrationTestBase): """Base class for source minion pool integration tests. Skips the entire test class when the export provider does not advertise - ``PROVIDER_TYPE_SOURCE_MINION_POOL`` support. + ``PROVIDER_TYPE_SOURCE_MINION_POOL`` support. Use for tests that *need* + source pools (a hard skip, rather than the silent fallback to temporary + workers that ReplicaIntegrationTestBase's opt-in pool usage gives you). """ @classmethod def setUpClass(cls): + # Check before super(), so that ReplicaIntegrationTestBase.setUpClass + # does not attempt pool creation against a provider that doesn't + # support it. h = harness._IntegrationHarness.get() available = providers_factory.get_available_providers() exp_types = available.get(h.exp_provider_platform, {}).get("types", []) @@ -704,7 +850,9 @@ class DestinationMinionPoolTestBase(CoriolisIntegrationTestBase): """Base class for minion pool integration tests. Skips the entire test class when the import provider does not advertise - ``PROVIDER_TYPE_DESTINATION_MINION_POOL`` support. + ``PROVIDER_TYPE_DESTINATION_MINION_POOL`` support. Use for tests that *need* + destination pools (a hard skip, rather than the silent fallback to + temporary workers that ReplicaIntegrationTestBase's opt-in pool usage gives you). """ @classmethod @@ -724,50 +872,63 @@ def setUpClass(cls): super().setUpClass() -class MinionPoolReplicaTestBase( - DestinationMinionPoolTestBase, ReplicaIntegrationTestBase -): - """Base class for replica integration tests using destination minion pools. +class AnyMinionPoolMixin: + """Skips the test class if neither the source nor the destination has minion pools. - Extends the assertions to also verify that the minions in the pool have - been used, and that the minions and the pool returns to an available state. + For tests that use minion pools on each side that supports them. Must come before + ReplicaIntegrationTestBase in the bases list. """ - _CREATE_DST_MINION_POOL = True + @classmethod + def setUpClass(cls): + super().setUpClass() - def _execute_and_wait(self, transfer_id, timeout=600): - super()._execute_and_wait(transfer_id, timeout=timeout) - self.assertPoolAllocated(self._dst_pool_id) - self.assertMachinesAvailable(self._dst_pool_id) + if not (cls._dst_pool_id or cls._src_pool_id): + raise unittest.SkipTest( + "Neither the source nor the destination provider supports minion pools" + ) - def assertExecutionCompleted(self, execution_id, timeout=600): - 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) - self.assertMachinesAvailable(self._dst_pool_id) +class DedicatedMinionPoolsMixin: + """Gives a ReplicaIntegrationTestBase subclass its own pool(s). + For tests that are disruptive to the pool or need a sizing other than the + shared pool's configurations. Must come before the test base in the bases list. + """ -class SourceMinionPoolReplicaTestBase( - SourceMinionPoolTestBase, ReplicaIntegrationTestBase -): - """Base class for replica integration tests using source minion pools. + # Overridable params for the dedicated pool(s) created by _create_dedicated_pool, + # e.g.: to exercise pool machine power-cycling via small idle time and the + # "poweroff" strategy. + _POOL_MINIMUM_MINIONS = _DEFAULT_POOL_MINIMUM_MINIONS + _POOL_MAXIMUM_MINIONS = _DEFAULT_POOL_MAXIMUM_MINIONS + _POOL_MINION_MAX_IDLE_TIME = _DEFAULT_POOL_MINION_MAX_IDLE_TIME + _POOL_MINION_RETENTION_STRATEGY = _DEFAULT_POOL_MINION_RETENTION_STRATEGY - Extends the assertions to also verify that the minions in the pool have - been used, and that the minions and the pool returns to an available state. - """ + @classmethod + def _create_dedicated_pool(cls, platform): + """Create a class-scoped pool on the class's endpoint of *platform*. - _CREATE_SRC_MINION_POOL = True + Sized per the _POOL_* configurations. + """ + is_dst = platform == constants.PROVIDER_PLATFORM_DESTINATION + endpoint = cls._get_shared_endpoint(platform) + pool = cls._create_pool( + endpoint.id, + "dst-transfer-pool" if is_dst else "src-transfer-pool", + skip_allocation=False, + wait_for_allocation=True, + platform=platform, + 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, + ) + return pool.id - def _execute_and_wait(self, transfer_id, timeout=600): - super()._execute_and_wait(transfer_id, timeout=timeout) - self.assertPoolAllocated(self._src_pool_id) - self.assertMachinesAvailable(self._src_pool_id) - - 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) + @classmethod + def _get_src_pool_id(cls): + return cls._create_dedicated_pool(constants.PROVIDER_PLATFORM_SOURCE) + + @classmethod + def _get_dst_pool_id(cls): + return cls._create_dedicated_pool(constants.PROVIDER_PLATFORM_DESTINATION) diff --git a/coriolis/tests/integration/deployments/test_deployment.py b/coriolis/tests/integration/deployments/test_deployment.py index c2b10615a..d50a080ab 100644 --- a/coriolis/tests/integration/deployments/test_deployment.py +++ b/coriolis/tests/integration/deployments/test_deployment.py @@ -85,9 +85,3 @@ def test_cancel_deployment(self): constants.EXECUTION_STATUS_CANCELED_FOR_DEBUGGING, ], ) - - -class MinionPoolReplicaDeploymentTests( - base.MinionPoolReplicaTestBase, ReplicaDeploymentIntegrationTest -): - """Replica deployment that uses a pre-allocated destination minion pool.""" diff --git a/coriolis/tests/integration/deployments/test_luks_osmorphing.py b/coriolis/tests/integration/deployments/test_luks_osmorphing.py index 4a20fe018..24ee23e7d 100644 --- a/coriolis/tests/integration/deployments/test_luks_osmorphing.py +++ b/coriolis/tests/integration/deployments/test_luks_osmorphing.py @@ -50,6 +50,9 @@ class _LUKSOSMorphingMixin: _SRC_DEVICE_SIZE_MB = 512 _CONTAINER_IMAGE = "ubuntu:24.04" + # Exercises the temporary OS Morphing minions. + _CREATE_DST_MINION_POOL = False + @classmethod def setUpClass(cls): harness = integration_harness._IntegrationHarness.get() @@ -214,11 +217,21 @@ def _assert_firstboot_setup(self): ) +class LUKSOSMorphingMinionPoolDeploymentTest( + integration_base.DestinationMinionPoolTestBase, LUKSOSMorphingDeploymentTest +): + """Same as LUKSOSMorphingDeploymentTest, OS Morphing in a pool minion.""" + + _CREATE_DST_MINION_POOL = True + + class LUKSRockyLinuxOSMorphingDeploymentTest( _LUKSOSMorphingMixin, integration_base.ReplicaIntegrationTestBase ): """LUKS + dracut OS morphing test using Rocky Linux 9.""" + _CREATE_DST_MINION_POOL = True + # kernel-core (~150 MB installed) needs extra room on top of the base # container image and the other morphing packages. _SRC_DEVICE_SIZE_MB = 777 diff --git a/coriolis/tests/integration/deployments/test_osmorphing.py b/coriolis/tests/integration/deployments/test_osmorphing.py index e533401a8..08538633f 100644 --- a/coriolis/tests/integration/deployments/test_osmorphing.py +++ b/coriolis/tests/integration/deployments/test_osmorphing.py @@ -40,6 +40,9 @@ def setUp(self): class OsMorphingDeploymentTest(OsMorphingDeploymentTestBase): + # Exercises the temporary workers. + _CREATE_DST_MINION_POOL = False + def test_deployment_with_os_morphing(self): self.assertFalse( osmorphing_utils.path_exists_on_device(self._src_device, "usr/bin/jq"), @@ -53,6 +56,8 @@ def test_deployment_with_os_morphing(self): "jq was not found on the destination device after OS morphing", ) + +class _OsMorphingScriptTestsMixin: def test_os_morphing_global_script_basic_format(self): expected_string = str(uuid.uuid4()) user_scripts = { @@ -194,21 +199,17 @@ def test_os_morphing_global_script_first_boot(self): class OsMorphingMinionPoolDeploymentTest( - integration_base.DestinationMinionPoolTestBase, OsMorphingDeploymentTestBase + integration_base.DestinationMinionPoolTestBase, + _OsMorphingScriptTestsMixin, + OsMorphingDeploymentTestBase, ): - """OS morphing deployment using a minion pool for the OS morphing phase.""" + """OS morphing deployment using the shared destination pool. - @classmethod - def setUpClass(cls): - super().setUpClass() + OsMorphingDeploymentTest tests cover the temporary OS Morphing minion, while these + ones covers OS Morphing in a (reused) pool minion, including the user scripts. + """ - pool = cls._create_pool( - cls._dst_endpoint.id, - "osmorph-pool", - skip_allocation=False, - wait_for_allocation=True, - ) - cls._osmorph_pool_id = pool.id + _CREATE_DST_MINION_POOL = True def test_deployment_with_os_morphing(self): self.assertFalse( @@ -216,12 +217,7 @@ def test_deployment_with_os_morphing(self): "jq was found on the source device before OS morphing", ) - deployment_kwargs = { - "instance_osmorphing_minion_pool_mappings": { - self._instance_name: self._osmorph_pool_id, - }, - } - self._execute_transfer_and_deployment(deployment_kwargs) + self._execute_transfer_and_deployment() self.assertTrue( osmorphing_utils.path_exists_on_device(self._dst_device, "usr/bin/jq"), @@ -229,9 +225,7 @@ def test_deployment_with_os_morphing(self): ) ctxt = self._get_db_context() - pool = db_api.get_minion_pool( - ctxt, self._osmorph_pool_id, include_machines=True - ) + pool = db_api.get_minion_pool(ctxt, self._dst_pool_id, include_machines=True) self.assertTrue(pool.minion_machines, "OS morphing pool has no minion machines") for machine in pool.minion_machines: @@ -240,6 +234,28 @@ def test_deployment_with_os_morphing(self): "OS morphing minion machine %s was never used" % machine.id, ) + +class OsMorphingMinionPoolAllocationFailureTest( + integration_base.DestinationMinionPoolTestBase, OsMorphingDeploymentTestBase +): + """OS morphing minion pool allocation failure test. + + Deliberately breaks the pool's only machine, so it needs its own dedicated pool, + rather than the shared one. + """ + + @classmethod + def setUpClass(cls): + super().setUpClass() + + pool = cls._create_pool( + cls._dst_endpoint.id, + "osmorph-pool", + skip_allocation=False, + wait_for_allocation=True, + ) + cls._osmorph_pool_id = pool.id + def test_osmorphing_minion_allocation_failure_cleans_up(self): """OS morphing minion pool allocation fail test. diff --git a/coriolis/tests/integration/management/test_region.py b/coriolis/tests/integration/management/test_region.py index 6c26a3f2d..93ae73784 100644 --- a/coriolis/tests/integration/management/test_region.py +++ b/coriolis/tests/integration/management/test_region.py @@ -50,6 +50,33 @@ def test_region_crud(self): class RegionSchedulingTests(base.ReplicaIntegrationTestBase): + """Tests scheduling for endpoints with mapped regions. + + Mutates the endpoints' mapped_regions, so it can't use the shared ones. Since it's + a new endpoint, we can't reuse the existing minion pool. + """ + + _CREATE_DST_MINION_POOL = False + _CREATE_SRC_MINION_POOL = False + + @classmethod + def _get_src_endpoint(cls): + return cls._create_endpoint( + name="test-src", + endpoint_type=cls._exp_platform, + description="integration source endpoint", + connection_info=cls._exp_conn_info, + ) + + @classmethod + def _get_dst_endpoint(cls): + return cls._create_endpoint( + name="test-dest", + endpoint_type=cls._imp_platform, + description="integration destination endpoint", + connection_info=cls._imp_conn_info, + ) + @classmethod def setUpClass(cls): super().setUpClass() diff --git a/coriolis/tests/integration/test_failure_recovery.py b/coriolis/tests/integration/test_failure_recovery.py index 1d7ee0f48..e558ff679 100644 --- a/coriolis/tests/integration/test_failure_recovery.py +++ b/coriolis/tests/integration/test_failure_recovery.py @@ -30,6 +30,11 @@ class TransferFailureIntegrationTest(base.ReplicaIntegrationTestBase): """Error path and resource cleanup.""" + # Patches deploy_replica_{target,source}_resources, which are not used for + # pool-backed transfers. + _CREATE_DST_MINION_POOL = False + _CREATE_SRC_MINION_POOL = False + def _assertResourcesCleaned(self, execution_id, task_type, resource_key): ctxt = self._get_db_context() execution = db_api.get_tasks_execution(ctxt, execution_id) @@ -127,8 +132,18 @@ def _slow_then_fail(self_provider, *args, **kwargs): self.assertTargetResourcesCleaned(execution.id) -class MinionPoolAllocationFailureTest(base.MinionPoolReplicaTestBase): - """Transfer minion pool allocation failure tests.""" +class MinionPoolAllocationFailureTest( + base.DedicatedMinionPoolsMixin, + base.DestinationMinionPoolTestBase, + base.ReplicaIntegrationTestBase, +): + """Transfer minion pool allocation failure tests. + + Deliberately breaks its pool's only machine, so it needs a dedicated pool + rather than the shared one. + """ + + _CREATE_SRC_MINION_POOL = False def test_transfer_minion_allocation_failure_cleans_up(self): """Transfer minion pool allocation fail test. diff --git a/coriolis/tests/integration/test_minion_pools.py b/coriolis/tests/integration/test_minion_pools.py index 13c003fcd..2722e64a2 100644 --- a/coriolis/tests/integration/test_minion_pools.py +++ b/coriolis/tests/integration/test_minion_pools.py @@ -197,7 +197,7 @@ def _create_pool(self, endpoint_id, **kwargs): ) -class _MinionPoolPowerCycleTestMixin: +class _MinionPoolPowerCycleTestMixin(base.DedicatedMinionPoolsMixin): """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 @@ -305,26 +305,34 @@ def test_transfer_after_pool_machine_power_cycle(self): class MinionPoolPowerCycleTransferTest( - _MinionPoolPowerCycleTestMixin, base.MinionPoolReplicaTestBase + _MinionPoolPowerCycleTestMixin, + base.DestinationMinionPoolTestBase, + base.ReplicaIntegrationTestBase, ): """Power-cycle test exercising a destination minion pool.""" + _CREATE_SRC_MINION_POOL = False + @property def _pool_id(self): return self._dst_pool_id class SourceMinionPoolPowerCycleTransferTest( - _MinionPoolPowerCycleTestMixin, base.SourceMinionPoolReplicaTestBase + _MinionPoolPowerCycleTestMixin, + base.SourceMinionPoolTestBase, + base.ReplicaIntegrationTestBase, ): """Power-cycle test exercising a source minion pool.""" + _CREATE_DST_MINION_POOL = False + @property def _pool_id(self): return self._src_pool_id -class _MinionPoolRefreshDeallocationTestMixin: +class _MinionPoolRefreshDeallocationTestMixin(base.DedicatedMinionPoolsMixin): """Excess pool machine gets deleted on refresh. Mirrors _MinionPoolPowerCycleTestMixin but with the default "delete" retention @@ -395,20 +403,28 @@ def test_excess_pool_machine_deleted_on_refresh(self): class MinionPoolRefreshDeallocationTransferTest( - _MinionPoolRefreshDeallocationTestMixin, base.MinionPoolReplicaTestBase + _MinionPoolRefreshDeallocationTestMixin, + base.DestinationMinionPoolTestBase, + base.ReplicaIntegrationTestBase, ): """Deletion-on-refresh test exercising a destination minion pool.""" + _CREATE_SRC_MINION_POOL = False + @property def _pool_id(self): return self._dst_pool_id class SourceMinionPoolRefreshDeallocationTransferTest( - _MinionPoolRefreshDeallocationTestMixin, base.SourceMinionPoolReplicaTestBase + _MinionPoolRefreshDeallocationTestMixin, + base.SourceMinionPoolTestBase, + base.ReplicaIntegrationTestBase, ): """Deletion-on-refresh test exercising a source minion pool.""" + _CREATE_DST_MINION_POOL = False + @property def _pool_id(self): return self._src_pool_id diff --git a/coriolis/tests/integration/test_provider/imp.py b/coriolis/tests/integration/test_provider/imp.py index fa45f3b3c..b6a9121e9 100644 --- a/coriolis/tests/integration/test_provider/imp.py +++ b/coriolis/tests/integration/test_provider/imp.py @@ -445,11 +445,18 @@ def create_minion( # # Mount the host's /lib/modules tree so that modprobe can # resolve built-in modules. + # + # A pool minion is created once and reused for arbitrary future + # transfers / deployments, so unlike deploy_os_morphing_resources (which + # only adds /dev/mapper/control when it already knows the deployment + # is LUKS-encrypted), this must always include it: luksOpen needs it, + # and there's no way to add it retroactively to an already-running + # container if a later LUKS-encrypted instance gets mapped to it. volumes = ["/lib/modules:/lib/modules:ro"] result = self._create_minion( "coriolis-pool-minion", connection_info, - [], + ["/dev/mapper/control"], volumes, device_cgroup_rules=["b *:* rwm"], ) @@ -484,5 +491,8 @@ def get_additional_os_morphing_info( "os_type": instance_deployment_info.get("os_type", "linux"), "ignore_devices": ignore_devices, "_include_loop_devices": True, + constants.ENCRYPTED_DISKS_PASS: instance_deployment_info.get( + constants.ENCRYPTED_DISKS_PASS + ), } } diff --git a/coriolis/tests/integration/transfers/test_executions.py b/coriolis/tests/integration/transfers/test_executions.py index 7bed60eab..b79db89e5 100644 --- a/coriolis/tests/integration/transfers/test_executions.py +++ b/coriolis/tests/integration/transfers/test_executions.py @@ -11,58 +11,22 @@ from coriolis.tests.integration import base -class TransferExecutionsTests(base.ReplicaIntegrationTestBase): - # Provider method to fail in test_execution_auto_deploy_transfer_failure. - # Plain transfers deploy fresh target resources via deploy_replica_target_resources, - # minion-pool-backed transfers instead attach volumes to a pre-allocated minion via - # attach_volumes_to_minion, so deploy_replica_target_resources is never called and - # would not inject any failure there. - _AUTO_DEPLOY_FAILURE_METHOD = "deploy_replica_target_resources" - - def test_executions(self): - # We didn't start the execution yet. - executions = self._client.transfer_executions.list(self._transfer.id) - self.assertIsInstance(executions, list) - self.assertEqual(0, len(executions)) - - # Start the execution. - execution = self._client.transfer_executions.create( - self._transfer.id, shutdown_instances=False - ) - self.addCleanup( - self._cleanup_execution, - self._transfer.id, - execution.id, - ) - - self.assertExecutionCompleted(execution.id) - executions = self._client.transfer_executions.list(self._transfer.id) - ids = [e.id for e in executions] - self.assertIn(execution.id, ids) - - # Get the execution. - fetched = self._client.transfer_executions.get(self._transfer.id, execution.id) - self.assertEqual(execution.id, fetched.id) - - # Delete the execution. - self._client.transfer_executions.delete(self._transfer.id, execution.id) - - executions = self._client.transfer_executions.list(self._transfer.id) - ids = [e.id for e in executions] - self.assertNotIn(execution.id, ids) - - def test_shutdown_instances(self): - # shutdown_instances=True calls provider.shutdown_instance(). - execution = self._client.transfer_executions.create( - self._transfer.id, shutdown_instances=True - ) - self.addCleanup( - self._cleanup_execution, - self._transfer.id, - execution.id, - ) - - self.assertExecutionCompleted(execution.id) +class _TransferExecutionsPathTestsMixin: + """Tests that depend on the type of minions the transfer uses.""" + + @property + def _deploy_resources_method(self): + """Provider method called by the transfers to prepare their destination. + + Plain transfers deploy fresh target resources via + deploy_replica_target_resources, minion-pool-backed transfers instead attach + volumes to a pre-allocated minion via attach_volumes_to_minion, so + deploy_replica_target_resources is never called and would not inject any + failure or delay there. + """ + if self._dst_pool_id: + return "attach_volumes_to_minion" + return "deploy_replica_target_resources" def _get_transfer_deployment(self): deployments = self._client.deployments.list() @@ -73,32 +37,6 @@ def _get_transfer_deployment(self): return transfer_deployments[0] - def test_execution_auto_deploy(self): - """auto_deploy=True hands the deployment off to the deployer manager. - - Exercises the deployer_manager -> conductor.confirm_deployer_completed handoff: - the deployer_manager service polls the PENDING deployment, notices the - underlying transfer execution (the "deployer") completed, and calls back into - the conductor to kick off the actual deployment. - """ - execution = self._client.transfer_executions.create( - self._transfer.id, - shutdown_instances=False, - auto_deploy=True, - ) - self.addCleanup( - self._cleanup_execution, - self._transfer.id, - execution.id, - ) - - self.assertExecutionCompleted(execution.id) - - deployment = self._get_transfer_deployment() - self.addCleanup(self._cleanup_deployment, deployment.id) - - self.assertDeploymentCompleted(deployment.id) - def test_execution_auto_deploy_transfer_failure(self): """A failed "deployer" transfer execution errors out the deployment. @@ -112,7 +50,7 @@ def test_execution_auto_deploy_transfer_failure(self): with mock.patch.object( self._harness.imp_provider_class, - self._AUTO_DEPLOY_FAILURE_METHOD, + self._deploy_resources_method, side_effect=injected_error, ): execution = self._client.transfer_executions.create( @@ -156,7 +94,7 @@ def _test_cancel_running_execution(self, force): # Artificially bump the execution time of a transfer. self._patch_add_delay( self._harness.imp_provider_class, - "deploy_replica_target_resources", + self._deploy_resources_method, ) execution = self._client.transfer_executions.create( @@ -189,9 +127,92 @@ def _test_cancel_running_execution(self, force): ) +class TransferExecutionsTests( + _TransferExecutionsPathTestsMixin, base.ReplicaIntegrationTestBase +): + """Transfer executions tests that use temporary workers.""" + + _CREATE_DST_MINION_POOL = False + _CREATE_SRC_MINION_POOL = False + + class MinionPoolTransferExecutionsTests( - base.MinionPoolReplicaTestBase, TransferExecutionsTests + base.AnyMinionPoolMixin, + _TransferExecutionsPathTestsMixin, + base.ReplicaIntegrationTestBase, ): - """Transfer executions that use a pre-allocated destination minion pool.""" + """Transfer executions tests that use minion pools. + + Also contains the tests which don't depend on the type of minion the transfer uses. + """ + + def test_executions(self): + # We didn't start the execution yet. + executions = self._client.transfer_executions.list(self._transfer.id) + self.assertIsInstance(executions, list) + self.assertEqual(0, len(executions)) + + # Start the execution. + execution = self._client.transfer_executions.create( + self._transfer.id, shutdown_instances=False + ) + self.addCleanup( + self._cleanup_execution, + self._transfer.id, + execution.id, + ) - _AUTO_DEPLOY_FAILURE_METHOD = "attach_volumes_to_minion" + self.assertExecutionCompleted(execution.id) + executions = self._client.transfer_executions.list(self._transfer.id) + ids = [e.id for e in executions] + self.assertIn(execution.id, ids) + + # Get the execution. + fetched = self._client.transfer_executions.get(self._transfer.id, execution.id) + self.assertEqual(execution.id, fetched.id) + + # Delete the execution. + self._client.transfer_executions.delete(self._transfer.id, execution.id) + + executions = self._client.transfer_executions.list(self._transfer.id) + ids = [e.id for e in executions] + self.assertNotIn(execution.id, ids) + + def test_shutdown_instances(self): + # shutdown_instances=True calls provider.shutdown_instance(). + execution = self._client.transfer_executions.create( + self._transfer.id, shutdown_instances=True + ) + self.addCleanup( + self._cleanup_execution, + self._transfer.id, + execution.id, + ) + + self.assertExecutionCompleted(execution.id) + + def test_execution_auto_deploy(self): + """auto_deploy=True hands the deployment off to the deployer manager. + + Exercises the deployer_manager -> conductor.confirm_deployer_completed handoff: + the deployer_manager service polls the PENDING deployment, notices the + underlying transfer execution (the "deployer") completed, and calls back into + the conductor to kick off the actual deployment. + """ + execution = self._client.transfer_executions.create( + self._transfer.id, + shutdown_instances=False, + auto_deploy=True, + ) + self.addCleanup( + self._cleanup_execution, + self._transfer.id, + execution.id, + ) + + self.assertExecutionCompleted(execution.id) + + deployment = self._get_transfer_deployment() + self.addCleanup(self._cleanup_deployment, deployment.id) + + self.assertDeploymentCompleted(deployment.id) diff --git a/coriolis/tests/integration/transfers/test_transfer.py b/coriolis/tests/integration/transfers/test_transfer.py index 1eb1dfa28..1feed4cab 100644 --- a/coriolis/tests/integration/transfers/test_transfer.py +++ b/coriolis/tests/integration/transfers/test_transfer.py @@ -150,6 +150,10 @@ class ReplicaTransferIntegrationTest( ): """Full-pipeline replica transfer integration tests.""" + # Exercises the temporary workers code path. + _CREATE_DST_MINION_POOL = False + _CREATE_SRC_MINION_POOL = False + def test_transfer_with_ssh_backup_writer(self): # NOTE: for minion pools, updating the data_transfer_mechanism will not # set up the new transfer mechanism into existing minions. @@ -241,6 +245,10 @@ class ClusteredTransferIntegrationTest(base.ReplicaIntegrationTestBase): them (same disk id in both instances' export_info). """ + # Exercises the temporary workers code path. + _CREATE_DST_MINION_POOL = False + _CREATE_SRC_MINION_POOL = False + @classmethod def setUpClass(cls): h = integration_harness._IntegrationHarness.get() @@ -430,20 +438,24 @@ def _fail_for_instance_a( class MinionPoolTransferTest( - base.MinionPoolReplicaTestBase, _ReplicaTransferTestsMixin + base.AnyMinionPoolMixin, + _ReplicaTransferTestsMixin, + base.ReplicaIntegrationTestBase, ): - """Transfer execution that uses a pre-allocated destination minion pool.""" + """Transfer execution that uses pre-allocated source and destination minion pools. - def test_transfer(self): - super().test_transfer() - self.assertPoolAllocated(self._dst_pool_id) - self.assertMachinesAvailable(self._dst_pool_id) + The pools are used on each side that supports them; the pools' health is checked + by assertExecutionCompleted. + """ class ReplicaTransferViaSSHTunnelTest(base.ReplicaIntegrationTestBase): """Transfer tests using an SSH tunneled replicator client.""" _EXTRA_SOURCE_ENVIRONMENT = {"use_tunnel": True} + # Exercises the temporary workers code path. + _CREATE_DST_MINION_POOL = False + _CREATE_SRC_MINION_POOL = False @classmethod def setUpClass(cls): @@ -480,14 +492,3 @@ def _spy_get_ssh_tunnel(client_self): test_utils.devices_match(self._src_device, self._dst_device), "Devices do not match after transfer via SSH tunnel", ) - - -class SourceMinionPoolTransferTest( - base.SourceMinionPoolReplicaTestBase, _ReplicaTransferTestsMixin -): - """Transfer execution that uses a pre-allocated source minion pool.""" - - def test_transfer(self): - super().test_transfer() - self.assertPoolAllocated(self._src_pool_id) - self.assertMachinesAvailable(self._src_pool_id)