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