From 0d39b7667666d6577fd01ce4852bf1f5c23fab8d Mon Sep 17 00:00:00 2001 From: Subham Sangwan Date: Sat, 2 May 2026 23:41:49 +0530 Subject: [PATCH 1/3] Fix depends_on_past broken for manual/asset-triggered runs without logical_date Fall back to run_after when logical_date is None to ensure depends_on_past works for all run types. Fixes #65578 --- airflow-core/src/airflow/models/dag.py | 12 +++++++++--- airflow-core/src/airflow/models/dagrun.py | 23 +++++++++++++++++------ 2 files changed, 26 insertions(+), 9 deletions(-) diff --git a/airflow-core/src/airflow/models/dag.py b/airflow-core/src/airflow/models/dag.py index ca84b7047b435..c98c35b819ca8 100644 --- a/airflow-core/src/airflow/models/dag.py +++ b/airflow-core/src/airflow/models/dag.py @@ -240,14 +240,20 @@ def get_last_dagrun(dag_id: str, session: Session, include_manually_triggered: b """ Return the last dag run for a dag, None if there was none. - Last dag run can be any type of run e.g. scheduled or backfilled. + Last dag run can be any type of run e.g. scheduled, backfilled, manual, or asset-triggered. Overridden DagRuns are ignored. + + For scheduled runs with logical_date, ordering uses logical_date. + For runs without logical_date (manual/asset-triggered), ordering falls back to run_after. + This ensures correct behavior for depends_on_past with all run types (AIP-39). """ DR = DagRun - query = select(DR).where(DR.dag_id == dag_id, DR.logical_date.is_not(None)) + query = select(DR).where(DR.dag_id == dag_id) if not include_manually_triggered: query = query.where(DR.run_type != DagRunType.MANUAL) - query = query.order_by(DR.logical_date.desc()) + # Order by logical_date if available, otherwise by run_after + # This ensures all run types are included and ordered correctly + query = query.order_by(func.coalesce(DR.logical_date, DR.run_after).desc()) return session.scalar(query.limit(1)) diff --git a/airflow-core/src/airflow/models/dagrun.py b/airflow-core/src/airflow/models/dagrun.py index 603041ef56b26..929f00a25a46e 100644 --- a/airflow-core/src/airflow/models/dagrun.py +++ b/airflow-core/src/airflow/models/dagrun.py @@ -983,19 +983,30 @@ def get_previous_dagrun( """ Return the previous DagRun, if there is one. + For scheduled runs with a logical_date, uses logical_date for ordering. + For runs without a logical_date (manual/asset-triggered), falls back to run_after. + This ensures depends_on_past works correctly for all run types (AIP-39). + :param dag_run: the dag run :param session: SQLAlchemy ORM Session :param state: the dag run state """ - if not dag_run or dag_run.logical_date is None: + if not dag_run: return None - filters = [ - DagRun.dag_id == dag_run.dag_id, - DagRun.logical_date < dag_run.logical_date, - ] + + filters = [DagRun.dag_id == dag_run.dag_id] if state is not None: filters.append(DagRun.state == state) - return session.scalar(select(DagRun).where(*filters).order_by(DagRun.logical_date.desc()).limit(1)) + + # For scheduled runs with logical_date, use logical_date for ordering + if dag_run.logical_date is not None: + filters.append(DagRun.logical_date < dag_run.logical_date) + return session.scalar(select(DagRun).where(*filters).order_by(DagRun.logical_date.desc()).limit(1)) + + # For runs without logical_date (manual/asset-triggered), fall back to run_after + # This ensures depends_on_past checks work for non-scheduled runs + filters.append(DagRun.run_after < dag_run.run_after) + return session.scalar(select(DagRun).where(*filters).order_by(DagRun.run_after.desc()).limit(1)) @staticmethod @provide_session From b31f9b3e7cebd1b46e36f35235515f855deaadd9 Mon Sep 17 00:00:00 2001 From: Subham Sangwan Date: Sun, 3 May 2026 00:57:40 +0530 Subject: [PATCH 2/3] Fix depends_on_past and run lookups for nullable logical_date --- airflow-core/src/airflow/models/dag.py | 11 ++--- airflow-core/src/airflow/models/dagrun.py | 13 +----- .../airflow/ti_deps/deps/prev_dagrun_dep.py | 32 +++++++++++--- .../unit/ti_deps/deps/test_prev_dagrun_dep.py | 44 +++++++++++++++++++ 4 files changed, 73 insertions(+), 27 deletions(-) diff --git a/airflow-core/src/airflow/models/dag.py b/airflow-core/src/airflow/models/dag.py index c98c35b819ca8..12242223f45f0 100644 --- a/airflow-core/src/airflow/models/dag.py +++ b/airflow-core/src/airflow/models/dag.py @@ -240,19 +240,14 @@ def get_last_dagrun(dag_id: str, session: Session, include_manually_triggered: b """ Return the last dag run for a dag, None if there was none. - Last dag run can be any type of run e.g. scheduled, backfilled, manual, or asset-triggered. - Overridden DagRuns are ignored. - - For scheduled runs with logical_date, ordering uses logical_date. - For runs without logical_date (manual/asset-triggered), ordering falls back to run_after. - This ensures correct behavior for depends_on_past with all run types (AIP-39). + Last dag run can be any type of run (e.g. scheduled, manual, asset-triggered). + Overridden DagRuns are ignored (AIP-39). """ DR = DagRun query = select(DR).where(DR.dag_id == dag_id) if not include_manually_triggered: query = query.where(DR.run_type != DagRunType.MANUAL) - # Order by logical_date if available, otherwise by run_after - # This ensures all run types are included and ordered correctly + # Order by logical_date if available, otherwise fallback to run_after (AIP-39) query = query.order_by(func.coalesce(DR.logical_date, DR.run_after).desc()) return session.scalar(query.limit(1)) diff --git a/airflow-core/src/airflow/models/dagrun.py b/airflow-core/src/airflow/models/dagrun.py index 929f00a25a46e..a3485962f54d5 100644 --- a/airflow-core/src/airflow/models/dagrun.py +++ b/airflow-core/src/airflow/models/dagrun.py @@ -981,11 +981,7 @@ def get_previous_dagrun( dag_run: DagRun, state: DagRunState | None = None, session: Session = NEW_SESSION ) -> DagRun | None: """ - Return the previous DagRun, if there is one. - - For scheduled runs with a logical_date, uses logical_date for ordering. - For runs without a logical_date (manual/asset-triggered), falls back to run_after. - This ensures depends_on_past works correctly for all run types (AIP-39). + Return the previous DagRun, if there is one (AIP-39). :param dag_run: the dag run :param session: SQLAlchemy ORM Session @@ -998,13 +994,6 @@ def get_previous_dagrun( if state is not None: filters.append(DagRun.state == state) - # For scheduled runs with logical_date, use logical_date for ordering - if dag_run.logical_date is not None: - filters.append(DagRun.logical_date < dag_run.logical_date) - return session.scalar(select(DagRun).where(*filters).order_by(DagRun.logical_date.desc()).limit(1)) - - # For runs without logical_date (manual/asset-triggered), fall back to run_after - # This ensures depends_on_past checks work for non-scheduled runs filters.append(DagRun.run_after < dag_run.run_after) return session.scalar(select(DagRun).where(*filters).order_by(DagRun.run_after.desc()).limit(1)) diff --git a/airflow-core/src/airflow/ti_deps/deps/prev_dagrun_dep.py b/airflow-core/src/airflow/ti_deps/deps/prev_dagrun_dep.py index 6cb2390a5dcae..5b239b680a9b9 100644 --- a/airflow-core/src/airflow/ti_deps/deps/prev_dagrun_dep.py +++ b/airflow-core/src/airflow/ti_deps/deps/prev_dagrun_dep.py @@ -19,7 +19,7 @@ from typing import TYPE_CHECKING -from sqlalchemy import func, or_, select +from sqlalchemy import func, literal, or_, select from airflow.models.backfill import BackfillDagRun from airflow.models.dagrun import DagRun @@ -75,12 +75,30 @@ def _has_any_prior_tis(ti: TI, *, session: Session) -> bool: This function exists for easy mocking in tests. """ - query = exists_query( - TI.dag_id == ti.dag_id, - TI.task_id == ti.task_id, - TI.logical_date < ti.logical_date, - session=session, - ) + if ti.logical_date is not None: + query = exists_query( + TI.dag_id == ti.dag_id, + TI.task_id == ti.task_id, + TI.logical_date < ti.logical_date, + session=session, + ) + else: + # Fallback to run_after for manual/asset runs (AIP-39) + dr = ti.get_dagrun(session=session) + query = ( + session.scalar( + select(literal(1)) + .select_from(TI) + .join(DagRun, TI.run_id == DagRun.run_id) + .where( + TI.dag_id == ti.dag_id, + TI.task_id == ti.task_id, + DagRun.run_after < dr.run_after, + ) + .limit(1) + ) + is not None + ) return query @staticmethod diff --git a/airflow-core/tests/unit/ti_deps/deps/test_prev_dagrun_dep.py b/airflow-core/tests/unit/ti_deps/deps/test_prev_dagrun_dep.py index 59a2274c5feb8..2376992eac87a 100644 --- a/airflow-core/tests/unit/ti_deps/deps/test_prev_dagrun_dep.py +++ b/airflow-core/tests/unit/ti_deps/deps/test_prev_dagrun_dep.py @@ -96,6 +96,50 @@ def test_first_task_run_of_new_task(self, testing_dag_bundle): assert dep.is_met(ti=ti, dep_context=dep_context) mock_has_any_prior_tis.assert_called_once_with(ti, session=ANY) + def test_prev_dagrun_with_manual_run_logical_date_null(self, session): + dag = DAG( + "test_depends_on_past_with_null_logical_date", + schedule=None, + start_date=START_DATE, + ) + task = BaseOperator( + task_id="test_task", + dag=dag, + depends_on_past=True, + start_date=START_DATE, + wait_for_downstream=False, + ) + scheduler_dag = sync_dag_to_db(dag, session=session) + + first_run_after = convert_to_utc(datetime(2016, 1, 1)) + first_dr = scheduler_dag.create_dagrun( + run_id="manual__2016-01-01T00:00:00+00:00", + state=DagRunState.RUNNING, + logical_date=None, + run_type=DagRunType.MANUAL, + data_interval=None, + run_after=first_run_after, + triggered_by=DagRunTriggeredByType.TEST, + ) + first_ti = first_dr.get_task_instance(task.task_id, session=session) + first_ti.set_state(TaskInstanceState.FAILED, session=session) + + second_run_after = convert_to_utc(datetime(2016, 1, 2)) + second_dr = scheduler_dag.create_dagrun( + run_id="manual__2016-01-02T00:00:00+00:00", + state=DagRunState.RUNNING, + logical_date=None, + run_type=DagRunType.MANUAL, + data_interval=None, + run_after=second_run_after, + triggered_by=DagRunTriggeredByType.TEST, + ) + second_ti = second_dr.get_task_instance(task.task_id, session=session) + second_ti.task = task + + dep_context = DepContext(ignore_depends_on_past=False) + assert not PrevDagrunDep().is_met(ti=second_ti, dep_context=dep_context) + @pytest.mark.parametrize( "kwargs", From 1a85677538c6b82ebd26042472aa0db60301835a Mon Sep 17 00:00:00 2001 From: Subham Sangwan Date: Sun, 3 May 2026 08:26:46 +0530 Subject: [PATCH 3/3] Fix get_previous_dagrun tiebreak when run_after values are equal --- airflow-core/src/airflow/models/dagrun.py | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/airflow-core/src/airflow/models/dagrun.py b/airflow-core/src/airflow/models/dagrun.py index a3485962f54d5..c1a308cbbe2c2 100644 --- a/airflow-core/src/airflow/models/dagrun.py +++ b/airflow-core/src/airflow/models/dagrun.py @@ -994,8 +994,19 @@ def get_previous_dagrun( if state is not None: filters.append(DagRun.state == state) - filters.append(DagRun.run_after < dag_run.run_after) - return session.scalar(select(DagRun).where(*filters).order_by(DagRun.run_after.desc()).limit(1)) + # Use (run_after, id) to correctly order runs when run_after values are equal + # (e.g. two manual runs triggered at the same time, or a scheduled run whose + # run_after equals the next run's run_after). id is a monotonically-increasing + # surrogate key so it gives a stable, deterministic tiebreak. + filters.append( + or_( + DagRun.run_after < dag_run.run_after, + and_(DagRun.run_after == dag_run.run_after, DagRun.id < dag_run.id), + ) + ) + return session.scalar( + select(DagRun).where(*filters).order_by(DagRun.run_after.desc(), DagRun.id.desc()).limit(1) + ) @staticmethod @provide_session