This is an automated email from the ASF dual-hosted git repository.
vincbeck 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 94b9fe33284 UI: Add team column and filter to Dag Run and Task
Instance lists (#70241)
94b9fe33284 is described below
commit 94b9fe33284b3067205d753916af118d18daf0c2
Author: Vincent <[email protected]>
AuthorDate: Thu Jul 23 13:31:26 2026 -0400
UI: Add team column and filter to Dag Run and Task Instance lists (#70241)
When multi-team mode is enabled, operators triaging failures across teams
need to see which team owns each Dag Run and Task Instance, and scope the
lists to a single team. This adds:
- A "Team" column on the Dag Run list (after Run Type) and the Task Instance
list (after State), each linking to the Dag list filtered by that team
- A team filter in the existing filter bar on both pages
- Backend support: a nullable `team_name` field on `DAGRunResponse` and
`TaskInstanceResponse`, populated (via bundle association) only when
multi-team is enabled, plus a `teams` query parameter on the list
endpoints that scopes results by team
The column and filter render only when the `multi_team` configuration is
enabled, so single-team deployments are unaffected.
---
.../src/airflow/api_fastapi/common/db/dags.py | 20 +++++
.../src/airflow/api_fastapi/common/parameters.py | 46 ++++++++++
.../api_fastapi/core_api/datamodels/dag_run.py | 1 +
.../core_api/datamodels/task_instances.py | 1 +
.../api_fastapi/core_api/openapi/_private_ui.yaml | 5 ++
.../core_api/openapi/v2-rest-api-generated.yaml | 26 ++++++
.../api_fastapi/core_api/routes/public/dag_run.py | 7 ++
.../core_api/routes/public/task_instances.py | 8 ++
.../src/airflow/ui/openapi-gen/queries/common.ts | 10 ++-
.../ui/openapi-gen/queries/ensureQueryData.ts | 12 ++-
.../src/airflow/ui/openapi-gen/queries/prefetch.ts | 12 ++-
.../src/airflow/ui/openapi-gen/queries/queries.ts | 12 ++-
.../src/airflow/ui/openapi-gen/queries/suspense.ts | 12 ++-
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 22 +++++
.../ui/openapi-gen/requests/services.gen.ts | 4 +
.../airflow/ui/openapi-gen/requests/types.gen.ts | 4 +
.../airflow/ui/public/i18n/locales/en/common.json | 1 +
.../src/airflow/ui/src/constants/filterConfigs.tsx | 17 +++-
.../src/airflow/ui/src/pages/DagRuns/DagRuns.tsx | 24 +++++-
.../ui/src/pages/DagRuns/DagRunsFilters.tsx | 6 ++
.../ui/src/pages/TaskInstances/TaskInstances.tsx | 23 ++++-
.../pages/TaskInstances/TaskInstancesFilter.tsx | 7 ++
.../src/airflow/ui/src/utils/useFiltersHandler.ts | 1 +
.../core_api/routes/public/test_assets.py | 1 +
.../core_api/routes/public/test_dag_run.py | 83 +++++++++++++++++-
.../core_api/routes/public/test_hitl.py | 1 +
.../core_api/routes/public/test_task_instances.py | 98 ++++++++++++++++++++++
.../tests/unit/cli/commands/test_asset_command.py | 2 +
.../src/airflowctl/api/datamodels/generated.py | 2 +
29 files changed, 443 insertions(+), 25 deletions(-)
diff --git a/airflow-core/src/airflow/api_fastapi/common/db/dags.py
b/airflow-core/src/airflow/api_fastapi/common/db/dags.py
index 09cb9afc8cc..7113e104666 100644
--- a/airflow-core/src/airflow/api_fastapi/common/db/dags.py
+++ b/airflow-core/src/airflow/api_fastapi/common/db/dags.py
@@ -17,6 +17,7 @@
from __future__ import annotations
+from collections.abc import Sequence
from typing import TYPE_CHECKING
from sqlalchemy import func, select
@@ -26,10 +27,12 @@ from airflow.api_fastapi.common.db.common import (
apply_filters_to_select,
)
from airflow.api_fastapi.common.parameters import BaseParam, RangeFilter,
SortParam
+from airflow.configuration import conf
from airflow.models import DagModel
from airflow.models.dagrun import DagRun
if TYPE_CHECKING:
+ from sqlalchemy.orm import Session
from sqlalchemy.sql import Select
@@ -93,3 +96,20 @@ def generate_dag_with_latest_run_query(
)
return query
+
+
+def attach_team_names(objects: Sequence, *, session: Session) -> None:
+ """
+ Attach the owning team name to each object exposing a ``dag_id``.
+
+ Only performs a database lookup when multi-team mode is enabled; otherwise
every
+ object keeps its default ``team_name`` of ``None``. The resolved name is
set as the
+ ``team_name`` attribute on each object so the response serializer can read
it.
+ """
+ if not objects or not conf.getboolean("core", "multi_team"):
+ return
+
+ dag_ids = list({obj.dag_id for obj in objects})
+ team_names_by_dag_id = DagModel.get_dag_id_to_team_name_mapping(dag_ids,
session=session)
+ for obj in objects:
+ obj.team_name = team_names_by_dag_id.get(obj.dag_id)
diff --git a/airflow-core/src/airflow/api_fastapi/common/parameters.py
b/airflow-core/src/airflow/api_fastapi/common/parameters.py
index 18445acfb04..1f940dcfd5a 100644
--- a/airflow-core/src/airflow/api_fastapi/common/parameters.py
+++ b/airflow-core/src/airflow/api_fastapi/common/parameters.py
@@ -1044,6 +1044,52 @@ class _TeamsFilter(BaseParam[list[str]]):
return cls().set_value(teams)
+class _DagIdTeamsFilter(BaseParam[list[str]]):
+ """Filter rows by team name through their ``dag_id`` (via bundle
association)."""
+
+ def __init__(
+ self,
+ dag_id_attribute: ColumnElement | InstrumentedAttribute,
+ value: list[str] | None = None,
+ skip_none: bool = True,
+ ) -> None:
+ super().__init__(value, skip_none)
+ self.dag_id_attribute = dag_id_attribute
+
+ def to_orm(self, select: Select) -> Select:
+ if self.skip_none is False:
+ raise ValueError(f"Cannot set 'skip_none' to False on a
{type(self)}")
+
+ if not self.value:
+ return select
+
+ from airflow.models.team import Team
+
+ return select.where(
+ self.dag_id_attribute.in_(
+ sql_select(DagModel.dag_id)
+ .join(DagBundleModel, DagModel.bundle_name ==
DagBundleModel.name)
+ .join(DagBundleModel.teams)
+ .where(Team.name.in_(self.value))
+ )
+ )
+
+ @classmethod
+ def depends(cls, *args: Any, **kwargs: Any) -> Self:
+ raise NotImplementedError("Use teams_filter_factory instead, depends
is not implemented.")
+
+
+def teams_filter_factory(
+ dag_id_attribute: ColumnElement | InstrumentedAttribute,
+) -> Callable[[list[str]], _DagIdTeamsFilter]:
+ """Build a ``teams`` filter that scopes rows by team through the given
``dag_id`` column."""
+
+ def depends_teams_filter(teams: list[str] = Query(default_factory=list))
-> _DagIdTeamsFilter:
+ return _DagIdTeamsFilter(dag_id_attribute).set_value(teams)
+
+ return depends_teams_filter
+
+
def _safe_parse_datetime(date_to_check: str) -> datetime:
"""
Parse datetime and raise error for invalid dates.
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_run.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_run.py
index cae9d3a8339..d728f8fd0b0 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_run.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_run.py
@@ -184,6 +184,7 @@ class DAGRunResponse(BaseModel):
dag_display_name: str = Field(validation_alias=AliasPath("dag_model",
"dag_display_name"))
partition_key: str | None
partition_date: datetime | None
+ team_name: str | None = None
class DAGRunCollectionResponse(BaseModel):
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
index 85e14bc3baf..4ea8a3c1f72 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
@@ -89,6 +89,7 @@ class TaskInstanceResponse(BaseModel):
trigger: TriggerResponse | None
queued_by_job: JobResponse | None = Field(alias="triggerer_job")
dag_version: DagVersionResponse | None
+ team_name: str | None = None
class TaskInstanceCollectionResponse(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 6dc04dde73b..ca8f7299982 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
@@ -4090,6 +4090,11 @@ components:
anyOf:
- $ref: '#/components/schemas/DagVersionResponse'
- type: 'null'
+ team_name:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Team Name
type: object
required:
- id
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index 4185e7541d5..1d5c1b6abf7 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -2681,6 +2681,14 @@ paths:
- type: string
- type: 'null'
title: Bundle Version
+ - name: teams
+ in: query
+ required: false
+ schema:
+ type: array
+ items:
+ type: string
+ title: Teams
- name: order_by
in: query
required: false
@@ -9169,6 +9177,14 @@ paths:
items:
type: integer
title: Version Number
+ - name: teams
+ in: query
+ required: false
+ schema:
+ type: array
+ items:
+ type: string
+ title: Teams
- name: try_number
in: query
required: false
@@ -14230,6 +14246,11 @@ components:
format: date-time
- type: 'null'
title: Partition Date
+ team_name:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Team Name
type: object
required:
- dag_run_id
@@ -16414,6 +16435,11 @@ components:
anyOf:
- $ref: '#/components/schemas/DagVersionResponse'
- type: 'null'
+ team_name:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Team Name
type: object
required:
- id
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py
index 7b03c7b3ece..ea650963e35 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py
@@ -42,6 +42,7 @@ from airflow.api_fastapi.common.db.dag_runs import (
attach_dag_versions_to_runs,
eager_load_dag_run_for_list,
)
+from airflow.api_fastapi.common.db.dags import attach_team_names
from airflow.api_fastapi.common.parameters import (
FilterOptionEnum,
FilterParam,
@@ -58,6 +59,7 @@ from airflow.api_fastapi.common.parameters import (
Range,
RangeFilter,
SortParam,
+ _DagIdTeamsFilter,
_PrefixSearchParam,
_SearchParam,
datetime_range_filter_factory,
@@ -65,6 +67,7 @@ from airflow.api_fastapi.common.parameters import (
float_range_filter_factory,
prefix_search_param_factory,
search_param_factory,
+ teams_filter_factory,
)
from airflow.api_fastapi.common.router import AirflowRouter
from airflow.api_fastapi.common.types import Mimetype
@@ -500,6 +503,7 @@ def get_dag_runs(
bundle_version: Annotated[
FilterParam[str | None],
Depends(filter_param_factory(DagRun.bundle_version, str | None))
],
+ teams: Annotated[_DagIdTeamsFilter,
Depends(teams_filter_factory(DagRun.dag_id))],
order_by: Annotated[
SortParam,
Depends(
@@ -655,6 +659,7 @@ def get_dag_runs(
partition_key_pattern,
partition_key_prefix_pattern,
consuming_asset_pattern,
+ teams,
]
if use_cursor:
@@ -688,6 +693,7 @@ def get_dag_runs(
has_next = has_more
attach_dag_versions_to_runs(dag_runs, session=session)
+ attach_team_names(dag_runs, session=session)
return DAGRunCollectionResponse(
dag_runs=dag_runs,
@@ -707,6 +713,7 @@ def get_dag_runs(
)
dag_runs = list(session.scalars(dag_run_select))
attach_dag_versions_to_runs(dag_runs, session=session)
+ attach_team_names(dag_runs, session=session)
return DAGRunCollectionResponse(
dag_runs=dag_runs,
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
index 7305e4e89be..214179cdca9 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
+++
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
@@ -41,6 +41,7 @@ from airflow.api_fastapi.common.dagbag import (
resolve_run_on_latest_version,
)
from airflow.api_fastapi.common.db.common import SessionDep,
apply_filters_to_select, paginated_select
+from airflow.api_fastapi.common.db.dags import attach_team_names
from airflow.api_fastapi.common.db.task_instances import
eager_load_TI_and_TIH_for_validation
from airflow.api_fastapi.common.parameters import (
FilterOptionEnum,
@@ -71,6 +72,7 @@ from airflow.api_fastapi.common.parameters import (
Range,
RangeFilter,
SortParam,
+ _DagIdTeamsFilter,
_PrefixSearchParam,
_SearchParam,
datetime_range_filter_factory,
@@ -78,6 +80,7 @@ from airflow.api_fastapi.common.parameters import (
float_range_filter_factory,
prefix_search_param_factory,
search_param_factory,
+ teams_filter_factory,
)
from airflow.api_fastapi.common.router import AirflowRouter
from airflow.api_fastapi.core_api.base import OrmClause
@@ -481,6 +484,7 @@ def get_task_instances(
queue_name_prefix_pattern: QueryTIQueueNamePrefixPatternSearch,
executor: QueryTIExecutorFilter,
version_number: QueryTIDagVersionFilter,
+ teams: Annotated[_DagIdTeamsFilter,
Depends(teams_filter_factory(TI.dag_id))],
try_number: QueryTITryNumberFilter,
operator: QueryTIOperatorFilter,
operator_name_pattern: QueryTIOperatorNamePatternSearch,
@@ -599,6 +603,7 @@ def get_task_instances(
map_index,
rendered_map_index_pattern,
rendered_map_index_prefix_pattern,
+ teams,
]
if use_cursor:
@@ -635,6 +640,8 @@ def get_task_instances(
has_prev = bool(cursor)
has_next = has_more
+ attach_team_names(task_instances, session=session)
+
return TaskInstanceCollectionResponse(
task_instances=task_instances,
next_cursor=(
@@ -656,6 +663,7 @@ def get_task_instances(
session=session,
)
task_instances = list(session.scalars(task_instance_select))
+ attach_team_names(task_instances, session=session)
return TaskInstanceCollectionResponse(
task_instances=task_instances,
total_entries=total_entries,
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 d27c22d162d..4f485202973 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -160,7 +160,7 @@ export const UseDagRunServiceGetDagRunKeyFn = ({ dagId,
dagRunId }: {
export type DagRunServiceGetDagRunsDefaultResponse = Awaited<ReturnType<typeof
DagRunService.getDagRuns>>;
export type DagRunServiceGetDagRunsQueryResult<TData =
DagRunServiceGetDagRunsDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
export const useDagRunServiceGetDagRunsKey = "DagRunServiceGetDagRuns";
-export const UseDagRunServiceGetDagRunsKeyFn = ({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLte, partitionKeyPattern, partitionKeyPrefixPattern, runAfterGt,
runAfterGte, runAfterLt, runAfterLte, runIdPattern, [...]
+export const UseDagRunServiceGetDagRunsKeyFn = ({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLte, partitionKeyPattern, partitionKeyPrefixPattern, runAfterGt,
runAfterGte, runAfterLt, runAfterLte, runIdPattern, [...]
bundleVersion?: string;
confContains?: string;
consumingAssetPattern?: string;
@@ -200,13 +200,14 @@ export const UseDagRunServiceGetDagRunsKeyFn = ({
bundleVersion, confContains, c
startDateLt?: string;
startDateLte?: string;
state?: string[];
+ teams?: string[];
triggeringUserNamePattern?: string;
triggeringUserNamePrefixPattern?: string;
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
-}, queryKey?: Array<unknown>) => [useDagRunServiceGetDagRunsKey, ...(queryKey
?? [{ bundleVersion, confContains, consumingAssetPattern, cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagVersion, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte, limit,
logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte, offset, orderBy,
partitionDateGte, partitionDateLte, partitionKeyPattern,
partitionKeyPrefixPattern, runAfterGt, runAfterGte, runA [...]
+}, queryKey?: Array<unknown>) => [useDagRunServiceGetDagRunsKey, ...(queryKey
?? [{ bundleVersion, confContains, consumingAssetPattern, cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagVersion, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte, limit,
logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte, offset, orderBy,
partitionDateGte, partitionDateLte, partitionKeyPattern,
partitionKeyPrefixPattern, runAfterGt, runAfterGte, runA [...]
export type DagRunServiceGetUpstreamAssetEventsDefaultResponse =
Awaited<ReturnType<typeof DagRunService.getUpstreamAssetEvents>>;
export type DagRunServiceGetUpstreamAssetEventsQueryResult<TData =
DagRunServiceGetUpstreamAssetEventsDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
export const useDagRunServiceGetUpstreamAssetEventsKey =
"DagRunServiceGetUpstreamAssetEvents";
@@ -541,7 +542,7 @@ export const
UseTaskInstanceServiceGetMappedTaskInstanceKeyFn = ({ dagId, dagRun
export type TaskInstanceServiceGetTaskInstancesDefaultResponse =
Awaited<ReturnType<typeof TaskInstanceService.getTaskInstances>>;
export type TaskInstanceServiceGetTaskInstancesQueryResult<TData =
TaskInstanceServiceGetTaskInstancesDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
export const useTaskInstanceServiceGetTaskInstancesKey =
"TaskInstanceServiceGetTaskInstances";
-export const UseTaskInstanceServiceGetTaskInstancesKeyFn = ({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPattern,
orderBy, pool, poolNamePattern, poolNamePrefixPattern, queue, queueNamePattern,
queueNamePrefixPattern, renderedMapIndex [...]
+export const UseTaskInstanceServiceGetTaskInstancesKeyFn = ({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPattern,
orderBy, pool, poolNamePattern, poolNamePrefixPattern, queue, queueNamePattern,
queueNamePrefixPattern, renderedMapIndex [...]
cursor?: string;
dagId: string;
dagIdPattern?: string;
@@ -590,13 +591,14 @@ export const UseTaskInstanceServiceGetTaskInstancesKeyFn
= ({ cursor, dagId, dag
taskDisplayNamePrefixPattern?: string;
taskGroupId?: string;
taskId?: string;
+ teams?: string[];
tryNumber?: number[];
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
versionNumber?: number[];
-}, queryKey?: Array<unknown>) => [useTaskInstanceServiceGetTaskInstancesKey,
...(queryKey ?? [{ cursor, dagId, dagIdPattern, dagIdPrefixPattern, dagRunId,
durationGt, durationGte, durationLt, durationLte, endDateGt, endDateGte,
endDateLt, endDateLte, executor, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, mapIndex, offset, operator, operatorNamePattern,
operatorNamePrefixPattern, orderBy, pool, poolNamePattern,
poolNamePrefixPattern, queue, queueNamePattern, queueN [...]
+}, queryKey?: Array<unknown>) => [useTaskInstanceServiceGetTaskInstancesKey,
...(queryKey ?? [{ cursor, dagId, dagIdPattern, dagIdPrefixPattern, dagRunId,
durationGt, durationGte, durationLt, durationLte, endDateGt, endDateGte,
endDateLt, endDateLte, executor, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, mapIndex, offset, operator, operatorNamePattern,
operatorNamePrefixPattern, orderBy, pool, poolNamePattern,
poolNamePrefixPattern, queue, queueNamePattern, queueN [...]
export type TaskInstanceServiceGetTaskInstanceTryDetailsDefaultResponse =
Awaited<ReturnType<typeof TaskInstanceService.getTaskInstanceTryDetails>>;
export type TaskInstanceServiceGetTaskInstanceTryDetailsQueryResult<TData =
TaskInstanceServiceGetTaskInstanceTryDetailsDefaultResponse, TError = unknown>
= UseQueryResult<TData, TError>;
export const useTaskInstanceServiceGetTaskInstanceTryDetailsKey =
"TaskInstanceServiceGetTaskInstanceTryDetails";
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 45999cf181b..08ecd2f13f3 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
@@ -349,6 +349,7 @@ export const ensureUseDagRunServiceGetDagRunData =
(queryClient: QueryClient, {
* @param data.state
* @param data.dagVersion
* @param data.bundleVersion
+* @param data.teams
* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
state, dag_id, run_id, logical_date, partition_date, run_after, start_date,
end_date, updated_at, conf, duration, dag_run_id`
* @param data.runIdPattern SQL LIKE expression — use `%` / `_` wildcards (e.g.
`%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 | dag2`).
Regular expressions are **not** supported.
*
@@ -370,7 +371,7 @@ export const ensureUseDagRunServiceGetDagRunData =
(queryClient: QueryClient, {
* @returns DAGRunCollectionResponse Successful Response
* @throws ApiError
*/
-export const ensureUseDagRunServiceGetDagRunsData = (queryClient: QueryClient,
{ bundleVersion, confContains, consumingAssetPattern, cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagVersion, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte, limit,
logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte, offset, orderBy,
partitionDateGte, partitionDateLte, partitionKeyPattern,
partitionKeyPrefixPattern, runAfterGt, runAfterGte, runAfte [...]
+export const ensureUseDagRunServiceGetDagRunsData = (queryClient: QueryClient,
{ bundleVersion, confContains, consumingAssetPattern, cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagVersion, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte, limit,
logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte, offset, orderBy,
partitionDateGte, partitionDateLte, partitionKeyPattern,
partitionKeyPrefixPattern, runAfterGt, runAfterGte, runAfte [...]
bundleVersion?: string;
confContains?: string;
consumingAssetPattern?: string;
@@ -410,13 +411,14 @@ export const ensureUseDagRunServiceGetDagRunsData =
(queryClient: QueryClient, {
startDateLt?: string;
startDateLte?: string;
state?: string[];
+ teams?: string[];
triggeringUserNamePattern?: string;
triggeringUserNamePrefixPattern?: string;
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
-}) => queryClient.ensureQueryData({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLte, partitionKeyPattern, partitionKeyPrefixPattern, runAfterGt,
runAfterGte, r [...]
+}) => queryClient.ensureQueryData({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLte, partitionKeyPattern, partitionKeyPrefixPattern, runAfterGt,
runAfterGte, r [...]
/**
* Get Upstream Asset Events
* If dag run is asset-triggered, return the asset events that triggered it.
@@ -1145,6 +1147,7 @@ export const
ensureUseTaskInstanceServiceGetMappedTaskInstanceData = (queryClien
* @param data.queueNamePrefixPattern Prefix match — returns items whose value
starts with the given string (case-sensitive, index-friendly). Use the pipe `|`
operator for OR logic (e.g. `dag1|dag2`). Use `~` to match all. Wildcard
characters (`%`, `_`) are treated as literal characters. Trailing
non-alphanumeric characters in the prefix are stripped before matching so the
range scan stays index-compatible under locale-aware collations — e.g. `test_`
effectively matches items starting wit [...]
* @param data.executor
* @param data.versionNumber
+* @param data.teams
* @param data.tryNumber
* @param data.operator
* @param data.operatorNamePattern SQL LIKE expression — use `%` / `_`
wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g.
`dag1 | dag2`). Regular expressions are **not** supported.
@@ -1162,7 +1165,7 @@ export const
ensureUseTaskInstanceServiceGetMappedTaskInstanceData = (queryClien
* @returns TaskInstanceCollectionResponse Successful Response
* @throws ApiError
*/
-export const ensureUseTaskInstanceServiceGetTaskInstancesData = (queryClient:
QueryClient, { cursor, dagId, dagIdPattern, dagIdPrefixPattern, dagRunId,
durationGt, durationGte, durationLt, durationLte, endDateGt, endDateGte,
endDateLt, endDateLte, executor, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, mapIndex, offset, operator, operatorNamePattern,
operatorNamePrefixPattern, orderBy, pool, poolNamePattern,
poolNamePrefixPattern, queue, queueNamePattern, queueName [...]
+export const ensureUseTaskInstanceServiceGetTaskInstancesData = (queryClient:
QueryClient, { cursor, dagId, dagIdPattern, dagIdPrefixPattern, dagRunId,
durationGt, durationGte, durationLt, durationLte, endDateGt, endDateGte,
endDateLt, endDateLte, executor, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, mapIndex, offset, operator, operatorNamePattern,
operatorNamePrefixPattern, orderBy, pool, poolNamePattern,
poolNamePrefixPattern, queue, queueNamePattern, queueName [...]
cursor?: string;
dagId: string;
dagIdPattern?: string;
@@ -1211,13 +1214,14 @@ export const
ensureUseTaskInstanceServiceGetTaskInstancesData = (queryClient: Qu
taskDisplayNamePrefixPattern?: string;
taskGroupId?: string;
taskId?: string;
+ teams?: string[];
tryNumber?: number[];
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
versionNumber?: number[];
-}) => queryClient.ensureQueryData({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPattern,
orderBy, pool, poolNamePattern, poolNamePrefixPattern, queue, queueNamePattern,
que [...]
+}) => queryClient.ensureQueryData({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPattern,
orderBy, pool, poolNamePattern, poolNamePrefixPattern, queue, queueNamePattern,
que [...]
/**
* Get Task Instance Try Details
* Get task instance details by try number.
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 4a44a68a8f3..0b0d7b46eeb 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
@@ -349,6 +349,7 @@ export const prefetchUseDagRunServiceGetDagRun =
(queryClient: QueryClient, { da
* @param data.state
* @param data.dagVersion
* @param data.bundleVersion
+* @param data.teams
* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
state, dag_id, run_id, logical_date, partition_date, run_after, start_date,
end_date, updated_at, conf, duration, dag_run_id`
* @param data.runIdPattern SQL LIKE expression — use `%` / `_` wildcards (e.g.
`%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 | dag2`).
Regular expressions are **not** supported.
*
@@ -370,7 +371,7 @@ export const prefetchUseDagRunServiceGetDagRun =
(queryClient: QueryClient, { da
* @returns DAGRunCollectionResponse Successful Response
* @throws ApiError
*/
-export const prefetchUseDagRunServiceGetDagRuns = (queryClient: QueryClient, {
bundleVersion, confContains, consumingAssetPattern, cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagVersion, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte, limit,
logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte, offset, orderBy,
partitionDateGte, partitionDateLte, partitionKeyPattern,
partitionKeyPrefixPattern, runAfterGt, runAfterGte, runAfterL [...]
+export const prefetchUseDagRunServiceGetDagRuns = (queryClient: QueryClient, {
bundleVersion, confContains, consumingAssetPattern, cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagVersion, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte, limit,
logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte, offset, orderBy,
partitionDateGte, partitionDateLte, partitionKeyPattern,
partitionKeyPrefixPattern, runAfterGt, runAfterGte, runAfterL [...]
bundleVersion?: string;
confContains?: string;
consumingAssetPattern?: string;
@@ -410,13 +411,14 @@ export const prefetchUseDagRunServiceGetDagRuns =
(queryClient: QueryClient, { b
startDateLt?: string;
startDateLte?: string;
state?: string[];
+ teams?: string[];
triggeringUserNamePattern?: string;
triggeringUserNamePrefixPattern?: string;
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
-}) => queryClient.prefetchQuery({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLte, partitionKeyPattern, partitionKeyPrefixPattern, runAfterGt,
runAfterGte, run [...]
+}) => queryClient.prefetchQuery({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLte, partitionKeyPattern, partitionKeyPrefixPattern, runAfterGt,
runAfterGte, run [...]
/**
* Get Upstream Asset Events
* If dag run is asset-triggered, return the asset events that triggered it.
@@ -1145,6 +1147,7 @@ export const
prefetchUseTaskInstanceServiceGetMappedTaskInstance = (queryClient:
* @param data.queueNamePrefixPattern Prefix match — returns items whose value
starts with the given string (case-sensitive, index-friendly). Use the pipe `|`
operator for OR logic (e.g. `dag1|dag2`). Use `~` to match all. Wildcard
characters (`%`, `_`) are treated as literal characters. Trailing
non-alphanumeric characters in the prefix are stripped before matching so the
range scan stays index-compatible under locale-aware collations — e.g. `test_`
effectively matches items starting wit [...]
* @param data.executor
* @param data.versionNumber
+* @param data.teams
* @param data.tryNumber
* @param data.operator
* @param data.operatorNamePattern SQL LIKE expression — use `%` / `_`
wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g.
`dag1 | dag2`). Regular expressions are **not** supported.
@@ -1162,7 +1165,7 @@ export const
prefetchUseTaskInstanceServiceGetMappedTaskInstance = (queryClient:
* @returns TaskInstanceCollectionResponse Successful Response
* @throws ApiError
*/
-export const prefetchUseTaskInstanceServiceGetTaskInstances = (queryClient:
QueryClient, { cursor, dagId, dagIdPattern, dagIdPrefixPattern, dagRunId,
durationGt, durationGte, durationLt, durationLte, endDateGt, endDateGte,
endDateLt, endDateLte, executor, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, mapIndex, offset, operator, operatorNamePattern,
operatorNamePrefixPattern, orderBy, pool, poolNamePattern,
poolNamePrefixPattern, queue, queueNamePattern, queueNamePr [...]
+export const prefetchUseTaskInstanceServiceGetTaskInstances = (queryClient:
QueryClient, { cursor, dagId, dagIdPattern, dagIdPrefixPattern, dagRunId,
durationGt, durationGte, durationLt, durationLte, endDateGt, endDateGte,
endDateLt, endDateLte, executor, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, mapIndex, offset, operator, operatorNamePattern,
operatorNamePrefixPattern, orderBy, pool, poolNamePattern,
poolNamePrefixPattern, queue, queueNamePattern, queueNamePr [...]
cursor?: string;
dagId: string;
dagIdPattern?: string;
@@ -1211,13 +1214,14 @@ export const
prefetchUseTaskInstanceServiceGetTaskInstances = (queryClient: Quer
taskDisplayNamePrefixPattern?: string;
taskGroupId?: string;
taskId?: string;
+ teams?: string[];
tryNumber?: number[];
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
versionNumber?: number[];
-}) => queryClient.prefetchQuery({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPattern,
orderBy, pool, poolNamePattern, poolNamePrefixPattern, queue, queueNamePattern,
queue [...]
+}) => queryClient.prefetchQuery({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPattern,
orderBy, pool, poolNamePattern, poolNamePrefixPattern, queue, queueNamePattern,
queue [...]
/**
* Get Task Instance Try Details
* Get task instance details by try number.
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 e3536fa0809..59d5c2dce19 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -349,6 +349,7 @@ export const useDagRunServiceGetDagRun = <TData =
Common.DagRunServiceGetDagRunD
* @param data.state
* @param data.dagVersion
* @param data.bundleVersion
+* @param data.teams
* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
state, dag_id, run_id, logical_date, partition_date, run_after, start_date,
end_date, updated_at, conf, duration, dag_run_id`
* @param data.runIdPattern SQL LIKE expression — use `%` / `_` wildcards (e.g.
`%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 | dag2`).
Regular expressions are **not** supported.
*
@@ -370,7 +371,7 @@ export const useDagRunServiceGetDagRun = <TData =
Common.DagRunServiceGetDagRunD
* @returns DAGRunCollectionResponse Successful Response
* @throws ApiError
*/
-export const useDagRunServiceGetDagRuns = <TData =
Common.DagRunServiceGetDagRunsDefaultResponse, TError = unknown, TQueryKey
extends Array<unknown> = unknown[]>({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLt [...]
+export const useDagRunServiceGetDagRuns = <TData =
Common.DagRunServiceGetDagRunsDefaultResponse, TError = unknown, TQueryKey
extends Array<unknown> = unknown[]>({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte,
partitionDateLt [...]
bundleVersion?: string;
confContains?: string;
consumingAssetPattern?: string;
@@ -410,13 +411,14 @@ export const useDagRunServiceGetDagRuns = <TData =
Common.DagRunServiceGetDagRun
startDateLt?: string;
startDateLte?: string;
state?: string[];
+ teams?: string[];
triggeringUserNamePattern?: string;
triggeringUserNamePrefixPattern?: string;
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
-}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte, pa [...]
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte, pa [...]
/**
* Get Upstream Asset Events
* If dag run is asset-triggered, return the asset events that triggered it.
@@ -1145,6 +1147,7 @@ export const useTaskInstanceServiceGetMappedTaskInstance
= <TData = Common.TaskI
* @param data.queueNamePrefixPattern Prefix match — returns items whose value
starts with the given string (case-sensitive, index-friendly). Use the pipe `|`
operator for OR logic (e.g. `dag1|dag2`). Use `~` to match all. Wildcard
characters (`%`, `_`) are treated as literal characters. Trailing
non-alphanumeric characters in the prefix are stripped before matching so the
range scan stays index-compatible under locale-aware collations — e.g. `test_`
effectively matches items starting wit [...]
* @param data.executor
* @param data.versionNumber
+* @param data.teams
* @param data.tryNumber
* @param data.operator
* @param data.operatorNamePattern SQL LIKE expression — use `%` / `_`
wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g.
`dag1 | dag2`). Regular expressions are **not** supported.
@@ -1162,7 +1165,7 @@ export const useTaskInstanceServiceGetMappedTaskInstance
= <TData = Common.TaskI
* @returns TaskInstanceCollectionResponse Successful Response
* @throws ApiError
*/
-export const useTaskInstanceServiceGetTaskInstances = <TData =
Common.TaskInstanceServiceGetTaskInstancesDefaultResponse, TError = unknown,
TQueryKey extends Array<unknown> = unknown[]>({ cursor, dagId, dagIdPattern,
dagIdPrefixPattern, dagRunId, durationGt, durationGte, durationLt, durationLte,
endDateGt, endDateGte, endDateLt, endDateLte, executor, limit, logicalDateGt,
logicalDateGte, logicalDateLt, logicalDateLte, mapIndex, offset, operator,
operatorNamePattern, operatorNamePrefixPat [...]
+export const useTaskInstanceServiceGetTaskInstances = <TData =
Common.TaskInstanceServiceGetTaskInstancesDefaultResponse, TError = unknown,
TQueryKey extends Array<unknown> = unknown[]>({ cursor, dagId, dagIdPattern,
dagIdPrefixPattern, dagRunId, durationGt, durationGte, durationLt, durationLte,
endDateGt, endDateGte, endDateLt, endDateLte, executor, limit, logicalDateGt,
logicalDateGte, logicalDateLt, logicalDateLte, mapIndex, offset, operator,
operatorNamePattern, operatorNamePrefixPat [...]
cursor?: string;
dagId: string;
dagIdPattern?: string;
@@ -1211,13 +1214,14 @@ export const useTaskInstanceServiceGetTaskInstances =
<TData = Common.TaskInstan
taskDisplayNamePrefixPattern?: string;
taskGroupId?: string;
taskId?: string;
+ teams?: string[];
tryNumber?: number[];
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
versionNumber?: number[];
-}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPa [...]
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorNamePrefixPa [...]
/**
* Get Task Instance Try Details
* Get task instance details by try number.
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 a43a834efd2..26965f1a286 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
@@ -349,6 +349,7 @@ export const useDagRunServiceGetDagRunSuspense = <TData =
Common.DagRunServiceGe
* @param data.state
* @param data.dagVersion
* @param data.bundleVersion
+* @param data.teams
* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
state, dag_id, run_id, logical_date, partition_date, run_after, start_date,
end_date, updated_at, conf, duration, dag_run_id`
* @param data.runIdPattern SQL LIKE expression — use `%` / `_` wildcards (e.g.
`%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 | dag2`).
Regular expressions are **not** supported.
*
@@ -370,7 +371,7 @@ export const useDagRunServiceGetDagRunSuspense = <TData =
Common.DagRunServiceGe
* @returns DAGRunCollectionResponse Successful Response
* @throws ApiError
*/
-export const useDagRunServiceGetDagRunsSuspense = <TData =
Common.DagRunServiceGetDagRunsDefaultResponse, TError = unknown, TQueryKey
extends Array<unknown> = unknown[]>({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte, partiti [...]
+export const useDagRunServiceGetDagRunsSuspense = <TData =
Common.DagRunServiceGetDagRunsDefaultResponse, TError = unknown, TQueryKey
extends Array<unknown> = unknown[]>({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDateGte, partiti [...]
bundleVersion?: string;
confContains?: string;
consumingAssetPattern?: string;
@@ -410,13 +411,14 @@ export const useDagRunServiceGetDagRunsSuspense = <TData
= Common.DagRunServiceG
startDateLt?: string;
startDateLte?: string;
state?: string[];
+ teams?: string[];
triggeringUserNamePattern?: string;
triggeringUserNamePrefixPattern?: string;
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
-}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDat [...]
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseDagRunServiceGetDagRunsKeyFn({ bundleVersion, confContains,
consumingAssetPattern, cursor, dagId, dagIdPattern, dagIdPrefixPattern,
dagVersion, durationGt, durationGte, durationLt, durationLte, endDateGt,
endDateGte, endDateLt, endDateLte, limit, logicalDateGt, logicalDateGte,
logicalDateLt, logicalDateLte, offset, orderBy, partitionDat [...]
/**
* Get Upstream Asset Events
* If dag run is asset-triggered, return the asset events that triggered it.
@@ -1145,6 +1147,7 @@ export const
useTaskInstanceServiceGetMappedTaskInstanceSuspense = <TData = Comm
* @param data.queueNamePrefixPattern Prefix match — returns items whose value
starts with the given string (case-sensitive, index-friendly). Use the pipe `|`
operator for OR logic (e.g. `dag1|dag2`). Use `~` to match all. Wildcard
characters (`%`, `_`) are treated as literal characters. Trailing
non-alphanumeric characters in the prefix are stripped before matching so the
range scan stays index-compatible under locale-aware collations — e.g. `test_`
effectively matches items starting wit [...]
* @param data.executor
* @param data.versionNumber
+* @param data.teams
* @param data.tryNumber
* @param data.operator
* @param data.operatorNamePattern SQL LIKE expression — use `%` / `_`
wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g.
`dag1 | dag2`). Regular expressions are **not** supported.
@@ -1162,7 +1165,7 @@ export const
useTaskInstanceServiceGetMappedTaskInstanceSuspense = <TData = Comm
* @returns TaskInstanceCollectionResponse Successful Response
* @throws ApiError
*/
-export const useTaskInstanceServiceGetTaskInstancesSuspense = <TData =
Common.TaskInstanceServiceGetTaskInstancesDefaultResponse, TError = unknown,
TQueryKey extends Array<unknown> = unknown[]>({ cursor, dagId, dagIdPattern,
dagIdPrefixPattern, dagRunId, durationGt, durationGte, durationLt, durationLte,
endDateGt, endDateGte, endDateLt, endDateLte, executor, limit, logicalDateGt,
logicalDateGte, logicalDateLt, logicalDateLte, mapIndex, offset, operator,
operatorNamePattern, operatorNameP [...]
+export const useTaskInstanceServiceGetTaskInstancesSuspense = <TData =
Common.TaskInstanceServiceGetTaskInstancesDefaultResponse, TError = unknown,
TQueryKey extends Array<unknown> = unknown[]>({ cursor, dagId, dagIdPattern,
dagIdPrefixPattern, dagRunId, durationGt, durationGte, durationLt, durationLte,
endDateGt, endDateGte, endDateLt, endDateLte, executor, limit, logicalDateGt,
logicalDateGte, logicalDateLt, logicalDateLte, mapIndex, offset, operator,
operatorNamePattern, operatorNameP [...]
cursor?: string;
dagId: string;
dagIdPattern?: string;
@@ -1211,13 +1214,14 @@ export const
useTaskInstanceServiceGetTaskInstancesSuspense = <TData = Common.Ta
taskDisplayNamePrefixPattern?: string;
taskGroupId?: string;
taskId?: string;
+ teams?: string[];
tryNumber?: number[];
updatedAtGt?: string;
updatedAtGte?: string;
updatedAtLt?: string;
updatedAtLte?: string;
versionNumber?: number[];
-}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorName [...]
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseTaskInstanceServiceGetTaskInstancesKeyFn({ cursor, dagId,
dagIdPattern, dagIdPrefixPattern, dagRunId, durationGt, durationGte,
durationLt, durationLte, endDateGt, endDateGte, endDateLt, endDateLte,
executor, limit, logicalDateGt, logicalDateGte, logicalDateLt, logicalDateLte,
mapIndex, offset, operator, operatorNamePattern, operatorName [...]
/**
* Get Task Instance Try Details
* Get task instance details by try number.
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 fc4f1711c40..2e4c24e1170 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
@@ -3904,6 +3904,17 @@ export const $DAGRunResponse = {
}
],
title: 'Partition Date'
+ },
+ team_name: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Team Name'
}
},
type: 'object',
@@ -7115,6 +7126,17 @@ export const $TaskInstanceResponse = {
type: 'null'
}
]
+ },
+ team_name: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Team Name'
}
},
type: 'object',
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 1316da09da6..b2e288a93f1 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
@@ -1101,6 +1101,7 @@ export class DagRunService {
* @param data.state
* @param data.dagVersion
* @param data.bundleVersion
+ * @param data.teams
* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
state, dag_id, run_id, logical_date, partition_date, run_after, start_date,
end_date, updated_at, conf, duration, dag_run_id`
* @param data.runIdPattern SQL LIKE expression — use `%` / `_` wildcards
(e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g. `dag1 |
dag2`). Regular expressions are **not** supported.
*
@@ -1164,6 +1165,7 @@ export class DagRunService {
state: data.state,
dag_version: data.dagVersion,
bundle_version: data.bundleVersion,
+ teams: data.teams,
order_by: data.orderBy,
run_id_pattern: data.runIdPattern,
run_id_prefix_pattern: data.runIdPrefixPattern,
@@ -2722,6 +2724,7 @@ export class TaskInstanceService {
* @param data.queueNamePrefixPattern Prefix match — returns items whose
value starts with the given string (case-sensitive, index-friendly). Use the
pipe `|` operator for OR logic (e.g. `dag1|dag2`). Use `~` to match all.
Wildcard characters (`%`, `_`) are treated as literal characters. Trailing
non-alphanumeric characters in the prefix are stripped before matching so the
range scan stays index-compatible under locale-aware collations — e.g. `test_`
effectively matches items startin [...]
* @param data.executor
* @param data.versionNumber
+ * @param data.teams
* @param data.tryNumber
* @param data.operator
* @param data.operatorNamePattern SQL LIKE expression — use `%` / `_`
wildcards (e.g. `%customer_%`). Use the pipe `|` operator for OR logic (e.g.
`dag1 | dag2`). Regular expressions are **not** supported.
@@ -2790,6 +2793,7 @@ export class TaskInstanceService {
queue_name_prefix_pattern: data.queueNamePrefixPattern,
executor: data.executor,
version_number: data.versionNumber,
+ teams: data.teams,
try_number: data.tryNumber,
operator: data.operator,
operator_name_pattern: data.operatorNamePattern,
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 d275bc8c3a3..0d36910280a 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
@@ -1050,6 +1050,7 @@ export type DAGRunResponse = {
dag_display_name: string;
partition_key: string | null;
partition_date: string | null;
+ team_name?: string | null;
};
/**
@@ -1848,6 +1849,7 @@ export type TaskInstanceResponse = {
trigger: TriggerResponse | null;
triggerer_job: JobResponse | null;
dag_version: DagVersionResponse | null;
+ team_name?: string | null;
};
/**
@@ -3248,6 +3250,7 @@ export type GetDagRunsData = {
startDateLt?: string | null;
startDateLte?: string | null;
state?: Array<(string)>;
+ teams?: Array<(string)>;
/**
* SQL LIKE expression — use `%` / `_` wildcards (e.g. `%customer_%`). Use
the pipe `|` operator for OR logic (e.g. `dag1 | dag2`). Regular expressions
are **not** supported.
*
@@ -3964,6 +3967,7 @@ export type GetTaskInstancesData = {
*/
taskGroupId?: string | null;
taskId?: string | null;
+ teams?: Array<(string)>;
tryNumber?: Array<(number)>;
updatedAtGt?: string | null;
updatedAtGte?: string | null;
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
index 9beb02245b6..f182cd88b68 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
@@ -8,6 +8,7 @@
"Variables": "Variables"
},
"allOperators": "All Operators",
+ "allTeams": "All Teams",
"appearance": {
"appearance": "Appearance",
"darkMode": "Dark Mode",
diff --git a/airflow-core/src/airflow/ui/src/constants/filterConfigs.tsx
b/airflow-core/src/airflow/ui/src/constants/filterConfigs.tsx
index 3da64282ff5..7c29e276d7b 100644
--- a/airflow-core/src/airflow/ui/src/constants/filterConfigs.tsx
+++ b/airflow-core/src/airflow/ui/src/constants/filterConfigs.tsx
@@ -19,7 +19,7 @@
import { Flex } from "@chakra-ui/react";
import { useTranslation } from "react-i18next";
import { BiTargetLock } from "react-icons/bi";
-import { FiBarChart, FiUser, FiDatabase } from "react-icons/fi";
+import { FiBarChart, FiUser, FiUsers, FiDatabase } from "react-icons/fi";
import { LuBrackets } from "react-icons/lu";
import {
MdDateRange,
@@ -34,6 +34,7 @@ import {
} from "react-icons/md";
import { PiQueue } from "react-icons/pi";
+import { useTeamsServiceListTeams } from "openapi/queries";
import type { DagRunState, DagRunType, TaskInstanceState } from
"openapi/requests/types.gen";
import { DagIcon } from "src/assets/DagIcon";
import { TaskIcon } from "src/assets/TaskIcon";
@@ -47,6 +48,7 @@ import {
jobTypeOptions,
taskInstanceStateOptions,
} from "src/constants/stateOptions";
+import { useConfig } from "src/queries/useConfig";
import { SearchParamsKeys } from "./searchParams";
@@ -60,6 +62,10 @@ export enum FilterTypes {
export const useFilterConfigs = () => {
const { t: translate } = useTranslation(["browse", "common", "components",
"admin", "hitl"]);
+ const multiTeamEnabled = Boolean(useConfig("multi_team"));
+ const { data: teamsData } = useTeamsServiceListTeams({ orderBy: ["name"] },
undefined, {
+ enabled: multiTeamEnabled,
+ });
const filterConfigMap = {
[SearchParamsKeys.ASSET_EVENT_DATE_RANGE]: {
@@ -382,6 +388,15 @@ export const useFilterConfigs = () => {
})),
type: FilterTypes.SELECT,
},
+ [SearchParamsKeys.TEAMS]: {
+ icon: <FiUsers />,
+ label: translate("common:dagDetails.team"),
+ options: [
+ { label: translate("common:allTeams"), value: "" },
+ ...(teamsData?.teams ?? []).map((team) => ({ label: team.name, value:
team.name })),
+ ],
+ type: FilterTypes.SELECT,
+ },
[SearchParamsKeys.TRIGGERING_USER_NAME_PATTERN]: {
hotkeyDisabled: true,
icon: <FiUser />,
diff --git a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.tsx
b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.tsx
index d9efe994719..183306a396c 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.tsx
@@ -48,6 +48,7 @@ import { RouterLink } from "src/components/ui";
import { ActionBar } from "src/components/ui/ActionBar";
import { SearchParamsKeys, type SearchParamsKeysType } from
"src/constants/searchParams";
import { useAdvancedSearchArg } from "src/hooks/useAdvancedSearch";
+import { useConfig } from "src/queries/useConfig";
import { renderDuration, useAutoRefresh, isStatePending, useDocumentTitle }
from "src/utils";
import BulkClearDagRunsButton from "./BulkClearDagRunsButton";
@@ -82,6 +83,7 @@ const {
START_DATE_GTE: START_DATE_GTE_PARAM,
START_DATE_LTE: START_DATE_LTE_PARAM,
STATE: STATE_PARAM,
+ TEAMS: TEAMS_PARAM,
TRIGGERING_USER_NAME_PATTERN: TRIGGERING_USER_NAME_PATTERN_PARAM,
}: SearchParamsKeysType = SearchParamsKeys;
@@ -91,7 +93,7 @@ type ColumnProps = {
readonly translate: TFunction;
} & GetColumnsParams;
-const runColumns = ({ dagId, open, translate }: ColumnProps):
Array<ColumnDef<DAGRunResponse>> => [
+const runColumns = ({ dagId, multiTeam, open, translate }: ColumnProps):
Array<ColumnDef<DAGRunResponse>> => [
{
accessorKey: "select",
cell: ({ row }) => <SelectionRowCheckbox colorPalette="brand"
rowKey={getRowKey(row.original)} />,
@@ -154,6 +156,21 @@ const runColumns = ({ dagId, open, translate }:
ColumnProps): Array<ColumnDef<DA
enableSorting: false,
header: translate("dagRun.runType"),
},
+ ...(multiTeam
+ ? [
+ {
+ accessorKey: "team_name",
+ cell: ({ row: { original } }: DagRunRow) =>
+ original.team_name !== undefined && original.team_name !== null ? (
+ <RouterLink
to={`/dags?teams=${encodeURIComponent(original.team_name)}`}>
+ {original.team_name}
+ </RouterLink>
+ ) : undefined,
+ enableSorting: false,
+ header: translate("dagDetails.team"),
+ },
+ ]
+ : []),
{
accessorKey: "triggering_user_name",
cell: ({ row: { original } }) => <Text>{original.triggering_user_name ??
""}</Text>,
@@ -229,6 +246,7 @@ export const DagRuns = () => {
const [searchParams] = useSearchParams();
const { onClose, onOpen, open } = useDisclosure();
+ const multiTeamEnabled = Boolean(useConfig("multi_team"));
const { setTableURLState, tableURLState } = useTableURLState({
columnVisibility: {
@@ -262,6 +280,7 @@ export const DagRuns = () => {
const durationLte = searchParams.get(DURATION_LTE_PARAM);
const confContains = searchParams.get(CONF_CONTAINS_PARAM);
const partitionKeyPattern = searchParams.get(PARTITION_KEY_PATTERN_PARAM);
+ const teams = searchParams.getAll(TEAMS_PARAM);
const refetchInterval = useAutoRefresh({});
@@ -316,6 +335,7 @@ export const DagRuns = () => {
startDateGte: startDateGte ?? undefined,
startDateLte: startDateLte ?? undefined,
state: filteredState === null ? undefined : [filteredState],
+ teams: teams.length > 0 ? teams : undefined,
...triggeringUserArg,
},
undefined,
@@ -339,7 +359,7 @@ export const DagRuns = () => {
const columns = runColumns({
dagId,
- multiTeam: false,
+ multiTeam: multiTeamEnabled,
open,
translate,
});
diff --git a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRunsFilters.tsx
b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRunsFilters.tsx
index f39e66cab41..fa7abf528b3 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRunsFilters.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRunsFilters.tsx
@@ -20,6 +20,7 @@ import { VStack } from "@chakra-ui/react";
import { FilterBar } from "src/components/FilterBar";
import { SearchParamsKeys } from "src/constants/searchParams";
+import { useConfig } from "src/queries/useConfig";
import { useFiltersHandler, type FilterableSearchParamsKeys } from "src/utils";
type DagRunsFiltersProps = {
@@ -27,6 +28,7 @@ type DagRunsFiltersProps = {
};
export const DagRunsFilters = ({ dagId }: DagRunsFiltersProps) => {
+ const multiTeamEnabled = Boolean(useConfig("multi_team"));
const searchParamKeys: Array<FilterableSearchParamsKeys> = [
SearchParamsKeys.RUN_ID_PATTERN,
SearchParamsKeys.STATE,
@@ -45,6 +47,10 @@ export const DagRunsFilters = ({ dagId }:
DagRunsFiltersProps) => {
SearchParamsKeys.CONSUMING_ASSET_PATTERN,
];
+ if (multiTeamEnabled) {
+ searchParamKeys.push(SearchParamsKeys.TEAMS);
+ }
+
if (dagId === undefined) {
searchParamKeys.unshift(SearchParamsKeys.DAG_ID_PATTERN);
}
diff --git
a/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstances.tsx
b/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstances.tsx
index f09d7f07a22..c1b5f583467 100644
--- a/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstances.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstances.tsx
@@ -44,6 +44,7 @@ import { RouterLink } from "src/components/ui";
import { ActionBar } from "src/components/ui/ActionBar";
import { SearchParamsKeys, type SearchParamsKeysType } from
"src/constants/searchParams";
import { useAdvancedSearchArg } from "src/hooks/useAdvancedSearch";
+import { useConfig } from "src/queries/useConfig";
import { useAutoRefresh, isStatePending, renderDuration, useDocumentTitle }
from "src/utils";
import { getTaskInstanceLink } from "src/utils/links";
@@ -78,6 +79,7 @@ const {
RUN_ID_PATTERN: RUN_ID_PATTERN_PARAM,
START_DATE: START_DATE_PARAM,
TASK_STATE: STATE_PARAM,
+ TEAMS: TEAMS_PARAM,
TRY_NUMBER: TRY_NUMBER_PARAM,
}: SearchParamsKeysType = SearchParamsKeys;
@@ -90,6 +92,7 @@ type ColumnProps = {
const taskInstanceColumns = ({
dagId,
+ multiTeam,
runId,
taskId,
translate,
@@ -174,6 +177,21 @@ const taskInstanceColumns = ({
),
header: () => translate("state"),
},
+ ...(multiTeam
+ ? [
+ {
+ accessorKey: "team_name",
+ cell: ({ row: { original } }: TaskInstanceRow) =>
+ original.team_name !== undefined && original.team_name !== null ? (
+ <RouterLink
to={`/dags?teams=${encodeURIComponent(original.team_name)}`}>
+ {original.team_name}
+ </RouterLink>
+ ) : undefined,
+ enableSorting: false,
+ header: translate("dagDetails.team"),
+ },
+ ]
+ : []),
{
accessorKey: "start_date",
cell: ({ row: { original } }) =>
@@ -258,6 +276,7 @@ export const TaskInstances = () => {
useDocumentTitle(dagId === undefined ?
translate("common:taskInstance_other") : undefined);
const [searchParams] = useSearchParams();
+ const multiTeamEnabled = Boolean(useConfig("multi_team"));
const { setTableURLState, tableURLState } = useTableURLState({
columnVisibility: {
@@ -291,6 +310,7 @@ export const TaskInstances = () => {
const filteredRunId = searchParams.get(RUN_ID_PATTERN_PARAM);
const hasFilteredState = filteredState.length > 0;
const taskDisplayNamePattern = searchParams.get(NAME_PATTERN_PARAM);
+ const teams = searchParams.getAll(TEAMS_PARAM);
const refetchInterval = useAutoRefresh({});
@@ -361,6 +381,7 @@ export const TaskInstances = () => {
...taskDisplayNameArg,
taskGroupId: groupId ?? undefined,
taskId: Boolean(groupId) ? undefined : taskId,
+ teams: teams.length > 0 ? teams : undefined,
tryNumber: tryNumberFilter !== null && tryNumberFilter !== "" ?
[Number(tryNumberFilter)] : undefined,
versionNumber:
filteredDagVersion !== null && filteredDagVersion !== "" ?
[Number(filteredDagVersion)] : undefined,
@@ -386,7 +407,7 @@ export const TaskInstances = () => {
const columns = taskInstanceColumns({
dagId,
- multiTeam: false,
+ multiTeam: multiTeamEnabled,
runId,
taskId: Boolean(groupId) ? undefined : taskId,
translate,
diff --git
a/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstancesFilter.tsx
b/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstancesFilter.tsx
index 2e14d51c7b3..a59c2ceadd5 100644
---
a/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstancesFilter.tsx
+++
b/airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstancesFilter.tsx
@@ -21,6 +21,7 @@ import { useSearchParams, useParams } from "react-router-dom";
import { FilterBar, type FilterValue } from "src/components/FilterBar";
import { SearchParamsKeys, type SearchParamsKeysType } from
"src/constants/searchParams";
+import { useConfig } from "src/queries/useConfig";
import { useFiltersHandler, type FilterableSearchParamsKeys } from "src/utils";
const {
@@ -38,11 +39,13 @@ const {
RENDERED_MAP_INDEX: RENDERED_MAP_INDEX_PARAM,
RUN_ID_PATTERN: RUN_ID_PATTERN_PARAM,
TASK_STATE: STATE_PARAM,
+ TEAMS: TEAMS_PARAM,
TRY_NUMBER: TRY_NUMBER_PARAM,
}: SearchParamsKeysType = SearchParamsKeys;
export const TaskInstancesFilter = () => {
const { dagId, runId } = useParams();
+ const multiTeamEnabled = Boolean(useConfig("multi_team"));
const paramKeys: Array<FilterableSearchParamsKeys> = [
NAME_PATTERN_PARAM as FilterableSearchParamsKeys,
LOGICAL_DATE_RANGE_PARAM as FilterableSearchParamsKeys,
@@ -59,6 +62,10 @@ export const TaskInstancesFilter = () => {
STATE_PARAM as FilterableSearchParamsKeys,
];
+ if (multiTeamEnabled) {
+ paramKeys.push(TEAMS_PARAM as FilterableSearchParamsKeys);
+ }
+
if (runId === undefined) {
paramKeys.unshift(RUN_ID_PATTERN_PARAM as FilterableSearchParamsKeys);
}
diff --git a/airflow-core/src/airflow/ui/src/utils/useFiltersHandler.ts
b/airflow-core/src/airflow/ui/src/utils/useFiltersHandler.ts
index ae5f48e5ef1..54f71369370 100644
--- a/airflow-core/src/airflow/ui/src/utils/useFiltersHandler.ts
+++ b/airflow-core/src/airflow/ui/src/utils/useFiltersHandler.ts
@@ -97,6 +97,7 @@ export type FilterableSearchParamsKeys =
| SearchParamsKeys.SUBJECT_SEARCH
| SearchParamsKeys.TASK_ID
| SearchParamsKeys.TASK_ID_PATTERN
+ | SearchParamsKeys.TEAMS
| SearchParamsKeys.TRIGGERING_USER_NAME_PATTERN
| SearchParamsKeys.TRY_NUMBER
| SearchParamsKeys.USER;
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
index 835a8ccb6dc..c8e6b65ccaf 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
@@ -1827,6 +1827,7 @@ class TestPostAssetMaterialize(TestAssets):
"triggering_user_name": "test",
"conf": {},
"note": None,
+ "team_name": None,
}
@pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py
index d8b7eeae5e1..f9893a3caaa 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py
@@ -25,7 +25,7 @@ from unittest import mock
import pytest
import time_machine
from fastapi.testclient import TestClient
-from sqlalchemy import func, select, update
+from sqlalchemy import delete, func, select, update
from airflow import plugins_manager
from airflow._shared.module_loading import qualname
@@ -53,6 +53,7 @@ from airflow.utils.types import DagRunTriggeredByType,
DagRunType
from tests_common.test_utils.api_fastapi import _check_dag_run_note,
_check_last_log
from tests_common.test_utils.asserts import assert_queries_count
+from tests_common.test_utils.config import conf_vars
from tests_common.test_utils.db import (
clear_db_assets,
clear_db_connections,
@@ -329,9 +330,39 @@ def get_dag_run_dict(run: DagRun):
"partition_date": (
from_datetime_to_zulu_without_ms(run.partition_date) if
run.partition_date else None
),
+ "team_name": None,
}
+def _attach_dag_to_team(session, dag_id: str, *, bundle_name: str, team_name:
str) -> str:
+ """
+ Associate a Dag with a team via a team-scoped bundle for multi-team tests.
+
+ Returns the Dag's original bundle name so the caller can restore it during
cleanup
+ (``DagModel.bundle_name`` is a foreign key with no ``ON DELETE`` action).
+ """
+ original_bundle_name =
session.scalar(select(DagModel.bundle_name).where(DagModel.dag_id == dag_id))
+ bundle = DagBundleModel(name=bundle_name)
+ bundle.teams.append(Team(name=team_name))
+ session.add(bundle)
+ session.flush()
+ session.execute(update(DagModel).where(DagModel.dag_id ==
dag_id).values(bundle_name=bundle_name))
+ session.commit()
+ return original_bundle_name
+
+
+def _detach_dag_from_team(
+ session, dag_id: str, *, bundle_name: str, team_name: str,
original_bundle_name: str
+) -> None:
+ """Undo :func:`_attach_dag_to_team`, restoring the Dag's original
bundle."""
+ session.execute(
+ update(DagModel).where(DagModel.dag_id ==
dag_id).values(bundle_name=original_bundle_name)
+ )
+ session.execute(delete(DagBundleModel).where(DagBundleModel.name ==
bundle_name))
+ session.execute(delete(Team).where(Team.name == team_name))
+ session.commit()
+
+
class TestGetDagRun:
@pytest.mark.parametrize(
("dag_id", "run_id", "state", "run_type", "triggered_by",
"dag_run_note"),
@@ -436,6 +467,53 @@ class TestGetDagRuns:
"partition_date_gte and partition_date_lte are not supported."
)
+ @conf_vars({("core", "multi_team"): "True"})
+ @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
+ def test_get_dag_runs_includes_team_name(self, test_client, session):
+ original_bundle_name = _attach_dag_to_team(
+ session, DAG1_ID, bundle_name="team-bundle-runs",
team_name="team-runs"
+ )
+ try:
+ response = test_client.get(f"/dags/{DAG1_ID}/dagRuns")
+ assert response.status_code == 200
+ body = response.json()
+ assert body["dag_runs"]
+ assert all(run["team_name"] == "team-runs" for run in
body["dag_runs"])
+ finally:
+ _detach_dag_from_team(
+ session,
+ DAG1_ID,
+ bundle_name="team-bundle-runs",
+ team_name="team-runs",
+ original_bundle_name=original_bundle_name,
+ )
+
+ @conf_vars({("core", "multi_team"): "True"})
+ @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
+ def test_get_dag_runs_filtered_by_team(self, test_client, session):
+ original_bundle_name = _attach_dag_to_team(
+ session, DAG1_ID, bundle_name="team-bundle-filter",
team_name="team-filter"
+ )
+ try:
+ response = test_client.get("/dags/~/dagRuns", params={"teams":
["team-filter"]})
+ assert response.status_code == 200
+ body = response.json()
+ assert body["total_entries"] == 2
+ assert {run["dag_id"] for run in body["dag_runs"]} == {DAG1_ID}
+
+ # A team with no Dags returns nothing.
+ response = test_client.get("/dags/~/dagRuns", params={"teams":
["nonexistent-team"]})
+ assert response.status_code == 200
+ assert response.json()["total_entries"] == 0
+ finally:
+ _detach_dag_from_team(
+ session,
+ DAG1_ID,
+ bundle_name="team-bundle-filter",
+ team_name="team-filter",
+ original_bundle_name=original_bundle_name,
+ )
+
def test_invalid_order_by_raises_400(self, test_client):
response = test_client.get("/dags/test_dag1/dagRuns?order_by=invalid")
assert response.status_code == 400
@@ -3308,6 +3386,7 @@ class TestTriggerDagRun:
"triggering_user_name": "test",
"partition_key": None,
"partition_date": None,
+ "team_name": None,
}
assert response.json() == expected_response_json
@@ -3549,6 +3628,7 @@ class TestTriggerDagRun:
"note": note,
"partition_key": None,
"partition_date": None,
+ "team_name": None,
}
assert response_2.status_code == 409
@@ -3639,6 +3719,7 @@ class TestTriggerDagRun:
"note": None,
"partition_key": None,
"partition_date": None,
+ "team_name": None,
}
@pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
index 6d21cc8c6c4..240c9334e80 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
@@ -266,6 +266,7 @@ def expected_sample_hitl_detail_dict(sample_ti:
TaskInstance) -> dict[str, Any]:
"state": None,
"task_display_name": "sample_task_hitl",
"task_id": TASK_ID,
+ "team_name": None,
"trigger": None,
"triggerer_job": None,
"try_number": 0,
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
index ae3a8388fea..f68acf75d60 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
@@ -76,6 +76,35 @@ DEFAULT_DATETIME_1 =
dt.datetime.fromisoformat(DEFAULT_DATETIME_STR_1)
DEFAULT_DATETIME_2 = dt.datetime.fromisoformat(DEFAULT_DATETIME_STR_2)
+def _attach_dag_to_team(session, dag_id: str, *, bundle_name: str, team_name:
str) -> str:
+ """
+ Associate a Dag with a team via a team-scoped bundle for multi-team tests.
+
+ Returns the Dag's original bundle name so the caller can restore it during
cleanup
+ (``DagModel.bundle_name`` is a foreign key with no ``ON DELETE`` action).
+ """
+ original_bundle_name =
session.scalar(select(DagModel.bundle_name).where(DagModel.dag_id == dag_id))
+ bundle = DagBundleModel(name=bundle_name)
+ bundle.teams.append(Team(name=team_name))
+ session.add(bundle)
+ session.flush()
+ session.execute(update(DagModel).where(DagModel.dag_id ==
dag_id).values(bundle_name=bundle_name))
+ session.commit()
+ return original_bundle_name
+
+
+def _detach_dag_from_team(
+ session, dag_id: str, *, bundle_name: str, team_name: str,
original_bundle_name: str
+) -> None:
+ """Undo :func:`_attach_dag_to_team`, restoring the Dag's original
bundle."""
+ session.execute(
+ update(DagModel).where(DagModel.dag_id ==
dag_id).values(bundle_name=original_bundle_name)
+ )
+ session.execute(delete(DagBundleModel).where(DagBundleModel.name ==
bundle_name))
+ session.execute(delete(Team).where(Team.name == team_name))
+ session.commit()
+
+
class TestTaskInstanceEndpoint:
@staticmethod
def clear_db():
@@ -242,6 +271,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
def test_should_respond_200_with_decorator(self, test_client, session):
@@ -318,6 +348,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"run_after": mock.ANY,
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
"dag_version": {
"id": response_data["dag_version"]["id"],
"version_number": expected_version_number,
@@ -411,6 +442,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"state": "running",
"unixname": getuser(),
},
+ "team_name": None,
}
def test_should_respond_200_with_task_state_in_removed(self, test_client,
session):
@@ -464,6 +496,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
def test_should_respond_200_task_instance_with_rendered(self, test_client,
session):
@@ -520,6 +553,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
def test_raises_404_for_nonexistent_task_instance(self, test_client):
@@ -640,6 +674,7 @@ class TestGetMappedTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
def test_should_respond_401(self, unauthenticated_test_client):
@@ -2255,6 +2290,60 @@ class TestGetTaskInstances(TestTaskInstanceEndpoint):
all_backward = backward_ids + [ti["id"] for ti in
forward_pages[-1]["task_instances"]]
assert all_backward == forward_ids, "Backward walk + last page must
match the forward walk exactly"
+ @conf_vars({("core", "multi_team"): "True"})
+ def test_should_include_team_name(self, test_client, session):
+ self.create_task_instances(session)
+ original_bundle_name = _attach_dag_to_team(
+ session, "example_python_operator", bundle_name="team-bundle-tis",
team_name="team-tis"
+ )
+ try:
+ response =
test_client.get(f"/dags/{'example_python_operator'}/dagRuns/~/taskInstances")
+ assert response.status_code == 200
+ body = response.json()
+ assert body["task_instances"]
+ assert all(ti["team_name"] == "team-tis" for ti in
body["task_instances"])
+ finally:
+ _detach_dag_from_team(
+ session,
+ "example_python_operator",
+ bundle_name="team-bundle-tis",
+ team_name="team-tis",
+ original_bundle_name=original_bundle_name,
+ )
+
+ @conf_vars({("core", "multi_team"): "True"})
+ def test_should_filter_by_team(self, test_client, session):
+ self.create_task_instances(session)
+ original_bundle_name = _attach_dag_to_team(
+ session,
+ "example_python_operator",
+ bundle_name="team-bundle-tis-filter",
+ team_name="team-tis-filter",
+ )
+ try:
+ response = test_client.get(
+ "/dags/~/dagRuns/~/taskInstances", params={"teams":
["team-tis-filter"]}
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["total_entries"] > 0
+ assert all(ti["dag_id"] == "example_python_operator" for ti in
body["task_instances"])
+
+ # A team with no Dags returns nothing.
+ response = test_client.get(
+ "/dags/~/dagRuns/~/taskInstances", params={"teams":
["nonexistent-team"]}
+ )
+ assert response.status_code == 200
+ assert response.json()["total_entries"] == 0
+ finally:
+ _detach_dag_from_team(
+ session,
+ "example_python_operator",
+ bundle_name="team-bundle-tis-filter",
+ team_name="team-tis-filter",
+ original_bundle_name=original_bundle_name,
+ )
+
class TestGetTaskDependencies(TestTaskInstanceEndpoint):
def setup_method(self):
@@ -3817,6 +3906,7 @@ class
TestPostClearTaskInstances(TestTaskInstanceEndpoint):
"task_display_name": "print_the_context",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
"try_number": 0,
"unixname": getuser(),
},
@@ -4757,6 +4847,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
],
"total_entries": 1,
@@ -5033,6 +5124,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
],
"total_entries": 1,
@@ -5171,6 +5263,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
],
"total_entries": 1,
@@ -5234,6 +5327,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
],
"total_entries": 1,
@@ -5329,6 +5423,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
],
"total_entries": 1,
@@ -5412,6 +5507,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
_check_task_instance_note(
@@ -5607,6 +5703,7 @@ class
TestPatchTaskInstanceDryRun(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
],
"total_entries": 1,
@@ -5895,6 +5992,7 @@ class
TestPatchTaskInstanceDryRun(TestTaskInstanceEndpoint):
"run_after": "2020-01-01T00:00:00Z",
"trigger": None,
"triggerer_job": None,
+ "team_name": None,
}
],
"total_entries": 1,
diff --git a/airflow-core/tests/unit/cli/commands/test_asset_command.py
b/airflow-core/tests/unit/cli/commands/test_asset_command.py
index a7329a96251..b6b7fd793fa 100644
--- a/airflow-core/tests/unit/cli/commands/test_asset_command.py
+++ b/airflow-core/tests/unit/cli/commands/test_asset_command.py
@@ -164,6 +164,7 @@ def test_cli_assets_materialize(mock_hasattr, parser:
ArgumentParser, stdout_cap
"run_type": "asset_materialization",
"start_date": None,
"state": "queued",
+ "team_name": None,
"triggered_by": "cli",
"triggering_user_name": "root",
"run_after": "2025-02-12T19:27:59.066046Z",
@@ -204,6 +205,7 @@ def
test_cli_assets_materialize_with_view_url_template(parser: ArgumentParser, s
"run_type": "asset_materialization",
"start_date": None,
"state": "queued",
+ "team_name": None,
"triggered_by": "cli",
"triggering_user_name": "root",
"run_after": "2025-02-12T19:27:59.066046Z",
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 3d3fea1fc43..d7b63613a86 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -1858,6 +1858,7 @@ class DAGRunResponse(BaseModel):
dag_display_name: Annotated[str, Field(title="Dag Display Name")]
partition_key: Annotated[str | None, Field(title="Partition Key")] = None
partition_date: Annotated[datetime | None, Field(title="Partition Date")]
= None
+ team_name: Annotated[str | None, Field(title="Team Name")] = None
class DAGRunsBatchBody(BaseModel):
@@ -2151,6 +2152,7 @@ class TaskInstanceResponse(BaseModel):
trigger: TriggerResponse | None = None
triggerer_job: JobResponse | None = None
dag_version: DagVersionResponse | None = None
+ team_name: Annotated[str | None, Field(title="Team Name")] = None
class TaskResponse(BaseModel):