Coverage for python/lsst/ctrl/bps/htcondor/prepare_utils.py: 82%
506 statements
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-23 10:18 +0000
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-23 10:18 +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):
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 arguments = ""
340 if gwjob.arguments:
341 arguments = gwjob.arguments
342 arguments = _replace_cmd_vars(arguments, gwjob)
343 arguments = _replace_wms_vars(arguments)
344 arguments = _replace_file_vars(cached_vals["bpsUseShared"], arguments, generic_workflow, gwjob)
346 if cached_vals.get("bpsMakeCommand", True):
347 # Way to have fallback to previous behavior as well as
348 # a way forward to centralize logic in bps.
350 jobcmds["getenv"] = "True"
352 if gwjob.executable.transfer_executable:
353 jobcmds["transfer_executable"] = "True"
354 jobcmds["executable"] = gwjob.executable.src_uri
355 else:
356 jobcmds["executable"] = _fix_env_var_syntax(gwjob.executable.src_uri)
358 if arguments:
359 arguments = _fix_env_var_syntax(arguments)
360 jobcmds["arguments"] = arguments
362 else:
363 # Instead of making a bash script, run /bin/bash -c <commands>
364 # HTCondor v25 has a job command called shell that can replace the
365 # /bin/bash when we get to that version.
367 # Don't set getenv as setting up the environment is assumed to be
368 # part of the payloadCommand.
370 if arguments:
371 arguments = _fix_env_var_syntax_shell(arguments)
373 if gwjob.executable.transfer_executable:
374 # Since replacing executable need to add this executable to the
375 # file transfer list.
376 gwfile = GenericWorkflowFile(
377 name=gwjob.executable.name, src_uri=gwjob.executable.src_uri, wms_transfer=True
378 )
379 generic_workflow.add_job_inputs(gwjob.name, [gwfile])
380 exec_name = os.path.basename(gwjob.executable.src_uri)
381 # Ensure the executable copy is executable.
382 gwjob_command = f"chmod u+x {exec_name}; ./{exec_name} {arguments}"
383 else:
384 exec_name = _fix_env_var_syntax_shell(gwjob.executable.src_uri)
385 gwjob_command = f"{exec_name} {arguments}"
387 payload_command = cached_vals["payloadCommand"]
388 _LOG.debug("%s payload_command pre-format: %s", gwjob.label, payload_command)
389 payload_command = re.sub("{gwjobCommand}", gwjob_command, payload_command)
390 payload_command = re.sub("{gwjobExports}", job_exports, payload_command)
392 # Remove newlines
393 payload_command = re.sub("\n", "", payload_command)
395 _LOG.debug("%s payload_command post-format: %s", gwjob.label, payload_command)
397 jobcmds["arguments"] = f"-c '{payload_command}'"
399 jobcmds["executable"] = "/bin/bash"
400 # Don't need to transfer /bin/bash
401 jobcmds["transfer_executable"] = "False"
403 return jobcmds
406def _translate_dag_cmds(gwjob):
407 """Translate job values into DAGMan commands.
409 Parameters
410 ----------
411 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
412 Job containing values to be translated.
414 Returns
415 -------
416 dagcmds : `dict` [`str`, `~typing.Any`]
417 DAGMan commands for the job.
418 """
419 # Values in the dag script that just are name mappings.
420 dag_translation = {
421 "abort_on_value": "abort_dag_on",
422 "abort_return_value": "abort_exit",
423 "priority": "priority",
424 }
426 dagcmds = {}
427 for gwkey, htckey in dag_translation.items():
428 dagcmds[htckey] = getattr(gwjob, gwkey, None)
430 # Still to be coded: vars "pre_cmdline", "post_cmdline"
431 return dagcmds
434def _fix_env_var_syntax(oldstr):
435 """Change ENV place holders to HTCondor Env var syntax.
437 Parameters
438 ----------
439 oldstr : `str`
440 String in which environment variable syntax is to be fixed.
442 Returns
443 -------
444 newstr : `str`
445 Given string with environment variable syntax fixed.
446 """
447 newstr = oldstr
448 for key in re.findall(r"<ENV:([^>]+)>", oldstr):
449 newstr = newstr.replace(rf"<ENV:{key}>", f"$ENV({key})")
450 return newstr
453def _fix_env_var_syntax_shell(oldstr):
454 """Change ENV place holders to shell var syntax.
456 Parameters
457 ----------
458 oldstr : `str`
459 String in which environment variable syntax is to be fixed.
461 Returns
462 -------
463 newstr : `str`
464 Given string with environment variable syntax fixed.
465 """
466 newstr = oldstr
467 for key in re.findall(r"<ENV:([^>]+)>", oldstr):
468 newstr = newstr.replace(rf"<ENV:{key}>", f"${{{key}}}")
469 return newstr
472def _replace_file_vars(use_shared, arguments, workflow, gwjob):
473 """Replace file placeholders in command line arguments with correct
474 physical file names.
476 Parameters
477 ----------
478 use_shared : `bool`
479 Whether HTCondor can assume shared filesystem.
480 arguments : `str`
481 Arguments string in which to replace file placeholders.
482 workflow : `lsst.ctrl.bps.GenericWorkflow`
483 Generic workflow that contains file information.
484 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
485 The job corresponding to the arguments.
487 Returns
488 -------
489 arguments : `str`
490 Given arguments string with file placeholders replaced.
491 """
492 # Replace input file placeholders with paths.
493 for gwfile in workflow.get_job_inputs(gwjob.name, data=True, transfer_only=False):
494 if not gwfile.wms_transfer: 494 ↛ 497line 494 didn't jump to line 497 because the condition on line 494 was never true
495 # Must assume full URI if in command line and told WMS is not
496 # responsible for transferring file.
497 uri = gwfile.src_uri
498 elif use_shared: 498 ↛ 505line 498 didn't jump to line 505 because the condition on line 498 was always true
499 if gwfile.job_shared: 499 ↛ 503line 499 didn't jump to line 503 because the condition on line 499 was always true
500 # Have shared filesystems and jobs can share file.
501 uri = gwfile.src_uri
502 else:
503 uri = os.path.basename(gwfile.src_uri)
504 else: # Using push transfer
505 uri = os.path.basename(gwfile.src_uri)
506 arguments = arguments.replace(f"<FILE:{gwfile.name}>", uri)
508 # Replace output file placeholders with paths.
509 for gwfile in workflow.get_job_outputs(gwjob.name, data=True, transfer_only=False): 509 ↛ 510line 509 didn't jump to line 510 because the loop on line 509 never started
510 if not gwfile.wms_transfer:
511 # Must assume full URI if in command line and told WMS is not
512 # responsible for transferring file.
513 uri = gwfile.src_uri
514 elif use_shared:
515 if gwfile.job_shared:
516 # Have shared filesystems and jobs can share file.
517 uri = gwfile.src_uri
518 else:
519 uri = os.path.basename(gwfile.src_uri)
520 else: # Using push transfer
521 uri = os.path.basename(gwfile.src_uri)
522 arguments = arguments.replace(f"<FILE:{gwfile.name}>", uri)
523 return arguments
526def _replace_cmd_vars(arguments, gwjob):
527 """Replace format-style placeholders in arguments.
529 Parameters
530 ----------
531 arguments : `str`
532 Arguments string in which to replace placeholders.
533 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
534 Job containing values to be used to replace placeholders
535 (in particular gwjob.cmdvals).
537 Returns
538 -------
539 arguments : `str`
540 Given arguments string with placeholders replaced.
541 """
542 replacements = gwjob.cmdvals if gwjob.cmdvals is not None else {}
543 try:
544 arguments = arguments.format(**replacements)
545 except (KeyError, TypeError) as exc: # TypeError in case None instead of {}
546 _LOG.error(
547 "Could not replace command variables for job %s: replacement for %s not provided",
548 gwjob.name,
549 str(exc),
550 )
551 _LOG.debug("arguments: %s\ncmdvals: %s", arguments, replacements)
552 raise
553 return arguments
556def _replace_wms_vars(orig_string: str) -> str:
557 """Replace special wms placeholders in given string.
559 Parameters
560 ----------
561 orig_string : `str`
562 String in which to replace wms placeholders.
564 Returns
565 -------
566 updated_string : `str`
567 Given string with wms placeholders replaced.
568 """
569 values = {"attemptNum": "$$([NumJobStarts])"}
570 updated_string = orig_string
571 for key in re.findall(r"<WMS:([^>]+)>", orig_string):
572 try:
573 updated_string = updated_string.replace(rf"<WMS:{key}>", values[key])
574 except KeyError:
575 _LOG.error("Unrecognized WMS placeholder: %s in %s", key, orig_string)
576 raise
577 return updated_string
580def _handle_job_inputs(
581 generic_workflow: GenericWorkflow, job_name: str, use_shared: bool, out_prefix: str
582) -> dict[str, str]:
583 """Add job input files from generic workflow to job.
585 Parameters
586 ----------
587 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
588 The generic workflow (e.g., has executable name and arguments).
589 job_name : `str`
590 Unique name for the job.
591 use_shared : `bool`
592 Whether job has access to files via shared filesystem.
593 out_prefix : `str`
594 The root directory into which all WMS-specific files are written.
596 Returns
597 -------
598 htc_commands : `dict` [`str`, `str`]
599 HTCondor commands for the job submission script.
600 """
601 inputs = []
602 for gwf_file in generic_workflow.get_job_inputs(job_name, data=True, transfer_only=True): 602 ↛ 603line 602 didn't jump to line 603 because the loop on line 602 never started
603 _LOG.debug("src_uri=%s", gwf_file.src_uri)
605 uri = Path(gwf_file.src_uri)
607 # Note if use_shared and job_shared, don't need to transfer file.
609 if not use_shared: # Copy file using push to job
610 inputs.append(str(uri))
611 elif not gwf_file.job_shared: # Jobs require own copy
612 # if using shared filesystem, but still need copy in job. Use
613 # HTCondor's curl plugin for a local copy.
614 if uri.is_dir():
615 raise RuntimeError(
616 f"HTCondor plugin cannot transfer directories locally within job {gwf_file.src_uri}"
617 )
618 inputs.append(f"file://{uri}")
620 htc_commands = {}
621 if inputs: 621 ↛ 622line 621 didn't jump to line 622 because the condition on line 621 was never true
622 htc_commands["transfer_input_files"] = ",".join(inputs)
623 _LOG.debug("transfer_input_files=%s", htc_commands["transfer_input_files"])
624 return htc_commands
627def _handle_job_outputs(
628 generic_workflow: GenericWorkflow, job_name: str, use_shared: bool, out_prefix: str
629) -> dict[str, str]:
630 """Add job output files from generic workflow to the job if any.
632 Parameters
633 ----------
634 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
635 The generic workflow (e.g., has executable name and arguments).
636 job_name : `str`
637 Unique name for the job.
638 use_shared : `bool`
639 Whether job has access to files via shared filesystem.
640 out_prefix : `str`
641 The root directory into which all WMS-specific files are written.
643 Returns
644 -------
645 htc_commands : `dict` [`str`, `str`]
646 HTCondor commands for the job submission script.
647 """
648 outputs = []
649 output_remaps = []
650 for gwf_file in generic_workflow.get_job_outputs(job_name, data=True, transfer_only=True):
651 _LOG.debug("src_uri=%s", gwf_file.src_uri)
653 uri = Path(gwf_file.src_uri)
654 if not use_shared:
655 outputs.append(uri.name)
656 output_remaps.append(f"{uri.name}={str(uri)}")
658 # Set to an empty string to disable and only update if there are output
659 # files to transfer. Otherwise, HTCondor will transfer back all files in
660 # the job’s temporary working directory that have been modified or created
661 # by the job.
662 htc_commands = {"transfer_output_files": '""'}
663 if outputs:
664 htc_commands["transfer_output_files"] = ",".join(outputs)
665 _LOG.debug("transfer_output_files=%s", htc_commands["transfer_output_files"])
667 htc_commands["transfer_output_remaps"] = f'"{";".join(output_remaps)}"'
668 _LOG.debug("transfer_output_remaps=%s", htc_commands["transfer_output_remaps"])
669 return htc_commands
672def _create_periodic_release_expr(
673 memory: int, multiplier: float | None, limit: int, additional_expr: str = ""
674) -> str:
675 """Construct an HTCondorAd expression for releasing held jobs.
677 Parameters
678 ----------
679 memory : `int`
680 Requested memory in MB.
681 multiplier : `float` or None
682 Memory growth rate between retries.
683 limit : `int`
684 Memory limit.
685 additional_expr : `str`, optional
686 Expression to add to periodic_release. Defaults to empty string.
688 Returns
689 -------
690 expr : `str`
691 A string representing an HTCondor ClassAd expression for releasing job.
692 """
693 _LOG.debug(
694 "periodic_release: memory: %s, multiplier: %s, limit: %s, additional_expr: %s",
695 memory,
696 multiplier,
697 limit,
698 additional_expr,
699 )
701 # ctrl_bps sets multiplier to None in the GenericWorkflow if
702 # memoryMultiplier <= 1, but checking value just in case.
703 if (not multiplier or multiplier <= 1) and not additional_expr:
704 return ""
706 # Job ClassAds attributes 'HoldReasonCode' and 'HoldReasonSubCode' are
707 # UNDEFINED if job is not HELD (i.e. when 'JobStatus' is not 5).
708 # The special comparison operators ensure that all comparisons below will
709 # evaluate to FALSE in this case.
710 #
711 # Note:
712 # May not be strictly necessary. Operators '&&' and '||' are not strict so
713 # the entire expression should evaluate to FALSE when the job is not HELD.
714 # According to ClassAd evaluation semantics FALSE && UNDEFINED is FALSE,
715 # but better safe than sorry.
716 is_held = "JobStatus == 5"
717 is_retry_allowed = "NumJobStarts <= JobMaxRetries"
719 mem_expr = ""
720 if memory and multiplier and multiplier > 1 and limit:
721 was_mem_exceeded = (
722 "(HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 "
723 "|| HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34)"
724 )
725 was_below_limit = f"min({{int({memory} * pow({multiplier}, NumJobStarts - 1)), {limit}}}) < {limit}"
726 mem_expr = f"{was_mem_exceeded} && {was_below_limit}"
728 user_expr = ""
729 if additional_expr:
730 # Never auto release a job held by user.
731 user_expr = f"HoldReasonCode =!= 1 && {additional_expr}"
733 # Automatically release job if held because output file not found
734 # (e.g., job failed so didn't produce output file).
735 transfer_expr = "HoldReasonCode =?= 12"
737 expr = f"{is_held} && {is_retry_allowed}"
738 if user_expr and mem_expr:
739 expr += f" && ({transfer_expr} || {mem_expr} || {user_expr})"
740 elif user_expr:
741 expr += f" && ({transfer_expr} || {user_expr})"
742 elif mem_expr: 742 ↛ 745line 742 didn't jump to line 745 because the condition on line 742 was always true
743 expr += f" && ({transfer_expr} || {mem_expr})"
745 return expr
748def _create_periodic_remove_expr(memory, multiplier, limit):
749 """Construct an HTCondorAd expression for removing jobs from the queue.
751 Parameters
752 ----------
753 memory : `int`
754 Requested memory in MB.
755 multiplier : `float`
756 Memory growth rate between retries.
757 limit : `int`
758 Memory limit.
760 Returns
761 -------
762 expr : `str`
763 A string representing an HTCondor ClassAd expression for removing jobs.
764 """
765 # Job ClassAds attributes 'HoldReasonCode' and 'HoldReasonSubCode'
766 # are UNDEFINED if job is not HELD (i.e. when 'JobStatus' is not 5).
767 # The special comparison operators ensure that all comparisons below
768 # will evaluate to FALSE in this case.
769 #
770 # Note:
771 # May not be strictly necessary. Operators '&&' and '||' are not
772 # strict so the entire expression should evaluate to FALSE when the
773 # job is not HELD. According to ClassAd evaluation semantics
774 # FALSE && UNDEFINED is FALSE, but better safe than sorry.
775 is_held = "JobStatus == 5"
776 is_retry_disallowed = "NumJobStarts > JobMaxRetries"
778 mem_expr = ""
779 if memory and multiplier and multiplier > 1 and limit:
780 mem_limit_expr = f"min({{int({memory} * pow({multiplier}, NumJobStarts - 1)), {limit}}}) == {limit}"
782 mem_expr = ( # Add || here so only added if adding memory expr
783 " || ((HoldReasonCode =?= 34 && HoldReasonSubCode =?= 0 "
784 f"|| HoldReasonCode =?= 3 && HoldReasonSubCode =?= 34) && {mem_limit_expr})"
785 )
787 expr = f"{is_held} && ({is_retry_disallowed}{mem_expr})"
788 return expr
791def _create_request_memory_expr(memory, multiplier, limit):
792 """Construct an HTCondor ClassAd expression for safe memory scaling.
794 Parameters
795 ----------
796 memory : `int`
797 Requested memory in MB.
798 multiplier : `float`
799 Memory growth rate between retries.
800 limit : `int`
801 Memory limit.
803 Returns
804 -------
805 expr : `str`
806 A string representing an HTCondor ClassAd expression enabling safe
807 memory scaling between job retries.
808 """
809 # The check if the job was held due to exceeding memory requirements
810 # will be made *after* job was released back to the job queue (is in
811 # the IDLE state), hence the need to use `Last*` job ClassAds instead of
812 # the ones describing job's current state.
813 #
814 # Also, 'Last*' job ClassAds attributes are UNDEFINED when a job is
815 # initially put in the job queue. The special comparison operators ensure
816 # that all comparisons below will evaluate to FALSE in this case.
817 was_mem_exceeded = (
818 "LastJobStatus =?= 5 "
819 "&& (LastHoldReasonCode =?= 34 && LastHoldReasonSubCode =?= 0 "
820 "|| LastHoldReasonCode =?= 3 && LastHoldReasonSubCode =?= 34)"
821 )
823 # If job runs the first time or was held for reasons other than exceeding
824 # the memory, set the required memory to the requested value or use
825 # the memory value measured by HTCondor (MemoryUsage) depending on
826 # whichever is greater.
827 expr = (
828 f"({was_mem_exceeded}) "
829 f"? min({{int({memory} * pow({multiplier}, NumJobStarts)), {limit}}}) "
830 f": min({{max({{{memory}, MemoryUsage ?: 0}}), {limit}}})"
831 )
832 return expr
835def _gather_site_values(config, compute_site):
836 """Gather values specific to given site.
838 Parameters
839 ----------
840 config : `lsst.ctrl.bps.BpsConfig`
841 BPS configuration that includes necessary submit/runtime
842 information.
843 compute_site : `str`
844 Compute site name.
846 Returns
847 -------
848 site_values : `dict` [`str`, `~typing.Any`]
849 Values specific to the given site.
850 """
851 site_values = {"attrs": {}, "profile": {}}
852 search_opts = {}
853 if compute_site: 853 ↛ 857line 853 didn't jump to line 857 because the condition on line 853 was always true
854 search_opts["curvals"] = {"curr_site": compute_site}
856 # Determine the hard limit for the memory requirement.
857 found, limit = config.search("memoryLimit", opt=search_opts)
858 if not found: 858 ↛ 859line 858 didn't jump to line 859 because the condition on line 858 was never true
859 search_opts["default"] = DEFAULT_HTC_EXEC_PATT
860 _, patt = config.search("executeMachinesPattern", opt=search_opts)
861 del search_opts["default"]
863 # To reduce the amount of data, ignore dynamic slots (if any) as,
864 # by definition, they cannot have more memory than
865 # the partitionable slot they are the part of.
866 constraint = f'SlotType != "Dynamic" && regexp("{patt}", Machine)'
867 pool_info = condor_status(constraint=constraint)
868 try:
869 limit = max(int(info["TotalSlotMemory"]) for info in pool_info.values())
870 except ValueError:
871 _LOG.debug("No execute machine in the pool matches %s", patt)
872 if limit: 872 ↛ 874line 872 didn't jump to line 874 because the condition on line 872 was always true
873 config[".bps_defined.memory_limit"] = limit
874 site_values["memoryLimit"] = limit
876 _, site_values["bpsUseShared"] = config.search("bpsUseShared", opt={"default": False})
878 searchobj = config[f".site.{compute_site}.profile.condor"]
879 if searchobj:
880 search_opts["searchobj"] = searchobj
881 search_opts["replaceVars"] = True
882 for key in searchobj:
883 if key.startswith("+"):
884 _, val = config.search(key, opt=search_opts)
885 site_values["attrs"][key[1:]] = val
886 else:
887 _, val = config.search(key, opt=search_opts)
888 site_values["profile"][key] = val
890 searchobj = config[f".site.{compute_site}"]
891 if searchobj:
892 for key, value in searchobj.items():
893 if key not in site_values and key not in ["attrs", "profile"]: 893 ↛ 894line 893 didn't jump to line 894 because the condition on line 893 was never true
894 site_values[key] = value
896 _LOG.debug("site_values = %s", site_values)
897 return site_values
900def _gather_label_values(config: BpsConfig, label: str) -> dict[str, Any]:
901 """Gather values specific to given job label.
903 Parameters
904 ----------
905 config : `lsst.ctrl.bps.BpsConfig`
906 BPS configuration that includes necessary submit/runtime
907 information.
908 label : `str`
909 GenericWorkflowJob label.
911 Returns
912 -------
913 values : `dict` [`str`, `~typing.Any`]
914 Values specific to the given job label.
915 """
916 _LOG.debug("_gather_label_values: label = %s", label)
917 values: dict[str, Any] = {"attrs": {}, "profile": {}}
919 search_opts = config.get_search_opts(label)
921 # Determine the hard limit for the memory requirement.
922 found, limit = config.search("memoryLimit", opt=search_opts)
923 if not found: 923 ↛ 924line 923 didn't jump to line 924 because the condition on line 923 was never true
924 search_opts["default"] = DEFAULT_HTC_EXEC_PATT
925 _, patt = config.search("executeMachinesPattern", opt=search_opts)
926 del search_opts["default"]
928 # To reduce the amount of data, ignore dynamic slots (if any) as,
929 # by definition, they cannot have more memory than
930 # the partitionable slot they are the part of.
931 constraint = f'SlotType != "Dynamic" && regexp("{patt}", Machine)'
932 pool_info = condor_status(constraint=constraint)
933 try:
934 limit = max(int(info["TotalSlotMemory"]) for info in pool_info.values())
935 except ValueError:
936 _LOG.debug("No execute machine in the pool matches %s", patt)
938 if limit: 938 ↛ 942line 938 didn't jump to line 942 because the condition on line 938 was always true
939 config[".bps_defined.memory_limit"] = limit
940 values["memoryLimit"] = limit
942 values["bpsUseShared"] = False
943 found, value = config.search("bpsUseShared", opt=search_opts)
944 if found:
945 values["bpsUseShared"] = value
947 found, value = config.search("releaseExpr", opt=search_opts)
948 if found:
949 values["releaseExpr"] = value
951 values["overwriteJobFiles"] = True
952 found, value = config.search("overwriteJobFiles", opt=search_opts)
953 if found:
954 values["overwriteJobFiles"] = value
956 found, value = config.search("releaseExpr", opt=search_opts)
957 if found:
958 values["releaseExpr"] = value
960 found, value = config.search("bpsMakeCommand", opt=search_opts)
961 values["bpsMakeCommand"] = value if found else True
962 if found and not value:
963 search_opts["skipNames"] = {"gwjobCommand", "gwjobExports"}
964 _LOG.debug("_gather_label_values: search_opts = %s", search_opts)
965 found, value = config.search("payloadCommand", opt=search_opts)
966 if found: 966 ↛ 970line 966 didn't jump to line 970 because the condition on line 966 was always true
967 values["payloadCommand"] = value
968 _LOG.debug("payloadCommand = %s", value)
970 found, value = config.search("bpsUseHTCEnvironment", opt=search_opts)
971 values["bpsUseHTCEnvironment"] = value if found else values["bpsMakeCommand"]
973 found, value = config.search("nodeset", opt=search_opts)
974 if found:
975 values["nodeset"] = value
977 found, profile_sect = config.search("profile", opt=search_opts)
978 if found and "condor" in profile_sect:
979 for subkey, val in profile_sect["condor"].items():
980 if subkey.startswith("+"):
981 values["attrs"][subkey[1:]] = val
982 else:
983 values["profile"][subkey] = val
985 # Copy all of the site values.
986 if "curr_site" in search_opts["curvals"]:
987 site_obj = config[f".site.{search_opts['curvals']['curr_site']}"]
988 if site_obj: 988 ↛ 993line 988 didn't jump to line 993 because the condition on line 988 was always true
989 for key, value in site_obj.items():
990 if key not in values and key not in ["attrs", "profile"]:
991 values[key] = value
993 _LOG.debug("_gather_label_values: label = %s, values = %s", label, values)
995 return values
998def _group_to_subdag(
999 config: BpsConfig, generic_workflow_group: GenericWorkflowGroup, out_prefix: str
1000) -> HTCJob:
1001 """Convert a generic workflow group to an HTCondor dag.
1003 Parameters
1004 ----------
1005 config : `lsst.ctrl.bps.BpsConfig`
1006 Workflow configuration.
1007 generic_workflow_group : `lsst.ctrl.bps.GenericWorkflowGroup`
1008 The generic workflow group to convert.
1009 out_prefix : `str`
1010 Location prefix to be used when creating jobs.
1012 Returns
1013 -------
1014 htc_job : `lsst.ctrl.bps.htcondor.HTCJob`
1015 Job for running the HTCondor dag.
1016 """
1017 jobname = f"wms_{generic_workflow_group.name}"
1018 htc_job = HTCJob(name=jobname, label=generic_workflow_group.label)
1019 htc_job.add_dag_cmds({"dir": f"subdags/{jobname}"})
1020 htc_job.subdag = _generic_workflow_to_htcondor_dag(config, generic_workflow_group, out_prefix)
1021 if not generic_workflow_group.blocking: 1021 ↛ 1027line 1021 didn't jump to line 1027 because the condition on line 1021 was always true
1022 htc_job.dagcmds["post"] = {
1023 "defer": "",
1024 "executable": f"{os.path.dirname(__file__)}/subdag_post.sh",
1025 "arguments": f"{jobname} $RETURN",
1026 }
1027 return htc_job
1030def _create_check_job(group_job_name: str, job_label: str, site_values: dict) -> HTCJob:
1031 """Create a job to check status of a group job.
1033 Parameters
1034 ----------
1035 group_job_name : `str`
1036 Name of the group job.
1037 job_label : `str`
1038 Label to use for the check status job.
1039 site_values : `dict`
1040 Site specific values.
1042 Returns
1043 -------
1044 htc_job : `lsst.ctrl.bps.htcondor.HTCJob`
1045 Job description for the job to check group job status.
1046 """
1047 htc_job = HTCJob(name=f"wms_check_status_{group_job_name}", label=job_label)
1048 htc_job.subfile = "${CTRL_BPS_HTCONDOR_DIR}/python/lsst/ctrl/bps/htcondor/check_group_status.sub"
1049 # ADD nodeset to VARS
1050 job_vars = {"group_job_name": group_job_name}
1051 if "nodeset" in site_values and site_values["nodeset"]:
1052 job_vars["job_nodeset"] = site_values["nodeset"]
1053 htc_job.add_dag_cmds({"dir": f"subdags/{group_job_name}", "vars": job_vars})
1055 return htc_job
1058def _generic_workflow_to_htcondor_dag(
1059 config: BpsConfig, generic_workflow: GenericWorkflow, out_prefix: str
1060) -> HTCDag:
1061 """Convert a GenericWorkflow to a HTCDag.
1063 Parameters
1064 ----------
1065 config : `lsst.ctrl.bps.BpsConfig`
1066 Workflow configuration.
1067 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
1068 The GenericWorkflow to convert.
1069 out_prefix : `str`
1070 Location prefix where the HTCondor files will be written.
1072 Returns
1073 -------
1074 dag : `lsst.ctrl.bps.htcondor.HTCDag`
1075 The HTCDag representation of the given GenericWorkflow.
1076 """
1077 dag = HTCDag(name=generic_workflow.name)
1079 _LOG.debug("htcondor dag attribs %s", generic_workflow.run_attrs)
1080 dag.add_attribs(generic_workflow.run_attrs)
1081 dag.add_attribs(
1082 {
1083 "bps_run_quanta": create_count_summary(generic_workflow.quanta_counts),
1084 "bps_job_summary": create_count_summary(generic_workflow.job_counts),
1085 }
1086 )
1087 _, tmp_template = config.search("subDirTemplate", opt={"replaceVars": False, "default": ""})
1089 _, save_htc_dot = config.search("saveHTCdot", opt={"default": False})
1090 dag.graph["write_dot"] = save_htc_dot
1092 if isinstance(tmp_template, str): 1092 ↛ 1095line 1092 didn't jump to line 1095 because the condition on line 1092 was always true
1093 subdir_template = defaultdict(lambda: tmp_template)
1094 else:
1095 subdir_template = tmp_template
1097 # Save list of lazy group jobs for later extra handling.
1098 lazy_groups = []
1100 # Create all DAG jobs
1101 cached_values = {} # Cache label-specific values to reduce config lookups.
1102 # Note: Can't use get_job_by_label because those only include payload jobs.
1103 for job_name in generic_workflow:
1104 gwjob = generic_workflow.get_job(job_name)
1105 if gwjob.node_type in [ 1105 ↛ 1120line 1105 didn't jump to line 1120 because the condition on line 1105 was always true
1106 GenericWorkflowNodeType.PAYLOAD,
1107 GenericWorkflowNodeType.LAZY_GROUP,
1108 ]:
1109 gwjob = cast(GenericWorkflowJob, gwjob)
1110 if gwjob.label not in cached_values:
1111 cached_values[gwjob.label] = _gather_label_values(config, gwjob.label)
1112 _LOG.debug("cached: %s= %s", gwjob.label, cached_values[gwjob.label])
1113 htc_job = _create_job(
1114 subdir_template[gwjob.label],
1115 cached_values[gwjob.label],
1116 generic_workflow,
1117 gwjob,
1118 out_prefix,
1119 )
1120 elif gwjob.node_type == GenericWorkflowNodeType.NOOP:
1121 gwjob = cast(GenericWorkflowNoopJob, gwjob)
1122 htc_job = HTCJob(f"wms_{gwjob.name}", label=gwjob.label)
1123 htc_job.subfile = "${CTRL_BPS_HTCONDOR_DIR}/python/lsst/ctrl/bps/htcondor/noop.sub"
1124 htc_job.add_job_attrs({"bps_job_name": gwjob.name, "bps_job_label": gwjob.label})
1125 htc_job.add_dag_cmds({"noop": True})
1126 elif gwjob.node_type == GenericWorkflowNodeType.GROUP:
1127 gwjob = cast(GenericWorkflowGroup, gwjob)
1128 cached_values[gwjob.label] = _gather_label_values(config, gwjob.label)
1129 htc_job = _group_to_subdag(config, gwjob, out_prefix)
1130 else:
1131 raise RuntimeError(f"Unsupported generic workflow node type {gwjob.node_type} ({gwjob.name})")
1132 _LOG.debug("Calling adding job %s %s", htc_job.name, htc_job.label)
1133 dag.add_job(htc_job)
1135 # Have to add the placeholder job for the workflow
1136 if gwjob.node_type == GenericWorkflowNodeType.LAZY_GROUP:
1137 lazy_groups.append(gwjob.name)
1139 # Add job dependencies to the DAG (be careful with wms_ jobs)
1140 for job_name in generic_workflow:
1141 gwjob = generic_workflow.get_job(job_name)
1142 parent_name = (
1143 gwjob.name
1144 if gwjob.node_type in [GenericWorkflowNodeType.PAYLOAD, GenericWorkflowNodeType.LAZY_GROUP]
1145 else f"wms_{gwjob.name}"
1146 )
1147 successor_jobs = [generic_workflow.get_job(j) for j in generic_workflow.successors(job_name)]
1148 children_names = []
1149 if gwjob.node_type == GenericWorkflowNodeType.GROUP: 1149 ↛ 1150line 1149 didn't jump to line 1150 because the condition on line 1149 was never true
1150 gwjob = cast(GenericWorkflowGroup, gwjob)
1151 group_children = [] # Dependencies between same group jobs
1152 for sjob in successor_jobs:
1153 if sjob.node_type == GenericWorkflowNodeType.GROUP and sjob.label == gwjob.label:
1154 group_children.append(f"wms_{sjob.name}")
1155 elif sjob.node_type == GenericWorkflowNodeType.PAYLOAD:
1156 children_names.append(sjob.name)
1157 else:
1158 children_names.append(f"wms_{sjob.name}")
1159 if group_children:
1160 dag.add_job_relationships([parent_name], group_children)
1161 if not gwjob.blocking:
1162 # Since subdag will always succeed, need to add a special
1163 # job that fails if group failed to block payload children.
1164 check_job = _create_check_job(
1165 f"wms_{gwjob.name}", gwjob.label, cached_values.get(gwjob.label, {})
1166 )
1167 dag.add_job(check_job)
1168 dag.add_job_relationships([f"wms_{gwjob.name}"], [check_job.name])
1169 parent_name = check_job.name
1170 else:
1171 for sjob in successor_jobs:
1172 if sjob.node_type in [ 1172 ↛ 1178line 1172 didn't jump to line 1178 because the condition on line 1172 was always true
1173 GenericWorkflowNodeType.PAYLOAD,
1174 GenericWorkflowNodeType.LAZY_GROUP,
1175 ]:
1176 children_names.append(sjob.name)
1177 else:
1178 children_names.append(f"wms_{sjob.name}")
1180 dag.add_job_relationships([parent_name], children_names)
1182 # Go back and add placeholder jobs for the lazy group dags
1183 for lazy_group_name in lazy_groups:
1184 _add_lazy_placeholder(lazy_group_name, generic_workflow, dag, out_prefix)
1186 # If final job exists in generic workflow, create DAG final job
1187 final = generic_workflow.get_final()
1188 if final and isinstance(final, GenericWorkflowJob):
1189 if final.label not in cached_values: 1189 ↛ 1191line 1189 didn't jump to line 1191 because the condition on line 1189 was always true
1190 cached_values[final.label] = _gather_label_values(config, final.label)
1191 final_htjob = _create_job(
1192 subdir_template[final.label],
1193 cached_values[final.label],
1194 generic_workflow,
1195 final,
1196 out_prefix,
1197 )
1198 if "post" not in final_htjob.dagcmds: 1198 ↛ 1204line 1198 didn't jump to line 1204 because the condition on line 1198 was always true
1199 final_htjob.dagcmds["post"] = {
1200 "defer": "",
1201 "executable": f"{os.path.dirname(__file__)}/final_post.sh",
1202 "arguments": f"{final.name} $DAG_STATUS $RETURN",
1203 }
1204 dag.add_final_job(final_htjob)
1205 elif final and isinstance(final, GenericWorkflow):
1206 raise NotImplementedError("HTCondor plugin does not support a workflow as the final job")
1207 elif final: 1207 ↛ 1208line 1207 didn't jump to line 1208 because the condition on line 1207 was never true
1208 raise TypeError(f"Invalid type for GenericWorkflow.get_final() results ({type(final)})")
1210 return dag
1213def _add_lazy_placeholder(
1214 prepare_job_name: str,
1215 generic_workflow: GenericWorkflow,
1216 dag: HTCDag,
1217 out_prefix: str,
1218):
1219 _LOG.debug("prepare_job_name = %s", prepare_job_name)
1221 # Make a fake job for the placeholder workflow
1222 job = HTCJob("placeholder")
1223 job.add_job_cmds(
1224 {
1225 "executable": "/usr/bin/echo",
1226 "arguments": '"BPS internal error - placeholder DAG was not replaced."',
1227 }
1228 )
1229 job.add_job_attrs({"bps_job_name": "placeholder", "bps_job_label": "placeholder"})
1230 job.add_job_attrs(generic_workflow.run_attrs)
1232 # Placeholder dag name needs to be the run name because that's
1233 # currently what the bps code will name the workflow.
1234 placeholder_dag_name = generic_workflow.name.replace("_ctrl", "")
1235 placeholder_dag = HTCDag(name=placeholder_dag_name)
1236 _LOG.debug("dag name = %s", placeholder_dag.graph["name"])
1237 placeholder_dag.add_attribs(generic_workflow.run_attrs)
1238 placeholder_dag.add_job(job)
1240 # To help ordering in bps report, save info so can figure out where
1241 # to put jobs in lazy dag in bps_job_label
1242 dag.add_attribs({"bps_lazy_mapping": f"{placeholder_dag.name}:{prepare_job_name}"})
1244 # The subdag job to be added to the control dag
1245 dag_job = HTCJob("wms_lazy_payload", "wms_lazy_payload")
1246 dag_job.subfile = f"{placeholder_dag.name}.condor.sub"
1247 dag_job.subdag = placeholder_dag
1249 dag.add_job(dag_job)
1251 # Update edges inserting between prepare job and any following jobs.
1252 successor_jobs = list(dag.successors(prepare_job_name))
1253 for job in successor_jobs:
1254 dag.add_edge(dag_job.name, job)
1255 dag.remove_edge(prepare_job_name, job)
1257 dag.add_edge(prepare_job_name, dag_job.name)
1259 _LOG.debug("_add_lazy_placeholder: edges = %s", dag.edges)
1262def _update_job_summary(subworkflow_name: str, subworkflow_summary: str, submit_path: str) -> None:
1263 """Add summary for jobs in given subworkflow to parent workflow's job
1264 summary.
1266 Parameters
1267 ----------
1268 subworkflow_name : `str`
1269 Name of subworkflow.
1270 subworkflow_summary : `str`
1271 Job summary for the subworkflow.
1272 submit_path : `str`
1273 Directory in which to find the DAG info file.
1274 """
1275 _LOG.debug("submit_path = %s", submit_path)
1276 filename, dag_info = read_dag_info(submit_path)
1277 _LOG.debug("dag_info = %s", dag_info)
1279 schedd_name = next(iter(dag_info))
1280 dag_values = next(iter(dag_info[schedd_name].values()))
1281 _LOG.debug("dag_values = %s", dag_values)
1283 # Get lazy mapping and the job that generated this dag
1284 generator_name = None
1285 lazy_mapping = dag_values.get("bps_lazy_mapping", None)
1286 if lazy_mapping: # find this workflow's generator job
1287 for part in lazy_mapping.split(";"):
1288 info = part.split(":")
1289 if info[0] == subworkflow_name:
1290 generator_name = info[1]
1291 break
1293 # Update bps_job_summary
1294 _LOG.debug(
1295 "Before replace, name = %s, bps_job_summary = %s, add summary = %s, generator_name = %s",
1296 subworkflow_name,
1297 dag_values["bps_job_summary"],
1298 subworkflow_summary,
1299 generator_name,
1300 )
1301 if generator_name:
1302 generator_summary = f"{generator_name}:1"
1303 dag_values["bps_job_summary"] = dag_values["bps_job_summary"].replace(
1304 generator_summary, f"{generator_summary};{subworkflow_summary}"
1305 )
1306 else:
1307 # just append to end of bps_job_summary
1308 dag_values["bps_job_summary"] += f";{subworkflow_summary}"
1310 _LOG.debug(
1311 "After replace, name = %s, bps_job_summary = %s",
1312 subworkflow_name,
1313 dag_values["bps_job_summary"],
1314 )
1316 # Save updated bps_job_summary
1317 write_dag_info(filename, dag_info)