diff --git a/CHANGELOG.md b/CHANGELOG.md index f15c72dc..5fe0397d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,9 @@ ### Fixed - (`core`) Fix a RuntimeException error when automatically correcting deprecated data paths in task args +### Changed +- (`core`) The `max_cores` parameter limits also the number of CPU cores pre-allocated for each training run. + ## 11.0.1.0 - 2026-07-02 ### Added diff --git a/khiops/core/api.py b/khiops/core/api.py index 3774c4c0..62071282 100644 --- a/khiops/core/api.py +++ b/khiops/core/api.py @@ -29,7 +29,7 @@ is_string_like, type_error_message, ) -from khiops.core.internals.runner import get_runner +from khiops.core.internals.runner import KhiopsLocalRunner, get_runner, set_runner from khiops.core.internals.task import get_task_registry # Construction rules @@ -143,6 +143,17 @@ def _run_task(task_name, task_args): # Obtain the api function from the registry task = get_task_registry().get_task(task_name, get_khiops_version()) + # Instantiate a new runner with a specific environment if max_cores is defined + if system_settings.max_cores is not None: + # Ensure the previous runner is a KhiopsLocalRunner + # to avoid overwriting a mocked runner + if isinstance(get_runner(), KhiopsLocalRunner): + khiops_env = os.environ.copy() + # The max pre-allocated cpu cores must have the same value + khiops_env["KHIOPS_PROC_NUMBER"] = str(system_settings.max_cores) + khiops_runner = KhiopsLocalRunner(khiops_env) + set_runner(khiops_runner) + # Execute the Khiops task and cleanup when necessary try: get_runner().run( @@ -199,6 +210,12 @@ def _preprocess_arguments(args): if arg == "max_cores": max_cores = args[arg] if max_cores is not None: + # This `max_cores` system setting will be used in the khiops scenario + # to limit the CPU cores to use for the training. + # An additional environment variable (local to this specific run) + # MUST also be set in the runner to avoid allocating + # all the available CPU cores. + # Thus, allocated CPU cores = max number of CPU cores used system_settings.max_cores = int(max_cores) elif arg == "memory_limit_mb": memory_limit_mb = args[arg] diff --git a/khiops/core/internals/runner.py b/khiops/core/internals/runner.py index e210a3c5..6f235fc1 100644 --- a/khiops/core/internals/runner.py +++ b/khiops/core/internals/runner.py @@ -50,7 +50,14 @@ def _isdir_without_all_perms(dir_path): ) -def get_default_samples_dir(): +def _current_environment(environment=None): + """Returns the provided environment or the process environment.""" + if environment is not None: + return environment + return os.environ + + +def get_default_samples_dir(environment=None): """Returns the default samples directory The default samples directory is computed according to the following priorities: @@ -60,16 +67,20 @@ def get_default_samples_dir(): - `%USERPROFILE%\\khiops_data\\samples` otherwise - Linux/macOS: `$HOME/khiops_data/samples` """ - if "KHIOPS_SAMPLES_DIR" in os.environ and os.environ["KHIOPS_SAMPLES_DIR"]: - samples_dir = os.environ["KHIOPS_SAMPLES_DIR"] - elif platform.system() == "Windows" and "PUBLIC" in os.environ: - samples_dir = os.path.join(os.environ["PUBLIC"], "khiops_data", "samples") + environment = _current_environment(environment) + if "KHIOPS_SAMPLES_DIR" in environment and environment["KHIOPS_SAMPLES_DIR"]: + samples_dir = environment["KHIOPS_SAMPLES_DIR"] + elif platform.system() == "Windows" and "PUBLIC" in environment: + samples_dir = os.path.join(environment["PUBLIC"], "khiops_data", "samples") else: # The filesystem abstract layer is used here # as the path can be either local or remote - samples_dir = fs.get_child_path( - fs.get_child_path(os.environ["HOME"], "khiops_data"), "samples" - ) + if "HOME" in environment: + samples_dir = fs.get_child_path( + fs.get_child_path(environment["HOME"], "khiops_data"), "samples" + ) + else: + raise KeyError("HOME") return samples_dir @@ -237,7 +248,7 @@ def _check_conda_env_bin_dir(conda_env_bin_dir): return is_conda_env_bin_dir -def _infer_khiops_installation_method(trace=False): +def _infer_khiops_installation_method(environment=None, trace=False): """Returns the Khiops installation method Definitions : @@ -255,6 +266,8 @@ def _infer_khiops_installation_method(trace=False): - or in a classical virtual environment (highly encouraged) """ + environment = _current_environment(environment) + # We are in a Conda environment if # - the CONDA_PREFIX environment variable exists and, # - the khiops_env script exists within: @@ -263,8 +276,8 @@ def _infer_khiops_installation_method(trace=False): # Note: The check that the Khiops binaries are actually executable is done # afterwards by the initializations method. installation_method = "unknown" - if "CONDA_PREFIX" in os.environ: - conda_env_dir = os.environ["CONDA_PREFIX"] + if "CONDA_PREFIX" in environment: + conda_env_dir = environment["CONDA_PREFIX"] if platform.system() == "Windows": conda_binary_dir = os.path.join(conda_env_dir, "Library", "bin") else: @@ -335,13 +348,22 @@ def _get_current_library_installer(): return "unknown" -def _build_khiops_process_environment(): +def _build_khiops_process_environment(environment=None): """Build a specific environment used for the execution of khiops in a process This environment can be modified freely without interfering with the global one. + + Parameters + ---------- + + environment: dict + optional source of the environment """ - khiops_env = os.environ.copy() + if environment is not None: + khiops_env = environment.copy() + else: + khiops_env = os.environ.copy() # Ensure HOME is always set for OpenMPI 5+ # (using KHIOPS_MPI_HOME if it exists) @@ -946,8 +968,18 @@ class KhiopsLocalRunner(KhiopsRunner): """ - def __init__(self): + def __init__(self, environment=None): + """Initialize a local runner. + + Parameters + ---------- + environment : dict, optional + Environment owned by this runner. If omitted, initialization keeps + using the process environment for compatibility. + """ + # Define specific attributes + self._environment = environment.copy() if environment is not None else None self._mpi_command_args = None self._khiops_path = None self._khiops_coclustering_path = None @@ -962,7 +994,10 @@ def __init__(self): self._initialize_khiops_environment() def _initialize_khiops_environment(self): - installation_method = _infer_khiops_installation_method() + runner_environment = _current_environment(self._environment) + installation_method = _infer_khiops_installation_method( + environment=runner_environment + ) match installation_method: # In conda-based environments, khiops_env is not in PATH; # its location must be inferred from the conda env directory. @@ -974,7 +1009,9 @@ def _initialize_khiops_environment(self): khiops_env_path += ".cmd" # In an activated conda environment, khiops_env is in PATH. case "conda": - khiops_env_path = self._infer_khiops_env_from_path(installation_method) + khiops_env_path = self._infer_khiops_env_from_path( + installation_method, runner_environment + ) case "pip": # Ensure the binary dependency is still installed. try: @@ -1011,18 +1048,23 @@ def _initialize_khiops_environment(self): # which is in PATH. else: khiops_env_path = self._infer_khiops_env_from_path( - installation_method + installation_method, runner_environment ) case _: raise KhiopsEnvironmentError( f"Unknown installation method '{installation_method}'." ) + khiops_env_process_options = { + "stdout": subprocess.PIPE, + "stderr": subprocess.PIPE, + "universal_newlines": True, + } + if self._environment is not None: + khiops_env_process_options["env"] = runner_environment.copy() + with subprocess.Popen( - [khiops_env_path, "--env"], - stdout=subprocess.PIPE, - stderr=subprocess.PIPE, - universal_newlines=True, + [khiops_env_path, "--env"], **khiops_env_process_options ) as khiops_env_process: stdout, stderr = khiops_env_process.communicate() if khiops_env_process.returncode != 0: @@ -1054,32 +1096,33 @@ def _initialize_khiops_environment(self): # to prepend to or to set to HOME for OpenMPI 5+ # when running khiops core if var_name == "KHIOPS_MPI_HOME": - os.environ["KHIOPS_MPI_HOME"] = var_value + runner_environment["KHIOPS_MPI_HOME"] = var_value # Set paths to Khiops binaries elif var_name == "KHIOPS_PATH": self.khiops_path = var_value - os.environ["KHIOPS_PATH"] = var_value + runner_environment["KHIOPS_PATH"] = var_value elif var_name == "KHIOPS_COCLUSTERING_PATH": self.khiops_coclustering_path = var_value - os.environ["KHIOPS_COCLUSTERING_PATH"] = var_value + runner_environment["KHIOPS_COCLUSTERING_PATH"] = var_value # Set MPI command elif var_name == "KHIOPS_MPI_COMMAND": self._mpi_command_args = shlex.split(var_value) - os.environ["KHIOPS_MPI_COMMAND"] = var_value + runner_environment["KHIOPS_MPI_COMMAND"] = var_value # On Windows, in 'pip' installations, # "KHIOPS_MPI_DLL_PATH" (containing the Intel MPI Library) # must be added to "PATH" otherwise Khiops wouldn't find it # and fail immediately elif installation_method == "pip" and var_name == "KHIOPS_MPI_DLL_PATH": - os.environ["PATH"] = os.pathsep.join( - [var_value, os.environ.get("PATH")] + runner_environment["PATH"] = os.pathsep.join( + [var_value, runner_environment.get("PATH")] ) + runner_environment["KHIOPS_MPI_DLL_PATH"] = var_value # Propagate all the other environment variables to Khiops binaries else: - os.environ[var_name] = var_value + runner_environment[var_name] = var_value # Set KHIOPS_API_MODE to `true` - os.environ["KHIOPS_API_MODE"] = "true" + runner_environment["KHIOPS_API_MODE"] = "true" # Check the tools exist and are executable self._check_tools() @@ -1087,8 +1130,11 @@ def _initialize_khiops_environment(self): # Initialize the default samples dir self._initialize_default_samples_dir() - def _infer_khiops_env_from_path(self, installation_method): - khiops_env_path = shutil.which("khiops_env") + def _infer_khiops_env_from_path(self, installation_method, environment=None): + # Fall back to searching `khiops_env` according to the `PATH` set in + # `os.environ`. + search_path = environment.get("PATH") if environment is not None else None + khiops_env_path = shutil.which("khiops_env", path=search_path) if khiops_env_path is None: raise KhiopsEnvironmentError( "The 'khiops_env' script not found for the current " @@ -1100,7 +1146,7 @@ def _infer_khiops_env_from_path(self, installation_method): def _initialize_default_samples_dir(self): """See class docstring""" - samples_dir = get_default_samples_dir() + samples_dir = get_default_samples_dir(self._environment) _check_samples_dir(samples_dir) self._samples_dir = samples_dir assert self._samples_dir is not None @@ -1162,7 +1208,10 @@ def _detect_library_installation_incompatibilities(self, library_root_dir_path): error_list = [] warning_list = [] - installation_method = _infer_khiops_installation_method() + runner_environment = _current_environment(self._environment) + installation_method = _infer_khiops_installation_method( + environment=runner_environment + ) # activated 'conda' installation if installation_method == "conda": @@ -1182,14 +1231,14 @@ def _detect_library_installation_incompatibilities(self, library_root_dir_path): ) warning_list.append(warning) - conda_prefix_path = Path(os.environ["CONDA_PREFIX"]) + conda_prefix_path = Path(runner_environment["CONDA_PREFIX"]) # the conda environment must match the library installation if not library_root_dir_path.is_relative_to(conda_prefix_path): error = ( "Khiops Python library installation " f"path '{library_root_dir_path}' " "does not match the current Conda environment " - f"'{os.environ['CONDA_PREFIX']}'. " + f"'{runner_environment['CONDA_PREFIX']}'. " "Either deactivate the current Conda environment " "or use the Khiops Python library " "belonging to the current Conda environment. " @@ -1203,7 +1252,7 @@ def _detect_library_installation_incompatibilities(self, library_root_dir_path): error = ( f"Khiops binary path '{self.khiops_path}' " "does not match the current Conda environment " - f"'{os.environ['CONDA_PREFIX']}'. " + f"'{runner_environment['CONDA_PREFIX']}'. " "We recommend installing the Khiops binary " "in the current Conda environment. " "Go to https://khiops.org for instructions.\n" @@ -1335,7 +1384,10 @@ def _build_status_message(self): ) # Build the messages for install type and mpi - install_type_msg = _infer_khiops_installation_method() + runner_environment = _current_environment(self._environment) + install_type_msg = _infer_khiops_installation_method( + environment=runner_environment + ) if self._mpi_command_args: mpi_command_args_msg = " ".join(self._mpi_command_args) else: @@ -1466,7 +1518,13 @@ def _get_samples_dir(self): self._samples_dir_checked = True return self._samples_dir - def raw_run(self, tool_name, command_line_args=None, use_mpi=True, trace=False): + def raw_run( + self, + tool_name, + command_line_args=None, + use_mpi=True, + trace=False, + ): """Execute a Khiops tool with given command line arguments Parameters @@ -1521,9 +1579,9 @@ def raw_run(self, tool_name, command_line_args=None, use_mpi=True, trace=False): print(f"Khiops execution call: {khiops_call}") # Build custom Khiops process environment - # which makes sure HOME is defined and set + # which makes sure for example HOME is defined and set # according to khiops_env's KHIOPS_MPI_HOME - khiops_env = _build_khiops_process_environment() + khiops_env = _build_khiops_process_environment(self._environment) # Execute the process with subprocess.Popen( @@ -1550,7 +1608,9 @@ def _run( # Execute the tool khiops_args = command_line_options.build_command_line_options(scenario_path) stdout, stderr, return_code = self.raw_run( - tool_name, command_line_args=khiops_args, trace=trace + tool_name, + command_line_args=khiops_args, + trace=trace, ) return return_code, stdout, stderr diff --git a/tests/test_core.py b/tests/test_core.py index facdc667..fdb64186 100644 --- a/tests/test_core.py +++ b/tests/test_core.py @@ -30,6 +30,7 @@ from khiops.core.internals.runner import KhiopsLocalRunner, KhiopsRunner from khiops.core.internals.scenario import ConfigurableKhiopsScenario from khiops.core.internals.version import KhiopsVersion +from tests.test_helper import KhiopsTestHelper # Disable warning about access to protected member: These are tests # pylint: disable=protected-access @@ -2673,9 +2674,11 @@ def run( self, task, task_args, - command_line_options, + command_line_options=None, trace=False, system_settings=None, + stdout_file_path="", + stderr_file_path="", force_ansi_scenario=False, **kwargs, ): @@ -3233,6 +3236,58 @@ def test_raise_exception_on_error_case_without_a_message(self): output_msg = str(context.exception) self.assertEqual(output_msg, expected_msg) + def test_max_cores_param_sets_mpi_proc_number(self): + # Prepare to collect the runner attributes + runner_attributes_trace = KhiopsTestHelper.create_parameter_trace() + KhiopsTestHelper.wrap_with_runner_attributes_trace( + "khiops.core.internals.runner", "KhiopsRunner.run", runner_attributes_trace + ) + # Set the file paths + dictionary_file_path = os.path.join(kh.get_samples_dir(), "Adult", "Adult.kdic") + data_table_path = os.path.join(kh.get_samples_dir(), "Adult", "Adult.txt") + report_file_path = os.path.join( + "kh_samples", "train_predictor_file_paths", "AnalysisResults.khj" + ) + # The existing MockedRunnerContext class is not used here + # as it creates a new KhiopsLocalRunner instance that does not fit our needs + with ( + mock.patch.object( + KhiopsLocalRunner, + "raw_run", + create_mocked_raw_run( + stdout=False, # ask for an empty stdout + stderr=False, # ask for an empty stderr + return_code=0, # zero error code (success) + ), + ), + mock.patch.object( + KhiopsLocalRunner, + "_get_khiops_version", + return_value=KhiopsVersion("11.0.1-rc1.0"), + ), + ): + # Train the predictor + kh.train_predictor( + dictionary_file_path_or_domain=dictionary_file_path, + dictionary_name="Adult", + data_table_path=data_table_path, + target_variable="class", + analysis_report_file_path=report_file_path, + trace=True, + max_cores=17, + ) + index_of_core_flag = runner_attributes_trace["_mpi_command_args"].index("-n") + self.assertTrue(index_of_core_flag >= 0, msg="The cpu cores flag must be set") + self.assertTrue( + index_of_core_flag < len(runner_attributes_trace["_mpi_command_args"]), + msg="The cpu cores flag must be followed by its value", + ) + self.assertEqual( + "17", + runner_attributes_trace["_mpi_command_args"][index_of_core_flag + 1], + msg="The cpu cores value must match the max_cores input parameter", + ) + class LocalFileSystemTests(unittest.TestCase): """Test the methods of the `LocalFileSystem`""" diff --git a/tests/test_helper.py b/tests/test_helper.py index f328bac1..64c346f6 100644 --- a/tests/test_helper.py +++ b/tests/test_helper.py @@ -285,6 +285,24 @@ def wrapper(wrapped, _instance, args, kwargs): return wrapper + @staticmethod + def wrap_with_runner_attributes_trace(module, function, runner_attributes_trace): + """Wrap function with runner attributes trace""" + + @wrapt.patch_function_wrapper(module, function) + def wrapper(wrapped, _instance, args, kwargs): + # mutate runner_attributes_trace as previously bound / initialized in the + # outer scope of the `wrap_with_parameter_trace` method by its + # caller: + nonlocal runner_attributes_trace + + # Collect all the updated runner attributes + runner_attributes_trace.update(vars(_instance)) + + return wrapped(*args, **kwargs) + + return wrapper + @staticmethod def get_resources_dir(): """Helper to get the directory containing the fixtures"""