seanmuth commented on code in PR #73689:
URL: https://github.com/apache/airflow/pull/73689#discussion_r4124440437
##########
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:
You're right — confirmed `_set_exceeds_max_active_runs` doesn't log anything
on `main`. I'd conflated it with #72403, which was closed without merging.
Reworded the docstring to drop that reference; the reasoning for keeping this
out of the scheduler (its own creation paths run every loop) stands fine
without it. Also caught the same stale #72403-as-merged framing in the tracking
issue's body and corrected it there too.
---
Drafted-by: Claude Sonnet 5; reviewed by @seanmuth before posting
##########
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:
Done — consolidated into one parametrized test (also added a 4th case
covering the backfill fix above) and switched the negative assertions to the
`not in caplog` membership form instead of `caplog.records`.
---
Drafted-by: Claude Sonnet 5; reviewed by @seanmuth 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]