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


##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py:
##########
@@ -265,18 +265,51 @@ def get_dag_details(
     )
 
     # Count only running Dag runs: this stat shows runs that are actually 
executing right now.
+    # Excludes backfill runs -- those are gated by Backfill.max_active_runs, 
not the Dag's own
+    # max_active_runs, so a backfill with a higher limit would otherwise 
render as exceeding it.
     active_runs_count = (
         session.scalar(
             select(func.count())
             .select_from(DagRun)
-            .where(DagRun.dag_id == dag_id, DagRun.state == 
DagRunState.RUNNING)
+            .where(
+                DagRun.dag_id == dag_id,
+                DagRun.state == DagRunState.RUNNING,
+                DagRun.backfill_id.is_(None),
+            )
+        )
+        or 0
+    )
+
+    # Count queued Dag runs: these are waiting for an active run to finish 
before they can start.
+    # Excludes backfill runs for the same reason as active_runs_count above.
+    queued_runs_count = (
+        session.scalar(
+            select(func.count())
+            .select_from(DagRun)
+            .where(
+                DagRun.dag_id == dag_id,
+                DagRun.state == DagRunState.QUEUED,
+                DagRun.backfill_id.is_(None),
+            )
         )
         or 0
     )
 
-    # Add is_favorite and active_runs_count fields to the Dag model
+    # Computed fresh here rather than read from 
DagModel.exceeds_max_non_backfill: that column
+    # is a scheduler-side cache written on parse and on a handful of scheduler 
events, but never
+    # by a manual/API/operator trigger -- so for a schedule=None Dag it can 
stay stale (wrong in
+    # either direction) for as long as min_file_process_interval, or for a 
run's entire lifetime
+    # if it finishes before the next parse.
+    non_backfill_active_runs_count = DagRun.active_runs_of_dags(
+        dag_ids=[dag_id], exclude_backfill=True, session=session
+    ).get(dag_id, 0)
+    is_at_max_active_runs = non_backfill_active_runs_count >= 
(dag_model.max_active_runs or 0)

Review Comment:
   With both counts now excluding backfill, 
`active_runs_of_dags(exclude_backfill=True)` is just `active_runs_count + 
queued_runs_count`, so this third SELECT can go: one `select(DagRun.state, 
func.count())...group_by(DagRun.state)` gives both counts, and 
`is_at_max_active_runs` is their sum against the limit. That also stops the 
three values disagreeing when a run is triggered between the separate queries, 
and leaves one backfill filter instead of two (`run_type != BACKFILL_JOB` here 
vs `backfill_id IS NULL` above).



##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py:
##########
@@ -221,6 +222,19 @@ class DAGDetailsResponse(DAGResponse):
     owner_links: dict[str, str] | None = None
     is_favorite: bool = False
     active_runs_count: int = 0
+    queued_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 "
+            "active_runs_count counts RUNNING runs only and includes backfill 
runs. A Dag with "

Review Comment:
   Now that `active_runs_count` excludes backfill runs, this description is 
wrong: it still says `active_runs_count` "includes backfill runs", and the 
one-running-backfill-run example now returns `active_runs_count: 0`. It also 
calls RUNNING + QUEUED the scheduler's promotion check, but 
`get_queued_dag_runs_to_set_running` gates on RUNNING only; RUNNING + QUEUED is 
the run-creation gate (`exceeds_max_non_backfill` in `dags_needing_dagruns`). 
Since `active_runs_count` has now meant three things (RUNNING + QUEUED in 
3.3.0, RUNNING with backfill in 3.3.2, RUNNING without backfill here), could it 
and `queued_runs_count` get a `Field(description=...)` of their own, with the 
spec, UI client and airflowctl models regenerated?



##########
airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx:
##########
@@ -93,14 +99,77 @@ describe("Header", () => {
     expect(screen.getByText("2 of 2")).toBeInTheDocument();
   });
 
+  it("does not show an info icon or queued count when nothing is queued", () 
=> {
+    render(
+      <Wrapper>
+        <Header dag={{ ...mockDag, active_runs_count: 1, max_active_runs: 2 }} 
/>
+      </Wrapper>,
+    );
+
+    
expect(screen.queryByTestId("active-runs-exceeds-max-info")).not.toBeInTheDocument();
+    expect(screen.getByText("1 of 2")).toBeInTheDocument();
+  });
+
+  it("does not show the icon when exactly at capacity with nothing queued", () 
=> {
+    render(
+      <Wrapper>
+        <Header dag={{ ...mockDag, active_runs_count: 2, max_active_runs: 2 }} 
/>
+      </Wrapper>,
+    );
+
+    
expect(screen.queryByTestId("active-runs-exceeds-max-info")).not.toBeInTheDocument();
+    expect(screen.getByText("2 of 2")).toBeInTheDocument();
+  });
+
+  it("shows an info icon and the queued count when runs are queued behind the 
maximum", () => {
+    render(
+      <Wrapper>
+        <Header dag={{ ...mockDag, active_runs_count: 1, max_active_runs: 1, 
queued_runs_count: 2 }} />

Review Comment:
   Every icon-shown case is exactly at capacity, so changing the `>=` in 
`isBlockedByMaxActiveRuns` to `===` still passes all of these. Over capacity is 
reachable (lower `max_active_runs` while runs are active), so an 
`active_runs_count: 3, max_active_runs: 1, queued_runs_count: 2` case asserting 
the icon and `"3 of 1 (2 queued)"` would pin it.



##########
airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx:
##########
@@ -113,11 +121,28 @@ export const Header = ({
     },
     ...nextRunStat,
     {
-      label: translate("dagDetails.activeRuns"),
+      key: "activeRuns",
+      label:
+        isBlockedByMaxActiveRuns && hasQueuedRuns ? (
+          <HStack gap={1}>
+            {translate("dagDetails.activeRuns")}
+            <Tooltip 
content={translate("dagDetails.activeRunsExceedsMaxTooltip")} portalled>
+              <FiInfo data-testid="active-runs-exceeds-max-info" />
+            </Tooltip>
+          </HStack>
+        ) : (
+          translate("dagDetails.activeRuns")
+        ),
       value:
         dag?.max_active_runs === undefined
           ? undefined
-          : `${dag.active_runs_count ?? 0} of ${dag.max_active_runs}`,
+          : hasQueuedRuns
+            ? translate("dagDetails.activeRunsWithQueued", {
+                activeRuns: dag.active_runs_count ?? 0,
+                maxActiveRuns: dag.max_active_runs,
+                queuedRuns: dag.queued_runs_count,
+              })
+            : `${dag.active_runs_count ?? 0} of ${dag.max_active_runs}`,

Review Comment:
   Only the queued branch goes through i18n now; this one still hard-codes `of` 
in English, so in a translated locale the stat switches language as the queue 
drains. An `activeRunsOfMax` key next to `activeRunsWithQueued` would keep the 
two consistent.



##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py:
##########
@@ -1498,19 +1513,177 @@ def test_dag_details_includes_active_runs_count(self, 
session, test_client):
         assert response.status_code == 200
         body = response.json()
 
-        # Verify active_runs_count field is present and correct
-        assert "active_runs_count" in body
-        assert isinstance(body["active_runs_count"], int)
-        assert body["active_runs_count"] == 1  # only running counts, queued 
does not
+        assert body["active_runs_count"] == 1  # only running counts, 
queued/success do not
+        assert body["queued_runs_count"] == 2  # only queued counts, 
running/success do not
 
         # Test with DAG that has no active runs
         response = test_client.get(f"/dags/{DAG1_ID}/details")
         assert response.status_code == 200
         body = response.json()
 
-        assert "active_runs_count" in body
-        assert isinstance(body["active_runs_count"], int)
         assert body["active_runs_count"] == 0
+        assert body["queued_runs_count"] == 0
+
+    def 
test_dag_details_active_runs_count_and_queued_runs_count_exclude_backfill_runs(
+        self, session, test_client
+    ):
+        """Backfill runs don't count against the Dag's own max_active_runs, so 
they're excluded."""
+        from airflow.models.backfill import Backfill

Review Comment:
   This import can go at the top of the module, as `test_backfills.py` and 
`models/test_backfill.py` already do; there's no cycle to avoid here. The 
`Backfill` row committed here also outlives the test: `_clear_db` doesn't call 
`clear_db_backfills()`, and `Backfill.dag_id` has no FK, so `clear_db_dags` 
won't sweep it.



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