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