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

henry3260 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 6e5f3136b9b Fix dateless dag.test and tasks test destroying unrelated 
Dag runs (#73813)
6e5f3136b9b is described below

commit 6e5f3136b9bc66524d17f813684b87e6b5a092cd
Author: Y-C <[email protected]>
AuthorDate: Tue Oct 6 18:56:48 2026 +0800

    Fix dateless dag.test and tasks test destroying unrelated Dag runs (#73813)
    
    Airflow 3 allows Dag runs without a logical date, which the helper behind
    these two debugging commands never accounted for. An absent date widened
    both its collision lookup and its pre-test clear into match-everything
    queries, so testing a single task could delete a run the user had
    triggered from the UI, along with its task instances and their history.
    
    A command whose whole purpose is local debugging must not touch data it
    was never pointed at.
    
    Co-authored-by: Eason09053360 
<[email protected]>
    Co-authored-by: Henry Chen <[email protected]>
---
 airflow-core/src/airflow/models/dagrun.py     | 20 +++++++------
 airflow-core/tests/unit/models/test_dag.py    | 27 ++++++++++++++++++
 airflow-core/tests/unit/models/test_dagrun.py | 41 ++++++++++++++++++++++++++-
 task-sdk/src/airflow/sdk/definitions/dag.py   | 18 ++++++------
 4 files changed, 89 insertions(+), 17 deletions(-)

diff --git a/airflow-core/src/airflow/models/dagrun.py 
b/airflow-core/src/airflow/models/dagrun.py
index fc26cbccd16..db4c38db5d1 100644
--- a/airflow-core/src/airflow/models/dagrun.py
+++ b/airflow-core/src/airflow/models/dagrun.py
@@ -2586,24 +2586,28 @@ def get_or_create_dagrun(
     """
     Create a DAG run, replacing an existing instance if needed to prevent 
collisions.
 
-    This function is only meant to be used by :meth:`DAG.test` as a helper 
function.
+    This function is only meant to be used by :meth:`DAG.test` and ``airflow 
tasks test``
+    as a helper function.
 
     :param dag: DAG to be used to find run.
     :param conf: Configuration to pass to newly created run.
     :param start_date: Start date of new run.
-    :param logical_date: Logical date for finding an existing run.
+    :param logical_date: Logical date for finding an existing run to replace. 
``None`` skips
+        the lookup, since NULL dates cannot violate the ``(dag_id, 
logical_date)`` unique key.
     :param run_id: Run ID for the new DAG run.
     :param triggered_by: the entity which triggers the dag_run
     :param triggering_user_name: the user name who triggers the dag_run
 
     :return: The newly created DAG run.
     """
-    dr = session.scalar(
-        select(DagRun).where(DagRun.dag_id == dag.dag_id, DagRun.logical_date 
== logical_date)
-    )
-    if dr:
-        session.delete(dr)
-        session.commit()
+    # ``== None`` compiles to ``IS NULL``, which would match an unrelated 
dateless run.
+    if logical_date is not None:
+        dr = session.scalar(
+            select(DagRun).where(DagRun.dag_id == dag.dag_id, 
DagRun.logical_date == logical_date)
+        )
+        if dr:
+            session.delete(dr)
+            session.flush()
     dr = dag.create_dagrun(
         run_id=run_id,
         logical_date=logical_date,
diff --git a/airflow-core/tests/unit/models/test_dag.py 
b/airflow-core/tests/unit/models/test_dag.py
index e5d0939a182..7fa1200b028 100644
--- a/airflow-core/tests/unit/models/test_dag.py
+++ b/airflow-core/tests/unit/models/test_dag.py
@@ -1854,6 +1854,33 @@ class TestDag:
         assert ti.state == (TaskInstanceState.SUCCESS if succeed_on_last_try 
else TaskInstanceState.FAILED)
         assert dr.state == (DagRunState.SUCCESS if succeed_on_last_try else 
DagRunState.FAILED)
 
+    def 
test_dag_test_without_logical_date_keeps_other_runs_task_instances(self, 
testing_dag_bundle, session):
+        dag = DAG(dag_id="test_dateless_dag_test", schedule=None, 
start_date=DEFAULT_DATE)
+
+        @task_decorator
+        def check_task():
+            pass
+
+        with dag:
+            check_task()
+
+        _create_dagrun(
+            dag,
+            logical_date=DEFAULT_DATE,
+            data_interval=(DEFAULT_DATE, DEFAULT_DATE),
+            run_type=DagRunType.SCHEDULED,
+            state=DagRunState.SUCCESS,
+        )
+        run_id = session.scalar(select(DagRun.run_id).where(DagRun.dag_id == 
dag.dag_id))
+        session.execute(update(TI).where(TI.run_id == 
run_id).values(state=TaskInstanceState.SUCCESS))
+        session.commit()
+
+        dag.test(logical_date=None)
+
+        session.expire_all()
+        states = session.scalars(select(TI.state).where(TI.run_id == 
run_id)).all()
+        assert states == [TaskInstanceState.SUCCESS]
+
     def test_dag_test_with_dependencies(self, testing_dag_bundle):
         dag = DAG(dag_id="test_local_testing_conn_file", schedule=None, 
start_date=DEFAULT_DATE)
         sync_dag_to_db(dag)
diff --git a/airflow-core/tests/unit/models/test_dagrun.py 
b/airflow-core/tests/unit/models/test_dagrun.py
index e85713943d6..4221ca746d3 100644
--- a/airflow-core/tests/unit/models/test_dagrun.py
+++ b/airflow-core/tests/unit/models/test_dagrun.py
@@ -53,7 +53,7 @@ from airflow._shared.timezones import timezone
 from airflow.callbacks.callback_requests import DagCallbackRequest, 
DagRunContext
 from airflow.models.dag import DagModel, infer_automated_data_interval
 from airflow.models.dag_version import DagVersion
-from airflow.models.dagrun import DagRun, DagRunNote, clear_partition_runs
+from airflow.models.dagrun import DagRun, DagRunNote, clear_partition_runs, 
get_or_create_dagrun
 from airflow.models.deadline import Deadline
 from airflow.models.deadline_alert import DeadlineAlert as DeadlineAlertModel
 from airflow.models.serialized_dag import SerializedDagModel
@@ -5442,3 +5442,42 @@ class TestApplyPartitionDateWindowSubDay:
             session=session,
         )
         assert cleared == 3
+
+
+class TestGetOrCreateDagrun:
+    @pytest.fixture(autouse=True)
+    def _clean_db(self):
+        clear_db_runs()
+        clear_db_dags()
+        yield
+        clear_db_runs()
+        clear_db_dags()
+
+    def test_none_logical_date_keeps_unrelated_runs(self, dag_maker, session):
+        with dag_maker("test_get_or_create_dagrun", serialized=True):
+            EmptyOperator(task_id="t1")
+        dag_maker.create_dagrun(run_id="manual__existing", logical_date=None, 
state=DagRunState.SUCCESS)
+        now = timezone.utcnow()
+
+        created = get_or_create_dagrun(
+            dag=dag_maker.serialized_dag,
+            run_id="__airflow_temporary_run_x",
+            logical_date=None,
+            data_interval=None,
+            run_after=now,
+            conf=None,
+            triggered_by=DagRunTriggeredByType.CLI,
+            triggering_user_name=None,
+            start_date=now,
+            session=session,
+        )
+
+        assert created.run_id == "__airflow_temporary_run_x"
+        run_ids = set(
+            session.scalars(select(DagRun.run_id).where(DagRun.dag_id == 
"test_get_or_create_dagrun"))
+        )
+        assert run_ids == {"manual__existing", "__airflow_temporary_run_x"}
+        existing_ti_count = session.scalar(
+            
select(func.count()).select_from(TaskInstance).where(TaskInstance.run_id == 
"manual__existing")
+        )
+        assert existing_ti_count == 1
diff --git a/task-sdk/src/airflow/sdk/definitions/dag.py 
b/task-sdk/src/airflow/sdk/definitions/dag.py
index 4a5764e800e..2dd62ff8e87 100644
--- a/task-sdk/src/airflow/sdk/definitions/dag.py
+++ b/task-sdk/src/airflow/sdk/definitions/dag.py
@@ -1262,14 +1262,16 @@ class DAG:
             # Allow users to explicitly pass None. If it isn't set, we default 
to current time.
             logical_date = logical_date if is_arg_set(logical_date) else 
timezone.utcnow()
 
-            log.debug("Clearing existing task instances for logical date %s", 
logical_date)
-            # TODO: Replace with calling client.dag_run.clear in Execution API 
at some point
-            SerializedDAG.clear_dags(
-                dags=[scheduler_dag],
-                start_date=logical_date,
-                end_date=logical_date,
-                dag_run_state=False,
-            )
+            # Unset bounds match every run, so a dateless test would clear 
unrelated runs.
+            if logical_date is not None:
+                log.debug("Clearing existing task instances for logical date 
%s", logical_date)
+                # TODO: Replace with calling client.dag_run.clear in Execution 
API at some point
+                SerializedDAG.clear_dags(
+                    dags=[scheduler_dag],
+                    start_date=logical_date,
+                    end_date=logical_date,
+                    dag_run_state=False,
+                )
 
             log.debug("Getting dagrun for dag %s", self.dag_id)
             logical_date = timezone.coerce_datetime(logical_date)

Reply via email to