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()