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

508 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-09-06 01:56 -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): 325 ↛ 329line 325 didn't jump to line 329 because the condition on line 325 was always true

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 if cached_vals.get("bpsMakeCommand", True): 

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

341 # a way forward to centralize logic in bps. 

342 

343 jobcmds["getenv"] = "True" 

344 

345 if gwjob.executable.transfer_executable: 

346 jobcmds["transfer_executable"] = "True" 

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

348 else: 

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

350 

351 if gwjob.arguments: 

352 arguments = gwjob.arguments 

353 arguments = _replace_cmd_vars(arguments, gwjob) 

354 arguments = _replace_wms_vars(arguments) 

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

356 arguments = _fix_env_var_syntax(arguments) 

357 jobcmds["arguments"] = arguments 

358 

359 else: 

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

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

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

363 

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

365 # part of the payloadCommand. 

366 

367 if gwjob.arguments: 367 ↛ 374line 367 didn't jump to line 374 because the condition on line 367 was always true

368 arguments = gwjob.arguments 

369 arguments = _replace_cmd_vars(arguments, gwjob) 

370 arguments = _replace_wms_vars(arguments) 

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

372 arguments = _fix_env_var_syntax_shell(arguments) 

373 

374 if gwjob.executable.transfer_executable: 

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

376 # file transfer list. 

377 gwfile = GenericWorkflowFile( 

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

379 ) 

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

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

382 # Ensure the executable copy is executable. 

383 gwjobCommand = f"chmod u+x {exec_name}; ./{exec_name} {arguments}" 

384 else: 

385 exec_name = _fix_env_var_syntax_shell(gwjob.executable.src_uri) 

386 gwjobCommand = f"{exec_name} {arguments}" 

387 

388 payloadCommand = cached_vals["payloadCommand"] 

389 _LOG.debug("%s payloadCommand pre-format: %s", gwjob.label, payloadCommand) 

390 payloadCommand = re.sub("{gwjobCommand}", gwjobCommand, payloadCommand) 

391 payloadCommand = re.sub("{gwjobExports}", job_exports, payloadCommand) 

392 

393 # Remove newlines 

394 payloadCommand = re.sub("\n", "", payloadCommand) 

395 

396 _LOG.debug("%s payloadCommand post-format: %s", gwjob.label, payloadCommand) 

397 

398 # jobcmds["arguments"] = htc_escape(f"-c '{payloadCommand}'") 

399 jobcmds["arguments"] = f"-c '{payloadCommand}'" 

400 

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

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

403 jobcmds["transfer_executable"] = "False" 

404 

405 return jobcmds 

406 

407 

408def _translate_dag_cmds(gwjob): 

409 """Translate job values into DAGMan commands. 

410 

411 Parameters 

412 ---------- 

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

414 Job containing values to be translated. 

415 

416 Returns 

417 ------- 

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

419 DAGMan commands for the job. 

420 """ 

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

422 dag_translation = { 

423 "abort_on_value": "abort_dag_on", 

424 "abort_return_value": "abort_exit", 

425 "priority": "priority", 

426 } 

427 

428 dagcmds = {} 

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

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

431 

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

433 return dagcmds 

434 

435 

436def _fix_env_var_syntax(oldstr): 

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

438 

439 Parameters 

440 ---------- 

441 oldstr : `str` 

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

443 

444 Returns 

445 ------- 

446 newstr : `str` 

447 Given string with environment variable syntax fixed. 

448 """ 

449 newstr = oldstr 

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

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

452 return newstr 

453 

454 

455def _fix_env_var_syntax_shell(oldstr): 

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

457 

458 Parameters 

459 ---------- 

460 oldstr : `str` 

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

462 

463 Returns 

464 ------- 

465 newstr : `str` 

466 Given string with environment variable syntax fixed. 

467 """ 

468 newstr = oldstr 

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

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

471 return newstr 

472 

473 

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

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

476 physical file names. 

477 

478 Parameters 

479 ---------- 

480 use_shared : `bool` 

481 Whether HTCondor can assume shared filesystem. 

482 arguments : `str` 

483 Arguments string in which to replace file placeholders. 

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

485 Generic workflow that contains file information. 

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

487 The job corresponding to the arguments. 

488 

489 Returns 

490 ------- 

491 arguments : `str` 

492 Given arguments string with file placeholders replaced. 

493 """ 

494 # Replace input file placeholders with paths. 

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

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

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

498 # responsible for transferring file. 

499 uri = gwfile.src_uri 

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

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

502 # Have shared filesystems and jobs can share file. 

503 uri = gwfile.src_uri 

504 else: 

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

506 else: # Using push transfer 

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

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

509 

510 # Replace output file placeholders with paths. 

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

512 if not gwfile.wms_transfer: 

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

514 # responsible for transferring file. 

515 uri = gwfile.src_uri 

516 elif use_shared: 

517 if gwfile.job_shared: 

518 # Have shared filesystems and jobs can share file. 

519 uri = gwfile.src_uri 

520 else: 

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

522 else: # Using push transfer 

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

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

525 return arguments 

526 

527 

528def _replace_cmd_vars(arguments, gwjob): 

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

530 

531 Parameters 

532 ---------- 

533 arguments : `str` 

534 Arguments string in which to replace placeholders. 

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

536 Job containing values to be used to replace placeholders 

537 (in particular gwjob.cmdvals). 

538 

539 Returns 

540 ------- 

541 arguments : `str` 

542 Given arguments string with placeholders replaced. 

543 """ 

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

545 try: 

546 arguments = arguments.format(**replacements) 

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

548 _LOG.error( 

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

550 gwjob.name, 

551 str(exc), 

552 ) 

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

554 raise 

555 return arguments 

556 

557 

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

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

560 

561 Parameters 

562 ---------- 

563 orig_string : `str` 

564 String in which to replace wms placeholders. 

565 

566 Returns 

567 ------- 

568 updated_string : `str` 

569 Given string with wms placeholders replaced. 

570 """ 

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

572 updated_string = orig_string 

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

574 try: 

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

576 except KeyError: 

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

578 raise 

579 return updated_string 

580 

581 

582def _handle_job_inputs( 

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

584) -> dict[str, str]: 

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

586 

587 Parameters 

588 ---------- 

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

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

591 job_name : `str` 

592 Unique name for the job. 

593 use_shared : `bool` 

594 Whether job has access to files via shared filesystem. 

595 out_prefix : `str` 

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

597 

598 Returns 

599 ------- 

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

601 HTCondor commands for the job submission script. 

602 """ 

603 inputs = [] 

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

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

606 

607 uri = Path(gwf_file.src_uri) 

608 

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

610 

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

612 inputs.append(str(uri)) 

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

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

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

616 if uri.is_dir(): 

617 raise RuntimeError( 

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

619 ) 

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

621 

622 htc_commands = {} 

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

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

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

626 return htc_commands 

627 

628 

629def _handle_job_outputs( 

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

631) -> dict[str, str]: 

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

633 

634 Parameters 

635 ---------- 

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

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

638 job_name : `str` 

639 Unique name for the job. 

640 use_shared : `bool` 

641 Whether job has access to files via shared filesystem. 

642 out_prefix : `str` 

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

644 

645 Returns 

646 ------- 

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

648 HTCondor commands for the job submission script. 

649 """ 

650 outputs = [] 

651 output_remaps = [] 

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

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

654 

655 uri = Path(gwf_file.src_uri) 

656 if not use_shared: 

657 outputs.append(uri.name) 

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

659 

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

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

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

663 # by the job. 

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

665 if outputs: 

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

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

668 

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

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

671 return htc_commands 

672 

673 

674def _create_periodic_release_expr( 

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

676) -> str: 

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

678 

679 Parameters 

680 ---------- 

681 memory : `int` 

682 Requested memory in MB. 

683 multiplier : `float` or None 

684 Memory growth rate between retries. 

685 limit : `int` 

686 Memory limit. 

687 additional_expr : `str`, optional 

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

689 

690 Returns 

691 ------- 

692 expr : `str` 

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

694 """ 

695 _LOG.debug( 

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

697 memory, 

698 multiplier, 

699 limit, 

700 additional_expr, 

701 ) 

702 

703 # ctrl_bps sets multiplier to None in the GenericWorkflow if 

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

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

706 return "" 

707 

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

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

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

711 # evaluate to FALSE in this case. 

712 # 

713 # Note: 

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

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

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

717 # but better safe than sorry. 

718 is_held = "JobStatus == 5" 

719 is_retry_allowed = "NumJobStarts <= JobMaxRetries" 

720 

721 mem_expr = "" 

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

723 was_mem_exceeded = ( 

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

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

726 ) 

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

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

729 

730 user_expr = "" 

731 if additional_expr: 

732 # Never auto release a job held by user. 

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

734 

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

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

737 transfer_expr = "HoldReasonCode =?= 12" 

738 

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

740 if user_expr and mem_expr: 

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

742 elif user_expr: 

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

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

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

746 

747 return expr 

748 

749 

750def _create_periodic_remove_expr(memory, multiplier, limit): 

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

752 

753 Parameters 

754 ---------- 

755 memory : `int` 

756 Requested memory in MB. 

757 multiplier : `float` 

758 Memory growth rate between retries. 

759 limit : `int` 

760 Memory limit. 

761 

762 Returns 

763 ------- 

764 expr : `str` 

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

766 """ 

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

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

769 # The special comparison operators ensure that all comparisons below 

770 # will evaluate to FALSE in this case. 

771 # 

772 # Note: 

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

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

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

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

777 is_held = "JobStatus == 5" 

778 is_retry_disallowed = "NumJobStarts > JobMaxRetries" 

779 

780 mem_expr = "" 

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

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

783 

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

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

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

787 ) 

788 

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

790 return expr 

791 

792 

793def _create_request_memory_expr(memory, multiplier, limit): 

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

795 

796 Parameters 

797 ---------- 

798 memory : `int` 

799 Requested memory in MB. 

800 multiplier : `float` 

801 Memory growth rate between retries. 

802 limit : `int` 

803 Memory limit. 

804 

805 Returns 

806 ------- 

807 expr : `str` 

808 A string representing an HTCondor ClassAd expression enabling safe 

809 memory scaling between job retries. 

810 """ 

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

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

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

814 # the ones describing job's current state. 

815 # 

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

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

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

819 was_mem_exceeded = ( 

820 "LastJobStatus =?= 5 " 

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

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

823 ) 

824 

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

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

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

828 # whichever is greater. 

829 expr = ( 

830 f"({was_mem_exceeded}) " 

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

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

833 ) 

834 return expr 

835 

836 

837def _gather_site_values(config, compute_site): 

838 """Gather values specific to given site. 

839 

840 Parameters 

841 ---------- 

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

843 BPS configuration that includes necessary submit/runtime 

844 information. 

845 compute_site : `str` 

846 Compute site name. 

847 

848 Returns 

849 ------- 

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

851 Values specific to the given site. 

852 """ 

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

854 search_opts = {} 

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

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

857 

858 # Determine the hard limit for the memory requirement. 

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

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

861 search_opts["default"] = DEFAULT_HTC_EXEC_PATT 

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

863 del search_opts["default"] 

864 

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

866 # by definition, they cannot have more memory than 

867 # the partitionable slot they are the part of. 

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

869 pool_info = condor_status(constraint=constraint) 

870 try: 

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

872 except ValueError: 

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

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

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

876 site_values["memoryLimit"] = limit 

877 

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

879 

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

881 if searchobj: 

882 search_opts["searchobj"] = searchobj 

883 search_opts["replaceVars"] = True 

884 for key in searchobj: 

885 if key.startswith("+"): 

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

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

888 else: 

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

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

891 

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

893 if searchobj: 

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

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

896 site_values[key] = value 

897 

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

899 return site_values 

900 

901 

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

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

904 

905 Parameters 

906 ---------- 

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

908 BPS configuration that includes necessary submit/runtime 

909 information. 

910 label : `str` 

911 GenericWorkflowJob label. 

912 

913 Returns 

914 ------- 

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

916 Values specific to the given job label. 

917 """ 

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

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

920 

921 search_opts = config.get_search_opts(label) 

922 

923 # Determine the hard limit for the memory requirement. 

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

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

926 search_opts["default"] = DEFAULT_HTC_EXEC_PATT 

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

928 del search_opts["default"] 

929 

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

931 # by definition, they cannot have more memory than 

932 # the partitionable slot they are the part of. 

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

934 pool_info = condor_status(constraint=constraint) 

935 try: 

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

937 except ValueError: 

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

939 

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

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

942 values["memoryLimit"] = limit 

943 

944 values["bpsUseShared"] = False 

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

946 if found: 

947 values["bpsUseShared"] = value 

948 

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

950 if found: 

951 values["releaseExpr"] = value 

952 

953 values["overwriteJobFiles"] = True 

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

955 if found: 

956 values["overwriteJobFiles"] = value 

957 

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

959 if found: 

960 values["releaseExpr"] = value 

961 

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

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

964 if found and not value: 

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

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

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

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

969 values["payloadCommand"] = value 

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

971 

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

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

974 

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

976 if found: 

977 values["nodeset"] = value 

978 

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

980 if found and "condor" in profile_sect: 

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

982 if subkey.startswith("+"): 

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

984 else: 

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

986 

987 # Copy all of the site values. 

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

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

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

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

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

993 values[key] = value 

994 

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

996 

997 return values 

998 

999 

1000def _group_to_subdag( 

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

1002) -> HTCJob: 

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

1004 

1005 Parameters 

1006 ---------- 

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

1008 Workflow configuration. 

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

1010 The generic workflow group to convert. 

1011 out_prefix : `str` 

1012 Location prefix to be used when creating jobs. 

1013 

1014 Returns 

1015 ------- 

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

1017 Job for running the HTCondor dag. 

1018 """ 

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

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

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

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

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

1024 htc_job.dagcmds["post"] = { 

1025 "defer": "", 

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

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

1028 } 

1029 return htc_job 

1030 

1031 

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

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

1034 

1035 Parameters 

1036 ---------- 

1037 group_job_name : `str` 

1038 Name of the group job. 

1039 job_label : `str` 

1040 Label to use for the check status job. 

1041 site_values : `dict` 

1042 Site specific values. 

1043 

1044 Returns 

1045 ------- 

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

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

1048 """ 

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

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

1051 # ADD nodeset to VARS 

1052 job_vars = {"group_job_name": group_job_name} 

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

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

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

1056 

1057 return htc_job 

1058 

1059 

1060def _generic_workflow_to_htcondor_dag( 

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

1062) -> HTCDag: 

1063 """Convert a GenericWorkflow to a HTCDag. 

1064 

1065 Parameters 

1066 ---------- 

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

1068 Workflow configuration. 

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

1070 The GenericWorkflow to convert. 

1071 out_prefix : `str` 

1072 Location prefix where the HTCondor files will be written. 

1073 

1074 Returns 

1075 ------- 

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

1077 The HTCDag representation of the given GenericWorkflow. 

1078 """ 

1079 dag = HTCDag(name=generic_workflow.name) 

1080 

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

1082 dag.add_attribs(generic_workflow.run_attrs) 

1083 dag.add_attribs( 

1084 { 

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

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

1087 } 

1088 ) 

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

1090 

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

1092 dag.graph["write_dot"] = save_htc_dot 

1093 

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

1095 subdir_template = defaultdict(lambda: tmp_template) 

1096 else: 

1097 subdir_template = tmp_template 

1098 

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

1100 lazy_groups = [] 

1101 

1102 # Create all DAG jobs 

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

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

1105 for job_name in generic_workflow: 

1106 gwjob = generic_workflow.get_job(job_name) 

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

1108 GenericWorkflowNodeType.PAYLOAD, 

1109 GenericWorkflowNodeType.LAZY_GROUP, 

1110 ]: 

1111 gwjob = cast(GenericWorkflowJob, gwjob) 

1112 if gwjob.label not in cached_values: 

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

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

1115 htc_job = _create_job( 

1116 subdir_template[gwjob.label], 

1117 cached_values[gwjob.label], 

1118 generic_workflow, 

1119 gwjob, 

1120 out_prefix, 

1121 ) 

1122 elif gwjob.node_type == GenericWorkflowNodeType.NOOP: 

1123 gwjob = cast(GenericWorkflowNoopJob, gwjob) 

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

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

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

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

1128 elif gwjob.node_type == GenericWorkflowNodeType.GROUP: 

1129 gwjob = cast(GenericWorkflowGroup, gwjob) 

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

1131 htc_job = _group_to_subdag(config, gwjob, out_prefix) 

1132 else: 

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

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

1135 dag.add_job(htc_job) 

1136 

1137 # Have to add the placeholder job for the workflow 

1138 if gwjob.node_type == GenericWorkflowNodeType.LAZY_GROUP: 

1139 lazy_groups.append(gwjob.name) 

1140 

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

1142 for job_name in generic_workflow: 

1143 gwjob = generic_workflow.get_job(job_name) 

1144 parent_name = ( 

1145 gwjob.name 

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

1147 else f"wms_{gwjob.name}" 

1148 ) 

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

1150 children_names = [] 

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

1152 gwjob = cast(GenericWorkflowGroup, gwjob) 

1153 group_children = [] # Dependencies between same group jobs 

1154 for sjob in successor_jobs: 

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

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

1157 elif sjob.node_type == GenericWorkflowNodeType.PAYLOAD: 

1158 children_names.append(sjob.name) 

1159 else: 

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

1161 if group_children: 

1162 dag.add_job_relationships([parent_name], group_children) 

1163 if not gwjob.blocking: 

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

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

1166 check_job = _create_check_job( 

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

1168 ) 

1169 dag.add_job(check_job) 

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

1171 parent_name = check_job.name 

1172 else: 

1173 for sjob in successor_jobs: 

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

1175 GenericWorkflowNodeType.PAYLOAD, 

1176 GenericWorkflowNodeType.LAZY_GROUP, 

1177 ]: 

1178 children_names.append(sjob.name) 

1179 else: 

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

1181 

1182 dag.add_job_relationships([parent_name], children_names) 

1183 

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

1185 for lazy_group_name in lazy_groups: 

1186 _add_lazy_placeholder(lazy_group_name, generic_workflow, dag, out_prefix) 

1187 

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

1189 final = generic_workflow.get_final() 

1190 if final and isinstance(final, GenericWorkflowJob): 

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

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

1193 final_htjob = _create_job( 

1194 subdir_template[final.label], 

1195 cached_values[final.label], 

1196 generic_workflow, 

1197 final, 

1198 out_prefix, 

1199 ) 

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

1201 final_htjob.dagcmds["post"] = { 

1202 "defer": "", 

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

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

1205 } 

1206 dag.add_final_job(final_htjob) 

1207 elif final and isinstance(final, GenericWorkflow): 

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

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

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

1211 

1212 return dag 

1213 

1214 

1215def _add_lazy_placeholder( 

1216 prepare_job_name: str, 

1217 generic_workflow: GenericWorkflow, 

1218 dag: HTCDag, 

1219 out_prefix: str, 

1220): 

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

1222 

1223 # Make a fake job for the placeholder workflow 

1224 job = HTCJob("placeholder") 

1225 job.add_job_cmds( 

1226 { 

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

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

1229 } 

1230 ) 

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

1232 job.add_job_attrs(generic_workflow.run_attrs) 

1233 

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

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

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

1237 placeholder_dag = HTCDag(name=placeholder_dag_name) 

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

1239 placeholder_dag.add_attribs(generic_workflow.run_attrs) 

1240 placeholder_dag.add_job(job) 

1241 

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

1243 # to put jobs in lazy dag in bps_job_label 

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

1245 

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

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

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

1249 dag_job.subdag = placeholder_dag 

1250 

1251 dag.add_job(dag_job) 

1252 

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

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

1255 for job in successor_jobs: 

1256 dag.add_edge(dag_job.name, job) 

1257 dag.remove_edge(prepare_job_name, job) 

1258 

1259 dag.add_edge(prepare_job_name, dag_job.name) 

1260 

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

1262 

1263 

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

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

1266 summary. 

1267 

1268 Parameters 

1269 ---------- 

1270 subworkflow_name : `str` 

1271 Name of subworkflow. 

1272 subworkflow_summary : `str` 

1273 Job summary for the subworkflow. 

1274 submit_path : `str` 

1275 Directory in which to find the DAG info file. 

1276 """ 

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

1278 filename, dag_info = read_dag_info(submit_path) 

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

1280 

1281 schedd_name = next(iter(dag_info)) 

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

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

1284 

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

1286 generator_name = None 

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

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

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

1290 info = part.split(":") 

1291 if info[0] == subworkflow_name: 

1292 generator_name = info[1] 

1293 break 

1294 

1295 # Update bps_job_summary 

1296 _LOG.debug( 

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

1298 subworkflow_name, 

1299 dag_values["bps_job_summary"], 

1300 subworkflow_summary, 

1301 generator_name, 

1302 ) 

1303 if generator_name: 

1304 generator_summary = f"{generator_name}:1" 

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

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

1307 ) 

1308 else: 

1309 # just append to end of bps_job_summary 

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

1311 

1312 _LOG.debug( 

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

1314 subworkflow_name, 

1315 dag_values["bps_job_summary"], 

1316 ) 

1317 

1318 # Save updated bps_job_summary 

1319 write_dag_info(filename, dag_info)