Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
50ac973
Expose exceeds_max_active_runs on the Dag details API response
seanmuth Sep 24, 2026
a989d94
Fix exceeds_max_non_backfill always false for non-schedulable Dags
seanmuth Sep 24, 2026
91504eb
Rename exceeds_max_active_runs to is_at_max_active_runs, batch its query
seanmuth Sep 25, 2026
7fa35fb
Fix test asserting active_runs_of_dags is skipped for non-schedulable…
seanmuth Sep 25, 2026
a6cee5c
Update the statement-budget test for the batched active-run-count query
seanmuth Sep 29, 2026
c1d3055
Compute is_at_max_active_runs at request time instead of aliasing a s…
seanmuth Sep 29, 2026
d3cb9fd
Shorten is_at_max_active_runs description and drop stale design notes
seanmuth Sep 30, 2026
7cb55d4
UI: Warn when a Dag's active runs exceed max_active_runs
seanmuth Sep 24, 2026
237bba9
Show queued run count and switch to an info icon on the Dag header
seanmuth Sep 24, 2026
48d2efd
UI: Trim the queued-runs tooltip wording
seanmuth Sep 24, 2026
9773a8e
Regenerate OpenAPI spec and UI client after rebase
seanmuth Sep 25, 2026
11e088f
Regenerate airflowctl datamodels after rebase
seanmuth Sep 25, 2026
16f1570
Fix Header.test.tsx assertion for the interpolated queued-count string
seanmuth Sep 25, 2026
cc82d49
Exclude backfill runs from active_runs_count and queued_runs_count
seanmuth Sep 29, 2026
c54f61c
Fix the queued-runs icon's false positive, portal, and key hack
seanmuth Sep 29, 2026
1906b01
Regenerate OpenAPI spec and UI client after rebase
seanmuth Sep 29, 2026
b578b39
Regenerate airflowctl datamodels after rebase
seanmuth Sep 29, 2026
b74476a
Fix backfill-exclusion test for this branch's own active_runs_count fix
seanmuth Sep 29, 2026
f49dbc9
Count active and queued runs in one query, describe both fields
seanmuth Sep 30, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
from pydantic import (
AliasGenerator,
ConfigDict,
Field,
computed_field,
field_serializer,
field_validator,
Expand Down Expand Up @@ -220,7 +221,18 @@ class DAGDetailsResponse(DAGResponse):
rerun_with_latest_version: bool | None = None
owner_links: dict[str, str] | None = None
is_favorite: bool = False
active_runs_count: int = 0
active_runs_count: int = Field(
default=0, description="Number of currently running runs (backfill runs excluded)."
)
queued_runs_count: int = Field(
default=0, description="Number of currently queued runs (backfill runs excluded)."
)
is_at_max_active_runs: bool = Field(
description=(
"Whether this Dag is currently at its max_active_runs limit, counting running and "
"queued runs (backfill runs excluded)."
)
)
team_name: str | None = None

@field_validator("timezone", mode="before")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13437,7 +13437,18 @@ components:
active_runs_count:
type: integer
title: Active Runs Count
description: Number of currently running runs (backfill runs excluded).
default: 0
queued_runs_count:
type: integer
title: Queued Runs Count
description: Number of currently queued runs (backfill runs excluded).
default: 0
is_at_max_active_runs:
type: boolean
title: Is At Max Active Runs
description: Whether this Dag is currently at its max_active_runs limit,
counting running and queued runs (backfill runs excluded).
team_name:
anyOf:
- type: string
Expand Down Expand Up @@ -13511,6 +13522,7 @@ components:
- timezone
- last_parsed
- default_args
- is_at_max_active_runs
- is_backfillable
- file_token
- concurrency
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,7 @@
from airflow.models.dagrun import DagRun
from airflow.utils.sqlalchemy import with_row_locks
from airflow.utils.state import DagRunState, DagSchedulingState
from airflow.utils.types import DagRunType

dags_router = AirflowRouter(tags=["DAG"], prefix="/dags")

Expand Down Expand Up @@ -264,19 +265,33 @@ def get_dag_details(
is not None
)

# Count only running Dag runs: this stat shows runs that are actually executing right now.
active_runs_count = (
session.scalar(
select(func.count())
.select_from(DagRun)
.where(DagRun.dag_id == dag_id, DagRun.state == DagRunState.RUNNING)
# One query for both counts so they can't disagree with each other if a run changes state
# in between. Backfill runs are excluded: they're gated by Backfill.max_active_runs, not
# the Dag's own max_active_runs.
active_runs_count = queued_runs_count = 0
for state, count in session.execute(
select(DagRun.state, func.count())
.where(
DagRun.dag_id == dag_id,
DagRun.state.in_((DagRunState.RUNNING, DagRunState.QUEUED)),
DagRun.run_type != DagRunType.BACKFILL_JOB,
)
or 0
)
.group_by(DagRun.state)
):
if state == DagRunState.RUNNING:
active_runs_count = count
else:
queued_runs_count = count

# Computed here rather than read from DagModel.exceeds_max_non_backfill: that scheduler-side
# cache isn't updated when a run is triggered manually or through the API.
is_at_max_active_runs = active_runs_count + queued_runs_count >= (dag_model.max_active_runs or 0)

# Add is_favorite and active_runs_count fields to the Dag model
# Add is_favorite, active_runs_count, queued_runs_count, and is_at_max_active_runs fields
setattr(dag_model, "is_favorite", is_favorite)
setattr(dag_model, "active_runs_count", active_runs_count)
setattr(dag_model, "queued_runs_count", queued_runs_count)
setattr(dag_model, "is_at_max_active_runs", is_at_max_active_runs)

return DAGDetailsResponse.model_validate(dag_model)

Expand Down
26 changes: 14 additions & 12 deletions airflow-core/src/airflow/dag_processing/collection.py
Original file line number Diff line number Diff line change
Expand Up @@ -167,18 +167,19 @@ def _get_latest_runs_stmt_partitioned(dag_id: str) -> Select:

class _RunInfo(NamedTuple):
latest_run: DagRun | None
num_active_runs: int

@classmethod
def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self:
"""
Query the run counts from the db.
Query the latest-run info from the db.

:param dags: dict of dags to query
:param dag: the Dag to calculate run info for
:param session: DB session
"""
# Skip these queries entirely if no Dags can be scheduled to save time.
# Only used to calculate the *next* scheduled run, so it can be skipped for Dags that can
# never be scheduled in the first place.
if not dag.timetable.can_be_scheduled:
return cls(None, 0)
return cls(None)

if dag.timetable.partitioned:
log.debug("Getting latest run for partitioned Dag", dag_id=dag.dag_id)
Expand All @@ -195,12 +196,7 @@ def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self:
)
else:
log.debug("no latest run found", dag_id=dag.dag_id)
active_run_counts = DagRun.active_runs_of_dags(
dag_ids=[dag.dag_id],
exclude_backfill=True,
session=session,
)
return cls(latest_run, active_run_counts.get(dag.dag_id, 0))
return cls(latest_run)


def _update_dag_tags(tag_names: set[str], dm: DagModel, *, session: Session) -> None:
Expand Down Expand Up @@ -623,6 +619,12 @@ def update_dags(
session: Session,
) -> None:
# we exclude backfill from active run counts since their concurrency is separate
# Batched once across every Dag in this update rather than inside the loop below,
# since DagRun.active_runs_of_dags already accepts a list of dag_ids -- calling it
# per-Dag would turn this into an N+1 query.
active_run_counts = DagRun.active_runs_of_dags(
dag_ids=list(orm_dags), exclude_backfill=True, session=session
)
for dag_id, dm in sorted(orm_dags.items()):
run_info = _RunInfo.calculate(dag=self.dags[dag_id], session=session)
dag = self.dags[dag_id]
Expand Down Expand Up @@ -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
dm.calculate_dagrun_date_fields(dag, reference_run=reference_run)
if not dag.timetable.asset_condition:
dm.schedule_asset_references = []
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3283,8 +3283,20 @@ export const $DAGDetailsResponse = {
active_runs_count: {
type: 'integer',
title: 'Active Runs Count',
description: 'Number of currently running runs (backfill runs excluded).',
default: 0
},
queued_runs_count: {
Comment thread
seanmuth marked this conversation as resolved.
type: 'integer',
title: 'Queued Runs Count',
description: 'Number of currently queued runs (backfill runs excluded).',
default: 0
},
is_at_max_active_runs: {
type: 'boolean',
title: 'Is At Max Active Runs',
description: 'Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).'
},
team_name: {
anyOf: [
{
Expand Down Expand Up @@ -3331,7 +3343,7 @@ Deprecated: Use max_active_tasks instead.`,
}
},
type: 'object',
required: ['dag_id', 'dag_display_name', 'is_paused', 'is_stale', 'last_parsed_time', 'last_parse_duration', 'last_expired', 'bundle_name', 'bundle_version', 'relative_fileloc', 'fileloc', 'description', 'timetable_summary', 'timetable_description', 'timetable_partitioned', 'timetable_periodic', 'tags', 'max_active_tasks', 'max_active_runs', 'max_consecutive_failed_dag_runs', 'has_task_concurrency_limits', 'has_import_errors', 'next_dagrun_logical_date', 'next_dagrun_data_interval_start', 'next_dagrun_data_interval_end', 'next_dagrun_run_after', 'allowed_run_types', 'owners', 'catchup', 'dag_run_timeout', 'asset_expression', 'doc_md', 'start_date', 'end_date', 'is_paused_upon_creation', 'params', 'render_template_as_native_obj', 'template_search_path', 'timezone', 'last_parsed', 'default_args', 'is_backfillable', 'file_token', 'concurrency', 'latest_dag_version'],
required: ['dag_id', 'dag_display_name', 'is_paused', 'is_stale', 'last_parsed_time', 'last_parse_duration', 'last_expired', 'bundle_name', 'bundle_version', 'relative_fileloc', 'fileloc', 'description', 'timetable_summary', 'timetable_description', 'timetable_partitioned', 'timetable_periodic', 'tags', 'max_active_tasks', 'max_active_runs', 'max_consecutive_failed_dag_runs', 'has_task_concurrency_limits', 'has_import_errors', 'next_dagrun_logical_date', 'next_dagrun_data_interval_start', 'next_dagrun_data_interval_end', 'next_dagrun_run_after', 'allowed_run_types', 'owners', 'catchup', 'dag_run_timeout', 'asset_expression', 'doc_md', 'start_date', 'end_date', 'is_paused_upon_creation', 'params', 'render_template_as_native_obj', 'template_search_path', 'timezone', 'last_parsed', 'default_args', 'is_at_max_active_runs', 'is_backfillable', 'file_token', 'concurrency', 'latest_dag_version'],
title: 'DAGDetailsResponse',
description: 'Specific serializer for Dag Details responses.'
} as const;
Expand Down
11 changes: 11 additions & 0 deletions airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -910,7 +910,18 @@ export type DAGDetailsResponse = {
[key: string]: (string);
} | null;
is_favorite?: boolean;
/**
* Number of currently running runs (backfill runs excluded).
*/
active_runs_count?: number;
/**
* Number of currently queued runs (backfill runs excluded).
*/
queued_runs_count?: number;
/**
* Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).
*/
is_at_max_active_runs: boolean;
team_name?: string | null;
/**
* Whether this Dag's schedule supports backfilling.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,9 @@
"dagBundle_other": "Dag Bundles",
"dagDetails": {
"activeRuns": "Active Runs",
"activeRunsExceedsMaxTooltip": "Runs beyond the limit will not start until an existing active run completes.",
"activeRunsOfMax": "{{activeRuns}} of {{maxActiveRuns}}",
"activeRunsWithQueued": "{{activeRuns}} of {{maxActiveRuns}} ({{queuedRuns}} queued)",
"catchup": "Catchup",
"dagRunTimeout": "Dag Run Timeout",
"defaultArgs": "Default Args",
Expand Down
13 changes: 11 additions & 2 deletions airflow-core/src/airflow/ui/src/components/HeaderCard.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,17 @@ import { StateBadge } from "src/components/StateBadge";

import { DagDeactivatedBanner } from "./DagDeactivatedBanner";

// A stat needs a stable React key. A string label already is one, so key is optional there;
// any other label (e.g. a label with an inline icon/tooltip) must supply its own key explicitly.
type Stat =
| { readonly key: string; readonly label: ReactNode; readonly value: ReactNode | string }
| { readonly key?: string; readonly label: string; readonly value: ReactNode | string };

type Props = {
readonly actions?: ReactNode;
readonly icon: ReactNode;
readonly state?: TaskInstanceState | null;
readonly stats: Array<{ key?: string; label: string; value: ReactNode | string }>;
readonly stats: Array<Stat>;
readonly subTitle?: ReactNode | string;
readonly title: ReactNode | string;
readonly type: "asset" | "dag" | "dagBundle" | "dagRun" | "task" | "taskGroup" | "taskInstance";
Expand Down Expand Up @@ -78,7 +84,10 @@ export const HeaderCard = ({ actions, icon, state, stats, subTitle, title, type

<HStack alignItems="flex-start" flexWrap="wrap" gap={6} my={3}>
{stats.map((stat) => (
<Box data-testid="stat" key={stat.key ?? stat.label}>
<Box
data-testid="stat"
key={stat.key ?? (typeof stat.label === "string" ? stat.label : undefined)}
>
<Box
color="fg.muted"
fontSize="xs"
Expand Down
88 changes: 84 additions & 4 deletions airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,13 @@
import "@testing-library/jest-dom";
import { render, screen } from "@testing-library/react";
import type { DAGDetailsResponse } from "openapi-gen/requests/types.gen";
import { afterEach, describe, expect, it, vi } from "vitest";
import { afterEach, beforeAll, describe, expect, it, vi } from "vitest";

import i18n from "src/i18n/config";
import { MOCK_DAG } from "src/mocks/handlers/dag";
import { Wrapper } from "src/utils/Wrapper";

import commonLocale from "../../../public/i18n/locales/en/common.json";
import { Header } from "./Header";

const mockConfig: Record<string, unknown> = { multi_team: false };
Expand All @@ -49,13 +50,18 @@ const mockDag = {
next_dagrun_logical_date: "2024-08-22T00:00:00+00:00",
next_dagrun_run_after: "2024-08-22T19:00:00+00:00",
owner_links: {},
queued_runs_count: 0,
relative_fileloc: "stale_dag.py",
tags: [],
timetable_partitioned: false,
timetable_summary: "* * * * *",
} as unknown as DAGDetailsResponse;

describe("Header", () => {
beforeAll(() => {
i18n.addResourceBundle("en", "common", commonLocale, true, true);
});

afterEach(() => {
mockConfig.multi_team = false;
});
Expand All @@ -67,7 +73,7 @@ describe("Header", () => {
</Wrapper>,
);

expect(screen.queryByText(i18n.t("dag:dagDetails.nextRun"))).not.toBeInTheDocument();
expect(screen.queryByText(i18n.t("common:dagDetails.nextRun"))).not.toBeInTheDocument();
expect(screen.queryByRole("button", { name: "Reparse Dag" })).not.toBeInTheDocument();
});

Expand All @@ -78,7 +84,7 @@ describe("Header", () => {
</Wrapper>,
);

expect(screen.getByText(i18n.t("dag:dagDetails.nextRun"))).toBeInTheDocument();
expect(screen.getByText(i18n.t("common:dagDetails.nextRun"))).toBeInTheDocument();
expect(screen.queryByText("2024-08-22 19:00:00")).not.toBeInTheDocument();
});

Expand All @@ -93,14 +99,88 @@ 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 }} />

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added the 3 of 1 (2 queued) over-capacity case.


Drafted-by: Claude Code (Opus 5.5); reviewed by @seanmuth before posting

</Wrapper>,
);

expect(screen.getByTestId("active-runs-exceeds-max-info")).toBeInTheDocument();
expect(screen.getByText("1 of 1 (2 queued)")).toBeInTheDocument();
});

it("shows the info icon when over capacity, e.g. after max_active_runs was lowered", () => {
render(
<Wrapper>
<Header dag={{ ...mockDag, active_runs_count: 3, max_active_runs: 1, queued_runs_count: 2 }} />
</Wrapper>,
);

expect(screen.getByTestId("active-runs-exceeds-max-info")).toBeInTheDocument();
expect(screen.getByText("3 of 1 (2 queued)")).toBeInTheDocument();
});

it("shows the queued count without the icon when below capacity (scheduler hasn't caught up yet)", () => {
render(
<Wrapper>
<Header dag={{ ...mockDag, active_runs_count: 0, max_active_runs: 2, queued_runs_count: 1 }} />
</Wrapper>,
);

expect(screen.queryByTestId("active-runs-exceeds-max-info")).not.toBeInTheDocument();
expect(screen.getByText("0 of 2 (1 queued)")).toBeInTheDocument();
});

it("shows the queued count without the icon for a paused Dag even at capacity", () => {
render(
<Wrapper>
<Header
dag={{
...mockDag,
active_runs_count: 2,
is_paused: true,
max_active_runs: 2,
queued_runs_count: 1,
}}
/>
</Wrapper>,
);

expect(screen.queryByTestId("active-runs-exceeds-max-info")).not.toBeInTheDocument();
expect(screen.getByText("2 of 2 (1 queued)")).toBeInTheDocument();
});

it("renders the draining badge instead of the next run timestamp for a draining Dag", () => {
render(
<Wrapper>
<Header dag={{ ...mockDag, is_stale: false, scheduling_state: "draining" }} />
</Wrapper>,
);

expect(screen.getByText(i18n.t("dag:dagDetails.nextRun"))).toBeInTheDocument();
expect(screen.getByText(i18n.t("common:dagDetails.nextRun"))).toBeInTheDocument();
expect(screen.queryByText("2024-08-22 19:00:00")).not.toBeInTheDocument();
expect(screen.getByTestId("draining-badge")).toBeInTheDocument();
});
Expand Down
Loading
Loading