kaxil commented on code in PR #73692:
URL: https://github.com/apache/airflow/pull/73692#discussion_r4139349353
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py:
##########
@@ -221,6 +222,18 @@ class DAGDetailsResponse(DAGResponse):
owner_links: dict[str, str] | None = None
is_favorite: bool = False
active_runs_count: int = 0
+ is_at_max_active_runs: bool = Field(
+ description=(
+ "Whether this Dag currently has as many active runs as its
max_active_runs allows. "
+ "Counted differently from active_runs_count above: this counts
RUNNING and QUEUED "
+ "runs (excluding backfill runs), matching the scheduler's own
promotion check, while "
Review Comment:
"Matching the scheduler's own promotion check" isn't right. The QUEUED to
RUNNING gate (`get_queued_dag_runs_to_set_running` in dagrun.py, and
`_start_queued_dagruns`) counts RUNNING runs only. RUNNING plus QUEUED minus
backfill is the check the scheduler makes before creating another scheduled run
(`_set_exceeds_max_active_runs`, read by `dags_needing_dagruns`), and the
second example here shows the gap: one queued run with `max_active_runs=1`
reads `true`, yet the scheduler promotes that run on its next loop. Could you
describe it as the run-creation check? This text is copied into the YAML, both
TS files and airflowctl, so those need regenerating too.
Separately, #73693 doesn't read this field. Its tooltip in `Header.tsx`
works out its own answer from `active_runs_count`, `max_active_runs` and
`queued_runs_count`, which is closer to the question #73686 asks (are queued
runs held back by the limit). So this adds a required public field that nothing
in-tree reads, and it answers a different question. Does it need to ship in
this PR, or could it wait until there's a caller that wants the run-creation
answer?
##########
airflow-core/src/airflow/dag_processing/collection.py:
##########
@@ -690,7 +692,7 @@ def update_dags(
dm.bundle_version = self.bundle_version
reference_run: DagRun | None = run_info.latest_run
- dm.exceeds_max_non_backfill = run_info.num_active_runs >=
dm.max_active_runs
+ dm.exceeds_max_non_backfill = active_run_counts.get(dag_id, 0) >=
dm.max_active_runs
Review Comment:
Now that the route computes `is_at_max_active_runs` itself, nothing reads
`exceeds_max_non_backfill` for a Dag with `can_be_scheduled=False`. Its only
readers are `dags_needing_dagruns` and `_create_dag_runs`, and a
`schedule=None` Dag never gets `next_dagrun_create_after` set or an asset
trigger. The batching is still a good cut in per-Dag queries, but the PR body
(the `DAG_ALIAS_MAPPING` alias, "permanently stuck False", #73693 depending on
it) and the docstrings in test_collection.py and test_dag.py (the removed
`num_active_runs`, "the flag is read regardless") still describe the old
design. Could you update them to say this is a batching change with no
scheduling effect for `schedule=None` Dags?
##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py:
##########
@@ -1512,6 +1513,96 @@ def test_dag_details_includes_active_runs_count(self,
session, test_client):
assert isinstance(body["active_runs_count"], int)
assert body["active_runs_count"] == 0
+ def test_dag_details_includes_is_at_max_active_runs(self, session,
test_client):
+ """is_at_max_active_runs is computed fresh from real DagRuns, not a
stale cached column."""
+ dag_model = session.get(DagModel, DAG2_ID)
+ dag_model.max_active_runs = 1
+ session.add(
+ DagRun(
+ dag_id=DAG2_ID,
+ run_id="is_at_max_active_runs_running",
+ logical_date=datetime(2021, 6, 15, 4, 0, 0,
tzinfo=timezone.utc),
+ start_date=datetime(2021, 6, 15, 4, 0, 0, tzinfo=timezone.utc),
+ run_type=DagRunType.MANUAL,
+ state=DagRunState.RUNNING,
+ triggered_by=DagRunTriggeredByType.TEST,
+ )
+ )
+ session.commit()
+
+ response = test_client.get(f"/dags/{DAG2_ID}/details")
+ assert response.status_code == 200
+ body = response.json()
+
+ assert body["is_at_max_active_runs"] is True
+
+ # Test with a DAG that has not hit its max_active_runs
+ response = test_client.get(f"/dags/{DAG1_ID}/details")
+ assert response.status_code == 200
+ body = response.json()
+
+ assert body["is_at_max_active_runs"] is False
+
+ def test_dag_details_is_at_max_active_runs_counts_queued_runs_too(self,
session, test_client):
+ """A queued (non-backfill) run counts toward is_at_max_active_runs
even with 0 running."""
+ dag_model = session.get(DagModel, DAG2_ID)
+ dag_model.max_active_runs = 1
+ session.add(
+ DagRun(
+ dag_id=DAG2_ID,
+ run_id="is_at_max_active_runs_queued",
+ logical_date=datetime(2021, 6, 15, 4, 0, 0,
tzinfo=timezone.utc),
+ start_date=datetime(2021, 6, 15, 4, 0, 0, tzinfo=timezone.utc),
+ run_type=DagRunType.MANUAL,
+ state=DagRunState.QUEUED,
+ triggered_by=DagRunTriggeredByType.TEST,
+ )
+ )
+ session.commit()
+
+ response = test_client.get(f"/dags/{DAG2_ID}/details")
+ assert response.status_code == 200
+ body = response.json()
+
+ assert body["active_runs_count"] == 0
+ assert body["is_at_max_active_runs"] is True
+
+ def test_dag_details_is_at_max_active_runs_excludes_backfill_runs(self,
session, test_client):
+ """A running backfill run doesn't count toward the Dag's own
is_at_max_active_runs."""
+ from airflow.models.backfill import Backfill
Review Comment:
This import can go at the top of the file. The `Backfill` row also outlives
the test, because `_clear_db` doesn't call `clear_db_backfills()`. You don't
need the row at all: `active_runs_of_dags` filters on `run_type` alone and
`backfill_id` is nullable, so a `DagRun` with
`run_type=DagRunType.BACKFILL_JOB` covers the same case.
--
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]