Coverage for tests/test_butler.py: 97%
1901 statements
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-22 09:33 +0000
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-22 09:33 +0000
1# This file is part of daf_butler.
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 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 <http://www.gnu.org/licenses/>.
28"""Tests for Butler."""
30from __future__ import annotations
32import json
33import logging
34import os
35import pathlib
36import pickle
37import re
38import shutil
39import tempfile
40import unittest
41import unittest.mock
42import uuid
43import warnings
44import weakref
45from collections.abc import Callable, Mapping
46from typing import TYPE_CHECKING, Any, cast
48import astropy.time
49from sqlalchemy.exc import IntegrityError
51from lsst.daf.butler import (
52 Butler,
53 ButlerConfig,
54 ButlerMetrics,
55 ButlerRepoIndex,
56 CollectionCycleError,
57 CollectionType,
58 Config,
59 DataCoordinate,
60 DatasetExistence,
61 DatasetNotFoundError,
62 DatasetProvenance,
63 DatasetRef,
64 DatasetType,
65 DimensionRecord,
66 FileDataset,
67 NoDefaultCollectionError,
68 StorageClassFactory,
69 ValidationError,
70 script,
71)
72from lsst.daf.butler._rubin.file_datasets import transfer_datasets_to_datastore
73from lsst.daf.butler._rubin.temporary_for_ingest import TemporaryForIngest
74from lsst.daf.butler._rubin.transfer_datasets_in_place import transfer_datasets_in_place
75from lsst.daf.butler.datastore import NullDatastore
76from lsst.daf.butler.datastore.file_templates import FileTemplate, FileTemplateValidationError
77from lsst.daf.butler.datastores.file_datastore.retrieve_artifacts import ZipIndex
78from lsst.daf.butler.datastores.fileDatastore import FileDatastore
79from lsst.daf.butler.direct_butler import DirectButler
80from lsst.daf.butler.registry import (
81 CollectionError,
82 CollectionTypeError,
83 ConflictingDefinitionError,
84 DataIdValueError,
85 DatasetTypeExpressionError,
86 MissingCollectionError,
87 OrphanedRecordError,
88)
89from lsst.daf.butler.registry.sql_registry import SqlRegistry
90from lsst.daf.butler.repo_relocation import BUTLER_ROOT_TAG
91from lsst.daf.butler.tests import MetricsExample, MetricsExampleModel, MultiDetectorFormatter
92from lsst.daf.butler.tests._repo_template_cache import make_repo_for_test
93from lsst.daf.butler.tests.dict_convertible_model import DictConvertibleModel
94from lsst.daf.butler.tests.postgresql import TemporaryPostgresInstance, setup_postgres_test_db
95from lsst.daf.butler.tests.server_available import butler_server_import_error, butler_server_is_available
96from lsst.daf.butler.tests.utils import (
97 MetricTestRepo,
98 TestCaseMixin,
99 create_populated_sqlite_registry,
100 makeTestTempDir,
101 removeTestTempDir,
102 safeTestTempDir,
103)
104from lsst.resources import ResourcePath
105from lsst.resources.http import HttpResourcePath
106from lsst.resources.tests import make_remote_test_uri
107from lsst.utils import doImportType
108from lsst.utils.introspection import get_full_type_name
110if butler_server_is_available:
111 from lsst.daf.butler.tests.server import create_test_server
114if TYPE_CHECKING:
115 import types
117 from lsst.daf.butler import DimensionGroup, Registry, StorageClass
119TESTDIR = os.path.abspath(os.path.dirname(__file__))
122def clean_environment() -> None:
123 """Remove external environment variables that affect the tests."""
124 for k in ("DAF_BUTLER_REPOSITORY_INDEX",):
125 os.environ.pop(k, None)
128def makeExampleMetrics() -> MetricsExample:
129 """Return example dataset suitable for tests."""
130 return MetricsExample(
131 {"AM1": 5.2, "AM2": 30.6},
132 {"a": [1, 2, 3], "b": {"blue": 5, "red": "green"}},
133 [563, 234, 456.7, 752, 8, 9, 27],
134 )
137class TransactionTestError(Exception):
138 """Specific error for testing transactions, to prevent misdiagnosing
139 that might otherwise occur when a standard exception is used.
140 """
142 pass
145class ButlerConfigTests(unittest.TestCase):
146 """Simple tests for ButlerConfig that are not tested in any other test
147 cases.
148 """
150 def testSearchPath(self) -> None:
151 configFile = os.path.join(TESTDIR, "config", "basic", "butler.yaml")
152 with self.assertLogs("lsst.daf.butler", level="DEBUG") as cm:
153 config1 = ButlerConfig(configFile)
154 self.assertNotIn("testConfigs", "\n".join(cm.output))
156 overrideDirectory = os.path.join(TESTDIR, "config", "testConfigs")
157 with self.assertLogs("lsst.daf.butler", level="DEBUG") as cm:
158 config2 = ButlerConfig(configFile, searchPaths=[overrideDirectory])
159 self.assertIn("testConfigs", "\n".join(cm.output))
161 key = ("datastore", "records", "table")
162 self.assertNotEqual(config1[key], config2[key])
163 self.assertEqual(config2[key], "override_record")
166class ButlerPutGetTests(TestCaseMixin):
167 """Helper method for running a suite of put/get tests from different
168 butler configurations.
169 """
171 root: str
172 default_run = "ingésτ😺"
173 storageClassFactory: StorageClassFactory
174 configFile: str | None
175 tmpConfigFile: str
177 @staticmethod
178 def addDatasetType(
179 datasetTypeName: str, dimensions: DimensionGroup, storageClass: StorageClass | str, registry: Registry
180 ) -> DatasetType:
181 """Create a DatasetType and register it"""
182 datasetType = DatasetType(datasetTypeName, dimensions, storageClass)
183 registry.registerDatasetType(datasetType)
184 return datasetType
186 @classmethod
187 def setUpClass(cls) -> None:
188 cls.storageClassFactory = StorageClassFactory()
189 if cls.configFile is not None: 189 ↛ exitline 189 didn't return from function 'setUpClass' because the condition on line 189 was always true
190 cls.storageClassFactory.addFromConfig(cls.configFile)
192 def assertGetComponents(
193 self,
194 butler: Butler,
195 datasetRef: DatasetRef,
196 components: tuple[str, ...],
197 reference: Any,
198 collections: Any = None,
199 ) -> None:
200 datasetType = datasetRef.datasetType
201 dataId = datasetRef.dataId
202 deferred = butler.getDeferred(datasetRef)
204 for component in components:
205 compTypeName = datasetType.componentTypeName(component)
206 result = butler.get(compTypeName, dataId, collections=collections)
207 self.assertEqual(result, getattr(reference, component))
208 result_deferred = deferred.get(component=component)
209 self.assertEqual(result_deferred, result)
211 def tearDown(self) -> None:
212 if self.root is not None: 212 ↛ exitline 212 didn't return from function 'tearDown' because the condition on line 212 was always true
213 removeTestTempDir(self.root)
215 def create_empty_butler(
216 self,
217 run: str | None = None,
218 writeable: bool | None = None,
219 metrics: ButlerMetrics | None = None,
220 cleanup: bool = True,
221 ):
222 """Create a Butler for the test repository, without inserting test
223 data.
224 """
225 butler = Butler.from_config(self.tmpConfigFile, run=run, writeable=writeable, metrics=metrics)
226 if cleanup:
227 self.enterContext(butler)
228 assert isinstance(butler, DirectButler), "Expect DirectButler in configuration"
229 return butler
231 def create_butler(
232 self,
233 run: str,
234 storageClass: StorageClass | str,
235 datasetTypeName: str,
236 metrics: ButlerMetrics | None = None,
237 ) -> tuple[Butler, DatasetType]:
238 """Create a Butler for the test repository and insert some test data
239 into it.
240 """
241 butler = self.create_empty_butler(run=run, metrics=metrics)
243 collections = set(butler.collections.query("*"))
244 self.assertEqual(collections, {run})
245 # Create and register a DatasetType
246 dimensions = butler.dimensions.conform(["instrument", "visit"])
248 datasetType = self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry)
250 # Add needed Dimensions
251 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"})
252 butler.registry.insertDimensionData(
253 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"}
254 )
255 butler.registry.insertDimensionData(
256 "visit_system", {"instrument": "DummyCamComp", "id": 1, "name": "default"}
257 )
258 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20200101})
259 visit_start = astropy.time.Time("2020-01-01 08:00:00.123456789", scale="tai")
260 visit_end = astropy.time.Time("2020-01-01 08:00:36.66", scale="tai")
261 butler.registry.insertDimensionData(
262 "visit",
263 {
264 "instrument": "DummyCamComp",
265 "id": 423,
266 "name": "fourtwentythree",
267 "physical_filter": "d-r",
268 "datetime_begin": visit_start,
269 "datetime_end": visit_end,
270 "day_obs": 20200101,
271 },
272 )
274 # Add more visits for some later tests
275 for visit_id in (424, 425):
276 butler.registry.insertDimensionData(
277 "visit",
278 {
279 "instrument": "DummyCamComp",
280 "id": visit_id,
281 "name": f"fourtwentyfour_{visit_id}",
282 "physical_filter": "d-r",
283 "day_obs": 20200101,
284 },
285 )
286 return butler, datasetType
288 def runPutGetTest(self, storageClass: StorageClass, datasetTypeName: str) -> Butler:
289 # New datasets will be added to run and tag, but we will only look in
290 # tag when looking up datasets.
291 run = self.default_run
292 butler, datasetType = self.create_butler(run, storageClass, datasetTypeName)
293 assert butler.run is not None
295 # Create and store a dataset
296 metric = makeExampleMetrics()
297 dataId = butler.registry.expandDataId({"instrument": "DummyCamComp", "visit": 423})
299 # Dataset should not exist if we haven't added it
300 with self.assertRaises(DatasetNotFoundError):
301 butler.get(datasetTypeName, dataId)
303 # Put and remove the dataset once as a DatasetRef, once as a dataId,
304 # and once with a DatasetType
306 # Keep track of any collections we add and do not clean up
307 expected_collections = {run}
309 counter = 0
310 ref = DatasetRef(datasetType, dataId, id=uuid.UUID(int=1), run="put_run_1")
311 args = tuple[DatasetRef] | tuple[str | DatasetType, DataCoordinate]
312 for args in ((ref,), (datasetTypeName, dataId), (datasetType, dataId)):
313 # Since we are using subTest we can get cascading failures
314 # here with the first attempt failing and the others failing
315 # immediately because the dataset already exists. Work around
316 # this by using a distinct run collection each time
317 counter += 1
318 this_run = f"put_run_{counter}"
319 butler.collections.register(this_run)
320 expected_collections.update({this_run})
322 with self.subTest(args=repr(args)):
323 kwargs: dict[str, Any] = {}
324 if not isinstance(args[0], DatasetRef): # type: ignore
325 kwargs["run"] = this_run
326 ref = butler.put(metric, *args, **kwargs)
327 self.assertIsInstance(ref, DatasetRef)
329 # Test get of a ref.
330 metricOut = butler.get(ref)
331 self.assertEqual(metric, metricOut)
332 # Test get
333 metricOut = butler.get(ref.datasetType.name, dataId, collections=this_run)
334 self.assertEqual(metric, metricOut)
335 # Test get with a datasetRef
336 metricOut = butler.get(ref)
337 self.assertEqual(metric, metricOut)
338 # Test getDeferred with dataId
339 metricOut = butler.getDeferred(ref.datasetType.name, dataId, collections=this_run).get()
340 self.assertEqual(metric, metricOut)
341 # Test getDeferred with a ref
342 metricOut = butler.getDeferred(ref).get()
343 self.assertEqual(metric, metricOut)
345 # Check we can get components
346 if storageClass.isComposite():
347 self.assertGetComponents(
348 butler, ref, ("summary", "data", "output"), metric, collections=this_run
349 )
351 primary_uri, secondary_uris = butler.getURIs(ref)
352 n_uris = len(secondary_uris)
353 if primary_uri:
354 n_uris += 1
356 # Can the artifacts themselves be retrieved?
357 if not butler._datastore.isEphemeral:
358 # Create a temporary directory to hold the retrieved
359 # artifacts.
360 with tempfile.TemporaryDirectory(
361 prefix="butler-artifacts-", ignore_cleanup_errors=True
362 ) as artifact_root:
363 root_uri = ResourcePath(artifact_root, forceDirectory=True)
365 for preserve_path in (True, False):
366 destination = root_uri.join(f"{preserve_path}_{counter}/")
367 log = logging.getLogger("lsst.x")
368 log.debug("Using destination %s for args %s", destination, args)
369 # Use copy so that we can test that overwrite
370 # protection works (using "auto" for File URIs
371 # would use hard links and subsequent transfer
372 # would work because it knows they are the same
373 # file).
374 transferred = butler.retrieveArtifacts(
375 [ref], destination, preserve_path=preserve_path, transfer="copy"
376 )
377 self.assertGreater(len(transferred), 0)
378 artifacts = list(ResourcePath.findFileResources([destination]))
379 # Filter out the index file.
380 artifacts = [a for a in artifacts if a.basename() != ZipIndex.index_name]
381 self.assertEqual(set(transferred), set(artifacts))
383 for artifact in transferred:
384 path_in_destination = artifact.relative_to(destination)
385 self.assertIsNotNone(path_in_destination)
386 assert path_in_destination is not None
388 # When path is not preserved there should not
389 # be any path separators.
390 num_seps = path_in_destination.count("/")
391 if preserve_path:
392 self.assertGreater(num_seps, 0)
393 else:
394 self.assertEqual(num_seps, 0)
396 self.assertEqual(
397 len(artifacts),
398 n_uris,
399 "Comparing expected artifacts vs actual:"
400 f" {artifacts} vs {primary_uri} and {secondary_uris}",
401 )
403 if preserve_path:
404 # No need to run these twice
405 with self.assertRaises(ValueError):
406 butler.retrieveArtifacts([ref], destination, transfer="move")
408 with self.assertRaisesRegex(
409 ValueError, "^Destination location must refer to a directory"
410 ):
411 butler.retrieveArtifacts(
412 [ref], ResourcePath("/some/file.txt", forceDirectory=False)
413 )
415 with self.assertRaises(FileExistsError):
416 butler.retrieveArtifacts([ref], destination)
418 transferred_again = butler.retrieveArtifacts(
419 [ref], destination, preserve_path=preserve_path, overwrite=True
420 )
421 self.assertEqual(set(transferred_again), set(transferred))
423 # Now remove the dataset completely.
424 butler.pruneDatasets([ref], purge=True, unstore=True)
425 # Lookup with original args should still fail.
426 kwargs = {"collections": this_run}
427 if isinstance(args[0], DatasetRef):
428 kwargs = {} # Prevent warning from being issued.
429 self.assertFalse(butler.exists(*args, **kwargs))
430 # get() should still fail.
431 with self.assertRaises((FileNotFoundError, DatasetNotFoundError)):
432 butler.get(ref)
433 # Registry shouldn't be able to find it by dataset_id anymore.
434 self.assertIsNone(butler.get_dataset(ref.id))
436 # Do explicit registry removal since we know they are
437 # empty
438 butler.collections.x_remove(this_run)
439 expected_collections.remove(this_run)
441 # Create DatasetRef for put using default run.
442 refIn = DatasetRef(datasetType, dataId, id=uuid.UUID(int=1), run=butler.run)
444 # Check that getDeferred fails with standalone ref.
445 with self.assertRaises(LookupError):
446 butler.getDeferred(refIn)
448 # Put the dataset again, since the last thing we did was remove it
449 # and we want to use the default collection.
450 ref = butler.put(metric, refIn)
452 # Get with parameters
453 stop = 4
454 sliced = butler.get(ref, parameters={"slice": slice(stop)})
455 self.assertNotEqual(metric, sliced)
456 self.assertEqual(metric.summary, sliced.summary)
457 self.assertEqual(metric.output, sliced.output)
458 assert metric.data is not None # for mypy
459 self.assertEqual(metric.data[:stop], sliced.data)
460 # getDeferred with parameters
461 sliced = butler.getDeferred(ref, parameters={"slice": slice(stop)}).get()
462 self.assertNotEqual(metric, sliced)
463 self.assertEqual(metric.summary, sliced.summary)
464 self.assertEqual(metric.output, sliced.output)
465 self.assertEqual(metric.data[:stop], sliced.data)
466 # getDeferred with deferred parameters
467 sliced = butler.getDeferred(ref).get(parameters={"slice": slice(stop)})
468 self.assertNotEqual(metric, sliced)
469 self.assertEqual(metric.summary, sliced.summary)
470 self.assertEqual(metric.output, sliced.output)
471 self.assertEqual(metric.data[:stop], sliced.data)
473 if storageClass.isComposite():
474 # Check that components can be retrieved
475 metricOut = butler.get(ref.datasetType.name, dataId)
476 compNameS = ref.datasetType.componentTypeName("summary")
477 compNameD = ref.datasetType.componentTypeName("data")
478 summary = butler.get(compNameS, dataId)
479 self.assertEqual(summary, metric.summary)
480 data = butler.get(compNameD, dataId)
481 self.assertEqual(data, metric.data)
483 if "counter" in storageClass.derivedComponents:
484 count = butler.get(ref.datasetType.componentTypeName("counter"), dataId)
485 self.assertEqual(count, len(data))
487 count = butler.get(
488 ref.datasetType.componentTypeName("counter"), dataId, parameters={"slice": slice(stop)}
489 )
490 self.assertEqual(count, stop)
492 compRef = butler.find_dataset(compNameS, dataId, collections=butler.collections.defaults)
493 assert compRef is not None
494 summary = butler.get(compRef)
495 self.assertEqual(summary, metric.summary)
497 # Create a Dataset type that has the same name but is inconsistent.
498 inconsistentDatasetType = DatasetType(
499 datasetTypeName, datasetType.dimensions, self.storageClassFactory.getStorageClass("Config")
500 )
502 # Getting with a dataset type that does not match registry fails
503 with self.assertRaisesRegex(
504 ValueError,
505 "(Supplied dataset type .* inconsistent with registry)"
506 "|(The new storage class .* is not compatible with the existing storage class)",
507 ):
508 butler.get(inconsistentDatasetType, dataId)
510 # Combining a DatasetRef with a dataId should fail
511 with self.assertRaisesRegex(ValueError, "DatasetRef given, cannot use dataId as well"):
512 butler.get(ref, dataId)
513 # Getting with an explicit ref should fail if the id doesn't match.
514 with self.assertRaises((FileNotFoundError, DatasetNotFoundError)):
515 butler.get(DatasetRef(ref.datasetType, ref.dataId, id=uuid.UUID(int=101), run=butler.run))
517 # Getting a dataset with unknown parameters should fail
518 with self.assertRaisesRegex(KeyError, "Parameter 'unsupported' not understood"):
519 butler.get(ref, parameters={"unsupported": True})
521 # Check we have a collection
522 collections = set(butler.collections.query("*"))
523 self.assertEqual(collections, expected_collections)
525 # Clean up to check that we can remove something that may have
526 # already had a component removed
527 butler.pruneDatasets([ref], unstore=True, purge=True)
529 # Add the same ref again, so we can check that duplicate put fails.
530 ref = butler.put(metric, datasetType, dataId)
532 # Repeat put will fail.
533 with self.assertRaisesRegex(
534 ConflictingDefinitionError, "A database constraint failure was triggered"
535 ):
536 butler.put(metric, datasetType, dataId)
538 # Remove the datastore entry.
539 butler.pruneDatasets([ref], unstore=True, purge=False, disassociate=False)
541 # Put will still fail
542 with self.assertRaisesRegex(
543 ConflictingDefinitionError, "A database constraint failure was triggered"
544 ):
545 butler.put(metric, datasetType, dataId)
547 # Repeat the same sequence with resolved ref.
548 butler.pruneDatasets([ref], unstore=True, purge=True)
549 ref = butler.put(metric, refIn)
551 # Repeat put will fail.
552 with self.assertRaisesRegex(ConflictingDefinitionError, "Datastore already contains dataset"):
553 butler.put(metric, refIn)
555 # Remove the datastore entry.
556 butler.pruneDatasets([ref], unstore=True, purge=False, disassociate=False)
558 # In case of resolved ref this write will succeed.
559 ref = butler.put(metric, refIn)
561 # Leave the dataset in place since some downstream tests require
562 # something to be present
564 return butler
566 def testDeferredCollectionPassing(self) -> None:
567 # Construct a butler with no run or collection, but make it writeable.
568 butler = self.create_empty_butler(writeable=True)
569 # Create and register a DatasetType
570 dimensions = butler.dimensions.conform(["instrument", "visit"])
571 datasetType = self.addDatasetType(
572 "example", dimensions, self.storageClassFactory.getStorageClass("StructuredData"), butler.registry
573 )
574 # Add needed Dimensions
575 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"})
576 butler.registry.insertDimensionData(
577 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"}
578 )
579 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20250101})
580 butler.registry.insertDimensionData(
581 "visit",
582 {
583 "instrument": "DummyCamComp",
584 "id": 423,
585 "name": "fourtwentythree",
586 "physical_filter": "d-r",
587 "day_obs": 20250101,
588 },
589 )
590 dataId = {"instrument": "DummyCamComp", "visit": 423}
591 # Create dataset.
592 metric = makeExampleMetrics()
593 # Register a new run and put dataset.
594 run = "deferred"
595 self.assertTrue(butler.collections.register(run))
596 # Second time it will be allowed but indicate no-op
597 self.assertFalse(butler.collections.register(run))
598 ref = butler.put(metric, datasetType, dataId, run=run)
599 # Putting with no run should fail with TypeError.
600 with self.assertRaises(CollectionError):
601 butler.put(metric, datasetType, dataId)
602 # Dataset should exist.
603 self.assertTrue(butler.exists(datasetType, dataId, collections=[run]))
604 # We should be able to get the dataset back, but with and without
605 # a deferred dataset handle.
606 self.assertEqual(metric, butler.get(datasetType, dataId, collections=[run]))
607 self.assertEqual(metric, butler.getDeferred(datasetType, dataId, collections=[run]).get())
608 # Trying to find the dataset without any collection is an error.
609 with self.assertRaises(NoDefaultCollectionError):
610 butler.exists(datasetType, dataId)
611 with self.assertRaises(CollectionError):
612 butler.get(datasetType, dataId)
613 # Associate the dataset with a different collection.
614 butler.collections.register("tagged", type=CollectionType.TAGGED)
615 butler.registry.associate("tagged", [ref])
616 # Deleting the dataset from the new collection should make it findable
617 # in the original collection.
618 butler.pruneDatasets([ref], tags=["tagged"])
619 self.assertTrue(butler.exists(datasetType, dataId, collections=[run]))
622class ButlerTests(ButlerPutGetTests):
623 """Tests for Butler."""
625 useTempRoot = True
626 validationCanFail: bool
627 fullConfigKey: str | None
628 registryStr: str | None
629 datastoreName: list[str] | None
630 datastoreStr: list[str]
631 predictionSupported = True
632 """Does getURIs support 'prediction mode'?"""
634 def setUp(self) -> None:
635 """Create a new butler root for each test."""
636 self.root = makeTestTempDir(TESTDIR)
637 make_repo_for_test(self.root, config=Config(self.configFile))
638 self.tmpConfigFile = os.path.join(self.root, "butler.yaml")
640 def are_uris_equivalent(self, uri1: ResourcePath, uri2: ResourcePath) -> bool:
641 """Return True if two URIs refer to the same resource.
643 Subclasses may override to handle unique requirements.
644 """
645 return uri1 == uri2
647 def testConstructor(self) -> None:
648 """Independent test of constructor."""
649 butler = Butler.from_config(self.tmpConfigFile, run=self.default_run)
650 self.enterContext(butler)
651 self.assertIsInstance(butler, Butler)
653 # Check that butler.yaml is added automatically.
654 if self.tmpConfigFile.endswith(end := "/butler.yaml"):
655 config_dir = self.tmpConfigFile[: -len(end)]
656 butler = Butler.from_config(config_dir, run=self.default_run)
657 self.enterContext(butler)
658 self.assertIsInstance(butler, Butler)
660 # Even with a ResourcePath.
661 butler = Butler.from_config(ResourcePath(config_dir, forceDirectory=True), run=self.default_run)
662 self.enterContext(butler)
663 self.assertIsInstance(butler, Butler)
665 collections = set(butler.collections.query("*"))
666 self.assertEqual(collections, {self.default_run})
668 # Check that some special characters can be included in run name.
669 special_run = "u@b.c-A"
670 butler_special = Butler.from_config(butler=butler, run=special_run)
671 self.enterContext(butler_special)
672 collections = set(butler_special.registry.queryCollections("*@*"))
673 self.assertEqual(collections, {special_run})
675 butler2 = Butler.from_config(butler=butler, collections=["other"])
676 self.enterContext(butler2)
677 self.assertEqual(butler2.collections.defaults, ("other",))
678 self.assertIsNone(butler2.run)
679 self.assertEqual(type(butler._datastore), type(butler2._datastore))
680 self.assertEqual(butler._datastore.config, butler2._datastore.config)
682 # Test that we can use an environment variable to find this
683 # repository.
684 butler_index = Config()
685 butler_index["label"] = self.tmpConfigFile
686 for suffix in (".yaml", ".json"):
687 # Ensure that the content differs so that we know that
688 # we aren't reusing the cache.
689 bad_label = f"file://bucket/not_real{suffix}"
690 butler_index["bad_label"] = bad_label
691 with ResourcePath.temporary_uri(suffix=suffix) as temp_file:
692 butler_index.dumpToUri(temp_file)
693 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": str(temp_file)}):
694 self.assertEqual(Butler.get_known_repos(), {"label", "bad_label"})
695 uri = Butler.get_repo_uri("bad_label")
696 self.assertEqual(uri, ResourcePath(bad_label))
697 uri = Butler.get_repo_uri("label")
698 butler = Butler.from_config(uri, writeable=False)
699 self.assertIsInstance(butler, Butler)
700 butler.close()
701 butler = Butler.from_config("label", writeable=False)
702 self.assertIsInstance(butler, Butler)
703 butler.close()
704 with self.assertRaisesRegex(FileNotFoundError, "aliases:.*bad_label"):
705 Butler.from_config("not_there", writeable=False)
706 with self.assertRaisesRegex(FileNotFoundError, "resolved from alias 'bad_label'"):
707 Butler.from_config("bad_label")
708 with self.assertRaises(FileNotFoundError):
709 # Should ignore aliases.
710 Butler.from_config(ResourcePath("label", forceAbsolute=False))
711 with self.assertRaises(KeyError) as cm:
712 Butler.get_repo_uri("missing")
713 self.assertEqual(
714 Butler.get_repo_uri("missing", True), ResourcePath("missing", forceAbsolute=False)
715 )
716 self.assertIn("not known to", str(cm.exception))
717 # Should report no failure.
718 self.assertEqual(ButlerRepoIndex.get_failure_reason(), "")
719 with ResourcePath.temporary_uri(suffix=suffix) as temp_file:
720 # Now with empty configuration.
721 butler_index = Config()
722 butler_index.dumpToUri(temp_file)
723 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": str(temp_file)}):
724 with self.assertRaisesRegex(FileNotFoundError, "(no known aliases)"):
725 Butler.from_config("label")
726 with ResourcePath.temporary_uri(suffix=suffix) as temp_file:
727 # Now with bad contents.
728 with open(temp_file.ospath, "w") as fh:
729 print("'", file=fh)
730 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": str(temp_file)}):
731 with self.assertRaisesRegex(FileNotFoundError, "(no known aliases:.*could not be read)"):
732 Butler.from_config("label")
733 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": "file://not_found/x.yaml"}):
734 with self.assertRaises(FileNotFoundError):
735 Butler.get_repo_uri("label")
736 self.assertEqual(Butler.get_known_repos(), set())
738 with self.assertRaisesRegex(FileNotFoundError, "index file not found"):
739 Butler.from_config("label")
741 # Check that we can create Butler when the alias file is not found.
742 butler = Butler.from_config(self.tmpConfigFile, writeable=False)
743 self.enterContext(butler)
744 self.assertIsInstance(butler, Butler)
745 with self.assertRaises(RuntimeError) as cm:
746 # No environment variable set.
747 Butler.get_repo_uri("label")
748 self.assertEqual(Butler.get_repo_uri("label", True), ResourcePath("label", forceAbsolute=False))
749 self.assertIn("No repository index defined", str(cm.exception))
750 with self.assertRaisesRegex(FileNotFoundError, "no known aliases.*No repository index"):
751 # No aliases registered.
752 Butler.from_config("not_there")
753 self.assertEqual(Butler.get_known_repos(), set())
755 def testClose(self):
756 butler = self.create_empty_butler(cleanup=False)
757 is_direct_butler = isinstance(butler, DirectButler)
758 if is_direct_butler: 758 ↛ 761line 758 didn't jump to line 761 because the condition on line 758 was always true
759 self.assertFalse(butler._closed)
761 with butler as butler_from_context_manager:
762 self.assertIs(butler, butler_from_context_manager)
763 if is_direct_butler: 763 ↛ 769line 763 didn't jump to line 769 because the condition on line 763 was always true
764 self.assertTrue(butler._closed)
765 with self.assertRaisesRegex(RuntimeError, "has been closed"):
766 butler.get_dataset_type("raw")
768 # Close may be called multiple times.
769 butler.close()
770 if is_direct_butler: 770 ↛ exitline 770 didn't return from function 'testClose' because the condition on line 770 was always true
771 self.assertTrue(butler._closed)
773 def testGarbageCollection(self):
774 """Test that Butler does not have any circular references that prevent
775 it from being garbage collected immediately when it goes out of scope.
776 """
777 butler = self.create_empty_butler(cleanup=False)
778 is_direct_butler = isinstance(butler, DirectButler)
779 butler_ref = weakref.ref(butler)
780 if is_direct_butler: 780 ↛ 787line 780 didn't jump to line 787 because the condition on line 780 was always true
781 registry_ref = weakref.ref(butler._registry)
782 managers_ref = weakref.ref(butler._registry._managers)
783 datastore_ref = weakref.ref(butler._datastore)
784 db_ref = weakref.ref(butler._registry._db)
785 engine_ref = weakref.ref(butler._registry._db._engine)
787 with warnings.catch_warnings():
788 # Hide warnings from unclosed database handles.
789 warnings.simplefilter("ignore", ResourceWarning)
790 del butler
791 self.assertIsNone(butler_ref(), "Butler should have been garbage collected")
792 if is_direct_butler: 792 ↛ 800line 792 didn't jump to line 800 because the condition on line 792 was always true
793 self.assertIsNone(registry_ref(), "SqlRegistry should have been garbage collected")
794 self.assertIsNone(managers_ref(), "Registry managers should have been garbage collected")
795 self.assertIsNone(datastore_ref(), "Datastore should have been garbage collected")
796 self.assertIsNone(db_ref(), "Database should have been garbage collected")
797 # SQLAlchemy has internal reference cycles, so the Engine instance
798 # is not cleaned up promptly even if we release our reference to
799 # it. Explicitly clean it up here to avoid file handles leaking.
800 if is_direct_butler: 800 ↛ exitline 800 didn't jump to the function exit
801 engine = engine_ref()
802 if engine is not None: 802 ↛ exitline 802 didn't jump to the function exit
803 engine.dispose()
805 def testDafButlerRepositories(self):
806 with unittest.mock.patch.dict(
807 os.environ,
808 {"DAF_BUTLER_REPOSITORIES": "label: 'https://someuri.com'\notherLabel: 'https://otheruri.com'\n"},
809 ):
810 self.assertEqual(str(Butler.get_repo_uri("label")), "https://someuri.com")
812 with unittest.mock.patch.dict(
813 os.environ,
814 {
815 "DAF_BUTLER_REPOSITORIES": "label: https://someuri.com",
816 "DAF_BUTLER_REPOSITORY_INDEX": "https://someuri.com",
817 },
818 ):
819 with self.assertRaisesRegex(RuntimeError, "Only one of the environment variables"):
820 Butler.get_repo_uri("label")
822 with unittest.mock.patch.dict(
823 os.environ,
824 {"DAF_BUTLER_REPOSITORIES": "invalid"},
825 ):
826 with self.assertRaisesRegex(ValueError, "Repository index not in expected format"):
827 Butler.get_repo_uri("label")
829 def testBasicPutGet(self) -> None:
830 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
831 self.runPutGetTest(storageClass, "test_metric")
833 def testCompositePutGetConcrete(self) -> None:
834 storageClass = self.storageClassFactory.getStorageClass("StructuredCompositeReadCompNoDisassembly")
835 butler = self.runPutGetTest(storageClass, "test_metric")
837 # Should *not* be disassembled
838 datasets = list(butler.registry.queryDatasets(..., collections=self.default_run))
839 self.assertEqual(len(datasets), 1)
840 uri, components = butler.getURIs(datasets[0])
841 self.assertIsInstance(uri, ResourcePath)
842 self.assertFalse(components)
843 self.assertEqual(uri.fragment, "", f"Checking absence of fragment in {uri}")
844 self.assertIn("423", str(uri), f"Checking visit is in URI {uri}")
846 # Predicted dataset
847 if self.predictionSupported: 847 ↛ exitline 847 didn't return from function 'testCompositePutGetConcrete' because the condition on line 847 was always true
848 dataId = {"instrument": "DummyCamComp", "visit": 424}
849 uri, components = butler.getURIs(datasets[0].datasetType, dataId=dataId, predict=True)
850 self.assertFalse(components)
851 self.assertIsInstance(uri, ResourcePath)
852 self.assertIn("424", str(uri), f"Checking visit is in URI {uri}")
853 self.assertEqual(uri.fragment, "predicted", f"Checking for fragment in {uri}")
854 # Repeat with a DatasetRef to test that code path.
855 ref = DatasetRef(
856 datasets[0].datasetType,
857 dataId=DataCoordinate.standardize(dataId, universe=butler.dimensions),
858 run=self.default_run,
859 )
860 uri2, components2 = butler.getURIs(ref, predict=True)
861 self.assertFalse(components2)
862 self.assertEqual(uri, uri2)
864 def testCompositePutGetVirtual(self) -> None:
865 storageClass = self.storageClassFactory.getStorageClass("StructuredCompositeReadComp")
866 butler = self.runPutGetTest(storageClass, "test_metric_comp")
868 # Should be disassembled
869 datasets = list(butler.registry.queryDatasets(..., collections=self.default_run))
870 self.assertEqual(len(datasets), 1)
871 uri, components = butler.getURIs(datasets[0])
873 if butler._datastore.isEphemeral:
874 # Never disassemble in-memory datastore
875 self.assertIsInstance(uri, ResourcePath)
876 self.assertFalse(components)
877 self.assertEqual(uri.fragment, "", f"Checking absence of fragment in {uri}")
878 self.assertIn("423", str(uri), f"Checking visit is in URI {uri}")
879 else:
880 self.assertIsNone(uri)
881 self.assertEqual(set(components), set(storageClass.components))
882 for compuri in components.values():
883 self.assertIsInstance(compuri, ResourcePath)
884 self.assertIn("423", str(compuri), f"Checking visit is in URI {compuri}")
885 self.assertEqual(compuri.fragment, "", f"Checking absence of fragment in {compuri}")
887 if self.predictionSupported: 887 ↛ exitline 887 didn't return from function 'testCompositePutGetVirtual' because the condition on line 887 was always true
888 # Predicted dataset
889 dataId = {"instrument": "DummyCamComp", "visit": 424}
890 uri, components = butler.getURIs(datasets[0].datasetType, dataId=dataId, predict=True)
892 if butler._datastore.isEphemeral:
893 # Never disassembled
894 self.assertIsInstance(uri, ResourcePath)
895 self.assertFalse(components)
896 self.assertIn("424", str(uri), f"Checking visit is in URI {uri}")
897 self.assertEqual(uri.fragment, "predicted", f"Checking for fragment in {uri}")
898 else:
899 self.assertIsNone(uri)
900 self.assertEqual(set(components), set(storageClass.components))
901 for compuri in components.values():
902 self.assertIsInstance(compuri, ResourcePath)
903 self.assertIn("424", str(compuri), f"Checking visit is in URI {compuri}")
904 self.assertEqual(compuri.fragment, "predicted", f"Checking for fragment in {compuri}")
906 def testStorageClassOverrideGet(self) -> None:
907 """Test storage class conversion on get with override."""
908 storageClass = self.storageClassFactory.getStorageClass("StructuredData")
909 datasetTypeName = "anything"
910 run = self.default_run
912 butler, datasetType = self.create_butler(run, storageClass, datasetTypeName)
914 # Create and store a dataset.
915 metric = makeExampleMetrics()
916 dataId = {"instrument": "DummyCamComp", "visit": 423}
918 ref = butler.put(metric, datasetType, dataId)
920 # Return native type.
921 retrieved = butler.get(ref)
922 self.assertEqual(retrieved, metric)
924 # Specify an override.
925 new_sc = self.storageClassFactory.getStorageClass("MetricsConversion")
926 model = butler.get(ref, storageClass=new_sc)
927 self.assertNotEqual(type(model), type(retrieved))
928 self.assertIs(type(model), new_sc.pytype)
929 self.assertEqual(retrieved, model)
931 # Defer but override later.
932 deferred = butler.getDeferred(ref)
933 model = deferred.get(storageClass=new_sc)
934 self.assertIs(type(model), new_sc.pytype)
935 self.assertEqual(retrieved, model)
937 # Defer but override up front.
938 deferred = butler.getDeferred(ref, storageClass=new_sc)
939 model = deferred.get()
940 self.assertIs(type(model), new_sc.pytype)
941 self.assertEqual(retrieved, model)
943 # Retrieve a component. Should be a tuple.
944 data = butler.get("anything.data", dataId, storageClass="StructuredDataDataTestTuple")
945 self.assertIs(type(data), tuple)
946 self.assertEqual(data, tuple(retrieved.data))
948 # Parameter on the write storage class should work regardless
949 # of read storage class.
950 data = butler.get(
951 "anything.data",
952 dataId,
953 storageClass="StructuredDataDataTestTuple",
954 parameters={"slice": slice(2, 4)},
955 )
956 self.assertEqual(len(data), 2)
958 # Try a parameter that is known to the read storage class but not
959 # the write storage class.
960 with self.assertRaises(KeyError):
961 butler.get(
962 "anything.data",
963 dataId,
964 storageClass="StructuredDataDataTestTuple",
965 parameters={"xslice": slice(2, 4)},
966 )
968 def testComponentFromOverriddenStorageClass(self) -> None:
969 """Test component get where the component is only defined by the
970 read storage class and not by the storage class used to write.
971 """
972 # StructuredDataNoComponents defines no components at all, whereas
973 # MetricsConversion (which it can be converted to) defines several.
974 write_sc = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
975 read_sc = self.storageClassFactory.getStorageClass("MetricsConversion")
976 self.assertFalse(write_sc.allComponents())
977 self.assertIn("summary", read_sc.allComponents())
979 butler, datasetType = self.create_butler(self.default_run, write_sc, "unstructured")
981 metric = makeExampleMetrics()
982 dataId = {"instrument": "DummyCamComp", "visit": 423}
983 ref = butler.put(metric, datasetType, dataId)
985 # The composite conversion on its own must work.
986 self.assertIs(type(butler.get(ref, storageClass=read_sc)), read_sc.pytype)
988 # A component of the converted composite, requested via a ref.
989 component_ref = ref.overrideStorageClass(read_sc).makeComponentRef("summary")
990 self.assertEqual(butler.get(component_ref), metric.summary)
992 # The same component, requested via a deferred handle that was given
993 # the storage class override up front.
994 deferred = butler.getDeferred(ref, storageClass=read_sc)
995 self.assertEqual(deferred.get(component="summary"), metric.summary)
997 # A component whose storage class is also overridden, on top of the
998 # storage class the read composite declares for it.
999 converted = butler.get(component_ref, storageClass="DictConvertibleModel")
1000 self.assertIsInstance(converted, DictConvertibleModel)
1001 self.assertEqual(converted.content, metric.summary)
1003 # The handle storage class applies to the composite and so selects the
1004 # component, while the one given to get() applies to the component.
1005 converted = deferred.get(component="summary", storageClass="DictConvertibleModel")
1006 self.assertIsInstance(converted, DictConvertibleModel)
1007 self.assertEqual(converted.content, metric.summary)
1009 def testPytypePutCoercion(self) -> None:
1010 """Test python type coercion on Butler.get and put."""
1011 # Store some data with the normal example storage class.
1012 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
1013 datasetTypeName = "test_metric"
1014 butler, _ = self.create_butler(self.default_run, storageClass, datasetTypeName)
1016 dataId = {"instrument": "DummyCamComp", "visit": 423}
1018 # Put a dict and this should coerce to a MetricsExample
1019 test_dict = {"summary": {"a": 1}, "output": {"b": 2}}
1020 metric_ref = butler.put(test_dict, datasetTypeName, dataId=dataId, visit=424)
1021 test_metric = butler.get(metric_ref)
1022 self.assertEqual(get_full_type_name(test_metric), "lsst.daf.butler.tests.MetricsExample")
1023 self.assertEqual(test_metric.summary, test_dict["summary"])
1024 self.assertEqual(test_metric.output, test_dict["output"])
1026 # Check that the put still works if a DatasetType is given with
1027 # a definition matching this python type.
1028 registry_type = butler.get_dataset_type(datasetTypeName)
1029 this_type = DatasetType(datasetTypeName, registry_type.dimensions, "StructuredDataDictJson")
1030 metric2_ref = butler.put(test_dict, this_type, dataId=dataId, visit=425)
1031 self.assertEqual(metric2_ref.datasetType, registry_type)
1033 # The get will return the type expected by registry.
1034 test_metric2 = butler.get(metric2_ref)
1035 self.assertEqual(get_full_type_name(test_metric2), "lsst.daf.butler.tests.MetricsExample")
1037 # Make a new DatasetRef with the compatible but different DatasetType.
1038 # This should now return a dict.
1039 new_ref = DatasetRef(this_type, metric2_ref.dataId, id=metric2_ref.id, run=metric2_ref.run)
1040 test_dict2 = butler.get(new_ref)
1041 self.assertEqual(get_full_type_name(test_dict2), "dict")
1043 # Get it again with the wrong dataset type definition using get()
1044 # rather than get(). This should be consistent with get()
1045 # behavior and return the type of the DatasetType.
1046 test_dict3 = butler.get(this_type, dataId=dataId, visit=425)
1047 self.assertEqual(get_full_type_name(test_dict3), "dict")
1049 def test_ingest_zip(self) -> None:
1050 """Create butler, export data, delete data, import from Zip."""
1051 butler, dataset_type = self.create_butler(
1052 run=self.default_run, storageClass="StructuredData", datasetTypeName="metrics"
1053 )
1055 metric = makeExampleMetrics()
1056 refs = []
1057 for visit in (423, 424, 425):
1058 ref = butler.put(metric, dataset_type, instrument="DummyCamComp", visit=visit)
1059 refs.append(ref)
1061 # Retrieve a Zip file.
1062 with tempfile.TemporaryDirectory(ignore_cleanup_errors=True) as tmpdir:
1063 zip = butler.retrieve_artifacts_zip(refs, destination=tmpdir)
1065 # Ingest will fail.
1066 with self.assertRaises(ConflictingDefinitionError):
1067 butler.ingest_zip(zip)
1069 # Clear out the collection.
1070 butler.removeRuns([self.default_run])
1071 self.assertFalse(butler.exists(refs[0]))
1073 butler.ingest_zip(zip, transfer="copy")
1074 self.assertGreater(butler._metrics.time_in_ingest, 0.0)
1075 self.assertEqual(butler._metrics.n_ingest, len(refs))
1077 # Check that it fails if we try it again.
1078 with self.assertRaises(ConflictingDefinitionError):
1079 butler.ingest_zip(zip, transfer="copy")
1081 # This will be a no-op.
1082 butler.ingest_zip(zip, transfer="copy", skip_existing=True)
1084 # Create an entirely new local file butler in this temp directory.
1085 new_butler_cfg = make_repo_for_test(tmpdir)
1086 new_butler = Butler.from_config(new_butler_cfg, writeable=True)
1087 self.enterContext(new_butler)
1089 # This will fail since dimensions records are missing.
1090 with self.assertRaises(ConflictingDefinitionError):
1091 new_butler.ingest_zip(zip, transfer="copy")
1093 # Dry run should work.
1094 new_butler.ingest_zip(zip, transfer="copy", dry_run=True)
1096 new_butler.ingest_zip(zip, transfer="copy", transfer_dimensions=True)
1097 self.assertTrue(butler.exists(refs[0]))
1099 # Check that the refs can be read again.
1100 _ = [butler.get(ref) for ref in refs]
1102 uri = butler.getURI(refs[2])
1103 self.assertTrue(uri.exists())
1105 # Delete one dataset. The Zip file should still exist and allow
1106 # remaining refs to be read.
1107 butler.pruneDatasets([refs[0]], purge=True, unstore=True)
1108 self.assertTrue(uri.exists())
1110 metric2 = butler.get(refs[1])
1111 self.assertEqual(metric2, metric, msg=f"{metric2} != {metric}")
1113 butler.removeRuns([self.default_run])
1114 self.assertFalse(uri.exists())
1115 self.assertFalse(butler.exists(refs[-1]))
1117 with self.assertRaises(ValueError):
1118 butler.retrieve_artifacts_zip([], destination=".")
1120 def testIngest(self) -> None:
1121 butler = self.create_empty_butler(run=self.default_run)
1123 # Create and register a DatasetType
1124 dimensions = butler.dimensions.conform(["instrument", "visit", "detector"])
1126 storageClass = self.storageClassFactory.getStorageClass("StructuredDataDictYaml")
1127 datasetTypeName = "metric"
1129 datasetType = self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry)
1131 # Add needed Dimensions
1132 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"})
1133 butler.registry.insertDimensionData(
1134 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"}
1135 )
1136 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20250101})
1137 for detector in (1, 2):
1138 butler.registry.insertDimensionData(
1139 "detector", {"instrument": "DummyCamComp", "id": detector, "full_name": f"detector{detector}"}
1140 )
1142 butler.registry.insertDimensionData(
1143 "visit",
1144 {
1145 "instrument": "DummyCamComp",
1146 "id": 423,
1147 "name": "fourtwentythree",
1148 "physical_filter": "d-r",
1149 "day_obs": 20250101,
1150 },
1151 {
1152 "instrument": "DummyCamComp",
1153 "id": 424,
1154 "name": "fourtwentyfour",
1155 "physical_filter": "d-r",
1156 "day_obs": 20250101,
1157 },
1158 )
1160 formatter = doImportType("lsst.daf.butler.formatters.yaml.YamlFormatter")
1161 dataRoot = os.path.join(TESTDIR, "data", "basic")
1162 datasets = []
1163 # Test one DatasetRef with a run that exists, and the other with a run
1164 # that doesn't exist, to verify that run collections are created when
1165 # required.
1166 runs = {1: self.default_run, 2: "a/new/run"}
1167 for detector in (1, 2):
1168 detector_name = f"detector_{detector}"
1169 metricFile = os.path.join(dataRoot, f"{detector_name}.yaml")
1170 dataId = butler.registry.expandDataId(
1171 {"instrument": "DummyCamComp", "visit": 423, "detector": detector}
1172 )
1173 # Create a DatasetRef for ingest
1174 refIn = DatasetRef(datasetType, dataId, run=runs[detector])
1176 datasets.append(FileDataset(path=metricFile, refs=[refIn], formatter=formatter))
1178 butler.ingest(*datasets, transfer="copy")
1180 dataId1 = {"instrument": "DummyCamComp", "detector": 1, "visit": 423}
1181 dataId2 = {"instrument": "DummyCamComp", "detector": 2, "visit": 423}
1183 metrics1 = butler.get(datasetTypeName, dataId1)
1184 metrics2 = butler.get(datasetTypeName, dataId2, collections="a/new/run")
1185 self.assertNotEqual(metrics1, metrics2)
1187 # Compare URIs
1188 uri1 = butler.getURI(datasetTypeName, dataId1)
1189 uri2 = butler.getURI(datasetTypeName, dataId2, collections="a/new/run")
1190 self.assertFalse(self.are_uris_equivalent(uri1, uri2), f"Cf. {uri1} with {uri2}")
1192 # Re-ingesting the same datasets raises an error with
1193 # skip_existing=False.
1194 with self.assertRaises(ConflictingDefinitionError):
1195 butler.ingest(*datasets, transfer="copy")
1196 # skip_existing=True makes it a no-op to re-ingest the same datasets.
1197 butler.ingest(*datasets, transfer="copy", skip_existing=True)
1199 # Now do a multi-dataset but single file ingest
1200 metricFile = os.path.join(dataRoot, "detectors.yaml")
1201 refs = []
1202 for detector in (1, 2):
1203 detector_name = f"detector_{detector}"
1204 dataId = butler.registry.expandDataId(
1205 {"instrument": "DummyCamComp", "visit": 424, "detector": detector}
1206 )
1207 # Create a DatasetRef for ingest
1208 refs.append(DatasetRef(datasetType, dataId, run=self.default_run))
1210 # Test "move" transfer to ensure that the files themselves
1211 # have disappeared following ingest.
1212 with ResourcePath.temporary_uri(suffix=".yaml") as tempFile:
1213 tempFile.transfer_from(ResourcePath(metricFile), transfer="copy")
1215 datasets = []
1216 datasets.append(FileDataset(path=tempFile, refs=refs, formatter=MultiDetectorFormatter))
1218 # For first ingest use copy.
1219 butler.ingest(*datasets, transfer="copy", record_validation_info=False)
1221 # Now try to ingest again in "execution butler" mode where
1222 # the registry entries exist but the datastore does not have
1223 # the files. We also need to strip the dimension records to ensure
1224 # that they will be re-added by the ingest.
1225 ref = datasets[0].refs[0]
1226 datasets[0].refs = [
1227 cast(
1228 DatasetRef,
1229 butler.find_dataset(ref.datasetType, data_id=ref.dataId, collections=ref.run),
1230 )
1231 for ref in datasets[0].refs
1232 ]
1233 all_refs = []
1234 for dataset in datasets:
1235 refs = []
1236 for ref in dataset.refs:
1237 # Create a dict from the dataId to drop the records.
1238 new_data_id = dict(ref.dataId.required)
1239 new_ref = butler.find_dataset(ref.datasetType, new_data_id, collections=ref.run)
1240 assert new_ref is not None
1241 self.assertFalse(new_ref.dataId.hasRecords())
1242 refs.append(new_ref)
1243 dataset.refs = refs
1244 all_refs.extend(dataset.refs)
1245 butler.pruneDatasets(all_refs, disassociate=False, unstore=True, purge=False)
1247 # Use move mode to test that the file is deleted. Also
1248 # disable recording of file size.
1249 butler.ingest(*datasets, transfer="move", record_validation_info=False)
1251 # Check that every ref now has records.
1252 for dataset in datasets:
1253 for ref in dataset.refs:
1254 self.assertTrue(ref.dataId.hasRecords())
1256 # Ensure that the file has disappeared.
1257 self.assertFalse(tempFile.exists())
1259 # Check that the datastore recorded no file size.
1260 # Not all datastores can support this.
1261 try:
1262 infos = butler._datastore.getStoredItemsInfo(datasets[0].refs[0]) # type: ignore[attr-defined]
1263 self.assertEqual(infos[0].file_size, -1)
1264 except AttributeError:
1265 pass
1267 dataId1 = {"instrument": "DummyCamComp", "detector": 1, "visit": 424}
1268 dataId2 = {"instrument": "DummyCamComp", "detector": 2, "visit": 424}
1270 multi1 = butler.get(datasetTypeName, dataId1)
1271 multi2 = butler.get(datasetTypeName, dataId2)
1273 self.assertEqual(multi1, metrics1)
1274 self.assertEqual(multi2, metrics2)
1276 # Compare URIs
1277 uri1 = butler.getURI(datasetTypeName, dataId1)
1278 uri2 = butler.getURI(datasetTypeName, dataId2)
1279 self.assertTrue(self.are_uris_equivalent(uri1, uri2), f"Cf. {uri1} with {uri2}")
1281 # Test that removing one does not break the second
1282 # This line will issue a warning log message for a ChainedDatastore
1283 # that uses an InMemoryDatastore since in-memory can not ingest
1284 # files.
1285 butler.pruneDatasets([datasets[0].refs[0]], unstore=True, disassociate=False)
1286 self.assertFalse(butler.exists(datasetTypeName, dataId1))
1287 self.assertTrue(butler.exists(datasetTypeName, dataId2))
1288 multi2b = butler.get(datasetTypeName, dataId2)
1289 self.assertEqual(multi2, multi2b)
1291 # Ensure we can ingest 0 datasets
1292 datasets = []
1293 butler.ingest(*datasets)
1295 def testPickle(self) -> None:
1296 """Test pickle support."""
1297 butler = self.create_empty_butler(run=self.default_run)
1298 assert isinstance(butler, DirectButler), "Expect DirectButler in configuration"
1299 butlerOut = pickle.loads(pickle.dumps(butler))
1300 self.enterContext(butlerOut)
1301 self.assertIsInstance(butlerOut, Butler)
1302 self.assertEqual(butlerOut._config, butler._config)
1303 self.assertEqual(list(butlerOut.collections.defaults), list(butler.collections.defaults))
1304 self.assertEqual(butlerOut.run, butler.run)
1306 def testGetDatasetTypes(self) -> None:
1307 butler = self.create_empty_butler(run=self.default_run)
1308 dimensions = butler.dimensions.conform(["instrument", "visit", "physical_filter"])
1309 dimensionEntries: list[tuple[str, list[Mapping[str, Any]]]] = [
1310 (
1311 "instrument",
1312 [
1313 {"instrument": "DummyCam"},
1314 {"instrument": "DummyHSC"},
1315 {"instrument": "DummyCamComp"},
1316 ],
1317 ),
1318 ("physical_filter", [{"instrument": "DummyCam", "name": "d-r", "band": "R"}]),
1319 ("day_obs", [{"instrument": "DummyCam", "id": 20250101}]),
1320 (
1321 "visit",
1322 [
1323 {
1324 "instrument": "DummyCam",
1325 "id": 42,
1326 "name": "fortytwo",
1327 "physical_filter": "d-r",
1328 "day_obs": 20250101,
1329 }
1330 ],
1331 ),
1332 ]
1333 storageClass = self.storageClassFactory.getStorageClass("StructuredData")
1334 # Add needed Dimensions
1335 for element, data in dimensionEntries:
1336 butler.registry.insertDimensionData(element, *data)
1338 # When a DatasetType is added to the registry entries are not created
1339 # for components but querying them can return the components.
1340 datasetTypeNames = {"metric", "metric2", "metric4", "metric33", "pvi", "paramtest"}
1341 components = set()
1342 for datasetTypeName in datasetTypeNames:
1343 # Create and register a DatasetType
1344 self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry)
1346 for componentName in storageClass.components:
1347 components.add(DatasetType.nameWithComponent(datasetTypeName, componentName))
1349 fromRegistry: set[DatasetType] = set()
1350 for parent_dataset_type in butler.registry.queryDatasetTypes():
1351 fromRegistry.add(parent_dataset_type)
1352 fromRegistry.update(parent_dataset_type.makeAllComponentDatasetTypes())
1353 self.assertEqual({d.name for d in fromRegistry}, datasetTypeNames | components)
1355 # Query with wildcard.
1356 dataset_types = butler.registry.queryDatasetTypes("metric*")
1357 self.assertEqual(len(dataset_types), 4, f"Got: {dataset_types}")
1358 # but not regex.
1359 with self.assertRaises(DatasetTypeExpressionError):
1360 butler.registry.queryDatasetTypes(["pvi", re.compile("metric.*")])
1362 # Now that we have some dataset types registered, validate them
1363 butler.validateConfiguration(
1364 ignore=[
1365 "test_metric_comp",
1366 "metric3",
1367 "metric5",
1368 "calexp",
1369 "DummySC",
1370 "datasetType.component",
1371 "random_data",
1372 "random_data_2",
1373 ]
1374 )
1376 # Add a new datasetType that will fail template validation
1377 self.addDatasetType("test_metric_comp", dimensions, storageClass, butler.registry)
1378 if self.validationCanFail:
1379 with self.assertRaises(ValidationError):
1380 butler.validateConfiguration()
1382 # Rerun validation but with a subset of dataset type names
1383 butler.validateConfiguration(datasetTypeNames=["metric4"])
1385 # Rerun validation but ignore the bad datasetType
1386 butler.validateConfiguration(
1387 ignore=[
1388 "test_metric_comp",
1389 "metric3",
1390 "metric5",
1391 "calexp",
1392 "DummySC",
1393 "datasetType.component",
1394 "random_data",
1395 "random_data_2",
1396 ]
1397 )
1399 def testTransaction(self) -> None:
1400 butler = self.create_empty_butler(run=self.default_run)
1401 datasetTypeName = "test_metric"
1402 dimensions = butler.dimensions.conform(["instrument", "visit"])
1403 dimensionEntries: tuple[tuple[str, Mapping[str, Any]], ...] = (
1404 ("instrument", {"instrument": "DummyCam"}),
1405 ("physical_filter", {"instrument": "DummyCam", "name": "d-r", "band": "R"}),
1406 ("day_obs", {"instrument": "DummyCam", "id": 20250101}),
1407 (
1408 "visit",
1409 {
1410 "instrument": "DummyCam",
1411 "id": 42,
1412 "name": "fortytwo",
1413 "physical_filter": "d-r",
1414 "day_obs": 20250101,
1415 },
1416 ),
1417 )
1418 storageClass = self.storageClassFactory.getStorageClass("StructuredData")
1419 metric = makeExampleMetrics()
1420 dataId = {"instrument": "DummyCam", "visit": 42}
1421 # Create and register a DatasetType
1422 datasetType = self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry)
1423 with self.assertRaises(TransactionTestError):
1424 with butler.transaction():
1425 # Add needed Dimensions
1426 for args in dimensionEntries:
1427 butler.registry.insertDimensionData(*args)
1428 # Store a dataset
1429 ref = butler.put(metric, datasetTypeName, dataId)
1430 self.assertIsInstance(ref, DatasetRef)
1431 # Test get of a ref.
1432 metricOut = butler.get(ref)
1433 self.assertEqual(metric, metricOut)
1434 # Test get
1435 metricOut = butler.get(datasetTypeName, dataId)
1436 self.assertEqual(metric, metricOut)
1437 # Check we can get components
1438 self.assertGetComponents(butler, ref, ("summary", "data", "output"), metric)
1439 raise TransactionTestError("This should roll back the entire transaction")
1440 with self.assertRaises(DataIdValueError, msg=f"Check can't expand DataId {dataId}"):
1441 butler.registry.expandDataId(dataId)
1442 # Should raise LookupError for missing data ID value
1443 with self.assertRaises(LookupError, msg=f"Check can't get by {datasetTypeName} and {dataId}"):
1444 butler.get(datasetTypeName, dataId)
1445 # Also check explicitly if Dataset entry is missing
1446 self.assertIsNone(butler.find_dataset(datasetType, dataId, collections=butler.collections.defaults))
1447 # Direct retrieval should not find the file in the Datastore
1448 with self.assertRaises(FileNotFoundError, msg=f"Check {ref} can't be retrieved directly"):
1449 butler.get(ref)
1451 def testMakeRepo(self) -> None:
1452 """Test that we can write butler configuration to a new repository via
1453 the Butler.makeRepo interface and then instantiate a butler from the
1454 repo root.
1455 """
1456 # Do not run the test if we know this datastore configuration does
1457 # not support a file system root
1458 if self.fullConfigKey is None:
1459 return
1461 # create two separate directories
1462 root1 = tempfile.mkdtemp(dir=self.root)
1463 root2 = tempfile.mkdtemp(dir=self.root)
1465 # This test asserts on repository creation itself, so it must not go
1466 # through the caching helper.
1467 self.assertFalse(Butler.has_repo_config(root1))
1468 butlerConfig = Butler.makeRepo(root1, config=Config(self.configFile))
1469 self.assertTrue(Butler.has_repo_config(root1))
1470 limited = Config(self.configFile)
1471 butler1 = Butler.from_config(butlerConfig)
1472 self.enterContext(butler1)
1473 assert isinstance(butler1, DirectButler), "Expect DirectButler in configuration"
1474 butlerConfig = Butler.makeRepo(root2, standalone=True, config=Config(self.configFile))
1475 full = Config(self.tmpConfigFile)
1476 butler2 = Butler.from_config(butlerConfig)
1477 self.enterContext(butler2)
1478 assert isinstance(butler2, DirectButler), "Expect DirectButler in configuration"
1479 # Butlers should have the same configuration regardless of whether
1480 # defaults were expanded.
1481 self.assertEqual(butler1._config, butler2._config)
1482 # Config files loaded directly should not be the same.
1483 self.assertNotEqual(limited, full)
1484 # Make sure "limited" doesn't have a few keys we know it should be
1485 # inheriting from defaults.
1486 self.assertIn(self.fullConfigKey, full)
1487 self.assertNotIn(self.fullConfigKey, limited)
1489 # Collections don't appear until something is put in them
1490 collections1 = set(butler1.registry.queryCollections())
1491 self.assertEqual(collections1, set())
1492 self.assertEqual(set(butler2.registry.queryCollections()), collections1)
1494 # Check that a config with no associated file name will not
1495 # work properly with relocatable Butler repo
1496 butlerConfig.configFile = None
1497 with self.assertRaises(ValueError):
1498 Butler.from_config(butlerConfig)
1500 with self.assertRaises(FileExistsError):
1501 Butler.makeRepo(self.root, standalone=True, config=Config(self.configFile), overwrite=False)
1503 def testStringification(self) -> None:
1504 butler = Butler.from_config(self.tmpConfigFile, run=self.default_run)
1505 self.enterContext(butler)
1506 butlerStr = str(butler)
1508 if self.datastoreStr is not None: 1508 ↛ 1511line 1508 didn't jump to line 1511 because the condition on line 1508 was always true
1509 for testStr in self.datastoreStr:
1510 self.assertIn(testStr, butlerStr)
1511 if self.registryStr is not None: 1511 ↛ 1514line 1511 didn't jump to line 1514 because the condition on line 1511 was always true
1512 self.assertIn(self.registryStr, butlerStr)
1514 datastoreName = butler._datastore.name
1515 if self.datastoreName is not None: 1515 ↛ exitline 1515 didn't return from function 'testStringification' because the condition on line 1515 was always true
1516 for testStr in self.datastoreName:
1517 self.assertIn(testStr, datastoreName)
1519 def testButlerRewriteDataId(self) -> None:
1520 """Test that dataIds can be rewritten based on dimension records."""
1521 butler = self.create_empty_butler(run=self.default_run)
1523 storageClass = self.storageClassFactory.getStorageClass("StructuredDataDict")
1524 datasetTypeName = "random_data"
1526 # Create dimension records.
1527 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"})
1528 butler.registry.insertDimensionData(
1529 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"}
1530 )
1531 butler.registry.insertDimensionData(
1532 "detector", {"instrument": "DummyCamComp", "id": 1, "full_name": "det1"}
1533 )
1535 dimensions = butler.dimensions.conform(["instrument", "exposure"])
1536 datasetType = DatasetType(datasetTypeName, dimensions, storageClass)
1537 butler.registry.registerDatasetType(datasetType)
1539 n_exposures = 5
1540 dayobs = 20210530
1542 # Create records for multiple day_obs but same seq_num to test that
1543 # we are constraining gets properly when day_obs/seq_num is used
1544 # for an exposure. Second day is year in future but is not used.
1545 for day_obs in (dayobs, dayobs + 1_00_00):
1546 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": day_obs})
1548 for i in range(n_exposures):
1549 group_name = f"group_{day_obs}_{i}"
1550 butler.registry.insertDimensionData(
1551 "group", {"instrument": "DummyCamComp", "name": group_name}
1552 )
1553 butler.registry.insertDimensionData(
1554 "exposure",
1555 {
1556 "instrument": "DummyCamComp",
1557 "id": day_obs + i,
1558 "obs_id": f"exp_{day_obs}_{i}",
1559 "seq_num": i,
1560 "day_obs": day_obs,
1561 "physical_filter": "d-r",
1562 "group": group_name,
1563 },
1564 )
1566 # Write some data.
1567 for i in range(n_exposures):
1568 metric = {"something": i, "other": "metric", "list": [2 * x for x in range(i)]}
1570 # Use the seq_num for the put to test rewriting.
1571 dataId = {"seq_num": i, "day_obs": dayobs, "instrument": "DummyCamComp", "physical_filter": "d-r"}
1572 ref = butler.put(metric, datasetTypeName, dataId=dataId)
1574 # Check that the exposure is correct in the dataId
1575 self.assertEqual(ref.dataId["exposure"], dayobs + i)
1577 # and check that we can get the dataset back with the same dataId
1578 new_metric = butler.get(datasetTypeName, dataId=dataId)
1579 self.assertEqual(new_metric, metric)
1581 # Check that we can find the datasets using the day_obs or the
1582 # exposure.day_obs.
1583 datasets_1 = list(
1584 butler.registry.queryDatasets(
1585 datasetType,
1586 collections=self.default_run,
1587 where="day_obs = :dayObs AND instrument = :instr",
1588 bind={"dayObs": dayobs, "instr": "DummyCamComp"},
1589 )
1590 )
1591 datasets_2 = list(
1592 butler.registry.queryDatasets(
1593 datasetType,
1594 collections=self.default_run,
1595 where="exposure.day_obs = :dayObs AND instrument = :instr",
1596 bind={"dayObs": dayobs, "instr": "DummyCamComp"},
1597 )
1598 )
1599 self.assertEqual(datasets_1, datasets_2)
1601 def testGetDatasetCollectionCaching(self):
1602 # Prior to DM-41117, there was a bug where get_dataset would throw
1603 # MissingCollectionError if you tried to fetch a dataset that was added
1604 # after the collection cache was last updated.
1605 reader_butler, datasetType = self.create_butler(self.default_run, "int", "datasettypename")
1606 writer_butler = self.create_empty_butler(writeable=True, run="new_run")
1607 dataId = {"instrument": "DummyCamComp", "visit": 423}
1608 put_ref = writer_butler.put(123, datasetType, dataId)
1609 get_ref = reader_butler.get_dataset(put_ref.id)
1610 self.assertEqual(get_ref.id, put_ref.id)
1611 # Also works when looking up via a hexadecimal string instead of a UUID
1612 # instance.
1613 hex_ref = reader_butler.get_dataset(put_ref.id.hex)
1614 self.assertEqual(hex_ref.id, put_ref.id)
1616 def testCollectionChainRedefine(self):
1617 butler = self._setup_to_test_collection_chain()
1619 butler.collections.redefine_chain("chain", "a")
1620 self._check_chain(butler, ["a"])
1622 # Duplicates are removed from the list of children
1623 butler.collections.redefine_chain("chain", ["c", "b", "c"])
1624 self._check_chain(butler, ["c", "b"])
1626 # Empty list clears the chain
1627 butler.collections.redefine_chain("chain", [])
1628 self._check_chain(butler, [])
1630 self._test_common_chain_functionality(butler, butler.collections.redefine_chain)
1632 def testCollectionChainPrepend(self):
1633 butler = self._setup_to_test_collection_chain()
1635 # Duplicates are removed from the list of children
1636 butler.collections.prepend_chain("chain", ["c", "b", "c"])
1637 self._check_chain(butler, ["c", "b"])
1639 # Prepend goes on the front of existing chain
1640 butler.collections.prepend_chain("chain", ["a"])
1641 self._check_chain(butler, ["a", "c", "b"])
1643 # Empty prepend does nothing
1644 butler.collections.prepend_chain("chain", [])
1645 self._check_chain(butler, ["a", "c", "b"])
1647 # Prepending children that already exist in the chain removes them from
1648 # their current position.
1649 butler.collections.prepend_chain("chain", ["d", "b", "c"])
1650 self._check_chain(butler, ["d", "b", "c", "a"])
1652 self._test_common_chain_functionality(butler, butler.collections.prepend_chain)
1654 def testCollectionChainExtend(self):
1655 butler = self._setup_to_test_collection_chain()
1657 # Duplicates are removed from the list of children
1658 butler.collections.extend_chain("chain", ["c", "b", "c"])
1659 self._check_chain(butler, ["c", "b"])
1661 # Extend goes on the end of existing chain
1662 butler.collections.extend_chain("chain", ["a"])
1663 self._check_chain(butler, ["c", "b", "a"])
1665 # Empty extend does nothing
1666 butler.collections.extend_chain("chain", [])
1667 self._check_chain(butler, ["c", "b", "a"])
1669 # Extending children that already exist in the chain removes them from
1670 # their current position.
1671 butler.collections.extend_chain("chain", ["d", "b", "c"])
1672 self._check_chain(butler, ["a", "d", "b", "c"])
1674 self._test_common_chain_functionality(butler, butler.collections.extend_chain)
1676 def testCollectionChainRemove(self) -> None:
1677 butler = self._setup_to_test_collection_chain()
1679 butler.collections.redefine_chain("chain", ["a", "b", "c", "d"])
1681 butler.collections.remove_from_chain("chain", "c")
1682 self._check_chain(butler, ["a", "b", "d"])
1684 # Duplicates are allowed in the list of children
1685 butler.collections.remove_from_chain("chain", ["b", "b", "a"])
1686 self._check_chain(butler, ["d"])
1688 # Empty remove does nothing
1689 butler.collections.remove_from_chain("chain", [])
1690 self._check_chain(butler, ["d"])
1692 # Removing children that aren't in the chain does nothing
1693 butler.collections.remove_from_chain("chain", ["a", "chain"])
1694 self._check_chain(butler, ["d"])
1696 self._test_common_chain_functionality(
1697 butler, butler.collections.remove_from_chain, skip_cycle_check=True
1698 )
1700 def _setup_to_test_collection_chain(self) -> Butler:
1701 butler = self.create_empty_butler(writeable=True)
1703 butler.collections.register("chain", CollectionType.CHAINED)
1705 runs = ["a", "b", "c", "d"]
1706 for run in runs:
1707 butler.collections.register(run)
1709 butler.collections.register("staticchain", CollectionType.CHAINED)
1710 butler.collections.redefine_chain("staticchain", ["a", "b"])
1712 return butler
1714 def _check_chain(self, butler: Butler, expected: list[str]) -> None:
1715 children = butler.collections.get_info("chain").children
1716 self.assertEqual(expected, list(children))
1718 def _test_common_chain_functionality(
1719 self, butler, func: Callable[[str, str | list[str]], Any], *, skip_cycle_check=False
1720 ) -> None:
1721 # Missing parent collection
1722 with self.assertRaises(MissingCollectionError):
1723 func("doesnotexist", [])
1724 # Missing child collection
1725 with self.assertRaises(MissingCollectionError):
1726 func("chain", ["doesnotexist"])
1727 # Forbid operations on non-chained collections
1728 with self.assertRaises(CollectionTypeError):
1729 func("d", ["a"])
1731 # Prevent collection cycles
1732 if not skip_cycle_check:
1733 butler.collections.register("chain2", CollectionType.CHAINED)
1734 func("chain2", "chain")
1735 with self.assertRaises(CollectionCycleError):
1736 func("chain", "chain2")
1738 # Make sure none of the earlier operations interfered with unrelated
1739 # chains.
1740 self.assertEqual(["a", "b"], list(butler.collections.get_info("staticchain").children))
1742 with butler._caching_context():
1743 with self.assertRaisesRegex(RuntimeError, "Chained collection modification not permitted"):
1744 func("chain", "a")
1746 def test_transfer_dimension_records_from(self) -> None:
1747 source_butler = self.create_empty_butler(writeable=True)
1748 source_butler.import_(filename=_get_test_data_path("lsstcam-subset.yaml"))
1750 visit_id = 2025120200439
1751 exposure_id = visit_id
1752 target_butler = self.enterContext(create_populated_sqlite_registry())
1753 target_butler.transfer_dimension_records_from(
1754 source_butler,
1755 [
1756 # Should trigger the lookup of visit and all its associated
1757 # "populated_by" records (visit_detector_region,
1758 # visit_definition, etc.)
1759 DataCoordinate.standardize(
1760 {"instrument": "LSSTCam", "visit": visit_id, "detector": 10},
1761 universe=source_butler.dimensions,
1762 ),
1763 # Shouldn't add any records to the lookup.
1764 DataCoordinate.make_empty(source_butler.dimensions),
1765 ],
1766 )
1768 def _fetch_record(dimension: str) -> DimensionRecord:
1769 records = target_butler.query_dimension_records(dimension)
1770 self.assertEqual(len(records), 1)
1771 return records[0]
1773 visit = _fetch_record("visit")
1774 self.assertEqual(visit.id, visit_id)
1775 self.assertEqual(visit.day_obs, 20251202)
1776 self.assertEqual(visit.target_name, "lowdust")
1777 self.assertEqual(visit.seq_num, 439)
1778 original_visit = source_butler.query_dimension_records("visit", instrument="LSSTCam", visit=visit_id)[
1779 0
1780 ]
1781 self.assertEqual(visit.region, original_visit.region)
1782 self.assertEqual(visit.timespan, original_visit.timespan)
1784 visit_detector_region = _fetch_record("visit_detector_region")
1785 self.assertEqual(visit_detector_region.instrument, "LSSTCam")
1786 self.assertEqual(visit_detector_region.detector, 10)
1787 self.assertEqual(visit_detector_region.visit, visit_id)
1788 original_visit_detector_region = source_butler.query_dimension_records(
1789 "visit_detector_region", instrument="LSSTCam", visit=visit_id, detector=10
1790 )[0]
1791 self.assertEqual(visit_detector_region.region, original_visit_detector_region.region)
1793 visit_definition = _fetch_record("visit_definition")
1794 self.assertEqual(visit_definition.instrument, "LSSTCam")
1795 self.assertEqual(visit_definition.exposure, 2025120200439)
1796 self.assertEqual(visit_definition.visit, visit_id)
1798 # The matching exposure record should have been pulled in via
1799 # visit -> visit_definition.
1800 exposure = _fetch_record("exposure")
1801 self.assertEqual(exposure.instrument, "LSSTCam")
1802 self.assertEqual(exposure.id, 2025120200439)
1803 self.assertEqual(exposure.obs_id, "MC_O_20251202_000439")
1804 original_exposure = source_butler.query_dimension_records(
1805 "exposure", instrument="LSSTCam", exposure=exposure_id
1806 )[0]
1807 self.assertEqual(exposure.timespan, original_exposure.timespan)
1809 group = _fetch_record("group")
1810 self.assertEqual(group.instrument, "LSSTCam")
1811 self.assertEqual(group.name, "2025-12-03T07:58:10.858")
1813 visit_system_memberships = target_butler.query_dimension_records("visit_system_membership")
1814 visit_system_memberships.sort(key=lambda record: record.visit_system)
1815 self.assertEqual(len(visit_system_memberships), 2)
1816 self.assertEqual(visit_system_memberships[0].visit_system, 0)
1817 self.assertEqual(visit_system_memberships[1].visit_system, 2)
1818 self.assertEqual(visit_system_memberships[0].visit, visit_id)
1819 self.assertEqual(visit_system_memberships[1].visit, visit_id)
1821 visit_systems = target_butler.query_dimension_records("visit_system")
1822 visit_systems.sort(key=lambda record: record.id)
1823 visit_system_memberships.sort(key=lambda record: record.visit_system)
1824 self.assertEqual(visit_systems[0].id, 0)
1825 self.assertEqual(visit_systems[1].id, 2)
1826 self.assertEqual(visit_systems[0].name, "one-to-one")
1827 self.assertEqual(visit_systems[1].name, "by-seq-start-end")
1830class FileDatastoreButlerTests(ButlerTests):
1831 """Common tests and specialization of ButlerTests for butlers backed
1832 by datastores that inherit from FileDatastore.
1833 """
1835 trustModeSupported = True
1837 def testComponentFromOverriddenStorageClassWarns(self) -> None:
1838 """Test that getting a component that only the read storage class
1839 defines warns, since the whole dataset has to be retrieved and
1840 converted before the component can be extracted.
1841 """
1842 write_sc = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
1843 read_sc = self.storageClassFactory.getStorageClass("MetricsConversion")
1844 butler, datasetType = self.create_butler(self.default_run, write_sc, "unstructured")
1845 metric = makeExampleMetrics()
1846 dataId = {"instrument": "DummyCamComp", "visit": 423}
1847 ref = butler.put(metric, datasetType, dataId)
1848 component_ref = ref.overrideStorageClass(read_sc).makeComponentRef("summary")
1850 logger = "lsst.daf.butler.datastores.file_datastore.get"
1851 with self.assertLogs(logger, level="WARNING") as cm:
1852 self.assertEqual(butler.get(component_ref), metric.summary)
1853 message = "\n".join(cm.output)
1854 # The message must name the component, the storage class that lacks it
1855 # along with the components it does have, and the storage class the
1856 # dataset has to be converted to.
1857 self.assertIn("summary", message)
1858 self.assertIn(write_sc.name, message)
1859 self.assertIn("components it does define: none", message)
1860 self.assertIn(read_sc.name, message)
1861 self.assertIn("less efficient", message)
1863 # Reading a component that the write storage class does define must not
1864 # warn.
1865 composite_type = self.addDatasetType(
1866 "composite", datasetType.dimensions, "StructuredData", butler.registry
1867 )
1868 composite_ref = butler.put(metric, composite_type, dataId)
1869 with self.assertNoLogs(logger, level="WARNING"):
1870 self.assertEqual(butler.get(composite_ref.makeComponentRef("summary")), metric.summary)
1872 def checkFileExists(self, root: str | ResourcePath, relpath: str | ResourcePath) -> bool:
1873 """Check if file exists at a given path (relative to root).
1875 Test testPutTemplates verifies actual physical existance of the files
1876 in the requested location.
1877 """
1878 uri = ResourcePath(root, forceDirectory=True)
1879 return uri.join(relpath).exists()
1881 def testPutTemplates(self) -> None:
1882 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
1883 butler = self.create_empty_butler(run=self.default_run)
1885 # Add needed Dimensions
1886 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"})
1887 butler.registry.insertDimensionData(
1888 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"}
1889 )
1890 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20250101})
1891 butler.registry.insertDimensionData(
1892 "visit",
1893 {
1894 "instrument": "DummyCamComp",
1895 "id": 423,
1896 "name": "v423",
1897 "physical_filter": "d-r",
1898 "day_obs": 20250101,
1899 },
1900 )
1901 butler.registry.insertDimensionData(
1902 "visit",
1903 {
1904 "instrument": "DummyCamComp",
1905 "id": 425,
1906 "name": "v425",
1907 "physical_filter": "d-r",
1908 "day_obs": 20250101,
1909 },
1910 )
1912 # Create and store a dataset
1913 metric = makeExampleMetrics()
1915 # Create two almost-identical DatasetTypes (both will use default
1916 # template)
1917 dimensions = butler.dimensions.conform(["instrument", "visit"])
1918 butler.registry.registerDatasetType(DatasetType("metric1", dimensions, storageClass))
1919 butler.registry.registerDatasetType(DatasetType("metric2", dimensions, storageClass))
1920 butler.registry.registerDatasetType(DatasetType("metric3", dimensions, storageClass))
1922 dataId1 = {"instrument": "DummyCamComp", "visit": 423}
1923 dataId2 = {"instrument": "DummyCamComp", "visit": 423, "physical_filter": "d-r"}
1925 # Put with exactly the data ID keys needed
1926 ref = butler.put(metric, "metric1", dataId1)
1927 uri = butler.getURI(ref)
1928 self.assertTrue(uri.exists())
1929 self.assertTrue(
1930 uri.unquoted_path.endswith(f"{self.default_run}/metric1/??#?/d-r/DummyCamComp_423.pickle")
1931 )
1933 # Check the template based on dimensions
1934 if hasattr(butler._datastore, "templates"):
1935 butler._datastore.templates.validateTemplates([ref])
1937 # Put with extra data ID keys (physical_filter is an optional
1938 # dependency); should not change template (at least the way we're
1939 # defining them to behave now; the important thing is that they
1940 # must be consistent).
1941 ref = butler.put(metric, "metric2", dataId2)
1942 uri = butler.getURI(ref)
1943 self.assertTrue(uri.exists())
1944 self.assertTrue(
1945 uri.unquoted_path.endswith(f"{self.default_run}/metric2/d-r/DummyCamComp_v423.pickle")
1946 )
1948 # Check the template based on dimensions
1949 if hasattr(butler._datastore, "templates"):
1950 butler._datastore.templates.validateTemplates([ref])
1952 # Use a template that has a typo in dimension record metadata.
1953 # Easier to test with a butler that has a ref with records attached.
1954 template = FileTemplate("a/{visit.name}/{id}_{visit.namex:?}.fits")
1955 with self.assertLogs("lsst.daf.butler.datastore.file_templates", "INFO"):
1956 path = template.format(ref)
1957 self.assertEqual(path, f"a/v423/{ref.id}_fits")
1959 template = FileTemplate("a/{visit.name}/{id}_{visit.namex}.fits")
1960 with self.assertRaises(KeyError):
1961 with self.assertLogs("lsst.daf.butler.datastore.file_templates", "INFO"):
1962 template.format(ref)
1964 # Now use a file template that will not result in unique filenames
1965 with self.assertRaises(FileTemplateValidationError):
1966 butler.put(metric, "metric3", dataId1)
1968 def testImportExport(self) -> None:
1969 # Run put/get tests just to create and populate a repo.
1970 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
1971 self.runImportExportTest(storageClass)
1973 @unittest.expectedFailure
1974 def testImportExportVirtualComposite(self) -> None:
1975 # Run put/get tests just to create and populate a repo.
1976 storageClass = self.storageClassFactory.getStorageClass("StructuredComposite")
1977 self.runImportExportTest(storageClass)
1979 def runImportExportTest(self, storageClass: StorageClass) -> None:
1980 """Test exporting and importing.
1982 This test does an export to a temp directory and an import back
1983 into a new temp directory repo. It does not assume a posix datastore.
1984 """
1985 exportButler = self.runPutGetTest(storageClass, "test_metric")
1987 # Test that we must have a file extension.
1988 with self.assertRaises(ValueError):
1989 with exportButler.export(filename="dump", directory=".") as export:
1990 pass
1992 # Test that unknown format is not allowed.
1993 with self.assertRaises(ValueError):
1994 with exportButler.export(filename="dump.fits", directory=".") as export:
1995 pass
1997 # Test that the repo actually has at least one dataset.
1998 datasets = list(exportButler.registry.queryDatasets(..., collections=...))
1999 self.assertGreater(len(datasets), 0)
2000 # Add a DimensionRecord that's unused by those datasets.
2001 skymapRecord = {"name": "example_skymap", "hash": (50).to_bytes(8, byteorder="little")}
2002 exportButler.registry.insertDimensionData("skymap", skymapRecord)
2003 # Export and then import datasets.
2004 with safeTestTempDir(TESTDIR) as exportDir:
2005 exportFile = os.path.join(exportDir, "exports.yaml")
2006 with exportButler.export(filename=exportFile, directory=exportDir, transfer="auto") as export:
2007 export.saveDatasets(datasets)
2008 # Export the same datasets again. This should quietly do
2009 # nothing because of internal deduplication, and it shouldn't
2010 # complain about being asked to export the "htm7" elements even
2011 # though there aren't any in these datasets or in the database.
2012 export.saveDatasets(datasets, elements=["htm7"])
2013 # Save one of the data IDs again; this should be harmless
2014 # because of internal deduplication.
2015 export.saveDataIds([datasets[0].dataId])
2016 # Save some dimension records directly.
2017 export.saveDimensionData("skymap", [skymapRecord])
2018 self.assertTrue(os.path.exists(exportFile))
2019 with safeTestTempDir(TESTDIR) as importDir:
2020 # We always want this to be a local posix butler
2021 make_repo_for_test(
2022 importDir, config=Config(os.path.join(TESTDIR, "config/basic/butler.yaml"))
2023 )
2024 # Calling script.butlerImport tests the implementation of the
2025 # butler command line interface "import" subcommand. Functions
2026 # in the script folder are generally considered protected and
2027 # should not be used as public api.
2028 with open(exportFile) as f:
2029 script.butlerImport(
2030 importDir,
2031 export_file=f,
2032 directory=exportDir,
2033 transfer="auto",
2034 skip_dimensions=None,
2035 )
2036 importButler = Butler.from_config(importDir, run=self.default_run)
2037 self.enterContext(importButler)
2038 for ref in datasets:
2039 with self.subTest(ref=repr(ref)):
2040 # Test for existence by passing in the DatasetType and
2041 # data ID separately, to avoid lookup by dataset_id.
2042 self.assertTrue(importButler.exists(ref.datasetType, ref.dataId))
2043 self.assertEqual(
2044 list(importButler.registry.queryDimensionRecords("skymap")),
2045 [importButler.dimensions["skymap"].RecordClass(**skymapRecord)],
2046 )
2048 def testRemoveRuns(self) -> None:
2049 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
2050 butler = self.create_empty_butler(writeable=True)
2051 # Load registry data with dimensions to hang datasets off of.
2052 butler.import_(filename=ResourcePath("resource://lsst.daf.butler/tests/registry_data/base.yaml"))
2053 # Add some RUN-type collection.
2054 run1 = "run1"
2055 butler.collections.register(run1)
2056 run2 = "run2"
2057 butler.collections.register(run2)
2058 # put a dataset in each
2059 metric = makeExampleMetrics()
2060 dimensions = butler.dimensions.conform(["instrument", "physical_filter"])
2061 datasetType = self.addDatasetType(
2062 "prune_collections_test_dataset", dimensions, storageClass, butler.registry
2063 )
2064 ref1 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run1)
2065 ref2 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run2)
2066 uri1 = butler.getURI(ref1)
2067 uri2 = butler.getURI(ref2)
2069 # Put one of the runs in a chain.
2070 butler.collections.register("Chain", CollectionType.CHAINED)
2071 butler.collections.extend_chain("Chain", run1)
2073 with self.assertRaises(OrphanedRecordError):
2074 butler.registry.removeDatasetType(datasetType.name)
2076 # Remove a non-run.
2077 with self.assertRaises(TypeError):
2078 butler.removeRuns(["Chain"])
2080 # Remove without unlinking from chain should fail.
2081 with self.assertRaises(IntegrityError):
2082 butler.removeRuns([run1])
2084 # Remove from both runs. No longer use unstore parameter since it
2085 # always purges.
2086 butler.removeRuns([run1, run2], unlink_from_chains=True)
2088 # Should be nothing in registry for either one, and datastore should
2089 # not think either exists.
2090 with self.assertRaises(MissingCollectionError):
2091 butler.collections.get_info(run1)
2092 with self.assertRaises(MissingCollectionError):
2093 butler.collections.get_info(run1)
2094 self.assertFalse(butler.stored(ref1))
2095 self.assertFalse(butler.stored(ref2))
2096 # We always unstore so both URIs should be gone.
2097 self.assertFalse(uri1.exists())
2098 self.assertFalse(uri2.exists())
2100 # Now that the collections have been pruned we can remove the
2101 # dataset type
2102 butler.registry.removeDatasetType(datasetType.name)
2104 with self.assertLogs("lsst.daf.butler.registry", "INFO") as cm:
2105 butler.registry.removeDatasetType(("test*", "test*"))
2106 self.assertIn("not defined", "\n".join(cm.output))
2108 def remove_dataset_out_of_band(self, butler: Butler, ref: DatasetRef) -> None:
2109 """Simulate an external actor removing a file outside of Butler's
2110 knowledge.
2112 Subclasses may override to handle more complicated datastore
2113 configurations.
2114 """
2115 uri = butler.getURI(ref)
2116 uri.remove()
2117 datastore = cast(FileDatastore, butler._datastore)
2118 datastore.cacheManager.remove_from_cache(ref)
2120 def testPruneDatasets(self) -> None:
2121 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
2122 butler = self.create_empty_butler(writeable=True)
2123 # Load registry data with dimensions to hang datasets off of.
2124 butler.import_(filename=_get_test_data_path("base.yaml"))
2125 # Add some RUN-type collections.
2126 run1 = "run1"
2127 butler.collections.register(run1)
2128 run2 = "run2"
2129 butler.collections.register(run2)
2130 # put some datasets. ref1 and ref2 have the same data ID, and are in
2131 # different runs. ref3 has a different data ID.
2132 metric = makeExampleMetrics()
2133 dimensions = butler.dimensions.conform(["instrument", "physical_filter"])
2134 datasetType = self.addDatasetType(
2135 "prune_collections_test_dataset", dimensions, storageClass, butler.registry
2136 )
2137 ref1 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run1)
2138 ref2 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run2)
2139 ref3 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-R1"}, run=run1)
2141 many_stored = butler.stored_many([ref1, ref2, ref3])
2142 for ref, stored in many_stored.items():
2143 self.assertTrue(stored, f"Ref {ref} should be stored")
2145 many_exists = butler._exists_many([ref1, ref2, ref3])
2146 for ref, exists in many_exists.items():
2147 self.assertTrue(exists, f"Checking ref {ref} exists.")
2148 self.assertEqual(exists, DatasetExistence.VERIFIED, f"Ref {ref} should be stored")
2150 # Simple prune.
2151 butler.pruneDatasets([ref1, ref2, ref3], purge=True, unstore=True)
2152 self.assertFalse(butler.exists(ref1.datasetType, ref1.dataId, collections=run1))
2154 many_stored = butler.stored_many([ref1, ref2, ref3])
2155 for ref, stored in many_stored.items():
2156 self.assertFalse(stored, f"Ref {ref} should not be stored")
2158 many_exists = butler._exists_many([ref1, ref2, ref3])
2159 for ref, exists in many_exists.items():
2160 self.assertEqual(exists, DatasetExistence.UNRECOGNIZED, f"Ref {ref} should not be stored")
2162 # Put data back.
2163 ref1_new = butler.put(metric, ref1)
2164 self.assertEqual(ref1_new, ref1) # Reuses original ID.
2165 ref2 = butler.put(metric, ref2)
2167 many_stored = butler.stored_many([ref1, ref2, ref3])
2168 self.assertTrue(many_stored[ref1])
2169 self.assertTrue(many_stored[ref2])
2170 self.assertFalse(many_stored[ref3])
2172 ref3 = butler.put(metric, ref3)
2174 many_exists = butler._exists_many([ref1, ref2, ref3])
2175 for ref, exists in many_exists.items():
2176 self.assertTrue(exists, f"Ref {ref} should not be stored")
2178 # Clear out the datasets from registry and start again.
2179 refs = [ref1, ref2, ref3]
2180 butler.pruneDatasets(refs, purge=True, unstore=True)
2181 for ref in refs:
2182 butler.put(metric, ref)
2184 # Confirm we can retrieve deferred.
2185 dref1 = butler.getDeferred(ref1) # known and exists
2186 metric1 = dref1.get()
2187 self.assertEqual(metric1, metric)
2189 # Test different forms of file availability.
2190 # Need to be in a state where:
2191 # - one ref just has registry record.
2192 # - one ref has a missing file but a datastore record.
2193 # - one ref has a missing datastore record but file is there.
2194 # - one ref does not exist anywhere.
2195 # Do not need to test a ref that has everything since that is tested
2196 # above.
2197 ref0 = DatasetRef(
2198 datasetType,
2199 DataCoordinate.standardize(
2200 {"instrument": "Cam1", "physical_filter": "Cam1-G"}, universe=butler.dimensions
2201 ),
2202 run=run1,
2203 )
2205 # Delete from datastore and retain in Registry.
2206 butler.pruneDatasets([ref1], purge=False, unstore=True, disassociate=False)
2208 # File has been removed.
2209 self.remove_dataset_out_of_band(butler, ref2)
2211 # Datastore has lost track.
2212 butler._datastore.forget([ref3])
2214 # First test with a standard butler.
2215 exists_many = butler._exists_many([ref0, ref1, ref2, ref3], full_check=True)
2216 self.assertEqual(exists_many[ref0], DatasetExistence.UNRECOGNIZED)
2217 self.assertEqual(exists_many[ref1], DatasetExistence.RECORDED)
2218 self.assertEqual(exists_many[ref2], DatasetExistence.RECORDED | DatasetExistence.DATASTORE)
2219 self.assertEqual(exists_many[ref3], DatasetExistence.RECORDED)
2221 exists_many = butler._exists_many([ref0, ref1, ref2, ref3], full_check=False)
2222 self.assertEqual(exists_many[ref0], DatasetExistence.UNRECOGNIZED)
2223 self.assertEqual(exists_many[ref1], DatasetExistence.RECORDED | DatasetExistence._ASSUMED)
2224 self.assertEqual(exists_many[ref2], DatasetExistence.KNOWN)
2225 self.assertEqual(exists_many[ref3], DatasetExistence.RECORDED | DatasetExistence._ASSUMED)
2226 self.assertTrue(exists_many[ref2])
2228 # Check that per-ref query gives the same answer as many query.
2229 for ref, exists in exists_many.items():
2230 self.assertEqual(butler.exists(ref, full_check=False), exists)
2232 # Get deferred checks for existence before it allows it to be
2233 # retrieved.
2234 with self.assertRaises(LookupError):
2235 butler.getDeferred(ref3) # not known, file exists
2236 dref2 = butler.getDeferred(ref2) # known but file missing
2237 with self.assertRaises(FileNotFoundError):
2238 dref2.get()
2240 # Test again with a trusting butler.
2241 if self.trustModeSupported: 2241 ↛ exitline 2241 didn't return from function 'testPruneDatasets' because the condition on line 2241 was always true
2242 butler._datastore.trustGetRequest = True
2243 exists_many = butler._exists_many([ref0, ref1, ref2, ref3], full_check=True)
2244 self.assertEqual(exists_many[ref0], DatasetExistence.UNRECOGNIZED)
2245 self.assertEqual(exists_many[ref1], DatasetExistence.RECORDED)
2246 self.assertEqual(exists_many[ref2], DatasetExistence.RECORDED | DatasetExistence.DATASTORE)
2247 self.assertEqual(exists_many[ref3], DatasetExistence.RECORDED | DatasetExistence._ARTIFACT)
2249 # When trusting we can get a deferred dataset handle that is not
2250 # known but does exist.
2251 dref3 = butler.getDeferred(ref3)
2252 metric3 = dref3.get()
2253 self.assertEqual(metric3, metric)
2255 # Check that per-ref query gives the same answer as many query.
2256 for ref, exists in exists_many.items():
2257 self.assertEqual(butler.exists(ref, full_check=True), exists)
2259 # Create a ref that surprisingly has the UUID of an existing ref
2260 # but is not the same.
2261 ref_bad = DatasetRef(datasetType, dataId=ref3.dataId, run=ref3.run, id=ref2.id)
2262 with self.assertRaises(ValueError):
2263 butler.exists(ref_bad)
2265 # Create a ref that has a compatible storage class.
2266 ref_compat = ref2.overrideStorageClass("StructuredDataDict")
2267 exists = butler.exists(ref_compat)
2268 self.assertEqual(exists, exists_many[ref2])
2270 # Remove everything and start from scratch.
2271 butler._datastore.trustGetRequest = False
2272 butler.pruneDatasets(refs, purge=True, unstore=True)
2273 for ref in refs:
2274 butler.put(metric, ref)
2276 # These tests mess directly with the trash table and can leave the
2277 # datastore in an odd state. Do them at the end.
2278 # Check that in normal mode, deleting the record will lead to
2279 # trash not touching the file.
2280 uri1 = butler.getURI(ref1)
2281 butler._datastore.bridge.moveToTrash(
2282 [ref1], transaction=None
2283 ) # Update the dataset_location table
2284 butler._datastore.forget([ref1])
2285 butler._datastore.trash(ref1)
2286 butler._datastore.emptyTrash()
2287 self.assertTrue(uri1.exists())
2288 uri1.remove() # Clean it up.
2290 # Simulate execution butler setup by deleting the datastore
2291 # record but keeping the file around and trusting.
2292 butler._datastore.trustGetRequest = True
2293 uris = butler.get_many_uris([ref2, ref3])
2294 uri2 = uris[ref2].primaryURI
2295 uri3 = uris[ref3].primaryURI
2296 self.assertTrue(uri2.exists())
2297 self.assertTrue(uri3.exists())
2299 # Remove the datastore record.
2300 butler._datastore.bridge.moveToTrash(
2301 [ref2], transaction=None
2302 ) # Update the dataset_location table
2303 butler._datastore.forget([ref2])
2304 self.assertTrue(uri2.exists())
2305 butler._datastore.trash([ref2, ref3])
2306 # Immediate removal for ref2 file
2307 self.assertFalse(uri2.exists())
2308 # But ref3 has to wait for the empty.
2309 self.assertTrue(uri3.exists())
2310 butler._datastore.emptyTrash()
2311 self.assertFalse(uri3.exists())
2313 # Clear out the datasets from registry.
2314 butler.pruneDatasets([ref1, ref2, ref3], purge=True, unstore=True)
2316 def test_butler_metrics(self):
2317 """Test that metrics are collected."""
2318 run = "test_run"
2319 metrics = ButlerMetrics()
2320 butler, datasetType = self.create_butler(
2321 run, "MetricsExampleModelProvenance", "prov_metric", metrics=metrics
2322 )
2323 data = MetricsExampleModel(
2324 summary={"AM1": 5.2, "AM2": 30.6},
2325 output={"a": [1, 2, 3], "b": {"blue": 5, "red": "green"}},
2326 data=[563, 234, 456.7, 752, 8, 9, 27],
2327 )
2329 data_ref = butler.put(data, datasetType, visit=424, instrument="DummyCamComp")
2330 butler.get(data_ref)
2331 butler.get(data_ref)
2332 self.assertEqual(metrics.n_get, 2)
2333 self.assertGreater(metrics.time_in_get, 0.0)
2334 self.assertEqual(metrics.n_put, 1)
2335 self.assertGreater(metrics.time_in_put, 0.0)
2337 deferred = butler.getDeferred(data_ref)
2338 deferred.get()
2339 self.assertEqual(metrics.n_get, 3)
2341 with butler.record_metrics() as new:
2342 data_ref_2 = butler.put(data, datasetType, visit=425, instrument="DummyCamComp")
2343 butler.get(data_ref)
2345 butler.pruneDatasets([data_ref, data_ref_2], purge=True, unstore=True)
2346 with ResourcePath.temporary_uri(suffix=".json") as tmpFile:
2347 tmpFile.write(data.model_dump_json().encode())
2348 refs = [
2349 DatasetRef(datasetType, data_ref_2.dataId, run),
2350 DatasetRef(datasetType, data_ref.dataId, run),
2351 ]
2352 datasets = [FileDataset(path=tmpFile, refs=refs)]
2353 butler.ingest(*datasets, transfer="copy")
2355 self.assertEqual(new.n_get, 1)
2356 self.assertEqual(new.n_put, 1)
2357 self.assertEqual(new.n_ingest, 2)
2360class PosixDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase):
2361 """PosixDatastore specialization of a butler"""
2363 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
2364 fullConfigKey: str | None = ".datastore.formatters"
2365 validationCanFail = True
2366 datastoreStr = ["/tmp"]
2367 datastoreName = [f"FileDatastore@{BUTLER_ROOT_TAG}"]
2368 registryStr = "/gen3.sqlite3"
2370 def testPathConstructor(self) -> None:
2371 """Independent test of constructor using PathLike."""
2372 butler = Butler.from_config(self.tmpConfigFile, run=self.default_run)
2373 self.enterContext(butler)
2374 self.assertIsInstance(butler, Butler)
2376 # And again with a Path object with the butler yaml
2377 path = pathlib.Path(self.tmpConfigFile)
2378 butler = Butler.from_config(path, writeable=False)
2379 self.enterContext(butler)
2380 self.assertIsInstance(butler, Butler)
2382 # And again with a Path object without the butler yaml
2383 # (making sure we skip it if the tmp config doesn't end
2384 # in butler.yaml -- which is the case for a subclass)
2385 if self.tmpConfigFile.endswith("butler.yaml"):
2386 path = pathlib.Path(os.path.dirname(self.tmpConfigFile))
2387 butler = Butler.from_config(path, writeable=False)
2388 self.enterContext(butler)
2389 self.assertIsInstance(butler, Butler)
2391 def testExportTransferCopy(self) -> None:
2392 """Test local export using all transfer modes"""
2393 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
2394 exportButler = self.runPutGetTest(storageClass, "test_metric")
2395 # Test that the repo actually has at least one dataset.
2396 datasets = list(exportButler.registry.queryDatasets(..., collections=...))
2397 self.assertGreater(len(datasets), 0)
2398 uris = [exportButler.getURI(d) for d in datasets]
2399 assert isinstance(exportButler._datastore, FileDatastore)
2400 datastoreRoot = exportButler.get_datastore_roots()[exportButler.get_datastore_names()[0]]
2402 pathsInStore = [uri.relative_to(datastoreRoot) for uri in uris]
2404 for path in pathsInStore:
2405 # Assume local file system
2406 assert path is not None
2407 self.assertTrue(self.checkFileExists(datastoreRoot, path), f"Checking path {path}")
2409 for transfer in ("copy", "link", "symlink", "relsymlink"):
2410 with safeTestTempDir(TESTDIR) as exportDir:
2411 with exportButler.export(directory=exportDir, format="yaml", transfer=transfer) as export:
2412 export.saveDatasets(datasets)
2413 for path in pathsInStore:
2414 assert path is not None
2415 self.assertTrue(
2416 self.checkFileExists(exportDir, path),
2417 f"Check that mode {transfer} exported files",
2418 )
2420 def testPytypeCoercion(self) -> None:
2421 """Test python type coercion on Butler.get and put."""
2422 # Store some data with the normal example storage class.
2423 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
2424 datasetTypeName = "test_metric"
2425 butler = self.runPutGetTest(storageClass, datasetTypeName)
2427 dataId = {"instrument": "DummyCamComp", "visit": 423}
2428 metric = butler.get(datasetTypeName, dataId=dataId)
2429 self.assertEqual(get_full_type_name(metric), "lsst.daf.butler.tests.MetricsExample")
2431 datasetType_ori = butler.get_dataset_type(datasetTypeName)
2432 self.assertEqual(datasetType_ori.storageClass.name, "StructuredDataNoComponents")
2434 # Now need to hack the registry dataset type definition.
2435 # There is no API for this.
2436 assert isinstance(butler._registry, SqlRegistry)
2437 manager = butler._registry._managers.datasets
2438 assert hasattr(manager, "_db") and hasattr(manager, "_static")
2439 manager._db.update(
2440 manager._static.dataset_type,
2441 {"name": datasetTypeName},
2442 {datasetTypeName: datasetTypeName, "storage_class": "StructuredDataNoComponentsModel"},
2443 )
2445 # Force reset of dataset type cache
2446 butler.registry.refresh()
2448 datasetType_new = butler.get_dataset_type(datasetTypeName)
2449 self.assertEqual(datasetType_new.name, datasetType_ori.name)
2450 self.assertEqual(datasetType_new.storageClass.name, "StructuredDataNoComponentsModel")
2452 metric_model = butler.get(datasetTypeName, dataId=dataId)
2453 self.assertNotEqual(type(metric_model), type(metric))
2454 self.assertEqual(get_full_type_name(metric_model), "lsst.daf.butler.tests.MetricsExampleModel")
2456 # Put the model and read it back to show that everything now
2457 # works as normal.
2458 metric_ref = butler.put(metric_model, datasetTypeName, dataId=dataId, visit=424)
2459 metric_model_new = butler.get(metric_ref)
2460 self.assertEqual(metric_model_new, metric_model)
2462 # Hack the storage class again to something that will fail on the
2463 # get with no conversion class.
2464 manager._db.update(
2465 manager._static.dataset_type,
2466 {"name": datasetTypeName},
2467 {datasetTypeName: datasetTypeName, "storage_class": "StructuredDataListYaml"},
2468 )
2469 butler.registry.refresh()
2471 with self.assertRaises(ValueError):
2472 butler.get(datasetTypeName, dataId=dataId)
2474 def test_provenance(self):
2475 """Test that provenance is attached on put."""
2476 run = "test_run"
2477 butler, datasetType = self.create_butler(run, "MetricsExampleModelProvenance", "prov_metric")
2478 metric = MetricsExampleModel(
2479 summary={"AM1": 5.2, "AM2": 30.6},
2480 output={"a": [1, 2, 3], "b": {"blue": 5, "red": "green"}},
2481 data=[563, 234, 456.7, 752, 8, 9, 27],
2482 )
2483 # Provenance can be attached to the object being put. Whether
2484 # it is or not is dependent on the formatter. For this test we
2485 # copy on adding provenance to ensure they differ.
2486 self.assertIsNone(metric.dataset_id)
2487 metric_ref = butler.put(metric, datasetType, visit=424, instrument="DummyCamComp")
2488 self.assertIsNone(metric.dataset_id)
2489 metric_2 = butler.get(metric_ref)
2490 self.assertEqual(metric_2.data, metric.data)
2491 self.assertEqual(metric_2.dataset_id, metric_ref.id)
2492 self.assertIsNone(metric_2.provenance)
2494 # Put with provenance.
2495 prov = DatasetProvenance(quantum_id=uuid.uuid4())
2496 prov.add_input(metric_ref)
2497 prov.add_extra_provenance(metric_ref.id, {"answer": 42})
2498 metric_ref2 = butler.put(metric, datasetType, visit=423, instrument="DummyCamComp", provenance=prov)
2499 metric_3 = butler.get(metric_ref2)
2500 self.assertEqual(metric_3.provenance, prov)
2502 # Check that we can extract provenance from dict form.
2503 prov_dict = prov.to_flat_dict(metric_ref2)
2504 prov_from_prov, ref_from_prov = DatasetProvenance.from_flat_dict(prov_dict, butler)
2505 self.assertEqual(ref_from_prov, metric_ref2)
2506 # Direct __eq__ of the provenance does not work because one side
2507 # includes dimension records.
2508 self.assertEqual({ref.id for ref in prov_from_prov.inputs}, {ref.id for ref in prov.inputs})
2509 self.assertEqual(prov_from_prov.quantum_id, prov.quantum_id)
2510 self.assertEqual(prov_from_prov.extras, prov.extras)
2512 # Force a bad ID into the dict.
2513 prov_dict["id"] = uuid.uuid4()
2514 with self.assertRaises(ValueError):
2515 DatasetProvenance.from_flat_dict(prov_dict, butler)
2516 del prov_dict["id"]
2517 prov_dict["input 0 id"] = uuid.uuid4()
2518 with self.assertRaises(ValueError):
2519 DatasetProvenance.from_flat_dict(prov_dict, butler)
2521 # Check that simple types can be reconstructed with non-standard
2522 # separators.
2523 prov_dict = prov.to_flat_dict(metric_ref2, prefix="XYZ", sep="😎", simple_types=True)
2524 prov_from_prov, ref_from_prov = DatasetProvenance.from_flat_dict(prov_dict, butler)
2525 self.assertEqual(ref_from_prov, metric_ref2)
2526 self.assertEqual({ref.id for ref in prov_from_prov.inputs}, {ref.id for ref in prov.inputs})
2528 with self.assertRaises(ValueError):
2529 DatasetProvenance.from_flat_dict({"unknown": 42}, butler)
2531 def test_specialized_file_datasets_functions(self):
2532 """Test a workflow used in Prompt Processing where we export datasets
2533 from one repository and write them in-place to the datastore of
2534 another, without immediately inserting registry entries for the
2535 datasets.
2536 """
2537 repo = MetricTestRepo.create_from_butler(
2538 self.create_empty_butler(writeable=True),
2539 self.tmpConfigFile,
2540 "StructuredCompositeReadCompNoDisassembly",
2541 )
2542 source_butler = repo.butler
2544 # Test writing outputs to a FileDatastore.
2545 with tempfile.TemporaryDirectory() as tempdir:
2546 target_repo_config = make_repo_for_test(tempdir)
2547 refs = [repo.ref1, repo.ref2]
2548 datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), refs)
2549 self.assertEqual(len(datasets), 2)
2550 self.assertEqual({ref.id for ref in refs}, {dataset.refs[0].id for dataset in datasets})
2551 for dataset in datasets:
2552 path = ResourcePath(dataset.path, forceAbsolute=False)
2553 # Paths should be relative paths to the target datastore.
2554 self.assertFalse(path.isabs())
2555 # Files should have been copied into the target datastore
2556 self.assertTrue(ResourcePath(tempdir).join(path).exists())
2558 # Make sure the target Butler can ingest the datasets.
2559 target_butler = Butler(target_repo_config, writeable=True)
2560 self.enterContext(target_butler)
2561 target_butler.transfer_dimension_records_from(source_butler, refs)
2562 target_butler.ingest(*datasets, transfer=None)
2563 self.assertIsNotNone(target_butler.get(repo.ref1))
2564 self.assertIsNotNone(target_butler.get(repo.ref2))
2566 # Giving an empty list of files is a no-op.
2567 no_datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), [])
2568 self.assertEqual(len(no_datasets), 0)
2570 # Test writing outputs to a ChainedDatastore.
2571 with tempfile.TemporaryDirectory() as tempdir:
2572 # Set up a second dataset type, so we can split the files across
2573 # multiple datastore roots.
2574 dt1 = repo.datasetType
2575 dt2 = DatasetType("other", dt1.dimensions, dt1.storageClass)
2576 source_butler.registry.registerDatasetType(dt2)
2577 other_ref = repo.addDataset(repo.ref1.dataId, datasetType=dt2)
2578 config = Config.fromString(
2579 f"""
2580 datastore:
2581 cls: lsst.daf.butler.datastores.chainedDatastore.ChainedDatastore
2582 datastore_constraints:
2583 - constraints:
2584 accept:
2585 - {dt1.name}
2586 - constraints:
2587 accept:
2588 - {dt2.name}
2589 datastores:
2590 - datastore:
2591 cls: lsst.daf.butler.datastores.fileDatastore.FileDatastore
2592 root: <butlerRoot>/FileDatastore_0
2593 - datastore:
2594 cls: lsst.daf.butler.datastores.fileDatastore.FileDatastore
2595 root: <butlerRoot>/FileDatastore_1
2596 """
2597 )
2598 target_repo_config = make_repo_for_test(tempdir, config)
2599 refs = [repo.ref1, repo.ref2, other_ref]
2600 datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), refs)
2601 self.assertEqual(len(datasets), 3)
2602 self.assertEqual({ref.id for ref in refs}, {dataset.refs[0].id for dataset in datasets})
2603 for dataset in datasets:
2604 path = ResourcePath(dataset.path, forceAbsolute=False)
2605 # Paths should be relative paths to the target datastore.
2606 self.assertFalse(path.isabs())
2607 # Files should have been split up between the two datastores
2608 # in the chain.
2609 datastore_root = ResourcePath(tempdir)
2610 if dataset.refs[0].datasetType.name == dt1.name:
2611 datastore_root = datastore_root.join("FileDatastore_0")
2612 else:
2613 datastore_root = datastore_root.join("FileDatastore_1")
2614 self.assertTrue(datastore_root.join(path).exists())
2616 # Make sure the target Butler can ingest the datasets.
2617 target_butler = Butler(target_repo_config, writeable=True)
2618 self.enterContext(target_butler)
2619 target_butler.transfer_dimension_records_from(source_butler, refs)
2620 target_butler.ingest(*datasets, transfer=None)
2621 self.assertIsNotNone(target_butler.get(repo.ref1))
2622 self.assertIsNotNone(target_butler.get(repo.ref2))
2623 self.assertIsNotNone(target_butler.get(other_ref))
2625 def test_temporary_for_ingest(self) -> None:
2626 """Test the `lsst.daf.butler._rubin.ingest_from_temporary` module."""
2627 with self.create_empty_butler("example_run") as butler:
2628 dataset_type = DatasetType("example", butler.dimensions.empty, "StructuredDataDict")
2629 butler.registry.registerDatasetType(dataset_type)
2630 ref = DatasetRef(dataset_type, DataCoordinate.make_empty(butler.dimensions), "example_run")
2631 with TemporaryForIngest(butler, ref) as temporary:
2632 temporary.path.write(b"three: 3")
2633 found = TemporaryForIngest.find_orphaned_temporaries_by_ref(ref, butler)
2634 self.assertEqual(found, [temporary.path])
2635 self.assertIn(".tmp", temporary.ospath)
2636 temporary.ingest()
2637 loaded = butler.get(ref)
2638 self.assertEqual(loaded, {"three": 3})
2641class PostgresPosixDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase):
2642 """PosixDatastore specialization of a butler using Postgres"""
2644 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
2645 fullConfigKey = ".datastore.formatters"
2646 validationCanFail = True
2647 datastoreStr = ["/tmp"]
2648 datastoreName = [f"FileDatastore@{BUTLER_ROOT_TAG}"]
2649 registryStr = "PostgreSQL@test"
2651 @classmethod
2652 def setUpClass(cls) -> None:
2653 cls.postgresql = cls.enterClassContext(setup_postgres_test_db())
2654 super().setUpClass()
2656 def setUp(self) -> None:
2657 # Need to add a registry section to the config.
2658 self._temp_config = False
2659 config = Config(self.configFile)
2660 self.postgresql.patch_butler_config(config)
2661 with tempfile.NamedTemporaryFile("w", suffix=".yaml", delete=False) as fh:
2662 config.dump(fh)
2663 self.configFile = fh.name
2664 self._temp_config = True
2665 super().setUp()
2667 def tearDown(self) -> None:
2668 if self._temp_config and os.path.exists(self.configFile):
2669 os.remove(self.configFile)
2670 super().tearDown()
2672 def testMakeRepo(self) -> None:
2673 # The base class test assumes that it's using sqlite and assumes
2674 # the config file is acceptable to sqlite.
2675 raise unittest.SkipTest("Postgres config is not compatible with this test.")
2678class ClonedPostgresPosixDatastoreButlerTestCase(PostgresPosixDatastoreButlerTestCase, unittest.TestCase):
2679 """Test that Butler with a Postgres registry still works after cloning."""
2681 def create_butler(
2682 self,
2683 run: str,
2684 storageClass: StorageClass | str,
2685 datasetTypeName: str,
2686 metrics: ButlerMetrics | None = None,
2687 ) -> tuple[DirectButler, DatasetType]:
2688 butler, datasetType = super().create_butler(run, storageClass, datasetTypeName, metrics=metrics)
2689 return butler.clone(run=run, metrics=metrics), datasetType
2692class InMemoryDatastoreButlerTestCase(ButlerTests, unittest.TestCase):
2693 """InMemoryDatastore specialization of a butler"""
2695 configFile = os.path.join(TESTDIR, "config/basic/butler-inmemory.yaml")
2696 fullConfigKey = None
2697 useTempRoot = False
2698 validationCanFail = False
2699 datastoreStr = ["datastore='InMemory"]
2700 datastoreName = ["InMemoryDatastore@"]
2701 registryStr = "/gen3.sqlite3"
2703 def testIngest(self) -> None:
2704 pass
2706 def test_ingest_zip(self) -> None:
2707 pass
2710class ClonedSqliteButlerTestCase(InMemoryDatastoreButlerTestCase, unittest.TestCase):
2711 """Test that a Butler with a Sqlite registry still works after cloning."""
2713 def create_butler(
2714 self,
2715 run: str,
2716 storageClass: StorageClass | str,
2717 datasetTypeName: str,
2718 metrics: ButlerMetrics | None = None,
2719 ) -> tuple[DirectButler, DatasetType]:
2720 butler, datasetType = super().create_butler(run, storageClass, datasetTypeName, metrics=metrics)
2721 return butler.clone(run=run), datasetType
2724class ChainedDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase):
2725 """PosixDatastore specialization"""
2727 configFile = os.path.join(TESTDIR, "config/basic/butler-chained.yaml")
2728 fullConfigKey = ".datastore.datastores.1.formatters"
2729 validationCanFail = True
2730 datastoreStr = ["datastore='InMemory", "/FileDatastore_1/,", "/FileDatastore_2/'"]
2731 datastoreName = [
2732 "InMemoryDatastore@",
2733 f"FileDatastore@{BUTLER_ROOT_TAG}/FileDatastore_1",
2734 "SecondDatastore",
2735 ]
2736 registryStr = "/gen3.sqlite3"
2738 def testPruneDatasets(self) -> None:
2739 # This test relies on manipulating files out-of-band, which is
2740 # impossible for this configuration because of the InMemoryDatastore in
2741 # the ChainedDatastore.
2742 pass
2744 def testComponentFromOverriddenStorageClassWarns(self) -> None:
2745 # The InMemoryDatastore in the ChainedDatastore satisfies the get, so
2746 # the FileDatastore warning about having to read the whole dataset to
2747 # extract the component is never issued.
2748 pass
2751class ButlerExplicitRootTestCase(PosixDatastoreButlerTestCase):
2752 """Test that a yaml file in one location can refer to a root in another."""
2754 datastoreStr = ["dir1"]
2755 # Disable the makeRepo test since we are deliberately not using
2756 # butler.yaml as the config name.
2757 fullConfigKey = None
2759 def setUp(self) -> None:
2760 self.root = makeTestTempDir(TESTDIR)
2762 # Make a new repository in one place
2763 self.dir1 = os.path.join(self.root, "dir1")
2764 make_repo_for_test(self.dir1, config=Config(self.configFile))
2766 # Move the yaml file to a different place and add a "root"
2767 self.dir2 = os.path.join(self.root, "dir2")
2768 os.makedirs(self.dir2, exist_ok=True)
2769 configFile1 = os.path.join(self.dir1, "butler.yaml")
2770 config = Config(configFile1)
2771 config["root"] = self.dir1
2772 configFile2 = os.path.join(self.dir2, "butler2.yaml")
2773 config.dumpToUri(configFile2)
2774 os.remove(configFile1)
2775 self.tmpConfigFile = configFile2
2777 def testFileLocations(self) -> None:
2778 self.assertNotEqual(self.dir1, self.dir2)
2779 self.assertTrue(os.path.exists(os.path.join(self.dir2, "butler2.yaml")))
2780 self.assertFalse(os.path.exists(os.path.join(self.dir1, "butler.yaml")))
2781 self.assertTrue(os.path.exists(os.path.join(self.dir1, "gen3.sqlite3")))
2784class ButlerMakeRepoOutfileTestCase(ButlerPutGetTests, unittest.TestCase):
2785 """Test that a config file created by makeRepo outside of repo works."""
2787 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
2789 def setUp(self) -> None:
2790 self.root = makeTestTempDir(TESTDIR)
2791 self.root2 = makeTestTempDir(TESTDIR)
2793 self.tmpConfigFile = os.path.join(self.root2, "different.yaml")
2794 make_repo_for_test(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile)
2796 def tearDown(self) -> None:
2797 if os.path.exists(self.root2): 2797 ↛ 2799line 2797 didn't jump to line 2799 because the condition on line 2797 was always true
2798 shutil.rmtree(self.root2, ignore_errors=True)
2799 super().tearDown()
2801 def testConfigExistence(self) -> None:
2802 c = Config(self.tmpConfigFile)
2803 uri_config = ResourcePath(c["root"])
2804 uri_expected = ResourcePath(self.root, forceDirectory=True)
2805 self.assertEqual(uri_config.geturl(), uri_expected.geturl())
2806 self.assertNotIn(":", uri_config.path, "Check for URI concatenated with normal path")
2808 def testPutGet(self) -> None:
2809 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents")
2810 self.runPutGetTest(storageClass, "test_metric")
2813class ButlerMakeRepoOutfileDirTestCase(ButlerMakeRepoOutfileTestCase):
2814 """Test that a config file created by makeRepo outside of repo works."""
2816 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
2818 def setUp(self) -> None:
2819 self.root = makeTestTempDir(TESTDIR)
2820 self.root2 = makeTestTempDir(TESTDIR)
2822 self.tmpConfigFile = self.root2
2823 make_repo_for_test(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile)
2825 def testConfigExistence(self) -> None:
2826 # Append the yaml file else Config constructor does not know the file
2827 # type.
2828 self.tmpConfigFile = os.path.join(self.tmpConfigFile, "butler.yaml")
2829 super().testConfigExistence()
2832class ButlerMakeRepoOutfileUriTestCase(ButlerMakeRepoOutfileTestCase):
2833 """Test that a config file created by makeRepo outside of repo works."""
2835 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
2837 def setUp(self) -> None:
2838 self.root = makeTestTempDir(TESTDIR)
2839 self.root2 = makeTestTempDir(TESTDIR)
2841 self.tmpConfigFile = ResourcePath(os.path.join(self.root2, "something.yaml")).geturl()
2842 make_repo_for_test(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile)
2845class RemoteTestDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase):
2846 """Specialization of a butler using a datastore root that reports itself
2847 as not local; a remote file datastore + a local SqlRegistry.
2848 """
2850 configFile = os.path.join(TESTDIR, "config/basic/butler-remotetest-store.yaml")
2851 fullConfigKey = None
2852 validationCanFail = True
2854 registryStr = "/gen3.sqlite3"
2855 """Expected format of the Registry string."""
2857 def setUp(self) -> None:
2858 config = Config(self.configFile)
2860 self.root = makeTestTempDir(TESTDIR)
2861 # The space in the directory name is deliberate. It ensures the URI
2862 # has to be percent-encoded correctly on the way in and decoded on
2863 # the way out.
2864 root_path = os.path.join(self.root, "butler root")
2865 os.makedirs(root_path)
2866 rooturi = make_remote_test_uri(root_path)
2867 config.update({"datastore": {"datastore": {"root": str(rooturi)}}})
2869 # The registry database has to live on a real local file system.
2870 self.reg_dir = makeTestTempDir(TESTDIR)
2871 config["registry", "db"] = f"sqlite:///{self.reg_dir}/gen3.sqlite3"
2873 self.datastoreStr = [f"datastore='{rooturi}'"]
2874 self.datastoreName = [f"FileDatastore@{rooturi}"]
2875 make_repo_for_test(rooturi, config=config, forceConfigRoot=False)
2876 self.tmpConfigFile = str(rooturi.join("butler.yaml", forceDirectory=False))
2878 def tearDown(self) -> None:
2879 removeTestTempDir(self.reg_dir)
2880 # The base class removes self.root, which contains the datastore.
2881 super().tearDown()
2884class DatastoreTransfers(TestCaseMixin):
2885 """Base test setup for data transfers between butlers. The concrete tests
2886 for specific configurations are in other classes, below.
2887 """
2889 storageClassFactory: StorageClassFactory
2891 @classmethod
2892 def setUpClass(cls) -> None:
2893 cls.storageClassFactory = StorageClassFactory()
2895 def setUp(self) -> None:
2896 self.root = makeTestTempDir(TESTDIR)
2897 self.config = Config(self.configFile)
2899 # Some tests cause convertors to be replaced so ensure
2900 # the storage class factory is reset each time.
2901 self.storageClassFactory.reset()
2902 self.storageClassFactory.addFromConfig(self.configFile)
2904 def tearDown(self) -> None:
2905 removeTestTempDir(self.root)
2907 def create_butler(self, manager: str | None, label: str, config_file: str | None = None) -> Butler:
2908 if manager is None: 2908 ↛ 2912line 2908 didn't jump to line 2912 because the condition on line 2908 was always true
2909 manager = (
2910 "lsst.daf.butler.registry.datasets.byDimensions.ByDimensionsDatasetRecordStorageManagerUUID"
2911 )
2912 config = Config(config_file if config_file is not None else self.configFile)
2913 config["registry", "managers", "datasets"] = manager
2914 butler = Butler.from_config(
2915 make_repo_for_test(f"{self.root}/butler{label}", config=config), writeable=True
2916 )
2917 self.enterContext(butler)
2918 return butler
2920 def assertButlerTransfers(
2921 self,
2922 purge: bool = False,
2923 storageClassName: str = "StructuredData",
2924 storageClassNameTarget: str | None = None,
2925 ) -> None:
2926 """Test that a run can be transferred to another butler."""
2927 storageClass = self.storageClassFactory.getStorageClass(storageClassName)
2928 if storageClassNameTarget is not None:
2929 storageClassTarget = self.storageClassFactory.getStorageClass(storageClassNameTarget)
2930 else:
2931 storageClassTarget = storageClass
2933 datasetTypeName = "random_data"
2935 # Test will create 3 collections and we will want to transfer
2936 # two of those three.
2937 runs = ["run1", "run2", "other"]
2939 # Also want to use two different dataset types to ensure that
2940 # grouping works.
2941 datasetTypeNames = ["random_data", "random_data_2"]
2943 # Create the run collections in the source butler.
2944 for run in runs:
2945 self.source_butler.collections.register(run)
2947 # Create dimensions in source butler.
2948 n_exposures = 30
2949 self.source_butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"})
2950 self.source_butler.registry.insertDimensionData(
2951 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"}
2952 )
2953 self.source_butler.registry.insertDimensionData(
2954 "detector", {"instrument": "DummyCamComp", "id": 1, "full_name": "det1"}
2955 )
2956 self.source_butler.registry.insertDimensionData(
2957 "day_obs",
2958 {
2959 "instrument": "DummyCamComp",
2960 "id": 20250101,
2961 },
2962 )
2964 for i in range(n_exposures):
2965 self.source_butler.registry.insertDimensionData(
2966 "group", {"instrument": "DummyCamComp", "name": f"group{i}"}
2967 )
2968 self.source_butler.registry.insertDimensionData(
2969 "exposure",
2970 {
2971 "instrument": "DummyCamComp",
2972 "id": i,
2973 "obs_id": f"exp{i}",
2974 "physical_filter": "d-r",
2975 "group": f"group{i}",
2976 "day_obs": 20250101,
2977 },
2978 )
2980 # Create dataset types in the source butler.
2981 dimensions = self.source_butler.dimensions.conform(["instrument", "exposure"])
2982 for datasetTypeName in datasetTypeNames:
2983 datasetType = DatasetType(datasetTypeName, dimensions, storageClass)
2984 self.source_butler.registry.registerDatasetType(datasetType)
2986 # Write a dataset to an unrelated run -- this will ensure that
2987 # we are rewriting integer dataset ids in the target if necessary.
2988 # Will not be relevant for UUID.
2989 run = "distraction"
2990 butler = Butler.from_config(butler=self.source_butler, run=run)
2991 self.enterContext(butler)
2992 butler.put(
2993 makeExampleMetrics(),
2994 datasetTypeName,
2995 exposure=1,
2996 instrument="DummyCamComp",
2997 physical_filter="d-r",
2998 )
3000 # Write some example metrics to the source
3001 butler = Butler.from_config(butler=self.source_butler)
3002 self.enterContext(butler)
3004 # Set of DatasetRefs that should be in the list of refs to transfer
3005 # but which will not be transferred.
3006 deleted: set[DatasetRef] = set()
3008 n_expected = 20 # Number of datasets expected to be transferred
3009 source_refs = []
3010 for i in range(n_exposures):
3011 # Put a third of datasets into each collection, only retain
3012 # two thirds.
3013 index = i % 3
3014 run = runs[index]
3015 datasetTypeName = datasetTypeNames[i % 2]
3017 metric = MetricsExample(
3018 summary={"counter": i}, output={"text": "metric"}, data=[2 * x for x in range(i)]
3019 )
3020 dataId = {"exposure": i, "instrument": "DummyCamComp", "physical_filter": "d-r"}
3021 ref = butler.put(metric, datasetTypeName, dataId=dataId, run=run)
3023 # Remove the datastore record using low-level API, but only
3024 # for a specific index.
3025 if purge and index == 1:
3026 # For one of these delete the file as well.
3027 # This allows the "missing" code to filter the
3028 # file out.
3029 # Access the individual datastores.
3030 datastores = []
3031 if hasattr(butler._datastore, "datastores"):
3032 datastores.extend(butler._datastore.datastores)
3033 else:
3034 datastores.append(butler._datastore)
3036 if not deleted:
3037 # For a chained datastore we need to remove
3038 # files in each chain.
3039 for datastore in datastores:
3040 # The file might not be known to the datastore
3041 # if constraints are used.
3042 try:
3043 primary, uris = datastore.getURIs(ref)
3044 except FileNotFoundError:
3045 continue
3046 if primary and primary.scheme != "mem":
3047 primary.remove()
3048 for uri in uris.values():
3049 if uri.scheme != "mem": 3049 ↛ 3048line 3049 didn't jump to line 3048 because the condition on line 3049 was always true
3050 uri.remove()
3051 n_expected -= 1
3052 deleted.add(ref)
3054 # Remove the datastore record.
3055 for datastore in datastores:
3056 if hasattr(datastore, "removeStoredItemInfo"): 3056 ↛ 3055line 3056 didn't jump to line 3055 because the condition on line 3056 was always true
3057 datastore.removeStoredItemInfo(ref)
3059 if index < 2:
3060 source_refs.append(ref)
3061 if ref not in deleted:
3062 new_metric = butler.get(ref)
3063 self.assertEqual(new_metric, metric)
3065 # Create some bad dataset types to ensure we check for inconsistent
3066 # definitions.
3067 badStorageClass = self.storageClassFactory.getStorageClass("StructuredDataList")
3068 for datasetTypeName in datasetTypeNames:
3069 datasetType = DatasetType(datasetTypeName, dimensions, badStorageClass)
3070 self.target_butler.registry.registerDatasetType(datasetType)
3071 with self.assertRaises(ConflictingDefinitionError) as cm:
3072 self.target_butler.transfer_from(self.source_butler, source_refs)
3073 self.assertIn("dataset type differs", str(cm.exception))
3075 # And remove the bad definitions.
3076 for datasetTypeName in datasetTypeNames:
3077 self.target_butler.registry.removeDatasetType(datasetTypeName)
3079 # Transfer without creating dataset types should fail.
3080 with self.assertRaises(KeyError):
3081 self.target_butler.transfer_from(self.source_butler, source_refs)
3083 # Transfer without creating dimensions should fail.
3084 with self.assertRaises(ConflictingDefinitionError) as cm:
3085 self.target_butler.transfer_from(self.source_butler, source_refs, register_dataset_types=True)
3086 self.assertIn("dimension", str(cm.exception))
3088 # The dry run test requires dataset types to exist. If we have
3089 # been given distinct storage classes for the target we have
3090 # to redefine at least one of the dataset types in the target butler.
3091 if storageClass != storageClassTarget:
3092 self.target_butler.registry.removeDatasetType(datasetTypeNames[0])
3093 datasetType = DatasetType(datasetTypeNames[0], dimensions, storageClassTarget)
3094 self.target_butler.registry.registerDatasetType(datasetType)
3096 # The failed transfer above leaves registry in an inconsistent
3097 # state because the run is created but then rolled back without
3098 # the collection cache being cleared. For now force a refresh.
3099 # Can remove with DM-35498.
3100 self.target_butler.registry.refresh()
3102 # Do a dry run -- this should not have any effect on the target butler.
3103 self.target_butler.transfer_from(self.source_butler, source_refs, dry_run=True)
3105 # Transfer the records for one ref to test the alternative API.
3106 with self.assertLogs(logger="lsst", level=logging.DEBUG) as log_cm:
3107 self.target_butler.transfer_dimension_records_from(self.source_butler, [source_refs[0]])
3108 self.assertIn("number of records transferred: 1", ";".join(log_cm.output))
3110 # Now transfer them to the second butler, including dimensions.
3111 with self.assertLogs(logger="lsst", level=logging.DEBUG) as log_cm:
3112 transferred = self.target_butler.transfer_from(
3113 self.source_butler,
3114 source_refs,
3115 register_dataset_types=True,
3116 transfer_dimensions=True,
3117 )
3118 self.assertEqual(len(transferred), n_expected)
3119 log_output = ";".join(log_cm.output)
3121 # A ChainedDatastore will use the in-memory datastore for mexists
3122 # so we can not rely on the mexists log message.
3123 self.assertIn("Number of datastore records found in source", log_output)
3124 self.assertIn("Creating output run", log_output)
3126 # Do the transfer twice to ensure that it will do nothing extra.
3127 # Only do this if purge=True because it does not work for int
3128 # dataset_id.
3129 if purge:
3130 # This should not need to register dataset types.
3131 transferred = self.target_butler.transfer_from(self.source_butler, source_refs)
3132 self.assertEqual(len(transferred), n_expected)
3134 with self.assertRaises((TypeError, AttributeError)):
3135 self.target_butler._datastore.transfer_from(self.source_butler, source_refs) # type: ignore
3137 with self.assertRaises(ValueError):
3138 self.target_butler._datastore.transfer_from(
3139 self.source_butler._datastore, source_refs, transfer="split"
3140 )
3142 # Now try to get the same refs from the new butler.
3143 for ref in source_refs:
3144 if ref not in deleted:
3145 new_metric = self.target_butler.get(ref)
3146 old_metric = self.source_butler.get(ref)
3147 self.assertEqual(new_metric, old_metric)
3149 # Try again without implicit storage class conversion
3150 # triggered by using the source ref. This will do conversion
3151 # since the formatter will be returning the source python type.
3152 target_ref = self.target_butler.get_dataset(ref.id)
3153 if target_ref.datasetType.storageClass != ref.datasetType.storageClass:
3154 new_metric = self.target_butler.get(target_ref)
3155 self.assertNotEqual(type(new_metric), type(old_metric))
3157 # Remove the dataset from the target and put it again
3158 # as if it was the right type all along for this butler.
3159 self.target_butler.pruneDatasets(
3160 [target_ref], unstore=True, purge=True, disassociate=True
3161 )
3162 self.target_butler.put(new_metric, target_ref)
3163 new_new_metric = self.target_butler.get(target_ref)
3164 new_old_metric = self.target_butler.get(
3165 target_ref, storageClass=ref.datasetType.storageClass
3166 )
3167 self.assertEqual(new_new_metric, new_metric)
3168 self.assertEqual(new_old_metric, old_metric)
3170 # Now prune run2 collection and create instead a CHAINED collection.
3171 # This should block the transfer.
3172 self.target_butler.removeRuns(["run2"])
3173 self.target_butler.collections.register("run2", CollectionType.CHAINED)
3174 with self.assertRaises(CollectionTypeError):
3175 # Re-importing the run1 datasets can be problematic if they
3176 # use integer IDs so filter those out.
3177 to_transfer = [ref for ref in source_refs if ref.run == "run2"]
3178 self.target_butler.transfer_from(self.source_butler, to_transfer)
3181class PosixDatastoreTransfers(DatastoreTransfers, unittest.TestCase):
3182 """Test data transfers between butlers.
3184 Test for different managers. UUID to UUID and integer to integer are
3185 tested. UUID to integer is not supported since we do not currently
3186 want to allow that. Integer to UUID is supported with the caveat
3187 that UUID4 will be generated and this will be incorrect for raw
3188 dataset types. The test ignores that.
3189 """
3191 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
3193 def create_butlers(
3194 self, manager1: str | None = None, manager2: str | None = None, source_config: str | None = None
3195 ) -> None:
3196 self.source_butler = self.create_butler(manager1, "1", config_file=source_config)
3197 self.target_butler = self.create_butler(manager2, "2")
3199 def testTransferUuidToUuid(self) -> None:
3200 self.create_butlers()
3201 self.assertButlerTransfers()
3203 def testTransferFromChainedUuidToUuid(self) -> None:
3204 """Force the source butler to be a ChainedDatastore."""
3205 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-chained.yaml"))
3206 self.assertButlerTransfers()
3208 def testTransferFromIncompatibleUuidToUuid(self) -> None:
3209 """Force the source butler to be a incompatible datastore."""
3210 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-inmemory.yaml"))
3211 with self.assertRaises(NotImplementedError):
3212 self.assertButlerTransfers()
3214 def testTransferFromIncompatibleChainUuidToUuid(self) -> None:
3215 """Force the source butler to be a incompatible datastore."""
3216 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-inmemory-chain.yaml"))
3217 with self.assertRaises(TypeError):
3218 self.assertButlerTransfers()
3220 def testTransferFromFileUuidToUuid(self) -> None:
3221 """Force the source butler to be a FileDatastore."""
3222 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler.yaml"))
3223 self.assertButlerTransfers()
3225 def testTransferMissing(self) -> None:
3226 """Test transfers where datastore records are missing.
3228 This is how execution butler works.
3229 """
3230 self.create_butlers()
3232 # Configure the source butler to allow trust.
3233 self.source_butler._datastore._set_trust_mode(True)
3235 self.assertButlerTransfers(purge=True)
3237 def testTransferMissingDisassembly(self) -> None:
3238 """Test transfers where datastore records are missing.
3240 This is how execution butler works.
3241 """
3242 self.create_butlers()
3244 # Configure the source butler to allow trust.
3245 self.source_butler._datastore._set_trust_mode(True)
3247 # Test disassembly.
3248 self.assertButlerTransfers(purge=True, storageClassName="StructuredComposite")
3250 def testTransferDifferingStorageClasses(self) -> None:
3251 """Test transfers when the source butler dataset type has a different
3252 but compatible storage class.
3253 """
3254 self.create_butlers()
3256 self.assertButlerTransfers(storageClassNameTarget="MetricsConversion")
3258 def testTransferDifferingStorageClassesDisassembly(self) -> None:
3259 """Test transfers when the source butler dataset type has a different
3260 but compatible storage class and where the source butler has
3261 disassembled.
3262 """
3263 self.create_butlers()
3265 self.assertButlerTransfers(
3266 storageClassName="StructuredComposite", storageClassNameTarget="MetricsConversion"
3267 )
3269 def testUnsafeDirectTransfer(self) -> None:
3270 """Test that transfer='unsafe_direct' records the absolute URI of
3271 source files in the target datastore.
3272 """
3273 self.create_butlers()
3274 dataset_type = DatasetType("dt", [], "int", universe=self.source_butler.dimensions)
3275 self.source_butler.registry.registerDatasetType(dataset_type)
3276 self.source_butler.collections.register("run")
3277 ref = self.source_butler.put(123, "dt", [], run="run")
3278 self.target_butler.transfer_from(
3279 self.source_butler, [ref], transfer="unsafe_direct", register_dataset_types=True
3280 )
3281 self.assertEqual(self.target_butler.get(ref), 123)
3282 self.assertEqual(self.source_butler.getURI(ref), self.target_butler.getURI(ref))
3284 def testAbsoluteURITransferDirect(self) -> None:
3285 """Test transfer using an absolute URI."""
3286 self._absolute_transfer("auto")
3288 def testAbsoluteURITransferUnsafeDirect(self) -> None:
3289 """Test transfer using an absolute URI."""
3290 self._absolute_transfer("unsafe_direct")
3292 def testAbsoluteURITransferCopy(self) -> None:
3293 """Test transfer using an absolute URI."""
3294 self._absolute_transfer("copy")
3296 def _absolute_transfer(self, transfer: str) -> None:
3297 self.create_butlers()
3299 storageClassName = "StructuredData"
3300 storageClass = self.storageClassFactory.getStorageClass(storageClassName)
3301 datasetTypeName = "random_data"
3302 run = "run1"
3303 self.source_butler.collections.register(run)
3305 dimensions = self.source_butler.dimensions.conform(())
3306 datasetType = DatasetType(datasetTypeName, dimensions, storageClass)
3307 self.source_butler.registry.registerDatasetType(datasetType)
3309 metrics = makeExampleMetrics()
3310 # Ingest from a URI that reports itself as not local, so that the test
3311 # distinguishes "the absolute URI was preserved" from "a local path
3312 # happened to work".
3313 source_dir = os.path.join(self.root, "source data")
3314 os.makedirs(source_dir)
3315 with ResourcePath.temporary_uri(prefix=make_remote_test_uri(source_dir), suffix=".json") as temp:
3316 self.assertFalse(temp.isLocal)
3317 dataId = DataCoordinate.make_empty(self.source_butler.dimensions)
3318 source_refs = [DatasetRef(datasetType, dataId, run=run)]
3319 temp.write(json.dumps(metrics.exportAsDict()).encode())
3320 dataset = FileDataset(path=temp, refs=source_refs)
3321 self.source_butler.ingest(dataset, transfer="direct")
3323 self.target_butler.transfer_from(
3324 self.source_butler, dataset.refs, register_dataset_types=True, transfer=transfer
3325 )
3327 uri = self.target_butler.getURI(dataset.refs[0])
3328 if transfer == "auto" or transfer == "unsafe_direct":
3329 self.assertEqual(uri, temp)
3330 else:
3331 self.assertNotEqual(uri, temp)
3333 def test_shared_dimension_group(self):
3334 """Test internal logic that divides dataset types by dimension group
3335 when doing registry updates.
3336 """
3337 self.create_butlers()
3338 self.source_butler.import_(filename=_get_test_data_path("base.yaml"), without_datastore=True)
3339 self.source_butler.import_(filename=_get_test_data_path("datasets.yaml"), without_datastore=True)
3341 source_butler = self.source_butler
3342 target_butler = self.target_butler
3344 # Create a dataset type with the same dimensions as the 'bias' dataset
3345 # type from base.yaml
3346 dataset_type = DatasetType(
3347 "test_type", ["instrument", "detector"], "int", universe=source_butler.dimensions
3348 )
3349 source_butler.registry.registerDatasetType(dataset_type)
3350 # This has the same data ID as one of the bias datasets in
3351 # datasets.yaml.
3352 test_ref = source_butler.registry.insertDatasets(
3353 "test_type", [{"instrument": "Cam1", "detector": 2}], run="imported_g"
3354 )[0]
3356 biases = source_butler.query_datasets("bias", ["imported_g", "imported_r"])
3357 flats = source_butler.query_datasets("flat", ["imported_g", "imported_r"])
3358 refs = [test_ref, *biases, *flats]
3360 # Test setup will be even more convoluted if we want the datastore to
3361 # actually transfer files. For testing the dimension group behavior,
3362 # we really only care about the registry.
3363 with unittest.mock.patch.object(target_butler._datastore, "transfer_from") as mock:
3364 mock.return_value = (set(refs), set())
3365 target_butler.transfer_from(
3366 source_butler,
3367 refs,
3368 transfer=None,
3369 register_dataset_types=True,
3370 skip_missing=False,
3371 transfer_dimensions=True,
3372 )
3374 transferred_test_ref = target_butler.find_dataset(
3375 "test_type", {"instrument": "Cam1", "detector": 2}, collections="imported_g"
3376 )
3377 self.assertEqual(transferred_test_ref.id, test_ref.id)
3379 transferred_bias = target_butler.find_dataset(
3380 "bias", {"instrument": "Cam1", "detector": 2}, collections="imported_g"
3381 )
3382 self.assertEqual(transferred_bias.id, uuid.UUID("51352db4-a47a-447c-b12d-a50b206b17cd"))
3384 transferred_flat = target_butler.find_dataset(
3385 "flat",
3386 {"instrument": "Cam1", "detector": 2, "physical_filter": "Cam1-R1", "band": "r"},
3387 collections="imported_r",
3388 )
3389 self.assertEqual(transferred_flat.id, uuid.UUID("c1296796-56c5-4acf-9b49-40d920c6f840"))
3392class ChainedDatastoreTransfers(PosixDatastoreTransfers):
3393 """Test transfers using a chained datastore."""
3395 configFile = os.path.join(TESTDIR, "config/basic/butler-chained.yaml")
3398@unittest.skipIf(not butler_server_is_available, butler_server_import_error)
3399class ButlerServerDatastoreTransfers(DatastoreTransfers, unittest.TestCase):
3400 """Test ``transfer_from`` involving Butler server."""
3402 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
3404 def test_transfers_from_remote_to_direct(self) -> None:
3405 from lsst.daf.butler.remote_butler._remote_file_transfer_source import (
3406 mock_file_transfer_uris_for_unit_test,
3407 )
3409 self.target_butler = self.create_butler(None, "2")
3410 with create_test_server(TESTDIR) as server:
3411 self.source_butler = server.hybrid_butler
3413 def _remap_transfer_url(path: HttpResourcePath) -> HttpResourcePath:
3414 # The Butler server returns HTTP URIs with a domain name that
3415 # is not resolvable because there is no actual HTTP server
3416 # involved in these tests. Strip this first layer of
3417 # indirection, and return the target of the redirect instead.
3418 response = server.client.get(str(path), follow_redirects=False, headers=path._extra_headers)
3419 return ResourcePath(str(response.next_request.url))
3421 with mock_file_transfer_uris_for_unit_test(_remap_transfer_url):
3422 self.assertButlerTransfers()
3425class TransferDatasetsInPlace(unittest.TestCase):
3426 """Test behavior of transfer_datasets_in_place() specialty function used by
3427 Prompt Publication service.
3428 """
3430 def test_file_datastore(self) -> None:
3431 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
3432 with (
3433 tempfile.TemporaryDirectory() as datastore_root,
3434 tempfile.TemporaryDirectory() as other_repo_root,
3435 ):
3436 config = Config(configFile)
3437 config["datastore", "datastore", "name"] = "file_datastore"
3438 make_repo_for_test(datastore_root, config=config)
3439 config["datastore", "datastore", "root"] = datastore_root
3440 make_repo_for_test(other_repo_root, config, forceConfigRoot=False)
3441 with (
3442 Butler(datastore_root, writeable=True) as source_butler,
3443 Butler(other_repo_root, writeable=True) as target_butler,
3444 ):
3445 self._test_transfer_datasets_in_place(source_butler, target_butler)
3447 def test_chained_datastore(self) -> None:
3448 configFile = os.path.join(TESTDIR, "config/basic/butler-chained-posix.yaml")
3449 with (
3450 tempfile.TemporaryDirectory() as datastore_root,
3451 tempfile.TemporaryDirectory() as other_repo_root,
3452 ):
3453 config = Config(configFile)
3454 config["datastore", "datastore", "datastores", 0, "datastore", "root"] = (
3455 f"{datastore_root}/butler_test_repository"
3456 )
3457 config["datastore", "datastore", "datastores", 1, "datastore", "root"] = (
3458 f"{datastore_root}/butler_test_repository2"
3459 )
3460 make_repo_for_test(datastore_root, config=config, forceConfigRoot=False)
3461 make_repo_for_test(other_repo_root, config=config, forceConfigRoot=False)
3462 with (
3463 Butler(datastore_root, writeable=True) as source_butler,
3464 Butler(other_repo_root, writeable=True) as target_butler,
3465 ):
3466 self._test_transfer_datasets_in_place(source_butler, target_butler)
3468 def _test_transfer_datasets_in_place(
3469 self, source_butler: DirectButler, target_butler: DirectButler
3470 ) -> None:
3471 metric_repo = MetricTestRepo.create_from_butler(
3472 source_butler,
3473 source_butler._config,
3474 )
3475 target_butler.transfer_dimension_records_from(source_butler, [metric_repo.ref1, metric_repo.ref2])
3476 # Verify that the setup was correct and the two repos have
3477 # independent registries.
3478 self.assertIsNone(target_butler.get_dataset(metric_repo.ref1.id))
3479 # Copy one dataset, and make sure we can load it from the
3480 # target repo.
3481 self.assertEqual(
3482 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1]),
3483 [metric_repo.ref1],
3484 )
3485 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1))
3486 self.assertIsNone(target_butler.get_dataset(metric_repo.ref2.id))
3487 self.assertEqual(source_butler.getURIs(metric_repo.ref1), target_butler.getURIs(metric_repo.ref1))
3488 # Trying to copy the same dataset again is a no-op.
3489 self.assertEqual(
3490 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1]),
3491 [],
3492 )
3493 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1))
3494 # A mix of existing and non-existing datasets.
3495 self.assertEqual(
3496 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1, metric_repo.ref2]),
3497 [metric_repo.ref2],
3498 )
3499 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1))
3500 self.assertEqual(target_butler.get(metric_repo.ref2), source_butler.get(metric_repo.ref2))
3502 # For testing datastore chaining, set up a dataset that is only
3503 # accepted by one of the datastores.
3504 source_butler.registry.registerDatasetType(
3505 DatasetType("rejected_by_first", source_butler.dimensions.conform([]), "int")
3506 )
3507 source_butler.registry.registerRun("run")
3508 ref = source_butler.put(1, "rejected_by_first", dataId={}, run="run")
3509 self.assertEqual(
3510 transfer_datasets_in_place(source_butler, target_butler, [ref]),
3511 [ref],
3512 )
3513 self.assertEqual(1, target_butler.get(ref))
3516class NullDatastoreTestCase(unittest.TestCase):
3517 """Test that we can fall back to a null datastore."""
3519 # Need a good config to create the repo.
3520 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml")
3521 storageClassFactory: StorageClassFactory
3523 @classmethod
3524 def setUpClass(cls) -> None:
3525 cls.storageClassFactory = StorageClassFactory()
3526 cls.storageClassFactory.addFromConfig(cls.configFile)
3528 def setUp(self) -> None:
3529 """Create a new butler root for each test."""
3530 self.root = makeTestTempDir(TESTDIR)
3531 make_repo_for_test(self.root, config=Config(self.configFile))
3533 def tearDown(self) -> None:
3534 removeTestTempDir(self.root)
3536 def test_fallback(self) -> None:
3537 # Read the butler config and mess with the datastore section.
3538 config_path = os.path.join(self.root, "butler.yaml")
3539 bad_config = Config(config_path)
3540 bad_config["datastore", "cls"] = "lsst.not.a.datastore.Datastore"
3541 bad_config.dumpToUri(config_path)
3543 with self.assertRaises(RuntimeError):
3544 Butler(self.root, without_datastore=False)
3546 with self.assertRaises(RuntimeError):
3547 Butler.from_config(self.root, without_datastore=False)
3549 butler = Butler.from_config(self.root, writeable=True, without_datastore=True)
3550 self.enterContext(butler)
3551 self.assertIsInstance(butler._datastore, NullDatastore)
3553 # Check that registry is working.
3554 butler.collections.register("MYRUN")
3555 collections = butler.collections.query("*")
3556 self.assertIn("MYRUN", set(collections))
3558 # Create a ref.
3559 dimensions = butler.dimensions.conform([])
3560 storageClass = self.storageClassFactory.getStorageClass("StructuredDataDict")
3561 datasetTypeName = "metric"
3562 datasetType = DatasetType(datasetTypeName, dimensions, storageClass)
3563 butler.registry.registerDatasetType(datasetType)
3564 ref = DatasetRef(datasetType, {}, run="MYRUN")
3566 # Check that datastore will complain.
3567 with self.assertRaises(FileNotFoundError):
3568 butler.get(ref)
3569 with self.assertRaises(FileNotFoundError):
3570 butler.getURI(ref)
3573@unittest.skipIf(not butler_server_is_available, butler_server_import_error)
3574class ButlerServerTests(FileDatastoreButlerTests):
3575 """Test RemoteButler and Butler server."""
3577 configFile = None
3578 predictionSupported = False
3579 trustModeSupported = False
3581 postgres: TemporaryPostgresInstance | None
3583 def setUp(self):
3584 self.server_instance = self.enterContext(create_test_server(TESTDIR))
3586 def tearDown(self):
3587 pass
3589 def are_uris_equivalent(self, uri1: ResourcePath, uri2: ResourcePath) -> bool:
3590 # S3 pre-signed URLs may end up with differing expiration times in the
3591 # query parameters, so ignore query parameters when comparing.
3592 return uri1.scheme == uri2.scheme and uri1.netloc == uri2.netloc and uri1.path == uri2.path
3594 def create_empty_butler(
3595 self,
3596 run: str | None = None,
3597 writeable: bool | None = None,
3598 metrics: ButlerMetrics | None = None,
3599 cleanup: bool = True,
3600 ) -> Butler:
3601 return self.server_instance.hybrid_butler.clone(run=run, metrics=metrics)
3603 def remove_dataset_out_of_band(self, butler: Butler, ref: DatasetRef) -> None:
3604 # Can't delete a file via S3 signed URLs, so we need to reach in
3605 # through DirectButler to delete the dataset.
3606 uri = self.server_instance.direct_butler.getURI(ref)
3607 uri.remove()
3609 def testConstructor(self):
3610 # RemoteButler constructor is tested in test_server.py and
3611 # test_remote_butler.py.
3612 pass
3614 def testDafButlerRepositories(self):
3615 # Loading of RemoteButler via repository index is tested in
3616 # test_server.py.
3617 pass
3619 def testGetDatasetTypes(self) -> None:
3620 # This is mostly a test of validateConfiguration, which is for
3621 # validating Datastore configuration and thus isn't relevant to
3622 # RemoteButler.
3623 pass
3625 def testMakeRepo(self) -> None:
3626 # Only applies to DirectButler.
3627 pass
3629 # Pickling not yet implemented for RemoteButler/HybridButler.
3630 @unittest.expectedFailure
3631 def testPickle(self) -> None:
3632 return super().testPickle()
3634 def testStringification(self) -> None:
3635 self.assertEqual(
3636 str(self.server_instance.remote_butler),
3637 "RemoteButler(https://test.example/api/butler/repo/testrepo/)",
3638 )
3640 def testTransaction(self) -> None:
3641 # Transactions will never be supported for RemoteButler.
3642 pass
3644 def testPutTemplates(self) -> None:
3645 # The Butler server instance is configured with different file naming
3646 # templates than this test is expecting.
3647 pass
3650@unittest.skipIf(not butler_server_is_available, butler_server_import_error)
3651class ButlerServerSqliteTests(ButlerServerTests, unittest.TestCase):
3652 """Tests for RemoteButler's registry shim, with a SQLite DB backing the
3653 server.
3654 """
3656 postgres = None
3659@unittest.skipIf(not butler_server_is_available, butler_server_import_error)
3660class ButlerServerPostgresTests(ButlerServerTests, unittest.TestCase):
3661 """Tests for RemoteButler's registry shim, with a Postgres DB backing the
3662 server.
3663 """
3665 @classmethod
3666 def setUpClass(cls):
3667 cls.postgres = cls.enterClassContext(setup_postgres_test_db())
3668 super().setUpClass()
3671def setup_module(module: types.ModuleType) -> None:
3672 """Set up the module for pytest."""
3673 clean_environment()
3676def _get_test_data_path(filename: str) -> ResourcePath:
3677 return ResourcePath(f"resource://lsst.daf.butler/tests/registry_data/{filename}")
3680if __name__ == "__main__":
3681 clean_environment()
3682 unittest.main()