Coverage for python/lsst/ctrl/bps/initialize.py: 96%
75 statements
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-26 09:18 +0000
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-26 09:18 +0000
1# This file is part of ctrl_bps.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This software is dual licensed under the GNU General Public License and also
10# under a 3-clause BSD license. Recipients may choose which of these licenses
11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt,
12# respectively. If you choose the GPL option then the following text applies
13# (but note that there is still no warranty even if you opt for BSD instead):
14#
15# This program is free software: you can redistribute it and/or modify
16# it under the terms of the GNU General Public License as published by
17# the Free Software Foundation, either version 3 of the License, or
18# (at your option) any later version.
19#
20# This program is distributed in the hope that it will be useful,
21# but WITHOUT ANY WARRANTY; without even the implied warranty of
22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
23# GNU General Public License for more details.
24#
25# You should have received a copy of the GNU General Public License
26# along with this program. If not, see <https://www.gnu.org/licenses/>.
28"""Driver for initializing a BPS submission."""
30__all__ = [
31 "custom_job_validator",
32 "init_submission",
33 "out_collection_validator",
34 "output_run_validator",
35 "submit_path_validator",
36 "translate_command_line_values",
37]
39import getpass
40import logging
41import re
42import shutil
43from collections.abc import Callable, Iterable
44from pathlib import Path
46from lsst.ctrl.bps import (
47 BPS_DEFAULTS,
48 BPS_SEARCH_ORDER,
49 DEFAULT_MEM_FMT,
50 DEFAULT_MEM_UNIT,
51 BpsConfig,
52)
53from lsst.ctrl.bps.bps_utils import _dump_env_info, _dump_pkg_info, mkdir
54from lsst.pipe.base import Instrument
55from lsst.utils import doImport
56from lsst.utils.timer import time_this
58_LOG = logging.getLogger(__name__)
61def translate_command_line_values(config: BpsConfig, **kwargs) -> None:
62 """Override config with command-line values.
64 Parameters
65 ----------
66 config : `lsst.ctrl.bps.BpsConfig`
67 BPS configuration.
68 **kwargs
69 Additional modifiers to the configuration from the command line.
70 """
71 _LOG.debug("translate_command_line_values: kwargs = %s", kwargs)
72 # Handle diffs between pipetask argument names vs bps yaml
73 translation = {
74 "input": "inCollection",
75 "output_run": "outputRun",
76 "qgraph": "qgraphFile",
77 "pipeline": "pipelineYaml",
78 "wms_service": "wmsServiceClass",
79 "compute_site": "computeSite",
80 }
81 for key, value in kwargs.items():
82 # Don't want to override config with None or empty string values.
83 if value:
84 # pipetask argument parser converts some values to list,
85 # but bps will want string.
86 if not isinstance(value, str) and isinstance(value, Iterable):
87 value = ",".join(value)
88 new_key = translation.get(key, re.sub(r"_(\S)", lambda match: match.group(1).upper(), key))
89 config[f".bps_cmdline.{new_key}"] = value
91 _LOG.debug("translate_command_line_values: .bps_cmdline = %s", config[".bps_cmdline"])
94def init_submission(
95 config_file: str, validators: Iterable[Callable[[BpsConfig], None]] = (), **kwargs
96) -> BpsConfig:
97 """Initialize BPS configuration and create submit directory.
99 Parameters
100 ----------
101 config_file : `str`
102 Name of the configuration file.
103 validators : `~collections.abc.Iterable` \
104 [`~collections.abc.Callable` [[`BpsConfig`], `None`]], optional
105 A list of functions performing checks on the given configuration.
106 Each function should take a single argument, a BpsConfig object, and
107 raise if the check fails. By default, no checks are performed.
108 **kwargs
109 Additional modifiers to the configuration.
111 Returns
112 -------
113 config : `lsst.ctrl.bps.BpsConfig`
114 Batch Processing Service configuration.
115 """
116 config = BpsConfig(
117 config_file,
118 search_order=BPS_SEARCH_ORDER,
119 defaults=BPS_DEFAULTS,
120 wms_service_class_fqn=kwargs.get("wms_service"),
121 )
123 with time_this(
124 log=_LOG,
125 level=logging.DEBUG,
126 prefix=None,
127 msg="Translating command line values completed",
128 mem_usage=True,
129 mem_unit=DEFAULT_MEM_UNIT,
130 mem_fmt=DEFAULT_MEM_FMT,
131 ):
132 translate_command_line_values(config, **kwargs)
134 with time_this(
135 log=_LOG,
136 level=logging.DEBUG,
137 prefix=None,
138 msg="Validation tests completed",
139 mem_usage=True,
140 mem_unit=DEFAULT_MEM_UNIT,
141 mem_fmt=DEFAULT_MEM_FMT,
142 ):
143 # Run validation tests on the given config if any.
144 for validator in validators:
145 validator(config)
147 # Set some initial values
148 config[".bps_defined.timestamp"] = Instrument.makeCollectionTimestamp()
150 if "operator" not in config:
151 config[".bps_defined.operator"] = getpass.getuser()
153 if "uniqProcName" not in config:
154 config[".bps_defined.uniqProcName"] = config["outputRun"].replace("/", "_")
156 with time_this(
157 log=_LOG,
158 level=logging.DEBUG,
159 prefix=None,
160 msg="Submission tests completed",
161 mem_usage=True,
162 mem_unit=DEFAULT_MEM_UNIT,
163 mem_fmt=DEFAULT_MEM_FMT,
164 ):
165 # If requested, run WMS plugin checks early in submission process to
166 # ensure WMS has what it will need for prepare() or submit().
167 if kwargs.get("runWmsSubmissionChecks", False):
168 found, wms_class = config.search("wmsServiceClass")
169 if not found:
170 raise KeyError("Missing wmsServiceClass in bps config. Aborting.")
172 # Check that can import wms service class.
173 wms_service_class = doImport(wms_class)
174 wms_service = wms_service_class(config)
176 try:
177 wms_service.run_submission_checks()
178 except NotImplementedError:
179 # Allow various plugins to implement only when needed to do
180 # extra checks.
181 _LOG.debug("run_submission_checks is not implemented in %s.", wms_class)
182 else:
183 _LOG.debug("Skipping submission checks.")
185 # Replace all bpsGenerateConfig
186 config.generate_config()
188 with time_this(
189 log=_LOG,
190 level=logging.DEBUG,
191 prefix=None,
192 msg="Submit path created",
193 mem_usage=True,
194 mem_unit=DEFAULT_MEM_UNIT,
195 mem_fmt=DEFAULT_MEM_FMT,
196 ):
197 # Make submit directory to contain all outputs.
198 submit_path = mkdir(config["submitPath"])
199 config[".bps_defined.submitPath"] = str(submit_path)
201 # Pre-make quantum graph filename
202 prefix = Path(submit_path)
203 qgraph_filename = prefix / config["qgraphFileTemplate"]
204 config[".bps_defined.runQgraphFile"] = str(qgraph_filename)
206 # save copy of configs (orig and expanded config)
207 shutil.copy2(config_file, submit_path)
208 expanded_config_file = f"{submit_path}/{config['uniqProcName']}_config.yaml"
209 config[".bps_defined.configFile"] = expanded_config_file
211 with open(expanded_config_file, "w") as fh:
212 config.dump(fh)
214 with time_this(
215 log=_LOG,
216 level=logging.DEBUG,
217 prefix=None,
218 msg="Saved environment and package information",
219 mem_usage=True,
220 mem_unit=DEFAULT_MEM_UNIT,
221 mem_fmt=DEFAULT_MEM_FMT,
222 ):
223 # Dump information about runtime environment and software versions.
224 _dump_env_info(f"{submit_path}/{config['uniqProcName']}.env.info.yaml")
225 _dump_pkg_info(f"{submit_path}/{config['uniqProcName']}.pkg.info.yaml")
227 return config
230def output_run_validator(config: BpsConfig) -> None:
231 """Check if 'outputRun' is specified in BPS config.
233 Parameters
234 ----------
235 config : `BpsConfig`
236 BPS configuration that needs to be validated.
238 Raises
239 ------
240 KeyError
241 Raised if 'outputRun' is not specified in the BPS configuration.
242 """
243 if "outputRun" not in config:
244 raise KeyError("Must specify the output run collection using 'outputRun'")
247def submit_path_validator(config: BpsConfig) -> None:
248 """Check if 'submitPath' is specified in BPS config.
250 Parameters
251 ----------
252 config : `BpsConfig`
253 BPS configuration that needs to be validated.
255 Raises
256 ------
257 KeyError
258 Raised if 'submitPath' is not specified in the BPS configuration.
259 """
260 if "submitPath" not in config:
261 raise KeyError("Must specify the submit-side run directory using 'submitPath'")
264def out_collection_validator(config: BpsConfig) -> None:
265 """Check if 'outCollection' is *not* specified in BPS config.
267 Parameters
268 ----------
269 config : `BpsConfig`
270 BPS configuration that needs to be validated.
272 Raises
273 ------
274 KeyError
275 Raised if 'outCollection' *is* specified in the BPS configuration.
276 """
277 if "outCollection" in config:
278 raise KeyError("'outCollection' is deprecated. Replace all references to it with 'outputRun'.")
281def custom_job_validator(config: BpsConfig) -> None:
282 """Check if 'customJob' is specified in BPS config.
284 Parameters
285 ----------
286 config : `BpsConfig`
287 BPS configuration that needs to be validated.
289 Raises
290 ------
291 KeyError
292 Raised if 'customJob' is not specified in the BPS configuration.
293 """
294 if "customJob" not in config and "executable" not in config["customJob"]:
295 raise KeyError("Must specify the details of the script to execute using 'customJob'.")