Coverage for tests/test_pipelines.py: 96%
123 statements
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-24 09:43 +0000
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-24 09:43 +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/>.
22import itertools
23import tempfile
24import unittest
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
32from lsst.resources import ResourcePath
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 }
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()
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)
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}'")
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"}
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}.")
186 def test_injection_ingredient(self):
187 """Test the source-injection post-processing ingredient pipeline.
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 )
208 def test_generated_pipeline_readiness(self):
209 """Test that the generated ApPipeWithFakes ingredient exists and is valid.
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 )
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.
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.")
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.")
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")
264 precon_isr = precon.to_graph().tasks["isr"]
265 base_isr = base.to_graph().tasks["isr"]
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 )
278 def test_inherited_subsets(self):
279 """Test that instrument-specific pipelines have all the subsets of their
280 generic counterparts.
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}.")
304class MemoryTester(lsst.utils.tests.MemoryTestCase):
305 pass
308def setup_module(module):
309 lsst.utils.tests.init()
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()