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 98e753847a3 Distinguish credential failures from AWS waiter exhaustion
(#73870)
98e753847a3 is described below
commit 98e753847a3332ee6e9935805a0f548edf84b617
Author: SameerMesiah97 <[email protected]>
AuthorDate: Mon Oct 5 15:33:08 2026 +0100
Distinguish credential failures from AWS waiter exhaustion (#73870)
* Distinguish credential failures from AWS waiter exhaustion
Raise a dedicated error when every waiter attempt fails due to missing
credentials instead of reporting the outcome as generic waiter
exhaustion. Apply the distinction to both synchronous and asynchronous
waiters and add test coverage for credential-only and mixed attempts.
* Include the underlying NoCredentialsError when credential attempts are
exhausted.
Preserve compatibility by subclassing WaiterMaxAttemptsError and update
tests accordingly.
* Raise WaiterMaxAttemptsError when the waiter makes no attempts
With waiter_max_attempts <= 0 the retry loop never runs, so the
all-attempts-failed-on-credentials flag kept its initial value and the
waiter blamed missing credentials that were never checked.
Generated-by: Claude Opus 5
---------
Co-authored-by: Sameer Mesiah <[email protected]>
Co-authored-by: Jarek Potiuk <[email protected]>
---
.../src/airflow/providers/amazon/aws/exceptions.py | 4 +
.../amazon/aws/utils/waiter_with_logging.py | 24 +++-
.../amazon/aws/utils/test_waiter_with_logging.py | 135 ++++++++++++++++++++-
3 files changed, 161 insertions(+), 2 deletions(-)
diff --git a/providers/amazon/src/airflow/providers/amazon/aws/exceptions.py
b/providers/amazon/src/airflow/providers/amazon/aws/exceptions.py
index 64c15f86189..7a59aa0fe2f 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/exceptions.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/exceptions.py
@@ -125,3 +125,7 @@ class WaiterTerminalFailure(AirflowException):
class WaiterMaxAttemptsError(AirflowException):
"""Raised when an AWS waiter exhausts its configured attempts."""
+
+
+class WaiterNoCredentialsError(WaiterMaxAttemptsError):
+ """Raised when an AWS waiter exhausts all attempts due to missing
credentials."""
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/utils/waiter_with_logging.py
b/providers/amazon/src/airflow/providers/amazon/aws/utils/waiter_with_logging.py
index 6ae8ecf7d33..aff06e9d225 100644
---
a/providers/amazon/src/airflow/providers/amazon/aws/utils/waiter_with_logging.py
+++
b/providers/amazon/src/airflow/providers/amazon/aws/utils/waiter_with_logging.py
@@ -25,7 +25,11 @@ from typing import TYPE_CHECKING, Any
import jmespath
from botocore.exceptions import NoCredentialsError, WaiterError
-from airflow.providers.amazon.aws.exceptions import WaiterMaxAttemptsError,
WaiterTerminalFailure
+from airflow.providers.amazon.aws.exceptions import (
+ WaiterMaxAttemptsError,
+ WaiterNoCredentialsError,
+ WaiterTerminalFailure,
+)
from airflow.providers.common.compat.sdk import AirflowException
if TYPE_CHECKING:
@@ -90,6 +94,9 @@ def wait(
log = logging.getLogger(__name__)
first_attempt = True
attempt = 0
+ all_attempts_no_credentials = True
+ last_no_credentials_error: NoCredentialsError | None = None
+
while attempt < waiter_max_attempts:
if not first_attempt:
time.sleep(waiter_delay)
@@ -98,9 +105,11 @@ def wait(
waiter.wait(**args, WaiterConfig={"MaxAttempts": 1})
except NoCredentialsError as error:
+ last_no_credentials_error = error
log.info(str(error))
except WaiterError as error:
+ all_attempts_no_credentials = False
error_reason = str(error)
last_response = error.last_response
@@ -134,6 +143,10 @@ def wait(
break
attempt += 1
else:
+ if all_attempts_no_credentials and last_no_credentials_error is not
None:
+ raise WaiterNoCredentialsError(
+ f"Waiter error: max attempts reached due to missing
credentials: {last_no_credentials_error}"
+ )
raise WaiterMaxAttemptsError("Waiter error: max attempts reached")
@@ -173,6 +186,9 @@ async def async_wait(
log = logging.getLogger(__name__)
first_attempt = True
attempt = 0
+ all_attempts_no_credentials = True
+ last_no_credentials_error: NoCredentialsError | None = None
+
while attempt < waiter_max_attempts:
if not first_attempt:
await asyncio.sleep(waiter_delay)
@@ -181,9 +197,11 @@ async def async_wait(
await waiter.wait(**args, WaiterConfig={"MaxAttempts": 1})
except NoCredentialsError as error:
+ last_no_credentials_error = error
log.info(str(error))
except WaiterError as error:
+ all_attempts_no_credentials = False
error_reason = str(error)
last_response = error.last_response
@@ -216,6 +234,10 @@ async def async_wait(
break
attempt += 1
else:
+ if all_attempts_no_credentials and last_no_credentials_error is not
None:
+ raise WaiterNoCredentialsError(
+ f"Waiter error: max attempts reached due to missing
credentials: {last_no_credentials_error}"
+ )
raise WaiterMaxAttemptsError("Waiter error: max attempts reached")
diff --git
a/providers/amazon/tests/unit/amazon/aws/utils/test_waiter_with_logging.py
b/providers/amazon/tests/unit/amazon/aws/utils/test_waiter_with_logging.py
index 817be27df2d..6f9267e5084 100644
--- a/providers/amazon/tests/unit/amazon/aws/utils/test_waiter_with_logging.py
+++ b/providers/amazon/tests/unit/amazon/aws/utils/test_waiter_with_logging.py
@@ -23,10 +23,11 @@ from unittest import mock
from unittest.mock import AsyncMock
import pytest
-from botocore.exceptions import WaiterError
+from botocore.exceptions import NoCredentialsError, WaiterError
from airflow.providers.amazon.aws.exceptions import (
WaiterMaxAttemptsError,
+ WaiterNoCredentialsError,
WaiterTerminalFailure,
)
from airflow.providers.amazon.aws.utils.waiter_with_logging import
_LazyStatusFormatter, async_wait, wait
@@ -190,6 +191,138 @@ class TestWaiter:
assert "Waiter error: max attempts reached" in str(exc.value)
assert mock_waiter.wait.call_count == 2
+ @mock.patch("time.sleep")
+ def test_wait_all_attempts_no_credentials(self, mock_sleep):
+ mock_waiter = mock.MagicMock()
+ error = NoCredentialsError()
+ mock_waiter.wait.side_effect = error
+
+ with pytest.raises(WaiterNoCredentialsError) as exc:
+ wait(
+ waiter=mock_waiter,
+ waiter_delay=123,
+ waiter_max_attempts=2,
+ args={"test_arg": "test_value"},
+ failure_message="test failure message",
+ status_message="test status message",
+ status_args=["Status.State"],
+ )
+
+ assert "Waiter error: max attempts reached due to missing credentials"
in str(exc.value)
+ assert str(error) in str(exc.value)
+ assert mock_waiter.wait.call_count == 2
+ mock_sleep.assert_called_once_with(123)
+
+ @pytest.mark.asyncio
+ async def test_async_wait_all_attempts_no_credentials(self):
+ mock_waiter = mock.MagicMock()
+ error = NoCredentialsError()
+ mock_waiter.wait = AsyncMock(side_effect=error)
+
+ with pytest.raises(WaiterNoCredentialsError) as exc:
+ await async_wait(
+ waiter=mock_waiter,
+ waiter_delay=0,
+ waiter_max_attempts=2,
+ args={"test_arg": "test_value"},
+ failure_message="test failure message",
+ status_message="test status message",
+ status_args=["Status.State"],
+ )
+
+ assert "Waiter error: max attempts reached due to missing credentials"
in str(exc.value)
+ assert str(error) in str(exc.value)
+ assert mock_waiter.wait.call_count == 2
+
+ @mock.patch("time.sleep")
+ def test_wait_mixed_attempts_raise_max_attempts_error(self, mock_sleep):
+ mock_waiter = mock.MagicMock()
+ error = WaiterError(
+ name="test_waiter",
+ reason="test_reason",
+ last_response=generate_response("Pending"),
+ )
+ mock_waiter.wait.side_effect = [error, NoCredentialsError()]
+
+ with pytest.raises(WaiterMaxAttemptsError) as exc:
+ wait(
+ waiter=mock_waiter,
+ waiter_delay=123,
+ waiter_max_attempts=2,
+ args={"test_arg": "test_value"},
+ failure_message="test failure message",
+ status_message="test status message",
+ status_args=["Status.State"],
+ )
+
+ assert "Waiter error: max attempts reached" in str(exc.value)
+ assert mock_waiter.wait.call_count == 2
+ assert type(exc.value) is WaiterMaxAttemptsError
+ mock_sleep.assert_called_once_with(123)
+
+ @pytest.mark.asyncio
+ async def test_async_wait_mixed_attempts_raise_max_attempts_error(self):
+ mock_waiter = mock.MagicMock()
+ error = WaiterError(
+ name="test_waiter",
+ reason="test_reason",
+ last_response=generate_response("Pending"),
+ )
+ mock_waiter.wait = AsyncMock(side_effect=[error, NoCredentialsError()])
+
+ with pytest.raises(WaiterMaxAttemptsError) as exc:
+ await async_wait(
+ waiter=mock_waiter,
+ waiter_delay=0,
+ waiter_max_attempts=2,
+ args={"test_arg": "test_value"},
+ failure_message="test failure message",
+ status_message="test status message",
+ status_args=["Status.State"],
+ )
+
+ assert "Waiter error: max attempts reached" in str(exc.value)
+ assert type(exc.value) is WaiterMaxAttemptsError
+ assert mock_waiter.wait.call_count == 2
+
+ @mock.patch("time.sleep")
+ def test_wait_zero_attempts_raise_max_attempts_error(self, mock_sleep):
+ mock_waiter = mock.MagicMock()
+
+ with pytest.raises(WaiterMaxAttemptsError) as exc:
+ wait(
+ waiter=mock_waiter,
+ waiter_delay=123,
+ waiter_max_attempts=0,
+ args={"test_arg": "test_value"},
+ failure_message="test failure message",
+ status_message="test status message",
+ status_args=["Status.State"],
+ )
+
+ assert type(exc.value) is WaiterMaxAttemptsError
+ mock_waiter.wait.assert_not_called()
+ mock_sleep.assert_not_called()
+
+ @pytest.mark.asyncio
+ async def test_async_wait_zero_attempts_raise_max_attempts_error(self):
+ mock_waiter = mock.MagicMock()
+ mock_waiter.wait = AsyncMock()
+
+ with pytest.raises(WaiterMaxAttemptsError) as exc:
+ await async_wait(
+ waiter=mock_waiter,
+ waiter_delay=0,
+ waiter_max_attempts=0,
+ args={"test_arg": "test_value"},
+ failure_message="test failure message",
+ status_message="test status message",
+ status_args=["Status.State"],
+ )
+
+ assert type(exc.value) is WaiterMaxAttemptsError
+ mock_waiter.wait.assert_not_called()
+
@mock.patch("time.sleep")
def test_wait_with_failure(self, mock_sleep):
mock_sleep.return_value = True