Coverage for python/lsst/ctrl/bps/pre_transform.py: 99%
98 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-09-06 01:54 -0700
« prev ^ index » next coverage.py v7.15.4, created at 2026-09-06 01:54 -0700
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 execute steps outside of BPS that need to be done first
29including running QuantumGraph generation and reading the QuantumGraph
30into memory.
31"""
33import logging
34import os
35import shlex
36import shutil
37import subprocess
38from pathlib import Path
40from lsst.ctrl.bps import BpsConfig, BpsSubprocessError
41from lsst.pipe.base import QuantumGraph
42from lsst.pipe.base.pipeline_graph import TaskImportMode
43from lsst.pipe.base.quantum_graph import PredictedQuantumGraph
44from lsst.resources import ResourcePath, ResourcePathExpression
45from lsst.utils import doImport
46from lsst.utils.logging import VERBOSE
47from lsst.utils.timer import time_this, timeMethod
49_LOG = logging.getLogger(__name__)
52@timeMethod(logger=_LOG, logLevel=VERBOSE)
53def acquire_quantum_graph(config: BpsConfig, out_prefix: str = "") -> str:
54 """Read a quantum graph from a file or create one from scratch.
56 Parameters
57 ----------
58 config : `lsst.ctrl.bps.BpsConfig`
59 Configuration values for BPS. In particular, looking for qgraphFile.
60 out_prefix : `str`, optional
61 Output path for the QuantumGraph and stdout/stderr from generating
62 the QuantumGraph. Default value is empty string.
64 Returns
65 -------
66 qgraph_filename : `str`
67 Name of file containing the quantum graph.
68 """
69 # Check to see if user provided pre-generated QuantumGraph.
70 found, input_qgraph_filename = config.search("qgraphFile")
71 if found and input_qgraph_filename:
72 if out_prefix is not None:
73 # Save a copy of the QuantumGraph file in out_prefix.
74 _LOG.info("Copying quantum graph from '%s'", input_qgraph_filename)
75 with time_this(log=_LOG, level=logging.INFO, prefix=None, msg="Completed copying quantum graph"):
76 qgraph_filename = os.path.join(out_prefix, os.path.basename(input_qgraph_filename))
77 shutil.copy2(input_qgraph_filename, qgraph_filename)
78 else:
79 # Use QuantumGraph file in original given location.
80 qgraph_filename = input_qgraph_filename
82 # Update the output run in the user provided quantum graph.
83 if "finalJob" in config:
84 update_quantum_graph(config, qgraph_filename, out_prefix)
85 else:
86 # Run command to create the QuantumGraph.
87 _LOG.info("Creating quantum graph")
88 with time_this(log=_LOG, level=logging.INFO, prefix=None, msg="Completed creating quantum graph"):
89 qgraph_filename = create_quantum_graph(config, out_prefix)
91 return qgraph_filename
94@timeMethod(logger=_LOG, logLevel=VERBOSE)
95def read_quantum_graph(qgraph_filename: ResourcePathExpression) -> PredictedQuantumGraph:
96 """Read a quantum graph from a file.
98 Parameters
99 ----------
100 qgraph_filename : `lsst.resources.ResourcePathExpression`
101 Name of file containing PredictedQuantumGraph to be read.
103 Returns
104 -------
105 qgraph : `lsst.pipe.base.quantum_graph.PredictedQuantumGraph`
106 A quantum graph read in from pre-generated file or one that is the
107 result of running code that generates it.
108 """
109 _LOG.info("Reading quantum graph from '%s'", qgraph_filename)
110 with time_this(log=_LOG, level=logging.INFO, prefix=None, msg="Completed reading quantum graph"):
111 qgraph_path = ResourcePath(qgraph_filename)
112 if qgraph_path.getExtension() == ".qg":
113 with PredictedQuantumGraph.open(qgraph_path, import_mode=TaskImportMode.DO_NOT_IMPORT) as reader:
114 reader.read_thin_graph()
115 qgraph = reader.finish()
116 elif qgraph_path.getExtension() == ".qgraph":
117 qgraph = PredictedQuantumGraph.from_old_quantum_graph(QuantumGraph.loadUri(qgraph_path))
118 else:
119 raise ValueError(f"Unrecognized extension for quantum graph file: {qgraph_filename}.")
120 return qgraph
123def execute(command: str, filename: str, write_buffering: int = 1) -> int:
124 """Execute a command.
126 Parameters
127 ----------
128 command : `str`
129 String representing the command to execute.
130 filename : `str`
131 A file to which both stderr and stdout will be written to.
132 write_buffering : `int`, optional
133 Buffering policy passed to open for the stdout/stderr file.
134 0 - not allowed here because writing in text mode.
135 1 - line buffering (default).
136 > 1 - size in bytes for a chunk buffer.
138 Returns
139 -------
140 exit_code : `int`
141 The exit code the command being executed finished with.
143 Raises
144 ------
145 ValueError
146 Raised if write_buffering is 0.
147 """
148 buffer_size = 5000
149 with open(filename, "w", write_buffering) as fh:
150 print(command, file=fh)
151 print("\n", file=fh) # Note: want a blank line
152 process = subprocess.Popen(
153 shlex.split(command), shell=False, stdout=subprocess.PIPE, stderr=subprocess.STDOUT
154 )
155 assert process.stdout is not None # for mypy
156 buffer = os.read(process.stdout.fileno(), buffer_size).decode()
157 while process.poll is None or buffer:
158 stripped_buffer = buffer.rstrip()
159 print(stripped_buffer, file=fh)
160 _LOG.info(stripped_buffer)
161 buffer = os.read(process.stdout.fileno(), buffer_size).decode()
162 process.stdout.close()
163 process.wait()
164 return process.returncode
167def create_quantum_graph(config: BpsConfig, out_prefix: str = "") -> str:
168 """Create QuantumGraph from pipeline definition.
170 Parameters
171 ----------
172 config : `lsst.ctrl.bps.BpsConfig`
173 BPS configuration.
174 out_prefix : `str`, optional
175 Path in which to output QuantumGraph as well as the stdout/stderr
176 from generating the QuantumGraph. Defaults to empty string so
177 code will write the QuantumGraph and stdout/stderr to the current
178 directory.
180 Returns
181 -------
182 qgraph_filename : `str`
183 Name of file containing generated QuantumGraph.
185 Raises
186 ------
187 KeyError
188 Raised if the command for generating the QuantumGraph was not found
189 in the provided configuration.
190 BpsSubprocessError
191 Raised if the command for generating the QuantumGraph failed.
192 """
193 found, qgraph_filename = config.search(".bps_defined.runQgraphFile")
194 if not found: 194 ↛ 199line 194 didn't jump to line 199 because the condition on line 194 was always true
195 # Create name of file to store QuantumGraph.
196 qgraph_filename = os.path.join(out_prefix, config["qgraphFileTemplate"])
198 # Get QuantumGraph generation command.
199 search_opt = {"curvals": {"qgraphFile": qgraph_filename}}
200 found, cmd = config.search("createQuantumGraph", opt=search_opt)
201 if not found:
202 raise KeyError("command for generating QuantumGraph not found")
203 _LOG.info(cmd)
205 # Run QuantumGraph generation.
206 out = os.path.join(out_prefix, "quantumGraphGeneration.out")
207 status = execute(cmd, out)
208 if status != 0:
209 raise BpsSubprocessError(
210 status,
211 f"Generating quantum graph failed with non-zero exit code ({status})\n"
212 f"Check {out} for more details.",
213 )
214 return qgraph_filename
217def update_quantum_graph(
218 config: BpsConfig, qgraph_filename: str, out_prefix: str = "", inplace: bool = False
219) -> None:
220 """Update output run in an existing quantum graph.
222 Parameters
223 ----------
224 config : `BpsConfig`
225 BPS configuration.
226 qgraph_filename : `str`
227 Name of file containing the quantum graph that needs to be updated.
228 out_prefix : `str`, optional
229 Path in which to output QuantumGraph as well as the stdout/stderr
230 from generating the QuantumGraph. Defaults to empty string so
231 code will write the QuantumGraph and stdout/stderr to the current
232 directory.
233 inplace : `bool`, optional
234 If set to True, all updates of the graph will be done in place without
235 creating a backup copy. Defaults to False.
236 """
237 src_qgraph = Path(qgraph_filename)
238 dest_qgraph = Path(qgraph_filename)
240 # If requested, create a backup copy of the quantum graph by adding
241 # '_orig' suffix to its stem (the filename without the extension).
242 if not inplace:
243 _LOG.info("Backing up quantum graph from '%s'", qgraph_filename)
244 src_qgraph = src_qgraph.parent / f"{src_qgraph.stem}_orig{src_qgraph.suffix}"
245 with time_this(log=_LOG, level=logging.INFO, prefix=None, msg="Completed backing up quantum graph"):
246 shutil.copy2(qgraph_filename, src_qgraph)
248 # Get the command for updating the quantum graph.
249 search_opt = {"curvals": {"inputQgraphFile": str(src_qgraph), "qgraphFile": str(dest_qgraph)}}
250 found, cmd = config.search("updateQuantumGraph", opt=search_opt)
251 if not found:
252 raise KeyError("command for updating quantum graph not found")
253 _LOG.info(cmd)
255 # Run the command to update the quantum graph.
256 out = os.path.join(out_prefix, "quantumGraphUpdate.out")
257 status = execute(cmd, out)
258 if status != 0:
259 raise BpsSubprocessError(
260 status,
261 f"Updating quantum graph failed with non-zero exit code ({status})\n"
262 f"Check {out} for more details.",
263 )
266def cluster_quanta(config, qgraph, name):
267 """Call specified function to group quanta into clusters to be run
268 together.
270 Parameters
271 ----------
272 config : `lsst.ctrl.bps.BpsConfig`
273 BPS configuration.
274 qgraph : `lsst.pipe.base.quantum_graph.PredictedQuantumGraph`
275 Original full quantum graph for the run.
276 name : `str`
277 Name for the ClusteredQuantumGraph that will be generated.
279 Returns
280 -------
281 cqgraph : `lsst.ctrl.bps.ClusteredQuantumGraph`
282 Generated ClusteredQuantumGraph.
284 Raises
285 ------
286 RuntimeError
287 If asked to validate and generated ClusteredQuantumGraph fails a test.
288 """
289 cluster_func = doImport(config["clusterAlgorithm"])
290 cqgraph = cluster_func(config, qgraph, name)
291 _, validate = config.search("validateClusteredQgraph", opt={"default": False})
292 if validate:
293 cqgraph.validate()
294 return cqgraph