Coverage for python/lsst/ctrl/bps/batch_submit.py: 99%

121 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-22 02:46 -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 to run submit stages as batch jobs.""" 

29 

30__all__ = ["batch_payload_prepare", "batch_submit", "create_batch_stages"] 

31 

32import logging 

33import os 

34from pathlib import Path 

35 

36from lsst.resources import ResourcePath, ResourcePathExpression 

37from lsst.utils.logging import VERBOSE 

38from lsst.utils.timer import time_this, timeMethod 

39 

40from . import ( 

41 DEFAULT_MEM_FMT, 

42 DEFAULT_MEM_UNIT, 

43 BpsConfig, 

44 GenericWorkflow, 

45 GenericWorkflowFile, 

46 GenericWorkflowJob, 

47 GenericWorkflowLazyGroup, 

48) 

49from .bps_utils import _make_id_link 

50from .pre_transform import cluster_quanta, read_quantum_graph 

51from .prepare import prepare 

52from .submit import submit 

53from .transform import _enhance_command, _get_job_values, transform 

54 

55_LOG = logging.getLogger(__name__) 

56 

57 

58@timeMethod(logger=_LOG, logLevel=VERBOSE) 

59def create_batch_stages( 

60 config: BpsConfig, prefix: ResourcePathExpression 

61) -> tuple[GenericWorkflow, BpsConfig]: 

62 """Create a GenericWorkflow that performs the submit stages as a workflow. 

63 

64 Parameters 

65 ---------- 

66 config : `lsst.ctrl.bps.BpsConfig` 

67 BPS configuration. 

68 prefix : `lsst.resources.ResourcePathExpression` 

69 Root path for any output files. 

70 

71 Returns 

72 ------- 

73 generic_workflow : `lsst.ctrl.bps.GenericWorkflow` 

74 The generic workflow transformed from the clustered quantum graph. 

75 generic_workflow_config : `lsst.ctrl.bps.BpsConfig` 

76 Configuration to accompany GenericWorkflow. 

77 """ 

78 prefix = ResourcePath(prefix) 

79 generic_workflow: GenericWorkflow = GenericWorkflow(name=f"{config['uniqProcName']}_ctrl") 

80 generic_workflow.run_attrs.update( 

81 { 

82 "bps_isjob": "True", 

83 "bps_project": config["project"], 

84 "bps_campaign": config["campaign"], 

85 "bps_run": config["uniqProcName"], 

86 "bps_runsite": config["computeSite"], 

87 "bps_operator": config["operator"], 

88 "bps_payload": config["payloadName"], 

89 } 

90 ) 

91 

92 # Save full run QuantumGraph for use by jobs 

93 qgraph_file = GenericWorkflowFile( 

94 "runQgraphFile", 

95 src_uri=config["runQgraphFile"], 

96 wms_transfer=True, 

97 job_access_remote=True, 

98 job_shared=True, 

99 ) 

100 generic_workflow.add_file(qgraph_file) 

101 

102 # Save config file for use by jobs 

103 config_file = GenericWorkflowFile( 

104 "configFile", 

105 src_uri=config["configFile"], 

106 wms_transfer=True, 

107 job_access_remote=False, 

108 job_shared=True, 

109 ) 

110 generic_workflow.add_file(config_file) 

111 

112 # Build QuantumGraph job 

113 build_job = GenericWorkflowJob( 

114 name="buildQuantumGraph", 

115 label="buildQuantumGraph", 

116 ) 

117 search_opt = config.get_search_opts(build_job.label) 

118 search_opt.update( 

119 { 

120 "replaceVars": False, 

121 "expandEnvVars": False, 

122 "replaceEnvVars": True, 

123 "required": False, 

124 } 

125 ) 

126 cmd_line_key = "jobCommand" 

127 job_values = _get_job_values(config, search_opt, cmd_line_key) 

128 if not job_values["executable"]: 

129 raise RuntimeError( 

130 f"Missing executable for buildQuantumGraph. Double check submit yaml for {cmd_line_key}" 

131 ) 

132 for key, value in job_values.items(): 

133 if key not in {"name", "label"}: 

134 setattr(build_job, key, value) 

135 

136 generic_workflow.add_job(build_job) 

137 _LOG.debug("build job's arguments: %s", build_job.arguments) 

138 

139 generic_workflow.add_file(config_file) 

140 generic_workflow.add_job_outputs(build_job.name, [qgraph_file]) 

141 generic_workflow.add_job_inputs(build_job.name, [config_file]) 

142 _enhance_command(config, generic_workflow, build_job, job_values) 

143 _LOG.debug("build job's arguments: %s", build_job.arguments) 

144 

145 # Build cluster/transform/prepare job 

146 prepare_job = GenericWorkflowLazyGroup( 

147 name="preparePayloadWorkflow", 

148 label="preparePayloadWorkflow", 

149 ) 

150 search_opt = config.get_search_opts(prepare_job.label) 

151 search_opt.update( 

152 { 

153 "replaceVars": False, 

154 "expandEnvVars": False, 

155 "replaceEnvVars": True, 

156 "required": False, 

157 } 

158 ) 

159 _LOG.debug("preparePayloadWorkflow search_opt = %s", search_opt) 

160 cmd_line_key = "jobCommand" 

161 job_values = _get_job_values(config, search_opt, cmd_line_key) 

162 if not job_values["executable"]: 

163 raise RuntimeError( 

164 f"Missing executable for preparePayloadWorkflow. Double check submit yaml for {cmd_line_key}" 

165 ) 

166 for key, value in job_values.items(): 

167 if key not in {"name", "label"}: 

168 setattr(prepare_job, key, value) 

169 

170 generic_workflow.add_job(prepare_job, parent_names=["buildQuantumGraph"]) 

171 generic_workflow.add_job_inputs(prepare_job.name, [qgraph_file, config_file]) 

172 _enhance_command(config, generic_workflow, prepare_job, job_values) 

173 

174 _, save_workflow = config.search("saveGenericWorkflow", opt={"default": False}) 

175 if save_workflow: 

176 with prefix.join("bps_stages_generic_workflow.pickle").open("wb") as outfh: 

177 generic_workflow.save(outfh, "pickle") 

178 

179 return generic_workflow, config 

180 

181 

182@timeMethod(logger=_LOG, logLevel=VERBOSE) 

183def batch_payload_prepare(config: BpsConfig, prefix: ResourcePathExpression) -> None: 

184 """Create a GenericWorkflow that performs the submit stages as a workflow. 

185 

186 Parameters 

187 ---------- 

188 config : `lsst.ctrl.bps.BpsConfig` 

189 BPS configuration. 

190 prefix : `lsst.resources.ResourcePathExpression` 

191 Root path for any output files. 

192 

193 Returns 

194 ------- 

195 generic_workflow : `lsst.ctrl.bps.GenericWorkflow` 

196 The generic workflow transformed from the clustered quantum graph. 

197 generic_workflow_config : `lsst.ctrl.bps.BpsConfig` 

198 Configuration to accompany GenericWorkflow. 

199 """ 

200 # Read existing QuantumGraph 

201 qgraph_uri = config["runQgraphFile"] 

202 qgraph = read_quantum_graph(qgraph_uri) 

203 

204 # Cluster 

205 _LOG.info("Starting cluster stage (grouping quanta into jobs)") 

206 with time_this( 

207 log=_LOG, 

208 level=logging.INFO, 

209 prefix=None, 

210 msg="Cluster stage completed", 

211 mem_usage=True, 

212 mem_unit=DEFAULT_MEM_UNIT, 

213 mem_fmt=DEFAULT_MEM_FMT, 

214 ): 

215 clustered_qgraph = cluster_quanta(config, qgraph, config["uniqProcName"]) 

216 

217 _LOG.info("ClusteredQuantumGraph contains %d cluster(s)", len(clustered_qgraph)) 

218 

219 submit_path = config[".bps_defined.submitPath"] 

220 _, save_clustered_qgraph = config.search("saveClusteredQgraph", opt={"default": False}) 

221 if save_clustered_qgraph: 

222 clustered_qgraph.save(os.path.join(submit_path, "bps_clustered_qgraph.pickle")) 

223 _, save_dot = config.search("saveDot", opt={"default": False}) 

224 if save_dot: 

225 clustered_qgraph.draw(os.path.join(submit_path, "bps_clustered_qgraph.dot")) 

226 

227 # Transform 

228 _LOG.info("Starting transform stage (creating generic workflow)") 

229 with time_this( 

230 log=_LOG, 

231 level=logging.INFO, 

232 prefix=None, 

233 msg="Transform stage completed", 

234 mem_usage=True, 

235 mem_unit=DEFAULT_MEM_UNIT, 

236 mem_fmt=DEFAULT_MEM_FMT, 

237 ): 

238 generic_workflow, generic_workflow_config = transform(config, clustered_qgraph, submit_path) 

239 _LOG.info("Generic workflow name '%s'", generic_workflow.name) 

240 

241 num_jobs = sum(generic_workflow.job_counts.values()) 

242 _LOG.info("GenericWorkflow contains %d job(s) (including final)", num_jobs) 

243 

244 # Want to read quantum graph from temp space if told to use it. 

245 gwfile = generic_workflow.get_file("runQgraphFile") 

246 gwfile.wms_transfer = False 

247 # Root computeSite is currently how to specify where the pipeline 

248 # will actually be run. 

249 found, use_run_temp_space = config.search( 

250 "bpsUseRunTempSpace", opt={"curvals": {"curr_site": config[".computeSite"]}} 

251 ) 

252 if found and use_run_temp_space: 

253 found, run_temp_space = config.search( 

254 "fileDistributionEndpoint", opt={"curvals": {"curr_site": config[".computeSite"]}} 

255 ) 

256 if found: 

257 _LOG.debug("run_temp_space = %s", run_temp_space) 

258 gwfile.src_uri = str(Path(run_temp_space) / config["qgraphFileTemplate"]) 

259 else: 

260 raise KeyError("Config is missing fileDistributionEndpoint.") 

261 elif not found: 261 ↛ 264line 261 didn't jump to line 264 because the condition on line 261 was always true

262 _LOG.debug("Config is missing bpsUseRunTempSpace") 

263 

264 _, save_workflow = config.search("saveGenericWorkflow", opt={"default": False}) 

265 if save_workflow: 

266 with open(os.path.join(submit_path, "bps_generic_workflow.pickle"), "wb") as outfh: 

267 generic_workflow.save(outfh, "pickle") 

268 _, save_dot = config.search("saveDot", opt={"default": False}) 

269 if save_dot: 

270 with open(os.path.join(submit_path, "bps_generic_workflow.dot"), "w") as outfh: 

271 generic_workflow.draw(outfh, "dot") 

272 

273 # Prepare 

274 _LOG.info("Starting prepare stage (creating specific implementation of workflow)") 

275 with time_this( 

276 log=_LOG, 

277 level=logging.INFO, 

278 prefix=None, 

279 msg="Prepare stage completed", 

280 mem_usage=True, 

281 mem_unit=DEFAULT_MEM_UNIT, 

282 mem_fmt=DEFAULT_MEM_FMT, 

283 ): 

284 wms_workflow = prepare(generic_workflow_config, generic_workflow, submit_path) 

285 

286 # Add payload workflow to currently running workflow 

287 _LOG.info("Starting update workflow") 

288 with time_this( 

289 log=_LOG, 

290 level=logging.INFO, 

291 prefix=None, 

292 msg="Workflow update completed", 

293 mem_usage=True, 

294 mem_unit=DEFAULT_MEM_UNIT, 

295 mem_fmt=DEFAULT_MEM_FMT, 

296 ): 

297 # Assuming submit_path for ctrl workflow is visible by this job. 

298 wms_workflow.add_to_parent_workflow(generic_workflow_config) 

299 

300 

301def batch_submit(config: BpsConfig): 

302 """Submit a workflow for execution with preparation done in batch jobs. 

303 

304 Parameters 

305 ---------- 

306 config : `lsst.ctrl.bps.BpsConfig` 

307 BPS configuration. 

308 

309 Returns 

310 ------- 

311 wms_workflow : `lsst.ctrl.bps.BaseWmsWorkflow` 

312 Submitted workflow. 

313 """ 

314 submit_path = config[".bps_defined.submitPath"] 

315 

316 _LOG.info("Starting to create control workflow") 

317 with time_this( 

318 log=_LOG, 

319 level=logging.INFO, 

320 prefix=None, 

321 msg="Creation completed", 

322 mem_usage=True, 

323 mem_unit=DEFAULT_MEM_UNIT, 

324 mem_fmt=DEFAULT_MEM_FMT, 

325 ): 

326 generic_workflow, config = create_batch_stages(config, submit_path) 

327 

328 _LOG.info("Starting to prepare control workflow") 

329 with time_this( 

330 log=_LOG, 

331 level=logging.INFO, 

332 prefix=None, 

333 msg="Preparation completed", 

334 mem_usage=True, 

335 mem_unit=DEFAULT_MEM_UNIT, 

336 mem_fmt=DEFAULT_MEM_FMT, 

337 ): 

338 wms_workflow = prepare(config, generic_workflow, submit_path) 

339 

340 _, dry_run = config.search("dryRun", opt={"default": False}) 

341 if not dry_run: 

342 _LOG.info("Starting to submit control workflow") 

343 with time_this( 

344 log=_LOG, 

345 level=logging.INFO, 

346 prefix=None, 

347 msg="Submission completed", 

348 mem_usage=True, 

349 mem_unit=DEFAULT_MEM_UNIT, 

350 mem_fmt=DEFAULT_MEM_FMT, 

351 ): 

352 submit(config, wms_workflow) 

353 _LOG.info("Run '%s' submitted for execution with id '%s'", wms_workflow.name, wms_workflow.run_id) 

354 

355 _make_id_link(config, wms_workflow.run_id) 

356 

357 return wms_workflow