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..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 @@ -28,6 +28,7 @@ from pydantic import ( AliasGenerator, ConfigDict, + Field, computed_field, field_serializer, field_validator, @@ -220,7 +221,18 @@ class DAGDetailsResponse(DAGResponse): rerun_with_latest_version: bool | None = None owner_links: dict[str, str] | None = None is_favorite: bool = False - active_runs_count: int = 0 + active_runs_count: int = Field( + default=0, description="Number of currently running runs (backfill runs excluded)." + ) + queued_runs_count: int = Field( + default=0, description="Number of currently queued runs (backfill runs excluded)." + ) + is_at_max_active_runs: bool = Field( + description=( + "Whether this Dag is currently at its max_active_runs limit, counting running and " + "queued runs (backfill runs excluded)." + ) + ) team_name: str | None = None @field_validator("timezone", mode="before") 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..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,7 +13437,18 @@ components: active_runs_count: type: integer title: Active Runs Count + description: Number of currently running runs (backfill runs excluded). default: 0 + queued_runs_count: + type: integer + title: Queued Runs Count + description: Number of currently queued runs (backfill runs excluded). + default: 0 + is_at_max_active_runs: + type: boolean + title: Is At Max Active Runs + description: Whether this Dag is currently at its max_active_runs limit, + counting running and queued runs (backfill runs excluded). team_name: anyOf: - type: string @@ -13511,6 +13522,7 @@ components: - timezone - last_parsed - default_args + - is_at_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 3b1f275a8b5da..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,19 +265,33 @@ def get_dag_details( is not None ) - # Count only running Dag runs: this stat shows runs that are actually executing right now. - active_runs_count = ( - session.scalar( - select(func.count()) - .select_from(DagRun) - .where(DagRun.dag_id == dag_id, DagRun.state == DagRunState.RUNNING) + # One query for both counts so they can't disagree with each other if a run changes state + # in between. Backfill runs are excluded: they're gated by Backfill.max_active_runs, not + # the Dag's own max_active_runs. + active_runs_count = queued_runs_count = 0 + for state, count in session.execute( + select(DagRun.state, func.count()) + .where( + DagRun.dag_id == dag_id, + DagRun.state.in_((DagRunState.RUNNING, DagRunState.QUEUED)), + DagRun.run_type != DagRunType.BACKFILL_JOB, ) - or 0 - ) + .group_by(DagRun.state) + ): + if state == DagRunState.RUNNING: + active_runs_count = count + else: + queued_runs_count = count + + # Computed here rather than read from DagModel.exceeds_max_non_backfill: that scheduler-side + # cache isn't updated when a run is triggered manually or through the API. + is_at_max_active_runs = active_runs_count + queued_runs_count >= (dag_model.max_active_runs or 0) - # Add is_favorite and active_runs_count fields to the Dag model + # Add is_favorite, active_runs_count, queued_runs_count, and is_at_max_active_runs fields setattr(dag_model, "is_favorite", is_favorite) setattr(dag_model, "active_runs_count", active_runs_count) + setattr(dag_model, "queued_runs_count", queued_runs_count) + setattr(dag_model, "is_at_max_active_runs", is_at_max_active_runs) return DAGDetailsResponse.model_validate(dag_model) diff --git a/airflow-core/src/airflow/dag_processing/collection.py b/airflow-core/src/airflow/dag_processing/collection.py index 36b948abad6c2..8c9c81b6fcfd7 100644 --- a/airflow-core/src/airflow/dag_processing/collection.py +++ b/airflow-core/src/airflow/dag_processing/collection.py @@ -167,18 +167,19 @@ def _get_latest_runs_stmt_partitioned(dag_id: str) -> Select: class _RunInfo(NamedTuple): latest_run: DagRun | None - num_active_runs: int @classmethod def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self: """ - Query the run counts from the db. + Query the latest-run info from the db. - :param dags: dict of dags to query + :param dag: the Dag to calculate run info for + :param session: DB session """ - # Skip these queries entirely if no Dags can be scheduled to save time. + # Only used to calculate the *next* scheduled run, so it can be skipped for Dags that can + # never be scheduled in the first place. if not dag.timetable.can_be_scheduled: - return cls(None, 0) + return cls(None) if dag.timetable.partitioned: log.debug("Getting latest run for partitioned Dag", dag_id=dag.dag_id) @@ -195,12 +196,7 @@ def calculate(cls, dag: LazyDeserializedDAG, *, session: Session) -> Self: ) else: log.debug("no latest run found", dag_id=dag.dag_id) - active_run_counts = DagRun.active_runs_of_dags( - dag_ids=[dag.dag_id], - exclude_backfill=True, - session=session, - ) - return cls(latest_run, active_run_counts.get(dag.dag_id, 0)) + return cls(latest_run) def _update_dag_tags(tag_names: set[str], dm: DagModel, *, session: Session) -> None: @@ -623,6 +619,12 @@ def update_dags( session: Session, ) -> None: # we exclude backfill from active run counts since their concurrency is separate + # Batched once across every Dag in this update rather than inside the loop below, + # since DagRun.active_runs_of_dags already accepts a list of dag_ids -- calling it + # per-Dag would turn this into an N+1 query. + active_run_counts = DagRun.active_runs_of_dags( + dag_ids=list(orm_dags), exclude_backfill=True, session=session + ) for dag_id, dm in sorted(orm_dags.items()): run_info = _RunInfo.calculate(dag=self.dags[dag_id], session=session) dag = self.dags[dag_id] @@ -690,7 +692,7 @@ def update_dags( dm.bundle_version = self.bundle_version reference_run: DagRun | None = run_info.latest_run - dm.exceeds_max_non_backfill = run_info.num_active_runs >= dm.max_active_runs + dm.exceeds_max_non_backfill = active_run_counts.get(dag_id, 0) >= dm.max_active_runs dm.calculate_dagrun_date_fields(dag, reference_run=reference_run) if not dag.timetable.asset_condition: dm.schedule_asset_references = [] 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..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,8 +3283,20 @@ export const $DAGDetailsResponse = { active_runs_count: { type: 'integer', title: 'Active Runs Count', + description: 'Number of currently running runs (backfill runs excluded).', default: 0 }, + queued_runs_count: { + type: 'integer', + title: 'Queued Runs Count', + description: 'Number of currently queued runs (backfill runs excluded).', + default: 0 + }, + is_at_max_active_runs: { + type: 'boolean', + title: 'Is At Max Active Runs', + description: 'Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).' + }, team_name: { anyOf: [ { @@ -3331,7 +3343,7 @@ Deprecated: Use max_active_tasks instead.`, } }, type: 'object', - required: ['dag_id', 'dag_display_name', 'is_paused', 'is_stale', 'last_parsed_time', 'last_parse_duration', 'last_expired', 'bundle_name', 'bundle_version', 'relative_fileloc', 'fileloc', 'description', 'timetable_summary', 'timetable_description', 'timetable_partitioned', 'timetable_periodic', 'tags', 'max_active_tasks', 'max_active_runs', 'max_consecutive_failed_dag_runs', 'has_task_concurrency_limits', 'has_import_errors', 'next_dagrun_logical_date', 'next_dagrun_data_interval_start', 'next_dagrun_data_interval_end', 'next_dagrun_run_after', 'allowed_run_types', 'owners', 'catchup', 'dag_run_timeout', 'asset_expression', 'doc_md', 'start_date', 'end_date', 'is_paused_upon_creation', 'params', 'render_template_as_native_obj', 'template_search_path', 'timezone', 'last_parsed', 'default_args', 'is_backfillable', 'file_token', 'concurrency', 'latest_dag_version'], + required: ['dag_id', 'dag_display_name', 'is_paused', 'is_stale', 'last_parsed_time', 'last_parse_duration', 'last_expired', 'bundle_name', 'bundle_version', 'relative_fileloc', 'fileloc', 'description', 'timetable_summary', 'timetable_description', 'timetable_partitioned', 'timetable_periodic', 'tags', 'max_active_tasks', 'max_active_runs', 'max_consecutive_failed_dag_runs', 'has_task_concurrency_limits', 'has_import_errors', 'next_dagrun_logical_date', 'next_dagrun_data_interval_start', 'next_dagrun_data_interval_end', 'next_dagrun_run_after', 'allowed_run_types', 'owners', 'catchup', 'dag_run_timeout', 'asset_expression', 'doc_md', 'start_date', 'end_date', 'is_paused_upon_creation', 'params', 'render_template_as_native_obj', 'template_search_path', 'timezone', 'last_parsed', 'default_args', 'is_at_max_active_runs', 'is_backfillable', 'file_token', 'concurrency', 'latest_dag_version'], title: 'DAGDetailsResponse', description: 'Specific serializer for Dag Details responses.' } as const; 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..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,7 +910,18 @@ export type DAGDetailsResponse = { [key: string]: (string); } | null; is_favorite?: boolean; + /** + * Number of currently running runs (backfill runs excluded). + */ active_runs_count?: number; + /** + * Number of currently queued runs (backfill runs excluded). + */ + queued_runs_count?: number; + /** + * Whether this Dag is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded). + */ + is_at_max_active_runs: boolean; team_name?: string | null; /** * Whether this Dag's schedule supports backfilling. 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..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 @@ -44,6 +44,9 @@ "dagBundle_other": "Dag Bundles", "dagDetails": { "activeRuns": "Active Runs", + "activeRunsExceedsMaxTooltip": "Runs beyond the limit will not start until an existing active run completes.", + "activeRunsOfMax": "{{activeRuns}} of {{maxActiveRuns}}", + "activeRunsWithQueued": "{{activeRuns}} of {{maxActiveRuns}} ({{queuedRuns}} queued)", "catchup": "Catchup", "dagRunTimeout": "Dag Run Timeout", "defaultArgs": "Default Args", diff --git a/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx b/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx index 5a6dce6312a23..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: string; value: ReactNode | string }>; + readonly stats: Array; readonly subTitle?: ReactNode | string; readonly title: ReactNode | string; readonly type: "asset" | "dag" | "dagBundle" | "dagRun" | "task" | "taskGroup" | "taskInstance"; @@ -78,7 +84,10 @@ export const HeaderCard = ({ actions, icon, state, stats, subTitle, title, type {stats.map((stat) => ( - + = { multi_team: false }; @@ -49,6 +50,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, @@ -56,6 +58,10 @@ const mockDag = { } as unknown as DAGDetailsResponse; describe("Header", () => { + beforeAll(() => { + i18n.addResourceBundle("en", "common", commonLocale, true, true); + }); + afterEach(() => { mockConfig.multi_team = false; }); @@ -67,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(); }); @@ -78,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(); }); @@ -93,6 +99,80 @@ describe("Header", () => { expect(screen.getByText("2 of 2")).toBeInTheDocument(); }); + it("does not show an info icon or queued count when nothing is queued", () => { + render( + +
+ , + ); + + expect(screen.queryByTestId("active-runs-exceeds-max-info")).not.toBeInTheDocument(); + expect(screen.getByText("1 of 2")).toBeInTheDocument(); + }); + + it("does not show the icon when exactly at capacity with nothing queued", () => { + render( + +
+ , + ); + + 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( + +
+ , + ); + + expect(screen.getByTestId("active-runs-exceeds-max-info")).toBeInTheDocument(); + expect(screen.getByText("1 of 1 (2 queued)")).toBeInTheDocument(); + }); + + it("shows the info icon when over capacity, e.g. after max_active_runs was lowered", () => { + render( + +
+ , + ); + + 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( + +
+ , + ); + + 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", () => { render( @@ -100,7 +180,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 7687488fcfc19..1377de0161a1d 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 { FiBookOpen, FiInfo } 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"; @@ -65,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 ? [] @@ -113,11 +121,31 @@ export const Header = ({ }, ...nextRunStat, { - label: translate("dagDetails.activeRuns"), + key: "activeRuns", + label: + isBlockedByMaxActiveRuns && hasQueuedRuns ? ( + + {translate("dagDetails.activeRuns")} + + + + + ) : ( + translate("dagDetails.activeRuns") + ), value: dag?.max_active_runs === undefined ? undefined - : `${dag.active_runs_count ?? 0} of ${dag.max_active_runs}`, + : hasQueuedRuns + ? translate("dagDetails.activeRunsWithQueued", { + activeRuns: dag.active_runs_count ?? 0, + maxActiveRuns: dag.max_active_runs, + queuedRuns: dag.queued_runs_count, + }) + : 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 73d2d08913cb7..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 @@ -1356,6 +1356,7 @@ def test_dag_details( "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, @@ -1394,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, @@ -1455,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): - """Test that DAG details include the active_runs_count field.""" - # Create running and queued DAG runs for DAG2 + 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.""" + # 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, @@ -1480,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( @@ -1498,19 +1513,153 @@ def test_dag_details_includes_active_runs_count(self, session, test_client): assert response.status_code == 200 body = response.json() - # Verify active_runs_count field is present and correct - assert "active_runs_count" in body - assert isinstance(body["active_runs_count"], int) - assert body["active_runs_count"] == 1 # only running counts, queued does not + assert body["active_runs_count"] == 1 # only running counts, queued/success do not + assert body["queued_runs_count"] == 2 # only queued counts, running/success do not # Test with DAG that has no active runs response = test_client.get(f"/dags/{DAG1_ID}/details") assert response.status_code == 200 body = response.json() - assert "active_runs_count" in body - assert isinstance(body["active_runs_count"], int) assert body["active_runs_count"] == 0 + assert body["queued_runs_count"] == 0 + + def test_dag_details_active_runs_count_and_queued_runs_count_exclude_backfill_runs( + self, session, test_client + ): + """Backfill runs don't count against the Dag's own max_active_runs, so they're excluded.""" + 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, + 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, + 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) + dag_model.max_active_runs = 1 + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="is_at_max_active_runs_running", + logical_date=datetime(2021, 6, 15, 4, 0, 0, tzinfo=timezone.utc), + start_date=datetime(2021, 6, 15, 4, 0, 0, tzinfo=timezone.utc), + run_type=DagRunType.MANUAL, + state=DagRunState.RUNNING, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.commit() + + response = test_client.get(f"/dags/{DAG2_ID}/details") + assert response.status_code == 200 + body = response.json() + + assert body["is_at_max_active_runs"] is True + + # Test with a DAG that has not hit its max_active_runs + response = test_client.get(f"/dags/{DAG1_ID}/details") + assert response.status_code == 200 + body = response.json() + + assert body["is_at_max_active_runs"] is False + + def test_dag_details_is_at_max_active_runs_counts_queued_runs_too(self, session, test_client): + """A queued (non-backfill) run counts toward is_at_max_active_runs even with 0 running.""" + dag_model = session.get(DagModel, DAG2_ID) + dag_model.max_active_runs = 1 + session.add( + DagRun( + dag_id=DAG2_ID, + run_id="is_at_max_active_runs_queued", + logical_date=datetime(2021, 6, 15, 4, 0, 0, tzinfo=timezone.utc), + start_date=datetime(2021, 6, 15, 4, 0, 0, tzinfo=timezone.utc), + run_type=DagRunType.MANUAL, + state=DagRunState.QUEUED, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.commit() + + response = test_client.get(f"/dags/{DAG2_ID}/details") + assert response.status_code == 200 + body = response.json() + + assert body["active_runs_count"] == 0 + assert body["is_at_max_active_runs"] is True + + def test_dag_details_is_at_max_active_runs_excludes_backfill_runs(self, session, test_client): + """A running backfill run doesn't count toward the Dag's own active-run stats.""" + dag_model = session.get(DagModel, DAG2_ID) + 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, + 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, + triggered_by=DagRunTriggeredByType.TEST, + ) + ) + session.commit() + + response = test_client.get(f"/dags/{DAG2_ID}/details") + assert response.status_code == 200 + body = response.json() + + # 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 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-core/tests/unit/dag_processing/test_collection.py b/airflow-core/tests/unit/dag_processing/test_collection.py index be36290abcece..e2aba871c9a2a 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,72 @@ 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 reflects its real active-run count. + + ``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: + ... + + 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 ): 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 diff --git a/airflow-core/tests/unit/models/test_dag.py b/airflow-core/tests/unit/models/test_dag.py index b22227b695be4..1af4ff7b26276 100644 --- a/airflow-core/tests/unit/models/test_dag.py +++ b/airflow-core/tests/unit/models/test_dag.py @@ -685,7 +685,8 @@ 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.""" 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 +694,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"), diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 8938ad8d0534b..88f97d33b2f72 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -2925,7 +2925,26 @@ 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 + 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 is currently at its max_active_runs limit, counting running and queued runs (backfill runs excluded).", + 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 f1eaf2e89efe1..748d05e49c4a0 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), + is_at_max_active_runs=False, is_paused_upon_creation=False, params={}, render_template_as_native_obj=True,