diff --git a/providers/standard/docs/sensors/datetime.rst b/providers/standard/docs/sensors/datetime.rst index 2a45ecb6bd869..3cbe87ae01c3b 100644 --- a/providers/standard/docs/sensors/datetime.rst +++ b/providers/standard/docs/sensors/datetime.rst @@ -39,9 +39,14 @@ To run the sensor in deferrable mode, set ``deferrable=True``. See :ref:`deferri TimeSensor ========== -Use the :class:`~airflow.providers.standard.sensors.time_sensor.TimeSensor` to end sensing after time specified. ``TimeSensor`` can be run in deferrable mode, if a Triggerer is available. +Use the :class:`~airflow.providers.standard.sensors.time.TimeSensor` to end sensing after time specified. ``TimeSensor`` can be run in deferrable mode, if a Triggerer is available. -Time will be evaluated against ``data_interval_end`` if present for the Dag run, otherwise ``run_after`` will be used. +Time is evaluated against the wall-clock date in the Dag's timezone at execution time (poke, deferral, or trigger start), not against ``data_interval_end`` or ``run_after``. For interval-relative behavior, use :class:`~airflow.providers.standard.sensors.time_delta.TimeDeltaSensor`. + +When ``start_from_trigger=True``, the sensor starts on the triggerer via +:class:`~airflow.providers.standard.triggers.temporal.TimeOfDayTrigger`. The trigger +stores only the parse-stable ``target_time`` and ``tz``; the concrete moment is +resolved when the trigger starts, so repeated Dag parses do not create new Dag versions. .. exampleinclude:: /../src/airflow/providers/standard/example_dags/example_sensors.py :language: python diff --git a/providers/standard/src/airflow/providers/standard/sensors/time.py b/providers/standard/src/airflow/providers/standard/sensors/time.py index 2e3b19eb89940..9135de22c3fbd 100644 --- a/providers/standard/src/airflow/providers/standard/sensors/time.py +++ b/providers/standard/src/airflow/providers/standard/sensors/time.py @@ -17,14 +17,17 @@ # under the License. from __future__ import annotations -import dataclasses import datetime import warnings from typing import TYPE_CHECKING, Any from airflow.exceptions import AirflowProviderDeprecationWarning from airflow.providers.common.compat.sdk import BaseSensorOperator, conf, timezone -from airflow.providers.standard.triggers.temporal import DateTimeTrigger +from airflow.providers.standard.triggers.temporal import ( + DateTimeTrigger, + resolve_time_of_day_moment, + serializable_timezone, +) from airflow.triggers.base import StartTriggerArgs if TYPE_CHECKING: @@ -35,8 +38,24 @@ class TimeSensor(BaseSensorOperator): """ Waits until the specified time of the day. + The time is evaluated against the wall-clock date in the Dag's timezone at + execution time (poke / deferral / trigger start), not at Dag-parse time. + This avoids dag_version churn from baking an absolute ``moment`` into + serialized ``start_trigger_args``. + + When ``start_from_trigger=True``, the sensor starts directly on the triggerer + via :class:`~airflow.providers.standard.triggers.temporal.TimeOfDayTrigger`, + which stores only the parse-stable ``target_time`` + tz and resolves + the concrete moment when the trigger actually starts. + :param target_time: time after which the job succeeds :param deferrable: whether to defer execution + :param start_from_trigger: Start the task directly from the triggerer without + going into the worker. + :param end_from_trigger: End the task directly from the triggerer without + going into the worker. + :param trigger_kwargs: Accepted for API compatibility with other sensors that + support dynamic task mapping into start-from-trigger; not used by TimeSensor. .. seealso:: For more information on how to use this sensor, take a look at the guide: @@ -44,13 +63,7 @@ class TimeSensor(BaseSensorOperator): """ - start_trigger_args = StartTriggerArgs( - trigger_cls="airflow.providers.standard.triggers.temporal.DateTimeTrigger", - trigger_kwargs={"moment": "", "end_from_trigger": False}, - next_method="execute_complete", - next_kwargs=None, - timeout=None, - ) + start_trigger_args = None start_from_trigger = False def __init__( @@ -64,37 +77,68 @@ def __init__( **kwargs, ) -> None: super().__init__(**kwargs) - - # Create a "date-aware" timestamp that will be used as the "target_datetime". This is a requirement - # of the DateTimeTrigger - - # Get date considering dag.timezone - aware_time = timezone.coerce_datetime( - datetime.datetime.combine( - datetime.datetime.now(self.dag.timezone), target_time, self.dag.timezone - ) - ) - - # Now that the dag's timezone has made the datetime timezone aware, we need to convert to UTC - self.target_datetime = timezone.convert_to_utc(aware_time) + # Wall-clock only; tzinfo is stripped so serialized target_time is deterministic. + if isinstance(target_time, datetime.time) and target_time.tzinfo is not None: + self.target_time = target_time.replace(tzinfo=None) + else: + self.target_time = target_time self.deferrable = deferrable self.start_from_trigger = start_from_trigger self.end_from_trigger = end_from_trigger + # Cached for this task attempt so a local-date rollover does not change the target. + self._cached_target_datetime: datetime.datetime | None = None if self.start_from_trigger: - # Replaced rather than mutated: ``start_trigger_args`` is a class attribute, so - # assigning through it would overwrite the arguments of every other task built - # from this operator. - self.start_trigger_args = dataclasses.replace( - self.start_trigger_args, - trigger_kwargs=dict(moment=self.target_datetime, end_from_trigger=self.end_from_trigger), + dag = self._dag + if dag is None: + raise ValueError( + "TimeSensor(start_from_trigger=True) requires the sensor to be attached to a Dag " + "so the timezone is known." + ) + # Parse-stable kwargs only (no datetime.now()); moment is resolved when the trigger starts. + self.start_trigger_args = StartTriggerArgs( + trigger_cls="airflow.providers.standard.triggers.temporal.TimeOfDayTrigger", + trigger_kwargs={ + "target_time": self.target_time.isoformat(), + "tz": serializable_timezone(dag.timezone), + "end_from_trigger": self.end_from_trigger, + }, + next_method="execute_complete", + next_kwargs=None, + timeout=None, ) + def _resolve_target_datetime(self) -> datetime.datetime: + """Compute the UTC moment for target_time on today's date in the Dag timezone.""" + dag = self._dag + # Unattached sensors (unit tests / early construction) use UTC. + tz: str | int | datetime.tzinfo = "UTC" if dag is None else dag.timezone + return resolve_time_of_day_moment(self.target_time, tz=tz) + + def _get_target_datetime(self) -> datetime.datetime: + """Return the target moment, computing and caching it once per attempt.""" + if self._cached_target_datetime is None: + self._cached_target_datetime = self._resolve_target_datetime() + return self._cached_target_datetime + + @property + def target_datetime(self) -> datetime.datetime: + """ + Resolved target datetime in UTC. + + Computed on first access (or first execute/poke) from ``target_time`` and + the Dag timezone, then cached for the life of this instance. Two reads on + either side of midnight therefore return the *same* date for a given + attempt; a fresh task instance re-resolves against "today". + """ + return self._get_target_datetime() + def execute(self, context: Context) -> None: + moment = self._get_target_datetime() if self.deferrable: self.defer( trigger=DateTimeTrigger( - moment=self.target_datetime, # This needs to be an aware timestamp + moment=moment, end_from_trigger=self.end_from_trigger, ), method_name="execute_complete", @@ -106,10 +150,9 @@ def execute_complete(self, context: Context, event: Any = None) -> None: return None def poke(self, context: Context) -> bool: - self.log.info("Checking if the time (%s) has come", self.target_datetime) - - # self.target_date has been converted to UTC, so we do not need to convert timezone - return timezone.utcnow() > self.target_datetime + target_datetime = self._get_target_datetime() + self.log.info("Checking if the time (%s) has come", target_datetime) + return timezone.utcnow() > target_datetime class TimeSensorAsync(TimeSensor): diff --git a/providers/standard/src/airflow/providers/standard/triggers/temporal.py b/providers/standard/src/airflow/providers/standard/triggers/temporal.py index cfc206ff3131d..a3897f19f0868 100644 --- a/providers/standard/src/airflow/providers/standard/triggers/temporal.py +++ b/providers/standard/src/airflow/providers/standard/triggers/temporal.py @@ -22,11 +22,130 @@ from typing import Any import pendulum +from pendulum.tz.timezone import FixedTimezone, Timezone from airflow.providers.common.compat.sdk import timezone from airflow.triggers.base import BaseTrigger, TaskSuccessEvent, TriggerEvent +def _parse_timezone(value: str | int | datetime.tzinfo) -> Timezone | FixedTimezone: + """Return a pendulum timezone from an IANA name, fixed offset (seconds), or existing tzinfo.""" + if isinstance(value, (Timezone, FixedTimezone)): + return value + if isinstance(value, datetime.tzinfo): + # Generic tzinfo (zoneinfo, datetime.timezone): rebuild a pendulum zone from its + # IANA name or fixed offset so pendulum's DST arithmetic gets a native zone type. + return pendulum.timezone(serializable_timezone(value)) + return pendulum.timezone(value) + + +def serializable_timezone(tzinfo: datetime.tzinfo | None) -> str | int: + """ + Encode a tzinfo as a value that round-trips through ``pendulum.timezone`` / parse_timezone. + + Named zones become their IANA name (e.g. ``Asia/Singapore``). Fixed-offset zones become + the offset in seconds (int), which is what Airflow's own timezone serializer uses. + UTC / zero-offset is always the string ``UTC`` for stable serialization. + """ + if tzinfo is None: + return "UTC" + if isinstance(tzinfo, FixedTimezone): + if tzinfo.offset == 0: + return "UTC" + return tzinfo.offset + name = getattr(tzinfo, "name", None) or getattr(tzinfo, "key", None) or getattr(tzinfo, "zone", None) + if name: + if name in ("UTC", "utc", "+00:00"): + return "UTC" + return name + offset = tzinfo.utcoffset(None) + if offset is not None: + total = int(offset.total_seconds()) + return "UTC" if total == 0 else total + return "UTC" + + +def _coerce_target_time(target_time: datetime.time | str) -> datetime.time: + """Accept ``datetime.time`` or an ISO time string (Airflow serializes time as str).""" + if isinstance(target_time, str): + return datetime.time.fromisoformat(target_time) + if isinstance(target_time, datetime.time): + # Drop tzinfo so storage / comparison is wall-clock-only; zone is separate. + if target_time.tzinfo is not None: + return target_time.replace(tzinfo=None) + return target_time + raise TypeError(f"Expected datetime.time or str for target_time. Got {type(target_time)}") + + +def resolve_time_of_day_moment( + target_time: datetime.time | str, + *, + tz: str | int | datetime.tzinfo = "UTC", + as_of: datetime.datetime | None = None, +) -> pendulum.DateTime: + """ + Resolve ``target_time`` on "today" in ``tz`` to a UTC-aware moment. + + Semantics: + + - **Already passed today**: still returns today's occurrence (caller succeeds immediately). + Does *not* roll forward to the next day. + - **Non-existent local time** (spring-forward gap, e.g. 02:30 America/New_York on DST start): + shifts forward to the next valid local time (e.g. 03:30). + - **Ambiguous local time** (fall-back overlap, e.g. 01:30 on DST end): uses ``fold=0`` + (the first occurrence). + - Moment is computed from ``as_of`` (default: now) so callers can cache per attempt and + avoid midnight drift when re-checking within the same run. + """ + wall_time = _coerce_target_time(target_time) + tzinfo = _parse_timezone(tz) + + if as_of is None: + as_of = timezone.utcnow() + local_now = pendulum.instance(as_of).in_timezone(tzinfo) + + # pendulum.datetime shifts non-existent (gap) times forward to the next valid wall time. + moment_local = pendulum.datetime( + local_now.year, + local_now.month, + local_now.day, + wall_time.hour, + wall_time.minute, + wall_time.second, + wall_time.microsecond, + tz=tzinfo, + ) + + # If the wall clock was preserved, check for ambiguous (fold) times and prefer fold=0. + if ( + moment_local.hour, + moment_local.minute, + moment_local.second, + moment_local.microsecond, + ) == ( + wall_time.hour, + wall_time.minute, + wall_time.second, + wall_time.microsecond, + ): + dt0 = datetime.datetime( + local_now.year, + local_now.month, + local_now.day, + wall_time.hour, + wall_time.minute, + wall_time.second, + wall_time.microsecond, + tzinfo=tzinfo, + fold=0, + ) + dt1 = dt0.replace(fold=1) + if dt0.utcoffset() != dt1.utcoffset(): + moment_local = pendulum.instance(dt0) + + return timezone.convert_to_utc(moment_local) + + class DateTimeTrigger(BaseTrigger): """ Trigger based on a datetime. @@ -86,6 +205,58 @@ async def run(self) -> AsyncIterator[TriggerEvent]: yield TriggerEvent(self.moment) +class TimeOfDayTrigger(DateTimeTrigger): + """ + Trigger that fires once the wall-clock reaches ``target_time`` in ``tz``. + + The concrete moment is resolved at construction (triggerer start), not at + Dag-parse time. That keeps ``start_trigger_args`` parse-stable while + preserving ``TimeSensor(start_from_trigger=True)``. + + ``serialize()`` includes the resolved ``moment`` so a reconstruct cannot + re-resolve "today" after midnight. ``start_trigger_args`` still omit + ``moment`` so Dag serialization stays parse-stable. + + ``target_time`` is accepted as an ISO time string so trigger kwargs remain + JSON/serde-safe (``datetime.time`` is not accepted by Airflow's trigger serde). + + :param target_time: wall-clock time of day (``datetime.time`` or ISO time string) + :param tz: IANA name (str) or fixed offset in seconds (int); must round-trip + through ``pendulum.timezone``. Named ``tz`` to avoid shadowing the + ``timezone`` module imported from the SDK compat layer. + :param end_from_trigger: whether the trigger should mark the task successful after + the time condition is reached + :param moment: optional pre-resolved UTC moment; used on serialize reconstruct + """ + + def __init__( + self, + target_time: datetime.time | str, + *, + tz: str | int = "UTC", + end_from_trigger: bool = False, + moment: datetime.datetime | None = None, + ) -> None: + wall = _coerce_target_time(target_time) + super().__init__( + moment=moment if moment is not None else resolve_time_of_day_moment(wall, tz=tz), + end_from_trigger=end_from_trigger, + ) + self.target_time: str = wall.isoformat() + self.tz: str | int = tz + + def serialize(self) -> tuple[str, dict[str, Any]]: + return ( + "airflow.providers.standard.triggers.temporal.TimeOfDayTrigger", + { + "target_time": self.target_time, + "tz": self.tz, + "end_from_trigger": self.end_from_trigger, + "moment": self.moment, + }, + ) + + class TimeDeltaTrigger(DateTimeTrigger): """ Create DateTimeTriggers based on delays. diff --git a/providers/standard/tests/unit/standard/sensors/test_time.py b/providers/standard/tests/unit/standard/sensors/test_time.py index 2b44c61d474b2..329230bfab24e 100644 --- a/providers/standard/tests/unit/standard/sensors/test_time.py +++ b/providers/standard/tests/unit/standard/sensors/test_time.py @@ -22,11 +22,18 @@ import pendulum import pytest import time_machine +from pendulum.tz.timezone import FixedTimezone +from airflow.exceptions import AirflowProviderDeprecationWarning from airflow.models.dag import DAG from airflow.providers.common.compat.sdk import TaskDeferred -from airflow.providers.standard.sensors.time import TimeSensor -from airflow.providers.standard.triggers.temporal import DateTimeTrigger +from airflow.providers.standard.sensors.time import TimeSensor, TimeSensorAsync +from airflow.providers.standard.triggers.temporal import ( + DateTimeTrigger, + TimeOfDayTrigger, + resolve_time_of_day_moment, +) +from airflow.triggers.base import StartTriggerArgs from tests_common.test_utils.compat import timezone @@ -110,8 +117,6 @@ def test_task_is_deferred(self): ): op = TimeSensor(task_id="test", target_time=time(10, 0), deferrable=True) - # This should be converted to UTC in the __init__. Since there is no default timezone, it will become - # aware, but note changed assert op.target_datetime.utcoffset() is not None with pytest.raises(TaskDeferred) as exc_info: @@ -127,7 +132,7 @@ def test_execute_complete_accepts_event(self): with DAG( dag_id="test_execute_complete_accepts_event", schedule=None, - start_date=datetime(2020, 1, 1), # Matches above + start_date=datetime(2020, 1, 1), ): op = TimeSensor(task_id="test", target_time=time(10, 0), deferrable=True) @@ -136,21 +141,255 @@ def test_execute_complete_accepts_event(self): except TypeError as e: pytest.fail(f"TypeError raised: {e}") + @time_machine.travel("2020-07-07 15:00:00", tick=False) + def test_target_already_passed_today_succeeds(self): + with DAG("e1", schedule=None, start_date=datetime(2020, 1, 1)): + op = TimeSensor(task_id="t", target_time=time(10, 0)) + assert op.poke({}) is True + assert op.target_datetime.date() == pendulum.date(2020, 7, 7) + + def test_target_datetime_cached_across_midnight(self): + with time_machine.travel("2020-07-07 23:59:00", tick=False): + with DAG("e2", schedule=None, start_date=datetime(2020, 1, 1)): + op = TimeSensor(task_id="t", target_time=time(10, 0)) + first = op.target_datetime + with time_machine.travel("2020-07-08 00:01:00", tick=False): + second = op.target_datetime + assert first == second + assert first.date() == pendulum.date(2020, 7, 7) + + @time_machine.travel("2024-03-10 12:00:00", tick=False) + def test_dst_spring_forward_shifts_forward(self): + ny = pendulum.timezone("America/New_York") + with DAG("e3", schedule=None, start_date=datetime(2024, 1, 1, tzinfo=ny)): + op = TimeSensor(task_id="t", target_time=time(2, 30)) + # 02:30 does not exist; pendulum shifts to 03:30 EDT = 07:30 UTC + assert op.target_datetime == pendulum.datetime(2024, 3, 10, 7, 30, tz="UTC") + + @time_machine.travel("2024-11-03 12:00:00", tick=False) + def test_dst_fall_back_uses_fold_zero(self): + ny = pendulum.timezone("America/New_York") + with DAG("e4", schedule=None, start_date=datetime(2024, 1, 1, tzinfo=ny)): + op = TimeSensor(task_id="t", target_time=time(1, 30)) + # fold=0 → EDT (UTC-4) → 01:30-04:00 = 05:30 UTC + assert op.target_datetime == pendulum.datetime(2024, 11, 3, 5, 30, tz="UTC") + + @time_machine.travel("2020-07-07 08:00:00", tick=False) + def test_no_dag_context_falls_back_to_utc(self): + op = TimeSensor(task_id="t", target_time=time(10, 0)) + assert op.target_datetime == pendulum.datetime(2020, 7, 7, 10, 0, tz="UTC") + assert op.poke({}) is False + + def test_start_from_trigger_without_dag_raises(self): + with pytest.raises(ValueError, match="attached to a Dag"): + TimeSensor(task_id="t", target_time=time(10, 0), start_from_trigger=True) + + @time_machine.travel("2020-07-07 00:00:00", tick=False) + def test_time_sensor_async_inherits_lazy_resolution(self): + with pytest.warns(AirflowProviderDeprecationWarning, match="TimeSensorAsync is deprecated"): + with DAG("e6", schedule=None, start_date=datetime(2020, 1, 1)): + op = TimeSensorAsync(task_id="t", target_time=time(10, 0)) + assert op.deferrable is True + with pytest.raises(TaskDeferred) as exc_info: + op.execute({}) + assert isinstance(exc_info.value.trigger, DateTimeTrigger) + assert exc_info.value.trigger.moment == pendulum.datetime(2020, 7, 7, 10) + + def test_fixed_and_named_tz_in_start_trigger_args(self): + fixed = FixedTimezone(19800) # +05:30 + with DAG("e7_fixed", schedule=None, start_date=datetime(2020, 1, 1, tzinfo=fixed)): + op_fixed = TimeSensor( + task_id="fixed", + target_time=time(10, 30), + start_from_trigger=True, + end_from_trigger=True, + ) + assert op_fixed.start_trigger_args.trigger_kwargs["tz"] == 19800 + assert op_fixed.start_trigger_args.trigger_kwargs["target_time"] == "10:30:00" + assert op_fixed.start_trigger_args.trigger_kwargs["end_from_trigger"] is True + + named = pendulum.timezone("Asia/Kolkata") + with DAG("e7_named", schedule=None, start_date=datetime(2020, 1, 1, tzinfo=named)): + op_named = TimeSensor( + task_id="named", + target_time=time(10, 30), + start_from_trigger=True, + ) + assert op_named.start_trigger_args.trigger_kwargs["tz"] == "Asia/Kolkata" + + for op in (op_fixed, op_named): + kwargs = op.start_trigger_args.trigger_kwargs + trigger = TimeOfDayTrigger(**kwargs) + assert trigger.target_time == kwargs["target_time"] + assert trigger.tz == kwargs["tz"] + classpath, ser = trigger.serialize() + assert classpath.endswith("TimeOfDayTrigger") + restored = TimeOfDayTrigger(**ser) + assert restored.moment == trigger.moment + assert restored.target_time == kwargs["target_time"] + assert restored.tz == kwargs["tz"] + def test_start_trigger_args_are_not_shared_between_tasks(self): - """Each task must carry its own trigger arguments. + with DAG("e8", schedule=None, start_date=datetime(2020, 1, 1)): + op_a = TimeSensor(task_id="a", target_time=time(9, 0), start_from_trigger=True) + op_b = TimeSensor(task_id="b", target_time=time(17, 30), start_from_trigger=True) - ``start_trigger_args`` is a class attribute, so assigning through it made every task - built from this operator advertise the moment of whichever was constructed last. + assert TimeSensor.start_trigger_args is None + assert op_a.start_trigger_args is not op_b.start_trigger_args + assert op_a.start_trigger_args.trigger_kwargs["target_time"] == "09:00:00" + assert op_b.start_trigger_args.trigger_kwargs["target_time"] == "17:30:00" + + @staticmethod + def _serialized_dag_payload(dag) -> dict: + """Serialized-Dag dict on every supported Airflow version. + + ``SerializedDAG.to_dict`` exists on 2.11 and 3.0-3.2 but was removed on main; + ``LazyDeserializedDAG.from_dag`` only exists on 3.3+. Parse-time churn shows up + in this payload first, whichever API produces it. """ - with DAG( - dag_id="test_start_trigger_args_not_shared", - schedule=None, - start_date=datetime(2020, 1, 1), - ): - early = TimeSensor(task_id="early", target_time=time(1, 0), start_from_trigger=True) - late = TimeSensor(task_id="late", target_time=time(23, 0), start_from_trigger=True) + from airflow.serialization.serialized_objects import SerializedDAG + + if hasattr(SerializedDAG, "to_dict"): + return SerializedDAG.to_dict(dag) + from airflow.serialization.serialized_objects import LazyDeserializedDAG + + return LazyDeserializedDAG.from_dag(dag).data + + @staticmethod + def _find_start_trigger_args(obj): + """Locate encoded start_trigger_args in a Serialized Dag payload.""" + if isinstance(obj, dict): + if "start_trigger_args" in obj: + return obj["start_trigger_args"] + if obj.get("__type") == "START_TRIGGER_ARGS": + return obj + for value in obj.values(): + found = TestTimeSensor._find_start_trigger_args(value) + if found is not None: + return found + elif isinstance(obj, list): + for item in obj: + found = TestTimeSensor._find_start_trigger_args(item) + if found is not None: + return found + return None - assert early.start_trigger_args is not late.start_trigger_args - assert early.start_trigger_args.trigger_kwargs["moment"] == early.target_datetime - assert late.start_trigger_args.trigger_kwargs["moment"] == late.target_datetime - assert early.target_datetime != late.target_datetime + @staticmethod + def _trigger_kwargs_from_encoded(start_trigger_args: dict) -> dict: + raw = start_trigger_args["trigger_kwargs"] + if isinstance(raw, dict) and "__var" in raw: + return raw["__var"] + return raw + + def test_start_from_trigger_keeps_serialized_dag_hash_stable(self): + import hashlib + import json + + def build_payload_and_hash() -> tuple[dict, str]: + with DAG( + dag_id="test_start_from_trigger_hash_stable", + schedule=None, + start_date=datetime(2020, 1, 1), + ) as dag: + TimeSensor( + task_id="test", + target_time=time(10, 0), + start_from_trigger=True, + end_from_trigger=True, + ) + payload = self._serialized_dag_payload(dag) + digest = hashlib.md5(json.dumps(payload, sort_keys=True, default=str).encode()).hexdigest() + return payload, digest + + with time_machine.travel("2025-01-01 00:00:00", tick=False): + first_payload, first_hash = build_payload_and_hash() + with time_machine.travel("2025-06-15 12:34:56", tick=False): + second_payload, second_hash = build_payload_and_hash() + + first_sta = self._find_start_trigger_args(first_payload) + assert first_sta is not None, "serialized payload must include start_trigger_args" + first_kwargs = self._trigger_kwargs_from_encoded(first_sta) + assert "TimeOfDayTrigger" in first_sta["trigger_cls"] + assert first_kwargs["target_time"] == "10:00:00" + assert "tz" in first_kwargs + assert "moment" not in first_kwargs + + assert first_hash == second_hash + assert first_payload == second_payload + + def test_start_from_trigger_serialized_dict_identical_across_parses(self): + """SerializedDAG payload must be byte-identical when now() differs (E9).""" + + def build_data(): + with DAG( + dag_id="test_start_from_trigger_dict_stable", + schedule=None, + start_date=datetime(2020, 1, 1), + ) as dag: + TimeSensor( + task_id="test", + target_time=time(10, 0), + start_from_trigger=True, + ) + return self._serialized_dag_payload(dag) + + with time_machine.travel("2025-01-01 00:00:00", tick=False): + first = build_data() + with time_machine.travel("2025-12-31 23:59:59", tick=False): + second = build_data() + + sta = self._find_start_trigger_args(first) + assert sta is not None, "serialized payload must include start_trigger_args" + kwargs = self._trigger_kwargs_from_encoded(sta) + assert "TimeOfDayTrigger" in sta["trigger_cls"] + assert "moment" not in kwargs + assert kwargs["target_time"] == "10:00:00" + + assert first == second + + @time_machine.travel("2020-07-07 00:00:00", tick=False) + def test_end_from_trigger_propagates(self): + with DAG("e10", schedule=None, start_date=datetime(2020, 1, 1)): + op = TimeSensor( + task_id="t", + target_time=time(10, 0), + start_from_trigger=True, + end_from_trigger=True, + deferrable=True, + ) + assert op.start_trigger_args.trigger_kwargs["end_from_trigger"] is True + + with pytest.raises(TaskDeferred) as exc_info: + op.execute({}) + assert exc_info.value.trigger.end_from_trigger is True + + def test_start_from_trigger_uses_time_of_day_trigger(self): + with DAG("sft", schedule=None, start_date=datetime(2020, 1, 1)): + op = TimeSensor(task_id="t", target_time=time(10, 0), start_from_trigger=True) + assert op.start_from_trigger is True + assert isinstance(op.start_trigger_args, StartTriggerArgs) + assert op.start_trigger_args.trigger_cls.endswith("TimeOfDayTrigger") + assert "moment" not in op.start_trigger_args.trigger_kwargs + assert "timezone" not in op.start_trigger_args.trigger_kwargs + assert "tz" in op.start_trigger_args.trigger_kwargs + + def test_start_from_trigger_still_in_signature(self): + import inspect + + sig = inspect.signature(TimeSensor.__init__) + assert "start_from_trigger" in sig.parameters + assert sig.parameters["start_from_trigger"].default is False + + @time_machine.travel("2020-07-07 08:00:00", tick=False) + def test_poke_and_defer_match_time_of_day_trigger_moment(self): + with DAG("same_moment", schedule=None, start_date=datetime(2020, 1, 1)): + poke_op = TimeSensor(task_id="poke", target_time=time(10, 0)) + defer_op = TimeSensor(task_id="defer", target_time=time(10, 0), deferrable=True) + sft_op = TimeSensor(task_id="sft", target_time=time(10, 0), start_from_trigger=True) + as_of = timezone.utcnow() + kwargs = sft_op.start_trigger_args.trigger_kwargs + trigger_moment = resolve_time_of_day_moment(kwargs["target_time"], tz=kwargs["tz"], as_of=as_of) + assert poke_op.target_datetime == trigger_moment + with pytest.raises(TaskDeferred) as exc_info: + defer_op.execute({}) + assert exc_info.value.trigger.moment == trigger_moment diff --git a/providers/standard/tests/unit/standard/triggers/test_temporal.py b/providers/standard/tests/unit/standard/triggers/test_temporal.py index 9a3b2c423d451..260a57ecf0b30 100644 --- a/providers/standard/tests/unit/standard/triggers/test_temporal.py +++ b/providers/standard/tests/unit/standard/triggers/test_temporal.py @@ -24,7 +24,13 @@ import pytest from airflow.providers.common.compat.sdk import timezone -from airflow.providers.standard.triggers.temporal import DateTimeTrigger, TimeDeltaTrigger +from airflow.providers.standard.triggers.temporal import ( + DateTimeTrigger, + TimeDeltaTrigger, + TimeOfDayTrigger, + resolve_time_of_day_moment, + serializable_timezone, +) from airflow.triggers.base import TriggerEvent from airflow.utils.state import TaskInstanceState @@ -145,3 +151,98 @@ async def test_datetime_trigger_mocked(mock_sleep, mock_utcnow): result = trigger_task.result() assert isinstance(result, TriggerEvent) assert result.payload == trigger_moment + + +def test_time_of_day_trigger_serialization(): + as_of = pendulum.datetime(2020, 7, 7, 8, 0, tz="UTC") + with mock.patch( + "airflow.providers.standard.triggers.temporal.timezone.utcnow", + return_value=as_of, + ): + trigger = TimeOfDayTrigger( + target_time=datetime.time(10, 30, 45), + tz="Asia/Singapore", + end_from_trigger=True, + ) + classpath, kwargs = trigger.serialize() + assert classpath == "airflow.providers.standard.triggers.temporal.TimeOfDayTrigger" + expected_moment = pendulum.datetime(2020, 7, 7, 2, 30, 45, tz="UTC") # 10:30:45 +08 + assert kwargs == { + "target_time": "10:30:45", + "tz": "Asia/Singapore", + "end_from_trigger": True, + "moment": expected_moment, + } + restored = TimeOfDayTrigger(**kwargs) + assert restored.serialize() == (classpath, kwargs) + + +def test_time_of_day_trigger_serialize_stable_across_midnight(): + """Reconstruct from serialize() must keep the original moment after local midnight.""" + as_of = pendulum.datetime(2020, 7, 7, 23, 0, tz="UTC") + with mock.patch( + "airflow.providers.standard.triggers.temporal.timezone.utcnow", + return_value=as_of, + ): + trigger = TimeOfDayTrigger(target_time=datetime.time(23, 30), tz="UTC") + classpath, kwargs = trigger.serialize() + assert classpath.endswith("TimeOfDayTrigger") + original_moment = kwargs["moment"] + assert original_moment == pendulum.datetime(2020, 7, 7, 23, 30, tz="UTC") + + later = pendulum.datetime(2020, 7, 8, 0, 5, tz="UTC") + with mock.patch( + "airflow.providers.standard.triggers.temporal.timezone.utcnow", + return_value=later, + ): + restored = TimeOfDayTrigger(**kwargs) + assert restored.moment == original_moment + + +def test_time_of_day_trigger_fixed_offset_tz(): + trigger = TimeOfDayTrigger(target_time="07:00:00", tz=19800, end_from_trigger=False) + assert trigger.tz == 19800 + assert trigger.target_time == "07:00:00" + _, kwargs = trigger.serialize() + assert kwargs["tz"] == 19800 + assert kwargs["target_time"] == "07:00:00" + assert "moment" in kwargs + + +def test_resolve_time_of_day_moment_already_passed(): + as_of = pendulum.datetime(2020, 7, 7, 15, 0, tz="UTC") + moment = resolve_time_of_day_moment(datetime.time(10, 0), tz="UTC", as_of=as_of) + assert moment == pendulum.datetime(2020, 7, 7, 10, 0, tz="UTC") + + +def test_resolve_time_of_day_moment_dst_spring_forward(): + as_of = pendulum.datetime(2024, 3, 10, 12, 0, tz="UTC") + moment = resolve_time_of_day_moment(datetime.time(2, 30), tz="America/New_York", as_of=as_of) + # Gap → shift forward to 03:30 EDT = 07:30 UTC + assert moment == pendulum.datetime(2024, 3, 10, 7, 30, tz="UTC") + + +def test_resolve_time_of_day_moment_dst_fall_back_fold_zero(): + as_of = pendulum.datetime(2024, 11, 3, 12, 0, tz="UTC") + moment = resolve_time_of_day_moment(datetime.time(1, 30), tz="America/New_York", as_of=as_of) + # fold=0 → first occurrence (EDT, UTC-4) = 05:30 UTC + assert moment == pendulum.datetime(2024, 11, 3, 5, 30, tz="UTC") + + +def test_serializable_timezone_named_and_fixed(): + assert serializable_timezone(pendulum.timezone("UTC")) == "UTC" + assert serializable_timezone(pendulum.timezone("Asia/Singapore")) == "Asia/Singapore" + assert serializable_timezone(pendulum.tz.timezone.FixedTimezone(19800)) == 19800 + assert serializable_timezone(None) == "UTC" + + +@pytest.mark.asyncio +async def test_time_of_day_trigger_fires_for_past_target(): + """If today's target is already past, the trigger should fire immediately.""" + past = (timezone.utcnow() - datetime.timedelta(hours=1)).time().replace(microsecond=0) + trigger = TimeOfDayTrigger(target_time=past, tz="UTC", end_from_trigger=False) + trigger_task = asyncio.create_task(trigger.run().__anext__()) + await asyncio.sleep(0.5) + assert trigger_task.done() is True + result = trigger_task.result() + assert isinstance(result, TriggerEvent)