From 7d50191b816c76db4c4da288f94cec2b820d90df Mon Sep 17 00:00:00 2001 From: PoAn Yang Date: Sat, 11 Jul 2026 10:02:50 +0900 Subject: [PATCH] Avoid deactivating configured Dag bundles on transient get bundle errors Signed-off-by: PoAn Yang --- .../airflow/dag_processing/bundles/manager.py | 7 +- .../bundles/test_dag_bundle_manager.py | 30 +++++++ .../tests/unit/dag_processing/test_manager.py | 84 +++++++++++++++++++ 3 files changed, 120 insertions(+), 1 deletion(-) diff --git a/airflow-core/src/airflow/dag_processing/bundles/manager.py b/airflow-core/src/airflow/dag_processing/bundles/manager.py index 7d966e0b5074f..be06a49efed9a 100644 --- a/airflow-core/src/airflow/dag_processing/bundles/manager.py +++ b/airflow-core/src/airflow/dag_processing/bundles/manager.py @@ -314,13 +314,18 @@ def _extract_and_sign_template(bundle_name: str) -> tuple[str | None, dict]: if not team: raise _bundle_item_exc(f"Team '{config.team_name}' does not exist") + # Pop configured bundles out of ``stored`` before ``_extract_and_sign_template``. + # A transient connection error from that call must not deactivate the bundle, + # which would make all Dags in it non-schedulable. + bundle = stored.pop(name, None) + try: new_template, new_params = _extract_and_sign_template(name) except Exception as e: self.log.exception("Error creating bundle '%s': %s", name, e) continue - if bundle := stored.pop(name, None): + if bundle: bundle.active = True if new_template != bundle.signed_url_template: bundle.signed_url_template = new_template diff --git a/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py b/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py index 0d17069c831dd..3180f116f70d5 100644 --- a/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py +++ b/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py @@ -199,6 +199,36 @@ def _get_bundle_names_and_active(): ] +@pytest.mark.db_test +@conf_vars({("core", "LOAD_EXAMPLES"): "False"}) +def test_sync_bundles_to_db_keeps_configured_bundle_on_transient_creation_error(clear_db, session): + """A configured bundle failing to instantiate with transient error must not be deactivated.""" + with patch.dict( + os.environ, {"AIRFLOW__DAG_PROCESSOR__DAG_BUNDLE_CONFIG_LIST": json.dumps(BASIC_BUNDLE_CONFIG)} + ): + manager = DagBundlesManager() + manager.sync_bundles_to_db() + + session.add( + ParseImportError( + bundle_name="my-test-bundle", + filename="some_file.py", + stacktrace="some error", + ) + ) + session.flush() + + manager = DagBundlesManager() + with patch.object( + DagBundlesManager, "get_bundle", autospec=True, side_effect=RuntimeError("transient error") + ): + manager.sync_bundles_to_db() + + bundle = session.scalar(select(DagBundleModel).where(DagBundleModel.name == "my-test-bundle")) + assert bundle.active is True + assert session.scalar(select(func.count(ParseImportError.id))) == 1 + + @pytest.mark.db_test @conf_vars({("core", "LOAD_EXAMPLES"): "False"}) def test_sync_bundles_to_db_does_not_log_removing_none_team(clear_db, caplog): diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index ae28f3317fb79..106cd6398df26 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -45,6 +45,7 @@ from airflow._shared.timezones import timezone from airflow.callbacks.callback_requests import DagCallbackRequest from airflow.dag_processing.bundles.base import BaseDagBundle +from airflow.dag_processing.bundles.local import LocalDagBundle from airflow.dag_processing.bundles.manager import DagBundlesManager from airflow.dag_processing.dagbag import DagBag from airflow.dag_processing.manager import ( @@ -83,6 +84,18 @@ logger = logging.getLogger(__name__) TEST_DAG_FOLDER = Path(__file__).parents[1].resolve() / "dags" + + +class _HealthControlledDagBundle(LocalDagBundle): + """LocalDagBundle that fails to instantiate unless its health file contains "ok".""" + + def __init__(self, *, health_file: str, **kwargs) -> None: + state = Path(health_file).read_text().strip() + if state != "ok": + raise RuntimeError(f"simulated {state}") + super().__init__(**kwargs) + + DEFAULT_DATE = timezone.datetime(2016, 1, 1) @@ -1072,6 +1085,77 @@ def test_deactivate_stale_dags_marks_dags_in_inactive_bundles(self, session): ) assert is_stale_by_dag == {"dag_in_inactive_bundle": True, "dag_in_active_bundle": False} + def test_transient_bundle_creation_error_keeps_bundle_active_and_dags_fresh(self, tmp_path): + """A bundle failing to instantiate with transient error must not be deactivated.""" + dags_dir = tmp_path / "dags" + dags_dir.mkdir() + (dags_dir / "good_dag.py").write_text( + textwrap.dedent( + """ + from airflow.providers.standard.operators.empty import EmptyOperator + from airflow.sdk import DAG + + with DAG(dag_id="dag_in_flaky_bundle", schedule=None): + EmptyOperator(task_id="noop") + """ + ) + ) + (dags_dir / "broken_dag.py").write_text("an invalid airflow DAG") + health_file = tmp_path / "health.txt" + + bundle_config = [ + { + "name": "flaky-bundle", + "classpath": f"{__name__}._HealthControlledDagBundle", + "kwargs": { + "path": str(dags_dir), + "health_file": str(health_file), + "refresh_interval": 0, + }, + } + ] + + def get_db_state(): + with create_session() as session: + active = session.scalar( + select(DagBundleModel.active).where(DagBundleModel.name == "flaky-bundle") + ) + fresh_dags = session.scalar( + select(func.count()) + .select_from(DagModel) + .where(DagModel.bundle_name == "flaky-bundle", ~DagModel.is_stale) + ) + import_errors = session.scalar(select(func.count()).select_from(ParseImportError)) + return active, fresh_dags, import_errors + + with conf_vars({("dag_processor", "dag_bundle_config_list"): json.dumps(bundle_config)}): + # Phase 1: misconfigured from day one — sync never succeeds, so no bundle row exists at all + health_file.write_text("misconfiguration") + DagFileProcessorManager(max_runs=1, processor_timeout=365 * 86_400).run() + assert get_db_state() == (None, 0, 0) + + # Phase 2: healthy restart — the bundle row appears and Dags get parsed + health_file.write_text("ok") + DagFileProcessorManager(max_runs=1, processor_timeout=365 * 86_400).run() + assert get_db_state() == (True, 1, 1) + + # Phase 3: restart with transient connection error — must not deactivate the bundle + health_file.write_text("transient connection error") + DagFileProcessorManager(max_runs=1, processor_timeout=365 * 86_400).run() + assert get_db_state() == (True, 1, 1) + + # Phase 4: healthy again + health_file.write_text("ok") + DagFileProcessorManager(max_runs=1, processor_timeout=365 * 86_400).run() + assert get_db_state() == (True, 1, 1) + + # Phase 5: removing the (again failing) bundle from config must follow the removal + # path: deactivate it, mark its Dags stale, and drop its import errors. + health_file.write_text("transient connection error") + with conf_vars({("dag_processor", "dag_bundle_config_list"): json.dumps([])}): + DagFileProcessorManager(max_runs=1, processor_timeout=365 * 86_400).run() + assert get_db_state() == (False, 0, 0) + @mock.patch("airflow.dag_processing.manager.is_lock_not_available_error") @pytest.mark.usefixtures("testing_dag_bundle") def test_deactivate_stale_dags_handles_lock_timeout(self, mock_is_lock_not_available, session, caplog):