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

505 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-26 09:28 +0000

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 for name, value in gwjob.environment.items(): 

313 if isinstance(value, str): 313 ↛ 317line 313 didn't jump to line 317 because the condition on line 313 was always true

314 value = _replace_wms_vars(value) 

315 value = _fix_env_var_syntax_shell(value) 

316 value = htc_escape(value) 

317 if use_htc_env: 317 ↛ 320line 317 didn't jump to line 320 because the condition on line 317 was always true

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

319 else: 

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

321 

322 # Process above added one trailing space 

323 if use_htc_env: 323 ↛ 327line 323 didn't jump to line 327 because the condition on line 323 was always true

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

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

326 

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

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

329 # a way forward to centralize logic in bps. 

330 

331 jobcmds["getenv"] = "True" 

332 

333 if gwjob.executable.transfer_executable: 

334 jobcmds["transfer_executable"] = "True" 

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

336 else: 

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

338 

339 if gwjob.arguments: 

340 arguments = gwjob.arguments 

341 arguments = _replace_cmd_vars(arguments, gwjob) 

342 arguments = _replace_wms_vars(arguments) 

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

344 arguments = _fix_env_var_syntax(arguments) 

345 jobcmds["arguments"] = arguments 

346 

347 else: 

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

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

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

351 

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

353 # part of the payloadCommand. 

354 

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

356 arguments = gwjob.arguments 

357 arguments = _replace_cmd_vars(arguments, gwjob) 

358 arguments = _replace_wms_vars(arguments) 

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

360 arguments = _fix_env_var_syntax_shell(arguments) 

361 

362 if gwjob.executable.transfer_executable: 

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

364 # file transfer list. 

365 gwfile = GenericWorkflowFile( 

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

367 ) 

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

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

370 # Ensure the executable copy is executable. 

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

372 else: 

373 exec_name = _fix_env_var_syntax_shell(gwjob.executable.src_uri) 

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

375 

376 payloadCommand = cached_vals["payloadCommand"] 

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

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

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

380 

381 # Remove newlines 

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

383 

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

385 

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

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

388 

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

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

391 jobcmds["transfer_executable"] = "False" 

392 

393 return jobcmds 

394 

395 

396def _translate_dag_cmds(gwjob): 

397 """Translate job values into DAGMan commands. 

398 

399 Parameters 

400 ---------- 

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

402 Job containing values to be translated. 

403 

404 Returns 

405 ------- 

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

407 DAGMan commands for the job. 

408 """ 

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

410 dag_translation = { 

411 "abort_on_value": "abort_dag_on", 

412 "abort_return_value": "abort_exit", 

413 "priority": "priority", 

414 } 

415 

416 dagcmds = {} 

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

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

419 

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

421 return dagcmds 

422 

423 

424def _fix_env_var_syntax(oldstr): 

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

426 

427 Parameters 

428 ---------- 

429 oldstr : `str` 

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

431 

432 Returns 

433 ------- 

434 newstr : `str` 

435 Given string with environment variable syntax fixed. 

436 """ 

437 newstr = oldstr 

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

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

440 return newstr 

441 

442 

443def _fix_env_var_syntax_shell(oldstr): 

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

445 

446 Parameters 

447 ---------- 

448 oldstr : `str` 

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

450 

451 Returns 

452 ------- 

453 newstr : `str` 

454 Given string with environment variable syntax fixed. 

455 """ 

456 newstr = oldstr 

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

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

459 return newstr 

460 

461 

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

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

464 physical file names. 

465 

466 Parameters 

467 ---------- 

468 use_shared : `bool` 

469 Whether HTCondor can assume shared filesystem. 

470 arguments : `str` 

471 Arguments string in which to replace file placeholders. 

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

473 Generic workflow that contains file information. 

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

475 The job corresponding to the arguments. 

476 

477 Returns 

478 ------- 

479 arguments : `str` 

480 Given arguments string with file placeholders replaced. 

481 """ 

482 # Replace input file placeholders with paths. 

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

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

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

486 # responsible for transferring file. 

487 uri = gwfile.src_uri 

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

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

490 # Have shared filesystems and jobs can share file. 

491 uri = gwfile.src_uri 

492 else: 

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

494 else: # Using push transfer 

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

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

497 

498 # Replace output file placeholders with paths. 

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

500 if not gwfile.wms_transfer: 

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

502 # responsible for transferring file. 

503 uri = gwfile.src_uri 

504 elif use_shared: 

505 if gwfile.job_shared: 

506 # Have shared filesystems and jobs can share file. 

507 uri = gwfile.src_uri 

508 else: 

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

510 else: # Using push transfer 

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

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

513 return arguments 

514 

515 

516def _replace_cmd_vars(arguments, gwjob): 

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

518 

519 Parameters 

520 ---------- 

521 arguments : `str` 

522 Arguments string in which to replace placeholders. 

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

524 Job containing values to be used to replace placeholders 

525 (in particular gwjob.cmdvals). 

526 

527 Returns 

528 ------- 

529 arguments : `str` 

530 Given arguments string with placeholders replaced. 

531 """ 

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

533 try: 

534 arguments = arguments.format(**replacements) 

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

536 _LOG.error( 

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

538 gwjob.name, 

539 str(exc), 

540 ) 

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

542 raise 

543 return arguments 

544 

545 

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

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

548 

549 Parameters 

550 ---------- 

551 orig_string : `str` 

552 String in which to replace wms placeholders. 

553 

554 Returns 

555 ------- 

556 updated_string : `str` 

557 Given string with wms placeholders replaced. 

558 """ 

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

560 updated_string = orig_string 

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

562 try: 

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

564 except KeyError: 

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

566 raise 

567 return updated_string 

568 

569 

570def _handle_job_inputs( 

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

572) -> dict[str, str]: 

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

574 

575 Parameters 

576 ---------- 

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

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

579 job_name : `str` 

580 Unique name for the job. 

581 use_shared : `bool` 

582 Whether job has access to files via shared filesystem. 

583 out_prefix : `str` 

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

585 

586 Returns 

587 ------- 

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

589 HTCondor commands for the job submission script. 

590 """ 

591 inputs = [] 

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

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

594 

595 uri = Path(gwf_file.src_uri) 

596 

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

598 

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

600 inputs.append(str(uri)) 

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

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

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

604 if uri.is_dir(): 

605 raise RuntimeError( 

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

607 ) 

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

609 

610 htc_commands = {} 

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

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

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

614 return htc_commands 

615 

616 

617def _handle_job_outputs( 

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

619) -> dict[str, str]: 

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

621 

622 Parameters 

623 ---------- 

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

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

626 job_name : `str` 

627 Unique name for the job. 

628 use_shared : `bool` 

629 Whether job has access to files via shared filesystem. 

630 out_prefix : `str` 

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

632 

633 Returns 

634 ------- 

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

636 HTCondor commands for the job submission script. 

637 """ 

638 outputs = [] 

639 output_remaps = [] 

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

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

642 

643 uri = Path(gwf_file.src_uri) 

644 if not use_shared: 

645 outputs.append(uri.name) 

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

647 

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

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

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

651 # by the job. 

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

653 if outputs: 

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

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

656 

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

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

659 return htc_commands 

660 

661 

662def _create_periodic_release_expr( 

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

664) -> str: 

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

666 

667 Parameters 

668 ---------- 

669 memory : `int` 

670 Requested memory in MB. 

671 multiplier : `float` or None 

672 Memory growth rate between retries. 

673 limit : `int` 

674 Memory limit. 

675 additional_expr : `str`, optional 

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

677 

678 Returns 

679 ------- 

680 expr : `str` 

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

682 """ 

683 _LOG.debug( 

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

685 memory, 

686 multiplier, 

687 limit, 

688 additional_expr, 

689 ) 

690 

691 # ctrl_bps sets multiplier to None in the GenericWorkflow if 

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

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

694 return "" 

695 

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

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

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

699 # evaluate to FALSE in this case. 

700 # 

701 # Note: 

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

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

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

705 # but better safe than sorry. 

706 is_held = "JobStatus == 5" 

707 is_retry_allowed = "NumJobStarts <= JobMaxRetries" 

708 

709 mem_expr = "" 

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

711 was_mem_exceeded = ( 

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

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

714 ) 

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

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

717 

718 user_expr = "" 

719 if additional_expr: 

720 # Never auto release a job held by user. 

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

722 

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

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

725 transfer_expr = "HoldReasonCode =?= 12" 

726 

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

728 if user_expr and mem_expr: 

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

730 elif user_expr: 

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

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

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

734 

735 return expr 

736 

737 

738def _create_periodic_remove_expr(memory, multiplier, limit): 

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

740 

741 Parameters 

742 ---------- 

743 memory : `int` 

744 Requested memory in MB. 

745 multiplier : `float` 

746 Memory growth rate between retries. 

747 limit : `int` 

748 Memory limit. 

749 

750 Returns 

751 ------- 

752 expr : `str` 

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

754 """ 

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

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

757 # The special comparison operators ensure that all comparisons below 

758 # will evaluate to FALSE in this case. 

759 # 

760 # Note: 

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

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

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

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

765 is_held = "JobStatus == 5" 

766 is_retry_disallowed = "NumJobStarts > JobMaxRetries" 

767 

768 mem_expr = "" 

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

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

771 

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

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

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

775 ) 

776 

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

778 return expr 

779 

780 

781def _create_request_memory_expr(memory, multiplier, limit): 

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

783 

784 Parameters 

785 ---------- 

786 memory : `int` 

787 Requested memory in MB. 

788 multiplier : `float` 

789 Memory growth rate between retries. 

790 limit : `int` 

791 Memory limit. 

792 

793 Returns 

794 ------- 

795 expr : `str` 

796 A string representing an HTCondor ClassAd expression enabling safe 

797 memory scaling between job retries. 

798 """ 

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

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

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

802 # the ones describing job's current state. 

803 # 

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

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

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

807 was_mem_exceeded = ( 

808 "LastJobStatus =?= 5 " 

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

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

811 ) 

812 

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

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

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

816 # whichever is greater. 

817 expr = ( 

818 f"({was_mem_exceeded}) " 

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

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

821 ) 

822 return expr 

823 

824 

825def _gather_site_values(config, compute_site): 

826 """Gather values specific to given site. 

827 

828 Parameters 

829 ---------- 

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

831 BPS configuration that includes necessary submit/runtime 

832 information. 

833 compute_site : `str` 

834 Compute site name. 

835 

836 Returns 

837 ------- 

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

839 Values specific to the given site. 

840 """ 

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

842 search_opts = {} 

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

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

845 

846 # Determine the hard limit for the memory requirement. 

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

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

849 search_opts["default"] = DEFAULT_HTC_EXEC_PATT 

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

851 del search_opts["default"] 

852 

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

854 # by definition, they cannot have more memory than 

855 # the partitionable slot they are the part of. 

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

857 pool_info = condor_status(constraint=constraint) 

858 try: 

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

860 except ValueError: 

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

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

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

864 site_values["memoryLimit"] = limit 

865 

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

867 

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

869 if searchobj: 

870 search_opts["searchobj"] = searchobj 

871 search_opts["replaceVars"] = True 

872 for key in searchobj: 

873 if key.startswith("+"): 

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

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

876 else: 

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

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

879 

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

881 if searchobj: 

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

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

884 site_values[key] = value 

885 

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

887 return site_values 

888 

889 

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

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

892 

893 Parameters 

894 ---------- 

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

896 BPS configuration that includes necessary submit/runtime 

897 information. 

898 label : `str` 

899 GenericWorkflowJob label. 

900 

901 Returns 

902 ------- 

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

904 Values specific to the given job label. 

905 """ 

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

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

908 

909 search_opts = config.get_search_opts(label) 

910 

911 # Determine the hard limit for the memory requirement. 

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

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

914 search_opts["default"] = DEFAULT_HTC_EXEC_PATT 

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

916 del search_opts["default"] 

917 

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

919 # by definition, they cannot have more memory than 

920 # the partitionable slot they are the part of. 

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

922 pool_info = condor_status(constraint=constraint) 

923 try: 

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

925 except ValueError: 

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

927 

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

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

930 values["memoryLimit"] = limit 

931 

932 values["bpsUseShared"] = False 

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

934 if found: 

935 values["bpsUseShared"] = value 

936 

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

938 if found: 

939 values["releaseExpr"] = value 

940 

941 values["overwriteJobFiles"] = True 

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

943 if found: 

944 values["overwriteJobFiles"] = value 

945 

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

947 if found: 

948 values["releaseExpr"] = value 

949 

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

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

952 if found and not value: 

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

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

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

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

957 values["payloadCommand"] = value 

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

959 

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

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

962 

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

964 if found: 

965 values["nodeset"] = value 

966 

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

968 if found and "condor" in profile_sect: 

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

970 if subkey.startswith("+"): 

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

972 else: 

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

974 

975 # Copy all of the site values. 

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

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

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

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

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

981 values[key] = value 

982 

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

984 

985 return values 

986 

987 

988def _group_to_subdag( 

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

990) -> HTCJob: 

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

992 

993 Parameters 

994 ---------- 

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

996 Workflow configuration. 

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

998 The generic workflow group to convert. 

999 out_prefix : `str` 

1000 Location prefix to be used when creating jobs. 

1001 

1002 Returns 

1003 ------- 

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

1005 Job for running the HTCondor dag. 

1006 """ 

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

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

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

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

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

1012 htc_job.dagcmds["post"] = { 

1013 "defer": "", 

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

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

1016 } 

1017 return htc_job 

1018 

1019 

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

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

1022 

1023 Parameters 

1024 ---------- 

1025 group_job_name : `str` 

1026 Name of the group job. 

1027 job_label : `str` 

1028 Label to use for the check status job. 

1029 site_values : `dict` 

1030 Site specific values. 

1031 

1032 Returns 

1033 ------- 

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

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

1036 """ 

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

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

1039 # ADD nodeset to VARS 

1040 job_vars = {"group_job_name": group_job_name} 

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

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

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

1044 

1045 return htc_job 

1046 

1047 

1048def _generic_workflow_to_htcondor_dag( 

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

1050) -> HTCDag: 

1051 """Convert a GenericWorkflow to a HTCDag. 

1052 

1053 Parameters 

1054 ---------- 

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

1056 Workflow configuration. 

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

1058 The GenericWorkflow to convert. 

1059 out_prefix : `str` 

1060 Location prefix where the HTCondor files will be written. 

1061 

1062 Returns 

1063 ------- 

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

1065 The HTCDag representation of the given GenericWorkflow. 

1066 """ 

1067 dag = HTCDag(name=generic_workflow.name) 

1068 

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

1070 dag.add_attribs(generic_workflow.run_attrs) 

1071 dag.add_attribs( 

1072 { 

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

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

1075 } 

1076 ) 

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

1078 

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

1080 dag.graph["write_dot"] = save_htc_dot 

1081 

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

1083 subdir_template = defaultdict(lambda: tmp_template) 

1084 else: 

1085 subdir_template = tmp_template 

1086 

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

1088 lazy_groups = [] 

1089 

1090 # Create all DAG jobs 

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

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

1093 for job_name in generic_workflow: 

1094 gwjob = generic_workflow.get_job(job_name) 

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

1096 GenericWorkflowNodeType.PAYLOAD, 

1097 GenericWorkflowNodeType.LAZY_GROUP, 

1098 ]: 

1099 gwjob = cast(GenericWorkflowJob, gwjob) 

1100 if gwjob.label not in cached_values: 

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

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

1103 htc_job = _create_job( 

1104 subdir_template[gwjob.label], 

1105 cached_values[gwjob.label], 

1106 generic_workflow, 

1107 gwjob, 

1108 out_prefix, 

1109 ) 

1110 elif gwjob.node_type == GenericWorkflowNodeType.NOOP: 

1111 gwjob = cast(GenericWorkflowNoopJob, gwjob) 

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

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

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

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

1116 elif gwjob.node_type == GenericWorkflowNodeType.GROUP: 

1117 gwjob = cast(GenericWorkflowGroup, gwjob) 

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

1119 htc_job = _group_to_subdag(config, gwjob, out_prefix) 

1120 else: 

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

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

1123 dag.add_job(htc_job) 

1124 

1125 # Have to add the placeholder job for the workflow 

1126 if gwjob.node_type == GenericWorkflowNodeType.LAZY_GROUP: 

1127 lazy_groups.append(gwjob.name) 

1128 

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

1130 for job_name in generic_workflow: 

1131 gwjob = generic_workflow.get_job(job_name) 

1132 parent_name = ( 

1133 gwjob.name 

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

1135 else f"wms_{gwjob.name}" 

1136 ) 

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

1138 children_names = [] 

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

1140 gwjob = cast(GenericWorkflowGroup, gwjob) 

1141 group_children = [] # Dependencies between same group jobs 

1142 for sjob in successor_jobs: 

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

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

1145 elif sjob.node_type == GenericWorkflowNodeType.PAYLOAD: 

1146 children_names.append(sjob.name) 

1147 else: 

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

1149 if group_children: 

1150 dag.add_job_relationships([parent_name], group_children) 

1151 if not gwjob.blocking: 

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

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

1154 check_job = _create_check_job( 

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

1156 ) 

1157 dag.add_job(check_job) 

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

1159 parent_name = check_job.name 

1160 else: 

1161 for sjob in successor_jobs: 

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

1163 GenericWorkflowNodeType.PAYLOAD, 

1164 GenericWorkflowNodeType.LAZY_GROUP, 

1165 ]: 

1166 children_names.append(sjob.name) 

1167 else: 

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

1169 

1170 dag.add_job_relationships([parent_name], children_names) 

1171 

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

1173 for lazy_group_name in lazy_groups: 

1174 _add_lazy_placeholder(lazy_group_name, generic_workflow, dag, out_prefix) 

1175 

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

1177 final = generic_workflow.get_final() 

1178 if final and isinstance(final, GenericWorkflowJob): 

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

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

1181 final_htjob = _create_job( 

1182 subdir_template[final.label], 

1183 cached_values[final.label], 

1184 generic_workflow, 

1185 final, 

1186 out_prefix, 

1187 ) 

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

1189 final_htjob.dagcmds["post"] = { 

1190 "defer": "", 

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

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

1193 } 

1194 dag.add_final_job(final_htjob) 

1195 elif final and isinstance(final, GenericWorkflow): 

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

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

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

1199 

1200 return dag 

1201 

1202 

1203def _add_lazy_placeholder( 

1204 prepare_job_name: str, 

1205 generic_workflow: GenericWorkflow, 

1206 dag: HTCDag, 

1207 out_prefix: str, 

1208): 

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

1210 

1211 # Make a fake job for the placeholder workflow 

1212 job = HTCJob("placeholder") 

1213 job.add_job_cmds( 

1214 { 

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

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

1217 } 

1218 ) 

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

1220 job.add_job_attrs(generic_workflow.run_attrs) 

1221 

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

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

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

1225 placeholder_dag = HTCDag(name=placeholder_dag_name) 

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

1227 placeholder_dag.add_attribs(generic_workflow.run_attrs) 

1228 placeholder_dag.add_job(job) 

1229 

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

1231 # to put jobs in lazy dag in bps_job_label 

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

1233 

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

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

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

1237 dag_job.subdag = placeholder_dag 

1238 

1239 dag.add_job(dag_job) 

1240 

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

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

1243 for job in successor_jobs: 

1244 dag.add_edge(dag_job.name, job) 

1245 dag.remove_edge(prepare_job_name, job) 

1246 

1247 dag.add_edge(prepare_job_name, dag_job.name) 

1248 

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

1250 

1251 

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

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

1254 summary. 

1255 

1256 Parameters 

1257 ---------- 

1258 subworkflow_name : `str` 

1259 Name of subworkflow. 

1260 subworkflow_summary : `str` 

1261 Job summary for the subworkflow. 

1262 submit_path : `str` 

1263 Directory in which to find the DAG info file. 

1264 """ 

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

1266 filename, dag_info = read_dag_info(submit_path) 

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

1268 

1269 schedd_name = next(iter(dag_info)) 

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

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

1272 

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

1274 generator_name = None 

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

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

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

1278 info = part.split(":") 

1279 if info[0] == subworkflow_name: 

1280 generator_name = info[1] 

1281 break 

1282 

1283 # Update bps_job_summary 

1284 _LOG.debug( 

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

1286 subworkflow_name, 

1287 dag_values["bps_job_summary"], 

1288 subworkflow_summary, 

1289 generator_name, 

1290 ) 

1291 if generator_name: 

1292 generator_summary = f"{generator_name}:1" 

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

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

1295 ) 

1296 else: 

1297 # just append to end of bps_job_summary 

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

1299 

1300 _LOG.debug( 

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

1302 subworkflow_name, 

1303 dag_values["bps_job_summary"], 

1304 ) 

1305 

1306 # Save updated bps_job_summary 

1307 write_dag_info(filename, dag_info)