kaxil commented on code in PR #73689:
URL: https://github.com/apache/airflow/pull/73689#discussion_r4127780586


##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -1649,6 +1649,65 @@ def 
test_dagrun_deadline_variable_interval_missing_variable_fails(self, _, sessi
                     session=session,
                 )
 
+    @pytest.mark.parametrize(
+        ("max_active_runs", "running_run_type", "expect_log"),
+        [
+            pytest.param(1, "manual", True, id="at_max"),
+            pytest.param(1, None, False, id="below_max"),
+            pytest.param(0, "manual", False, id="max_active_runs_unset"),
+            pytest.param(1, "backfill", False, 
id="running_backfill_run_does_not_count"),
+        ],
+    )
+    def test_log_if_new_run_blocked_by_max_active_runs(
+        self, session, caplog, max_active_runs, running_run_type, expect_log
+    ):
+        from airflow.models.backfill import Backfill

Review Comment:
   No circular import here (other test modules import `Backfill` at module 
level), so this can move up to the top-level `airflow.models` imports.



##########
airflow-core/src/airflow/models/dagrun.py:
##########
@@ -744,6 +744,54 @@ def active_runs_of_dags(
             query = query.where(cls.run_type != DagRunType.BACKFILL_JOB)
         return {dag_id: count for dag_id, count in session.execute(query)}
 
+    @classmethod
+    @provide_session
+    def log_if_new_run_blocked_by_max_active_runs(
+        cls,
+        *,
+        dag: SerializedDAG,
+        run_id: str,
+        session: Session = NEW_SESSION,
+    ) -> None:
+        """
+        Log if a just-created DagRun will not be scheduled yet because the Dag 
is at max_active_runs.
+
+        Meant to be called by request-driven trigger surfaces -- manual 
UI/REST-API triggers and
+        TriggerDagRunOperator/CLI (via 
:func:`airflow.api.common.trigger_dag.trigger_dag`) -- right
+        after :meth:`SerializedDAG.create_dagrun`. Deliberately not called 
from the scheduler's own
+        run-creation call sites: those run every scheduling loop, so calling 
this from there too
+        would turn a one-shot, per-trigger log into a per-loop one instead.
+
+        Counts running non-backfill runs only, to match the promotion queries
+        (:meth:`get_queued_dag_runs_to_set_running` and the scheduler's 
``_start_queued_dagruns``),
+        which both count running runs per ``(dag_id, backfill_id)`` -- a 
running backfill run
+        doesn't hold back a manual trigger's own max_active_runs slot.
+
+        :meta private:
+        """
+        if not dag.max_active_runs:

Review Comment:
   The backfill fix lines up the count, but the limit still comes from a 
different place than the gate uses. Both promotion paths read 
`DagModel.max_active_runs` (`_start_queued_dagruns` goes through 
`dag_run.max_active_runs`, with a comment about not trusting a stale serialized 
Dag), and the serialized write can lag `DagModel` by up to 
`min_serialized_dag_update_interval` after someone edits the value. 
`max_active_runs=0` is also the one case where the run is never promoted (the 
prefilter is `0 < 0`), and this early return skips logging it; the 
`max_active_runs_unset` test case locks that in. Making this an instance method 
that reads `self.max_active_runs or 0` off the just-created run would fix both 
and drop the `dag`/`run_id` params.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to