Skip to content
Draft
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
1 change: 1 addition & 0 deletions doc/changes/DM-53944.feature.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Made improvements to how the environment specifications in submit yaml are handled. See ``Job Environment`` section in docs for details and examples.
78 changes: 78 additions & 0 deletions doc/lsst.ctrl.bps/quickstart.rst
Original file line number Diff line number Diff line change
Expand Up @@ -1541,6 +1541,84 @@ invisible to the user. ``bps report`` will still show same labels and
total counts as without ordering. ``cancel`` and ``restart`` will still
work the same.


.. _job-environment:

Job Environment
---------------

One can specify environment values for jobs in the submit yaml. When creating
the job description for the WMS-plugins, BPS will search for environment
settings following normal scoping ordering but performs a special merge
instead of replacing the entire section. If an environment value has
another yaml environment variable, BPS will replace that yaml environment
variable with its value. This will allow normal path chaining
behavior across sections.

.. note::

To not include an environment variable from another section, use the special
string ``BPS_NONE`` as the environment variable value (case-insensitive).

Here are some examples, given the following BPS config yaml what
environment values BPS tells the WMS-plugin to set inside the job:

.. code-block:: YAML

var1: "root_val1"
environment:
VAR2: "root_val2"
VAR3: "root_val3"
VAR_PATH: "${PACKAGE_DIR}/root_dir:${VAR_PATH}"
TEST_VAR: "one {var1} three"
site:
site1:
var1: "site_val1"
environment:
VAR4: "site_val4"
VAR_PATH: "${PACKAGE_DIR}/site_dir:${VAR_PATH}"
cluster:
cl1:
var1: "cl1_val1"
environment:
VAR2: "BPS_NONE"
VAR3: "cl1_val3"
VAR4: "cl1_val4"

#. BPS environment for a cl1 job for site1 :

.. code-block::

VAR3="cl1_val3"
VAR4="cl1_val4"
VAR_PATH="${PACKAGE_DIR}/site_dir:${PACKAGE_DIR}/root_dir:${VAR_PATH}"
TEST_VAR="one cl1_val1 three"

#. BPS environment for a non-cl1 job for site1 :

.. code-block::

VAR2="root_val2"
VAR3="root_val3"
VAR4="site_val4"
VAR_PATH="${PACKAGE_DIR}/site_dir:${PACKAGE_DIR}/root_dir:${VAR_PATH}"
TEST_VAR="one site_val1 three"

#. BPS environment for a non-cli job for some site other than site1:

.. code-block::

VAR2="root_val2"
VAR3="root_val3"
VAR_PATH="${PACKAGE_DIR}/root_dir:${VAR_PATH}"
TEST_VAR="one root_val1 three"

Different WMS-plugins create the job environment via different mechanisms
(e.g., the HTCondor plugin copies the submit environment to the job) or may
not support setting the environment via submit yaml. If the WMS-plugin
supports setting the environment via submit yaml, the job environment must
at least include what is in the submit yaml, but it may include more.

.. _bps-config-generation:

Config Generation
Expand Down
10 changes: 5 additions & 5 deletions python/lsst/ctrl/bps/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,11 +29,7 @@

from astropy import units as u

__all__ = [
"DEFAULT_MEM_FMT",
"DEFAULT_MEM_RETRIES",
"DEFAULT_MEM_UNIT",
]
__all__ = ["BPS_NONE", "DEFAULT_MEM_FMT", "DEFAULT_MEM_RETRIES", "DEFAULT_MEM_UNIT"]


DEFAULT_MEM_RETRIES = 5
Expand All @@ -47,3 +43,7 @@
DEFAULT_MEM_FMT = ".3f"
"""Default format specifier to use when reporting memory consumption.
"""

BPS_NONE = "BPS_NONE"
"""Special yaml value for a None value that can be read and written.
"""
175 changes: 127 additions & 48 deletions python/lsst/ctrl/bps/transform.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,17 @@
import math
import os
import re
from typing import Any

from lsst.ctrl.bps import ClusteredQuantumGraph
from lsst.daf.butler import Config
from lsst.pipe.base import QuantumGraph
from lsst.utils.logging import VERBOSE
from lsst.utils.timer import timeMethod

from . import (
BPS_NONE,
BPS_SEARCH_ORDER,
DEFAULT_MEM_RETRIES,
BpsConfig,
GenericWorkflow,
Expand Down Expand Up @@ -411,45 +415,7 @@ def _get_job_values(config, search_opt, cmd_line_key):
else:
job_values[attr] = getattr(default_gwjob, attr)

# Need to replace all config variables in environment values.
# Also change env vars in environment values to bash syntax.
#
# Note: Because job_values["environment"] is a BpsConfig and
# currently cannot have 2 search objects, for each environment
# setting, we have to get the setting string as is and then
# separately use the overall config to replace values inside
# the setting string.
tmp_job_env = job_values.get("environment", None)
if tmp_job_env:
_LOG.debug("_get_job_values: job_values['environment'] = %s", tmp_job_env)

# Don't want to replace when getting environment setting string.
as_is_search_opt = {
"replaceVars": False,
"expandEnvVars": False,
"replaceEnvBps2Shell": False,
"replaceEnvShell2Bps": False,
}

# When updating environment string, use given search options,
# but ensure making the environment string using bash syntax.
env_search_opt = copy.copy(search_opt)
env_search_opt["replaceVars"] = True # Replace bps config variables.
env_search_opt["replaceEnvBps2Shell"] = False # Replace bps <ENV:var> syntax.
env_search_opt["replaceEnvShell2Bps"] = True # Do not replace shell env syntax.
env_search_opt["expandEnvVars"] = False # Do not replace with submission env value.

job_env = {} # While replacing variables, convert to plain dict.

for name in tmp_job_env:
# Get environment setting string as is.
value = tmp_job_env.search(name, as_is_search_opt)[1]
_LOG.debug("_get_job_values: as is value for %s = %s", name, value)
# Replace config vars and env placeholders
job_env[name] = config.modify_value(name, str(value), env_search_opt)
_LOG.debug("_get_job_values: new env value for %s = %s", name, job_env[name])
# Save new dictionary back with other job values.
job_values["environment"] = job_env
job_values["environment"] = gather_job_environment(config, search_opt)

# If the automatic memory scaling is enabled (i.e. the memory multiplier
# is set and it is a positive number greater than 1.0), adjust number
Expand Down Expand Up @@ -631,8 +597,6 @@ def create_generic_workflow(
_, when_save = config.search("whenSaveJobQgraph", {"default": WhenToSaveQuantumGraphs.TRANSFORM.name})
save_qgraph_per_job = WhenToSaveQuantumGraphs[when_save.upper()]

search_opt = {"replaceVars": False, "expandEnvVars": False, "replaceEnvVars": True, "required": False}

generic_workflow = GenericWorkflow(name)

# Save full run QuantumGraph for use by jobs
Expand Down Expand Up @@ -664,13 +628,15 @@ def create_generic_workflow(
gwjob = GenericWorkflowJob(cluster.name, cluster.label)

# First get job values from cluster or cluster config
search_opt["curvals"] = {"curr_cluster": cluster.label}
found, value = config.search("computeSite", opt=search_opt)
if found:
search_opt["curvals"]["curr_site"] = value
found, value = config.search("computeCloud", opt=search_opt)
if found:
search_opt["curvals"]["curr_cloud"] = value
search_opt = config.get_search_opts(cluster.label)
search_opt.update(
{
"replaceVars": False,
"expandEnvVars": False,
"replaceEnvVars": True,
"required": False,
}
)

# If some config values are set for this cluster
if cluster.label not in cached_job_values:
Expand All @@ -694,13 +660,23 @@ def create_generic_workflow(
_get_job_values(config["cluster"][cluster.label], search_opt, "runQuantumCommand")
)
cluster_job_values = copy.copy(cached_job_values[cluster.label])
# Environment is special because of the way it is merged.
# It needs to be set in the cached_pipetask_values, so
# don't include it here.
cluster_job_values.pop("environment", None)

cluster_job_values["name"] = cluster.name
cluster_job_values["label"] = cluster.label
cluster_job_values["quanta_counts"] = cluster.quanta_counts
cluster_job_values["tags"] = cluster.tags
_LOG.debug("cluster_job_values = %s", cluster_job_values)
_handle_job_values(cluster_job_values, gwjob, cluster_job_values.keys())
_LOG.debug(
"After _handle_job_values cluster %s: gwjob.environment = %s", cluster.label, gwjob.environment
)
_LOG.debug(
"After _handle_job_values cluster %s: gwjob.arguments = %s", cluster.label, gwjob.arguments
)

# For purposes of whether to continue searching for a value is whether
# the value evaluates to False.
Expand All @@ -719,7 +695,19 @@ def create_generic_workflow(
if task_label not in cached_pipetask_values:
search_opt["curvals"]["curr_pipetask"] = task_label
cached_pipetask_values[task_label] = _get_job_values(config, search_opt, "runQuantumCommand")
_LOG.debug(
"cached_pipetask_values[%s]['environment'] = %s",
task_label,
cached_pipetask_values[task_label].get("environment", None),
)
_LOG.debug(
"cached_pipetask_values[%s]['arguments'] = %s",
task_label,
cached_pipetask_values[task_label].get("arguments", None),
)
_handle_job_values(cached_pipetask_values[task_label], gwjob, unset_attributes)
_LOG.debug("After _handle_job_values pipetask: gwjob.environment = %s", gwjob.environment)
_LOG.debug("After _handle_job_values pipetask: gwjob.arguments = %s", gwjob.arguments)

# Update job with workflow attribute and profile values.
qgraph_gwfile = _get_qgraph_gwfile(
Expand All @@ -732,7 +720,10 @@ def create_generic_workflow(
gwjob.cmdvals["qgraphNodeId"] = ",".join(
sorted([f"{node_id}" for node_id in cluster.qgraph_node_ids])
)
_LOG.debug("Before _enhance_command: gwjob.arguments = %s", gwjob.arguments)
_LOG.debug("Before _enhance_command: gwjob.cmdvals = %s", gwjob.cmdvals)
_enhance_command(config, generic_workflow, gwjob, cached_job_values)
_LOG.debug("After _enhance_command: gwjob.environment = %s", gwjob.environment)

# If writing per-job QuantumGraph files during TRANSFORM stage,
# write it now while in memory.
Expand Down Expand Up @@ -948,3 +939,91 @@ def add_final_job_as_sink(generic_workflow, final_job):

generic_workflow.add_job(final_job)
generic_workflow.add_job_relationships(gw_sinks, final_job.name)


def gather_job_environment(config: BpsConfig, search_opt: dict[str, Any]) -> dict[str, str]:
"""Gather environment settings using given config search options.

Parameters
----------
config : `lsst.ctrl.bps.BpsConfig`
Bps configuration.
search_opt : `dict` [`str`, `~typing.Any`]
Config search options.

Returns
-------
environment : `dict` [`str`, `str`]
Dictionary of environment variable names and values.
"""
# Don't want to replace when getting environment setting string.
as_is_search_opt = {
"replaceVars": False,
"expandEnvVars": False,
"replaceEnvBps2Shell": True, # Merge needs it to be shell syntax
"replaceEnvShell2Bps": False,
}

def _update_env(env_updates: dict[str, str], environment: dict[str, str]):
"""Update environment with values from given environment
section.

Parameters
----------
env_updates : `dict` [`str`, `str`] or `~lsst.daf.butler.Config`
New values to use for update.
environment : `dict` [`str`, `str`]
Current environment values to update. Modified in place.
"""
for key, value in env_updates.items():
value = config.modify_value(key, str(value), as_is_search_opt)

for envkey in re.findall(r"\${([^}]+)}", value):
if envkey == key and envkey in environment:
oldval = environment[key]
value = re.sub(rf"\${{{envkey}}}", oldval, value)
environment[key] = value

environment = {}
if ".environment" in config:
_, root_env = config.search("environment", as_is_search_opt)

# Cast to Config to avoid replacing variables.
environment.update(Config(root_env))

curvals = search_opt.get("curvals", {})

# In order to do the concatenation correctly, must
# search in reverse order than normal config searches
for sect in reversed(BPS_SEARCH_ORDER):
sect_key = "curr_" + sect
if sect_key in curvals and sect in config and curvals[sect_key] in config[sect]:
search_sect = config[sect][curvals[sect_key]]
if "environment" in search_sect:
# Cast to Config to avoid replacing variables
_update_env(Config(search_sect["environment"]), environment)

# Because not using normal config search, also have to check the
# search object if given.
if "searchobj" in search_opt and "environment" in search_opt["searchobj"]:
# Cast to Config to avoid replacing variables
env_sect = Config(search_opt["searchobj"]["environment"])
_update_env(env_sect, environment)

# Need to replace all config variables in environment values.
# Also remove any environment variable where value is "BPS_NONE".
new_environment = {}
if environment:
env_search_opt = copy.copy(search_opt)
env_search_opt["replaceVars"] = True # Replace bps config variables.
env_search_opt["replaceEnvBps2Shell"] = False # Keep bps <ENV:var> syntax.
env_search_opt["replaceEnvShell2Bps"] = True # Replace shell env syntax.
env_search_opt["expandEnvVars"] = False # Do not replace with submission env values.
for name in environment:
# Ensure that environment values are strings
new_environment[name] = config.modify_value(name, str(environment[name]), env_search_opt)
if new_environment[name].upper() == BPS_NONE:
_LOG.debug("Removing %s from the merged environment", name)
del new_environment[name]

return new_environment
Loading
Loading