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)