Coverage for python/lsst/ctrl/bps/htcondor/prepare_utils.py: 82%
508 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-15 09:04 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-15 09:04 +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/>.
28"""Utility functions for preparing the HTCondor workflow."""
30import logging
31import os
32import re
33from collections import defaultdict
34from pathlib import Path
35from typing import Any, cast
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
48from .lssthtc import (
49 HTCDag,
50 HTCJob,
51 condor_status,
52 htc_escape,
53 read_dag_info,
54 write_dag_info,
55)
57_LOG = logging.getLogger(__name__)
59DEFAULT_HTC_EXEC_PATT = ".*worker.*"
60"""Default pattern for searching execute machines in an HTCondor pool.
61"""
64def _create_job(subdir_template, cached_values, generic_workflow, gwjob, out_prefix):
65 """Convert GenericWorkflow job nodes to DAG jobs.
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.
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)
88 curvals = defaultdict(str)
89 curvals["label"] = gwjob.label
90 if gwjob.tags:
91 curvals.update(gwjob.tags)
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})
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 }
117 htc_job_cmds.update(_translate_job_cmds(cached_values, generic_workflow, gwjob))
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])
127 key = "log"
128 htc_job_cmds[key] = f"{gwjob.name}.$(Cluster).{key}"
129 _LOG.debug("HTCondor %s = %s", key, htc_job_cmds[key])
131 htc_job_cmds.update(
132 _handle_job_inputs(generic_workflow, gwjob.name, cached_values["bpsUseShared"], out_prefix)
133 )
135 htc_job_cmds.update(
136 _handle_job_outputs(generic_workflow, gwjob.name, cached_values["bpsUseShared"], out_prefix)
137 )
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
148 # Add the job cmds dict to the job object.
149 htc_job.add_job_cmds(htc_job_cmds)
151 # Add job-related cmds to the DAG (e.g., VARS)
152 htc_job.add_dag_cmds(_translate_dag_cmds(gwjob))
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})
161 return htc_job
164def _translate_job_cmds(cached_vals, generic_workflow, gwjob):
165 """Translate the job data that are one to one mapping
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.
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 }
193 jobcmds = {}
194 for gwkey, htckey in job_translation.items():
195 jobcmds[htckey] = getattr(gwjob, gwkey, None)
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")
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.")
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"
217 if gwjob.request_memory:
218 jobcmds["request_memory"] = f"{gwjob.request_memory}"
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 )
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
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 )
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
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
262 jobcmds["periodic_remove"] = _create_periodic_remove_expr(
263 gwjob.request_memory, gwjob.memory_multiplier, memory_max
264 )
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
272 # Handle command line
273 new_job_cmds = _translate_command_line(cached_vals, generic_workflow, gwjob)
274 jobcmds.update(new_job_cmds)
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
283 return jobcmds
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.
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.
300 Returns
301 -------
302 jobcmds : `dict` [`str` `Any`]
303 Commands to add to HTC submit description.
304 """
305 jobcmds = {}
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}';"
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"])
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.
343 jobcmds["getenv"] = "True"
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)
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
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.
364 # Don't set getenv as setting up the environment is assumed to be
365 # part of the payloadCommand.
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)
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}"
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)
393 # Remove newlines
394 payloadCommand = re.sub("\n", "", payloadCommand)
396 _LOG.debug("%s payloadCommand post-format: %s", gwjob.label, payloadCommand)
398 # jobcmds["arguments"] = htc_escape(f"-c '{payloadCommand}'")
399 jobcmds["arguments"] = f"-c '{payloadCommand}'"
401 jobcmds["executable"] = "/bin/bash"
402 # Don't need to transfer /bin/bash
403 jobcmds["transfer_executable"] = "False"
405 return jobcmds
408def _translate_dag_cmds(gwjob):
409 """Translate job values into DAGMan commands.
411 Parameters
412 ----------
413 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
414 Job containing values to be translated.
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 }
428 dagcmds = {}
429 for gwkey, htckey in dag_translation.items():
430 dagcmds[htckey] = getattr(gwjob, gwkey, None)
432 # Still to be coded: vars "pre_cmdline", "post_cmdline"
433 return dagcmds
436def _fix_env_var_syntax(oldstr):
437 """Change ENV place holders to HTCondor Env var syntax.
439 Parameters
440 ----------
441 oldstr : `str`
442 String in which environment variable syntax is to be fixed.
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
455def _fix_env_var_syntax_shell(oldstr):
456 """Change ENV place holders to shell var syntax.
458 Parameters
459 ----------
460 oldstr : `str`
461 String in which environment variable syntax is to be fixed.
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
474def _replace_file_vars(use_shared, arguments, workflow, gwjob):
475 """Replace file placeholders in command line arguments with correct
476 physical file names.
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.
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)
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
528def _replace_cmd_vars(arguments, gwjob):
529 """Replace format-style placeholders in arguments.
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).
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
558def _replace_wms_vars(orig_string: str) -> str:
559 """Replace special wms placeholders in given string.
561 Parameters
562 ----------
563 orig_string : `str`
564 String in which to replace wms placeholders.
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
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.
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.
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)
607 uri = Path(gwf_file.src_uri)
609 # Note if use_shared and job_shared, don't need to transfer file.
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}")
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
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.
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.
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)
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)}")
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"])
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
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.
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.
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 )
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 ""
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"
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}"
730 user_expr = ""
731 if additional_expr:
732 # Never auto release a job held by user.
733 user_expr = f"HoldReasonCode =!= 1 && {additional_expr}"
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"
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})"
747 return expr
750def _create_periodic_remove_expr(memory, multiplier, limit):
751 """Construct an HTCondorAd expression for removing jobs from the queue.
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.
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"
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}"
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 )
789 expr = f"{is_held} && ({is_retry_disallowed}{mem_expr})"
790 return expr
793def _create_request_memory_expr(memory, multiplier, limit):
794 """Construct an HTCondor ClassAd expression for safe memory scaling.
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.
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 )
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
837def _gather_site_values(config, compute_site):
838 """Gather values specific to given site.
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.
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}
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"]
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
878 _, site_values["bpsUseShared"] = config.search("bpsUseShared", opt={"default": False})
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
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
898 _LOG.debug("site_values = %s", site_values)
899 return site_values
902def _gather_label_values(config: BpsConfig, label: str) -> dict[str, Any]:
903 """Gather values specific to given job label.
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.
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": {}}
921 search_opts = config.get_search_opts(label)
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"]
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)
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
944 values["bpsUseShared"] = False
945 found, value = config.search("bpsUseShared", opt=search_opts)
946 if found:
947 values["bpsUseShared"] = value
949 found, value = config.search("releaseExpr", opt=search_opts)
950 if found:
951 values["releaseExpr"] = value
953 values["overwriteJobFiles"] = True
954 found, value = config.search("overwriteJobFiles", opt=search_opts)
955 if found:
956 values["overwriteJobFiles"] = value
958 found, value = config.search("releaseExpr", opt=search_opts)
959 if found:
960 values["releaseExpr"] = value
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)
972 found, value = config.search("bpsUseHTCEnvironment", opt=search_opts)
973 values["bpsUseHTCEnvironment"] = value if found else values["bpsMakeCommand"]
975 found, value = config.search("nodeset", opt=search_opts)
976 if found:
977 values["nodeset"] = value
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
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
995 _LOG.debug("_gather_label_values: label = %s, values = %s", label, values)
997 return values
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.
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.
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
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.
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.
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})
1057 return htc_job
1060def _generic_workflow_to_htcondor_dag(
1061 config: BpsConfig, generic_workflow: GenericWorkflow, out_prefix: str
1062) -> HTCDag:
1063 """Convert a GenericWorkflow to a HTCDag.
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.
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)
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": ""})
1091 _, save_htc_dot = config.search("saveHTCdot", opt={"default": False})
1092 dag.graph["write_dot"] = save_htc_dot
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
1099 # Save list of lazy group jobs for later extra handling.
1100 lazy_groups = []
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)
1137 # Have to add the placeholder job for the workflow
1138 if gwjob.node_type == GenericWorkflowNodeType.LAZY_GROUP:
1139 lazy_groups.append(gwjob.name)
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}")
1182 dag.add_job_relationships([parent_name], children_names)
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)
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)})")
1212 return dag
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)
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)
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)
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}"})
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
1251 dag.add_job(dag_job)
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)
1259 dag.add_edge(prepare_job_name, dag_job.name)
1261 _LOG.debug("_add_lazy_placeholder: edges = %s", dag.edges)
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.
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)
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)
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
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}"
1312 _LOG.debug(
1313 "After replace, name = %s, bps_job_summary = %s",
1314 subworkflow_name,
1315 dag_values["bps_job_summary"],
1316 )
1318 # Save updated bps_job_summary
1319 write_dag_info(filename, dag_info)