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

vatsrahul1001 pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/v3-3-test by this push:
     new ed40a894ef2  allow deadline alert UUID references in serialized Dag 
schema (#70148) (#72738)
ed40a894ef2 is described below

commit ed40a894ef249446560726c4de967961078fb9da
Author: Rahul Vats <[email protected]>
AuthorDate: Wed Sep 9 19:09:32 2026 +0530

     allow deadline alert UUID references in serialized Dag schema (#70148) 
(#72738)
    
    Co-authored-by: Kevin Yang <[email protected]>
    (cherry picked from commit eaf5fc19ce7f799993f9e18fe5090a5718487e73)
    
    # Conflicts:
    #       devel-common/src/tests_common/test_utils/dag.py
    
    Co-authored-by: Jung-Hyun Andrew Kim <[email protected]>
---
 airflow-core/src/airflow/models/serialized_dag.py  |  6 ++
 airflow-core/src/airflow/serialization/schema.json |  7 ++
 .../tests/unit/models/test_serialized_dag.py       | 97 ++++++++++++++++++++++
 devel-common/src/tests_common/test_utils/dag.py    |  6 +-
 4 files changed, 115 insertions(+), 1 deletion(-)

diff --git a/airflow-core/src/airflow/models/serialized_dag.py 
b/airflow-core/src/airflow/models/serialized_dag.py
index 9698d8fb745..b0e57e79991 100644
--- a/airflow-core/src/airflow/models/serialized_dag.py
+++ b/airflow-core/src/airflow/models/serialized_dag.py
@@ -19,6 +19,7 @@
 
 from __future__ import annotations
 
+import copy
 import logging
 import zlib
 from collections.abc import Callable, Iterable, Iterator, Sequence
@@ -644,6 +645,11 @@ class SerializedDagModel(Base):
         name_updated = False
         reused_deadline_data: dict[str, dict] | None = None
         if dag.data.get("dag", {}).get("deadline"):
+            # The deadline handling below rewrites data["dag"]["deadline"] 
from a list of
+            # encoded dicts into a list of UUID references. Work on a copy so 
we never mutate
+            # the caller's LazyDeserializedDAG in place.
+
+            dag = dag.model_copy(update={"data": copy.deepcopy(dag.data)})
             # Try to reuse existing deadline UUIDs if the deadline definitions 
haven't changed.
             # This preserves the hash and avoids unnecessary 
SerializedDagModel recreations.
             existing_serialized_dag = session.scalar(
diff --git a/airflow-core/src/airflow/serialization/schema.json 
b/airflow-core/src/airflow/serialization/schema.json
index b9efe884480..e585d7e2dcd 100644
--- a/airflow-core/src/airflow/serialization/schema.json
+++ b/airflow-core/src/airflow/serialization/schema.json
@@ -200,6 +200,13 @@
                     "type": "array",
                     "items": { "$ref": "#/definitions/dict" }
                 },
+                {
+                    "$comment": "Once persisted, a Dag's deadline alerts live 
as rows in the deadline_alert table and the serialized Dag keeps only a list of 
UUID strings referencing them (see 
SerializedDagModel._generate_deadline_uuids). This branch lets the stored form 
validate at any lifecycle stage, not only before the dict->UUID rewrite.",
+                    "type": "array",
+                    "items": { "type": "string",
+                    "pattern": 
"^[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$"
+                  }
+                },
                 { "type": "null" }
             ]
         },
diff --git a/airflow-core/tests/unit/models/test_serialized_dag.py 
b/airflow-core/tests/unit/models/test_serialized_dag.py
index acbcb181962..326da2d4167 100644
--- a/airflow-core/tests/unit/models/test_serialized_dag.py
+++ b/airflow-core/tests/unit/models/test_serialized_dag.py
@@ -19,6 +19,7 @@
 
 from __future__ import annotations
 
+import copy
 import logging
 from datetime import timedelta
 from unittest import mock
@@ -1131,3 +1132,99 @@ class TestSerializedDagModel:
         alert = session.scalar(select(DAM).where(DAM.serialized_dag_id == 
orig_serdag.id))
         assert alert is not None
         assert alert.id == orig_alert.id
+
+    def test_write_dag_with_deadline_passes_schema_validation(self, 
testing_dag_bundle, session):
+        """The persisted serialized Dag for a deadline-bearing Dag must 
satisfy the JSON schema.
+
+        write_dag stores ``data["dag"]["deadline"]`` as a list of UUID strings 
referencing
+        deadline_alert rows, so the schema has to accept that persisted form 
and not only the
+        list-of-dicts form produced before the dict->UUID rewrite.
+        """
+        dag_id = "test_deadline_schema_valid"
+        dag = DAG(
+            dag_id=dag_id,
+            deadline=DeadlineAlert(
+                reference=DeadlineReference.DAGRUN_QUEUED_AT,
+                interval=timedelta(minutes=5),
+                callback=AsyncCallback(empty_callback_for_deadline),
+            ),
+        )
+        EmptyOperator(task_id="task1", dag=dag)
+        sync_dag_to_db(dag, session=session)
+        session.commit()
+
+        result = session.scalar(select(SDM).where(SDM.dag_id == dag_id))
+        persisted_deadline = result.data["dag"]["deadline"]
+        assert isinstance(persisted_deadline, list)
+        assert persisted_deadline
+        assert all(isinstance(ref, str) for ref in persisted_deadline)
+
+        # Must not raise: the stored UUID-reference form has to satisfy the 
serialized Dag schema.
+        DagSerialization.validate_schema(result.data)
+
+    def test_write_dag_does_not_mutate_caller_deadline_data(self, 
testing_dag_bundle, session):
+        """write_dag must not rewrite the caller's LazyDeserializedDAG 
deadline in place.
+
+        The dict->UUID replacement in ``_generate_deadline_uuids`` has to 
happen on a copy so a
+        LazyDeserializedDAG the caller still references keeps its original 
list-of-dicts deadline.
+        """
+        dag_id = "test_deadline_no_mutation"
+        dag = DAG(
+            dag_id=dag_id,
+            deadline=DeadlineAlert(
+                reference=DeadlineReference.DAGRUN_QUEUED_AT,
+                interval=timedelta(minutes=5),
+                callback=AsyncCallback(empty_callback_for_deadline),
+            ),
+        )
+        EmptyOperator(task_id="task1", dag=dag)
+        sync_dag_to_db(dag, session=session)
+        session.commit()
+
+        # Change the interval so write_dag regenerates UUIDs (the dict->UUID 
rewrite path)
+        # rather than reusing the existing ones.
+        dag.deadline = DeadlineAlert(
+            reference=DeadlineReference.DAGRUN_QUEUED_AT,
+            interval=timedelta(minutes=10),
+            callback=AsyncCallback(empty_callback_for_deadline),
+        )
+        lazy_dag = LazyDeserializedDAG.from_dag(dag)
+        original_deadline = copy.deepcopy(lazy_dag.data["dag"]["deadline"])
+        assert original_deadline
+        assert all(isinstance(item, dict) for item in original_deadline)
+
+        SDM.write_dag(lazy_dag, bundle_name="testing", session=session)
+        session.commit()
+
+        assert lazy_dag.data["dag"]["deadline"] == original_deadline
+
+    def test_sync_dag_to_db_returns_db_normalized_deadline_ids(self, 
testing_dag_bundle, session):
+        """sync_dag_to_db must return a SerializedDAG with the DB-normalized 
deadline UUIDs. Verify that the UUIDs returned by sync_dag_to_db match the 
persisted deadline_alert rows in the DB."""
+        dag_id = "test_sync_dag_to_db_deadline_ids"
+        dag = DAG(
+            dag_id=dag_id,
+            deadline=DeadlineAlert(
+                reference=DeadlineReference.DAGRUN_QUEUED_AT,
+                interval=timedelta(minutes=5),
+                callback=AsyncCallback(empty_callback_for_deadline),
+            ),
+        )
+        EmptyOperator(task_id="task1", dag=dag)
+
+        scheduler_dag = sync_dag_to_db(dag, session=session)
+        session.commit()
+
+        latest_serdag = session.scalar(
+            select(SDM).where(SDM.dag_id == 
dag_id).order_by(SDM.created_at.desc())
+        )
+        assert latest_serdag is not None
+
+        persisted_alerts = 
session.scalars(select(DAM).where(DAM.serialized_dag_id == 
latest_serdag.id)).all()
+
+        persisted_uuids = {str(alert.id) for alert in persisted_alerts}
+        returned_uuids = scheduler_dag.deadline or []
+
+        assert returned_uuids
+        assert all(isinstance(ref, str) for ref in returned_uuids)
+        assert len(returned_uuids) == len(set(returned_uuids))
+        assert set(returned_uuids) == persisted_uuids
diff --git a/devel-common/src/tests_common/test_utils/dag.py 
b/devel-common/src/tests_common/test_utils/dag.py
index 6e02ddf61ab..e042f53d0b2 100644
--- a/devel-common/src/tests_common/test_utils/dag.py
+++ b/devel-common/src/tests_common/test_utils/dag.py
@@ -69,7 +69,11 @@ def sync_dags_to_db(
     def _write_dag(dag: DAG) -> SerializedDAG:
         data = DagSerialization.to_dict(dag)
         SerializedDagModel.write_dag(LazyDeserializedDAG(data=data), 
bundle_name, session=session)
-        return DagSerialization.from_dict(data)
+        session.flush()
+        serialized_dag = SerializedDagModel.get_dag(dag.dag_id, 
session=session)
+        if serialized_dag is None:
+            raise RuntimeError(f"Serialized DAG {dag.dag_id!r} was not found 
after writing to the database")
+        return serialized_dag
 
     SerializedDAG.bulk_write_to_db(bundle_name, None, dags, session=session)
     scheduler_dags = [_write_dag(dag) for dag in dags]

Reply via email to