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 3abf9ca59c8 WaitSensor: make time_to_wait templated (#70480)
3abf9ca59c8 is described below
commit 3abf9ca59c8077aa6ebc151d48b8aa4203637831
Author: raphaelauv <[email protected]>
AuthorDate: Fri Jul 31 21:12:08 2026 +0200
WaitSensor: make time_to_wait templated (#70480)
---
.../providers/standard/sensors/time_delta.py | 21 ++++++++++++------
.../tests/unit/standard/sensors/test_time_delta.py | 25 ++++++++++++++++++++++
2 files changed, 39 insertions(+), 7 deletions(-)
diff --git
a/providers/standard/src/airflow/providers/standard/sensors/time_delta.py
b/providers/standard/src/airflow/providers/standard/sensors/time_delta.py
index e139499fec6..0f9bcf45d07 100644
--- a/providers/standard/src/airflow/providers/standard/sensors/time_delta.py
+++ b/providers/standard/src/airflow/providers/standard/sensors/time_delta.py
@@ -18,6 +18,7 @@
from __future__ import annotations
import warnings
+from collections.abc import Sequence
from datetime import datetime, timedelta
from time import sleep
from typing import TYPE_CHECKING, Any
@@ -177,6 +178,8 @@ class WaitSensor(BaseSensorOperator):
:param deferrable: Run sensor in deferrable mode
"""
+ template_fields: Sequence[str] = ("time_to_wait",)
+
def __init__(
self,
time_to_wait: timedelta | int,
@@ -185,20 +188,24 @@ class WaitSensor(BaseSensorOperator):
) -> None:
super().__init__(**kwargs)
self.deferrable = deferrable
- if isinstance(time_to_wait, int):
- self.time_to_wait = timedelta(minutes=time_to_wait)
- else:
- self.time_to_wait = time_to_wait
+ self.time_to_wait = time_to_wait
+
+ def _resolve_time_to_wait(self) -> timedelta:
+ value = self.time_to_wait
+ if isinstance(value, timedelta):
+ return value
+ return timedelta(minutes=int(value))
def execute(self, context: Context) -> None:
+ time_to_wait = self._resolve_time_to_wait()
if self.deferrable:
self.defer(
trigger=(
- TimeDeltaTrigger(self.time_to_wait, end_from_trigger=True)
+ TimeDeltaTrigger(time_to_wait, end_from_trigger=True)
if AIRFLOW_V_3_0_PLUS
- else TimeDeltaTrigger(self.time_to_wait)
+ else TimeDeltaTrigger(time_to_wait)
),
method_name="execute_complete",
)
else:
- sleep(int(self.time_to_wait.total_seconds()))
+ sleep(int(time_to_wait.total_seconds()))
diff --git a/providers/standard/tests/unit/standard/sensors/test_time_delta.py
b/providers/standard/tests/unit/standard/sensors/test_time_delta.py
index 5c4eeef2a8e..9ec8fc4233d 100644
--- a/providers/standard/tests/unit/standard/sensors/test_time_delta.py
+++ b/providers/standard/tests/unit/standard/sensors/test_time_delta.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import re
from datetime import timedelta
from typing import Any
@@ -253,3 +254,27 @@ class TestTimeDeltaSensorAsync:
op.execute(context)
assert caught.value.trigger.moment == expected_time
+
+ @pytest.mark.parametrize(
+ "time_to_wait",
+ [timedelta(minutes=1), 1, "{{ 1*2 }}"],
+ )
+ def test_wait_sensor_templating(self, mocker, time_to_wait):
+ defer_mock = mocker.patch(DEFER_PATH)
+ op = WaitSensor(task_id="wait_sensor_check",
time_to_wait=time_to_wait, dag=self.dag, deferrable=True)
+
+ with time_machine.travel(pendulum.datetime(year=2024, month=8, day=1,
tz="UTC"), tick=False):
+ context = op.render_template_fields({})
+ op.execute(context)
+ defer_mock.assert_called_once()
+
+ def test_wait_sensor_templating_error(self, mocker):
+ op = WaitSensor(
+ task_id="wait_sensor_check", time_to_wait="{{ 'nothing' }}",
dag=self.dag, deferrable=True
+ )
+ with time_machine.travel(pendulum.datetime(year=2024, month=8, day=1,
tz="UTC"), tick=False):
+ context = op.render_template_fields({})
+ with pytest.raises(
+ ValueError, match=re.escape("invalid literal for int() with
base 10: 'nothing'")
+ ):
+ op.execute(context)