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


##########
airflow-core/src/airflow/models/dagrun.py:
##########
@@ -744,6 +744,47 @@ 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, and the 
scheduler already has
+        separate periodic bookkeeping for this 
(``_set_exceeds_max_active_runs``) that intentionally
+        only logs once when a Dag *newly* becomes blocked, rather than every 
loop. Calling this from
+        there too would turn a one-shot, per-trigger log back into that same 
per-loop spam.
+
+        :meta private:
+        """
+        if not dag.max_active_runs:
+            return
+        num_running = (
+            session.scalar(
+                select(func.count())
+                .select_from(cls)
+                .where(cls.dag_id == dag.dag_id, cls.state == 
DagRunState.RUNNING)

Review Comment:
   This counts every running run of the Dag, backfill runs included, but the 
scheduler doesn't. Both `get_queued_dag_runs_to_set_running` and 
`_start_queued_dagruns` count running runs per `(dag_id, backfill_id)`, so a 
manual run is only held back by other non-backfill runs. With a backfill in 
flight on a `max_active_runs=1` Dag, a manual trigger would log "will not be 
scheduled yet" and then start on the next loop, which is the misleading signal 
this PR is trying to remove. Adding `cls.backfill_id.is_(None)` here, plus a 
test with a running backfill run, would line it up with the promotion query.



##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -1649,6 +1649,86 @@ def 
test_dagrun_deadline_variable_interval_missing_variable_fails(self, _, sessi
                     session=session,
                 )
 
+    def test_log_if_new_run_blocked_by_max_active_runs_logs_when_at_max(self, 
session, caplog):
+        dag = DAG(
+            
dag_id="test_log_if_new_run_blocked_by_max_active_runs_logs_when_at_max",
+            schedule=None,
+            max_active_runs=1,
+        )
+        scheduler_dag = sync_dag_to_db(dag, session=session)
+        scheduler_dag.create_dagrun(
+            run_id="running_run",
+            logical_date=DEFAULT_DATE,
+            data_interval=(DEFAULT_DATE, DEFAULT_DATE),
+            run_after=DEFAULT_DATE,
+            run_type=DagRunType.MANUAL,
+            state=DagRunState.RUNNING,
+            triggered_by=DagRunTriggeredByType.TEST,
+            session=session,
+        )
+
+        with caplog.at_level("INFO", logger="airflow.models.dagrun"):
+            DagRun.log_if_new_run_blocked_by_max_active_runs(
+                dag=scheduler_dag, run_id="queued_run", session=session
+            )
+
+        assert {
+            "event": "created DagRun will not be scheduled yet, dag is at 
max_active_runs",
+            "dag_id": 
"test_log_if_new_run_blocked_by_max_active_runs_logs_when_at_max",
+            "run_id": "queued_run",
+            "active_runs": 1,
+            "max_active_runs": 1,
+            "log_level": "info",
+        } in caplog
+
+    def 
test_log_if_new_run_blocked_by_max_active_runs_does_not_log_when_below_max(self,
 session, caplog):

Review Comment:
   These three tests share the same setup and differ only in `max_active_runs`, 
whether a running run exists, and the expected outcome, so one 
`pytest.mark.parametrize` test would cover them. The negative checks can also 
be `"created DagRun will not be scheduled yet, dag is at max_active_runs" not 
in caplog`, since `caplog.records` on the structlog capture is marked for 
removal in `tests_common`.



##########
airflow-core/src/airflow/models/dagrun.py:
##########
@@ -744,6 +744,47 @@ 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, and the 
scheduler already has
+        separate periodic bookkeeping for this 
(``_set_exceeds_max_active_runs``) that intentionally

Review Comment:
   `_set_exceeds_max_active_runs` doesn't log anything on main. The "only logs 
once when a Dag newly becomes blocked" behaviour came from #72403, which was 
closed without merging. The reason for keeping this out of the scheduler (its 
creation paths run every loop) holds on its own, so I'd drop the reference to 
that bookkeeping or reword it to describe what's actually there.



-- 
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