From 44a81445200818996eaf963780b7313a81fd127c Mon Sep 17 00:00:00 2001 From: justinpakzad <114518232+justinpakzad@users.noreply.github.com> Date: Sat, 21 Mar 2026 13:51:22 -0400 Subject: [PATCH 1/3] fixed pagination bug and updated docstring to clarify dag_id_pattern wildcard usage --- .../openapi/v2-rest-api-generated.yaml | 9 +++- .../core_api/routes/public/dags.py | 29 +++++++++++- .../airflow/ui/openapi-gen/queries/queries.ts | 4 ++ .../ui/openapi-gen/requests/services.gen.ts | 4 ++ .../core_api/routes/public/test_dags.py | 47 +++++++++++++++++++ 5 files changed, 90 insertions(+), 3 deletions(-) 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 b167e463a35f0..3e6564676fa4b 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 @@ -3291,7 +3291,14 @@ paths: tags: - DAG summary: Patch Dags - description: Patch multiple DAGs. + description: 'Patch multiple DAGs. + + + If `dag_id_pattern` is not provided, no DAGs will be matched regardless + + of other filters. To match all DAGs, pass a wildcard value such as `~` + + or `%` for `dag_id_pattern`.' operationId: patch_dags security: - OAuth2PasswordBearer: [] 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 6fbb7c831e9d9..34291ea5b5459 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 @@ -328,7 +328,13 @@ def patch_dags( session: SessionDep, update_mask: list[str] | None = Query(None), ) -> DAGCollectionResponse: - """Patch multiple DAGs.""" + """ + Patch multiple DAGs. + + If `dag_id_pattern` is not provided, no DAGs will be matched regardless + of other filters. To match all DAGs, pass a wildcard value such as `~` + or `%` for `dag_id_pattern`. + """ if update_mask: if update_mask != ["is_paused"]: raise HTTPException( @@ -356,7 +362,26 @@ def patch_dags( session=session, ) dags = session.scalars(dags_select).all() - dags_to_update = {dag.dag_id for dag in dags} + dags_to_update: set[str] = set() + + for current_offset in range(0, total_entries, limit.value): # type: ignore[arg-type] + dags_select, _ = paginated_select( + statement=select(DagModel.dag_id), + filters=[ + exclude_stale, + paused, + dag_id_pattern, + tags, + owners, + editable_dags_filter, + ], + order_by=None, + offset=QueryOffset.depends(current_offset), + limit=limit, + session=session, + ) + dags_to_update.update(session.scalars(dags_select).all()) + session.execute( update(DagModel) .where(DagModel.dag_id.in_(dags_to_update)) 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 dac7a198e59bd..2ac0c9986841b 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts @@ -2115,6 +2115,10 @@ export const useDagRunServicePatchDagRun = Date: Wed, 1 Apr 2026 18:42:38 -0400 Subject: [PATCH 2/3] removed batch loop to update all dags in one shot and added additional test case --- .../core_api/routes/public/dags.py | 38 +++++++------------ .../core_api/routes/public/test_dags.py | 13 +++++++ 2 files changed, 27 insertions(+), 24 deletions(-) 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 34291ea5b5459..417ba24dc2f49 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 @@ -26,10 +26,7 @@ from airflow.api.common import delete_dag as delete_dag_module from airflow.api_fastapi.common.dagbag import DagBagDep, get_latest_version_of_dag -from airflow.api_fastapi.common.db.common import ( - SessionDep, - paginated_select, -) +from airflow.api_fastapi.common.db.common import SessionDep, apply_filters_to_select, paginated_select from airflow.api_fastapi.common.db.dags import generate_dag_with_latest_run_query from airflow.api_fastapi.common.parameters import ( FilterOptionEnum, @@ -362,29 +359,22 @@ def patch_dags( session=session, ) dags = session.scalars(dags_select).all() - dags_to_update: set[str] = set() - - for current_offset in range(0, total_entries, limit.value): # type: ignore[arg-type] - dags_select, _ = paginated_select( - statement=select(DagModel.dag_id), - filters=[ - exclude_stale, - paused, - dag_id_pattern, - tags, - owners, - editable_dags_filter, - ], - order_by=None, - offset=QueryOffset.depends(current_offset), - limit=limit, - session=session, - ) - dags_to_update.update(session.scalars(dags_select).all()) + + filtered_dag_ids = apply_filters_to_select( + statement=select(DagModel.dag_id), + filters=[ + exclude_stale, + paused, + dag_id_pattern, + tags, + owners, + editable_dags_filter, + ], + ) session.execute( update(DagModel) - .where(DagModel.dag_id.in_(dags_to_update)) + .where(DagModel.dag_id.in_(filtered_dag_ids)) .values(is_paused=patch_body.is_paused) .execution_options(synchronize_session="fetch") ) 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 7ff928606e4eb..434cd9e07f3d3 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 @@ -837,6 +837,19 @@ def test_patch_dags_by_tags( paused_dag_ids = {dag["dag_id"] for dag in resp_body["dags"] if dag["is_paused"]} assert paused_dag_ids == set(expected_paused_ids) + def test_patch_dags_updates_all_beyond_limit(self, test_client, session): + response = test_client.patch( + "/dags", + json={"is_paused": True}, + params={"dag_id_pattern": "~", "limit": 1}, + ) + assert response.status_code == 200 + assert len(response.json()["dags"]) == 1 + paused_dags = session.scalars( + select(DagModel.dag_id).where(DagModel.is_paused, ~DagModel.is_stale) + ).all() + assert set(paused_dags) == {DAG1_ID, DAG2_ID} + @mock.patch("airflow.api_fastapi.auth.managers.base_auth_manager.BaseAuthManager.get_authorized_dag_ids") def test_patch_dags_should_call_authorized_dag_ids(self, mock_get_authorized_dag_ids, test_client): mock_get_authorized_dag_ids.return_value = {DAG1_ID, DAG2_ID} From 3fa2ec15911296502138e22d8f7e449ca9fbc22f Mon Sep 17 00:00:00 2001 From: justinpakzad <114518232+justinpakzad@users.noreply.github.com> Date: Wed, 1 Apr 2026 22:47:47 -0400 Subject: [PATCH 3/3] Fixed MySQL subquery issue --- .../src/airflow/api_fastapi/core_api/routes/public/dags.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 417ba24dc2f49..b83cc1ef2236c 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 @@ -370,11 +370,11 @@ def patch_dags( owners, editable_dags_filter, ], - ) + ).subquery() session.execute( update(DagModel) - .where(DagModel.dag_id.in_(filtered_dag_ids)) + .where(DagModel.dag_id.in_(select(filtered_dag_ids.c.dag_id))) .values(is_paused=patch_body.is_paused) .execution_options(synchronize_session="fetch") )