This is an automated email from the ASF dual-hosted git repository.

potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 0d5c54b6be9 Retry EKS resource deletion on transient 
ResourceInUseException (#70260)
0d5c54b6be9 is described below

commit 0d5c54b6be9823f9b4ba149ac06b1b9aeaa184b8
Author: Ramit Kataria <[email protected]>
AuthorDate: Fri Jul 31 18:28:02 2026 -0700

    Retry EKS resource deletion on transient ResourceInUseException (#70260)
    
    Deleting an EKS cluster, nodegroup, or Fargate profile can transiently
    fail with ResourceInUseException while the resource is still settling
    from a prior operation. The delete operators and the deferrable delete
    trigger surfaced that error immediately, failing the task and leaving
    the resource orphaned.
    
    Both the synchronous operators and the async trigger now retry these
    deletes with exponential backoff until the resource settles, giving up
    and re-raising only after a bounded timeout so a genuinely wedged
    resource still fails.
    
    With the operators handling this directly, the manual task-level retries
    in the EKS system test DAGs are no longer needed and are removed.
    
    Co-authored-by: Niko Oliveira <[email protected]>
---
 docs/spelling_wordlist.txt                         |   1 +
 .../airflow/providers/amazon/aws/operators/eks.py  |  44 +++++---
 .../airflow/providers/amazon/aws/triggers/eks.py   |  30 ++++--
 .../airflow/providers/amazon/aws/utils/__init__.py |  42 +++++++-
 .../aws/example_eks_with_fargate_in_one_step.py    |   5 -
 .../amazon/aws/example_eks_with_fargate_profile.py |   5 -
 .../aws/example_eks_with_nodegroup_in_one_step.py  |   4 -
 .../amazon/aws/example_eks_with_nodegroups.py      |   4 -
 .../tests/unit/amazon/aws/operators/test_eks.py    | 119 +++++++++++++++++++++
 .../tests/unit/amazon/aws/triggers/test_eks.py     |  44 ++++++++
 .../tests/unit/amazon/aws/utils/test_utils.py      |  66 ++++++++++++
 11 files changed, 324 insertions(+), 40 deletions(-)

diff --git a/docs/spelling_wordlist.txt b/docs/spelling_wordlist.txt
index 36baf354aa6..73212e3d679 100644
--- a/docs/spelling_wordlist.txt
+++ b/docs/spelling_wordlist.txt
@@ -1619,6 +1619,7 @@ StatsD
 statsd
 stderr
 stdin
+stdlib
 stdout
 StorageClass
 storages
diff --git a/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py 
b/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py
index 7ab58406ffb..c2efc42f8a5 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py
@@ -28,6 +28,7 @@ from collections.abc import Sequence
 from datetime import timedelta
 from typing import TYPE_CHECKING, Any, cast
 
+import tenacity
 from botocore.exceptions import ClientError, WaiterError
 
 from airflow.exceptions import AirflowProviderDeprecationWarning
@@ -42,7 +43,10 @@ from airflow.providers.amazon.aws.triggers.eks import (
     EksDeleteNodegroupTrigger,
     EksPodTrigger,
 )
-from airflow.providers.amazon.aws.utils import validate_execute_complete_event
+from airflow.providers.amazon.aws.utils import (
+    build_resource_in_use_retry_args,
+    validate_execute_complete_event,
+)
 from airflow.providers.amazon.aws.utils.mixins import aws_template_fields
 from airflow.providers.amazon.aws.utils.waiter_with_logging import wait
 from airflow.providers.cncf.kubernetes.utils.pod_manager import OnFinishAction
@@ -139,10 +143,12 @@ def _create_compute(
             # delete_nodegroup_on_failure defaults to True to prevent orphaned 
nodegroups.
             if delete_nodegroup_on_failure:
                 try:
-                    eks_hook.delete_nodegroup(
-                        clusterName=cluster_name,
-                        nodegroupName=nodegroup_name,
-                    )
+                    for attempt in 
tenacity.Retrying(**build_resource_in_use_retry_args(log)):
+                        with attempt:
+                            eks_hook.delete_nodegroup(
+                                clusterName=cluster_name,
+                                nodegroupName=nodegroup_name,
+                            )
                     log.info(
                         "Issued delete request for nodegroup '%s' in cluster 
'%s' after failure.",
                         nodegroup_name,
@@ -809,7 +815,9 @@ class EksDeleteClusterOperator(AwsBaseOperator[EksHook]):
             self.delete_any_nodegroups()
             self.delete_any_fargate_profiles()
 
-        self.hook.delete_cluster(name=self.cluster_name)
+        for attempt in 
tenacity.Retrying(**build_resource_in_use_retry_args(self.log)):
+            with attempt:
+                self.hook.delete_cluster(name=self.cluster_name)
 
         if self.wait_for_completion:
             self.log.info("Waiting for cluster to delete.  This will take some 
time.")
@@ -825,8 +833,11 @@ class EksDeleteClusterOperator(AwsBaseOperator[EksHook]):
         nodegroups = self.hook.list_nodegroups(clusterName=self.cluster_name)
         if nodegroups:
             
self.log.info(CAN_NOT_DELETE_MSG.format(compute=NODEGROUP_FULL_NAME, 
count=len(nodegroups)))
+            retry_args = build_resource_in_use_retry_args(self.log)
             for group in nodegroups:
-                self.hook.delete_nodegroup(clusterName=self.cluster_name, 
nodegroupName=group)
+                for attempt in tenacity.Retrying(**retry_args):
+                    with attempt:
+                        
self.hook.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=group)
             # Note this is a custom waiter so we're using hook.get_waiter(), 
not hook.conn.get_waiter().
             self.log.info("Waiting for all nodegroups to delete.  This will 
take some time.")
             
self.hook.get_waiter("all_nodegroups_deleted").wait(clusterName=self.cluster_name)
@@ -843,11 +854,16 @@ class EksDeleteClusterOperator(AwsBaseOperator[EksHook]):
         if fargate_profiles:
             self.log.info(CAN_NOT_DELETE_MSG.format(compute=FARGATE_FULL_NAME, 
count=len(fargate_profiles)))
             self.log.info("Waiting for Fargate profiles to delete.  This will 
take some time.")
+            retry_args = build_resource_in_use_retry_args(self.log)
             for profile in fargate_profiles:
                 # The API will return a (cluster) ResourceInUseException if 
you try
                 # to delete Fargate profiles in parallel the way we can with 
nodegroups,
                 # so each must be deleted sequentially
-                
self.hook.delete_fargate_profile(clusterName=self.cluster_name, 
fargateProfileName=profile)
+                for attempt in tenacity.Retrying(**retry_args):
+                    with attempt:
+                        self.hook.delete_fargate_profile(
+                            clusterName=self.cluster_name, 
fargateProfileName=profile
+                        )
                 self.hook.conn.get_waiter("fargate_profile_deleted").wait(
                     clusterName=self.cluster_name, fargateProfileName=profile
                 )
@@ -921,7 +937,9 @@ class EksDeleteNodegroupOperator(AwsBaseOperator[EksHook]):
         super().__init__(**kwargs)
 
     def execute(self, context: Context):
-        self.hook.delete_nodegroup(clusterName=self.cluster_name, 
nodegroupName=self.nodegroup_name)
+        for attempt in 
tenacity.Retrying(**build_resource_in_use_retry_args(self.log)):
+            with attempt:
+                self.hook.delete_nodegroup(clusterName=self.cluster_name, 
nodegroupName=self.nodegroup_name)
         if self.deferrable:
             self.defer(
                 trigger=EksDeleteNodegroupTrigger(
@@ -1009,9 +1027,11 @@ class 
EksDeleteFargateProfileOperator(AwsBaseOperator[EksHook]):
         super().__init__(**kwargs)
 
     def execute(self, context: Context):
-        self.hook.delete_fargate_profile(
-            clusterName=self.cluster_name, 
fargateProfileName=self.fargate_profile_name
-        )
+        for attempt in 
tenacity.Retrying(**build_resource_in_use_retry_args(self.log)):
+            with attempt:
+                self.hook.delete_fargate_profile(
+                    clusterName=self.cluster_name, 
fargateProfileName=self.fargate_profile_name
+                )
         if self.deferrable:
             self.defer(
                 trigger=EksDeleteFargateProfileTrigger(
diff --git a/providers/amazon/src/airflow/providers/amazon/aws/triggers/eks.py 
b/providers/amazon/src/airflow/providers/amazon/aws/triggers/eks.py
index a386f96da1b..18535d2344a 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/triggers/eks.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/triggers/eks.py
@@ -19,10 +19,12 @@ from __future__ import annotations
 import datetime
 from typing import TYPE_CHECKING, Any
 
+import tenacity
 from botocore.exceptions import ClientError
 
 from airflow.providers.amazon.aws.hooks.eks import EksHook
 from airflow.providers.amazon.aws.triggers.base import AwsBaseWaiterTrigger
+from airflow.providers.amazon.aws.utils import build_resource_in_use_retry_args
 from airflow.providers.amazon.aws.utils.waiter_with_logging import async_wait
 from airflow.providers.cncf.kubernetes.triggers.pod import KubernetesPodTrigger
 from airflow.providers.common.compat.sdk import AirflowException
@@ -275,13 +277,15 @@ class EksDeleteClusterTrigger(AwsBaseWaiterTrigger):
             if self.force_delete_compute:
                 await self.delete_any_nodegroups(client=client)
                 await self.delete_any_fargate_profiles(client=client)
-            try:
-                await client.delete_cluster(name=self.cluster_name)
-            except ClientError as ex:
-                if ex.response.get("Error").get("Code") == 
"ResourceNotFoundException":
-                    pass
-                else:
-                    raise
+            async for attempt in 
tenacity.AsyncRetrying(**build_resource_in_use_retry_args(self.log)):
+                with attempt:
+                    try:
+                        await client.delete_cluster(name=self.cluster_name)
+                    except ClientError as ex:
+                        # The cluster is already gone — nothing to wait on, so 
stop retrying.
+                        if ex.response.get("Error", {}).get("Code") == 
"ResourceNotFoundException":
+                            break
+                        raise
             await async_wait(
                 waiter=waiter,
                 waiter_delay=int(self.waiter_delay),
@@ -305,8 +309,11 @@ class EksDeleteClusterTrigger(AwsBaseWaiterTrigger):
         if nodegroups.get("nodegroups", None):
             self.log.info("Deleting nodegroups")
             waiter = self.hook().get_waiter("all_nodegroups_deleted", 
deferrable=True, client=client)
+            retry_args = build_resource_in_use_retry_args(self.log)
             for group in nodegroups["nodegroups"]:
-                await client.delete_nodegroup(clusterName=self.cluster_name, 
nodegroupName=group)
+                async for attempt in tenacity.AsyncRetrying(**retry_args):
+                    with attempt:
+                        await 
client.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=group)
             await async_wait(
                 waiter=waiter,
                 waiter_delay=int(self.waiter_delay),
@@ -330,8 +337,13 @@ class EksDeleteClusterTrigger(AwsBaseWaiterTrigger):
         fargate_profiles = await 
client.list_fargate_profiles(clusterName=self.cluster_name)
         if fargate_profiles.get("fargateProfileNames"):
             self.log.info("Waiting for Fargate profiles to delete.  This will 
take some time.")
+            retry_args = build_resource_in_use_retry_args(self.log)
             for profile in fargate_profiles["fargateProfileNames"]:
-                await 
client.delete_fargate_profile(clusterName=self.cluster_name, 
fargateProfileName=profile)
+                async for attempt in tenacity.AsyncRetrying(**retry_args):
+                    with attempt:
+                        await client.delete_fargate_profile(
+                            clusterName=self.cluster_name, 
fargateProfileName=profile
+                        )
                 await async_wait(
                     waiter=client.get_waiter("fargate_profile_deleted"),
                     waiter_delay=int(self.waiter_delay),
diff --git 
a/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py 
b/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py
index 59d93015802..b4bf0d4f7c3 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/utils/__init__.py
@@ -22,14 +22,54 @@ import re
 from datetime import datetime, timezone
 from enum import Enum
 from importlib import metadata
-from typing import Any
+from typing import TYPE_CHECKING, Any
+
+import tenacity
+from botocore.exceptions import ClientError
 
 from airflow.providers.common.compat.sdk import AirflowException
 from airflow.utils.helpers import prune_dict
 from airflow.version import version
 
+if TYPE_CHECKING:
+    from airflow.sdk.types import Logger
+
 log = logging.getLogger(__name__)
 
+# AWS briefly rejects a delete call with ResourceInUseException while the 
target resource is still
+# settling from a prior operation (e.g. EKS finalizing a nodegroup removal 
before the cluster can be
+# deleted). Retry with exponential backoff (1s, 2s, 4s, ... capped at 
RESOURCE_IN_USE_RETRY_MAX_WAIT
+# per wait) until RESOURCE_IN_USE_RETRY_TIMEOUT elapses, then give up and 
re-raise. This rides out the
+# settling window without hanging a genuinely wedged resource for long.
+RESOURCE_IN_USE_RETRY_TIMEOUT = 300
+RESOURCE_IN_USE_RETRY_MAX_WAIT = 60
+
+
+def is_resource_in_use_error(exception: BaseException) -> bool:
+    """Return True if the exception is a transient AWS 
``ResourceInUseException``."""
+    return (
+        isinstance(exception, ClientError)
+        and exception.response.get("Error", {}).get("Code") == 
"ResourceInUseException"
+    )
+
+
+def build_resource_in_use_retry_args(logger: Logger | logging.Logger) -> 
dict[str, Any]:
+    """
+    Build tenacity arguments for retrying a call on a transient 
``ResourceInUseException``.
+
+    Shared by synchronous operators (``tenacity.Retrying``) and deferrable 
triggers
+    (``tenacity.AsyncRetrying``) so both back off identically. 
``reraise=True`` keeps the
+    original error as the task failure once the retry timeout is exhausted. 
Accepts either an
+    Airflow structlog logger (``self.log``) or a stdlib ``logging.Logger`` 
(module-level helpers).
+    """
+    return {
+        "retry": tenacity.retry_if_exception(is_resource_in_use_error),
+        "wait": tenacity.wait_exponential(max=RESOURCE_IN_USE_RETRY_MAX_WAIT),
+        "stop": tenacity.stop_after_delay(RESOURCE_IN_USE_RETRY_TIMEOUT),
+        "before_sleep": tenacity.before_sleep_log(logger, logging.WARNING),
+        "reraise": True,
+    }
+
 
 def trim_none_values(obj: dict):
     return prune_dict(obj)
diff --git 
a/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_in_one_step.py
 
b/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_in_one_step.py
index e22fd0ec2bb..44a8df57e54 100644
--- 
a/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_in_one_step.py
+++ 
b/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_in_one_step.py
@@ -18,8 +18,6 @@ from __future__ import annotations
 
 from datetime import datetime
 
-from pendulum import duration
-
 from airflow.providers.amazon.aws.hooks.eks import ClusterStates, 
FargateProfileStates
 from airflow.providers.amazon.aws.operators.eks import (
     EksCreateClusterOperator,
@@ -132,9 +130,6 @@ with DAG(
         trigger_rule=TriggerRule.ALL_DONE,
         cluster_name=cluster_name,
         force_delete_compute=True,
-        retries=4,
-        retry_delay=duration(seconds=30),
-        retry_exponential_backoff=True,
     )
 
     await_delete_cluster = EksClusterStateSensor(
diff --git 
a/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_profile.py 
b/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_profile.py
index 8d2a3987f26..1da5b9c8e01 100644
--- 
a/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_profile.py
+++ 
b/providers/amazon/tests/system/amazon/aws/example_eks_with_fargate_profile.py
@@ -18,8 +18,6 @@ from __future__ import annotations
 
 from datetime import datetime
 
-from pendulum import duration
-
 from airflow.providers.amazon.aws.hooks.eks import ClusterStates, 
FargateProfileStates
 from airflow.providers.amazon.aws.operators.eks import (
     EksCreateClusterOperator,
@@ -146,9 +144,6 @@ with DAG(
         task_id="delete_eks_fargate_profile",
         cluster_name=cluster_name,
         fargate_profile_name=fargate_profile_name,
-        retries=4,
-        retry_delay=duration(seconds=30),
-        retry_exponential_backoff=True,
     )
     # [END howto_operator_eks_delete_fargate_profile]
     delete_fargate_profile.trigger_rule = TriggerRule.ALL_DONE
diff --git 
a/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroup_in_one_step.py
 
b/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroup_in_one_step.py
index 2cdc633f4b7..47a27687caf 100644
--- 
a/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroup_in_one_step.py
+++ 
b/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroup_in_one_step.py
@@ -19,7 +19,6 @@ from __future__ import annotations
 from datetime import datetime
 
 import boto3
-from pendulum import duration
 
 from airflow.providers.amazon.aws.hooks.eks import ClusterStates, 
NodegroupStates
 from airflow.providers.amazon.aws.operators.eks import (
@@ -149,9 +148,6 @@ with DAG(
         task_id="delete_nodegroup_and_cluster",
         cluster_name=cluster_name,
         force_delete_compute=True,
-        retries=4,
-        retry_delay=duration(seconds=30),
-        retry_exponential_backoff=True,
     )
     # [END howto_operator_eks_force_delete_cluster]
     delete_nodegroup_and_cluster.trigger_rule = TriggerRule.ALL_DONE
diff --git 
a/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroups.py 
b/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroups.py
index 79abe2b3919..018d49c0831 100644
--- a/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroups.py
+++ b/providers/amazon/tests/system/amazon/aws/example_eks_with_nodegroups.py
@@ -19,7 +19,6 @@ from __future__ import annotations
 from datetime import datetime
 
 import boto3
-from pendulum import duration
 
 from airflow.providers.amazon.aws.hooks.eks import ClusterStates, 
NodegroupStates
 from airflow.providers.amazon.aws.operators.eks import (
@@ -172,9 +171,6 @@ with DAG(
         task_id="delete_nodegroup",
         cluster_name=cluster_name,
         nodegroup_name=nodegroup_name,
-        retries=4,
-        retry_delay=duration(seconds=30),
-        retry_exponential_backoff=True,
     )
     # [END howto_operator_eks_delete_nodegroup]
     delete_nodegroup.trigger_rule = TriggerRule.ALL_DONE
diff --git a/providers/amazon/tests/unit/amazon/aws/operators/test_eks.py 
b/providers/amazon/tests/unit/amazon/aws/operators/test_eks.py
index 5003e56d76a..abd9b6c334d 100644
--- a/providers/amazon/tests/unit/amazon/aws/operators/test_eks.py
+++ b/providers/amazon/tests/unit/amazon/aws/operators/test_eks.py
@@ -72,6 +72,15 @@ CREATE_NODEGROUP_KWARGS = {
     "instanceTypes": "t4g.large",
 }
 
+RESOURCE_IN_USE_ERROR = ClientError(
+    error_response={"Error": {"Code": "ResourceInUseException", "Message": 
"update in progress"}},
+    operation_name="DeleteCluster",
+)
+RESOURCE_NOT_FOUND_ERROR = ClientError(
+    error_response={"Error": {"Code": "ResourceNotFoundException", "Message": 
"not found"}},
+    operation_name="DeleteCluster",
+)
+
 
 class ClusterParams(TypedDict):
     cluster_name: str
@@ -628,6 +637,39 @@ class TestEksCreateNodegroupOperator:
             nodegroupName=NODEGROUP_NAME,
         )
 
+    @mock.patch("time.sleep", return_value=None)
+    @mock.patch.object(EksHook, "delete_nodegroup")
+    @mock.patch("airflow.providers.amazon.aws.operators.eks.wait")
+    @mock.patch.object(EksHook, "create_nodegroup")
+    def test_nodegroup_cleanup_retries_on_resource_in_use(
+        self,
+        mock_create_nodegroup,
+        mock_waiter,
+        mock_delete_nodegroup,
+        mock_sleep,
+    ):
+        mock_waiter.side_effect = AirflowException("Nodegroup creation failed: 
Waiter NodegroupActive failed")
+        # A freshly-failed nodegroup may still be settling, so the cleanup 
delete rides out
+        # transient ResourceInUseException before succeeding.
+        mock_delete_nodegroup.side_effect = [RESOURCE_IN_USE_ERROR, 
RESOURCE_IN_USE_ERROR, None]
+
+        operator = EksCreateNodegroupOperator(
+            task_id=TASK_ID,
+            cluster_name=CLUSTER_NAME,
+            nodegroup_name=NODEGROUP_NAME,
+            nodegroup_subnets=SUBNET_IDS,
+            nodegroup_role_arn=NODEROLE_ARN[1],
+            wait_for_completion=True,
+            delete_nodegroup_on_failure=True,
+        )
+
+        # The original creation error is still raised once cleanup eventually 
succeeds.
+        with pytest.raises(AirflowException, match="Nodegroup creation 
failed"):
+            operator.execute({})
+
+        assert mock_delete_nodegroup.call_count == 3
+        mock_delete_nodegroup.assert_called_with(clusterName=CLUSTER_NAME, 
nodegroupName=NODEGROUP_NAME)
+
     @mock.patch.object(EksHook, "delete_nodegroup")
     @mock.patch("airflow.providers.amazon.aws.operators.eks.wait")
     @mock.patch.object(EksHook, "create_nodegroup")
@@ -718,6 +760,53 @@ class TestEksDeleteClusterOperator:
         with pytest.raises(TaskDeferred):
             self.delete_cluster_operator.execute({})
 
+    @mock.patch("time.sleep", return_value=None)
+    @mock.patch.object(Waiter, "wait")
+    @mock.patch.object(EksHook, "list_nodegroups")
+    @mock.patch.object(EksHook, "delete_cluster")
+    def test_delete_cluster_retries_on_resource_in_use(
+        self, mock_delete_cluster, mock_list_nodegroups, mock_waiter, 
mock_sleep
+    ):
+        mock_list_nodegroups.return_value = []
+        mock_delete_cluster.side_effect = [RESOURCE_IN_USE_ERROR, 
RESOURCE_IN_USE_ERROR, None]
+
+        self.delete_cluster_operator.execute({})
+
+        assert mock_delete_cluster.call_count == 3
+        mock_delete_cluster.assert_called_with(name=self.cluster_name)
+
+    @mock.patch("time.sleep", return_value=None)
+    @mock.patch.object(EksHook, "get_waiter")
+    @mock.patch.object(EksHook, "list_nodegroups")
+    @mock.patch.object(EksHook, "delete_nodegroup")
+    def test_delete_any_nodegroups_retries_on_resource_in_use(
+        self, mock_delete_nodegroup, mock_list_nodegroups, mock_get_waiter, 
mock_sleep
+    ):
+        mock_list_nodegroups.return_value = ["ng1"]
+        mock_delete_nodegroup.side_effect = [RESOURCE_IN_USE_ERROR, 
RESOURCE_IN_USE_ERROR, None]
+
+        self.delete_cluster_operator.delete_any_nodegroups()
+
+        assert mock_delete_nodegroup.call_count == 3
+        
mock_delete_nodegroup.assert_called_with(clusterName=self.cluster_name, 
nodegroupName="ng1")
+
+    @mock.patch("time.sleep", return_value=None)
+    @mock.patch.object(Waiter, "wait")
+    @mock.patch.object(EksHook, "list_fargate_profiles")
+    @mock.patch.object(EksHook, "delete_fargate_profile")
+    def test_delete_any_fargate_profiles_retries_on_resource_in_use(
+        self, mock_delete_fargate_profile, mock_list_fargate_profiles, 
mock_waiter, mock_sleep
+    ):
+        mock_list_fargate_profiles.return_value = ["fp1"]
+        mock_delete_fargate_profile.side_effect = [RESOURCE_IN_USE_ERROR, 
RESOURCE_IN_USE_ERROR, None]
+
+        self.delete_cluster_operator.delete_any_fargate_profiles()
+
+        assert mock_delete_fargate_profile.call_count == 3
+        mock_delete_fargate_profile.assert_called_with(
+            clusterName=self.cluster_name, fargateProfileName="fp1"
+        )
+
     def test_template_fields(self):
         validate_template_fields(self.delete_cluster_operator)
 
@@ -761,6 +850,21 @@ class TestEksDeleteNodegroupOperator:
         mock_waiter.assert_called_with(mock.ANY, clusterName=CLUSTER_NAME, 
nodegroupName=NODEGROUP_NAME)
         assert_expected_waiter_type(mock_waiter, "NodegroupDeleted")
 
+    @mock.patch("time.sleep", return_value=None)
+    @mock.patch.object(Waiter, "wait")
+    @mock.patch.object(EksHook, "delete_nodegroup")
+    def test_delete_nodegroup_retries_on_resource_in_use(
+        self, mock_delete_nodegroup, mock_waiter, mock_sleep
+    ):
+        mock_delete_nodegroup.side_effect = [RESOURCE_IN_USE_ERROR, 
RESOURCE_IN_USE_ERROR, None]
+
+        self.delete_nodegroup_operator.execute({})
+
+        assert mock_delete_nodegroup.call_count == 3
+        mock_delete_nodegroup.assert_called_with(
+            clusterName=self.cluster_name, nodegroupName=self.nodegroup_name
+        )
+
     def test_template_fields(self):
         validate_template_fields(self.delete_nodegroup_operator)
 
@@ -822,6 +926,21 @@ class TestEksDeleteFargateProfileOperator:
             "Trigger is not a EksDeleteFargateProfileTrigger"
         )
 
+    @mock.patch("time.sleep", return_value=None)
+    @mock.patch.object(Waiter, "wait")
+    @mock.patch.object(EksHook, "delete_fargate_profile")
+    def test_delete_fargate_profile_retries_on_resource_in_use(
+        self, mock_delete_fargate_profile, mock_waiter, mock_sleep
+    ):
+        mock_delete_fargate_profile.side_effect = [RESOURCE_IN_USE_ERROR, 
RESOURCE_IN_USE_ERROR, None]
+
+        self.delete_fargate_profile_operator.execute({})
+
+        assert mock_delete_fargate_profile.call_count == 3
+        mock_delete_fargate_profile.assert_called_with(
+            clusterName=self.cluster_name, 
fargateProfileName=self.fargate_profile_name
+        )
+
     def test_template_fields(self):
         validate_template_fields(self.delete_fargate_profile_operator)
 
diff --git a/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py 
b/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py
index fcdd712cc56..bc92c52fb1a 100644
--- a/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py
+++ b/providers/amazon/tests/unit/amazon/aws/triggers/test_eks.py
@@ -174,6 +174,20 @@ class TestEksDeleteClusterTriggerRun(TestEksTrigger):
         assert exception._excinfo[1].response == response
         assert exception._excinfo[1].operation_name == operation_name
 
+    @pytest.mark.asyncio
+    @patch("asyncio.sleep", return_value=None)
+    async def test_run_retries_on_resource_in_use(self, mock_sleep):
+        in_use = ClientError({"Error": {"Code": "ResourceInUseException"}}, 
"delete_eks_cluster")
+        delete_cluster_mock = AsyncMock(side_effect=[in_use, in_use, None])
+        self.mock_client.delete_cluster = delete_cluster_mock
+
+        generator = self.trigger.run()
+        response = await generator.asend(None)
+
+        assert delete_cluster_mock.call_count == 3
+        delete_cluster_mock.assert_called_with(name=CLUSTER_NAME)
+        assert response == TriggerEvent({"status": "deleted"})
+
     @pytest.mark.asyncio
     async def test_run_parameterizes_async_wait_correctly(self):
         self.mock_client.get_waiter = Mock(return_value="waiter")
@@ -247,6 +261,36 @@ class 
TestEksDeleteClusterTriggerDeleteNodegroupsAndFargateProfiles(TestEksTrigg
             status_args=["nodegroups"],
         )
 
+    @pytest.mark.asyncio
+    @patch("asyncio.sleep", return_value=None)
+    async def test_delete_nodegroups_retries_on_resource_in_use(self, 
mock_sleep):
+        in_use = ClientError({"Error": {"Code": "ResourceInUseException"}}, 
"DeleteNodegroup")
+        mock_list_node_groups = AsyncMock(return_value={"nodegroups": ["g1"]})
+        mock_delete_nodegroup = AsyncMock(side_effect=[in_use, in_use, None])
+        mock_client = AsyncMock(list_nodegroups=mock_list_node_groups, 
delete_nodegroup=mock_delete_nodegroup)
+
+        await self.trigger.delete_any_nodegroups(mock_client)
+
+        assert mock_delete_nodegroup.call_count == 3
+        mock_delete_nodegroup.assert_called_with(clusterName=CLUSTER_NAME, 
nodegroupName="g1")
+
+    @pytest.mark.asyncio
+    @patch("asyncio.sleep", return_value=None)
+    async def test_delete_fargate_profiles_retries_on_resource_in_use(self, 
mock_sleep):
+        in_use = ClientError({"Error": {"Code": "ResourceInUseException"}}, 
"DeleteFargateProfile")
+        mock_list_fargate_profiles = 
AsyncMock(return_value={"fargateProfileNames": ["p1"]})
+        mock_delete_fargate_profile = AsyncMock(side_effect=[in_use, in_use, 
None])
+        mock_client = AsyncMock(
+            list_fargate_profiles=mock_list_fargate_profiles,
+            delete_fargate_profile=mock_delete_fargate_profile,
+            get_waiter=self.mock_waiter,
+        )
+
+        await self.trigger.delete_any_fargate_profiles(mock_client)
+
+        assert mock_delete_fargate_profile.call_count == 3
+        
mock_delete_fargate_profile.assert_called_with(clusterName=CLUSTER_NAME, 
fargateProfileName="p1")
+
     @pytest.mark.asyncio
     async def 
test_when_there_are_no_nodegroups_it_should_only_log_message(self):
         mock_list_node_groups = AsyncMock(return_value={"nodegroups": []})
diff --git a/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py 
b/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py
index b9a04a8f54d..2eea945642b 100644
--- a/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py
+++ b/providers/amazon/tests/unit/amazon/aws/utils/test_utils.py
@@ -17,16 +17,22 @@
 from __future__ import annotations
 
 import datetime
+import logging
+from unittest import mock
 
 import pytest
+import tenacity
+from botocore.exceptions import ClientError
 
 from airflow.providers.amazon.aws.utils import (
     _StringCompareEnum,
+    build_resource_in_use_retry_args,
     datetime_to_epoch,
     datetime_to_epoch_ms,
     datetime_to_epoch_us,
     get_airflow_version,
     get_botocore_version,
+    is_resource_in_use_error,
 )
 
 DT = datetime.datetime(2000, 1, 1, tzinfo=datetime.timezone.utc)
@@ -66,3 +72,63 @@ def test_botocore_version():
     assert isinstance(botocore_version[0], int), "botocore major version 
expected to be an integer"
     assert isinstance(botocore_version[1], int), "botocore minor version 
expected to be an integer"
     assert isinstance(botocore_version[2], int), "botocore patch version 
expected to be an integer"
+
+
+def _run_with_retry(call, **overrides):
+    retry_args = {**build_resource_in_use_retry_args(logging.getLogger()), 
**overrides}
+    for attempt in tenacity.Retrying(**retry_args):
+        with attempt:
+            return call()
+
+
[email protected](
+    ("exception", "expected"),
+    [
+        pytest.param(
+            ClientError({"Error": {"Code": "ResourceInUseException"}}, 
"DeleteCluster"),
+            True,
+            id="resource_in_use",
+        ),
+        pytest.param(
+            ClientError({"Error": {"Code": "ResourceNotFoundException"}}, 
"DeleteCluster"),
+            False,
+            id="other_client_error",
+        ),
+        pytest.param(ValueError("boom"), False, id="non_client_error"),
+    ],
+)
+def test_is_resource_in_use_error(exception, expected):
+    assert is_resource_in_use_error(exception) is expected
+
+
[email protected]("time.sleep", return_value=None)
+def test_build_resource_in_use_retry_args_retries_then_succeeds(mock_sleep):
+    in_use = ClientError({"Error": {"Code": "ResourceInUseException"}}, 
"DeleteCluster")
+    call = mock.Mock(side_effect=[in_use, in_use, "ok"])
+
+    assert _run_with_retry(call) == "ok"
+    assert call.call_count == 3
+
+
[email protected]("time.sleep", return_value=None)
+def 
test_build_resource_in_use_retry_args_reraises_when_stop_reached(mock_sleep):
+    in_use = ClientError({"Error": {"Code": "ResourceInUseException"}}, 
"DeleteCluster")
+    call = mock.Mock(side_effect=in_use)
+
+    # Override the stop so exhaustion is reached without waiting out the real 
timeout.
+    with pytest.raises(ClientError) as exc_info:
+        _run_with_retry(call, stop=tenacity.stop_after_attempt(3))
+
+    assert exc_info.value is in_use
+    assert call.call_count == 3
+
+
[email protected]("time.sleep", return_value=None)
+def 
test_build_resource_in_use_retry_args_does_not_retry_other_errors(mock_sleep):
+    other = ClientError({"Error": {"Code": "ResourceNotFoundException"}}, 
"DeleteCluster")
+    call = mock.Mock(side_effect=other)
+
+    with pytest.raises(ClientError):
+        _run_with_retry(call)
+
+    assert call.call_count == 1

Reply via email to