This is an automated email from the ASF dual-hosted git repository.
henry3260 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 4c0ba9e41da Make Edge executor adoption test immune to executor_loader
reloads (#74122)
4c0ba9e41da is described below
commit 4c0ba9e41da99d583b1906bb7309e3057828daa1
Author: Shahar Epstein <[email protected]>
AuthorDate: Sat Oct 3 05:33:08 2026 +0300
Make Edge executor adoption test immune to executor_loader reloads (#74122)
---
.../edge3/tests/unit/edge3/executors/test_edge_executor.py | 12 +++++++++++-
1 file changed, 11 insertions(+), 1 deletion(-)
diff --git a/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
b/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
index bb636de8382..2f5ab17c0c2 100644
--- a/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
+++ b/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
@@ -51,6 +51,16 @@ if AIRFLOW_V_3_3_PLUS:
pytestmark = pytest.mark.db_test
+# Patch the ExecutorLoader class object bound at import time by the module
that calls init_executors()
+# (Job before Airflow 3.2, SchedulerJobRunner since), not the one currently in
+# airflow.executors.executor_loader: tests in other providers (e.g.
cncf.kubernetes, celery) reload() that
+# module, which replaces the class there, so a patch via the module path never
reaches the scheduler.
+SCHEDULER_EXECUTOR_LOADER = (
+ "airflow.jobs.scheduler_job_runner.ExecutorLoader"
+ if AIRFLOW_V_3_2_PLUS
+ else "airflow.jobs.job.ExecutorLoader"
+)
+
class TestEdgeExecutor:
@pytest.fixture(autouse=True)
@@ -440,7 +450,7 @@ class TestEdgeExecutor:
assert job in session
assert executor.running == {key}
-
@mock.patch("airflow.executors.executor_loader.ExecutorLoader.init_executors",
autospec=True)
+ @mock.patch(f"{SCHEDULER_EXECUTOR_LOADER}.init_executors", autospec=True)
def test_scheduler_restart_adopts_queued_edge_task(self,
mock_init_executors, dag_maker, session):
with dag_maker("test_dag", session=session):
EmptyOperator(task_id="test_task")