This is an automated email from the ASF dual-hosted git repository.

bbovenzi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 9e67831f2bb Allow filtering the Dag runs list by Dag tags (#70861)
9e67831f2bb is described below

commit 9e67831f2bb64ec89a0860400322a2826a45712f
Author: Aaryan Mahajan <[email protected]>
AuthorDate: Thu Oct 1 01:14:53 2026 +0530

    Allow filtering the Dag runs list by Dag tags (#70861)
    
    * Allow filtering the Dag runs list by Dag tags
    
    Finding runs for Dags with a specific tag currently requires filtering
    by dag_id one at a time, since the cross-Dag runs table has no way to
    scope by tag. Add a tags filter (with an any/all match mode) to the
    dag runs API and wire it into the UI's runs table filter bar.
    
    * Pass the tags match mode from the Dag runs filter to the API
    
    The tags filter now comes from the shared typeahead editor, which offers an
    any/all match mode, so the runs table has to forward it as it does on the
    Dags list. Without that the mode chosen in the filter bar was silently
    ignored and every selection behaved as "any".
---
 .../api_fastapi/common/parameters/__init__.py      |  2 +
 .../airflow/api_fastapi/common/parameters/dag.py   | 46 ++++++++++++++++++
 .../core_api/openapi/v2-rest-api-generated.yaml    | 19 ++++++++
 .../api_fastapi/core_api/routes/public/dag_run.py  |  4 ++
 .../src/airflow/ui/openapi-gen/queries/common.ts   |  6 ++-
 .../ui/openapi-gen/queries/ensureQueryData.ts      |  8 +++-
 .../src/airflow/ui/openapi-gen/queries/prefetch.ts |  8 +++-
 .../src/airflow/ui/openapi-gen/queries/queries.ts  |  8 +++-
 .../src/airflow/ui/openapi-gen/queries/suspense.ts |  8 +++-
 .../ui/openapi-gen/requests/services.gen.ts        |  4 ++
 .../airflow/ui/openapi-gen/requests/types.gen.ts   |  2 +
 .../src/airflow/ui/src/mocks/handlers/dag_runs.ts  | 49 ++++++++++++++++++-
 .../airflow/ui/src/pages/DagRuns/DagRuns.test.tsx  | 27 +++++++++++
 .../src/airflow/ui/src/pages/DagRuns/DagRuns.tsx   |  6 +++
 .../ui/src/pages/DagRuns/DagRunsFilters.tsx        |  1 +
 .../core_api/routes/public/test_dag_run.py         | 56 +++++++++++++++++++++-
 16 files changed, 241 insertions(+), 13 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 8ec274d0078..a3fca665e9a 100644
--- a/airflow-core/src/airflow/api_fastapi/common/parameters/__init__.py
+++ b/airflow-core/src/airflow/api_fastapi/common/parameters/__init__.py
@@ -66,7 +66,9 @@ from airflow.api_fastapi.common.parameters.dag import (
     QueryTagsFilter as QueryTagsFilter,
     QueryTeamsFilter as QueryTeamsFilter,
     QueryTimetableTypePrefixPatternSearch as 
QueryTimetableTypePrefixPatternSearch,
+    _DagIdTagsFilter as _DagIdTagsFilter,
     _DagIdTeamsFilter as _DagIdTeamsFilter,
+    tags_filter_factory as tags_filter_factory,
     teams_filter_factory as teams_filter_factory,
 )
 from airflow.api_fastapi.common.parameters.dag_run import (
diff --git a/airflow-core/src/airflow/api_fastapi/common/parameters/dag.py 
b/airflow-core/src/airflow/api_fastapi/common/parameters/dag.py
index f6b0ecaa34c..e5c469e4eb4 100644
--- a/airflow-core/src/airflow/api_fastapi/common/parameters/dag.py
+++ b/airflow-core/src/airflow/api_fastapi/common/parameters/dag.py
@@ -216,6 +216,52 @@ def teams_filter_factory(
     return depends_teams_filter
 
 
+class _DagIdTagsFilter(BaseParam[_TagFilterModel]):
+    """Filter rows by Dag tags through their ``dag_id``."""
+
+    def __init__(
+        self,
+        dag_id_attribute: ColumnElement | InstrumentedAttribute,
+        value: _TagFilterModel | 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 or not self.value.tags:
+            return select
+
+        conditions = [DagModel.tags.any(DagTag.name == tag) for tag in 
self.value.tags]
+        operator = or_ if not self.value.tags_match_mode or 
self.value.tags_match_mode == "any" else and_
+        return select.where(
+            
self.dag_id_attribute.in_(sql_select(DagModel.dag_id).where(operator(*conditions)))
+        )
+
+    @classmethod
+    def depends(cls, *args: Any, **kwargs: Any) -> Self:
+        raise NotImplementedError("Use tags_filter_factory instead, depends is 
not implemented.")
+
+
+def tags_filter_factory(
+    dag_id_attribute: ColumnElement | InstrumentedAttribute,
+) -> Callable[[list[str], Literal["any", "all"] | None], _DagIdTagsFilter]:
+    """Build a ``tags`` filter that scopes rows by Dag tags through the given 
``dag_id`` column."""
+
+    def depends_tags_filter(
+        tags: list[str] = Query(default_factory=list),
+        tags_match_mode: Literal["any", "all"] | None = None,
+    ) -> _DagIdTagsFilter:
+        return _DagIdTagsFilter(dag_id_attribute).set_value(
+            _TagFilterModel(tags=tags, tags_match_mode=tags_match_mode)
+        )
+
+    return depends_tags_filter
+
+
 QueryPausedFilter = Annotated[
     FilterParam[bool | None],
     Depends(filter_param_factory(DagModel.is_paused, bool | None, 
filter_name="paused")),
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 c336ff3afa9..d73ae639a93 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
@@ -2625,6 +2625,25 @@ paths:
           items:
             type: string
           title: Teams
+      - name: tags
+        in: query
+        required: false
+        schema:
+          type: array
+          items:
+            type: string
+          title: Tags
+      - name: tags_match_mode
+        in: query
+        required: false
+        schema:
+          anyOf:
+          - enum:
+            - any
+            - all
+            type: string
+          - type: 'null'
+          title: Tags Match Mode
       - name: order_by
         in: query
         required: false
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 2515344c50d..458ef2fd7ce 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
@@ -64,6 +64,7 @@ from airflow.api_fastapi.common.parameters import (
     Range,
     RangeFilter,
     SortParam,
+    _DagIdTagsFilter,
     _DagIdTeamsFilter,
     _PrefixSearchParam,
     _SearchParam,
@@ -72,6 +73,7 @@ from airflow.api_fastapi.common.parameters import (
     float_range_filter_factory,
     prefix_search_param_factory,
     search_param_factory,
+    tags_filter_factory,
     teams_filter_factory,
 )
 from airflow.api_fastapi.common.router import AirflowRouter
@@ -521,6 +523,7 @@ def get_dag_runs(
         FilterParam[str | None], 
Depends(filter_param_factory(DagRun.bundle_version, str | None))
     ],
     teams: Annotated[_DagIdTeamsFilter, 
Depends(teams_filter_factory(DagRun.dag_id))],
+    tags: Annotated[_DagIdTagsFilter, 
Depends(tags_filter_factory(DagRun.dag_id))],
     order_by: Annotated[
         SortParam,
         Depends(
@@ -678,6 +681,7 @@ def get_dag_runs(
         partition_key_prefix_pattern,
         consuming_asset_pattern,
         teams,
+        tags,
     ]
 
     if use_cursor:
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 a9118eae857..5f604612c0a 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -202,7 +202,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;
@@ -242,6 +242,8 @@ export const UseDagRunServiceGetDagRunsKeyFn = ({ 
bundleVersion, confContains, c
   startDateLt?: string;
   startDateLte?: string;
   state?: string[];
+  tags?: string[];
+  tagsMatchMode?: "any" | "all";
   teams?: string[];
   triggeringUserNamePattern?: string;
   triggeringUserNamePrefixPattern?: string;
@@ -249,7 +251,7 @@ export const UseDagRunServiceGetDagRunsKeyFn = ({ 
bundleVersion, confContains, c
   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";
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 1f4cb527d80..05bbc8be639 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
@@ -424,6 +424,8 @@ export const ensureUseDagRunServiceGetDagRunData = 
(queryClient: QueryClient, {
 * @param data.dagVersion
 * @param data.bundleVersion
 * @param data.teams
+* @param data.tags
+* @param data.tagsMatchMode
 * @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 Case-insensitive substring match (SQL `ILIKE`). 
Slower than `run_id_prefix_pattern` on large tables — see "Filtering with 
pattern parameters".
 * @param data.runIdPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
@@ -437,7 +439,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;
@@ -477,6 +479,8 @@ export const ensureUseDagRunServiceGetDagRunsData = 
(queryClient: QueryClient, {
   startDateLt?: string;
   startDateLte?: string;
   state?: string[];
+  tags?: string[];
+  tagsMatchMode?: "any" | "all";
   teams?: string[];
   triggeringUserNamePattern?: string;
   triggeringUserNamePrefixPattern?: string;
@@ -484,7 +488,7 @@ export const ensureUseDagRunServiceGetDagRunsData = 
(queryClient: QueryClient, {
   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.
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 3092c1c8212..c1c988437f7 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
@@ -424,6 +424,8 @@ export const prefetchUseDagRunServiceGetDagRun = 
(queryClient: QueryClient, { da
 * @param data.dagVersion
 * @param data.bundleVersion
 * @param data.teams
+* @param data.tags
+* @param data.tagsMatchMode
 * @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 Case-insensitive substring match (SQL `ILIKE`). 
Slower than `run_id_prefix_pattern` on large tables — see "Filtering with 
pattern parameters".
 * @param data.runIdPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
@@ -437,7 +439,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;
@@ -477,6 +479,8 @@ export const prefetchUseDagRunServiceGetDagRuns = 
(queryClient: QueryClient, { b
   startDateLt?: string;
   startDateLte?: string;
   state?: string[];
+  tags?: string[];
+  tagsMatchMode?: "any" | "all";
   teams?: string[];
   triggeringUserNamePattern?: string;
   triggeringUserNamePrefixPattern?: string;
@@ -484,7 +488,7 @@ export const prefetchUseDagRunServiceGetDagRuns = 
(queryClient: QueryClient, { b
   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.
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 917a275bc78..8c816d0aa72 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -424,6 +424,8 @@ export const useDagRunServiceGetDagRun = <TData = 
Common.DagRunServiceGetDagRunD
 * @param data.dagVersion
 * @param data.bundleVersion
 * @param data.teams
+* @param data.tags
+* @param data.tagsMatchMode
 * @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 Case-insensitive substring match (SQL `ILIKE`). 
Slower than `run_id_prefix_pattern` on large tables — see "Filtering with 
pattern parameters".
 * @param data.runIdPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
@@ -437,7 +439,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;
@@ -477,6 +479,8 @@ export const useDagRunServiceGetDagRuns = <TData = 
Common.DagRunServiceGetDagRun
   startDateLt?: string;
   startDateLte?: string;
   state?: string[];
+  tags?: string[];
+  tagsMatchMode?: "any" | "all";
   teams?: string[];
   triggeringUserNamePattern?: string;
   triggeringUserNamePrefixPattern?: string;
@@ -484,7 +488,7 @@ export const useDagRunServiceGetDagRuns = <TData = 
Common.DagRunServiceGetDagRun
   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.
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 58d17ee7f74..646259af57b 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
@@ -424,6 +424,8 @@ export const useDagRunServiceGetDagRunSuspense = <TData = 
Common.DagRunServiceGe
 * @param data.dagVersion
 * @param data.bundleVersion
 * @param data.teams
+* @param data.tags
+* @param data.tagsMatchMode
 * @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 Case-insensitive substring match (SQL `ILIKE`). 
Slower than `run_id_prefix_pattern` on large tables — see "Filtering with 
pattern parameters".
 * @param data.runIdPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
@@ -437,7 +439,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;
@@ -477,6 +479,8 @@ export const useDagRunServiceGetDagRunsSuspense = <TData = 
Common.DagRunServiceG
   startDateLt?: string;
   startDateLte?: string;
   state?: string[];
+  tags?: string[];
+  tagsMatchMode?: "any" | "all";
   teams?: string[];
   triggeringUserNamePattern?: string;
   triggeringUserNamePrefixPattern?: string;
@@ -484,7 +488,7 @@ export const useDagRunServiceGetDagRunsSuspense = <TData = 
Common.DagRunServiceG
   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.
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 cab47101f2d..38227629a2c 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
@@ -1190,6 +1190,8 @@ export class DagRunService {
      * @param data.dagVersion
      * @param data.bundleVersion
      * @param data.teams
+     * @param data.tags
+     * @param data.tagsMatchMode
      * @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 Case-insensitive substring match (SQL 
`ILIKE`). Slower than `run_id_prefix_pattern` on large tables — see "Filtering 
with pattern parameters".
      * @param data.runIdPrefixPattern Case-sensitive, index-friendly prefix 
match. See "Filtering with pattern parameters".
@@ -1246,6 +1248,8 @@ export class DagRunService {
                 dag_version: data.dagVersion,
                 bundle_version: data.bundleVersion,
                 teams: data.teams,
+                tags: data.tags,
+                tags_match_mode: data.tagsMatchMode,
                 order_by: data.orderBy,
                 run_id_pattern: data.runIdPattern,
                 run_id_prefix_pattern: data.runIdPrefixPattern,
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 afc4e5835ce..f279f7d2cea 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
@@ -3587,6 +3587,8 @@ export type GetDagRunsData = {
     startDateLt?: string | null;
     startDateLte?: string | null;
     state?: Array<(string)>;
+    tags?: Array<(string)>;
+    tagsMatchMode?: 'any' | 'all' | null;
     teams?: Array<(string)>;
     /**
      * Case-insensitive substring match (SQL `ILIKE`). Slower than 
`triggering_user_name_prefix_pattern` on large tables — see "Filtering with 
pattern parameters".
diff --git a/airflow-core/src/airflow/ui/src/mocks/handlers/dag_runs.ts 
b/airflow-core/src/airflow/ui/src/mocks/handlers/dag_runs.ts
index 7fcb78ebfc9..2854fff56d4 100644
--- a/airflow-core/src/airflow/ui/src/mocks/handlers/dag_runs.ts
+++ b/airflow-core/src/airflow/ui/src/mocks/handlers/dag_runs.ts
@@ -56,13 +56,49 @@ const dagRunInRange = {
   triggering_user_name: "admin",
 };
 
+const dagRunTaggedDag = {
+  conf: null,
+  dag_display_name: "tagged_dag",
+  dag_id: "tagged_dag",
+  dag_run_id: "run_tagged_dag",
+  dag_versions: [],
+  data_interval_end: null,
+  data_interval_start: null,
+  duration: 1.0,
+  end_date: "2025-01-15T00:00:01Z",
+  logical_date: "2025-01-15T00:00:00Z",
+  partition_key: null,
+  run_after: "2025-01-15T00:00:00Z",
+  run_type: "manual",
+  start_date: "2025-01-15T00:00:00Z",
+  state: "success",
+  triggering_user_name: "admin",
+};
+
+const dagRunMultiTaggedDag = {
+  ...dagRunTaggedDag,
+  dag_display_name: "multi_tagged_dag",
+  dag_id: "multi_tagged_dag",
+  dag_run_id: "run_multi_tagged_dag",
+};
+
+// Maps dag_id to the tags its Dag is annotated with, so the "tags" query
+// param can be simulated without a real Dags table backing the mock.
+const dagIdTags: Record<string, Array<string>> = {
+  multi_tagged_dag: ["example_tag", "other_tag"],
+  tagged_dag: ["example_tag"],
+  test_dag: [],
+};
+
 export const handlers: Array<HttpHandler> = [
   http.get("/api/v2/dags/:dagId/dagRuns", ({ request }) => {
     const url = new URL(request.url);
     const logicalDateGte = url.searchParams.get("logical_date_gte");
     const logicalDateLte = url.searchParams.get("logical_date_lte");
+    const tags = url.searchParams.getAll("tags");
+    const tagsMatchMode = url.searchParams.get("tags_match_mode");
 
-    const allRuns = [dagRunBeforeFilter, dagRunInRange];
+    const allRuns = [dagRunBeforeFilter, dagRunInRange, dagRunTaggedDag, 
dagRunMultiTaggedDag];
 
     const filtered = allRuns.filter((run) => {
       const logicalDate = new Date(run.logical_date);
@@ -73,6 +109,17 @@ export const handlers: Array<HttpHandler> = [
       if (logicalDateLte !== null && logicalDate > new Date(logicalDateLte)) {
         return false;
       }
+      if (tags.length > 0) {
+        const runTags = dagIdTags[run.dag_id] ?? [];
+        const matches =
+          tagsMatchMode === "all"
+            ? tags.every((tag) => runTags.includes(tag))
+            : tags.some((tag) => runTags.includes(tag));
+
+        if (!matches) {
+          return false;
+        }
+      }
 
       return true;
     });
diff --git a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.test.tsx 
b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.test.tsx
index ca2ad5ada3a..24ee1bc1847 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.test.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.test.tsx
@@ -55,6 +55,33 @@ describe("DagRuns logical date filter", () => {
   });
 });
 
+// dag_runs mock handler (see src/mocks/handlers/dag_runs.ts) tags 
"tagged_dag" with
+// "example_tag" and "multi_tagged_dag" with "example_tag" and "other_tag"; 
"test_dag" has no tags.
+describe("DagRuns tags filter", () => {
+  it("filters runs by the tags query param", async () => {
+    render(<AppWrapper initialEntries={["/dag_runs?tags=example_tag"]} />);
+
+    await waitFor(() => 
expect(screen.getByText("run_tagged_dag")).toBeInTheDocument());
+    expect(screen.getByText("run_multi_tagged_dag")).toBeInTheDocument();
+    expect(screen.queryByText("run_in_range")).not.toBeInTheDocument();
+    expect(screen.queryByText("run_before_filter")).not.toBeInTheDocument();
+  });
+
+  it("matches runs of Dags with any of the tags by default", async () => {
+    render(<AppWrapper 
initialEntries={["/dag_runs?tags=example_tag&tags=other_tag"]} />);
+
+    await waitFor(() => 
expect(screen.getByText("run_tagged_dag")).toBeInTheDocument());
+    expect(screen.getByText("run_multi_tagged_dag")).toBeInTheDocument();
+  });
+
+  it("matches only runs of Dags with all of the tags when tags_match_mode is 
all", async () => {
+    render(<AppWrapper 
initialEntries={["/dag_runs?tags=example_tag&tags=other_tag&tags_match_mode=all"]}
 />);
+
+    await waitFor(() => 
expect(screen.getByText("run_multi_tagged_dag")).toBeInTheDocument());
+    expect(screen.queryByText("run_tagged_dag")).not.toBeInTheDocument();
+  });
+});
+
 describe("DagRuns conf expand/collapse", () => {
   // Relies on the conf column being visible by default, which is what renders 
the JSON viewer.
   // useTableURLState persists sorting to localStorage, so clear it between 
cases.
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 b7fb5c347fe..c826120fa88 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRuns.tsx
@@ -87,6 +87,8 @@ const {
   START_DATE_GTE: START_DATE_GTE_PARAM,
   START_DATE_LTE: START_DATE_LTE_PARAM,
   STATE: STATE_PARAM,
+  TAGS: TAGS_PARAM,
+  TAGS_MATCH_MODE: TAGS_MATCH_MODE_PARAM,
   TEAMS: TEAMS_PARAM,
   TRIGGERING_USER_NAME_PATTERN: TRIGGERING_USER_NAME_PATTERN_PARAM,
 }: SearchParamsKeysType = SearchParamsKeys;
@@ -279,6 +281,8 @@ 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 tags = searchParams.getAll(TAGS_PARAM);
+  const tagsMatchMode = searchParams.get(TAGS_MATCH_MODE_PARAM) === "all" ? 
"all" : "any";
   const teams = searchParams.getAll(TEAMS_PARAM);
 
   const refetchInterval = useAutoRefresh({});
@@ -334,6 +338,8 @@ export const DagRuns = () => {
       startDateGte: startDateGte ?? undefined,
       startDateLte: startDateLte ?? undefined,
       state: filteredState === null ? undefined : [filteredState],
+      tags: tags.length > 0 ? tags : undefined,
+      tagsMatchMode: tags.length > 0 ? tagsMatchMode : undefined,
       teams: teams.length > 0 ? teams : undefined,
       ...triggeringUserArg,
     },
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 8245729329d..9722b065d21 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRunsFilters.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagRuns/DagRunsFilters.tsx
@@ -52,6 +52,7 @@ export const DagRunsFilters = ({ dagId }: 
DagRunsFiltersProps) => {
 
   if (dagId === undefined) {
     searchParamKeys.unshift(SearchParamsKeys.DAG_ID_PATTERN);
+    searchParamKeys.push(SearchParamsKeys.TAGS);
   }
 
   const { filterConfigs, handleFiltersChange, initialValues } = 
useFiltersHandler(searchParamKeys);
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 8edfcb27816..075bfb0ede2 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
@@ -35,7 +35,7 @@ from airflow.api_fastapi.auth.managers.simple.user import 
SimpleAuthManagerUser
 from airflow.api_fastapi.common.dagbag import resolve_run_on_latest_version
 from airflow.api_fastapi.core_api.datamodels.dag_versions import 
DagVersionResponse
 from airflow.exceptions import ParamValidationError
-from airflow.models import DagModel, DagRun, Log
+from airflow.models import DagModel, DagRun, DagTag, Log
 from airflow.models.asset import AssetEvent, AssetModel
 from airflow.models.dagbundle import DagBundleModel
 from airflow.models.taskinstance import TaskInstance
@@ -335,6 +335,19 @@ def get_dag_run_dict(run: DagRun):
     }
 
 
+def _attach_tags_to_dag(session, dag_id: str, tag_names: list[str]) -> None:
+    """Assign Dag tags for tag-filter tests."""
+    for tag_name in tag_names:
+        session.add(DagTag(dag_id=dag_id, name=tag_name))
+    session.commit()
+
+
+def _detach_tags_from_dag(session, dag_id: str) -> None:
+    """Undo :func:`_attach_tags_to_dag`."""
+    session.execute(delete(DagTag).where(DagTag.dag_id == dag_id))
+    session.commit()
+
+
 class TestGetDagRun:
     @pytest.mark.parametrize(
         ("dag_id", "run_id", "state", "run_type", "triggered_by", 
"dag_run_note"),
@@ -472,6 +485,45 @@ class TestGetDagRuns:
             assert response.status_code == 200
             assert response.json()["total_entries"] == 0
 
+    def test_get_dag_runs_filtered_by_tag(self, test_client, session):
+        _attach_tags_to_dag(session, DAG1_ID, ["tag-filter-only"])
+        try:
+            response = test_client.get("/dags/~/dagRuns", params={"tags": 
["tag-filter-only"]})
+            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 tag with no Dags returns nothing.
+            response = test_client.get("/dags/~/dagRuns", params={"tags": 
["nonexistent-tag"]})
+            assert response.status_code == 200
+            assert response.json()["total_entries"] == 0
+        finally:
+            _detach_tags_from_dag(session, DAG1_ID)
+
+    def test_get_dag_runs_filtered_by_tags_match_mode(self, test_client, 
session):
+        _attach_tags_to_dag(session, DAG1_ID, ["tag-filter-a"])
+        _attach_tags_to_dag(session, DAG2_ID, ["tag-filter-a", "tag-filter-b"])
+        try:
+            response = test_client.get(
+                "/dags/~/dagRuns",
+                params={"tags": ["tag-filter-a", "tag-filter-b"], 
"tags_match_mode": "any"},
+            )
+            assert response.status_code == 200
+            assert response.json()["total_entries"] == 4
+
+            response = test_client.get(
+                "/dags/~/dagRuns",
+                params={"tags": ["tag-filter-a", "tag-filter-b"], 
"tags_match_mode": "all"},
+            )
+            assert response.status_code == 200
+            body = response.json()
+            assert body["total_entries"] == 2
+            assert {run["dag_id"] for run in body["dag_runs"]} == {DAG2_ID}
+        finally:
+            _detach_tags_from_dag(session, DAG1_ID)
+            _detach_tags_from_dag(session, DAG2_ID)
+
     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

Reply via email to