From 5e1de001072e49c2e00ecf8bbb452e3743ec22ef Mon Sep 17 00:00:00 2001 From: jlnav Date: Wed, 9 Sep 2026 10:32:50 -0500 Subject: [PATCH] decrease some of the exectuor_hworld and cancel_in_alloc pause durations. construct a FakePopen, _fake_launch, and _fake_cancel to test executor code without relying on flakiness of actual launches, and faster results. apply new fast_launch fixture (using these fakes) to all test_executor.py unit tests except those marked with a pytest.real_launch mark. --- libensemble/sim_funcs/executor_hworld.py | 6 +- .../test_cancel_in_alloc.py | 2 +- libensemble/tests/unit_tests/conftest.py | 154 ++++++++++++++++++ libensemble/tests/unit_tests/test_executor.py | 45 +++-- 4 files changed, 186 insertions(+), 21 deletions(-) diff --git a/libensemble/sim_funcs/executor_hworld.py b/libensemble/sim_funcs/executor_hworld.py index cb7e16a355..1ba638198b 100644 --- a/libensemble/sim_funcs/executor_hworld.py +++ b/libensemble/sim_funcs/executor_hworld.py @@ -88,8 +88,8 @@ def executor_hworld(H, _, sim_specs, info): print("sim_ended_count", sim_ended_count, flush=True) if ELAPSED_TIMEOUT: - args_for_sim = "sleep 60" # Manager kill - if signal received else completes - timeout = 65.0 + args_for_sim = "sleep 20" # Manager kill - if signal received else completes + timeout = 25.0 else: timeout = 6.0 @@ -104,7 +104,7 @@ def executor_hworld(H, _, sim_specs, info): args_for_sim = "sleep 1" # Should finish launch_shc = True elif sim_ended_count == 4: - args_for_sim = "sleep 8" # Worker kill on timeout + args_for_sim = "sleep 3" # Worker kill on timeout timeout = 1.0 elif sim_ended_count == 5: args_for_sim = "sleep 2 Fail" # Manager kill - if signal received else completes diff --git a/libensemble/tests/functionality_tests/test_cancel_in_alloc.py b/libensemble/tests/functionality_tests/test_cancel_in_alloc.py index 72953cdfd8..45e5f8ad28 100644 --- a/libensemble/tests/functionality_tests/test_cancel_in_alloc.py +++ b/libensemble/tests/functionality_tests/test_cancel_in_alloc.py @@ -37,7 +37,7 @@ "sim_f": sim_f, "in": ["x"], "out": [("f", float)], - "user": {"uniform_random_pause_ub": 10}, # long sleep ensures sims are still running when cancel fires + "user": {"uniform_random_pause_ub": 2}, # sleep ensures sims are still running when cancel fires } vocs = VOCS(variables={"x0": [-3, 3], "x1": [-2, 2]}, objectives={"f": "EXPLORE"}) diff --git a/libensemble/tests/unit_tests/conftest.py b/libensemble/tests/unit_tests/conftest.py index 5f6cc562e4..54ef5d6d3b 100644 --- a/libensemble/tests/unit_tests/conftest.py +++ b/libensemble/tests/unit_tests/conftest.py @@ -1,11 +1,162 @@ # https://stackoverflow.com/questions/47559524/pytest-how-to-skip-tests-unless-you-declare-an-option-flag/61193490#61193490 +import itertools +import os +import shutil +import subprocess +import time import warnings import pytest +import libensemble.utils.launcher as launcher + warnings.simplefilter("ignore", ResourceWarning) +# Speedup factor applied to simulated sim-app sleep durations by FakePopen. +# Chosen so that every kill/timeout ordering asserted by the tests is preserved +# (e.g. a "sleep 5" task still outlives a 0.5s wait() timeout). +FAST_LAUNCH_SPEEDUP = 5.0 + +# Simulated lifetime for tasks that must stay alive until killed (e.g. tasks +# whose stdout is polled for an "Error" marker before being killed). +FAKE_ERROR_TASK_DURATION = 10.0 + + +class FakePopen: + """Test double for :class:`subprocess.Popen` returned by the launcher. + + Emulates the ``my_simtask``/``my_serialtask`` argument protocol understood + by the compiled sim apps (``sleep [Error|Fail]``), scaled down by + ``FAST_LAUNCH_SPEEDUP`` so executor tests run fast and deterministically + without spawning real (MPI) subprocesses. + + Only :func:`libensemble.utils.launcher.launch` and + :func:`libensemble.utils.launcher.cancel` are patched; the remaining + launcher helpers (``wait``, ``killpg``, ``terminatepg``, ...) operate on + this object unchanged. + """ + + _pid_counter = itertools.count(100000) + + def __init__(self, cmd, stdout=None, stderr=None, **kwargs): + exe = cmd[0] + # Emulate Popen executable resolution so launch-failure paths + # (e.g. non-existent MPI runners exercising task retries) still work. + if "/" in exe or exe.startswith("."): + found = os.path.isfile(exe) + else: + found = shutil.which(exe) is not None + if not found: + raise FileNotFoundError(f"No such file or directory: {exe!r}") + self.cmd = list(cmd) + self.args = self.cmd[1:] + self.pid = next(FakePopen._pid_counter) + self.returncode = None + self._start = time.monotonic() + self._duration, self._exit_code = self._plan() + self._write_stdout(stdout) + + def _plan(self): + """Determine simulated (scaled) lifetime and exit code from app args.""" + cmd_str = " ".join(self.cmd) + if "c_startup" in cmd_str or "py_startup" in cmd_str: + return 0.0, 0 + duration = 3.0 # default sleep of the compiled sim apps + if "sleep" in self.args: + try: + duration = float(self.args[self.args.index("sleep") + 1]) + except (IndexError, ValueError): + pass + if "Error" in self.args: + # Must outlive any polling loop so the "Error" marker is always + # found in stdout while the task is still running. + return max(duration, FAKE_ERROR_TASK_DURATION) / FAST_LAUNCH_SPEEDUP, 0 + if "Fail" in self.args: + return duration / FAST_LAUNCH_SPEEDUP, 1 + return duration / FAST_LAUNCH_SPEEDUP, 0 + + def _write_stdout(self, stdout): + """Write the app output the tests expect to find in the stdout file.""" + if stdout is None: + return + cmd_str = " ".join(self.cmd) + if "c_startup" in cmd_str or "py_startup" in cmd_str: + stdout.write(f"{time.time()}\n") + stdout.flush() + return + if "my_simtask" in cmd_str or "my_serialtask" in cmd_str: + stdout.write("Hello world sleeping (simulated by FakePopen)\n") + if "Error" in self.args: + stdout.write("Oh Dear! An non-fatal Error seems to have occurred\n") + stdout.flush() + + def _complete(self): + if self.returncode is None: + self.returncode = self._exit_code + return self.returncode + + def poll(self): + """Emulate Popen.poll: None while running, returncode once complete.""" + if self.returncode is None and time.monotonic() - self._start >= self._duration: + self._complete() + return self.returncode + + def wait(self, timeout=None): + """Emulate Popen.wait, raising TimeoutExpired like the real thing.""" + if self.returncode is not None: + return self.returncode + remaining = self._duration - (time.monotonic() - self._start) + if remaining <= 0: + return self._complete() + if timeout is not None and remaining > timeout: + time.sleep(timeout) + if self.returncode is None and time.monotonic() - self._start < self._duration: + raise subprocess.TimeoutExpired(self.cmd, timeout) + return self._complete() + time.sleep(max(remaining, 0)) + return self._complete() + + def terminate(self): + if self.returncode is None: + self.returncode = -15 # as if SIGTERM was delivered + return self.returncode + + def kill(self): + if self.returncode is None: + self.returncode = -9 # as if SIGKILL was delivered + return self.returncode + + +def _fake_launch(cmd_template, specs=None, **kwargs): + """Drop-in replacement for launcher.launch returning a FakePopen.""" + cmd = launcher.form_command(cmd_template, specs) if specs is not None else cmd_template + return FakePopen(cmd, **kwargs) + + +def _fake_cancel(process, timeout=0): + """Drop-in replacement for launcher.cancel: terminate the fake at once.""" + if process.returncode is None: + process.terminate() + return process.wait() + + +@pytest.fixture +def fast_launch(request, monkeypatch): + """Replace real task subprocesses with fast, deterministic FakePopen. + + Patching :func:`libensemble.utils.launcher.launch` covers both the serial + :class:`Executor` and the :class:`MPIExecutor`, which both resolve it from + the shared launcher module at call time. Tests marked ``real_launch`` skip + the patching so a small set of genuine integration tests is retained. + """ + if request.node.get_closest_marker("real_launch"): + yield None + return + monkeypatch.setattr(launcher, "launch", _fake_launch) + monkeypatch.setattr(launcher, "cancel", _fake_cancel) + yield FakePopen + def pytest_addoption(parser): parser.addoption("--runextra", action="store_true", default=False, help="run extra tests") @@ -13,6 +164,9 @@ def pytest_addoption(parser): def pytest_configure(config): config.addinivalue_line("markers", "extra: mark test as extra to run") + config.addinivalue_line( + "markers", "real_launch: mark test as requiring real subprocess launches (skips fast_launch patching)" + ) def pytest_collection_modifyitems(config, items): diff --git a/libensemble/tests/unit_tests/test_executor.py b/libensemble/tests/unit_tests/test_executor.py index bf90464e84..e3160185b2 100644 --- a/libensemble/tests/unit_tests/test_executor.py +++ b/libensemble/tests/unit_tests/test_executor.py @@ -14,6 +14,11 @@ from libensemble.message_numbers import STOP_TAG, TASK_FAILED, UNSET_TAG from libensemble.resources.mpi_resources import MPIResourcesException +# Run all tests in this module with simulated task subprocesses (FakePopen via +# the fast_launch fixture) unless marked real_launch. This removes real MPI +# spawn overhead and wall-clock sleeps while exercising the executor logic. +pytestmark = pytest.mark.usefixtures("fast_launch") + NCORES = 1 build_sims = ["my_simtask.c", "my_serialtask.c", "c_startup.c"] @@ -166,7 +171,7 @@ def is_ompi(): # ----------------------------------------------------------------------------- # The following would typically be in the user sim_func. -def polling_loop(exctr, task, timeout_sec=2, delay=0.1): +def polling_loop(exctr, task, timeout_sec=1.0, delay=0.02): """Iterate over a loop, polling for an exit condition""" start = time.time() @@ -198,7 +203,7 @@ def polling_loop(exctr, task, timeout_sec=2, delay=0.1): return task -def polling_loop_multitask(exctr, task_list, timeout_sec=4.0, delay=0.1): +def polling_loop_multitask(exctr, task_list, timeout_sec=1.0, delay=0.02): """Iterate over a loop, polling for exit conditions on multiple tasks""" start = time.time() @@ -342,7 +347,7 @@ def test_kill_on_timeout(): cores = NCORES args_for_sim = "sleep 10" task = exctr.submit(calc_type="sim", num_procs=cores, app_args=args_for_sim) - task = polling_loop(exctr, task) + task = polling_loop(exctr, task, timeout_sec=0.5) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "USER_KILLED", "task.state should be USER_KILLED. Returned " + str(task.state) @@ -354,7 +359,7 @@ def test_kill_on_timeout_polling_loop_method(): cores = NCORES args_for_sim = "sleep 10" task = exctr.submit(calc_type="sim", num_procs=cores, app_args=args_for_sim) - exctr.polling_loop(task, timeout=1) + exctr.polling_loop(task, timeout=0.5) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "USER_KILLED", "task.state should be USER_KILLED. Returned " + str(task.state) @@ -425,7 +430,7 @@ def test_procs_and_machinefile_logic(): f.write(socket.gethostname() + "\n") task = exctr.submit(calc_type="sim", machinefile=machinefilename, app_args=args_for_sim) - task = polling_loop(exctr, task, timeout_sec=4, delay=0.1) + task = polling_loop(exctr, task, timeout_sec=0.5, delay=0.02) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FINISHED", "task.state should be FINISHED. Returned " + str(task.state) @@ -441,8 +446,8 @@ def test_procs_and_machinefile_logic(): ) else: task = exctr.submit(calc_type="sim", num_procs=6, num_nodes=2, procs_per_node=3, app_args=args_for_sim) - task = polling_loop(exctr, task, timeout_sec=4, delay=0.1) - time.sleep(0.25) + task = polling_loop(exctr, task, timeout_sec=0.5, delay=0.02) + time.sleep(0.05) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FINISHED", "task.state should be FINISHED. Returned " + str(task.state) @@ -466,8 +471,8 @@ def test_procs_and_machinefile_logic(): else: task = exctr.submit(calc_type="sim", num_nodes=2, procs_per_node=3, app_args=args_for_sim) assert 1 - task = polling_loop(exctr, task, timeout_sec=4, delay=0.1) - time.sleep(0.25) + task = polling_loop(exctr, task, timeout_sec=0.5, delay=0.02) + time.sleep(0.05) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FINISHED", "task.state should be FINISHED. Returned " + str(task.state) @@ -482,14 +487,14 @@ def test_procs_and_machinefile_logic(): # Testing no num_nodes (should not fail). task = exctr.submit(calc_type="sim", num_procs=2, procs_per_node=2, app_args=args_for_sim) assert 1 - task = polling_loop(exctr, task, timeout_sec=4, delay=0.1) + task = polling_loop(exctr, task, timeout_sec=0.5, delay=0.02) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FINISHED", "task.state should be FINISHED. Returned " + str(task.state) # Testing no procs_per_node (shouldn't fail) task = exctr.submit(calc_type="sim", num_nodes=1, num_procs=2, app_args=args_for_sim) assert 1 - task = polling_loop(exctr, task, timeout_sec=4, delay=0.1) + task = polling_loop(exctr, task, timeout_sec=0.5, delay=0.02) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FINISHED", "task.state should be FINISHED. Returned " + str(task.state) @@ -536,7 +541,7 @@ def test_finish_and_kill(): args_for_sim = "sleep 0.1" task = exctr.submit(calc_type="sim", num_procs=cores, app_args=args_for_sim) while not task.finished: - time.sleep(0.1) + time.sleep(0.02) task.poll() assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FINISHED", "task.state should be FINISHED. Returned " + str(task.state) @@ -676,7 +681,7 @@ def test_task_failure(): cores = NCORES args_for_sim = "sleep 1.0 Fail" task = exctr.submit(calc_type="sim", num_procs=cores, app_args=args_for_sim) - task = polling_loop(exctr, task, timeout_sec=3) + task = polling_loop(exctr, task, timeout_sec=0.5) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FAILED", "task.state should be FAILED. Returned " + str(task.state) @@ -711,7 +716,7 @@ def test_retries_launch_fail(): print(f"\nTest: {sys._getframe().f_code.co_name}\n") setup_executor_fakerunner() exctr = Executor.executor - exctr.retry_delay_incr = 0.05 + exctr.retry_delay_incr = 0.02 cores = NCORES args_for_sim = "sleep 0" task = exctr.submit(calc_type="sim", num_procs=cores, app_args=args_for_sim) @@ -724,11 +729,11 @@ def test_retries_before_polling_loop_method(): print(f"\nTest: {sys._getframe().f_code.co_name}\n") setup_executor_fakerunner() exctr = Executor.executor - exctr.retry_delay_incr = 0.05 + exctr.retry_delay_incr = 0.02 cores = NCORES args_for_sim = "sleep 0" task = exctr.submit(calc_type="sim", num_procs=cores, app_args=args_for_sim) - exctr.polling_loop(task, timeout=1) + exctr.polling_loop(task, timeout=0.5) assert task.finished, "task.finished should be True. Returned " + str(task.finished) assert task.state == "FAILED_TO_START", "task.state should be FAILED_TO_START. Returned " + str(task.state) assert task.run_attempts == 5, "task.run_attempts should be 5. Returned " + str(task.run_attempts) @@ -738,7 +743,7 @@ def test_retries_run_fail(): print(f"\nTest: {sys._getframe().f_code.co_name}\n") setup_executor() exctr = Executor.executor - exctr.retry_delay_incr = 0.05 + exctr.retry_delay_incr = 0.02 cores = NCORES args_for_sim = "sleep 0 Fail" task = exctr.submit(calc_type="sim", num_procs=cores, app_args=args_for_sim, wait_on_start=True) @@ -811,6 +816,7 @@ def test_serial_exe_exception(): pytest.fail("Expected exception") +@pytest.mark.real_launch def test_serial_exe_env_script(): env_script_path = os.path.join(os.getcwd(), "./env_script_in.sh") setup_serial_executor() @@ -835,6 +841,7 @@ def test_serial_exe_dryrun(): assert task.state == "FINISHED", "task.state should be FINISHED. Returned " + str(task.state) +@pytest.mark.real_launch def test_serial_startup_times(): print(f"\nTest: {sys._getframe().f_code.co_name}\n") setup_executor_startups() @@ -864,6 +871,8 @@ def test_serial_startup_times(): def test_futures_interface(): print(f"\nTest: {sys._getframe().f_code.co_name}\n") setup_executor() + exctr = Executor.executor + exctr.fail_time = 0.3 # Don't wait out full fail_time for early-failure check cores = NCORES args_for_sim = "sleep 3" with Executor.executor as exctr: @@ -877,6 +886,8 @@ def test_futures_interface(): def test_futures_interface_cancel(): print(f"\nTest: {sys._getframe().f_code.co_name}\n") setup_executor() + exctr = Executor.executor + exctr.fail_time = 0.3 # Don't wait out full fail_time for early-failure check cores = NCORES args_for_sim = "sleep 3" with Executor.executor as exctr: