This is an automated email from the ASF dual-hosted git repository.
bbovenzi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new b316afb44dc Show larger dag run and task instance counts on the
dashboard (#70892)
b316afb44dc is described below
commit b316afb44dc9686f5c39aab678a9a89cb21d53d2
Author: Jed Cunningham <[email protected]>
AuthorDate: Mon Aug 3 11:07:09 2026 -0600
Show larger dag run and task instance counts on the dashboard (#70892)
The dashboard counted each state separately with a limit of 1000, so any
busy state
showed "1000+" instead of a real number. On a large install almost every
state sat at
the cap, leaving the panel with no usable figures at all.
Count the window in a single scan instead. When it fits, every count is
exact. When it
does not, report what was read as a lower bound, rounded down, which stays
far closer
to the real volume than a fixed cap. Dag runs and task instances are judged
separately,
since a window often holds few enough dag runs to count exactly while
holding far too
many task instances.
Counting a bounded number of rows is also cheaper than the old per-state
limits, which
scanned the whole window for any state that could not fill its own limit
(postgres, 12M
task instances):
window before after
15min 15 ms 4 ms
24h 3,270 ms 45 ms
7d 20,752 ms 46 ms
MySQL and SQLite show the same pattern.
The Dags list keeps its previous capped counts and now owns that constant.
---
.../core_api/datamodels/ui/dashboard.py | 4 +-
.../api_fastapi/core_api/openapi/_private_ui.yaml | 12 ++-
.../airflow/api_fastapi/core_api/routes/ui/dags.py | 4 +-
.../api_fastapi/core_api/routes/ui/dashboard.py | 88 +++++++++++++---------
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 14 +++-
.../airflow/ui/openapi-gen/requests/types.gen.ts | 3 +-
.../Dashboard/HistoricalMetrics/DagRunMetrics.tsx | 19 +++--
.../HistoricalMetrics/HistoricalMetrics.tsx | 4 +-
.../Dashboard/HistoricalMetrics/MetricSection.tsx | 5 +-
.../HistoricalMetrics/TaskInstanceMetrics.tsx | 14 ++--
.../core_api/routes/ui/test_dashboard.py | 66 +++++++++++-----
11 files changed, 149 insertions(+), 84 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dashboard.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dashboard.py
index 46fc051daf3..1235b8d76aa 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dashboard.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dashboard.py
@@ -52,7 +52,9 @@ class HistoricalMetricDataResponse(BaseModel):
dag_run_states: DAGRunStates
task_instance_states: TaskInstanceStateCount
- state_count_limit: int
+ # True when the counts above are floors on the real values rather than
exact figures.
+ dag_run_counts_are_lower_bounds: bool = False
+ task_instance_counts_are_lower_bounds: bool = False
class DashboardDagStatsResponse(BaseModel):
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
index 58d097827ea..32f0f6f1735 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
+++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
@@ -3771,14 +3771,18 @@ components:
$ref: '#/components/schemas/DAGRunStates'
task_instance_states:
$ref: '#/components/schemas/TaskInstanceStateCount'
- state_count_limit:
- type: integer
- title: State Count Limit
+ dag_run_counts_are_lower_bounds:
+ type: boolean
+ title: Dag Run Counts Are Lower Bounds
+ default: false
+ task_instance_counts_are_lower_bounds:
+ type: boolean
+ title: Task Instance Counts Are Lower Bounds
+ default: false
type: object
required:
- dag_run_states
- task_instance_states
- - state_count_limit
title: HistoricalMetricDataResponse
description: Historical Metric Data serializer for responses.
JobResponse:
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dags.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dags.py
index 962196e51ef..2700cb1e0b5 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dags.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dags.py
@@ -68,7 +68,6 @@ from airflow.api_fastapi.core_api.datamodels.ui.dags import (
DAGWithLatestDagRunsResponse,
)
from airflow.api_fastapi.core_api.openapi.exceptions import
create_openapi_http_exception_doc
-from airflow.api_fastapi.core_api.routes.ui.dashboard import STATE_COUNT_CAP
from airflow.api_fastapi.core_api.security import (
GetUserDep,
ReadableDagsFilterDep,
@@ -83,6 +82,9 @@ from airflow.utils.state import DagRunState, TaskInstanceState
dags_router = AirflowRouter(prefix="/dags", tags=["DAG"])
+# Per-dag run counts read at most this many rows per state; the UI shows "N+"
at the cap.
+STATE_COUNT_CAP = 1000
+
@dags_router.get(
"",
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dashboard.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dashboard.py
index 0498082a4ed..62c51c79d2e 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dashboard.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dashboard.py
@@ -16,10 +16,11 @@
# under the License.
from __future__ import annotations
-from typing import cast
+from decimal import ROUND_FLOOR, Context
+from typing import TYPE_CHECKING, cast
from fastapi import Depends, status
-from sqlalchemy import func, literal, select, union_all
+from sqlalchemy import func, select
from sqlalchemy.sql.expression import case, false
from airflow._shared.timezones import timezone
@@ -38,11 +39,45 @@ from airflow.models.dagrun import DagRun
from airflow.models.taskinstance import TaskInstance
from airflow.utils.state import DagRunState, TaskInstanceState
+if TYPE_CHECKING:
+ from sqlalchemy.orm import Session
+
dashboard_router = AirflowRouter(tags=["Dashboard"], prefix="/dashboard")
-# Cap for state counts — avoids counting millions of rows.
-# The UI shows "N+" when the returned count equals this value.
-STATE_COUNT_CAP = 1000
+# Rows a single scan reads. Windows that fit are counted exactly; wider ones
report a floor.
+EXACT_COUNT_LIMIT = 50_000
+
+
+_ROUNDING = Context(prec=2, rounding=ROUND_FLOOR)
+
+
+def _round_down(value: int) -> int:
+ """Round to two significant digits, never upwards: that would claim
uncounted rows."""
+ return int(_ROUNDING.create_decimal(value))
+
+
+def _compute_state_counts(
+ model, filters, *, session: Session, join=None, null_label: str | None =
None
+) -> tuple[dict[str, int], bool]:
+ """
+ Per-state counts for the window, and whether they are lower bounds rather
than exact.
+
+ A scan that stopped early counted only some of the rows, but each count is
still a floor
+ on the real value.
+ """
+ stmt = select(model.state.label("state")).select_from(model)
+ if join is not None:
+ stmt = stmt.join(join)
+ window = stmt.where(*filters).limit(EXACT_COUNT_LIMIT + 1).subquery()
+ rows = session.execute(select(window.c.state,
func.count().label("cnt")).group_by(window.c.state)).all()
+ are_lower_bounds = sum(row.cnt for row in rows) > EXACT_COUNT_LIMIT
+ counts: dict[str, int] = {}
+ for row in rows:
+ label = row.state or null_label
+ if label is None:
+ continue
+ counts[label] = _round_down(row.cnt) if are_lower_bounds else row.cnt
+ return counts, are_lower_bounds
@dashboard_router.get(
@@ -70,50 +105,31 @@ def historical_metrics(
DagRun.dag_id.in_(permitted_dag_ids),
]
- # Build one LIMIT-capped subquery per state, then UNION ALL them into a
- # single query. Every state gets the same treatment: at most
STATE_COUNT_CAP
- # rows are read from the index, so even states with millions of rows
- # (typically "success") are counted in single-digit milliseconds.
- # Each branch is wrapped in a subquery so LIMIT works on all backends
- # (SQLite rejects LIMIT inside bare UNION ALL arms).
- def _capped_state_counts(model, states, label_fn, join=None):
- branches = []
- for state in states:
- stmt =
select(literal(label_fn(state)).label("state")).select_from(model)
- if join is not None:
- stmt = stmt.join(join)
- branch = (
- stmt.where(*dag_run_filters)
- .where(model.state == state if state else
model.state.is_(None))
- .limit(STATE_COUNT_CAP)
- .subquery()
- )
- branches.append(select(branch.c.state))
- capped = union_all(*branches).subquery()
- return session.execute(
- select(capped.c.state,
func.count().label("cnt")).group_by(capped.c.state)
- ).all()
-
- dag_run_state_counts = _capped_state_counts(DagRun, list(DagRunState),
lambda s: s.value)
- ti_state_counts = _capped_state_counts(
+ # Judged separately: dag runs often fit when task instances do not.
+ dag_run_states, dag_runs_are_lower_bounds = _compute_state_counts(
+ DagRun, dag_run_filters, session=session
+ )
+ task_instance_states, task_instances_are_lower_bounds =
_compute_state_counts(
TaskInstance,
- [None, *TaskInstanceState],
- lambda s: s.value if s else "no_status",
+ dag_run_filters,
+ session=session,
join=TaskInstance.dag_run,
+ null_label="no_status",
)
return HistoricalMetricDataResponse.model_validate(
{
"dag_run_states": {
**{dag_run_state.value: 0 for dag_run_state in DagRunState},
- **{row.state: row.cnt for row in dag_run_state_counts},
+ **dag_run_states,
},
"task_instance_states": {
"no_status": 0,
**{ti_state.value: 0 for ti_state in TaskInstanceState},
- **{row.state: row.cnt for row in ti_state_counts},
+ **task_instance_states,
},
- "state_count_limit": STATE_COUNT_CAP,
+ "dag_run_counts_are_lower_bounds": dag_runs_are_lower_bounds,
+ "task_instance_counts_are_lower_bounds":
task_instances_are_lower_bounds,
}
)
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 d14b118651c..e1f10b5e91e 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
@@ -10165,13 +10165,19 @@ export const $HistoricalMetricDataResponse = {
task_instance_states: {
'$ref': '#/components/schemas/TaskInstanceStateCount'
},
- state_count_limit: {
- type: 'integer',
- title: 'State Count Limit'
+ dag_run_counts_are_lower_bounds: {
+ type: 'boolean',
+ title: 'Dag Run Counts Are Lower Bounds',
+ default: false
+ },
+ task_instance_counts_are_lower_bounds: {
+ type: 'boolean',
+ title: 'Task Instance Counts Are Lower Bounds',
+ default: false
}
},
type: 'object',
- required: ['dag_run_states', 'task_instance_states', 'state_count_limit'],
+ required: ['dag_run_states', 'task_instance_states'],
title: 'HistoricalMetricDataResponse',
description: 'Historical Metric Data serializer for 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 f19669c72fe..c0b4e673d86 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
@@ -2586,7 +2586,8 @@ export type GridTISummaries = {
export type HistoricalMetricDataResponse = {
dag_run_states: DAGRunStates;
task_instance_states: TaskInstanceStateCount;
- state_count_limit: number;
+ dag_run_counts_are_lower_bounds?: boolean;
+ task_instance_counts_are_lower_bounds?: boolean;
};
/**
diff --git
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/DagRunMetrics.tsx
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/DagRunMetrics.tsx
index 7d0bec32d76..6908564e504 100644
---
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/DagRunMetrics.tsx
+++
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/DagRunMetrics.tsx
@@ -24,21 +24,26 @@ import { FiBarChart } from "react-icons/fi";
import { MetricSection } from "./MetricSection";
type DagRunMetricsProps = {
+ readonly countsAreLowerBounds: boolean;
readonly dagRunStates: DAGRunStates;
readonly endDate?: string;
readonly startDate: string;
- readonly stateCountLimit: number;
};
const DAGRUN_STATES: Array<keyof DAGRunStates> = ["queued", "running",
"success", "failed"];
-export const DagRunMetrics = ({ dagRunStates, endDate, startDate,
stateCountLimit }: DagRunMetricsProps) => {
+export const DagRunMetrics = ({
+ countsAreLowerBounds,
+ dagRunStates,
+ endDate,
+ startDate,
+}: DagRunMetricsProps) => {
const { t: translate } = useTranslation();
const total = Object.values(dagRunStates).reduce((sum, count) => sum +
count, 0);
- // When any state hit the API's STATE_COUNT_CAP, the summed total is only a
- // lower bound, so per-state percentages computed from it are wrong (#67336).
- // Suppress percentages for the whole group in that case.
- const isTotalTruncated = Object.values(dagRunStates).some((count) => count
>= stateCountLimit);
+ // The total is only a lower bound when the counts are, so percentages would
be wrong.
+ const isTotalTruncated = countsAreLowerBounds;
+ // "0+" would be meaningless.
+ const isLowerBound = (count: number) => countsAreLowerBounds && count > 0;
return (
<Box borderRadius={5} borderWidth={1} p={4}>
@@ -50,7 +55,7 @@ export const DagRunMetrics = ({ dagRunStates, endDate,
startDate, stateCountLimi
<Stack gap={4}>
{DAGRUN_STATES.map((state) => (
<MetricSection
- capped={dagRunStates[state] >= stateCountLimit}
+ capped={isLowerBound(dagRunStates[state])}
endDate={endDate}
isTotalTruncated={isTotalTruncated}
key={state}
diff --git
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/HistoricalMetrics.tsx
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/HistoricalMetrics.tsx
index d193522beeb..24174bc9b03 100644
---
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/HistoricalMetrics.tsx
+++
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/HistoricalMetrics.tsx
@@ -74,13 +74,13 @@ export const HistoricalMetrics = ({ endDate, startDate }:
HistoricalMetricsProps
{!isLoading && data !== undefined && (
<Box>
<DagRunMetrics
+ countsAreLowerBounds={data.dag_run_counts_are_lower_bounds
?? false}
dagRunStates={data.dag_run_states}
startDate={startDate}
- stateCountLimit={data.state_count_limit}
/>
<TaskInstanceMetrics
+
countsAreLowerBounds={data.task_instance_counts_are_lower_bounds ?? false}
startDate={startDate}
- stateCountLimit={data.state_count_limit}
taskInstanceStates={data.task_instance_states}
/>
</Box>
diff --git
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/MetricSection.tsx
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/MetricSection.tsx
index 20aae19444d..1cc585ffa3f 100644
---
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/MetricSection.tsx
+++
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/MetricSection.tsx
@@ -48,6 +48,7 @@ export const MetricSection = ({
state,
total,
}: MetricSectionProps) => {
+ // A lower bound has no known proportion, so it deliberately fills the bar.
const stateWidth = capped ? BAR_WIDTH : total === 0 ? 0 : (runs / total) *
BAR_WIDTH;
const remainingWidth = BAR_WIDTH - stateWidth;
const hidePercent = isTotalTruncated;
@@ -57,7 +58,7 @@ export const MetricSection = ({
const searchParams = new URLSearchParams(
`?${stateParam}=${state}&${SearchParamsKeys.START_DATE_GTE}=${startDate}`,
);
- const { t: translate } = useTranslation();
+ const { i18n, t: translate } = useTranslation();
if (endDate !== undefined) {
searchParams.append(SearchParamsKeys.END_DATE, endDate);
@@ -70,7 +71,7 @@ export const MetricSection = ({
<RouterLink to={`/${kind}?${searchParams.toString()}`}>
<StateBadge fontSize="md" state={state === "no_status" ? null :
state}>
{}
- {capped ? `${runs}+` : runs}
+ {`${runs.toLocaleString(i18n.language)}${capped ? "+" : ""}`}
</StateBadge>
</RouterLink>
<Text>{translate(`states.${state}`)}</Text>
diff --git
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/TaskInstanceMetrics.tsx
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/TaskInstanceMetrics.tsx
index 5cbf5811e17..fa839ed6014 100644
---
a/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/TaskInstanceMetrics.tsx
+++
b/airflow-core/src/airflow/ui/src/pages/Dashboard/HistoricalMetrics/TaskInstanceMetrics.tsx
@@ -24,9 +24,9 @@ import { MdOutlineTask } from "react-icons/md";
import { MetricSection } from "./MetricSection";
type TaskInstanceMetricsProps = {
+ readonly countsAreLowerBounds: boolean;
readonly endDate?: string;
readonly startDate: string;
- readonly stateCountLimit: number;
readonly taskInstanceStates: TaskInstanceStateCount;
};
@@ -48,17 +48,17 @@ const TASK_STATES: Array<keyof TaskInstanceStateCount> = [
];
export const TaskInstanceMetrics = ({
+ countsAreLowerBounds,
endDate,
startDate,
- stateCountLimit,
taskInstanceStates,
}: TaskInstanceMetricsProps) => {
const { t: translate } = useTranslation();
const total = Object.values(taskInstanceStates).reduce((sum, count) => sum +
count, 0);
- // When any state hit the API's STATE_COUNT_CAP, the summed total is only a
- // lower bound, so per-state percentages computed from it are wrong (#67336).
- // Suppress percentages for the whole group in that case.
- const isTotalTruncated = Object.values(taskInstanceStates).some((count) =>
count >= stateCountLimit);
+ // The total is only a lower bound when the counts are, so percentages would
be wrong.
+ const isTotalTruncated = countsAreLowerBounds;
+ // "0+" would be meaningless.
+ const isLowerBound = (count: number) => countsAreLowerBounds && count > 0;
return (
<Box borderRadius={5} borderWidth={1} mt={2} p={4}>
@@ -73,7 +73,7 @@ export const TaskInstanceMetrics = ({
).map((state) =>
taskInstanceStates[state] > 0 ? (
<MetricSection
- capped={taskInstanceStates[state] >= stateCountLimit}
+ capped={isLowerBound(taskInstanceStates[state])}
endDate={endDate}
isTotalTruncated={isTotalTruncated}
key={state}
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py
index 50869e0ee9f..6ed4eada4b3 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py
@@ -18,13 +18,16 @@
from __future__ import annotations
from datetime import timedelta
+from types import SimpleNamespace
from unittest import mock
import pendulum
import pytest
+from airflow.api_fastapi.core_api.routes.ui import dashboard
from airflow.models.dag import DagModel
from airflow.models.dagbag import DBDagBag
+from airflow.models.dagrun import DagRun
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.utils.state import DagRunState, TaskInstanceState
from airflow.utils.types import DagRunType
@@ -265,7 +268,8 @@ class TestHistoricalMetricsDataEndpoint:
"up_for_retry": 0,
"upstream_failed": 0,
},
- "state_count_limit": 1000,
+ "dag_run_counts_are_lower_bounds": False,
+ "task_instance_counts_are_lower_bounds": False,
},
),
(
@@ -288,7 +292,8 @@ class TestHistoricalMetricsDataEndpoint:
"up_for_retry": 0,
"upstream_failed": 0,
},
- "state_count_limit": 1000,
+ "dag_run_counts_are_lower_bounds": False,
+ "task_instance_counts_are_lower_bounds": False,
},
),
(
@@ -311,7 +316,8 @@ class TestHistoricalMetricsDataEndpoint:
"up_for_retry": 0,
"upstream_failed": 0,
},
- "state_count_limit": 1000,
+ "dag_run_counts_are_lower_bounds": False,
+ "task_instance_counts_are_lower_bounds": False,
},
),
],
@@ -323,29 +329,30 @@ class TestHistoricalMetricsDataEndpoint:
assert response.status_code == 200
assert response.json() == expected
+ @pytest.mark.parametrize(
+ ("exact_limit", "dag_runs_bounded", "task_instances_bounded"),
+ [
+ pytest.param(8, False, False, id="both-counted-exactly"),
+ pytest.param(4, False, True,
id="dag-runs-exact-task-instances-bounded"),
+ pytest.param(3, True, True, id="neither-counted-exactly"),
+ ],
+ )
@pytest.mark.usefixtures("freeze_time_for_dagruns", "make_dag_runs")
- def test_state_counts_are_capped(self, test_client):
- """State counts are capped at STATE_COUNT_CAP; fixture creates 4 dag
runs and 8 TIs."""
- with
mock.patch("airflow.api_fastapi.core_api.routes.ui.dashboard.STATE_COUNT_CAP",
1):
+ def test_exact_limit_applies_per_group(
+ self, test_client, exact_limit, dag_runs_bounded,
task_instances_bounded
+ ):
+ """The fixture's 4 dag runs and 8 task instances cross
EXACT_COUNT_LIMIT independently."""
+ with mock.patch.object(dashboard, "EXACT_COUNT_LIMIT", exact_limit):
response = test_client.get(
"/dashboard/historical_metrics_data",
params={"start_date": "2023-01-01T00:00", "end_date":
"2023-08-02T00:00"},
)
assert response.status_code == 200
data = response.json()
-
- assert data["state_count_limit"] == 1
-
- dr_states = data["dag_run_states"]
- assert dr_states["success"] == 1
- assert dr_states["failed"] == 1
- assert dr_states["running"] == 1
- assert dr_states["queued"] == 1
-
- ti_states = data["task_instance_states"]
- assert ti_states["success"] == 1
- assert ti_states["failed"] == 1
- assert ti_states["no_status"] == 1
+ assert data["dag_run_counts_are_lower_bounds"] is dag_runs_bounded
+ assert data["task_instance_counts_are_lower_bounds"] is
task_instances_bounded
+ if not dag_runs_bounded:
+ assert data["dag_run_states"] == {"failed": 1, "queued": 1,
"running": 1, "success": 1}
def test_should_response_401(self, unauthenticated_test_client):
response = unauthenticated_test_client.get(
@@ -428,3 +435,24 @@ class TestDagStatsEndpoint:
def test_should_response_403(self, unauthorized_test_client):
response = unauthorized_test_client.get("/dashboard/dag_stats")
assert response.status_code == 403
+
+
[email protected](
+ ("limit", "expected_counts", "expected_bounded"),
+ [
+ pytest.param(1_000, {"success": 4_200, "failed": 130}, True,
id="bounded-counts-rounded-down"),
+ pytest.param(50_000, {"success": 4_250, "failed": 137}, False,
id="exact-counts-left-alone"),
+ ],
+)
+def test_compute_state_counts_rounds_only_lower_bounds(limit, expected_counts,
expected_bounded):
+ """Rounding applies to lower bounds and nothing else."""
+ session = mock.MagicMock()
+ session.execute.return_value.all.return_value = [
+ SimpleNamespace(state="success", cnt=4_250),
+ SimpleNamespace(state="failed", cnt=137),
+ ]
+ with mock.patch.object(dashboard, "EXACT_COUNT_LIMIT", limit):
+ counts, are_lower_bounds = dashboard._compute_state_counts(DagRun, [],
session=session)
+
+ assert are_lower_bounds is expected_bounded
+ assert counts == expected_counts