Coverage for python/lsst/ctrl/bps/transform.py: 80%

328 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-01 09:22 +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/>. 

27 

28"""Driver for the transformation of a QuantumGraph into a generic workflow.""" 

29 

30import copy 

31import dataclasses 

32import logging 

33import math 

34import os 

35import re 

36 

37from lsst.ctrl.bps import ClusteredQuantumGraph 

38from lsst.pipe.base import QuantumGraph 

39from lsst.utils.logging import VERBOSE 

40from lsst.utils.timer import timeMethod 

41 

42from . import ( 

43 DEFAULT_MEM_RETRIES, 

44 BpsConfig, 

45 GenericWorkflow, 

46 GenericWorkflowExec, 

47 GenericWorkflowFile, 

48 GenericWorkflowJob, 

49) 

50from .bps_utils import ( 

51 WhenToSaveQuantumGraphs, 

52 create_job_quantum_graph_filename, 

53 save_qg_subgraph, 

54) 

55 

56# All available job attributes. 

57_ATTRS_ALL = frozenset([field.name for field in dataclasses.fields(GenericWorkflowJob)]) 

58 

59# Job attributes that need to be set to their maximal value in the cluster. 

60_ATTRS_MAX = frozenset( 

61 { 

62 "memory_multiplier", 

63 "number_of_retries", 

64 "request_cpus", 

65 "request_memory", 

66 "request_memory_max", 

67 } 

68) 

69 

70# Job attributes that need to be set to sum of their values in the cluster. 

71_ATTRS_SUM = frozenset( 

72 { 

73 "request_disk", 

74 "request_walltime", 

75 } 

76) 

77 

78# Job attributes do not fall into a specific category 

79_ATTRS_MISC = frozenset( 

80 { 

81 "label", # taskDef labels aren't same in job and may not match job label 

82 "cmdvals", 

83 "profile", 

84 "attrs", 

85 } 

86) 

87 

88# Attributes that need to be the same for each quanta in the cluster. 

89_ATTRS_UNIVERSAL = frozenset(_ATTRS_ALL - (_ATTRS_MAX | _ATTRS_MISC | _ATTRS_SUM)) 

90 

91_LOG = logging.getLogger(__name__) 

92 

93 

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

95def transform( 

96 config: BpsConfig, cqgraph: ClusteredQuantumGraph, prefix: str 

97) -> tuple[GenericWorkflow, BpsConfig]: 

98 """Transform a ClusteredQuantumGraph to a GenericWorkflow. 

99 

100 Parameters 

101 ---------- 

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

103 BPS configuration. 

104 cqgraph : `lsst.ctrl.bps.ClusteredQuantumGraph` 

105 A clustered quantum graph to transform into a generic workflow. 

106 prefix : `str` 

107 Root path for any output files. 

108 

109 Returns 

110 ------- 

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

112 The generic workflow transformed from the clustered quantum graph. 

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

114 Configuration to accompany GenericWorkflow. 

115 """ 

116 if cqgraph.name is not None: 

117 name = cqgraph.name 

118 else: 

119 _, name = config.search("uniqProcName", opt={"required": True}) 

120 

121 generic_workflow = create_generic_workflow(config, cqgraph, name, prefix) 

122 generic_workflow_config = create_generic_workflow_config(config, prefix) 

123 

124 return generic_workflow, generic_workflow_config 

125 

126 

127def add_workflow_init_nodes(config, qgraph, generic_workflow): 

128 """Add nodes to workflow graph that perform initialization steps. 

129 

130 Assumes that all of the initialization should be executed prior to any 

131 of the current workflow. 

132 

133 Parameters 

134 ---------- 

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

136 BPS configuration. 

137 qgraph : `lsst.pipe.base.graph.QuantumGraph` 

138 The quantum graph the generic workflow represents. 

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

140 Generic workflow to which the initialization steps should be added. 

141 """ 

142 # Create a workflow graph that will have task and file nodes necessary for 

143 # initializing the pipeline execution 

144 init_workflow = create_init_workflow(config, qgraph, generic_workflow.get_file("runQgraphFile")) 

145 _LOG.debug("init_workflow nodes = %s", init_workflow.nodes()) 

146 generic_workflow.add_workflow_source(init_workflow) 

147 

148 

149def create_init_workflow( 

150 config: BpsConfig, qgraph: QuantumGraph, qgraph_gwfile: GenericWorkflowFile 

151) -> GenericWorkflow: 

152 """Create workflow for running initialization job(s). 

153 

154 Parameters 

155 ---------- 

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

157 BPS configuration. 

158 qgraph : `lsst.pipe.base.graph.QuantumGraph` 

159 The quantum graph the generic workflow represents. 

160 qgraph_gwfile : `lsst.ctrl.bps.GenericWorkflowFile` 

161 File object for the full run QuantumGraph file. 

162 

163 Returns 

164 ------- 

165 init_workflow : `lsst.ctrl.bps.GenericWorkflow` 

166 GenericWorkflow consisting of job(s) to initialize workflow. 

167 """ 

168 _LOG.debug("creating init subgraph") 

169 _LOG.debug("creating init task input(s)") 

170 search_opt = { 

171 "curvals": {"curr_pipetask": "pipetaskInit"}, 

172 "replaceVars": False, 

173 "expandEnvVars": False, 

174 "replaceEnvVars": True, 

175 "required": False, 

176 } 

177 found, value = config.search("computeSite", opt=search_opt) 

178 if found: 178 ↛ 180line 178 didn't jump to line 180 because the condition on line 178 was always true

179 search_opt["curvals"]["curr_site"] = value 

180 found, value = config.search("computeCloud", opt=search_opt) 

181 if found: 

182 search_opt["curvals"]["curr_cloud"] = value 

183 

184 init_workflow = GenericWorkflow("init") 

185 init_workflow.add_file(qgraph_gwfile) 

186 

187 # create job for executing --init-only 

188 gwjob = GenericWorkflowJob("pipetaskInit", "pipetaskInit") 

189 

190 job_values = _get_job_values(config, search_opt, "runQuantumCommand") 

191 job_values["name"] = "pipetaskInit" 

192 job_values["label"] = "pipetaskInit" 

193 

194 # Adjust job attributes values if necessary. 

195 _handle_job_values(job_values, gwjob) 

196 

197 init_workflow.add_job(gwjob) 

198 init_workflow.add_job_inputs(gwjob.name, [qgraph_gwfile]) 

199 _enhance_command(config, init_workflow, gwjob, {}) 

200 

201 return init_workflow 

202 

203 

204def _enhance_command(config, generic_workflow, gwjob, cached_job_values): 

205 """Enhance command line with env and file placeholders 

206 and gather command line values. 

207 

208 Parameters 

209 ---------- 

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

211 BPS configuration. 

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

213 Generic workflow that contains the job. 

214 gwjob : `lsst.ctrl.bps.GenericWorkflowJob` 

215 Generic workflow job to which the updated executable, arguments, 

216 and values should be saved. 

217 cached_job_values : `dict` [`str`, dict[`str`, `~typing.Any`]] 

218 Cached values common across jobs with same label. Updated if values 

219 aren't already saved for given gwjob's label. 

220 """ 

221 _LOG.debug("gwjob given to _enhance_command: %s", gwjob) 

222 

223 search_opt = config.get_search_opts(gwjob.label) 

224 

225 search_opt["curvals"]["jobName"] = gwjob.name 

226 search_opt["curvals"]["jobLabel"] = gwjob.label 

227 for key, value in gwjob.tags.items(): 

228 search_opt["curvals"][key] = value 

229 

230 search_opt.update( 

231 { 

232 "replaceVars": False, 

233 "expandEnvVars": False, 

234 "replaceEnvVars": True, 

235 "required": False, 

236 } 

237 ) 

238 

239 if gwjob.label not in cached_job_values: 

240 cached_job_values[gwjob.label] = {} 

241 # Allowing whenSaveJobQgraph and useLazyCommands per pipetask label. 

242 key = "whenSaveJobQgraph" 

243 _, when_save = config.search(key, opt=search_opt) 

244 cached_job_values[gwjob.label][key] = WhenToSaveQuantumGraphs[when_save.upper()] 

245 

246 key = "useLazyCommands" 

247 search_opt["default"] = True 

248 _, cached_job_values[gwjob.label][key] = config.search(key, opt=search_opt) 

249 del search_opt["default"] 

250 

251 # Change qgraph variable to match whether using run or per-job qgraph 

252 # Note: these are lookup keys, not actual physical filenames. 

253 if cached_job_values[gwjob.label]["whenSaveJobQgraph"] == WhenToSaveQuantumGraphs.NEVER: 

254 gwjob.arguments = gwjob.arguments.replace("{qgraphFile}", "{runQgraphFile}") 

255 elif gwjob.name == "pipetaskInit": 255 ↛ 256line 255 didn't jump to line 256 because the condition on line 255 was never true

256 gwjob.arguments = gwjob.arguments.replace("{qgraphFile}", "{runQgraphFile}") 

257 else: # Needed unique file keys for per-job QuantumGraphs 

258 gwjob.arguments = gwjob.arguments.replace("{qgraphFile}", f"{{qgraphFile_{gwjob.name}}}") 

259 

260 # Replace files with special placeholders 

261 for gwfile in generic_workflow.get_job_inputs(gwjob.name): 

262 gwjob.arguments = gwjob.arguments.replace(f"{{{gwfile.name}}}", f"<FILE:{gwfile.name}>") 

263 for gwfile in generic_workflow.get_job_outputs(gwjob.name): 

264 gwjob.arguments = gwjob.arguments.replace(f"{{{gwfile.name}}}", f"<FILE:{gwfile.name}>") 

265 

266 # Replace wms variables with wms placeholders. 

267 gwjob.arguments = re.sub( 

268 r"{wms([^}]+)}", lambda x: f"<WMS:{x[1][0].lower() + x[1][1:]}>", gwjob.arguments 

269 ) 

270 

271 # Save dict of other values needed to complete command line. 

272 # (Be careful to not replace env variables as they may 

273 # be different in compute job.) 

274 search_opt["replaceVars"] = True 

275 _LOG.debug("before cmdvals = %s (search_opt = %s)", gwjob.cmdvals, search_opt) 

276 for key in re.findall(r"{([^}]+)}", gwjob.arguments): 

277 _LOG.debug("looking for %s in cmdvals", key) 

278 if key in gwjob.cmdvals: 

279 continue 

280 elif key in cached_job_values[gwjob.label]: 

281 gwjob.cmdvals[key] = cached_job_values[gwjob.label][key] 

282 else: 

283 _, gwjob.cmdvals[key] = config.search(key, opt=search_opt) 

284 _LOG.debug("after cmdvals = %s", gwjob.cmdvals) 

285 

286 # backwards compatibility 

287 if not cached_job_values[gwjob.label]["useLazyCommands"]: 287 ↛ 288line 287 didn't jump to line 288 because the condition on line 287 was never true

288 if "bpsUseShared" not in cached_job_values[gwjob.label]: 

289 key = "bpsUseShared" 

290 search_opt["default"] = True 

291 _, cached_job_values[gwjob.label][key] = config.search(key, opt=search_opt) 

292 del search_opt["default"] 

293 

294 gwjob.arguments = _fill_arguments( 

295 cached_job_values[gwjob.label]["bpsUseShared"], generic_workflow, gwjob.arguments, gwjob.cmdvals 

296 ) 

297 

298 

299def _fill_arguments(use_shared, generic_workflow, arguments, cmdvals): 

300 """Replace placeholders in command line string in job. 

301 

302 Parameters 

303 ---------- 

304 use_shared : `bool` 

305 Whether using shared filesystem. 

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

307 Generic workflow containing the job. 

308 arguments : `str` 

309 String containing placeholders. 

310 cmdvals : `dict` [`str`, `~typing.Any`] 

311 Any command line values that can be used to replace placeholders. 

312 

313 Returns 

314 ------- 

315 arguments : `str` 

316 Command line with FILE and ENV placeholders replaced. 

317 """ 

318 # Replace file placeholders 

319 for file_key in re.findall(r"<FILE:([^>]+)>", arguments): 

320 gwfile = generic_workflow.get_file(file_key) 

321 if not gwfile.wms_transfer: 

322 # Must assume full URI if in command line and told WMS is not 

323 # responsible for transferring file. 

324 uri = gwfile.src_uri 

325 elif use_shared: 

326 if gwfile.job_shared: 

327 # Have shared filesystems and jobs can share file. 

328 uri = gwfile.src_uri 

329 else: 

330 uri = os.path.basename(gwfile.src_uri) 

331 else: # Using push transfer 

332 uri = os.path.basename(gwfile.src_uri) 

333 

334 arguments = arguments.replace(f"<FILE:{file_key}>", uri) 

335 

336 # Replace env placeholder with submit-side values 

337 arguments = re.sub(r"<ENV:([^>]+)>", r"$\1", arguments) 

338 arguments = os.path.expandvars(arguments) 

339 

340 # Replace remaining vars 

341 arguments = arguments.format(**cmdvals) 

342 

343 return arguments 

344 

345 

346def _get_qgraph_gwfile(config, save_qgraph_per_job, gwjob, run_qgraph_file, prefix): 

347 """Get qgraph location to be used by job. 

348 

349 Parameters 

350 ---------- 

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

352 Bps configuration. 

353 save_qgraph_per_job : `lsst.ctrl.bps.bps_utils.WhenToSaveQuantumGraphs` 

354 What submission stage to save per-job qgraph files (or NEVER) 

355 gwjob : `lsst.ctrl.bps.GenericWorkflowJob` 

356 Job for which determining QuantumGraph file. 

357 run_qgraph_file : `lsst.ctrl.bps.GenericWorkflowFile` 

358 File representation of the full run QuantumGraph. 

359 prefix : `str` 

360 Path prefix for any files written. 

361 

362 Returns 

363 ------- 

364 gwfile : `lsst.ctrl.bps.GenericWorkflowFile` 

365 Representation of butler location (may not include filename). 

366 """ 

367 qgraph_gwfile = None 

368 if save_qgraph_per_job != WhenToSaveQuantumGraphs.NEVER: 368 ↛ 369line 368 didn't jump to line 369 because the condition on line 368 was never true

369 qgraph_gwfile = GenericWorkflowFile( 

370 f"qgraphFile_{gwjob.name}", 

371 src_uri=create_job_quantum_graph_filename(config, gwjob, prefix), 

372 wms_transfer=True, 

373 job_access_remote=True, 

374 job_shared=True, 

375 ) 

376 else: 

377 qgraph_gwfile = run_qgraph_file 

378 

379 return qgraph_gwfile 

380 

381 

382def _get_job_values(config, search_opt, cmd_line_key): 

383 """Gather generic workflow job values from the bps config. 

384 

385 Parameters 

386 ---------- 

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

388 Bps configuration. 

389 search_opt : `dict` [`str`, `~typing.Any`] 

390 Search options to be used when searching config. 

391 cmd_line_key : `str` or None 

392 Which command line key to search for (e.g., "runQuantumCommand"). 

393 

394 Returns 

395 ------- 

396 job_values : `dict` [ `str`, `~typing.Any` ]` 

397 A mapping between job attributes and their values. 

398 """ 

399 _LOG.debug("cmd_line_key=%s, search_opt=%s", cmd_line_key, search_opt) 

400 

401 # Create a dummy job to easily access the default values. 

402 default_gwjob = GenericWorkflowJob("default_job", "default_label") 

403 

404 job_values = {} 

405 for attr in _ATTRS_ALL: 

406 # Variable names in yaml are camel case instead of snake case. 

407 yaml_name = re.sub(r"_(\S)", lambda match: match.group(1).upper(), attr) 

408 found, value = config.search(yaml_name, opt=search_opt) 

409 if found: 

410 job_values[attr] = value 

411 else: 

412 job_values[attr] = getattr(default_gwjob, attr) 

413 

414 # Need to replace all config variables in environment values. 

415 # Also change env vars in environment values to bash syntax. 

416 # 

417 # Note: Because job_values["environment"] is a BpsConfig and 

418 # currently cannot have 2 search objects, for each environment 

419 # setting, we have to get the setting string as is and then 

420 # separately use the overall config to replace values inside 

421 # the setting string. 

422 tmp_job_env = job_values.get("environment", None) 

423 if tmp_job_env: 

424 _LOG.debug("_get_job_values: job_values['environment'] = %s", tmp_job_env) 

425 

426 # Don't want to replace when getting environment setting string. 

427 as_is_search_opt = { 

428 "replaceVars": False, 

429 "expandEnvVars": False, 

430 "replaceEnvBps2Shell": False, 

431 "replaceEnvShell2Bps": False, 

432 } 

433 

434 # When updating environment string, use given search options, 

435 # but ensure making the environment string using bash syntax. 

436 env_search_opt = copy.copy(search_opt) 

437 env_search_opt["replaceVars"] = True # Replace bps config variables. 

438 env_search_opt["replaceEnvBps2Shell"] = False # Replace bps <ENV:var> syntax. 

439 env_search_opt["replaceEnvShell2Bps"] = True # Do not replace shell env syntax. 

440 env_search_opt["expandEnvVars"] = False # Do not replace with submission env value. 

441 

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

443 

444 for name in tmp_job_env: 

445 # Get environment setting string as is. 

446 value = tmp_job_env.search(name, as_is_search_opt)[1] 

447 _LOG.debug("_get_job_values: as is value for %s = %s", name, value) 

448 # Replace config vars and env placeholders 

449 job_env[name] = config.modify_value(name, str(value), env_search_opt) 

450 _LOG.debug("_get_job_values: new env value for %s = %s", name, job_env[name]) 

451 # Save new dictionary back with other job values. 

452 job_values["environment"] = job_env 

453 

454 # If the automatic memory scaling is enabled (i.e. the memory multiplier 

455 # is set and it is a positive number greater than 1.0), adjust number 

456 # of retries when necessary. If the memory multiplier is invalid, disable 

457 # automatic memory scaling. 

458 if job_values["memory_multiplier"] is not None: 

459 if math.ceil(float(job_values["memory_multiplier"])) > 1: 

460 if job_values["number_of_retries"] is None: 460 ↛ 465line 460 didn't jump to line 465 because the condition on line 460 was always true

461 job_values["number_of_retries"] = DEFAULT_MEM_RETRIES 

462 else: 

463 job_values["memory_multiplier"] = None 

464 

465 if cmd_line_key: 

466 found, cmdline = config.search(cmd_line_key, opt=search_opt) 

467 # Make sure cmdline isn't None as that could be sent in as a 

468 # default value in search_opt. 

469 if found and cmdline: 

470 cmd, args = cmdline.split(" ", 1) 

471 job_values["executable"] = GenericWorkflowExec(os.path.basename(cmd), cmd, False) 

472 if args: 472 ↛ 475line 472 didn't jump to line 475 because the condition on line 472 was always true

473 job_values["arguments"] = args 

474 

475 return job_values 

476 

477 

478def _handle_job_values(quantum_job_values, gwjob, attributes=_ATTRS_ALL): 

479 """Set the job attributes in the cluster to their correct values. 

480 

481 Parameters 

482 ---------- 

483 quantum_job_values : `dict` [`str`, Any] 

484 Job values for running single Quantum. 

485 gwjob : `lsst.ctrl.bps.GenericWorkflowJob` 

486 Generic workflow job in which to store the universal values. 

487 attributes : `~collections.abc.Iterable` [`str`], optional 

488 Job attributes to be set in the job following different rules. 

489 The default value is _ATTRS_ALL. 

490 """ 

491 _LOG.debug("Call to _handle_job_values") 

492 _handle_job_values_universal(quantum_job_values, gwjob, attributes) 

493 _handle_job_values_max(quantum_job_values, gwjob, attributes) 

494 _handle_job_values_sum(quantum_job_values, gwjob, attributes) 

495 

496 

497def _handle_job_values_universal(quantum_job_values, gwjob, attributes=_ATTRS_UNIVERSAL): 

498 """Handle job attributes that must have the same value for every quantum 

499 in the cluster. 

500 

501 Parameters 

502 ---------- 

503 quantum_job_values : `dict` [`str`, Any] 

504 Job values for running single Quantum. 

505 gwjob : `lsst.ctrl.bps.GenericWorkflowJob` 

506 Generic workflow job in which to store the universal values. 

507 attributes : `~collections.abc.Iterable` [`str`], optional 

508 Job attributes to be set in the job following different rules. 

509 The default value is _ATTRS_UNIVERSAL. 

510 """ 

511 for attr in _ATTRS_UNIVERSAL & set(attributes): 

512 _LOG.debug( 

513 "Handling job %s (job=%s, quantum=%s)", 

514 attr, 

515 getattr(gwjob, attr), 

516 quantum_job_values.get(attr, "MISSING"), 

517 ) 

518 current_value = getattr(gwjob, attr) 

519 try: 

520 quantum_value = quantum_job_values[attr] 

521 except KeyError: 

522 continue 

523 else: 

524 if not current_value: 

525 setattr(gwjob, attr, quantum_value) 

526 elif current_value != quantum_value: 526 ↛ 527line 526 didn't jump to line 527 because the condition on line 526 was never true

527 _LOG.error( 

528 "Inconsistent value for %s in Cluster %s Quantum Number %s\n" 

529 "Current cluster value: %s\n" 

530 "Quantum value: %s", 

531 attr, 

532 gwjob.name, 

533 quantum_job_values.get("qgraphNodeId", "MISSING"), 

534 current_value, 

535 quantum_value, 

536 ) 

537 raise RuntimeError(f"Inconsistent value for {attr} in cluster {gwjob.name}.") 

538 

539 

540def _handle_job_values_max(quantum_job_values, gwjob, attributes=_ATTRS_MAX): 

541 """Handle job attributes that should be set to their maximum value in 

542 the in cluster. 

543 

544 Parameters 

545 ---------- 

546 quantum_job_values : `dict` [`str`, `~typing.Any`] 

547 Job values for running single Quantum. 

548 gwjob : `lsst.ctrl.bps.GenericWorkflowJob` 

549 Generic workflow job in which to store the aggregate values. 

550 attributes : `~collections.abc.Iterable` [`str`], optional 

551 Job attributes to be set in the job following different rules. 

552 The default value is _ATTR_MAX. 

553 """ 

554 for attr in _ATTRS_MAX & set(attributes): 

555 current_value = getattr(gwjob, attr) 

556 try: 

557 quantum_value = quantum_job_values[attr] 

558 except KeyError: 

559 continue 

560 else: 

561 needs_update = False 

562 if current_value is None: 562 ↛ 566line 562 didn't jump to line 566 because the condition on line 562 was always true

563 if quantum_value is not None: 563 ↛ 564line 563 didn't jump to line 564 because the condition on line 563 was never true

564 needs_update = True 

565 else: 

566 if quantum_value is not None and current_value < quantum_value: 

567 needs_update = True 

568 if needs_update: 568 ↛ 569line 568 didn't jump to line 569 because the condition on line 568 was never true

569 setattr(gwjob, attr, quantum_value) 

570 

571 # When updating memory requirements for a job, check if memory 

572 # autoscaling is enabled. If it is, always use the memory 

573 # multiplier and the number of retries which comes with the 

574 # quantum. 

575 # 

576 # Note that as a result, the quantum with the biggest memory 

577 # requirements will determine whether the memory autoscaling 

578 # will be enabled (or disabled) depending on the value of its 

579 # memory multiplier. 

580 if attr == "request_memory": 

581 gwjob.memory_multiplier = quantum_job_values["memory_multiplier"] 

582 if gwjob.memory_multiplier is not None: 

583 gwjob.number_of_retries = quantum_job_values["number_of_retries"] 

584 

585 

586def _handle_job_values_sum(quantum_job_values, gwjob, attributes=_ATTRS_SUM): 

587 """Handle job attributes that are the sum of their values in the cluster. 

588 

589 Parameters 

590 ---------- 

591 quantum_job_values : `dict` [`str`, `~typing.Any`] 

592 Job values for running single Quantum. 

593 gwjob : `lsst.ctrl.bps.GenericWorkflowJob` 

594 Generic workflow job in which to store the aggregate values. 

595 attributes : `~collections.abc.Iterable` [`str`], optional 

596 Job attributes to be set in the job following different rules. 

597 The default value is _ATTRS_SUM. 

598 """ 

599 for attr in _ATTRS_SUM & set(attributes): 

600 current_value = getattr(gwjob, attr) 

601 if not current_value: 601 ↛ 604line 601 didn't jump to line 604 because the condition on line 601 was always true

602 setattr(gwjob, attr, quantum_job_values[attr]) 

603 else: 

604 setattr(gwjob, attr, current_value + quantum_job_values[attr]) 

605 

606 

607def create_generic_workflow( 

608 config: BpsConfig, cqgraph: ClusteredQuantumGraph, name: str, prefix: str 

609) -> GenericWorkflow: 

610 """Create a generic workflow from a ClusteredQuantumGraph such that it 

611 has information needed for WMS (e.g., command lines). 

612 

613 Parameters 

614 ---------- 

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

616 BPS configuration. 

617 cqgraph : `lsst.ctrl.bps.ClusteredQuantumGraph` 

618 ClusteredQuantumGraph for running a specific pipeline on a specific 

619 payload. 

620 name : `str` 

621 Name for the workflow (typically unique). 

622 prefix : `str` 

623 Root path for any output files. 

624 

625 Returns 

626 ------- 

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

628 Generic workflow for the given ClusteredQuantumGraph + config. 

629 """ 

630 # Determine whether saving per-job QuantumGraph files in the loop. 

631 _, when_save = config.search("whenSaveJobQgraph", {"default": WhenToSaveQuantumGraphs.TRANSFORM.name}) 

632 save_qgraph_per_job = WhenToSaveQuantumGraphs[when_save.upper()] 

633 

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

635 

636 generic_workflow = GenericWorkflow(name) 

637 

638 # Save full run QuantumGraph for use by jobs 

639 generic_workflow.add_file( 

640 GenericWorkflowFile( 

641 "runQgraphFile", 

642 src_uri=config["runQgraphFile"], 

643 wms_transfer=True, 

644 job_access_remote=True, 

645 job_shared=True, 

646 ) 

647 ) 

648 

649 # Cache pipetask specific or more generic job values to minimize number 

650 # on config searches. 

651 cached_job_values = {} 

652 cached_pipetask_values = {} 

653 

654 for cluster in cqgraph.clusters(): 

655 _LOG.debug("Loop over clusters: %s, %s", cluster, type(cluster)) 

656 _LOG.debug( 

657 "cqgraph: name=%s, len=%s, label=%s, ids=%s", 

658 cluster.name, 

659 len(cluster.qgraph_node_ids), 

660 cluster.label, 

661 cluster.qgraph_node_ids, 

662 ) 

663 

664 gwjob = GenericWorkflowJob(cluster.name, cluster.label) 

665 

666 # First get job values from cluster or cluster config 

667 search_opt["curvals"] = {"curr_cluster": cluster.label} 

668 found, value = config.search("computeSite", opt=search_opt) 

669 if found: 669 ↛ 671line 669 didn't jump to line 671 because the condition on line 669 was always true

670 search_opt["curvals"]["curr_site"] = value 

671 found, value = config.search("computeCloud", opt=search_opt) 

672 if found: 

673 search_opt["curvals"]["curr_cloud"] = value 

674 

675 # If some config values are set for this cluster 

676 if cluster.label not in cached_job_values: 

677 _LOG.debug("config['cluster'][%s] = %s", cluster.label, config["cluster"][cluster.label]) 

678 cached_job_values[cluster.label] = {} 

679 

680 # Allowing whenSaveJobQgraph and useLazyCommands per cluster label. 

681 key = "whenSaveJobQgraph" 

682 _, when_save = config.search(key, opt=search_opt) 

683 cached_job_values[cluster.label][key] = WhenToSaveQuantumGraphs[when_save.upper()] 

684 

685 key = "useLazyCommands" 

686 search_opt["default"] = True 

687 _, cached_job_values[cluster.label][key] = config.search(key, opt=search_opt) 

688 del search_opt["default"] 

689 

690 if cluster.label in config["cluster"]: 690 ↛ 693line 690 didn't jump to line 693 because the condition on line 690 was never true

691 # Don't want to get global defaults here so only look in 

692 # cluster section. 

693 cached_job_values[cluster.label].update( 

694 _get_job_values(config["cluster"][cluster.label], search_opt, "runQuantumCommand") 

695 ) 

696 cluster_job_values = copy.copy(cached_job_values[cluster.label]) 

697 

698 cluster_job_values["name"] = cluster.name 

699 cluster_job_values["label"] = cluster.label 

700 cluster_job_values["quanta_counts"] = cluster.quanta_counts 

701 cluster_job_values["tags"] = cluster.tags 

702 _LOG.debug("cluster_job_values = %s", cluster_job_values) 

703 _handle_job_values(cluster_job_values, gwjob, cluster_job_values.keys()) 

704 

705 # For purposes of whether to continue searching for a value is whether 

706 # the value evaluates to False. 

707 unset_attributes = {attr for attr in _ATTRS_ALL if not getattr(gwjob, attr)} 

708 

709 _LOG.debug("unset_attributes=%s", unset_attributes) 

710 _LOG.debug("set=%s", _ATTRS_ALL - unset_attributes) 

711 

712 # For job info not defined at cluster level, attempt to get job info 

713 # either common or aggregate for all Quanta in cluster. 

714 for node_id in iter(cluster.qgraph_node_ids): 

715 _LOG.debug("node_id=%s", node_id) 

716 quantum_info = cqgraph.get_quantum_info(node_id) 

717 

718 task_label = quantum_info["task_label"] 

719 if task_label not in cached_pipetask_values: 

720 search_opt["curvals"]["curr_pipetask"] = task_label 

721 cached_pipetask_values[task_label] = _get_job_values(config, search_opt, "runQuantumCommand") 

722 _handle_job_values(cached_pipetask_values[task_label], gwjob, unset_attributes) 

723 

724 # Update job with workflow attribute and profile values. 

725 qgraph_gwfile = _get_qgraph_gwfile( 

726 config, save_qgraph_per_job, gwjob, generic_workflow.get_file("runQgraphFile"), prefix 

727 ) 

728 

729 generic_workflow.add_job(gwjob) 

730 generic_workflow.add_job_inputs(gwjob.name, [qgraph_gwfile]) 

731 

732 gwjob.cmdvals["qgraphNodeId"] = ",".join( 

733 sorted([f"{node_id}" for node_id in cluster.qgraph_node_ids]) 

734 ) 

735 _enhance_command(config, generic_workflow, gwjob, cached_job_values) 

736 

737 # If writing per-job QuantumGraph files during TRANSFORM stage, 

738 # write it now while in memory. 

739 if save_qgraph_per_job == WhenToSaveQuantumGraphs.TRANSFORM: 739 ↛ 740line 739 didn't jump to line 740 because the condition on line 739 was never true

740 save_qg_subgraph(cqgraph.qgraph, qgraph_gwfile.src_uri, cluster.qgraph_node_ids) 

741 

742 # Create job dependencies. 

743 for parent in cqgraph.clusters(): 

744 for child in cqgraph.successors(parent): 

745 generic_workflow.add_job_relationships(parent.name, child.name) 

746 

747 # Add initial workflow. 

748 if config.get("runInit", "{default: False}"): 748 ↛ 751line 748 didn't jump to line 751 because the condition on line 748 was always true

749 add_workflow_init_nodes(config, cqgraph.qgraph, generic_workflow) 

750 

751 generic_workflow.run_attrs.update( 

752 { 

753 "bps_isjob": "True", 

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

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

756 "bps_run": generic_workflow.name, 

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

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

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

760 } 

761 ) 

762 

763 # Add final job 

764 add_final_job(config, generic_workflow, prefix) 

765 

766 if "ordering" in config: 766 ↛ 767line 766 didn't jump to line 767 because the condition on line 766 was never true

767 generic_workflow.add_special_job_ordering(config["ordering"]) 

768 

769 return generic_workflow 

770 

771 

772def create_generic_workflow_config(config, prefix): 

773 """Create generic workflow configuration. 

774 

775 Parameters 

776 ---------- 

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

778 Bps configuration. 

779 prefix : `str` 

780 Root path for any output files. 

781 

782 Returns 

783 ------- 

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

785 Configuration accompanying the GenericWorkflow. 

786 """ 

787 generic_workflow_config = BpsConfig(config) 

788 generic_workflow_config["workflowName"] = config["uniqProcName"] 

789 generic_workflow_config["workflowPath"] = prefix 

790 return generic_workflow_config 

791 

792 

793def add_final_job(config: BpsConfig, generic_workflow: GenericWorkflow, prefix: str) -> None: 

794 """Add final workflow job depending upon configuration. 

795 

796 Depending on configuration, the final job will be added as a special job 

797 which will always run regardless of the exit status of the workflow or 

798 a regular sink node which will only run if the workflow execution finished 

799 with no errors. 

800 

801 Parameters 

802 ---------- 

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

804 Bps configuration. 

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

806 Generic workflow to which attributes should be added. 

807 prefix : `str` 

808 Directory in which to output final script. 

809 """ 

810 _, when_run = config.search(".finalJob.whenRun") 

811 if when_run.upper() != "NEVER": 811 ↛ exitline 811 didn't return from function 'add_final_job' because the condition on line 811 was always true

812 gwjob = create_final_job(config, generic_workflow, prefix) 

813 if when_run.upper() == "ALWAYS": 813 ↛ 815line 813 didn't jump to line 815 because the condition on line 813 was always true

814 generic_workflow.add_final(gwjob) 

815 elif when_run.upper() == "SUCCESS": 

816 add_final_job_as_sink(generic_workflow, gwjob) 

817 else: 

818 raise ValueError(f"Invalid value for finalJob.whenRun: {when_run}") 

819 

820 

821def create_final_job(config: BpsConfig, generic_workflow: GenericWorkflow, prefix: str) -> GenericWorkflowJob: 

822 """Create the final workflow job. 

823 

824 Parameters 

825 ---------- 

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

827 Bps configuration. 

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

829 Generic workflow to which attributes should be added. 

830 prefix : `str` 

831 Directory in which to output final script. 

832 

833 Returns 

834 ------- 

835 final_job : `lsst.ctrl.bps.GenericWorkflowJob` 

836 Final workflow job. 

837 """ 

838 job_name = "finalJob" 

839 gwjob = GenericWorkflowJob(job_name, job_name) 

840 

841 search_opt = {"searchobj": config[job_name], "curvals": {}, "default": None} 

842 found, value = config.search("computeSite", opt=search_opt) 

843 if found: 843 ↛ 845line 843 didn't jump to line 845 because the condition on line 843 was always true

844 search_opt["curvals"]["curr_site"] = value 

845 found, value = config.search("computeCloud", opt=search_opt) 

846 if found: 846 ↛ 855line 846 didn't jump to line 855 because the condition on line 846 was always true

847 search_opt["curvals"]["curr_cloud"] = value 

848 

849 # Set job attributes based on the values find in the config excluding 

850 # the ones in the _ATTRS_MISC group. The attributes in this group are 

851 # somewhat "special": 

852 # * HTCondor plugin, which uses 'attrs' and 'profile', has its own 

853 # mechanism for setting them, 

854 # * 'cmdvals' is being set internally, not via config. 

855 job_values = _get_job_values(config, search_opt, None) 

856 for attr in _ATTRS_ALL - _ATTRS_MISC: 

857 if not getattr(gwjob, attr) and job_values.get(attr, None): 

858 setattr(gwjob, attr, job_values[attr]) 

859 

860 # Create script and add command line to job. 

861 gwjob.executable, gwjob.arguments = create_final_command(config, prefix) 

862 

863 # Determine inputs from command line. 

864 for file_key in re.findall(r"<FILE:([^>]+)>", gwjob.arguments): 

865 gwfile = generic_workflow.get_file(file_key) 

866 generic_workflow.add_job_inputs(gwjob.name, gwfile) 

867 

868 _enhance_command(config, generic_workflow, gwjob, {}) 

869 return gwjob 

870 

871 

872def create_final_command(config: BpsConfig, prefix: str) -> tuple[GenericWorkflowExec, str]: 

873 """Create the command and shell script for the final job. 

874 

875 Parameters 

876 ---------- 

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

878 Bps configuration. 

879 prefix : `str` 

880 Directory in which to output final script. 

881 

882 Returns 

883 ------- 

884 executable : `lsst.ctrl.bps.GenericWorkflowExec` 

885 Executable object for the final script. 

886 arguments : `str` 

887 Command line needed to call the final script. 

888 

889 Raises 

890 ------ 

891 RuntimeError if no commands found. 

892 """ 

893 search_opt = { 

894 "replaceVars": True, 

895 "skipNames": ["butlerConfig", "qgraphFile"], 

896 "replaceEnvVars": False, 

897 "expandEnvVars": False, 

898 "searchobj": config["finalJob"], 

899 } 

900 

901 script_file = os.path.join(prefix, "final_job.bash") 

902 with open(script_file, "w", encoding="utf8") as fh: 

903 print("#!/bin/bash\n", file=fh) 

904 print("set -e", file=fh) 

905 print("set -x", file=fh) 

906 

907 print("qgraphFile=$1", file=fh) 

908 print("butlerConfig=$2", file=fh) 

909 

910 command_len = 0 # Make sure at least write one actual command 

911 i = 1 

912 found, command = config.search(f"command{i}", opt=search_opt) 

913 while found: 

914 # The files will be args to script, so change to shell vars 

915 command = command.replace("{qgraphFile}", "${qgraphFile}") 

916 command = command.replace("{butlerConfig}", "${butlerConfig}") 

917 

918 print(command, file=fh) 

919 command_len += len(command.strip()) 

920 

921 # Search for next command 

922 i += 1 

923 found, command = config.search(f"command{i}", opt=search_opt) 

924 if command_len == 0: 

925 raise RuntimeError( 

926 "No finalJob commands were found. Use NEVER for finalJob.whenRun to turn off finalJob" 

927 ) 

928 os.chmod(script_file, 0o755) 

929 executable = GenericWorkflowExec(os.path.basename(script_file), script_file, True) 

930 

931 _, orig_butler = config.search("butlerConfig") 

932 return executable, f"<FILE:runQgraphFile> {orig_butler}" 

933 

934 

935def add_final_job_as_sink(generic_workflow, final_job): 

936 """Add final job as the single sink for the workflow. 

937 

938 Parameters 

939 ---------- 

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

941 Generic workflow to which attributes should be added. 

942 final_job : `lsst.ctrl.bps.GenericWorkflowJob` 

943 Job to add as new sink node depending upon all previous sink nodes. 

944 """ 

945 # Find sink nodes of generic workflow graph. 

946 gw_sinks = [n for n in generic_workflow if generic_workflow.out_degree(n) == 0] 

947 _LOG.debug("gw_sinks = %s", gw_sinks) 

948 

949 generic_workflow.add_job(final_job) 

950 generic_workflow.add_job_relationships(gw_sinks, final_job.name)