From 48872f8e32234fb75992678c0bf038e97a089890 Mon Sep 17 00:00:00 2001 From: Yuseok Jo Date: Tue, 7 Jul 2026 21:49:16 +0900 Subject: [PATCH] Show task state counts for the latest Dag run on the Dags list page --- .../core_api/datamodels/ui/dags.py | 19 +++ .../core_api/openapi/_private_ui.yaml | 77 +++++++++++ .../api_fastapi/core_api/routes/ui/dags.py | 74 +++++++++++ .../airflow/ui/openapi-gen/queries/common.ts | 6 + .../ui/openapi-gen/queries/ensureQueryData.ts | 13 ++ .../ui/openapi-gen/queries/prefetch.ts | 13 ++ .../airflow/ui/openapi-gen/queries/queries.ts | 13 ++ .../ui/openapi-gen/queries/suspense.ts | 13 ++ .../ui/openapi-gen/requests/schemas.gen.ts | 43 ++++++ .../ui/openapi-gen/requests/services.gen.ts | 25 +++- .../ui/openapi-gen/requests/types.gen.ts | 42 ++++++ .../ui/public/i18n/locales/en/dags.json | 5 + .../ui/src/pages/DagsList/DagCard.test.tsx | 16 ++- .../airflow/ui/src/pages/DagsList/DagCard.tsx | 24 +++- .../ui/src/pages/DagsList/DagsList.tsx | 47 ++++++- .../LatestRunTaskStateCounts.test.tsx | 91 +++++++++++++ .../DagsList/LatestRunTaskStateCounts.tsx | 96 ++++++++++++++ .../queries/useLatestRunTaskStateCounts.tsx | 45 +++++++ .../core_api/routes/ui/test_dags.py | 123 ++++++++++++++++++ 19 files changed, 775 insertions(+), 10 deletions(-) create mode 100644 airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.test.tsx create mode 100644 airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.tsx create mode 100644 airflow-core/src/airflow/ui/src/queries/useLatestRunTaskStateCounts.tsx diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dags.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dags.py index 590084cc2c5bf..68267c9a02346 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dags.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dags.py @@ -48,6 +48,25 @@ class DAGRunStateCountsResponse(BaseModel): state_counts: dict[DagRunState, int] +class DAGLatestRunTaskInstanceStateCountsResponse(BaseModel): + """ + Task-instance state counts for a Dag's latest run. + + ``state_counts`` only carries states present in the run; task instances without a + state yet are keyed as ``no_status``. + """ + + dag_id: str + run_id: str + state_counts: dict[str, int] + + +class DAGsLatestRunTaskInstanceStateCountsCollectionResponse(BaseModel): + """Collection of per-Dag latest-run task-instance state counts for the Dag list page.""" + + dags: list[DAGLatestRunTaskInstanceStateCountsResponse] + + class DAGsRunStateCountsCollectionResponse(BaseModel): """Collection of per-Dag DagRun-state counts for the Dag list page.""" 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 3defcad70f9fd..bfaffcb35e2df 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 @@ -640,6 +640,44 @@ paths: application/json: schema: $ref: '#/components/schemas/HTTPValidationError' + /ui/dags/latest_run_task_instance_state_counts: + get: + tags: + - DAG + summary: Get Latest Run Task Instance State Counts + description: 'Return task-instance state counts for each Dag''s latest run, + for the Dag list page. + + + Dags without any run are omitted from the response.' + operationId: get_latest_run_task_instance_state_counts_ui + security: + - OAuth2PasswordBearer: [] + - HTTPBearer: [] + parameters: + - name: dag_ids + in: query + required: true + schema: + type: array + items: + type: string + minItems: 1 + maxItems: 100 + title: Dag Ids + responses: + '200': + description: Successful Response + content: + application/json: + schema: + $ref: '#/components/schemas/DAGsLatestRunTaskInstanceStateCountsCollectionResponse' + '422': + description: Validation Error + content: + application/json: + schema: + $ref: '#/components/schemas/HTTPValidationError' /ui/dependencies: get: tags: @@ -2366,6 +2404,32 @@ components: that the API server/Web UI can use this data to render connection form UI.' + DAGLatestRunTaskInstanceStateCountsResponse: + properties: + dag_id: + type: string + title: Dag Id + run_id: + type: string + title: Run Id + state_counts: + additionalProperties: + type: integer + type: object + title: State Counts + type: object + required: + - dag_id + - run_id + - state_counts + title: DAGLatestRunTaskInstanceStateCountsResponse + description: 'Task-instance state counts for a Dag''s latest run. + + + ``state_counts`` only carries states present in the run; task instances without + a + + state yet are keyed as ``no_status``.' DAGRunLightResponse: properties: id: @@ -2677,6 +2741,19 @@ components: - file_token title: DAGWithLatestDagRunsResponse description: DAG with latest dag runs response serializer. + DAGsLatestRunTaskInstanceStateCountsCollectionResponse: + properties: + dags: + items: + $ref: '#/components/schemas/DAGLatestRunTaskInstanceStateCountsResponse' + type: array + title: Dags + type: object + required: + - dags + title: DAGsLatestRunTaskInstanceStateCountsCollectionResponse + description: Collection of per-Dag latest-run task-instance state counts for + the Dag list page. DAGsRunStateCountsCollectionResponse: properties: dags: 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 38d0c93cef88d..af7ebe9599969 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 @@ -58,7 +58,9 @@ from airflow.api_fastapi.core_api.datamodels.dags import DAG_ALIAS_MAPPING, DAGResponse from airflow.api_fastapi.core_api.datamodels.ui.dag_runs import DAGRunLightResponse from airflow.api_fastapi.core_api.datamodels.ui.dags import ( + DAGLatestRunTaskInstanceStateCountsResponse, DAGRunStateCountsResponse, + DAGsLatestRunTaskInstanceStateCountsCollectionResponse, DAGsRunStateCountsCollectionResponse, DAGWithLatestDagRunsCollectionResponse, DAGWithLatestDagRunsResponse, @@ -343,3 +345,75 @@ def get_dag_run_state_counts( ], state_count_limit=STATE_COUNT_CAP, ) + + +@dags_router.get( + "/latest_run_task_instance_state_counts", + dependencies=[ + Depends(requires_access_dag(method="GET")), + Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK_INSTANCE)), + ], + operation_id="get_latest_run_task_instance_state_counts_ui", +) +def get_latest_run_task_instance_state_counts( + session: SessionDep, + readable_dags_filter: ReadableDagsFilterDep, + dag_ids: Annotated[list[str], Query(min_length=1, max_length=conf.getint("api", "maximum_page_limit"))], +) -> DAGsLatestRunTaskInstanceStateCountsCollectionResponse: + """ + Return task-instance state counts for each Dag's latest run, for the Dag list page. + + Dags without any run are omitted from the response. + """ + permitted_dag_ids = readable_dags_filter.value or set() + requested_dag_ids = sorted(set(dag_ids) & permitted_dag_ids) + + dags: list[DAGLatestRunTaskInstanceStateCountsResponse] = [] + if not requested_dag_ids: + return DAGsLatestRunTaskInstanceStateCountsCollectionResponse(dags=dags) + + latest_run_branches = [ + select(DagRun.dag_id, DagRun.run_id) + .where(DagRun.dag_id == dag_id) + .order_by(DagRun.run_after.desc()) + .limit(1) + .subquery() + for dag_id in requested_dag_ids + ] + latest_runs_union = union_all(*(select(branch) for branch in latest_run_branches)).subquery() + latest_run_id_by_dag: dict[str, str] = { + row.dag_id: row.run_id for row in session.execute(select(latest_runs_union)) + } + + if latest_run_id_by_dag: + # Each branch filters on (dag_id, run_id) equality, which the ti_dag_run index + # covers. A run's task instances are bounded by the Dag's task structure, so the + # per-state counts are exact (no cap needed here, unlike the cross-run counts in + # get_dag_run_state_counts). + ti_branches = [ + select(literal(dag_id).label("dag_id"), TaskInstance.state.label("state")) + .where(TaskInstance.dag_id == dag_id, TaskInstance.run_id == run_id) + .subquery() + for dag_id, run_id in latest_run_id_by_dag.items() + ] + tis_union = union_all(*(select(branch) for branch in ti_branches)).subquery() + counts_by_dag: dict[str, dict[str, int]] = {dag_id: {} for dag_id in latest_run_id_by_dag} + for row in session.execute( + select(tis_union.c.dag_id, tis_union.c.state, func.count().label("cnt")).group_by( + tis_union.c.dag_id, tis_union.c.state + ) + ): + state_key = row.state if row.state is not None else "no_status" + counts_by_dag[row.dag_id][state_key] = row.cnt + + dags = [ + DAGLatestRunTaskInstanceStateCountsResponse( + dag_id=dag_id, + run_id=latest_run_id_by_dag[dag_id], + state_counts=counts_by_dag[dag_id], + ) + for dag_id in requested_dag_ids + if dag_id in latest_run_id_by_dag + ] + + return DAGsLatestRunTaskInstanceStateCountsCollectionResponse(dags=dags) diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts index 2d6dc7ffe7db3..0063a9627c4b3 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts @@ -364,6 +364,12 @@ export const useDagServiceGetDagRunStateCountsUiKey = "DagServiceGetDagRunStateC export const UseDagServiceGetDagRunStateCountsUiKeyFn = ({ dagIds }: { dagIds: string[]; }, queryKey?: Array) => [useDagServiceGetDagRunStateCountsUiKey, ...(queryKey ?? [{ dagIds }])]; +export type DagServiceGetLatestRunTaskInstanceStateCountsUiDefaultResponse = Awaited>; +export type DagServiceGetLatestRunTaskInstanceStateCountsUiQueryResult = UseQueryResult; +export const useDagServiceGetLatestRunTaskInstanceStateCountsUiKey = "DagServiceGetLatestRunTaskInstanceStateCountsUi"; +export const UseDagServiceGetLatestRunTaskInstanceStateCountsUiKeyFn = ({ dagIds }: { + dagIds: string[]; +}, queryKey?: Array) => [useDagServiceGetLatestRunTaskInstanceStateCountsUiKey, ...(queryKey ?? [{ dagIds }])]; export type EventLogServiceGetEventLogDefaultResponse = Awaited>; export type EventLogServiceGetEventLogQueryResult = UseQueryResult; export const useEventLogServiceGetEventLogKey = "EventLogServiceGetEventLog"; diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts index 6eba7e9a9db90..4f55965136d8d 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts @@ -724,6 +724,19 @@ export const ensureUseDagServiceGetDagRunStateCountsUiData = (queryClient: Query dagIds: string[]; }) => queryClient.ensureQueryData({ queryKey: Common.UseDagServiceGetDagRunStateCountsUiKeyFn({ dagIds }), queryFn: () => DagService.getDagRunStateCountsUi({ dagIds }) }); /** +* Get Latest Run Task Instance State Counts +* Return task-instance state counts for each Dag's latest run, for the Dag list page. +* +* Dags without any run are omitted from the response. +* @param data The data for the request. +* @param data.dagIds +* @returns DAGsLatestRunTaskInstanceStateCountsCollectionResponse Successful Response +* @throws ApiError +*/ +export const ensureUseDagServiceGetLatestRunTaskInstanceStateCountsUiData = (queryClient: QueryClient, { dagIds }: { + dagIds: string[]; +}) => queryClient.ensureQueryData({ queryKey: Common.UseDagServiceGetLatestRunTaskInstanceStateCountsUiKeyFn({ dagIds }), queryFn: () => DagService.getLatestRunTaskInstanceStateCountsUi({ dagIds }) }); +/** * Get Event Log * @param data The data for the request. * @param data.eventLogId diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts index 037f6d01a396c..a3d7068f4bbc3 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts @@ -724,6 +724,19 @@ export const prefetchUseDagServiceGetDagRunStateCountsUi = (queryClient: QueryCl dagIds: string[]; }) => queryClient.prefetchQuery({ queryKey: Common.UseDagServiceGetDagRunStateCountsUiKeyFn({ dagIds }), queryFn: () => DagService.getDagRunStateCountsUi({ dagIds }) }); /** +* Get Latest Run Task Instance State Counts +* Return task-instance state counts for each Dag's latest run, for the Dag list page. +* +* Dags without any run are omitted from the response. +* @param data The data for the request. +* @param data.dagIds +* @returns DAGsLatestRunTaskInstanceStateCountsCollectionResponse Successful Response +* @throws ApiError +*/ +export const prefetchUseDagServiceGetLatestRunTaskInstanceStateCountsUi = (queryClient: QueryClient, { dagIds }: { + dagIds: string[]; +}) => queryClient.prefetchQuery({ queryKey: Common.UseDagServiceGetLatestRunTaskInstanceStateCountsUiKeyFn({ dagIds }), queryFn: () => DagService.getLatestRunTaskInstanceStateCountsUi({ dagIds }) }); +/** * Get Event Log * @param data The data for the request. * @param data.eventLogId diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts index 09ddcdd3837e6..c5cff5dbff5f2 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts @@ -724,6 +724,19 @@ export const useDagServiceGetDagRunStateCountsUi = , "queryKey" | "queryFn">) => useQuery({ queryKey: Common.UseDagServiceGetDagRunStateCountsUiKeyFn({ dagIds }, queryKey), queryFn: () => DagService.getDagRunStateCountsUi({ dagIds }) as TData, ...options }); /** +* Get Latest Run Task Instance State Counts +* Return task-instance state counts for each Dag's latest run, for the Dag list page. +* +* Dags without any run are omitted from the response. +* @param data The data for the request. +* @param data.dagIds +* @returns DAGsLatestRunTaskInstanceStateCountsCollectionResponse Successful Response +* @throws ApiError +*/ +export const useDagServiceGetLatestRunTaskInstanceStateCountsUi = = unknown[]>({ dagIds }: { + dagIds: string[]; +}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useQuery({ queryKey: Common.UseDagServiceGetLatestRunTaskInstanceStateCountsUiKeyFn({ dagIds }, queryKey), queryFn: () => DagService.getLatestRunTaskInstanceStateCountsUi({ dagIds }) as TData, ...options }); +/** * Get Event Log * @param data The data for the request. * @param data.eventLogId diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts index e3ec883630999..e1232d54fb6d4 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts @@ -724,6 +724,19 @@ export const useDagServiceGetDagRunStateCountsUiSuspense = , "queryKey" | "queryFn">) => useSuspenseQuery({ queryKey: Common.UseDagServiceGetDagRunStateCountsUiKeyFn({ dagIds }, queryKey), queryFn: () => DagService.getDagRunStateCountsUi({ dagIds }) as TData, ...options }); /** +* Get Latest Run Task Instance State Counts +* Return task-instance state counts for each Dag's latest run, for the Dag list page. +* +* Dags without any run are omitted from the response. +* @param data The data for the request. +* @param data.dagIds +* @returns DAGsLatestRunTaskInstanceStateCountsCollectionResponse Successful Response +* @throws ApiError +*/ +export const useDagServiceGetLatestRunTaskInstanceStateCountsUiSuspense = = unknown[]>({ dagIds }: { + dagIds: string[]; +}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useSuspenseQuery({ queryKey: Common.UseDagServiceGetLatestRunTaskInstanceStateCountsUiKeyFn({ dagIds }, queryKey), queryFn: () => DagService.getLatestRunTaskInstanceStateCountsUi({ dagIds }) as TData, ...options }); +/** * Get Event Log * @param data The data for the request. * @param data.eventLogId 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 55b25b2c9753f..bdab027cb2dbe 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 @@ -8865,6 +8865,33 @@ It is used to transfer providers information loaded by providers_manager such th the API server/Web UI can use this data to render connection form UI.` } as const; +export const $DAGLatestRunTaskInstanceStateCountsResponse = { + properties: { + dag_id: { + type: 'string', + title: 'Dag Id' + }, + run_id: { + type: 'string', + title: 'Run Id' + }, + state_counts: { + additionalProperties: { + type: 'integer' + }, + type: 'object', + title: 'State Counts' + } + }, + type: 'object', + required: ['dag_id', 'run_id', 'state_counts'], + title: 'DAGLatestRunTaskInstanceStateCountsResponse', + description: `Task-instance state counts for a Dag's latest run. + +\`\`state_counts\`\` only carries states present in the run; task instances without a +state yet are keyed as \`\`no_status\`\`.` +} as const; + export const $DAGRunLightResponse = { properties: { id: { @@ -9315,6 +9342,22 @@ export const $DAGWithLatestDagRunsResponse = { description: 'DAG with latest dag runs response serializer.' } as const; +export const $DAGsLatestRunTaskInstanceStateCountsCollectionResponse = { + properties: { + dags: { + items: { + '$ref': '#/components/schemas/DAGLatestRunTaskInstanceStateCountsResponse' + }, + type: 'array', + title: 'Dags' + } + }, + type: 'object', + required: ['dags'], + title: 'DAGsLatestRunTaskInstanceStateCountsCollectionResponse', + description: 'Collection of per-Dag latest-run task-instance state counts for the Dag list page.' +} as const; + export const $DAGsRunStateCountsCollectionResponse = { properties: { dags: { diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts index 782417965463a..8230fbc06ebc7 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts @@ -3,7 +3,7 @@ import type { CancelablePromise } from './core/CancelablePromise'; import { OpenAPI } from './core/OpenAPI'; import { request as __request } from './core/request'; -import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData, GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse, GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData, CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse, GetAssetQueuedEventsData, GetAssetQueuedEventsResponse, DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData, GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse, DeleteDagAssetQueuedEventsData, DeleteDagAssetQueuedEventsResponse, GetDagAssetQueuedEventData, GetDagAssetQueuedEventResponse, DeleteDagAssetQueuedEventData, DeleteDagAssetQueuedEventResponse, NextRunAssetsData, NextRunAssetsResponse2, ListBackfillsData, ListBackfillsResponse, CreateBackfillData, CreateBackfillResponse, GetBackfillData, GetBackfillResponse, PauseBackfillData, PauseBackfillResponse, UnpauseBackfillData, UnpauseBackfillResponse, CancelBackfillData, CancelBackfillResponse, CreateBackfillDryRunData, CreateBackfillDryRunResponse, ListBackfillsUiData, ListBackfillsUiResponse, DeleteConnectionData, DeleteConnectionResponse, GetConnectionData, GetConnectionResponse, PatchConnectionData, PatchConnectionResponse, GetConnectionTestData, GetConnectionTestResponse, EnqueueConnectionTestData, EnqueueConnectionTestResponse, GetConnectionsData, GetConnectionsResponse, PostConnectionData, PostConnectionResponse, BulkConnectionsData, BulkConnectionsResponse, TestConnectionData, TestConnectionResponse, CreateDefaultConnectionsResponse, HookMetaDataResponse, GetDagRunData, GetDagRunResponse, DeleteDagRunData, DeleteDagRunResponse, PatchDagRunData, PatchDagRunResponse, BulkDagRunsData, BulkDagRunsResponse, GetDagRunsData, GetDagRunsResponse, TriggerDagRunData, TriggerDagRunResponse, GetUpstreamAssetEventsData, GetUpstreamAssetEventsResponse, ClearDagRunData, ClearDagRunResponse, WaitDagRunUntilFinishedData, WaitDagRunUntilFinishedResponse, GetListDagRunsBatchData, GetListDagRunsBatchResponse, ClearDagRunsData, ClearDagRunsResponse, ClearDagRunPartitionsData, ClearDagRunPartitionsResponse, GetDagRunStatsData, GetDagRunStatsResponse, GetDagSourceData, GetDagSourceResponse, GetDagStatsData, GetDagStatsResponse, GetConfigData, GetConfigResponse, GetConfigValueData, GetConfigValueResponse, GetConfigsResponse, ListDagWarningsData, ListDagWarningsResponse, GetDagsData, GetDagsResponse, PatchDagsData, PatchDagsResponse, GetDagData, GetDagResponse, PatchDagData, PatchDagResponse, DeleteDagData, DeleteDagResponse, GetDagDetailsData, GetDagDetailsResponse, FavoriteDagData, FavoriteDagResponse, UnfavoriteDagData, UnfavoriteDagResponse, GetDagTagsData, GetDagTagsResponse, GetDagsUiData, GetDagsUiResponse, GetLatestRunInfoData, GetLatestRunInfoResponse, GetDagRunStateCountsUiData, GetDagRunStateCountsUiResponse, GetEventLogData, GetEventLogResponse, GetEventLogsData, GetEventLogsResponse, GetExtraLinksData, GetExtraLinksResponse, GetTaskInstanceData, GetTaskInstanceResponse, PatchTaskInstanceData, PatchTaskInstanceResponse, DeleteTaskInstanceData, DeleteTaskInstanceResponse, GetMappedTaskInstancesData, GetMappedTaskInstancesResponse, GetTaskInstanceDependenciesByMapIndexData, GetTaskInstanceDependenciesByMapIndexResponse, GetTaskInstanceDependenciesData, GetTaskInstanceDependenciesResponse, GetTaskInstanceTriesData, GetTaskInstanceTriesResponse, GetMappedTaskInstanceTriesData, GetMappedTaskInstanceTriesResponse, GetMappedTaskInstanceData, GetMappedTaskInstanceResponse, PatchTaskInstanceByMapIndexData, PatchTaskInstanceByMapIndexResponse, GetTaskInstancesData, GetTaskInstancesResponse, BulkTaskInstancesData, BulkTaskInstancesResponse, GetTaskInstancesBatchData, GetTaskInstancesBatchResponse, GetTaskInstanceTryDetailsData, GetTaskInstanceTryDetailsResponse, GetMappedTaskInstanceTryDetailsData, GetMappedTaskInstanceTryDetailsResponse, PostClearTaskInstancesData, PostClearTaskInstancesResponse, PatchTaskGroupInstancesData, PatchTaskGroupInstancesResponse, PatchTaskGroupInstancesDryRunData, PatchTaskGroupInstancesDryRunResponse, PatchTaskInstanceDryRunByMapIndexData, PatchTaskInstanceDryRunByMapIndexResponse, PatchTaskInstanceDryRunData, PatchTaskInstanceDryRunResponse, GetLogData, GetLogResponse, GetExternalLogUrlData, GetExternalLogUrlResponse, UpdateHitlDetailData, UpdateHitlDetailResponse, GetHitlDetailData, GetHitlDetailResponse, GetHitlDetailTryDetailData, GetHitlDetailTryDetailResponse, GetHitlDetailsData, GetHitlDetailsResponse, GetImportErrorData, GetImportErrorResponse, GetImportErrorsData, GetImportErrorsResponse, GetJobsData, GetJobsResponse, GetPluginsData, GetPluginsResponse, ImportErrorsResponse, DeletePoolData, DeletePoolResponse, GetPoolData, GetPoolResponse, PatchPoolData, PatchPoolResponse, GetPoolsData, GetPoolsResponse, PostPoolData, PostPoolResponse, BulkPoolsData, BulkPoolsResponse, GetProvidersData, GetProvidersResponse, ListAssetStateStoreData, ListAssetStateStoreResponse, ClearAssetStateStoreData, ClearAssetStateStoreResponse, GetAssetStateStoreData, GetAssetStateStoreResponse, SetAssetStateStoreData, SetAssetStateStoreResponse, DeleteAssetStateStoreData, DeleteAssetStateStoreResponse, ListTaskStateStoreData, ListTaskStateStoreResponse, ClearTaskStateStoreData, ClearTaskStateStoreResponse, GetTaskStateStoreData, GetTaskStateStoreResponse, SetTaskStateStoreData, SetTaskStateStoreResponse, PatchTaskStateStoreData, PatchTaskStateStoreResponse, DeleteTaskStateStoreData, DeleteTaskStateStoreResponse, GetXcomEntryData, GetXcomEntryResponse, UpdateXcomEntryData, UpdateXcomEntryResponse, DeleteXcomEntryData, DeleteXcomEntryResponse, GetXcomEntriesData, GetXcomEntriesResponse, CreateXcomEntryData, CreateXcomEntryResponse, GetTasksData, GetTasksResponse, GetTaskData, GetTaskResponse, DeleteVariableData, DeleteVariableResponse, GetVariableData, GetVariableResponse, PatchVariableData, PatchVariableResponse, GetVariablesData, GetVariablesResponse, PostVariableData, PostVariableResponse, BulkVariablesData, BulkVariablesResponse, ReparseDagFileData, ReparseDagFileResponse, GetDagVersionData, GetDagVersionResponse, GetDagVersionsData, GetDagVersionsResponse, GetHealthResponse, GetVersionResponse, LoginData, LoginResponse, LogoutResponse, GetAuthMenusResponse, GetCurrentUserInfoResponse, GenerateTokenData, GenerateTokenResponse2, GetPartitionedDagRunsData, GetPartitionedDagRunsResponse, GetPendingPartitionedDagRunData, GetPendingPartitionedDagRunResponse, GetDependenciesData, GetDependenciesResponse, HistoricalMetricsData, HistoricalMetricsResponse, DagStatsResponse2, GetDeadlinesData, GetDeadlinesResponse, GetDagDeadlineAlertsData, GetDagDeadlineAlertsResponse, StructureDataData, StructureDataResponse2, GetDagStructureData, GetDagStructureResponse, GetGridRunsData, GetGridRunsResponse, GetGridTiSummariesStreamData, GetGridTiSummariesStreamResponse, GetGanttDataData, GetGanttDataResponse, GetCalendarData, GetCalendarResponse, GetCalendarDeadlinesData, GetCalendarDeadlinesResponse, ListTeamsData, ListTeamsResponse } from './types.gen'; +import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData, GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse, GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData, CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse, GetAssetQueuedEventsData, GetAssetQueuedEventsResponse, DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData, GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse, DeleteDagAssetQueuedEventsData, DeleteDagAssetQueuedEventsResponse, GetDagAssetQueuedEventData, GetDagAssetQueuedEventResponse, DeleteDagAssetQueuedEventData, DeleteDagAssetQueuedEventResponse, NextRunAssetsData, NextRunAssetsResponse2, ListBackfillsData, ListBackfillsResponse, CreateBackfillData, CreateBackfillResponse, GetBackfillData, GetBackfillResponse, PauseBackfillData, PauseBackfillResponse, UnpauseBackfillData, UnpauseBackfillResponse, CancelBackfillData, CancelBackfillResponse, CreateBackfillDryRunData, CreateBackfillDryRunResponse, ListBackfillsUiData, ListBackfillsUiResponse, DeleteConnectionData, DeleteConnectionResponse, GetConnectionData, GetConnectionResponse, PatchConnectionData, PatchConnectionResponse, GetConnectionTestData, GetConnectionTestResponse, EnqueueConnectionTestData, EnqueueConnectionTestResponse, GetConnectionsData, GetConnectionsResponse, PostConnectionData, PostConnectionResponse, BulkConnectionsData, BulkConnectionsResponse, TestConnectionData, TestConnectionResponse, CreateDefaultConnectionsResponse, HookMetaDataResponse, GetDagRunData, GetDagRunResponse, DeleteDagRunData, DeleteDagRunResponse, PatchDagRunData, PatchDagRunResponse, BulkDagRunsData, BulkDagRunsResponse, GetDagRunsData, GetDagRunsResponse, TriggerDagRunData, TriggerDagRunResponse, GetUpstreamAssetEventsData, GetUpstreamAssetEventsResponse, ClearDagRunData, ClearDagRunResponse, WaitDagRunUntilFinishedData, WaitDagRunUntilFinishedResponse, GetListDagRunsBatchData, GetListDagRunsBatchResponse, ClearDagRunsData, ClearDagRunsResponse, ClearDagRunPartitionsData, ClearDagRunPartitionsResponse, GetDagRunStatsData, GetDagRunStatsResponse, GetDagSourceData, GetDagSourceResponse, GetDagStatsData, GetDagStatsResponse, GetConfigData, GetConfigResponse, GetConfigValueData, GetConfigValueResponse, GetConfigsResponse, ListDagWarningsData, ListDagWarningsResponse, GetDagsData, GetDagsResponse, PatchDagsData, PatchDagsResponse, GetDagData, GetDagResponse, PatchDagData, PatchDagResponse, DeleteDagData, DeleteDagResponse, GetDagDetailsData, GetDagDetailsResponse, FavoriteDagData, FavoriteDagResponse, UnfavoriteDagData, UnfavoriteDagResponse, GetDagTagsData, GetDagTagsResponse, GetDagsUiData, GetDagsUiResponse, GetLatestRunInfoData, GetLatestRunInfoResponse, GetDagRunStateCountsUiData, GetDagRunStateCountsUiResponse, GetLatestRunTaskInstanceStateCountsUiData, GetLatestRunTaskInstanceStateCountsUiResponse, GetEventLogData, GetEventLogResponse, GetEventLogsData, GetEventLogsResponse, GetExtraLinksData, GetExtraLinksResponse, GetTaskInstanceData, GetTaskInstanceResponse, PatchTaskInstanceData, PatchTaskInstanceResponse, DeleteTaskInstanceData, DeleteTaskInstanceResponse, GetMappedTaskInstancesData, GetMappedTaskInstancesResponse, GetTaskInstanceDependenciesByMapIndexData, GetTaskInstanceDependenciesByMapIndexResponse, GetTaskInstanceDependenciesData, GetTaskInstanceDependenciesResponse, GetTaskInstanceTriesData, GetTaskInstanceTriesResponse, GetMappedTaskInstanceTriesData, GetMappedTaskInstanceTriesResponse, GetMappedTaskInstanceData, GetMappedTaskInstanceResponse, PatchTaskInstanceByMapIndexData, PatchTaskInstanceByMapIndexResponse, GetTaskInstancesData, GetTaskInstancesResponse, BulkTaskInstancesData, BulkTaskInstancesResponse, GetTaskInstancesBatchData, GetTaskInstancesBatchResponse, GetTaskInstanceTryDetailsData, GetTaskInstanceTryDetailsResponse, GetMappedTaskInstanceTryDetailsData, GetMappedTaskInstanceTryDetailsResponse, PostClearTaskInstancesData, PostClearTaskInstancesResponse, PatchTaskGroupInstancesData, PatchTaskGroupInstancesResponse, PatchTaskGroupInstancesDryRunData, PatchTaskGroupInstancesDryRunResponse, PatchTaskInstanceDryRunByMapIndexData, PatchTaskInstanceDryRunByMapIndexResponse, PatchTaskInstanceDryRunData, PatchTaskInstanceDryRunResponse, GetLogData, GetLogResponse, GetExternalLogUrlData, GetExternalLogUrlResponse, UpdateHitlDetailData, UpdateHitlDetailResponse, GetHitlDetailData, GetHitlDetailResponse, GetHitlDetailTryDetailData, GetHitlDetailTryDetailResponse, GetHitlDetailsData, GetHitlDetailsResponse, GetImportErrorData, GetImportErrorResponse, GetImportErrorsData, GetImportErrorsResponse, GetJobsData, GetJobsResponse, GetPluginsData, GetPluginsResponse, ImportErrorsResponse, DeletePoolData, DeletePoolResponse, GetPoolData, GetPoolResponse, PatchPoolData, PatchPoolResponse, GetPoolsData, GetPoolsResponse, PostPoolData, PostPoolResponse, BulkPoolsData, BulkPoolsResponse, GetProvidersData, GetProvidersResponse, ListAssetStateStoreData, ListAssetStateStoreResponse, ClearAssetStateStoreData, ClearAssetStateStoreResponse, GetAssetStateStoreData, GetAssetStateStoreResponse, SetAssetStateStoreData, SetAssetStateStoreResponse, DeleteAssetStateStoreData, DeleteAssetStateStoreResponse, ListTaskStateStoreData, ListTaskStateStoreResponse, ClearTaskStateStoreData, ClearTaskStateStoreResponse, GetTaskStateStoreData, GetTaskStateStoreResponse, SetTaskStateStoreData, SetTaskStateStoreResponse, PatchTaskStateStoreData, PatchTaskStateStoreResponse, DeleteTaskStateStoreData, DeleteTaskStateStoreResponse, GetXcomEntryData, GetXcomEntryResponse, UpdateXcomEntryData, UpdateXcomEntryResponse, DeleteXcomEntryData, DeleteXcomEntryResponse, GetXcomEntriesData, GetXcomEntriesResponse, CreateXcomEntryData, CreateXcomEntryResponse, GetTasksData, GetTasksResponse, GetTaskData, GetTaskResponse, DeleteVariableData, DeleteVariableResponse, GetVariableData, GetVariableResponse, PatchVariableData, PatchVariableResponse, GetVariablesData, GetVariablesResponse, PostVariableData, PostVariableResponse, BulkVariablesData, BulkVariablesResponse, ReparseDagFileData, ReparseDagFileResponse, GetDagVersionData, GetDagVersionResponse, GetDagVersionsData, GetDagVersionsResponse, GetHealthResponse, GetVersionResponse, LoginData, LoginResponse, LogoutResponse, GetAuthMenusResponse, GetCurrentUserInfoResponse, GenerateTokenData, GenerateTokenResponse2, GetPartitionedDagRunsData, GetPartitionedDagRunsResponse, GetPendingPartitionedDagRunData, GetPendingPartitionedDagRunResponse, GetDependenciesData, GetDependenciesResponse, HistoricalMetricsData, HistoricalMetricsResponse, DagStatsResponse2, GetDeadlinesData, GetDeadlinesResponse, GetDagDeadlineAlertsData, GetDagDeadlineAlertsResponse, StructureDataData, StructureDataResponse2, GetDagStructureData, GetDagStructureResponse, GetGridRunsData, GetGridRunsResponse, GetGridTiSummariesStreamData, GetGridTiSummariesStreamResponse, GetGanttDataData, GetGanttDataResponse, GetCalendarData, GetCalendarResponse, GetCalendarDeadlinesData, GetCalendarDeadlinesResponse, ListTeamsData, ListTeamsResponse } from './types.gen'; export class AssetService { /** @@ -2012,6 +2012,29 @@ export class DagService { }); } + /** + * Get Latest Run Task Instance State Counts + * Return task-instance state counts for each Dag's latest run, for the Dag list page. + * + * Dags without any run are omitted from the response. + * @param data The data for the request. + * @param data.dagIds + * @returns DAGsLatestRunTaskInstanceStateCountsCollectionResponse Successful Response + * @throws ApiError + */ + public static getLatestRunTaskInstanceStateCountsUi(data: GetLatestRunTaskInstanceStateCountsUiData): CancelablePromise { + return __request(OpenAPI, { + method: 'GET', + url: '/ui/dags/latest_run_task_instance_state_counts', + query: { + dag_ids: data.dagIds + }, + errors: { + 422: 'Validation Error' + } + }); + } + } export class EventLogService { 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 e1bf23b160208..50e3a187e5110 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 @@ -2273,6 +2273,20 @@ export type ConnectionHookMetaData = { } | null; }; +/** + * Task-instance state counts for a Dag's latest run. + * + * ``state_counts`` only carries states present in the run; task instances without a + * state yet are keyed as ``no_status``. + */ +export type DAGLatestRunTaskInstanceStateCountsResponse = { + dag_id: string; + run_id: string; + state_counts: { + [key: string]: (number); + }; +}; + /** * DAG Run serializer for responses. */ @@ -2362,6 +2376,13 @@ export type DAGWithLatestDagRunsResponse = { readonly file_token: string; }; +/** + * Collection of per-Dag latest-run task-instance state counts for the Dag list page. + */ +export type DAGsLatestRunTaskInstanceStateCountsCollectionResponse = { + dags: Array; +}; + /** * Collection of per-Dag DagRun-state counts for the Dag list page. */ @@ -3532,6 +3553,12 @@ export type GetDagRunStateCountsUiData = { export type GetDagRunStateCountsUiResponse = DAGsRunStateCountsCollectionResponse; +export type GetLatestRunTaskInstanceStateCountsUiData = { + dagIds: Array<(string)>; +}; + +export type GetLatestRunTaskInstanceStateCountsUiResponse = DAGsLatestRunTaskInstanceStateCountsCollectionResponse; + export type GetEventLogData = { eventLogId: number; }; @@ -6416,6 +6443,21 @@ export type $OpenApiTs = { }; }; }; + '/ui/dags/latest_run_task_instance_state_counts': { + get: { + req: GetLatestRunTaskInstanceStateCountsUiData; + res: { + /** + * Successful Response + */ + 200: DAGsLatestRunTaskInstanceStateCountsCollectionResponse; + /** + * Validation Error + */ + 422: HTTPValidationError; + }; + }; + }; '/api/v2/eventLogs/{event_log_id}': { get: { req: GetEventLogData; diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json index c8b195c93b8cd..a336824c15d58 100644 --- a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json +++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json @@ -24,6 +24,11 @@ }, "runIdPatternFilter": "Search Dag Runs" }, + "latestRunTaskStateCounts": { + "label": "Latest run tasks", + "loading": "Loading latest run task state counts", + "tooltip": "{{state}}: {{formattedCount}} — click to filter tasks" + }, "ownerLink": "Owner link for {{owner}}", "runAndTaskActions": { "affectedTasks": { diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.test.tsx b/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.test.tsx index ad163ae19f1e6..1da462e349dbc 100644 --- a/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.test.tsx +++ b/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.test.tsx @@ -58,9 +58,19 @@ const GMTWrapper = ({ children }: PropsWithChildren) => ( // 4 extra "state-badge" testids and break getByTestId assertions in tests // that target the latest-run badge. const renderCard = (dag: DAGWithLatestDagRunsResponse) => - render(, { - wrapper: GMTWrapper, - }); + render( + , + { + wrapper: GMTWrapper, + }, + ); const mockDag = { allowed_run_types: null, diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.tsx b/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.tsx index cd320d417d0ea..4678af9939b22 100644 --- a/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.tsx +++ b/airflow-core/src/airflow/ui/src/pages/DagsList/DagCard.tsx @@ -19,7 +19,10 @@ import { Box, Flex, Grid, GridItem, HStack, Spinner } from "@chakra-ui/react"; import { useTranslation } from "react-i18next"; -import type { DAGWithLatestDagRunsResponse } from "openapi/requests/types.gen"; +import type { + DAGLatestRunTaskInstanceStateCountsResponse, + DAGWithLatestDagRunsResponse, +} from "openapi/requests/types.gen"; import { DeleteDagButton } from "src/components/DagActions/DeleteDagButton"; import { FavoriteDagButton } from "src/components/DagActions/FavoriteDagButton"; import DagRunInfo from "src/components/DagRunInfo"; @@ -32,17 +35,27 @@ import { isStatePending, useAutoRefresh } from "src/utils"; import { DagRunStateCounts } from "./DagRunStateCounts"; import { DagTags } from "./DagTags"; +import { LatestRunTaskStateCounts } from "./LatestRunTaskStateCounts"; import { RecentRuns } from "./RecentRuns"; import { Schedule } from "./Schedule"; type Props = { readonly dag: DAGWithLatestDagRunsResponse; + readonly latestRunTaskStateCounts: DAGLatestRunTaskInstanceStateCountsResponse | undefined; + readonly latestRunTaskStateCountsLoading: boolean; readonly runStateCounts: Record | undefined; readonly runStateCountsLoading: boolean; readonly stateCountLimit: number | undefined; }; -export const DagCard = ({ dag, runStateCounts, runStateCountsLoading, stateCountLimit }: Props) => { +export const DagCard = ({ + dag, + latestRunTaskStateCounts, + latestRunTaskStateCountsLoading, + runStateCounts, + runStateCountsLoading, + stateCountLimit, +}: Props) => { const { t: translate } = useTranslation(["common", "dag"]); const [latestRun] = dag.latest_dag_runs; @@ -135,6 +148,13 @@ export const DagCard = ({ dag, runStateCounts, runStateCountsLoading, stateCount stateCountLimit={stateCountLimit} /> + + + ); diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx b/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx index e9c5215cc5a06..66eff61e50cf3 100644 --- a/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx +++ b/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx @@ -22,7 +22,11 @@ import { useTranslation } from "react-i18next"; import { useSearchParams } from "react-router-dom"; import { useLocalStorage } from "usehooks-ts"; -import type { DagRunState, DAGWithLatestDagRunsResponse } from "openapi/requests/types.gen"; +import type { + DAGLatestRunTaskInstanceStateCountsResponse, + DagRunState, + DAGWithLatestDagRunsResponse, +} from "openapi/requests/types.gen"; import { DeleteDagButton } from "src/components/DagActions/DeleteDagButton"; import { FavoriteDagButton } from "src/components/DagActions/FavoriteDagButton"; import DagRunInfo from "src/components/DagRunInfo"; @@ -42,6 +46,7 @@ import { DagsLayout } from "src/layouts/DagsLayout"; import { useConfig } from "src/queries/useConfig"; import { useDagRunStateCounts } from "src/queries/useDagRunStateCounts"; import { useDags } from "src/queries/useDags"; +import { useLatestRunTaskStateCounts } from "src/queries/useLatestRunTaskStateCounts"; import { useDocumentTitle } from "src/utils"; import { DagImportErrors } from "../Dashboard/Stats/DagImportErrors"; @@ -49,6 +54,7 @@ import { DagCard } from "./DagCard"; import { DagRunStateCounts } from "./DagRunStateCounts"; import { DagTags } from "./DagTags"; import { DagsFilters } from "./DagsFilters"; +import { LatestRunTaskStateCounts } from "./LatestRunTaskStateCounts"; import { Schedule } from "./Schedule"; import { SortSelect } from "./SortSelect"; import { useTagFilter } from "./useTagFilter"; @@ -59,9 +65,15 @@ type RunStateCountsContext = { readonly stateCountLimit: number | undefined; }; +type LatestRunTaskStateCountsContext = { + readonly entriesByDag: Record; + readonly isLoading: boolean; +}; + const createColumns = ( translate: (key: string, options?: Record) => string, runStateContext: RunStateCountsContext, + taskStateContext: LatestRunTaskStateCountsContext, ): Array> => [ { accessorKey: "is_paused", @@ -145,6 +157,19 @@ const createColumns = ( enableSorting: false, header: () => translate("dags:runStateCounts.label"), }, + { + accessorKey: "latest_run_task_state_counts", + cell: ({ row: { original } }) => ( + + ), + enableSorting: false, + header: () => translate("dags:latestRunTaskStateCounts.label"), + }, { accessorKey: "tags", cell: ({ @@ -204,10 +229,15 @@ const { PAUSED, }: SearchParamsKeysType = SearchParamsKeys; -const createCardDef = (runStateContext: RunStateCountsContext): CardDef => ({ +const createCardDef = ( + runStateContext: RunStateCountsContext, + taskStateContext: LatestRunTaskStateCountsContext, +): CardDef => ({ card: ({ row }) => ( { stateCountLimit: runStateCountsData?.state_count_limit, }; - const columns = createColumns(translate, runStateContext); - const cardDef = createCardDef(runStateContext); + const { data: taskStateCountsData, isLoading: taskStateCountsLoading } = useLatestRunTaskStateCounts({ + dagIds: data?.dags.map((dag) => dag.dag_id) ?? [], + dags: data?.dags, + }); + const taskStateContext: LatestRunTaskStateCountsContext = { + entriesByDag: Object.fromEntries((taskStateCountsData?.dags ?? []).map((entry) => [entry.dag_id, entry])), + isLoading: taskStateCountsLoading, + }; + + const columns = createColumns(translate, runStateContext, taskStateContext); + const cardDef = createCardDef(runStateContext, taskStateContext); const handleSortChange = ({ value }: SelectValueChangeDetails>) => { setTableURLState({ diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.test.tsx b/airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.test.tsx new file mode 100644 index 0000000000000..f443712ed0caa --- /dev/null +++ b/airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.test.tsx @@ -0,0 +1,91 @@ +/*! + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +import "@testing-library/jest-dom/vitest"; +import { render, screen } from "@testing-library/react"; +import { MemoryRouter } from "react-router-dom"; +import { describe, expect, it } from "vitest"; + +import type { DAGLatestRunTaskInstanceStateCountsResponse } from "openapi/requests/types.gen"; +import { BaseWrapper } from "src/utils/Wrapper"; + +import "../../i18n/config"; +import { LatestRunTaskStateCounts } from "./LatestRunTaskStateCounts"; + +const renderCounts = ( + entry: DAGLatestRunTaskInstanceStateCountsResponse | undefined, + options: { isLoading?: boolean } = {}, +) => + render(, { + wrapper: ({ children }) => ( + + {children} + + ), + }); + +const makeEntry = (stateCounts: Record): DAGLatestRunTaskInstanceStateCountsResponse => ({ + dag_id: "my_dag", + run_id: "run_1", + state_counts: stateCounts, +}); + +describe("LatestRunTaskStateCounts", () => { + it("renders skeleton placeholders while loading", () => { + renderCounts(undefined, { isLoading: true }); + expect(screen.getByTestId("latest-run-task-state-counts-loading-my_dag")).toBeInTheDocument(); + expect(screen.queryByTestId("latest-run-task-state-counts-my_dag")).toBeNull(); + }); + + it("renders nothing for a Dag without runs", () => { + renderCounts(undefined); + expect(screen.queryByTestId("latest-run-task-state-counts-my_dag")).toBeNull(); + }); + + it("renders one clickable badge per present state, linking to the run's filtered task list", () => { + renderCounts(makeEntry({ failed: 2, running: 1, success: 7 })); + expect(screen.getByTestId("latest-run-task-state-counts-my_dag")).toBeInTheDocument(); + + const failedLink = screen.getByTestId("latest-run-task-state-count-failed-my_dag"); + const runningLink = screen.getByTestId("latest-run-task-state-count-running-my_dag"); + const successLink = screen.getByTestId("latest-run-task-state-count-success-my_dag"); + + expect(failedLink).toHaveAttribute("href", "/dags/my_dag/runs/run_1?task_state=failed"); + expect(runningLink).toHaveAttribute("href", "/dags/my_dag/runs/run_1?task_state=running"); + expect(successLink).toHaveAttribute("href", "/dags/my_dag/runs/run_1?task_state=success"); + + expect(failedLink).toHaveTextContent("2"); + expect(runningLink).toHaveTextContent("1"); + expect(successLink).toHaveTextContent("7"); + }); + + it("omits absent and zero-count states instead of rendering empty badges", () => { + renderCounts(makeEntry({ queued: 0, success: 3 })); + expect(screen.getByTestId("latest-run-task-state-count-success-my_dag")).toBeInTheDocument(); + expect(screen.queryByTestId("latest-run-task-state-count-queued-my_dag")).toBeNull(); + expect(screen.queryByTestId("latest-run-task-state-count-failed-my_dag")).toBeNull(); + }); + + it("maps no_status to the 'none' task filter value", () => { + renderCounts(makeEntry({ no_status: 2, success: 1 })); + const noStatusLink = screen.getByTestId("latest-run-task-state-count-no_status-my_dag"); + + expect(noStatusLink).toHaveAttribute("href", "/dags/my_dag/runs/run_1?task_state=none"); + expect(noStatusLink).toHaveTextContent("2"); + }); +}); diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.tsx b/airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.tsx new file mode 100644 index 0000000000000..487b9888313d7 --- /dev/null +++ b/airflow-core/src/airflow/ui/src/pages/DagsList/LatestRunTaskStateCounts.tsx @@ -0,0 +1,96 @@ +/*! + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +import { HStack, Skeleton } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; + +import type { + DAGLatestRunTaskInstanceStateCountsResponse, + TaskInstanceState, +} from "openapi/requests/types.gen"; +import { StateBadge } from "src/components/StateBadge"; +import { RouterLink, Tooltip } from "src/components/ui"; +import { SearchParamsKeys } from "src/constants/searchParams"; +import { sortStateEntries } from "src/utils"; + +type Props = { + readonly compact?: boolean; + readonly dagId: string; + readonly entry: DAGLatestRunTaskInstanceStateCountsResponse | undefined; + readonly isLoading: boolean; +}; + +export const LatestRunTaskStateCounts = ({ compact = false, dagId, entry, isLoading }: Props) => { + const { t: translate } = useTranslation(["dags", "common"]); + const gap = compact ? 0.5 : 1; + const fontSize = compact ? "xs" : "sm"; + + if (isLoading) { + // Badges are dynamic (only states present in the run), so the final count is + // unknown while loading; three pills approximate a typical row without reflow. + return ( + + {[1, 2, 3].map((idx) => ( + + ))} + + ); + } + + if (entry === undefined) { + return undefined; + } + + const stateEntries = sortStateEntries(entry.state_counts); + + return ( + + {stateEntries.map(([state, count]) => { + const translatedState = translate(`common:states.${state}` as const); + const tooltipContent = translate("latestRunTaskStateCounts.tooltip", { + formattedCount: `${count}`, + state: translatedState, + }); + // Task instances without a state are keyed "no_status"; the task list + // filters them with the "none" value and StateBadge renders them as null. + const filterValue = state === "no_status" ? "none" : state; + + return ( + + + + {count} + + + + ); + })} + + ); +}; diff --git a/airflow-core/src/airflow/ui/src/queries/useLatestRunTaskStateCounts.tsx b/airflow-core/src/airflow/ui/src/queries/useLatestRunTaskStateCounts.tsx new file mode 100644 index 0000000000000..2c32c437c1816 --- /dev/null +++ b/airflow-core/src/airflow/ui/src/queries/useLatestRunTaskStateCounts.tsx @@ -0,0 +1,45 @@ +/*! + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +import { useDagServiceGetLatestRunTaskInstanceStateCountsUi } from "openapi/queries"; +import type { DAGWithLatestDagRunsResponse } from "openapi/requests/types.gen"; +import { isStatePending, useAutoRefresh } from "src/utils"; + +export const useLatestRunTaskStateCounts = ({ + dagIds, + dags, +}: { + readonly dagIds: ReadonlyArray; + // Refresh predicate is derived from useDags' data so the counts query doesn't + // need to be loaded before it knows whether to poll — avoids a chicken-and-egg. + readonly dags: ReadonlyArray | undefined; +}) => { + const refetchInterval = useAutoRefresh({}); + const hasPendingRun = + dags?.some((dag) => !dag.is_paused && dag.latest_dag_runs.some((run) => isStatePending(run.state))) ?? + false; + + // Stable key: sort the dag_ids so pagination/sort order changes don't churn the cache. + const sortedDagIds = [...dagIds].sort(); + + return useDagServiceGetLatestRunTaskInstanceStateCountsUi({ dagIds: sortedDagIds }, undefined, { + enabled: sortedDagIds.length > 0, + placeholderData: (prev) => prev, + refetchInterval: hasPendingRun ? refetchInterval : false, + }); +}; diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py index 4e561fd343664..ddd8dc0c75795 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py @@ -32,6 +32,7 @@ from airflow.models.dag import DagModel, DagTag from airflow.models.dag_favorite import DagFavorite from airflow.models.hitl import HITLDetail +from airflow.providers.standard.operators.empty import EmptyOperator from airflow.sdk.timezone import utcnow from airflow.utils.session import NEW_SESSION, provide_session from airflow.utils.state import DagRunState, TaskInstanceState @@ -587,3 +588,125 @@ def test_should_response_401(self, unauthenticated_test_client): def test_should_response_403(self, unauthorized_test_client): response = unauthorized_test_client.get("/dags/run_state_counts", params={"dag_ids": [DAG1_ID]}) assert response.status_code == 403 + + +TI_COUNTS_DAG_ID = "test_dag_latest_run_ti_counts" +LATEST_RUN_TI_COUNTS_ENDPOINT = "/dags/latest_run_task_instance_state_counts" + + +class TestGetLatestRunTaskInstanceStateCounts(TestPublicDagEndpoint): + """Tests for ``GET /ui/dags/latest_run_task_instance_state_counts``.""" + + @pytest.fixture(autouse=True) + def seed_runs_with_task_instances(self, setup, dag_maker, session) -> None: + # A dedicated Dag with two runs: only the latest run (by run_after) may be + # counted. Its four tasks cover three distinct states plus the null-state + # ("no_status") case; the older run is all-success noise the endpoint must skip. + with dag_maker(TI_COUNTS_DAG_ID, schedule=None, session=session): + for idx in range(4): + EmptyOperator(task_id=f"task_{idx}") + + base = utcnow() - pendulum.duration(days=1) + older_run = dag_maker.create_dagrun( + run_id="older_run", state=DagRunState.SUCCESS, logical_date=base, run_after=base + ) + for ti in older_run.task_instances: + ti.state = TaskInstanceState.SUCCESS + latest = base + pendulum.duration(hours=1) + latest_run = dag_maker.create_dagrun( + run_id="latest_run", state=DagRunState.RUNNING, logical_date=latest, run_after=latest + ) + latest_states = [ + TaskInstanceState.SUCCESS, + TaskInstanceState.FAILED, + TaskInstanceState.RUNNING, + None, + ] + tis = sorted(latest_run.task_instances, key=lambda ti: ti.task_id) + for ti, state in zip(tis, latest_states, strict=True): + ti.state = state + dag_maker.sync_dagbag_to_db() + session.commit() + + @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + def test_counts_only_the_latest_run(self, test_client): + response = test_client.get(LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": [TI_COUNTS_DAG_ID]}) + assert response.status_code == 200 + # Only latest-run states may appear: the older all-success run must not leak in + # (it would push success to 5), and unset states surface as "no_status". + assert response.json()["dags"] == [ + { + "dag_id": TI_COUNTS_DAG_ID, + "run_id": "latest_run", + "state_counts": {"success": 1, "failed": 1, "running": 1, "no_status": 1}, + } + ] + + @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + def test_omits_dags_without_runs(self, test_client): + # DAG2 exists but has no runs in the parent fixture; it must simply be absent. + response = test_client.get( + LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": [TI_COUNTS_DAG_ID, DAG2_ID]} + ) + assert response.status_code == 200 + dag_ids = [entry["dag_id"] for entry in response.json()["dags"]] + assert dag_ids == [TI_COUNTS_DAG_ID] + + @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + def test_counts_multiple_dags_independently(self, test_client): + response = test_client.get( + LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": [TI_COUNTS_DAG_ID, DAG1_ID]} + ) + assert response.status_code == 200 + by_dag = {entry["dag_id"]: entry for entry in response.json()["dags"]} + assert set(by_dag) == {TI_COUNTS_DAG_ID, DAG1_ID} + # DAG1's only run comes from the parent fixture: a single task instance that was + # never scheduled, so it surfaces under "no_status". + assert by_dag[DAG1_ID]["state_counts"] == {"no_status": 1} + assert by_dag[TI_COUNTS_DAG_ID]["state_counts"] == { + "success": 1, + "failed": 1, + "running": 1, + "no_status": 1, + } + + def test_deduplicates_dag_ids(self, test_client): + response = test_client.get( + LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": [TI_COUNTS_DAG_ID, TI_COUNTS_DAG_ID]} + ) + assert response.status_code == 200 + dag_ids = [entry["dag_id"] for entry in response.json()["dags"]] + assert dag_ids == [TI_COUNTS_DAG_ID] + + def test_rejects_too_many_dag_ids(self, test_client): + # The page never sends more than maximum_page_limit dag_ids; a direct call with a + # larger list is rejected so the per-Dag UNION ALL width stays bounded. + too_many = [f"dag_{idx}" for idx in range(conf.getint("api", "maximum_page_limit") + 1)] + response = test_client.get(LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": too_many}) + assert response.status_code == 422 + + @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + def test_permission_filter_hides_disallowed_dags(self, test_client): + with mock.patch.object( + SimpleAuthManager, + "get_authorized_dag_ids", + return_value={TI_COUNTS_DAG_ID}, + ): + response = test_client.get( + LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": [TI_COUNTS_DAG_ID, DAG1_ID]} + ) + assert response.status_code == 200 + dag_ids = [entry["dag_id"] for entry in response.json()["dags"]] + assert dag_ids == [TI_COUNTS_DAG_ID] + + def test_should_response_401(self, unauthenticated_test_client): + response = unauthenticated_test_client.get( + LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": [TI_COUNTS_DAG_ID]} + ) + assert response.status_code == 401 + + def test_should_response_403(self, unauthorized_test_client): + response = unauthorized_test_client.get( + LATEST_RUN_TI_COUNTS_ENDPOINT, params={"dag_ids": [TI_COUNTS_DAG_ID]} + ) + assert response.status_code == 403