Coverage for python/lsst/ctrl/bps/transform.py: 80%
328 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-19 09:38 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-19 09:38 +0000
1# This file is part of ctrl_bps.
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"""Driver for the transformation of a QuantumGraph into a generic workflow."""
30import copy
31import dataclasses
32import logging
33import math
34import os
35import re
37from lsst.ctrl.bps import ClusteredQuantumGraph
38from lsst.pipe.base import QuantumGraph
39from lsst.utils.logging import VERBOSE
40from lsst.utils.timer import timeMethod
42from . import (
43 DEFAULT_MEM_RETRIES,
44 BpsConfig,
45 GenericWorkflow,
46 GenericWorkflowExec,
47 GenericWorkflowFile,
48 GenericWorkflowJob,
49)
50from .bps_utils import (
51 WhenToSaveQuantumGraphs,
52 create_job_quantum_graph_filename,
53 save_qg_subgraph,
54)
56# All available job attributes.
57_ATTRS_ALL = frozenset([field.name for field in dataclasses.fields(GenericWorkflowJob)])
59# Job attributes that need to be set to their maximal value in the cluster.
60_ATTRS_MAX = frozenset(
61 {
62 "memory_multiplier",
63 "number_of_retries",
64 "request_cpus",
65 "request_memory",
66 "request_memory_max",
67 }
68)
70# Job attributes that need to be set to sum of their values in the cluster.
71_ATTRS_SUM = frozenset(
72 {
73 "request_disk",
74 "request_walltime",
75 }
76)
78# Job attributes do not fall into a specific category
79_ATTRS_MISC = frozenset(
80 {
81 "label", # taskDef labels aren't same in job and may not match job label
82 "cmdvals",
83 "profile",
84 "attrs",
85 }
86)
88# Attributes that need to be the same for each quanta in the cluster.
89_ATTRS_UNIVERSAL = frozenset(_ATTRS_ALL - (_ATTRS_MAX | _ATTRS_MISC | _ATTRS_SUM))
91_LOG = logging.getLogger(__name__)
94@timeMethod(logger=_LOG, logLevel=VERBOSE)
95def transform(
96 config: BpsConfig, cqgraph: ClusteredQuantumGraph, prefix: str
97) -> tuple[GenericWorkflow, BpsConfig]:
98 """Transform a ClusteredQuantumGraph to a GenericWorkflow.
100 Parameters
101 ----------
102 config : `lsst.ctrl.bps.BpsConfig`
103 BPS configuration.
104 cqgraph : `lsst.ctrl.bps.ClusteredQuantumGraph`
105 A clustered quantum graph to transform into a generic workflow.
106 prefix : `str`
107 Root path for any output files.
109 Returns
110 -------
111 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
112 The generic workflow transformed from the clustered quantum graph.
113 generic_workflow_config : `lsst.ctrl.bps.BpsConfig`
114 Configuration to accompany GenericWorkflow.
115 """
116 if cqgraph.name is not None:
117 name = cqgraph.name
118 else:
119 _, name = config.search("uniqProcName", opt={"required": True})
121 generic_workflow = create_generic_workflow(config, cqgraph, name, prefix)
122 generic_workflow_config = create_generic_workflow_config(config, prefix)
124 return generic_workflow, generic_workflow_config
127def add_workflow_init_nodes(config, qgraph, generic_workflow):
128 """Add nodes to workflow graph that perform initialization steps.
130 Assumes that all of the initialization should be executed prior to any
131 of the current workflow.
133 Parameters
134 ----------
135 config : `lsst.ctrl.bps.BpsConfig`
136 BPS configuration.
137 qgraph : `lsst.pipe.base.graph.QuantumGraph`
138 The quantum graph the generic workflow represents.
139 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
140 Generic workflow to which the initialization steps should be added.
141 """
142 # Create a workflow graph that will have task and file nodes necessary for
143 # initializing the pipeline execution
144 init_workflow = create_init_workflow(config, qgraph, generic_workflow.get_file("runQgraphFile"))
145 _LOG.debug("init_workflow nodes = %s", init_workflow.nodes())
146 generic_workflow.add_workflow_source(init_workflow)
149def create_init_workflow(
150 config: BpsConfig, qgraph: QuantumGraph, qgraph_gwfile: GenericWorkflowFile
151) -> GenericWorkflow:
152 """Create workflow for running initialization job(s).
154 Parameters
155 ----------
156 config : `lsst.ctrl.bps.BpsConfig`
157 BPS configuration.
158 qgraph : `lsst.pipe.base.graph.QuantumGraph`
159 The quantum graph the generic workflow represents.
160 qgraph_gwfile : `lsst.ctrl.bps.GenericWorkflowFile`
161 File object for the full run QuantumGraph file.
163 Returns
164 -------
165 init_workflow : `lsst.ctrl.bps.GenericWorkflow`
166 GenericWorkflow consisting of job(s) to initialize workflow.
167 """
168 _LOG.debug("creating init subgraph")
169 _LOG.debug("creating init task input(s)")
170 search_opt = {
171 "curvals": {"curr_pipetask": "pipetaskInit"},
172 "replaceVars": False,
173 "expandEnvVars": False,
174 "replaceEnvVars": True,
175 "required": False,
176 }
177 found, value = config.search("computeSite", opt=search_opt)
178 if found: 178 ↛ 180line 178 didn't jump to line 180 because the condition on line 178 was always true
179 search_opt["curvals"]["curr_site"] = value
180 found, value = config.search("computeCloud", opt=search_opt)
181 if found:
182 search_opt["curvals"]["curr_cloud"] = value
184 init_workflow = GenericWorkflow("init")
185 init_workflow.add_file(qgraph_gwfile)
187 # create job for executing --init-only
188 gwjob = GenericWorkflowJob("pipetaskInit", "pipetaskInit")
190 job_values = _get_job_values(config, search_opt, "runQuantumCommand")
191 job_values["name"] = "pipetaskInit"
192 job_values["label"] = "pipetaskInit"
194 # Adjust job attributes values if necessary.
195 _handle_job_values(job_values, gwjob)
197 init_workflow.add_job(gwjob)
198 init_workflow.add_job_inputs(gwjob.name, [qgraph_gwfile])
199 _enhance_command(config, init_workflow, gwjob, {})
201 return init_workflow
204def _enhance_command(config, generic_workflow, gwjob, cached_job_values):
205 """Enhance command line with env and file placeholders
206 and gather command line values.
208 Parameters
209 ----------
210 config : `lsst.ctrl.bps.BpsConfig`
211 BPS configuration.
212 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
213 Generic workflow that contains the job.
214 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
215 Generic workflow job to which the updated executable, arguments,
216 and values should be saved.
217 cached_job_values : `dict` [`str`, dict[`str`, `~typing.Any`]]
218 Cached values common across jobs with same label. Updated if values
219 aren't already saved for given gwjob's label.
220 """
221 _LOG.debug("gwjob given to _enhance_command: %s", gwjob)
223 search_opt = config.get_search_opts(gwjob.label)
225 search_opt["curvals"]["jobName"] = gwjob.name
226 search_opt["curvals"]["jobLabel"] = gwjob.label
227 for key, value in gwjob.tags.items():
228 search_opt["curvals"][key] = value
230 search_opt.update(
231 {
232 "replaceVars": False,
233 "expandEnvVars": False,
234 "replaceEnvVars": True,
235 "required": False,
236 }
237 )
239 if gwjob.label not in cached_job_values:
240 cached_job_values[gwjob.label] = {}
241 # Allowing whenSaveJobQgraph and useLazyCommands per pipetask label.
242 key = "whenSaveJobQgraph"
243 _, when_save = config.search(key, opt=search_opt)
244 cached_job_values[gwjob.label][key] = WhenToSaveQuantumGraphs[when_save.upper()]
246 key = "useLazyCommands"
247 search_opt["default"] = True
248 _, cached_job_values[gwjob.label][key] = config.search(key, opt=search_opt)
249 del search_opt["default"]
251 # Change qgraph variable to match whether using run or per-job qgraph
252 # Note: these are lookup keys, not actual physical filenames.
253 if cached_job_values[gwjob.label]["whenSaveJobQgraph"] == WhenToSaveQuantumGraphs.NEVER:
254 gwjob.arguments = gwjob.arguments.replace("{qgraphFile}", "{runQgraphFile}")
255 elif gwjob.name == "pipetaskInit": 255 ↛ 256line 255 didn't jump to line 256 because the condition on line 255 was never true
256 gwjob.arguments = gwjob.arguments.replace("{qgraphFile}", "{runQgraphFile}")
257 else: # Needed unique file keys for per-job QuantumGraphs
258 gwjob.arguments = gwjob.arguments.replace("{qgraphFile}", f"{{qgraphFile_{gwjob.name}}}")
260 # Replace files with special placeholders
261 for gwfile in generic_workflow.get_job_inputs(gwjob.name):
262 gwjob.arguments = gwjob.arguments.replace(f"{{{gwfile.name}}}", f"<FILE:{gwfile.name}>")
263 for gwfile in generic_workflow.get_job_outputs(gwjob.name):
264 gwjob.arguments = gwjob.arguments.replace(f"{{{gwfile.name}}}", f"<FILE:{gwfile.name}>")
266 # Replace wms variables with wms placeholders.
267 gwjob.arguments = re.sub(
268 r"{wms([^}]+)}", lambda x: f"<WMS:{x[1][0].lower() + x[1][1:]}>", gwjob.arguments
269 )
271 # Save dict of other values needed to complete command line.
272 # (Be careful to not replace env variables as they may
273 # be different in compute job.)
274 search_opt["replaceVars"] = True
275 _LOG.debug("before cmdvals = %s (search_opt = %s)", gwjob.cmdvals, search_opt)
276 for key in re.findall(r"{([^}]+)}", gwjob.arguments):
277 _LOG.debug("looking for %s in cmdvals", key)
278 if key in gwjob.cmdvals:
279 continue
280 elif key in cached_job_values[gwjob.label]:
281 gwjob.cmdvals[key] = cached_job_values[gwjob.label][key]
282 else:
283 _, gwjob.cmdvals[key] = config.search(key, opt=search_opt)
284 _LOG.debug("after cmdvals = %s", gwjob.cmdvals)
286 # backwards compatibility
287 if not cached_job_values[gwjob.label]["useLazyCommands"]: 287 ↛ 288line 287 didn't jump to line 288 because the condition on line 287 was never true
288 if "bpsUseShared" not in cached_job_values[gwjob.label]:
289 key = "bpsUseShared"
290 search_opt["default"] = True
291 _, cached_job_values[gwjob.label][key] = config.search(key, opt=search_opt)
292 del search_opt["default"]
294 gwjob.arguments = _fill_arguments(
295 cached_job_values[gwjob.label]["bpsUseShared"], generic_workflow, gwjob.arguments, gwjob.cmdvals
296 )
299def _fill_arguments(use_shared, generic_workflow, arguments, cmdvals):
300 """Replace placeholders in command line string in job.
302 Parameters
303 ----------
304 use_shared : `bool`
305 Whether using shared filesystem.
306 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
307 Generic workflow containing the job.
308 arguments : `str`
309 String containing placeholders.
310 cmdvals : `dict` [`str`, `~typing.Any`]
311 Any command line values that can be used to replace placeholders.
313 Returns
314 -------
315 arguments : `str`
316 Command line with FILE and ENV placeholders replaced.
317 """
318 # Replace file placeholders
319 for file_key in re.findall(r"<FILE:([^>]+)>", arguments):
320 gwfile = generic_workflow.get_file(file_key)
321 if not gwfile.wms_transfer:
322 # Must assume full URI if in command line and told WMS is not
323 # responsible for transferring file.
324 uri = gwfile.src_uri
325 elif use_shared:
326 if gwfile.job_shared:
327 # Have shared filesystems and jobs can share file.
328 uri = gwfile.src_uri
329 else:
330 uri = os.path.basename(gwfile.src_uri)
331 else: # Using push transfer
332 uri = os.path.basename(gwfile.src_uri)
334 arguments = arguments.replace(f"<FILE:{file_key}>", uri)
336 # Replace env placeholder with submit-side values
337 arguments = re.sub(r"<ENV:([^>]+)>", r"$\1", arguments)
338 arguments = os.path.expandvars(arguments)
340 # Replace remaining vars
341 arguments = arguments.format(**cmdvals)
343 return arguments
346def _get_qgraph_gwfile(config, save_qgraph_per_job, gwjob, run_qgraph_file, prefix):
347 """Get qgraph location to be used by job.
349 Parameters
350 ----------
351 config : `lsst.ctrl.bps.BpsConfig`
352 Bps configuration.
353 save_qgraph_per_job : `lsst.ctrl.bps.bps_utils.WhenToSaveQuantumGraphs`
354 What submission stage to save per-job qgraph files (or NEVER)
355 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
356 Job for which determining QuantumGraph file.
357 run_qgraph_file : `lsst.ctrl.bps.GenericWorkflowFile`
358 File representation of the full run QuantumGraph.
359 prefix : `str`
360 Path prefix for any files written.
362 Returns
363 -------
364 gwfile : `lsst.ctrl.bps.GenericWorkflowFile`
365 Representation of butler location (may not include filename).
366 """
367 qgraph_gwfile = None
368 if save_qgraph_per_job != WhenToSaveQuantumGraphs.NEVER: 368 ↛ 369line 368 didn't jump to line 369 because the condition on line 368 was never true
369 qgraph_gwfile = GenericWorkflowFile(
370 f"qgraphFile_{gwjob.name}",
371 src_uri=create_job_quantum_graph_filename(config, gwjob, prefix),
372 wms_transfer=True,
373 job_access_remote=True,
374 job_shared=True,
375 )
376 else:
377 qgraph_gwfile = run_qgraph_file
379 return qgraph_gwfile
382def _get_job_values(config, search_opt, cmd_line_key):
383 """Gather generic workflow job values from the bps config.
385 Parameters
386 ----------
387 config : `lsst.ctrl.bps.BpsConfig`
388 Bps configuration.
389 search_opt : `dict` [`str`, `~typing.Any`]
390 Search options to be used when searching config.
391 cmd_line_key : `str` or None
392 Which command line key to search for (e.g., "runQuantumCommand").
394 Returns
395 -------
396 job_values : `dict` [ `str`, `~typing.Any` ]`
397 A mapping between job attributes and their values.
398 """
399 _LOG.debug("cmd_line_key=%s, search_opt=%s", cmd_line_key, search_opt)
401 # Create a dummy job to easily access the default values.
402 default_gwjob = GenericWorkflowJob("default_job", "default_label")
404 job_values = {}
405 for attr in _ATTRS_ALL:
406 # Variable names in yaml are camel case instead of snake case.
407 yaml_name = re.sub(r"_(\S)", lambda match: match.group(1).upper(), attr)
408 found, value = config.search(yaml_name, opt=search_opt)
409 if found:
410 job_values[attr] = value
411 else:
412 job_values[attr] = getattr(default_gwjob, attr)
414 # Need to replace all config variables in environment values.
415 # Also change env vars in environment values to bash syntax.
416 #
417 # Note: Because job_values["environment"] is a BpsConfig and
418 # currently cannot have 2 search objects, for each environment
419 # setting, we have to get the setting string as is and then
420 # separately use the overall config to replace values inside
421 # the setting string.
422 tmp_job_env = job_values.get("environment", None)
423 if tmp_job_env:
424 _LOG.debug("_get_job_values: job_values['environment'] = %s", tmp_job_env)
426 # Don't want to replace when getting environment setting string.
427 as_is_search_opt = {
428 "replaceVars": False,
429 "expandEnvVars": False,
430 "replaceEnvBps2Shell": False,
431 "replaceEnvShell2Bps": False,
432 }
434 # When updating environment string, use given search options,
435 # but ensure making the environment string using bash syntax.
436 env_search_opt = copy.copy(search_opt)
437 env_search_opt["replaceVars"] = True # Replace bps config variables.
438 env_search_opt["replaceEnvBps2Shell"] = False # Replace bps <ENV:var> syntax.
439 env_search_opt["replaceEnvShell2Bps"] = True # Do not replace shell env syntax.
440 env_search_opt["expandEnvVars"] = False # Do not replace with submission env value.
442 job_env = {} # While replacing variables, convert to plain dict.
444 for name in tmp_job_env:
445 # Get environment setting string as is.
446 value = tmp_job_env.search(name, as_is_search_opt)[1]
447 _LOG.debug("_get_job_values: as is value for %s = %s", name, value)
448 # Replace config vars and env placeholders
449 job_env[name] = config.modify_value(name, str(value), env_search_opt)
450 _LOG.debug("_get_job_values: new env value for %s = %s", name, job_env[name])
451 # Save new dictionary back with other job values.
452 job_values["environment"] = job_env
454 # If the automatic memory scaling is enabled (i.e. the memory multiplier
455 # is set and it is a positive number greater than 1.0), adjust number
456 # of retries when necessary. If the memory multiplier is invalid, disable
457 # automatic memory scaling.
458 if job_values["memory_multiplier"] is not None:
459 if math.ceil(float(job_values["memory_multiplier"])) > 1:
460 if job_values["number_of_retries"] is None: 460 ↛ 465line 460 didn't jump to line 465 because the condition on line 460 was always true
461 job_values["number_of_retries"] = DEFAULT_MEM_RETRIES
462 else:
463 job_values["memory_multiplier"] = None
465 if cmd_line_key:
466 found, cmdline = config.search(cmd_line_key, opt=search_opt)
467 # Make sure cmdline isn't None as that could be sent in as a
468 # default value in search_opt.
469 if found and cmdline:
470 cmd, args = cmdline.split(" ", 1)
471 job_values["executable"] = GenericWorkflowExec(os.path.basename(cmd), cmd, False)
472 if args: 472 ↛ 475line 472 didn't jump to line 475 because the condition on line 472 was always true
473 job_values["arguments"] = args
475 return job_values
478def _handle_job_values(quantum_job_values, gwjob, attributes=_ATTRS_ALL):
479 """Set the job attributes in the cluster to their correct values.
481 Parameters
482 ----------
483 quantum_job_values : `dict` [`str`, Any]
484 Job values for running single Quantum.
485 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
486 Generic workflow job in which to store the universal values.
487 attributes : `~collections.abc.Iterable` [`str`], optional
488 Job attributes to be set in the job following different rules.
489 The default value is _ATTRS_ALL.
490 """
491 _LOG.debug("Call to _handle_job_values")
492 _handle_job_values_universal(quantum_job_values, gwjob, attributes)
493 _handle_job_values_max(quantum_job_values, gwjob, attributes)
494 _handle_job_values_sum(quantum_job_values, gwjob, attributes)
497def _handle_job_values_universal(quantum_job_values, gwjob, attributes=_ATTRS_UNIVERSAL):
498 """Handle job attributes that must have the same value for every quantum
499 in the cluster.
501 Parameters
502 ----------
503 quantum_job_values : `dict` [`str`, Any]
504 Job values for running single Quantum.
505 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
506 Generic workflow job in which to store the universal values.
507 attributes : `~collections.abc.Iterable` [`str`], optional
508 Job attributes to be set in the job following different rules.
509 The default value is _ATTRS_UNIVERSAL.
510 """
511 for attr in _ATTRS_UNIVERSAL & set(attributes):
512 _LOG.debug(
513 "Handling job %s (job=%s, quantum=%s)",
514 attr,
515 getattr(gwjob, attr),
516 quantum_job_values.get(attr, "MISSING"),
517 )
518 current_value = getattr(gwjob, attr)
519 try:
520 quantum_value = quantum_job_values[attr]
521 except KeyError:
522 continue
523 else:
524 if not current_value:
525 setattr(gwjob, attr, quantum_value)
526 elif current_value != quantum_value: 526 ↛ 527line 526 didn't jump to line 527 because the condition on line 526 was never true
527 _LOG.error(
528 "Inconsistent value for %s in Cluster %s Quantum Number %s\n"
529 "Current cluster value: %s\n"
530 "Quantum value: %s",
531 attr,
532 gwjob.name,
533 quantum_job_values.get("qgraphNodeId", "MISSING"),
534 current_value,
535 quantum_value,
536 )
537 raise RuntimeError(f"Inconsistent value for {attr} in cluster {gwjob.name}.")
540def _handle_job_values_max(quantum_job_values, gwjob, attributes=_ATTRS_MAX):
541 """Handle job attributes that should be set to their maximum value in
542 the in cluster.
544 Parameters
545 ----------
546 quantum_job_values : `dict` [`str`, `~typing.Any`]
547 Job values for running single Quantum.
548 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
549 Generic workflow job in which to store the aggregate values.
550 attributes : `~collections.abc.Iterable` [`str`], optional
551 Job attributes to be set in the job following different rules.
552 The default value is _ATTR_MAX.
553 """
554 for attr in _ATTRS_MAX & set(attributes):
555 current_value = getattr(gwjob, attr)
556 try:
557 quantum_value = quantum_job_values[attr]
558 except KeyError:
559 continue
560 else:
561 needs_update = False
562 if current_value is None: 562 ↛ 566line 562 didn't jump to line 566 because the condition on line 562 was always true
563 if quantum_value is not None: 563 ↛ 564line 563 didn't jump to line 564 because the condition on line 563 was never true
564 needs_update = True
565 else:
566 if quantum_value is not None and current_value < quantum_value:
567 needs_update = True
568 if needs_update: 568 ↛ 569line 568 didn't jump to line 569 because the condition on line 568 was never true
569 setattr(gwjob, attr, quantum_value)
571 # When updating memory requirements for a job, check if memory
572 # autoscaling is enabled. If it is, always use the memory
573 # multiplier and the number of retries which comes with the
574 # quantum.
575 #
576 # Note that as a result, the quantum with the biggest memory
577 # requirements will determine whether the memory autoscaling
578 # will be enabled (or disabled) depending on the value of its
579 # memory multiplier.
580 if attr == "request_memory":
581 gwjob.memory_multiplier = quantum_job_values["memory_multiplier"]
582 if gwjob.memory_multiplier is not None:
583 gwjob.number_of_retries = quantum_job_values["number_of_retries"]
586def _handle_job_values_sum(quantum_job_values, gwjob, attributes=_ATTRS_SUM):
587 """Handle job attributes that are the sum of their values in the cluster.
589 Parameters
590 ----------
591 quantum_job_values : `dict` [`str`, `~typing.Any`]
592 Job values for running single Quantum.
593 gwjob : `lsst.ctrl.bps.GenericWorkflowJob`
594 Generic workflow job in which to store the aggregate values.
595 attributes : `~collections.abc.Iterable` [`str`], optional
596 Job attributes to be set in the job following different rules.
597 The default value is _ATTRS_SUM.
598 """
599 for attr in _ATTRS_SUM & set(attributes):
600 current_value = getattr(gwjob, attr)
601 if not current_value: 601 ↛ 604line 601 didn't jump to line 604 because the condition on line 601 was always true
602 setattr(gwjob, attr, quantum_job_values[attr])
603 else:
604 setattr(gwjob, attr, current_value + quantum_job_values[attr])
607def create_generic_workflow(
608 config: BpsConfig, cqgraph: ClusteredQuantumGraph, name: str, prefix: str
609) -> GenericWorkflow:
610 """Create a generic workflow from a ClusteredQuantumGraph such that it
611 has information needed for WMS (e.g., command lines).
613 Parameters
614 ----------
615 config : `lsst.ctrl.bps.BpsConfig`
616 BPS configuration.
617 cqgraph : `lsst.ctrl.bps.ClusteredQuantumGraph`
618 ClusteredQuantumGraph for running a specific pipeline on a specific
619 payload.
620 name : `str`
621 Name for the workflow (typically unique).
622 prefix : `str`
623 Root path for any output files.
625 Returns
626 -------
627 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
628 Generic workflow for the given ClusteredQuantumGraph + config.
629 """
630 # Determine whether saving per-job QuantumGraph files in the loop.
631 _, when_save = config.search("whenSaveJobQgraph", {"default": WhenToSaveQuantumGraphs.TRANSFORM.name})
632 save_qgraph_per_job = WhenToSaveQuantumGraphs[when_save.upper()]
634 search_opt = {"replaceVars": False, "expandEnvVars": False, "replaceEnvVars": True, "required": False}
636 generic_workflow = GenericWorkflow(name)
638 # Save full run QuantumGraph for use by jobs
639 generic_workflow.add_file(
640 GenericWorkflowFile(
641 "runQgraphFile",
642 src_uri=config["runQgraphFile"],
643 wms_transfer=True,
644 job_access_remote=True,
645 job_shared=True,
646 )
647 )
649 # Cache pipetask specific or more generic job values to minimize number
650 # on config searches.
651 cached_job_values = {}
652 cached_pipetask_values = {}
654 for cluster in cqgraph.clusters():
655 _LOG.debug("Loop over clusters: %s, %s", cluster, type(cluster))
656 _LOG.debug(
657 "cqgraph: name=%s, len=%s, label=%s, ids=%s",
658 cluster.name,
659 len(cluster.qgraph_node_ids),
660 cluster.label,
661 cluster.qgraph_node_ids,
662 )
664 gwjob = GenericWorkflowJob(cluster.name, cluster.label)
666 # First get job values from cluster or cluster config
667 search_opt["curvals"] = {"curr_cluster": cluster.label}
668 found, value = config.search("computeSite", opt=search_opt)
669 if found: 669 ↛ 671line 669 didn't jump to line 671 because the condition on line 669 was always true
670 search_opt["curvals"]["curr_site"] = value
671 found, value = config.search("computeCloud", opt=search_opt)
672 if found:
673 search_opt["curvals"]["curr_cloud"] = value
675 # If some config values are set for this cluster
676 if cluster.label not in cached_job_values:
677 _LOG.debug("config['cluster'][%s] = %s", cluster.label, config["cluster"][cluster.label])
678 cached_job_values[cluster.label] = {}
680 # Allowing whenSaveJobQgraph and useLazyCommands per cluster label.
681 key = "whenSaveJobQgraph"
682 _, when_save = config.search(key, opt=search_opt)
683 cached_job_values[cluster.label][key] = WhenToSaveQuantumGraphs[when_save.upper()]
685 key = "useLazyCommands"
686 search_opt["default"] = True
687 _, cached_job_values[cluster.label][key] = config.search(key, opt=search_opt)
688 del search_opt["default"]
690 if cluster.label in config["cluster"]: 690 ↛ 693line 690 didn't jump to line 693 because the condition on line 690 was never true
691 # Don't want to get global defaults here so only look in
692 # cluster section.
693 cached_job_values[cluster.label].update(
694 _get_job_values(config["cluster"][cluster.label], search_opt, "runQuantumCommand")
695 )
696 cluster_job_values = copy.copy(cached_job_values[cluster.label])
698 cluster_job_values["name"] = cluster.name
699 cluster_job_values["label"] = cluster.label
700 cluster_job_values["quanta_counts"] = cluster.quanta_counts
701 cluster_job_values["tags"] = cluster.tags
702 _LOG.debug("cluster_job_values = %s", cluster_job_values)
703 _handle_job_values(cluster_job_values, gwjob, cluster_job_values.keys())
705 # For purposes of whether to continue searching for a value is whether
706 # the value evaluates to False.
707 unset_attributes = {attr for attr in _ATTRS_ALL if not getattr(gwjob, attr)}
709 _LOG.debug("unset_attributes=%s", unset_attributes)
710 _LOG.debug("set=%s", _ATTRS_ALL - unset_attributes)
712 # For job info not defined at cluster level, attempt to get job info
713 # either common or aggregate for all Quanta in cluster.
714 for node_id in iter(cluster.qgraph_node_ids):
715 _LOG.debug("node_id=%s", node_id)
716 quantum_info = cqgraph.get_quantum_info(node_id)
718 task_label = quantum_info["task_label"]
719 if task_label not in cached_pipetask_values:
720 search_opt["curvals"]["curr_pipetask"] = task_label
721 cached_pipetask_values[task_label] = _get_job_values(config, search_opt, "runQuantumCommand")
722 _handle_job_values(cached_pipetask_values[task_label], gwjob, unset_attributes)
724 # Update job with workflow attribute and profile values.
725 qgraph_gwfile = _get_qgraph_gwfile(
726 config, save_qgraph_per_job, gwjob, generic_workflow.get_file("runQgraphFile"), prefix
727 )
729 generic_workflow.add_job(gwjob)
730 generic_workflow.add_job_inputs(gwjob.name, [qgraph_gwfile])
732 gwjob.cmdvals["qgraphNodeId"] = ",".join(
733 sorted([f"{node_id}" for node_id in cluster.qgraph_node_ids])
734 )
735 _enhance_command(config, generic_workflow, gwjob, cached_job_values)
737 # If writing per-job QuantumGraph files during TRANSFORM stage,
738 # write it now while in memory.
739 if save_qgraph_per_job == WhenToSaveQuantumGraphs.TRANSFORM: 739 ↛ 740line 739 didn't jump to line 740 because the condition on line 739 was never true
740 save_qg_subgraph(cqgraph.qgraph, qgraph_gwfile.src_uri, cluster.qgraph_node_ids)
742 # Create job dependencies.
743 for parent in cqgraph.clusters():
744 for child in cqgraph.successors(parent):
745 generic_workflow.add_job_relationships(parent.name, child.name)
747 # Add initial workflow.
748 if config.get("runInit", "{default: False}"): 748 ↛ 751line 748 didn't jump to line 751 because the condition on line 748 was always true
749 add_workflow_init_nodes(config, cqgraph.qgraph, generic_workflow)
751 generic_workflow.run_attrs.update(
752 {
753 "bps_isjob": "True",
754 "bps_project": config["project"],
755 "bps_campaign": config["campaign"],
756 "bps_run": generic_workflow.name,
757 "bps_operator": config["operator"],
758 "bps_payload": config["payloadName"],
759 "bps_runsite": config["computeSite"],
760 }
761 )
763 # Add final job
764 add_final_job(config, generic_workflow, prefix)
766 if "ordering" in config: 766 ↛ 767line 766 didn't jump to line 767 because the condition on line 766 was never true
767 generic_workflow.add_special_job_ordering(config["ordering"])
769 return generic_workflow
772def create_generic_workflow_config(config, prefix):
773 """Create generic workflow configuration.
775 Parameters
776 ----------
777 config : `lsst.ctrl.bps.BpsConfig`
778 Bps configuration.
779 prefix : `str`
780 Root path for any output files.
782 Returns
783 -------
784 generic_workflow_config : `lsst.ctrl.bps.BpsConfig`
785 Configuration accompanying the GenericWorkflow.
786 """
787 generic_workflow_config = BpsConfig(config)
788 generic_workflow_config["workflowName"] = config["uniqProcName"]
789 generic_workflow_config["workflowPath"] = prefix
790 return generic_workflow_config
793def add_final_job(config: BpsConfig, generic_workflow: GenericWorkflow, prefix: str) -> None:
794 """Add final workflow job depending upon configuration.
796 Depending on configuration, the final job will be added as a special job
797 which will always run regardless of the exit status of the workflow or
798 a regular sink node which will only run if the workflow execution finished
799 with no errors.
801 Parameters
802 ----------
803 config : `lsst.ctrl.bps.BpsConfig`
804 Bps configuration.
805 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
806 Generic workflow to which attributes should be added.
807 prefix : `str`
808 Directory in which to output final script.
809 """
810 _, when_run = config.search(".finalJob.whenRun")
811 if when_run.upper() != "NEVER": 811 ↛ exitline 811 didn't return from function 'add_final_job' because the condition on line 811 was always true
812 gwjob = create_final_job(config, generic_workflow, prefix)
813 if when_run.upper() == "ALWAYS": 813 ↛ 815line 813 didn't jump to line 815 because the condition on line 813 was always true
814 generic_workflow.add_final(gwjob)
815 elif when_run.upper() == "SUCCESS":
816 add_final_job_as_sink(generic_workflow, gwjob)
817 else:
818 raise ValueError(f"Invalid value for finalJob.whenRun: {when_run}")
821def create_final_job(config: BpsConfig, generic_workflow: GenericWorkflow, prefix: str) -> GenericWorkflowJob:
822 """Create the final workflow job.
824 Parameters
825 ----------
826 config : `lsst.ctrl.bps.BpsConfig`
827 Bps configuration.
828 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
829 Generic workflow to which attributes should be added.
830 prefix : `str`
831 Directory in which to output final script.
833 Returns
834 -------
835 final_job : `lsst.ctrl.bps.GenericWorkflowJob`
836 Final workflow job.
837 """
838 job_name = "finalJob"
839 gwjob = GenericWorkflowJob(job_name, job_name)
841 search_opt = {"searchobj": config[job_name], "curvals": {}, "default": None}
842 found, value = config.search("computeSite", opt=search_opt)
843 if found: 843 ↛ 845line 843 didn't jump to line 845 because the condition on line 843 was always true
844 search_opt["curvals"]["curr_site"] = value
845 found, value = config.search("computeCloud", opt=search_opt)
846 if found: 846 ↛ 855line 846 didn't jump to line 855 because the condition on line 846 was always true
847 search_opt["curvals"]["curr_cloud"] = value
849 # Set job attributes based on the values find in the config excluding
850 # the ones in the _ATTRS_MISC group. The attributes in this group are
851 # somewhat "special":
852 # * HTCondor plugin, which uses 'attrs' and 'profile', has its own
853 # mechanism for setting them,
854 # * 'cmdvals' is being set internally, not via config.
855 job_values = _get_job_values(config, search_opt, None)
856 for attr in _ATTRS_ALL - _ATTRS_MISC:
857 if not getattr(gwjob, attr) and job_values.get(attr, None):
858 setattr(gwjob, attr, job_values[attr])
860 # Create script and add command line to job.
861 gwjob.executable, gwjob.arguments = create_final_command(config, prefix)
863 # Determine inputs from command line.
864 for file_key in re.findall(r"<FILE:([^>]+)>", gwjob.arguments):
865 gwfile = generic_workflow.get_file(file_key)
866 generic_workflow.add_job_inputs(gwjob.name, gwfile)
868 _enhance_command(config, generic_workflow, gwjob, {})
869 return gwjob
872def create_final_command(config: BpsConfig, prefix: str) -> tuple[GenericWorkflowExec, str]:
873 """Create the command and shell script for the final job.
875 Parameters
876 ----------
877 config : `lsst.ctrl.bps.BpsConfig`
878 Bps configuration.
879 prefix : `str`
880 Directory in which to output final script.
882 Returns
883 -------
884 executable : `lsst.ctrl.bps.GenericWorkflowExec`
885 Executable object for the final script.
886 arguments : `str`
887 Command line needed to call the final script.
889 Raises
890 ------
891 RuntimeError if no commands found.
892 """
893 search_opt = {
894 "replaceVars": True,
895 "skipNames": ["butlerConfig", "qgraphFile"],
896 "replaceEnvVars": False,
897 "expandEnvVars": False,
898 "searchobj": config["finalJob"],
899 }
901 script_file = os.path.join(prefix, "final_job.bash")
902 with open(script_file, "w", encoding="utf8") as fh:
903 print("#!/bin/bash\n", file=fh)
904 print("set -e", file=fh)
905 print("set -x", file=fh)
907 print("qgraphFile=$1", file=fh)
908 print("butlerConfig=$2", file=fh)
910 command_len = 0 # Make sure at least write one actual command
911 i = 1
912 found, command = config.search(f"command{i}", opt=search_opt)
913 while found:
914 # The files will be args to script, so change to shell vars
915 command = command.replace("{qgraphFile}", "${qgraphFile}")
916 command = command.replace("{butlerConfig}", "${butlerConfig}")
918 print(command, file=fh)
919 command_len += len(command.strip())
921 # Search for next command
922 i += 1
923 found, command = config.search(f"command{i}", opt=search_opt)
924 if command_len == 0:
925 raise RuntimeError(
926 "No finalJob commands were found. Use NEVER for finalJob.whenRun to turn off finalJob"
927 )
928 os.chmod(script_file, 0o755)
929 executable = GenericWorkflowExec(os.path.basename(script_file), script_file, True)
931 _, orig_butler = config.search("butlerConfig")
932 return executable, f"<FILE:runQgraphFile> {orig_butler}"
935def add_final_job_as_sink(generic_workflow, final_job):
936 """Add final job as the single sink for the workflow.
938 Parameters
939 ----------
940 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
941 Generic workflow to which attributes should be added.
942 final_job : `lsst.ctrl.bps.GenericWorkflowJob`
943 Job to add as new sink node depending upon all previous sink nodes.
944 """
945 # Find sink nodes of generic workflow graph.
946 gw_sinks = [n for n in generic_workflow if generic_workflow.out_degree(n) == 0]
947 _LOG.debug("gw_sinks = %s", gw_sinks)
949 generic_workflow.add_job(final_job)
950 generic_workflow.add_job_relationships(gw_sinks, final_job.name)