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

kaxil 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 370606c1065 Fix Dags staying blocked by `max_active_runs` until the 
next parse (#73032)
370606c1065 is described below

commit 370606c1065c7a2f75869d8ecfafdfd9a53a0942
Author: PoAn Yang <[email protected]>
AuthorDate: Tue Oct 6 20:47:50 2026 +0800

    Fix Dags staying blocked by `max_active_runs` until the next parse (#73032)
---
 airflow-core/src/airflow/jobs/scheduler_job_runner.py | 12 ++----------
 airflow-core/tests/unit/jobs/test_scheduler_job.py    | 12 ++++++++++--
 2 files changed, 12 insertions(+), 12 deletions(-)

diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py 
b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
index 5e07bf3cf70..1f794f41058 100644
--- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
+++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
@@ -3156,11 +3156,7 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
             session.flush()
             self.log.info("Run %s of %s has timed-out", dag_run.run_id, 
dag_run.dag_id)
 
-            if dag_run.state in State.finished_dr_states and dag_run.run_type 
in (
-                DagRunType.SCHEDULED,
-                DagRunType.MANUAL,
-                DagRunType.ASSET_TRIGGERED,
-            ):
+            if dag_run.state in State.finished_dr_states and dag_run.run_type 
!= DagRunType.BACKFILL_JOB:
                 self._set_exceeds_max_active_runs(dag_model=dag_model, 
session=session)
 
             callback_to_execute: DagCallbackRequest | None = None
@@ -3224,11 +3220,7 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
         # TODO[HA]: Rename update_state -> schedule_dag_run, ?? something else?
         schedulable_tis, callback_to_run = 
dag_run.update_state(session=session, execute_callbacks=False)
 
-        if dag_run.state in State.finished_dr_states and dag_run.run_type in (
-            DagRunType.SCHEDULED,
-            DagRunType.MANUAL,
-            DagRunType.ASSET_TRIGGERED,
-        ):
+        if dag_run.state in State.finished_dr_states and dag_run.run_type != 
DagRunType.BACKFILL_JOB:
             self._set_exceeds_max_active_runs(dag_model=dag_model, 
session=session)
 
         # This will do one query per dag run. We "could" build up a complex
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py 
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index eafe68660aa..7211bd68fe4 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -6039,6 +6039,7 @@ class TestSchedulerJob:
         )
         assert actual == expected
 
+    @pytest.mark.parametrize("timed_out", [False, True], ids=["finished", 
"timed_out"])
     @pytest.mark.parametrize(
         ("run_type", "expected"),
         [
@@ -6046,19 +6047,26 @@ class TestSchedulerJob:
             (DagRunType.SCHEDULED, True),
             (DagRunType.BACKFILL_JOB, False),
             (DagRunType.ASSET_TRIGGERED, True),
+            (DagRunType.OPERATOR_TRIGGERED, True),
+            (DagRunType.ASSET_MATERIALIZATION, True),
         ],
         ids=[
             DagRunType.MANUAL.name,
             DagRunType.SCHEDULED.name,
             DagRunType.BACKFILL_JOB.name,
             DagRunType.ASSET_TRIGGERED.name,
+            DagRunType.OPERATOR_TRIGGERED.name,
+            DagRunType.ASSET_MATERIALIZATION.name,
         ],
     )
-    def test_should_update_dag_next_dagruns_after_run_type(self, run_type, 
expected, session, dag_maker):
+    def test_should_update_dag_next_dagruns_after_run_type(
+        self, run_type, expected, timed_out, session, dag_maker
+    ):
         """Test that whether next dag run is updated depends on run type"""
         with dag_maker(
             schedule="*/1 * * * *",
             max_active_runs=3,
+            dagrun_timeout=datetime.timedelta(seconds=60),
         ):
             EmptyOperator(task_id="dummy")
 
@@ -6066,7 +6074,7 @@ class TestSchedulerJob:
             run_id="run",
             run_type=run_type,
             logical_date=DEFAULT_DATE,
-            start_date=timezone.utcnow(),
+            start_date=timezone.utcnow() - datetime.timedelta(days=1) if 
timed_out else timezone.utcnow(),
             state=State.SUCCESS,
             session=session,
         )

Reply via email to