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]