diff --git a/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_commands.py b/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_commands.py index 243dc8ea95f..41efe797a22 100644 --- a/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_commands.py +++ b/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_commands.py @@ -6,6 +6,7 @@ import json import os import random +import re import subprocess import tempfile import time @@ -115,6 +116,28 @@ def _is_transient_operation_conflict(ex): "ProvisioningState of extension: Updating" in message ) + @staticmethod + def _is_resource_already_exists_conflict(ex): + return "already exists" in str(ex).casefold() + + @staticmethod + def _extract_cli_option(command, *option_names): + for option_name in option_names: + match = re.search(rf"{re.escape(option_name)}(?:=|\s+)(\S+)", command) + if match: + return match.group(1).strip("\"'") + return None + + @classmethod + def _build_show_command_for_already_existing_resource(cls, command): + if not re.match(r"^aks\s+create\b", command.strip()): + return None + resource_group = cls._extract_cli_option(command, "--resource-group", "-g") + name = cls._extract_cli_option(command, "--name", "-n") + if not resource_group or not name: + return None + return f"aks show --resource-group {resource_group} --name {name}" + def _execute_with_transient_conflict_retry(self, command, expect_failure): from azure.cli.testsdk.base import execute import logging @@ -127,6 +150,15 @@ def _execute_with_transient_conflict_retry(self, command, expect_failure): try: return execute(self.cli_ctx, command, expect_failure=expect_failure) except (HttpResponseError, CLIError) as ex: + if ( + not expect_failure and + attempt > 0 and + getattr(self, "_allow_retried_create_recovery", False) and + self._is_resource_already_exists_conflict(ex) + ): + show_command = self._build_show_command_for_already_existing_resource(command) + if show_command: + return execute(self.cli_ctx, show_command, expect_failure=False) if ( expect_failure or not self._is_transient_operation_conflict(ex) or @@ -144,6 +176,14 @@ def _execute_with_transient_conflict_retry(self, command, expect_failure): raise AssertionError("unreachable") + def _cmd_with_retried_create_recovery(self, command, checks=None): + previous_value = getattr(self, "_allow_retried_create_recovery", False) + self._allow_retried_create_recovery = True + try: + return self.cmd(command, checks=checks) + finally: + self._allow_retried_create_recovery = previous_value + def _refetch_settled_aks_result(self, resource_id, fallback_result): from azure.cli.testsdk.base import execute @@ -383,6 +423,38 @@ def _wait_for_cluster_update(self): checks=[self.is_empty()], ) + def _wait_for_cluster_property(self, query, expected, attempts=20, delay=30): + if not (self.is_live or self.in_recording): + return + last_value = None + for attempt in range(attempts): + last_value = self.cmd( + f'aks show --resource-group={{resource_group}} --name={{name}} ' + f'--query "{query}" -o json' + ).get_output_in_json() + if str(last_value).casefold() == str(expected).casefold(): + return + if attempt < attempts - 1: + time.sleep(delay) + raise AssertionError( + f"Cluster property '{query}' did not reach {expected!r}; last value: {last_value!r}" + ) + + def _cmd_or_skip_if_artifact_streaming_unavailable(self, command, checks=None): + try: + return self.cmd(command, checks=checks) + except Exception as ex: # pylint: disable=broad-except + message = str(ex) + if ( + "UnmarshalError" in message and + 'unknown field "artifactStreamingProfile"' in message + ): + self.skipTest( + "The stable AKS API used by Azure CLI does not currently expose " + "artifactStreamingProfile; coverage remains in aks-preview." + ) + raise + # Substrings identifying an "unsupported/unavailable" condition (as opposed to e.g. a # value/quota/permission validation error). On their own these are too generic to trigger a # skip safely, since many unrelated errors (invalid VM size, bad SKU, etc.) also contain @@ -1619,10 +1691,8 @@ def test_aks_create_default_service_with_monitoring_addon(self, resource_group, ]) # disable monitoring add-on - disable_addon_output = self.cmd('aks disable-addons -a monitoring -g {resource_group} -n {name}', checks=[ - self.check('addonProfiles.omsagent.enabled', False), - ]).get_output_in_json() - assert bool(disable_addon_output["addonProfiles"]["omsagent"]["config"]) == False + self.cmd('aks disable-addons -a monitoring -g {resource_group} -n {name}') + self._wait_for_cluster_property('addonProfiles.omsagent.enabled', False) # show again show_output = self.cmd('aks show -g {resource_group} -n {name}', checks=[ @@ -4697,7 +4767,7 @@ def test_aks_nodepool_add_with_artifact_streaming( ) # nodepool add - self.cmd( + self._cmd_or_skip_if_artifact_streaming_unavailable( "aks nodepool add --resource-group={resource_group} --cluster-name={name} --name={nodepool2_name} " "--node-vm-size={node_vm_size} " "--enable-artifact-streaming --aks-custom-headers=AKSHTTPCustomFeatures=Microsoft.ContainerService/ArtifactStreamingPreview", @@ -4756,7 +4826,7 @@ def test_aks_nodepool_update_with_artifact_streaming( ) # enable artifact streaming - self.cmd( + self._cmd_or_skip_if_artifact_streaming_unavailable( "aks nodepool update " "--resource-group={resource_group} " "--cluster-name={name} " @@ -8841,13 +8911,21 @@ def test_aks_create_with_control_plane_metrics(self, resource_group, resource_gr # the final state via ``aks show`` after the cluster settles. # Control Plane Metrics availability is still rolling out per-subscription/region; if the # service reports the feature/toggle as unsupported here, skip rather than fail the test. - self._cmd_or_skip_if_unsupported( - create_cmd, - checks=[ - self.check('provisioningState', 'Succeeded'), - ], - skip_reason="Control Plane Metrics toggle is not yet available in this subscription/region", - ) + try: + self._cmd_with_retried_create_recovery( + create_cmd, + checks=[self.check('provisioningState', 'Succeeded')], + ) + except Exception as ex: # pylint: disable=broad-except + message = str(ex).casefold() + if ( + any(marker in message for marker in self._CONTROL_PLANE_METRICS_CONTEXT_MARKERS) and + any(marker in message for marker in self._UNSUPPORTED_CONDITION_MARKERS) + ): + self.skipTest( + "Control Plane Metrics toggle is not yet available in this subscription/region" + ) + raise wait_cmd = 'aks wait --resource-group={resource_group} --name={name} --created ' \ '--interval 60 --timeout 1800' @@ -8893,9 +8971,10 @@ def test_aks_update_with_control_plane_metrics(self, resource_group, resource_gr create_cmd = 'aks create --resource-group={resource_group} --name={name} --location={location} ' \ '--ssh-key-value={ssh_key_value} --node-vm-size={node_vm_size} --enable-managed-identity ' \ '--enable-azure-monitor-metrics --azure-monitor-workspace-resource-id={amw_id} --output=json' - self.cmd(create_cmd, checks=[ - self.check('provisioningState', 'Succeeded'), - ]) + self._cmd_with_retried_create_recovery( + create_cmd, + checks=[self.check('provisioningState', 'Succeeded')], + ) # wait for AMW background setup to complete before issuing update wait_cmd = 'aks wait --resource-group={resource_group} --name={name} --updated --timeout=1800' @@ -10577,12 +10656,14 @@ def test_aks_kubenet_to_cni_overlay_migration(self, resource_group, resource_gro 'name': aks_name, 'location': cluster_location, 'k8s_version': create_version, + 'node_vm_size': 'Standard_D2s_v3', 'ssh_key_value': self.generate_ssh_keys(), }) # create create_cmd = 'aks create --resource-group={resource_group} --name={name} --location={location} ' \ '--network-plugin kubenet --ssh-key-value={ssh_key_value} --kubernetes-version {k8s_version} ' \ + '--node-vm-size {node_vm_size} ' \ '--service-cidr 172.56.0.0/16 --dns-service-ip 172.56.0.10 --pod-cidr 100.112.0.0/12 ' \ '--aks-custom-headers AKSHTTPCustomFeatures=Microsoft.ContainerService/AzureOverlayPreview' self.cmd(create_cmd, checks=[ @@ -14503,9 +14584,11 @@ def test_aks_create_acns_with_flow_logs( checks=[ self.check("provisioningState", "Succeeded"), self.check("addonProfiles.omsagent.enabled", True), - self.check("addonProfiles.omsagent.config.enableRetinaNetworkFlags", "False"), ], ) + self._wait_for_cluster_property( + "addonProfiles.omsagent.config.enableRetinaNetworkFlags", "False" + ) # update: enable high log scale mode independently via aks update self.cmd( @@ -14527,9 +14610,11 @@ def test_aks_create_acns_with_flow_logs( checks=[ self.check("provisioningState", "Succeeded"), self.check("addonProfiles.omsagent.enabled", True), - self.check("addonProfiles.omsagent.config.enableRetinaNetworkFlags", "True"), ], ) + self._wait_for_cluster_property( + "addonProfiles.omsagent.config.enableRetinaNetworkFlags", "True" + ) # delete self.cmd( diff --git a/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_provisioning_retry.py b/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_provisioning_retry.py index 85dbc866013..55d3887a4fb 100644 --- a/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_provisioning_retry.py +++ b/src/azure-cli/azure/cli/command_modules/acs/tests/latest/test_aks_provisioning_retry.py @@ -200,6 +200,192 @@ def test_live_run_waits_for_cluster_update(self): ) +class TestWaitForClusterProperty(unittest.TestCase): + + @staticmethod + def _make_instance(values, is_live=True, in_recording=False): + from azure.cli.command_modules.acs.tests.latest.test_aks_commands import ( + AzureKubernetesServiceScenarioTest, + ) + instance = object.__new__(AzureKubernetesServiceScenarioTest) + instance.is_live = is_live + instance.in_recording = in_recording + instance.cmd = MagicMock( + side_effect=[MockExecutionResult(value) for value in values] + ) + return instance + + @patch('time.sleep', return_value=None) + def test_polls_until_property_matches_case_insensitively(self, mock_sleep): + instance = self._make_instance(['true', 'False']) + + instance._wait_for_cluster_property('addonProfiles.omsagent.enabled', False) + + self.assertEqual(instance.cmd.call_count, 2) + mock_sleep.assert_called_once_with(30) + + def test_replay_does_not_issue_poll_requests(self): + instance = self._make_instance([], is_live=False) + + instance._wait_for_cluster_property('addonProfiles.omsagent.enabled', False) + + instance.cmd.assert_not_called() + + @patch('time.sleep', return_value=None) + def test_raises_when_property_never_matches(self, _mock_sleep): + instance = self._make_instance(['true', 'true']) + + with self.assertRaisesRegex(AssertionError, 'last value'): + instance._wait_for_cluster_property( + 'addonProfiles.omsagent.enabled', False, attempts=2 + ) + + +class TestAlreadyExistsConflictHandling(unittest.TestCase): + + @staticmethod + def _make_instance(): + from azure.cli.command_modules.acs.tests.latest.test_aks_commands import ( + AzureKubernetesServiceScenarioTest, + ) + instance = object.__new__(AzureKubernetesServiceScenarioTest) + instance.cli_ctx = MagicMock() + instance._allow_retried_create_recovery = False + return instance + + def test_builds_show_command_for_create(self): + instance = self._make_instance() + + command = instance._build_show_command_for_already_existing_resource( + 'aks create -g rg --name cluster --node-count 1' + ) + + self.assertEqual(command, 'aks show --resource-group rg --name cluster') + + def test_does_not_translate_non_create_command(self): + instance = self._make_instance() + + self.assertIsNone( + instance._build_show_command_for_already_existing_resource( + 'aks update -g rg -n cluster' + ) + ) + + @patch.dict(os.environ, { + 'AZURE_CLI_TEST_OPERATION_MAX_RETRIES': '3', + 'AZURE_CLI_TEST_OPERATION_BASE_DELAY': '0.01', + }) + @patch('time.sleep', return_value=None) + @patch('random.uniform', return_value=0) + @patch('azure.cli.testsdk.base.execute') + def test_already_exists_after_transient_retry_uses_show( + self, mock_execute, _mock_random, mock_sleep + ): + expected = MockExecutionResult({'provisioningState': 'Succeeded'}) + mock_execute.side_effect = [ + CLIError('Another operation is in progress.'), + CLIError("The cluster 'cluster' already exists."), + expected, + ] + instance = self._make_instance() + instance._allow_retried_create_recovery = True + + result = instance._execute_with_transient_conflict_retry( + 'aks create --resource-group rg --name cluster', False + ) + + self.assertIs(result, expected) + mock_execute.assert_called_with( + instance.cli_ctx, + 'aks show --resource-group rg --name cluster', + expect_failure=False, + ) + mock_sleep.assert_called_once() + + @patch.dict(os.environ, {'AZURE_CLI_TEST_OPERATION_MAX_RETRIES': '3'}) + @patch('time.sleep', return_value=None) + @patch('azure.cli.testsdk.base.execute') + def test_first_attempt_already_exists_still_raises( + self, mock_execute, mock_sleep + ): + mock_execute.side_effect = CLIError("The cluster 'cluster' already exists.") + + with self.assertRaisesRegex(CLIError, 'already exists'): + self._make_instance()._execute_with_transient_conflict_retry( + 'aks create --resource-group rg --name cluster', False + ) + + mock_execute.assert_called_once() + mock_sleep.assert_not_called() + + @patch.dict(os.environ, { + 'AZURE_CLI_TEST_OPERATION_MAX_RETRIES': '3', + 'AZURE_CLI_TEST_OPERATION_BASE_DELAY': '0.01', + }) + @patch('time.sleep', return_value=None) + @patch('random.uniform', return_value=0) + @patch('azure.cli.testsdk.base.execute') + def test_already_exists_without_explicit_recovery_still_raises( + self, mock_execute, _mock_random, _mock_sleep + ): + mock_execute.side_effect = [ + CLIError('Another operation is in progress.'), + CLIError("The cluster 'cluster' already exists."), + ] + + with self.assertRaisesRegex(CLIError, 'already exists'): + self._make_instance()._execute_with_transient_conflict_retry( + 'aks create --resource-group rg --name cluster', False + ) + + def test_recovery_wrapper_restores_previous_state(self): + instance = self._make_instance() + instance.cmd = MagicMock(return_value='result') + + result = instance._cmd_with_retried_create_recovery( + 'aks create -g rg -n cluster', checks=['check'] + ) + + self.assertEqual(result, 'result') + instance.cmd.assert_called_once_with( + 'aks create -g rg -n cluster', checks=['check'] + ) + self.assertFalse(instance._allow_retried_create_recovery) + + +class TestArtifactStreamingStableApiHandling(unittest.TestCase): + + def test_skips_only_exact_stable_api_unmarshal_error(self): + from azure.cli.command_modules.acs.tests.latest.test_aks_commands import ( + AzureKubernetesServiceScenarioTest, + ) + instance = object.__new__(AzureKubernetesServiceScenarioTest) + instance.cmd = MagicMock(side_effect=CLIError( + 'UnmarshalError: json: unknown field "artifactStreamingProfile"' + )) + instance.skipTest = MagicMock(side_effect=unittest.SkipTest('unsupported')) + + with self.assertRaises(unittest.SkipTest): + instance._cmd_or_skip_if_artifact_streaming_unavailable('aks nodepool update') + + def test_unrelated_unmarshal_error_propagates(self): + from azure.cli.command_modules.acs.tests.latest.test_aks_commands import ( + AzureKubernetesServiceScenarioTest, + ) + instance = object.__new__(AzureKubernetesServiceScenarioTest) + instance.cmd = MagicMock( + side_effect=CLIError('UnmarshalError: json: unknown field "other"') + ) + instance.skipTest = MagicMock() + + with self.assertRaisesRegex(CLIError, 'unknown field'): + instance._cmd_or_skip_if_artifact_streaming_unavailable( + 'aks nodepool update' + ) + + instance.skipTest.assert_not_called() + + class TestCmdWithRetry(unittest.TestCase): def _make_instance(self):