Coverage for python/lsst/ctrl/bps/pre_transform.py: 99%

98 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-22 10:08 +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/>. 

27 

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""" 

32 

33import logging 

34import os 

35import shlex 

36import shutil 

37import subprocess 

38from pathlib import Path 

39 

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 

48 

49_LOG = logging.getLogger(__name__) 

50 

51 

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. 

55 

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. 

63 

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 

81 

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) 

90 

91 return qgraph_filename 

92 

93 

94@timeMethod(logger=_LOG, logLevel=VERBOSE) 

95def read_quantum_graph(qgraph_filename: ResourcePathExpression) -> PredictedQuantumGraph: 

96 """Read a quantum graph from a file. 

97 

98 Parameters 

99 ---------- 

100 qgraph_filename : `lsst.resources.ResourcePathExpression` 

101 Name of file containing PredictedQuantumGraph to be read. 

102 

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 

121 

122 

123def execute(command: str, filename: str, write_buffering: int = 1) -> int: 

124 """Execute a command. 

125 

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. 

137 

138 Returns 

139 ------- 

140 exit_code : `int` 

141 The exit code the command being executed finished with. 

142 

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 

165 

166 

167def create_quantum_graph(config: BpsConfig, out_prefix: str = "") -> str: 

168 """Create QuantumGraph from pipeline definition. 

169 

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. 

179 

180 Returns 

181 ------- 

182 qgraph_filename : `str` 

183 Name of file containing generated QuantumGraph. 

184 

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"]) 

197 

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) 

204 

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 

215 

216 

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. 

221 

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) 

239 

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) 

247 

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) 

254 

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 ) 

264 

265 

266def cluster_quanta(config, qgraph, name): 

267 """Call specified function to group quanta into clusters to be run 

268 together. 

269 

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. 

278 

279 Returns 

280 ------- 

281 cqgraph : `lsst.ctrl.bps.ClusteredQuantumGraph` 

282 Generated ClusteredQuantumGraph. 

283 

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