Coverage for python/lsst/ctrl/bps/htcondor/prepare_utils.py: 82%
505 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-08-29 09:20 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-08-29 09:20 +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 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}';"
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"])
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.
331 jobcmds["getenv"] = "True"
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)
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
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.
352 # Don't set getenv as setting up the environment is assumed to be
353 # part of the payloadCommand.
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)
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}"
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)
381 # Remove newlines
382 payloadCommand = re.sub("\n", "", payloadCommand)
384 _LOG.debug("%s payloadCommand post-format: %s", gwjob.label, payloadCommand)
386 # jobcmds["arguments"] = htc_escape(f"-c '{payloadCommand}'")
387 jobcmds["arguments"] = f"-c '{payloadCommand}'"
389 jobcmds["executable"] = "/bin/bash"
390 # Don't need to transfer /bin/bash
391 jobcmds["transfer_executable"] = "False"
393 return jobcmds
396def _translate_dag_cmds(gwjob):
397 """Translate job values into DAGMan commands.
399 Parameters
400 ----------
401 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
402 Job containing values to be translated.
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 }
416 dagcmds = {}
417 for gwkey, htckey in dag_translation.items():
418 dagcmds[htckey] = getattr(gwjob, gwkey, None)
420 # Still to be coded: vars "pre_cmdline", "post_cmdline"
421 return dagcmds
424def _fix_env_var_syntax(oldstr):
425 """Change ENV place holders to HTCondor Env var syntax.
427 Parameters
428 ----------
429 oldstr : `str`
430 String in which environment variable syntax is to be fixed.
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
443def _fix_env_var_syntax_shell(oldstr):
444 """Change ENV place holders to shell var syntax.
446 Parameters
447 ----------
448 oldstr : `str`
449 String in which environment variable syntax is to be fixed.
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
462def _replace_file_vars(use_shared, arguments, workflow, gwjob):
463 """Replace file placeholders in command line arguments with correct
464 physical file names.
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.
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)
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
516def _replace_cmd_vars(arguments, gwjob):
517 """Replace format-style placeholders in arguments.
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).
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
546def _replace_wms_vars(orig_string: str) -> str:
547 """Replace special wms placeholders in given string.
549 Parameters
550 ----------
551 orig_string : `str`
552 String in which to replace wms placeholders.
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
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.
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.
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)
595 uri = Path(gwf_file.src_uri)
597 # Note if use_shared and job_shared, don't need to transfer file.
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}")
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
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.
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.
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)
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)}")
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"])
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
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.
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.
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 )
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 ""
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"
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}"
718 user_expr = ""
719 if additional_expr:
720 # Never auto release a job held by user.
721 user_expr = f"HoldReasonCode =!= 1 && {additional_expr}"
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"
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})"
735 return expr
738def _create_periodic_remove_expr(memory, multiplier, limit):
739 """Construct an HTCondorAd expression for removing jobs from the queue.
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.
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"
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}"
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 )
777 expr = f"{is_held} && ({is_retry_disallowed}{mem_expr})"
778 return expr
781def _create_request_memory_expr(memory, multiplier, limit):
782 """Construct an HTCondor ClassAd expression for safe memory scaling.
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.
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 )
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
825def _gather_site_values(config, compute_site):
826 """Gather values specific to given site.
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.
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}
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"]
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
866 _, site_values["bpsUseShared"] = config.search("bpsUseShared", opt={"default": False})
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
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
886 _LOG.debug("site_values = %s", site_values)
887 return site_values
890def _gather_label_values(config: BpsConfig, label: str) -> dict[str, Any]:
891 """Gather values specific to given job label.
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.
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": {}}
909 search_opts = config.get_search_opts(label)
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"]
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)
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
932 values["bpsUseShared"] = False
933 found, value = config.search("bpsUseShared", opt=search_opts)
934 if found:
935 values["bpsUseShared"] = value
937 found, value = config.search("releaseExpr", opt=search_opts)
938 if found:
939 values["releaseExpr"] = value
941 values["overwriteJobFiles"] = True
942 found, value = config.search("overwriteJobFiles", opt=search_opts)
943 if found:
944 values["overwriteJobFiles"] = value
946 found, value = config.search("releaseExpr", opt=search_opts)
947 if found:
948 values["releaseExpr"] = value
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)
960 found, value = config.search("bpsUseHTCEnvironment", opt=search_opts)
961 values["bpsUseHTCEnvironment"] = value if found else values["bpsMakeCommand"]
963 found, value = config.search("nodeset", opt=search_opts)
964 if found:
965 values["nodeset"] = value
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
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
983 _LOG.debug("_gather_label_values: label = %s, values = %s", label, values)
985 return values
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.
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.
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
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.
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.
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})
1045 return htc_job
1048def _generic_workflow_to_htcondor_dag(
1049 config: BpsConfig, generic_workflow: GenericWorkflow, out_prefix: str
1050) -> HTCDag:
1051 """Convert a GenericWorkflow to a HTCDag.
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.
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)
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": ""})
1079 _, save_htc_dot = config.search("saveHTCdot", opt={"default": False})
1080 dag.graph["write_dot"] = save_htc_dot
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
1087 # Save list of lazy group jobs for later extra handling.
1088 lazy_groups = []
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)
1125 # Have to add the placeholder job for the workflow
1126 if gwjob.node_type == GenericWorkflowNodeType.LAZY_GROUP:
1127 lazy_groups.append(gwjob.name)
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}")
1170 dag.add_job_relationships([parent_name], children_names)
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)
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)})")
1200 return dag
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)
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)
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)
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}"})
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
1239 dag.add_job(dag_job)
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)
1247 dag.add_edge(prepare_job_name, dag_job.name)
1249 _LOG.debug("_add_lazy_placeholder: edges = %s", dag.edges)
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.
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)
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)
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
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}"
1300 _LOG.debug(
1301 "After replace, name = %s, bps_job_summary = %s",
1302 subworkflow_name,
1303 dag_values["bps_job_summary"],
1304 )
1306 # Save updated bps_job_summary
1307 write_dag_info(filename, dag_info)