This is an automated email from the ASF dual-hosted git repository.
o-nikolas 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 f03b3bb5ee3 Fix deadline serialization, repr, and prune edge cases
(#70421)
f03b3bb5ee3 is described below
commit f03b3bb5ee35c42be8000325ee3b7d21b8998ef5
Author: Sean Ghaeli <[email protected]>
AuthorDate: Mon Aug 17 09:56:25 2026 -0700
Fix deadline serialization, repr, and prune edge cases (#70421)
Split out of #68919 per review: decoder __class_path routing, clear
error for missing __class_path, repr guards for severed dagrun and
dict-shaped interval, and prune guard for missed deadlines.
---
airflow-core/src/airflow/models/deadline.py | 11 +-
airflow-core/src/airflow/models/deadline_alert.py | 7 +-
airflow-core/src/airflow/serialization/decoders.py | 7 +-
.../airflow/serialization/definitions/deadline.py | 10 +-
airflow-core/tests/unit/models/test_deadline.py | 17 +++
.../tests/unit/models/test_deadline_alert.py | 30 +++++
.../tests/unit/models/test_prune_deadlines.py | 126 +++++++++++++++++++++
.../test_deadline_reference_registry.py | 82 ++++++++++++++
8 files changed, 282 insertions(+), 8 deletions(-)
diff --git a/airflow-core/src/airflow/models/deadline.py
b/airflow-core/src/airflow/models/deadline.py
index 67438ce428a..77e3c682e7b 100644
--- a/airflow-core/src/airflow/models/deadline.py
+++ b/airflow-core/src/airflow/models/deadline.py
@@ -147,8 +147,10 @@ class Deadline(Base):
def _determine_resource() -> tuple[str, str]:
"""Determine the type of resource based on which values are
present."""
if self.dagrun_id:
- # The deadline is for a Dag run:
- return "DagRun", f"Dag: {self.dagrun.dag_id} Run:
{self.dagrun_id}"
+ # Guard the relationship: the FK can be set while ``dagrun``
resolves to None (e.g.
+ # after a cascade delete). __repr__ must not raise, so fall
back to the id-only form.
+ dag_id = self.dagrun.dag_id if self.dagrun is not None else
"<unknown>"
+ return "DagRun", f"Dag: {dag_id} Run: {self.dagrun_id}"
return "Unknown", ""
@@ -183,9 +185,10 @@ class Deadline(Base):
return 0
try:
- # Get deadlines which match the provided conditions and their
associated DagRuns.
+ # Exclude deadlines already marked ``missed``: the scheduler owns
their (queued)
+ # callbacks, so prune must never cascade-delete them.
deadline_dagrun_pairs = session.execute(
- select(Deadline,
DagRun).join(DagRun).where(and_(*filter_conditions))
+ select(Deadline,
DagRun).join(DagRun).where(and_(*filter_conditions)).where(~Deadline.missed)
).all()
except AttributeError as e:
diff --git a/airflow-core/src/airflow/models/deadline_alert.py
b/airflow-core/src/airflow/models/deadline_alert.py
index 20bfed459ee..f9c8f2a6879 100644
--- a/airflow-core/src/airflow/models/deadline_alert.py
+++ b/airflow-core/src/airflow/models/deadline_alert.py
@@ -57,11 +57,14 @@ class DeadlineAlert(Base):
interval_seconds = None
+ # Legacy rows store a bare number instead of a serialized dict.
if isinstance(self.interval, (int, float)):
interval_seconds = int(self.interval)
- elif isinstance(self.interval, datetime.timedelta):
- interval_seconds = int(self.interval.total_seconds())
+ elif isinstance(self.interval, dict):
+ data = self.interval.get("__data__")
+ if isinstance(data, (int, float)):
+ interval_seconds = int(data)
if interval_seconds is None:
interval_display = "dynamic"
diff --git a/airflow-core/src/airflow/serialization/decoders.py
b/airflow-core/src/airflow/serialization/decoders.py
index 78ae241813c..0a8f703de30 100644
--- a/airflow-core/src/airflow/serialization/decoders.py
+++ b/airflow-core/src/airflow/serialization/decoders.py
@@ -167,7 +167,12 @@ def decode_deadline_reference(reference_data: dict):
"""Decode a previously serialized deadline reference."""
ref_name =
reference_data.get(SerializedReferenceModels.REFERENCE_TYPE_FIELD)
- if ref_name and SerializedReferenceModels.is_builtin_reference(ref_name):
+ # A custom reference may share a name with a builtin, so ``__class_path``
wins.
+ if "__class_path" in reference_data:
+ reference_class:
type[SerializedReferenceModels.SerializedBaseDeadlineReference] = (
+ SerializedReferenceModels.SerializedCustomReference
+ )
+ elif ref_name and SerializedReferenceModels.is_builtin_reference(ref_name):
reference_class =
SerializedReferenceModels.get_reference_class(ref_name)
else:
reference_class = SerializedReferenceModels.SerializedCustomReference
diff --git a/airflow-core/src/airflow/serialization/definitions/deadline.py
b/airflow-core/src/airflow/serialization/definitions/deadline.py
index 20fac2b54e8..968d69e3037 100644
--- a/airflow-core/src/airflow/serialization/definitions/deadline.py
+++ b/airflow-core/src/airflow/serialization/definitions/deadline.py
@@ -320,7 +320,15 @@ class SerializedReferenceModels:
def deserialize_reference(cls, reference_data: dict):
from airflow.serialization.helpers import
find_registered_custom_deadline_reference
- custom_class =
find_registered_custom_deadline_reference(reference_data["__class_path"])
+ class_path = reference_data.get("__class_path")
+ if not class_path:
+ raise ValueError(
+ "Cannot deserialize deadline reference: unrecognized
reference_type "
+
f"{reference_data.get(SerializedReferenceModels.REFERENCE_TYPE_FIELD)!r} with
no "
+ "'__class_path' to import. The stored reference is
corrupt, from a newer "
+ "Airflow version, or references a custom class whose
plugin is no longer installed."
+ )
+ custom_class =
find_registered_custom_deadline_reference(class_path)
inner_ref = custom_class.deserialize_reference(reference_data)
return cls(inner_ref)
diff --git a/airflow-core/tests/unit/models/test_deadline.py
b/airflow-core/tests/unit/models/test_deadline.py
index b52bbd5120b..1771269b548 100644
--- a/airflow-core/tests/unit/models/test_deadline.py
+++ b/airflow-core/tests/unit/models/test_deadline.py
@@ -209,6 +209,23 @@ class TestDeadline:
assert f"needed by {DEFAULT_DATE}" in repr_str
assert TEST_CALLBACK_PATH in repr_str
+ def test_repr_with_dagrun_id_but_no_dagrun_relationship(self,
deadline_orm):
+ """__repr__ must NOT raise when dagrun_id is set but the dagrun
relationship is None.
+
+ The FK (dagrun_id) can be set while the relationship resolves to None
— e.g. the DagRun
+ was deleted (ondelete=CASCADE) and this is a stale/expired in-memory
Deadline. A __repr__
+ that raised AttributeError here would break log lines, tracebacks, and
debugger displays
+ exactly when something is already going wrong. The repr falls back to
an id-only form.
+ """
+ # Sever the relationship while keeping the FK id (simulates
deleted/detached DagRun).
+ deadline_orm.dagrun = None
+ assert deadline_orm.dagrun_id is not None
+
+ repr_str = repr(deadline_orm) # must not raise
+ assert "[DagRun Deadline]" in repr_str
+ assert f"Run: {deadline_orm.dagrun_id}" in repr_str
+ assert "Dag: <unknown>" in repr_str
+
@pytest.mark.db_test
def test_bundle_name_propagated_to_callback(self, dagrun, session):
"""The bundle name is forwarded to the callback so the triggerer can
resolve its team."""
diff --git a/airflow-core/tests/unit/models/test_deadline_alert.py
b/airflow-core/tests/unit/models/test_deadline_alert.py
index a9b1854f6ab..56ed2f8cc86 100644
--- a/airflow-core/tests/unit/models/test_deadline_alert.py
+++ b/airflow-core/tests/unit/models/test_deadline_alert.py
@@ -117,6 +117,36 @@ class TestDeadlineAlert:
assert "interval=1m" in repr_str
assert repr(deadline_alert_orm.callback_def) in repr_str
+ @pytest.mark.parametrize(
+ ("interval", "expected"),
+ [
+ # Post-0117 shape: interval is the serialized dict, not a bare
number.
+ pytest.param(
+ {"__classname__": "datetime.timedelta", "__data__": 7200.0},
"interval=2h", id="timedelta_2h"
+ ),
+ # A corrupted dict without ``__data__`` must still render (no
raise) as dynamic.
+ pytest.param({"unexpected": "shape"}, "interval=dynamic",
id="corrupted_dict_dynamic"),
+ # A VariableInterval serializes with a dict ``__data__`` (its
key), not a number,
+ # so it renders as dynamic rather than a fixed duration.
+ pytest.param(
+ {
+ "__classname__":
"airflow.sdk.definitions.deadline.VariableInterval",
+ "__data__": {"key": "deadline_seconds"},
+ },
+ "interval=dynamic",
+ id="variable_interval_dynamic",
+ ),
+ ],
+ )
+ def test_deadline_alert_repr_does_not_raise_on_json_dict_interval(
+ self, deadline_alert_orm, interval, expected
+ ):
+ """``DeadlineAlert.__repr__`` must not raise for the production
JSON-dict interval shape."""
+ deadline_alert_orm.interval = interval
+ repr_str = repr(deadline_alert_orm) # must not raise
+ assert "[DeadlineAlert]" in repr_str
+ assert expected in repr_str
+
def test_deadline_alert_matches_definition(self, session,
deadline_reference):
alert1 = DeadlineAlert(
serialized_dag_id=SERIALIZED_DAG_ID,
diff --git a/airflow-core/tests/unit/models/test_prune_deadlines.py
b/airflow-core/tests/unit/models/test_prune_deadlines.py
new file mode 100644
index 00000000000..ba1712e3357
--- /dev/null
+++ b/airflow-core/tests/unit/models/test_prune_deadlines.py
@@ -0,0 +1,126 @@
+# 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.
+"""Coverage for ``Deadline.prune_deadlines``: on-time/overdue/pending
selection,
+callback cascade behaviour, and concurrent-mutation edge cases."""
+
+from __future__ import annotations
+
+from datetime import timedelta
+from typing import TYPE_CHECKING
+
+import pytest
+import time_machine
+from sqlalchemy import select
+
+from airflow.models import DagRun
+from airflow.models.callback import Callback
+from airflow.models.deadline import Deadline
+from airflow.providers.standard.operators.empty import EmptyOperator
+from airflow.sdk.definitions.callback import AsyncCallback
+from airflow.utils.state import DagRunState
+
+from tests_common.test_utils import db
+from unit.models import DEFAULT_DATE
+
+if TYPE_CHECKING:
+ from sqlalchemy.orm import Session
+
+DAG_ID = "prune_deadlines_test_dag"
+
+
+async def _prune_test_callback():
+ pass
+
+
+CALLBACK_PATH = f"{__name__}.{_prune_test_callback.__name__}"
+
+
+def _clean_db():
+ db.clear_db_dags()
+ db.clear_db_runs()
+ db.clear_db_deadline()
+
+
[email protected]
+def dagrun(session, dag_maker):
+ with dag_maker(DAG_ID):
+ EmptyOperator(task_id="task_id")
+ with time_machine.travel(DEFAULT_DATE):
+ dag_maker.create_dagrun(state=DagRunState.QUEUED,
logical_date=DEFAULT_DATE)
+ session.commit()
+ return session.scalars(select(DagRun)).one()
+
+
+def _make_deadline(session: Session, *, dagrun_id: int, deadline_time,
state=None) -> Deadline:
+ deadline = Deadline(
+ deadline_time=deadline_time,
+ callback=AsyncCallback(CALLBACK_PATH),
+ dagrun_id=dagrun_id,
+ dag_id=DAG_ID,
+ deadline_alert_id=None,
+ )
+ session.add(deadline)
+ session.flush()
+ if state is not None:
+ deadline.callback.state = state
+ session.add(deadline.callback)
+ session.flush()
+ return deadline
+
+
[email protected]_test
+class TestPruneDeadlines:
+ @staticmethod
+ def teardown_method():
+ _clean_db()
+
+ def test_handle_miss_then_prune_does_not_delete_missed(self, dagrun,
session):
+ """
+ Inverse race: handle_miss marks the deadline missed and queues the
callback
+ BEFORE the on-time prune runs. prune must NOT delete a deadline
already marked
+ ``missed`` — its callback is owned by the scheduler/triggerer, and
cascade-deleting
+ the deadline would silently drop that queued callback (a lost-callback
window).
+
+ prune_deadlines explicitly filters ``~Deadline.missed`` so a missed
deadline (and its
+ queued callback) survives even if the DagRun's end_date would
otherwise match the
+ on-time predicate. (In the real scheduler a missed deadline always has
+ ``deadline_time < now <= end_date`` so it can't match the on-time
predicate anyway;
+ the explicit filter makes the invariant robust to future callers and
clock skew.)
+ """
+ d = _make_deadline(session, dagrun_id=dagrun.id,
deadline_time=DEFAULT_DATE + timedelta(hours=1))
+ # Set missed + queued callback directly to test the prune ``~missed``
guard in isolation.
+ d.callback.queue(session=session)
+ d.missed = True
+ session.add_all([d, d.callback])
+ session.flush()
+ assert d.missed is True
+ deadline_id = d.id
+ callback_id = d.callback.id
+
+ # DagRun reports on-time completion (end_date <= deadline_time) — the
exact condition
+ # that would have pruned the row before the ~missed guard was added.
+ dagrun.end_date = DEFAULT_DATE
+ session.add(dagrun)
+ session.flush()
+
+ deleted = Deadline.prune_deadlines(session=session,
conditions={DagRun.id: dagrun.id})
+ session.flush()
+
+ # The missed deadline and its queued callback both survive.
+ assert deleted == 0
+ assert session.get(Deadline, deadline_id) is not None
+ assert session.get(Callback, callback_id) is not None
diff --git
a/airflow-core/tests/unit/serialization/test_deadline_reference_registry.py
b/airflow-core/tests/unit/serialization/test_deadline_reference_registry.py
index feae74bf980..ff54a6b2569 100644
--- a/airflow-core/tests/unit/serialization/test_deadline_reference_registry.py
+++ b/airflow-core/tests/unit/serialization/test_deadline_reference_registry.py
@@ -98,3 +98,85 @@ def
test_serialized_custom_reference_rejects_unregistered(monkeypatch):
SerializedReferenceModels.SerializedCustomReference.deserialize_reference(
{"__class_path": "some.other.module.UnregisteredReference"}
)
+
+
[email protected](
+ "reference_data",
+ [
+ pytest.param({"reference_type": "TotallyUnknownReference"},
id="unknown_type_no_class_path"),
+ pytest.param({"reference_type": "X", "__class_path": ""},
id="empty_class_path"),
+ ],
+)
+def
test_serialized_custom_reference_missing_class_path_raises_clear_error(reference_data):
+ """A reference routed to SerializedCustomReference but lacking a usable
``__class_path``
+ (corrupt / hand-edited row, blob from a newer version, or a custom ref
whose plugin is gone)
+ must raise a clear ValueError — NOT a bare ``KeyError: '__class_path'``."""
+ with pytest.raises(ValueError, match="unrecognized reference_type"):
+
SerializedReferenceModels.SerializedCustomReference.deserialize_reference(reference_data)
+
+
+class FixedDatetimeDeadline(ReferenceModels.BaseDeadlineReference):
+ """Custom reference whose class name deliberately collides with a builtin
reference name."""
+
+ required_kwargs: set[str] = set()
+
+ def serialize_reference(self) -> dict:
+ return {"reference_type": self.reference_name, "marker": "i-am-custom"}
+
+ def _evaluate_with(self, *, session, **kwargs):
+ raise AssertionError("custom evaluate should not be exercised in this
test")
+
+
+class DagRunLogicalDateDeadline(ReferenceModels.BaseDeadlineReference):
+ """Custom reference colliding with a builtin that has no required
deserialize fields."""
+
+ required_kwargs: set[str] = set()
+
+ def serialize_reference(self) -> dict:
+ return {"reference_type": self.reference_name}
+
+ def _evaluate_with(self, *, session, **kwargs):
+ raise AssertionError("custom evaluate should not be exercised in this
test")
+
+
+_COLLIDING_REFS = {
+ f"{FixedDatetimeDeadline.__module__}.FixedDatetimeDeadline":
FixedDatetimeDeadline,
+ f"{DagRunLogicalDateDeadline.__module__}.DagRunLogicalDateDeadline":
DagRunLogicalDateDeadline,
+}
+
+
[email protected]
+def colliding_plugin_registry(monkeypatch):
+ """Advertise custom references whose names collide with builtin reference
names."""
+ monkeypatch.setattr(
+ plugins_manager,
+ "get_deadline_references_plugins",
+ lambda: _COLLIDING_REFS,
+ )
+ return _COLLIDING_REFS
+
+
[email protected](
+ "custom_cls",
+ [FixedDatetimeDeadline],
+)
+def
test_custom_reference_name_collision_routes_to_custom(colliding_plugin_registry,
custom_cls):
+ """
+ A custom reference whose class name collides with a builtin must
round-trip as the custom
+ class, not silently decode as the builtin (which loses the user's
evaluation logic or raises
+ a spurious KeyError on builtin-only fields).
+
+ Regression test: ``decode_deadline_reference`` previously routed solely by
the
+ ``reference_type`` name string, ignoring the authoritative
``__class_path`` key.
+ """
+ from airflow.serialization.decoders import decode_deadline_reference
+ from airflow.serialization.encoders import encode_deadline_reference
+
+ encoded = encode_deadline_reference(custom_cls())
+ # The encoder stamps __class_path for custom references regardless of name
collision.
+ assert encoded["__class_path"] ==
f"{custom_cls.__module__}.{custom_cls.__name__}"
+
+ decoded = decode_deadline_reference(encoded)
+
+ assert isinstance(decoded,
SerializedReferenceModels.SerializedCustomReference)
+ assert isinstance(decoded.inner_ref, custom_cls)