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 18fe1c60986 UI: Add bulk pause, drain and unpause actions to the Dags
list (#73055)
18fe1c60986 is described below
commit 18fe1c6098605e089dcb7571b3839183f7b32752
Author: Dheeraj Turaga <[email protected]>
AuthorDate: Sat Sep 26 16:32:16 2026 -0500
UI: Add bulk pause, drain and unpause actions to the Dags list (#73055)
* Add a bulk pause/drain action to the Dags list
Draining a Dag with unfinished runs now asks whether to drain or pause
immediately, which turned stopping many Dags into two clicks per Dag
instead of one. Selecting several Dags in the table view and pausing
or draining them together needs at most one such choice for the whole
batch.
Follow-up to the draining state added in #72407.
related: #22006
* Mention bulk pause/drain in the Dag draining release note
The bulk action in the Dags list is part of the graceful draining feature,
so the release note should credit it alongside the core and airflowctl changes.
* Keep already-paused Dags paused when bulk draining
Draining clears is_paused, so a paused Dag caught up in a bulk drain let
the scheduler resume its unfinished runs until the drain completed, the
opposite of what the user asked for. Such Dags also forced the
drain-vs-pause prompt even when every unpaused Dag in the selection was
idle.
* Add a bulk unpause action to the Dags list
A selection could be paused or drained in one go but had to be resumed
one toggle at a time. Unpausing always asks for confirmation because
resuming many Dags at once can start a burst of catchup runs.
---
airflow-core/docs/security/api_permissions_ref.rst | 4 +
airflow-core/newsfragments/72407.significant.rst | 5 +-
.../api_fastapi/core_api/datamodels/dags.py | 6 +
.../core_api/openapi/v2-rest-api-generated.yaml | 154 +++++++++++++++
.../api_fastapi/core_api/routes/public/dags.py | 26 ++-
.../src/airflow/api_fastapi/core_api/security.py | 34 ++++
.../api_fastapi/core_api/services/public/dags.py | 96 ++++++++++
.../src/airflow/ui/openapi-gen/queries/common.ts | 1 +
.../src/airflow/ui/openapi-gen/queries/queries.ts | 15 +-
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 164 ++++++++++++++++
.../ui/openapi-gen/requests/services.gen.ts | 24 ++-
.../airflow/ui/openapi-gen/requests/types.gen.ts | 82 ++++++++
.../airflow/ui/public/i18n/locales/en/dags.json | 3 +-
.../ui/src/components/PauseOrDrainChoiceModal.tsx | 66 +++++++
.../src/airflow/ui/src/components/TogglePause.tsx | 44 ++---
.../src/airflow/ui/src/mocks/handlers/dags.ts | 6 +-
.../pages/DagsList/BulkPauseDrainDagsButton.tsx | 105 +++++++++++
.../src/pages/DagsList/BulkUnpauseDagsButton.tsx | 85 +++++++++
.../ui/src/pages/DagsList/DagsList.test.tsx | 143 ++++++++++++++
.../src/airflow/ui/src/pages/DagsList/DagsList.tsx | 123 ++++++++----
.../ui/src/queries/useBulkSetDagSchedulingState.ts | 89 +++++++++
.../core_api/routes/public/test_dags.py | 210 +++++++++++++++++++++
.../src/airflowctl/api/datamodels/generated.py | 70 +++++++
providers/fab/docs/auth-manager/access-control.rst | 4 +
scripts/ci/prek/extract_permissions.py | 1 +
scripts/ci/prek/fab_permissions_doc.py | 8 +
scripts/tests/ci/prek/test_extract_permissions.py | 3 +
27 files changed, 1488 insertions(+), 83 deletions(-)
diff --git a/airflow-core/docs/security/api_permissions_ref.rst
b/airflow-core/docs/security/api_permissions_ref.rst
index c930fdb9f07..c6a86c1dc84 100644
--- a/airflow-core/docs/security/api_permissions_ref.rst
+++ b/airflow-core/docs/security/api_permissions_ref.rst
@@ -218,6 +218,10 @@ source code so it stays up to date as endpoints are added
or changed.
- ``/api/v2/dags``
- ``DAG``
- ``PUT``
+ * - ``PATCH``
+ - ``/api/v2/dags/bulk``
+ - ``DAG``
+ - ``multi``
* - ``DELETE``
- ``/api/v2/dags/{dag_id}``
- ``DAG``
diff --git a/airflow-core/newsfragments/72407.significant.rst
b/airflow-core/newsfragments/72407.significant.rst
index 77abb598897..90355b11c8f 100644
--- a/airflow-core/newsfragments/72407.significant.rst
+++ b/airflow-core/newsfragments/72407.significant.rst
@@ -1,4 +1,4 @@
-Gracefully pause Dags by draining active runs (#72407, #73225)
+Gracefully pause Dags by draining active runs (#72407, #73225, #73055)
Airflow now supports a ``draining`` scheduling state for Dags. Draining stops
the scheduler from creating new scheduled or asset-triggered Dag runs while
@@ -9,6 +9,7 @@ After outstanding Dag runs finish, the Dag automatically
transitions to
``paused``. Users can cancel draining to return the Dag to ``active``.
The UI, REST API, and ``airflowctl`` expose the ``active``, ``draining``, and
-``paused`` scheduling states. Existing ``is_paused`` API usage remains
+``paused`` scheduling states, and the Dags list can pause, drain, or unpause
+several selected Dags at once. Existing ``is_paused`` API usage remains
supported; use ``scheduling_state`` when the draining state must be selected or
reported explicitly.
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py
index db3b013976f..1030d6e3d40 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dags.py
@@ -182,6 +182,12 @@ class DAGPatchBody(StrictBaseModel):
DAGPatchBodyPartial = make_partial_model(DAGPatchBody)
+class BulkDAGBody(DAGPatchBody):
+ """Request body for bulk update of Dags."""
+
+ dag_id: str
+
+
class DAGCollectionResponse(BaseModel):
"""Dag Collection serializer for responses."""
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 49694cd5d03..d57c05b6812 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
@@ -4580,6 +4580,47 @@ paths:
application/json:
schema:
$ref: '#/components/schemas/HTTPValidationError'
+ /api/v2/dags/bulk:
+ patch:
+ tags:
+ - DAG
+ summary: Bulk Dags
+ description: Bulk pause, resume, or drain Dags by id.
+ operationId: bulk_dags
+ requestBody:
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/BulkBody_BulkDAGBody_'
+ required: true
+ responses:
+ '200':
+ description: Successful Response
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/BulkResponse'
+ '401':
+ description: Unauthorized
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ '403':
+ description: Forbidden
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ '422':
+ description: Validation Error
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPValidationError'
+ security:
+ - OAuth2PasswordBearer: []
+ - HTTPBearer: []
/api/v2/dags/{dag_id}/favorite:
post:
tags:
@@ -12039,6 +12080,21 @@ components:
This structure helps users understand which key actions succeeded and
which
failed.'
+ BulkBody_BulkDAGBody_:
+ properties:
+ actions:
+ items:
+ oneOf:
+ - $ref: '#/components/schemas/BulkCreateAction_BulkDAGBody_'
+ - $ref: '#/components/schemas/BulkUpdateAction_BulkDAGBody_'
+ - $ref: '#/components/schemas/BulkDeleteAction_BulkDAGBody_'
+ type: array
+ title: Actions
+ additionalProperties: false
+ type: object
+ required:
+ - actions
+ title: BulkBody[BulkDAGBody]
BulkBody_BulkDAGRunBody_:
properties:
actions:
@@ -12114,6 +12170,28 @@ components:
required:
- actions
title: BulkBody[VariableBody]
+ BulkCreateAction_BulkDAGBody_:
+ properties:
+ action:
+ type: string
+ const: create
+ title: Action
+ description: The action to be performed on the entities.
+ entities:
+ items:
+ $ref: '#/components/schemas/BulkDAGBody'
+ type: array
+ title: Entities
+ description: A list of entities to be created.
+ action_on_existence:
+ $ref: '#/components/schemas/BulkActionOnExistence'
+ default: fail
+ additionalProperties: false
+ type: object
+ required:
+ - action
+ - entities
+ title: BulkCreateAction[BulkDAGBody]
BulkCreateAction_BulkDAGRunBody_:
properties:
action:
@@ -12224,6 +12302,26 @@ components:
- action
- entities
title: BulkCreateAction[VariableBody]
+ BulkDAGBody:
+ properties:
+ is_paused:
+ anyOf:
+ - type: boolean
+ - type: 'null'
+ title: Is Paused
+ scheduling_state:
+ anyOf:
+ - $ref: '#/components/schemas/DagSchedulingState'
+ - type: 'null'
+ dag_id:
+ type: string
+ title: Dag Id
+ additionalProperties: false
+ type: object
+ required:
+ - dag_id
+ title: BulkDAGBody
+ description: Request body for bulk update of Dags.
BulkDAGRunBody:
properties:
dag_run_id:
@@ -12315,6 +12413,30 @@ components:
type: object
title: BulkDAGRunClearBody
description: Request body for the bulk clear Dag Runs endpoint.
+ BulkDeleteAction_BulkDAGBody_:
+ properties:
+ action:
+ type: string
+ const: delete
+ title: Action
+ description: The action to be performed on the entities.
+ entities:
+ items:
+ anyOf:
+ - type: string
+ - $ref: '#/components/schemas/BulkDAGBody'
+ type: array
+ title: Entities
+ description: A list of entity id/key or entity objects to be deleted.
+ action_on_non_existence:
+ $ref: '#/components/schemas/BulkActionNotOnExistence'
+ default: fail
+ additionalProperties: false
+ type: object
+ required:
+ - action
+ - entities
+ title: BulkDeleteAction[BulkDAGBody]
BulkDeleteAction_BulkDAGRunBody_:
properties:
action:
@@ -12520,6 +12642,38 @@ components:
- task_id
title: BulkTaskInstanceBody
description: Request body for bulk update, and delete task instances.
+ BulkUpdateAction_BulkDAGBody_:
+ properties:
+ action:
+ type: string
+ const: update
+ title: Action
+ description: The action to be performed on the entities.
+ entities:
+ items:
+ $ref: '#/components/schemas/BulkDAGBody'
+ type: array
+ title: Entities
+ description: A list of entities to be updated.
+ update_mask:
+ anyOf:
+ - items:
+ type: string
+ type: array
+ - type: 'null'
+ title: Update Mask
+ description: A list of field names to update for each entity.Only
these
+ fields will be applied from the request body to the database
model.Any
+ extra fields provided will be ignored.
+ action_on_non_existence:
+ $ref: '#/components/schemas/BulkActionNotOnExistence'
+ default: fail
+ additionalProperties: false
+ type: object
+ required:
+ - action
+ - entities
+ title: BulkUpdateAction[BulkDAGBody]
BulkUpdateAction_BulkDAGRunBody_:
properties:
action:
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py
index 3b1f275a8b5..ad6531a0ed5 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dags.py
@@ -58,7 +58,9 @@ from airflow.api_fastapi.common.parameters import (
)
from airflow.api_fastapi.common.router import AirflowRouter
from airflow.api_fastapi.compat import HTTP_422_UNPROCESSABLE_CONTENT
+from airflow.api_fastapi.core_api.datamodels.common import BulkBody,
BulkResponse
from airflow.api_fastapi.core_api.datamodels.dags import (
+ BulkDAGBody,
DAGCollectionResponse,
DAGDetailsResponse,
DAGPatchBody,
@@ -71,7 +73,9 @@ from airflow.api_fastapi.core_api.security import (
GetUserDep,
ReadableDagsFilterDep,
requires_access_dag,
+ requires_access_dag_bulk,
)
+from airflow.api_fastapi.core_api.services.public.dags import BulkDagService,
get_scheduling_state
from airflow.api_fastapi.logging.decorators import action_logging
from airflow.exceptions import AirflowException, DagNotFound
from airflow.models import DagModel
@@ -83,14 +87,6 @@ from airflow.utils.state import DagRunState,
DagSchedulingState
dags_router = AirflowRouter(tags=["DAG"], prefix="/dags")
-def _get_scheduling_state(patch_body: DAGPatchBody) -> DagSchedulingState:
- if patch_body.scheduling_state is not None:
- return patch_body.scheduling_state
- if patch_body.is_paused is True:
- return DagSchedulingState.PAUSED
- return DagSchedulingState.ACTIVE
-
-
@dags_router.get("", dependencies=[Depends(requires_access_dag(method="GET"))])
def get_dags(
limit: QueryLimit,
@@ -281,6 +277,16 @@ def get_dag_details(
return DAGDetailsResponse.model_validate(dag_model)
+@dags_router.patch(
+ # Declared before "/{dag_id}" so this literal segment isn't shadowed by
that path param.
+ "/bulk",
+ dependencies=[Depends(requires_access_dag_bulk()),
Depends(action_logging())],
+)
+def bulk_dags(request: BulkBody[BulkDAGBody], session: SessionDep) ->
BulkResponse:
+ """Bulk pause, resume, or drain Dags by id."""
+ return BulkDagService(session=session, request=request).handle_request()
+
+
@dags_router.patch(
"/{dag_id}",
responses=create_openapi_http_exception_doc(
@@ -327,7 +333,7 @@ def patch_dag(
except ValidationError as e:
raise RequestValidationError(errors=e.errors())
- dag.set_scheduling_state(_get_scheduling_state(patch_body))
+ dag.set_scheduling_state(get_scheduling_state(patch_body))
return dag
@@ -410,7 +416,7 @@ def patch_dags(
],
).subquery()
- scheduling_state = _get_scheduling_state(patch_body)
+ scheduling_state = get_scheduling_state(patch_body)
session.execute(
update(DagModel)
.where(DagModel.dag_id.in_(select(filtered_dag_ids.c.dag_id)))
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/security.py
b/airflow-core/src/airflow/api_fastapi/core_api/security.py
index 61c7444bde2..1937cf1e7cc 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/security.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/security.py
@@ -67,6 +67,7 @@ from airflow.api_fastapi.core_api.datamodels.common import (
)
from airflow.api_fastapi.core_api.datamodels.connections import ConnectionBody
from airflow.api_fastapi.core_api.datamodels.dag_run import BulkDAGRunBody,
BulkDAGRunClearBody
+from airflow.api_fastapi.core_api.datamodels.dags import BulkDAGBody
from airflow.api_fastapi.core_api.datamodels.pools import PoolBody
from airflow.api_fastapi.core_api.datamodels.variables import VariableBody
from airflow.configuration import conf
@@ -1101,6 +1102,39 @@ def requires_access_dag_run_bulk() ->
Callable[[BulkBody[BulkDAGRunBody], BaseUs
return inner
+def requires_access_dag_bulk() -> Callable[[BulkBody[BulkDAGBody], BaseUser],
None]:
+ def inner(
+ request: BulkBody[BulkDAGBody],
+ user: GetUserDep,
+ ) -> None:
+ entity_methods: list[tuple[str, ResourceMethod]] = []
+ for action in request.actions:
+ methods = _get_resource_methods_from_bulk_request(action)
+ for entity in action.entities:
+ entity_dag_id = entity if isinstance(entity, str) else
entity.dag_id
+ for method in methods:
+ entity_methods.append((entity_dag_id, method))
+
+ if not entity_methods:
+ return
+
+ dag_id_to_team = DagModel.get_dag_id_to_team_name_mapping(
+ list({dag_id for dag_id, _ in entity_methods})
+ )
+ requests: list[IsAuthorizedDagRequest] = [
+ {"method": method, "details": DagDetails(id=dag_id,
team_name=dag_id_to_team.get(dag_id))}
+ for dag_id, method in entity_methods
+ ]
+ _requires_access(
+ is_authorized_callback=lambda:
get_auth_manager().batch_is_authorized_dag(
+ requests=requests,
+ user=user,
+ )
+ )
+
+ return inner
+
+
def requires_access_dag_run_clear_bulk() -> Callable[[BulkDAGRunClearBody,
BaseUser, str], None]:
def inner(
body: BulkDAGRunClearBody,
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/services/public/dags.py
b/airflow-core/src/airflow/api_fastapi/core_api/services/public/dags.py
new file mode 100644
index 00000000000..2671273298c
--- /dev/null
+++ b/airflow-core/src/airflow/api_fastapi/core_api/services/public/dags.py
@@ -0,0 +1,96 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+from __future__ import annotations
+
+from fastapi import HTTPException, status
+from sqlalchemy import select
+
+from airflow.api_fastapi.core_api.datamodels.common import (
+ BulkActionNotOnExistence,
+ BulkActionResponse,
+ BulkCreateAction,
+ BulkDeleteAction,
+ BulkUpdateAction,
+)
+from airflow.api_fastapi.core_api.datamodels.dags import BulkDAGBody,
DAGPatchBody
+from airflow.api_fastapi.core_api.services.public.common import BulkService
+from airflow.models.dag import DagModel
+from airflow.utils.state import DagSchedulingState
+
+
+def get_scheduling_state(patch_body: DAGPatchBody) -> DagSchedulingState:
+ """Resolve the target scheduling state from a legacy `is_paused` or a
`scheduling_state` patch body."""
+ if patch_body.scheduling_state is not None:
+ return patch_body.scheduling_state
+ if patch_body.is_paused is True:
+ return DagSchedulingState.PAUSED
+ return DagSchedulingState.ACTIVE
+
+
+class BulkDagService(BulkService[BulkDAGBody]):
+ """Service for handling bulk operations on Dags."""
+
+ def handle_bulk_create(self, action: BulkCreateAction[BulkDAGBody],
results: BulkActionResponse) -> None:
+ results.errors.append(
+ {
+ "error": "Dags bulk create is not supported.",
+ "status_code": status.HTTP_405_METHOD_NOT_ALLOWED,
+ }
+ )
+
+ def handle_bulk_delete(self, action: BulkDeleteAction[BulkDAGBody],
results: BulkActionResponse) -> None:
+ results.errors.append(
+ {
+ "error": "Dags bulk delete is not supported. Use the delete
Dag endpoint instead.",
+ "status_code": status.HTTP_405_METHOD_NOT_ALLOWED,
+ }
+ )
+
+ def handle_bulk_update(self, action: BulkUpdateAction[BulkDAGBody],
results: BulkActionResponse) -> None:
+ """Bulk update Dags (pause, resume, or drain)."""
+ entities_by_id = {entity.dag_id: entity for entity in action.entities}
+ if not entities_by_id:
+ return
+
+ dag_map = {
+ dag.dag_id: dag
+ for dag in self.session.scalars(
+
select(DagModel).where(DagModel.dag_id.in_(entities_by_id.keys()))
+ )
+ }
+ not_found_ids = set(entities_by_id) - set(dag_map)
+
+ try:
+ if action.action_on_non_existence == BulkActionNotOnExistence.FAIL
and not_found_ids:
+ raise HTTPException(
+ status.HTTP_404_NOT_FOUND,
+ f"The Dags with these ids: {sorted(not_found_ids)} were
not found",
+ )
+ update_ids = (
+ set(dag_map)
+ if action.action_on_non_existence ==
BulkActionNotOnExistence.SKIP
+ else set(entities_by_id)
+ )
+ for dag_id in update_ids:
+ dag = dag_map.get(dag_id)
+ if dag is None:
+ continue
+
dag.set_scheduling_state(get_scheduling_state(entities_by_id[dag_id]))
+ results.success.append(dag_id)
+ except HTTPException as e:
+ results.errors.append({"error": f"{e.detail}", "status_code":
e.status_code})
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 05479e1f51d..c6c5d28aff6 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -1191,6 +1191,7 @@ export type DagRunServicePatchDagRunMutationResult =
Awaited<ReturnType<typeof D
export type DagRunServiceBulkDagRunsMutationResult = Awaited<ReturnType<typeof
DagRunService.bulkDagRuns>>;
export type DagServicePatchDagsMutationResult = Awaited<ReturnType<typeof
DagService.patchDags>>;
export type DagServicePatchDagMutationResult = Awaited<ReturnType<typeof
DagService.patchDag>>;
+export type DagServiceBulkDagsMutationResult = Awaited<ReturnType<typeof
DagService.bulkDags>>;
export type TaskInstanceServicePatchTaskInstanceMutationResult =
Awaited<ReturnType<typeof TaskInstanceService.patchTaskInstance>>;
export type TaskInstanceServicePatchTaskInstanceByMapIndexMutationResult =
Awaited<ReturnType<typeof TaskInstanceService.patchTaskInstanceByMapIndex>>;
export type TaskInstanceServiceBulkTaskInstancesMutationResult =
Awaited<ReturnType<typeof TaskInstanceService.bulkTaskInstances>>;
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 2fd053cb65c..aaa024a638c 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -2,7 +2,7 @@
import { UseMutationOptions, UseQueryOptions, useMutation, useQuery } from
"@tanstack/react-query";
import { AssetService, AssetStateStoreService, AuthLinksService,
BackfillService, CalendarService, ConfigService, ConnectionService,
DagBundleService, DagParsingService, DagRunService, DagService,
DagSourceService, DagStatsService, DagVersionService, DagWarningService,
DashboardService, DeadlinesService, DependenciesService, EventLogService,
ExperimentalService, ExtraLinksService, GanttService, GridService,
ImportErrorService, JobService, LoginService, MonitorService,
PartitionedDagRunSe [...]
-import { AssetStateStoreBody, BackfillPostBody, BulkBody_BulkDAGRunBody_,
BulkBody_BulkTaskInstanceBody_, BulkBody_ConnectionBody_, BulkBody_PoolBody_,
BulkBody_VariableBody_, BulkDAGRunClearBody, ClearPartitionsBody,
ClearTaskInstancesBody, ConnectionBody, ConnectionTestRequestBody,
CreateAssetEventsBody, DAGPatchBody, DAGRunClearBody, DAGRunPatchBody,
DAGRunsBatchBody, DagRunState, DagSchedulingState, DagWarningType,
GenerateTokenBody, MaterializeAssetBody, PatchTaskInstanceBody, PoolB [...]
+import { AssetStateStoreBody, BackfillPostBody, BulkBody_BulkDAGBody_,
BulkBody_BulkDAGRunBody_, BulkBody_BulkTaskInstanceBody_,
BulkBody_ConnectionBody_, BulkBody_PoolBody_, BulkBody_VariableBody_,
BulkDAGRunClearBody, ClearPartitionsBody, ClearTaskInstancesBody,
ConnectionBody, ConnectionTestRequestBody, CreateAssetEventsBody, DAGPatchBody,
DAGRunClearBody, DAGRunPatchBody, DAGRunsBatchBody, DagRunState,
DagSchedulingState, DagWarningType, GenerateTokenBody, MaterializeAssetBody,
Patch [...]
import * as Common from "./common";
/**
* Get Assets
@@ -2801,6 +2801,19 @@ export const useDagServicePatchDag = <TData =
Common.DagServicePatchDagMutationR
updateMask?: string[];
}, TContext>({ mutationFn: ({ dagId, requestBody, updateMask }) =>
DagService.patchDag({ dagId, requestBody, updateMask }) as unknown as
Promise<TData>, ...options });
/**
+* Bulk Dags
+* Bulk pause, resume, or drain Dags by id.
+* @param data The data for the request.
+* @param data.requestBody
+* @returns BulkResponse Successful Response
+* @throws ApiError
+*/
+export const useDagServiceBulkDags = <TData =
Common.DagServiceBulkDagsMutationResult, TError = unknown, TContext =
unknown>(options?: Omit<UseMutationOptions<TData, TError, {
+ requestBody: BulkBody_BulkDAGBody_;
+}, TContext>, "mutationFn">) => useMutation<TData, TError, {
+ requestBody: BulkBody_BulkDAGBody_;
+}, TContext>({ mutationFn: ({ requestBody }) => DagService.bulkDags({
requestBody }) as unknown as Promise<TData>, ...options });
+/**
* Patch Task Instance
* Update a task instance.
* @param data The data for the request.
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 71d21c125c5..11bd522c626 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
@@ -1105,6 +1105,32 @@ The response includes a list of successful keys and any
errors encountered durin
This structure helps users understand which key actions succeeded and which
failed.`
} as const;
+export const $BulkBody_BulkDAGBody_ = {
+ properties: {
+ actions: {
+ items: {
+ oneOf: [
+ {
+ '$ref':
'#/components/schemas/BulkCreateAction_BulkDAGBody_'
+ },
+ {
+ '$ref':
'#/components/schemas/BulkUpdateAction_BulkDAGBody_'
+ },
+ {
+ '$ref':
'#/components/schemas/BulkDeleteAction_BulkDAGBody_'
+ }
+ ]
+ },
+ type: 'array',
+ title: 'Actions'
+ }
+ },
+ additionalProperties: false,
+ type: 'object',
+ required: ['actions'],
+ title: 'BulkBody[BulkDAGBody]'
+} as const;
+
export const $BulkBody_BulkDAGRunBody_ = {
properties: {
actions: {
@@ -1235,6 +1261,33 @@ export const $BulkBody_VariableBody_ = {
title: 'BulkBody[VariableBody]'
} as const;
+export const $BulkCreateAction_BulkDAGBody_ = {
+ properties: {
+ action: {
+ type: 'string',
+ const: 'create',
+ title: 'Action',
+ description: 'The action to be performed on the entities.'
+ },
+ entities: {
+ items: {
+ '$ref': '#/components/schemas/BulkDAGBody'
+ },
+ type: 'array',
+ title: 'Entities',
+ description: 'A list of entities to be created.'
+ },
+ action_on_existence: {
+ '$ref': '#/components/schemas/BulkActionOnExistence',
+ default: 'fail'
+ }
+ },
+ additionalProperties: false,
+ type: 'object',
+ required: ['action', 'entities'],
+ title: 'BulkCreateAction[BulkDAGBody]'
+} as const;
+
export const $BulkCreateAction_BulkDAGRunBody_ = {
properties: {
action: {
@@ -1370,6 +1423,41 @@ export const $BulkCreateAction_VariableBody_ = {
title: 'BulkCreateAction[VariableBody]'
} as const;
+export const $BulkDAGBody = {
+ properties: {
+ is_paused: {
+ anyOf: [
+ {
+ type: 'boolean'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Is Paused'
+ },
+ scheduling_state: {
+ anyOf: [
+ {
+ '$ref': '#/components/schemas/DagSchedulingState'
+ },
+ {
+ type: 'null'
+ }
+ ]
+ },
+ dag_id: {
+ type: 'string',
+ title: 'Dag Id'
+ }
+ },
+ additionalProperties: false,
+ type: 'object',
+ required: ['dag_id'],
+ title: 'BulkDAGBody',
+ description: 'Request body for bulk update of Dags.'
+} as const;
+
export const $BulkDAGRunBody = {
properties: {
dag_run_id: {
@@ -1511,6 +1599,40 @@ export const $BulkDAGRunClearBody = {
description: 'Request body for the bulk clear Dag Runs endpoint.'
} as const;
+export const $BulkDeleteAction_BulkDAGBody_ = {
+ properties: {
+ action: {
+ type: 'string',
+ const: 'delete',
+ title: 'Action',
+ description: 'The action to be performed on the entities.'
+ },
+ entities: {
+ items: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ '$ref': '#/components/schemas/BulkDAGBody'
+ }
+ ]
+ },
+ type: 'array',
+ title: 'Entities',
+ description: 'A list of entity id/key or entity objects to be
deleted.'
+ },
+ action_on_non_existence: {
+ '$ref': '#/components/schemas/BulkActionNotOnExistence',
+ default: 'fail'
+ }
+ },
+ additionalProperties: false,
+ type: 'object',
+ required: ['action', 'entities'],
+ title: 'BulkDeleteAction[BulkDAGBody]'
+} as const;
+
export const $BulkDeleteAction_BulkDAGRunBody_ = {
properties: {
action: {
@@ -1815,6 +1937,48 @@ export const $BulkTaskInstanceBody = {
description: 'Request body for bulk update, and delete task instances.'
} as const;
+export const $BulkUpdateAction_BulkDAGBody_ = {
+ properties: {
+ action: {
+ type: 'string',
+ const: 'update',
+ title: 'Action',
+ description: 'The action to be performed on the entities.'
+ },
+ entities: {
+ items: {
+ '$ref': '#/components/schemas/BulkDAGBody'
+ },
+ type: 'array',
+ title: 'Entities',
+ description: 'A list of entities to be updated.'
+ },
+ update_mask: {
+ anyOf: [
+ {
+ items: {
+ type: 'string'
+ },
+ type: 'array'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Update Mask',
+ description: 'A list of field names to update for each entity.Only
these fields will be applied from the request body to the database model.Any
extra fields provided will be ignored.'
+ },
+ action_on_non_existence: {
+ '$ref': '#/components/schemas/BulkActionNotOnExistence',
+ default: 'fail'
+ }
+ },
+ additionalProperties: false,
+ type: 'object',
+ required: ['action', 'entities'],
+ title: 'BulkUpdateAction[BulkDAGBody]'
+} as const;
+
export const $BulkUpdateAction_BulkDAGRunBody_ = {
properties: {
action: {
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 7fe097a0bba..f796d225c6e 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
@@ -3,7 +3,7 @@
import type { CancelablePromise } from './core/CancelablePromise';
import { OpenAPI } from './core/OpenAPI';
import { request as __request } from './core/request';
-import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData,
GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse,
GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData,
CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse,
GetAssetQueuedEventsData, GetAssetQueuedEventsResponse,
DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData,
GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse,
Dele [...]
+import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData,
GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse,
GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData,
CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse,
GetAssetQueuedEventsData, GetAssetQueuedEventsResponse,
DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData,
GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse,
Dele [...]
export class AssetService {
/**
@@ -2047,6 +2047,28 @@ export class DagService {
});
}
+ /**
+ * Bulk Dags
+ * Bulk pause, resume, or drain Dags by id.
+ * @param data The data for the request.
+ * @param data.requestBody
+ * @returns BulkResponse Successful Response
+ * @throws ApiError
+ */
+ public static bulkDags(data: BulkDagsData):
CancelablePromise<BulkDagsResponse> {
+ return __request(OpenAPI, {
+ method: 'PATCH',
+ url: '/api/v2/dags/bulk',
+ body: data.requestBody,
+ mediaType: 'application/json',
+ errors: {
+ 401: 'Unauthorized',
+ 403: 'Forbidden',
+ 422: 'Validation Error'
+ }
+ });
+ }
+
/**
* Favorite Dag
* Mark the Dag as favorite.
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 776009c9afc..8f70583bf01 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
@@ -341,6 +341,10 @@ export type BulkActionResponse = {
}>;
};
+export type BulkBody_BulkDAGBody_ = {
+ actions: Array<(BulkCreateAction_BulkDAGBody_ |
BulkUpdateAction_BulkDAGBody_ | BulkDeleteAction_BulkDAGBody_)>;
+};
+
export type BulkBody_BulkDAGRunBody_ = {
actions: Array<(BulkCreateAction_BulkDAGRunBody_ |
BulkUpdateAction_BulkDAGRunBody_ | BulkDeleteAction_BulkDAGRunBody_)>;
};
@@ -361,6 +365,18 @@ export type BulkBody_VariableBody_ = {
actions: Array<(BulkCreateAction_VariableBody_ |
BulkUpdateAction_VariableBody_ | BulkDeleteAction_VariableBody_)>;
};
+export type BulkCreateAction_BulkDAGBody_ = {
+ /**
+ * The action to be performed on the entities.
+ */
+ action: "create";
+ /**
+ * A list of entities to be created.
+ */
+ entities: Array<BulkDAGBody>;
+ action_on_existence?: BulkActionOnExistence;
+};
+
export type BulkCreateAction_BulkDAGRunBody_ = {
/**
* The action to be performed on the entities.
@@ -421,6 +437,15 @@ export type BulkCreateAction_VariableBody_ = {
action_on_existence?: BulkActionOnExistence;
};
+/**
+ * Request body for bulk update of Dags.
+ */
+export type BulkDAGBody = {
+ is_paused?: boolean | null;
+ scheduling_state?: DagSchedulingState | null;
+ dag_id: string;
+};
+
/**
* Request body for bulk operations on Dag Runs.
*/
@@ -461,6 +486,18 @@ export type BulkDAGRunClearBody = {
dag_runs?: Array<BulkDAGRunBody>;
};
+export type BulkDeleteAction_BulkDAGBody_ = {
+ /**
+ * The action to be performed on the entities.
+ */
+ action: "delete";
+ /**
+ * A list of entity id/key or entity objects to be deleted.
+ */
+ entities: Array<(string | BulkDAGBody)>;
+ action_on_non_existence?: BulkActionNotOnExistence;
+};
+
export type BulkDeleteAction_BulkDAGRunBody_ = {
/**
* The action to be performed on the entities.
@@ -559,6 +596,22 @@ export type BulkTaskInstanceBody = {
dag_run_id?: string | null;
};
+export type BulkUpdateAction_BulkDAGBody_ = {
+ /**
+ * The action to be performed on the entities.
+ */
+ action: "update";
+ /**
+ * A list of entities to be updated.
+ */
+ entities: Array<BulkDAGBody>;
+ /**
+ * A list of field names to update for each entity.Only these fields will
be applied from the request body to the database model.Any extra fields
provided will be ignored.
+ */
+ update_mask?: Array<(string)> | null;
+ action_on_non_existence?: BulkActionNotOnExistence;
+};
+
export type BulkUpdateAction_BulkDAGRunBody_ = {
/**
* The action to be performed on the entities.
@@ -3780,6 +3833,12 @@ export type GetDagDetailsData = {
export type GetDagDetailsResponse = DAGDetailsResponse;
+export type BulkDagsData = {
+ requestBody: BulkBody_BulkDAGBody_;
+};
+
+export type BulkDagsResponse = BulkResponse;
+
export type FavoriteDagData = {
dagId: string;
};
@@ -6760,6 +6819,29 @@ export type $OpenApiTs = {
};
};
};
+ '/api/v2/dags/bulk': {
+ patch: {
+ req: BulkDagsData;
+ res: {
+ /**
+ * Successful Response
+ */
+ 200: BulkResponse;
+ /**
+ * Unauthorized
+ */
+ 401: HTTPExceptionResponse;
+ /**
+ * Forbidden
+ */
+ 403: HTTPExceptionResponse;
+ /**
+ * Validation Error
+ */
+ 422: HTTPValidationError;
+ };
+ };
+ };
'/api/v2/dags/{dag_id}/favorite': {
post: {
req: FavoriteDagData;
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
index b0a1d8c86b3..d77b200dbb6 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
@@ -97,7 +97,8 @@
"cancelDrain": "Cancel drain",
"drain": "Drain then pause",
"drainPrompt": "This Dag has Dag runs in progress. Draining lets them
finish before the Dag pauses; pausing now stops scheduling immediately without
waiting for them.",
- "pauseNow": "Pause now"
+ "pauseNow": "Pause now",
+ "pauseSelected": "Pause / Drain"
},
"schedulingBanner": {
"message": "This Dag is draining. It will pause once its running Dag runs
finish."
diff --git
a/airflow-core/src/airflow/ui/src/components/PauseOrDrainChoiceModal.tsx
b/airflow-core/src/airflow/ui/src/components/PauseOrDrainChoiceModal.tsx
new file mode 100644
index 00000000000..7ab678e797e
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/components/PauseOrDrainChoiceModal.tsx
@@ -0,0 +1,66 @@
+/*!
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+import type { ReactNode } from "react";
+
+import { Button } from "@chakra-ui/react";
+import { useTranslation } from "react-i18next";
+
+import { Modal } from "src/system-components";
+
+type Props = {
+ readonly children?: ReactNode;
+ readonly displayName: string;
+ readonly onChooseDrain: () => void;
+ readonly onChoosePause: () => void;
+ readonly onOpenChange: () => void;
+ readonly open: boolean;
+};
+
+/** The "drain or pause now?" choice, shared by the single-Dag toggle and the
bulk pause/drain action. */
+export const PauseOrDrainChoiceModal = ({
+ children,
+ displayName,
+ onChooseDrain,
+ onChoosePause,
+ onOpenChange,
+ open,
+}: Props) => {
+ const { t: translate } = useTranslation(["common", "dags"]);
+
+ return (
+ <Modal
+ footerActions={
+ <>
+ <Button data-testid="drain-dag" onClick={onChooseDrain}>
+ {translate("dags:schedulingActions.drain")}
+ </Button>
+ <Button data-testid="pause-dag-now" onClick={onChoosePause}
variant="outline">
+ {translate("dags:schedulingActions.pauseNow")}
+ </Button>
+ </>
+ }
+ onOpenChange={onOpenChange}
+ open={open}
+ title={`${translate("common:pause")} ${displayName}?`}
+ >
+ {translate("dags:schedulingActions.drainPrompt")}
+ {children}
+ </Modal>
+ );
+};
diff --git a/airflow-core/src/airflow/ui/src/components/TogglePause.tsx
b/airflow-core/src/airflow/ui/src/components/TogglePause.tsx
index 2cb0ff52d91..892048ecdb5 100644
--- a/airflow-core/src/airflow/ui/src/components/TogglePause.tsx
+++ b/airflow-core/src/airflow/ui/src/components/TogglePause.tsx
@@ -18,18 +18,19 @@
*/
import { useState } from "react";
-import { Button, useDisclosure } from "@chakra-ui/react";
+import { useDisclosure } from "@chakra-ui/react";
import { useTranslation } from "react-i18next";
import { MdHourglassTop } from "react-icons/md";
import type { DagSchedulingState } from "openapi/requests/types.gen";
-import { Modal, Switch, Tooltip, type SwitchProps } from
"src/system-components";
+import { Switch, Tooltip, type SwitchProps } from "src/system-components";
import { useConfig } from "src/queries/useConfig";
import { useTogglePause } from "src/queries/useTogglePause";
import { ConfirmationModal } from "./ConfirmationModal";
+import { PauseOrDrainChoiceModal } from "./PauseOrDrainChoiceModal";
type Props = {
readonly dagDisplayName?: string;
@@ -124,36 +125,19 @@ export const TogglePause = ({
}}
open={open}
/>
- <Modal
- footerActions={
- <>
- <Button
- data-testid="drain-dag"
- onClick={() => {
- setSchedulingState("draining");
- onChoiceClose();
- }}
- >
- {translate("dags:schedulingActions.drain")}
- </Button>
- <Button
- data-testid="pause-dag-now"
- onClick={() => {
- setSchedulingState("paused");
- onChoiceClose();
- }}
- variant="outline"
- >
- {translate("dags:schedulingActions.pauseNow")}
- </Button>
- </>
- }
+ <PauseOrDrainChoiceModal
+ displayName={displayName}
+ onChooseDrain={() => {
+ setSchedulingState("draining");
+ onChoiceClose();
+ }}
+ onChoosePause={() => {
+ setSchedulingState("paused");
+ onChoiceClose();
+ }}
onOpenChange={onChoiceClose}
open={choiceOpen}
- title={`${translate("common:pause")} ${displayName}?`}
- >
- {translate("dags:schedulingActions.drainPrompt")}
- </Modal>
+ />
</>
);
};
diff --git a/airflow-core/src/airflow/ui/src/mocks/handlers/dags.ts
b/airflow-core/src/airflow/ui/src/mocks/handlers/dags.ts
index e56cb05a3e5..b035a2c1f6a 100644
--- a/airflow-core/src/airflow/ui/src/mocks/handlers/dags.ts
+++ b/airflow-core/src/airflow/ui/src/mocks/handlers/dags.ts
@@ -18,7 +18,7 @@
*/
import { http, HttpResponse, type HttpHandler } from "msw";
-const successDag = {
+export const successDag = {
dag_display_name: "tutorial_taskflow_api_success",
dag_id: "tutorial_taskflow_api_success",
file_token:
@@ -53,7 +53,7 @@ const successDag = {
timetable_type: "NullTimetable",
};
-const failedDag = {
+export const failedDag = {
dag_display_name: "tutorial_taskflow_api_failed",
dag_id: "tutorial_taskflow_api_failed",
file_token:
@@ -88,7 +88,7 @@ const failedDag = {
timetable_type: "CronTriggerTimetable",
};
-const pausedDag = {
+export const pausedDag = {
dag_display_name: "paused_dag",
dag_id: "paused_dag",
file_token:
diff --git
a/airflow-core/src/airflow/ui/src/pages/DagsList/BulkPauseDrainDagsButton.tsx
b/airflow-core/src/airflow/ui/src/pages/DagsList/BulkPauseDrainDagsButton.tsx
new file mode 100644
index 00000000000..cabb36c2777
--- /dev/null
+++
b/airflow-core/src/airflow/ui/src/pages/DagsList/BulkPauseDrainDagsButton.tsx
@@ -0,0 +1,105 @@
+/*!
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+import { Button, useDisclosure } from "@chakra-ui/react";
+import { useTranslation } from "react-i18next";
+
+import type { DAGWithLatestDagRunsResponse, DagSchedulingState } from
"openapi/requests/types.gen";
+
+import { ActionErrors } from "src/components/ActionErrors";
+import { ConfirmationModal } from "src/components/ConfirmationModal";
+import { PauseOrDrainChoiceModal } from
"src/components/PauseOrDrainChoiceModal";
+
+import { useBulkSetDagSchedulingState } from
"src/queries/useBulkSetDagSchedulingState";
+
+type Props = {
+ readonly deselectKeys: (keys: Array<string>) => void;
+ readonly selectedDags: Array<DAGWithLatestDagRunsResponse>;
+};
+
+const BulkPauseDrainDagsButton = ({ deselectKeys, selectedDags }: Props) => {
+ const { t: translate } = useTranslation(["common", "dags"]);
+ const { onClose, onOpen, open } = useDisclosure();
+ const { bulkAction, data, error, isPending, reset } =
useBulkSetDagSchedulingState({
+ deselectKeys,
+ onSuccessConfirm: onClose,
+ });
+
+ // Nothing is running in any selected unpaused Dag, so draining and pausing
now are equivalent
+ // for the whole batch — skip the drain-vs-pause choice, same as the
single-Dag toggle does.
+ // Already-paused Dags are never drained (see runBulkAction), so their runs
don't count.
+ const allIdle = selectedDags.every((dag) => dag.is_paused ||
!dag.has_unfinished_runs);
+ const displayName = `${selectedDags.length} ${translate("dag", { count:
selectedDags.length })}`;
+
+ const runBulkAction = (schedulingState: DagSchedulingState) => {
+ bulkAction({
+ actions: [
+ {
+ action: "update",
+ action_on_non_existence: "skip",
+ entities: selectedDags.map((dag) => ({
+ dag_id: dag.dag_id,
+ // Draining a paused Dag would unpause it and let the scheduler
resume its unfinished runs.
+ scheduling_state: dag.is_paused ? "paused" : schedulingState,
+ })),
+ },
+ ],
+ });
+ };
+
+ const handleOpen = () => {
+ reset();
+ onOpen();
+ };
+
+ return (
+ <>
+ <Button
+ data-testid="bulk-pause-drain-dags"
+ disabled={selectedDags.every((dag) => dag.is_paused)}
+ loading={isPending}
+ onClick={handleOpen}
+ variant="outline"
+ >
+ {translate("dags:schedulingActions.pauseSelected")}
+ </Button>
+ {allIdle ? (
+ <ConfirmationModal
+ header={`${translate("common:pause")} ${displayName}?`}
+ onConfirm={() => runBulkAction("paused")}
+ onOpenChange={onClose}
+ open={open}
+ >
+ <ActionErrors actionResponse={data?.update} error={error} />
+ </ConfirmationModal>
+ ) : (
+ <PauseOrDrainChoiceModal
+ displayName={displayName}
+ onChooseDrain={() => runBulkAction("draining")}
+ onChoosePause={() => runBulkAction("paused")}
+ onOpenChange={onClose}
+ open={open}
+ >
+ <ActionErrors actionResponse={data?.update} error={error} />
+ </PauseOrDrainChoiceModal>
+ )}
+ </>
+ );
+};
+
+export default BulkPauseDrainDagsButton;
diff --git
a/airflow-core/src/airflow/ui/src/pages/DagsList/BulkUnpauseDagsButton.tsx
b/airflow-core/src/airflow/ui/src/pages/DagsList/BulkUnpauseDagsButton.tsx
new file mode 100644
index 00000000000..1b673ac6faf
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/DagsList/BulkUnpauseDagsButton.tsx
@@ -0,0 +1,85 @@
+/*!
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+import { Button, useDisclosure } from "@chakra-ui/react";
+import { useTranslation } from "react-i18next";
+
+import type { DAGWithLatestDagRunsResponse } from "openapi/requests/types.gen";
+
+import { ActionErrors } from "src/components/ActionErrors";
+import { ConfirmationModal } from "src/components/ConfirmationModal";
+
+import { useBulkSetDagSchedulingState } from
"src/queries/useBulkSetDagSchedulingState";
+
+type Props = {
+ readonly deselectKeys: (keys: Array<string>) => void;
+ readonly selectedDags: Array<DAGWithLatestDagRunsResponse>;
+};
+
+const BulkUnpauseDagsButton = ({ deselectKeys, selectedDags }: Props) => {
+ const { t: translate } = useTranslation(["common", "dags"]);
+ const { onClose, onOpen, open } = useDisclosure();
+ const { bulkAction, data, error, isPending, reset } =
useBulkSetDagSchedulingState({
+ deselectKeys,
+ onSuccessConfirm: onClose,
+ });
+
+ const displayName = `${selectedDags.length} ${translate("dag", { count:
selectedDags.length })}`;
+
+ const handleOpen = () => {
+ reset();
+ onOpen();
+ };
+
+ // Unlike the single-Dag toggle, always confirm: unpausing many Dags at once
can trigger
+ // a burst of catchup runs.
+ return (
+ <>
+ <Button
+ data-testid="bulk-unpause-dags"
+ disabled={selectedDags.every((dag) => !dag.is_paused &&
dag.scheduling_state !== "draining")}
+ loading={isPending}
+ onClick={handleOpen}
+ variant="outline"
+ >
+ {translate("common:unpause")}
+ </Button>
+ <ConfirmationModal
+ header={`${translate("common:unpause")} ${displayName}?`}
+ onConfirm={() =>
+ bulkAction({
+ actions: [
+ {
+ action: "update",
+ action_on_non_existence: "skip",
+ // Draining Dags go back to active too, the same as clicking a
draining Dag's toggle.
+ entities: selectedDags.map((dag) => ({ dag_id: dag.dag_id,
scheduling_state: "active" })),
+ },
+ ],
+ })
+ }
+ onOpenChange={onClose}
+ open={open}
+ >
+ <ActionErrors actionResponse={data?.update} error={error} />
+ </ConfirmationModal>
+ </>
+ );
+};
+
+export default BulkUnpauseDagsButton;
diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.test.tsx
b/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.test.tsx
index 9e55cfcb008..810a34f5ad6 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.test.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.test.tsx
@@ -22,8 +22,11 @@ import { delay, http, HttpResponse } from "msw";
import { setupServer, type SetupServer } from "msw/node";
import { afterAll, afterEach, beforeAll, describe, expect, it } from "vitest";
+import type { DagSchedulingState } from "openapi/requests/types.gen";
+
import { DAGS_LIST_DISPLAY_KEY } from "src/constants/localStorage";
import { handlers } from "src/mocks/handlers";
+import { failedDag, pausedDag, successDag } from "src/mocks/handlers/dags";
import { AppWrapper } from "src/utils/AppWrapper";
let server: SetupServer;
@@ -102,6 +105,146 @@ describe("Dag Filters", () => {
});
});
+type BulkRequestBody = {
+ actions: Array<{ entities: Array<{ dag_id: string; scheduling_state: string
}> }>;
+};
+
+const renderDagsAndSelectAll = async (
+ dags: Array<{ scheduling_state?: DagSchedulingState } & typeof successDag>,
+) => {
+ const captured: { body?: BulkRequestBody } = {};
+
+ server.use(
+ http.get("/ui/dags", () => HttpResponse.json({ dags, total_entries:
dags.length })),
+ http.patch("/api/v2/dags/bulk", async ({ request }) => {
+ captured.body = (await request.json()) as BulkRequestBody;
+
+ return HttpResponse.json({
+ update: { errors: [], success: dags.map((dag) => dag.dag_id) },
+ });
+ }),
+ );
+
+ localStorage.setItem(DAGS_LIST_DISPLAY_KEY, JSON.stringify("table"));
+ render(<AppWrapper initialEntries={["/dags"]} />);
+
+ await waitFor(() => expect(screen.getByText(dags[0]?.dag_display_name ??
"")).toBeInTheDocument());
+
+
fireEvent.click(within(screen.getByTestId("table-list")).getAllByRole("checkbox")[0]
as HTMLInputElement);
+
+ return captured;
+};
+
+const getRequestedSchedulingStates = async (captured: { body?: BulkRequestBody
}) => {
+ await waitFor(() => expect(captured.body).toBeDefined());
+
+ return Object.fromEntries(
+ (captured.body?.actions[0]?.entities ?? []).map((entity) =>
[entity.dag_id, entity.scheduling_state]),
+ );
+};
+
+describe("Bulk pause/drain Dags", () => {
+ it.each([
+ { pausedDagHasUnfinishedRuns: false, scenario: "no selected Dag has
unfinished runs" },
+ { pausedDagHasUnfinishedRuns: true, scenario: "only an already-paused Dag
has unfinished runs" },
+ ])(
+ "skips the drain choice and pauses every selected Dag when $scenario",
+ async ({ pausedDagHasUnfinishedRuns }) => {
+ const captured = await renderDagsAndSelectAll([
+ successDag,
+ failedDag,
+ { ...pausedDag, has_unfinished_runs: pausedDagHasUnfinishedRuns },
+ ]);
+
+ fireEvent.click(await screen.findByTestId("bulk-pause-drain-dags"));
+
+ const confirmButton = await
screen.findByTestId("confirmation-confirm-button");
+
+ expect(screen.queryByTestId("drain-dag")).not.toBeInTheDocument();
+ fireEvent.click(confirmButton);
+
+ expect(await getRequestedSchedulingStates(captured)).toEqual({
+ paused_dag: "paused",
+ tutorial_taskflow_api_failed: "paused",
+ tutorial_taskflow_api_success: "paused",
+ });
+ },
+ );
+
+ it.each([
+ { choice: "drain-dag", expectedState: "draining" },
+ { choice: "pause-dag-now", expectedState: "paused" },
+ ])(
+ "offers the drain choice once for a mix of idle and running Dags and
applies $choice to the unpaused ones",
+ async ({ choice, expectedState }) => {
+ const captured = await renderDagsAndSelectAll([
+ { ...successDag, has_unfinished_runs: true },
+ failedDag,
+ { ...pausedDag, has_unfinished_runs: true },
+ ]);
+
+ fireEvent.click(await screen.findByTestId("bulk-pause-drain-dags"));
+
+ const choiceButton = await screen.findByTestId(choice);
+
+
expect(screen.queryByTestId("confirmation-confirm-button")).not.toBeInTheDocument();
+ fireEvent.click(choiceButton);
+
+ expect(await getRequestedSchedulingStates(captured)).toEqual({
+ paused_dag: "paused",
+ tutorial_taskflow_api_failed: expectedState,
+ tutorial_taskflow_api_success: expectedState,
+ });
+ },
+ );
+
+ it("unpauses every selected Dag, including cancelling a drain, in one bulk
request", async () => {
+ const captured = await renderDagsAndSelectAll([
+ successDag,
+ { ...failedDag, scheduling_state: "draining" },
+ pausedDag,
+ ]);
+
+ fireEvent.click(await screen.findByTestId("bulk-unpause-dags"));
+ fireEvent.click(await screen.findByTestId("confirmation-confirm-button"));
+
+ expect(await getRequestedSchedulingStates(captured)).toEqual({
+ paused_dag: "active",
+ tutorial_taskflow_api_failed: "active",
+ tutorial_taskflow_api_success: "active",
+ });
+ });
+
+ it.each([
+ {
+ dags: [pausedDag],
+ pauseDisabled: true,
+ scenario: "every selected Dag is paused",
+ unpauseDisabled: false,
+ },
+ {
+ dags: [successDag, failedDag],
+ pauseDisabled: false,
+ scenario: "every selected Dag is active",
+ unpauseDisabled: true,
+ },
+ {
+ dags: [{ ...failedDag, scheduling_state: "draining" as const }],
+ pauseDisabled: false,
+ scenario: "a selected Dag is draining",
+ unpauseDisabled: false,
+ },
+ ])(
+ "disables the bulk actions that would change nothing when $scenario",
+ async ({ dags, pauseDisabled, unpauseDisabled }) => {
+ await renderDagsAndSelectAll(dags);
+
+ expect(await
screen.findByTestId("bulk-pause-drain-dags")).toHaveProperty("disabled",
pauseDisabled);
+
expect(screen.getByTestId("bulk-unpause-dags")).toHaveProperty("disabled",
unpauseDisabled);
+ },
+ );
+});
+
describe("Dag sorting", () => {
it("sorts cards by latest run after", async () => {
render(<AppWrapper initialEntries={["/dags"]} />);
diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx
b/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx
index 8d505c47bcc..2f0fb4495b1 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagsList/DagsList.tsx
@@ -28,7 +28,7 @@ import type {
DAGWithLatestDagRunsResponse,
} from "openapi/requests/types.gen";
-import { RouterLink } from "src/system-components";
+import { ActionBar, RouterLink } from "src/system-components";
import { DagsLayout } from "src/layouts/DagsLayout";
@@ -37,6 +37,12 @@ import { FavoriteDagButton } from
"src/components/DagActions/FavoriteDagButton";
import DagRunInfo from "src/components/DagRunInfo";
import { DataTable } from "src/components/DataTable";
import type { CardDef } from "src/components/DataTable/types";
+import {
+ SelectionHeaderCheckbox,
+ SelectionProvider,
+ SelectionRowCheckbox,
+ useRowSelection,
+} from "src/components/DataTable/useRowSelection";
import { useTableURLState } from "src/components/DataTable/useTableUrlState";
import { DrainingBadge } from "src/components/DrainingBadge";
import { ErrorAlert } from "src/components/ErrorAlert";
@@ -55,6 +61,8 @@ import { useDags } from "src/queries/useDags";
import { useDocumentTitle } from "src/utils";
import { DagImportErrors } from "../Dashboard/Stats/DagImportErrors";
+import BulkPauseDrainDagsButton from "./BulkPauseDrainDagsButton";
+import BulkUnpauseDagsButton from "./BulkUnpauseDagsButton";
import { DagCard } from "./DagCard";
import { DagRunStateCounts } from "./DagRunStateCounts";
import { DagTags } from "./DagTags";
@@ -62,6 +70,8 @@ import { DagsFilters } from "./DagsFilters";
import { Schedule } from "./Schedule";
import { SortSelect } from "./SortSelect";
+const getRowKey = (dag: DAGWithLatestDagRunsResponse) => dag.dag_id;
+
type GetColumnsParams = {
readonly multiTeam: boolean;
};
@@ -77,6 +87,16 @@ const createColumns = (
runStateContext: RunStateCountsContext,
{ multiTeam }: GetColumnsParams,
): Array<ColumnDef<DAGWithLatestDagRunsResponse>> => [
+ {
+ accessorKey: "select",
+ cell: ({ row }) => <SelectionRowCheckbox colorPalette="brand"
rowKey={getRowKey(row.original)} />,
+ enableHiding: false,
+ enableSorting: false,
+ header: () => <SelectionHeaderCheckbox colorPalette="brand" />,
+ meta: {
+ skeletonWidth: 10,
+ },
+ },
{
accessorKey: "is_paused",
cell: ({ row: { original } }) => (
@@ -360,6 +380,13 @@ export const DagsList = () => {
const columns = createColumns(translate, runStateContext, { multiTeam:
multiTeamEnabled });
const cardDef = createCardDef(runStateContext);
+ const { allRowsSelected, clearSelections, deselectKeys, handleRowSelect,
handleSelectAll, selectedRows } =
+ useRowSelection({
+ data: data?.dags,
+ getKey: getRowKey,
+ });
+ const selectedDags = (data?.dags ?? []).filter((dag) =>
selectedRows.has(getRowKey(dag)));
+
const handleSortChange = ({ value }:
SelectValueChangeDetails<Array<string>>) => {
setTableURLState({
pagination,
@@ -370,45 +397,71 @@ export const DagsList = () => {
});
};
+ const handleDisplayToggleChange = (nextDisplay: "card" | "table") => {
+ setDisplay(nextDisplay);
+ if (nextDisplay !== "table") {
+ // The card view has no selection affordance, so drop any stale
selection made in table view.
+ clearSelections();
+ }
+ };
+
const totalEntries = data?.total_entries ?? 0;
return (
<DagsLayout>
<Box pb={8}>
- <DataTable
- cardDef={cardDef}
- columns={columns}
- data={data?.dags ?? []}
- displayMode={display}
- enableMultiSort
- errorMessage={<ErrorAlert error={error} />}
- filterActions={
- <VStack alignItems="flex-start" gap={2} w="100%">
- <SearchBar
- advancedSearch={advancedSearch}
- defaultValue={dagDisplayNamePattern}
- onChange={handleSearchChange}
- placeholder={translate("dags:search.dags")}
- />
- <DagsFilters />
- </VStack>
- }
- headingExtra={<DagImportErrors iconOnly />}
- initialState={tableURLState}
- isFetching={isFetching}
- isLoading={isLoading}
- modelName="common:dag"
- onDisplayToggleChange={setDisplay}
- onStateChange={setTableURLState}
- presentationActions={
- display === "card" ? (
- <SortSelect handleSortChange={handleSortChange}
orderBy={orderBy[0]} />
- ) : undefined
- }
- showDisplayToggle
- skeletonCount={display === "card" ? 5 : undefined}
- total={totalEntries}
- />
+ <SelectionProvider
+ allRowsSelected={allRowsSelected}
+ onRowSelect={handleRowSelect}
+ onSelectAll={handleSelectAll}
+ selectedRows={selectedRows}
+ >
+ <DataTable
+ cardDef={cardDef}
+ columns={columns}
+ data={data?.dags ?? []}
+ displayMode={display}
+ enableMultiSort
+ errorMessage={<ErrorAlert error={error} />}
+ filterActions={
+ <VStack alignItems="flex-start" gap={2} w="100%">
+ <SearchBar
+ advancedSearch={advancedSearch}
+ defaultValue={dagDisplayNamePattern}
+ onChange={handleSearchChange}
+ placeholder={translate("dags:search.dags")}
+ />
+ <DagsFilters />
+ </VStack>
+ }
+ headingExtra={<DagImportErrors iconOnly />}
+ initialState={tableURLState}
+ isFetching={isFetching}
+ isLoading={isLoading}
+ modelName="common:dag"
+ onDisplayToggleChange={handleDisplayToggleChange}
+ onStateChange={setTableURLState}
+ presentationActions={
+ display === "card" ? (
+ <SortSelect handleSortChange={handleSortChange}
orderBy={orderBy[0]} />
+ ) : undefined
+ }
+ showDisplayToggle
+ skeletonCount={display === "card" ? 5 : undefined}
+ total={totalEntries}
+ />
+ <ActionBar.Root closeOnInteractOutside={false} open={display ===
"table" && selectedRows.size > 0}>
+ <ActionBar.Content>
+ <ActionBar.SelectionTrigger>
+ {selectedRows.size} {translate("selected")}
+ </ActionBar.SelectionTrigger>
+ <ActionBar.Separator />
+ <BulkPauseDrainDagsButton deselectKeys={deselectKeys}
selectedDags={selectedDags} />
+ <BulkUnpauseDagsButton deselectKeys={deselectKeys}
selectedDags={selectedDags} />
+ <ActionBar.CloseTrigger onClick={clearSelections} />
+ </ActionBar.Content>
+ </ActionBar.Root>
+ </SelectionProvider>
</Box>
</DagsLayout>
);
diff --git
a/airflow-core/src/airflow/ui/src/queries/useBulkSetDagSchedulingState.ts
b/airflow-core/src/airflow/ui/src/queries/useBulkSetDagSchedulingState.ts
new file mode 100644
index 00000000000..11d5607b48c
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/queries/useBulkSetDagSchedulingState.ts
@@ -0,0 +1,89 @@
+/*!
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+import { useQueryClient } from "@tanstack/react-query";
+import { useTranslation } from "react-i18next";
+
+import {
+ useDagServiceBulkDags,
+ useDagServiceGetDagDetailsKey,
+ useDagServiceGetDagKey,
+ useDagServiceGetDagsKey,
+ useDagServiceGetDagsUiKey,
+} from "openapi/queries";
+import type { BulkBody_BulkDAGBody_, BulkResponse } from
"openapi/requests/types.gen";
+
+import { toaster } from "src/system-components";
+
+type Props = {
+ readonly deselectKeys: (keys: Array<string>) => void;
+ readonly onSuccessConfirm: VoidFunction;
+};
+
+export const useBulkSetDagSchedulingState = ({ deselectKeys, onSuccessConfirm
}: Props) => {
+ const queryClient = useQueryClient();
+ const { t: translate } = useTranslation(["common", "dags"]);
+
+ const onSuccess = async (responseData: BulkResponse) => {
+ await Promise.all([
+ queryClient.invalidateQueries({ queryKey: [useDagServiceGetDagsKey] }),
+ queryClient.invalidateQueries({ queryKey: [useDagServiceGetDagsUiKey] }),
+ queryClient.invalidateQueries({ queryKey: [useDagServiceGetDagKey] }),
+ queryClient.invalidateQueries({ queryKey:
[useDagServiceGetDagDetailsKey] }),
+ ]);
+
+ const updateResult = responseData.update;
+
+ if (!updateResult) {
+ return;
+ }
+
+ const successKeys = updateResult.success ?? [];
+ const actionErrors = updateResult.errors ?? [];
+
+ if (successKeys.length > 0) {
+ toaster.create({
+ description: translate("toaster.bulkUpdate.success.description", {
+ count: successKeys.length,
+ keys: successKeys.join(", "),
+ resourceName: translate("dag_other"),
+ }),
+ title: translate("toaster.bulkUpdate.success.title", {
+ resourceName: translate("dag_other"),
+ }),
+ type: "success",
+ });
+ deselectKeys(successKeys);
+ }
+
+ // Per-entity failures (status 200 with items in ``errors``) keep the
dialog open
+ // so the user can see what failed; the consumer renders
``data.update.errors``.
+ if (actionErrors.length === 0) {
+ onSuccessConfirm();
+ }
+ };
+
+ const { data, error, isPending, mutate, reset } = useDagServiceBulkDags({
onSuccess });
+
+ const bulkAction = (requestBody: BulkBody_BulkDAGBody_) => {
+ reset();
+ mutate({ requestBody });
+ };
+
+ return { bulkAction, data, error, isPending, reset };
+};
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py
index 73d2d08913c..3f658470717 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py
@@ -22,8 +22,10 @@ from unittest import mock
import pendulum
import pytest
+from fastapi.testclient import TestClient
from sqlalchemy import delete, insert, select, update
+from airflow.api_fastapi.auth.managers.simple.user import SimpleAuthManagerUser
from airflow.models.asset import AssetModel, DagScheduleAssetReference
from airflow.models.dag import DagModel, DagTag
from airflow.models.dag_favorite import DagFavorite
@@ -1186,6 +1188,214 @@ class TestPatchDags(TestDagEndpoint):
assert response.status_code == 403
+class TestBulkDags(TestDagEndpoint):
+ """Unit tests for bulk pause/resume/drain of Dags."""
+
+ def test_bulk_update_pauses_dags(self, test_client, session):
+ response = test_client.patch(
+ "/dags/bulk",
+ json={
+ "actions": [
+ {
+ "action": "update",
+ "entities": [
+ {"dag_id": DAG1_ID, "scheduling_state":
DagSchedulingState.PAUSED},
+ {"dag_id": DAG2_ID, "scheduling_state":
DagSchedulingState.PAUSED},
+ ],
+ }
+ ]
+ },
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert sorted(body["update"]["success"]) == [DAG1_ID, DAG2_ID]
+ assert body["update"]["errors"] == []
+ session.expire_all()
+ assert session.scalar(select(DagModel.is_paused).where(DagModel.dag_id
== DAG1_ID)) is True
+ assert session.scalar(select(DagModel.is_paused).where(DagModel.dag_id
== DAG2_ID)) is True
+ check_last_log(session, dag_id=None, event="bulk_dags",
logical_date=None)
+
+ def test_bulk_update_drains_dags(self, test_client, session):
+ response = test_client.patch(
+ "/dags/bulk",
+ json={
+ "actions": [
+ {
+ "action": "update",
+ "entities": [
+ {"dag_id": DAG1_ID, "scheduling_state":
DagSchedulingState.DRAINING},
+ ],
+ }
+ ]
+ },
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["update"]["success"] == [DAG1_ID]
+ session.expire_all()
+ dag = session.scalar(select(DagModel).where(DagModel.dag_id ==
DAG1_ID))
+ assert dag.is_paused is False
+ assert dag.is_draining is True
+
+ def test_bulk_update_unpauses_paused_and_draining_dags(self, test_client,
session):
+ session.execute(update(DagModel).where(DagModel.dag_id ==
DAG1_ID).values(is_paused=True))
+ session.execute(update(DagModel).where(DagModel.dag_id ==
DAG2_ID).values(is_draining=True))
+ session.commit()
+
+ response = test_client.patch(
+ "/dags/bulk",
+ json={
+ "actions": [
+ {
+ "action": "update",
+ "entities": [
+ {"dag_id": DAG1_ID, "scheduling_state":
DagSchedulingState.ACTIVE},
+ {"dag_id": DAG2_ID, "scheduling_state":
DagSchedulingState.ACTIVE},
+ ],
+ }
+ ]
+ },
+ )
+ assert response.status_code == 200
+ assert sorted(response.json()["update"]["success"]) == [DAG1_ID,
DAG2_ID]
+ session.expire_all()
+ for dag_id in (DAG1_ID, DAG2_ID):
+ dag = session.get(DagModel, dag_id)
+ assert dag.is_paused is False
+ assert dag.is_draining is False
+
+ def test_bulk_update_accepts_legacy_is_paused(self, test_client, session):
+ response = test_client.patch(
+ "/dags/bulk",
+ json={"actions": [{"action": "update", "entities": [{"dag_id":
DAG1_ID, "is_paused": True}]}]},
+ )
+ assert response.status_code == 200
+ session.expire_all()
+ assert session.scalar(select(DagModel.is_paused).where(DagModel.dag_id
== DAG1_ID)) is True
+
+ def test_bulk_update_not_found_fails(self, test_client, session):
+ """FAIL semantics: an unknown dag_id fails the whole action and
nothing is updated."""
+ response = test_client.patch(
+ "/dags/bulk",
+ json={
+ "actions": [
+ {
+ "action": "update",
+ "entities": [
+ {"dag_id": DAG1_ID, "scheduling_state":
DagSchedulingState.PAUSED},
+ {"dag_id": "does_not_exist", "scheduling_state":
DagSchedulingState.PAUSED},
+ ],
+ }
+ ]
+ },
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["update"]["success"] == []
+ errors = body["update"]["errors"]
+ assert len(errors) == 1
+ assert errors[0]["status_code"] == 404
+ assert "does_not_exist" in errors[0]["error"]
+ session.expire_all()
+ assert session.scalar(select(DagModel.is_paused).where(DagModel.dag_id
== DAG1_ID)) is False
+
+ def test_bulk_update_not_found_skip(self, test_client, session):
+ response = test_client.patch(
+ "/dags/bulk",
+ json={
+ "actions": [
+ {
+ "action": "update",
+ "action_on_non_existence": "skip",
+ "entities": [
+ {"dag_id": DAG1_ID, "scheduling_state":
DagSchedulingState.PAUSED},
+ {"dag_id": "does_not_exist", "scheduling_state":
DagSchedulingState.PAUSED},
+ ],
+ }
+ ]
+ },
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["update"]["success"] == [DAG1_ID]
+ assert body["update"]["errors"] == []
+
+ @pytest.mark.parametrize("action", ["create", "delete"])
+ def test_bulk_create_and_delete_not_supported(self, test_client, action):
+ entity = (
+ {"dag_id": DAG1_ID, "scheduling_state": DagSchedulingState.PAUSED}
+ if action == "create"
+ else DAG1_ID
+ )
+ response = test_client.patch(
+ "/dags/bulk",
+ json={"actions": [{"action": action, "entities": [entity]}]},
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body[action]["success"] == []
+ assert body[action]["errors"][0]["status_code"] == 405
+
+ def test_bulk_update_should_response_401(self,
unauthenticated_test_client):
+ response = unauthenticated_test_client.patch(
+ "/dags/bulk",
+ json={"actions": [{"action": "update", "entities": [{"dag_id":
DAG1_ID, "is_paused": True}]}]},
+ )
+ assert response.status_code == 401
+
+ def test_bulk_update_should_response_403(self, unauthorized_test_client):
+ response = unauthorized_test_client.patch(
+ "/dags/bulk",
+ json={"actions": [{"action": "update", "entities": [{"dag_id":
DAG1_ID, "is_paused": True}]}]},
+ )
+ assert response.status_code == 403
+
+ def test_bulk_update_rejects_unauthorized_dag_ids(self, test_client,
session):
+ """A 403 if any entity references a Dag the user can't access; nothing
is updated."""
+ restricted_bundle = DagBundleModel(name="restricted-bundle-bulk")
+ restricted_team = Team(name="restricted-team-bulk")
+ restricted_bundle.teams.append(restricted_team)
+ session.add_all([restricted_bundle, restricted_team])
+ session.flush()
+ session.execute(
+ update(DagModel).where(DagModel.dag_id ==
DAG2_ID).values(bundle_name="restricted-bundle-bulk")
+ )
+ session.commit()
+
+ auth_manager = test_client.app.state.auth_manager
+ token = auth_manager._get_token_signer().generate(
+ auth_manager.serialize_user(
+ SimpleAuthManagerUser(username="limited-user", role="user",
teams=[]),
+ )
+ )
+ with (
+ mock.patch("airflow.models.revoked_token.RevokedToken.is_revoked",
return_value=False),
+ TestClient(
+ test_client.app,
+ headers={"Authorization": f"Bearer {token}"},
+ base_url=str(test_client.base_url),
+ ) as limited_test_client,
+ ):
+ response = limited_test_client.patch(
+ "/dags/bulk",
+ json={
+ "actions": [
+ {
+ "action": "update",
+ "entities": [
+ {"dag_id": DAG1_ID, "is_paused": True},
+ {"dag_id": DAG2_ID, "is_paused": True},
+ ],
+ }
+ ]
+ },
+ )
+
+ assert response.status_code == 403
+ session.expire_all()
+ assert session.scalar(select(DagModel.is_paused).where(DagModel.dag_id
== DAG1_ID)) is False
+
+
class TestFavoriteDag(TestDagEndpoint):
"""Unit tests for favoriting a DAG."""
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 8938ad8d053..dcd041c7f51 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -1704,6 +1704,19 @@ class BulkCreateActionVariableBody(BaseModel):
action_on_existence: BulkActionOnExistence | None = "fail"
+class BulkDAGBody(BaseModel):
+ """
+ Request body for bulk update of Dags.
+ """
+
+ model_config = ConfigDict(
+ extra="forbid",
+ )
+ is_paused: Annotated[bool | None, Field(title="Is Paused")] = None
+ scheduling_state: DagSchedulingState | None = None
+ dag_id: Annotated[str, Field(title="Dag Id")]
+
+
class BulkDAGRunBody(BaseModel):
"""
Request body for bulk operations on Dag Runs.
@@ -1767,6 +1780,20 @@ class BulkDAGRunClearBody(BaseModel):
dag_runs: Annotated[list[BulkDAGRunBody] | None, Field(title="Dag Runs")]
= None
+class BulkDeleteActionBulkDAGBody(BaseModel):
+ model_config = ConfigDict(
+ extra="forbid",
+ )
+ action: Annotated[
+ Literal["delete"], Field(description="The action to be performed on
the entities.", title="Action")
+ ]
+ entities: Annotated[
+ list[str | BulkDAGBody],
+ Field(description="A list of entity id/key or entity objects to be
deleted.", title="Entities"),
+ ]
+ action_on_non_existence: BulkActionNotOnExistence | None = "fail"
+
+
class BulkDeleteActionBulkDAGRunBody(BaseModel):
model_config = ConfigDict(
extra="forbid",
@@ -1843,6 +1870,26 @@ class BulkTaskInstanceBody(BaseModel):
dag_run_id: Annotated[str | None, Field(title="Dag Run Id")] = None
+class BulkUpdateActionBulkDAGBody(BaseModel):
+ model_config = ConfigDict(
+ extra="forbid",
+ )
+ action: Annotated[
+ Literal["update"], Field(description="The action to be performed on
the entities.", title="Action")
+ ]
+ entities: Annotated[
+ list[BulkDAGBody], Field(description="A list of entities to be
updated.", title="Entities")
+ ]
+ update_mask: Annotated[
+ list[str] | None,
+ Field(
+ description="A list of field names to update for each entity.Only
these fields will be applied from the request body to the database model.Any
extra fields provided will be ignored.",
+ title="Update Mask",
+ ),
+ ] = None
+ action_on_non_existence: BulkActionNotOnExistence | None = "fail"
+
+
class BulkUpdateActionBulkDAGRunBody(BaseModel):
model_config = ConfigDict(
extra="forbid",
@@ -2554,6 +2601,19 @@ class BulkBodyVariableBody(BaseModel):
]
+class BulkCreateActionBulkDAGBody(BaseModel):
+ model_config = ConfigDict(
+ extra="forbid",
+ )
+ action: Annotated[
+ Literal["create"], Field(description="The action to be performed on
the entities.", title="Action")
+ ]
+ entities: Annotated[
+ list[BulkDAGBody], Field(description="A list of entities to be
created.", title="Entities")
+ ]
+ action_on_existence: BulkActionOnExistence | None = "fail"
+
+
class BulkCreateActionBulkDAGRunBody(BaseModel):
model_config = ConfigDict(
extra="forbid",
@@ -2805,6 +2865,16 @@ class TaskInstanceHistoryCollectionResponse(BaseModel):
total_entries: Annotated[int, Field(title="Total Entries")]
+class BulkBodyBulkDAGBody(BaseModel):
+ model_config = ConfigDict(
+ extra="forbid",
+ )
+ actions: Annotated[
+ list[BulkCreateActionBulkDAGBody | BulkUpdateActionBulkDAGBody |
BulkDeleteActionBulkDAGBody],
+ Field(title="Actions"),
+ ]
+
+
class BulkBodyBulkDAGRunBody(BaseModel):
model_config = ConfigDict(
extra="forbid",
diff --git a/providers/fab/docs/auth-manager/access-control.rst
b/providers/fab/docs/auth-manager/access-control.rst
index ffb1849aae9..b57fd34b5a0 100644
--- a/providers/fab/docs/auth-manager/access-control.rst
+++ b/providers/fab/docs/auth-manager/access-control.rst
@@ -392,6 +392,10 @@ Stable API Permissions
- PATCH
- DAGs.can_edit
- User
+ * - ``/api/v2/dags/bulk``
+ - PATCH
+ - DAGs.can_edit
+ - User
* - ``/api/v2/dags/{dag_id}``
- DELETE
- DAGs.can_delete
diff --git a/scripts/ci/prek/extract_permissions.py
b/scripts/ci/prek/extract_permissions.py
index 1f0c68d313e..5ade72296bf 100644
--- a/scripts/ci/prek/extract_permissions.py
+++ b/scripts/ci/prek/extract_permissions.py
@@ -225,6 +225,7 @@ def _extract_entity_arg(call_node: ast.Call) -> str | None:
# Map from requires_access_* function name → (resource base name, forced
entity or None)
_FN_TO_RESOURCE_INFO: dict[str, tuple[str, str | None]] = {
"requires_access_dag": ("DAG", None),
+ "requires_access_dag_bulk": ("DAG", None),
"requires_access_dag_from_file_token": ("DAG", None), # reparse
authorizes the file_token's Dags
"requires_access_backfill": ("DAG", "RUN"), # backfill is a DAG.RUN alias
"requires_access_dag_run_bulk": ("DAG", "RUN"), # dag_run bulk is a
DAG.RUN alias
diff --git a/scripts/ci/prek/fab_permissions_doc.py
b/scripts/ci/prek/fab_permissions_doc.py
index 53e81dba467..cfabeb581cf 100755
--- a/scripts/ci/prek/fab_permissions_doc.py
+++ b/scripts/ci/prek/fab_permissions_doc.py
@@ -230,6 +230,14 @@ def fab_permissions_for(entry: PermissionEntry) ->
list[str]:
auth_method = entry.required_permission
action = _METHOD_TO_ACTION.get(auth_method, "can_read")
+ if entry.resource == "DAG" and auth_method == "multi":
+ # A bulk `requires_access_dag_*` call with no method arg (e.g.
+ # ``requires_access_dag_bulk``) always authorizes a write, same as the
base Dag
+ # check below for a ``DAG.<entity>`` bulk call --
``_METHOD_TO_ACTION`` has no
+ # "multi" entry, so it would otherwise fall back to the generic
"can_read".
+ dag_name = display.get("RESOURCE_DAG", "DAGs")
+ return [f"{dag_name}.can_edit"]
+
if entry.resource.startswith("DAG."):
entity = entry.resource.split(".", 1)[1]
consts = dag_map.get(entity, ())
diff --git a/scripts/tests/ci/prek/test_extract_permissions.py
b/scripts/tests/ci/prek/test_extract_permissions.py
index 2900c66f603..763578f4697 100644
--- a/scripts/tests/ci/prek/test_extract_permissions.py
+++ b/scripts/tests/ci/prek/test_extract_permissions.py
@@ -274,6 +274,9 @@ class TestBuildResourceLabel:
def test_alias_dag_run_bulk_forces_run_entity(self):
assert _build_resource_label("requires_access_dag_run_bulk", None) ==
"DAG.RUN"
+ def test_dag_bulk_with_no_entity_returns_base(self):
+ assert _build_resource_label("requires_access_dag_bulk", None) == "DAG"
+
def test_alias_event_log_forces_audit_log(self):
assert _build_resource_label("requires_access_event_log", None) ==
"DAG.AUDIT_LOG"