This is an automated email from the ASF dual-hosted git repository.
o-nikolas 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 acf0ad84cec Use the configured region for deferred Neptune cluster
tasks (#71646)
acf0ad84cec is described below
commit acf0ad84cec08c98143024984da60482e43a5d81
Author: Sepuri Sai Krishna <[email protected]>
AuthorDate: Thu Aug 20 05:14:56 2026 +0530
Use the configured region for deferred Neptune cluster tasks (#71646)
The triggers accepted a region_name but never handed it to the base trigger
that assigns the attribute, and two of the operator defer sites never
passed
the hook configuration at all. A deferred start or stop therefore polled
the
default region, and the task failed only after exhausting its waiter
attempts
with an error that never mentioned the region.
---
.../providers/amazon/aws/operators/neptune.py | 6 +++
.../providers/amazon/aws/triggers/neptune.py | 3 ++
.../unit/amazon/aws/operators/test_neptune.py | 46 ++++++++++++++++++++++
.../tests/unit/amazon/aws/triggers/test_neptune.py | 38 ++++++++++++++++++
4 files changed, 93 insertions(+)
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/operators/neptune.py
b/providers/amazon/src/airflow/providers/amazon/aws/operators/neptune.py
index 9e3916da5eb..93678c9a7cc 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/operators/neptune.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/operators/neptune.py
@@ -175,6 +175,9 @@ class
NeptuneStartDbClusterOperator(AwsBaseOperator[NeptuneHook]):
db_cluster_id=self.cluster_id,
waiter_delay=self.waiter_delay,
waiter_max_attempts=self.waiter_max_attempts,
+ region_name=self.region_name,
+ botocore_config=self.botocore_config,
+ verify=self.verify,
),
method_name="execute_complete",
)
@@ -302,6 +305,9 @@ class
NeptuneStopDbClusterOperator(AwsBaseOperator[NeptuneHook]):
db_cluster_id=self.cluster_id,
waiter_delay=self.waiter_delay,
waiter_max_attempts=self.waiter_max_attempts,
+ region_name=self.region_name,
+ botocore_config=self.botocore_config,
+ verify=self.verify,
),
method_name="execute_complete",
)
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/triggers/neptune.py
b/providers/amazon/src/airflow/providers/amazon/aws/triggers/neptune.py
index 2c7a3a9b4f1..e7b19fc75c1 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/triggers/neptune.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/triggers/neptune.py
@@ -58,6 +58,7 @@ class NeptuneClusterAvailableTrigger(AwsBaseWaiterTrigger):
waiter_delay=waiter_delay,
waiter_max_attempts=waiter_max_attempts,
aws_conn_id=aws_conn_id,
+ region_name=region_name,
**kwargs,
)
@@ -103,6 +104,7 @@ class NeptuneClusterStoppedTrigger(AwsBaseWaiterTrigger):
waiter_delay=waiter_delay,
waiter_max_attempts=waiter_max_attempts,
aws_conn_id=aws_conn_id,
+ region_name=region_name,
**kwargs,
)
@@ -148,6 +150,7 @@ class
NeptuneClusterInstancesAvailableTrigger(AwsBaseWaiterTrigger):
waiter_delay=waiter_delay,
waiter_max_attempts=waiter_max_attempts,
aws_conn_id=aws_conn_id,
+ region_name=region_name,
**kwargs,
)
diff --git a/providers/amazon/tests/unit/amazon/aws/operators/test_neptune.py
b/providers/amazon/tests/unit/amazon/aws/operators/test_neptune.py
index 7ca801b2fbc..dc16e6ec994 100644
--- a/providers/amazon/tests/unit/amazon/aws/operators/test_neptune.py
+++ b/providers/amazon/tests/unit/amazon/aws/operators/test_neptune.py
@@ -29,11 +29,20 @@ from airflow.providers.amazon.aws.operators.neptune import (
NeptuneStartDbClusterOperator,
NeptuneStopDbClusterOperator,
)
+from airflow.providers.amazon.aws.triggers.neptune import (
+ NeptuneClusterAvailableTrigger,
+ NeptuneClusterStoppedTrigger,
+)
from airflow.providers.common.compat.sdk import AirflowException, TaskDeferred
from unit.amazon.aws.utils.test_template_fields import validate_template_fields
CLUSTER_ID = "test_cluster"
+REGION_NAME = "eu-west-2"
+VERIFY = False
+BOTOCORE_CONFIG = {"read_timeout": 42}
+WAITER_DELAY = 12
+WAITER_MAX_ATTEMPTS = 34
EXPECTED_RESPONSE = {"db_cluster_id": CLUSTER_ID}
@@ -54,6 +63,43 @@ def _create_cluster(hook: NeptuneHook):
raise ValueError("AWS not properly mocked")
[email protected](
+ ("operator_class", "trigger_class"),
+ [
+ (NeptuneStartDbClusterOperator, NeptuneClusterAvailableTrigger),
+ (NeptuneStopDbClusterOperator, NeptuneClusterStoppedTrigger),
+ ],
+)
[email protected](NeptuneHook, "conn")
+def test_deferred_trigger_receives_the_operator_configuration(mock_conn,
operator_class, trigger_class):
+ operator = operator_class(
+ task_id="task_test",
+ db_cluster_id=CLUSTER_ID,
+ deferrable=True,
+ wait_for_completion=False,
+ aws_conn_id="aws_default",
+ region_name=REGION_NAME,
+ verify=VERIFY,
+ botocore_config=BOTOCORE_CONFIG,
+ waiter_delay=WAITER_DELAY,
+ waiter_max_attempts=WAITER_MAX_ATTEMPTS,
+ )
+
+ with pytest.raises(TaskDeferred) as deferred:
+ operator.execute(None)
+
+ assert isinstance(deferred.value.trigger, trigger_class)
+ assert deferred.value.trigger.serialize()[1] == {
+ "db_cluster_id": CLUSTER_ID,
+ "aws_conn_id": "aws_default",
+ "region_name": REGION_NAME,
+ "verify": VERIFY,
+ "botocore_config": BOTOCORE_CONFIG,
+ "waiter_delay": WAITER_DELAY,
+ "waiter_max_attempts": WAITER_MAX_ATTEMPTS,
+ }
+
+
class TestNeptuneStartClusterOperator:
@mock.patch.object(NeptuneHook, "conn")
@mock.patch.object(NeptuneHook, "get_waiter")
diff --git a/providers/amazon/tests/unit/amazon/aws/triggers/test_neptune.py
b/providers/amazon/tests/unit/amazon/aws/triggers/test_neptune.py
index 4fa9d0ed38f..2329fafb633 100644
--- a/providers/amazon/tests/unit/amazon/aws/triggers/test_neptune.py
+++ b/providers/amazon/tests/unit/amazon/aws/triggers/test_neptune.py
@@ -30,6 +30,44 @@ from airflow.providers.amazon.aws.triggers.neptune import (
from airflow.triggers.base import TriggerEvent
CLUSTER_ID = "test-cluster"
+REGION_NAME = "eu-west-2"
+VERIFY = False
+BOTOCORE_CONFIG = {"read_timeout": 42}
+WAITER_DELAY = 12
+WAITER_MAX_ATTEMPTS = 34
+
+
[email protected](
+ "trigger_class",
+ [
+ NeptuneClusterAvailableTrigger,
+ NeptuneClusterStoppedTrigger,
+ NeptuneClusterInstancesAvailableTrigger,
+ ],
+)
+def test_hook_configuration_survives_serialization(trigger_class):
+ trigger = trigger_class(
+ db_cluster_id=CLUSTER_ID,
+ aws_conn_id="aws_default",
+ region_name=REGION_NAME,
+ verify=VERIFY,
+ botocore_config=BOTOCORE_CONFIG,
+ waiter_delay=WAITER_DELAY,
+ waiter_max_attempts=WAITER_MAX_ATTEMPTS,
+ )
+
+ assert trigger.region_name == REGION_NAME
+ assert trigger.verify == VERIFY
+ assert trigger.botocore_config == BOTOCORE_CONFIG
+ assert trigger.serialize()[1] == {
+ "db_cluster_id": CLUSTER_ID,
+ "aws_conn_id": "aws_default",
+ "region_name": REGION_NAME,
+ "verify": VERIFY,
+ "botocore_config": BOTOCORE_CONFIG,
+ "waiter_delay": WAITER_DELAY,
+ "waiter_max_attempts": WAITER_MAX_ATTEMPTS,
+ }
class TestNeptuneClusterAvailableTrigger: