From 50ac973f5594c253b13c41d0cc17925b1df9a5c8 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Thu, 24 Sep 2026 15:08:44 -0500 Subject: [PATCH 01/19] Expose exceeds_max_active_runs on the Dag details API response DagModel.exceeds_max_non_backfill has been persisted since 3.2.0 (migration 0099) and is already used internally by the scheduler to avoid re-evaluating Dags that are already at their concurrency limit, but it was never surfaced anywhere in the public API. There is currently no way for a client to tell "is this Dag currently blocked on max_active_runs" without independently computing active-run counts against the limit. This exposes it on GET /dags/{dag_id}/details as exceeds_max_active_runs, aliased to the underlying model column via the existing DAG_ALIAS_MAPPING mechanism. Part of #73686. --- .../api_fastapi/core_api/datamodels/dags.py | 2 ++ .../openapi/v2-rest-api-generated.yaml | 4 ++++ .../ui/openapi-gen/requests/schemas.gen.ts | 6 ++++- .../ui/openapi-gen/requests/types.gen.ts | 1 + .../core_api/routes/public/test_dags.py | 22 +++++++++++++++++++ .../airflowctl/api/datamodels/generated.py | 1 + .../tests/airflow_ctl/api/test_operations.py | 1 + 7 files changed, 36 insertions(+), 1 deletion(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py index db3b013976f6f..bf613f432b08f 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py @@ -199,6 +199,7 @@ class DAGDetailsResponse(DAGResponse): "dag_run_timeout": "dagrun_timeout", "last_parsed": "last_loaded", "template_search_path": "template_searchpath", + "exceeds_max_active_runs": "exceeds_max_non_backfill", **DAG_ALIAS_MAPPING, }.get(field_name, field_name), ), @@ -221,6 +222,7 @@ class DAGDetailsResponse(DAGResponse): owner_links: dict[str, str] | None = None is_favorite: bool = False active_runs_count: int = 0 + exceeds_max_active_runs: bool team_name: str | None = None @field_validator("timezone", mode="before") diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 49694cd5d035c..00c44e439d03e 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13438,6 +13438,9 @@ components: type: integer title: Active Runs Count default: 0 + exceeds_max_active_runs: + type: boolean + title: Exceeds Max Active Runs team_name: anyOf: - type: string @@ -13511,6 +13514,7 @@ components: - timezone - last_parsed - default_args + - exceeds_max_active_runs - is_backfillable - file_token - concurrency diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index 71d21c125c53e..16d53f231fd6d 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3285,6 +3285,10 @@ export const $DAGDetailsResponse = { title: 'Active Runs Count', default: 0 }, + exceeds_max_active_runs: { + type: 'boolean', + title: 'Exceeds Max Active Runs' + }, team_name: { anyOf: [ { @@ -3331,7 +3335,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', 'exceeds_max_active_runs', 'is_backfillable', 'file_token', 'concurrency', 'latest_dag_version'], title: 'DAGDetailsResponse', description: 'Specific serializer for Dag Details responses.' } as const; diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index 776009c9afc25..f32036fb41e3a 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -911,6 +911,7 @@ export type DAGDetailsResponse = { } | null; is_favorite?: boolean; active_runs_count?: number; + exceeds_max_active_runs: boolean; team_name?: string | null; /** * Whether this Dag's schedule supports backfilling. diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index 73d2d08913cb7..0a5f1e27b924d 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1352,6 +1352,7 @@ def test_dag_details( "description": None, "doc_md": "details", "end_date": None, + "exceeds_max_active_runs": False, "fileloc": __file__, "file_token": file_token, "has_import_errors": False, @@ -1512,6 +1513,27 @@ 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_exceeds_max_active_runs(self, session, test_client): + """Test that DAG details include the exceeds_max_active_runs field.""" + dag_model = session.get(DagModel, DAG2_ID) + dag_model.exceeds_max_non_backfill = True + session.commit() + + response = test_client.get(f"/dags/{DAG2_ID}/details") + assert response.status_code == 200 + body = response.json() + + assert "exceeds_max_active_runs" in body + assert body["exceeds_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 "exceeds_max_active_runs" in body + assert body["exceeds_max_active_runs"] is False + def test_dag_details_team_name_none_without_multi_team(self, test_client): """Without multi-team enabled, ``team_name`` stays ``None`` and no lookup happens.""" response = test_client.get(f"/dags/{DAG1_ID}/details") diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 8938ad8d0534b..748aabf35843d 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2926,6 +2926,7 @@ class DAGDetailsResponse(BaseModel): owner_links: Annotated[dict[str, str] | None, Field(title="Owner Links")] = None is_favorite: Annotated[bool | None, Field(title="Is Favorite")] = False active_runs_count: Annotated[int | None, Field(title="Active Runs Count")] = 0 + exceeds_max_active_runs: Annotated[bool, Field(title="Exceeds Max Active Runs")] team_name: Annotated[str | None, Field(title="Team Name")] = None is_backfillable: Annotated[ bool, Field(description="Whether this Dag's schedule supports backfilling.", title="Is Backfillable") diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py b/airflow-ctl/tests/airflow_ctl/api/test_operations.py index f1eaf2e89efe1..8cab3ed9945ed 100644 --- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py +++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py @@ -1120,6 +1120,7 @@ class TestDagOperations: doc_md=None, start_date=datetime.datetime(2024, 12, 31, 23, 59, 59), end_date=datetime.datetime(2025, 1, 1, 0, 0, 0), + exceeds_max_active_runs=False, is_paused_upon_creation=False, params={}, render_template_as_native_obj=True, From a989d947cb5fea50a753332863e9eb17662538c5 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Thu, 24 Sep 2026 16:12:54 -0500 Subject: [PATCH 02/19] Fix exceeds_max_non_backfill always false for non-schedulable Dags The dag-processor's active-run-count calculation short-circuited to 0 for any Dag whose timetable can't be scheduled (e.g. schedule=None, triggered manually or via the API), regardless of how many runs were actually active. Since exceeds_max_non_backfill is recomputed on every parse cycle, this made it permanently unreliable for exactly the kind of Dag that most needs it: one triggered manually more often than its max_active_runs allows never has a schedule to be "scheduled" against, but still has a real concurrency limit. Only the latest-run lookup (used for scheduling the next run) is skippable for such Dags; the active-run count is not. --- .../src/airflow/dag_processing/collection.py | 21 +++--- .../unit/dag_processing/test_collection.py | 68 ++++++++++++++++++- 2 files changed, 80 insertions(+), 9 deletions(-) diff --git a/airflow-core/src/airflow/dag_processing/collection.py b/airflow-core/src/airflow/dag_processing/collection.py index 36b948abad6c2..fc206974e23ed 100644 --- a/airflow-core/src/airflow/dag_processing/collection.py +++ b/airflow-core/src/airflow/dag_processing/collection.py @@ -176,9 +176,19 @@ def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self: :param dags: dict of dags to query """ - # Skip these queries entirely if no Dags can be scheduled to save time. + active_run_counts = DagRun.active_runs_of_dags( + dag_ids=[dag.dag_id], + exclude_backfill=True, + session=session, + ) + num_active_runs = active_run_counts.get(dag.dag_id, 0) + + # The latest-run lookup below is only used to calculate the *next* scheduled run, so it + # can be skipped for Dags that can never be scheduled in the first place. num_active_runs + # is still meaningful for such Dags (e.g. a schedule=None Dag triggered manually more + # often than its max_active_runs allows) and must always be computed. if not dag.timetable.can_be_scheduled: - return cls(None, 0) + return cls(None, num_active_runs) if dag.timetable.partitioned: log.debug("Getting latest run for partitioned Dag", dag_id=dag.dag_id) @@ -195,12 +205,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, num_active_runs) def _update_dag_tags(tag_names: set[str], dm: DagModel, *, session: Session) -> None: diff --git a/airflow-core/tests/unit/dag_processing/test_collection.py b/airflow-core/tests/unit/dag_processing/test_collection.py index be36290abcece..d1585bc9f9213 100644 --- a/airflow-core/tests/unit/dag_processing/test_collection.py +++ b/airflow-core/tests/unit/dag_processing/test_collection.py @@ -68,7 +68,8 @@ from airflow.serialization.serialized_objects import LazyDeserializedDAG from airflow.timetables.simple import PartitionedAtRuntime from airflow.triggers.base import BaseEventTrigger -from airflow.utils.types import DagRunType +from airflow.utils.state import DagRunState +from airflow.utils.types import DagRunTriggeredByType, DagRunType from tests_common.test_utils.config import conf_vars from tests_common.test_utils.db import ( @@ -1448,6 +1449,71 @@ def test_max_active_runs_defaults_from_conf_when_none(self, testing_dag_bundle, orm_dag = session.get(DagModel, "dag_max_runs_default") assert orm_dag.max_active_runs == 4 + def test_exceeds_max_non_backfill_reflects_real_active_runs_for_non_schedulable_dag( + self, testing_dag_bundle, session, dag_maker + ): + """A schedule=None Dag's exceeds_max_non_backfill must reflect real active-run counts. + + ``can_be_scheduled`` is False for schedule=None Dags, but that must not make the + parser hardcode num_active_runs to 0 -- the flag is read regardless of whether the + Dag can be automatically scheduled. + """ + with dag_maker("dag_schedule_none_exceeds_max", schedule=None, max_active_runs=1) as dag: + ... + + session.add( + DagRun( + dag_id=dag.dag_id, + run_id="running_run", + logical_date=tz.datetime(2024, 1, 1), + start_date=tz.utcnow(), + run_type=DagRunType.MANUAL, + state=DagRunState.RUNNING, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.add( + DagRun( + dag_id=dag.dag_id, + run_id="queued_run", + logical_date=tz.datetime(2024, 1, 2), + start_date=tz.utcnow(), + run_type=DagRunType.MANUAL, + state=DagRunState.QUEUED, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.commit() + + update_dag_parsing_results_in_db("testing", None, [dag], dict(), 0.1, set(), session) + + orm_dag = session.get(DagModel, "dag_schedule_none_exceeds_max") + assert orm_dag.exceeds_max_non_backfill is True + + def test_exceeds_max_non_backfill_false_within_limit_for_non_schedulable_dag( + self, testing_dag_bundle, session, dag_maker + ): + with dag_maker("dag_schedule_none_within_max", schedule=None, max_active_runs=2) as dag: + ... + + session.add( + DagRun( + dag_id=dag.dag_id, + run_id="running_run", + logical_date=tz.datetime(2024, 1, 1), + start_date=tz.utcnow(), + run_type=DagRunType.MANUAL, + state=DagRunState.RUNNING, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.commit() + + update_dag_parsing_results_in_db("testing", None, [dag], dict(), 0.1, set(), session) + + orm_dag = session.get(DagModel, "dag_schedule_none_within_max") + assert orm_dag.exceeds_max_non_backfill is False + def test_max_consecutive_failed_dag_runs_explicit_value_is_used( self, testing_dag_bundle, session, dag_maker ): From 91504eba37ee7a4d7b02fce13a80cf85c1d66658 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Fri, 25 Sep 2026 11:01:59 -0500 Subject: [PATCH 03/19] Rename exceeds_max_active_runs to is_at_max_active_runs, batch its query Review feedback from #73692: - "exceeds" was inaccurate: the flag is true once a Dag is at or above its limit, not only when strictly over it. Renamed the exposed field via the existing DAG_ALIAS_MAPPING mechanism; the underlying exceeds_max_non_backfill column is unchanged, since renaming a persisted column is a separate, more deliberate change. - The dag-processor's active-run-count calculation was already an N+1 (one query per Dag per parse cycle) before this PR's fix to the non-schedulable-Dag case; that fix just extended the same pattern to more Dags. DagRun.active_runs_of_dags already accepts a batch of dag_ids, so hoist the call out of the per-Dag loop in update_dags into a single call across every Dag in the update. --- .../api_fastapi/core_api/datamodels/dags.py | 4 +-- .../openapi/v2-rest-api-generated.yaml | 6 ++-- .../src/airflow/dag_processing/collection.py | 28 +++++++++++-------- .../ui/openapi-gen/requests/schemas.gen.ts | 6 ++-- .../ui/openapi-gen/requests/types.gen.ts | 2 +- .../core_api/routes/public/test_dags.py | 14 +++++----- .../airflowctl/api/datamodels/generated.py | 2 +- .../tests/airflow_ctl/api/test_operations.py | 2 +- 8 files changed, 35 insertions(+), 29 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py index bf613f432b08f..d1913abcb7c46 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py @@ -199,7 +199,7 @@ class DAGDetailsResponse(DAGResponse): "dag_run_timeout": "dagrun_timeout", "last_parsed": "last_loaded", "template_search_path": "template_searchpath", - "exceeds_max_active_runs": "exceeds_max_non_backfill", + "is_at_max_active_runs": "exceeds_max_non_backfill", **DAG_ALIAS_MAPPING, }.get(field_name, field_name), ), @@ -222,7 +222,7 @@ class DAGDetailsResponse(DAGResponse): owner_links: dict[str, str] | None = None is_favorite: bool = False active_runs_count: int = 0 - exceeds_max_active_runs: bool + is_at_max_active_runs: bool team_name: str | None = None @field_validator("timezone", mode="before") diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 00c44e439d03e..538fdf736b6b5 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13438,9 +13438,9 @@ components: type: integer title: Active Runs Count default: 0 - exceeds_max_active_runs: + is_at_max_active_runs: type: boolean - title: Exceeds Max Active Runs + title: Is At Max Active Runs team_name: anyOf: - type: string @@ -13514,7 +13514,7 @@ components: - timezone - last_parsed - default_args - - exceeds_max_active_runs + - is_at_max_active_runs - is_backfillable - file_token - concurrency diff --git a/airflow-core/src/airflow/dag_processing/collection.py b/airflow-core/src/airflow/dag_processing/collection.py index fc206974e23ed..c87375482b314 100644 --- a/airflow-core/src/airflow/dag_processing/collection.py +++ b/airflow-core/src/airflow/dag_processing/collection.py @@ -170,19 +170,15 @@ class _RunInfo(NamedTuple): num_active_runs: int @classmethod - def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self: + def calculate(cls, dag: LazyDeserializedDAG, *, num_active_runs: int, 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 num_active_runs: active run count for this Dag, pre-computed by the caller in one + batched query across every Dag in the current update so this doesn't become an N+1. + :param session: DB session """ - active_run_counts = DagRun.active_runs_of_dags( - dag_ids=[dag.dag_id], - exclude_backfill=True, - session=session, - ) - num_active_runs = active_run_counts.get(dag.dag_id, 0) - # The latest-run lookup below is only used to calculate the *next* scheduled run, so it # can be skipped for Dags that can never be scheduled in the first place. num_active_runs # is still meaningful for such Dags (e.g. a schedule=None Dag triggered manually more @@ -628,8 +624,18 @@ 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) + run_info = _RunInfo.calculate( + dag=self.dags[dag_id], + num_active_runs=active_run_counts.get(dag_id, 0), + session=session, + ) dag = self.dags[dag_id] dm.fileloc = dag.fileloc dm.relative_fileloc = dag.relative_fileloc diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index 16d53f231fd6d..2c20b295fb10a 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3285,9 +3285,9 @@ export const $DAGDetailsResponse = { title: 'Active Runs Count', default: 0 }, - exceeds_max_active_runs: { + is_at_max_active_runs: { type: 'boolean', - title: 'Exceeds Max Active Runs' + title: 'Is At Max Active Runs' }, team_name: { anyOf: [ @@ -3335,7 +3335,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', 'exceeds_max_active_runs', '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; diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index f32036fb41e3a..4e0dafda75bc2 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -911,7 +911,7 @@ export type DAGDetailsResponse = { } | null; is_favorite?: boolean; active_runs_count?: number; - exceeds_max_active_runs: boolean; + is_at_max_active_runs: boolean; team_name?: string | null; /** * Whether this Dag's schedule supports backfilling. diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index 0a5f1e27b924d..f4332331b70aa 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1352,11 +1352,11 @@ def test_dag_details( "description": None, "doc_md": "details", "end_date": None, - "exceeds_max_active_runs": False, "fileloc": __file__, "file_token": file_token, "has_import_errors": False, "has_task_concurrency_limits": True, + "is_at_max_active_runs": False, "is_backfillable": False, "is_favorite": False, "is_paused": False, @@ -1513,8 +1513,8 @@ 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_exceeds_max_active_runs(self, session, test_client): - """Test that DAG details include the exceeds_max_active_runs field.""" + def test_dag_details_includes_is_at_max_active_runs(self, session, test_client): + """Test that DAG details include the is_at_max_active_runs field.""" dag_model = session.get(DagModel, DAG2_ID) dag_model.exceeds_max_non_backfill = True session.commit() @@ -1523,16 +1523,16 @@ def test_dag_details_includes_exceeds_max_active_runs(self, session, test_client assert response.status_code == 200 body = response.json() - assert "exceeds_max_active_runs" in body - assert body["exceeds_max_active_runs"] is True + assert "is_at_max_active_runs" in body + 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 "exceeds_max_active_runs" in body - assert body["exceeds_max_active_runs"] is False + assert "is_at_max_active_runs" in body + assert body["is_at_max_active_runs"] is False def test_dag_details_team_name_none_without_multi_team(self, test_client): """Without multi-team enabled, ``team_name`` stays ``None`` and no lookup happens.""" diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 748aabf35843d..642091635f67c 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2926,7 +2926,7 @@ class DAGDetailsResponse(BaseModel): owner_links: Annotated[dict[str, str] | None, Field(title="Owner Links")] = None is_favorite: Annotated[bool | None, Field(title="Is Favorite")] = False active_runs_count: Annotated[int | None, Field(title="Active Runs Count")] = 0 - exceeds_max_active_runs: Annotated[bool, Field(title="Exceeds Max Active Runs")] + is_at_max_active_runs: Annotated[bool, Field(title="Is At Max Active Runs")] team_name: Annotated[str | None, Field(title="Team Name")] = None is_backfillable: Annotated[ bool, Field(description="Whether this Dag's schedule supports backfilling.", title="Is Backfillable") diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py b/airflow-ctl/tests/airflow_ctl/api/test_operations.py index 8cab3ed9945ed..748d05e49c4a0 100644 --- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py +++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py @@ -1120,7 +1120,7 @@ class TestDagOperations: doc_md=None, start_date=datetime.datetime(2024, 12, 31, 23, 59, 59), end_date=datetime.datetime(2025, 1, 1, 0, 0, 0), - exceeds_max_active_runs=False, + is_at_max_active_runs=False, is_paused_upon_creation=False, params={}, render_template_as_native_obj=True, From 7fa35fb1e07c09e2a6b7c0da584374921b5b8187 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Fri, 25 Sep 2026 11:55:10 -0500 Subject: [PATCH 04/19] Fix test asserting active_runs_of_dags is skipped for non-schedulable Dags test_bulk_write_to_db_interval_save_runtime encoded the old, buggy short-circuit this PR removes: active_runs_of_dags is now always batched once per update, regardless of whether any Dag in the batch can be scheduled, so exceeds_max_non_backfill stays accurate for schedule=None Dags too. Co-Authored-By: Claude --- airflow-core/tests/unit/models/test_dag.py | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/airflow-core/tests/unit/models/test_dag.py b/airflow-core/tests/unit/models/test_dag.py index b22227b695be4..2364541319a76 100644 --- a/airflow-core/tests/unit/models/test_dag.py +++ b/airflow-core/tests/unit/models/test_dag.py @@ -685,7 +685,14 @@ def test_bulk_write_to_db_multiple_dags(self, testing_dag_bundle): SerializedDAG.bulk_write_to_db("testing", None, dags) @pytest.mark.parametrize("interval", [None, "@daily"]) - def test_bulk_write_to_db_interval_save_runtime(self, testing_dag_bundle, interval): + def test_bulk_write_to_db_interval_batches_active_run_count_query(self, testing_dag_bundle, interval): + """active_runs_of_dags is called exactly once per update, batched across every Dag. + + This holds regardless of whether any Dag in the batch has a schedule: the active-run + count is needed even for schedule=None Dags (to keep exceeds_max_non_backfill accurate + for Dags that are only ever triggered manually/via the API), so this must not be skipped + just because none of the Dags being updated can be scheduled. + """ mock_active_runs_of_dags = mock.MagicMock(side_effect=DagRun.active_runs_of_dags) with mock.patch.object(DagRun, "active_runs_of_dags", mock_active_runs_of_dags): dags_null_timetable = [ @@ -693,10 +700,7 @@ def test_bulk_write_to_db_interval_save_runtime(self, testing_dag_bundle, interv create_scheduler_dag(DAG("dag-interval-test", schedule=interval, start_date=TEST_DATE)), ] SerializedDAG.bulk_write_to_db("testing", None, dags_null_timetable) - if interval: - mock_active_runs_of_dags.assert_called_once() - else: - mock_active_runs_of_dags.assert_not_called() + mock_active_runs_of_dags.assert_called_once() @pytest.mark.parametrize( ("state", "catchup", "expected_next_dagrun"), From a6cee5c2223ab9e3f11c0c120141bf9456ed7041 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Tue, 29 Sep 2026 12:12:49 -0400 Subject: [PATCH 05/19] Update the statement-budget test for the batched active-run-count query The N+1 fix moved DagRun.active_runs_of_dags from once per Dag to once per persistence call, so it now counts toward FIXED_PER_CALL instead of UNCHANGED_PER_DAG/REWRITE_PER_DAG. Co-Authored-By: Claude --- airflow-core/tests/unit/dag_processing/test_manager.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index fbdc9bc9b7529..6d885adc32c6d 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -270,10 +270,11 @@ def _statement_breakdown(counts: Counter[tuple[str, str]]) -> str: # content is unchanged; once the hash has moved and [core] min_serialized_dag_update_interval has # lapsed it rewrites it, which costs two more statements per Dag and nothing extra per call. The # per-call price is a file that parsed cleanly: one reporting import errors also looks up whichever -# of them are already recorded. -FIXED_PER_CALL = 9 -UNCHANGED_PER_DAG = 3 -REWRITE_PER_DAG = 5 +# of them are already recorded. Includes one active-run-count SELECT batched once per call across +# every Dag in the file, rather than once per Dag. +FIXED_PER_CALL = 10 +UNCHANGED_PER_DAG = 2 +REWRITE_PER_DAG = 4 # A file that failed to parse and so defines no Dags. Two of the five are import_error SELECTs: the # bounded lookup, and the listener re-reading the row the update beside it already had. IMPORT_ERROR_PER_CALL = 5 From c1d30558f6f52b88b5430467c3459a26c27fd00a Mon Sep 17 00:00:00 2001 From: seanmuth Date: Tue, 29 Sep 2026 12:25:07 -0400 Subject: [PATCH 06/19] Compute is_at_max_active_runs at request time instead of aliasing a stale scheduler cache exceeds_max_non_backfill is a scheduler-side cache written on parse and by a handful of scheduler events (scheduled-run creation, dagrun timeout, run finish) -- but never by a manual/API/operator trigger. For a schedule=None Dag, triggering a run left is_at_max_active_runs false until the next parse cycle (up to min_file_process_interval, 30s by default), or for the run's entire lifetime if it finished before that parse -- exactly the case this field exists to cover. Marking a run failed through the API had the reverse problem. get_dag_details now computes it fresh, the same way it already does for active_runs_count/queued_runs_count, via DagRun.active_runs_of_dags(exclude_backfill=True) compared against max_active_runs. Removed the now-unused alias to exceeds_max_non_backfill from DAGDetailsResponse; that column stays as the scheduler's own internal optimization, just no longer exposed through this field. Documented via Field(description=...) that this counts differently than active_runs_count: RUNNING+QUEUED with backfill runs excluded (matching the scheduler's own promotion check), vs active_runs_count's RUNNING-only that includes backfill runs. Also simplified _RunInfo/_RunInfo.calculate in the dag-processor: num_active_runs was only ever passed through unchanged since the N+1 fix moved its computation to the caller, so update_dags now reads it directly from the batched query result instead of round-tripping it through calculate's parameter and return value. Co-Authored-By: Claude --- .../api_fastapi/core_api/datamodels/dags.py | 15 +++- .../openapi/v2-rest-api-generated.yaml | 8 ++ .../core_api/routes/public/dags.py | 13 +++- .../src/airflow/dag_processing/collection.py | 23 ++---- .../ui/openapi-gen/requests/schemas.gen.ts | 3 +- .../ui/openapi-gen/requests/types.gen.ts | 3 + .../core_api/routes/public/test_dags.py | 77 ++++++++++++++++++- .../airflowctl/api/datamodels/generated.py | 8 +- 8 files changed, 125 insertions(+), 25 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py index d1913abcb7c46..6d670a9f769b3 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py @@ -28,6 +28,7 @@ from pydantic import ( AliasGenerator, ConfigDict, + Field, computed_field, field_serializer, field_validator, @@ -199,7 +200,6 @@ class DAGDetailsResponse(DAGResponse): "dag_run_timeout": "dagrun_timeout", "last_parsed": "last_loaded", "template_search_path": "template_searchpath", - "is_at_max_active_runs": "exceeds_max_non_backfill", **DAG_ALIAS_MAPPING, }.get(field_name, field_name), ), @@ -222,7 +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 + 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 " + "one running backfill run and no others can show active_runs_count: 1 alongside " + "is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and " + "otherwise no active runs can show active_runs_count: 0 alongside " + "is_at_max_active_runs: true." + ) + ) team_name: str | None = None @field_validator("timezone", mode="before") diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 538fdf736b6b5..403b365f1e216 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13441,6 +13441,14 @@ components: is_at_max_active_runs: type: boolean title: Is At Max Active Runs + 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 one running backfill + run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: + false, and a Dag with one queued (non-backfill) run and otherwise no active + runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.' team_name: anyOf: - type: string diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py index 3b1f275a8b5da..00ee0cb85a55b 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py @@ -274,9 +274,20 @@ def get_dag_details( 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) + + # Add is_favorite, active_runs_count, and is_at_max_active_runs fields to the Dag model setattr(dag_model, "is_favorite", is_favorite) setattr(dag_model, "active_runs_count", active_runs_count) + setattr(dag_model, "is_at_max_active_runs", is_at_max_active_runs) return DAGDetailsResponse.model_validate(dag_model) diff --git a/airflow-core/src/airflow/dag_processing/collection.py b/airflow-core/src/airflow/dag_processing/collection.py index c87375482b314..8c9c81b6fcfd7 100644 --- a/airflow-core/src/airflow/dag_processing/collection.py +++ b/airflow-core/src/airflow/dag_processing/collection.py @@ -167,24 +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, *, num_active_runs: int, session: Session) -> Self: + def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self: """ Query the latest-run info from the db. :param dag: the Dag to calculate run info for - :param num_active_runs: active run count for this Dag, pre-computed by the caller in one - batched query across every Dag in the current update so this doesn't become an N+1. :param session: DB session """ - # The latest-run lookup below is only used to calculate the *next* scheduled run, so it - # can be skipped for Dags that can never be scheduled in the first place. num_active_runs - # is still meaningful for such Dags (e.g. a schedule=None Dag triggered manually more - # often than its max_active_runs allows) and must always be computed. + # 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, num_active_runs) + return cls(None) if dag.timetable.partitioned: log.debug("Getting latest run for partitioned Dag", dag_id=dag.dag_id) @@ -201,7 +196,7 @@ def calculate(cls, dag: LazyDeserializedDAG, *, num_active_runs: int, session: S ) else: log.debug("no latest run found", dag_id=dag.dag_id) - return cls(latest_run, num_active_runs) + return cls(latest_run) def _update_dag_tags(tag_names: set[str], dm: DagModel, *, session: Session) -> None: @@ -631,11 +626,7 @@ def update_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], - num_active_runs=active_run_counts.get(dag_id, 0), - session=session, - ) + run_info = _RunInfo.calculate(dag=self.dags[dag_id], session=session) dag = self.dags[dag_id] dm.fileloc = dag.fileloc dm.relative_fileloc = dag.relative_fileloc @@ -701,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 = [] diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index 2c20b295fb10a..fd4266ba4061a 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3287,7 +3287,8 @@ export const $DAGDetailsResponse = { }, is_at_max_active_runs: { type: 'boolean', - title: 'Is At Max Active Runs' + title: 'Is At Max Active Runs', + 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true." }, team_name: { anyOf: [ diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index 4e0dafda75bc2..07df5073899c3 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -911,6 +911,9 @@ export type DAGDetailsResponse = { } | null; is_favorite?: boolean; active_runs_count?: number; + /** + * 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true. + */ is_at_max_active_runs: boolean; team_name?: string | null; /** diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index f4332331b70aa..80eb20413785e 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1514,16 +1514,26 @@ def test_dag_details_includes_active_runs_count(self, session, test_client): assert body["active_runs_count"] == 0 def test_dag_details_includes_is_at_max_active_runs(self, session, test_client): - """Test that DAG details include the is_at_max_active_runs field.""" + """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.exceeds_max_non_backfill = True + 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 "is_at_max_active_runs" in body assert body["is_at_max_active_runs"] is True # Test with a DAG that has not hit its max_active_runs @@ -1531,7 +1541,66 @@ def test_dag_details_includes_is_at_max_active_runs(self, session, test_client): assert response.status_code == 200 body = response.json() - assert "is_at_max_active_runs" in body + 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 + + dag_model = session.get(DagModel, DAG2_ID) + dag_model.max_active_runs = 1 + backfill = Backfill( + dag_id=DAG2_ID, + from_date=datetime(2021, 6, 15, tzinfo=timezone.utc), + to_date=datetime(2021, 6, 16, tzinfo=timezone.utc), + dag_run_conf=None, + ) + session.add(backfill) + session.flush() + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="is_at_max_active_runs_backfill_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.BACKFILL_JOB, + state=DagRunState.RUNNING, + backfill_id=backfill.id, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.commit() + + response = test_client.get(f"/dags/{DAG2_ID}/details") + assert response.status_code == 200 + body = response.json() + + # active_runs_count includes the backfill run; is_at_max_active_runs excludes it. + assert body["active_runs_count"] == 1 assert body["is_at_max_active_runs"] is False def test_dag_details_team_name_none_without_multi_team(self, test_client): diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 642091635f67c..3b47487e06173 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2926,7 +2926,13 @@ class DAGDetailsResponse(BaseModel): owner_links: Annotated[dict[str, str] | None, Field(title="Owner Links")] = None is_favorite: Annotated[bool | None, Field(title="Is Favorite")] = False active_runs_count: Annotated[int | None, Field(title="Active Runs Count")] = 0 - is_at_max_active_runs: Annotated[bool, Field(title="Is At Max Active Runs")] + is_at_max_active_runs: Annotated[ + 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.", + title="Is At Max Active Runs", + ), + ] team_name: Annotated[str | None, Field(title="Team Name")] = None is_backfillable: Annotated[ bool, Field(description="Whether this Dag's schedule supports backfilling.", title="Is Backfillable") From d3cb9fd42409e626cc916f989446da3b4a0596e2 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Wed, 30 Sep 2026 10:23:55 -0400 Subject: [PATCH 07/19] Shorten is_at_max_active_runs description and drop stale design notes The previous description wrongly called RUNNING + QUEUED the scheduler's promotion check (that gate counts RUNNING only). State what the field is instead: whether the Dag is currently at its limit, counting running and queued non-backfill runs. Test docstrings still described the old design, where the API aliased exceeds_max_non_backfill and so depended on the dag-processor keeping it accurate for schedule=None Dags. The backfill test no longer needs a Backfill row: active_runs_of_dags filters on run_type alone. Co-Authored-By: Claude --- .../airflow/api_fastapi/core_api/datamodels/dags.py | 10 ++-------- .../core_api/openapi/v2-rest-api-generated.yaml | 10 ++-------- .../airflow/ui/openapi-gen/requests/schemas.gen.ts | 2 +- .../src/airflow/ui/openapi-gen/requests/types.gen.ts | 2 +- .../api_fastapi/core_api/routes/public/test_dags.py | 11 ----------- .../tests/unit/dag_processing/test_collection.py | 9 +++++---- airflow-core/tests/unit/models/test_dag.py | 8 +------- .../src/airflowctl/api/datamodels/generated.py | 2 +- 8 files changed, 13 insertions(+), 41 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py index 6d670a9f769b3..d50dbf3966d3c 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py @@ -224,14 +224,8 @@ class DAGDetailsResponse(DAGResponse): 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 " - "active_runs_count counts RUNNING runs only and includes backfill runs. A Dag with " - "one running backfill run and no others can show active_runs_count: 1 alongside " - "is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and " - "otherwise no active runs can show active_runs_count: 0 alongside " - "is_at_max_active_runs: true." + "Whether this Dag is currently at its max_active_runs limit, counting running and " + "queued runs (backfill runs excluded)." ) ) team_name: str | None = None diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 403b365f1e216..0f93a56ac0f03 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13441,14 +13441,8 @@ components: is_at_max_active_runs: type: boolean title: Is At Max Active Runs - 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 one running backfill - run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: - false, and a Dag with one queued (non-backfill) run and otherwise no active - runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.' + 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 diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index fd4266ba4061a..b79e848e31de8 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3288,7 +3288,7 @@ export const $DAGDetailsResponse = { is_at_max_active_runs: { type: 'boolean', title: 'Is At Max Active Runs', - 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true." + description: 'Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).' }, team_name: { anyOf: [ diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index 07df5073899c3..0f926b82d2630 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -912,7 +912,7 @@ export type DAGDetailsResponse = { is_favorite?: boolean; active_runs_count?: number; /** - * 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true. + * 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; diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index 80eb20413785e..75fcc5a341f1d 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1569,18 +1569,8 @@ def test_dag_details_is_at_max_active_runs_counts_queued_runs_too(self, session, 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 - dag_model = session.get(DagModel, DAG2_ID) dag_model.max_active_runs = 1 - backfill = Backfill( - dag_id=DAG2_ID, - from_date=datetime(2021, 6, 15, tzinfo=timezone.utc), - to_date=datetime(2021, 6, 16, tzinfo=timezone.utc), - dag_run_conf=None, - ) - session.add(backfill) - session.flush() session.add( DagRun( dag_id=DAG2_ID, @@ -1589,7 +1579,6 @@ def test_dag_details_is_at_max_active_runs_excludes_backfill_runs(self, session, start_date=datetime(2021, 6, 15, 4, 0, 0, tzinfo=timezone.utc), run_type=DagRunType.BACKFILL_JOB, state=DagRunState.RUNNING, - backfill_id=backfill.id, triggered_by=DagRunTriggeredByType.TEST, ) ) diff --git a/airflow-core/tests/unit/dag_processing/test_collection.py b/airflow-core/tests/unit/dag_processing/test_collection.py index d1585bc9f9213..e2aba871c9a2a 100644 --- a/airflow-core/tests/unit/dag_processing/test_collection.py +++ b/airflow-core/tests/unit/dag_processing/test_collection.py @@ -1452,11 +1452,12 @@ def test_max_active_runs_defaults_from_conf_when_none(self, testing_dag_bundle, def test_exceeds_max_non_backfill_reflects_real_active_runs_for_non_schedulable_dag( self, testing_dag_bundle, session, dag_maker ): - """A schedule=None Dag's exceeds_max_non_backfill must reflect real active-run counts. + """A schedule=None Dag's exceeds_max_non_backfill reflects its real active-run count. - ``can_be_scheduled`` is False for schedule=None Dags, but that must not make the - parser hardcode num_active_runs to 0 -- the flag is read regardless of whether the - Dag can be automatically scheduled. + ``can_be_scheduled`` is False for schedule=None Dags; the active-run count comes from the + batched query in ``update_dags`` for every Dag either way. Nothing reads the flag for + such a Dag (it never gets ``next_dagrun_create_after`` set), so this only keeps the cached + column accurate rather than changing scheduling. """ with dag_maker("dag_schedule_none_exceeds_max", schedule=None, max_active_runs=1) as dag: ... diff --git a/airflow-core/tests/unit/models/test_dag.py b/airflow-core/tests/unit/models/test_dag.py index 2364541319a76..1af4ff7b26276 100644 --- a/airflow-core/tests/unit/models/test_dag.py +++ b/airflow-core/tests/unit/models/test_dag.py @@ -686,13 +686,7 @@ def test_bulk_write_to_db_multiple_dags(self, testing_dag_bundle): @pytest.mark.parametrize("interval", [None, "@daily"]) def test_bulk_write_to_db_interval_batches_active_run_count_query(self, testing_dag_bundle, interval): - """active_runs_of_dags is called exactly once per update, batched across every Dag. - - This holds regardless of whether any Dag in the batch has a schedule: the active-run - count is needed even for schedule=None Dags (to keep exceeds_max_non_backfill accurate - for Dags that are only ever triggered manually/via the API), so this must not be skipped - just because none of the Dags being updated can be scheduled. - """ + """active_runs_of_dags is called exactly once per update, batched across every Dag.""" mock_active_runs_of_dags = mock.MagicMock(side_effect=DagRun.active_runs_of_dags) with mock.patch.object(DagRun, "active_runs_of_dags", mock_active_runs_of_dags): dags_null_timetable = [ diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 3b47487e06173..1c70d5b6000f8 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2929,7 +2929,7 @@ class DAGDetailsResponse(BaseModel): is_at_max_active_runs: Annotated[ 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.", + description="Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).", title="Is At Max Active Runs", ), ] From 7cb55d489e803838a9958c5725a16cfbd3ff9645 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Thu, 24 Sep 2026 15:22:07 -0500 Subject: [PATCH 08/19] UI: Warn when a Dag's active runs exceed max_active_runs The Dag header's "Active Runs" stat shows "X of Y" once a Dag is at or over its max_active_runs limit, but nothing explains what that means or what happens next. A newly created run beyond the limit simply sits queued with no explanation visible in the UI, and "3 of 1" reads as a plain oddity rather than a Dag waiting on capacity. This adds a warning-triangle icon with a tooltip next to the stat label, shown only when the Dag has exceeded max_active_runs, explaining that additional runs will not start until an existing active run completes. Builds on #73692, which adds the exceeds_max_active_runs field this consumes on DAGDetailsResponse. Part of #73686. --- .../ui/public/i18n/locales/en/common.json | 1 + .../airflow/ui/src/components/HeaderCard.tsx | 6 ++--- .../airflow/ui/src/pages/Dag/Header.test.tsx | 23 +++++++++++++++++++ .../src/airflow/ui/src/pages/Dag/Header.tsx | 17 +++++++++++--- 4 files changed, 41 insertions(+), 6 deletions(-) diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json index fe8b8b93d0d60..34fda3629bf11 100644 --- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json +++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json @@ -44,6 +44,7 @@ "dagBundle_other": "Dag Bundles", "dagDetails": { "activeRuns": "Active Runs", + "activeRunsExceedsMaxTooltip": "This Dag has more active runs than its Max Active Runs limit allows. Runs beyond the limit will not start until an existing active run completes.", "catchup": "Catchup", "dagRunTimeout": "Dag Run Timeout", "defaultArgs": "Default Args", diff --git a/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx b/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx index 5a6dce6312a23..418ba50142408 100644 --- a/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx +++ b/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx @@ -31,7 +31,7 @@ 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<{ key?: string; label: ReactNode | string; value: ReactNode | string }>; readonly subTitle?: ReactNode | string; readonly title: ReactNode | string; readonly type: "asset" | "dag" | "dagBundle" | "dagRun" | "task" | "taskGroup" | "taskInstance"; @@ -77,8 +77,8 @@ export const HeaderCard = ({ actions, icon, state, stats, subTitle, title, type - {stats.map((stat) => ( - + {stats.map((stat, index) => ( + { expect(screen.getByText("2 of 2")).toBeInTheDocument(); }); + it("does not show a warning icon when active runs are within the maximum", () => { + render( + +
+ , + ); + + expect(screen.queryByTestId("active-runs-exceeds-max-warning")).not.toBeInTheDocument(); + }); + + it("shows a warning icon when active runs exceed the maximum", () => { + render( + +
+ , + ); + + expect(screen.getByTestId("active-runs-exceeds-max-warning")).toBeInTheDocument(); + }); + it("renders the draining badge instead of the next run timestamp for a draining Dag", () => { render( diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx index 7687488fcfc19..1dde412141fe8 100644 --- a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx +++ b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx @@ -16,13 +16,14 @@ * specific language governing permissions and limitations * under the License. */ +import { HStack } from "@chakra-ui/react"; import { useTranslation } from "react-i18next"; -import { FiBookOpen } from "react-icons/fi"; +import { FiAlertTriangle, FiBookOpen } from "react-icons/fi"; import { useParams } from "react-router-dom"; import type { DAGDetailsResponse, DagRunState } from "openapi/requests/types.gen"; -import { RouterLink } from "src/system-components"; +import { RouterLink, Tooltip } from "src/system-components"; import { DeleteDagButton } from "src/components/DagActions/DeleteDagButton"; import { FavoriteDagButton } from "src/components/DagActions/FavoriteDagButton"; @@ -113,7 +114,17 @@ export const Header = ({ }, ...nextRunStat, { - label: translate("dagDetails.activeRuns"), + label: + dag?.exceeds_max_active_runs === true ? ( + + {translate("dagDetails.activeRuns")} + + + + + ) : ( + translate("dagDetails.activeRuns") + ), value: dag?.max_active_runs === undefined ? undefined From 237bba961dc567d1799bd793895bd95aac41d05a Mon Sep 17 00:00:00 2001 From: seanmuth Date: Thu, 24 Sep 2026 16:22:38 -0500 Subject: [PATCH 09/19] Show queued run count and switch to an info icon on the Dag header Two refinements based on testing this against a live reproduction: - The number displayed can never actually show more active runs than the limit allows (RUNNING is capped by the scheduler's promotion gate), so a warning-severity icon overstated the situation. Switch to a plain info icon. - Key the tooltip off a new, live-computed queued_runs_count field instead of exceeds_max_active_runs. The flag answers "is this Dag at or over capacity" (useful on its own, via #73692), which is a slightly different question from "are there runs actually waiting right now" -- the latter is what the UI needs, and computing it fresh on every request sidesteps any staleness in the persisted flag entirely. Also display the queued count directly ("1 of 1 (2 queued)"), so the information doesn't require a hover at all. --- .../api_fastapi/core_api/datamodels/dags.py | 1 + .../core_api/openapi/v2-rest-api-generated.yaml | 12 +++++++----- .../api_fastapi/core_api/routes/public/dags.py | 13 ++++++++++++- .../ui/openapi-gen/requests/schemas.gen.ts | 12 ++++++++---- .../airflow/ui/openapi-gen/requests/types.gen.ts | 6 ++---- .../ui/public/i18n/locales/en/common.json | 1 + .../src/airflow/ui/src/pages/Dag/Header.test.tsx | 16 ++++++++-------- .../src/airflow/ui/src/pages/Dag/Header.tsx | 14 ++++++++++---- .../core_api/routes/public/test_dags.py | 12 +++++++++++- .../src/airflowctl/api/datamodels/generated.py | 9 ++------- 10 files changed, 62 insertions(+), 34 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py index d50dbf3966d3c..7e468a0b5c1d9 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py @@ -222,6 +222,7 @@ 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 is currently at its max_active_runs limit, counting running and " diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 0f93a56ac0f03..ed9b97887b60c 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13438,11 +13438,13 @@ components: type: integer title: Active Runs Count default: 0 - is_at_max_active_runs: + queued_runs_count: + type: integer + title: Queued Runs Count + default: 0 + exceeds_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). + title: Exceeds Max Active Runs team_name: anyOf: - type: string @@ -13516,7 +13518,7 @@ components: - timezone - last_parsed - default_args - - is_at_max_active_runs + - exceeds_max_active_runs - is_backfillable - file_token - concurrency diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py index 00ee0cb85a55b..143353eba1719 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py @@ -274,6 +274,16 @@ def get_dag_details( or 0 ) + # Count queued Dag runs: these are waiting for an active run to finish before they can start. + queued_runs_count = ( + session.scalar( + select(func.count()) + .select_from(DagRun) + .where(DagRun.dag_id == dag_id, DagRun.state == DagRunState.QUEUED) + ) + or 0 + ) + # 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 @@ -284,9 +294,10 @@ def get_dag_details( ).get(dag_id, 0) is_at_max_active_runs = non_backfill_active_runs_count >= (dag_model.max_active_runs or 0) - # Add is_favorite, active_runs_count, and is_at_max_active_runs 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) diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index b79e848e31de8..6d2792e9e2890 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3285,10 +3285,14 @@ export const $DAGDetailsResponse = { title: 'Active Runs Count', default: 0 }, - is_at_max_active_runs: { + queued_runs_count: { + type: 'integer', + title: 'Queued Runs Count', + default: 0 + }, + exceeds_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).' + title: 'Exceeds Max Active Runs' }, team_name: { anyOf: [ @@ -3336,7 +3340,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_at_max_active_runs', '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', 'exceeds_max_active_runs', 'is_backfillable', 'file_token', 'concurrency', 'latest_dag_version'], title: 'DAGDetailsResponse', description: 'Specific serializer for Dag Details responses.' } as const; diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index 0f926b82d2630..5c96ff7b0077a 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -911,10 +911,8 @@ export type DAGDetailsResponse = { } | null; is_favorite?: boolean; active_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; + queued_runs_count?: number; + exceeds_max_active_runs: boolean; team_name?: string | null; /** * Whether this Dag's schedule supports backfilling. diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json index 34fda3629bf11..a2c3dc2ff6e22 100644 --- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json +++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json @@ -45,6 +45,7 @@ "dagDetails": { "activeRuns": "Active Runs", "activeRunsExceedsMaxTooltip": "This Dag has more active runs than its Max Active Runs limit allows. Runs beyond the limit will not start until an existing active run completes.", + "activeRunsWithQueued": "{{activeRuns}} of {{maxActiveRuns}} ({{queuedRuns}} queued)", "catchup": "Catchup", "dagRunTimeout": "Dag Run Timeout", "defaultArgs": "Default Args", diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx b/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx index 8a53c0451dc84..71d5457128f3e 100644 --- a/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx +++ b/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx @@ -40,7 +40,6 @@ const mockDag = { bundle_name: "dags-folder", bundle_version: "1", default_args: {}, - exceeds_max_active_runs: false, fileloc: "/files/dags/stale_dag.py", is_favorite: false, is_stale: true, @@ -50,6 +49,7 @@ 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, @@ -94,26 +94,26 @@ describe("Header", () => { expect(screen.getByText("2 of 2")).toBeInTheDocument(); }); - it("does not show a warning icon when active runs are within the maximum", () => { + it("does not show an info icon or queued count when nothing is queued", () => { render(
, ); - expect(screen.queryByTestId("active-runs-exceeds-max-warning")).not.toBeInTheDocument(); + expect(screen.queryByTestId("active-runs-exceeds-max-info")).not.toBeInTheDocument(); + expect(screen.getByText("1 of 2")).toBeInTheDocument(); }); - it("shows a warning icon when active runs exceed the maximum", () => { + it("shows an info icon and the queued count when runs are queued behind the maximum", () => { render( -
+
, ); - expect(screen.getByTestId("active-runs-exceeds-max-warning")).toBeInTheDocument(); + expect(screen.getByTestId("active-runs-exceeds-max-info")).toBeInTheDocument(); + expect(screen.getByText("1 of 1 (2 queued)")).toBeInTheDocument(); }); it("renders the draining badge instead of the next run timestamp for a draining Dag", () => { diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx index 1dde412141fe8..2a7c6a7372071 100644 --- a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx +++ b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx @@ -18,7 +18,7 @@ */ import { HStack } from "@chakra-ui/react"; import { useTranslation } from "react-i18next"; -import { FiAlertTriangle, FiBookOpen } from "react-icons/fi"; +import { FiBookOpen, FiInfo } from "react-icons/fi"; import { useParams } from "react-router-dom"; import type { DAGDetailsResponse, DagRunState } from "openapi/requests/types.gen"; @@ -115,11 +115,11 @@ export const Header = ({ ...nextRunStat, { label: - dag?.exceeds_max_active_runs === true ? ( + (dag?.queued_runs_count ?? 0) > 0 ? ( {translate("dagDetails.activeRuns")} - + ) : ( @@ -128,7 +128,13 @@ export const Header = ({ value: dag?.max_active_runs === undefined ? undefined - : `${dag.active_runs_count ?? 0} of ${dag.max_active_runs}`, + : (dag.queued_runs_count ?? 0) > 0 + ? 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}`, }, { label: translate("dagDetails.owner"), diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index 75fcc5a341f1d..f2be265fa4ab1 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1395,6 +1395,7 @@ def test_dag_details( "value": 1, } }, + "queued_runs_count": 0, "relative_fileloc": "test_dags.py", "render_template_as_native_obj": False, "rerun_with_latest_version": None, @@ -1457,7 +1458,7 @@ def test_dag_details_serves_legacy_asset_expression_as_null(self, session, test_ assert response.json()["asset_expression"] is None def test_dag_details_includes_active_runs_count(self, session, test_client): - """Test that DAG details include the active_runs_count field.""" + """Test that DAG details include the active_runs_count and queued_runs_count fields.""" # Create running and queued DAG runs for DAG2 session.add( DagRun( @@ -1504,6 +1505,11 @@ def test_dag_details_includes_active_runs_count(self, session, test_client): assert isinstance(body["active_runs_count"], int) assert body["active_runs_count"] == 1 # only running counts, queued does not + # Verify queued_runs_count field is present and correct + assert "queued_runs_count" in body + assert isinstance(body["queued_runs_count"], int) + assert body["queued_runs_count"] == 1 # 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 @@ -1513,6 +1519,10 @@ 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 + assert "queued_runs_count" in body + assert isinstance(body["queued_runs_count"], int) + assert body["queued_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) diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 1c70d5b6000f8..7abb27e721c65 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2926,13 +2926,8 @@ class DAGDetailsResponse(BaseModel): owner_links: Annotated[dict[str, str] | None, Field(title="Owner Links")] = None is_favorite: Annotated[bool | None, Field(title="Is Favorite")] = False active_runs_count: Annotated[int | None, Field(title="Active Runs Count")] = 0 - is_at_max_active_runs: Annotated[ - bool, - Field( - description="Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).", - title="Is At Max Active Runs", - ), - ] + queued_runs_count: Annotated[int | None, Field(title="Queued Runs Count")] = 0 + exceeds_max_active_runs: Annotated[bool, Field(title="Exceeds Max Active Runs")] team_name: Annotated[str | None, Field(title="Team Name")] = None is_backfillable: Annotated[ bool, Field(description="Whether this Dag's schedule supports backfilling.", title="Is Backfillable") From 48d2efde6ff3339baa4036c292d3890083b4c2ed Mon Sep 17 00:00:00 2001 From: seanmuth Date: Thu, 24 Sep 2026 16:28:58 -0500 Subject: [PATCH 10/19] UI: Trim the queued-runs tooltip wording Drop the opening sentence -- now that the tooltip is keyed off queued_runs_count rather than an "active vs max" comparison, "more active runs than its limit allows" no longer accurately describes the condition, and the second sentence already says what matters. --- airflow-core/src/airflow/ui/public/i18n/locales/en/common.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json index a2c3dc2ff6e22..598bdd58d0547 100644 --- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json +++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json @@ -44,7 +44,7 @@ "dagBundle_other": "Dag Bundles", "dagDetails": { "activeRuns": "Active Runs", - "activeRunsExceedsMaxTooltip": "This Dag has more active runs than its Max Active Runs limit allows. Runs beyond the limit will not start until an existing active run completes.", + "activeRunsExceedsMaxTooltip": "Runs beyond the limit will not start until an existing active run completes.", "activeRunsWithQueued": "{{activeRuns}} of {{maxActiveRuns}} ({{queuedRuns}} queued)", "catchup": "Catchup", "dagRunTimeout": "Dag Run Timeout", From 9773a8ef048feee35e7fefa66df7cd2bee33042b Mon Sep 17 00:00:00 2001 From: seanmuth Date: Fri, 25 Sep 2026 11:08:48 -0500 Subject: [PATCH 11/19] Regenerate OpenAPI spec and UI client after rebase The rebase conflict resolution took a placeholder version of these generated files; regenerate them fresh so they reflect the is_at_max_active_runs rename and queued_runs_count addition. Co-Authored-By: Claude --- .../api_fastapi/core_api/openapi/v2-rest-api-generated.yaml | 6 +++--- .../src/airflow/ui/openapi-gen/requests/schemas.gen.ts | 6 +++--- .../src/airflow/ui/openapi-gen/requests/types.gen.ts | 2 +- 3 files changed, 7 insertions(+), 7 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index ed9b97887b60c..8567e6b21a466 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13442,9 +13442,9 @@ components: type: integer title: Queued Runs Count default: 0 - exceeds_max_active_runs: + is_at_max_active_runs: type: boolean - title: Exceeds Max Active Runs + title: Is At Max Active Runs team_name: anyOf: - type: string @@ -13518,7 +13518,7 @@ components: - timezone - last_parsed - default_args - - exceeds_max_active_runs + - is_at_max_active_runs - is_backfillable - file_token - concurrency diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index 6d2792e9e2890..189ddb4c6d2a6 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3290,9 +3290,9 @@ export const $DAGDetailsResponse = { title: 'Queued Runs Count', default: 0 }, - exceeds_max_active_runs: { + is_at_max_active_runs: { type: 'boolean', - title: 'Exceeds Max Active Runs' + title: 'Is At Max Active Runs' }, team_name: { anyOf: [ @@ -3340,7 +3340,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', 'exceeds_max_active_runs', '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; diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index 5c96ff7b0077a..3aee03220f2da 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -912,7 +912,7 @@ export type DAGDetailsResponse = { is_favorite?: boolean; active_runs_count?: number; queued_runs_count?: number; - exceeds_max_active_runs: boolean; + is_at_max_active_runs: boolean; team_name?: string | null; /** * Whether this Dag's schedule supports backfilling. From 11e088fa9cba90adbb31b411ccb284962df4073b Mon Sep 17 00:00:00 2001 From: seanmuth Date: Fri, 25 Sep 2026 11:09:33 -0500 Subject: [PATCH 12/19] Regenerate airflowctl datamodels after rebase Same placeholder-conflict-resolution issue as the previous commit -- this file still had exceeds_max_active_runs from before the rename. Co-Authored-By: Claude --- airflow-ctl/src/airflowctl/api/datamodels/generated.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 7abb27e721c65..3e155849b1b65 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2927,7 +2927,7 @@ class DAGDetailsResponse(BaseModel): is_favorite: Annotated[bool | None, Field(title="Is Favorite")] = False active_runs_count: Annotated[int | None, Field(title="Active Runs Count")] = 0 queued_runs_count: Annotated[int | None, Field(title="Queued Runs Count")] = 0 - exceeds_max_active_runs: Annotated[bool, Field(title="Exceeds Max Active Runs")] + is_at_max_active_runs: Annotated[bool, Field(title="Is At Max Active Runs")] team_name: Annotated[str | None, Field(title="Team Name")] = None is_backfillable: Annotated[ bool, Field(description="Whether this Dag's schedule supports backfilling.", title="Is Backfillable") From 16f157083e1442370a21e350862178062470d084 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Fri, 25 Sep 2026 12:01:02 -0500 Subject: [PATCH 13/19] Fix Header.test.tsx assertion for the interpolated queued-count string This test suite's i18n setup doesn't resolve real translations -- every translate() call renders its raw key, as every other assertion in this file already accounts for by comparing against i18n.t(...) rather than a hardcoded final string. The new queued-count assertion missed that pattern and compared against the literal English output, which never matches in CI. Co-Authored-By: Claude --- airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx b/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx index 71d5457128f3e..0623cd047d854 100644 --- a/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx +++ b/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx @@ -113,7 +113,11 @@ describe("Header", () => { ); expect(screen.getByTestId("active-runs-exceeds-max-info")).toBeInTheDocument(); - expect(screen.getByText("1 of 1 (2 queued)")).toBeInTheDocument(); + expect( + screen.getByText( + i18n.t("common:dagDetails.activeRunsWithQueued", { activeRuns: 1, maxActiveRuns: 1, queuedRuns: 2 }), + ), + ).toBeInTheDocument(); }); it("renders the draining badge instead of the next run timestamp for a draining Dag", () => { From cc82d492ddd04b9e491c7b022d1a075fa60f16a9 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Tue, 29 Sep 2026 09:45:52 -0400 Subject: [PATCH 14/19] Exclude backfill runs from active_runs_count and queued_runs_count Both counts were counting every run regardless of backfill_id, but a backfill's own runs are gated by Backfill.max_active_runs, not the Dag's. A 100-date backfill on a max_active_runs=16 Dag would render "10 of 16 (90 queued)" pointing at the wrong limit. Filters both queries on DagRun.backfill_id.is_(None) to match how the promotion queries (get_queued_dag_runs_to_set_running, _start_queued_dagruns) already scope concurrency per (dag_id, backfill_id). Also strengthened the existing count test to seed two queued runs instead of one, so it can't pass by coincidence if queued_runs_count accidentally queried RUNNING instead of QUEUED. Co-Authored-By: Claude --- .../core_api/routes/public/dags.py | 15 ++- .../core_api/routes/public/test_dags.py | 92 +++++++++++++++---- 2 files changed, 89 insertions(+), 18 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py index 143353eba1719..0f216686fc48e 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py @@ -265,21 +265,32 @@ 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) + .where( + DagRun.dag_id == dag_id, + DagRun.state == DagRunState.QUEUED, + DagRun.backfill_id.is_(None), + ) ) or 0 ) diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index f2be265fa4ab1..7a5d09ab0ede1 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1457,9 +1457,11 @@ def test_dag_details_serves_legacy_asset_expression_as_null(self, session, test_ assert response.status_code == 200 assert response.json()["asset_expression"] is None - def test_dag_details_includes_active_runs_count(self, session, test_client): + def test_dag_details_includes_active_runs_count_and_queued_runs_count(self, session, test_client): """Test that DAG details include the active_runs_count and queued_runs_count fields.""" - # Create running and queued DAG runs for DAG2 + # One running, two queued, and one successful (uncounted) run for DAG2 -- two queued + # runs (not one) so a query that accidentally filters on RUNNING instead of QUEUED for + # queued_runs_count can't coincidentally match active_runs_count's value of 1. session.add( DagRun( dag_id=DAG2_ID, @@ -1482,6 +1484,17 @@ def test_dag_details_includes_active_runs_count(self, session, test_client): triggered_by=DagRunTriggeredByType.TEST, ) ) + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="queued_run_2", + logical_date=datetime(2021, 6, 15, 2, 30, 0, tzinfo=timezone.utc), + start_date=datetime(2021, 6, 15, 2, 30, 0, tzinfo=timezone.utc), + run_type=DagRunType.MANUAL, + state=DagRunState.QUEUED, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) # Add a successful DAG run (should not be counted) session.add( DagRun( @@ -1500,29 +1513,76 @@ 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 - - # Verify queued_runs_count field is present and correct - assert "queued_runs_count" in body - assert isinstance(body["queued_runs_count"], int) - assert body["queued_runs_count"] == 1 # only queued counts, running/success do 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 "queued_runs_count" in body - assert isinstance(body["queued_runs_count"], int) 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 + + backfill = Backfill( + dag_id=DAG2_ID, + from_date=datetime(2021, 6, 15, tzinfo=timezone.utc), + to_date=datetime(2021, 6, 16, tzinfo=timezone.utc), + dag_run_conf=None, + ) + session.add(backfill) + session.flush() + + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="manual_running_run", + logical_date=datetime(2021, 6, 15, 1, 0, 0, tzinfo=timezone.utc), + start_date=datetime(2021, 6, 15, 1, 0, 0, tzinfo=timezone.utc), + run_type=DagRunType.MANUAL, + state=DagRunState.RUNNING, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="backfill_running_run", + logical_date=datetime(2021, 6, 15, 2, 0, 0, tzinfo=timezone.utc), + start_date=datetime(2021, 6, 15, 2, 0, 0, tzinfo=timezone.utc), + run_type=DagRunType.BACKFILL_JOB, + state=DagRunState.RUNNING, + backfill_id=backfill.id, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="backfill_queued_run", + logical_date=datetime(2021, 6, 15, 3, 0, 0, tzinfo=timezone.utc), + start_date=datetime(2021, 6, 15, 3, 0, 0, tzinfo=timezone.utc), + run_type=DagRunType.BACKFILL_JOB, + state=DagRunState.QUEUED, + backfill_id=backfill.id, + 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"] == 1 # the backfill's running run is excluded + assert body["queued_runs_count"] == 0 # the backfill's queued run is excluded + 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) From c54f61c9b5026845bef73a860acaf8120852de1a Mon Sep 17 00:00:00 2001 From: seanmuth Date: Tue, 29 Sep 2026 09:50:31 -0400 Subject: [PATCH 15/19] Fix the queued-runs icon's false positive, portal, and key hack The info icon showed whenever queued_runs_count > 0, but the tooltip claims a run is "waiting for an active run to complete" -- not true for a paused Dag (queued runs never promote at all while paused) or for the brief window where the scheduler just hasn't picked up a queued run yet on a busy deployment. Now only shows the icon when active_runs_count is actually at or over max_active_runs and the Dag isn't paused; the "(N queued)" text still always shows whenever queued_runs_count > 0, regardless of the icon. Added `portalled` to the Tooltip: without it, the tooltip's content rendered inline inside HeaderCard's stat-label Box, which sets textTransform="uppercase" -- so the tooltip text appeared in caps, overlapping the stat's value. Reverted HeaderCard's key={stat.key ?? index} back to key={stat.key ?? stat.label}: the array-index fallback slipped past react/no-array-index-key (which doesn't inspect ?? expressions) and silently switched every other HeaderCard caller from a label-keyed to an index-keyed list. The active-runs stat now passes key: "activeRuns" explicitly instead, and the stats prop is typed as a discriminated union so a non-string label always requires an explicit key going forward. The key expression itself narrows label with typeof, since label's type still includes non-Key ReactNode values like false. Rewrote Header.test.tsx's i18n-dependent assertions to load the real en/common locale bundle (matching RenderedJsonField.test.tsx) instead of comparing against i18n.t(...) with no bundle loaded, which just compared the raw translation key to itself and would pass regardless of the actual interpolated values. Also fixed two pre-existing assertions that were querying the "dag" namespace for a key (dagDetails.nextRun) that actually lives in "common" -- previously invisible because both sides always fell back to the same raw key. Added coverage for the corrected icon condition: exactly at capacity with nothing queued, below capacity with something queued, and a paused Dag at capacity with something queued -- all cases where the icon must stay hidden even though it previously would have shown (or, for the first case, was already covered). Co-Authored-By: Claude --- .../airflow/ui/src/components/HeaderCard.tsx | 15 ++++- .../airflow/ui/src/pages/Dag/Header.test.tsx | 60 ++++++++++++++++--- .../src/airflow/ui/src/pages/Dag/Header.tsx | 14 ++++- 3 files changed, 74 insertions(+), 15 deletions(-) diff --git a/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx b/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx index 418ba50142408..05f38f488c228 100644 --- a/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx +++ b/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx @@ -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: ReactNode | string; value: ReactNode | string }>; + readonly stats: Array; readonly subTitle?: ReactNode | string; readonly title: ReactNode | string; readonly type: "asset" | "dag" | "dagBundle" | "dagRun" | "task" | "taskGroup" | "taskInstance"; @@ -77,8 +83,11 @@ export const HeaderCard = ({ actions, icon, state, stats, subTitle, title, type - {stats.map((stat, index) => ( - + {stats.map((stat) => ( + = { multi_team: false }; @@ -57,6 +58,10 @@ const mockDag = { } as unknown as DAGDetailsResponse; describe("Header", () => { + beforeAll(() => { + i18n.addResourceBundle("en", "common", commonLocale, true, true); + }); + afterEach(() => { mockConfig.multi_team = false; }); @@ -68,7 +73,7 @@ describe("Header", () => { , ); - 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(); }); @@ -79,7 +84,7 @@ describe("Header", () => { , ); - 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(); }); @@ -105,6 +110,17 @@ describe("Header", () => { expect(screen.getByText("1 of 2")).toBeInTheDocument(); }); + it("does not show the icon when exactly at capacity with nothing queued", () => { + render( + +
+ , + ); + + 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( @@ -113,11 +129,37 @@ describe("Header", () => { ); expect(screen.getByTestId("active-runs-exceeds-max-info")).toBeInTheDocument(); - expect( - screen.getByText( - i18n.t("common:dagDetails.activeRunsWithQueued", { activeRuns: 1, maxActiveRuns: 1, queuedRuns: 2 }), - ), - ).toBeInTheDocument(); + expect(screen.getByText("1 of 1 (2 queued)")).toBeInTheDocument(); + }); + + it("shows the queued count without the icon when below capacity (scheduler hasn't caught up yet)", () => { + render( + +
+ , + ); + + 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( + +
+ , + ); + + 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", () => { @@ -127,7 +169,7 @@ describe("Header", () => { , ); - 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(); }); diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx index 2a7c6a7372071..44d683b7d9122 100644 --- a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx +++ b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx @@ -66,6 +66,13 @@ export const Header = ({ const { dagId } = useParams(); const showTeam = useShowTeam(dag?.team_name); const isStale = dag?.is_stale; + const hasQueuedRuns = (dag?.queued_runs_count ?? 0) > 0; + // Queued runs alone don't mean anything is stuck: on a busy deployment the scheduler may + // simply not have picked them up yet, and a paused Dag never promotes queued runs at all. + // Only warn once active runs are actually at (or over) the limit. + const isBlockedByMaxActiveRuns = + (dag?.active_runs_count ?? 0) >= (dag?.max_active_runs ?? Number.POSITIVE_INFINITY) && + dag?.is_paused !== true; const nextRunStat = isStale ? [] @@ -114,11 +121,12 @@ export const Header = ({ }, ...nextRunStat, { + key: "activeRuns", label: - (dag?.queued_runs_count ?? 0) > 0 ? ( + isBlockedByMaxActiveRuns && hasQueuedRuns ? ( {translate("dagDetails.activeRuns")} - + @@ -128,7 +136,7 @@ export const Header = ({ value: dag?.max_active_runs === undefined ? undefined - : (dag.queued_runs_count ?? 0) > 0 + : hasQueuedRuns ? translate("dagDetails.activeRunsWithQueued", { activeRuns: dag.active_runs_count ?? 0, maxActiveRuns: dag.max_active_runs, From 1906b01808156bfd90890d074d908e78658c63c2 Mon Sep 17 00:00:00 2001 From: seanmuth Date: Tue, 29 Sep 2026 12:29:14 -0400 Subject: [PATCH 16/19] Regenerate OpenAPI spec and UI client after rebase The rebase conflict resolution took a placeholder version of these generated files; regenerate them fresh so they reflect the is_at_max_active_runs description added on #73692. Co-Authored-By: Claude --- .../core_api/openapi/v2-rest-api-generated.yaml | 8 ++++++++ .../src/airflow/ui/openapi-gen/requests/schemas.gen.ts | 3 ++- .../src/airflow/ui/openapi-gen/requests/types.gen.ts | 3 +++ 3 files changed, 13 insertions(+), 1 deletion(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 8567e6b21a466..3aa0e8c6c971e 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13445,6 +13445,14 @@ components: is_at_max_active_runs: type: boolean title: Is At Max Active Runs + 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 one running backfill + run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: + false, and a Dag with one queued (non-backfill) run and otherwise no active + runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.' team_name: anyOf: - type: string diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index 189ddb4c6d2a6..b474f28917164 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3292,7 +3292,8 @@ export const $DAGDetailsResponse = { }, is_at_max_active_runs: { type: 'boolean', - title: 'Is At Max Active Runs' + title: 'Is At Max Active Runs', + 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true." }, team_name: { anyOf: [ diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index 3aee03220f2da..dfcc2a420e645 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -912,6 +912,9 @@ export type DAGDetailsResponse = { is_favorite?: boolean; active_runs_count?: number; queued_runs_count?: number; + /** + * 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true. + */ is_at_max_active_runs: boolean; team_name?: string | null; /** From b578b391cf15a7b604495f679d54ba4c225777fc Mon Sep 17 00:00:00 2001 From: seanmuth Date: Tue, 29 Sep 2026 12:29:53 -0400 Subject: [PATCH 17/19] Regenerate airflowctl datamodels after rebase Same placeholder-conflict-resolution issue as the previous commit -- this file was still missing is_at_max_active_runs' description field. Co-Authored-By: Claude --- airflow-ctl/src/airflowctl/api/datamodels/generated.py | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 3e155849b1b65..fb8cd54f56ffe 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2927,7 +2927,13 @@ class DAGDetailsResponse(BaseModel): is_favorite: Annotated[bool | None, Field(title="Is Favorite")] = False active_runs_count: Annotated[int | None, Field(title="Active Runs Count")] = 0 queued_runs_count: Annotated[int | None, Field(title="Queued Runs Count")] = 0 - is_at_max_active_runs: Annotated[bool, Field(title="Is At Max Active Runs")] + is_at_max_active_runs: Annotated[ + 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.", + title="Is At Max Active Runs", + ), + ] team_name: Annotated[str | None, Field(title="Team Name")] = None is_backfillable: Annotated[ bool, Field(description="Whether this Dag's schedule supports backfilling.", title="Is Backfillable") From b74476a817f830ac1f2ad56c3dbd5e90356676de Mon Sep 17 00:00:00 2001 From: seanmuth Date: Tue, 29 Sep 2026 12:34:56 -0400 Subject: [PATCH 18/19] Fix backfill-exclusion test for this branch's own active_runs_count fix The test inherited from #73692's branch expected active_runs_count to still include the backfill run (true there, since that branch doesn't have this branch's own backfill-exclusion fix for active_runs_count). On this branch both fields exclude it, so rewrote the test to add a manual running run alongside the backfill one and assert the backfill run doesn't push either count from 1 to 2 -- falsifiable regardless of which of the two fields' exclusion logic might regress. Co-Authored-By: Claude --- .../core_api/routes/public/test_dags.py | 18 +++++++++++++++--- 1 file changed, 15 insertions(+), 3 deletions(-) diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index 7a5d09ab0ede1..3cc4c3fb90141 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1638,9 +1638,20 @@ def test_dag_details_is_at_max_active_runs_counts_queued_runs_too(self, session, 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.""" + """A running backfill run doesn't count toward the Dag's own active-run stats.""" dag_model = session.get(DagModel, DAG2_ID) - dag_model.max_active_runs = 1 + dag_model.max_active_runs = 2 + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="is_at_max_active_runs_manual_running", + logical_date=datetime(2021, 6, 15, 3, 0, 0, tzinfo=timezone.utc), + start_date=datetime(2021, 6, 15, 3, 0, 0, tzinfo=timezone.utc), + run_type=DagRunType.MANUAL, + state=DagRunState.RUNNING, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) session.add( DagRun( dag_id=DAG2_ID, @@ -1658,7 +1669,8 @@ def test_dag_details_is_at_max_active_runs_excludes_backfill_runs(self, session, assert response.status_code == 200 body = response.json() - # active_runs_count includes the backfill run; is_at_max_active_runs excludes it. + # The manual run counts; the backfill run on top of it doesn't push either field to 2 + # (which would make is_at_max_active_runs true, since max_active_runs is 2). assert body["active_runs_count"] == 1 assert body["is_at_max_active_runs"] is False From f49dbc90f7ad73b156ccfcaf0f9106ab1b92f52b Mon Sep 17 00:00:00 2001 From: seanmuth Date: Wed, 30 Sep 2026 10:31:05 -0400 Subject: [PATCH 19/19] Count active and queued runs in one query, describe both fields active_runs_count, queued_runs_count and is_at_max_active_runs came from three separate SELECTs with two different backfill filters, so they could disagree if a run changed state between them. One query grouped by state now gives both counts, and is_at_max_active_runs is their sum against the limit. active_runs_count has meant something different in each of the last few releases, so both counts now carry a description of what they count. The non-queued "X of Y" stat now goes through i18n like the queued one, so a translated UI doesn't switch language as the queue drains. Added an over-capacity test so the icon condition's >= is pinned, and dropped the Backfill row from the backfill test (the filter is on run_type, so the row isn't needed and it outlived the test). Co-Authored-By: Claude --- .../api_fastapi/core_api/datamodels/dags.py | 8 ++- .../openapi/v2-rest-api-generated.yaml | 12 ++-- .../core_api/routes/public/dags.py | 58 +++++++------------ .../ui/openapi-gen/requests/schemas.gen.ts | 4 +- .../ui/openapi-gen/requests/types.gen.ts | 8 ++- .../ui/public/i18n/locales/en/common.json | 1 + .../airflow/ui/src/pages/Dag/Header.test.tsx | 11 ++++ .../src/airflow/ui/src/pages/Dag/Header.tsx | 5 +- .../core_api/routes/public/test_dags.py | 13 ----- .../airflowctl/api/datamodels/generated.py | 17 +++++- 10 files changed, 70 insertions(+), 67 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py index 7e468a0b5c1d9..7d153e78092a4 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py @@ -221,8 +221,12 @@ 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 - queued_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 " diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 3aa0e8c6c971e..f32cdb89ccada 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -13437,22 +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 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 one running backfill - run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: - false, and a Dag with one queued (non-backfill) run and otherwise no active - runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.' + 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 diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py index 0f216686fc48e..7e2eb89d82522 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py @@ -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") @@ -264,46 +265,27 @@ def get_dag_details( is not None ) - # 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, - 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), - ) + # 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 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) + # 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, active_runs_count, queued_runs_count, and is_at_max_active_runs fields setattr(dag_model, "is_favorite", is_favorite) diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index b474f28917164..51cba31ecd060 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -3283,17 +3283,19 @@ 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: { 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 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true." + description: 'Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).' }, team_name: { anyOf: [ diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index dfcc2a420e645..0be3dafe38933 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -910,10 +910,16 @@ 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 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true. + * 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; diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json index 598bdd58d0547..fc10152aa1ea5 100644 --- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json +++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json @@ -45,6 +45,7 @@ "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", diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx b/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx index 0801a550d0e40..014c6fa49d827 100644 --- a/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx +++ b/airflow-core/src/airflow/ui/src/pages/Dag/Header.test.tsx @@ -132,6 +132,17 @@ describe("Header", () => { 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( + +
+ , + ); + + 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( diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx index 44d683b7d9122..1377de0161a1d 100644 --- a/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx +++ b/airflow-core/src/airflow/ui/src/pages/Dag/Header.tsx @@ -142,7 +142,10 @@ export const Header = ({ maxActiveRuns: dag.max_active_runs, queuedRuns: dag.queued_runs_count, }) - : `${dag.active_runs_count ?? 0} of ${dag.max_active_runs}`, + : translate("dagDetails.activeRunsOfMax", { + activeRuns: dag.active_runs_count ?? 0, + maxActiveRuns: dag.max_active_runs, + }), }, { label: translate("dagDetails.owner"), diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index 3cc4c3fb90141..1ecfdeb57a8e6 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1528,17 +1528,6 @@ def test_dag_details_active_runs_count_and_queued_runs_count_exclude_backfill_ru 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 - - backfill = Backfill( - dag_id=DAG2_ID, - from_date=datetime(2021, 6, 15, tzinfo=timezone.utc), - to_date=datetime(2021, 6, 16, tzinfo=timezone.utc), - dag_run_conf=None, - ) - session.add(backfill) - session.flush() - session.add( DagRun( dag_id=DAG2_ID, @@ -1558,7 +1547,6 @@ def test_dag_details_active_runs_count_and_queued_runs_count_exclude_backfill_ru start_date=datetime(2021, 6, 15, 2, 0, 0, tzinfo=timezone.utc), run_type=DagRunType.BACKFILL_JOB, state=DagRunState.RUNNING, - backfill_id=backfill.id, triggered_by=DagRunTriggeredByType.TEST, ) ) @@ -1570,7 +1558,6 @@ def test_dag_details_active_runs_count_and_queued_runs_count_exclude_backfill_ru start_date=datetime(2021, 6, 15, 3, 0, 0, tzinfo=timezone.utc), run_type=DagRunType.BACKFILL_JOB, state=DagRunState.QUEUED, - backfill_id=backfill.id, triggered_by=DagRunTriggeredByType.TEST, ) ) diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index fb8cd54f56ffe..88f97d33b2f72 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2925,12 +2925,23 @@ class DAGDetailsResponse(BaseModel): rerun_with_latest_version: Annotated[bool | None, Field(title="Rerun With Latest Version")] = None owner_links: Annotated[dict[str, str] | None, Field(title="Owner Links")] = None is_favorite: Annotated[bool | None, Field(title="Is Favorite")] = False - active_runs_count: Annotated[int | None, Field(title="Active Runs Count")] = 0 - queued_runs_count: Annotated[int | None, Field(title="Queued Runs Count")] = 0 + active_runs_count: Annotated[ + int | None, + Field( + description="Number of currently running runs (backfill runs excluded).", + title="Active Runs Count", + ), + ] = 0 + queued_runs_count: Annotated[ + int | None, + Field( + description="Number of currently queued runs (backfill runs excluded).", title="Queued Runs Count" + ), + ] = 0 is_at_max_active_runs: Annotated[ 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 one running backfill run and no others can show active_runs_count: 1 alongside is_at_max_active_runs: false, and a Dag with one queued (non-backfill) run and otherwise no active runs can show active_runs_count: 0 alongside is_at_max_active_runs: true.", + description="Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).", title="Is At Max Active Runs", ), ]