seanmuth commented on code in PR #73693:
URL: https://github.com/apache/airflow/pull/73693#discussion_r4145876843
##########
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:
Fixed. `is_at_max_active_runs` now has the short description from #73692,
and `active_runs_count` and `queued_runs_count` each have their own (running /
queued, backfill runs excluded). The spec, UI client and airflowctl are
regenerated.
---
Drafted-by: Claude Code (Opus 5.5); reviewed by @seanmuth before posting
##########
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:
Done. One query grouped by state gives both counts, and
`is_at_max_active_runs` is their sum against the limit. There's a single
backfill filter, `run_type != BACKFILL_JOB`, matching `active_runs_of_dags`.
---
Drafted-by: Claude Code (Opus 5.5); reviewed by @seanmuth before posting
##########
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:
Added `activeRunsOfMax` next to `activeRunsWithQueued`, so both branches go
through i18n.
---
Drafted-by: Claude Code (Opus 5.5); reviewed by @seanmuth before posting
##########
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:
Added the `3 of 1 (2 queued)` over-capacity case.
---
Drafted-by: Claude Code (Opus 5.5); reviewed by @seanmuth before posting
##########
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:
Dropped the `Backfill` rows and the import instead. With the filter on
`run_type`, a `DagRun` with `run_type=BACKFILL_JOB` covers it and nothing
outlives the test.
---
Drafted-by: Claude Code (Opus 5.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]