Skip to content

Commit 063d5f1

Browse files
authored
UI: Add team column and filter to Dag list page (#70028)
When multi-team mode is enabled, operators need to identify which team owns each Dag and filter the list by team. This adds: - A "Team" column (after tags) in the Dag list table that links to a filtered view of that team's Dags - A team multi-select filter in the header controls - Backend support: `team_name` field in the DAGWithLatestDagRunsResponse and a `teams` query parameter for filtering Dags by team (via bundle association) Both the column and filter are conditionally rendered only when the `multi_team` configuration is enabled.
1 parent 15a4478 commit 063d5f1

19 files changed

Lines changed: 232 additions & 14 deletions

File tree

airflow-core/src/airflow/api_fastapi/common/parameters.py

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,7 @@
6262
from airflow.models.dag import DagModel, DagTag
6363
from airflow.models.dag_favorite import DagFavorite
6464
from airflow.models.dag_version import DagVersion
65+
from airflow.models.dagbundle import DagBundleModel
6566
from airflow.models.dagrun import DagRun
6667
from airflow.models.errors import ParseImportError
6768
from airflow.models.hitl import HITLDetail
@@ -1020,6 +1021,29 @@ def depends(cls, owners: list[str] = Query(default_factory=list)) -> _OwnersFilt
10201021
return cls().set_value(owners)
10211022

10221023

1024+
class _TeamsFilter(BaseParam[list[str]]):
1025+
"""Filter Dags by team name (via bundle association)."""
1026+
1027+
def to_orm(self, select: Select) -> Select:
1028+
if self.skip_none is False:
1029+
raise ValueError(f"Cannot set 'skip_none' to False on a {type(self)}")
1030+
1031+
if not self.value:
1032+
return select
1033+
1034+
from airflow.models.team import Team
1035+
1036+
return select.where(
1037+
DagModel.bundle_name.in_(
1038+
sql_select(DagBundleModel.name).join(DagBundleModel.teams).where(Team.name.in_(self.value))
1039+
)
1040+
)
1041+
1042+
@classmethod
1043+
def depends(cls, teams: list[str] = Query(default_factory=list)) -> _TeamsFilter:
1044+
return cls().set_value(teams)
1045+
1046+
10231047
def _safe_parse_datetime(date_to_check: str) -> datetime:
10241048
"""
10251049
Parse datetime and raise error for invalid dates.
@@ -1268,6 +1292,7 @@ def depends_float(
12681292
]
12691293
QueryTagsFilter = Annotated[_TagsFilter, Depends(_TagsFilter.depends)]
12701294
QueryOwnersFilter = Annotated[_OwnersFilter, Depends(_OwnersFilter.depends)]
1295+
QueryTeamsFilter = Annotated[_TeamsFilter, Depends(_TeamsFilter.depends)]
12711296

12721297

12731298
class _HasAssetScheduleFilter(BaseParam[bool]):

airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/dags.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ class DAGWithLatestDagRunsResponse(DAGResponse):
3232
latest_dag_runs: list[DAGRunLightResponse]
3333
pending_actions: list[HITLDetail]
3434
is_favorite: bool
35+
team_name: str | None = None
3536

3637

3738
class DAGWithLatestDagRunsCollectionResponse(BaseModel):

airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -338,6 +338,14 @@ paths:
338338
items:
339339
type: string
340340
title: Owners
341+
- name: teams
342+
in: query
343+
required: false
344+
schema:
345+
type: array
346+
items:
347+
type: string
348+
title: Teams
341349
- name: dag_ids
342350
in: query
343351
required: false
@@ -2629,6 +2637,11 @@ components:
26292637
is_favorite:
26302638
type: boolean
26312639
title: Is Favorite
2640+
team_name:
2641+
anyOf:
2642+
- type: string
2643+
- type: 'null'
2644+
title: Team Name
26322645
is_backfillable:
26332646
type: boolean
26342647
title: Is Backfillable

airflow-core/src/airflow/api_fastapi/core_api/routes/ui/dags.py

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@
5151
QueryPausedFilter,
5252
QueryPendingActionsFilter,
5353
QueryTagsFilter,
54+
QueryTeamsFilter,
5455
SortParam,
5556
filter_param_factory,
5657
)
@@ -96,6 +97,7 @@ def get_dags(
9697
offset: QueryOffset,
9798
tags: QueryTagsFilter,
9899
owners: QueryOwnersFilter,
100+
teams: QueryTeamsFilter,
99101
dag_ids: Annotated[
100102
FilterParam[list[str] | None],
101103
Depends(filter_param_factory(DagModel.dag_id, list[str] | None, FilterOptionEnum.IN, "dag_ids")),
@@ -157,6 +159,7 @@ def get_dags(
157159
dag_display_name_prefix_pattern,
158160
tags,
159161
owners,
162+
teams,
160163
last_dag_run_state,
161164
dag_run_state,
162165
is_favorite,
@@ -232,6 +235,13 @@ def get_dags(
232235
for dag_id, hitl_detail in pending_actions:
233236
pending_actions_by_dag_id[dag_id].append(hitl_detail)
234237

238+
# Fetch team names when multi-team is enabled
239+
team_names_by_dag_id: dict[str, str | None] = {}
240+
if conf.getboolean("core", "multi_team") and dags:
241+
team_names_by_dag_id = DagModel.get_dag_id_to_team_name_mapping(
242+
[dag.dag_id for dag in dags], session=session
243+
)
244+
235245
# aggregate rows by dag_id
236246
# Build the dict dynamically from DAGResponse.model_fields so that new fields
237247
# added to DAGResponse are picked up automatically without code changes here.
@@ -249,6 +259,7 @@ def get_dags(
249259
"latest_dag_runs": [],
250260
"pending_actions": pending_actions_by_dag_id[dag.dag_id],
251261
"is_favorite": dag.dag_id in favorite_dag_ids,
262+
"team_name": team_names_by_dag_id.get(dag.dag_id),
252263
}
253264
)
254265
dag_runs_by_dag_id[dag.dag_id] = DAGWithLatestDagRunsResponse.model_validate(dag_data)

airflow-core/src/airflow/ui/openapi-gen/queries/common.ts

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -339,7 +339,7 @@ export const UseDagServiceGetDagTagsKeyFn = ({ limit, offset, orderBy, tagNamePa
339339
export type DagServiceGetDagsUiDefaultResponse = Awaited<ReturnType<typeof DagService.getDagsUi>>;
340340
export type DagServiceGetDagsUiQueryResult<TData = DagServiceGetDagsUiDefaultResponse, TError = unknown> = UseQueryResult<TData, TError>;
341341
export const useDagServiceGetDagsUiKey = "DagServiceGetDagsUi";
342-
export const UseDagServiceGetDagsUiKeyFn = ({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }: {
342+
export const UseDagServiceGetDagsUiKeyFn = ({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }: {
343343
assetDependency?: string;
344344
bundleName?: string;
345345
bundleVersion?: string;
@@ -363,7 +363,8 @@ export const UseDagServiceGetDagsUiKeyFn = ({ assetDependency, bundleName, bundl
363363
paused?: boolean;
364364
tags?: string[];
365365
tagsMatchMode?: "any" | "all";
366-
} = {}, queryKey?: Array<unknown>) => [useDagServiceGetDagsUiKey, ...(queryKey ?? [{ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }])];
366+
teams?: string[];
367+
} = {}, queryKey?: Array<unknown>) => [useDagServiceGetDagsUiKey, ...(queryKey ?? [{ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }])];
367368
export type DagServiceGetLatestRunInfoDefaultResponse = Awaited<ReturnType<typeof DagService.getLatestRunInfo>>;
368369
export type DagServiceGetLatestRunInfoQueryResult<TData = DagServiceGetLatestRunInfoDefaultResponse, TError = unknown> = UseQueryResult<TData, TError>;
369370
export const useDagServiceGetLatestRunInfoKey = "DagServiceGetLatestRunInfo";

airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -675,6 +675,7 @@ export const ensureUseDagServiceGetDagTagsData = (queryClient: QueryClient, { li
675675
* @param data.tags
676676
* @param data.tagsMatchMode
677677
* @param data.owners
678+
* @param data.teams
678679
* @param data.dagIds
679680
* @param data.dagIdPattern SQL LIKE expression — use `%` / `_` wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 | dag2`). Regular expressions are **not** supported.
680681
*
@@ -699,7 +700,7 @@ export const ensureUseDagServiceGetDagTagsData = (queryClient: QueryClient, { li
699700
* @returns DAGWithLatestDagRunsCollectionResponse Successful Response
700701
* @throws ApiError
701702
*/
702-
export const ensureUseDagServiceGetDagsUiData = (queryClient: QueryClient, { assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }: {
703+
export const ensureUseDagServiceGetDagsUiData = (queryClient: QueryClient, { assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }: {
703704
assetDependency?: string;
704705
bundleName?: string;
705706
bundleVersion?: string;
@@ -723,7 +724,8 @@ export const ensureUseDagServiceGetDagsUiData = (queryClient: QueryClient, { ass
723724
paused?: boolean;
724725
tags?: string[];
725726
tagsMatchMode?: "any" | "all";
726-
} = {}) => queryClient.ensureQueryData({ queryKey: Common.UseDagServiceGetDagsUiKeyFn({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }), queryFn: () => DagService.getDagsUi({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }) });
727+
teams?: string[];
728+
} = {}) => queryClient.ensureQueryData({ queryKey: Common.UseDagServiceGetDagsUiKeyFn({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }), queryFn: () => DagService.getDagsUi({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }) });
727729
/**
728730
* Get Latest Run Info
729731
* Get latest run.

airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -675,6 +675,7 @@ export const prefetchUseDagServiceGetDagTags = (queryClient: QueryClient, { limi
675675
* @param data.tags
676676
* @param data.tagsMatchMode
677677
* @param data.owners
678+
* @param data.teams
678679
* @param data.dagIds
679680
* @param data.dagIdPattern SQL LIKE expression — use `%` / `_` wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 | dag2`). Regular expressions are **not** supported.
680681
*
@@ -699,7 +700,7 @@ export const prefetchUseDagServiceGetDagTags = (queryClient: QueryClient, { limi
699700
* @returns DAGWithLatestDagRunsCollectionResponse Successful Response
700701
* @throws ApiError
701702
*/
702-
export const prefetchUseDagServiceGetDagsUi = (queryClient: QueryClient, { assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }: {
703+
export const prefetchUseDagServiceGetDagsUi = (queryClient: QueryClient, { assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }: {
703704
assetDependency?: string;
704705
bundleName?: string;
705706
bundleVersion?: string;
@@ -723,7 +724,8 @@ export const prefetchUseDagServiceGetDagsUi = (queryClient: QueryClient, { asset
723724
paused?: boolean;
724725
tags?: string[];
725726
tagsMatchMode?: "any" | "all";
726-
} = {}) => queryClient.prefetchQuery({ queryKey: Common.UseDagServiceGetDagsUiKeyFn({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }), queryFn: () => DagService.getDagsUi({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }) });
727+
teams?: string[];
728+
} = {}) => queryClient.prefetchQuery({ queryKey: Common.UseDagServiceGetDagsUiKeyFn({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }), queryFn: () => DagService.getDagsUi({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }) });
727729
/**
728730
* Get Latest Run Info
729731
* Get latest run.

airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -675,6 +675,7 @@ export const useDagServiceGetDagTags = <TData = Common.DagServiceGetDagTagsDefau
675675
* @param data.tags
676676
* @param data.tagsMatchMode
677677
* @param data.owners
678+
* @param data.teams
678679
* @param data.dagIds
679680
* @param data.dagIdPattern SQL LIKE expression — use `%` / `_` wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 | dag2`). Regular expressions are **not** supported.
680681
*
@@ -699,7 +700,7 @@ export const useDagServiceGetDagTags = <TData = Common.DagServiceGetDagTagsDefau
699700
* @returns DAGWithLatestDagRunsCollectionResponse Successful Response
700701
* @throws ApiError
701702
*/
702-
export const useDagServiceGetDagsUi = <TData = Common.DagServiceGetDagsUiDefaultResponse, TError = unknown, TQueryKey extends Array<unknown> = unknown[]>({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }: {
703+
export const useDagServiceGetDagsUi = <TData = Common.DagServiceGetDagsUiDefaultResponse, TError = unknown, TQueryKey extends Array<unknown> = unknown[]>({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }: {
703704
assetDependency?: string;
704705
bundleName?: string;
705706
bundleVersion?: string;
@@ -723,7 +724,8 @@ export const useDagServiceGetDagsUi = <TData = Common.DagServiceGetDagsUiDefault
723724
paused?: boolean;
724725
tags?: string[];
725726
tagsMatchMode?: "any" | "all";
726-
} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, "queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey: Common.UseDagServiceGetDagsUiKeyFn({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }, queryKey), queryFn: () => DagService.getDagsUi({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode }) as TData, ...options });
727+
teams?: string[];
728+
} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, "queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey: Common.UseDagServiceGetDagsUiKeyFn({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }, queryKey), queryFn: () => DagService.getDagsUi({ assetDependency, bundleName, bundleVersion, dagDisplayNamePattern, dagDisplayNamePrefixPattern, dagIdPattern, dagIdPrefixPattern, dagIds, dagRunsLimit, dagRunState, excludeStale, hasAssetSchedule, hasImportErrors, hasPendingActions, isFavorite, lastDagRunState, limit, offset, orderBy, owners, paused, tags, tagsMatchMode, teams }) as TData, ...options });
727729
/**
728730
* Get Latest Run Info
729731
* Get latest run.

0 commit comments

Comments
 (0)