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)

Reply via email to