Coverage for python/lsst/ctrl/bps/htcondor/prepare_utils.py: 82%

506 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-28 02:24 -0700

1# This file is part of ctrl_bps_htcondor. 

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"""Utility functions for preparing the HTCondor workflow.""" 

29 

30import logging 

31import os 

32import re 

33from collections import defaultdict 

34from pathlib import Path 

35from typing import Any, cast 

36 

37from lsst.ctrl.bps import ( 

38 BpsConfig, 

39 GenericWorkflow, 

40 GenericWorkflowFile, 

41 GenericWorkflowGroup, 

42 GenericWorkflowJob, 

43 GenericWorkflowNodeType, 

44 GenericWorkflowNoopJob, 

45) 

46from lsst.ctrl.bps.bps_utils import create_count_summary 

47 

48from .lssthtc import ( 

49 HTCDag, 

50 HTCJob, 

51 condor_status, 

52 htc_escape, 

53 read_dag_info, 

54 write_dag_info, 

55) 

56 

57_LOG = logging.getLogger(__name__) 

58 

59DEFAULT_HTC_EXEC_PATT = ".*worker.*" 

60"""Default pattern for searching execute machines in an HTCondor pool. 

61""" 

62 

63 

64def _create_job(subdir_template, cached_values, generic_workflow, gwjob, out_prefix): 

65 """Convert GenericWorkflow job nodes to DAG jobs. 

66 

67 Parameters 

68 ---------- 

69 subdir_template : `str` 

70 Template for making subdirs. 

71 cached_values : `dict` 

72 Site and label specific values. 

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

74 Generic workflow that is being converted. 

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

76 The generic job to convert to a HTCondor job. 

77 out_prefix : `str` 

78 Directory prefix for HTCondor files. 

79 

80 Returns 

81 ------- 

82 htc_job : `lsst.ctrl.bps.wms.htcondor.HTCJob` 

83 The HTCondor job equivalent to the given generic job. 

84 """ 

85 _LOG.debug("_create_job: cached_values = %s", cached_values) 

86 htc_job = HTCJob(gwjob.name, label=gwjob.label) 

87 

88 curvals = defaultdict(str) 

89 curvals["label"] = gwjob.label 

90 if gwjob.tags: 

91 curvals.update(gwjob.tags) 

92 

93 subdir = Path("jobs") / subdir_template.format_map(curvals) 

94 htc_job.subdir = subdir 

95 htc_job.subfile = f"{gwjob.name}.sub" 

96 htc_job.add_dag_cmds({"dir": subdir}) 

97 

98 htc_job_cmds = { 

99 "universe": "vanilla", 

100 "should_transfer_files": "YES", 

101 "when_to_transfer_output": "ON_EXIT_OR_EVICT", 

102 "transfer_output_files": '""', # Set to empty string to disable 

103 "transfer_executable": "False", 

104 # Exceeding memory sometimes triggers SIGBUS or SIGSEGV error. Tell 

105 # htcondor to put on hold any jobs which exited by a signal. If 

106 # executed in a bash script, like finalJob, the signals will become 

107 # exit codes above 128 (exit code = 128 + signal number). 

108 "on_exit_hold": "ExitBySignal == true || ExitCode > 128", 

109 "on_exit_hold_reason": "ExitBySignal == true ? " 

110 'strcat("Job raised a signal ", string(ExitSignal), ' 

111 '". Handling job as if it has gone over memory limit.") : ' 

112 'strcat("Job exit code (", string(ExitCode), ") > 128. ' 

113 'Handling job as if it has gone over memory limit.")', 

114 "on_exit_hold_subcode": "34", 

115 } 

116 

117 htc_job_cmds.update(_translate_job_cmds(cached_values, generic_workflow, gwjob)) 

118 

119 # Combine stdout and stderr to reduce the number of files. 

120 for key in ("output", "error"): 

121 if cached_values["overwriteJobFiles"]: 

122 htc_job_cmds[key] = f"{gwjob.name}.$(Cluster).out" 

123 else: 

124 htc_job_cmds[key] = f"{gwjob.name}.$(Cluster).$$([NumJobStarts ?: 0]).out" 

125 _LOG.debug("HTCondor %s = %s", key, htc_job_cmds[key]) 

126 

127 key = "log" 

128 htc_job_cmds[key] = f"{gwjob.name}.$(Cluster).{key}" 

129 _LOG.debug("HTCondor %s = %s", key, htc_job_cmds[key]) 

130 

131 htc_job_cmds.update( 

132 _handle_job_inputs(generic_workflow, gwjob.name, cached_values["bpsUseShared"], out_prefix) 

133 ) 

134 

135 htc_job_cmds.update( 

136 _handle_job_outputs(generic_workflow, gwjob.name, cached_values["bpsUseShared"], out_prefix) 

137 ) 

138 

139 # If specified, add nodeset to the job 

140 if "nodeset" in cached_values: 

141 htc_job.add_job_attrs({"JobNodeset": cached_values["nodeset"]}) 

142 clause = f'( Target.Nodeset == "{cached_values["nodeset"]}" )' 

143 if "requirements" in htc_job_cmds: 

144 htc_job_cmds["requirements"] = f"({htc_job_cmds['requirements']}) && {clause}" 

145 else: 

146 htc_job_cmds["requirements"] = clause 

147 

148 # Add the job cmds dict to the job object. 

149 htc_job.add_job_cmds(htc_job_cmds) 

150 

151 # Add job-related cmds to the DAG (e.g., VARS) 

152 htc_job.add_dag_cmds(_translate_dag_cmds(gwjob)) 

153 

154 # Add job attributes to job. 

155 _LOG.debug("gwjob.attrs = %s", gwjob.attrs) 

156 htc_job.add_job_attrs(gwjob.attrs) 

157 htc_job.add_job_attrs(cached_values["attrs"]) 

158 htc_job.add_job_attrs({"bps_job_quanta": create_count_summary(gwjob.quanta_counts)}) 

159 htc_job.add_job_attrs({"bps_job_name": gwjob.name, "bps_job_label": gwjob.label}) 

160 

161 return htc_job 

162 

163 

164def _translate_job_cmds(cached_vals, generic_workflow, gwjob): 

165 """Translate the job data that are one to one mapping 

166 

167 Parameters 

168 ---------- 

169 cached_vals : `dict` [`str`, `~typing.Any`] 

170 Config values common to jobs with same site or label. 

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

172 Generic workflow that contains job to being converted. 

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

174 Generic workflow job to be converted. 

175 

176 Returns 

177 ------- 

178 htc_job_commands : `dict` [`str`, `~typing.Any`] 

179 Contains commands which can appear in the HTCondor submit description 

180 file. 

181 """ 

182 # Values in the job script that just are name mappings. 

183 job_translation = { 

184 "mail_to": "notify_user", 

185 "when_to_mail": "notification", 

186 "request_cpus": "request_cpus", 

187 "priority": "priority", 

188 "category": "category", 

189 "accounting_group": "accounting_group", 

190 "accounting_user": "accounting_group_user", 

191 } 

192 

193 jobcmds = {} 

194 for gwkey, htckey in job_translation.items(): 

195 jobcmds[htckey] = getattr(gwjob, gwkey, None) 

196 

197 # If accounting info was not set explicitly, use site settings if any. 

198 if not gwjob.accounting_group: 198 ↛ 200line 198 didn't jump to line 200 because the condition on line 198 was always true

199 jobcmds["accounting_group"] = cached_vals.get("accountingGroup") 

200 if not gwjob.accounting_user: 200 ↛ 204line 200 didn't jump to line 204 because the condition on line 200 was always true

201 jobcmds["accounting_group_user"] = cached_vals.get("accountingUser") 

202 

203 # job commands that need modification 

204 if gwjob.retry_unless_exit: 

205 if isinstance(gwjob.retry_unless_exit, int): 

206 jobcmds["retry_until"] = f"{gwjob.retry_unless_exit}" 

207 elif isinstance(gwjob.retry_unless_exit, list): 

208 jobcmds["retry_until"] = ( 

209 f"member(ExitCode, {{{','.join([str(x) for x in gwjob.retry_unless_exit])}}})" 

210 ) 

211 else: 

212 raise ValueError("retryUnlessExit must be an integer or a list of integers.") 

213 

214 if gwjob.request_disk: 214 ↛ 215line 214 didn't jump to line 215 because the condition on line 214 was never true

215 jobcmds["request_disk"] = f"{gwjob.request_disk}MB" 

216 

217 if gwjob.request_memory: 

218 jobcmds["request_memory"] = f"{gwjob.request_memory}" 

219 

220 memory_max = 0 

221 if gwjob.memory_multiplier: 

222 # Do not use try-except! At the moment, BpsConfig returns an empty 

223 # string if it does not contain the key. 

224 memory_limit = cached_vals["memoryLimit"] 

225 if not memory_limit: 225 ↛ 226line 225 didn't jump to line 226 because the condition on line 225 was never true

226 raise RuntimeError( 

227 "Memory autoscaling enabled, but automatic detection of the memory limit " 

228 "failed; setting it explicitly with 'memoryLimit' or changing worker node " 

229 "search pattern 'executeMachinesPattern' might help." 

230 ) 

231 

232 # Set maximal amount of memory job can ask for. 

233 # 

234 # The check below assumes that 'memory_limit' was set to a value which 

235 # realistically reflects actual physical limitations of a given compute 

236 # resource. 

237 memory_max = memory_limit 

238 if gwjob.request_memory_max and gwjob.request_memory_max < memory_limit: 238 ↛ 239line 238 didn't jump to line 239 because the condition on line 238 was never true

239 memory_max = gwjob.request_memory_max 

240 

241 # Make job ask for more memory each time it failed due to insufficient 

242 # memory requirements. 

243 jobcmds["request_memory"] = _create_request_memory_expr( 

244 gwjob.request_memory, gwjob.memory_multiplier, memory_max 

245 ) 

246 

247 user_release_expr = cached_vals.get("releaseExpr", "") 

248 if gwjob.number_of_retries is not None and gwjob.number_of_retries >= 0: 

249 jobcmds["max_retries"] = gwjob.number_of_retries 

250 

251 # No point in adding periodic_release if 0 retries 

252 if gwjob.number_of_retries > 0: 

253 periodic_release = _create_periodic_release_expr( 

254 gwjob.request_memory, 

255 gwjob.memory_multiplier, 

256 memory_max, 

257 user_release_expr, 

258 ) 

259 if periodic_release: 259 ↛ 262line 259 didn't jump to line 262 because the condition on line 259 was always true

260 jobcmds["periodic_release"] = periodic_release 

261 

262 jobcmds["periodic_remove"] = _create_periodic_remove_expr( 

263 gwjob.request_memory, gwjob.memory_multiplier, memory_max 

264 ) 

265 

266 # Assume concurrency_limit implemented using HTCondor concurrency limits. 

267 # May need to move to special site-specific implementation if sites use 

268 # other mechanisms. 

269 if gwjob.concurrency_limit: 269 ↛ 270line 269 didn't jump to line 270 because the condition on line 269 was never true

270 jobcmds["concurrency_limit"] = gwjob.concurrency_limit 

271 

272 # Handle command line 

273 new_job_cmds = _translate_command_line(cached_vals, generic_workflow, gwjob) 

274 jobcmds.update(new_job_cmds) 

275 

276 # Add extra "pass-thru" job commands 

277 if gwjob.profile: 

278 for key, val in gwjob.profile.items(): 

279 jobcmds[key] = val 

280 for key, val in cached_vals["profile"].items(): 

281 jobcmds[key] = val 

282 

283 return jobcmds 

284 

285 

286def _translate_command_line( 

287 cached_vals: dict[str, Any], generic_workflow: GenericWorkflow, gwjob: GenericWorkflowJob 

288) -> dict[str, Any]: 

289 """Make environment, executable and argument settings for job. 

290 

291 Parameters 

292 ---------- 

293 cached_vals : `dict` [`str`, `~typing.Any`] 

294 Config values common to jobs with same site or label. 

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

296 Generic workflow that contains job to being converted. 

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

298 Generic workflow job to be converted. 

299 

300 Returns 

301 ------- 

302 jobcmds : `dict` [`str` `Any`] 

303 Commands to add to HTC submit description. 

304 """ 

305 jobcmds = {} 

306 

307 job_exports = "" 

308 htc_envs = "" 

309 if gwjob.environment: 

310 use_htc_env = cached_vals.get("bpsUseHTCEnvironment", False) 

311 _LOG.debug("_translate_command_line: use_htc_env = %s", use_htc_env) 

312 if use_htc_env: 

313 # Even though it seems like just using getenv environment will 

314 # work, we must use HTCondor env syntax to get submit side value. 

315 # An environment variable defined in the job description just 

316 # overrides any value from getenv. So we can't ask the job 

317 # description to prepend/append to the value from getenv. 

318 _fix_env = _fix_env_var_syntax 

319 else: 

320 # If not using environment in the job description, setting the 

321 # environment is implemented as exports in the commands run via 

322 # the shell. 

323 _fix_env = _fix_env_var_syntax_shell 

324 for name, value in gwjob.environment.items(): 

325 if isinstance(value, str): 

326 value = _replace_wms_vars(value) 

327 value = _fix_env(value) 

328 value = htc_escape(value) 

329 if use_htc_env: 

330 htc_envs += f"{name}='{value}' " # Add single quotes to allow internal spaces 

331 else: 

332 job_exports += f"export {name}='{value}';" 

333 

334 # Process above added one trailing space 

335 if use_htc_env: 

336 jobcmds["environment"] = htc_envs.rstrip() 

337 _LOG.debug("_translate_command_line: saving htc environment = %s", jobcmds["environment"]) 

338 

339 arguments = "" 

340 if gwjob.arguments: 

341 arguments = gwjob.arguments 

342 arguments = _replace_cmd_vars(arguments, gwjob) 

343 arguments = _replace_wms_vars(arguments) 

344 arguments = _replace_file_vars(cached_vals["bpsUseShared"], arguments, generic_workflow, gwjob) 

345 

346 if cached_vals.get("bpsMakeCommand", True): 

347 # Way to have fallback to previous behavior as well as 

348 # a way forward to centralize logic in bps. 

349 

350 jobcmds["getenv"] = "True" 

351 

352 if gwjob.executable.transfer_executable: 

353 jobcmds["transfer_executable"] = "True" 

354 jobcmds["executable"] = gwjob.executable.src_uri 

355 else: 

356 jobcmds["executable"] = _fix_env_var_syntax(gwjob.executable.src_uri) 

357 

358 if arguments: 

359 arguments = _fix_env_var_syntax(arguments) 

360 jobcmds["arguments"] = arguments 

361 

362 else: 

363 # Instead of making a bash script, run /bin/bash -c <commands> 

364 # HTCondor v25 has a job command called shell that can replace the 

365 # /bin/bash when we get to that version. 

366 

367 # Don't set getenv as setting up the environment is assumed to be 

368 # part of the payloadCommand. 

369 

370 if arguments: 

371 arguments = _fix_env_var_syntax_shell(arguments) 

372 

373 if gwjob.executable.transfer_executable: 

374 # Since replacing executable need to add this executable to the 

375 # file transfer list. 

376 gwfile = GenericWorkflowFile( 

377 name=gwjob.executable.name, src_uri=gwjob.executable.src_uri, wms_transfer=True 

378 ) 

379 generic_workflow.add_job_inputs(gwjob.name, [gwfile]) 

380 exec_name = os.path.basename(gwjob.executable.src_uri) 

381 # Ensure the executable copy is executable. 

382 gwjob_command = f"chmod u+x {exec_name}; ./{exec_name} {arguments}" 

383 else: 

384 exec_name = _fix_env_var_syntax_shell(gwjob.executable.src_uri) 

385 gwjob_command = f"{exec_name} {arguments}" 

386 

387 payload_command = cached_vals["payloadCommand"] 

388 _LOG.debug("%s payload_command pre-format: %s", gwjob.label, payload_command) 

389 payload_command = re.sub("{gwjobCommand}", gwjob_command, payload_command) 

390 payload_command = re.sub("{gwjobExports}", job_exports, payload_command) 

391 

392 # Remove newlines 

393 payload_command = re.sub("\n", "", payload_command) 

394 

395 _LOG.debug("%s payload_command post-format: %s", gwjob.label, payload_command) 

396 

397 jobcmds["arguments"] = f"-c '{payload_command}'" 

398 

399 jobcmds["executable"] = "/bin/bash" 

400 # Don't need to transfer /bin/bash 

401 jobcmds["transfer_executable"] = "False" 

402 

403 return jobcmds 

404 

405 

406def _translate_dag_cmds(gwjob): 

407 """Translate job values into DAGMan commands. 

408 

409 Parameters 

410 ---------- 

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

412 Job containing values to be translated. 

413 

414 Returns 

415 ------- 

416 dagcmds : `dict` [`str`, `~typing.Any`] 

417 DAGMan commands for the job. 

418 """ 

419 # Values in the dag script that just are name mappings. 

420 dag_translation = { 

421 "abort_on_value": "abort_dag_on", 

422 "abort_return_value": "abort_exit", 

423 "priority": "priority", 

424 } 

425 

426 dagcmds = {} 

427 for gwkey, htckey in dag_translation.items(): 

428 dagcmds[htckey] = getattr(gwjob, gwkey, None) 

429 

430 # Still to be coded: vars "pre_cmdline", "post_cmdline" 

431 return dagcmds 

432 

433 

434def _fix_env_var_syntax(oldstr): 

435 """Change ENV place holders to HTCondor Env var syntax. 

436 

437 Parameters 

438 ---------- 

439 oldstr : `str` 

440 String in which environment variable syntax is to be fixed. 

441 

442 Returns 

443 ------- 

444 newstr : `str` 

445 Given string with environment variable syntax fixed. 

446 """ 

447 newstr = oldstr 

448 for key in re.findall(r"<ENV:([^>]+)>", oldstr): 

449 newstr = newstr.replace(rf"<ENV:{key}>", f"$ENV({key})") 

450 return newstr 

451 

452 

453def _fix_env_var_syntax_shell(oldstr): 

454 """Change ENV place holders to shell var syntax. 

455 

456 Parameters 

457 ---------- 

458 oldstr : `str` 

459 String in which environment variable syntax is to be fixed. 

460 

461 Returns 

462 ------- 

463 newstr : `str` 

464 Given string with environment variable syntax fixed. 

465 """ 

466 newstr = oldstr 

467 for key in re.findall(r"<ENV:([^>]+)>", oldstr): 

468 newstr = newstr.replace(rf"<ENV:{key}>", f"${{{key}}}") 

469 return newstr 

470 

471 

472def _replace_file_vars(use_shared, arguments, workflow, gwjob): 

473 """Replace file placeholders in command line arguments with correct 

474 physical file names. 

475 

476 Parameters 

477 ---------- 

478 use_shared : `bool` 

479 Whether HTCondor can assume shared filesystem. 

480 arguments : `str` 

481 Arguments string in which to replace file placeholders. 

482 workflow : `lsst.ctrl.bps.GenericWorkflow` 

483 Generic workflow that contains file information. 

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

485 The job corresponding to the arguments. 

486 

487 Returns 

488 ------- 

489 arguments : `str` 

490 Given arguments string with file placeholders replaced. 

491 """ 

492 # Replace input file placeholders with paths. 

493 for gwfile in workflow.get_job_inputs(gwjob.name, data=True, transfer_only=False): 

494 if not gwfile.wms_transfer: 494 ↛ 497line 494 didn't jump to line 497 because the condition on line 494 was never true

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

496 # responsible for transferring file. 

497 uri = gwfile.src_uri 

498 elif use_shared: 498 ↛ 505line 498 didn't jump to line 505 because the condition on line 498 was always true

499 if gwfile.job_shared: 499 ↛ 503line 499 didn't jump to line 503 because the condition on line 499 was always true

500 # Have shared filesystems and jobs can share file. 

501 uri = gwfile.src_uri 

502 else: 

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

504 else: # Using push transfer 

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

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

507 

508 # Replace output file placeholders with paths. 

509 for gwfile in workflow.get_job_outputs(gwjob.name, data=True, transfer_only=False): 509 ↛ 510line 509 didn't jump to line 510 because the loop on line 509 never started

510 if not gwfile.wms_transfer: 

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

512 # responsible for transferring file. 

513 uri = gwfile.src_uri 

514 elif use_shared: 

515 if gwfile.job_shared: 

516 # Have shared filesystems and jobs can share file. 

517 uri = gwfile.src_uri 

518 else: 

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

520 else: # Using push transfer 

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

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

523 return arguments 

524 

525 

526def _replace_cmd_vars(arguments, gwjob): 

527 """Replace format-style placeholders in arguments. 

528 

529 Parameters 

530 ---------- 

531 arguments : `str` 

532 Arguments string in which to replace placeholders. 

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

534 Job containing values to be used to replace placeholders 

535 (in particular gwjob.cmdvals). 

536 

537 Returns 

538 ------- 

539 arguments : `str` 

540 Given arguments string with placeholders replaced. 

541 """ 

542 replacements = gwjob.cmdvals if gwjob.cmdvals is not None else {} 

543 try: 

544 arguments = arguments.format(**replacements) 

545 except (KeyError, TypeError) as exc: # TypeError in case None instead of {} 

546 _LOG.error( 

547 "Could not replace command variables for job %s: replacement for %s not provided", 

548 gwjob.name, 

549 str(exc), 

550 ) 

551 _LOG.debug("arguments: %s\ncmdvals: %s", arguments, replacements) 

552 raise 

553 return arguments 

554 

555 

556def _replace_wms_vars(orig_string: str) -> str: 

557 """Replace special wms placeholders in given string. 

558 

559 Parameters 

560 ---------- 

561 orig_string : `str` 

562 String in which to replace wms placeholders. 

563 

564 Returns 

565 ------- 

566 updated_string : `str` 

567 Given string with wms placeholders replaced. 

568 """ 

569 values = {"attemptNum": "$$([NumJobStarts])"} 

570 updated_string = orig_string 

571 for key in re.findall(r"<WMS:([^>]+)>", orig_string): 

572 try: 

573 updated_string = updated_string.replace(rf"<WMS:{key}>", values[key]) 

574 except KeyError: 

575 _LOG.error("Unrecognized WMS placeholder: %s in %s", key, orig_string) 

576 raise 

577 return updated_string 

578 

579 

580def _handle_job_inputs( 

581 generic_workflow: GenericWorkflow, job_name: str, use_shared: bool, out_prefix: str 

582) -> dict[str, str]: 

583 """Add job input files from generic workflow to job. 

584 

585 Parameters 

586 ---------- 

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

588 The generic workflow (e.g., has executable name and arguments). 

589 job_name : `str` 

590 Unique name for the job. 

591 use_shared : `bool` 

592 Whether job has access to files via shared filesystem. 

593 out_prefix : `str` 

594 The root directory into which all WMS-specific files are written. 

595 

596 Returns 

597 ------- 

598 htc_commands : `dict` [`str`, `str`] 

599 HTCondor commands for the job submission script. 

600 """ 

601 inputs = [] 

602 for gwf_file in generic_workflow.get_job_inputs(job_name, data=True, transfer_only=True): 602 ↛ 603line 602 didn't jump to line 603 because the loop on line 602 never started

603 _LOG.debug("src_uri=%s", gwf_file.src_uri) 

604 

605 uri = Path(gwf_file.src_uri) 

606 

607 # Note if use_shared and job_shared, don't need to transfer file. 

608 

609 if not use_shared: # Copy file using push to job 

610 inputs.append(str(uri)) 

611 elif not gwf_file.job_shared: # Jobs require own copy 

612 # if using shared filesystem, but still need copy in job. Use 

613 # HTCondor's curl plugin for a local copy. 

614 if uri.is_dir(): 

615 raise RuntimeError( 

616 f"HTCondor plugin cannot transfer directories locally within job {gwf_file.src_uri}" 

617 ) 

618 inputs.append(f"file://{uri}") 

619 

620 htc_commands = {} 

621 if inputs: 621 ↛ 622line 621 didn't jump to line 622 because the condition on line 621 was never true

622 htc_commands["transfer_input_files"] = ",".join(inputs) 

623 _LOG.debug("transfer_input_files=%s", htc_commands["transfer_input_files"]) 

624 return htc_commands 

625 

626 

627def _handle_job_outputs( 

628 generic_workflow: GenericWorkflow, job_name: str, use_shared: bool, out_prefix: str 

629) -> dict[str, str]: 

630 """Add job output files from generic workflow to the job if any. 

631 

632 Parameters 

633 ---------- 

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

635 The generic workflow (e.g., has executable name and arguments). 

636 job_name : `str` 

637 Unique name for the job. 

638 use_shared : `bool` 

639 Whether job has access to files via shared filesystem. 

640 out_prefix : `str` 

641 The root directory into which all WMS-specific files are written. 

642 

643 Returns 

644 ------- 

645 htc_commands : `dict` [`str`, `str`] 

646 HTCondor commands for the job submission script. 

647 """ 

648 outputs = [] 

649 output_remaps = [] 

650 for gwf_file in generic_workflow.get_job_outputs(job_name, data=True, transfer_only=True): 

651 _LOG.debug("src_uri=%s", gwf_file.src_uri) 

652 

653 uri = Path(gwf_file.src_uri) 

654 if not use_shared: 

655 outputs.append(uri.name) 

656 output_remaps.append(f"{uri.name}={str(uri)}") 

657 

658 # Set to an empty string to disable and only update if there are output 

659 # files to transfer. Otherwise, HTCondor will transfer back all files in 

660 # the job’s temporary working directory that have been modified or created 

661 # by the job. 

662 htc_commands = {"transfer_output_files": '""'} 

663 if outputs: 

664 htc_commands["transfer_output_files"] = ",".join(outputs) 

665 _LOG.debug("transfer_output_files=%s", htc_commands["transfer_output_files"]) 

666 

667 htc_commands["transfer_output_remaps"] = f'"{";".join(output_remaps)}"' 

668 _LOG.debug("transfer_output_remaps=%s", htc_commands["transfer_output_remaps"]) 

669 return htc_commands 

670 

671 

672def _create_periodic_release_expr( 

673 memory: int, multiplier: float | None, limit: int, additional_expr: str = "" 

674) -> str: 

675 """Construct an HTCondorAd expression for releasing held jobs. 

676 

677 Parameters 

678 ---------- 

679 memory : `int` 

680 Requested memory in MB. 

681 multiplier : `float` or None 

682 Memory growth rate between retries. 

683 limit : `int` 

684 Memory limit. 

685 additional_expr : `str`, optional 

686 Expression to add to periodic_release. Defaults to empty string. 

687 

688 Returns 

689 ------- 

690 expr : `str` 

691 A string representing an HTCondor ClassAd expression for releasing job. 

692 """ 

693 _LOG.debug( 

694 "periodic_release: memory: %s, multiplier: %s, limit: %s, additional_expr: %s", 

695 memory, 

696 multiplier, 

697 limit, 

698 additional_expr, 

699 ) 

700 

701 # ctrl_bps sets multiplier to None in the GenericWorkflow if 

702 # memoryMultiplier <= 1, but checking value just in case. 

703 if (not multiplier or multiplier <= 1) and not additional_expr: 

704 return "" 

705 

706 # Job ClassAds attributes 'HoldReasonCode' and 'HoldReasonSubCode' are 

707 # UNDEFINED if job is not HELD (i.e. when 'JobStatus' is not 5). 

708 # The special comparison operators ensure that all comparisons below will 

709 # evaluate to FALSE in this case. 

710 # 

711 # Note: 

712 # May not be strictly necessary. Operators '&&' and '||' are not strict so 

713 # the entire expression should evaluate to FALSE when the job is not HELD. 

714 # According to ClassAd evaluation semantics FALSE && UNDEFINED is FALSE, 

715 # but better safe than sorry. 

716 is_held = "JobStatus == 5" 

717 is_retry_allowed = "NumJobStarts <= JobMaxRetries" 

718 

719 mem_expr = "" 

720 if memory and multiplier and multiplier > 1 and limit: 

721 was_mem_exceeded = ( 

722 "(HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 " 

723 "|| HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34)" 

724 ) 

725 was_below_limit = f"min({{int({memory} * pow({multiplier}, NumJobStarts - 1)), {limit}}}) < {limit}" 

726 mem_expr = f"{was_mem_exceeded} && {was_below_limit}" 

727 

728 user_expr = "" 

729 if additional_expr: 

730 # Never auto release a job held by user. 

731 user_expr = f"HoldReasonCode =!= 1 && {additional_expr}" 

732 

733 # Automatically release job if held because output file not found 

734 # (e.g., job failed so didn't produce output file). 

735 transfer_expr = "HoldReasonCode =?= 12" 

736 

737 expr = f"{is_held} && {is_retry_allowed}" 

738 if user_expr and mem_expr: 

739 expr += f" && ({transfer_expr} || {mem_expr} || {user_expr})" 

740 elif user_expr: 

741 expr += f" && ({transfer_expr} || {user_expr})" 

742 elif mem_expr: 742 ↛ 745line 742 didn't jump to line 745 because the condition on line 742 was always true

743 expr += f" && ({transfer_expr} || {mem_expr})" 

744 

745 return expr 

746 

747 

748def _create_periodic_remove_expr(memory, multiplier, limit): 

749 """Construct an HTCondorAd expression for removing jobs from the queue. 

750 

751 Parameters 

752 ---------- 

753 memory : `int` 

754 Requested memory in MB. 

755 multiplier : `float` 

756 Memory growth rate between retries. 

757 limit : `int` 

758 Memory limit. 

759 

760 Returns 

761 ------- 

762 expr : `str` 

763 A string representing an HTCondor ClassAd expression for removing jobs. 

764 """ 

765 # Job ClassAds attributes 'HoldReasonCode' and 'HoldReasonSubCode' 

766 # are UNDEFINED if job is not HELD (i.e. when 'JobStatus' is not 5). 

767 # The special comparison operators ensure that all comparisons below 

768 # will evaluate to FALSE in this case. 

769 # 

770 # Note: 

771 # May not be strictly necessary. Operators '&&' and '||' are not 

772 # strict so the entire expression should evaluate to FALSE when the 

773 # job is not HELD. According to ClassAd evaluation semantics 

774 # FALSE && UNDEFINED is FALSE, but better safe than sorry. 

775 is_held = "JobStatus == 5" 

776 is_retry_disallowed = "NumJobStarts > JobMaxRetries" 

777 

778 mem_expr = "" 

779 if memory and multiplier and multiplier > 1 and limit: 

780 mem_limit_expr = f"min({{int({memory} * pow({multiplier}, NumJobStarts - 1)), {limit}}}) == {limit}" 

781 

782 mem_expr = ( # Add || here so only added if adding memory expr 

783 " || ((HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 " 

784 f"|| HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && {mem_limit_expr})" 

785 ) 

786 

787 expr = f"{is_held} && ({is_retry_disallowed}{mem_expr})" 

788 return expr 

789 

790 

791def _create_request_memory_expr(memory, multiplier, limit): 

792 """Construct an HTCondor ClassAd expression for safe memory scaling. 

793 

794 Parameters 

795 ---------- 

796 memory : `int` 

797 Requested memory in MB. 

798 multiplier : `float` 

799 Memory growth rate between retries. 

800 limit : `int` 

801 Memory limit. 

802 

803 Returns 

804 ------- 

805 expr : `str` 

806 A string representing an HTCondor ClassAd expression enabling safe 

807 memory scaling between job retries. 

808 """ 

809 # The check if the job was held due to exceeding memory requirements 

810 # will be made *after* job was released back to the job queue (is in 

811 # the IDLE state), hence the need to use `Last*` job ClassAds instead of 

812 # the ones describing job's current state. 

813 # 

814 # Also, 'Last*' job ClassAds attributes are UNDEFINED when a job is 

815 # initially put in the job queue. The special comparison operators ensure 

816 # that all comparisons below will evaluate to FALSE in this case. 

817 was_mem_exceeded = ( 

818 "LastJobStatus =?= 5 " 

819 "&& (LastHoldReasonCode =?= 34 && LastHoldReasonSubCode =?= 0 " 

820 "|| LastHoldReasonCode =?= 3 && LastHoldReasonSubCode =?= 34)" 

821 ) 

822 

823 # If job runs the first time or was held for reasons other than exceeding 

824 # the memory, set the required memory to the requested value or use 

825 # the memory value measured by HTCondor (MemoryUsage) depending on 

826 # whichever is greater. 

827 expr = ( 

828 f"({was_mem_exceeded}) " 

829 f"? min({{int({memory} * pow({multiplier}, NumJobStarts)), {limit}}}) " 

830 f": min({{max({{{memory}, MemoryUsage ?: 0}}), {limit}}})" 

831 ) 

832 return expr 

833 

834 

835def _gather_site_values(config, compute_site): 

836 """Gather values specific to given site. 

837 

838 Parameters 

839 ---------- 

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

841 BPS configuration that includes necessary submit/runtime 

842 information. 

843 compute_site : `str` 

844 Compute site name. 

845 

846 Returns 

847 ------- 

848 site_values : `dict` [`str`, `~typing.Any`] 

849 Values specific to the given site. 

850 """ 

851 site_values = {"attrs": {}, "profile": {}} 

852 search_opts = {} 

853 if compute_site: 853 ↛ 857line 853 didn't jump to line 857 because the condition on line 853 was always true

854 search_opts["curvals"] = {"curr_site": compute_site} 

855 

856 # Determine the hard limit for the memory requirement. 

857 found, limit = config.search("memoryLimit", opt=search_opts) 

858 if not found: 858 ↛ 859line 858 didn't jump to line 859 because the condition on line 858 was never true

859 search_opts["default"] = DEFAULT_HTC_EXEC_PATT 

860 _, patt = config.search("executeMachinesPattern", opt=search_opts) 

861 del search_opts["default"] 

862 

863 # To reduce the amount of data, ignore dynamic slots (if any) as, 

864 # by definition, they cannot have more memory than 

865 # the partitionable slot they are the part of. 

866 constraint = f'SlotType != "Dynamic" && regexp("{patt}", Machine)' 

867 pool_info = condor_status(constraint=constraint) 

868 try: 

869 limit = max(int(info["TotalSlotMemory"]) for info in pool_info.values()) 

870 except ValueError: 

871 _LOG.debug("No execute machine in the pool matches %s", patt) 

872 if limit: 872 ↛ 874line 872 didn't jump to line 874 because the condition on line 872 was always true

873 config[".bps_defined.memory_limit"] = limit 

874 site_values["memoryLimit"] = limit 

875 

876 _, site_values["bpsUseShared"] = config.search("bpsUseShared", opt={"default": False}) 

877 

878 searchobj = config[f".site.{compute_site}.profile.condor"] 

879 if searchobj: 

880 search_opts["searchobj"] = searchobj 

881 search_opts["replaceVars"] = True 

882 for key in searchobj: 

883 if key.startswith("+"): 

884 _, val = config.search(key, opt=search_opts) 

885 site_values["attrs"][key[1:]] = val 

886 else: 

887 _, val = config.search(key, opt=search_opts) 

888 site_values["profile"][key] = val 

889 

890 searchobj = config[f".site.{compute_site}"] 

891 if searchobj: 

892 for key, value in searchobj.items(): 

893 if key not in site_values and key not in ["attrs", "profile"]: 893 ↛ 894line 893 didn't jump to line 894 because the condition on line 893 was never true

894 site_values[key] = value 

895 

896 _LOG.debug("site_values = %s", site_values) 

897 return site_values 

898 

899 

900def _gather_label_values(config: BpsConfig, label: str) -> dict[str, Any]: 

901 """Gather values specific to given job label. 

902 

903 Parameters 

904 ---------- 

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

906 BPS configuration that includes necessary submit/runtime 

907 information. 

908 label : `str` 

909 GenericWorkflowJob label. 

910 

911 Returns 

912 ------- 

913 values : `dict` [`str`, `~typing.Any`] 

914 Values specific to the given job label. 

915 """ 

916 _LOG.debug("_gather_label_values: label = %s", label) 

917 values: dict[str, Any] = {"attrs": {}, "profile": {}} 

918 

919 search_opts = config.get_search_opts(label) 

920 

921 # Determine the hard limit for the memory requirement. 

922 found, limit = config.search("memoryLimit", opt=search_opts) 

923 if not found: 923 ↛ 924line 923 didn't jump to line 924 because the condition on line 923 was never true

924 search_opts["default"] = DEFAULT_HTC_EXEC_PATT 

925 _, patt = config.search("executeMachinesPattern", opt=search_opts) 

926 del search_opts["default"] 

927 

928 # To reduce the amount of data, ignore dynamic slots (if any) as, 

929 # by definition, they cannot have more memory than 

930 # the partitionable slot they are the part of. 

931 constraint = f'SlotType != "Dynamic" && regexp("{patt}", Machine)' 

932 pool_info = condor_status(constraint=constraint) 

933 try: 

934 limit = max(int(info["TotalSlotMemory"]) for info in pool_info.values()) 

935 except ValueError: 

936 _LOG.debug("No execute machine in the pool matches %s", patt) 

937 

938 if limit: 938 ↛ 942line 938 didn't jump to line 942 because the condition on line 938 was always true

939 config[".bps_defined.memory_limit"] = limit 

940 values["memoryLimit"] = limit 

941 

942 values["bpsUseShared"] = False 

943 found, value = config.search("bpsUseShared", opt=search_opts) 

944 if found: 

945 values["bpsUseShared"] = value 

946 

947 found, value = config.search("releaseExpr", opt=search_opts) 

948 if found: 

949 values["releaseExpr"] = value 

950 

951 values["overwriteJobFiles"] = True 

952 found, value = config.search("overwriteJobFiles", opt=search_opts) 

953 if found: 

954 values["overwriteJobFiles"] = value 

955 

956 found, value = config.search("releaseExpr", opt=search_opts) 

957 if found: 

958 values["releaseExpr"] = value 

959 

960 found, value = config.search("bpsMakeCommand", opt=search_opts) 

961 values["bpsMakeCommand"] = value if found else True 

962 if found and not value: 

963 search_opts["skipNames"] = {"gwjobCommand", "gwjobExports"} 

964 _LOG.debug("_gather_label_values: search_opts = %s", search_opts) 

965 found, value = config.search("payloadCommand", opt=search_opts) 

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

967 values["payloadCommand"] = value 

968 _LOG.debug("payloadCommand = %s", value) 

969 

970 found, value = config.search("bpsUseHTCEnvironment", opt=search_opts) 

971 values["bpsUseHTCEnvironment"] = value if found else values["bpsMakeCommand"] 

972 

973 found, value = config.search("nodeset", opt=search_opts) 

974 if found: 

975 values["nodeset"] = value 

976 

977 found, profile_sect = config.search("profile", opt=search_opts) 

978 if found and "condor" in profile_sect: 

979 for subkey, val in profile_sect["condor"].items(): 

980 if subkey.startswith("+"): 

981 values["attrs"][subkey[1:]] = val 

982 else: 

983 values["profile"][subkey] = val 

984 

985 # Copy all of the site values. 

986 if "curr_site" in search_opts["curvals"]: 

987 site_obj = config[f".site.{search_opts['curvals']['curr_site']}"] 

988 if site_obj: 988 ↛ 993line 988 didn't jump to line 993 because the condition on line 988 was always true

989 for key, value in site_obj.items(): 

990 if key not in values and key not in ["attrs", "profile"]: 

991 values[key] = value 

992 

993 _LOG.debug("_gather_label_values: label = %s, values = %s", label, values) 

994 

995 return values 

996 

997 

998def _group_to_subdag( 

999 config: BpsConfig, generic_workflow_group: GenericWorkflowGroup, out_prefix: str 

1000) -> HTCJob: 

1001 """Convert a generic workflow group to an HTCondor dag. 

1002 

1003 Parameters 

1004 ---------- 

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

1006 Workflow configuration. 

1007 generic_workflow_group : `lsst.ctrl.bps.GenericWorkflowGroup` 

1008 The generic workflow group to convert. 

1009 out_prefix : `str` 

1010 Location prefix to be used when creating jobs. 

1011 

1012 Returns 

1013 ------- 

1014 htc_job : `lsst.ctrl.bps.htcondor.HTCJob` 

1015 Job for running the HTCondor dag. 

1016 """ 

1017 jobname = f"wms_{generic_workflow_group.name}" 

1018 htc_job = HTCJob(name=jobname, label=generic_workflow_group.label) 

1019 htc_job.add_dag_cmds({"dir": f"subdags/{jobname}"}) 

1020 htc_job.subdag = _generic_workflow_to_htcondor_dag(config, generic_workflow_group, out_prefix) 

1021 if not generic_workflow_group.blocking: 1021 ↛ 1027line 1021 didn't jump to line 1027 because the condition on line 1021 was always true

1022 htc_job.dagcmds["post"] = { 

1023 "defer": "", 

1024 "executable": f"{os.path.dirname(__file__)}/subdag_post.sh", 

1025 "arguments": f"{jobname} $RETURN", 

1026 } 

1027 return htc_job 

1028 

1029 

1030def _create_check_job(group_job_name: str, job_label: str, site_values: dict) -> HTCJob: 

1031 """Create a job to check status of a group job. 

1032 

1033 Parameters 

1034 ---------- 

1035 group_job_name : `str` 

1036 Name of the group job. 

1037 job_label : `str` 

1038 Label to use for the check status job. 

1039 site_values : `dict` 

1040 Site specific values. 

1041 

1042 Returns 

1043 ------- 

1044 htc_job : `lsst.ctrl.bps.htcondor.HTCJob` 

1045 Job description for the job to check group job status. 

1046 """ 

1047 htc_job = HTCJob(name=f"wms_check_status_{group_job_name}", label=job_label) 

1048 htc_job.subfile = "${CTRL_BPS_HTCONDOR_DIR}/python/lsst/ctrl/bps/htcondor/check_group_status.sub" 

1049 # ADD nodeset to VARS 

1050 job_vars = {"group_job_name": group_job_name} 

1051 if "nodeset" in site_values and site_values["nodeset"]: 

1052 job_vars["job_nodeset"] = site_values["nodeset"] 

1053 htc_job.add_dag_cmds({"dir": f"subdags/{group_job_name}", "vars": job_vars}) 

1054 

1055 return htc_job 

1056 

1057 

1058def _generic_workflow_to_htcondor_dag( 

1059 config: BpsConfig, generic_workflow: GenericWorkflow, out_prefix: str 

1060) -> HTCDag: 

1061 """Convert a GenericWorkflow to a HTCDag. 

1062 

1063 Parameters 

1064 ---------- 

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

1066 Workflow configuration. 

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

1068 The GenericWorkflow to convert. 

1069 out_prefix : `str` 

1070 Location prefix where the HTCondor files will be written. 

1071 

1072 Returns 

1073 ------- 

1074 dag : `lsst.ctrl.bps.htcondor.HTCDag` 

1075 The HTCDag representation of the given GenericWorkflow. 

1076 """ 

1077 dag = HTCDag(name=generic_workflow.name) 

1078 

1079 _LOG.debug("htcondor dag attribs %s", generic_workflow.run_attrs) 

1080 dag.add_attribs(generic_workflow.run_attrs) 

1081 dag.add_attribs( 

1082 { 

1083 "bps_run_quanta": create_count_summary(generic_workflow.quanta_counts), 

1084 "bps_job_summary": create_count_summary(generic_workflow.job_counts), 

1085 } 

1086 ) 

1087 _, tmp_template = config.search("subDirTemplate", opt={"replaceVars": False, "default": ""}) 

1088 

1089 _, save_htc_dot = config.search("saveHTCdot", opt={"default": False}) 

1090 dag.graph["write_dot"] = save_htc_dot 

1091 

1092 if isinstance(tmp_template, str): 1092 ↛ 1095line 1092 didn't jump to line 1095 because the condition on line 1092 was always true

1093 subdir_template = defaultdict(lambda: tmp_template) 

1094 else: 

1095 subdir_template = tmp_template 

1096 

1097 # Save list of lazy group jobs for later extra handling. 

1098 lazy_groups = [] 

1099 

1100 # Create all DAG jobs 

1101 cached_values = {} # Cache label-specific values to reduce config lookups. 

1102 # Note: Can't use get_job_by_label because those only include payload jobs. 

1103 for job_name in generic_workflow: 

1104 gwjob = generic_workflow.get_job(job_name) 

1105 if gwjob.node_type in [ 1105 ↛ 1120line 1105 didn't jump to line 1120 because the condition on line 1105 was always true

1106 GenericWorkflowNodeType.PAYLOAD, 

1107 GenericWorkflowNodeType.LAZY_GROUP, 

1108 ]: 

1109 gwjob = cast(GenericWorkflowJob, gwjob) 

1110 if gwjob.label not in cached_values: 

1111 cached_values[gwjob.label] = _gather_label_values(config, gwjob.label) 

1112 _LOG.debug("cached: %s= %s", gwjob.label, cached_values[gwjob.label]) 

1113 htc_job = _create_job( 

1114 subdir_template[gwjob.label], 

1115 cached_values[gwjob.label], 

1116 generic_workflow, 

1117 gwjob, 

1118 out_prefix, 

1119 ) 

1120 elif gwjob.node_type == GenericWorkflowNodeType.NOOP: 

1121 gwjob = cast(GenericWorkflowNoopJob, gwjob) 

1122 htc_job = HTCJob(f"wms_{gwjob.name}", label=gwjob.label) 

1123 htc_job.subfile = "${CTRL_BPS_HTCONDOR_DIR}/python/lsst/ctrl/bps/htcondor/noop.sub" 

1124 htc_job.add_job_attrs({"bps_job_name": gwjob.name, "bps_job_label": gwjob.label}) 

1125 htc_job.add_dag_cmds({"noop": True}) 

1126 elif gwjob.node_type == GenericWorkflowNodeType.GROUP: 

1127 gwjob = cast(GenericWorkflowGroup, gwjob) 

1128 cached_values[gwjob.label] = _gather_label_values(config, gwjob.label) 

1129 htc_job = _group_to_subdag(config, gwjob, out_prefix) 

1130 else: 

1131 raise RuntimeError(f"Unsupported generic workflow node type {gwjob.node_type} ({gwjob.name})") 

1132 _LOG.debug("Calling adding job %s %s", htc_job.name, htc_job.label) 

1133 dag.add_job(htc_job) 

1134 

1135 # Have to add the placeholder job for the workflow 

1136 if gwjob.node_type == GenericWorkflowNodeType.LAZY_GROUP: 

1137 lazy_groups.append(gwjob.name) 

1138 

1139 # Add job dependencies to the DAG (be careful with wms_ jobs) 

1140 for job_name in generic_workflow: 

1141 gwjob = generic_workflow.get_job(job_name) 

1142 parent_name = ( 

1143 gwjob.name 

1144 if gwjob.node_type in [GenericWorkflowNodeType.PAYLOAD, GenericWorkflowNodeType.LAZY_GROUP] 

1145 else f"wms_{gwjob.name}" 

1146 ) 

1147 successor_jobs = [generic_workflow.get_job(j) for j in generic_workflow.successors(job_name)] 

1148 children_names = [] 

1149 if gwjob.node_type == GenericWorkflowNodeType.GROUP: 1149 ↛ 1150line 1149 didn't jump to line 1150 because the condition on line 1149 was never true

1150 gwjob = cast(GenericWorkflowGroup, gwjob) 

1151 group_children = [] # Dependencies between same group jobs 

1152 for sjob in successor_jobs: 

1153 if sjob.node_type == GenericWorkflowNodeType.GROUP and sjob.label == gwjob.label: 

1154 group_children.append(f"wms_{sjob.name}") 

1155 elif sjob.node_type == GenericWorkflowNodeType.PAYLOAD: 

1156 children_names.append(sjob.name) 

1157 else: 

1158 children_names.append(f"wms_{sjob.name}") 

1159 if group_children: 

1160 dag.add_job_relationships([parent_name], group_children) 

1161 if not gwjob.blocking: 

1162 # Since subdag will always succeed, need to add a special 

1163 # job that fails if group failed to block payload children. 

1164 check_job = _create_check_job( 

1165 f"wms_{gwjob.name}", gwjob.label, cached_values.get(gwjob.label, {}) 

1166 ) 

1167 dag.add_job(check_job) 

1168 dag.add_job_relationships([f"wms_{gwjob.name}"], [check_job.name]) 

1169 parent_name = check_job.name 

1170 else: 

1171 for sjob in successor_jobs: 

1172 if sjob.node_type in [ 1172 ↛ 1178line 1172 didn't jump to line 1178 because the condition on line 1172 was always true

1173 GenericWorkflowNodeType.PAYLOAD, 

1174 GenericWorkflowNodeType.LAZY_GROUP, 

1175 ]: 

1176 children_names.append(sjob.name) 

1177 else: 

1178 children_names.append(f"wms_{sjob.name}") 

1179 

1180 dag.add_job_relationships([parent_name], children_names) 

1181 

1182 # Go back and add placeholder jobs for the lazy group dags 

1183 for lazy_group_name in lazy_groups: 

1184 _add_lazy_placeholder(lazy_group_name, generic_workflow, dag, out_prefix) 

1185 

1186 # If final job exists in generic workflow, create DAG final job 

1187 final = generic_workflow.get_final() 

1188 if final and isinstance(final, GenericWorkflowJob): 

1189 if final.label not in cached_values: 1189 ↛ 1191line 1189 didn't jump to line 1191 because the condition on line 1189 was always true

1190 cached_values[final.label] = _gather_label_values(config, final.label) 

1191 final_htjob = _create_job( 

1192 subdir_template[final.label], 

1193 cached_values[final.label], 

1194 generic_workflow, 

1195 final, 

1196 out_prefix, 

1197 ) 

1198 if "post" not in final_htjob.dagcmds: 1198 ↛ 1204line 1198 didn't jump to line 1204 because the condition on line 1198 was always true

1199 final_htjob.dagcmds["post"] = { 

1200 "defer": "", 

1201 "executable": f"{os.path.dirname(__file__)}/final_post.sh", 

1202 "arguments": f"{final.name} $DAG_STATUS $RETURN", 

1203 } 

1204 dag.add_final_job(final_htjob) 

1205 elif final and isinstance(final, GenericWorkflow): 

1206 raise NotImplementedError("HTCondor plugin does not support a workflow as the final job") 

1207 elif final: 1207 ↛ 1208line 1207 didn't jump to line 1208 because the condition on line 1207 was never true

1208 raise TypeError(f"Invalid type for GenericWorkflow.get_final() results ({type(final)})") 

1209 

1210 return dag 

1211 

1212 

1213def _add_lazy_placeholder( 

1214 prepare_job_name: str, 

1215 generic_workflow: GenericWorkflow, 

1216 dag: HTCDag, 

1217 out_prefix: str, 

1218): 

1219 _LOG.debug("prepare_job_name = %s", prepare_job_name) 

1220 

1221 # Make a fake job for the placeholder workflow 

1222 job = HTCJob("placeholder") 

1223 job.add_job_cmds( 

1224 { 

1225 "executable": "/usr/bin/echo", 

1226 "arguments": '"BPS internal error - placeholder DAG was not replaced."', 

1227 } 

1228 ) 

1229 job.add_job_attrs({"bps_job_name": "placeholder", "bps_job_label": "placeholder"}) 

1230 job.add_job_attrs(generic_workflow.run_attrs) 

1231 

1232 # Placeholder dag name needs to be the run name because that's 

1233 # currently what the bps code will name the workflow. 

1234 placeholder_dag_name = generic_workflow.name.replace("_ctrl", "") 

1235 placeholder_dag = HTCDag(name=placeholder_dag_name) 

1236 _LOG.debug("dag name = %s", placeholder_dag.graph["name"]) 

1237 placeholder_dag.add_attribs(generic_workflow.run_attrs) 

1238 placeholder_dag.add_job(job) 

1239 

1240 # To help ordering in bps report, save info so can figure out where 

1241 # to put jobs in lazy dag in bps_job_label 

1242 dag.add_attribs({"bps_lazy_mapping": f"{placeholder_dag.name}:{prepare_job_name}"}) 

1243 

1244 # The subdag job to be added to the control dag 

1245 dag_job = HTCJob("wms_lazy_payload", "wms_lazy_payload") 

1246 dag_job.subfile = f"{placeholder_dag.name}.condor.sub" 

1247 dag_job.subdag = placeholder_dag 

1248 

1249 dag.add_job(dag_job) 

1250 

1251 # Update edges inserting between prepare job and any following jobs. 

1252 successor_jobs = list(dag.successors(prepare_job_name)) 

1253 for job in successor_jobs: 

1254 dag.add_edge(dag_job.name, job) 

1255 dag.remove_edge(prepare_job_name, job) 

1256 

1257 dag.add_edge(prepare_job_name, dag_job.name) 

1258 

1259 _LOG.debug("_add_lazy_placeholder: edges = %s", dag.edges) 

1260 

1261 

1262def _update_job_summary(subworkflow_name: str, subworkflow_summary: str, submit_path: str) -> None: 

1263 """Add summary for jobs in given subworkflow to parent workflow's job 

1264 summary. 

1265 

1266 Parameters 

1267 ---------- 

1268 subworkflow_name : `str` 

1269 Name of subworkflow. 

1270 subworkflow_summary : `str` 

1271 Job summary for the subworkflow. 

1272 submit_path : `str` 

1273 Directory in which to find the DAG info file. 

1274 """ 

1275 _LOG.debug("submit_path = %s", submit_path) 

1276 filename, dag_info = read_dag_info(submit_path) 

1277 _LOG.debug("dag_info = %s", dag_info) 

1278 

1279 schedd_name = next(iter(dag_info)) 

1280 dag_values = next(iter(dag_info[schedd_name].values())) 

1281 _LOG.debug("dag_values = %s", dag_values) 

1282 

1283 # Get lazy mapping and the job that generated this dag 

1284 generator_name = None 

1285 lazy_mapping = dag_values.get("bps_lazy_mapping", None) 

1286 if lazy_mapping: # find this workflow's generator job 

1287 for part in lazy_mapping.split(";"): 

1288 info = part.split(":") 

1289 if info[0] == subworkflow_name: 

1290 generator_name = info[1] 

1291 break 

1292 

1293 # Update bps_job_summary 

1294 _LOG.debug( 

1295 "Before replace, name = %s, bps_job_summary = %s, add summary = %s, generator_name = %s", 

1296 subworkflow_name, 

1297 dag_values["bps_job_summary"], 

1298 subworkflow_summary, 

1299 generator_name, 

1300 ) 

1301 if generator_name: 

1302 generator_summary = f"{generator_name}:1" 

1303 dag_values["bps_job_summary"] = dag_values["bps_job_summary"].replace( 

1304 generator_summary, f"{generator_summary};{subworkflow_summary}" 

1305 ) 

1306 else: 

1307 # just append to end of bps_job_summary 

1308 dag_values["bps_job_summary"] += f";{subworkflow_summary}" 

1309 

1310 _LOG.debug( 

1311 "After replace, name = %s, bps_job_summary = %s", 

1312 subworkflow_name, 

1313 dag_values["bps_job_summary"], 

1314 ) 

1315 

1316 # Save updated bps_job_summary 

1317 write_dag_info(filename, dag_info)