This is an automated email from the ASF dual-hosted git repository.

ferruzzi 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 e00e53257e9 Add SerializedVariableInterval for deadline alerts (#71802)
e00e53257e9 is described below

commit e00e53257e90a5380731cf5e148d2db57961f417
Author: SameerMesiah97 <[email protected]>
AuthorDate: Sat Aug 29 04:48:16 2026 +0100

    Add SerializedVariableInterval for deadline alerts (#71802)
---
 airflow-core/newsfragments/71802.bugfix.rst        |  1 +
 airflow-core/src/airflow/serialization/decoders.py | 10 +++-
 .../src/airflow/serialization/definitions/dag.py   |  9 ++-
 .../airflow/serialization/definitions/deadline.py  | 28 ++++++++-
 airflow-core/tests/unit/models/test_dagrun.py      |  7 +--
 .../serialization/definitions/test_deadline.py     | 67 ++++++++++++++++++++++
 .../unit/serialization/test_serialized_objects.py  | 33 +++++++++--
 generated/known_sdk_imports_in_core.txt            |  3 +-
 task-sdk/src/airflow/sdk/definitions/deadline.py   | 11 +++-
 .../tests/task_sdk/definitions/test_deadline.py    | 15 ++---
 10 files changed, 156 insertions(+), 28 deletions(-)

diff --git a/airflow-core/newsfragments/71802.bugfix.rst 
b/airflow-core/newsfragments/71802.bugfix.rst
new file mode 100644
index 00000000000..0c173c40e4d
--- /dev/null
+++ b/airflow-core/newsfragments/71802.bugfix.rst
@@ -0,0 +1 @@
+Fixed ``VariableInterval`` deadline intervals to allow zero and negative 
offsets, matching the existing ``timedelta`` interval semantics. 
``VariableInterval`` is now converted to the core-side 
``SerializedVariableInterval`` representation during deadline deserialization 
and resolved during deadline evaluation.
diff --git a/airflow-core/src/airflow/serialization/decoders.py 
b/airflow-core/src/airflow/serialization/decoders.py
index 49f88480b57..be5812ac24b 100644
--- a/airflow-core/src/airflow/serialization/decoders.py
+++ b/airflow-core/src/airflow/serialization/decoders.py
@@ -38,6 +38,7 @@ from airflow.serialization.definitions.deadline import (
     DeadlineAlertFields,
     SerializedDeadlineAlert,
     SerializedReferenceModels,
+    SerializedVariableInterval,
 )
 from airflow.serialization.enums import DagAttributeTypes as DAT, Encoding
 from airflow.serialization.helpers import (
@@ -203,7 +204,7 @@ def decode_deadline_alert(encoded_data: dict):
             "from a version that supports VariableInterval. Downgrade is not 
fully reversible."
         )
 
-    interval: datetime.timedelta | VariableInterval
+    interval: datetime.timedelta | SerializedVariableInterval
 
     # Backward compatibility: previously interval was stored as 
total_seconds() (float/int).
     # Handle numeric values by converting to timedelta.
@@ -211,8 +212,13 @@ def decode_deadline_alert(encoded_data: dict):
         interval = datetime.timedelta(seconds=raw_interval)
     else:
         deserialized = deserialize(raw_interval)
-        if isinstance(deserialized, (datetime.timedelta, VariableInterval)):
+
+        if isinstance(deserialized, datetime.timedelta):
+            interval = deserialized
+        elif isinstance(deserialized, SerializedVariableInterval):
             interval = deserialized
+        elif isinstance(deserialized, VariableInterval):
+            interval = SerializedVariableInterval(key=deserialized.key)
         else:
             raise TypeError(f"Invalid interval type: 
{type(deserialized).__name__}")
 
diff --git a/airflow-core/src/airflow/serialization/definitions/dag.py 
b/airflow-core/src/airflow/serialization/definitions/dag.py
index 8ed0fee2cca..5b3a5352430 100644
--- a/airflow-core/src/airflow/serialization/definitions/dag.py
+++ b/airflow-core/src/airflow/serialization/definitions/dag.py
@@ -50,9 +50,12 @@ from airflow.models.deadline import Deadline
 from airflow.models.deadline_alert import DeadlineAlert as DeadlineAlertModel
 from airflow.models.taskinstancekey import TaskInstanceKey
 from airflow.models.tasklog import LogTemplate
-from airflow.sdk.definitions.deadline import VariableInterval
 from airflow.serialization.decoders import decode_deadline_alert
-from airflow.serialization.definitions.deadline import DeadlineAlertFields, 
SerializedReferenceModels
+from airflow.serialization.definitions.deadline import (
+    DeadlineAlertFields,
+    SerializedReferenceModels,
+    SerializedVariableInterval,
+)
 from airflow.serialization.definitions.param import SerializedParamsDict
 from airflow.serialization.enums import DagAttributeTypes as DAT, Encoding
 from airflow.timetables.base import DagRunInfo, DataInterval, TimeRestriction
@@ -754,7 +757,7 @@ class SerializedDAG:
 
             interval = deserialized_deadline_alert.interval
 
-            if isinstance(interval, VariableInterval):
+            if isinstance(interval, SerializedVariableInterval):
                 interval = interval.resolve()
 
             if isinstance(deserialized_deadline_alert.reference, 
SerializedReferenceModels.TYPES.DAGRUN):
diff --git a/airflow-core/src/airflow/serialization/definitions/deadline.py 
b/airflow-core/src/airflow/serialization/definitions/deadline.py
index 968d69e3037..e5108c6267c 100644
--- a/airflow-core/src/airflow/serialization/definitions/deadline.py
+++ b/airflow-core/src/airflow/serialization/definitions/deadline.py
@@ -28,6 +28,7 @@ from sqlalchemy import select
 
 from airflow._shared.timezones import timezone
 from airflow.models.deadline import classproperty
+from airflow.models.variable import Variable
 from airflow.utils.log.logging_mixin import LoggingMixin
 from airflow.utils.session import provide_session
 from airflow.utils.sqlalchemy import get_dialect_name
@@ -38,8 +39,6 @@ if TYPE_CHECKING:
     from sqlalchemy import ColumnElement
     from sqlalchemy.orm import Session
 
-    from airflow.sdk.definitions.deadline import VariableInterval
-
 logger = logging.getLogger(__name__)
 
 
@@ -380,11 +379,34 @@ def _fetch_from_db(column, *, session: Session, dag_id: 
str, run_id: str) -> dat
     return result
 
 
[email protected](frozen=True)
+class SerializedVariableInterval:
+    """Core-side serialized representation of a variable-backed deadline 
interval."""
+
+    key: str
+
+    def resolve(self) -> timedelta:
+
+        try:
+            value = Variable.get(self.key)
+        except KeyError as e:
+            raise ValueError(f"VariableInterval '{self.key}' not found") from e
+
+        try:
+            seconds = int(value)
+        except (TypeError, ValueError) as e:
+            raise ValueError(
+                f"VariableInterval '{self.key}' must be an integer (seconds), 
got: {value!r}"
+            ) from e
+
+        return timedelta(seconds=seconds)
+
+
 @attrs.define
 class SerializedDeadlineAlert:
     """Serialized representation of a deadline alert."""
 
     reference: SerializedReferenceModels.SerializedBaseDeadlineReference
-    interval: timedelta | VariableInterval
+    interval: timedelta | SerializedVariableInterval
     callback: Any
     name: str | None = None
diff --git a/airflow-core/tests/unit/models/test_dagrun.py 
b/airflow-core/tests/unit/models/test_dagrun.py
index c46c30d1194..c30c1c7c94b 100644
--- a/airflow-core/tests/unit/models/test_dagrun.py
+++ b/airflow-core/tests/unit/models/test_dagrun.py
@@ -61,6 +61,7 @@ from airflow.models.taskinstance import TaskInstance, 
TaskInstanceNote, clear_ta
 from airflow.models.taskmap import TaskMap
 from airflow.models.taskreschedule import TaskReschedule
 from airflow.models.trigger import Trigger
+from airflow.models.variable import Variable
 from airflow.providers.standard.operators.bash import BashOperator
 from airflow.providers.standard.operators.empty import EmptyOperator
 from airflow.providers.standard.operators.python import PythonOperator, 
ShortCircuitOperator
@@ -76,8 +77,6 @@ from airflow.sdk import (
 )
 from airflow.sdk.definitions.callback import AsyncCallback
 from airflow.sdk.definitions.deadline import DeadlineAlert, DeadlineReference, 
VariableInterval
-from airflow.sdk.definitions.variable import Variable
-from airflow.sdk.exceptions import AirflowRuntimeError
 from airflow.serialization.definitions.deadline import 
SerializedReferenceModels
 from airflow.serialization.serialized_objects import LazyDeserializedDAG
 from airflow.settings import get_policy_plugin_manager
@@ -1532,7 +1531,7 @@ class TestDagRun:
         )
         dag_run.dag = scheduler_dag
 
-        # First update resolve interval to "5".
+        # First update resolves interval to "60".
         dag_run.update_state(session=session)
 
         deadline = session.execute(select(Deadline)).scalars().one_or_none()
@@ -1556,7 +1555,7 @@ class TestDagRun:
         with mock.patch.object(
             Variable,
             "get",
-            side_effect=AirflowRuntimeError(mock_err),
+            side_effect=KeyError(mock_err),
         ):
             future_date = datetime.datetime.now() + 
datetime.timedelta(days=365)
 
diff --git a/airflow-core/tests/unit/serialization/definitions/test_deadline.py 
b/airflow-core/tests/unit/serialization/definitions/test_deadline.py
new file mode 100644
index 00000000000..3e12fe0830c
--- /dev/null
+++ b/airflow-core/tests/unit/serialization/definitions/test_deadline.py
@@ -0,0 +1,67 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+from __future__ import annotations
+
+from datetime import timedelta
+
+import pytest
+
+from airflow.models.variable import Variable
+from airflow.serialization.definitions.deadline import 
SerializedVariableInterval
+
+
+class TestVariableInterval:
+    @pytest.mark.parametrize(
+        ("value", "expected"),
+        [
+            ("3", timedelta(seconds=3)),
+            ("10", timedelta(seconds=10)),
+            ("05", timedelta(seconds=5)),
+            ("0", timedelta(0)),
+            ("-5", timedelta(seconds=-5)),
+        ],
+    )
+    def test_resolve_valid(self, mocker, value, expected):
+        mocker.patch.object(Variable, "get", return_value=value)
+
+        interval = SerializedVariableInterval(key="test_interval")
+
+        assert interval.resolve() == expected
+
+    @pytest.mark.parametrize(
+        ("value", "raise_missing", "match"),
+        [
+            (None, True, "not found"),
+            ("abc", False, "must be an integer"),
+            ("", False, "must be an integer"),
+        ],
+    )
+    def test_resolve_invalid(self, mocker, value, raise_missing, match):
+        if raise_missing:
+            mocker.patch.object(
+                Variable,
+                "get",
+                side_effect=KeyError("test_interval"),
+            )
+        else:
+            mocker.patch.object(Variable, "get", return_value=value)
+
+        interval = SerializedVariableInterval(key="test_interval")
+
+        with pytest.raises(ValueError, match=match):
+            interval.resolve()
diff --git a/airflow-core/tests/unit/serialization/test_serialized_objects.py 
b/airflow-core/tests/unit/serialization/test_serialized_objects.py
index da9a8316585..3a84c70d29b 100644
--- a/airflow-core/tests/unit/serialization/test_serialized_objects.py
+++ b/airflow-core/tests/unit/serialization/test_serialized_objects.py
@@ -79,6 +79,7 @@ from airflow.sdk.definitions.deadline import (
     AsyncCallback,
     DeadlineAlert,
     DeadlineReference,
+    VariableInterval,
 )
 from airflow.sdk.definitions.decorators import task
 from airflow.sdk.definitions.operator_resources import Resources
@@ -95,7 +96,11 @@ from airflow.serialization.definitions.assets import (
     SerializedAssetBase,
     SerializedAssetRef,
 )
-from airflow.serialization.definitions.deadline import DeadlineAlertFields, 
SerializedDeadlineAlert
+from airflow.serialization.definitions.deadline import (
+    DeadlineAlertFields,
+    SerializedDeadlineAlert,
+    SerializedVariableInterval,
+)
 from airflow.serialization.encoders import ensure_serialized_asset, 
ensure_serialized_deadline_alert
 from airflow.serialization.enums import DagAttributeTypes as DAT, Encoding
 from airflow.serialization.helpers import PartitionMapperNotFound
@@ -492,13 +497,33 @@ def test_serialize_deserialize_connection():
 
 
 @pytest.mark.parametrize("reference", REFERENCE_TYPES)
-def test_serialize_deserialize_deadline_alert(reference):
[email protected](
+    ("interval", "expected_interval"),
+    [
+        pytest.param(
+            timedelta(hours=1),
+            timedelta(hours=1),
+            id="timedelta",
+        ),
+        pytest.param(
+            VariableInterval("deadline_seconds"),
+            SerializedVariableInterval("deadline_seconds"),
+            id="sdk_variable_interval",
+        ),
+        pytest.param(
+            SerializedVariableInterval("deadline_seconds"),
+            SerializedVariableInterval("deadline_seconds"),
+            id="serialized_variable_interval",
+        ),
+    ],
+)
+def test_serialize_deserialize_deadline_alert(reference, interval, 
expected_interval):
     public_deadline_alert_fields = {
         field.lower() for field in vars(DeadlineAlertFields) if not 
field.startswith("_")
     }
     original = DeadlineAlert(
         reference=reference,
-        interval=timedelta(hours=1),
+        interval=interval,
         callback=AsyncCallback(empty_callback_for_deadline, 
kwargs=TEST_CALLBACK_KWARGS),
     )
 
@@ -509,7 +534,7 @@ def test_serialize_deserialize_deadline_alert(reference):
 
     deserialized = BaseSerialization.deserialize(serialized)
     assert deserialized.reference.serialize_reference() == 
reference.serialize_reference()
-    assert deserialized.interval == original.interval
+    assert deserialized.interval == expected_interval
     assert deserialized.callback == original.callback
 
 
diff --git a/generated/known_sdk_imports_in_core.txt 
b/generated/known_sdk_imports_in_core.txt
index f0aaae24ac2..85bd4863e64 100644
--- a/generated/known_sdk_imports_in_core.txt
+++ b/generated/known_sdk_imports_in_core.txt
@@ -25,8 +25,7 @@ airflow-core/src/airflow/providers_manager.py::5
 airflow-core/src/airflow/secrets/__init__.py::1
 airflow-core/src/airflow/serialization/decoders.py::2
 airflow-core/src/airflow/serialization/definitions/baseoperator.py::1
-airflow-core/src/airflow/serialization/definitions/dag.py::2
-airflow-core/src/airflow/serialization/definitions/deadline.py::1
+airflow-core/src/airflow/serialization/definitions/dag.py::1
 airflow-core/src/airflow/serialization/definitions/mappedoperator.py::5
 airflow-core/src/airflow/serialization/encoders.py::11
 airflow-core/src/airflow/serialization/serialized_objects.py::17
diff --git a/task-sdk/src/airflow/sdk/definitions/deadline.py 
b/task-sdk/src/airflow/sdk/definitions/deadline.py
index f3a14aeb9cc..4aabc493fb4 100644
--- a/task-sdk/src/airflow/sdk/definitions/deadline.py
+++ b/task-sdk/src/airflow/sdk/definitions/deadline.py
@@ -17,6 +17,7 @@
 from __future__ import annotations
 
 import logging
+import warnings
 from abc import ABC
 from dataclasses import dataclass
 from datetime import datetime, timedelta
@@ -438,6 +439,13 @@ class VariableInterval:
     key: str
 
     def resolve(self) -> timedelta:
+        warnings.warn(
+            "VariableInterval.resolve() is deprecated and will be removed in a 
future release. "
+            "Deadline interval resolution is handled internally during 
deadline evaluation.",
+            DeprecationWarning,
+            stacklevel=2,
+        )
+
         try:
             value = Variable.get(self.key)
         except AirflowRuntimeError as e:
@@ -450,7 +458,4 @@ class VariableInterval:
                 f"VariableInterval '{self.key}' must be an integer (seconds), 
got: {value!r}"
             ) from e
 
-        if seconds <= 0:
-            raise ValueError(f"VariableInterval '{self.key}' must be > 0, got: 
{seconds}")
-
         return timedelta(seconds=seconds)
diff --git a/task-sdk/tests/task_sdk/definitions/test_deadline.py 
b/task-sdk/tests/task_sdk/definitions/test_deadline.py
index b104980e4c9..47cc6a874d8 100644
--- a/task-sdk/tests/task_sdk/definitions/test_deadline.py
+++ b/task-sdk/tests/task_sdk/definitions/test_deadline.py
@@ -173,7 +173,9 @@ class TestVariableInterval:
         [
             ("3", timedelta(seconds=3)),
             ("10", timedelta(seconds=10)),
-            ("05", timedelta(seconds=5)),  # leading zero
+            ("05", timedelta(seconds=5)),
+            ("0", timedelta(0)),
+            ("-5", timedelta(seconds=-5)),
         ],
     )
     def test_resolve_valid(self, mocker, value, expected):
@@ -181,7 +183,8 @@ class TestVariableInterval:
 
         interval = VariableInterval(key="test_interval")
 
-        assert interval.resolve() == expected
+        with pytest.warns(DeprecationWarning, 
match="VariableInterval.resolve"):
+            assert interval.resolve() == expected
 
     @pytest.mark.parametrize(
         ("value", "raise_runtime", "match"),
@@ -189,12 +192,9 @@ class TestVariableInterval:
             (None, True, "not found"),
             ("abc", False, "must be an integer"),
             ("", False, "must be an integer"),
-            ("0", False, "must be > 0"),
-            ("-5", False, "must be > 0"),
         ],
     )
     def test_resolve_invalid(self, mocker, value, raise_runtime, match):
-
         if raise_runtime:
             mock_err = mock.Mock()
             mock_err.error.value = "MISSING"
@@ -210,5 +210,6 @@ class TestVariableInterval:
 
         interval = VariableInterval(key="test_interval")
 
-        with pytest.raises(ValueError, match=match):
-            interval.resolve()
+        with pytest.warns(DeprecationWarning, 
match="VariableInterval.resolve"):
+            with pytest.raises(ValueError, match=match):
+                interval.resolve()

Reply via email to