Coverage for python/lsst/ctrl/bps/batch_submit.py: 99%
121 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-14 09:34 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-14 09:34 +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 to run submit stages as batch jobs."""
30__all__ = ["batch_payload_prepare", "batch_submit", "create_batch_stages"]
32import logging
33import os
34from pathlib import Path
36from lsst.resources import ResourcePath, ResourcePathExpression
37from lsst.utils.logging import VERBOSE
38from lsst.utils.timer import time_this, timeMethod
40from . import (
41 DEFAULT_MEM_FMT,
42 DEFAULT_MEM_UNIT,
43 BpsConfig,
44 GenericWorkflow,
45 GenericWorkflowFile,
46 GenericWorkflowJob,
47 GenericWorkflowLazyGroup,
48)
49from .bps_utils import _make_id_link
50from .pre_transform import cluster_quanta, read_quantum_graph
51from .prepare import prepare
52from .submit import submit
53from .transform import _enhance_command, _get_job_values, transform
55_LOG = logging.getLogger(__name__)
58@timeMethod(logger=_LOG, logLevel=VERBOSE)
59def create_batch_stages(
60 config: BpsConfig, prefix: ResourcePathExpression
61) -> tuple[GenericWorkflow, BpsConfig]:
62 """Create a GenericWorkflow that performs the submit stages as a workflow.
64 Parameters
65 ----------
66 config : `lsst.ctrl.bps.BpsConfig`
67 BPS configuration.
68 prefix : `lsst.resources.ResourcePathExpression`
69 Root path for any output files.
71 Returns
72 -------
73 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
74 The generic workflow transformed from the clustered quantum graph.
75 generic_workflow_config : `lsst.ctrl.bps.BpsConfig`
76 Configuration to accompany GenericWorkflow.
77 """
78 prefix = ResourcePath(prefix)
79 generic_workflow: GenericWorkflow = GenericWorkflow(name=f"{config['uniqProcName']}_ctrl")
80 generic_workflow.run_attrs.update(
81 {
82 "bps_isjob": "True",
83 "bps_project": config["project"],
84 "bps_campaign": config["campaign"],
85 "bps_run": config["uniqProcName"],
86 "bps_runsite": config["computeSite"],
87 "bps_operator": config["operator"],
88 "bps_payload": config["payloadName"],
89 }
90 )
92 # Save full run QuantumGraph for use by jobs
93 qgraph_file = GenericWorkflowFile(
94 "runQgraphFile",
95 src_uri=config["runQgraphFile"],
96 wms_transfer=True,
97 job_access_remote=True,
98 job_shared=True,
99 )
100 generic_workflow.add_file(qgraph_file)
102 # Save config file for use by jobs
103 config_file = GenericWorkflowFile(
104 "configFile",
105 src_uri=config["configFile"],
106 wms_transfer=True,
107 job_access_remote=False,
108 job_shared=True,
109 )
110 generic_workflow.add_file(config_file)
112 # Build QuantumGraph job
113 build_job = GenericWorkflowJob(
114 name="buildQuantumGraph",
115 label="buildQuantumGraph",
116 )
117 search_opt = config.get_search_opts(build_job.label)
118 search_opt.update(
119 {
120 "replaceVars": False,
121 "expandEnvVars": False,
122 "replaceEnvVars": True,
123 "required": False,
124 }
125 )
126 cmd_line_key = "jobCommand"
127 job_values = _get_job_values(config, search_opt, cmd_line_key)
128 if not job_values["executable"]:
129 raise RuntimeError(
130 f"Missing executable for buildQuantumGraph. Double check submit yaml for {cmd_line_key}"
131 )
132 for key, value in job_values.items():
133 if key not in {"name", "label"}:
134 setattr(build_job, key, value)
136 generic_workflow.add_job(build_job)
137 _LOG.debug("build job's arguments: %s", build_job.arguments)
139 generic_workflow.add_file(config_file)
140 generic_workflow.add_job_outputs(build_job.name, [qgraph_file])
141 generic_workflow.add_job_inputs(build_job.name, [config_file])
142 _enhance_command(config, generic_workflow, build_job, job_values)
143 _LOG.debug("build job's arguments: %s", build_job.arguments)
145 # Build cluster/transform/prepare job
146 prepare_job = GenericWorkflowLazyGroup(
147 name="preparePayloadWorkflow",
148 label="preparePayloadWorkflow",
149 )
150 search_opt = config.get_search_opts(prepare_job.label)
151 search_opt.update(
152 {
153 "replaceVars": False,
154 "expandEnvVars": False,
155 "replaceEnvVars": True,
156 "required": False,
157 }
158 )
159 _LOG.debug("preparePayloadWorkflow search_opt = %s", search_opt)
160 cmd_line_key = "jobCommand"
161 job_values = _get_job_values(config, search_opt, cmd_line_key)
162 if not job_values["executable"]:
163 raise RuntimeError(
164 f"Missing executable for preparePayloadWorkflow. Double check submit yaml for {cmd_line_key}"
165 )
166 for key, value in job_values.items():
167 if key not in {"name", "label"}:
168 setattr(prepare_job, key, value)
170 generic_workflow.add_job(prepare_job, parent_names=["buildQuantumGraph"])
171 generic_workflow.add_job_inputs(prepare_job.name, [qgraph_file, config_file])
172 _enhance_command(config, generic_workflow, prepare_job, job_values)
174 _, save_workflow = config.search("saveGenericWorkflow", opt={"default": False})
175 if save_workflow:
176 with prefix.join("bps_stages_generic_workflow.pickle").open("wb") as outfh:
177 generic_workflow.save(outfh, "pickle")
179 return generic_workflow, config
182@timeMethod(logger=_LOG, logLevel=VERBOSE)
183def batch_payload_prepare(config: BpsConfig, prefix: ResourcePathExpression) -> None:
184 """Create a GenericWorkflow that performs the submit stages as a workflow.
186 Parameters
187 ----------
188 config : `lsst.ctrl.bps.BpsConfig`
189 BPS configuration.
190 prefix : `lsst.resources.ResourcePathExpression`
191 Root path for any output files.
193 Returns
194 -------
195 generic_workflow : `lsst.ctrl.bps.GenericWorkflow`
196 The generic workflow transformed from the clustered quantum graph.
197 generic_workflow_config : `lsst.ctrl.bps.BpsConfig`
198 Configuration to accompany GenericWorkflow.
199 """
200 # Read existing QuantumGraph
201 qgraph_uri = config["runQgraphFile"]
202 qgraph = read_quantum_graph(qgraph_uri)
204 # Cluster
205 _LOG.info("Starting cluster stage (grouping quanta into jobs)")
206 with time_this(
207 log=_LOG,
208 level=logging.INFO,
209 prefix=None,
210 msg="Cluster stage completed",
211 mem_usage=True,
212 mem_unit=DEFAULT_MEM_UNIT,
213 mem_fmt=DEFAULT_MEM_FMT,
214 ):
215 clustered_qgraph = cluster_quanta(config, qgraph, config["uniqProcName"])
217 _LOG.info("ClusteredQuantumGraph contains %d cluster(s)", len(clustered_qgraph))
219 submit_path = config[".bps_defined.submitPath"]
220 _, save_clustered_qgraph = config.search("saveClusteredQgraph", opt={"default": False})
221 if save_clustered_qgraph:
222 clustered_qgraph.save(os.path.join(submit_path, "bps_clustered_qgraph.pickle"))
223 _, save_dot = config.search("saveDot", opt={"default": False})
224 if save_dot:
225 clustered_qgraph.draw(os.path.join(submit_path, "bps_clustered_qgraph.dot"))
227 # Transform
228 _LOG.info("Starting transform stage (creating generic workflow)")
229 with time_this(
230 log=_LOG,
231 level=logging.INFO,
232 prefix=None,
233 msg="Transform stage completed",
234 mem_usage=True,
235 mem_unit=DEFAULT_MEM_UNIT,
236 mem_fmt=DEFAULT_MEM_FMT,
237 ):
238 generic_workflow, generic_workflow_config = transform(config, clustered_qgraph, submit_path)
239 _LOG.info("Generic workflow name '%s'", generic_workflow.name)
241 num_jobs = sum(generic_workflow.job_counts.values())
242 _LOG.info("GenericWorkflow contains %d job(s) (including final)", num_jobs)
244 # Want to read quantum graph from temp space if told to use it.
245 gwfile = generic_workflow.get_file("runQgraphFile")
246 gwfile.wms_transfer = False
247 # Root computeSite is currently how to specify where the pipeline
248 # will actually be run.
249 found, use_run_temp_space = config.search(
250 "bpsUseRunTempSpace", opt={"curvals": {"curr_site": config[".computeSite"]}}
251 )
252 if found and use_run_temp_space:
253 found, run_temp_space = config.search(
254 "fileDistributionEndpoint", opt={"curvals": {"curr_site": config[".computeSite"]}}
255 )
256 if found:
257 _LOG.debug("run_temp_space = %s", run_temp_space)
258 gwfile.src_uri = str(Path(run_temp_space) / config["qgraphFileTemplate"])
259 else:
260 raise KeyError("Config is missing fileDistributionEndpoint.")
261 elif not found: 261 ↛ 264line 261 didn't jump to line 264 because the condition on line 261 was always true
262 _LOG.debug("Config is missing bpsUseRunTempSpace")
264 _, save_workflow = config.search("saveGenericWorkflow", opt={"default": False})
265 if save_workflow:
266 with open(os.path.join(submit_path, "bps_generic_workflow.pickle"), "wb") as outfh:
267 generic_workflow.save(outfh, "pickle")
268 _, save_dot = config.search("saveDot", opt={"default": False})
269 if save_dot:
270 with open(os.path.join(submit_path, "bps_generic_workflow.dot"), "w") as outfh:
271 generic_workflow.draw(outfh, "dot")
273 # Prepare
274 _LOG.info("Starting prepare stage (creating specific implementation of workflow)")
275 with time_this(
276 log=_LOG,
277 level=logging.INFO,
278 prefix=None,
279 msg="Prepare stage completed",
280 mem_usage=True,
281 mem_unit=DEFAULT_MEM_UNIT,
282 mem_fmt=DEFAULT_MEM_FMT,
283 ):
284 wms_workflow = prepare(generic_workflow_config, generic_workflow, submit_path)
286 # Add payload workflow to currently running workflow
287 _LOG.info("Starting update workflow")
288 with time_this(
289 log=_LOG,
290 level=logging.INFO,
291 prefix=None,
292 msg="Workflow update completed",
293 mem_usage=True,
294 mem_unit=DEFAULT_MEM_UNIT,
295 mem_fmt=DEFAULT_MEM_FMT,
296 ):
297 # Assuming submit_path for ctrl workflow is visible by this job.
298 wms_workflow.add_to_parent_workflow(generic_workflow_config)
301def batch_submit(config: BpsConfig):
302 """Submit a workflow for execution with preparation done in batch jobs.
304 Parameters
305 ----------
306 config : `lsst.ctrl.bps.BpsConfig`
307 BPS configuration.
309 Returns
310 -------
311 wms_workflow : `lsst.ctrl.bps.BaseWmsWorkflow`
312 Submitted workflow.
313 """
314 submit_path = config[".bps_defined.submitPath"]
316 _LOG.info("Starting to create control workflow")
317 with time_this(
318 log=_LOG,
319 level=logging.INFO,
320 prefix=None,
321 msg="Creation completed",
322 mem_usage=True,
323 mem_unit=DEFAULT_MEM_UNIT,
324 mem_fmt=DEFAULT_MEM_FMT,
325 ):
326 generic_workflow, config = create_batch_stages(config, submit_path)
328 _LOG.info("Starting to prepare control workflow")
329 with time_this(
330 log=_LOG,
331 level=logging.INFO,
332 prefix=None,
333 msg="Preparation completed",
334 mem_usage=True,
335 mem_unit=DEFAULT_MEM_UNIT,
336 mem_fmt=DEFAULT_MEM_FMT,
337 ):
338 wms_workflow = prepare(config, generic_workflow, submit_path)
340 _, dry_run = config.search("dryRun", opt={"default": False})
341 if not dry_run:
342 _LOG.info("Starting to submit control workflow")
343 with time_this(
344 log=_LOG,
345 level=logging.INFO,
346 prefix=None,
347 msg="Submission completed",
348 mem_usage=True,
349 mem_unit=DEFAULT_MEM_UNIT,
350 mem_fmt=DEFAULT_MEM_FMT,
351 ):
352 submit(config, wms_workflow)
353 _LOG.info("Run '%s' submitted for execution with id '%s'", wms_workflow.name, wms_workflow.run_id)
355 _make_id_link(config, wms_workflow.run_id)
357 return wms_workflow