Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
19 changes: 18 additions & 1 deletion khiops/core/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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]
Expand Down
142 changes: 101 additions & 41 deletions khiops/core/internals/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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


Expand Down Expand Up @@ -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 :
Expand All @@ -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:
Expand All @@ -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:
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand All @@ -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:
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -1054,41 +1096,45 @@ 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()

# 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 "
Expand All @@ -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
Expand Down Expand Up @@ -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":

Expand All @@ -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. "
Expand All @@ -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"
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand All @@ -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
Expand Down
Loading
Loading