diff --git a/airflow-core/tests/unit/api_fastapi/conftest.py b/airflow-core/tests/unit/api_fastapi/conftest.py index f275d48f967be..93261560cb339 100644 --- a/airflow-core/tests/unit/api_fastapi/conftest.py +++ b/airflow-core/tests/unit/api_fastapi/conftest.py @@ -17,6 +17,7 @@ from __future__ import annotations import datetime +import json import os from typing import TYPE_CHECKING from unittest import mock @@ -30,12 +31,10 @@ from airflow.api_fastapi.app import create_app from airflow.api_fastapi.auth.managers.simple.user import SimpleAuthManagerUser from airflow.dag_processing.bundles.manager import DagBundlesManager -from airflow.models import Connection -from airflow.providers.git.bundles.git import GitDagBundle from airflow.providers.standard.operators.empty import EmptyOperator from tests_common.test_utils.config import conf_vars -from tests_common.test_utils.db import clear_db_connections, parse_and_sync_to_db +from tests_common.test_utils.db import parse_and_sync_to_db if TYPE_CHECKING: from airflow.api_fastapi.auth.managers.simple.simple_auth_manager import SimpleAuthManager @@ -197,43 +196,23 @@ def create_test_client(apps="all"): @pytest.fixture -def configure_git_connection_for_dag_bundle(session): - clear_db_connections(False) - # Git connection is required for the bundles to have a url. - connection = Connection( - conn_id="git_default", - conn_type="git", - description="default git connection", - host="http://test_host.github.com", - port=8081, - login="", - ) - session.add(connection) - with ( - conf_vars( - { - ( - "dag_processor", - "dag_bundle_config_list", - ): '[{ "name": "dag_maker", "classpath": "airflow.providers.git.bundles.git.GitDagBundle", "kwargs": {"subdir": "dags", "tracking_ref": "main", "refresh_interval": 0}}, { "name": "another_bundle_name", "classpath": "airflow.providers.git.bundles.git.GitDagBundle", "kwargs": {"subdir": "dags", "tracking_ref": "main", "refresh_interval": 0}}]' - } - ), - mock.patch("airflow.providers.git.bundles.git.GitHook") as mock_git_hook, - mock.patch.object(GitDagBundle, "get_current_version") as mock_get_current_version, - ): - mock_get_current_version.return_value = "some_commit_hash" - mock_git_hook.return_value.repo_url = connection.host +def configure_dag_bundles_with_view_url(): + # Rendering a per-version link only needs the URL template a remote bundle would supply. + bundle_config = [ + { + "name": name, + "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", + "kwargs": {"view_url_template": "http://test_host.github.com/tree/{version}/dags"}, + } + for name in ("dag_maker", "another_bundle_name") + ] + with conf_vars({("dag_processor", "dag_bundle_config_list"): json.dumps(bundle_config)}): DagBundlesManager().sync_bundles_to_db() yield - # in case no flush or commit was executed after the "session.add" above, we need to flush the session - # manually here to make sure that the added connection will be deleted by query(Connection).delete() - # in the`clear_db_connections` function below - session.flush() - clear_db_connections(False) @pytest.fixture -def make_dag_with_multiple_versions(dag_maker, configure_git_connection_for_dag_bundle, session): +def make_dag_with_multiple_versions(dag_maker, configure_dag_bundles_with_view_url, session): """ Create DAG with multiple versions diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py index 33cc8f9fbd159..cdb5fc033a4bf 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py @@ -2262,7 +2262,7 @@ def create_dags(self, setup, dag_maker, session): EmptyOperator(task_id="task") session.commit() - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") @mock.patch( "airflow.api_fastapi.auth.managers.simple.user.SimpleAuthManagerUser.get_display_name", return_value="Jane Doe", @@ -2272,7 +2272,7 @@ def test_materialize_records_triggering_user_display_name(self, mock_display_nam assert response.status_code == 200 assert response.json()["triggering_user_name"] == "Jane Doe" - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200(self, test_client): response = test_client.post("/assets/1/materialize") assert response.status_code == 200 @@ -2302,14 +2302,14 @@ def test_should_respond_200(self, test_client): "team_name": None, } - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_with_partition_key(self, test_client): partition_key = "2026-03-23" response = test_client.post("/assets/1/materialize", json={"partition_key": partition_key}) assert response.status_code == 200 assert response.json()["partition_key"] == partition_key - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_with_trigger_fields(self, test_client): payload = { "conf": {"foo": "bar"}, @@ -2332,7 +2332,7 @@ def test_should_respond_200_with_trigger_fields(self, test_client): assert response.json()["partition_key"] == "2026-03-23" assert response.json()["run_type"] == "asset_materialization" - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_with_trigger_fields_without_dag_run_id(self, test_client): payload = { "conf": {"foo": "bar"}, @@ -2429,7 +2429,7 @@ def test_materialize_allowed_run_types_from_requested_version(self, test_client, == f"Dag with dag_id: '{self.DAG_ASSET1_ID}' does not allow asset materialization runs" ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_403_when_user_cannot_trigger_dag(self, test_client): with mock.patch( "airflow.api_fastapi.core_api.routes.public.assets.get_auth_manager", @@ -2501,7 +2501,7 @@ def test_should_respond_with_bundle_version(self, test_client, session, dag_make in response.json()["detail"] ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_400_on_invalid_dag_run_id(self, test_client): """A dag_run_id containing '..' triggers ValueError in DagRun.validate_run_id. @@ -2514,7 +2514,7 @@ def test_should_respond_400_on_invalid_dag_run_id(self, test_client): assert response.status_code == 400 assert "must not contain '..'" in response.json()["detail"] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_with_partition_date_for_partitioned_dag( self, test_client, dag_maker, session ): diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py index 8edfcb278161e..b1b8d0576132d 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py @@ -373,7 +373,7 @@ class TestGetDagRun: ), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_get_dag_run(self, test_client, dag_id, run_id, state, run_type, triggered_by, dag_run_note): response = test_client.get(f"/dags/{dag_id}/dagRuns/{run_id}") assert response.status_code == 200 @@ -386,7 +386,7 @@ def test_get_dag_run(self, test_client, dag_id, run_id, state, run_type, trigger assert body["note"] == dag_run_note @conf_vars({("core", "multi_team"): "True"}) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_get_dag_run_includes_team_name(self, test_client, session): with attach_dag_to_team(session, DAG1_ID, bundle_name="team-bundle-run", team_name="team-run"): response = test_client.get(f"/dags/{DAG1_ID}/dagRuns/{DAG1_RUN1_ID}") @@ -413,7 +413,7 @@ class TestGetDagRuns: ("dag_id", "total_entries"), [(DAG1_ID, 2), (DAG2_ID, 2), ("~", 4)], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_get_dag_runs(self, test_client, session, dag_id, total_entries): response = test_client.get(f"/dags/{dag_id}/dagRuns") assert response.status_code == 200 @@ -425,7 +425,7 @@ def test_get_dag_runs(self, test_client, session, dag_id, total_entries): ).one() assert each == get_dag_run_dict(run) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_get_dag_runs_not_found(self, test_client): response = test_client.get("/dags/invalid/dagRuns") assert response.status_code == 404 @@ -448,7 +448,7 @@ def test_partition_date_day_filters_reject_non_partitioned_dag(self, test_client ) @conf_vars({("core", "multi_team"): "True"}) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_get_dag_runs_includes_team_name(self, test_client, session): with attach_dag_to_team(session, DAG1_ID, bundle_name="team-bundle-runs", team_name="team-runs"): response = test_client.get(f"/dags/{DAG1_ID}/dagRuns") @@ -458,7 +458,7 @@ def test_get_dag_runs_includes_team_name(self, test_client, session): assert all(run["team_name"] == "team-runs" for run in body["dag_runs"]) @conf_vars({("core", "multi_team"): "True"}) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_get_dag_runs_filtered_by_team(self, test_client, session): with attach_dag_to_team(session, DAG1_ID, bundle_name="team-bundle-filter", team_name="team-filter"): response = test_client.get("/dags/~/dagRuns", params={"teams": ["team-filter"]}) @@ -505,7 +505,7 @@ def test_should_respond_403(self, unauthorized_test_client): pytest.param("duration", [DAG1_RUN1_ID, DAG1_RUN2_ID], id="order_by_duration"), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_return_correct_results_with_order_by(self, test_client, order_by, expected_order): # Test ascending order @@ -536,7 +536,7 @@ def test_return_correct_results_with_order_by(self, test_client, order_by, expec ({"limit": 1, "offset": 2}, []), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_limit_and_offset(self, test_client, query_params, expected_dag_id_order): response = test_client.get("/dags/test_dag1/dagRuns", params=query_params) assert response.status_code == 200 @@ -606,7 +606,7 @@ def test_bad_limit_and_offset(self, test_client, query_params, expected_detail): "-run_after", ], # test with multiple ordering fields (alias, non-alias, datetime, non-datetime) ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_cursor_pagination_first_two_page(self, test_client, order_by): """First page with cursor='' and second page fetched via the returned next_cursor.""" response = test_client.get( @@ -639,7 +639,7 @@ def test_cursor_pagination_first_two_page(self, test_client, order_by): "order_by", ["id", "dag_run_id", "logical_date", "-run_after"], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_cursor_pagination_returns_cursor_response(self, test_client, order_by): """When cursor param is provided, response has cursor fields and a bounded total_entries.""" response1 = test_client.get( @@ -670,7 +670,7 @@ def test_cursor_pagination_returns_cursor_response(self, test_client, order_by): "order_by", ["id", "dag_run_id", "logical_date", "-run_after"], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_cursor_pagination_forward_and_backward_consistency(self, test_client, order_by): """Walk all pages forward via next_cursor, then backward via previous_cursor, and compare.""" total_runs = 4 # 4 dag runs are created by the setup fixture @@ -724,7 +724,7 @@ def test_cursor_pagination_forward_and_backward_consistency(self, test_client, o assert all_backward == forward_ids @mock.patch("airflow.api_fastapi.common.db.common.EXACT_COUNT_LIMIT", 2) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_cursor_pagination_total_entries_capped(self, test_client): """With more matching runs than the cap, total_entries is reported as the cap.""" response = test_client.get( @@ -737,7 +737,7 @@ def test_cursor_pagination_total_entries_capped(self, test_client): assert body["total_entries"] == 2 assert body["total_entries_limit"] == 2 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_cursor_pagination_invalid_token(self, test_client): response = test_client.get( "/dags/~/dagRuns", @@ -745,7 +745,7 @@ def test_cursor_pagination_invalid_token(self, test_client): ) assert response.status_code == 400 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_cursor_pagination_nullable_sort_column_returns_all_rows(self, test_client, session): """Cursor pagination sorted by a nullable column must not silently drop rows. @@ -1076,14 +1076,14 @@ def test_cursor_pagination_nullable_sort_column_returns_all_rows(self, test_clie ("~", {"partition_key_prefix_pattern": "2026"}, [DAG1_RUN1_ID]), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_filters(self, test_client, dag_id, query_params, expected_dag_id_list): response = test_client.get(f"/dags/{dag_id}/dagRuns", params=query_params) assert response.status_code == 200 body = response.json() assert [each["dag_run_id"] for each in body["dag_runs"]] == expected_dag_id_list - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_partition_date_day_filters_use_timetable_timezone(self, test_client, dag_maker, session): dag_id = "test_partition_date_local_day" with dag_maker( @@ -1232,7 +1232,7 @@ def test_invalid_dag_version(self, test_client): class TestListDagRunsBatch: - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_list_dag_runs_return_200(self, test_client, session): with assert_queries_count(5): response = test_client.post("/dags/~/dagRuns/list", json={}) @@ -1274,7 +1274,7 @@ def test_list_dag_runs_with_invalid_dag_id(self, test_client): [["invalid"], 200, []], ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_list_dag_runs_with_dag_ids_filter(self, test_client, dag_ids, status_code, expected_dag_id_list): with assert_queries_count(5): response = test_client.post("/dags/~/dagRuns/list", json={"dag_ids": dag_ids}) @@ -1307,7 +1307,7 @@ def test_invalid_order_by_raises_400(self, test_client): pytest.param("conf", DAG_RUNS_LIST, id="order_by_conf"), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_dag_runs_ordering(self, test_client, order_by, expected_order): # Test ascending order response = test_client.post("/dags/~/dagRuns/list", json={"order_by": order_by}) @@ -1335,7 +1335,7 @@ def test_dag_runs_ordering(self, test_client, order_by, expected_order): ({"page_limit": 1, "page_offset": 2}, DAG_RUNS_LIST[2:3]), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_limit_and_offset(self, test_client, post_body, expected_dag_id_order): response = test_client.post("/dags/~/dagRuns/list", json=post_body) assert response.status_code == 200 @@ -1455,7 +1455,7 @@ def test_bad_limit_and_offset(self, test_client, post_body, expected_detail): ), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_filters(self, test_client, post_body, expected_dag_id_list): response = test_client.post("/dags/~/dagRuns/list", json=post_body) assert response.status_code == 200 @@ -1616,7 +1616,7 @@ class TestPatchDagRun: ), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_patch_dag_run(self, test_client, dag_id, run_id, patch_body, response_body, note_data, session): response = test_client.patch(f"/dags/{dag_id}/dagRuns/{run_id}", json=patch_body) assert response.status_code == 200 @@ -1684,7 +1684,7 @@ def test_should_respond_403(self, unauthorized_test_client): ), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_patch_dag_run_with_update_mask( self, test_client, query_params, patch_body, response_body, expected_status_code, note_data, session ): @@ -1722,7 +1722,7 @@ def test_patch_dag_run_bad_request(self, test_client): ("failed", [DagRunState.FAILED], "Dag Run's state was manually set to `failed`."), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_patch_dag_run_notifies_listeners( self, test_client, state, expected_dagrun_state, expected_msg, listener_manager ): @@ -1735,7 +1735,7 @@ def test_patch_dag_run_notifies_listeners( assert listener.dag_run_msg == expected_msg assert listener.dag_run_has_dag_attr is True - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_patch_dag_run_listener_sees_note_when_note_and_state_both_patched( self, test_client, listener_manager ): @@ -1755,7 +1755,7 @@ def test_patch_dag_run_listener_sees_note_when_note_and_state_both_patched( ("failed", TaskInstanceState.FAILED), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_patch_dag_run_notifies_ti_listeners_for_running_tasks( self, test_client, @@ -1791,7 +1791,7 @@ def test_patch_dag_run_notifies_ti_listeners_for_running_tasks( assert listener.state[1] is DagRunState(dag_run_state) assert len(listener.state) == 2 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_patch_dag_run_does_not_notify_ti_listeners_for_non_running_tasks( self, test_client, @@ -1825,7 +1825,7 @@ def test_patch_dag_run_does_not_notify_ti_listeners_for_non_running_tasks( # The list length check distinguishes "only dagrun fired" from "both fired". assert listener.state == [DagRunState.SUCCESS] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_patch_dag_run_does_not_notify_ti_listeners_for_running_teardown_tasks( self, test_client, @@ -1903,7 +1903,7 @@ def test_should_respond_403(self, unauthorized_test_client): class TestGetDagRunAssetTriggerEvents: - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") @pytest.mark.parametrize( "partition_key", ["test_partition_key", None], @@ -2008,7 +2008,7 @@ def test_should_respond_200( } assert response.json() == expected_response - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") @pytest.mark.parametrize("num_events", [1, 5]) def test_query_count_does_not_scale_with_consumed_events( self, num_events, test_client, dag_maker, session @@ -2084,7 +2084,7 @@ def test_should_respond_404(self, test_client): class TestClearDagRun: - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run(self, test_client, session): response = test_client.post( f"/dags/{DAG1_ID}/dagRuns/{DAG1_RUN1_ID}/clear", @@ -2102,7 +2102,7 @@ def test_clear_dag_run(self, test_client, session): logical_date=None, ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_whose_dag_version_was_deleted(self, test_client, session): """A run that kept its bundle version after ``airflow db clean`` removed its Dag version.""" session.execute( @@ -2148,7 +2148,7 @@ def test_should_respond_403(self, unauthorized_test_client): "omit-leaves-existing", ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_applies_note(self, test_client, session, body, expected_note): """``note`` in the clear body writes to the Dag Run; ``None`` / unset leaves it alone.""" response = test_client.post(f"/dags/{DAG1_ID}/dagRuns/{DAG1_RUN1_ID}/clear", json=body) @@ -2159,7 +2159,7 @@ def test_clear_dag_run_applies_note(self, test_client, session, body, expected_n ) assert dag_run.note == expected_note - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_dry_run_does_not_apply_note(self, test_client, session): """``note`` is ignored on dry-run (no side effects).""" response = test_client.post( @@ -2181,7 +2181,7 @@ def test_clear_dag_run_dry_run_does_not_apply_note(self, test_client, session): [{"only_failed": True}, DAG1_RUN2_ID, ["failed"]], ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_dry_run(self, test_client, session, body, dag_run_id, expected_state): response = test_client.post(f"/dags/{DAG1_ID}/dagRuns/{dag_run_id}/clear", json=body) assert response.status_code == 200 @@ -2201,7 +2201,7 @@ def test_clear_dag_run_dry_run(self, test_client, session, body, dag_run_id, exp ) assert logs == 0 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_dry_run_response_has_full_task_instance_fields(self, test_client): """Regression test: dry-run response must include all TaskInstanceResponse fields. @@ -2229,7 +2229,7 @@ def test_clear_dag_run_dry_run_response_has_full_task_instance_fields(self, test # rendered_fields must be present (defaults to {}) assert "rendered_fields" in ti - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_dry_run_only_failed_returns_only_failed_tasks_with_full_fields(self, test_client): """Regression test: only_failed=True dry-run must return only failed TIs with full fields. @@ -2267,7 +2267,7 @@ def test_clear_dag_run_unprocessable_entity(self, test_client): assert body["detail"][0]["msg"] == "Field required" assert body["detail"][0]["loc"][0] == "body" - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_only_new_dry_run(self, test_client, session): """Test that only_new dry_run returns 0 new tasks when all tasks already have TIs. @@ -2291,7 +2291,7 @@ def test_clear_dag_run_only_new_dry_run(self, test_client, session): assert logs == 0 @mock.patch("airflow.serialization.definitions.dag.SerializedDAG.clear") - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_dag_run_only_new_non_dry_run(self, mock_clear, test_client, session): """Test that only_new non-dry_run clears and returns a DAGRunResponse.""" mock_clear.return_value = 2 @@ -2328,7 +2328,7 @@ def test_clear_dag_run_only_new_and_only_failed_mutually_exclusive(self, test_cl class TestBulkClearDagRuns: - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_specific_dag(self, test_client, session): """Specific dag_id in URL, dag_run_id in body — clears both runs and queues them.""" response = test_client.post( @@ -2347,7 +2347,7 @@ def test_bulk_clear_specific_dag(self, test_client, session): assert run["state"] == "queued" assert run["dag_id"] == DAG1_ID - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_wildcard_across_dags(self, test_client, session): """``~`` URL with per-entity dag_id — clears runs across Dags in one call.""" response = test_client.post( @@ -2368,7 +2368,7 @@ def test_bulk_clear_wildcard_across_dags(self, test_client, session): for run in body["dag_runs"]: assert run["state"] == "queued" - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_dry_run_collects_affected_tis_across_runs(self, test_client, session): """Dry-run returns the union of affected TIs across the listed runs without mutating state.""" response = test_client.post( @@ -2390,7 +2390,7 @@ def test_bulk_clear_dry_run_collects_affected_tis_across_runs(self, test_client, ) assert dag_run.state == DAG1_RUN1_STATE - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_dry_run_only_failed_filters(self, test_client): """``only_failed=True`` shrinks the dry-run preview to failed TIs only.""" response = test_client.post( @@ -2406,7 +2406,7 @@ def test_bulk_clear_dry_run_only_failed_filters(self, test_client): assert all(ti["state"] == "failed" for ti in body["task_instances"]) assert body["total_entries"] == 1 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_applies_note_to_each_run(self, test_client, session): """``note`` in the body is applied to every cleared run in the same transaction.""" response = test_client.post( @@ -2422,7 +2422,7 @@ def test_bulk_clear_applies_note_to_each_run(self, test_client, session): dag_run = session.scalar(select(DagRun).where(DagRun.dag_id == DAG1_ID, DagRun.run_id == run_id)) assert dag_run.note == "bulk cleared by test" - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_wildcard_rejects_missing_dag_id(self, test_client): """``~`` URL requires every entry to carry a concrete dag_id; 400 otherwise.""" response = test_client.post( @@ -2435,7 +2435,7 @@ def test_bulk_clear_wildcard_rejects_missing_dag_id(self, test_client): assert response.status_code == 400 assert DAG1_RUN1_ID in response.json()["detail"] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_specific_url_rejects_mismatched_dag_id(self, test_client): """When the URL has a specific dag_id, mismatched per-entity dag_id is rejected.""" response = test_client.post( @@ -2447,7 +2447,7 @@ def test_bulk_clear_specific_url_rejects_mismatched_dag_id(self, test_client): ) assert response.status_code == 400 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_missing_run_returns_404(self, test_client): response = test_client.post( f"/dags/{DAG1_ID}/clearDagRuns", @@ -2458,7 +2458,7 @@ def test_bulk_clear_missing_run_returns_404(self, test_client): ) assert response.status_code == 404 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_rejects_only_new_with_only_failed(self, test_client): """``only_new`` and ``only_failed`` are mutually exclusive at the body validator level.""" response = test_client.post( @@ -2486,7 +2486,7 @@ def test_bulk_clear_unauthorized_returns_403(self, unauthorized_test_client): ) assert response.status_code == 403 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_bulk_clear_rejects_unauthorized_dag_ids_from_request_body(self, test_client, session): """A 403 at the route level if any entry references a Dag the user can't access; nothing is cleared.""" restricted_bundle_name = "restricted-bundle-clear" @@ -2548,7 +2548,7 @@ class TestClearDagRunOnlyNew: """ @pytest.fixture - def dag_two_versions(self, dag_maker, configure_git_connection_for_dag_bundle, session): + def dag_two_versions(self, dag_maker, configure_dag_bundles_with_view_url, session): """ Two-version DAG with one run on v1. @@ -2660,7 +2660,7 @@ class TestBulkClearDagRunsPartitionSelector: """ @pytest.fixture - def partition_dag(self, dag_maker, configure_git_connection_for_dag_bundle, session): + def partition_dag(self, dag_maker, configure_dag_bundles_with_view_url, session): """Dag with two runs carrying partition_key and partition_date.""" with dag_maker( PARTITION_DAG_ID, @@ -2899,7 +2899,7 @@ class TestClearPartitions: """ @pytest.fixture - def partitioned_dag_with_runs(self, dag_maker, configure_git_connection_for_dag_bundle, session): + def partitioned_dag_with_runs(self, dag_maker, configure_dag_bundles_with_view_url, session): """Dag with three runs carrying partition fields and task instances.""" dag_id = "clear_partitions_test_dag" with dag_maker( @@ -3166,7 +3166,7 @@ def test_partition_date_single_bound_window( assert run.partition_key is not None def test_cross_dag_run_id_collision_does_not_clear_other_dag( - self, test_client, session, dag_maker, configure_git_connection_for_dag_bundle + self, test_client, session, dag_maker, configure_dag_bundles_with_view_url ): """Clearing by run_id only affects the target Dag; a second Dag with the same run_id is untouched.""" shared_run_id = "shared_run_id" @@ -3254,7 +3254,7 @@ def test_cross_dag_run_id_collision_does_not_clear_other_dag( assert ti_b_after.state == State.SUCCESS def test_cross_dag_run_id_collision_dry_run_counts_only_target_dag( - self, test_client, session, dag_maker, configure_git_connection_for_dag_bundle + self, test_client, session, dag_maker, configure_dag_bundles_with_view_url ): """dry_run=True TI count only includes the target Dag's TIs, not a bystander sharing the same run_id.""" shared_run_id = "shared_dry_run_id" @@ -3390,7 +3390,7 @@ def _dags_for_trigger_tests(self, session=None): (None, None, None, None, None), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200( self, test_client, dag_run_id, note, data_interval_start, data_interval_end, note_data, session ): @@ -3652,7 +3652,7 @@ def test_should_respond_400_if_manual_runs_denied(self, test_client, session, da assert response.json()["detail"] == f"Dag with dag_id: '{dag_id}' does not allow manual runs" @time_machine.travel(timezone.utcnow(), tick=False) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_response_409_for_duplicate_logical_date(self, test_client): RUN_ID_1 = "random_1" RUN_ID_2 = "random_2" @@ -3752,7 +3752,7 @@ def test_response_409(self, test_client): assert "detail" in response_json assert list(response_json["detail"].keys()) == ["reason", "statement", "orig_error", "message"] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_with_null_logical_date(self, test_client): response = test_client.post( f"/dags/{DAG1_ID}/dagRuns", @@ -3785,7 +3785,7 @@ def test_should_respond_200_with_null_logical_date(self, test_client): "team_name": None, } - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_generate_unique_run_id_for_scheduled_dag(self, dag_maker, test_client, session): "Ensure manual triggers on scheduled DAGs don't conflict on run_id" scheduled_dag_id = "test_scheduled_dag" @@ -3877,7 +3877,7 @@ def test_custom_timetable_generate_run_id_for_manual_trigger(self, dag_maker, te run = session.scalars(select(DagRun).where(DagRun.run_id == run_id_without_logical_date)).one() assert run.dag_id == custom_dag_id - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_trigger_dag_run_with_bundle_version(self, test_client, session, dag_maker): """Test triggering a DAG run with a specific bundle version.""" from tests_common.test_utils.dag import sync_dag_to_db @@ -3991,7 +3991,7 @@ def test_trigger_dag_run_bundle_version_validates_against_old_param_schema( ) assert response.status_code == 200 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_trigger_dag_run_bundle_version_uses_v1_timetable(self, test_client, session, dag_maker): """Triggering with bundle_version='v1' must derive data_interval from v1's timetable, not v2's.""" from tests_common.test_utils.dag import sync_dag_to_db @@ -4031,7 +4031,7 @@ def test_trigger_dag_run_bundle_version_uses_v1_timetable(self, test_client, ses assert data["data_interval_start"] == "2023-12-31T00:00:00Z" assert data["data_interval_end"] == "2024-01-01T00:00:00Z" - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_trigger_dag_run_allowed_run_types_from_requested_version(self, test_client, session, dag_maker): """allowed_run_types is enforced from the requested bundle version, not the latest.""" from tests_common.test_utils.dag import sync_dag_to_db @@ -4090,7 +4090,7 @@ def test_should_respond_400_when_partition_key_given_for_non_partitioned_dag(sel assert response.status_code == 400 assert "not a partitioned Dag" in response.json()["detail"] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_when_partition_key_given_for_partitioned_dag( self, dag_maker, test_client, session ): @@ -4117,7 +4117,7 @@ def test_should_respond_200_when_partition_key_given_for_partitioned_dag( ) assert response.status_code == 200 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_when_partition_key_given_for_partitioned_at_runtime_dag( self, dag_maker, test_client, session ): @@ -4157,7 +4157,7 @@ def test_should_respond_400_when_empty_partition_key(self, test_client): == "Dag 'test_dag1' is not a partitioned Dag and does not accept a partition_key." ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_400_when_over_length_partition_key(self, dag_maker, test_client, session): """A partition_key exceeding 250 characters must return 400, not 500.""" partitioned_dag_id = "test_over_length_partition_key" @@ -4179,7 +4179,7 @@ def test_should_respond_400_when_over_length_partition_key(self, dag_maker, test assert response.status_code == 400 assert "at most 250 characters" in response.json()["detail"] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_respond_200_when_exactly_max_length_partition_key(self, dag_maker, test_client, session): """A partition_key of exactly 250 characters must be accepted.""" partitioned_dag_id = "test_max_length_partition_key" @@ -4200,7 +4200,7 @@ def test_should_respond_200_when_exactly_max_length_partition_key(self, dag_make ) assert response.status_code == 200 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_trigger_partitioned_dag_populates_partition_date(self, dag_maker, test_client, session): """Triggering a CronPartitionTimetable Dag with a valid key populates partition_date on the run. @@ -4231,7 +4231,7 @@ def test_trigger_partitioned_dag_populates_partition_date(self, dag_maker, test_ assert dag_run.partition_key == "2025-06-01T00:00:00" assert dag_run.partition_date == datetime(2025, 6, 1, 0, 0, 0, tzinfo=timezone.utc) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_trigger_partitioned_dag_invalid_key_returns_400(self, dag_maker, test_client, session): """An invalid partition_key for a CronPartitionTimetable Dag must return HTTP 400.""" partitioned_dag_id = "test_trigger_invalid_partition_key" @@ -4252,7 +4252,7 @@ def test_trigger_partitioned_dag_invalid_key_returns_400(self, dag_maker, test_c ) assert response.status_code == 400 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_trigger_partitioned_at_runtime_dag_leaves_partition_date_none( self, dag_maker, test_client, session ): @@ -4344,7 +4344,7 @@ def test_dag_level_false_overrides_fallback_true(self, dag_maker, session): result = resolve_run_on_latest_version(None, "test_resolver_dag_false", session, fallback=True) assert result is False - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_clear_endpoint_invokes_resolver_when_field_omitted(self, test_client): """Clearing without run_on_latest_version triggers the server-side resolver.""" with mock.patch( diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py index b37ccf1a1302f..af7d820691119 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dags.py @@ -1220,7 +1220,7 @@ class TestDagDetails(TestDagEndpoint): ({}, DAG2_ID, 200, DAG2_ID, "2021-06-15T00:00:00Z", {}, 0.24), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_dag_details( self, test_client, @@ -1609,7 +1609,7 @@ def _create_dag_for_deletion( (DAG5_ID, DAG5_DISPLAY_NAME, 409, 200, True, True), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_delete_dag( self, dag_maker, diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py index 33f4a5441d542..eacc77878ad7d 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dags.py @@ -109,7 +109,7 @@ def setup_dag_runs(self, *, session: Session = NEW_SESSION) -> None: ({"bundle_name": "wrong_bundle", "bundle_version": "some_commit_hash"}, [], 0), ], ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_should_return_200(self, test_client, query_params, expected_ids, expected_total_dag_runs): response = test_client.get("/dags", params=query_params) assert response.status_code == 200 @@ -135,7 +135,7 @@ def test_should_return_200(self, test_client, query_params, expected_ids, expect assert previous_run_after > dag_run["run_after"] previous_run_after = dag_run["run_after"] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") @pytest.mark.parametrize( ("query_params", "expected_ids"), [ @@ -149,7 +149,7 @@ def test_timetable_type_filter(self, test_client: TestClient, query_params, expe assert response.status_code == 200 assert [dag["dag_id"] for dag in response.json()["dags"]] == expected_ids - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") @pytest.mark.parametrize("state", ["queued", "running", "failed", "success"]) def test_dag_run_state_matches_any_run_not_only_latest(self, test_client, session, state): # Give DAG1 an older run in the probed state while its latest run ends in a @@ -175,7 +175,7 @@ def test_dag_run_state_matches_any_run_not_only_latest(self, test_client, sessio assert any_state.status_code == 200 assert [dag["dag_id"] for dag in any_state.json()["dags"]] == [DAG1_ID] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_last_and_any_run_state_filters_combined(self, test_client, session): # Regression: combining the last-run and any-run state filters must return the # intersection, not raise. The any-run EXISTS subquery must not correlate to the @@ -322,7 +322,7 @@ def restricted_is_authorized_dag(self, *, method, user, access_entity=None, deta response = test_client.get("/dags") assert response.status_code == 403 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_orders_latest_runs_by_run_after(self, test_client, session): running_run_after = pendulum.datetime(2026, 1, 1, tz="UTC") queued_run_after = pendulum.datetime(2026, 1, 2, tz="UTC") @@ -428,7 +428,7 @@ def test_get_dags_no_n_plus_one_queries(self, session, test_client): f"({first_query_count} → {second_query_count}), suggesting n+1 queries for tags" ) - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_latest_run_should_return_200(self, test_client): response = test_client.get(f"/dags/{DAG1_ID}/latest_run") assert response.status_code == 200 @@ -647,7 +647,7 @@ def seed_dag_runs(self, setup, session) -> None: ) session.commit() - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_returns_zero_filled_counts_per_requested_dag(self, test_client): response = test_client.get( "/dags/run_state_counts", @@ -686,7 +686,7 @@ def test_rejects_too_many_dag_ids(self, test_client): response = test_client.get("/dags/run_state_counts", params={"dag_ids": too_many}) assert response.status_code == 422 - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") def test_permission_filter_hides_disallowed_dags(self, test_client, session): # The test user is granted read on DAG1, DAG2, DAG4, DAG5 (DAG3 is paused/stale in # the parent fixture). Asking for a dag that exists but the caller cannot read @@ -704,7 +704,7 @@ def test_permission_filter_hides_disallowed_dags(self, test_client, session): dag_ids = [entry["dag_id"] for entry in response.json()["dags"]] assert dag_ids == [DAG1_ID] - @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle") + @pytest.mark.usefixtures("configure_dag_bundles_with_view_url") @mock.patch("airflow.api_fastapi.core_api.routes.ui.dags.STATE_COUNT_CAP", 3) def test_caps_counts_at_state_count_cap(self, test_client, session): # Push SUCCESS over the (patched) cap while FAILED stays under it: states cap independently.