Coverage for python/lsst/ctrl/bps/initialize.py: 96%

75 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-25 15:17 -0700

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/>. 

27 

28"""Driver for initializing a BPS submission.""" 

29 

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] 

38 

39import getpass 

40import logging 

41import re 

42import shutil 

43from collections.abc import Callable, Iterable 

44from pathlib import Path 

45 

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 

57 

58_LOG = logging.getLogger(__name__) 

59 

60 

61def translate_command_line_values(config: BpsConfig, **kwargs) -> None: 

62 """Override config with command-line values. 

63 

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 

90 

91 _LOG.debug("translate_command_line_values: .bps_cmdline = %s", config[".bps_cmdline"]) 

92 

93 

94def init_submission( 

95 config_file: str, validators: Iterable[Callable[[BpsConfig], None]] = (), **kwargs 

96) -> BpsConfig: 

97 """Initialize BPS configuration and create submit directory. 

98 

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. 

110 

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 ) 

122 

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) 

133 

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) 

146 

147 # Set some initial values 

148 config[".bps_defined.timestamp"] = Instrument.makeCollectionTimestamp() 

149 

150 if "operator" not in config: 

151 config[".bps_defined.operator"] = getpass.getuser() 

152 

153 if "uniqProcName" not in config: 

154 config[".bps_defined.uniqProcName"] = config["outputRun"].replace("/", "_") 

155 

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.") 

171 

172 # Check that can import wms service class. 

173 wms_service_class = doImport(wms_class) 

174 wms_service = wms_service_class(config) 

175 

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.") 

184 

185 # Replace all bpsGenerateConfig 

186 config.generate_config() 

187 

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) 

200 

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) 

205 

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 

210 

211 with open(expanded_config_file, "w") as fh: 

212 config.dump(fh) 

213 

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") 

226 

227 return config 

228 

229 

230def output_run_validator(config: BpsConfig) -> None: 

231 """Check if 'outputRun' is specified in BPS config. 

232 

233 Parameters 

234 ---------- 

235 config : `BpsConfig` 

236 BPS configuration that needs to be validated. 

237 

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'") 

245 

246 

247def submit_path_validator(config: BpsConfig) -> None: 

248 """Check if 'submitPath' is specified in BPS config. 

249 

250 Parameters 

251 ---------- 

252 config : `BpsConfig` 

253 BPS configuration that needs to be validated. 

254 

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'") 

262 

263 

264def out_collection_validator(config: BpsConfig) -> None: 

265 """Check if 'outCollection' is *not* specified in BPS config. 

266 

267 Parameters 

268 ---------- 

269 config : `BpsConfig` 

270 BPS configuration that needs to be validated. 

271 

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'.") 

279 

280 

281def custom_job_validator(config: BpsConfig) -> None: 

282 """Check if 'customJob' is specified in BPS config. 

283 

284 Parameters 

285 ---------- 

286 config : `BpsConfig` 

287 BPS configuration that needs to be validated. 

288 

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'.")