Coverage for tests/test_pipelines.py: 96%

123 statements  

« prev     ^ index     » next       coverage.py v7.16.2, created at 2026-09-30 12:07 +0000

1# This file is part of ap_pipe. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

5# (http://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 program is free software: you can redistribute it and/or modify 

10# it under the terms of the GNU General Public License as published by 

11# the Free Software Foundation, either version 3 of the License, or 

12# (at your option) any later version. 

13# 

14# This program is distributed in the hope that it will be useful, 

15# but WITHOUT ANY WARRANTY; without even the implied warranty of 

16# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the 

17# GNU General Public License for more details. 

18# 

19# You should have received a copy of the GNU General Public License 

20# along with this program. If not, see <http://www.gnu.org/licenses/>. 

21 

22import itertools 

23import tempfile 

24import unittest 

25 

26import lsst.daf.butler.tests as butlerTests 

27import lsst.pipe.base 

28from lsst.pipe.base.tests.pipelineStepTester import PipelineStepTester # Can't use fully-qualified name 

29import lsst.utils 

30import lsst.utils.tests 

31 

32from lsst.resources import ResourcePath 

33 

34 

35class PipelineDefintionsTestSuite(lsst.utils.tests.TestCase): 

36 """Tests of the self-consistency of our pipeline definitions. 

37 """ 

38 def setUp(self): 

39 self.path = ResourcePath("eups://ap_pipe/pipelines/", forceDirectory=True) 

40 # Each pipeline file should have a subset that represents it in 

41 # higher-level pipelines. 

42 self.synonyms = {"ApPipe.yaml": "apPipe", 

43 "ApPipeDaytime.yaml": "apPipe", 

44 "ApPipeWithIsrTaskLSST.yaml": "apPipe", 

45 "ApPipeWithPreconvolution.yaml": "apPipe", 

46 "ApPipeWithFakes.yaml": "apPipe", 

47 "SingleFrame.yaml": "singleFrame", 

48 "SingleFrameWithIsrTaskLSST.yaml": "singleFrame", 

49 "RunIsrWithoutInterChipCrosstalk.yaml": "runIsr", 

50 "RunIsrForCrosstalkSources.yaml": "runOverscan", 

51 "CreateInjectionCatalogs.yaml": "fakeCreation", 

52 } 

53 

54 def test_graph_build(self): 

55 """Test that each pipeline definition file can be 

56 used to build a graph. 

57 """ 

58 files = ResourcePath.findFileResources([self.path], file_filter=r".*\.yaml$") 

59 for file in files: 

60 if "QuickTemplate" in file.path: 

61 # Our QuickTemplate definition cannot be tested here because it 

62 # depends on drp_tasks, which we cannot make a dependency here. 

63 continue 

64 if "PromptTemplate" in file.path: 

65 # Our PromptTemplate definition cannot be tested here because it 

66 # depends on drp_tasks, which we cannot make a dependency here. 

67 continue 

68 with self.subTest(file=str(file)): 

69 pipeline = lsst.pipe.base.Pipeline.from_uri(file) 

70 pipeline.addConfigOverride("parameters", "apdb_config", "some/file/path.yaml") 

71 # If this fails, it will produce a useful error message. 

72 pipeline.to_graph() 

73 

74 def test_datasets(self): 

75 files = ResourcePath.findFileResources( 

76 [self.path.join("_ingredients", forceDirectory=True)], file_filter=r".*\.yaml$" 

77 ) 

78 for file in files: 

79 if "QuickTemplate" in file.path: 

80 # Our QuickTemplate definition cannot be tested here because it 

81 # depends on drp_tasks, which we cannot make a dependency here. 

82 continue 

83 if "injection/" in file.path: 

84 # The source-injection post-processing ingredient is a partial 

85 # pipeline merged into full AP pipelines at build time; 

86 # it is validated separately by test_injection_ingredient. 

87 continue 

88 with self.subTest(file=str(file)): 

89 expected_inputs = { 

90 # ISR 

91 "raw", "camera", "crosstalk", "crosstalkSources", "bias", "dark", "flat", "ptc", 

92 "fringe", "straylightData", "bfKernel", "newBFKernel", "defects", "linearizer", 

93 "opticsTransmission", "filterTransmission", "atmosphereTransmission", 

94 "illumMaskedImage", "deferredChargeCalib", 

95 # ISR-LSST 

96 "bfk", "cti", "dnlLUT", "gain_correction", 

97 # Everything else 

98 "skyMap", "gaia_dr3_20230707", "gaia_dr2_20200414", "ps1_pv3_3pi_20170110", 

99 "template_coadd", "pretrainedModelPackage", "dia_source_apdb" 

100 } 

101 # Detect source-injection pipelines by task label rather than 

102 # relying on filename conventions. 

103 temp_pipeline = lsst.pipe.base.Pipeline.from_uri(file) 

104 temp_pipeline.addConfigOverride("parameters", "apdb_config", "some/file/path.yaml") 

105 if "injectVisit" in temp_pipeline.task_labels: 

106 expected_inputs.add("injection_catalog") 

107 expected_inputs.add("VisitDetectorFakeSourceCat") 

108 tester = PipelineStepTester( 

109 filename=file, 

110 step_suffixes=[""], # Test full pipeline 

111 initial_dataset_types=[("ps1_pv3_3pi_20170110", {"htm7"}, "SimpleCatalog", False), 

112 ("gaia_dr2_20200414", {"htm7"}, "SimpleCatalog", False), 

113 ("gaia_dr3_20230707", {"htm7"}, "SimpleCatalog", False), 

114 ], 

115 expected_inputs=expected_inputs, 

116 # Pipeline outputs highly in flux, don't test 

117 expected_outputs=set(), 

118 pipeline_patches={"parameters:apdb_config": "some/file/path.yaml", 

119 }, 

120 ) 

121 # Tester modifies Butler registry, so need a fresh repo every time 

122 with tempfile.TemporaryDirectory() as tempRepo: 

123 butler = butlerTests.makeTestRepo(tempRepo) 

124 tester.run(butler, self) 

125 

126 def test_whole_subset(self): 

127 """Test that each pipeline's synonymous subset includes all tasks, 

128 including those imported from other files. 

129 """ 

130 files = ResourcePath.findFileResources([self.path], file_filter=r".*\.yaml$") 

131 for file in files: 

132 if "QuickTemplate" in file.path: 

133 # Our QuickTemplate definition cannot be tested here because it 

134 # depends on drp_tasks, which we cannot make a dependency here. 

135 continue 

136 elif "injection/" in file.path: 

137 # PostInjectedTasksApPipe is not actually an AP pipeline 

138 continue 

139 elif "ApdbDeduplication" in file.path: 

140 # The task to export catalogs from the APDB and re-run 

141 # association is not intended to be part of Prompt Processing 

142 # or batch AP pipeline runs. 

143 continue 

144 elif "PromptTemplate" in file.path: 

145 # Our PromptTemplate definition cannot be tested here because it 

146 # depends on drp_tasks, which we cannot make a dependency here. 

147 continue 

148 with self.subTest(file=str(file)): 

149 pipeline = lsst.pipe.base.Pipeline.from_uri(file) 

150 subset = self.synonyms.get(file.basename(), "<unknown_synonym>") 

151 self.assertEqual(pipeline.subsets.get(subset, "<missing>"), set(pipeline.task_labels), 

152 msg=f"These tasks are missing from subset '{subset}'") 

153 

154 def test_ap_pipe_subsets(self): 

155 """Test the unique subsets of ApPipe. 

156 """ 

157 files = ResourcePath.findFileResources([self.path], file_filter=r"^ApPipe.*\.yaml$") 

158 required_subsets = {"preload", "prompt", "afterburner"} 

159 # getRegionTimeFromVisit is part of no subset besides apPipe. This is a 

160 # very deliberate exception; see RFC-997. 

161 no_subset_wanted = {"getRegionTimeFromVisit"} 

162 

163 for file in files: 

164 if "injection/" in file.path: 164 ↛ 166line 164 didn't jump to line 166 because the condition on line 164 was never true

165 # PostInjectedTasksApPipe is not actually an AP pipeline 

166 continue 

167 with self.subTest(file=str(file)): 

168 pipeline = lsst.pipe.base.Pipeline.from_uri(file) 

169 # Do all steps exist? 

170 self.assertGreaterEqual(pipeline.subsets.keys(), required_subsets, 

171 msg="An AP pipeline is missing subsets " 

172 f"{required_subsets - pipeline.subsets.keys()}.") 

173 # Is each task part of exactly one step? 

174 for set1, set2 in itertools.product(required_subsets, required_subsets): 

175 if set1 == set2: 

176 continue 

177 tasks1 = pipeline.subsets[set1] 

178 tasks2 = pipeline.subsets[set2] 

179 self.assertTrue(tasks1.isdisjoint(tasks2), 

180 msg=f"Subsets '{set1}' and '{set2}' share tasks " 

181 f"{tasks1.intersection(tasks2)}.") 

182 subsetted = set().union(*[pipeline.subsets[s] for s in required_subsets]) 

183 self.assertEqual(subsetted, set(pipeline.task_labels) - no_subset_wanted, 

184 msg=f"These tasks are not in any of the subsets {required_subsets}.") 

185 

186 def test_injection_ingredient(self): 

187 """Test the source-injection post-processing ingredient pipeline. 

188 

189 PostInjectedTasksApPipe is a partial pipeline merged into full AP 

190 pipelines at build time by make_injection_pipeline. This test 

191 validates that it can build a graph and contains the expected tasks. 

192 """ 

193 ingredient = self.path.join("_ingredients").join("injection").join("PostInjectedTasksApPipe.yaml") 

194 with self.subTest(file=str(ingredient)): 

195 pipeline = lsst.pipe.base.Pipeline.from_uri(ingredient) 

196 expected_tasks = { 

197 "injectedMatchDiaSrc", 

198 "injectedMatchAssocDiaSrc", 

199 "consolidateMatchDiaSrc", 

200 "consolidateMatchAssocDiaSrc", 

201 } 

202 self.assertGreaterEqual( 

203 set(pipeline.task_labels), 

204 expected_tasks, 

205 msg="Source-injection post-processing ingredient is missing expected tasks.", 

206 ) 

207 

208 def test_generated_pipeline_readiness(self): 

209 """Test that the generated ApPipeWithFakes ingredient exists and is valid. 

210 

211 pipelines/_ingredients/ApPipeWithFakes.yaml is generated at build time 

212 by make_injection_pipeline (invoked via scons). This test verifies that 

213 generation occurred before pipeline tests run, and that the generated 

214 pipeline includes the source-injection task. If this test is skipped, 

215 run 'scons' in the ap_pipe root directory first. 

216 """ 

217 generated = self.path.join("_ingredients/ApPipeWithFakes.yaml") 

218 if not generated.exists(): 218 ↛ 221line 218 didn't jump to line 221 because the condition on line 218 was never true

219 # fail the test with a message that explains how to fix the problem, 

220 # rather than silently skipping it 

221 self.fail( 

222 f"{generated} has not been generated yet. " 

223 "Run 'scons' in the ap_pipe root directory to generate it." 

224 ) 

225 with self.subTest(file=str(generated)): 

226 pipeline = lsst.pipe.base.Pipeline.from_uri(generated) 

227 pipeline.addConfigOverride("parameters", "apdb_config", "some/file/path.yaml") 

228 self.assertIn( 

229 "injectVisit", 

230 pipeline.task_labels, 

231 msg="Generated ApPipeWithFakes.yaml is missing the 'injectVisit' task.", 

232 ) 

233 

234 def test_preconvolution_isr_matches_ap_pipe(self): 

235 """Test that, for each instrument, ApPipeWithPreconvolution defines 

236 the same isr task (class and config) as the corresponding ApPipe. 

237 

238 Preconvolution changes only image subtraction and DIA-source 

239 detection; instrument signature removal must be unaffected. 

240 """ 

241 files = [ 

242 f for f in ResourcePath.findFileResources( 

243 [self.path], file_filter=r"^ApPipeWithPreconvolution\.yaml$" 

244 ) 

245 if "_ingredients" not in f.path 

246 ] 

247 # Sanity-check that this test actually has cameras to compare. 

248 self.assertGreater(len(files), 0, 

249 msg="No camera-specific ApPipeWithPreconvolution.yaml files found.") 

250 

251 for precon_file in files: 

252 with self.subTest(file=str(precon_file)): 

253 base_file = precon_file.dirname().join("ApPipe.yaml") 

254 self.assertTrue(base_file.exists(), 

255 msg=f"Expected sibling ApPipe.yaml next to {precon_file}: " 

256 f"{base_file} does not exist.") 

257 

258 precon = lsst.pipe.base.Pipeline.from_uri(precon_file) 

259 base = lsst.pipe.base.Pipeline.from_uri(base_file) 

260 # apdb_config has no default and must be set before to_graph(). 

261 precon.addConfigOverride("parameters", "apdb_config", "some/file/path.yaml") 

262 base.addConfigOverride("parameters", "apdb_config", "some/file/path.yaml") 

263 

264 precon_isr = precon.to_graph().tasks["isr"] 

265 base_isr = base.to_graph().tasks["isr"] 

266 

267 self.assertEqual(precon_isr.task_class_name, base_isr.task_class_name, 

268 msg=f"isr task class differs between ApPipe.yaml and " 

269 f"ApPipeWithPreconvolution.yaml in {precon_file.dirname()}.") 

270 # Can't just do `assertEqual(precon_isr, base_isr)` since 

271 # Task nodes are intentionally not equality comparable. 

272 self.assertTrue( 

273 base_isr.config.compare(precon_isr.config, shortcut=False), 

274 msg=f"isr task config differs between ApPipe.yaml and " 

275 f"ApPipeWithPreconvolution.yaml in {precon_file.dirname()}." 

276 ) 

277 

278 def test_inherited_subsets(self): 

279 """Test that instrument-specific pipelines have all the subsets of their 

280 generic counterparts. 

281 

282 Note that this does not check inheritance *within* `_ingredients`! 

283 """ 

284 files = [ 

285 f for f in ResourcePath.findFileResources([self.path], file_filter=r".*\.yaml$") 

286 if "_ingredients" not in f.path 

287 ] 

288 for file in files: 

289 if "QuickTemplate" in file.path: 

290 # Our QuickTemplate definition cannot be tested here because it 

291 # depends on drp_tasks, which we cannot make a dependency here. 

292 continue 

293 with self.subTest(file=str(file)): 

294 generic = self.path.join("_ingredients/", forceDirectory=True).join(file.basename()) 

295 if not generic.exists(): 

296 continue 

297 special_subsets = lsst.pipe.base.Pipeline.from_uri(file).subsets.keys() 

298 generic_subsets = lsst.pipe.base.Pipeline.from_uri(generic).subsets.keys() 

299 self.assertGreaterEqual(special_subsets, generic_subsets, 

300 msg="The instrument-specific pipeline is missing subsets " 

301 f"{generic_subsets - special_subsets}.") 

302 

303 

304class MemoryTester(lsst.utils.tests.MemoryTestCase): 

305 pass 

306 

307 

308def setup_module(module): 

309 lsst.utils.tests.init() 

310 

311 

312if __name__ == "__main__": 312 ↛ 313line 312 didn't jump to line 313 because the condition on line 312 was never true

313 lsst.utils.tests.init() 

314 unittest.main()