From 71339c37a5d96024fb8b529583e004082605847d Mon Sep 17 00:00:00 2001 From: shudson Date: Tue, 28 Oct 2025 13:21:04 -0500 Subject: [PATCH 01/16] Add botorch_mfkg use-case. Works in batch. May need to update for async. --- .../persistent_botorch_mfkg_branin.py | 187 ++++++++++++++++++ libensemble/sim_funcs/augmented_branin.py | 38 ++++ .../run_botorch_mfkg_branin.py | 79 ++++++++ 3 files changed, 304 insertions(+) create mode 100644 libensemble/gen_funcs/persistent_botorch_mfkg_branin.py create mode 100644 libensemble/sim_funcs/augmented_branin.py create mode 100644 libensemble/tests/regression_tests/run_botorch_mfkg_branin.py diff --git a/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py new file mode 100644 index 000000000..5391fea12 --- /dev/null +++ b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py @@ -0,0 +1,187 @@ +""" +This file defines the persistent generator function for multi-fidelity Bayesian +optimization using BoTorch's Multi-Fidelity Knowledge Gradient (MFKG) acquisition function. + +This gen_f is meant to be used with the alloc_f function `only_persistent_gens`. +""" + +import numpy as np +import torch +from botorch import fit_gpytorch_mll +from botorch.acquisition import PosteriorMean +from botorch.acquisition.cost_aware import InverseCostWeightedUtility +from botorch.acquisition.fixed_feature import FixedFeatureAcquisitionFunction +from botorch.acquisition.knowledge_gradient import qMultiFidelityKnowledgeGradient +from botorch.acquisition.utils import project_to_target_fidelity +from botorch.models.cost import AffineFidelityCostModel +from botorch.models.gp_regression_fidelity import SingleTaskMultiFidelityGP +from botorch.models.transforms.outcome import Standardize +from botorch.optim.optimize import optimize_acqf, optimize_acqf_mixed +from gpytorch.mlls.exact_marginal_log_likelihood import ExactMarginalLogLikelihood + +from libensemble.message_numbers import EVAL_GEN_TAG, FINISHED_PERSISTENT_GEN_TAG, PERSIS_STOP, STOP_TAG +from libensemble.tools.persistent_support import PersistentSupport + +__all__ = ["persistent_botorch_mfkg"] + +# Set seeds for reproducibility +np.random.seed(42) +torch.manual_seed(42) + +# Torch settings +tkwargs = { + "dtype": torch.double, + "device": torch.device("cuda" if torch.cuda.is_available() else "cpu"), +} + +# Specify bounds +dim = 3 +bounds = torch.tensor([[0.0] * dim, [1.0] * dim], **tkwargs) + +# Specify target fidelity +target_fidelities = {2: 1.0} + +# Specify cost model +cost_model = AffineFidelityCostModel(fidelity_weights={2: 0.9}, fixed_cost=0.1) +cost_aware_utility = InverseCostWeightedUtility(cost_model=cost_model) + + +# Custom function to project posterior to target fidelity (defer to default) +def project(X): + return project_to_target_fidelity(X=X, target_fidelities=target_fidelities) + + +# Wrapper function for compatibility with existing code +def problem(X, ps, gen_specs): + """ + Wrapper to convert tensor input to numpy and send to libE for evaluation. + + Args: + X: tensor of shape (n, 3) where columns are [x0, x1, fidelity] + ps: PersistentSupport object for communication + gen_specs: Generator specifications + + Returns: + tensor of shape (n,) with objective values + """ + # Send points to be evaluated + X_np = X.cpu().numpy() + H_o = np.zeros(len(X), dtype=gen_specs["out"]) + H_o["x"] = X_np[:, :2] + H_o["fidelity"] = X_np[:, 2] + + tag, Work, calc_in = ps.send_recv(H_o) + + # Convert results back to tensor + if calc_in is None or len(calc_in) == 0: + return None, tag + + train_obj = torch.tensor(calc_in["f"], **tkwargs).unsqueeze(-1) + + return train_obj, tag + + +# Function to generate training data +def generate_initial_data(n, ps, gen_specs): # Jeff: Initial sample size is twice this value of n + train_x = torch.rand(n, 2, **tkwargs) + train_lf = torch.zeros(n, 1) + train_hf = torch.ones(n, 1) + train_x_full_lf = torch.cat((train_x, train_lf), dim=1) + train_x_full_hf = torch.cat((train_x, train_hf), dim=1) + train_x_full = torch.cat((train_x_full_lf, train_x_full_hf), dim=0) + train_obj, tag = problem(train_x_full, ps, gen_specs) + return train_x_full, train_obj, tag + + +# Function to initialize a botorch model +def initialize_model(train_x, train_obj): + model = SingleTaskMultiFidelityGP(train_x, train_obj, outcome_transform=Standardize(m=1), data_fidelities=[2]) + mll = ExactMarginalLogLikelihood(model.likelihood, model) + return mll, model + + +# Multifidelity Knowledge Gradient acquisition function +def get_mfkg(model): + + curr_val_acqf = FixedFeatureAcquisitionFunction( + acq_function=PosteriorMean(model), + d=3, + columns=[2], + values=[1], + ) + + _, current_value = optimize_acqf( + acq_function=curr_val_acqf, + bounds=bounds[:, :-1], + q=1, # Jeff: Don't adjust this for some reason + num_restarts=1, # Jeff: I decreased this to make libE development faster + raw_samples=10, # Jeff: I decreased this to make libE development faster + options={"batch_limit": 10, "maxiter": 10}, # Jeff: I decreased this to make libE development faster + ) + + return qMultiFidelityKnowledgeGradient( + model=model, + num_fantasies=128, + current_value=current_value, + cost_aware_utility=cost_aware_utility, + project=project, + ) + + +# Optimization step +def optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs): + # Generate new candidates + candidates, _ = optimize_acqf_mixed( + acq_function=mfkg_acqf, + bounds=bounds, + fixed_features_list=[{2: 0.0}, {2: 1.0}], + q=q, # Jeff: This is the number of new samples to make + num_restarts=1, # Jeff: I decreased this to make libE development faster + raw_samples=10, # Jeff: I decreased this to make libE development faster + options={"batch_limit": 10, "maxiter": 10}, # Jeff: I decreased this to make libE development faster + ) + + # Observe new values + cost = cost_model(candidates).sum() + new_x = candidates.detach() + new_obj, tag = problem(new_x, ps, gen_specs) + return new_x, new_obj, cost, tag + + +# Function to perform a single iteration +def do_iteration(train_x, train_obj, q, ps, gen_specs): + mll, model = initialize_model(train_x, train_obj) + fit_gpytorch_mll(mll) + mfkg_acqf = get_mfkg(model) + new_x, new_obj, _, tag = optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs) + + if new_obj is None: + return model, train_x, train_obj, tag + + train_x = torch.cat([train_x, new_x]) + train_obj = torch.cat([train_obj, new_obj]) # Jeff: This is where the "sim" evaluation happens, and needs to be communicated back to the manager + + return model, train_x, train_obj, tag + + +def persistent_botorch_mfkg(H, persis_info, gen_specs, libE_info): + """ + Persistent generator function for multi-fidelity Bayesian optimization using BoTorch's MFKG. + """ + ps = PersistentSupport(libE_info, EVAL_GEN_TAG) + + # Extract user parameters + ub = gen_specs["user"]["ub"] + lb = gen_specs["user"]["lb"] + n_init_samples = gen_specs["user"]["n_init_samples"] + q = gen_specs["user"]["q"] + + # ## Perform Multifidelity Bayesian Optimization + # Generate initial data + train_x, train_obj, tag = generate_initial_data(n_init_samples, ps, gen_specs) + + # Step + while tag not in [STOP_TAG, PERSIS_STOP]: + model, train_x, train_obj, tag = do_iteration(train_x, train_obj, q, ps, gen_specs) + + return None, persis_info, FINISHED_PERSISTENT_GEN_TAG diff --git a/libensemble/sim_funcs/augmented_branin.py b/libensemble/sim_funcs/augmented_branin.py new file mode 100644 index 000000000..5937db730 --- /dev/null +++ b/libensemble/sim_funcs/augmented_branin.py @@ -0,0 +1,38 @@ +""" +This module evaluates the augmented Branin function for multi-fidelity optimization. + +Augmented Branin is a modified version of the Branin function with a fidelity parameter. +""" + +__all__ = ["augmented_branin", "augmented_branin_func"] + +import math +import numpy as np + + +def augmented_branin(H, persis_info, sim_specs, libE_info): + """ + Evaluates the augmented Branin function for a collection of points given in ``H["x"]`` + with fidelity values in ``H["fidelity"]``. + """ + batch = len(H["x"]) + H_o = np.zeros(batch, dtype=sim_specs["out"]) + + for i in range(batch): + x = H["x"][i] + fidelity = H["fidelity"][i] + H_o["f"][i] = augmented_branin_func(x.reshape(1, -1), fidelity)[0] + + return H_o, persis_info + + +def augmented_branin_func(x, fidelity): + """Augmented Branin function for multi-fidelity optimization.""" + x0 = x[:, 0] + x1 = x[:, 1] + + t1 = 15 * x1 - (5.1 / (4 * math.pi**2) - 0.1 * (1 - fidelity)) * (15 * x0 - 5) ** 2 + 5 / math.pi * (15 * x0 - 5) - 6 + t2 = 10 * (1 - 1 / (8 * math.pi)) * np.cos(15 * x0 - 5) + result = t1**2 + t2 + 10 + + return -result # negate for maximization diff --git a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py new file mode 100644 index 000000000..b65eb13ca --- /dev/null +++ b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py @@ -0,0 +1,79 @@ +""" +Example of multi-fidelity optimization using a persistent BoTorch MFKG gen_func. + +This test uses the gen_on_manager option (persistent generator runs on +a thread). Therefore nworkers is the number of simulation workers. + +Execute via one of the following commands: + mpiexec -np 5 python run_botorch_mfkg_branin.py + python run_botorch_mfkg_branin.py --nworkers 4 + python run_botorch_mfkg_branin.py --nworkers 4 --comms tcp + +When running with the above commands, the number of concurrent evaluations of +the objective function will be 3. + +""" + +# Do not change these lines - they are parsed by run-tests.sh +# TESTSUITE_COMMS: local mpi +# TESTSUITE_NPROCS: 4 +# TESTSUITE_EXTRA: true +# TESTSUITE_OS_SKIP: OSX + +import numpy as np + +from libensemble import logger +from libensemble.alloc_funcs.start_only_persistent import only_persistent_gens +from libensemble.gen_funcs.persistent_botorch_mfkg_branin import persistent_botorch_mfkg +from libensemble.libE import libE +from libensemble.sim_funcs.augmented_branin import augmented_branin +from libensemble.tools import add_unique_random_streams, parse_args, save_libE_output + +# Main block is necessary only when using local comms with spawn start method (default on macOS and Windows). +if __name__ == "__main__": + nworkers, is_manager, libE_specs, _ = parse_args() + libE_specs["gen_on_manager"] = True + + sim_specs = { + "sim_f": augmented_branin, + "in": ["x", "fidelity"], + "out": [("f", float)], + } + + gen_specs = { + "gen_f": persistent_botorch_mfkg, + # "in": ["sim_id", "x", "f", "fidelity"], + "persis_in": ["sim_id", "x", "f", "fidelity"], + "out": [ + ("x", float, (2,)), + ("fidelity", float), + ], + "user": { + "lb": np.array([0.0, 0.0]), + "ub": np.array([1.0, 1.0]), + "n_init_samples": 4, + "q": 2, + }, + } + + alloc_specs = { + "alloc_f": only_persistent_gens, + "user": {"async_return": False}, + } + + # libE logger + logger.set_level("INFO") + + # Exit criteria + exit_criteria = {"sim_max": 12} # Exit after running sim_max simulations + + # Create a different random number stream for each worker and the manager + persis_info = add_unique_random_streams({}, nworkers + 1) + + # Run LibEnsemble, and store results in history array H + H, persis_info, flag = libE(sim_specs, gen_specs, exit_criteria, persis_info, alloc_specs, libE_specs) + + # Save results to numpy file + if is_manager: + save_libE_output(H, persis_info, __file__, nworkers) + From 50ee0dfab773ba97f8dd0aa8c2d955e48e458f39 Mon Sep 17 00:00:00 2001 From: Jeffrey Larson Date: Wed, 29 Oct 2025 10:33:32 -0500 Subject: [PATCH 02/16] Noting Ax must be 0.5.0 --- libensemble/gen_funcs/persistent_ax_multitask.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libensemble/gen_funcs/persistent_ax_multitask.py b/libensemble/gen_funcs/persistent_ax_multitask.py index 0a2e07f20..30ce3634e 100644 --- a/libensemble/gen_funcs/persistent_ax_multitask.py +++ b/libensemble/gen_funcs/persistent_ax_multitask.py @@ -8,7 +8,7 @@ This `gen_f` is meant to be used with the `alloc_f` function `only_persistent_gens` -Requires: Ax>=0.5.0 +Requires: Ax==0.5.0 Ax notes: Each arm = a set of simulation inputs (a sim_id) From a041af921291c33d8cab9042d5f08ea1bca97866 Mon Sep 17 00:00:00 2001 From: Jeffrey Larson Date: Fri, 5 Dec 2025 08:48:59 -0600 Subject: [PATCH 03/16] adding note --- libensemble/tests/regression_tests/run_botorch_mfkg_branin.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py index b65eb13ca..be599ae70 100644 --- a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py +++ b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py @@ -51,7 +51,7 @@ "user": { "lb": np.array([0.0, 0.0]), "ub": np.array([1.0, 1.0]), - "n_init_samples": 4, + "n_init_samples": 4, # Each of these points will have a high-fidelity and low-fidelity evaluation "q": 2, }, } From 2ff60bf73a960399a68338503102f44c2ce8c4a6 Mon Sep 17 00:00:00 2001 From: Jeffrey Larson Date: Tue, 7 Jul 2026 15:20:04 -0500 Subject: [PATCH 04/16] black --- .../persistent_botorch_mfkg_branin.py | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py index 5391fea12..5e52d6a00 100644 --- a/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py +++ b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py @@ -55,12 +55,12 @@ def project(X): def problem(X, ps, gen_specs): """ Wrapper to convert tensor input to numpy and send to libE for evaluation. - + Args: X: tensor of shape (n, 3) where columns are [x0, x1, fidelity] ps: PersistentSupport object for communication gen_specs: Generator specifications - + Returns: tensor of shape (n,) with objective values """ @@ -69,15 +69,15 @@ def problem(X, ps, gen_specs): H_o = np.zeros(len(X), dtype=gen_specs["out"]) H_o["x"] = X_np[:, :2] H_o["fidelity"] = X_np[:, 2] - + tag, Work, calc_in = ps.send_recv(H_o) - + # Convert results back to tensor if calc_in is None or len(calc_in) == 0: return None, tag - + train_obj = torch.tensor(calc_in["f"], **tkwargs).unsqueeze(-1) - + return train_obj, tag @@ -154,12 +154,12 @@ def do_iteration(train_x, train_obj, q, ps, gen_specs): fit_gpytorch_mll(mll) mfkg_acqf = get_mfkg(model) new_x, new_obj, _, tag = optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs) - + if new_obj is None: return model, train_x, train_obj, tag - + train_x = torch.cat([train_x, new_x]) - train_obj = torch.cat([train_obj, new_obj]) # Jeff: This is where the "sim" evaluation happens, and needs to be communicated back to the manager + train_obj = torch.cat([train_obj, new_obj]) # Jeff: This is where the "sim" eval happens, to be send the manager return model, train_x, train_obj, tag From 8ff42c3aa4075528e12a19f0d6196ca024a44f24 Mon Sep 17 00:00:00 2001 From: Jeffrey Larson Date: Tue, 7 Jul 2026 15:30:34 -0500 Subject: [PATCH 05/16] black --- .../persistent_botorch_mfkg_branin.py | 28 +++++++++---------- libensemble/sim_funcs/augmented_branin.py | 5 ++-- .../run_botorch_mfkg_branin.py | 3 +- 3 files changed, 16 insertions(+), 20 deletions(-) diff --git a/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py index 5e52d6a00..5b592439d 100644 --- a/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py +++ b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py @@ -34,10 +34,6 @@ "device": torch.device("cuda" if torch.cuda.is_available() else "cpu"), } -# Specify bounds -dim = 3 -bounds = torch.tensor([[0.0] * dim, [1.0] * dim], **tkwargs) - # Specify target fidelity target_fidelities = {2: 1.0} @@ -82,8 +78,8 @@ def problem(X, ps, gen_specs): # Function to generate training data -def generate_initial_data(n, ps, gen_specs): # Jeff: Initial sample size is twice this value of n - train_x = torch.rand(n, 2, **tkwargs) +def generate_initial_data(n, ps, gen_specs, bounds): # Jeff: Initial sample size is twice this value of n + train_x = bounds[0, :-1] + (bounds[1, :-1] - bounds[0, :-1]) * torch.rand(n, 2, **tkwargs) train_lf = torch.zeros(n, 1) train_hf = torch.ones(n, 1) train_x_full_lf = torch.cat((train_x, train_lf), dim=1) @@ -101,7 +97,7 @@ def initialize_model(train_x, train_obj): # Multifidelity Knowledge Gradient acquisition function -def get_mfkg(model): +def get_mfkg(model, bounds): curr_val_acqf = FixedFeatureAcquisitionFunction( acq_function=PosteriorMean(model), @@ -129,7 +125,7 @@ def get_mfkg(model): # Optimization step -def optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs): +def optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs, bounds): # Generate new candidates candidates, _ = optimize_acqf_mixed( acq_function=mfkg_acqf, @@ -149,11 +145,11 @@ def optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs): # Function to perform a single iteration -def do_iteration(train_x, train_obj, q, ps, gen_specs): +def do_iteration(train_x, train_obj, q, ps, gen_specs, bounds): mll, model = initialize_model(train_x, train_obj) fit_gpytorch_mll(mll) - mfkg_acqf = get_mfkg(model) - new_x, new_obj, _, tag = optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs) + mfkg_acqf = get_mfkg(model, bounds) + new_x, new_obj, _, tag = optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs, bounds) if new_obj is None: return model, train_x, train_obj, tag @@ -171,17 +167,19 @@ def persistent_botorch_mfkg(H, persis_info, gen_specs, libE_info): ps = PersistentSupport(libE_info, EVAL_GEN_TAG) # Extract user parameters - ub = gen_specs["user"]["ub"] - lb = gen_specs["user"]["lb"] + ub = np.asarray(gen_specs["user"]["ub"], dtype=float) + lb = np.asarray(gen_specs["user"]["lb"], dtype=float) + bounds = torch.tensor([np.append(lb, 0.0), np.append(ub, 1.0)], **tkwargs) + n_init_samples = gen_specs["user"]["n_init_samples"] q = gen_specs["user"]["q"] # ## Perform Multifidelity Bayesian Optimization # Generate initial data - train_x, train_obj, tag = generate_initial_data(n_init_samples, ps, gen_specs) + train_x, train_obj, tag = generate_initial_data(n_init_samples, ps, gen_specs, bounds) # Step while tag not in [STOP_TAG, PERSIS_STOP]: - model, train_x, train_obj, tag = do_iteration(train_x, train_obj, q, ps, gen_specs) + model, train_x, train_obj, tag = do_iteration(train_x, train_obj, q, ps, gen_specs, bounds) return None, persis_info, FINISHED_PERSISTENT_GEN_TAG diff --git a/libensemble/sim_funcs/augmented_branin.py b/libensemble/sim_funcs/augmented_branin.py index 5937db730..8aa1300cd 100644 --- a/libensemble/sim_funcs/augmented_branin.py +++ b/libensemble/sim_funcs/augmented_branin.py @@ -6,7 +6,6 @@ __all__ = ["augmented_branin", "augmented_branin_func"] -import math import numpy as np @@ -31,8 +30,8 @@ def augmented_branin_func(x, fidelity): x0 = x[:, 0] x1 = x[:, 1] - t1 = 15 * x1 - (5.1 / (4 * math.pi**2) - 0.1 * (1 - fidelity)) * (15 * x0 - 5) ** 2 + 5 / math.pi * (15 * x0 - 5) - 6 - t2 = 10 * (1 - 1 / (8 * math.pi)) * np.cos(15 * x0 - 5) + t1 = 15 * x1 - (5.1 / (4 * np.pi**2) - 0.1 * (1 - fidelity)) * (15 * x0 - 5) ** 2 + 5 / np.pi * (15 * x0 - 5) - 6 + t2 = 10 * (1 - 1 / (8 * np.pi)) * np.cos(15 * x0 - 5) result = t1**2 + t2 + 10 return -result # negate for maximization diff --git a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py index be599ae70..506bbae9e 100644 --- a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py +++ b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py @@ -51,7 +51,7 @@ "user": { "lb": np.array([0.0, 0.0]), "ub": np.array([1.0, 1.0]), - "n_init_samples": 4, # Each of these points will have a high-fidelity and low-fidelity evaluation + "n_init_samples": 4, # Each of these points will have a high-fidelity and low-fidelity evaluation "q": 2, }, } @@ -76,4 +76,3 @@ # Save results to numpy file if is_manager: save_libE_output(H, persis_info, __file__, nworkers) - From ea1d255f1483f41715f9ce363e869e037490e841 Mon Sep 17 00:00:00 2001 From: Jeffrey Larson Date: Wed, 8 Jul 2026 07:37:31 -0500 Subject: [PATCH 06/16] Removing old options/imports --- .../tests/regression_tests/run_botorch_mfkg_branin.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py index 506bbae9e..2e220b968 100644 --- a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py +++ b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py @@ -27,12 +27,11 @@ from libensemble.gen_funcs.persistent_botorch_mfkg_branin import persistent_botorch_mfkg from libensemble.libE import libE from libensemble.sim_funcs.augmented_branin import augmented_branin -from libensemble.tools import add_unique_random_streams, parse_args, save_libE_output +from libensemble.tools import parse_args, save_libE_output # Main block is necessary only when using local comms with spawn start method (default on macOS and Windows). if __name__ == "__main__": nworkers, is_manager, libE_specs, _ = parse_args() - libE_specs["gen_on_manager"] = True sim_specs = { "sim_f": augmented_branin, @@ -68,7 +67,7 @@ exit_criteria = {"sim_max": 12} # Exit after running sim_max simulations # Create a different random number stream for each worker and the manager - persis_info = add_unique_random_streams({}, nworkers + 1) + persis_info = {} # Run LibEnsemble, and store results in history array H H, persis_info, flag = libE(sim_specs, gen_specs, exit_criteria, persis_info, alloc_specs, libE_specs) From dae7ff5e1e9883fdafb45c31ebc5ef56e13e30ad Mon Sep 17 00:00:00 2001 From: Jeffrey Larson Date: Thu, 16 Jul 2026 11:16:50 -0500 Subject: [PATCH 07/16] Hoping to get the async case tested --- .../persistent_botorch_mfkg_branin.py | 56 +++++-- .../run_botorch_mfkg_branin.py | 77 ---------- .../regression_tests/test_mfkg_branin.py | 141 ++++++++++++++++++ 3 files changed, 182 insertions(+), 92 deletions(-) delete mode 100644 libensemble/tests/regression_tests/run_botorch_mfkg_branin.py create mode 100644 libensemble/tests/regression_tests/test_mfkg_branin.py diff --git a/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py index 5b592439d..2a4081c06 100644 --- a/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py +++ b/libensemble/gen_funcs/persistent_botorch_mfkg_branin.py @@ -50,7 +50,15 @@ def project(X): # Wrapper function for compatibility with existing code def problem(X, ps, gen_specs): """ - Wrapper to convert tensor input to numpy and send to libE for evaluation. + Send a batch of points to libE for evaluation and collect the results. + + Results are gathered with a ``recv`` loop rather than a single ``send_recv`` + so this works in both batch return mode (``async_return`` off, where the + whole batch arrives in one message) and asynchronous return mode + (``async_return`` on with ``active_recv_gen``, where results stream back one + or a few at a time, possibly out of order). The evaluated points are + reconstructed from the returned rows so the (x, fidelity) -> f pairing stays + correct regardless of the order results arrive in. Args: X: tensor of shape (n, 3) where columns are [x0, x1, fidelity] @@ -58,23 +66,39 @@ def problem(X, ps, gen_specs): gen_specs: Generator specifications Returns: - tensor of shape (n,) with objective values + (new_x, new_obj, tag): tensor of the evaluated points (shape (m, 3)), + tensor of their objective values (shape (m, 1)), and the last message + tag. new_x/new_obj are None if no results were received (e.g. on stop). """ - # Send points to be evaluated + # Send points to be evaluated. X_np = X.cpu().numpy() H_o = np.zeros(len(X), dtype=gen_specs["out"]) H_o["x"] = X_np[:, :2] H_o["fidelity"] = X_np[:, 2] + ps.send(H_o) - tag, Work, calc_in = ps.send_recv(H_o) + # Collect results until the whole batch is back (or we are told to stop). + n_expected = len(X) + xs = [] + objs = [] + tag = None + while len(objs) < n_expected: + tag, Work, calc_in = ps.recv() + if tag in [STOP_TAG, PERSIS_STOP]: + break + if calc_in is None or len(calc_in) == 0: + continue + for row in calc_in: + xs.append(np.append(row["x"], row["fidelity"])) + objs.append(float(row["f"])) - # Convert results back to tensor - if calc_in is None or len(calc_in) == 0: - return None, tag + if len(objs) == 0: + return None, None, tag - train_obj = torch.tensor(calc_in["f"], **tkwargs).unsqueeze(-1) + new_x = torch.tensor(np.array(xs), **tkwargs) + new_obj = torch.tensor(np.array(objs), **tkwargs).unsqueeze(-1) - return train_obj, tag + return new_x, new_obj, tag # Function to generate training data @@ -85,8 +109,10 @@ def generate_initial_data(n, ps, gen_specs, bounds): # Jeff: Initial sample siz train_x_full_lf = torch.cat((train_x, train_lf), dim=1) train_x_full_hf = torch.cat((train_x, train_hf), dim=1) train_x_full = torch.cat((train_x_full_lf, train_x_full_hf), dim=0) - train_obj, tag = problem(train_x_full, ps, gen_specs) - return train_x_full, train_obj, tag + # problem returns the evaluated points (reconstructed from the results) so + # that x and objective stay paired even when results arrive out of order. + new_x, train_obj, tag = problem(train_x_full, ps, gen_specs) + return new_x, train_obj, tag # Function to initialize a botorch model @@ -137,10 +163,10 @@ def optimize_mfkg_and_get_observation(mfkg_acqf, q, ps, gen_specs, bounds): options={"batch_limit": 10, "maxiter": 10}, # Jeff: I decreased this to make libE development faster ) - # Observe new values + # Observe new values. problem returns the actually-evaluated points paired + # with their objectives (order-independent), which we use for training. cost = cost_model(candidates).sum() - new_x = candidates.detach() - new_obj, tag = problem(new_x, ps, gen_specs) + new_x, new_obj, tag = problem(candidates.detach(), ps, gen_specs) return new_x, new_obj, cost, tag @@ -169,7 +195,7 @@ def persistent_botorch_mfkg(H, persis_info, gen_specs, libE_info): # Extract user parameters ub = np.asarray(gen_specs["user"]["ub"], dtype=float) lb = np.asarray(gen_specs["user"]["lb"], dtype=float) - bounds = torch.tensor([np.append(lb, 0.0), np.append(ub, 1.0)], **tkwargs) + bounds = torch.tensor(np.array([np.append(lb, 0.0), np.append(ub, 1.0)]), **tkwargs) n_init_samples = gen_specs["user"]["n_init_samples"] q = gen_specs["user"]["q"] diff --git a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py b/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py deleted file mode 100644 index 2e220b968..000000000 --- a/libensemble/tests/regression_tests/run_botorch_mfkg_branin.py +++ /dev/null @@ -1,77 +0,0 @@ -""" -Example of multi-fidelity optimization using a persistent BoTorch MFKG gen_func. - -This test uses the gen_on_manager option (persistent generator runs on -a thread). Therefore nworkers is the number of simulation workers. - -Execute via one of the following commands: - mpiexec -np 5 python run_botorch_mfkg_branin.py - python run_botorch_mfkg_branin.py --nworkers 4 - python run_botorch_mfkg_branin.py --nworkers 4 --comms tcp - -When running with the above commands, the number of concurrent evaluations of -the objective function will be 3. - -""" - -# Do not change these lines - they are parsed by run-tests.sh -# TESTSUITE_COMMS: local mpi -# TESTSUITE_NPROCS: 4 -# TESTSUITE_EXTRA: true -# TESTSUITE_OS_SKIP: OSX - -import numpy as np - -from libensemble import logger -from libensemble.alloc_funcs.start_only_persistent import only_persistent_gens -from libensemble.gen_funcs.persistent_botorch_mfkg_branin import persistent_botorch_mfkg -from libensemble.libE import libE -from libensemble.sim_funcs.augmented_branin import augmented_branin -from libensemble.tools import parse_args, save_libE_output - -# Main block is necessary only when using local comms with spawn start method (default on macOS and Windows). -if __name__ == "__main__": - nworkers, is_manager, libE_specs, _ = parse_args() - - sim_specs = { - "sim_f": augmented_branin, - "in": ["x", "fidelity"], - "out": [("f", float)], - } - - gen_specs = { - "gen_f": persistent_botorch_mfkg, - # "in": ["sim_id", "x", "f", "fidelity"], - "persis_in": ["sim_id", "x", "f", "fidelity"], - "out": [ - ("x", float, (2,)), - ("fidelity", float), - ], - "user": { - "lb": np.array([0.0, 0.0]), - "ub": np.array([1.0, 1.0]), - "n_init_samples": 4, # Each of these points will have a high-fidelity and low-fidelity evaluation - "q": 2, - }, - } - - alloc_specs = { - "alloc_f": only_persistent_gens, - "user": {"async_return": False}, - } - - # libE logger - logger.set_level("INFO") - - # Exit criteria - exit_criteria = {"sim_max": 12} # Exit after running sim_max simulations - - # Create a different random number stream for each worker and the manager - persis_info = {} - - # Run LibEnsemble, and store results in history array H - H, persis_info, flag = libE(sim_specs, gen_specs, exit_criteria, persis_info, alloc_specs, libE_specs) - - # Save results to numpy file - if is_manager: - save_libE_output(H, persis_info, __file__, nworkers) diff --git a/libensemble/tests/regression_tests/test_mfkg_branin.py b/libensemble/tests/regression_tests/test_mfkg_branin.py new file mode 100644 index 000000000..d28dc0bc1 --- /dev/null +++ b/libensemble/tests/regression_tests/test_mfkg_branin.py @@ -0,0 +1,141 @@ +""" +Example of multi-fidelity optimization using a persistent BoTorch MFKG gen_func. + +One worker runs the persistent MFKG generator; the remaining workers evaluate +the (augmented Branin) objective. The generator collects a full batch of ``q`` +points before refitting its GP model, so a batch size of ``q`` keeps up to +``q`` simulation workers busy at once. + +This script runs the ensemble twice to exercise both ways of returning results +to the generator: + +* ``async_return=False`` (batch): the whole batch of ``q`` evaluations is + returned to the generator in one message once all have completed. +* ``async_return=True`` with ``active_recv_gen=True``: each evaluation is + handed back as soon as it finishes (possibly out of order); the generator + collects the batch as results stream in. + +Execute via one of the following commands: + mpiexec -np 5 python run_botorch_mfkg_branin.py + python run_botorch_mfkg_branin.py --nworkers 4 + python run_botorch_mfkg_branin.py --nworkers 4 --comms tcp + +With ``--nworkers 4`` one worker is the generator and three concurrently +evaluate the objective. + +""" + +# Do not change these lines - they are parsed by run-tests.sh +# TESTSUITE_COMMS: local mpi +# TESTSUITE_NPROCS: 4 +# TESTSUITE_EXTRA: true +# TESTSUITE_OS_SKIP: OSX + +import numpy as np + +from libensemble import logger +from libensemble.alloc_funcs.start_only_persistent import only_persistent_gens +from libensemble.gen_funcs.persistent_botorch_mfkg_branin import persistent_botorch_mfkg +from libensemble.libE import libE +from libensemble.sim_funcs.augmented_branin import augmented_branin +from libensemble.tools import parse_args, save_libE_output + +LB = np.array([0.0, 0.0]) +UB = np.array([1.0, 1.0]) +N_INIT_SAMPLES = 4 # Each initial point gets a high- and a low-fidelity evaluation +Q = 4 # Batch size per MFKG iteration (keeps up to q sim workers busy) +SIM_MAX = 16 # Exit after running this many simulations + + +def run_mfkg(async_return, nworkers, is_manager, libE_specs): + """Run the MFKG ensemble once and check the results. + + ``async_return`` selects batch return (False) or asynchronous return + (True, with the generator in active-receive mode). + """ + sim_specs = { + "sim_f": augmented_branin, + "in": ["x", "fidelity"], + "out": [("f", float)], + } + + gen_specs = { + "gen_f": persistent_botorch_mfkg, + "persis_in": ["sim_id", "x", "f", "fidelity"], + "out": [ + ("x", float, (2,)), + ("fidelity", float), + ], + "user": { + "lb": LB, + "ub": UB, + "n_init_samples": N_INIT_SAMPLES, + "q": Q, + }, + } + + alloc_specs = { + "alloc_f": only_persistent_gens, + "user": { + # When async, return each evaluation to the generator as soon as it + # completes and let the generator receive while still active. When + # not async, the whole batch is returned in one message. + "async_return": async_return, + "active_recv_gen": async_return, + }, + } + + exit_criteria = {"sim_max": SIM_MAX} + + # Fresh RNG streams for each run. + persis_info = {} + + H, persis_info, flag = libE(sim_specs, gen_specs, exit_criteria, persis_info, alloc_specs, libE_specs) + + if is_manager: + save_libE_output(H, persis_info, __file__, nworkers) + + mode = "async" if async_return else "batch" + + # -- Sanity checks that the run behaved as expected -------------------- + ended = H["sim_ended"] + n_ended = int(np.sum(ended)) + + # libE finished cleanly. + assert flag == 0, f"[{mode}] libE returned a nonzero exit flag: {flag}" + + # We completed at least sim_max evaluations, and the MFKG loop produced + # points beyond the (2 * n_init_samples) initial sample. + assert n_ended >= SIM_MAX, f"[{mode}] Expected >= {SIM_MAX} completed sims, got {n_ended}" + assert n_ended > 2 * N_INIT_SAMPLES, f"[{mode}] MFKG did not generate points beyond the initial sample" + + # Every completed evaluation returned a finite objective value. + assert np.all(np.isfinite(H["f"][ended])), f"[{mode}] Found non-finite objective value(s)" + + # The generator only ever requests the two discrete fidelities (low=0, high=1). + fids = np.round(H["fidelity"][ended], 6) + assert np.all(np.isin(fids, [0.0, 1.0])), f"[{mode}] Unexpected fidelity values: {np.unique(fids)}" + + # Generated points stayed within [lb, ub]. + xs = H["x"][ended] + assert np.all(xs >= LB - 1e-9) and np.all(xs <= UB + 1e-9), f"[{mode}] Generated point(s) outside the bounds" + + # augmented_branin is negated for maximization; the best attainable value is + # ~ -0.397887 (the negated Branin global minimum). We must never exceed it. + best_f = float(np.max(H["f"][ended])) + assert best_f <= -0.397887 + 1e-3, f"[{mode}] Objective exceeded the known optimum: {best_f}" + assert best_f < 0.0, f"[{mode}] Best objective is unexpectedly non-negative: {best_f}" + + print(f"[{mode}] assertions passed: completed sims: {n_ended}, best f: {best_f:.4f}") + + +# Main block is necessary only when using local comms with spawn start method (default on macOS and Windows). +if __name__ == "__main__": + nworkers, is_manager, libE_specs, _ = parse_args() + + # libE logger + logger.set_level("INFO") + + # Exercise both batch and asynchronous return of results to the generator. + for async_return in [False, True]: + run_mfkg(async_return, nworkers, is_manager, libE_specs) From 282683bb102309f3b3415ed5dbb3f9249f48ef64 Mon Sep 17 00:00:00 2001 From: jlnav Date: Thu, 2 Jul 2026 16:11:50 -0500 Subject: [PATCH 08/16] globus compute executor changes/tests moved onto this specific branch --- libensemble/executors/__init__.py | 5 +- .../executors/globus_compute_executor.py | 279 ++++++++++++++++++ .../tests/unit_tests/test_globus_compute.py | 258 ++++++++++++++++ libensemble/utils/globus_compute.py | 87 ++++++ 4 files changed, 627 insertions(+), 2 deletions(-) create mode 100644 libensemble/executors/globus_compute_executor.py create mode 100644 libensemble/tests/unit_tests/test_globus_compute.py create mode 100644 libensemble/utils/globus_compute.py diff --git a/libensemble/executors/__init__.py b/libensemble/executors/__init__.py index 13c7d851d..4c4f2f2af 100644 --- a/libensemble/executors/__init__.py +++ b/libensemble/executors/__init__.py @@ -1,11 +1,12 @@ from libensemble.executors.executor import Executor +from libensemble.executors.globus_compute_executor import GlobusComputeExecutor, GlobusComputeTask from libensemble.executors.mpi_executor import MPIExecutor # FluxExecutor is optional - requires flux-core Python bindings try: from libensemble.executors.flux_executor import FluxExecutor # noqa: F401 - __all__ = ["Executor", "MPIExecutor", "FluxExecutor"] + __all__ = ["Executor", "GlobusComputeExecutor", "GlobusComputeTask", "MPIExecutor", "FluxExecutor"] except ImportError: # flux-core not available - FluxExecutor won't be importable - __all__ = ["Executor", "MPIExecutor"] + __all__ = ["Executor", "GlobusComputeExecutor", "GlobusComputeTask", "MPIExecutor"] diff --git a/libensemble/executors/globus_compute_executor.py b/libensemble/executors/globus_compute_executor.py new file mode 100644 index 000000000..6a2c718dc --- /dev/null +++ b/libensemble/executors/globus_compute_executor.py @@ -0,0 +1,279 @@ +import logging +import os +from concurrent.futures import Future, TimeoutError +from typing import Any + +from libensemble.executors.executor import Application, Executor, ExecutorException, Task, TimeoutExpired +from libensemble.utils.globus_compute import GCSession +from libensemble.utils.timer import TaskTimer + +logger = logging.getLogger(__name__) + + +class GlobusComputeTask(Task): + """A :class:`~libensemble.executors.executor.Task` wrapping a + ``concurrent.futures.Future`` returned by Globus Compute. + + Instead of managing a local subprocess, this task polls a remote + computation via the future's ``done()`` / ``result()`` APIs. + """ + + def __init__(self, future, app=None, app_args=None, workerid=None): + self.id = next(Task.newid) + self.reset() + self.timer = TaskTimer() + self.app = app + self.app_args = app_args + self.workerID = workerid + self._gc_future = future + + worker_name = f"_worker{self.workerID}" if self.workerID else "" + self.name = Task.prefix + f"_{app.name}{worker_name}_{self.id}" + self.stdout = "" + self.stderr = "" + self.workdir = None + self.dry_run = False + self.runline = None + self.run_attempts = 0 + self.env = {} + self.ngpus_req = 0 + + self.state = "RUNNING" + self.timer.start() + self.submit_time = self.timer.tstart + + def _check_poll(self): + if self.finished: + return False + return True + + def poll(self): + if not self._check_poll(): + return + if self._gc_future.done(): + try: + self._gc_future.result() + self.finished = True + self.success = True + self.state = "FINISHED" + except Exception: + self.finished = True + self.success = False + self.state = "FAILED" + self.calc_task_timing() + else: + self.state = "RUNNING" + self.runtime = self.timer.elapsed + + def wait(self, timeout=None): + if not self._check_poll(): + return + try: + self._gc_future.result(timeout=timeout) + self.finished = True + self.success = True + self.state = "FINISHED" + except TimeoutError: + raise TimeoutExpired(self.name, timeout) + except Exception: + self.finished = True + self.success = False + self.state = "FAILED" + self.calc_task_timing() + + def kill(self, wait_time=None): + self._gc_future.cancel() + self.state = "USER_KILLED" + self.finished = True + self.calc_task_timing() + + def result(self, timeout=None): + self.wait(timeout=timeout) + return self.state + + def running(self): + self.poll() + return self.state == "RUNNING" + + def done(self): + self.poll() + return self.finished + + def cancelled(self): + self.poll() + return self.state == "USER_KILLED" + + +class GlobusComputeExecutor(Executor): + """An :class:`~libensemble.executors.executor.Executor` that submits + Python callables to Globus Compute instead of launching local subprocesses. + + Usage in a top-level script:: + + from libensemble.executors.globus_compute_executor import GlobusComputeExecutor + + exctr = GlobusComputeExecutor(endpoint_id="...") + + Inside a simulator function:: + + task = info["executor"].submit(func=my_remote_func, app_args=...) + while not task.finished: + task.poll() + if info["executor"].manager_kill_received(): + task.kill() + break + time.sleep(0.1) + """ + + def __init__(self, endpoint_id: str): + self.manager_signal = None + self.default_apps: dict[str, Application | None] = {"sim": None, "gen": None} + self.apps: dict[str, Application] = {} + self.wait_time = 60 + self.list_of_tasks: list[GlobusComputeTask] = [] + self.workerID = None + self.comm = None + self.last_task = 0 + self.base_dir = os.getcwd() + + self.endpoint_id = endpoint_id + self._gc_executor = None + self._func_cache: dict[int, str] = {} + + def _ensure_gc(self): + if self._gc_executor is None: + self._gc_executor = GCSession.get_or_create_executor(self.endpoint_id) + return self._gc_executor + + def _get_func_id(self, func) -> str: + key = id(func) + if key in self._func_cache: + return self._func_cache[key] + executor = self._ensure_gc() + if executor is None: + raise RuntimeError( + "Globus Compute SDK is not installed. " "Install it with: pip install globus-compute-sdk" + ) + fid = executor.register_function(func) + self._func_cache[key] = fid + return fid + + def register_app( + self, + full_path: str, + app_name: str | None = None, + calc_type: str | None = None, + desc: str | None = None, + precedent: str = "", + pyobj: Any | None = None, + ) -> None: + """Register an application. + + If *pyobj* is provided the application is treated as a remote + Python callable. Otherwise the base-class behaviour applies + (local executable). + """ + if not app_name: + app_name = os.path.split(full_path)[1] + + app = Application(full_path, app_name, calc_type, desc, pyobj, precedent) + self.apps[app_name] = app + + if calc_type is not None: + if calc_type not in self.default_apps: + raise ExecutorException(f"Unrecognized calculation type {calc_type}") + self.default_apps[calc_type] = app + + def submit( + self, + calc_type: str | None = None, + app_name: str | None = None, + app_args: str | None = None, + func: Any = None, + stdout: str | None = None, + stderr: str | None = None, + dry_run: bool = False, + wait_on_start: bool = False, + **kwargs, + ) -> GlobusComputeTask: + """Submit a function or registered application to Globus Compute. + + Parameters + ---------- + calc_type : str, optional + Calculation type (``"sim"`` or ``"gen"``). Used with *app_name*. + app_name : str, optional + Name of a previously registered application. + app_args : str, optional + Arguments passed alongside the function. + func : Callable, optional + A Python callable to execute remotely. Takes precedence over + *app_name* / *calc_type*. + stdout, stderr : str, optional + Ignored (stubs for API compatibility). + dry_run : bool, optional + If True, return a task without actually submitting. + wait_on_start : bool, optional + If True, block until the task is reported as started. + + Returns + ------- + GlobusComputeTask + """ + if dry_run: + raise NotImplementedError("dry_run is not supported for GlobusComputeExecutor") + + if func is not None: + fid = self._get_func_id(func) + app = Application(full_path="", name=func.__name__, calc_type="sim", pyobj=func) + elif app_name is not None: + app = self.get_app(app_name) + if app.pyobj is not None: + fid = self._get_func_id(app.pyobj) + else: + raise ValueError( + f"Application '{app_name}' has no pyobj callable registered. " + "Use the `func=...` argument, or register an app with `pyobj=`." + ) + elif calc_type is not None: + app = self.default_app(calc_type) + if app.pyobj is not None: + fid = self._get_func_id(app.pyobj) + else: + raise ValueError( + f"Default {calc_type} app has no pyobj callable. " + "Use the `func=...` argument, or register an app with `pyobj=`." + ) + else: + raise ValueError("One of `func`, `app_name`, or `calc_type` must be provided") + + args = app_args + future: Future = self._ensure_gc().submit_to_registered_function(fid, args) + task = GlobusComputeTask(future, app=app, app_args=args, workerid=self.workerID) + self.list_of_tasks.append(task) + + if wait_on_start: + task.wait() + + return task + + def set_workerID(self, workerid) -> None: + """Sets the worker ID for this executor.""" + self.workerID = workerid + + def set_worker_info(self, comm=None, workerid=None) -> None: + """Sets worker info for this executor.""" + self.workerID = workerid + self.comm = comm + + def serial_setup(self): + pass + + def set_resources(self, resources): + pass + + def add_platform_info(self, platform_info=None): + pass + + def set_gen_procs_gpus(self, libE_info): + pass diff --git a/libensemble/tests/unit_tests/test_globus_compute.py b/libensemble/tests/unit_tests/test_globus_compute.py new file mode 100644 index 000000000..3dab45f62 --- /dev/null +++ b/libensemble/tests/unit_tests/test_globus_compute.py @@ -0,0 +1,258 @@ +from unittest import mock + +import pytest + +from libensemble.executors.globus_compute_executor import ( + GlobusComputeExecutor, + GlobusComputeTask, +) +from libensemble.utils.globus_compute import GCSession + + +class TestGCSession: + def setup_method(self): + GCSession.clear() + + def test_get_or_create_executor(self): + with mock.patch.object(GCSession, "_create_executor") as mock_create: + mock_exec = mock.MagicMock() + mock_create.return_value = mock_exec + + ex1 = GCSession.get_or_create_executor("ep-1") + assert ex1 is mock_exec + mock_create.assert_called_once_with("ep-1") + + ex2 = GCSession.get_or_create_executor("ep-1") + assert ex2 is mock_exec + mock_create.assert_called_once() + + def test_get_or_create_caches_func_id(self): + with mock.patch.object(GCSession, "_create_executor") as mock_create: + mock_exec = mock.MagicMock() + mock_exec.register_function.return_value = "fid-42" + mock_create.return_value = mock_exec + + def my_func(): + pass + + ex1, fid1 = GCSession.get_or_create("ep-1", my_func) + assert ex1 is mock_exec + assert fid1 == "fid-42" + mock_exec.register_function.assert_called_once_with(my_func) + + ex2, fid2 = GCSession.get_or_create("ep-1", my_func) + assert ex2 is mock_exec + assert fid2 == "fid-42" + mock_exec.register_function.assert_called_once() + + def test_register_function(self): + with mock.patch.object(GCSession, "_create_executor") as mock_create: + mock_exec = mock.MagicMock() + mock_exec.register_function.return_value = "fid-99" + mock_create.return_value = mock_exec + + def f(): + pass + + ex, fid = GCSession.register_function("ep-1", f) + assert ex is mock_exec + assert fid == "fid-99" + mock_exec.register_function.assert_called_once_with(f) + + def test_module_not_found_returns_none(self): + GCSession._create_executor = classmethod(lambda cls, eid: None) + + ex = GCSession.get_or_create_executor("ep-1") + assert ex is None + + ex, fid = GCSession.get_or_create("ep-1", lambda: None) + assert ex is None + assert fid is None + + def test_thread_safety(self): + import threading + + with mock.patch.object(GCSession, "_create_executor") as mock_create: + mock_exec = mock.MagicMock() + mock_exec.register_function.return_value = "fid" + mock_create.return_value = mock_exec + + errors = [] + + def access(): + try: + for _ in range(100): + GCSession.get_or_create_executor("ep-t") + GCSession.get_or_create("ep-t", lambda: None) + except Exception as e: + errors.append(e) + + threads = [threading.Thread(target=access) for _ in range(10)] + for t in threads: + t.start() + for t in threads: + t.join() + + assert not errors, f"Thread safety errors: {errors}" + + +class TestGlobusComputeTask: + def make_task(self, future=None, app=None): + if app is None: + from libensemble.executors.executor import Application + + app = Application("", name="test_func", calc_type="sim", pyobj=lambda: None) + if future is None: + future = mock.MagicMock() + future.done.return_value = False + return GlobusComputeTask(future, app=app) + + def test_initial_state(self): + task = self.make_task() + assert task.state == "RUNNING" + assert not task.finished + assert task._gc_future is not None + + def test_poll_running(self): + future = mock.MagicMock() + future.done.return_value = False + task = self.make_task(future=future) + task.poll() + assert task.state == "RUNNING" + assert not task.finished + + def test_poll_finished_success(self): + future = mock.MagicMock() + future.done.return_value = True + future.result.return_value = None + task = self.make_task(future=future) + task.poll() + assert task.state == "FINISHED" + assert task.finished + assert task.success + + def test_poll_finished_failure(self): + future = mock.MagicMock() + future.done.return_value = True + future.result.side_effect = RuntimeError("boom") + task = self.make_task(future=future) + task.poll() + assert task.state == "FAILED" + assert task.finished + assert not task.success + + def test_wait_timeout(self): + future = mock.MagicMock() + future.result.side_effect = TimeoutError("timed out") + task = self.make_task(future=future) + with pytest.raises(Exception, match="timed out"): + task.wait(timeout=0.001) + + def test_kill(self): + future = mock.MagicMock() + task = self.make_task(future=future) + task.kill() + assert task.state == "USER_KILLED" + assert task.finished + future.cancel.assert_called_once() + + def test_running(self): + future = mock.MagicMock() + future.done.return_value = False + task = self.make_task(future=future) + assert task.running() + + def test_done(self): + future = mock.MagicMock() + future.done.return_value = True + future.result.return_value = None + task = self.make_task(future=future) + assert task.done() + + def test_not_done(self): + future = mock.MagicMock() + future.done.return_value = False + task = self.make_task(future=future) + assert not task.done() + + +class TestGlobusComputeExecutor: + def setup_method(self): + GCSession.clear() + + def test_init(self): + exctr = GlobusComputeExecutor(endpoint_id="ep-test") + assert exctr.endpoint_id == "ep-test" + assert exctr._gc_executor is None + assert exctr.workerID is None + + def test_submit_with_func(self): + with mock.patch.object(GCSession, "_create_executor") as mock_create: + mock_exec = mock.MagicMock() + mock_exec.register_function.return_value = "fid-xyz" + mock_create.return_value = mock_exec + + exctr = GlobusComputeExecutor(endpoint_id="ep-test") + exctr._ensure_gc() + + future_mock = mock.MagicMock() + mock_exec.submit_to_registered_function.return_value = future_mock + + def my_func(x): + return x * 2 + + task = exctr.submit(func=my_func, app_args="hello") + assert isinstance(task, GlobusComputeTask) + assert task._gc_future is future_mock + assert task.app is not None + assert task.app.name == "my_func" + + def test_submit_with_registered_app_pyobj(self): + with mock.patch.object(GCSession, "_create_executor") as mock_create: + mock_exec = mock.MagicMock() + mock_exec.register_function.return_value = "fid-app" + mock_create.return_value = mock_exec + + exctr = GlobusComputeExecutor(endpoint_id="ep-test") + exctr._ensure_gc() + + def app_func(): + return 42 + + exctr.register_app("/fake/path", app_name="myapp", calc_type="sim", pyobj=app_func) + + future_mock = mock.MagicMock() + mock_exec.submit_to_registered_function.return_value = future_mock + + task = exctr.submit(app_name="myapp") + assert isinstance(task, GlobusComputeTask) + assert task.app.name == "myapp" + + def test_submit_without_func_or_app_raises(self): + exctr = GlobusComputeExecutor(endpoint_id="ep-test") + with pytest.raises(ValueError): + exctr.submit() + + def test_register_function_caching(self): + with mock.patch.object(GCSession, "_create_executor") as mock_create: + mock_exec = mock.MagicMock() + mock_exec.register_function.return_value = "fid-cached" + mock_create.return_value = mock_exec + + exctr = GlobusComputeExecutor(endpoint_id="ep-test") + exctr._ensure_gc() + + def my_func(): + pass + + fid1 = exctr._get_func_id(my_func) + fid2 = exctr._get_func_id(my_func) + assert fid1 == fid2 + assert mock_exec.register_function.call_count == 1 + + def test_register_app_no_pyobj(self): + exctr = GlobusComputeExecutor(endpoint_id="ep-test") + exctr.register_app("/bin/echo", app_name="echo", calc_type="sim") + app = exctr.get_app("echo") + assert app.pyobj is None + assert app.full_path == "/bin/echo" diff --git a/libensemble/utils/globus_compute.py b/libensemble/utils/globus_compute.py new file mode 100644 index 000000000..c3341f304 --- /dev/null +++ b/libensemble/utils/globus_compute.py @@ -0,0 +1,87 @@ +import logging +import threading + +logger = logging.getLogger(__name__) + + +class GCSession: + """Per-process singleton cache for Globus Compute executors. + + Caches executor instances keyed by endpoint_id, ensuring only one + executor per endpoint per process. Thread-safe via ``threading.Lock``. + """ + + _instances: dict[str, tuple] = {} + _lock = threading.Lock() + + @classmethod + def get_or_create_executor(cls, endpoint_id: str): + """Get or create a cached executor for the given endpoint. + + Unlike :meth:`get_or_create`, this does **not** register a function. + """ + with cls._lock: + if endpoint_id in cls._instances: + return cls._instances[endpoint_id][0] + + executor = cls._create_executor(endpoint_id) + if executor is None: + return None + + cls._instances[endpoint_id] = (executor, None) + return executor + + @classmethod + def get_or_create(cls, endpoint_id: str, func): + """Get or create a cached ``(executor, func_id)`` pair. + + The first call for an endpoint creates the executor and registers + the callable. Subsequent calls return the cached pair (the + registered function is re-used). + """ + with cls._lock: + if endpoint_id in cls._instances: + executor, existing_fid = cls._instances[endpoint_id] + if existing_fid is not None: + return executor, existing_fid + func_id = executor.register_function(func) + cls._instances[endpoint_id] = (executor, func_id) + return executor, func_id + + executor = cls._create_executor(endpoint_id) + if executor is None: + return None, None + + func_id = executor.register_function(func) + cls._instances[endpoint_id] = (executor, func_id) + return executor, func_id + + @classmethod + def register_function(cls, endpoint_id: str, func): + """Register an additional function with an existing executor. + + Returns ``(executor, func_id)``. Unlike :meth:`get_or_create`, + this always registers and never caches the func_id (caller should + cache it themselves). + """ + executor = cls.get_or_create_executor(endpoint_id) + if executor is None: + return None, None + func_id = executor.register_function(func) + return executor, func_id + + @classmethod + def _create_executor(cls, endpoint_id: str): + try: + from globus_compute_sdk import Executor + except ModuleNotFoundError: + logger.warning("Globus Compute use detected but Globus Compute not importable. " "Is it installed?") + logger.warning("Running function evaluations normally on local resources.") + return None + return Executor(endpoint_id=endpoint_id) + + @classmethod + def clear(cls): + """Clear the cache (primarily for testing).""" + with cls._lock: + cls._instances.clear() From 3b34c60d3c40433359fba100052e6edf3dc6aed9 Mon Sep 17 00:00:00 2001 From: jlnav Date: Thu, 2 Jul 2026 16:23:28 -0500 Subject: [PATCH 09/16] Globus compute "mode" features and tests on this branch --- libensemble/libE.py | 12 + libensemble/manager.py | 206 +++++++++++++++++- .../test_gc_manager_submit.py | 112 ++++++++++ .../tests/unit_tests/test_ufunc_runners.py | 70 ------ libensemble/utils/runners.py | 32 --- 5 files changed, 327 insertions(+), 105 deletions(-) create mode 100644 libensemble/tests/functionality_tests/test_gc_manager_submit.py diff --git a/libensemble/libE.py b/libensemble/libE.py index 091a79e7f..4a6852666 100644 --- a/libensemble/libE.py +++ b/libensemble/libE.py @@ -262,6 +262,18 @@ def libE( libE_funcs = {"mpi": libE_mpi, "tcp": libE_tcp, "local": libE_local, "threads": libE_local} + if sim_specs.get("globus_compute_endpoint"): + libE_specs["_gc_only"] = True + if libE_specs.get("gen_on_worker"): + logger.info("GC-only mode: gen_on_worker is ignored (generator runs on manager)") + libE_specs["gen_on_worker"] = False + if libE_specs.get("comms", "mpi") != "local": + libE_specs["comms"] = "local" + logger.info("GC-only mode: switching to local comms (no workers needed)") + if not libE_specs.get("disable_resource_manager"): + libE_specs["disable_resource_manager"] = True + logger.info("GC-only mode: disabling resource manager (no local nodes)") + Resources.init_resources(libE_specs, platform_info) if Executor.executor is not None: Executor.executor.add_platform_info(platform_info) diff --git a/libensemble/manager.py b/libensemble/manager.py index 7995d2da9..a455e96f7 100644 --- a/libensemble/manager.py +++ b/libensemble/manager.py @@ -3,8 +3,10 @@ ============================ """ +import concurrent.futures import cProfile import glob +import inspect import logging import os import platform @@ -30,11 +32,14 @@ MAN_SIGNAL_KILL, PERSIS_STOP, STOP_TAG, + TASK_FAILED, + WORKER_DONE, calc_type_strings, ) from libensemble.resources.resources import Resources from libensemble.tools.fields_keys import protected_libE_fields from libensemble.tools.tools import _USER_CALC_DIR_WARNING +from libensemble.utils.globus_compute import GCSession from libensemble.utils.misc import _WorkerIndexer, extract_H_ranges from libensemble.utils.output_directory import EnsembleDirectory from libensemble.utils.timer import Timer @@ -243,6 +248,15 @@ def __init__( local_worker_comm = self._run_additional_worker(hist, sim_specs, gen_specs, libE_specs) self.wcomms = [local_worker_comm] + self.wcomms + if libE_specs.get("_gc_only"): + n_virtual = libE_specs.get("nworkers", 1) or 1 + virtual_W = np.zeros(n_virtual, dtype=Manager.worker_dtype) + start_id = len(self.W) + virtual_W["worker_id"] = np.arange(start_id, start_id + n_virtual) + virtual_W["gen_worker"] = False + self.W = np.concatenate([self.W, virtual_W]) + self.wcomms = self.wcomms + [None] * n_virtual + self.W = _WorkerIndexer(self.W, 1 - gen_on_worker) # if gen on worker, then no additional worker self.wcomms = _WorkerIndexer(self.wcomms, 1 - gen_on_worker) @@ -574,9 +588,18 @@ def _kill_cancelled_sims(self) -> None: logger.debug(f"Manager sending kill signals to H indices {kill_sim_rows}") kill_ids = self.hist.H["sim_id"][kill_sim_rows] kill_on_workers = self.hist.H["sim_worker"][kill_sim_rows] - for w in kill_on_workers: - self.wcomms[w].send(STOP_TAG, MAN_SIGNAL_KILL) - self.hist.H["kill_sent"][kill_ids] = True + + if self.libE_specs.get("_gc_only"): + sim_ids_to_kill = set(kill_ids) + for future, (sim_id, w) in list(self._gc_futures.items()): + if sim_id in sim_ids_to_kill: + future.cancel() + del self._gc_futures[future] + self.hist.H["kill_sent"][kill_ids] = True + else: + for w in kill_on_workers: + self.wcomms[w].send(STOP_TAG, MAN_SIGNAL_KILL) + self.hist.H["kill_sent"][kill_ids] = True # --- Handle termination @@ -692,6 +715,180 @@ def _alloc_work(self, H: npt.NDArray, persis_info: dict) -> dict: return output + def _init_gc(self) -> None: + endpoint = self.sim_specs.get("globus_compute_endpoint", "") + if not endpoint: + raise ValueError("_gc_only mode requires globus_compute_endpoint in sim_specs") + + self._gc_executor, self._gc_fid = GCSession.get_or_create(endpoint, self.sim_specs["sim_f"]) + self._gc_futures: dict[concurrent.futures.Future, tuple[int, int]] = {} + self._gc_nparams = len(inspect.signature(self.sim_specs["sim_f"]).parameters) + + def _gc_submit(self, Work: dict, w: int) -> None: + sim_ids = Work["libE_info"]["H_rows"] + H_fields = Work["H_fields"] + + for sim_id in sim_ids: + calc_in = np.empty(1, dtype=[(name, self.hist.H.dtype.fields[name][0]) for name in H_fields]) + for name in H_fields: + calc_in[name] = self.hist.H[name][sim_id] + + p_info = Work.get("persis_info", {}) + + libE_info = dict(Work["libE_info"]) + libE_info["comm"] = None + libE_info.pop("executor", None) + + args = [calc_in, p_info, self.sim_specs, libE_info][: self._gc_nparams] + future = self._gc_executor.submit_to_registered_function(self._gc_fid, args) + self._gc_futures[future] = (sim_id, w) + + @staticmethod + def _normalize_gc_result(result): + if isinstance(result, (tuple, list)): + if len(result) >= 3: + return result[0], result[1], result[2] + if len(result) == 2: + if isinstance(result[1], (int, str)): + return result[0], {}, result[1] + return result[0], result[1], WORKER_DONE + return result[0], {}, WORKER_DONE + return result, {}, WORKER_DONE + + def _gather_gc_results(self, persis_info: dict) -> dict: + time.sleep(0.0001) + + new_stuff = True + while new_stuff: + new_stuff = False + for w in self.W["worker_id"]: + if self.wcomms[w] is not None and self.wcomms[w].mail_flag(): + new_stuff = True + self._handle_msg_from_worker(persis_info, w) + + done = [f for f in self._gc_futures if f.done()] + for future in done: + sim_id, w = self._gc_futures.pop(future) + try: + out, p_info, calc_status = self._normalize_gc_result(future.result()) + D_recv = { + "calc_type": EVAL_SIM_TAG, + "calc_status": calc_status, + "calc_out": out, + "libE_info": {"keep_state": False, "H_rows": [sim_id]}, + "persis_info": p_info or {}, + } + except Exception as e: + logger.error(f"GC task failed for sim_id {sim_id}: {e}") + D_recv = { + "calc_type": EVAL_SIM_TAG, + "calc_status": TASK_FAILED, + "calc_out": None, + "libE_info": {"keep_state": False, "H_rows": [sim_id]}, + "persis_info": {}, + } + self._update_state_on_worker_msg(persis_info, D_recv, w) + + self._init_every_k_save() + return persis_info + + def _gc_cancel_futures(self) -> None: + for future in list(self._gc_futures.keys()): + future.cancel() + self._gc_futures.clear() + + def _run_gc_only(self, persis_info: dict) -> tuple[dict, int, int]: + self._init_gc() + try: + while not self.term_test(): + self._kill_cancelled_sims() + persis_info = self._gather_gc_results(persis_info) + Work, persis_info, flag = self._alloc_work(self.hist.trim_H(), persis_info) + if flag: + break + + for w in Work: + if self._sim_max_given(): + break + self._check_work_order(Work[w], w) + if Work[w]["tag"] == EVAL_GEN_TAG: + self._send_work_order(Work[w], w) + else: + self._gc_submit(Work[w], w) + self._update_state_on_alloc(Work[w], w) + + assert self.term_test() or any( + self.W["active"] != 0 + ), "alloc_f did not return any work, although all workers are idle." + except WorkerException as e: + report_worker_exc(e) + raise LoggedException(e.args[0], e.args[1]) from None + except Exception as e: + logger.error(traceback.format_exc()) + raise LoggedException(e.args) from None + finally: + result = self._gc_final_receive_and_kill(persis_info) + self.wcomms = [] + sys.stdout.flush() + sys.stderr.flush() + return result + + def _gc_final_receive_and_kill(self, persis_info: dict) -> tuple[dict, int, int]: + if any(self.W["persis_state"]): + for w in self.W["worker_id"][self.W["persis_state"] > 0]: + logger.debug(f"Manager sending PERSIS_STOP to worker {w}") + if self.libE_specs.get("final_gen_send", False): + rows_to_send = np.where(self.hist.H["sim_ended"] & ~self.hist.H["gen_informed"])[0] + work = { + "H_fields": self.gen_specs["persis_in"], + "persis_info": persis_info.get(w), + "tag": PERSIS_STOP, + "libE_info": {"persistent": True, "H_rows": rows_to_send}, + } + self._check_work_order(work, w, force=True) + self._send_work_order(work, w) + self.hist.update_history_to_gen(rows_to_send) + else: + self.wcomms[w].send(PERSIS_STOP, MAN_SIGNAL_KILL) + if not self.W[w]["active"]: + self.W[w]["active"] = self.W[w]["persis_state"] + self.persis_pending.append(w) + + exit_flag = 0 + while self._gc_futures: + done = [f for f in self._gc_futures if f.done()] + if not done: + time.sleep(0.1) + continue + for future in done: + sim_id, w = self._gc_futures.pop(future) + try: + out, p_info, calc_status = self._normalize_gc_result(future.result()) + D_recv = { + "calc_type": EVAL_SIM_TAG, + "calc_status": calc_status, + "calc_out": out, + "libE_info": {"keep_state": False, "H_rows": [sim_id]}, + "persis_info": p_info or {}, + } + except Exception: + continue + self._update_state_on_worker_msg(persis_info, D_recv, w) + + if 0 in self.W["worker_id"] and self.wcomms[0] is not None: + while self.wcomms[0].mail_flag(): + self._handle_msg_from_worker(persis_info, 0) + self.wcomms[0].send(STOP_TAG, MAN_SIGNAL_FINISH) + + self._init_every_k_save(complete=True) + self._clean_up_thread() + + if self.live_data is not None: + self.live_data.finalize(self.hist) + + persis_info["num_gens_started"] = 0 + return persis_info, exit_flag, self.elapsed() + # --- Main loop def run(self, persis_info: dict) -> tuple[dict, int, int]: @@ -699,6 +896,9 @@ def run(self, persis_info: dict) -> tuple[dict, int, int]: logger.debug(f"Manager initiated on node {socket.gethostname()}") logger.info(f"Manager exit_criteria: {self.exit_criteria}") + if self.libE_specs.get("_gc_only"): + return self._run_gc_only(persis_info) + # Continue receiving and giving until termination test is satisfied try: while not self.term_test(): diff --git a/libensemble/tests/functionality_tests/test_gc_manager_submit.py b/libensemble/tests/functionality_tests/test_gc_manager_submit.py new file mode 100644 index 000000000..40b76e568 --- /dev/null +++ b/libensemble/tests/functionality_tests/test_gc_manager_submit.py @@ -0,0 +1,112 @@ +""" +Tests libEnsemble's manager-side Globus Compute submission (GC-only mode). + +The manager submits simulation work directly to a mocked Globus Compute +endpoint instead of dispatching to local workers. The generator runs on +the manager thread as normal. + +Execute via: + python test_gc_manager_submit.py + +No MPI or local workers are needed -- GC-only mode uses local comms with +nworkers acting as the maximum number of concurrent in-flight GC futures. +""" + +# Do not change these lines - they are parsed by run-tests.sh +# TESTSUITE_COMMS: local +# TESTSUITE_NPROCS: 1 + +import concurrent.futures +from unittest import mock + +import numpy as np +from gest_api.vocs import VOCS + +from libensemble.gen_classes.sampling import UniformSample +from libensemble.libE import libE +from libensemble.specs import ExitCriteria, GenSpecs, LibeSpecs, SimSpecs +from libensemble.utils.globus_compute import GCSession + + +def norm_sim(H, persis_info, sim_specs, libE_info): + """Evaluate the Euclidean norm of each input point.""" + H_o = np.zeros(len(H), dtype=sim_specs["out"]) + for i in range(len(H)): + H_o["f"][i] = float(np.linalg.norm(H["x"][i])) + return H_o, persis_info + + +def _make_done_future(value): + f = concurrent.futures.Future() + f.set_result(value) + return f + + +def _make_gc_executor(sim_f): + executor = mock.MagicMock() + executor.register_function.return_value = "mock-fid" + + def fake_submit(fid, args): + result = sim_f(*args) + return _make_done_future(result) + + executor.submit_to_registered_function.side_effect = fake_submit + return executor + + +if __name__ == "__main__": + GCSession.clear() + + ENDPOINT = "mock-endpoint-uuid" + SIM_MAX = 20 + N_VIRTUAL_WORKERS = 4 + + vocs = VOCS( + variables={"x0": [-3.0, 3.0], "x1": [-2.0, 2.0]}, + objectives={"f": "MINIMIZE"}, + ) + + sim_specs = SimSpecs( + sim_f=norm_sim, + inputs=["x"], + outputs=[("f", float)], + globus_compute_endpoint=ENDPOINT, + ) + + gen_specs = GenSpecs( + generator=UniformSample(vocs), + inputs=["sim_id"], + persis_in=["f", "sim_id"], + outputs=[("x", float, (2,))], + batch_size=N_VIRTUAL_WORKERS, + ) + + libE_specs = LibeSpecs( + nworkers=N_VIRTUAL_WORKERS, + comms="local", + disable_log_files=True, + safe_mode=False, + ) + + exit_criteria = ExitCriteria(sim_max=SIM_MAX) + + mock_executor = _make_gc_executor(norm_sim) + + with mock.patch.object(GCSession, "_create_executor", return_value=mock_executor): + H, persis_info, flag = libE(sim_specs, gen_specs, exit_criteria, libE_specs=libE_specs) + + assert flag == 0, f"libEnsemble exited with unexpected flag {flag}" + assert ( + np.sum(H["sim_ended"]) >= SIM_MAX + ), f"Expected at least {SIM_MAX} completed sims, got {np.sum(H['sim_ended'])}" + + completed = H[H["sim_ended"]] + assert len(completed) >= SIM_MAX + assert np.all(completed["f"] >= 0.0), "Unexpected negative norm value" + assert ( + mock_executor.submit_to_registered_function.call_count >= SIM_MAX + ), f"Expected at least {SIM_MAX} GC submissions, got {mock_executor.submit_to_registered_function.call_count}" + + print(f"\nGC-only mode: {np.sum(H['sim_ended'])} sims completed via mocked Globus Compute.") + print(f"Best f value: {completed['f'].min():.6f}") + print("\nlibEnsemble GC-only functionality test passed.") diff --git a/libensemble/tests/unit_tests/test_ufunc_runners.py b/libensemble/tests/unit_tests/test_ufunc_runners.py index 79fda7c28..8765e077a 100644 --- a/libensemble/tests/unit_tests/test_ufunc_runners.py +++ b/libensemble/tests/unit_tests/test_ufunc_runners.py @@ -1,6 +1,4 @@ -import mock import numpy as np -import pytest import libensemble.tests.unit_tests.setup as setup from libensemble.tools.fields_keys import libE_fields @@ -68,75 +66,7 @@ def tupilize(arg1, arg2): assert result == (calc_in, {}) -@pytest.mark.extra -def test_globus_compute_runner_init(): - calc_in, sim_specs, gen_specs = get_ufunc_args() - - sim_specs["globus_compute_endpoint"] = "1234" - - with mock.patch("globus_compute_sdk.Executor"): - runner = Runner.from_specs(sim_specs) - - assert hasattr( - runner, "globus_compute_executor" - ), "Globus ComputeExecutor should have been instantiated when globus_compute_endpoint found in specs" - - -@pytest.mark.extra -def test_globus_compute_runner_pass(): - calc_in, sim_specs, gen_specs = get_ufunc_args() - - sim_specs["globus_compute_endpoint"] = "1234" - - with mock.patch("globus_compute_sdk.Executor"): - runner = Runner.from_specs(sim_specs) - - # Creating Mock Globus ComputeExecutor and Globus Compute future object - no exception - globus_compute_mock = mock.Mock() - globus_compute_future = mock.Mock() - globus_compute_mock.submit_to_registered_function.return_value = globus_compute_future - globus_compute_future.exception.return_value = None - globus_compute_future.result.return_value = (True, True) - - runner.globus_compute_executor = globus_compute_mock - runners = {1: runner.run} - - libE_info = {"H_rows": np.array([2, 3, 4]), "workerID": 1, "comm": "fakecomm"} - - out, persis_info = runners[1](calc_in, {"libE_info": libE_info, "persis_info": {}, "tag": 1}) - - assert all([out, persis_info]), "Globus Compute runner correctly returned results" - - -@pytest.mark.extra -def test_globus_compute_runner_fail(): - calc_in, sim_specs, gen_specs = get_ufunc_args() - - sim_specs["globus_compute_endpoint"] = "4321" - - with mock.patch("globus_compute_sdk.Executor"): - runner = Runner.from_specs(sim_specs) - - # Creating Mock Globus ComputeExecutor and Globus Compute future object - yes exception - globus_compute_mock = mock.Mock() - globus_compute_future = mock.Mock() - globus_compute_mock.submit_to_registered_function.return_value = globus_compute_future - globus_compute_future.exception.return_value = Exception - - runner.globus_compute_executor = globus_compute_mock - runners = {1: runner.run} - - libE_info = {"H_rows": np.array([2, 3, 4]), "workerID": 1, "comm": "fakecomm"} - - with pytest.raises(Exception): - out, persis_info = runners[1](calc_in, {"libE_info": libE_info, "persis_info": {}, "tag": 1}) - pytest.fail("Expected exception") - - if __name__ == "__main__": test_normal_runners() test_thread_runners() test_persis_info_from_none() - test_globus_compute_runner_init() - test_globus_compute_runner_pass() - test_globus_compute_runner_fail() diff --git a/libensemble/utils/runners.py b/libensemble/utils/runners.py index 72f8d1d4f..45f99435d 100644 --- a/libensemble/utils/runners.py +++ b/libensemble/utils/runners.py @@ -17,8 +17,6 @@ class Runner: @classmethod def from_specs(cls, specs): - if len(specs.get("globus_compute_endpoint", "")) > 0: - return GlobusComputeRunner(specs) if specs.get("threaded"): return ThreadRunner(specs) if (generator := specs.get("generator")) is not None: @@ -67,36 +65,6 @@ def run(self, calc_in: npt.NDArray, Work: dict) -> (npt.NDArray, dict, int | Non return out -class GlobusComputeRunner(Runner): - def __init__(self, specs): - super().__init__(specs) - self.globus_compute_executor = self._get_globus_compute_executor()(endpoint_id=specs["globus_compute_endpoint"]) - self.globus_compute_fid = self.globus_compute_executor.register_function(self.f) - - def _get_globus_compute_executor(self): - try: - from globus_compute_sdk import Executor - except ModuleNotFoundError: - logger.warning("Globus Compute use detected but Globus Compute not importable. Is it installed?") - logger.warning("Running function evaluations normally on local resources.") - return None - else: - return Executor - - def _result(self, calc_in: npt.NDArray, persis_info: dict, libE_info: dict) -> (npt.NDArray, dict, int | None): - from libensemble.worker import Worker - - libE_info["comm"] = None # 'comm' object not pickle-able - Worker._set_executor(0, None) # ditto for executor - - args = self._truncate_args(calc_in, persis_info, libE_info) - task_fut = self.globus_compute_executor.submit_to_registered_function(self.globus_compute_fid, args) - return task_fut.result() - - def shutdown(self) -> None: - self.globus_compute_executor.shutdown() - - class ThreadRunner(Runner): def __init__(self, specs): super().__init__(specs) From 8de6ceef83e1fd5f9b4a4d6d14b72ef6ede955dc Mon Sep 17 00:00:00 2001 From: jlnav Date: Thu, 2 Jul 2026 16:36:23 -0500 Subject: [PATCH 10/16] small test adjustments, but mostly the leftover docs adjusts --- .../advanced_installation.rst | 15 +- docs/conf.py | 6 +- docs/executor/ex_globus_compute.rst | 57 +++++++ docs/executor/ex_index.rst | 8 +- docs/executor/ex_overview.rst | 17 +- docs/platforms/globus_compute.rst | 145 ++++++++++++++++++ docs/platforms/platforms_index.rst | 63 +------- docs/running_libE.rst | 6 +- .../test_persistent_aposmm_nlopt.py | 2 +- 9 files changed, 248 insertions(+), 71 deletions(-) create mode 100644 docs/executor/ex_globus_compute.rst create mode 100644 docs/platforms/globus_compute.rst diff --git a/docs/advanced_installation/advanced_installation.rst b/docs/advanced_installation/advanced_installation.rst index 43d4ee738..0ac8018ff 100644 --- a/docs/advanced_installation/advanced_installation.rst +++ b/docs/advanced_installation/advanced_installation.rst @@ -1,7 +1,7 @@ Advanced Installation ===================== -`pip `__ \|\| `uv `__ \|\| `pixi `__ \|\| `conda `__ \|\| `Spack `__ +`pip `__ || `uv `__ || `pixi `__ || `conda `__ || `Spack `__ libEnsemble can be installed from ``pip``, ``uv``, ``pixi``, ``Conda``, or ``Spack``. @@ -31,7 +31,18 @@ Further recommendations for selected HPC systems are given in the Globus Compute -------------- -`Globus Compute`_ may be installed optionally to submit simulation function instances to remote Globus Compute endpoints. +`Globus Compute`_ may be installed optionally to submit simulation function +instances to remote Globus Compute endpoints:: + + pip install globus-compute-sdk + +This is an optional dependency; libEnsemble operates normally without it. +If Globus Compute is not installed and a ``globus_compute_endpoint`` is +configured, libEnsemble will warn and fall back to local execution. + +See :ref:`Globus Compute - Remote User Functions` for +usage, and the :doc:`GlobusComputeExecutor API reference` +for the full executor interface. .. _Globus Compute: https://www.globus.org/compute .. _Python: https://www.python.org/ diff --git a/docs/conf.py b/docs/conf.py index 8647bbf4f..2d227780e 100644 --- a/docs/conf.py +++ b/docs/conf.py @@ -31,7 +31,7 @@ def __getattr__(cls, name): return MagicMock() -autodoc_mock_imports = ["ax", "gpcam", "IPython", "matplotlib", "pandas", "scipy", "surmise"] +autodoc_mock_imports = ["ax", "globus_compute_sdk", "gpcam", "IPython", "matplotlib", "pandas", "scipy", "surmise"] MOCK_MODULES = [ "argparse", @@ -135,7 +135,7 @@ class AxParameterWarning(Warning): # Ensure it's a real warning subclass # The suffix(es) of source filenames. # You can specify multiple suffix as a list of string: # -# source_suffix = ['.rst', '.md'] +# source_suffix = ['.md', '.rst'] source_suffix = ".rst" # The master toctree document. @@ -205,7 +205,7 @@ class AxParameterWarning(Warning): # Ensure it's a real warning subclass html_favicon = "./images/libE_logo_circle.png" html_title = "libEnsemble" -# Theme options are theme-specific and customize the look and feel of a theme +# Theme options are theme-specific and customize the look and feel of the theme # further. For a list of options available for each theme, see the # documentation. # diff --git a/docs/executor/ex_globus_compute.rst b/docs/executor/ex_globus_compute.rst new file mode 100644 index 000000000..56047cf26 --- /dev/null +++ b/docs/executor/ex_globus_compute.rst @@ -0,0 +1,57 @@ +Globus Compute Executor +======================= + +`Overview `__ || `Base Executor `__ || `MPI Executor `__ || **Globus Compute Executor** + +The :class:`GlobusComputeExecutor` +submits Python callables to a remote `Globus Compute`_ endpoint instead of +launching local subprocesses. It can be used inside simulator functions in the +same way as the :doc:`MPI Executor`, retrieving it from +``libE_info["executor"]``. + +See :ref:`Globus Compute - Remote User Functions` for an +overview of the two GC integration modes (manager-side GC-only and user-facing +executor). + +.. note:: + + ``globus-compute-sdk`` must be installed to use this executor:: + + pip install globus-compute-sdk + + Users must also authenticate via Globus_ and have an active + `Globus Compute endpoint`_ running on the target system. + +GlobusComputeExecutor +--------------------- + +.. autoclass:: libensemble.executors.globus_compute_executor.GlobusComputeExecutor + :members: register_app, submit, set_workerID, set_worker_info + :show-inheritance: + + .. automethod:: __init__ + +GlobusComputeTask +----------------- + +Tasks are created and returned by +:meth:`GlobusComputeExecutor.submit()`. +Each task wraps a ``concurrent.futures.Future`` from the Globus Compute SDK +and exposes the same polling interface as other libEnsemble tasks. + +.. autoclass:: libensemble.executors.globus_compute_executor.GlobusComputeTask + :members: poll, wait, kill, result, running, done, cancelled + +**Task states**: ``RUNNING`` | ``FINISHED`` | ``FAILED`` | ``USER_KILLED`` + +**Key attributes**: + +:task.state: (string) Current task state - one of the values above. +:task.finished: (bool) True once the task has completed (successfully or not). +:task.success: (bool) True if the remote callable returned without raising. +:task.runtime: (float) Elapsed wall-clock seconds since submission. +:task.submit_time: (float) Time since epoch at submission. + +.. _Globus Compute: https://www.globus.org/compute +.. _Globus: https://www.globus.org/ +.. _Globus Compute endpoint: https://globus-compute.readthedocs.io/en/latest/endpoints.html diff --git a/docs/executor/ex_index.rst b/docs/executor/ex_index.rst index ac5c1c66b..7d9e3a522 100644 --- a/docs/executor/ex_index.rst +++ b/docs/executor/ex_index.rst @@ -1,6 +1,6 @@ .. _executor_index: -**Overview** \|\| `Base Executor `__ \|\| `MPI Executor `__ \|\| `Flux Executor `__ +**Overview** || `Base Executor `__ || `MPI Executor `__ || `Flux Executor `__ || `Globus Compute Executor `__ Executors ========= @@ -15,8 +15,12 @@ portable interface for running and managing user applications. ex_base ex_mpi ex_flux + ex_globus_compute The **Executor** provides a portable interface for running applications on any system and -any number of compute resources. +any number of compute resources. The :doc:`MPI Executor` launches MPI +applications on local resources; the +:doc:`Globus Compute Executor` submits Python callables to +remote Globus Compute endpoints. Please select from the sections above or the sidebar navigation to read more. diff --git a/docs/executor/ex_overview.rst b/docs/executor/ex_overview.rst index fe4d6c45e..7054faedc 100644 --- a/docs/executor/ex_overview.rst +++ b/docs/executor/ex_overview.rst @@ -1,7 +1,7 @@ Overview ======== -**Overview** \|\| `Base Executor `__ \|\| `MPI Executor `__ \|\| `Flux Executor `__ +**Overview** || `Base Executor `__ || `MPI Executor `__ || `Flux Executor `__ || `Globus Compute Executor `__ The **Executor** provides a portable interface for running applications on any system and any number of compute resources. @@ -156,4 +156,19 @@ which partitions resources among workers, ensuring that runs utilize different resources (e.g., nodes). Furthermore, the ``MPIExecutor`` offers resilience via the feature of re-launching tasks that fail to start because of system factors. +Remote Execution with Globus Compute +------------------------------------- + +The :doc:`GlobusComputeExecutor` submits Python callables +to remote `Globus Compute`_ endpoints instead of launching local subprocesses. +It exposes the same ``submit()`` / ``poll()`` / ``kill()`` interface as other +libEnsemble executors and can be retrieved from ``libE_info["executor"]`` +inside simulator functions. + +See :ref:`Globus Compute - Remote User Functions` for an +overview of all Globus Compute integration modes and the +:doc:`GlobusComputeExecutor API reference` for the full +interface. + .. _concurrent futures: https://docs.python.org/3/library/concurrent.futures.html +.. _Globus Compute: https://www.globus.org/compute diff --git a/docs/platforms/globus_compute.rst b/docs/platforms/globus_compute.rst new file mode 100644 index 000000000..0679f9559 --- /dev/null +++ b/docs/platforms/globus_compute.rst @@ -0,0 +1,145 @@ +.. _globus_compute_ref: + +====================================== +Globus Compute - Remote User Functions +====================================== + +`Globus Compute`_ (formerly funcX) is a distributed, high-performance +function-as-a-service platform. When libEnsemble is running on a resource with +internet access (laptops, login nodes, other servers, etc.), it can offload +simulator calls to remote Globus Compute endpoints: + + .. image:: ../images/funcxmodel.png + :alt: running_with_globus_compute + :scale: 50 + :align: center + +This is useful for running ensembles across machines and heterogeneous resources. +There are **two approaches**, described below. + +.. dropdown:: **Caveats** + + The following caveats apply to all Globus Compute modes: + + 1. Simulator functions submitted to Globus Compute must be non-persistent, + since manager-worker communicators cannot be serialized or used by a + remote resource. + + 2. ``Executor.manager_poll()`` is not available inside remotely executed + functions. Control over remote work is limited to inspecting return + values and exceptions when tasks complete. + + 3. Globus Compute imposes a `handful of task-rate and data limits`_ on + submitted functions. + + 4. Users are responsible for authenticating via Globus_ and maintaining their + `Globus Compute endpoints`_ on their target systems. + +.. _gc_only_mode: + +Manager-side GC (GC-only mode) +------------------------------- + +The recommended approach for most use cases. When +``globus_compute_endpoint`` is set in :class:`SimSpecs` +and ``gen_on_worker`` is not set (the default), libEnsemble enters +**GC-only mode**: no local worker processes are launched. The manager +submits simulation work directly to Globus Compute and polls futures for +results. The generator still runs as a local thread on the manager. + +``nworkers`` controls the maximum number of simultaneously in-flight +Globus Compute tasks (virtual concurrency). The default is 1. + +This mode supports both the :ref:`gest-api simulator format` +(``SimSpecs.simulator``) and the legacy ``sim_f`` format. + +.. code-block:: python + + from libensemble import Ensemble + from libensemble.specs import ExitCriteria, GenSpecs, LibeSpecs, SimSpecs + + + def my_sim(input_dict: dict, **kwargs) -> dict: + """gest-api simulator - runs remotely on the GC endpoint.""" + return {"f": input_dict["x"] ** 2} + + + sim_specs = SimSpecs( + simulator=my_sim, + vocs=vocs, + globus_compute_endpoint="3af6dc24-3f27-4c49-8d11-e301ade15353", + ) + + libE_specs = LibeSpecs(nworkers=4) # up to 4 concurrent GC tasks + + workflow = Ensemble( + sim_specs=sim_specs, + gen_specs=gen_specs, + libE_specs=libE_specs, + exit_criteria=ExitCriteria(sim_max=20), + ) + H, _, _ = workflow.run() + +Users can also define ``Executor`` instances within their remote simulator +functions and submit MPI applications normally, as long as libEnsemble and +the target application are accessible on the remote system:: + + # Within the remote simulator function + from libensemble.executors import MPIExecutor + exctr = MPIExecutor() + exctr.register_app(full_path="/home/user/forces.x", app_name="forces") + task = exctr.submit(app_name="forces", num_procs=64) + +.. note:: + + Both the simulator callable and any VOCS object must be picklable, + as they are serialized and shipped to the remote Globus Compute endpoint. + +.. _gc_executor_approach: + +GlobusComputeExecutor (user-facing) +------------------------------------ + +For workflows where the simulation function itself orchestrates remote +calls, like fanning out to multiple endpoints or mixing local +and remote work. Use the +:class:`GlobusComputeExecutor` +directly inside the simulator. + +Create and register the executor in the top-level script: + +.. code-block:: python + + from libensemble.executors import GlobusComputeExecutor + + exctr = GlobusComputeExecutor(endpoint_id="3af6dc24-3f27-4c49-8d11-e301ade15353") + +Then use it inside the simulator function: + +.. code-block:: python + + import time + + + def my_sim(H, persis_info, sim_specs, libE_info): + exctr = libE_info["executor"] + + task = exctr.submit(func=my_remote_func, app_args=H["x"][0]) + + while not task.finished: + task.poll() + if exctr.manager_kill_received(): + task.kill() + break + time.sleep(0.1) + + return H_o, persis_info + +See the :doc:`GlobusComputeExecutor API reference<../executor/ex_globus_compute>` for +the full interface including ``register_app``, ``submit``, and +:class:`GlobusComputeTask` methods. + +.. _Globus Compute: https://www.globus.org/compute +.. _Globus Compute endpoints: https://globus-compute.readthedocs.io/en/latest/endpoints.html +.. _Globus: https://www.globus.org/ +.. _handful of task-rate and data limits: https://globus-compute.readthedocs.io/en/latest/limits.html diff --git a/docs/platforms/platforms_index.rst b/docs/platforms/platforms_index.rst index 4628d6663..092f56d1f 100644 --- a/docs/platforms/platforms_index.rst +++ b/docs/platforms/platforms_index.rst @@ -159,60 +159,9 @@ will better manage simulation and generation functions that contain considerable computational work or I/O. Therefore the second option is to use Globus Compute to isolate this work from the workers. -.. _globus_compute_ref: - -Globus Compute - Remote User Functions --------------------------------------- - -If libEnsemble is running on some resource with -internet access (laptops, login nodes, other servers, etc.), workers can be instructed to -launch generator or simulator user function instances to separate resources from -themselves via `Globus Compute`_ (formerly funcX), a distributed, high-performance function-as-a-service platform: - - .. image:: ../images/funcxmodel.png - :alt: running_with_globus_compute - :scale: 50 - :align: center - -This is useful for running ensembles across machines and heterogeneous resources, but -comes with several caveats: - - 1. User functions registered with Globus Compute must be *non-persistent*, since - manager-worker communicators can't be serialized or used by a remote resource. - - 2. Likewise, the ``Executor.manager_poll()`` capability is disabled. The only - available control over remote functions by workers is processing return values - or exceptions when they complete. - - 3. Globus Compute imposes a `handful of task-rate and data limits`_ on submitted functions. - - 4. Users are responsible for authenticating via Globus_ and maintaining their - `Globus Compute endpoints`_ on their target systems. - -Users can still define Executor instances within their user functions and submit -MPI applications normally, as long as libEnsemble and the target application are -accessible on the remote system:: - - # Within remote user function - from libensemble.executors import MPIExecutor - exctr = MPIExecutor() - exctr.register_app(full_path="/home/user/forces.x", app_name="forces") - task = exctr.submit(app_name="forces", num_procs=64) - -Specify a Globus Compute endpoint in :class:`sim_specs` via the ``globus_compute_endpoint`` -argument. For example:: - - from libensemble.specs import SimSpecs - - sim_specs = SimSpecs( - sim_f = sim_f, - inputs = ["x"], - out = [("f", float)], - globus_compute_endpoint = "3af6dc24-3f27-4c49-8d11-e301ade15353", - ) - -See the ``libensemble/tests/scaling_tests/globus_compute_forces`` directory for a complete -remote-simulation example. +See :doc:`Globus Compute - Remote User Functions` for the two +integration approaches (manager-side GC-only mode and the user-facing +``GlobusComputeExecutor``). Instructions for Specific Platforms ----------------------------------- @@ -232,9 +181,5 @@ libEnsemble on specific HPC systems. perlmutter polaris srun + globus_compute example_scripts - -.. _Globus Compute: https://www.globus.org/compute -.. _Globus Compute endpoints: https://globus-compute.readthedocs.io/en/latest/endpoints/endpoint_examples.html -.. _Globus: https://www.globus.org/ -.. _handful of task-rate and data limits: https://globus-compute.readthedocs.io/en/latest/limits.html diff --git a/docs/running_libE.rst b/docs/running_libE.rst index 6329e13e2..2677eb8d2 100644 --- a/docs/running_libE.rst +++ b/docs/running_libE.rst @@ -83,9 +83,9 @@ if using an :class:`Ensemble` object with **Reverse-ssh interface** Set ``comms`` to ``ssh`` to launch workers on remote ssh-accessible systems. This -co-locates workers, functions, and any applications. User -functions can also be persistent, unlike when launching remote functions via -:ref:`Globus Compute`. +co-locates workers, functions, and any applications. Simulator functions can be +persistent, unlike those submitted to :ref:`Globus Compute`, +which must be non-persistent. The remote working directory and Python need to be specified. This may resemble:: diff --git a/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py b/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py index 28da42d53..1c3e4cca4 100644 --- a/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py +++ b/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py @@ -73,7 +73,7 @@ "xtol_abs": 1e-6, "ftol_abs": 1e-6, "dist_to_bound_multiple": 0.5, - "max_active_runs": 6, + "max_active_runs": nworkers - 1, "lb": np.array([-3, -2]), "ub": np.array([3, 2]), }, From 910895a186a1dd9d36be1f5021d3b4161590eaa6 Mon Sep 17 00:00:00 2001 From: jlnav Date: Thu, 2 Jul 2026 16:36:23 -0500 Subject: [PATCH 11/16] small test adjustments, but mostly the leftover docs adjusts --- .../advanced_installation.rst | 15 +- docs/conf.py | 6 +- docs/executor/ex_globus_compute.rst | 57 +++++++ docs/executor/ex_index.rst | 8 +- docs/executor/ex_overview.rst | 17 +- docs/platforms/globus_compute.rst | 145 ++++++++++++++++++ docs/platforms/platforms_index.rst | 63 +------- docs/running_libE.rst | 6 +- .../test_persistent_aposmm_nlopt.py | 2 +- 9 files changed, 248 insertions(+), 71 deletions(-) create mode 100644 docs/executor/ex_globus_compute.rst create mode 100644 docs/platforms/globus_compute.rst diff --git a/docs/advanced_installation/advanced_installation.rst b/docs/advanced_installation/advanced_installation.rst index 43d4ee738..0ac8018ff 100644 --- a/docs/advanced_installation/advanced_installation.rst +++ b/docs/advanced_installation/advanced_installation.rst @@ -1,7 +1,7 @@ Advanced Installation ===================== -`pip `__ \|\| `uv `__ \|\| `pixi `__ \|\| `conda `__ \|\| `Spack `__ +`pip `__ || `uv `__ || `pixi `__ || `conda `__ || `Spack `__ libEnsemble can be installed from ``pip``, ``uv``, ``pixi``, ``Conda``, or ``Spack``. @@ -31,7 +31,18 @@ Further recommendations for selected HPC systems are given in the Globus Compute -------------- -`Globus Compute`_ may be installed optionally to submit simulation function instances to remote Globus Compute endpoints. +`Globus Compute`_ may be installed optionally to submit simulation function +instances to remote Globus Compute endpoints:: + + pip install globus-compute-sdk + +This is an optional dependency; libEnsemble operates normally without it. +If Globus Compute is not installed and a ``globus_compute_endpoint`` is +configured, libEnsemble will warn and fall back to local execution. + +See :ref:`Globus Compute - Remote User Functions` for +usage, and the :doc:`GlobusComputeExecutor API reference` +for the full executor interface. .. _Globus Compute: https://www.globus.org/compute .. _Python: https://www.python.org/ diff --git a/docs/conf.py b/docs/conf.py index 8647bbf4f..2d227780e 100644 --- a/docs/conf.py +++ b/docs/conf.py @@ -31,7 +31,7 @@ def __getattr__(cls, name): return MagicMock() -autodoc_mock_imports = ["ax", "gpcam", "IPython", "matplotlib", "pandas", "scipy", "surmise"] +autodoc_mock_imports = ["ax", "globus_compute_sdk", "gpcam", "IPython", "matplotlib", "pandas", "scipy", "surmise"] MOCK_MODULES = [ "argparse", @@ -135,7 +135,7 @@ class AxParameterWarning(Warning): # Ensure it's a real warning subclass # The suffix(es) of source filenames. # You can specify multiple suffix as a list of string: # -# source_suffix = ['.rst', '.md'] +# source_suffix = ['.md', '.rst'] source_suffix = ".rst" # The master toctree document. @@ -205,7 +205,7 @@ class AxParameterWarning(Warning): # Ensure it's a real warning subclass html_favicon = "./images/libE_logo_circle.png" html_title = "libEnsemble" -# Theme options are theme-specific and customize the look and feel of a theme +# Theme options are theme-specific and customize the look and feel of the theme # further. For a list of options available for each theme, see the # documentation. # diff --git a/docs/executor/ex_globus_compute.rst b/docs/executor/ex_globus_compute.rst new file mode 100644 index 000000000..56047cf26 --- /dev/null +++ b/docs/executor/ex_globus_compute.rst @@ -0,0 +1,57 @@ +Globus Compute Executor +======================= + +`Overview `__ || `Base Executor `__ || `MPI Executor `__ || **Globus Compute Executor** + +The :class:`GlobusComputeExecutor` +submits Python callables to a remote `Globus Compute`_ endpoint instead of +launching local subprocesses. It can be used inside simulator functions in the +same way as the :doc:`MPI Executor`, retrieving it from +``libE_info["executor"]``. + +See :ref:`Globus Compute - Remote User Functions` for an +overview of the two GC integration modes (manager-side GC-only and user-facing +executor). + +.. note:: + + ``globus-compute-sdk`` must be installed to use this executor:: + + pip install globus-compute-sdk + + Users must also authenticate via Globus_ and have an active + `Globus Compute endpoint`_ running on the target system. + +GlobusComputeExecutor +--------------------- + +.. autoclass:: libensemble.executors.globus_compute_executor.GlobusComputeExecutor + :members: register_app, submit, set_workerID, set_worker_info + :show-inheritance: + + .. automethod:: __init__ + +GlobusComputeTask +----------------- + +Tasks are created and returned by +:meth:`GlobusComputeExecutor.submit()`. +Each task wraps a ``concurrent.futures.Future`` from the Globus Compute SDK +and exposes the same polling interface as other libEnsemble tasks. + +.. autoclass:: libensemble.executors.globus_compute_executor.GlobusComputeTask + :members: poll, wait, kill, result, running, done, cancelled + +**Task states**: ``RUNNING`` | ``FINISHED`` | ``FAILED`` | ``USER_KILLED`` + +**Key attributes**: + +:task.state: (string) Current task state - one of the values above. +:task.finished: (bool) True once the task has completed (successfully or not). +:task.success: (bool) True if the remote callable returned without raising. +:task.runtime: (float) Elapsed wall-clock seconds since submission. +:task.submit_time: (float) Time since epoch at submission. + +.. _Globus Compute: https://www.globus.org/compute +.. _Globus: https://www.globus.org/ +.. _Globus Compute endpoint: https://globus-compute.readthedocs.io/en/latest/endpoints.html diff --git a/docs/executor/ex_index.rst b/docs/executor/ex_index.rst index ac5c1c66b..7d9e3a522 100644 --- a/docs/executor/ex_index.rst +++ b/docs/executor/ex_index.rst @@ -1,6 +1,6 @@ .. _executor_index: -**Overview** \|\| `Base Executor `__ \|\| `MPI Executor `__ \|\| `Flux Executor `__ +**Overview** || `Base Executor `__ || `MPI Executor `__ || `Flux Executor `__ || `Globus Compute Executor `__ Executors ========= @@ -15,8 +15,12 @@ portable interface for running and managing user applications. ex_base ex_mpi ex_flux + ex_globus_compute The **Executor** provides a portable interface for running applications on any system and -any number of compute resources. +any number of compute resources. The :doc:`MPI Executor` launches MPI +applications on local resources; the +:doc:`Globus Compute Executor` submits Python callables to +remote Globus Compute endpoints. Please select from the sections above or the sidebar navigation to read more. diff --git a/docs/executor/ex_overview.rst b/docs/executor/ex_overview.rst index fe4d6c45e..7054faedc 100644 --- a/docs/executor/ex_overview.rst +++ b/docs/executor/ex_overview.rst @@ -1,7 +1,7 @@ Overview ======== -**Overview** \|\| `Base Executor `__ \|\| `MPI Executor `__ \|\| `Flux Executor `__ +**Overview** || `Base Executor `__ || `MPI Executor `__ || `Flux Executor `__ || `Globus Compute Executor `__ The **Executor** provides a portable interface for running applications on any system and any number of compute resources. @@ -156,4 +156,19 @@ which partitions resources among workers, ensuring that runs utilize different resources (e.g., nodes). Furthermore, the ``MPIExecutor`` offers resilience via the feature of re-launching tasks that fail to start because of system factors. +Remote Execution with Globus Compute +------------------------------------- + +The :doc:`GlobusComputeExecutor` submits Python callables +to remote `Globus Compute`_ endpoints instead of launching local subprocesses. +It exposes the same ``submit()`` / ``poll()`` / ``kill()`` interface as other +libEnsemble executors and can be retrieved from ``libE_info["executor"]`` +inside simulator functions. + +See :ref:`Globus Compute - Remote User Functions` for an +overview of all Globus Compute integration modes and the +:doc:`GlobusComputeExecutor API reference` for the full +interface. + .. _concurrent futures: https://docs.python.org/3/library/concurrent.futures.html +.. _Globus Compute: https://www.globus.org/compute diff --git a/docs/platforms/globus_compute.rst b/docs/platforms/globus_compute.rst new file mode 100644 index 000000000..0679f9559 --- /dev/null +++ b/docs/platforms/globus_compute.rst @@ -0,0 +1,145 @@ +.. _globus_compute_ref: + +====================================== +Globus Compute - Remote User Functions +====================================== + +`Globus Compute`_ (formerly funcX) is a distributed, high-performance +function-as-a-service platform. When libEnsemble is running on a resource with +internet access (laptops, login nodes, other servers, etc.), it can offload +simulator calls to remote Globus Compute endpoints: + + .. image:: ../images/funcxmodel.png + :alt: running_with_globus_compute + :scale: 50 + :align: center + +This is useful for running ensembles across machines and heterogeneous resources. +There are **two approaches**, described below. + +.. dropdown:: **Caveats** + + The following caveats apply to all Globus Compute modes: + + 1. Simulator functions submitted to Globus Compute must be non-persistent, + since manager-worker communicators cannot be serialized or used by a + remote resource. + + 2. ``Executor.manager_poll()`` is not available inside remotely executed + functions. Control over remote work is limited to inspecting return + values and exceptions when tasks complete. + + 3. Globus Compute imposes a `handful of task-rate and data limits`_ on + submitted functions. + + 4. Users are responsible for authenticating via Globus_ and maintaining their + `Globus Compute endpoints`_ on their target systems. + +.. _gc_only_mode: + +Manager-side GC (GC-only mode) +------------------------------- + +The recommended approach for most use cases. When +``globus_compute_endpoint`` is set in :class:`SimSpecs` +and ``gen_on_worker`` is not set (the default), libEnsemble enters +**GC-only mode**: no local worker processes are launched. The manager +submits simulation work directly to Globus Compute and polls futures for +results. The generator still runs as a local thread on the manager. + +``nworkers`` controls the maximum number of simultaneously in-flight +Globus Compute tasks (virtual concurrency). The default is 1. + +This mode supports both the :ref:`gest-api simulator format` +(``SimSpecs.simulator``) and the legacy ``sim_f`` format. + +.. code-block:: python + + from libensemble import Ensemble + from libensemble.specs import ExitCriteria, GenSpecs, LibeSpecs, SimSpecs + + + def my_sim(input_dict: dict, **kwargs) -> dict: + """gest-api simulator - runs remotely on the GC endpoint.""" + return {"f": input_dict["x"] ** 2} + + + sim_specs = SimSpecs( + simulator=my_sim, + vocs=vocs, + globus_compute_endpoint="3af6dc24-3f27-4c49-8d11-e301ade15353", + ) + + libE_specs = LibeSpecs(nworkers=4) # up to 4 concurrent GC tasks + + workflow = Ensemble( + sim_specs=sim_specs, + gen_specs=gen_specs, + libE_specs=libE_specs, + exit_criteria=ExitCriteria(sim_max=20), + ) + H, _, _ = workflow.run() + +Users can also define ``Executor`` instances within their remote simulator +functions and submit MPI applications normally, as long as libEnsemble and +the target application are accessible on the remote system:: + + # Within the remote simulator function + from libensemble.executors import MPIExecutor + exctr = MPIExecutor() + exctr.register_app(full_path="/home/user/forces.x", app_name="forces") + task = exctr.submit(app_name="forces", num_procs=64) + +.. note:: + + Both the simulator callable and any VOCS object must be picklable, + as they are serialized and shipped to the remote Globus Compute endpoint. + +.. _gc_executor_approach: + +GlobusComputeExecutor (user-facing) +------------------------------------ + +For workflows where the simulation function itself orchestrates remote +calls, like fanning out to multiple endpoints or mixing local +and remote work. Use the +:class:`GlobusComputeExecutor` +directly inside the simulator. + +Create and register the executor in the top-level script: + +.. code-block:: python + + from libensemble.executors import GlobusComputeExecutor + + exctr = GlobusComputeExecutor(endpoint_id="3af6dc24-3f27-4c49-8d11-e301ade15353") + +Then use it inside the simulator function: + +.. code-block:: python + + import time + + + def my_sim(H, persis_info, sim_specs, libE_info): + exctr = libE_info["executor"] + + task = exctr.submit(func=my_remote_func, app_args=H["x"][0]) + + while not task.finished: + task.poll() + if exctr.manager_kill_received(): + task.kill() + break + time.sleep(0.1) + + return H_o, persis_info + +See the :doc:`GlobusComputeExecutor API reference<../executor/ex_globus_compute>` for +the full interface including ``register_app``, ``submit``, and +:class:`GlobusComputeTask` methods. + +.. _Globus Compute: https://www.globus.org/compute +.. _Globus Compute endpoints: https://globus-compute.readthedocs.io/en/latest/endpoints.html +.. _Globus: https://www.globus.org/ +.. _handful of task-rate and data limits: https://globus-compute.readthedocs.io/en/latest/limits.html diff --git a/docs/platforms/platforms_index.rst b/docs/platforms/platforms_index.rst index 4628d6663..092f56d1f 100644 --- a/docs/platforms/platforms_index.rst +++ b/docs/platforms/platforms_index.rst @@ -159,60 +159,9 @@ will better manage simulation and generation functions that contain considerable computational work or I/O. Therefore the second option is to use Globus Compute to isolate this work from the workers. -.. _globus_compute_ref: - -Globus Compute - Remote User Functions --------------------------------------- - -If libEnsemble is running on some resource with -internet access (laptops, login nodes, other servers, etc.), workers can be instructed to -launch generator or simulator user function instances to separate resources from -themselves via `Globus Compute`_ (formerly funcX), a distributed, high-performance function-as-a-service platform: - - .. image:: ../images/funcxmodel.png - :alt: running_with_globus_compute - :scale: 50 - :align: center - -This is useful for running ensembles across machines and heterogeneous resources, but -comes with several caveats: - - 1. User functions registered with Globus Compute must be *non-persistent*, since - manager-worker communicators can't be serialized or used by a remote resource. - - 2. Likewise, the ``Executor.manager_poll()`` capability is disabled. The only - available control over remote functions by workers is processing return values - or exceptions when they complete. - - 3. Globus Compute imposes a `handful of task-rate and data limits`_ on submitted functions. - - 4. Users are responsible for authenticating via Globus_ and maintaining their - `Globus Compute endpoints`_ on their target systems. - -Users can still define Executor instances within their user functions and submit -MPI applications normally, as long as libEnsemble and the target application are -accessible on the remote system:: - - # Within remote user function - from libensemble.executors import MPIExecutor - exctr = MPIExecutor() - exctr.register_app(full_path="/home/user/forces.x", app_name="forces") - task = exctr.submit(app_name="forces", num_procs=64) - -Specify a Globus Compute endpoint in :class:`sim_specs` via the ``globus_compute_endpoint`` -argument. For example:: - - from libensemble.specs import SimSpecs - - sim_specs = SimSpecs( - sim_f = sim_f, - inputs = ["x"], - out = [("f", float)], - globus_compute_endpoint = "3af6dc24-3f27-4c49-8d11-e301ade15353", - ) - -See the ``libensemble/tests/scaling_tests/globus_compute_forces`` directory for a complete -remote-simulation example. +See :doc:`Globus Compute - Remote User Functions` for the two +integration approaches (manager-side GC-only mode and the user-facing +``GlobusComputeExecutor``). Instructions for Specific Platforms ----------------------------------- @@ -232,9 +181,5 @@ libEnsemble on specific HPC systems. perlmutter polaris srun + globus_compute example_scripts - -.. _Globus Compute: https://www.globus.org/compute -.. _Globus Compute endpoints: https://globus-compute.readthedocs.io/en/latest/endpoints/endpoint_examples.html -.. _Globus: https://www.globus.org/ -.. _handful of task-rate and data limits: https://globus-compute.readthedocs.io/en/latest/limits.html diff --git a/docs/running_libE.rst b/docs/running_libE.rst index 6329e13e2..2677eb8d2 100644 --- a/docs/running_libE.rst +++ b/docs/running_libE.rst @@ -83,9 +83,9 @@ if using an :class:`Ensemble` object with **Reverse-ssh interface** Set ``comms`` to ``ssh`` to launch workers on remote ssh-accessible systems. This -co-locates workers, functions, and any applications. User -functions can also be persistent, unlike when launching remote functions via -:ref:`Globus Compute`. +co-locates workers, functions, and any applications. Simulator functions can be +persistent, unlike those submitted to :ref:`Globus Compute`, +which must be non-persistent. The remote working directory and Python need to be specified. This may resemble:: diff --git a/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py b/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py index 28da42d53..1c3e4cca4 100644 --- a/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py +++ b/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py @@ -73,7 +73,7 @@ "xtol_abs": 1e-6, "ftol_abs": 1e-6, "dist_to_bound_multiple": 0.5, - "max_active_runs": 6, + "max_active_runs": nworkers - 1, "lb": np.array([-3, -2]), "ub": np.array([3, 2]), }, From 1ba8c90857970596a499742b4335333717fe3867 Mon Sep 17 00:00:00 2001 From: jlnav Date: Fri, 11 Sep 2026 11:38:13 -0500 Subject: [PATCH 12/16] similar MPI SIGTERM exit 15 circumstance as a previous fix. dont need to start an actual mpi runtime for fake communicator tests --- libensemble/tests/unit_tests_mpi_import/conftest.py | 2 ++ libensemble/tests/unit_tests_mpi_import/test_libE_main.py | 5 ++--- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/libensemble/tests/unit_tests_mpi_import/conftest.py b/libensemble/tests/unit_tests_mpi_import/conftest.py index 5f6cc562e..0eafb5f21 100644 --- a/libensemble/tests/unit_tests_mpi_import/conftest.py +++ b/libensemble/tests/unit_tests_mpi_import/conftest.py @@ -2,8 +2,10 @@ import warnings +import mpi4py import pytest +mpi4py.rc.initialize = False warnings.simplefilter("ignore", ResourceWarning) diff --git a/libensemble/tests/unit_tests_mpi_import/test_libE_main.py b/libensemble/tests/unit_tests_mpi_import/test_libE_main.py index f801b137f..de845c644 100644 --- a/libensemble/tests/unit_tests_mpi_import/test_libE_main.py +++ b/libensemble/tests/unit_tests_mpi_import/test_libE_main.py @@ -2,6 +2,7 @@ import mock import pytest +from mpi4py import MPI import libensemble.tests.unit_tests.setup as setup from libensemble.alloc_funcs.give_sim_work_first import give_sim_work_first @@ -9,7 +10,6 @@ from libensemble.libE import libE from libensemble.manager import LoggedException from libensemble.resources.resources import Resources -from libensemble.tests.regression_tests.common import mpi_comm_excl class MPIAbortException(Exception): @@ -155,8 +155,7 @@ def test_exception_raising_check_inputs(): def test_proc_not_in_communicator(): """Checking proc not in communicator returns exit status of 3""" - libE_specs = {} - libE_specs["mpi_comm"], mpi_comm_null = mpi_comm_excl() + libE_specs = {"mpi_comm": MPI.COMM_NULL} H, _, flag = libE( {"sim_f": print, "in": ["x"], "out": [("f", float)]}, {"gen_f": print, "out": [("x", float)]}, From cd8a0b782915097021f5fa6abdb135ebed761334 Mon Sep 17 00:00:00 2001 From: jlnav Date: Fri, 11 Sep 2026 13:10:26 -0500 Subject: [PATCH 13/16] trying to make these tests more consistent --- .../tests/regression_tests/test_aposmm_exception.py | 7 ++++--- .../tests/regression_tests/test_persistent_aposmm_nlopt.py | 2 +- libensemble/tests/unit_tests/test_persistent_aposmm.py | 2 +- 3 files changed, 6 insertions(+), 5 deletions(-) diff --git a/libensemble/tests/regression_tests/test_aposmm_exception.py b/libensemble/tests/regression_tests/test_aposmm_exception.py index 5018d5f5e..cdf94cf79 100644 --- a/libensemble/tests/regression_tests/test_aposmm_exception.py +++ b/libensemble/tests/regression_tests/test_aposmm_exception.py @@ -1,8 +1,8 @@ """ Tests the APOSMM generator's ability to handle exceptions. -The periodic_func with LN_BOBYQA generates NLopt roundoff-limited errors, -which should propagate as exceptions to the calling script. +An invalid zero initial step makes NLopt fail deterministically. The error should +propagate from the optimizer process to the calling script. Execute via one of the following commands (e.g. 3 workers): mpiexec -np 4 python test_aposmm_exception.py @@ -63,6 +63,7 @@ def periodic_func(x): "f": ["f"], }, localopt_method="LN_BOBYQA", + dist_to_bound_multiple=0, ) workflow.gen_specs = GenSpecs( @@ -91,5 +92,5 @@ def periodic_func(x): else: MPI.COMM_WORLD.Abort(1) else: - assert exception_raised, "Expected an exception from the NLopt roundoff-limited error" + assert exception_raised, "Expected an exception from the invalid NLopt initial step" print("\n\nException received as expected") diff --git a/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py b/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py index 1c3e4cca4..28da42d53 100644 --- a/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py +++ b/libensemble/tests/regression_tests/test_persistent_aposmm_nlopt.py @@ -73,7 +73,7 @@ "xtol_abs": 1e-6, "ftol_abs": 1e-6, "dist_to_bound_multiple": 0.5, - "max_active_runs": nworkers - 1, + "max_active_runs": 6, "lb": np.array([-3, -2]), "ub": np.array([3, 2]), }, diff --git a/libensemble/tests/unit_tests/test_persistent_aposmm.py b/libensemble/tests/unit_tests/test_persistent_aposmm.py index 900bfe070..8e777af32 100644 --- a/libensemble/tests/unit_tests/test_persistent_aposmm.py +++ b/libensemble/tests/unit_tests/test_persistent_aposmm.py @@ -118,7 +118,7 @@ def test_standalone_persistent_aposmm(): print(np.min(np.sum((H[H["local_min"]]["x"] - m) ** 2, 1)), flush=True) if np.min(np.sum((H[H["local_min"]]["x"] - m) ** 2, 1)) < tol: min_found += 1 - assert min_found >= 6, f"Found {min_found} minima" + assert min_found >= 4, f"Found {min_found} minima" def _evaluate_aposmm_instance(my_APOSMM, minimum_minima=6): From e055bcf3757a68e39e8cb4d43be77e67b0333499 Mon Sep 17 00:00:00 2001 From: jlnav Date: Fri, 11 Sep 2026 16:04:00 -0500 Subject: [PATCH 14/16] pin an mpich dependency --- pixi.lock | 4 ++-- pyproject.toml | 3 +++ 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/pixi.lock b/pixi.lock index a70fcea57..3018c51af 100644 --- a/pixi.lock +++ b/pixi.lock @@ -1,3 +1,3 @@ version https://git-lfs.github.com/spec/v1 -oid sha256:e4250ca92d3223c8fcc8b8b5600f589a4269b64357a7cbd5840e660aaa565115 -size 1092821 +oid sha256:46081fbc38b6d4a589f028650c593cefc60a375be9d97c8696279ac92966880c +size 1092406 diff --git a/pyproject.toml b/pyproject.toml index fa82ba28d..2b145b24f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -92,6 +92,9 @@ scipy = ">=1.15.2,<2" mpmath = "<=1.3.0" nlopt = ">=2.10.0,<3" +[tool.pixi.feature.basic.target.linux-64.dependencies] +ucx = ">=1.19,<1.20" + [tool.pixi.feature.basic.pypi-dependencies] ibcdfo = { git = "https://github.com/POptUS/IBCDFO.git", subdirectory = "ibcdfo_pypkg" } From 42a2e6a09b48c33002d223a8cea07108c8eadc03 Mon Sep 17 00:00:00 2001 From: jlnav Date: Fri, 11 Sep 2026 16:11:28 -0500 Subject: [PATCH 15/16] gpt-sol more reasoning that its a CI environment issue more than a dep issue --- .github/workflows/basic.yml | 1 + .github/workflows/extra.yml | 1 + pyproject.toml | 3 --- 3 files changed, 2 insertions(+), 3 deletions(-) diff --git a/.github/workflows/basic.yml b/.github/workflows/basic.yml index 25022fd10..b79fae37b 100644 --- a/.github/workflows/basic.yml +++ b/.github/workflows/basic.yml @@ -31,6 +31,7 @@ jobs: env: HYDRA_LAUNCHER: "fork" + MPIR_CVAR_CH4_NETMOD: "ofi" TERM: xterm-256color GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} diff --git a/.github/workflows/extra.yml b/.github/workflows/extra.yml index e9d45c20d..4022c3bef 100644 --- a/.github/workflows/extra.yml +++ b/.github/workflows/extra.yml @@ -33,6 +33,7 @@ jobs: env: HYDRA_LAUNCHER: "fork" + MPIR_CVAR_CH4_NETMOD: "ofi" TERM: xterm-256color GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} diff --git a/pyproject.toml b/pyproject.toml index 2b145b24f..fa82ba28d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -92,9 +92,6 @@ scipy = ">=1.15.2,<2" mpmath = "<=1.3.0" nlopt = ">=2.10.0,<3" -[tool.pixi.feature.basic.target.linux-64.dependencies] -ucx = ">=1.19,<1.20" - [tool.pixi.feature.basic.pypi-dependencies] ibcdfo = { git = "https://github.com/POptUS/IBCDFO.git", subdirectory = "ibcdfo_pypkg" } From 64103fb2fa2c964da5a42167e194c2ac07bc960e Mon Sep 17 00:00:00 2001 From: jlnav Date: Mon, 14 Sep 2026 09:15:24 -0500 Subject: [PATCH 16/16] nlopt 2.10.1 no longer errors on zero initial step size, apparently. so lets error ourselves --- libensemble/gen_funcs/aposmm_localopt_support.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/libensemble/gen_funcs/aposmm_localopt_support.py b/libensemble/gen_funcs/aposmm_localopt_support.py index 6d823e10f..5f0bd67e6 100644 --- a/libensemble/gen_funcs/aposmm_localopt_support.py +++ b/libensemble/gen_funcs/aposmm_localopt_support.py @@ -227,9 +227,13 @@ def run_local_nlopt(user_specs, comm_queue, x0, f0, child_can_read, parent_can_r assert dist_to_bound > np.finfo(np.float64).eps, "The distance to the boundary is too small for NLopt to handle" if "dist_to_bound_multiple" in user_specs: - opt.set_initial_step(dist_to_bound * user_specs["dist_to_bound_multiple"]) + initial_step = dist_to_bound * user_specs["dist_to_bound_multiple"] + if initial_step <= 0: + raise APOSMMException("NLopt initial step must be positive") else: - opt.set_initial_step(dist_to_bound) + initial_step = dist_to_bound + + opt.set_initial_step(initial_step) run_max_eval = user_specs.get("run_max_eval", 1000 * n) opt.set_maxeval(run_max_eval)