seanmuth commented on code in PR #73689:
URL: https://github.com/apache/airflow/pull/73689#discussion_r4134738071
##########
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:
Good catch on both counts, fixed. Converted this to an instance method on
the just-created `DagRun` — it now reads `self.max_active_runs` (the
`dag_model`-backed association proxy, same source `_start_queued_dagruns` uses
via `dag_run.max_active_runs`) instead of the serialized Dag's value, and drops
the `dag`/`run_id` params entirely.
Also dropped the early `if not dag.max_active_runs: return` and let the
comparison run unconditionally as `num_running >= (self.max_active_runs or 0)`,
so `max_active_runs=0` now logs correctly instead of being treated as unset.
Updated the `max_active_runs_unset` test case to
`max_active_runs_zero_never_promoted` with `expect_log=True` to lock that in,
and the other three call sites (`trigger_dag.py`, and the two API routes) now
call `dag_run.log_if_new_run_blocked_by_max_active_runs()` on the run they just
created.
---
Drafted-by: Claude Sonnet 5 (no human review before posting)
##########
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:
Moved to the top-level imports, thanks.
---
Drafted-by: Claude Sonnet 5 (no human review before posting)
--
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]