From de93e086912aaacb2f809e958ac3f8927662ece7 Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Wed, 13 Aug 2025 07:12:36 +0800 Subject: [PATCH 1/8] Fix: Clean up FAB permissions when deleting DAGs (#50905) This commit adds cleanup logic to the delete_dag function to: - Remove DAG-specific resources from ab_view_menu table - Clean up associated permissions and role associations --- .../src/airflow/api/common/delete_dag.py | 103 ++++++++++++++++++ 1 file changed, 103 insertions(+) diff --git a/airflow-core/src/airflow/api/common/delete_dag.py b/airflow-core/src/airflow/api/common/delete_dag.py index a617bcdc5c61c..40eb9520bbd14 100644 --- a/airflow-core/src/airflow/api/common/delete_dag.py +++ b/airflow-core/src/airflow/api/common/delete_dag.py @@ -88,4 +88,107 @@ def delete_dag(dag_id: str, keep_records_in_log: bool = True, session: Session = .execution_options(synchronize_session="fetch") ) + # Clean up DAG-specific permissions from Flask-AppBuilder tables + _cleanup_dag_permissions(dag_id, session) + return count + + +def _cleanup_dag_permissions(dag_id: str, session: Session) -> None: + """ + Clean up DAG-specific permissions from Flask-AppBuilder tables. + + When a DAG is deleted, we need to clean up the corresponding permissions + to prevent orphaned entries in the ab_view_menu table. + + This addresses issue #50905: Deleted DAGs not removed from ab_view_menu table + and show up in permissions. + """ + from airflow.configuration import conf + + if "FabAuthManager" not in conf.get("core", "auth_manager"): + return + + # Try to import FAB models with version compatibility + def _get_fab_models(): + """Get FAB models with version compatibility handling.""" + try: + from airflow.providers.fab.auth_manager import models as fab_models + + return fab_models + except ImportError: + try: + # Handle Pre-airflow 2.9 case where FAB was part of the core airflow + from airflow.providers.fab.auth.managers.fab import models as fab_models + + return fab_models + except ImportError: + # If FAB provider is not available, skip cleanup + return None + except RuntimeError as e: + # Handle case where FAB provider is not even usable + if "needs Apache Airflow 2.9.0" in str(e): + try: + from airflow.providers.fab.auth.managers.fab import models as fab_models + + return fab_models + except ImportError: + return None + else: + return None + + fab_models = _get_fab_models() + if fab_models is None: + return + + Permission = fab_models.Permission + Resource = fab_models.Resource + assoc_permission_role = fab_models.assoc_permission_role + + from airflow.security.permissions import RESOURCE_DAG_PREFIX + + # Find all DAG-specific resources that match this dag_id + dag_resource_name = f"{RESOURCE_DAG_PREFIX}{dag_id}" + dag_resources = ( + session.query(Resource) + .filter( + Resource.name.in_( + [ + dag_resource_name, # DAG:dag_id + f"DAG Run:{dag_id}", # DAG_RUN:dag_id + f"Task Instance:{dag_id}", # TASK_INSTANCE:dag_id (if exists) + ] + ) + ) + .all() + ) + + if not dag_resources: + return + + dag_resource_ids = [resource.id for resource in dag_resources] + + # Find all permissions associated with these resources + dag_permissions = session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).all() + + if not dag_permissions: + # Delete resources even if no permissions exist + session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) + return + + dag_permission_ids = [permission.id for permission in dag_permissions] + + # Delete permission-role associations first (foreign key constraint) + session.query(assoc_permission_role).filter( + assoc_permission_role.c.permission_view_id.in_(dag_permission_ids) + ).delete(synchronize_session=False) + + # Delete permissions + session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).delete( + synchronize_session=False + ) + + # Delete resources (ab_view_menu entries) + session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) + + log.info("Cleaned up %d DAG-specific permissions for dag_id: %s", len(dag_permissions), dag_id) From 489d34e707638da246eae56313e635743917fe8d Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Wed, 20 Aug 2025 07:15:41 +0800 Subject: [PATCH 2/8] test: Add comprehensive tests for DAG permission cleanup Test Coverage: - Non-FAB auth manager scenarios (graceful skip) - FAB provider unavailable scenarios (import error handling) - Full cleanup process with mocked FAB models - Integration with delete_dag function - Edge cases (non-existent DAGs) --- .../tests/unit/api/common/test_delete_dag.py | 142 ++++++++++++++++++ 1 file changed, 142 insertions(+) create mode 100644 airflow-core/tests/unit/api/common/test_delete_dag.py diff --git a/airflow-core/tests/unit/api/common/test_delete_dag.py b/airflow-core/tests/unit/api/common/test_delete_dag.py new file mode 100644 index 0000000000000..47be46c951e8b --- /dev/null +++ b/airflow-core/tests/unit/api/common/test_delete_dag.py @@ -0,0 +1,142 @@ +# 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 unittest.mock import MagicMock, patch + +import pytest + +from airflow.api.common.delete_dag import _cleanup_dag_permissions, delete_dag +from airflow.exceptions import DagNotFound + +pytestmark = pytest.mark.db_test + + +class TestCleanupDagPermissions: + """Test cases for _cleanup_dag_permissions function.""" + + @patch("airflow.configuration.conf.get") + def test_cleanup_dag_permissions_non_fab_auth_manager(self, mock_conf_get, session): + """Test that cleanup is skipped when not using FabAuthManager.""" + mock_conf_get.return_value = "BasicAuthManager" + + # Should return early without any database operations + _cleanup_dag_permissions("test_dag", session) + + mock_conf_get.assert_called_once_with("core", "auth_manager") + + @patch("airflow.configuration.conf.get") + def test_cleanup_dag_permissions_fab_provider_not_available(self, mock_conf_get, session): + """Test that cleanup is skipped when FAB provider is not available.""" + mock_conf_get.return_value = "airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager" + + # Mock ImportError specifically for FAB provider imports + original_import = __builtins__["__import__"] + + def mock_import(name, *args, **kwargs): + if "airflow.providers.fab" in name: + raise ImportError(f"No module named '{name}'") + return original_import(name, *args, **kwargs) + + with patch("builtins.__import__", side_effect=mock_import): + # Mock session.query to track if it's called + with patch.object(session, "query") as mock_query: + # Should return early when FAB models are not available + _cleanup_dag_permissions("test_dag", session) + + # No database operations should be performed + mock_query.assert_not_called() + + @patch("airflow.configuration.conf.get") + @patch("airflow.api.common.delete_dag.log") + def test_cleanup_dag_permissions_full_cleanup(self, mock_log, mock_conf_get, session): + """Test full cleanup process when resources and permissions exist.""" + mock_conf_get.return_value = "airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager" + + # Mock successful FAB import + with patch("airflow.providers.fab.auth_manager.models") as mock_models: + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock resources and permissions + mock_dag_resource = MagicMock() + mock_dag_resource.id = 1 + mock_dag_permission = MagicMock() + mock_dag_permission.id = 1 + + def mock_query_side_effect(model): + mock_query = MagicMock() + if model == mock_resource: + mock_query.filter.return_value.all.return_value = [mock_dag_resource] + mock_query.filter.return_value.delete.return_value = None + elif model == mock_permission: + mock_query.filter.return_value.all.return_value = [mock_dag_permission] + mock_query.filter.return_value.delete.return_value = None + elif model == mock_assoc: + mock_query.filter.return_value.delete.return_value = None + return mock_query + + # Mock session.query method properly + mock_session_query = MagicMock(side_effect=mock_query_side_effect) + session.query = mock_session_query + + _cleanup_dag_permissions("test_dag", session) + + # Should log successful cleanup + mock_log.info.assert_called_once_with( + "Cleaned up %d DAG-specific permissions for dag_id: %s", 1, "test_dag" + ) + + +class TestDeleteDag: + """Test cases for delete_dag function integration with permission cleanup.""" + + def test_delete_dag_calls_cleanup_permissions(self, dag_maker, session): + """Test that delete_dag calls _cleanup_dag_permissions.""" + # Mock all the dependencies to avoid database constraints + with ( + patch("airflow.api.common.delete_dag._cleanup_dag_permissions") as mock_cleanup, + patch.object(session, "scalar") as mock_scalar, + patch.object(session, "execute") as mock_execute, + ): + # Mock that there are no running task instances + mock_scalar.side_effect = [ + None, + MagicMock(dag_id="test_dag"), + ] # First call for running TIs, second for DAG model + mock_execute.return_value.rowcount = 1 + + # This should call the cleanup function + result = delete_dag("test_dag", session=session) + + # Should call cleanup permissions with correct arguments + mock_cleanup.assert_called_once_with("test_dag", session) + # Should return count of deleted records + assert result >= 0 + + def test_delete_dag_nonexistent_dag_does_not_call_cleanup(self, session): + """Test that delete_dag does not call cleanup for non-existent DAG.""" + with patch("airflow.api.common.delete_dag._cleanup_dag_permissions") as mock_cleanup: + with pytest.raises(DagNotFound): + delete_dag("nonexistent_dag", session=session) + + # Should not call cleanup for non-existent DAG + mock_cleanup.assert_not_called() From f8b17bcf3a50780c62488360a39295452521228f Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Fri, 22 Aug 2025 07:06:13 +0800 Subject: [PATCH 3/8] fix: clean up DAG permissions when deleting DAGs - : Simplified core logic using auth manager delegation - : Added cleanup_dag_permissions method - : New specialized module for FAB permission cleanup - Enhanced test coverage for multiple auth manager scenarios --- .../src/airflow/api/common/delete_dag.py | 113 +++---------- .../tests/unit/api/common/test_delete_dag.py | 152 ++++++++--------- .../fab/auth_manager/dag_permissions.py | 104 ++++++++++++ .../fab/auth_manager/fab_auth_manager.py | 11 ++ .../fab/auth_manager/test_dag_permissions.py | 154 ++++++++++++++++++ 5 files changed, 359 insertions(+), 175 deletions(-) create mode 100644 providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py create mode 100644 providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py diff --git a/airflow-core/src/airflow/api/common/delete_dag.py b/airflow-core/src/airflow/api/common/delete_dag.py index 40eb9520bbd14..0b3f9f7787c50 100644 --- a/airflow-core/src/airflow/api/common/delete_dag.py +++ b/airflow-core/src/airflow/api/common/delete_dag.py @@ -88,7 +88,7 @@ def delete_dag(dag_id: str, keep_records_in_log: bool = True, session: Session = .execution_options(synchronize_session="fetch") ) - # Clean up DAG-specific permissions from Flask-AppBuilder tables + # Clean up DAG-specific permissions _cleanup_dag_permissions(dag_id, session) return count @@ -96,99 +96,26 @@ def delete_dag(dag_id: str, keep_records_in_log: bool = True, session: Session = def _cleanup_dag_permissions(dag_id: str, session: Session) -> None: """ - Clean up DAG-specific permissions from Flask-AppBuilder tables. + Clean up DAG-specific permissions from the auth manager. - When a DAG is deleted, we need to clean up the corresponding permissions - to prevent orphaned entries in the ab_view_menu table. - - This addresses issue #50905: Deleted DAGs not removed from ab_view_menu table - and show up in permissions. + This delegates the cleanup to the appropriate auth manager implementation. """ - from airflow.configuration import conf - - if "FabAuthManager" not in conf.get("core", "auth_manager"): - return - - # Try to import FAB models with version compatibility - def _get_fab_models(): - """Get FAB models with version compatibility handling.""" - try: - from airflow.providers.fab.auth_manager import models as fab_models - - return fab_models - except ImportError: - try: - # Handle Pre-airflow 2.9 case where FAB was part of the core airflow - from airflow.providers.fab.auth.managers.fab import models as fab_models - - return fab_models - except ImportError: - # If FAB provider is not available, skip cleanup - return None - except RuntimeError as e: - # Handle case where FAB provider is not even usable - if "needs Apache Airflow 2.9.0" in str(e): - try: - from airflow.providers.fab.auth.managers.fab import models as fab_models - - return fab_models - except ImportError: - return None - else: - return None - - fab_models = _get_fab_models() - if fab_models is None: - return - - Permission = fab_models.Permission - Resource = fab_models.Resource - assoc_permission_role = fab_models.assoc_permission_role - - from airflow.security.permissions import RESOURCE_DAG_PREFIX - - # Find all DAG-specific resources that match this dag_id - dag_resource_name = f"{RESOURCE_DAG_PREFIX}{dag_id}" - dag_resources = ( - session.query(Resource) - .filter( - Resource.name.in_( - [ - dag_resource_name, # DAG:dag_id - f"DAG Run:{dag_id}", # DAG_RUN:dag_id - f"Task Instance:{dag_id}", # TASK_INSTANCE:dag_id (if exists) - ] + try: + from airflow.api_fastapi.app import get_auth_manager + + auth_manager = get_auth_manager() + if hasattr(auth_manager, "cleanup_dag_permissions"): + auth_manager.cleanup_dag_permissions(dag_id, session) + log.info("Successfully cleaned up DAG permissions for dag_id: %s", dag_id) + else: + log.debug( + "Auth manager %s does not support DAG permission cleanup, skipping", + type(auth_manager).__name__, ) + except Exception as e: + # If auth manager is not available or fails, silently skip cleanup + # This ensures DAG deletion continues even if permission cleanup fails + log.warning( + "Failed to clean up DAG permissions for dag_id %s: %s. DAG deletion will continue.", dag_id, e ) - .all() - ) - - if not dag_resources: - return - - dag_resource_ids = [resource.id for resource in dag_resources] - - # Find all permissions associated with these resources - dag_permissions = session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).all() - - if not dag_permissions: - # Delete resources even if no permissions exist - session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) - return - - dag_permission_ids = [permission.id for permission in dag_permissions] - - # Delete permission-role associations first (foreign key constraint) - session.query(assoc_permission_role).filter( - assoc_permission_role.c.permission_view_id.in_(dag_permission_ids) - ).delete(synchronize_session=False) - - # Delete permissions - session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).delete( - synchronize_session=False - ) - - # Delete resources (ab_view_menu entries) - session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) - - log.info("Cleaned up %d DAG-specific permissions for dag_id: %s", len(dag_permissions), dag_id) + pass diff --git a/airflow-core/tests/unit/api/common/test_delete_dag.py b/airflow-core/tests/unit/api/common/test_delete_dag.py index 47be46c951e8b..83c2b9da1c9f4 100644 --- a/airflow-core/tests/unit/api/common/test_delete_dag.py +++ b/airflow-core/tests/unit/api/common/test_delete_dag.py @@ -29,81 +29,45 @@ class TestCleanupDagPermissions: """Test cases for _cleanup_dag_permissions function.""" - @patch("airflow.configuration.conf.get") - def test_cleanup_dag_permissions_non_fab_auth_manager(self, mock_conf_get, session): - """Test that cleanup is skipped when not using FabAuthManager.""" - mock_conf_get.return_value = "BasicAuthManager" + @patch("airflow.api_fastapi.app.get_auth_manager") + def test_cleanup_dag_permissions_auth_manager_without_cleanup_method( + self, mock_get_auth_manager, session + ): + """Test that cleanup is skipped when auth manager doesn't have cleanup_dag_permissions method.""" + mock_auth_manager = MagicMock() + # Mock auth manager without cleanup_dag_permissions method + del mock_auth_manager.cleanup_dag_permissions + mock_get_auth_manager.return_value = mock_auth_manager + + # Should return early without calling any methods + _cleanup_dag_permissions("test_dag", session) + + # Should still get the auth manager but not call cleanup + mock_get_auth_manager.assert_called_once() + + @patch("airflow.api_fastapi.app.get_auth_manager") + def test_cleanup_dag_permissions_auth_manager_with_cleanup_method(self, mock_get_auth_manager, session): + """Test that cleanup is called when auth manager has cleanup_dag_permissions method.""" + mock_auth_manager = MagicMock() + mock_auth_manager.cleanup_dag_permissions = MagicMock() + mock_get_auth_manager.return_value = mock_auth_manager - # Should return early without any database operations _cleanup_dag_permissions("test_dag", session) - mock_conf_get.assert_called_once_with("core", "auth_manager") - - @patch("airflow.configuration.conf.get") - def test_cleanup_dag_permissions_fab_provider_not_available(self, mock_conf_get, session): - """Test that cleanup is skipped when FAB provider is not available.""" - mock_conf_get.return_value = "airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager" - - # Mock ImportError specifically for FAB provider imports - original_import = __builtins__["__import__"] - - def mock_import(name, *args, **kwargs): - if "airflow.providers.fab" in name: - raise ImportError(f"No module named '{name}'") - return original_import(name, *args, **kwargs) - - with patch("builtins.__import__", side_effect=mock_import): - # Mock session.query to track if it's called - with patch.object(session, "query") as mock_query: - # Should return early when FAB models are not available - _cleanup_dag_permissions("test_dag", session) - - # No database operations should be performed - mock_query.assert_not_called() - - @patch("airflow.configuration.conf.get") - @patch("airflow.api.common.delete_dag.log") - def test_cleanup_dag_permissions_full_cleanup(self, mock_log, mock_conf_get, session): - """Test full cleanup process when resources and permissions exist.""" - mock_conf_get.return_value = "airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager" - - # Mock successful FAB import - with patch("airflow.providers.fab.auth_manager.models") as mock_models: - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock resources and permissions - mock_dag_resource = MagicMock() - mock_dag_resource.id = 1 - mock_dag_permission = MagicMock() - mock_dag_permission.id = 1 - - def mock_query_side_effect(model): - mock_query = MagicMock() - if model == mock_resource: - mock_query.filter.return_value.all.return_value = [mock_dag_resource] - mock_query.filter.return_value.delete.return_value = None - elif model == mock_permission: - mock_query.filter.return_value.all.return_value = [mock_dag_permission] - mock_query.filter.return_value.delete.return_value = None - elif model == mock_assoc: - mock_query.filter.return_value.delete.return_value = None - return mock_query - - # Mock session.query method properly - mock_session_query = MagicMock(side_effect=mock_query_side_effect) - session.query = mock_session_query + # Should call the auth manager's cleanup method + mock_get_auth_manager.assert_called_once() + mock_auth_manager.cleanup_dag_permissions.assert_called_once_with("test_dag", session) - _cleanup_dag_permissions("test_dag", session) + @patch("airflow.api_fastapi.app.get_auth_manager") + def test_cleanup_dag_permissions_handles_get_auth_manager_exception(self, mock_get_auth_manager, session): + """Test that cleanup handles exceptions from get_auth_manager gracefully.""" + mock_get_auth_manager.side_effect = Exception("Auth manager not available") - # Should log successful cleanup - mock_log.info.assert_called_once_with( - "Cleaned up %d DAG-specific permissions for dag_id: %s", 1, "test_dag" - ) + # Should not raise exception, should handle gracefully + try: + _cleanup_dag_permissions("test_dag", session) + except Exception: + pytest.fail("_cleanup_dag_permissions should handle get_auth_manager exceptions gracefully") class TestDeleteDag: @@ -111,19 +75,11 @@ class TestDeleteDag: def test_delete_dag_calls_cleanup_permissions(self, dag_maker, session): """Test that delete_dag calls _cleanup_dag_permissions.""" - # Mock all the dependencies to avoid database constraints - with ( - patch("airflow.api.common.delete_dag._cleanup_dag_permissions") as mock_cleanup, - patch.object(session, "scalar") as mock_scalar, - patch.object(session, "execute") as mock_execute, - ): - # Mock that there are no running task instances - mock_scalar.side_effect = [ - None, - MagicMock(dag_id="test_dag"), - ] # First call for running TIs, second for DAG model - mock_execute.return_value.rowcount = 1 + with dag_maker(dag_id="test_dag", session=session): + pass + # Mock the cleanup function to track if it's called + with patch("airflow.api.common.delete_dag._cleanup_dag_permissions") as mock_cleanup: # This should call the cleanup function result = delete_dag("test_dag", session=session) @@ -140,3 +96,35 @@ def test_delete_dag_nonexistent_dag_does_not_call_cleanup(self, session): # Should not call cleanup for non-existent DAG mock_cleanup.assert_not_called() + + @patch("airflow.api_fastapi.app.get_auth_manager") + def test_delete_dag_with_fab_auth_manager(self, mock_get_auth_manager, dag_maker, session): + """Test that delete_dag works correctly with FabAuthManager.""" + with dag_maker(dag_id="test_dag", session=session): + pass + + # Mock FabAuthManager with cleanup method + mock_auth_manager = MagicMock() + mock_auth_manager.cleanup_dag_permissions = MagicMock() + mock_get_auth_manager.return_value = mock_auth_manager + + result = delete_dag("test_dag", session=session) + + # Should call the auth manager's cleanup method + mock_auth_manager.cleanup_dag_permissions.assert_called_once_with("test_dag", session) + assert result >= 0 + + @patch("airflow.api_fastapi.app.get_auth_manager") + def test_delete_dag_with_simple_auth_manager(self, mock_get_auth_manager, dag_maker, session): + """Test that delete_dag works correctly with SimpleAuthManager (no cleanup method).""" + with dag_maker(dag_id="test_dag", session=session): + pass + + # Mock SimpleAuthManager without cleanup method + mock_auth_manager = MagicMock() + del mock_auth_manager.cleanup_dag_permissions + mock_get_auth_manager.return_value = mock_auth_manager + + # Should not raise any exception + result = delete_dag("test_dag", session=session) + assert result >= 0 diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py b/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py new file mode 100644 index 0000000000000..38e45c6837e3a --- /dev/null +++ b/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py @@ -0,0 +1,104 @@ +# +# 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. +"""DAG permissions management for FAB Auth Manager.""" + +from __future__ import annotations + +import logging +from typing import TYPE_CHECKING + +from airflow.security.permissions import RESOURCE_DAG_PREFIX + +if TYPE_CHECKING: + from sqlalchemy.orm import Session + +log = logging.getLogger(__name__) + + +def cleanup_dag_permissions(dag_id: str, session: Session | None = None) -> None: + """ + Clean up DAG-specific permissions from Flask-AppBuilder tables. + + When a DAG is deleted, we need to clean up the corresponding permissions + to prevent orphaned entries in the ab_view_menu table. + + This addresses issue #50905: Deleted DAGs not removed from ab_view_menu table + and show up in permissions. + + :param dag_id: Specific DAG ID to clean up. + :param session: Database session. If None, creates a new session. + """ + from airflow.utils.session import create_session + + if session is None: + with create_session() as session: + _cleanup_dag_permissions_impl(dag_id, session) + else: + _cleanup_dag_permissions_impl(dag_id, session) + + +def _cleanup_dag_permissions_impl(dag_id: str, session: Session) -> None: + """Implement DAG permissions cleanup.""" + from airflow.providers.fab.auth_manager.models import Permission, Resource, assoc_permission_role + + # Clean up specific DAG permissions + dag_resource_name = f"{RESOURCE_DAG_PREFIX}{dag_id}" + dag_resources = ( + session.query(Resource) + .filter( + Resource.name.in_( + [ + dag_resource_name, # DAG:dag_id + f"DAG Run:{dag_id}", # DAG_RUN:dag_id + f"Task Instance:{dag_id}", # TASK_INSTANCE:dag_id (if exists) + ] + ) + ) + .all() + ) + log.info("Cleaning up DAG-specific permissions for dag_id: %s", dag_id) + + if not dag_resources: + return + + dag_resource_ids = [resource.id for resource in dag_resources] + + # Find all permissions associated with these resources + dag_permissions = session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).all() + + if not dag_permissions: + # Delete resources even if no permissions exist + session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) + return + + dag_permission_ids = [permission.id for permission in dag_permissions] + + # Delete permission-role associations first (foreign key constraint) + session.query(assoc_permission_role).filter( + assoc_permission_role.c.permission_view_id.in_(dag_permission_ids) + ).delete(synchronize_session=False) + + # Delete permissions + session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).delete( + synchronize_session=False + ) + + # Delete resources (ab_view_menu entries) + session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) + + log.info("Cleaned up %d DAG-specific permissions", len(dag_permissions)) diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py b/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py index 1cd9590a508bf..3725c01ccbfaf 100644 --- a/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py +++ b/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py @@ -652,6 +652,17 @@ def _get_user_permissions(user: User): return [] return getattr(user, "perms") or [] + def cleanup_dag_permissions(self, dag_id: str, session: Session) -> None: + """ + Clean up DAG-specific permissions from Flask-AppBuilder tables. + + :param dag_id: the DAG ID to clean up permissions for + :param session: database session + """ + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + cleanup_dag_permissions(dag_id=dag_id, session=session) + def _sync_appbuilder_roles(self): """ Sync appbuilder roles to DB. diff --git a/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py b/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py new file mode 100644 index 0000000000000..4398625568e19 --- /dev/null +++ b/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py @@ -0,0 +1,154 @@ +# 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 unittest.mock import MagicMock, patch + + +class TestDagPermissions: + """Test cases for dag_permissions module.""" + + def test_cleanup_dag_permissions_with_session(self): + """Test cleanup_dag_permissions with provided session.""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models + with patch("airflow.providers.fab.auth_manager.models") as mock_models: + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session + mock_session = MagicMock() + + # Mock no resources found + mock_session.query.return_value.filter.return_value.all.return_value = [] + + cleanup_dag_permissions("test_dag", mock_session) + + # Should call session.query for Resource + mock_session.query.assert_called() + + def test_cleanup_dag_permissions_without_session(self): + """Test cleanup_dag_permissions without provided session (creates new session).""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models + with patch("airflow.providers.fab.auth_manager.models") as mock_models: + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock create_session + with patch("airflow.utils.session.create_session") as mock_create_session: + mock_session = MagicMock() + mock_create_session.return_value.__enter__.return_value = mock_session + + # Mock no resources found + mock_session.query.return_value.filter.return_value.all.return_value = [] + + cleanup_dag_permissions("test_dag") + + # Should create a new session + mock_create_session.assert_called_once() + + def test_cleanup_dag_permissions_with_resources_and_permissions(self): + """Test cleanup_dag_permissions with actual resources and permissions.""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models + with patch("airflow.providers.fab.auth_manager.models") as mock_models: + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session + mock_session = MagicMock() + + # Mock resources exist + mock_dag_resource = MagicMock() + mock_dag_resource.id = 1 + mock_session.query.return_value.filter.return_value.all.return_value = [mock_dag_resource] + + # Mock permissions exist + mock_dag_permission = MagicMock() + mock_dag_permission.id = 1 + + def mock_query_side_effect(model): + mock_query = MagicMock() + if model == mock_resource: + mock_query.filter.return_value.all.return_value = [mock_dag_resource] + mock_query.filter.return_value.delete.return_value = None + elif model == mock_permission: + mock_query.filter.return_value.all.return_value = [mock_dag_permission] + mock_query.filter.return_value.delete.return_value = None + elif model == mock_assoc: + mock_query.filter.return_value.delete.return_value = None + return mock_query + + mock_session.query.side_effect = mock_query_side_effect + + cleanup_dag_permissions("test_dag", mock_session) + + # Should perform deletion operations + assert mock_session.query.call_count >= 3 # Resource, Permission, assoc queries + + def test_cleanup_dag_permissions_with_resources_but_no_permissions(self): + """Test cleanup_dag_permissions with resources but no permissions.""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models + with patch("airflow.providers.fab.auth_manager.models") as mock_models: + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session + mock_session = MagicMock() + + # Mock resources exist but no permissions + mock_dag_resource = MagicMock() + mock_dag_resource.id = 1 + + def mock_query_side_effect(model): + mock_query = MagicMock() + if model == mock_resource: + mock_query.filter.return_value.all.return_value = [mock_dag_resource] + mock_query.filter.return_value.delete.return_value = None + elif model == mock_permission: + mock_query.filter.return_value.all.return_value = [] # No permissions + mock_query.filter.return_value.delete.return_value = None + return mock_query + + mock_session.query.side_effect = mock_query_side_effect + + cleanup_dag_permissions("test_dag", mock_session) + + # Should still delete resources even without permissions + assert mock_session.query.call_count >= 2 # Resource queries From 3af12529c75961759de7fffb6cb1820c9ae0445b Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Sat, 30 Aug 2025 15:52:18 +0800 Subject: [PATCH 4/8] feat(fab-auth-manager): replace automatic DAG permission cleanup with CLI command Replace automatic cleanup of DAG permissions during DAG deletion with a new CLI command approach as requested in code review. This gives operators explicit control over when and how DAG permissions are cleaned up. --- .../src/airflow/api/common/delete_dag.py | 30 ---- .../tests/unit/api/common/test_delete_dag.py | 130 ------------------ .../auth_manager/cli_commands/definition.py | 32 +++++ .../cli_commands/permissions_command.py | 122 ++++++++++++++++ .../fab/auth_manager/fab_auth_manager.py | 13 +- .../cli_commands/test_permissions_command.py | 103 ++++++++++++++ 6 files changed, 259 insertions(+), 171 deletions(-) delete mode 100644 airflow-core/tests/unit/api/common/test_delete_dag.py create mode 100644 providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py create mode 100644 providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py diff --git a/airflow-core/src/airflow/api/common/delete_dag.py b/airflow-core/src/airflow/api/common/delete_dag.py index 0b3f9f7787c50..a617bcdc5c61c 100644 --- a/airflow-core/src/airflow/api/common/delete_dag.py +++ b/airflow-core/src/airflow/api/common/delete_dag.py @@ -88,34 +88,4 @@ def delete_dag(dag_id: str, keep_records_in_log: bool = True, session: Session = .execution_options(synchronize_session="fetch") ) - # Clean up DAG-specific permissions - _cleanup_dag_permissions(dag_id, session) - return count - - -def _cleanup_dag_permissions(dag_id: str, session: Session) -> None: - """ - Clean up DAG-specific permissions from the auth manager. - - This delegates the cleanup to the appropriate auth manager implementation. - """ - try: - from airflow.api_fastapi.app import get_auth_manager - - auth_manager = get_auth_manager() - if hasattr(auth_manager, "cleanup_dag_permissions"): - auth_manager.cleanup_dag_permissions(dag_id, session) - log.info("Successfully cleaned up DAG permissions for dag_id: %s", dag_id) - else: - log.debug( - "Auth manager %s does not support DAG permission cleanup, skipping", - type(auth_manager).__name__, - ) - except Exception as e: - # If auth manager is not available or fails, silently skip cleanup - # This ensures DAG deletion continues even if permission cleanup fails - log.warning( - "Failed to clean up DAG permissions for dag_id %s: %s. DAG deletion will continue.", dag_id, e - ) - pass diff --git a/airflow-core/tests/unit/api/common/test_delete_dag.py b/airflow-core/tests/unit/api/common/test_delete_dag.py deleted file mode 100644 index 83c2b9da1c9f4..0000000000000 --- a/airflow-core/tests/unit/api/common/test_delete_dag.py +++ /dev/null @@ -1,130 +0,0 @@ -# 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 unittest.mock import MagicMock, patch - -import pytest - -from airflow.api.common.delete_dag import _cleanup_dag_permissions, delete_dag -from airflow.exceptions import DagNotFound - -pytestmark = pytest.mark.db_test - - -class TestCleanupDagPermissions: - """Test cases for _cleanup_dag_permissions function.""" - - @patch("airflow.api_fastapi.app.get_auth_manager") - def test_cleanup_dag_permissions_auth_manager_without_cleanup_method( - self, mock_get_auth_manager, session - ): - """Test that cleanup is skipped when auth manager doesn't have cleanup_dag_permissions method.""" - mock_auth_manager = MagicMock() - # Mock auth manager without cleanup_dag_permissions method - del mock_auth_manager.cleanup_dag_permissions - mock_get_auth_manager.return_value = mock_auth_manager - - # Should return early without calling any methods - _cleanup_dag_permissions("test_dag", session) - - # Should still get the auth manager but not call cleanup - mock_get_auth_manager.assert_called_once() - - @patch("airflow.api_fastapi.app.get_auth_manager") - def test_cleanup_dag_permissions_auth_manager_with_cleanup_method(self, mock_get_auth_manager, session): - """Test that cleanup is called when auth manager has cleanup_dag_permissions method.""" - mock_auth_manager = MagicMock() - mock_auth_manager.cleanup_dag_permissions = MagicMock() - mock_get_auth_manager.return_value = mock_auth_manager - - _cleanup_dag_permissions("test_dag", session) - - # Should call the auth manager's cleanup method - mock_get_auth_manager.assert_called_once() - mock_auth_manager.cleanup_dag_permissions.assert_called_once_with("test_dag", session) - - @patch("airflow.api_fastapi.app.get_auth_manager") - def test_cleanup_dag_permissions_handles_get_auth_manager_exception(self, mock_get_auth_manager, session): - """Test that cleanup handles exceptions from get_auth_manager gracefully.""" - mock_get_auth_manager.side_effect = Exception("Auth manager not available") - - # Should not raise exception, should handle gracefully - try: - _cleanup_dag_permissions("test_dag", session) - except Exception: - pytest.fail("_cleanup_dag_permissions should handle get_auth_manager exceptions gracefully") - - -class TestDeleteDag: - """Test cases for delete_dag function integration with permission cleanup.""" - - def test_delete_dag_calls_cleanup_permissions(self, dag_maker, session): - """Test that delete_dag calls _cleanup_dag_permissions.""" - with dag_maker(dag_id="test_dag", session=session): - pass - - # Mock the cleanup function to track if it's called - with patch("airflow.api.common.delete_dag._cleanup_dag_permissions") as mock_cleanup: - # This should call the cleanup function - result = delete_dag("test_dag", session=session) - - # Should call cleanup permissions with correct arguments - mock_cleanup.assert_called_once_with("test_dag", session) - # Should return count of deleted records - assert result >= 0 - - def test_delete_dag_nonexistent_dag_does_not_call_cleanup(self, session): - """Test that delete_dag does not call cleanup for non-existent DAG.""" - with patch("airflow.api.common.delete_dag._cleanup_dag_permissions") as mock_cleanup: - with pytest.raises(DagNotFound): - delete_dag("nonexistent_dag", session=session) - - # Should not call cleanup for non-existent DAG - mock_cleanup.assert_not_called() - - @patch("airflow.api_fastapi.app.get_auth_manager") - def test_delete_dag_with_fab_auth_manager(self, mock_get_auth_manager, dag_maker, session): - """Test that delete_dag works correctly with FabAuthManager.""" - with dag_maker(dag_id="test_dag", session=session): - pass - - # Mock FabAuthManager with cleanup method - mock_auth_manager = MagicMock() - mock_auth_manager.cleanup_dag_permissions = MagicMock() - mock_get_auth_manager.return_value = mock_auth_manager - - result = delete_dag("test_dag", session=session) - - # Should call the auth manager's cleanup method - mock_auth_manager.cleanup_dag_permissions.assert_called_once_with("test_dag", session) - assert result >= 0 - - @patch("airflow.api_fastapi.app.get_auth_manager") - def test_delete_dag_with_simple_auth_manager(self, mock_get_auth_manager, dag_maker, session): - """Test that delete_dag works correctly with SimpleAuthManager (no cleanup method).""" - with dag_maker(dag_id="test_dag", session=session): - pass - - # Mock SimpleAuthManager without cleanup method - mock_auth_manager = MagicMock() - del mock_auth_manager.cleanup_dag_permissions - mock_get_auth_manager.return_value = mock_auth_manager - - # Should not raise any exception - result = delete_dag("test_dag", session=session) - assert result >= 0 diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/definition.py b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/definition.py index 7f8f1e84e2798..eab5e5eedb449 100644 --- a/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/definition.py +++ b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/definition.py @@ -108,6 +108,14 @@ ("--include-dags",), help="If passed, DAG specific permissions will also be synced.", action="store_true" ) +# permissions cleanup +ARG_DRY_RUN = Arg( + ("--dry-run",), help="Show what would be cleaned up without making any changes.", action="store_true" +) +ARG_DAG_ID_OPTIONAL = Arg( + ("-d", "--dag-id"), help="Optional: Clean up permissions for specific DAG ID only", type=str +) + ################ # # COMMANDS # # ################ @@ -253,6 +261,30 @@ args=(ARG_INCLUDE_DAGS, ARG_VERBOSE), ) +PERMISSIONS_CLEANUP_COMMAND = ActionCommand( + name="permissions-cleanup", + help="Clean up DAG permissions in Flask-AppBuilder tables", + description=( + "Clean up DAG-specific permissions. By default, cleans up orphaned permissions " + "for deleted DAGs. Use --dag-id to clean up permissions for a specific DAG." + ), + func=lazy_load_command( + "airflow.providers.fab.auth_manager.cli_commands.permissions_command.permissions_cleanup" + ), + args=(ARG_DAG_ID_OPTIONAL, ARG_DRY_RUN, ARG_YES, ARG_VERBOSE), + epilog=( + "examples:\n" + "To see what orphaned permissions would be cleaned up:\n" + " $ airflow fab-auth-manager permissions-cleanup --dry-run\n" + "To clean up all orphaned permissions:\n" + " $ airflow fab-auth-manager permissions-cleanup\n" + "To clean up permissions for specific DAG:\n" + " $ airflow fab-auth-manager permissions-cleanup --dag-id my_dag\n" + "To clean up without confirmation:\n" + " $ airflow fab-auth-manager permissions-cleanup --yes" + ), +) + DB_COMMANDS = ( ActionCommand( name="migrate", diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py new file mode 100644 index 0000000000000..e35a27cb7d22a --- /dev/null +++ b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py @@ -0,0 +1,122 @@ +# +# 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. +"""Permissions cleanup command.""" + +from __future__ import annotations + +from airflow.utils import cli as cli_utils +from airflow.utils.providers_configuration_loader import providers_configuration_loaded + + +@cli_utils.action_cli +@providers_configuration_loaded +def permissions_cleanup(args): + """Clean up DAG permissions in Flask-AppBuilder tables.""" + from airflow.models import DagModel + from airflow.providers.fab.auth_manager.cli_commands.utils import get_application_builder + from airflow.providers.fab.auth_manager.models import Resource + from airflow.security.permissions import RESOURCE_DAG_PREFIX + from airflow.utils.session import create_session + + with get_application_builder() as _: + with create_session() as session: + # Get all existing DAG IDs from DagModel + existing_dag_ids = {dag.dag_id for dag in session.query(DagModel).all()} + + # Get all DAG-related resources from FAB tables + dag_resources = ( + session.query(Resource) + .filter( + Resource.name.like(f"{RESOURCE_DAG_PREFIX}%") + | Resource.name.like("DAG Run:%") + | Resource.name.like("Task Instance:%") + ) + .all() + ) + + orphaned_resources = [] + orphaned_dag_ids = set() + + for resource in dag_resources: + # Extract DAG ID from resource name + dag_id = None + if resource.name.startswith(RESOURCE_DAG_PREFIX): + dag_id = resource.name[len(RESOURCE_DAG_PREFIX) :] + elif resource.name.startswith("DAG Run:"): + dag_id = resource.name[len("DAG Run:") :] + elif resource.name.startswith("Task Instance:"): + dag_id = resource.name[len("Task Instance:") :] + + # Check if this DAG ID still exists + if dag_id and dag_id not in existing_dag_ids: + orphaned_resources.append(resource) + orphaned_dag_ids.add(dag_id) + + # Filter by specific DAG ID if provided + if args.dag_id: + if args.dag_id in orphaned_dag_ids: + orphaned_dag_ids = {args.dag_id} + print(f"Filtering to clean up permissions for DAG: {args.dag_id}") + else: + print( + f"DAG '{args.dag_id}' not found in orphaned permissions or still exists in database." + ) + return + + if not orphaned_dag_ids: + if args.dag_id: + print(f"No orphaned permissions found for DAG: {args.dag_id}") + else: + print("No orphaned DAG permissions found.") + return + + print(f"Found orphaned permissions for {len(orphaned_dag_ids)} deleted DAG(s):") + for dag_id in sorted(orphaned_dag_ids): + print(f" - {dag_id}") + + if args.dry_run: + print("\nDry run mode: No changes will be made.") + print(f"Would clean up permissions for {len(orphaned_dag_ids)} orphaned DAG(s).") + return + + # Perform cleanup if not in dry run mode + if not args.yes: + action = ( + f"clean up permissions for {len(orphaned_dag_ids)} DAG(s)" + if not args.dag_id + else f"clean up permissions for DAG '{args.dag_id}'" + ) + confirm = input(f"\nDo you want to {action}? [y/N]: ") + if confirm.lower() not in ("y", "yes"): + print("Cleanup cancelled.") + return + + # Perform the actual cleanup + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + cleanup_count = 0 + for dag_id in orphaned_dag_ids: + try: + cleanup_dag_permissions(dag_id, session) + cleanup_count += 1 + if args.verbose: + print(f"Cleaned up permissions for DAG: {dag_id}") + except Exception as e: + print(f"Failed to clean up permissions for DAG {dag_id}: {e}") + + print(f"\nSuccessfully cleaned up permissions for {cleanup_count} DAG(s).") diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py b/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py index 3725c01ccbfaf..a5436048e3630 100644 --- a/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py +++ b/providers/fab/src/airflow/providers/fab/auth_manager/fab_auth_manager.py @@ -60,6 +60,7 @@ from airflow.models import DagModel from airflow.providers.fab.auth_manager.cli_commands.definition import ( DB_COMMANDS, + PERMISSIONS_CLEANUP_COMMAND, ROLES_COMMANDS, SYNC_PERM_COMMAND, USERS_COMMANDS, @@ -211,6 +212,7 @@ def get_cli_commands() -> list[CLICommand]: subcommands=ROLES_COMMANDS, ), SYNC_PERM_COMMAND, # not in a command group + PERMISSIONS_CLEANUP_COMMAND, # single command for permissions cleanup ] # If Airflow version is 3.0.0 or higher, add the fab-db command group if packaging.version.parse( @@ -652,17 +654,6 @@ def _get_user_permissions(user: User): return [] return getattr(user, "perms") or [] - def cleanup_dag_permissions(self, dag_id: str, session: Session) -> None: - """ - Clean up DAG-specific permissions from Flask-AppBuilder tables. - - :param dag_id: the DAG ID to clean up permissions for - :param session: database session - """ - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - cleanup_dag_permissions(dag_id=dag_id, session=session) - def _sync_appbuilder_roles(self): """ Sync appbuilder roles to DB. diff --git a/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py new file mode 100644 index 0000000000000..91ba9f5c3cfbd --- /dev/null +++ b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py @@ -0,0 +1,103 @@ +# +# 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. +"""Test permissions command.""" + +from __future__ import annotations + +from unittest.mock import MagicMock, patch + +from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( + permissions_cleanup, +) + + +class TestPermissionsCommand: + """Test permissions cleanup CLI commands.""" + + @patch("airflow.providers.fab.auth_manager.dag_permissions.cleanup_dag_permissions") + @patch("airflow.utils.session.create_session") + @patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder") + def test_permissions_cleanup_success( + self, mock_get_application_builder, mock_create_session, mock_cleanup_dag_permissions + ): + """Test successful cleanup of DAG permissions.""" + # Mock session + mock_session = MagicMock() + mock_create_session.return_value.__enter__.return_value = mock_session + + # Mock application builder + mock_appbuilder = MagicMock() + mock_get_application_builder.return_value.__enter__.return_value = mock_appbuilder + + # Mock DAG models (existing DAGs) + mock_dag_model = MagicMock() + mock_dag_model.dag_id = "existing_dag" + mock_session.query.return_value.all.return_value = [mock_dag_model] + + # Mock orphaned resources + mock_resource = MagicMock() + mock_resource.name = "DAG:deleted_dag" + mock_session.query.return_value.filter.return_value.all.return_value = [mock_resource] + + # Mock args + import argparse + + args = argparse.Namespace() + args.dag_id = None + args.dry_run = False + args.yes = True + args.verbose = True + + # Execute command + permissions_cleanup(args) + + # Verify function calls + mock_cleanup_dag_permissions.assert_called() + + @patch("airflow.utils.session.create_session") + @patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder") + def test_permissions_cleanup_dry_run(self, mock_get_application_builder, mock_create_session): + """Test dry run mode for permissions cleanup.""" + # Mock session and data + mock_session = MagicMock() + mock_create_session.return_value.__enter__.return_value = mock_session + + # Mock DAG models (existing DAGs) + mock_dag_model = MagicMock() + mock_dag_model.dag_id = "existing_dag" + mock_session.query.return_value.all.return_value = [mock_dag_model] + + # Mock orphaned resources + mock_resource = MagicMock() + mock_resource.name = "DAG:deleted_dag" + mock_session.query.return_value.filter.return_value.all.return_value = [mock_resource] + + # Mock args + import argparse + + args = argparse.Namespace() + args.dag_id = None + args.dry_run = True + args.verbose = True + + # Mock application builder + mock_appbuilder = MagicMock() + mock_get_application_builder.return_value.__enter__.return_value = mock_appbuilder + + # Execute command + permissions_cleanup(args) From a72a8bfb44e51702b9be0f9d66a1d79c415344ba Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Mon, 1 Sep 2025 07:14:06 +0800 Subject: [PATCH 5/8] feat(fab): Update DAG permissions cleanup to SQLAlchemy 2.0 and improve tests - Replace deprecated session.query() with session.scalars(select()) and session.execute(delete()) - Update permissions_command.py to use SQLAlchemy 2.0 syntax for better future compatibility - Rewrite test_permissions_command.py to use integration tests instead of fragile mocks --- .../cli_commands/permissions_command.py | 12 +- .../fab/auth_manager/dag_permissions.py | 34 ++-- .../cli_commands/test_permissions_command.py | 140 +++++++++------- .../fab/auth_manager/test_dag_permissions.py | 154 ------------------ 4 files changed, 111 insertions(+), 229 deletions(-) delete mode 100644 providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py index e35a27cb7d22a..ff2f3beca5e05 100644 --- a/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py +++ b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py @@ -27,6 +27,8 @@ @providers_configuration_loaded def permissions_cleanup(args): """Clean up DAG permissions in Flask-AppBuilder tables.""" + from sqlalchemy import select + from airflow.models import DagModel from airflow.providers.fab.auth_manager.cli_commands.utils import get_application_builder from airflow.providers.fab.auth_manager.models import Resource @@ -36,18 +38,16 @@ def permissions_cleanup(args): with get_application_builder() as _: with create_session() as session: # Get all existing DAG IDs from DagModel - existing_dag_ids = {dag.dag_id for dag in session.query(DagModel).all()} + existing_dag_ids = {dag.dag_id for dag in session.scalars(select(DagModel)).all()} # Get all DAG-related resources from FAB tables - dag_resources = ( - session.query(Resource) - .filter( + dag_resources = session.scalars( + select(Resource).filter( Resource.name.like(f"{RESOURCE_DAG_PREFIX}%") | Resource.name.like("DAG Run:%") | Resource.name.like("Task Instance:%") ) - .all() - ) + ).all() orphaned_resources = [] orphaned_dag_ids = set() diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py b/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py index 38e45c6837e3a..a313badd4a177 100644 --- a/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py +++ b/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py @@ -54,13 +54,14 @@ def cleanup_dag_permissions(dag_id: str, session: Session | None = None) -> None def _cleanup_dag_permissions_impl(dag_id: str, session: Session) -> None: """Implement DAG permissions cleanup.""" + from sqlalchemy import select + from airflow.providers.fab.auth_manager.models import Permission, Resource, assoc_permission_role # Clean up specific DAG permissions dag_resource_name = f"{RESOURCE_DAG_PREFIX}{dag_id}" - dag_resources = ( - session.query(Resource) - .filter( + dag_resources = session.scalars( + select(Resource).filter( Resource.name.in_( [ dag_resource_name, # DAG:dag_id @@ -69,8 +70,7 @@ def _cleanup_dag_permissions_impl(dag_id: str, session: Session) -> None: ] ) ) - .all() - ) + ).all() log.info("Cleaning up DAG-specific permissions for dag_id: %s", dag_id) if not dag_resources: @@ -79,26 +79,32 @@ def _cleanup_dag_permissions_impl(dag_id: str, session: Session) -> None: dag_resource_ids = [resource.id for resource in dag_resources] # Find all permissions associated with these resources - dag_permissions = session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).all() + dag_permissions = session.scalars( + select(Permission).filter(Permission.resource_id.in_(dag_resource_ids)) + ).all() if not dag_permissions: # Delete resources even if no permissions exist - session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) + from sqlalchemy import delete + + session.execute(delete(Resource).where(Resource.id.in_(dag_resource_ids))) return dag_permission_ids = [permission.id for permission in dag_permissions] # Delete permission-role associations first (foreign key constraint) - session.query(assoc_permission_role).filter( - assoc_permission_role.c.permission_view_id.in_(dag_permission_ids) - ).delete(synchronize_session=False) + from sqlalchemy import delete - # Delete permissions - session.query(Permission).filter(Permission.resource_id.in_(dag_resource_ids)).delete( - synchronize_session=False + session.execute( + delete(assoc_permission_role).where( + assoc_permission_role.c.permission_view_id.in_(dag_permission_ids) + ) ) + # Delete permissions + session.execute(delete(Permission).where(Permission.resource_id.in_(dag_resource_ids))) + # Delete resources (ab_view_menu entries) - session.query(Resource).filter(Resource.id.in_(dag_resource_ids)).delete(synchronize_session=False) + session.execute(delete(Resource).where(Resource.id.in_(dag_resource_ids))) log.info("Cleaned up %d DAG-specific permissions", len(dag_permissions)) diff --git a/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py index 91ba9f5c3cfbd..d5536e029d15c 100644 --- a/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py +++ b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py @@ -19,85 +19,115 @@ from __future__ import annotations -from unittest.mock import MagicMock, patch +import argparse +from contextlib import redirect_stdout +from importlib import reload +from io import StringIO +from unittest.mock import patch -from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( - permissions_cleanup, -) +import pytest +from airflow.cli import cli_parser -class TestPermissionsCommand: - """Test permissions cleanup CLI commands.""" +from tests_common.test_utils.compat import ignore_provider_compatibility_error +from tests_common.test_utils.config import conf_vars - @patch("airflow.providers.fab.auth_manager.dag_permissions.cleanup_dag_permissions") - @patch("airflow.utils.session.create_session") - @patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder") - def test_permissions_cleanup_success( - self, mock_get_application_builder, mock_create_session, mock_cleanup_dag_permissions - ): - """Test successful cleanup of DAG permissions.""" - # Mock session - mock_session = MagicMock() - mock_create_session.return_value.__enter__.return_value = mock_session +with ignore_provider_compatibility_error("2.9.0+", __file__): + from airflow.providers.fab.auth_manager.cli_commands import permissions_command + from airflow.providers.fab.auth_manager.cli_commands.utils import get_application_builder - # Mock application builder - mock_appbuilder = MagicMock() - mock_get_application_builder.return_value.__enter__.return_value = mock_appbuilder +pytestmark = pytest.mark.db_test - # Mock DAG models (existing DAGs) - mock_dag_model = MagicMock() - mock_dag_model.dag_id = "existing_dag" - mock_session.query.return_value.all.return_value = [mock_dag_model] - # Mock orphaned resources - mock_resource = MagicMock() - mock_resource.name = "DAG:deleted_dag" - mock_session.query.return_value.filter.return_value.all.return_value = [mock_resource] +class TestPermissionsCommand: + """Test permissions cleanup CLI commands.""" - # Mock args - import argparse + @pytest.fixture(autouse=True) + def _set_attrs(self): + with conf_vars( + { + ( + "core", + "auth_manager", + ): "airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager", + } + ): + # Reload the module to use FAB auth manager + reload(cli_parser) + # Clearing the cache before calling it + cli_parser.get_parser.cache_clear() + self.parser = cli_parser.get_parser() + with get_application_builder() as appbuilder: + self.appbuilder = appbuilder + yield + @patch("airflow.providers.fab.auth_manager.dag_permissions.cleanup_dag_permissions") + def test_permissions_cleanup_success(self, mock_cleanup_dag_permissions): + """Test successful cleanup of DAG permissions.""" + # Mock args args = argparse.Namespace() args.dag_id = None args.dry_run = False args.yes = True args.verbose = True - # Execute command - permissions_cleanup(args) + with redirect_stdout(StringIO()) as stdout: + permissions_command.permissions_cleanup(args) - # Verify function calls - mock_cleanup_dag_permissions.assert_called() + # Verify function calls - it should be called for real DAGs + output = stdout.getvalue() + # Should either call cleanup or report no orphaned permissions found + assert ( + "Successfully cleaned up permissions" in output or "No orphaned DAG permissions found" in output + ) - @patch("airflow.utils.session.create_session") - @patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder") - def test_permissions_cleanup_dry_run(self, mock_get_application_builder, mock_create_session): + def test_permissions_cleanup_dry_run(self): """Test dry run mode for permissions cleanup.""" - # Mock session and data - mock_session = MagicMock() - mock_create_session.return_value.__enter__.return_value = mock_session + # Mock args + args = argparse.Namespace() + args.dag_id = None + args.dry_run = True + args.verbose = True - # Mock DAG models (existing DAGs) - mock_dag_model = MagicMock() - mock_dag_model.dag_id = "existing_dag" - mock_session.query.return_value.all.return_value = [mock_dag_model] + with redirect_stdout(StringIO()) as stdout: + permissions_command.permissions_cleanup(args) - # Mock orphaned resources - mock_resource = MagicMock() - mock_resource.name = "DAG:deleted_dag" - mock_session.query.return_value.filter.return_value.all.return_value = [mock_resource] + output = stdout.getvalue() + assert "Dry run mode" in output or "No orphaned DAG permissions found" in output + def test_permissions_cleanup_specific_dag(self): + """Test cleanup for a specific DAG.""" # Mock args - import argparse - args = argparse.Namespace() - args.dag_id = None + args.dag_id = "test_dag" args.dry_run = True args.verbose = True - # Mock application builder - mock_appbuilder = MagicMock() - mock_get_application_builder.return_value.__enter__.return_value = mock_appbuilder + with redirect_stdout(StringIO()) as stdout: + permissions_command.permissions_cleanup(args) + + output = stdout.getvalue() + # Check for appropriate output indicating DAG-specific operation + assert ( + "test_dag" in output + or "No orphaned permissions found for DAG" in output + or "not found in orphaned permissions" in output + or "No orphaned DAG permissions found" in output + ) + + @patch("builtins.input", return_value="n") + def test_permissions_cleanup_no_confirmation(self, mock_input): + """Test cleanup cancellation when user doesn't confirm.""" + # Mock args + args = argparse.Namespace() + args.dag_id = None + args.dry_run = False + args.yes = False + args.verbose = False + + with redirect_stdout(StringIO()) as stdout: + permissions_command.permissions_cleanup(args) - # Execute command - permissions_cleanup(args) + output = stdout.getvalue() + # Should not call cleanup if user declines or no orphaned permissions found + assert "Cleanup cancelled" in output or "No orphaned DAG permissions found" in output diff --git a/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py b/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py deleted file mode 100644 index 4398625568e19..0000000000000 --- a/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py +++ /dev/null @@ -1,154 +0,0 @@ -# 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 unittest.mock import MagicMock, patch - - -class TestDagPermissions: - """Test cases for dag_permissions module.""" - - def test_cleanup_dag_permissions_with_session(self): - """Test cleanup_dag_permissions with provided session.""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models - with patch("airflow.providers.fab.auth_manager.models") as mock_models: - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session - mock_session = MagicMock() - - # Mock no resources found - mock_session.query.return_value.filter.return_value.all.return_value = [] - - cleanup_dag_permissions("test_dag", mock_session) - - # Should call session.query for Resource - mock_session.query.assert_called() - - def test_cleanup_dag_permissions_without_session(self): - """Test cleanup_dag_permissions without provided session (creates new session).""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models - with patch("airflow.providers.fab.auth_manager.models") as mock_models: - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock create_session - with patch("airflow.utils.session.create_session") as mock_create_session: - mock_session = MagicMock() - mock_create_session.return_value.__enter__.return_value = mock_session - - # Mock no resources found - mock_session.query.return_value.filter.return_value.all.return_value = [] - - cleanup_dag_permissions("test_dag") - - # Should create a new session - mock_create_session.assert_called_once() - - def test_cleanup_dag_permissions_with_resources_and_permissions(self): - """Test cleanup_dag_permissions with actual resources and permissions.""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models - with patch("airflow.providers.fab.auth_manager.models") as mock_models: - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session - mock_session = MagicMock() - - # Mock resources exist - mock_dag_resource = MagicMock() - mock_dag_resource.id = 1 - mock_session.query.return_value.filter.return_value.all.return_value = [mock_dag_resource] - - # Mock permissions exist - mock_dag_permission = MagicMock() - mock_dag_permission.id = 1 - - def mock_query_side_effect(model): - mock_query = MagicMock() - if model == mock_resource: - mock_query.filter.return_value.all.return_value = [mock_dag_resource] - mock_query.filter.return_value.delete.return_value = None - elif model == mock_permission: - mock_query.filter.return_value.all.return_value = [mock_dag_permission] - mock_query.filter.return_value.delete.return_value = None - elif model == mock_assoc: - mock_query.filter.return_value.delete.return_value = None - return mock_query - - mock_session.query.side_effect = mock_query_side_effect - - cleanup_dag_permissions("test_dag", mock_session) - - # Should perform deletion operations - assert mock_session.query.call_count >= 3 # Resource, Permission, assoc queries - - def test_cleanup_dag_permissions_with_resources_but_no_permissions(self): - """Test cleanup_dag_permissions with resources but no permissions.""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models - with patch("airflow.providers.fab.auth_manager.models") as mock_models: - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session - mock_session = MagicMock() - - # Mock resources exist but no permissions - mock_dag_resource = MagicMock() - mock_dag_resource.id = 1 - - def mock_query_side_effect(model): - mock_query = MagicMock() - if model == mock_resource: - mock_query.filter.return_value.all.return_value = [mock_dag_resource] - mock_query.filter.return_value.delete.return_value = None - elif model == mock_permission: - mock_query.filter.return_value.all.return_value = [] # No permissions - mock_query.filter.return_value.delete.return_value = None - return mock_query - - mock_session.query.side_effect = mock_query_side_effect - - cleanup_dag_permissions("test_dag", mock_session) - - # Should still delete resources even without permissions - assert mock_session.query.call_count >= 2 # Resource queries From 7f083b5865ccd8e792d9b30ac2022ec44388ffc0 Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Mon, 1 Sep 2025 20:39:02 +0800 Subject: [PATCH 6/8] feat(fab): add missing test_dag_permissions.py Add unit tests for dag_permissions module to satisfy project structure requirements --- .../fab/auth_manager/test_dag_permissions.py | 160 ++++++++++++++++++ 1 file changed, 160 insertions(+) create mode 100644 providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py diff --git a/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py b/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py new file mode 100644 index 0000000000000..31584f3024cda --- /dev/null +++ b/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py @@ -0,0 +1,160 @@ +# 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 unittest.mock import MagicMock, patch + + +class TestDagPermissions: + """Test cases for dag_permissions module.""" + + def test_cleanup_dag_permissions_with_session(self): + """Test cleanup_dag_permissions with provided session.""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session with SQLAlchemy 2.0 methods + mock_session = MagicMock() + mock_session.scalars.return_value.all.return_value = [] + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag", mock_session) + + # Should call session.scalars and session.execute for SQLAlchemy 2.0 + assert mock_session.scalars.called or mock_session.execute.called + + def test_cleanup_dag_permissions_without_session(self): + """Test cleanup_dag_permissions without provided session (creates new session).""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock create_session + with patch("airflow.utils.session.create_session") as mock_create_session: + mock_session = MagicMock() + mock_create_session.return_value.__enter__.return_value = mock_session + mock_session.scalars.return_value.all.return_value = [] + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag") + + # Should create a new session + mock_create_session.assert_called_once() + + def test_cleanup_dag_permissions_with_resources_and_permissions(self): + """Test cleanup_dag_permissions with actual resources and permissions.""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session + mock_session = MagicMock() + + # Mock resources exist + mock_dag_resource = MagicMock() + mock_dag_resource.id = 1 + + # Mock permissions exist + mock_dag_permission = MagicMock() + mock_dag_permission.id = 1 + + # Setup mock returns for different select queries + def mock_scalars_side_effect(stmt): + mock_result = MagicMock() + # Return resources for resource queries + mock_result.all.return_value = [mock_dag_resource] + return mock_result + + mock_session.scalars.side_effect = mock_scalars_side_effect + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag", mock_session) + + # Should perform database operations + assert mock_session.scalars.called or mock_session.execute.called + + def test_cleanup_dag_permissions_with_resources_but_no_permissions(self): + """Test cleanup_dag_permissions with resources but no permissions.""" + from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session + mock_session = MagicMock() + + # Mock resources exist but no permissions + mock_dag_resource = MagicMock() + mock_dag_resource.id = 1 + + def mock_scalars_side_effect(stmt): + mock_result = MagicMock() + # Return resources for resource queries + mock_result.all.return_value = [mock_dag_resource] + return mock_result + + mock_session.scalars.side_effect = mock_scalars_side_effect + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag", mock_session) + + # Should still perform database operations + assert mock_session.scalars.called or mock_session.execute.called From a0f06d38a1dc99059da4ebe1c088639c47b08f7b Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Thu, 4 Sep 2025 07:23:34 +0800 Subject: [PATCH 7/8] refactor(fab): Consolidate DAG permissions code and modernize patterns - Replace hardcoded resource strings with constants in permissions_command.py - Use RESOURCE_DAG_PREFIX and RESOURCE_DETAILS_MAP for DAG Run - Remove deprecated Task Instance resource handling - Change to session handling with @provide_session decorator - Replace manual boolean parsing with airflow.utils.strings.to_boolean - Consolidate dag_permissions.py functionality into permissions_command.py - Consolidate test files for better maintainability --- .../cli_commands/permissions_command.py | 94 ++++++++-- .../fab/auth_manager/dag_permissions.py | 110 ------------ .../cli_commands/test_permissions_command.py | 153 ++++++++++++++++- .../fab/auth_manager/test_dag_permissions.py | 160 ------------------ 4 files changed, 235 insertions(+), 282 deletions(-) delete mode 100644 providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py delete mode 100644 providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py index ff2f3beca5e05..a535050750c64 100644 --- a/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py +++ b/providers/fab/src/airflow/providers/fab/auth_manager/cli_commands/permissions_command.py @@ -19,8 +19,83 @@ from __future__ import annotations +import logging +from typing import TYPE_CHECKING + from airflow.utils import cli as cli_utils from airflow.utils.providers_configuration_loader import providers_configuration_loaded +from airflow.utils.session import NEW_SESSION, provide_session +from airflow.utils.strings import to_boolean + +if TYPE_CHECKING: + from sqlalchemy.orm import Session + +log = logging.getLogger(__name__) + + +@provide_session +def cleanup_dag_permissions(dag_id: str, session: Session = NEW_SESSION) -> None: + """ + Clean up DAG-specific permissions from Flask-AppBuilder tables. + + When a DAG is deleted, we need to clean up the corresponding permissions + to prevent orphaned entries in the ab_view_menu table. + + This addresses issue #50905: Deleted DAGs not removed from ab_view_menu table + and show up in permissions. + + :param dag_id: Specific DAG ID to clean up. + :param session: Database session. + """ + from sqlalchemy import delete, select + + from airflow.providers.fab.auth_manager.models import Permission, Resource, assoc_permission_role + from airflow.security.permissions import RESOURCE_DAG_PREFIX, RESOURCE_DAG_RUN, RESOURCE_DETAILS_MAP + + # Clean up specific DAG permissions + dag_resources = session.scalars( + select(Resource).filter( + Resource.name.in_( + [ + f"{RESOURCE_DAG_PREFIX}{dag_id}", # DAG:dag_id + f"{RESOURCE_DETAILS_MAP[RESOURCE_DAG_RUN]['prefix']}{dag_id}", # DAG_RUN:dag_id + ] + ) + ) + ).all() + log.info("Cleaning up DAG-specific permissions for dag_id: %s", dag_id) + + if not dag_resources: + return + + dag_resource_ids = [resource.id for resource in dag_resources] + + # Find all permissions associated with these resources + dag_permissions = session.scalars( + select(Permission).filter(Permission.resource_id.in_(dag_resource_ids)) + ).all() + + if not dag_permissions: + # Delete resources even if no permissions exist + session.execute(delete(Resource).where(Resource.id.in_(dag_resource_ids))) + return + + dag_permission_ids = [permission.id for permission in dag_permissions] + + # Delete permission-role associations first (foreign key constraint) + session.execute( + delete(assoc_permission_role).where( + assoc_permission_role.c.permission_view_id.in_(dag_permission_ids) + ) + ) + + # Delete permissions + session.execute(delete(Permission).where(Permission.resource_id.in_(dag_resource_ids))) + + # Delete resources (ab_view_menu entries) + session.execute(delete(Resource).where(Resource.id.in_(dag_resource_ids))) + + log.info("Cleaned up %d DAG-specific permissions", len(dag_permissions)) @cli_utils.action_cli @@ -32,7 +107,11 @@ def permissions_cleanup(args): from airflow.models import DagModel from airflow.providers.fab.auth_manager.cli_commands.utils import get_application_builder from airflow.providers.fab.auth_manager.models import Resource - from airflow.security.permissions import RESOURCE_DAG_PREFIX + from airflow.security.permissions import ( + RESOURCE_DAG_PREFIX, + RESOURCE_DAG_RUN, + RESOURCE_DETAILS_MAP, + ) from airflow.utils.session import create_session with get_application_builder() as _: @@ -44,8 +123,7 @@ def permissions_cleanup(args): dag_resources = session.scalars( select(Resource).filter( Resource.name.like(f"{RESOURCE_DAG_PREFIX}%") - | Resource.name.like("DAG Run:%") - | Resource.name.like("Task Instance:%") + | Resource.name.like(f"{RESOURCE_DETAILS_MAP[RESOURCE_DAG_RUN]['prefix']}%") ) ).all() @@ -57,10 +135,8 @@ def permissions_cleanup(args): dag_id = None if resource.name.startswith(RESOURCE_DAG_PREFIX): dag_id = resource.name[len(RESOURCE_DAG_PREFIX) :] - elif resource.name.startswith("DAG Run:"): - dag_id = resource.name[len("DAG Run:") :] - elif resource.name.startswith("Task Instance:"): - dag_id = resource.name[len("Task Instance:") :] + elif resource.name.startswith(RESOURCE_DETAILS_MAP[RESOURCE_DAG_RUN]["prefix"]): + dag_id = resource.name[len(RESOURCE_DETAILS_MAP[RESOURCE_DAG_RUN]["prefix"]) :] # Check if this DAG ID still exists if dag_id and dag_id not in existing_dag_ids: @@ -102,13 +178,11 @@ def permissions_cleanup(args): else f"clean up permissions for DAG '{args.dag_id}'" ) confirm = input(f"\nDo you want to {action}? [y/N]: ") - if confirm.lower() not in ("y", "yes"): + if not to_boolean(confirm): print("Cleanup cancelled.") return # Perform the actual cleanup - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - cleanup_count = 0 for dag_id in orphaned_dag_ids: try: diff --git a/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py b/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py deleted file mode 100644 index a313badd4a177..0000000000000 --- a/providers/fab/src/airflow/providers/fab/auth_manager/dag_permissions.py +++ /dev/null @@ -1,110 +0,0 @@ -# -# 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. -"""DAG permissions management for FAB Auth Manager.""" - -from __future__ import annotations - -import logging -from typing import TYPE_CHECKING - -from airflow.security.permissions import RESOURCE_DAG_PREFIX - -if TYPE_CHECKING: - from sqlalchemy.orm import Session - -log = logging.getLogger(__name__) - - -def cleanup_dag_permissions(dag_id: str, session: Session | None = None) -> None: - """ - Clean up DAG-specific permissions from Flask-AppBuilder tables. - - When a DAG is deleted, we need to clean up the corresponding permissions - to prevent orphaned entries in the ab_view_menu table. - - This addresses issue #50905: Deleted DAGs not removed from ab_view_menu table - and show up in permissions. - - :param dag_id: Specific DAG ID to clean up. - :param session: Database session. If None, creates a new session. - """ - from airflow.utils.session import create_session - - if session is None: - with create_session() as session: - _cleanup_dag_permissions_impl(dag_id, session) - else: - _cleanup_dag_permissions_impl(dag_id, session) - - -def _cleanup_dag_permissions_impl(dag_id: str, session: Session) -> None: - """Implement DAG permissions cleanup.""" - from sqlalchemy import select - - from airflow.providers.fab.auth_manager.models import Permission, Resource, assoc_permission_role - - # Clean up specific DAG permissions - dag_resource_name = f"{RESOURCE_DAG_PREFIX}{dag_id}" - dag_resources = session.scalars( - select(Resource).filter( - Resource.name.in_( - [ - dag_resource_name, # DAG:dag_id - f"DAG Run:{dag_id}", # DAG_RUN:dag_id - f"Task Instance:{dag_id}", # TASK_INSTANCE:dag_id (if exists) - ] - ) - ) - ).all() - log.info("Cleaning up DAG-specific permissions for dag_id: %s", dag_id) - - if not dag_resources: - return - - dag_resource_ids = [resource.id for resource in dag_resources] - - # Find all permissions associated with these resources - dag_permissions = session.scalars( - select(Permission).filter(Permission.resource_id.in_(dag_resource_ids)) - ).all() - - if not dag_permissions: - # Delete resources even if no permissions exist - from sqlalchemy import delete - - session.execute(delete(Resource).where(Resource.id.in_(dag_resource_ids))) - return - - dag_permission_ids = [permission.id for permission in dag_permissions] - - # Delete permission-role associations first (foreign key constraint) - from sqlalchemy import delete - - session.execute( - delete(assoc_permission_role).where( - assoc_permission_role.c.permission_view_id.in_(dag_permission_ids) - ) - ) - - # Delete permissions - session.execute(delete(Permission).where(Permission.resource_id.in_(dag_resource_ids))) - - # Delete resources (ab_view_menu entries) - session.execute(delete(Resource).where(Resource.id.in_(dag_resource_ids))) - - log.info("Cleaned up %d DAG-specific permissions", len(dag_permissions)) diff --git a/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py index d5536e029d15c..963ff6decb7e2 100644 --- a/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py +++ b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py @@ -23,7 +23,7 @@ from contextlib import redirect_stdout from importlib import reload from io import StringIO -from unittest.mock import patch +from unittest.mock import MagicMock, patch import pytest @@ -61,7 +61,7 @@ def _set_attrs(self): self.appbuilder = appbuilder yield - @patch("airflow.providers.fab.auth_manager.dag_permissions.cleanup_dag_permissions") + @patch("airflow.providers.fab.auth_manager.cli_commands.permissions_command.cleanup_dag_permissions") def test_permissions_cleanup_success(self, mock_cleanup_dag_permissions): """Test successful cleanup of DAG permissions.""" # Mock args @@ -131,3 +131,152 @@ def test_permissions_cleanup_no_confirmation(self, mock_input): output = stdout.getvalue() # Should not call cleanup if user declines or no orphaned permissions found assert "Cleanup cancelled" in output or "No orphaned DAG permissions found" in output + + +class TestDagPermissions: + """Test cases for cleanup_dag_permissions function.""" + + def test_cleanup_dag_permissions_with_session(self): + """Test cleanup_dag_permissions with provided session.""" + from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( + cleanup_dag_permissions, + ) + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session with SQLAlchemy 2.0 methods + mock_session = MagicMock() + mock_session.scalars.return_value.all.return_value = [] + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag", mock_session) + + # Should call session.scalars and session.execute for SQLAlchemy 2.0 + assert mock_session.scalars.called or mock_session.execute.called + + def test_cleanup_dag_permissions_without_session(self): + """Test cleanup_dag_permissions without provided session (creates new session).""" + from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( + cleanup_dag_permissions, + ) + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock create_session + with patch("airflow.utils.session.create_session") as mock_create_session: + mock_session = MagicMock() + mock_create_session.return_value.__enter__.return_value = mock_session + mock_session.scalars.return_value.all.return_value = [] + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag") + + # Should create a new session + mock_create_session.assert_called_once() + + def test_cleanup_dag_permissions_with_resources_and_permissions(self): + """Test cleanup_dag_permissions with actual resources and permissions.""" + from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( + cleanup_dag_permissions, + ) + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session + mock_session = MagicMock() + + # Mock resources exist + mock_dag_resource = MagicMock() + mock_dag_resource.id = 1 + + # Mock permissions exist + mock_dag_permission = MagicMock() + mock_dag_permission.id = 1 + + # Setup mock returns for different select queries + def mock_scalars_side_effect(stmt): + mock_result = MagicMock() + # Return resources for resource queries + mock_result.all.return_value = [mock_dag_resource] + return mock_result + + mock_session.scalars.side_effect = mock_scalars_side_effect + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag", mock_session) + + # Should perform database operations + assert mock_session.scalars.called or mock_session.execute.called + + def test_cleanup_dag_permissions_with_resources_but_no_permissions(self): + """Test cleanup_dag_permissions with resources but no permissions.""" + from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( + cleanup_dag_permissions, + ) + + # Mock FAB models and select/delete + with ( + patch("airflow.providers.fab.auth_manager.models") as mock_models, + patch("sqlalchemy.select"), + patch("sqlalchemy.delete"), + ): + mock_resource = MagicMock() + mock_permission = MagicMock() + mock_assoc = MagicMock() + mock_models.Resource = mock_resource + mock_models.Permission = mock_permission + mock_models.assoc_permission_role = mock_assoc + + # Mock session + mock_session = MagicMock() + + # Mock resources exist but no permissions + mock_dag_resource = MagicMock() + mock_dag_resource.id = 1 + + def mock_scalars_side_effect(stmt): + mock_result = MagicMock() + # Return resources for resource queries + mock_result.all.return_value = [mock_dag_resource] + return mock_result + + mock_session.scalars.side_effect = mock_scalars_side_effect + mock_session.execute.return_value = None + + cleanup_dag_permissions("test_dag", mock_session) + + # Should still perform database operations + assert mock_session.scalars.called or mock_session.execute.called diff --git a/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py b/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py deleted file mode 100644 index 31584f3024cda..0000000000000 --- a/providers/fab/tests/unit/fab/auth_manager/test_dag_permissions.py +++ /dev/null @@ -1,160 +0,0 @@ -# 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 unittest.mock import MagicMock, patch - - -class TestDagPermissions: - """Test cases for dag_permissions module.""" - - def test_cleanup_dag_permissions_with_session(self): - """Test cleanup_dag_permissions with provided session.""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), - ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session with SQLAlchemy 2.0 methods - mock_session = MagicMock() - mock_session.scalars.return_value.all.return_value = [] - mock_session.execute.return_value = None - - cleanup_dag_permissions("test_dag", mock_session) - - # Should call session.scalars and session.execute for SQLAlchemy 2.0 - assert mock_session.scalars.called or mock_session.execute.called - - def test_cleanup_dag_permissions_without_session(self): - """Test cleanup_dag_permissions without provided session (creates new session).""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), - ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock create_session - with patch("airflow.utils.session.create_session") as mock_create_session: - mock_session = MagicMock() - mock_create_session.return_value.__enter__.return_value = mock_session - mock_session.scalars.return_value.all.return_value = [] - mock_session.execute.return_value = None - - cleanup_dag_permissions("test_dag") - - # Should create a new session - mock_create_session.assert_called_once() - - def test_cleanup_dag_permissions_with_resources_and_permissions(self): - """Test cleanup_dag_permissions with actual resources and permissions.""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), - ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session - mock_session = MagicMock() - - # Mock resources exist - mock_dag_resource = MagicMock() - mock_dag_resource.id = 1 - - # Mock permissions exist - mock_dag_permission = MagicMock() - mock_dag_permission.id = 1 - - # Setup mock returns for different select queries - def mock_scalars_side_effect(stmt): - mock_result = MagicMock() - # Return resources for resource queries - mock_result.all.return_value = [mock_dag_resource] - return mock_result - - mock_session.scalars.side_effect = mock_scalars_side_effect - mock_session.execute.return_value = None - - cleanup_dag_permissions("test_dag", mock_session) - - # Should perform database operations - assert mock_session.scalars.called or mock_session.execute.called - - def test_cleanup_dag_permissions_with_resources_but_no_permissions(self): - """Test cleanup_dag_permissions with resources but no permissions.""" - from airflow.providers.fab.auth_manager.dag_permissions import cleanup_dag_permissions - - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), - ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session - mock_session = MagicMock() - - # Mock resources exist but no permissions - mock_dag_resource = MagicMock() - mock_dag_resource.id = 1 - - def mock_scalars_side_effect(stmt): - mock_result = MagicMock() - # Return resources for resource queries - mock_result.all.return_value = [mock_dag_resource] - return mock_result - - mock_session.scalars.side_effect = mock_scalars_side_effect - mock_session.execute.return_value = None - - cleanup_dag_permissions("test_dag", mock_session) - - # Should still perform database operations - assert mock_session.scalars.called or mock_session.execute.called From 44bd7234db836bd76929e9c4f158870118be1bbe Mon Sep 17 00:00:00 2001 From: HsiuChuanHsu Date: Fri, 5 Sep 2025 07:45:28 +0800 Subject: [PATCH 8/8] test: enhance FAB permissions tests with function verification and real database operations 1. Function call verification (TestPermissionsCommand) - Replace stdout-only testing with precise function call verification - Add mock assertions for cleanup_dag_permissions with exact parameters 2. Real database operations (TestDagPermissions): - Replace mock-heavy approach with actual database interactions - Use real Resource, Action, and Permission entities with constraint handling --- .../cli_commands/test_permissions_command.py | 369 +++++++++++------- 1 file changed, 220 insertions(+), 149 deletions(-) diff --git a/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py index 963ff6decb7e2..cf4a0355647a6 100644 --- a/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py +++ b/providers/fab/tests/unit/fab/auth_manager/cli_commands/test_permissions_command.py @@ -62,7 +62,8 @@ def _set_attrs(self): yield @patch("airflow.providers.fab.auth_manager.cli_commands.permissions_command.cleanup_dag_permissions") - def test_permissions_cleanup_success(self, mock_cleanup_dag_permissions): + @patch("airflow.providers.fab.auth_manager.models.Resource") + def test_permissions_cleanup_success(self, mock_resource, mock_cleanup_dag_permissions): """Test successful cleanup of DAG permissions.""" # Mock args args = argparse.Namespace() @@ -71,17 +72,38 @@ def test_permissions_cleanup_success(self, mock_cleanup_dag_permissions): args.yes = True args.verbose = True - with redirect_stdout(StringIO()) as stdout: + # Mock orphaned resources + mock_orphaned_resource = MagicMock() + mock_orphaned_resource.name = "DAG:orphaned_dag" + + with ( + patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder"), + patch("airflow.utils.session.create_session") as mock_session_ctx, + patch("sqlalchemy.select"), + redirect_stdout(StringIO()), + ): + mock_session = MagicMock() + mock_session_ctx.return_value.__enter__.return_value = mock_session + + # Mock DagModel query - return existing DAGs + mock_dag_result = MagicMock() + mock_dag_result.all.return_value = [MagicMock(dag_id="existing_dag")] + + # Mock Resource query - return orphaned resources + mock_resource_result = MagicMock() + mock_resource_result.all.return_value = [mock_orphaned_resource] + + # Setup session.scalars to return different results for different queries + mock_session.scalars.side_effect = [mock_dag_result, mock_resource_result] + permissions_command.permissions_cleanup(args) - # Verify function calls - it should be called for real DAGs - output = stdout.getvalue() - # Should either call cleanup or report no orphaned permissions found - assert ( - "Successfully cleaned up permissions" in output or "No orphaned DAG permissions found" in output - ) + # Verify function calls - it should be called exactly once for the orphaned DAG + mock_cleanup_dag_permissions.assert_called_once_with("orphaned_dag", mock_session) - def test_permissions_cleanup_dry_run(self): + @patch("airflow.providers.fab.auth_manager.cli_commands.permissions_command.cleanup_dag_permissions") + @patch("airflow.providers.fab.auth_manager.models.Resource") + def test_permissions_cleanup_dry_run(self, mock_resource, mock_cleanup_dag_permissions): """Test dry run mode for permissions cleanup.""" # Mock args args = argparse.Namespace() @@ -89,34 +111,86 @@ def test_permissions_cleanup_dry_run(self): args.dry_run = True args.verbose = True - with redirect_stdout(StringIO()) as stdout: + # Mock orphaned resources + mock_orphaned_resource = MagicMock() + mock_orphaned_resource.name = "DAG:orphaned_dag" + + with ( + patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder"), + patch("airflow.utils.session.create_session") as mock_session_ctx, + patch("sqlalchemy.select"), + redirect_stdout(StringIO()) as stdout, + ): + mock_session = MagicMock() + mock_session_ctx.return_value.__enter__.return_value = mock_session + + # Mock DagModel query - return existing DAGs + mock_dag_result = MagicMock() + mock_dag_result.all.return_value = [MagicMock(dag_id="existing_dag")] + + # Mock Resource query - return orphaned resources + mock_resource_result = MagicMock() + mock_resource_result.all.return_value = [mock_orphaned_resource] + + # Setup session.scalars to return different results for different queries + mock_session.scalars.side_effect = [mock_dag_result, mock_resource_result] + permissions_command.permissions_cleanup(args) output = stdout.getvalue() assert "Dry run mode" in output or "No orphaned DAG permissions found" in output + # In dry run mode, cleanup_dag_permissions should NOT be called + mock_cleanup_dag_permissions.assert_not_called() - def test_permissions_cleanup_specific_dag(self): + @patch("airflow.providers.fab.auth_manager.cli_commands.permissions_command.cleanup_dag_permissions") + @patch("airflow.providers.fab.auth_manager.models.Resource") + def test_permissions_cleanup_specific_dag(self, mock_resource, mock_cleanup_dag_permissions): """Test cleanup for a specific DAG.""" # Mock args args = argparse.Namespace() args.dag_id = "test_dag" - args.dry_run = True + args.dry_run = False + args.yes = True args.verbose = True - with redirect_stdout(StringIO()) as stdout: + # Mock orphaned resource for the specific DAG + mock_orphaned_resource = MagicMock() + mock_orphaned_resource.name = "DAG:test_dag" + + with ( + patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder"), + patch("airflow.utils.session.create_session") as mock_session_ctx, + patch("sqlalchemy.select"), + redirect_stdout(StringIO()), + ): + mock_session = MagicMock() + mock_session_ctx.return_value.__enter__.return_value = mock_session + + # Mock DagModel query - return existing DAGs (NOT including the target DAG) + mock_dag_result = MagicMock() + mock_dag_result.all.return_value = [ + MagicMock(dag_id="existing_dag"), + MagicMock(dag_id="another_existing_dag"), + ] + + # Mock Resource query - return orphaned resources + mock_resource_result = MagicMock() + mock_resource_result.all.return_value = [mock_orphaned_resource] + + # Setup session.scalars to return different results for different queries + mock_session.scalars.side_effect = [mock_dag_result, mock_resource_result] + permissions_command.permissions_cleanup(args) - output = stdout.getvalue() - # Check for appropriate output indicating DAG-specific operation - assert ( - "test_dag" in output - or "No orphaned permissions found for DAG" in output - or "not found in orphaned permissions" in output - or "No orphaned DAG permissions found" in output - ) + # Should call cleanup_dag_permissions specifically for test_dag + mock_cleanup_dag_permissions.assert_called_once_with("test_dag", mock_session) + @patch("airflow.providers.fab.auth_manager.cli_commands.permissions_command.cleanup_dag_permissions") + @patch("airflow.providers.fab.auth_manager.models.Resource") @patch("builtins.input", return_value="n") - def test_permissions_cleanup_no_confirmation(self, mock_input): + def test_permissions_cleanup_no_confirmation( + self, mock_input, mock_resource, mock_cleanup_dag_permissions + ): """Test cleanup cancellation when user doesn't confirm.""" # Mock args args = argparse.Namespace() @@ -125,158 +199,155 @@ def test_permissions_cleanup_no_confirmation(self, mock_input): args.yes = False args.verbose = False - with redirect_stdout(StringIO()) as stdout: + # Mock orphaned resources + mock_orphaned_resource = MagicMock() + mock_orphaned_resource.name = "DAG:orphaned_dag" + + with ( + patch("airflow.providers.fab.auth_manager.cli_commands.utils.get_application_builder"), + patch("airflow.utils.session.create_session") as mock_session_ctx, + patch("sqlalchemy.select"), + redirect_stdout(StringIO()) as stdout, + ): + mock_session = MagicMock() + mock_session_ctx.return_value.__enter__.return_value = mock_session + + # Mock DagModel query - return existing DAGs + mock_dag_result = MagicMock() + mock_dag_result.all.return_value = [MagicMock(dag_id="existing_dag")] + + # Mock Resource query - return orphaned resources + mock_resource_result = MagicMock() + mock_resource_result.all.return_value = [mock_orphaned_resource] + + # Setup session.scalars to return different results for different queries + mock_session.scalars.side_effect = [mock_dag_result, mock_resource_result] + permissions_command.permissions_cleanup(args) output = stdout.getvalue() # Should not call cleanup if user declines or no orphaned permissions found assert "Cleanup cancelled" in output or "No orphaned DAG permissions found" in output + # cleanup_dag_permissions should NOT be called when user cancels + if "Cleanup cancelled" in output: + mock_cleanup_dag_permissions.assert_not_called() -class TestDagPermissions: - """Test cases for cleanup_dag_permissions function.""" - def test_cleanup_dag_permissions_with_session(self): - """Test cleanup_dag_permissions with provided session.""" - from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( - cleanup_dag_permissions, - ) +class TestDagPermissions: + """Test cases for cleanup_dag_permissions function with real database operations.""" - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), + @pytest.fixture(autouse=True) + def _setup_fab_test(self): + """Setup FAB for testing.""" + with conf_vars( + { + ( + "core", + "auth_manager", + ): "airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager", + } ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session with SQLAlchemy 2.0 methods - mock_session = MagicMock() - mock_session.scalars.return_value.all.return_value = [] - mock_session.execute.return_value = None - - cleanup_dag_permissions("test_dag", mock_session) + with get_application_builder(): + yield - # Should call session.scalars and session.execute for SQLAlchemy 2.0 - assert mock_session.scalars.called or mock_session.execute.called + def test_cleanup_dag_permissions_removes_specific_dag_resources(self): + """Test that cleanup_dag_permissions removes only the specified DAG resources.""" + from sqlalchemy import select - def test_cleanup_dag_permissions_without_session(self): - """Test cleanup_dag_permissions without provided session (creates new session).""" from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( cleanup_dag_permissions, ) + from airflow.providers.fab.auth_manager.models import Action, Permission, Resource + from airflow.security.permissions import RESOURCE_DAG_PREFIX + from airflow.utils.session import create_session + + with create_session() as session: + # Create resources for two different DAGs + target_resource = Resource(name=f"{RESOURCE_DAG_PREFIX}target_dag") + keep_resource = Resource(name=f"{RESOURCE_DAG_PREFIX}keep_dag") + session.add_all([target_resource, keep_resource]) + session.flush() + + # Get or create action + read_action = session.scalars(select(Action).where(Action.name == "can_read")).first() + if not read_action: + read_action = Action(name="can_read") + session.add(read_action) + session.flush() + + # Create permissions + target_perm = Permission(action=read_action, resource=target_resource) + keep_perm = Permission(action=read_action, resource=keep_resource) + session.add_all([target_perm, keep_perm]) + session.commit() + + # Execute cleanup + cleanup_dag_permissions("target_dag", session) + + # Verify: target resource deleted, keep resource remains + assert not session.get(Resource, target_resource.id) + assert session.get(Resource, keep_resource.id) + assert not session.get(Permission, target_perm.id) + assert session.get(Permission, keep_perm.id) + + def test_cleanup_dag_permissions_handles_no_matching_resources(self): + """Test that cleanup_dag_permissions handles DAGs with no matching resources gracefully.""" + from sqlalchemy import func, select - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), - ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock create_session - with patch("airflow.utils.session.create_session") as mock_create_session: - mock_session = MagicMock() - mock_create_session.return_value.__enter__.return_value = mock_session - mock_session.scalars.return_value.all.return_value = [] - mock_session.execute.return_value = None - - cleanup_dag_permissions("test_dag") - - # Should create a new session - mock_create_session.assert_called_once() - - def test_cleanup_dag_permissions_with_resources_and_permissions(self): - """Test cleanup_dag_permissions with actual resources and permissions.""" from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( cleanup_dag_permissions, ) + from airflow.providers.fab.auth_manager.models import Resource + from airflow.utils.session import create_session - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), - ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session - mock_session = MagicMock() - - # Mock resources exist - mock_dag_resource = MagicMock() - mock_dag_resource.id = 1 - - # Mock permissions exist - mock_dag_permission = MagicMock() - mock_dag_permission.id = 1 - - # Setup mock returns for different select queries - def mock_scalars_side_effect(stmt): - mock_result = MagicMock() - # Return resources for resource queries - mock_result.all.return_value = [mock_dag_resource] - return mock_result - - mock_session.scalars.side_effect = mock_scalars_side_effect - mock_session.execute.return_value = None - - cleanup_dag_permissions("test_dag", mock_session) + with create_session() as session: + initial_count = session.scalar(select(func.count(Resource.id))) + cleanup_dag_permissions("non_existent_dag", session) + assert session.scalar(select(func.count(Resource.id))) == initial_count - # Should perform database operations - assert mock_session.scalars.called or mock_session.execute.called - - def test_cleanup_dag_permissions_with_resources_but_no_permissions(self): - """Test cleanup_dag_permissions with resources but no permissions.""" + def test_cleanup_dag_permissions_handles_resources_without_permissions(self): + """Test cleanup when resources exist but have no permissions.""" from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( cleanup_dag_permissions, ) + from airflow.providers.fab.auth_manager.models import Resource + from airflow.security.permissions import RESOURCE_DAG_PREFIX + from airflow.utils.session import create_session - # Mock FAB models and select/delete - with ( - patch("airflow.providers.fab.auth_manager.models") as mock_models, - patch("sqlalchemy.select"), - patch("sqlalchemy.delete"), - ): - mock_resource = MagicMock() - mock_permission = MagicMock() - mock_assoc = MagicMock() - mock_models.Resource = mock_resource - mock_models.Permission = mock_permission - mock_models.assoc_permission_role = mock_assoc - - # Mock session - mock_session = MagicMock() - - # Mock resources exist but no permissions - mock_dag_resource = MagicMock() - mock_dag_resource.id = 1 - - def mock_scalars_side_effect(stmt): - mock_result = MagicMock() - # Return resources for resource queries - mock_result.all.return_value = [mock_dag_resource] - return mock_result + with create_session() as session: + # Create resource without permissions + resource = Resource(name=f"{RESOURCE_DAG_PREFIX}test_dag") + session.add(resource) + session.commit() + resource_id = resource.id - mock_session.scalars.side_effect = mock_scalars_side_effect - mock_session.execute.return_value = None + cleanup_dag_permissions("test_dag", session) + assert not session.get(Resource, resource_id) - cleanup_dag_permissions("test_dag", mock_session) + def test_cleanup_dag_permissions_with_default_session(self): + """Test cleanup_dag_permissions when no session is provided (uses default).""" + from sqlalchemy import func, select - # Should still perform database operations - assert mock_session.scalars.called or mock_session.execute.called + from airflow.providers.fab.auth_manager.cli_commands.permissions_command import ( + cleanup_dag_permissions, + ) + from airflow.providers.fab.auth_manager.models import Resource + from airflow.security.permissions import RESOURCE_DAG_PREFIX + from airflow.utils.session import create_session + + # Setup test data + with create_session() as session: + resource = Resource(name=f"{RESOURCE_DAG_PREFIX}test_dag") + session.add(resource) + session.commit() + + # Call cleanup without session parameter + cleanup_dag_permissions("test_dag") + + # Verify deletion + with create_session() as session: + count = session.scalar( + select(func.count(Resource.id)).where(Resource.name == f"{RESOURCE_DAG_PREFIX}test_dag") + ) + assert count == 0