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 d776f003b71 UI: Add team column and filter to the jobs list (#73295)
d776f003b71 is described below
commit d776f003b7186842f32245236a90dbb51f5eee6d
Author: Vincent <[email protected]>
AuthorDate: Tue Sep 29 10:20:56 2026 -0400
UI: Add team column and filter to the jobs list (#73295)
When multi-team mode is enabled, operators triaging scheduler, triggerer
and Dag
processor jobs need to see which team's workloads each job belongs to, and
to
scope the list to a single team when investigating that team's resource
usage.
The column and the filter render only when the `multi_team` configuration is
enabled, so single-team deployments are unaffected.
---
.../api_fastapi/common/parameters/__init__.py | 3 ++
.../airflow/api_fastapi/common/parameters/job.py | 50 +++++++++++++++++++
.../core_api/openapi/v2-rest-api-generated.yaml | 8 +++
.../api_fastapi/core_api/routes/public/job.py | 3 ++
.../airflow/cli/commands/dag_processor_command.py | 26 +++++++++-
.../src/airflow/ui/openapi-gen/queries/common.ts | 5 +-
.../ui/openapi-gen/queries/ensureQueryData.ts | 6 ++-
.../src/airflow/ui/openapi-gen/queries/prefetch.ts | 6 ++-
.../src/airflow/ui/openapi-gen/queries/queries.ts | 6 ++-
.../src/airflow/ui/openapi-gen/queries/suspense.ts | 6 ++-
.../ui/openapi-gen/requests/services.gen.ts | 4 +-
.../airflow/ui/openapi-gen/requests/types.gen.ts | 1 +
airflow-core/src/airflow/ui/src/pages/Jobs.tsx | 22 ++++++--
.../api_fastapi/core_api/routes/public/test_job.py | 17 +++++++
.../cli/commands/test_dag_processor_command.py | 58 ++++++++++++++++++++++
15 files changed, 206 insertions(+), 15 deletions(-)
diff --git a/airflow-core/src/airflow/api_fastapi/common/parameters/__init__.py
b/airflow-core/src/airflow/api_fastapi/common/parameters/__init__.py
index 39e748a2972..deae80a25ac 100644
--- a/airflow-core/src/airflow/api_fastapi/common/parameters/__init__.py
+++ b/airflow-core/src/airflow/api_fastapi/common/parameters/__init__.py
@@ -86,6 +86,9 @@ from airflow.api_fastapi.common.parameters.filter import (
FilterParam as FilterParam,
filter_param_factory as filter_param_factory,
)
+from airflow.api_fastapi.common.parameters.job import (
+ QueryJobTeamsFilter as QueryJobTeamsFilter,
+)
from airflow.api_fastapi.common.parameters.misc import (
QueryConnectionIdPatternSearch as QueryConnectionIdPatternSearch,
QueryConnectionIdPrefixPatternSearch as
QueryConnectionIdPrefixPatternSearch,
diff --git a/airflow-core/src/airflow/api_fastapi/common/parameters/job.py
b/airflow-core/src/airflow/api_fastapi/common/parameters/job.py
new file mode 100644
index 00000000000..6c418d2da48
--- /dev/null
+++ b/airflow-core/src/airflow/api_fastapi/common/parameters/job.py
@@ -0,0 +1,50 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+from __future__ import annotations
+
+from typing import TYPE_CHECKING, Annotated
+
+from fastapi import Depends, Query
+from sqlalchemy import select as sql_select
+
+from airflow.api_fastapi.common.parameters.base import BaseParam
+from airflow.jobs.job import Job
+from airflow.models.team import JobTeam
+
+if TYPE_CHECKING:
+ from sqlalchemy.sql import Select
+
+
+class _JobTeamsFilter(BaseParam[list[str]]):
+ """Filter jobs by the teams they serve (via the ``JobTeam``
association)."""
+
+ 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
+
+ return
select.where(Job.id.in_(sql_select(JobTeam.job_id).where(JobTeam.team_name.in_(self.value))))
+
+ @classmethod
+ def depends(cls, teams: list[str] = Query(default_factory=list)) ->
_JobTeamsFilter:
+ return cls().set_value(teams)
+
+
+QueryJobTeamsFilter = Annotated[_JobTeamsFilter,
Depends(_JobTeamsFilter.depends)]
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 415b36445d7..22f1afe6903 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
@@ -5506,6 +5506,14 @@ paths:
- type: string
- type: 'null'
title: Executor Class
+ - name: teams
+ in: query
+ required: false
+ schema:
+ type: array
+ items:
+ type: string
+ title: Teams
responses:
'200':
description: Successful Response
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py
index 3e7e16bb69a..3760b461b23 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py
@@ -28,6 +28,7 @@ from airflow.api_fastapi.common.db.common import (
)
from airflow.api_fastapi.common.parameters import (
FilterParam,
+ QueryJobTeamsFilter,
QueryLimit,
QueryOffset,
RangeFilter,
@@ -101,6 +102,7 @@ def get_jobs(
FilterParam[str | None],
Depends(filter_param_factory(Job.executor_class, str | None,
filter_name="executor_class")),
],
+ teams: QueryJobTeamsFilter,
is_alive: bool | None = None,
) -> JobCollectionResponse:
"""Get all jobs."""
@@ -116,6 +118,7 @@ def get_jobs(
job_type,
hostname,
executor_class,
+ teams,
],
order_by=order_by,
limit=limit,
diff --git a/airflow-core/src/airflow/cli/commands/dag_processor_command.py
b/airflow-core/src/airflow/cli/commands/dag_processor_command.py
index 104065d6d78..5bf28fb0c28 100644
--- a/airflow-core/src/airflow/cli/commands/dag_processor_command.py
+++ b/airflow-core/src/airflow/cli/commands/dag_processor_command.py
@@ -22,6 +22,8 @@ import logging
from typing import Any
from airflow.cli.commands.daemon_utils import run_command_with_daemon_option
+from airflow.configuration import conf
+from airflow.dag_processing.bundles.manager import
_get_configured_bundle_team_names
from airflow.dag_processing.manager import DagFileProcessorManager
from airflow.jobs.dag_processor_job_runner import DagProcessorJobRunner
from airflow.jobs.job import Job, run_job
@@ -33,12 +35,34 @@ from airflow.utils.providers_configuration_loader import
providers_configuration
log = logging.getLogger(__name__)
+def _get_team_names(bundle_names: list[str] | None) -> list[str]:
+ """
+ Return the teams this Dag processor serves, sorted and de-duplicated.
+
+ Teams are resolved from the bundle configuration rather than the metadata
DB: the job row is
+ written before ``sync_bundles()`` runs, so a DB lookup would see no rows
on a fresh deployment
+ and stale rows right after a bundle is reassigned in config. Config is the
source of truth, and
+ is what ``airflow_health.py`` resolves teams from too.
+
+ A processor started without ``--bundle-name`` parses every configured
bundle, so it serves
+ every configured team. A bundle mapped to no team contributes no team, so
a processor parsing
+ only team-less (or unknown) bundles is not team-scoped and serves the
empty list. Outside
+ multi-team mode team scoping is disabled entirely, matching
``DagFileProcessorManager``.
+ """
+ if not conf.getboolean("core", "multi_team"):
+ return []
+
+ configured = _get_configured_bundle_team_names()
+ names = bundle_names or list(configured)
+ return sorted({team for name in names if (team := configured.get(name)) is
not None})
+
+
def _create_dag_processor_job_runner(args: Any) -> DagProcessorJobRunner:
"""Create DagFileProcessorProcess instance."""
if args.bundle_name:
cli_utils.validate_dag_bundle_arg(args.bundle_name)
return DagProcessorJobRunner(
- job=Job(bundle_names=args.bundle_name),
+ job=Job(bundle_names=args.bundle_name,
team_names=_get_team_names(args.bundle_name)),
processor=DagFileProcessorManager(
max_runs=args.num_runs,
bundle_names_to_parse=args.bundle_name,
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 574871897ab..e1de3ccdd0d 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -786,7 +786,7 @@ export const UseImportErrorServiceGetImportErrorsKeyFn = ({
bundleName, filename
export type JobServiceGetJobsDefaultResponse = Awaited<ReturnType<typeof
JobService.getJobs>>;
export type JobServiceGetJobsQueryResult<TData =
JobServiceGetJobsDefaultResponse, TError = unknown> = UseQueryResult<TData,
TError>;
export const useJobServiceGetJobsKey = "JobServiceGetJobs";
-export const UseJobServiceGetJobsKeyFn = ({ dagId, endDateGt, endDateGte,
endDateLt, endDateLte, executorClass, hostname, isAlive, jobState, jobType,
limit, offset, orderBy, startDateGt, startDateGte, startDateLt, startDateLte }:
{
+export const UseJobServiceGetJobsKeyFn = ({ dagId, endDateGt, endDateGte,
endDateLt, endDateLte, executorClass, hostname, isAlive, jobState, jobType,
limit, offset, orderBy, startDateGt, startDateGte, startDateLt, startDateLte,
teams }: {
dagId?: string;
endDateGt?: string;
endDateGte?: string;
@@ -804,7 +804,8 @@ export const UseJobServiceGetJobsKeyFn = ({ dagId,
endDateGt, endDateGte, endDat
startDateGte?: string;
startDateLt?: string;
startDateLte?: string;
-} = {}, queryKey?: Array<unknown>) => [useJobServiceGetJobsKey, ...(queryKey
?? [{ dagId, endDateGt, endDateGte, endDateLt, endDateLte, executorClass,
hostname, isAlive, jobState, jobType, limit, offset, orderBy, startDateGt,
startDateGte, startDateLt, startDateLte }])];
+ teams?: string[];
+} = {}, queryKey?: Array<unknown>) => [useJobServiceGetJobsKey, ...(queryKey
?? [{ dagId, endDateGt, endDateGte, endDateLt, endDateLte, executorClass,
hostname, isAlive, jobState, jobType, limit, offset, orderBy, startDateGt,
startDateGte, startDateLt, startDateLte, teams }])];
export type PluginServiceGetPluginsDefaultResponse = Awaited<ReturnType<typeof
PluginService.getPlugins>>;
export type PluginServiceGetPluginsQueryResult<TData =
PluginServiceGetPluginsDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
export const usePluginServiceGetPluginsKey = "PluginServiceGetPlugins";
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 01983ce4e99..58c12d2a670 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
@@ -1561,10 +1561,11 @@ export const
ensureUseImportErrorServiceGetImportErrorsData = (queryClient: Quer
* @param data.jobType
* @param data.hostname
* @param data.executorClass
+* @param data.teams
* @returns JobCollectionResponse Successful Response
* @throws ApiError
*/
-export const ensureUseJobServiceGetJobsData = (queryClient: QueryClient, {
dagId, endDateGt, endDateGte, endDateLt, endDateLte, executorClass, hostname,
isAlive, jobState, jobType, limit, offset, orderBy, startDateGt, startDateGte,
startDateLt, startDateLte }: {
+export const ensureUseJobServiceGetJobsData = (queryClient: QueryClient, {
dagId, endDateGt, endDateGte, endDateLt, endDateLte, executorClass, hostname,
isAlive, jobState, jobType, limit, offset, orderBy, startDateGt, startDateGte,
startDateLt, startDateLte, teams }: {
dagId?: string;
endDateGt?: string;
endDateGte?: string;
@@ -1582,7 +1583,8 @@ export const ensureUseJobServiceGetJobsData =
(queryClient: QueryClient, { dagId
startDateGte?: string;
startDateLt?: string;
startDateLte?: string;
-} = {}) => queryClient.ensureQueryData({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte }), queryFn: ()
=> JobService.getJobs({ dagId, endDateGt, endDateGte, endDateLt, endDateLte,
executorClass, hostname, isAlive, jobState, jobType, limit, offset, orderBy,
startDateGt, startDateGte, startDateLt, startDateLte }) });
+ teams?: string[];
+} = {}) => queryClient.ensureQueryData({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte, teams }),
queryFn: () => JobService.getJobs({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startD [...]
/**
* Get Plugins
* @param data The data for the request.
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 501689dc55d..b04f0e549d1 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
@@ -1561,10 +1561,11 @@ export const
prefetchUseImportErrorServiceGetImportErrors = (queryClient: QueryC
* @param data.jobType
* @param data.hostname
* @param data.executorClass
+* @param data.teams
* @returns JobCollectionResponse Successful Response
* @throws ApiError
*/
-export const prefetchUseJobServiceGetJobs = (queryClient: QueryClient, {
dagId, endDateGt, endDateGte, endDateLt, endDateLte, executorClass, hostname,
isAlive, jobState, jobType, limit, offset, orderBy, startDateGt, startDateGte,
startDateLt, startDateLte }: {
+export const prefetchUseJobServiceGetJobs = (queryClient: QueryClient, {
dagId, endDateGt, endDateGte, endDateLt, endDateLte, executorClass, hostname,
isAlive, jobState, jobType, limit, offset, orderBy, startDateGt, startDateGte,
startDateLt, startDateLte, teams }: {
dagId?: string;
endDateGt?: string;
endDateGte?: string;
@@ -1582,7 +1583,8 @@ export const prefetchUseJobServiceGetJobs = (queryClient:
QueryClient, { dagId,
startDateGte?: string;
startDateLt?: string;
startDateLte?: string;
-} = {}) => queryClient.prefetchQuery({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte }), queryFn: ()
=> JobService.getJobs({ dagId, endDateGt, endDateGte, endDateLt, endDateLte,
executorClass, hostname, isAlive, jobState, jobType, limit, offset, orderBy,
startDateGt, startDateGte, startDateLt, startDateLte }) });
+ teams?: string[];
+} = {}) => queryClient.prefetchQuery({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte, teams }),
queryFn: () => JobService.getJobs({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDat [...]
/**
* Get Plugins
* @param data The data for the request.
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 4951fc96970..df1286bea25 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -1561,10 +1561,11 @@ export const useImportErrorServiceGetImportErrors =
<TData = Common.ImportErrorS
* @param data.jobType
* @param data.hostname
* @param data.executorClass
+* @param data.teams
* @returns JobCollectionResponse Successful Response
* @throws ApiError
*/
-export const useJobServiceGetJobs = <TData =
Common.JobServiceGetJobsDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte }: {
+export const useJobServiceGetJobs = <TData =
Common.JobServiceGetJobsDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte, teams }: {
dagId?: string;
endDateGt?: string;
endDateGte?: string;
@@ -1582,7 +1583,8 @@ export const useJobServiceGetJobs = <TData =
Common.JobServiceGetJobsDefaultResp
startDateGte?: string;
startDateLt?: string;
startDateLte?: string;
-} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte }, queryKey),
queryFn: () => JobService.getJobs({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAli [...]
+ teams?: string[];
+} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte, teams },
queryKey), queryFn: () => JobService.getJobs({ dagId, endDateGt, endDateGte,
endDateLt, endDateLte, executorClass, hostname [...]
/**
* Get Plugins
* @param data The data for the request.
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 6c91b91e143..31da0a3e582 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
@@ -1561,10 +1561,11 @@ export const
useImportErrorServiceGetImportErrorsSuspense = <TData = Common.Impo
* @param data.jobType
* @param data.hostname
* @param data.executorClass
+* @param data.teams
* @returns JobCollectionResponse Successful Response
* @throws ApiError
*/
-export const useJobServiceGetJobsSuspense = <TData =
Common.JobServiceGetJobsDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte }: {
+export const useJobServiceGetJobsSuspense = <TData =
Common.JobServiceGetJobsDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte, teams }: {
dagId?: string;
endDateGt?: string;
endDateGte?: string;
@@ -1582,7 +1583,8 @@ export const useJobServiceGetJobsSuspense = <TData =
Common.JobServiceGetJobsDef
startDateGte?: string;
startDateLt?: string;
startDateLte?: string;
-} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte }, queryKey),
queryFn: () => JobService.getJobs({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostnam [...]
+ teams?: string[];
+} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseJobServiceGetJobsKeyFn({ dagId, endDateGt, endDateGte, endDateLt,
endDateLte, executorClass, hostname, isAlive, jobState, jobType, limit, offset,
orderBy, startDateGt, startDateGte, startDateLt, startDateLte, teams },
queryKey), queryFn: () => JobService.getJobs({ dagId, endDateGt, endDateGte,
endDateLt, endDateLte, executorClass, [...]
/**
* Get Plugins
* @param data The data for the request.
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 d0a1a936ed9..86596b1b890 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
@@ -3641,6 +3641,7 @@ export class JobService {
* @param data.jobType
* @param data.hostname
* @param data.executorClass
+ * @param data.teams
* @returns JobCollectionResponse Successful Response
* @throws ApiError
*/
@@ -3665,7 +3666,8 @@ export class JobService {
dag_id: data.dagId,
job_type: data.jobType,
hostname: data.hostname,
- executor_class: data.executorClass
+ executor_class: data.executorClass,
+ teams: data.teams
},
errors: {
400: 'Bad Request',
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 7e33899b2ea..6dcc310a469 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
@@ -4558,6 +4558,7 @@ export type GetJobsData = {
startDateGte?: string | null;
startDateLt?: string | null;
startDateLte?: string | null;
+ teams?: Array<(string)>;
};
export type GetJobsResponse = JobCollectionResponse;
diff --git a/airflow-core/src/airflow/ui/src/pages/Jobs.tsx
b/airflow-core/src/airflow/ui/src/pages/Jobs.tsx
index 0db02b1e26c..7d66ab08fee 100644
--- a/airflow-core/src/airflow/ui/src/pages/Jobs.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Jobs.tsx
@@ -33,9 +33,10 @@ import { StateBadge } from "src/components/StateBadge";
import Time from "src/components/Time";
import { SearchParamsKeys } from "src/constants/searchParams";
+import { useConfig } from "src/queries/useConfig";
import { useDocumentTitle, useFiltersHandler, type FilterableSearchParamsKeys
} from "src/utils";
-const createColumns = (translate: TFunction): Array<ColumnDef<JobResponse>> =>
[
+const createColumns = (translate: TFunction, multiTeam: boolean):
Array<ColumnDef<JobResponse>> => [
{
accessorKey: "id",
header: translate("jobs.columns.id"),
@@ -44,6 +45,16 @@ const createColumns = (translate: TFunction):
Array<ColumnDef<JobResponse>> => [
accessorKey: "job_type",
header: translate("jobs.columns.jobType"),
},
+ ...(multiTeam
+ ? ([
+ {
+ accessorKey: "team_names",
+ cell: ({ row: { original } }) => original.team_names?.join(", "),
+ enableSorting: false,
+ header: translate("common:dagDetails.team"),
+ },
+ ] as Array<ColumnDef<JobResponse>>)
+ : []),
{
accessorKey: "state",
cell: ({
@@ -97,15 +108,18 @@ const jobsFilterKeys: Array<FilterableSearchParamsKeys> = [
export const Jobs = () => {
const { t: translate } = useTranslation(["admin", "common"]);
+ const multiTeamEnabled = Boolean(useConfig("multi_team"));
useDocumentTitle(translate("common:browse.jobs"));
const { setTableURLState, tableURLState } = useTableURLState();
const [searchParams] = useSearchParams();
- const { filterConfigs, handleFiltersChange, initialValues } =
useFiltersHandler(jobsFilterKeys);
+ const { filterConfigs, handleFiltersChange, initialValues } =
useFiltersHandler(
+ multiTeamEnabled ? [...jobsFilterKeys, SearchParamsKeys.TEAMS] :
jobsFilterKeys,
+ );
- const columns = createColumns(translate);
+ const columns = createColumns(translate, multiTeamEnabled);
const { pagination, sorting } = tableURLState;
const [sort] = sorting;
@@ -119,6 +133,7 @@ export const Jobs = () => {
const filteredStartDateLte =
searchParams.get(SearchParamsKeys.START_DATE_LTE);
const filteredEndDateGte = searchParams.get(SearchParamsKeys.END_DATE_GTE);
const filteredEndDateLte = searchParams.get(SearchParamsKeys.END_DATE_LTE);
+ const teams = searchParams.getAll(SearchParamsKeys.TEAMS);
const { data, error, isFetching, isLoading } = useJobServiceGetJobs({
endDateGte: filteredEndDateGte ?? undefined,
@@ -132,6 +147,7 @@ export const Jobs = () => {
orderBy,
startDateGte: filteredStartDateGte ?? undefined,
startDateLte: filteredStartDateLte ?? undefined,
+ teams: teams.length > 0 ? teams : undefined,
});
return (
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_job.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_job.py
index f969d48b305..7c51ac5c6f7 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_job.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_job.py
@@ -222,6 +222,23 @@ class TestGetJobs(TestJobEndpoint):
assert response_json["jobs"][0]["team_names"] ==
sorted([testing_team.name, extra_team.name])
assert response_json["jobs"][0]["bundle_names"] == ["bundle-a",
"bundle-b"]
+ def test_get_jobs_filters_by_teams(self, test_client, session: Session,
testing_team):
+ clear_db_jobs()
+ session.add_all(
+ [
+ Job(state=JobState.RUNNING, job_type="SchedulerJob",
team_names=[testing_team.name]),
+ Job(state=JobState.RUNNING, job_type="SchedulerJob"),
+ ]
+ )
+ session.commit()
+
+ response = test_client.get("/jobs", params={"teams":
[testing_team.name]})
+
+ assert response.status_code == 200
+ response_json = response.json()
+ assert response_json["total_entries"] == 1
+ assert response_json["jobs"][0]["team_names"] == [testing_team.name]
+
def test_should_raises_401_unauthenticated(self,
unauthenticated_test_client):
response = unauthenticated_test_client.get("/jobs")
assert response.status_code == 401
diff --git a/airflow-core/tests/unit/cli/commands/test_dag_processor_command.py
b/airflow-core/tests/unit/cli/commands/test_dag_processor_command.py
index 1631001b594..7f8d963e108 100644
--- a/airflow-core/tests/unit/cli/commands/test_dag_processor_command.py
+++ b/airflow-core/tests/unit/cli/commands/test_dag_processor_command.py
@@ -29,6 +29,26 @@ from tests_common.test_utils.config import conf_vars
pytestmark = pytest.mark.db_test
+# The Dag processor resolves teams from the configured bundle partition, not
the metadata DB, so
+# that a fresh or reassigned deployment gets the right team before
``sync_bundles()`` has run.
+BUNDLE_TEAMS = {
+ "bundle_a": "team_a",
+ "other_bundle_a": "team_a",
+ "bundle_b": "team_b",
+ "global_bundle": None,
+}
+
+
[email protected]
+def dag_bundles_with_teams():
+ with (
+ conf_vars({("core", "multi_team"): "True"}),
+ mock.patch.object(
+ dag_processor_command, "_get_configured_bundle_team_names",
return_value=dict(BUNDLE_TEAMS)
+ ),
+ ):
+ yield
+
class TestDagProcessorCommand:
"""
@@ -58,6 +78,44 @@ class TestDagProcessorCommand:
assert mock_runner.call_args.kwargs["processor"].bundle_names_to_parse
== ["testing"]
assert mock_runner.call_args.kwargs["job"].bundle_names == ["testing"]
+ @pytest.mark.usefixtures("dag_bundles_with_teams")
+ @pytest.mark.parametrize(
+ ("bundle_names", "expected_team_names"),
+ [
+ pytest.param(None, ["team_a", "team_b"],
id="every-configured-bundle"),
+ pytest.param(["bundle_a"], ["team_a"], id="one-bundle-of-a-team"),
+ pytest.param(["bundle_a", "other_bundle_a"], ["team_a"],
id="several-bundles-of-one-team"),
+ pytest.param(["bundle_a", "bundle_b"], ["team_a", "team_b"],
id="bundles-of-several-teams"),
+ pytest.param(["bundle_a", "global_bundle"], ["team_a"],
id="bundles-of-a-team-and-of-no-team"),
+ pytest.param(["global_bundle"], [], id="bundle-of-no-team"),
+ pytest.param(["unknown_bundle"], [], id="unknown-bundle"),
+ ],
+ )
+ def test_get_team_names(self, bundle_names, expected_team_names):
+ assert dag_processor_command._get_team_names(bundle_names) ==
expected_team_names
+
+ @conf_vars({("core", "multi_team"): "False"})
+ @mock.patch.object(
+ dag_processor_command, "_get_configured_bundle_team_names",
return_value=dict(BUNDLE_TEAMS)
+ )
+ def test_get_team_names_returns_empty_outside_multi_team(self,
mock_configured):
+ assert dag_processor_command._get_team_names(["bundle_a"]) == []
+ mock_configured.assert_not_called()
+
+ @conf_vars({("core", "load_examples"): "False"})
+
@mock.patch("airflow.cli.commands.dag_processor_command.DagProcessorJobRunner")
+ @mock.patch("airflow.utils.cli.validate_dag_bundle_arg")
+ @pytest.mark.usefixtures("dag_bundles_with_teams")
+ def test_job_records_the_teams_owning_the_parsed_bundles(self, _,
mock_runner):
+ mock_runner.return_value.job_type = "DagProcessorJob"
+ args = self.parser.parse_args(
+ ["dag-processor", "--bundle-name", "bundle_a", "--bundle-name",
"bundle_b"]
+ )
+
+ dag_processor_command.dag_processor(args)
+
+ assert mock_runner.call_args.kwargs["job"].team_names == ["team_a",
"team_b"]
+
@conf_vars({("core", "load_examples"): "False"})
@mock.patch("airflow.cli.commands.dag_processor_command.DagProcessorJobRunner")
@mock.patch("airflow.utils.cli.DagBundlesManager", autospec=True)