From b5441ea5a53659904e6ad7c353369503fa7f1c26 Mon Sep 17 00:00:00 2001 From: Seyeong Kim Date: Tue, 11 Aug 2026 16:35:40 +0000 Subject: [PATCH 1/7] fix(maas): retain machine FQDN Preserve the FQDN returned by MAAS so downstream removal steps can distinguish it from the short cluster node name. Signed-off-by: Seyeong Kim --- .../sunbeam/provider/maas/client.py | 1 + .../unit/sunbeam/provider/maas/test_maas.py | 20 +++++++++++++++++++ 2 files changed, 21 insertions(+) diff --git a/sunbeam-python/sunbeam/provider/maas/client.py b/sunbeam-python/sunbeam/provider/maas/client.py index 824fc8300..75e138642 100644 --- a/sunbeam-python/sunbeam/provider/maas/client.py +++ b/sunbeam-python/sunbeam/provider/maas/client.py @@ -341,6 +341,7 @@ def _convert_raw_machine(machine_raw: dict, root_disk: dict | None) -> dict: machine = { "system_id": machine_raw["system_id"], "hostname": machine_raw["hostname"], + "fqdn": machine_raw["fqdn"], "roles": list(set(tag_names).intersection(RoleTags.values())), "zone": machine_raw["zone"]["name"], "status": machine_raw["status_name"], diff --git a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py index fbdcc42b8..336b0a62b 100644 --- a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py +++ b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py @@ -17,6 +17,7 @@ from sunbeam.core.deployment import Networks from sunbeam.core.deployments import DeploymentsConfig from sunbeam.core.juju import ControllerNotFoundException +from sunbeam.provider.maas.client import _convert_raw_machine from sunbeam.provider.maas.commands import ( configure_cmd, remove_node, @@ -64,6 +65,25 @@ ) +class TestConvertRawMachine: + def test_preserves_fqdn(self): + machine_raw = { + "system_id": "sysid", + "hostname": "cloud-4", + "fqdn": "cloud-4.maas", + "blockdevice_set": [], + "interface_set": [], + "zone": {"name": "default"}, + "status_name": "Ready", + "cpu_count": 4, + "memory": 8192, + } + + machine = _convert_raw_machine(machine_raw, None) + + assert machine["fqdn"] == "cloud-4.maas" + + class TestMaasConfigureCommand: def test_network_agents_include_all_microovn_nodes( self, From 3f890f46b76d60f15abfe525c11558b20747e39f Mon Sep 17 00:00:00 2001 From: Seyeong Kim Date: Wed, 12 Aug 2026 12:34:29 +0000 Subject: [PATCH 2/7] fix(maas): remove stale OpenStack references Remove Nova services and Neutron agents by the MAAS FQDN after the target MicroOVN unit is gone. Verify that the control-plane records no longer exist so retrying node removal cannot silently leave stale state. Signed-off-by: Seyeong Kim --- .../sunbeam/provider/maas/commands.py | 13 +- sunbeam-python/sunbeam/steps/hypervisor.py | 93 ++++++- .../unit/sunbeam/provider/maas/test_maas.py | 64 ++++- .../unit/sunbeam/steps/test_hypervisor.py | 240 +++++++++++++++++- 4 files changed, 404 insertions(+), 6 deletions(-) diff --git a/sunbeam-python/sunbeam/provider/maas/commands.py b/sunbeam-python/sunbeam/provider/maas/commands.py index b456051fa..2d16e00f0 100644 --- a/sunbeam-python/sunbeam/provider/maas/commands.py +++ b/sunbeam-python/sunbeam/provider/maas/commands.py @@ -137,6 +137,7 @@ DeployHypervisorApplicationStep, DestroyHypervisorApplicationStep, ReapplyHypervisorTerraformPlanStep, + RemoveHypervisorReferencesStep, RemoveHypervisorUnitStep, ) from sunbeam.steps.juju import ( @@ -1693,6 +1694,9 @@ def remove_node(ctx: click.Context, name: str, force: bool, show_hints: bool) -> run_plan(check_plan, console, show_hints) + maas_client = MaasClient.from_deployment(deployment) + machine = get_machine(maas_client, name) + plan = [ MigrateK8SKubeconfigStep( client, name, jhelper, deployment.openstack_machines_model @@ -1701,7 +1705,7 @@ def remove_node(ctx: click.Context, name: str, force: bool, show_hints: bool) -> RemoveHypervisorUnitStep( client, jhelper, - deployment, + None, name, deployment.openstack_machines_model, force, @@ -1715,6 +1719,13 @@ def remove_node(ctx: click.Context, name: str, force: bool, show_hints: bool) -> RemoveMicroOVNUnitsStep( client, name, jhelper, deployment.openstack_machines_model ), + RemoveHypervisorReferencesStep( + jhelper, + deployment, + machine["hostname"], + machine["fqdn"], + force=force, + ), CordonK8SUnitStep(client, name, jhelper, deployment.openstack_machines_model), DrainK8SUnitStep( client, name, jhelper, deployment.openstack_machines_model, remove_pvc=True diff --git a/sunbeam-python/sunbeam/steps/hypervisor.py b/sunbeam-python/sunbeam/steps/hypervisor.py index 45cae9b2e..6d71a115b 100644 --- a/sunbeam-python/sunbeam/steps/hypervisor.py +++ b/sunbeam-python/sunbeam/steps/hypervisor.py @@ -5,6 +5,7 @@ import logging import typing +import click import tenacity from sunbeam.clusterd.client import Client @@ -32,9 +33,15 @@ ApplicationNotFoundException, JujuHelper, JujuStepHelper, + ModelNotFoundException, ) from sunbeam.core.manifest import Manifest -from sunbeam.core.openstack_api import remove_hypervisor +from sunbeam.core.openstack_api import ( + get_admin_connection, + remove_compute_service, + remove_hypervisor, + remove_network_service, +) from sunbeam.core.steps import ( DeployMachineApplicationStep, DestroyMachineApplicationStep, @@ -49,7 +56,9 @@ if typing.TYPE_CHECKING: import openstack + from keystoneauth1 import exceptions as keystoneauth_exceptions else: + keystoneauth_exceptions = LazyImport("keystoneauth1.exceptions") openstack = LazyImport("openstack") LOG = logging.getLogger(__name__) @@ -60,6 +69,10 @@ HYPERVISOR_UNIT_TIMEOUT = ( 1800 # 30 minutes, adding / removing units can take a long time ) +# Deleted Nova services and Neutron agents are normally gone on the first +# check. Poll every 10 seconds for up to 5 minutes before failing removal. +HYPERVISOR_REFERENCES_TIMEOUT = 300 +HYPERVISOR_REFERENCES_POLL_INTERVAL = 10 class DeployHypervisorApplicationStep(DeployMachineApplicationStep): @@ -330,6 +343,84 @@ def run(self, context: StepContext) -> Result: return Result(ResultType.COMPLETED) +class RemoveHypervisorReferencesStep(BaseStep): + """Remove Nova and Neutron references to a hypervisor.""" + + def __init__( + self, + jhelper: JujuHelper, + deployment: Deployment, + hostname: str, + fqdn: str, + force: bool = False, + ): + super().__init__( + "Remove openstack-hypervisor references", + "Remove openstack-hypervisor references from the control plane", + ) + self.jhelper = jhelper + self.deployment = deployment + self.force = force + self._hostnames = tuple(dict.fromkeys((hostname, fqdn))) + + def _remove_references(self) -> None: + """Remove references and raise while records remain.""" + conn = get_admin_connection(self.jhelper, self.deployment) + for hostname in self._hostnames: + remove_compute_service(hostname, conn) + remove_network_service(hostname, conn) + remaining_hosts = [] + for hostname in self._hostnames: + compute_services = list(conn.compute.services(host=hostname)) + network_agents = list(conn.network.agents(host=hostname)) + if compute_services or network_agents: + remaining_hosts.append(hostname) + if remaining_hosts: + raise tenacity.TryAgain( + f"Hypervisor references remain for {', '.join(remaining_hosts)}" + ) + + def run(self, context: StepContext) -> Result: + """Remove references until Nova and Neutron report none remain.""" + client_exceptions = ( + click.ClickException, + openstack.exceptions.SDKException, + keystoneauth_exceptions.ClientException, + ) + retry_exceptions: tuple[type[BaseException], ...] = (tenacity.TryAgain,) + if not self.force: + retry_exceptions += client_exceptions + try: + for attempt in tenacity.Retrying( + stop=tenacity.stop_after_delay(HYPERVISOR_REFERENCES_TIMEOUT), + wait=tenacity.wait_fixed(HYPERVISOR_REFERENCES_POLL_INTERVAL), + retry=tenacity.retry_if_exception_type(retry_exceptions), + reraise=True, + ): + with attempt: + self._remove_references() + except (ModelNotFoundException, ApplicationNotFoundException) as e: + LOG.debug("OpenStack control plane not deployed, skipping: %r", e) + return Result(ResultType.SKIPPED) + except client_exceptions as e: + LOG.error("Control plane error while removing hypervisor references") + if self.force: + LOG.warning( + "Force mode set, skipping hypervisor reference cleanup", + exc_info=True, + ) + return Result(ResultType.SKIPPED, str(e)) + return Result(ResultType.FAILED, str(e)) + except tenacity.TryAgain as e: + LOG.error( + "Hypervisor references still present after %ss", + HYPERVISOR_REFERENCES_TIMEOUT, + ) + return Result(ResultType.FAILED, str(e)) + + return Result(ResultType.COMPLETED) + + class ReapplyHypervisorTerraformPlanStep(BaseStep): """Reapply openstack-hyervisor terraform plan.""" diff --git a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py index 336b0a62b..99763b18b 100644 --- a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py +++ b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py @@ -57,8 +57,15 @@ ZoneBalanceCheck, ZonesCheck, ) +from sunbeam.steps.hypervisor import ( + RemoveHypervisorReferencesStep, + RemoveHypervisorUnitStep, +) from sunbeam.steps.juju import RemoveJujuMachineStep -from sunbeam.steps.microovn import ReapplyMicroOVNTerraformPlanStep +from sunbeam.steps.microovn import ( + ReapplyMicroOVNTerraformPlanStep, + RemoveMicroOVNUnitsStep, +) from sunbeam.steps.role_distributor import ( ReapplyRoleDistributorApplicationStep, RemoveRoleDistributorUnitsStep, @@ -2505,6 +2512,60 @@ def test_storage_ippool_label_with_different_names(self): class TestRemoveNodeRoleDistributor: + @pytest.fixture(autouse=True) + def mock_maas_machine(self, mocker): + maas_client = mocker.patch( + "sunbeam.provider.maas.commands.MaasClient.from_deployment" + ).return_value + get_machine = mocker.patch( + "sunbeam.provider.maas.commands.get_machine", + return_value={"hostname": "node-1", "fqdn": "node-1.maas"}, + ) + return maas_client, get_machine + + @patch("sunbeam.provider.maas.commands.JujuHelper") + @patch("sunbeam.provider.maas.commands.run_preflight_checks") + @patch("sunbeam.provider.maas.commands.run_plan") + def test_remove_orders_cleanup_steps( + self, + run_plan_cmd, + run_preflight, + juju_helper, + mock_maas_machine, + ): + maas_client, get_machine = mock_maas_machine + deployment = Mock() + deployment.openstack_machines_model = "openstack-machines" + deployment.get_ovn_manager.return_value.get_machines.return_value = [] + + result = CliRunner().invoke(remove_node, ["node-1"], obj=deployment) + + assert result.exit_code == 0, result.output + get_machine.assert_called_once_with(maas_client, "node-1") + plan = run_plan_cmd.call_args_list[1][0][0] + hypervisor_unit_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, RemoveHypervisorUnitStep) + ) + microovn_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, RemoveMicroOVNUnitsStep) + ) + references_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, RemoveHypervisorReferencesStep) + ) + hypervisor_unit = plan[hypervisor_unit_index] + references = plan[references_index] + + assert hypervisor_unit.deployment is None + assert references._hostnames == ("node-1", "node-1.maas") + assert references.force is False + assert microovn_index < references_index + @patch("sunbeam.provider.maas.commands.JujuHelper") @patch("sunbeam.provider.maas.commands.run_preflight_checks") @patch("sunbeam.provider.maas.commands.run_plan") @@ -2613,7 +2674,6 @@ def test_remove_skips_role_distributor_when_microovn_has_no_machines( isinstance(step, ReapplyMicroOVNTerraformPlanStep) for step in plan ) - class TestIsMaasDeployment: """Test is_maas_deployment TypeGuard function.""" diff --git a/sunbeam-python/tests/unit/sunbeam/steps/test_hypervisor.py b/sunbeam-python/tests/unit/sunbeam/steps/test_hypervisor.py index ae06ce217..505db400c 100644 --- a/sunbeam-python/tests/unit/sunbeam/steps/test_hypervisor.py +++ b/sunbeam-python/tests/unit/sunbeam/steps/test_hypervisor.py @@ -2,17 +2,22 @@ # SPDX-License-Identifier: Apache-2.0 import json -from unittest.mock import Mock, patch +from unittest.mock import Mock, call, patch +import click import pytest +from keystoneauth1.exceptions.catalog import EndpointNotFound +from keystoneauth1.exceptions.connection import ConnectFailure +from openstack.exceptions import SDKException from sunbeam.clusterd.service import NodeNotExistInClusterException from sunbeam.core.common import ResultType -from sunbeam.core.juju import ApplicationNotFoundException +from sunbeam.core.juju import ApplicationNotFoundException, ModelNotFoundException from sunbeam.core.terraform import TerraformException from sunbeam.steps.hypervisor import ( ReapplyHypervisorOptionalIntegrationsStep, ReapplyHypervisorTerraformPlanStep, + RemoveHypervisorReferencesStep, RemoveHypervisorUnitStep, ) @@ -311,6 +316,237 @@ def test_run_timeout( assert result.message == "timed out" +class TestRemoveHypervisorReferencesStep: + @patch("sunbeam.steps.hypervisor.get_admin_connection") + @patch("sunbeam.steps.hypervisor.remove_network_service") + @patch("sunbeam.steps.hypervisor.remove_compute_service") + def test_run_removes_short_and_fqdn_references( + self, + remove_compute_service, + remove_network_service, + get_admin_connection, + basic_jhelper, + basic_deployment, + step_context, + ): + conn = Mock() + conn.compute.services.return_value = [] + conn.network.agents.return_value = [] + get_admin_connection.return_value = conn + step = RemoveHypervisorReferencesStep( + basic_jhelper, + basic_deployment, + "cloud-4", + "cloud-4.maas", + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.COMPLETED + assert remove_compute_service.call_args_list == [ + call("cloud-4", conn), + call("cloud-4.maas", conn), + ] + assert remove_network_service.call_args_list == [ + call("cloud-4", conn), + call("cloud-4.maas", conn), + ] + assert conn.compute.services.call_args_list == [ + call(host="cloud-4"), + call(host="cloud-4.maas"), + ] + assert conn.network.agents.call_args_list == [ + call(host="cloud-4"), + call(host="cloud-4.maas"), + ] + + @patch("sunbeam.steps.hypervisor.HYPERVISOR_REFERENCES_POLL_INTERVAL", 0) + @patch("sunbeam.steps.hypervisor.get_admin_connection") + @patch("sunbeam.steps.hypervisor.remove_network_service") + @patch("sunbeam.steps.hypervisor.remove_compute_service") + def test_run_retries_until_references_are_gone( + self, + remove_compute_service, + remove_network_service, + get_admin_connection, + basic_jhelper, + basic_deployment, + step_context, + ): + conn = Mock() + conn.compute.services.side_effect = [[Mock()], [], [], []] + conn.network.agents.return_value = [] + get_admin_connection.return_value = conn + step = RemoveHypervisorReferencesStep( + basic_jhelper, + basic_deployment, + "cloud-4", + "cloud-4.maas", + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.COMPLETED + assert remove_compute_service.call_count == 4 + assert remove_network_service.call_count == 4 + + @patch("sunbeam.steps.hypervisor.get_admin_connection") + @patch("sunbeam.steps.hypervisor.remove_network_service") + @patch("sunbeam.steps.hypervisor.remove_compute_service") + def test_run_deduplicates_equal_hostnames( + self, + remove_compute_service, + remove_network_service, + get_admin_connection, + basic_jhelper, + basic_deployment, + step_context, + ): + conn = Mock() + conn.compute.services.return_value = [] + conn.network.agents.return_value = [] + get_admin_connection.return_value = conn + step = RemoveHypervisorReferencesStep( + basic_jhelper, + basic_deployment, + "cloud-4", + "cloud-4", + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.COMPLETED + remove_compute_service.assert_called_once_with("cloud-4", conn) + remove_network_service.assert_called_once_with("cloud-4", conn) + + @pytest.mark.parametrize( + "error", + [ + SDKException("control plane unavailable"), + EndpointNotFound("control plane unavailable"), + click.ClickException("Unable to get keystone leader"), + ], + ) + @patch("sunbeam.steps.hypervisor.get_admin_connection") + def test_run_force_skips_client_error( + self, + get_admin_connection, + basic_jhelper, + basic_deployment, + step_context, + error, + ): + get_admin_connection.side_effect = error + step = RemoveHypervisorReferencesStep( + basic_jhelper, + basic_deployment, + "cloud-4", + "cloud-4.maas", + force=True, + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.SKIPPED + assert result.message == str(error) + get_admin_connection.assert_called_once_with(basic_jhelper, basic_deployment) + + @pytest.mark.parametrize( + "error", + [ + ConnectFailure("control plane unavailable"), + click.ClickException("Unable to get keystone leader"), + ], + ) + @patch("sunbeam.steps.hypervisor.HYPERVISOR_REFERENCES_POLL_INTERVAL", 0) + @patch("sunbeam.steps.hypervisor.get_admin_connection") + def test_run_retries_client_error( + self, + get_admin_connection, + basic_jhelper, + basic_deployment, + step_context, + error, + ): + conn = Mock() + conn.compute.services.return_value = [] + conn.network.agents.return_value = [] + get_admin_connection.side_effect = [error, conn] + step = RemoveHypervisorReferencesStep( + basic_jhelper, + basic_deployment, + "cloud-4", + "cloud-4.maas", + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.COMPLETED + assert get_admin_connection.call_count == 2 + + @patch("sunbeam.steps.hypervisor.HYPERVISOR_REFERENCES_TIMEOUT", 0) + @patch("sunbeam.steps.hypervisor.get_admin_connection") + @patch("sunbeam.steps.hypervisor.remove_network_service") + @patch("sunbeam.steps.hypervisor.remove_compute_service") + def test_run_fails_when_references_persist( + self, + remove_compute_service, + remove_network_service, + get_admin_connection, + basic_jhelper, + basic_deployment, + step_context, + ): + conn = Mock() + conn.compute.services.return_value = [Mock()] + conn.network.agents.return_value = [Mock()] + get_admin_connection.return_value = conn + step = RemoveHypervisorReferencesStep( + basic_jhelper, + basic_deployment, + "cloud-4", + "cloud-4.maas", + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.FAILED + assert remove_compute_service.called + assert remove_network_service.called + + @pytest.mark.parametrize("force", [False, True]) + @pytest.mark.parametrize( + "error", + [ + ModelNotFoundException("Model 'openstack' not found"), + ApplicationNotFoundException("Application 'keystone' not found"), + ], + ) + @patch("sunbeam.steps.hypervisor.get_admin_connection") + def test_run_skips_without_control_plane( + self, + get_admin_connection, + basic_jhelper, + basic_deployment, + step_context, + force, + error, + ): + get_admin_connection.side_effect = error + step = RemoveHypervisorReferencesStep( + basic_jhelper, + basic_deployment, + "cloud-4", + "cloud-4.maas", + force=force, + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.SKIPPED + get_admin_connection.assert_called_once_with(basic_jhelper, basic_deployment) + + class TestReapplyHypervisorTerraformPlanStep: @pytest.fixture def get_network_config_patch(self): From e425f2316d3754f1a577f49388ce23d242d7cecf Mon Sep 17 00:00:00 2001 From: Seyeong Kim Date: Wed, 12 Aug 2026 12:44:36 +0000 Subject: [PATCH 3/7] fix(maas): remove stale Cinder services Removing a cinder-volume Juju unit leaves its service records in Cinder, where they still point to a node no longer in the cluster. Disable matching services before unit removal, then delete their records through a healthy Cinder control-plane unit and verify that none remain. Signed-off-by: Seyeong Kim --- sunbeam-python/sunbeam/core/juju.py | 11 +- .../sunbeam/provider/maas/commands.py | 14 + sunbeam-python/sunbeam/steps/cinder_volume.py | 238 ++++++++- .../tests/unit/sunbeam/core/test_juju.py | 18 + .../unit/sunbeam/provider/maas/test_maas.py | 22 + .../unit/sunbeam/steps/test_cinder_volume.py | 454 +++++++++++++++++- 6 files changed, 753 insertions(+), 4 deletions(-) diff --git a/sunbeam-python/sunbeam/core/juju.py b/sunbeam-python/sunbeam/core/juju.py index e1f095f8d..4a5f9bc58 100644 --- a/sunbeam-python/sunbeam/core/juju.py +++ b/sunbeam-python/sunbeam/core/juju.py @@ -820,12 +820,19 @@ def run_cmd_on_unit_payload( if env: args.extend(f"--env={k}={v}" for k, v in env.items()) + stderr = "" with self._model(model) as juju: try: stdout, _ = juju._cli(*args, "--", *(cmd.split()), log=False) except jubilant.CLIError as e: - stdout = e.stdout - return json.loads(stdout)[name]["results"] + stdout, stderr = e.stdout, e.stderr or "" + try: + return json.loads(stdout)[name]["results"] + except (json.JSONDecodeError, KeyError, TypeError) as e: + raise ExecFailedException( + f"Failed to parse command result for unit {name!r}: " + f"{stderr.strip() or e}" + ) from e def run_action( self, diff --git a/sunbeam-python/sunbeam/provider/maas/commands.py b/sunbeam-python/sunbeam/provider/maas/commands.py index 2d16e00f0..582ae946c 100644 --- a/sunbeam-python/sunbeam/provider/maas/commands.py +++ b/sunbeam-python/sunbeam/provider/maas/commands.py @@ -125,6 +125,8 @@ CheckCinderVolumeDistributionStep, DeployCinderVolumeApplicationStep, DestroyCinderVolumeApplicationStep, + DisableCinderVolumeServicesStep, + RemoveCinderVolumeServicesStep, RemoveCinderVolumeUnitsStep, ) from sunbeam.steps.clusterd import APPLICATION as CLUSTERD_APPLICATION @@ -1710,9 +1712,21 @@ def remove_node(ctx: click.Context, name: str, force: bool, show_hints: bool) -> deployment.openstack_machines_model, force, ), + DisableCinderVolumeServicesStep( + jhelper, + deployment, + machine["hostname"], + machine["fqdn"], + ), RemoveCinderVolumeUnitsStep( client, name, jhelper, deployment.openstack_machines_model ), + RemoveCinderVolumeServicesStep( + jhelper, + deployment, + machine["hostname"], + machine["fqdn"], + ), RemoveMicrocephUnitsStep( client, name, jhelper, deployment.openstack_machines_model ), diff --git a/sunbeam-python/sunbeam/steps/cinder_volume.py b/sunbeam-python/sunbeam/steps/cinder_volume.py index 53ceed0fe..38abb2ec5 100644 --- a/sunbeam-python/sunbeam/steps/cinder_volume.py +++ b/sunbeam-python/sunbeam/steps/cinder_volume.py @@ -2,6 +2,7 @@ # SPDX-License-Identifier: Apache-2.0 import logging +import typing from typing import Any import sunbeam.steps.microceph as microceph @@ -10,19 +11,38 @@ from sunbeam.clusterd.service import ( NodeNotExistInClusterException, ) -from sunbeam.core.common import BaseStep, Result, ResultType, Role, StepContext +from sunbeam.core.common import ( + BaseStep, + Result, + ResultType, + Role, + StepContext, + SunbeamException, +) from sunbeam.core.deployment import Deployment, Networks from sunbeam.core.juju import ( ApplicationNotFoundException, + JujuException, JujuHelper, + ModelNotFoundException, ) from sunbeam.core.manifest import CharmManifest, Manifest +from sunbeam.core.openstack import OPENSTACK_MODEL +from sunbeam.core.openstack_api import get_admin_connection from sunbeam.core.steps import ( DeployMachineApplicationStep, DestroyMachineApplicationStep, RemoveMachineUnitsStep, ) from sunbeam.core.terraform import TerraformException, TerraformHelper +from sunbeam.lazy import LazyImport + +if typing.TYPE_CHECKING: + import openstack + from keystoneauth1 import exceptions as keystoneauth_exceptions +else: + keystoneauth_exceptions = LazyImport("keystoneauth1.exceptions") + openstack = LazyImport("openstack") LOG = logging.getLogger(__name__) CONFIG_KEY = "TerraformVarsCinderVolumePlan" @@ -31,6 +51,31 @@ CINDER_VOLUME_UNIT_TIMEOUT = ( 1800 # 30 minutes, adding / removing units can take a long time ) +CINDER_APPLICATION = "cinder" +CINDER_API_CONTAINER = "cinder-api" +CINDER_VOLUME_BINARY = "cinder-volume" +CINDER_SERVICE_REMOVE_REASON = "Removing node from cluster" +# cinder-manage soft-deletes a single service row, which takes seconds. Five +# minutes also covers waiting for a hook already running on the unit and +# oslo.db's default connection retries (10 attempts, 10 seconds apart) when +# the database is unreachable. +CINDER_SERVICE_REMOVE_TIMEOUT = 300 + + +def get_cinder_volume_services( + conn: "openstack.connection.Connection", + hostname: str, + fqdn: str, +) -> list[Any]: + """Return cinder-volume services for the exact node hostnames.""" + expected_hosts = {hostname, fqdn} + + services = [] + for service in conn.block_storage.services(binary=CINDER_VOLUME_BINARY): + if service.host.split("@", 1)[0] in expected_hosts: + services.append(service) + + return sorted(services, key=lambda service: service.host) def get_mandatory_control_plane_offers( @@ -231,6 +276,197 @@ def get_unit_timeout(self) -> int: return CINDER_VOLUME_UNIT_TIMEOUT +class _CinderVolumeServiceStep(BaseStep): + """Shared service discovery for Cinder cleanup steps.""" + + def __init__( + self, + name: str, + description: str, + jhelper: JujuHelper, + deployment: Deployment, + hostname: str, + fqdn: str, + ): + super().__init__(name, description) + self.jhelper = jhelper + self.deployment = deployment + self.hostname = hostname + self.fqdn = fqdn + self.connection: Any | None = None + self.services: list[Any] = [] + + def _discover_services(self) -> Result: + """Discover Cinder service records for the target node.""" + try: + self.connection = get_admin_connection(self.jhelper, self.deployment) + self.services = get_cinder_volume_services( + self.connection, self.hostname, self.fqdn + ) + except (ModelNotFoundException, ApplicationNotFoundException) as e: + LOG.debug("OpenStack control plane not deployed, skipping: %r", e) + return Result(ResultType.SKIPPED) + except ( + openstack.exceptions.SDKException, + keystoneauth_exceptions.ClientException, + ) as e: + LOG.warning("Failed to discover Cinder volume services: %r", e) + return Result(ResultType.FAILED, str(e)) + + if not self.services: + return Result(ResultType.SKIPPED) + return Result(ResultType.COMPLETED) + + def _current_services(self) -> list[Any]: + """Return matching Cinder service records.""" + if self.connection is None: + raise SunbeamException("Cinder admin connection not found") + return get_cinder_volume_services(self.connection, self.hostname, self.fqdn) + + +class DisableCinderVolumeServicesStep(_CinderVolumeServiceStep): + """Disable matching Cinder volume services before unit removal.""" + + def __init__( + self, + jhelper: JujuHelper, + deployment: Deployment, + hostname: str, + fqdn: str, + ): + super().__init__( + "Disable Cinder Volume services", + "Disabling Cinder Volume services", + jhelper, + deployment, + hostname, + fqdn, + ) + + def is_skip(self, context: StepContext) -> Result: + """Determine whether matching Cinder services exist.""" + return self._discover_services() + + def run(self, context: StepContext) -> Result: + """Disable each enabled matching Cinder volume service.""" + if self.connection is None: + return Result(ResultType.FAILED, "Cinder admin connection not found") + try: + for service in self.services: + if service.status != "enabled": + continue + LOG.info("Disabling %s on %s", service.binary, service.host) + self.connection.block_storage.disable_service( + service, reason=CINDER_SERVICE_REMOVE_REASON + ) + except ( + openstack.exceptions.SDKException, + keystoneauth_exceptions.ClientException, + ) as e: + LOG.warning("Failed to disable Cinder volume service: %r", e) + return Result(ResultType.FAILED, str(e)) + + return Result(ResultType.COMPLETED) + + +class RemoveCinderVolumeServicesStep(_CinderVolumeServiceStep): + """Remove matching Cinder volume service records after unit removal.""" + + def __init__( + self, + jhelper: JujuHelper, + deployment: Deployment, + hostname: str, + fqdn: str, + ): + super().__init__( + "Remove Cinder Volume services", + "Removing Cinder Volume services", + jhelper, + deployment, + hostname, + fqdn, + ) + + def is_skip(self, context: StepContext) -> Result: + """Determine whether matching Cinder services exist.""" + return self._discover_services() + + def _healthy_units(self) -> list[str]: + """Return active Cinder API units with a leader first. + + Only idle or executing agents are used: juju exec waits for a + running hook to finish, while lost or errored agents never run it. + """ + try: + application = self.jhelper.get_application( + CINDER_APPLICATION, OPENSTACK_MODEL + ) + except JujuException as e: + LOG.warning("Failed to find Cinder control-plane units: %r", e) + return [] + + units = [ + (name, unit) + for name, unit in application.units.items() + if unit.workload_status.current == "active" + and unit.juju_status.current in ("idle", "executing") + ] + units.sort(key=lambda item: (not item[1].leader, item[0])) + return [name for name, _ in units] + + def _remove_service(self, unit: str, service: Any) -> None: + """Remove one Cinder service record through a healthy unit.""" + command = f"cinder-manage service remove {CINDER_VOLUME_BINARY} {service.host}" + result = self.jhelper.run_cmd_on_unit_payload( + unit, + OPENSTACK_MODEL, + command, + CINDER_API_CONTAINER, + timeout=CINDER_SERVICE_REMOVE_TIMEOUT, + ) + if result.get("return-code") != 0: + raise JujuException(f"Failed to remove Cinder service {service.host}") + + def run(self, context: StepContext) -> Result: + """Remove matching records and verify that none remain.""" + if self.connection is None: + return Result(ResultType.FAILED, "Cinder admin connection not found") + + try: + remaining = self.services + healthy_units = self._healthy_units() + if not healthy_units: + raise SunbeamException( + "No healthy Cinder control-plane units available" + ) + + for unit in healthy_units: + for service in remaining: + try: + self._remove_service(unit, service) + except JujuException as e: + LOG.warning( + "cinder-manage failed to remove %s on unit %s: %r", + service.host, + unit, + e, + ) + break + remaining = self._current_services() + if not remaining: + return Result(ResultType.COMPLETED) + + raise SunbeamException("Cinder service records remain after removal") + except ( + SunbeamException, + openstack.exceptions.SDKException, + keystoneauth_exceptions.ClientException, + ) as e: + LOG.warning("Failed to remove Cinder volume services: %r", e) + return Result(ResultType.FAILED, str(e)) + + class CheckCinderVolumeDistributionStep(BaseStep): _APPLICATION = APPLICATION diff --git a/sunbeam-python/tests/unit/sunbeam/core/test_juju.py b/sunbeam-python/tests/unit/sunbeam/core/test_juju.py index 7a248e4e7..57d31c7be 100644 --- a/sunbeam-python/tests/unit/sunbeam/core/test_juju.py +++ b/sunbeam-python/tests/unit/sunbeam/core/test_juju.py @@ -238,6 +238,24 @@ def test_run_cmd_on_unit_payload_cli_error(jhelper, juju): assert result["err"] == "fail" +@pytest.mark.parametrize( + "stdout", + ["", "{}", json.dumps({"app/0": {}})], +) +def test_run_cmd_on_unit_payload_normalizes_invalid_result(jhelper, juju, stdout): + juju._cli.return_value = (stdout, "") + + with pytest.raises(jujulib.ExecFailedException): + jhelper.run_cmd_on_unit_payload("app/0", "test-model", "ls", "container") + + +def test_run_cmd_on_unit_payload_normalizes_invalid_cli_error(jhelper, juju): + juju._cli.side_effect = jubilant.CLIError(1, ["exec"], "", "unit unavailable") + + with pytest.raises(jujulib.ExecFailedException, match="unit unavailable"): + jhelper.run_cmd_on_unit_payload("app/0", "test-model", "ls", "container") + + def test_set_model_config(jhelper, juju): jhelper.set_model_config("test-model", {"app": "bar"}) juju.model_config.assert_called() diff --git a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py index 99763b18b..4de766f30 100644 --- a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py +++ b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py @@ -57,6 +57,11 @@ ZoneBalanceCheck, ZonesCheck, ) +from sunbeam.steps.cinder_volume import ( + DisableCinderVolumeServicesStep, + RemoveCinderVolumeServicesStep, + RemoveCinderVolumeUnitsStep, +) from sunbeam.steps.hypervisor import ( RemoveHypervisorReferencesStep, RemoveHypervisorUnitStep, @@ -2566,6 +2571,23 @@ def test_remove_orders_cleanup_steps( assert references.force is False assert microovn_index < references_index + disable_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, DisableCinderVolumeServicesStep) + ) + cinder_units_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, RemoveCinderVolumeUnitsStep) + ) + cinder_services_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, RemoveCinderVolumeServicesStep) + ) + assert disable_index < cinder_units_index < cinder_services_index + @patch("sunbeam.provider.maas.commands.JujuHelper") @patch("sunbeam.provider.maas.commands.run_preflight_checks") @patch("sunbeam.provider.maas.commands.run_plan") diff --git a/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py b/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py index 930c10ad9..35242628e 100644 --- a/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py +++ b/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py @@ -4,12 +4,24 @@ from unittest.mock import MagicMock, Mock, patch import pytest - +from keystoneauth1.exceptions.catalog import EndpointNotFound +from keystoneauth1.exceptions.connection import ConnectFailure + +from sunbeam.core.common import ResultType +from sunbeam.core.juju import ( + ApplicationNotFoundException, + ExecFailedException, + JujuException, + ModelNotFoundException, +) from sunbeam.steps.cinder_volume import ( CINDER_VOLUME_APP_TIMEOUT, CINDER_VOLUME_UNIT_TIMEOUT, DeployCinderVolumeApplicationStep, + DisableCinderVolumeServicesStep, + RemoveCinderVolumeServicesStep, RemoveCinderVolumeUnitsStep, + get_cinder_volume_services, ) @@ -361,3 +373,443 @@ def test_get_unit_timeout(self, remove_cinder_volume_units_step): remove_cinder_volume_units_step.get_unit_timeout() == CINDER_VOLUME_UNIT_TIMEOUT ) + + +class TestCinderVolumeServiceCleanup: + @pytest.fixture + def connection(self): + return Mock() + + @pytest.fixture + def services(self): + return [ + Mock(binary="cinder-volume", host="cloud-4@backend-a", status="enabled"), + Mock( + binary="cinder-volume", host="cloud-4.maas@backend-b", status="disabled" + ), + ] + + def test_get_cinder_volume_services_matches_exact_hosts(self): + matching_short = Mock(binary="cinder-volume", host="cloud-4@backend-a") + matching_fqdn = Mock(binary="cinder-volume", host="cloud-4.maas@backend-b") + unrelated = Mock(binary="cinder-volume", host="cloud-40@backend-a") + conn = Mock() + conn.block_storage.services.return_value = [ + unrelated, + matching_fqdn, + matching_short, + ] + + services = get_cinder_volume_services(conn, "cloud-4", "cloud-4.maas") + + conn.block_storage.services.assert_called_once_with(binary="cinder-volume") + assert services == [matching_fqdn, matching_short] + + def test_disable_enabled_services( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + return_value=services, + ) + step = DisableCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.COMPLETED + connection.block_storage.disable_service.assert_called_once_with( + services[0], reason="Removing node from cluster" + ) + + def test_discovery_fails_on_keystoneauth_error( + self, + mocker, + basic_jhelper, + basic_deployment, + step_context, + ): + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + side_effect=EndpointNotFound("volume endpoint missing"), + ) + step = DisableCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + result = step.is_skip(step_context) + + assert result.result_type == ResultType.FAILED + + def test_disable_fails_on_keystoneauth_error( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + connection.block_storage.disable_service.side_effect = ConnectFailure( + "volume endpoint unavailable" + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + return_value=services, + ) + step = DisableCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.FAILED + + def test_remove_prefers_healthy_nonleader_when_leader_unhealthy( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + unhealthy_leader = Mock( + workload_status=Mock(current="blocked"), + juju_status=Mock(current="idle"), + leader=True, + ) + healthy_unit = Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current="idle"), + leader=False, + ) + application = Mock( + units={"cinder/0": unhealthy_leader, "cinder/1": healthy_unit} + ) + basic_jhelper.get_application.return_value = application + basic_jhelper.run_cmd_on_unit_payload.return_value = {"return-code": 0} + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + side_effect=[services, []], + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.COMPLETED + assert all( + call.args[0] == "cinder/1" + for call in basic_jhelper.run_cmd_on_unit_payload.call_args_list + ) + + def test_remove_succeeds_when_ambiguous_command_error_has_no_records( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + unit = Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current="idle"), + leader=True, + ) + basic_jhelper.get_application.return_value = Mock(units={"cinder/0": unit}) + basic_jhelper.run_cmd_on_unit_payload.side_effect = JujuException("unknown") + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + side_effect=[services, []], + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.COMPLETED + basic_jhelper.run_cmd_on_unit_payload.assert_called_once() + + def test_remove_falls_back_when_first_unit_disappears( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + healthy_units = { + name: Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current="idle"), + leader=False, + ) + for name in ("cinder/0", "cinder/1") + } + service = services[0] + basic_jhelper.get_application.return_value = Mock(units=healthy_units) + basic_jhelper.run_cmd_on_unit_payload.side_effect = [ + ExecFailedException("unit disappeared"), + {"return-code": 0}, + ] + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + side_effect=[[service], [service], []], + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.COMPLETED + assert [ + call.args[0] + for call in basic_jhelper.run_cmd_on_unit_payload.call_args_list + ] == ["cinder/0", "cinder/1"] + + def test_remove_fails_on_keystoneauth_error( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + unit = Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current="idle"), + leader=True, + ) + basic_jhelper.get_application.return_value = Mock(units={"cinder/0": unit}) + basic_jhelper.run_cmd_on_unit_payload.return_value = {"return-code": 0} + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + side_effect=[ + services, + ConnectFailure("volume endpoint unavailable"), + ], + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.FAILED + + def test_remove_fails_without_healthy_units( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + unhealthy = Mock( + workload_status=Mock(current="blocked"), + juju_status=Mock(current="idle"), + leader=True, + ) + basic_jhelper.get_application.return_value = Mock(units={"cinder/0": unhealthy}) + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + return_value=services, + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.FAILED + basic_jhelper.run_cmd_on_unit_payload.assert_not_called() + + def test_remove_fails_after_healthy_candidates_leave_records( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + healthy_units = { + name: Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current="idle"), + leader=False, + ) + for name in ("cinder/0", "cinder/1") + } + basic_jhelper.get_application.return_value = Mock(units=healthy_units) + basic_jhelper.run_cmd_on_unit_payload.return_value = {"return-code": 1} + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + return_value=services, + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.FAILED + assert [ + call.args[0] + for call in basic_jhelper.run_cmd_on_unit_payload.call_args_list + ] == ["cinder/0", "cinder/1"] + + def test_remove_skips_when_no_records( + self, mocker, basic_jhelper, basic_deployment, connection, step_context + ): + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", return_value=[] + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.SKIPPED + basic_jhelper.get_application.assert_not_called() + + @pytest.mark.parametrize( + "error", + [ + ModelNotFoundException("Model 'openstack' not found"), + ApplicationNotFoundException("Application 'keystone' not found"), + ], + ) + def test_discovery_skips_without_control_plane( + self, mocker, basic_jhelper, basic_deployment, step_context, error + ): + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + side_effect=error, + ) + step = DisableCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.SKIPPED + + def test_remove_uses_active_unit_while_it_runs_a_hook( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + ): + busy_unit = Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current="executing"), + leader=True, + ) + basic_jhelper.get_application.return_value = Mock(units={"cinder/0": busy_unit}) + basic_jhelper.run_cmd_on_unit_payload.return_value = {"return-code": 0} + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + side_effect=[services, []], + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.COMPLETED + assert all( + call.args[0] == "cinder/0" + for call in basic_jhelper.run_cmd_on_unit_payload.call_args_list + ) + + @pytest.mark.parametrize("agent_status", ["lost", "error"]) + def test_remove_skips_units_with_unavailable_agents( + self, + mocker, + basic_jhelper, + basic_deployment, + connection, + services, + step_context, + agent_status, + ): + unavailable_leader = Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current=agent_status), + leader=True, + ) + idle_unit = Mock( + workload_status=Mock(current="active"), + juju_status=Mock(current="idle"), + leader=False, + ) + basic_jhelper.get_application.return_value = Mock( + units={"cinder/0": unavailable_leader, "cinder/1": idle_unit} + ) + basic_jhelper.run_cmd_on_unit_payload.return_value = {"return-code": 0} + mocker.patch( + "sunbeam.steps.cinder_volume.get_admin_connection", + return_value=connection, + ) + mocker.patch( + "sunbeam.steps.cinder_volume.get_cinder_volume_services", + side_effect=[services, []], + ) + step = RemoveCinderVolumeServicesStep( + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + ) + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.COMPLETED + assert all( + call.args[0] == "cinder/1" + for call in basic_jhelper.run_cmd_on_unit_payload.call_args_list + ) From b3f53cdfd7dd9337bd72d8e2fd4de5c5e9d80c10 Mon Sep 17 00:00:00 2001 From: Seyeong Kim Date: Thu, 13 Aug 2026 02:45:12 +0000 Subject: [PATCH 4/7] fix(maas): remove OSDs before node removal Removing a MicroCeph Juju unit can leave its DB-backed OSDs in Ceph after the node is gone. Remove those OSDs before deleting the unit. Stop when CRUSH contains target OSDs absent from MicroCeph, because that orphaned state needs separate recovery. Signed-off-by: Seyeong Kim --- sunbeam-python/sunbeam/core/juju.py | 14 +- .../sunbeam/provider/maas/commands.py | 9 + sunbeam-python/sunbeam/steps/microceph.py | 155 +++++++++++++++++ .../tests/unit/sunbeam/core/test_juju.py | 15 ++ .../unit/sunbeam/provider/maas/test_maas.py | 29 +++- .../unit/sunbeam/steps/test_microceph.py | 160 +++++++++++++++++- 6 files changed, 370 insertions(+), 12 deletions(-) diff --git a/sunbeam-python/sunbeam/core/juju.py b/sunbeam-python/sunbeam/core/juju.py index 4a5f9bc58..e06ba43c9 100644 --- a/sunbeam-python/sunbeam/core/juju.py +++ b/sunbeam-python/sunbeam/core/juju.py @@ -759,22 +759,22 @@ def run_cmd_on_machine_unit_payload( ) -> "jubilant.Task": """Run a shell command on a machine unit. - Returns action results irrespective of the return-code - in action results. - :name: unit name :model: Name of the model where the application is located :cmd: Command to run :timeout: Timeout in seconds :returns: Command results - - Command execution failures are part of the results with - return-code, stdout, stderr. + :raises: ExecFailedException if command execution fails """ with self._model(model) as juju: try: task = juju.exec(cmd, unit=name, wait=timeout) - except jubilant.TaskError as e: + except ( + jubilant.TaskError, + jubilant.CLIError, + TimeoutError, + ValueError, + ) as e: raise ExecFailedException( f"Failed to run command {cmd!r} on unit" f" {name!r} in model {model!r}: {e}" diff --git a/sunbeam-python/sunbeam/provider/maas/commands.py b/sunbeam-python/sunbeam/provider/maas/commands.py index 582ae946c..028cb66d7 100644 --- a/sunbeam-python/sunbeam/provider/maas/commands.py +++ b/sunbeam-python/sunbeam/provider/maas/commands.py @@ -175,6 +175,7 @@ CheckMicrocephDistributionStep, DeployMicrocephApplicationStep, DestroyMicrocephApplicationStep, + RemoveMicrocephOSDsStep, RemoveMicrocephUnitsStep, SetCephMgrPoolSizeStep, ) @@ -1727,6 +1728,14 @@ def remove_node(ctx: click.Context, name: str, force: bool, show_hints: bool) -> machine["hostname"], machine["fqdn"], ), + RemoveMicrocephOSDsStep( + client, + name, + jhelper, + deployment.openstack_machines_model, + force=force, + hostnames=(machine["hostname"], machine["fqdn"]), + ), RemoveMicrocephUnitsStep( client, name, jhelper, deployment.openstack_machines_model ), diff --git a/sunbeam-python/sunbeam/steps/microceph.py b/sunbeam-python/sunbeam/steps/microceph.py index e87aefdce..2c504235b 100644 --- a/sunbeam-python/sunbeam/steps/microceph.py +++ b/sunbeam-python/sunbeam/steps/microceph.py @@ -2,6 +2,7 @@ # SPDX-License-Identifier: Apache-2.0 import ast +import json import logging from typing import Any @@ -24,6 +25,7 @@ from sunbeam.core.juju import ( ActionFailedException, ApplicationNotFoundException, + ExecFailedException, JujuHelper, LeaderNotFoundException, UnitNotFoundException, @@ -231,6 +233,159 @@ def get_unit_timeout(self) -> int: return MICROCEPH_UNIT_TIMEOUT +class RemoveMicrocephOSDsStep(BaseStep): + """Remove a node's MicroCeph OSDs before removing its unit.""" + + # Default timeout of `microceph disk remove`. + _OSD_REMOVE_TIMEOUT = 1800 + # Juju cancels the command when its wait expires, so let MicroCeph reach + # its own timeout first and report that error. + _COMMAND_TIMEOUT = _OSD_REMOVE_TIMEOUT + 60 + + def __init__( + self, + client: Client, + name: str, + jhelper: JujuHelper, + model: str, + force: bool = False, + hostnames: tuple[str, ...] = (), + ): + super().__init__( + "Remove MicroCeph OSDs", + "Removing MicroCeph OSDs", + ) + self.client = client + self.node = name + self.jhelper = jhelper + self.model = model + self.force = force + self._hostnames = tuple(dict.fromkeys((name, *hostnames))) + self.unit: str | None = None + + def _prepare(self) -> Result: + """Find the target unit, or the leader if the target is gone.""" + self.unit = None + try: + try: + node_info = self.client.cluster.get_node_info(self.node) + self.unit = self.jhelper.get_unit_from_machine( + APPLICATION, str(node_info["machineid"]), self.model + ) + except (NodeNotExistInClusterException, UnitNotFoundException): + self.unit = self.jhelper.get_leader_unit(APPLICATION, self.model) + except ApplicationNotFoundException: + LOG.debug("Failed to get application", exc_info=True) + return Result( + ResultType.SKIPPED, + f"Application {APPLICATION} has not been deployed yet", + ) + except LeaderNotFoundException as e: + return Result(ResultType.FAILED, str(e)) + + return Result(ResultType.COMPLETED) + + def is_skip(self, context: StepContext) -> Result: + """Determine whether cleanup can run and whether it is needed.""" + return self._prepare() + + def _run_command(self, command: str) -> str: + """Run a MicroCeph command and fail on transport or command errors.""" + if self.unit is None: + raise SunbeamException("MicroCeph cleanup unit is not available") + try: + result = self.jhelper.run_cmd_on_machine_unit_payload( + self.unit, + self.model, + command, + timeout=self._COMMAND_TIMEOUT, + ) + except ExecFailedException as e: + raise SunbeamException(f"Failed to run {command!r}: {e}") from e + + return result.stdout + + def _list_configured_osd_ids(self) -> list[int]: + """Return target OSD IDs that still exist in the MicroCeph database.""" + try: + disks = json.loads(self._run_command("microceph disk list --json"))[ + "ConfiguredDisks" + ] + return sorted( + {disk["osd"] for disk in disks if disk["location"] in self._hostnames} + ) + except (json.JSONDecodeError, KeyError, TypeError) as e: + raise SunbeamException( + f"Failed to parse configured disk listing: {e}" + ) from e + + def _list_crush_osd_ids(self) -> list[int]: + """Return OSD IDs under the target CRUSH host.""" + try: + nodes = json.loads( + self._run_command("microceph.ceph osd tree --format json") + )["nodes"] + return sorted( + { + osd_id + for node in nodes + if node["type"] == "host" and node["name"] in self._hostnames + for osd_id in node["children"] + } + ) + except (json.JSONDecodeError, KeyError, TypeError) as e: + raise SunbeamException(f"Failed to parse CRUSH tree: {e}") from e + + def _list_target_osds(self) -> tuple[list[int], list[int]]: + """Read both target OSD sources before changing either source.""" + return self._list_configured_osd_ids(), self._list_crush_osd_ids() + + def run(self, context: StepContext) -> Result: + """Remove DB-backed OSDs and verify both MicroCeph and CRUSH state.""" + if self.unit is None: + preparation = self._prepare() + if preparation.result_type != ResultType.COMPLETED: + return preparation + + try: + configured_osds, crush_osds = self._list_target_osds() + crush_only_osds = sorted(set(crush_osds) - set(configured_osds)) + if crush_only_osds: + return Result( + ResultType.FAILED, + f"CRUSH-only OSDs for {self.node}: {crush_only_osds}", + ) + + for osd_id in configured_osds: + command = ( + f"microceph disk remove osd.{osd_id} " + f"--timeout {self._OSD_REMOVE_TIMEOUT}" + ) + if self.force: + command += " --confirm-failure-domain-downgrade" + self._run_command(command) + + if not configured_osds: + return Result(ResultType.COMPLETED) + + remaining_configured, remaining_crush = self._list_target_osds() + if remaining_configured: + return Result( + ResultType.FAILED, + f"Configured OSDs remain for {self.node}: {remaining_configured}", + ) + if remaining_crush: + return Result( + ResultType.FAILED, + f"CRUSH OSDs remain for {self.node}: {remaining_crush}", + ) + except SunbeamException as e: + LOG.debug("Failed to clean up MicroCeph OSDs", exc_info=True) + return Result(ResultType.FAILED, str(e)) + + return Result(ResultType.COMPLETED) + + class ConfigureMicrocephOSDStep(BaseStep): """Configure Microceph OSD disks.""" diff --git a/sunbeam-python/tests/unit/sunbeam/core/test_juju.py b/sunbeam-python/tests/unit/sunbeam/core/test_juju.py index 57d31c7be..0822e4410 100644 --- a/sunbeam-python/tests/unit/sunbeam/core/test_juju.py +++ b/sunbeam-python/tests/unit/sunbeam/core/test_juju.py @@ -208,6 +208,21 @@ def test_run_cmd_on_machine_unit_payload_success(jhelper, juju): assert result.results["result"] == "ok" +@pytest.mark.parametrize( + "error", + [ + TimeoutError("timed out"), + jubilant.CLIError(1, ["exec"], "", "controller unavailable"), + ValueError("unit has no result"), + ], +) +def test_run_cmd_on_machine_unit_payload_normalizes_exec_errors(jhelper, juju, error): + juju.exec.side_effect = error + + with pytest.raises(jujulib.ExecFailedException): + jhelper.run_cmd_on_machine_unit_payload("app/0", "test-model", "ls") + + def test_run_action_success(jhelper, juju): juju.run = Mock(return_value=Mock(success=True, results={"app": "bar"})) diff --git a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py index 4de766f30..6c24a4c79 100644 --- a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py +++ b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py @@ -67,6 +67,7 @@ RemoveHypervisorUnitStep, ) from sunbeam.steps.juju import RemoveJujuMachineStep +from sunbeam.steps.microceph import RemoveMicrocephOSDsStep, RemoveMicrocephUnitsStep from sunbeam.steps.microovn import ( ReapplyMicroOVNTerraformPlanStep, RemoveMicroOVNUnitsStep, @@ -2516,7 +2517,7 @@ def test_storage_ippool_label_with_different_names(self): assert deployment.storage_ip_pool == expected_label -class TestRemoveNodeRoleDistributor: +class TestRemoveNode: @pytest.fixture(autouse=True) def mock_maas_machine(self, mocker): maas_client = mocker.patch( @@ -2531,19 +2532,25 @@ def mock_maas_machine(self, mocker): @patch("sunbeam.provider.maas.commands.JujuHelper") @patch("sunbeam.provider.maas.commands.run_preflight_checks") @patch("sunbeam.provider.maas.commands.run_plan") + @pytest.mark.parametrize( + ("cli_args", "force"), + [(["node-1"], False), (["--force", "node-1"], True)], + ) def test_remove_orders_cleanup_steps( self, run_plan_cmd, run_preflight, juju_helper, mock_maas_machine, + cli_args, + force, ): maas_client, get_machine = mock_maas_machine deployment = Mock() deployment.openstack_machines_model = "openstack-machines" deployment.get_ovn_manager.return_value.get_machines.return_value = [] - result = CliRunner().invoke(remove_node, ["node-1"], obj=deployment) + result = CliRunner().invoke(remove_node, cli_args, obj=deployment) assert result.exit_code == 0, result.output get_machine.assert_called_once_with(maas_client, "node-1") @@ -2568,7 +2575,7 @@ def test_remove_orders_cleanup_steps( assert hypervisor_unit.deployment is None assert references._hostnames == ("node-1", "node-1.maas") - assert references.force is False + assert references.force is force assert microovn_index < references_index disable_index = next( @@ -2588,6 +2595,21 @@ def test_remove_orders_cleanup_steps( ) assert disable_index < cinder_units_index < cinder_services_index + osd_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, RemoveMicrocephOSDsStep) + ) + unit_index = next( + i + for i, step in enumerate(plan) + if isinstance(step, RemoveMicrocephUnitsStep) + ) + assert osd_index < unit_index + assert plan[osd_index].node == "node-1" + assert plan[osd_index]._hostnames == ("node-1", "node-1.maas") + assert plan[osd_index].force is force + @patch("sunbeam.provider.maas.commands.JujuHelper") @patch("sunbeam.provider.maas.commands.run_preflight_checks") @patch("sunbeam.provider.maas.commands.run_plan") @@ -2696,6 +2718,7 @@ def test_remove_skips_role_distributor_when_microovn_has_no_machines( isinstance(step, ReapplyMicroOVNTerraformPlanStep) for step in plan ) + class TestIsMaasDeployment: """Test is_maas_deployment TypeGuard function.""" diff --git a/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py b/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py index 45710175c..d1c23b5db 100644 --- a/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py +++ b/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py @@ -1,11 +1,48 @@ # SPDX-FileCopyrightText: 2023 - Canonical Ltd # SPDX-License-Identifier: Apache-2.0 +import json from unittest.mock import Mock from sunbeam.core.common import ResultType -from sunbeam.core.juju import ActionFailedException -from sunbeam.steps.microceph import ConfigureMicrocephOSDStep, SetCephMgrPoolSizeStep +from sunbeam.core.juju import ( + ActionFailedException, + ExecFailedException, + UnitNotFoundException, +) +from sunbeam.steps.microceph import ( + ConfigureMicrocephOSDStep, + RemoveMicrocephOSDsStep, + SetCephMgrPoolSizeStep, +) + + +def _command_result(stdout=""): + return Mock(stdout=stdout) + + +def _configured_disks(disks): + return json.dumps({"ConfiguredDisks": disks}) + + +def _crush_tree(hostname=None, children=None): + nodes = [] + if hostname is not None: + nodes.append({"name": hostname, "type": "host", "children": children or []}) + return json.dumps({"nodes": nodes}) + + +def _cleanup_step(cclient, jhelper, force=False): + cclient.cluster.get_node_info.return_value = {"machineid": "1"} + jhelper.get_unit_from_machine.return_value = "microceph/0" + return RemoveMicrocephOSDsStep( + cclient, + "node-1", + jhelper, + "test-model", + hostnames=("node-1", "node-1.maas"), + force=force, + ) class TestConfigureMicrocephOSDStep: @@ -88,6 +125,125 @@ def test_run_with_wipe_false(self, cclient, jhelper, step_context): assert result.result_type == ResultType.COMPLETED +class TestRemoveMicrocephOSDsStep: + def test_removes_host_osds_and_verifies_state(self, cclient, jhelper, step_context): + step = _cleanup_step(cclient, jhelper) + jhelper.run_cmd_on_machine_unit_payload.side_effect = [ + _command_result( + _configured_disks( + [ + {"osd": 5, "location": "node-1.maas"}, + {"osd": 2, "location": "node-1"}, + {"osd": 9, "location": "node-2"}, + ] + ) + ), + _command_result(_crush_tree("node-1.maas", [5])), + _command_result(), + _command_result(), + _command_result(_configured_disks([{"osd": 9, "location": "node-2"}])), + _command_result(_crush_tree()), + ] + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + result = step.run(step_context) + + assert result.result_type == ResultType.COMPLETED + assert step.unit == "microceph/0" + jhelper.get_unit_from_machine.assert_called_once_with( + "microceph", "1", "test-model" + ) + jhelper.get_leader_unit.assert_not_called() + assert [ + call.args[2] + for call in jhelper.run_cmd_on_machine_unit_payload.call_args_list + ] == [ + "microceph disk list --json", + "microceph.ceph osd tree --format json", + "microceph disk remove osd.2 --timeout 1800", + "microceph disk remove osd.5 --timeout 1800", + "microceph disk list --json", + "microceph.ceph osd tree --format json", + ] + + def test_crush_only_osd_aborts_before_removal(self, cclient, jhelper, step_context): + jhelper.get_unit_from_machine.side_effect = UnitNotFoundException( + "unit is missing" + ) + jhelper.get_leader_unit.return_value = "microceph/3" + step = _cleanup_step(cclient, jhelper) + jhelper.run_cmd_on_machine_unit_payload.side_effect = [ + _command_result(_configured_disks([{"osd": 2, "location": "node-1"}])), + _command_result(_crush_tree("node-1", [2, 7])), + ] + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + result = step.run(step_context) + + assert result.result_type == ResultType.FAILED + assert step.unit == "microceph/3" + assert jhelper.run_cmd_on_machine_unit_payload.call_count == 2 + + def test_force_does_not_bypass_safety_checks(self, cclient, jhelper, step_context): + step = _cleanup_step(cclient, jhelper, force=True) + jhelper.run_cmd_on_machine_unit_payload.side_effect = [ + _command_result(_configured_disks([{"osd": 2, "location": "node-1"}])), + _command_result(_crush_tree("node-1", [2])), + _command_result(), + _command_result(_configured_disks([])), + _command_result(_crush_tree()), + ] + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + assert step.run(step_context).result_type == ResultType.COMPLETED + command = jhelper.run_cmd_on_machine_unit_payload.call_args_list[2].args[2] + assert command == ( + "microceph disk remove osd.2 --timeout 1800 " + "--confirm-failure-domain-downgrade" + ) + + def test_command_failure_aborts_cleanup(self, cclient, jhelper, step_context): + step = _cleanup_step(cclient, jhelper) + jhelper.run_cmd_on_machine_unit_payload.side_effect = [ + _command_result(_configured_disks([{"osd": 2, "location": "node-1"}])), + _command_result(_crush_tree("node-1", [2])), + ExecFailedException("remove failed"), + ] + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + result = step.run(step_context) + + assert result.result_type == ResultType.FAILED + assert jhelper.run_cmd_on_machine_unit_payload.call_count == 3 + + def test_invalid_listing_aborts_cleanup(self, cclient, jhelper, step_context): + step = _cleanup_step(cclient, jhelper) + jhelper.run_cmd_on_machine_unit_payload.return_value = _command_result("{}") + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + result = step.run(step_context) + + assert result.result_type == ResultType.FAILED + assert jhelper.run_cmd_on_machine_unit_payload.call_count == 1 + + def test_remaining_osd_fails_verification(self, cclient, jhelper, step_context): + step = _cleanup_step(cclient, jhelper) + configured = _configured_disks([{"osd": 2, "location": "node-1"}]) + tree = _crush_tree("node-1", [2]) + jhelper.run_cmd_on_machine_unit_payload.side_effect = [ + _command_result(configured), + _command_result(tree), + _command_result(), + _command_result(configured), + _command_result(tree), + ] + + assert step.is_skip(step_context).result_type == ResultType.COMPLETED + result = step.run(step_context) + + assert result.result_type == ResultType.FAILED + + class TestSetCephMgrPoolSizeStep: def test_is_skip(self, cclient, jhelper, step_context): cclient.cluster.list_nodes_by_role.return_value = [] From b8bf90328b603f49e6da198622c47b940349dfe1 Mon Sep 17 00:00:00 2001 From: Seyeong Kim Date: Tue, 8 Sep 2026 16:34:05 +0900 Subject: [PATCH 5/7] fix(maas): honor force during Cinder cleanup Forced node removal still aborts on Cinder client errors because the Cinder cleanup steps do not receive the force flag. Pass it to both steps so their client errors follow Nova and Neutron's force behavior, while confirmed remaining service records still fail cleanup. Signed-off-by: Seyeong Kim --- .../sunbeam/provider/maas/commands.py | 2 + sunbeam-python/sunbeam/steps/cinder_volume.py | 28 ++++++-- .../unit/sunbeam/provider/maas/test_maas.py | 2 + .../unit/sunbeam/steps/test_cinder_volume.py | 72 ++++++++++++++----- 4 files changed, 81 insertions(+), 23 deletions(-) diff --git a/sunbeam-python/sunbeam/provider/maas/commands.py b/sunbeam-python/sunbeam/provider/maas/commands.py index 028cb66d7..5c188b6d4 100644 --- a/sunbeam-python/sunbeam/provider/maas/commands.py +++ b/sunbeam-python/sunbeam/provider/maas/commands.py @@ -1718,6 +1718,7 @@ def remove_node(ctx: click.Context, name: str, force: bool, show_hints: bool) -> deployment, machine["hostname"], machine["fqdn"], + force=force, ), RemoveCinderVolumeUnitsStep( client, name, jhelper, deployment.openstack_machines_model @@ -1727,6 +1728,7 @@ def remove_node(ctx: click.Context, name: str, force: bool, show_hints: bool) -> deployment, machine["hostname"], machine["fqdn"], + force=force, ), RemoveMicrocephOSDsStep( client, diff --git a/sunbeam-python/sunbeam/steps/cinder_volume.py b/sunbeam-python/sunbeam/steps/cinder_volume.py index 38abb2ec5..240c668f5 100644 --- a/sunbeam-python/sunbeam/steps/cinder_volume.py +++ b/sunbeam-python/sunbeam/steps/cinder_volume.py @@ -5,6 +5,8 @@ import typing from typing import Any +import click + import sunbeam.steps.microceph as microceph from sunbeam import versions from sunbeam.clusterd.client import Client @@ -287,15 +289,24 @@ def __init__( deployment: Deployment, hostname: str, fqdn: str, + force: bool = False, ): super().__init__(name, description) self.jhelper = jhelper self.deployment = deployment self.hostname = hostname self.fqdn = fqdn + self.force = force self.connection: Any | None = None self.services: list[Any] = [] + def _client_error(self, error: Exception) -> Result: + """Skip OpenStack client errors in force mode.""" + if self.force: + LOG.warning("Force mode set, skipping %s: %s", self.name, error) + return Result(ResultType.SKIPPED, str(error)) + return Result(ResultType.FAILED, str(error)) + def _discover_services(self) -> Result: """Discover Cinder service records for the target node.""" try: @@ -307,11 +318,12 @@ def _discover_services(self) -> Result: LOG.debug("OpenStack control plane not deployed, skipping: %r", e) return Result(ResultType.SKIPPED) except ( + click.ClickException, openstack.exceptions.SDKException, keystoneauth_exceptions.ClientException, ) as e: LOG.warning("Failed to discover Cinder volume services: %r", e) - return Result(ResultType.FAILED, str(e)) + return self._client_error(e) if not self.services: return Result(ResultType.SKIPPED) @@ -333,6 +345,7 @@ def __init__( deployment: Deployment, hostname: str, fqdn: str, + force: bool = False, ): super().__init__( "Disable Cinder Volume services", @@ -341,6 +354,7 @@ def __init__( deployment, hostname, fqdn, + force=force, ) def is_skip(self, context: StepContext) -> Result: @@ -364,7 +378,7 @@ def run(self, context: StepContext) -> Result: keystoneauth_exceptions.ClientException, ) as e: LOG.warning("Failed to disable Cinder volume service: %r", e) - return Result(ResultType.FAILED, str(e)) + return self._client_error(e) return Result(ResultType.COMPLETED) @@ -378,6 +392,7 @@ def __init__( deployment: Deployment, hostname: str, fqdn: str, + force: bool = False, ): super().__init__( "Remove Cinder Volume services", @@ -386,6 +401,7 @@ def __init__( deployment, hostname, fqdn, + force=force, ) def is_skip(self, context: StepContext) -> Result: @@ -458,13 +474,15 @@ def run(self, context: StepContext) -> Result: return Result(ResultType.COMPLETED) raise SunbeamException("Cinder service records remain after removal") + except SunbeamException as e: + LOG.warning("Cinder volume service cleanup did not complete: %r", e) + return Result(ResultType.FAILED, str(e)) except ( - SunbeamException, openstack.exceptions.SDKException, keystoneauth_exceptions.ClientException, ) as e: - LOG.warning("Failed to remove Cinder volume services: %r", e) - return Result(ResultType.FAILED, str(e)) + LOG.warning("OpenStack API error during Cinder volume cleanup: %r", e) + return self._client_error(e) class CheckCinderVolumeDistributionStep(BaseStep): diff --git a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py index 6c24a4c79..5454f1161 100644 --- a/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py +++ b/sunbeam-python/tests/unit/sunbeam/provider/maas/test_maas.py @@ -2594,6 +2594,8 @@ def test_remove_orders_cleanup_steps( if isinstance(step, RemoveCinderVolumeServicesStep) ) assert disable_index < cinder_units_index < cinder_services_index + assert plan[disable_index].force is force + assert plan[cinder_services_index].force is force osd_index = next( i diff --git a/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py b/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py index 35242628e..1f444b1fd 100644 --- a/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py +++ b/sunbeam-python/tests/unit/sunbeam/steps/test_cinder_volume.py @@ -3,9 +3,10 @@ from unittest.mock import MagicMock, Mock, patch +import click import pytest -from keystoneauth1.exceptions.catalog import EndpointNotFound from keystoneauth1.exceptions.connection import ConnectFailure +from openstack.exceptions import ForbiddenException from sunbeam.core.common import ResultType from sunbeam.core.juju import ( @@ -432,26 +433,44 @@ def test_disable_enabled_services( services[0], reason="Removing node from cluster" ) - def test_discovery_fails_on_keystoneauth_error( + @pytest.mark.parametrize("force", [False, True]) + @pytest.mark.parametrize( + "error", + [ + ConnectFailure("connection refused"), + ForbiddenException(), + click.ClickException("Unable to get keystone leader"), + ], + ) + def test_discovery_client_error( self, mocker, basic_jhelper, basic_deployment, step_context, + force, + error, ): mocker.patch( "sunbeam.steps.cinder_volume.get_admin_connection", - side_effect=EndpointNotFound("volume endpoint missing"), + side_effect=error, ) step = DisableCinderVolumeServicesStep( - basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas", force=force ) result = step.is_skip(step_context) - assert result.result_type == ResultType.FAILED + assert result.result_type == ( + ResultType.SKIPPED if force else ResultType.FAILED + ) + assert result.message == str(error) - def test_disable_fails_on_keystoneauth_error( + @pytest.mark.parametrize("force", [False, True]) + @pytest.mark.parametrize( + "error", [ConnectFailure("volume endpoint unavailable"), ForbiddenException()] + ) + def test_disable_client_error( self, mocker, basic_jhelper, @@ -459,10 +478,10 @@ def test_disable_fails_on_keystoneauth_error( connection, services, step_context, + force, + error, ): - connection.block_storage.disable_service.side_effect = ConnectFailure( - "volume endpoint unavailable" - ) + connection.block_storage.disable_service.side_effect = error mocker.patch( "sunbeam.steps.cinder_volume.get_admin_connection", return_value=connection, @@ -472,11 +491,12 @@ def test_disable_fails_on_keystoneauth_error( return_value=services, ) step = DisableCinderVolumeServicesStep( - basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas", force=force ) assert step.is_skip(step_context).result_type == ResultType.COMPLETED - assert step.run(step_context).result_type == ResultType.FAILED + expected = ResultType.SKIPPED if force else ResultType.FAILED + assert step.run(step_context).result_type == expected def test_remove_prefers_healthy_nonleader_when_leader_unhealthy( self, @@ -595,7 +615,11 @@ def test_remove_falls_back_when_first_unit_disappears( for call in basic_jhelper.run_cmd_on_unit_payload.call_args_list ] == ["cinder/0", "cinder/1"] - def test_remove_fails_on_keystoneauth_error( + @pytest.mark.parametrize("force", [False, True]) + @pytest.mark.parametrize( + "error", [ConnectFailure("volume endpoint unavailable"), ForbiddenException()] + ) + def test_remove_verification_client_error( self, mocker, basic_jhelper, @@ -603,6 +627,8 @@ def test_remove_fails_on_keystoneauth_error( connection, services, step_context, + force, + error, ): unit = Mock( workload_status=Mock(current="active"), @@ -619,16 +645,18 @@ def test_remove_fails_on_keystoneauth_error( "sunbeam.steps.cinder_volume.get_cinder_volume_services", side_effect=[ services, - ConnectFailure("volume endpoint unavailable"), + error, ], ) step = RemoveCinderVolumeServicesStep( - basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas", force=force ) assert step.is_skip(step_context).result_type == ResultType.COMPLETED - assert step.run(step_context).result_type == ResultType.FAILED + expected = ResultType.SKIPPED if force else ResultType.FAILED + assert step.run(step_context).result_type == expected + @pytest.mark.parametrize("force", [False, True]) def test_remove_fails_without_healthy_units( self, mocker, @@ -637,6 +665,7 @@ def test_remove_fails_without_healthy_units( connection, services, step_context, + force, ): unhealthy = Mock( workload_status=Mock(current="blocked"), @@ -653,13 +682,15 @@ def test_remove_fails_without_healthy_units( return_value=services, ) step = RemoveCinderVolumeServicesStep( - basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas", force=force ) assert step.is_skip(step_context).result_type == ResultType.COMPLETED assert step.run(step_context).result_type == ResultType.FAILED basic_jhelper.run_cmd_on_unit_payload.assert_not_called() + @pytest.mark.parametrize("force", [False, True]) + @pytest.mark.parametrize("return_code", [0, 1]) def test_remove_fails_after_healthy_candidates_leave_records( self, mocker, @@ -668,7 +699,10 @@ def test_remove_fails_after_healthy_candidates_leave_records( connection, services, step_context, + force, + return_code, ): + services = services[:1] healthy_units = { name: Mock( workload_status=Mock(current="active"), @@ -678,7 +712,9 @@ def test_remove_fails_after_healthy_candidates_leave_records( for name in ("cinder/0", "cinder/1") } basic_jhelper.get_application.return_value = Mock(units=healthy_units) - basic_jhelper.run_cmd_on_unit_payload.return_value = {"return-code": 1} + basic_jhelper.run_cmd_on_unit_payload.return_value = { + "return-code": return_code + } mocker.patch( "sunbeam.steps.cinder_volume.get_admin_connection", return_value=connection, @@ -688,7 +724,7 @@ def test_remove_fails_after_healthy_candidates_leave_records( return_value=services, ) step = RemoveCinderVolumeServicesStep( - basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas" + basic_jhelper, basic_deployment, "cloud-4", "cloud-4.maas", force=force ) assert step.is_skip(step_context).result_type == ResultType.COMPLETED From bfb92822d5f2a2fc94d817751cfd4401a53456a1 Mon Sep 17 00:00:00 2001 From: Seyeong Kim Date: Sat, 19 Sep 2026 08:43:34 +0900 Subject: [PATCH 6/7] fix(maas): remove empty CRUSH hosts MicroCeph disk removal can leave an empty CRUSH host bucket, so a removed node still appears in Ceph's placement tree. Remove matching short-name and FQDN buckets after their OSDs are gone, and verify the buckets disappear even when retrying with no disks left. Signed-off-by: Seyeong Kim --- sunbeam-python/sunbeam/steps/microceph.py | 49 +++++++----- .../unit/sunbeam/steps/test_microceph.py | 80 +++++++++++++++++-- 2 files changed, 100 insertions(+), 29 deletions(-) diff --git a/sunbeam-python/sunbeam/steps/microceph.py b/sunbeam-python/sunbeam/steps/microceph.py index 2c504235b..bc722955d 100644 --- a/sunbeam-python/sunbeam/steps/microceph.py +++ b/sunbeam-python/sunbeam/steps/microceph.py @@ -319,26 +319,23 @@ def _list_configured_osd_ids(self) -> list[int]: f"Failed to parse configured disk listing: {e}" ) from e - def _list_crush_osd_ids(self) -> list[int]: - """Return OSD IDs under the target CRUSH host.""" + def _list_crush_hosts(self) -> dict[str, list[int]]: + """Return matching CRUSH hosts and their OSD IDs, including empty hosts.""" try: nodes = json.loads( self._run_command("microceph.ceph osd tree --format json") )["nodes"] - return sorted( - { - osd_id - for node in nodes - if node["type"] == "host" and node["name"] in self._hostnames - for osd_id in node["children"] - } - ) + return { + node["name"]: node["children"] + for node in nodes + if node["type"] == "host" and node["name"] in self._hostnames + } except (json.JSONDecodeError, KeyError, TypeError) as e: raise SunbeamException(f"Failed to parse CRUSH tree: {e}") from e - def _list_target_osds(self) -> tuple[list[int], list[int]]: + def _list_target_osds(self) -> tuple[list[int], dict[str, list[int]]]: """Read both target OSD sources before changing either source.""" - return self._list_configured_osd_ids(), self._list_crush_osd_ids() + return self._list_configured_osd_ids(), self._list_crush_hosts() def run(self, context: StepContext) -> Result: """Remove DB-backed OSDs and verify both MicroCeph and CRUSH state.""" @@ -348,8 +345,9 @@ def run(self, context: StepContext) -> Result: return preparation try: - configured_osds, crush_osds = self._list_target_osds() - crush_only_osds = sorted(set(crush_osds) - set(configured_osds)) + configured_osds, crush_hosts = self._list_target_osds() + crush_osds = {osd for children in crush_hosts.values() for osd in children} + crush_only_osds = sorted(crush_osds - set(configured_osds)) if crush_only_osds: return Result( ResultType.FAILED, @@ -365,19 +363,26 @@ def run(self, context: StepContext) -> Result: command += " --confirm-failure-domain-downgrade" self._run_command(command) - if not configured_osds: - return Result(ResultType.COMPLETED) - - remaining_configured, remaining_crush = self._list_target_osds() - if remaining_configured: + if configured_osds: + remaining_configured, crush_hosts = self._list_target_osds() + if remaining_configured: + return Result( + ResultType.FAILED, + f"Configured OSDs remain for {self.node}: " + f"{remaining_configured}", + ) + if any(crush_hosts.values()): return Result( ResultType.FAILED, - f"Configured OSDs remain for {self.node}: {remaining_configured}", + f"CRUSH OSDs remain for {self.node}: {crush_hosts}", ) - if remaining_crush: + + for hostname in crush_hosts: + self._run_command(f"microceph.ceph osd crush remove {hostname}") + if crush_hosts and (remaining_hosts := self._list_crush_hosts()): return Result( ResultType.FAILED, - f"CRUSH OSDs remain for {self.node}: {remaining_crush}", + f"CRUSH hosts remain for {self.node}: {sorted(remaining_hosts)}", ) except SunbeamException as e: LOG.debug("Failed to clean up MicroCeph OSDs", exc_info=True) diff --git a/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py b/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py index d1c23b5db..64a74d9aa 100644 --- a/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py +++ b/sunbeam-python/tests/unit/sunbeam/steps/test_microceph.py @@ -4,6 +4,8 @@ import json from unittest.mock import Mock +import pytest + from sunbeam.core.common import ResultType from sunbeam.core.juju import ( ActionFailedException, @@ -142,6 +144,8 @@ def test_removes_host_osds_and_verifies_state(self, cclient, jhelper, step_conte _command_result(), _command_result(), _command_result(_configured_disks([{"osd": 9, "location": "node-2"}])), + _command_result(_crush_tree("node-1.maas")), + _command_result(), _command_result(_crush_tree()), ] @@ -164,14 +168,19 @@ def test_removes_host_osds_and_verifies_state(self, cclient, jhelper, step_conte "microceph disk remove osd.5 --timeout 1800", "microceph disk list --json", "microceph.ceph osd tree --format json", + "microceph.ceph osd crush remove node-1.maas", + "microceph.ceph osd tree --format json", ] - def test_crush_only_osd_aborts_before_removal(self, cclient, jhelper, step_context): + @pytest.mark.parametrize("force", [False, True]) + def test_crush_only_osd_aborts_before_removal( + self, cclient, jhelper, step_context, force + ): jhelper.get_unit_from_machine.side_effect = UnitNotFoundException( "unit is missing" ) jhelper.get_leader_unit.return_value = "microceph/3" - step = _cleanup_step(cclient, jhelper) + step = _cleanup_step(cclient, jhelper, force=force) jhelper.run_cmd_on_machine_unit_payload.side_effect = [ _command_result(_configured_disks([{"osd": 2, "location": "node-1"}])), _command_result(_crush_tree("node-1", [2, 7])), @@ -184,6 +193,42 @@ def test_crush_only_osd_aborts_before_removal(self, cclient, jhelper, step_conte assert step.unit == "microceph/3" assert jhelper.run_cmd_on_machine_unit_payload.call_count == 2 + @pytest.mark.parametrize( + "hostname, remaining_hostname, expected", + [ + ("node-1.maas", None, ResultType.COMPLETED), + ("node-1", "node-1", ResultType.FAILED), + ], + ) + def test_verifies_empty_host_cleanup_on_retry( + self, cclient, jhelper, step_context, hostname, remaining_hostname, expected + ): + step = _cleanup_step(cclient, jhelper, force=True) + other_host = {"name": "node-10", "type": "host", "children": [9]} + tree = json.loads(_crush_tree(hostname)) + tree["nodes"].append(other_host) + remaining = json.loads(_crush_tree(remaining_hostname)) + remaining["nodes"].append(other_host) + jhelper.run_cmd_on_machine_unit_payload.side_effect = [ + _command_result(_configured_disks([{"osd": 9, "location": "node-10"}])), + _command_result(json.dumps(tree)), + _command_result(), + _command_result(json.dumps(remaining)), + ] + + result = step.run(step_context) + + assert result.result_type == expected + assert [ + call.args[2] + for call in jhelper.run_cmd_on_machine_unit_payload.call_args_list + ] == [ + "microceph disk list --json", + "microceph.ceph osd tree --format json", + f"microceph.ceph osd crush remove {hostname}", + "microceph.ceph osd tree --format json", + ] + def test_force_does_not_bypass_safety_checks(self, cclient, jhelper, step_context): step = _cleanup_step(cclient, jhelper, force=True) jhelper.run_cmd_on_machine_unit_payload.side_effect = [ @@ -202,11 +247,16 @@ def test_force_does_not_bypass_safety_checks(self, cclient, jhelper, step_contex "--confirm-failure-domain-downgrade" ) - def test_command_failure_aborts_cleanup(self, cclient, jhelper, step_context): + @pytest.mark.parametrize("hostname, osds", [("node-1", [2]), ("node-1.maas", [])]) + def test_command_failure_aborts_cleanup( + self, cclient, jhelper, step_context, hostname, osds + ): step = _cleanup_step(cclient, jhelper) jhelper.run_cmd_on_machine_unit_payload.side_effect = [ - _command_result(_configured_disks([{"osd": 2, "location": "node-1"}])), - _command_result(_crush_tree("node-1", [2])), + _command_result( + _configured_disks([{"osd": osd, "location": hostname} for osd in osds]) + ), + _command_result(_crush_tree(hostname, osds)), ExecFailedException("remove failed"), ] @@ -226,7 +276,12 @@ def test_invalid_listing_aborts_cleanup(self, cclient, jhelper, step_context): assert result.result_type == ResultType.FAILED assert jhelper.run_cmd_on_machine_unit_payload.call_count == 1 - def test_remaining_osd_fails_verification(self, cclient, jhelper, step_context): + @pytest.mark.parametrize( + "remaining_disks", [[], [{"osd": 2, "location": "node-1"}]] + ) + def test_remaining_osd_fails_verification( + self, cclient, jhelper, step_context, remaining_disks + ): step = _cleanup_step(cclient, jhelper) configured = _configured_disks([{"osd": 2, "location": "node-1"}]) tree = _crush_tree("node-1", [2]) @@ -234,7 +289,7 @@ def test_remaining_osd_fails_verification(self, cclient, jhelper, step_context): _command_result(configured), _command_result(tree), _command_result(), - _command_result(configured), + _command_result(_configured_disks(remaining_disks)), _command_result(tree), ] @@ -242,6 +297,17 @@ def test_remaining_osd_fails_verification(self, cclient, jhelper, step_context): result = step.run(step_context) assert result.result_type == ResultType.FAILED + assert jhelper.run_cmd_on_machine_unit_payload.call_count == 5 + + def test_already_removed_host_is_a_noop(self, cclient, jhelper, step_context): + step = _cleanup_step(cclient, jhelper) + jhelper.run_cmd_on_machine_unit_payload.side_effect = [ + _command_result(_configured_disks([])), + _command_result(_crush_tree("node-10", [9])), + ] + + assert step.run(step_context).result_type == ResultType.COMPLETED + assert jhelper.run_cmd_on_machine_unit_payload.call_count == 2 class TestSetCephMgrPoolSizeStep: From f7d1b0731bc664c156fae67b3e4b7bdaf3f0c138 Mon Sep 17 00:00:00 2001 From: Seyeong Kim Date: Sat, 19 Sep 2026 08:43:45 +0900 Subject: [PATCH 7/7] fix(maas): refresh MicroOVN placement on reapply Reapplying the MicroOVN plan after node removal reuses the machine IDs saved in the Terraform variables, so it can try to place MicroOVN units on a Juju machine that no longer exists. Recompute the placement from the remaining cluster nodes, as the deploy step does, and set both architecture lists so legacy variables cannot restore removed machines. Closes-Bug: #2163212 Signed-off-by: Seyeong Kim --- sunbeam-python/sunbeam/steps/microovn.py | 10 +++- .../tests/unit/sunbeam/steps/test_microovn.py | 46 +++++++++++++++++++ 2 files changed, 55 insertions(+), 1 deletion(-) diff --git a/sunbeam-python/sunbeam/steps/microovn.py b/sunbeam-python/sunbeam/steps/microovn.py index f6bc18aaa..de4bb98e3 100644 --- a/sunbeam-python/sunbeam/steps/microovn.py +++ b/sunbeam-python/sunbeam/steps/microovn.py @@ -327,6 +327,15 @@ def is_skip(self, context: StepContext) -> Result: ) def run(self, context: StepContext) -> Result: """Apply terraform configuration to deploy MicroOVN.""" + machines_by_arch = self.ovn_manager.get_machines_by_architecture() + self.extra_tfvars["microovn_machine_ids_by_architecture"] = { + ovn.DEFAULT_ARCHITECTURE: [], + ovn.ARM64_ARCHITECTURE: [], + **machines_by_arch, + } + distributor_ids = self.ovn_manager.get_token_distributor_machines() + self.extra_tfvars["token_distributor_machine_ids"] = distributor_ids[:1] + # Apply Network configs everytime reapply is called network_configs = get_external_network_configs(self.client) if "charm_openstack_network_agents_config" not in self.extra_tfvars: @@ -353,7 +362,6 @@ def run(self, context: StepContext) -> Result: except TerraformException as e: return Result(ResultType.FAILED, str(e)) - machines_by_arch = self.ovn_manager.get_machines_by_architecture() apps_to_wait = _microovn_applications_to_wait(machines_by_arch) try: for application in apps_to_wait: diff --git a/sunbeam-python/tests/unit/sunbeam/steps/test_microovn.py b/sunbeam-python/tests/unit/sunbeam/steps/test_microovn.py index 4c24b6f03..fd9cc0f6c 100644 --- a/sunbeam-python/tests/unit/sunbeam/steps/test_microovn.py +++ b/sunbeam-python/tests/unit/sunbeam/steps/test_microovn.py @@ -455,6 +455,7 @@ def reapply_microovn_terraform_step( ovn.DEFAULT_ARCHITECTURE: ["1"], ovn.ARM64_ARCHITECTURE: [], } + manager.get_token_distributor_machines.return_value = ["1"] return ReapplyMicroOVNTerraformPlanStep( basic_client, basic_tfhelper, @@ -530,3 +531,48 @@ def test_run_uses_microovn_accepted_statuses( accepted_status=["active", "unknown"], timeout=1200, ) + + @pytest.mark.parametrize( + "machines,distributors", + [ + ({"amd64": ["0"]}, ["0"]), + ({"amd64": ["2", "3"], "arm64": ["4"]}, ["2", "3"]), + ({"arm64": ["4"]}, []), + ], + ) + def test_run_refreshes_machine_placement( + self, + reapply_microovn_terraform_step, + basic_tfhelper, + step_context, + mocker, + machines, + distributors, + ): + step = reapply_microovn_terraform_step + step.ovn_manager.get_machines_by_architecture.return_value = machines + step.ovn_manager.get_token_distributor_machines.return_value = distributors + step.extra_tfvars.update( + { + "microovn_machine_ids_by_architecture": { + "amd64": ["0", "1"], + "arm64": ["4", "5"], + }, + "token_distributor_machine_ids": ["1"], + } + ) + mocker.patch( + "sunbeam.steps.microovn.get_external_network_configs", return_value={} + ) + + result = step.run(step_context) + + assert result.result_type == ResultType.COMPLETED + tfvars = basic_tfhelper.update_tfvars_and_apply_tf.call_args.kwargs[ + "override_tfvars" + ] + assert tfvars["microovn_machine_ids_by_architecture"] == { + "amd64": machines.get("amd64", []), + "arm64": machines.get("arm64", []), + } + assert tfvars["token_distributor_machine_ids"] == distributors[:1]