Coverage for tests/test_butler.py: 97%

1900 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-08-17 13:47 -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/>. 

27 

28"""Tests for Butler.""" 

29 

30from __future__ import annotations 

31 

32import json 

33import logging 

34import os 

35import pathlib 

36import pickle 

37import re 

38import shutil 

39import tempfile 

40import unittest 

41import unittest.mock 

42import uuid 

43import warnings 

44import weakref 

45from collections.abc import Callable, Mapping 

46from typing import TYPE_CHECKING, Any, cast 

47 

48import astropy.time 

49from sqlalchemy.exc import IntegrityError 

50 

51from lsst.daf.butler import ( 

52 Butler, 

53 ButlerConfig, 

54 ButlerMetrics, 

55 ButlerRepoIndex, 

56 CollectionCycleError, 

57 CollectionType, 

58 Config, 

59 DataCoordinate, 

60 DatasetExistence, 

61 DatasetNotFoundError, 

62 DatasetProvenance, 

63 DatasetRef, 

64 DatasetType, 

65 DimensionRecord, 

66 FileDataset, 

67 NoDefaultCollectionError, 

68 StorageClassFactory, 

69 ValidationError, 

70 script, 

71) 

72from lsst.daf.butler._rubin.file_datasets import transfer_datasets_to_datastore 

73from lsst.daf.butler._rubin.temporary_for_ingest import TemporaryForIngest 

74from lsst.daf.butler._rubin.transfer_datasets_in_place import transfer_datasets_in_place 

75from lsst.daf.butler.datastore import NullDatastore 

76from lsst.daf.butler.datastore.file_templates import FileTemplate, FileTemplateValidationError 

77from lsst.daf.butler.datastores.file_datastore.retrieve_artifacts import ZipIndex 

78from lsst.daf.butler.datastores.fileDatastore import FileDatastore 

79from lsst.daf.butler.direct_butler import DirectButler 

80from lsst.daf.butler.registry import ( 

81 CollectionError, 

82 CollectionTypeError, 

83 ConflictingDefinitionError, 

84 DataIdValueError, 

85 DatasetTypeExpressionError, 

86 MissingCollectionError, 

87 OrphanedRecordError, 

88) 

89from lsst.daf.butler.registry.sql_registry import SqlRegistry 

90from lsst.daf.butler.repo_relocation import BUTLER_ROOT_TAG 

91from lsst.daf.butler.tests import MetricsExample, MetricsExampleModel, MultiDetectorFormatter 

92from lsst.daf.butler.tests.dict_convertible_model import DictConvertibleModel 

93from lsst.daf.butler.tests.postgresql import TemporaryPostgresInstance, setup_postgres_test_db 

94from lsst.daf.butler.tests.server_available import butler_server_import_error, butler_server_is_available 

95from lsst.daf.butler.tests.utils import ( 

96 MetricTestRepo, 

97 TestCaseMixin, 

98 create_populated_sqlite_registry, 

99 makeTestTempDir, 

100 removeTestTempDir, 

101 safeTestTempDir, 

102) 

103from lsst.resources import ResourcePath 

104from lsst.resources.http import HttpResourcePath 

105from lsst.resources.tests import make_remote_test_uri 

106from lsst.utils import doImportType 

107from lsst.utils.introspection import get_full_type_name 

108 

109if butler_server_is_available: 

110 from lsst.daf.butler.tests.server import create_test_server 

111 

112 

113if TYPE_CHECKING: 

114 import types 

115 

116 from lsst.daf.butler import DimensionGroup, Registry, StorageClass 

117 

118TESTDIR = os.path.abspath(os.path.dirname(__file__)) 

119 

120 

121def clean_environment() -> None: 

122 """Remove external environment variables that affect the tests.""" 

123 for k in ("DAF_BUTLER_REPOSITORY_INDEX",): 

124 os.environ.pop(k, None) 

125 

126 

127def makeExampleMetrics() -> MetricsExample: 

128 """Return example dataset suitable for tests.""" 

129 return MetricsExample( 

130 {"AM1": 5.2, "AM2": 30.6}, 

131 {"a": [1, 2, 3], "b": {"blue": 5, "red": "green"}}, 

132 [563, 234, 456.7, 752, 8, 9, 27], 

133 ) 

134 

135 

136class TransactionTestError(Exception): 

137 """Specific error for testing transactions, to prevent misdiagnosing 

138 that might otherwise occur when a standard exception is used. 

139 """ 

140 

141 pass 

142 

143 

144class ButlerConfigTests(unittest.TestCase): 

145 """Simple tests for ButlerConfig that are not tested in any other test 

146 cases. 

147 """ 

148 

149 def testSearchPath(self) -> None: 

150 configFile = os.path.join(TESTDIR, "config", "basic", "butler.yaml") 

151 with self.assertLogs("lsst.daf.butler", level="DEBUG") as cm: 

152 config1 = ButlerConfig(configFile) 

153 self.assertNotIn("testConfigs", "\n".join(cm.output)) 

154 

155 overrideDirectory = os.path.join(TESTDIR, "config", "testConfigs") 

156 with self.assertLogs("lsst.daf.butler", level="DEBUG") as cm: 

157 config2 = ButlerConfig(configFile, searchPaths=[overrideDirectory]) 

158 self.assertIn("testConfigs", "\n".join(cm.output)) 

159 

160 key = ("datastore", "records", "table") 

161 self.assertNotEqual(config1[key], config2[key]) 

162 self.assertEqual(config2[key], "override_record") 

163 

164 

165class ButlerPutGetTests(TestCaseMixin): 

166 """Helper method for running a suite of put/get tests from different 

167 butler configurations. 

168 """ 

169 

170 root: str 

171 default_run = "ingésτ😺" 

172 storageClassFactory: StorageClassFactory 

173 configFile: str | None 

174 tmpConfigFile: str 

175 

176 @staticmethod 

177 def addDatasetType( 

178 datasetTypeName: str, dimensions: DimensionGroup, storageClass: StorageClass | str, registry: Registry 

179 ) -> DatasetType: 

180 """Create a DatasetType and register it""" 

181 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

182 registry.registerDatasetType(datasetType) 

183 return datasetType 

184 

185 @classmethod 

186 def setUpClass(cls) -> None: 

187 cls.storageClassFactory = StorageClassFactory() 

188 if cls.configFile is not None: 188 ↛ exitline 188 didn't return from function 'setUpClass' because the condition on line 188 was always true

189 cls.storageClassFactory.addFromConfig(cls.configFile) 

190 

191 def assertGetComponents( 

192 self, 

193 butler: Butler, 

194 datasetRef: DatasetRef, 

195 components: tuple[str, ...], 

196 reference: Any, 

197 collections: Any = None, 

198 ) -> None: 

199 datasetType = datasetRef.datasetType 

200 dataId = datasetRef.dataId 

201 deferred = butler.getDeferred(datasetRef) 

202 

203 for component in components: 

204 compTypeName = datasetType.componentTypeName(component) 

205 result = butler.get(compTypeName, dataId, collections=collections) 

206 self.assertEqual(result, getattr(reference, component)) 

207 result_deferred = deferred.get(component=component) 

208 self.assertEqual(result_deferred, result) 

209 

210 def tearDown(self) -> None: 

211 if self.root is not None: 211 ↛ exitline 211 didn't return from function 'tearDown' because the condition on line 211 was always true

212 removeTestTempDir(self.root) 

213 

214 def create_empty_butler( 

215 self, 

216 run: str | None = None, 

217 writeable: bool | None = None, 

218 metrics: ButlerMetrics | None = None, 

219 cleanup: bool = True, 

220 ): 

221 """Create a Butler for the test repository, without inserting test 

222 data. 

223 """ 

224 butler = Butler.from_config(self.tmpConfigFile, run=run, writeable=writeable, metrics=metrics) 

225 if cleanup: 

226 self.enterContext(butler) 

227 assert isinstance(butler, DirectButler), "Expect DirectButler in configuration" 

228 return butler 

229 

230 def create_butler( 

231 self, 

232 run: str, 

233 storageClass: StorageClass | str, 

234 datasetTypeName: str, 

235 metrics: ButlerMetrics | None = None, 

236 ) -> tuple[Butler, DatasetType]: 

237 """Create a Butler for the test repository and insert some test data 

238 into it. 

239 """ 

240 butler = self.create_empty_butler(run=run, metrics=metrics) 

241 

242 collections = set(butler.collections.query("*")) 

243 self.assertEqual(collections, {run}) 

244 # Create and register a DatasetType 

245 dimensions = butler.dimensions.conform(["instrument", "visit"]) 

246 

247 datasetType = self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry) 

248 

249 # Add needed Dimensions 

250 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"}) 

251 butler.registry.insertDimensionData( 

252 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"} 

253 ) 

254 butler.registry.insertDimensionData( 

255 "visit_system", {"instrument": "DummyCamComp", "id": 1, "name": "default"} 

256 ) 

257 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20200101}) 

258 visit_start = astropy.time.Time("2020-01-01 08:00:00.123456789", scale="tai") 

259 visit_end = astropy.time.Time("2020-01-01 08:00:36.66", scale="tai") 

260 butler.registry.insertDimensionData( 

261 "visit", 

262 { 

263 "instrument": "DummyCamComp", 

264 "id": 423, 

265 "name": "fourtwentythree", 

266 "physical_filter": "d-r", 

267 "datetime_begin": visit_start, 

268 "datetime_end": visit_end, 

269 "day_obs": 20200101, 

270 }, 

271 ) 

272 

273 # Add more visits for some later tests 

274 for visit_id in (424, 425): 

275 butler.registry.insertDimensionData( 

276 "visit", 

277 { 

278 "instrument": "DummyCamComp", 

279 "id": visit_id, 

280 "name": f"fourtwentyfour_{visit_id}", 

281 "physical_filter": "d-r", 

282 "day_obs": 20200101, 

283 }, 

284 ) 

285 return butler, datasetType 

286 

287 def runPutGetTest(self, storageClass: StorageClass, datasetTypeName: str) -> Butler: 

288 # New datasets will be added to run and tag, but we will only look in 

289 # tag when looking up datasets. 

290 run = self.default_run 

291 butler, datasetType = self.create_butler(run, storageClass, datasetTypeName) 

292 assert butler.run is not None 

293 

294 # Create and store a dataset 

295 metric = makeExampleMetrics() 

296 dataId = butler.registry.expandDataId({"instrument": "DummyCamComp", "visit": 423}) 

297 

298 # Dataset should not exist if we haven't added it 

299 with self.assertRaises(DatasetNotFoundError): 

300 butler.get(datasetTypeName, dataId) 

301 

302 # Put and remove the dataset once as a DatasetRef, once as a dataId, 

303 # and once with a DatasetType 

304 

305 # Keep track of any collections we add and do not clean up 

306 expected_collections = {run} 

307 

308 counter = 0 

309 ref = DatasetRef(datasetType, dataId, id=uuid.UUID(int=1), run="put_run_1") 

310 args = tuple[DatasetRef] | tuple[str | DatasetType, DataCoordinate] 

311 for args in ((ref,), (datasetTypeName, dataId), (datasetType, dataId)): 

312 # Since we are using subTest we can get cascading failures 

313 # here with the first attempt failing and the others failing 

314 # immediately because the dataset already exists. Work around 

315 # this by using a distinct run collection each time 

316 counter += 1 

317 this_run = f"put_run_{counter}" 

318 butler.collections.register(this_run) 

319 expected_collections.update({this_run}) 

320 

321 with self.subTest(args=repr(args)): 

322 kwargs: dict[str, Any] = {} 

323 if not isinstance(args[0], DatasetRef): # type: ignore 

324 kwargs["run"] = this_run 

325 ref = butler.put(metric, *args, **kwargs) 

326 self.assertIsInstance(ref, DatasetRef) 

327 

328 # Test get of a ref. 

329 metricOut = butler.get(ref) 

330 self.assertEqual(metric, metricOut) 

331 # Test get 

332 metricOut = butler.get(ref.datasetType.name, dataId, collections=this_run) 

333 self.assertEqual(metric, metricOut) 

334 # Test get with a datasetRef 

335 metricOut = butler.get(ref) 

336 self.assertEqual(metric, metricOut) 

337 # Test getDeferred with dataId 

338 metricOut = butler.getDeferred(ref.datasetType.name, dataId, collections=this_run).get() 

339 self.assertEqual(metric, metricOut) 

340 # Test getDeferred with a ref 

341 metricOut = butler.getDeferred(ref).get() 

342 self.assertEqual(metric, metricOut) 

343 

344 # Check we can get components 

345 if storageClass.isComposite(): 

346 self.assertGetComponents( 

347 butler, ref, ("summary", "data", "output"), metric, collections=this_run 

348 ) 

349 

350 primary_uri, secondary_uris = butler.getURIs(ref) 

351 n_uris = len(secondary_uris) 

352 if primary_uri: 

353 n_uris += 1 

354 

355 # Can the artifacts themselves be retrieved? 

356 if not butler._datastore.isEphemeral: 

357 # Create a temporary directory to hold the retrieved 

358 # artifacts. 

359 with tempfile.TemporaryDirectory( 

360 prefix="butler-artifacts-", ignore_cleanup_errors=True 

361 ) as artifact_root: 

362 root_uri = ResourcePath(artifact_root, forceDirectory=True) 

363 

364 for preserve_path in (True, False): 

365 destination = root_uri.join(f"{preserve_path}_{counter}/") 

366 log = logging.getLogger("lsst.x") 

367 log.debug("Using destination %s for args %s", destination, args) 

368 # Use copy so that we can test that overwrite 

369 # protection works (using "auto" for File URIs 

370 # would use hard links and subsequent transfer 

371 # would work because it knows they are the same 

372 # file). 

373 transferred = butler.retrieveArtifacts( 

374 [ref], destination, preserve_path=preserve_path, transfer="copy" 

375 ) 

376 self.assertGreater(len(transferred), 0) 

377 artifacts = list(ResourcePath.findFileResources([destination])) 

378 # Filter out the index file. 

379 artifacts = [a for a in artifacts if a.basename() != ZipIndex.index_name] 

380 self.assertEqual(set(transferred), set(artifacts)) 

381 

382 for artifact in transferred: 

383 path_in_destination = artifact.relative_to(destination) 

384 self.assertIsNotNone(path_in_destination) 

385 assert path_in_destination is not None 

386 

387 # When path is not preserved there should not 

388 # be any path separators. 

389 num_seps = path_in_destination.count("/") 

390 if preserve_path: 

391 self.assertGreater(num_seps, 0) 

392 else: 

393 self.assertEqual(num_seps, 0) 

394 

395 self.assertEqual( 

396 len(artifacts), 

397 n_uris, 

398 "Comparing expected artifacts vs actual:" 

399 f" {artifacts} vs {primary_uri} and {secondary_uris}", 

400 ) 

401 

402 if preserve_path: 

403 # No need to run these twice 

404 with self.assertRaises(ValueError): 

405 butler.retrieveArtifacts([ref], destination, transfer="move") 

406 

407 with self.assertRaisesRegex( 

408 ValueError, "^Destination location must refer to a directory" 

409 ): 

410 butler.retrieveArtifacts( 

411 [ref], ResourcePath("/some/file.txt", forceDirectory=False) 

412 ) 

413 

414 with self.assertRaises(FileExistsError): 

415 butler.retrieveArtifacts([ref], destination) 

416 

417 transferred_again = butler.retrieveArtifacts( 

418 [ref], destination, preserve_path=preserve_path, overwrite=True 

419 ) 

420 self.assertEqual(set(transferred_again), set(transferred)) 

421 

422 # Now remove the dataset completely. 

423 butler.pruneDatasets([ref], purge=True, unstore=True) 

424 # Lookup with original args should still fail. 

425 kwargs = {"collections": this_run} 

426 if isinstance(args[0], DatasetRef): 

427 kwargs = {} # Prevent warning from being issued. 

428 self.assertFalse(butler.exists(*args, **kwargs)) 

429 # get() should still fail. 

430 with self.assertRaises((FileNotFoundError, DatasetNotFoundError)): 

431 butler.get(ref) 

432 # Registry shouldn't be able to find it by dataset_id anymore. 

433 self.assertIsNone(butler.get_dataset(ref.id)) 

434 

435 # Do explicit registry removal since we know they are 

436 # empty 

437 butler.collections.x_remove(this_run) 

438 expected_collections.remove(this_run) 

439 

440 # Create DatasetRef for put using default run. 

441 refIn = DatasetRef(datasetType, dataId, id=uuid.UUID(int=1), run=butler.run) 

442 

443 # Check that getDeferred fails with standalone ref. 

444 with self.assertRaises(LookupError): 

445 butler.getDeferred(refIn) 

446 

447 # Put the dataset again, since the last thing we did was remove it 

448 # and we want to use the default collection. 

449 ref = butler.put(metric, refIn) 

450 

451 # Get with parameters 

452 stop = 4 

453 sliced = butler.get(ref, parameters={"slice": slice(stop)}) 

454 self.assertNotEqual(metric, sliced) 

455 self.assertEqual(metric.summary, sliced.summary) 

456 self.assertEqual(metric.output, sliced.output) 

457 assert metric.data is not None # for mypy 

458 self.assertEqual(metric.data[:stop], sliced.data) 

459 # getDeferred with parameters 

460 sliced = butler.getDeferred(ref, parameters={"slice": slice(stop)}).get() 

461 self.assertNotEqual(metric, sliced) 

462 self.assertEqual(metric.summary, sliced.summary) 

463 self.assertEqual(metric.output, sliced.output) 

464 self.assertEqual(metric.data[:stop], sliced.data) 

465 # getDeferred with deferred parameters 

466 sliced = butler.getDeferred(ref).get(parameters={"slice": slice(stop)}) 

467 self.assertNotEqual(metric, sliced) 

468 self.assertEqual(metric.summary, sliced.summary) 

469 self.assertEqual(metric.output, sliced.output) 

470 self.assertEqual(metric.data[:stop], sliced.data) 

471 

472 if storageClass.isComposite(): 

473 # Check that components can be retrieved 

474 metricOut = butler.get(ref.datasetType.name, dataId) 

475 compNameS = ref.datasetType.componentTypeName("summary") 

476 compNameD = ref.datasetType.componentTypeName("data") 

477 summary = butler.get(compNameS, dataId) 

478 self.assertEqual(summary, metric.summary) 

479 data = butler.get(compNameD, dataId) 

480 self.assertEqual(data, metric.data) 

481 

482 if "counter" in storageClass.derivedComponents: 

483 count = butler.get(ref.datasetType.componentTypeName("counter"), dataId) 

484 self.assertEqual(count, len(data)) 

485 

486 count = butler.get( 

487 ref.datasetType.componentTypeName("counter"), dataId, parameters={"slice": slice(stop)} 

488 ) 

489 self.assertEqual(count, stop) 

490 

491 compRef = butler.find_dataset(compNameS, dataId, collections=butler.collections.defaults) 

492 assert compRef is not None 

493 summary = butler.get(compRef) 

494 self.assertEqual(summary, metric.summary) 

495 

496 # Create a Dataset type that has the same name but is inconsistent. 

497 inconsistentDatasetType = DatasetType( 

498 datasetTypeName, datasetType.dimensions, self.storageClassFactory.getStorageClass("Config") 

499 ) 

500 

501 # Getting with a dataset type that does not match registry fails 

502 with self.assertRaisesRegex( 

503 ValueError, 

504 "(Supplied dataset type .* inconsistent with registry)" 

505 "|(The new storage class .* is not compatible with the existing storage class)", 

506 ): 

507 butler.get(inconsistentDatasetType, dataId) 

508 

509 # Combining a DatasetRef with a dataId should fail 

510 with self.assertRaisesRegex(ValueError, "DatasetRef given, cannot use dataId as well"): 

511 butler.get(ref, dataId) 

512 # Getting with an explicit ref should fail if the id doesn't match. 

513 with self.assertRaises((FileNotFoundError, DatasetNotFoundError)): 

514 butler.get(DatasetRef(ref.datasetType, ref.dataId, id=uuid.UUID(int=101), run=butler.run)) 

515 

516 # Getting a dataset with unknown parameters should fail 

517 with self.assertRaisesRegex(KeyError, "Parameter 'unsupported' not understood"): 

518 butler.get(ref, parameters={"unsupported": True}) 

519 

520 # Check we have a collection 

521 collections = set(butler.collections.query("*")) 

522 self.assertEqual(collections, expected_collections) 

523 

524 # Clean up to check that we can remove something that may have 

525 # already had a component removed 

526 butler.pruneDatasets([ref], unstore=True, purge=True) 

527 

528 # Add the same ref again, so we can check that duplicate put fails. 

529 ref = butler.put(metric, datasetType, dataId) 

530 

531 # Repeat put will fail. 

532 with self.assertRaisesRegex( 

533 ConflictingDefinitionError, "A database constraint failure was triggered" 

534 ): 

535 butler.put(metric, datasetType, dataId) 

536 

537 # Remove the datastore entry. 

538 butler.pruneDatasets([ref], unstore=True, purge=False, disassociate=False) 

539 

540 # Put will still fail 

541 with self.assertRaisesRegex( 

542 ConflictingDefinitionError, "A database constraint failure was triggered" 

543 ): 

544 butler.put(metric, datasetType, dataId) 

545 

546 # Repeat the same sequence with resolved ref. 

547 butler.pruneDatasets([ref], unstore=True, purge=True) 

548 ref = butler.put(metric, refIn) 

549 

550 # Repeat put will fail. 

551 with self.assertRaisesRegex(ConflictingDefinitionError, "Datastore already contains dataset"): 

552 butler.put(metric, refIn) 

553 

554 # Remove the datastore entry. 

555 butler.pruneDatasets([ref], unstore=True, purge=False, disassociate=False) 

556 

557 # In case of resolved ref this write will succeed. 

558 ref = butler.put(metric, refIn) 

559 

560 # Leave the dataset in place since some downstream tests require 

561 # something to be present 

562 

563 return butler 

564 

565 def testDeferredCollectionPassing(self) -> None: 

566 # Construct a butler with no run or collection, but make it writeable. 

567 butler = self.create_empty_butler(writeable=True) 

568 # Create and register a DatasetType 

569 dimensions = butler.dimensions.conform(["instrument", "visit"]) 

570 datasetType = self.addDatasetType( 

571 "example", dimensions, self.storageClassFactory.getStorageClass("StructuredData"), butler.registry 

572 ) 

573 # Add needed Dimensions 

574 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"}) 

575 butler.registry.insertDimensionData( 

576 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"} 

577 ) 

578 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20250101}) 

579 butler.registry.insertDimensionData( 

580 "visit", 

581 { 

582 "instrument": "DummyCamComp", 

583 "id": 423, 

584 "name": "fourtwentythree", 

585 "physical_filter": "d-r", 

586 "day_obs": 20250101, 

587 }, 

588 ) 

589 dataId = {"instrument": "DummyCamComp", "visit": 423} 

590 # Create dataset. 

591 metric = makeExampleMetrics() 

592 # Register a new run and put dataset. 

593 run = "deferred" 

594 self.assertTrue(butler.collections.register(run)) 

595 # Second time it will be allowed but indicate no-op 

596 self.assertFalse(butler.collections.register(run)) 

597 ref = butler.put(metric, datasetType, dataId, run=run) 

598 # Putting with no run should fail with TypeError. 

599 with self.assertRaises(CollectionError): 

600 butler.put(metric, datasetType, dataId) 

601 # Dataset should exist. 

602 self.assertTrue(butler.exists(datasetType, dataId, collections=[run])) 

603 # We should be able to get the dataset back, but with and without 

604 # a deferred dataset handle. 

605 self.assertEqual(metric, butler.get(datasetType, dataId, collections=[run])) 

606 self.assertEqual(metric, butler.getDeferred(datasetType, dataId, collections=[run]).get()) 

607 # Trying to find the dataset without any collection is an error. 

608 with self.assertRaises(NoDefaultCollectionError): 

609 butler.exists(datasetType, dataId) 

610 with self.assertRaises(CollectionError): 

611 butler.get(datasetType, dataId) 

612 # Associate the dataset with a different collection. 

613 butler.collections.register("tagged", type=CollectionType.TAGGED) 

614 butler.registry.associate("tagged", [ref]) 

615 # Deleting the dataset from the new collection should make it findable 

616 # in the original collection. 

617 butler.pruneDatasets([ref], tags=["tagged"]) 

618 self.assertTrue(butler.exists(datasetType, dataId, collections=[run])) 

619 

620 

621class ButlerTests(ButlerPutGetTests): 

622 """Tests for Butler.""" 

623 

624 useTempRoot = True 

625 validationCanFail: bool 

626 fullConfigKey: str | None 

627 registryStr: str | None 

628 datastoreName: list[str] | None 

629 datastoreStr: list[str] 

630 predictionSupported = True 

631 """Does getURIs support 'prediction mode'?""" 

632 

633 def setUp(self) -> None: 

634 """Create a new butler root for each test.""" 

635 self.root = makeTestTempDir(TESTDIR) 

636 Butler.makeRepo(self.root, config=Config(self.configFile)) 

637 self.tmpConfigFile = os.path.join(self.root, "butler.yaml") 

638 

639 def are_uris_equivalent(self, uri1: ResourcePath, uri2: ResourcePath) -> bool: 

640 """Return True if two URIs refer to the same resource. 

641 

642 Subclasses may override to handle unique requirements. 

643 """ 

644 return uri1 == uri2 

645 

646 def testConstructor(self) -> None: 

647 """Independent test of constructor.""" 

648 butler = Butler.from_config(self.tmpConfigFile, run=self.default_run) 

649 self.enterContext(butler) 

650 self.assertIsInstance(butler, Butler) 

651 

652 # Check that butler.yaml is added automatically. 

653 if self.tmpConfigFile.endswith(end := "/butler.yaml"): 

654 config_dir = self.tmpConfigFile[: -len(end)] 

655 butler = Butler.from_config(config_dir, run=self.default_run) 

656 self.enterContext(butler) 

657 self.assertIsInstance(butler, Butler) 

658 

659 # Even with a ResourcePath. 

660 butler = Butler.from_config(ResourcePath(config_dir, forceDirectory=True), run=self.default_run) 

661 self.enterContext(butler) 

662 self.assertIsInstance(butler, Butler) 

663 

664 collections = set(butler.collections.query("*")) 

665 self.assertEqual(collections, {self.default_run}) 

666 

667 # Check that some special characters can be included in run name. 

668 special_run = "u@b.c-A" 

669 butler_special = Butler.from_config(butler=butler, run=special_run) 

670 self.enterContext(butler_special) 

671 collections = set(butler_special.registry.queryCollections("*@*")) 

672 self.assertEqual(collections, {special_run}) 

673 

674 butler2 = Butler.from_config(butler=butler, collections=["other"]) 

675 self.enterContext(butler2) 

676 self.assertEqual(butler2.collections.defaults, ("other",)) 

677 self.assertIsNone(butler2.run) 

678 self.assertEqual(type(butler._datastore), type(butler2._datastore)) 

679 self.assertEqual(butler._datastore.config, butler2._datastore.config) 

680 

681 # Test that we can use an environment variable to find this 

682 # repository. 

683 butler_index = Config() 

684 butler_index["label"] = self.tmpConfigFile 

685 for suffix in (".yaml", ".json"): 

686 # Ensure that the content differs so that we know that 

687 # we aren't reusing the cache. 

688 bad_label = f"file://bucket/not_real{suffix}" 

689 butler_index["bad_label"] = bad_label 

690 with ResourcePath.temporary_uri(suffix=suffix) as temp_file: 

691 butler_index.dumpToUri(temp_file) 

692 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": str(temp_file)}): 

693 self.assertEqual(Butler.get_known_repos(), {"label", "bad_label"}) 

694 uri = Butler.get_repo_uri("bad_label") 

695 self.assertEqual(uri, ResourcePath(bad_label)) 

696 uri = Butler.get_repo_uri("label") 

697 butler = Butler.from_config(uri, writeable=False) 

698 self.assertIsInstance(butler, Butler) 

699 butler.close() 

700 butler = Butler.from_config("label", writeable=False) 

701 self.assertIsInstance(butler, Butler) 

702 butler.close() 

703 with self.assertRaisesRegex(FileNotFoundError, "aliases:.*bad_label"): 

704 Butler.from_config("not_there", writeable=False) 

705 with self.assertRaisesRegex(FileNotFoundError, "resolved from alias 'bad_label'"): 

706 Butler.from_config("bad_label") 

707 with self.assertRaises(FileNotFoundError): 

708 # Should ignore aliases. 

709 Butler.from_config(ResourcePath("label", forceAbsolute=False)) 

710 with self.assertRaises(KeyError) as cm: 

711 Butler.get_repo_uri("missing") 

712 self.assertEqual( 

713 Butler.get_repo_uri("missing", True), ResourcePath("missing", forceAbsolute=False) 

714 ) 

715 self.assertIn("not known to", str(cm.exception)) 

716 # Should report no failure. 

717 self.assertEqual(ButlerRepoIndex.get_failure_reason(), "") 

718 with ResourcePath.temporary_uri(suffix=suffix) as temp_file: 

719 # Now with empty configuration. 

720 butler_index = Config() 

721 butler_index.dumpToUri(temp_file) 

722 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": str(temp_file)}): 

723 with self.assertRaisesRegex(FileNotFoundError, "(no known aliases)"): 

724 Butler.from_config("label") 

725 with ResourcePath.temporary_uri(suffix=suffix) as temp_file: 

726 # Now with bad contents. 

727 with open(temp_file.ospath, "w") as fh: 

728 print("'", file=fh) 

729 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": str(temp_file)}): 

730 with self.assertRaisesRegex(FileNotFoundError, "(no known aliases:.*could not be read)"): 

731 Butler.from_config("label") 

732 with unittest.mock.patch.dict(os.environ, {"DAF_BUTLER_REPOSITORY_INDEX": "file://not_found/x.yaml"}): 

733 with self.assertRaises(FileNotFoundError): 

734 Butler.get_repo_uri("label") 

735 self.assertEqual(Butler.get_known_repos(), set()) 

736 

737 with self.assertRaisesRegex(FileNotFoundError, "index file not found"): 

738 Butler.from_config("label") 

739 

740 # Check that we can create Butler when the alias file is not found. 

741 butler = Butler.from_config(self.tmpConfigFile, writeable=False) 

742 self.enterContext(butler) 

743 self.assertIsInstance(butler, Butler) 

744 with self.assertRaises(RuntimeError) as cm: 

745 # No environment variable set. 

746 Butler.get_repo_uri("label") 

747 self.assertEqual(Butler.get_repo_uri("label", True), ResourcePath("label", forceAbsolute=False)) 

748 self.assertIn("No repository index defined", str(cm.exception)) 

749 with self.assertRaisesRegex(FileNotFoundError, "no known aliases.*No repository index"): 

750 # No aliases registered. 

751 Butler.from_config("not_there") 

752 self.assertEqual(Butler.get_known_repos(), set()) 

753 

754 def testClose(self): 

755 butler = self.create_empty_butler(cleanup=False) 

756 is_direct_butler = isinstance(butler, DirectButler) 

757 if is_direct_butler: 757 ↛ 760line 757 didn't jump to line 760 because the condition on line 757 was always true

758 self.assertFalse(butler._closed) 

759 

760 with butler as butler_from_context_manager: 

761 self.assertIs(butler, butler_from_context_manager) 

762 if is_direct_butler: 762 ↛ 768line 762 didn't jump to line 768 because the condition on line 762 was always true

763 self.assertTrue(butler._closed) 

764 with self.assertRaisesRegex(RuntimeError, "has been closed"): 

765 butler.get_dataset_type("raw") 

766 

767 # Close may be called multiple times. 

768 butler.close() 

769 if is_direct_butler: 769 ↛ exitline 769 didn't return from function 'testClose' because the condition on line 769 was always true

770 self.assertTrue(butler._closed) 

771 

772 def testGarbageCollection(self): 

773 """Test that Butler does not have any circular references that prevent 

774 it from being garbage collected immediately when it goes out of scope. 

775 """ 

776 butler = self.create_empty_butler(cleanup=False) 

777 is_direct_butler = isinstance(butler, DirectButler) 

778 butler_ref = weakref.ref(butler) 

779 if is_direct_butler: 779 ↛ 786line 779 didn't jump to line 786 because the condition on line 779 was always true

780 registry_ref = weakref.ref(butler._registry) 

781 managers_ref = weakref.ref(butler._registry._managers) 

782 datastore_ref = weakref.ref(butler._datastore) 

783 db_ref = weakref.ref(butler._registry._db) 

784 engine_ref = weakref.ref(butler._registry._db._engine) 

785 

786 with warnings.catch_warnings(): 

787 # Hide warnings from unclosed database handles. 

788 warnings.simplefilter("ignore", ResourceWarning) 

789 del butler 

790 self.assertIsNone(butler_ref(), "Butler should have been garbage collected") 

791 if is_direct_butler: 791 ↛ 799line 791 didn't jump to line 799 because the condition on line 791 was always true

792 self.assertIsNone(registry_ref(), "SqlRegistry should have been garbage collected") 

793 self.assertIsNone(managers_ref(), "Registry managers should have been garbage collected") 

794 self.assertIsNone(datastore_ref(), "Datastore should have been garbage collected") 

795 self.assertIsNone(db_ref(), "Database should have been garbage collected") 

796 # SQLAlchemy has internal reference cycles, so the Engine instance 

797 # is not cleaned up promptly even if we release our reference to 

798 # it. Explicitly clean it up here to avoid file handles leaking. 

799 if is_direct_butler: 799 ↛ exitline 799 didn't jump to the function exit

800 engine = engine_ref() 

801 if engine is not None: 801 ↛ exitline 801 didn't jump to the function exit

802 engine.dispose() 

803 

804 def testDafButlerRepositories(self): 

805 with unittest.mock.patch.dict( 

806 os.environ, 

807 {"DAF_BUTLER_REPOSITORIES": "label: 'https://someuri.com'\notherLabel: 'https://otheruri.com'\n"}, 

808 ): 

809 self.assertEqual(str(Butler.get_repo_uri("label")), "https://someuri.com") 

810 

811 with unittest.mock.patch.dict( 

812 os.environ, 

813 { 

814 "DAF_BUTLER_REPOSITORIES": "label: https://someuri.com", 

815 "DAF_BUTLER_REPOSITORY_INDEX": "https://someuri.com", 

816 }, 

817 ): 

818 with self.assertRaisesRegex(RuntimeError, "Only one of the environment variables"): 

819 Butler.get_repo_uri("label") 

820 

821 with unittest.mock.patch.dict( 

822 os.environ, 

823 {"DAF_BUTLER_REPOSITORIES": "invalid"}, 

824 ): 

825 with self.assertRaisesRegex(ValueError, "Repository index not in expected format"): 

826 Butler.get_repo_uri("label") 

827 

828 def testBasicPutGet(self) -> None: 

829 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

830 self.runPutGetTest(storageClass, "test_metric") 

831 

832 def testCompositePutGetConcrete(self) -> None: 

833 storageClass = self.storageClassFactory.getStorageClass("StructuredCompositeReadCompNoDisassembly") 

834 butler = self.runPutGetTest(storageClass, "test_metric") 

835 

836 # Should *not* be disassembled 

837 datasets = list(butler.registry.queryDatasets(..., collections=self.default_run)) 

838 self.assertEqual(len(datasets), 1) 

839 uri, components = butler.getURIs(datasets[0]) 

840 self.assertIsInstance(uri, ResourcePath) 

841 self.assertFalse(components) 

842 self.assertEqual(uri.fragment, "", f"Checking absence of fragment in {uri}") 

843 self.assertIn("423", str(uri), f"Checking visit is in URI {uri}") 

844 

845 # Predicted dataset 

846 if self.predictionSupported: 846 ↛ exitline 846 didn't return from function 'testCompositePutGetConcrete' because the condition on line 846 was always true

847 dataId = {"instrument": "DummyCamComp", "visit": 424} 

848 uri, components = butler.getURIs(datasets[0].datasetType, dataId=dataId, predict=True) 

849 self.assertFalse(components) 

850 self.assertIsInstance(uri, ResourcePath) 

851 self.assertIn("424", str(uri), f"Checking visit is in URI {uri}") 

852 self.assertEqual(uri.fragment, "predicted", f"Checking for fragment in {uri}") 

853 # Repeat with a DatasetRef to test that code path. 

854 ref = DatasetRef( 

855 datasets[0].datasetType, 

856 dataId=DataCoordinate.standardize(dataId, universe=butler.dimensions), 

857 run=self.default_run, 

858 ) 

859 uri2, components2 = butler.getURIs(ref, predict=True) 

860 self.assertFalse(components2) 

861 self.assertEqual(uri, uri2) 

862 

863 def testCompositePutGetVirtual(self) -> None: 

864 storageClass = self.storageClassFactory.getStorageClass("StructuredCompositeReadComp") 

865 butler = self.runPutGetTest(storageClass, "test_metric_comp") 

866 

867 # Should be disassembled 

868 datasets = list(butler.registry.queryDatasets(..., collections=self.default_run)) 

869 self.assertEqual(len(datasets), 1) 

870 uri, components = butler.getURIs(datasets[0]) 

871 

872 if butler._datastore.isEphemeral: 

873 # Never disassemble in-memory datastore 

874 self.assertIsInstance(uri, ResourcePath) 

875 self.assertFalse(components) 

876 self.assertEqual(uri.fragment, "", f"Checking absence of fragment in {uri}") 

877 self.assertIn("423", str(uri), f"Checking visit is in URI {uri}") 

878 else: 

879 self.assertIsNone(uri) 

880 self.assertEqual(set(components), set(storageClass.components)) 

881 for compuri in components.values(): 

882 self.assertIsInstance(compuri, ResourcePath) 

883 self.assertIn("423", str(compuri), f"Checking visit is in URI {compuri}") 

884 self.assertEqual(compuri.fragment, "", f"Checking absence of fragment in {compuri}") 

885 

886 if self.predictionSupported: 886 ↛ exitline 886 didn't return from function 'testCompositePutGetVirtual' because the condition on line 886 was always true

887 # Predicted dataset 

888 dataId = {"instrument": "DummyCamComp", "visit": 424} 

889 uri, components = butler.getURIs(datasets[0].datasetType, dataId=dataId, predict=True) 

890 

891 if butler._datastore.isEphemeral: 

892 # Never disassembled 

893 self.assertIsInstance(uri, ResourcePath) 

894 self.assertFalse(components) 

895 self.assertIn("424", str(uri), f"Checking visit is in URI {uri}") 

896 self.assertEqual(uri.fragment, "predicted", f"Checking for fragment in {uri}") 

897 else: 

898 self.assertIsNone(uri) 

899 self.assertEqual(set(components), set(storageClass.components)) 

900 for compuri in components.values(): 

901 self.assertIsInstance(compuri, ResourcePath) 

902 self.assertIn("424", str(compuri), f"Checking visit is in URI {compuri}") 

903 self.assertEqual(compuri.fragment, "predicted", f"Checking for fragment in {compuri}") 

904 

905 def testStorageClassOverrideGet(self) -> None: 

906 """Test storage class conversion on get with override.""" 

907 storageClass = self.storageClassFactory.getStorageClass("StructuredData") 

908 datasetTypeName = "anything" 

909 run = self.default_run 

910 

911 butler, datasetType = self.create_butler(run, storageClass, datasetTypeName) 

912 

913 # Create and store a dataset. 

914 metric = makeExampleMetrics() 

915 dataId = {"instrument": "DummyCamComp", "visit": 423} 

916 

917 ref = butler.put(metric, datasetType, dataId) 

918 

919 # Return native type. 

920 retrieved = butler.get(ref) 

921 self.assertEqual(retrieved, metric) 

922 

923 # Specify an override. 

924 new_sc = self.storageClassFactory.getStorageClass("MetricsConversion") 

925 model = butler.get(ref, storageClass=new_sc) 

926 self.assertNotEqual(type(model), type(retrieved)) 

927 self.assertIs(type(model), new_sc.pytype) 

928 self.assertEqual(retrieved, model) 

929 

930 # Defer but override later. 

931 deferred = butler.getDeferred(ref) 

932 model = deferred.get(storageClass=new_sc) 

933 self.assertIs(type(model), new_sc.pytype) 

934 self.assertEqual(retrieved, model) 

935 

936 # Defer but override up front. 

937 deferred = butler.getDeferred(ref, storageClass=new_sc) 

938 model = deferred.get() 

939 self.assertIs(type(model), new_sc.pytype) 

940 self.assertEqual(retrieved, model) 

941 

942 # Retrieve a component. Should be a tuple. 

943 data = butler.get("anything.data", dataId, storageClass="StructuredDataDataTestTuple") 

944 self.assertIs(type(data), tuple) 

945 self.assertEqual(data, tuple(retrieved.data)) 

946 

947 # Parameter on the write storage class should work regardless 

948 # of read storage class. 

949 data = butler.get( 

950 "anything.data", 

951 dataId, 

952 storageClass="StructuredDataDataTestTuple", 

953 parameters={"slice": slice(2, 4)}, 

954 ) 

955 self.assertEqual(len(data), 2) 

956 

957 # Try a parameter that is known to the read storage class but not 

958 # the write storage class. 

959 with self.assertRaises(KeyError): 

960 butler.get( 

961 "anything.data", 

962 dataId, 

963 storageClass="StructuredDataDataTestTuple", 

964 parameters={"xslice": slice(2, 4)}, 

965 ) 

966 

967 def testComponentFromOverriddenStorageClass(self) -> None: 

968 """Test component get where the component is only defined by the 

969 read storage class and not by the storage class used to write. 

970 """ 

971 # StructuredDataNoComponents defines no components at all, whereas 

972 # MetricsConversion (which it can be converted to) defines several. 

973 write_sc = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

974 read_sc = self.storageClassFactory.getStorageClass("MetricsConversion") 

975 self.assertFalse(write_sc.allComponents()) 

976 self.assertIn("summary", read_sc.allComponents()) 

977 

978 butler, datasetType = self.create_butler(self.default_run, write_sc, "unstructured") 

979 

980 metric = makeExampleMetrics() 

981 dataId = {"instrument": "DummyCamComp", "visit": 423} 

982 ref = butler.put(metric, datasetType, dataId) 

983 

984 # The composite conversion on its own must work. 

985 self.assertIs(type(butler.get(ref, storageClass=read_sc)), read_sc.pytype) 

986 

987 # A component of the converted composite, requested via a ref. 

988 component_ref = ref.overrideStorageClass(read_sc).makeComponentRef("summary") 

989 self.assertEqual(butler.get(component_ref), metric.summary) 

990 

991 # The same component, requested via a deferred handle that was given 

992 # the storage class override up front. 

993 deferred = butler.getDeferred(ref, storageClass=read_sc) 

994 self.assertEqual(deferred.get(component="summary"), metric.summary) 

995 

996 # A component whose storage class is also overridden, on top of the 

997 # storage class the read composite declares for it. 

998 converted = butler.get(component_ref, storageClass="DictConvertibleModel") 

999 self.assertIsInstance(converted, DictConvertibleModel) 

1000 self.assertEqual(converted.content, metric.summary) 

1001 

1002 # The handle storage class applies to the composite and so selects the 

1003 # component, while the one given to get() applies to the component. 

1004 converted = deferred.get(component="summary", storageClass="DictConvertibleModel") 

1005 self.assertIsInstance(converted, DictConvertibleModel) 

1006 self.assertEqual(converted.content, metric.summary) 

1007 

1008 def testPytypePutCoercion(self) -> None: 

1009 """Test python type coercion on Butler.get and put.""" 

1010 # Store some data with the normal example storage class. 

1011 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

1012 datasetTypeName = "test_metric" 

1013 butler, _ = self.create_butler(self.default_run, storageClass, datasetTypeName) 

1014 

1015 dataId = {"instrument": "DummyCamComp", "visit": 423} 

1016 

1017 # Put a dict and this should coerce to a MetricsExample 

1018 test_dict = {"summary": {"a": 1}, "output": {"b": 2}} 

1019 metric_ref = butler.put(test_dict, datasetTypeName, dataId=dataId, visit=424) 

1020 test_metric = butler.get(metric_ref) 

1021 self.assertEqual(get_full_type_name(test_metric), "lsst.daf.butler.tests.MetricsExample") 

1022 self.assertEqual(test_metric.summary, test_dict["summary"]) 

1023 self.assertEqual(test_metric.output, test_dict["output"]) 

1024 

1025 # Check that the put still works if a DatasetType is given with 

1026 # a definition matching this python type. 

1027 registry_type = butler.get_dataset_type(datasetTypeName) 

1028 this_type = DatasetType(datasetTypeName, registry_type.dimensions, "StructuredDataDictJson") 

1029 metric2_ref = butler.put(test_dict, this_type, dataId=dataId, visit=425) 

1030 self.assertEqual(metric2_ref.datasetType, registry_type) 

1031 

1032 # The get will return the type expected by registry. 

1033 test_metric2 = butler.get(metric2_ref) 

1034 self.assertEqual(get_full_type_name(test_metric2), "lsst.daf.butler.tests.MetricsExample") 

1035 

1036 # Make a new DatasetRef with the compatible but different DatasetType. 

1037 # This should now return a dict. 

1038 new_ref = DatasetRef(this_type, metric2_ref.dataId, id=metric2_ref.id, run=metric2_ref.run) 

1039 test_dict2 = butler.get(new_ref) 

1040 self.assertEqual(get_full_type_name(test_dict2), "dict") 

1041 

1042 # Get it again with the wrong dataset type definition using get() 

1043 # rather than get(). This should be consistent with get() 

1044 # behavior and return the type of the DatasetType. 

1045 test_dict3 = butler.get(this_type, dataId=dataId, visit=425) 

1046 self.assertEqual(get_full_type_name(test_dict3), "dict") 

1047 

1048 def test_ingest_zip(self) -> None: 

1049 """Create butler, export data, delete data, import from Zip.""" 

1050 butler, dataset_type = self.create_butler( 

1051 run=self.default_run, storageClass="StructuredData", datasetTypeName="metrics" 

1052 ) 

1053 

1054 metric = makeExampleMetrics() 

1055 refs = [] 

1056 for visit in (423, 424, 425): 

1057 ref = butler.put(metric, dataset_type, instrument="DummyCamComp", visit=visit) 

1058 refs.append(ref) 

1059 

1060 # Retrieve a Zip file. 

1061 with tempfile.TemporaryDirectory(ignore_cleanup_errors=True) as tmpdir: 

1062 zip = butler.retrieve_artifacts_zip(refs, destination=tmpdir) 

1063 

1064 # Ingest will fail. 

1065 with self.assertRaises(ConflictingDefinitionError): 

1066 butler.ingest_zip(zip) 

1067 

1068 # Clear out the collection. 

1069 butler.removeRuns([self.default_run]) 

1070 self.assertFalse(butler.exists(refs[0])) 

1071 

1072 butler.ingest_zip(zip, transfer="copy") 

1073 self.assertGreater(butler._metrics.time_in_ingest, 0.0) 

1074 self.assertEqual(butler._metrics.n_ingest, len(refs)) 

1075 

1076 # Check that it fails if we try it again. 

1077 with self.assertRaises(ConflictingDefinitionError): 

1078 butler.ingest_zip(zip, transfer="copy") 

1079 

1080 # This will be a no-op. 

1081 butler.ingest_zip(zip, transfer="copy", skip_existing=True) 

1082 

1083 # Create an entirely new local file butler in this temp directory. 

1084 new_butler_cfg = Butler.makeRepo(tmpdir) 

1085 new_butler = Butler.from_config(new_butler_cfg, writeable=True) 

1086 self.enterContext(new_butler) 

1087 

1088 # This will fail since dimensions records are missing. 

1089 with self.assertRaises(ConflictingDefinitionError): 

1090 new_butler.ingest_zip(zip, transfer="copy") 

1091 

1092 # Dry run should work. 

1093 new_butler.ingest_zip(zip, transfer="copy", dry_run=True) 

1094 

1095 new_butler.ingest_zip(zip, transfer="copy", transfer_dimensions=True) 

1096 self.assertTrue(butler.exists(refs[0])) 

1097 

1098 # Check that the refs can be read again. 

1099 _ = [butler.get(ref) for ref in refs] 

1100 

1101 uri = butler.getURI(refs[2]) 

1102 self.assertTrue(uri.exists()) 

1103 

1104 # Delete one dataset. The Zip file should still exist and allow 

1105 # remaining refs to be read. 

1106 butler.pruneDatasets([refs[0]], purge=True, unstore=True) 

1107 self.assertTrue(uri.exists()) 

1108 

1109 metric2 = butler.get(refs[1]) 

1110 self.assertEqual(metric2, metric, msg=f"{metric2} != {metric}") 

1111 

1112 butler.removeRuns([self.default_run]) 

1113 self.assertFalse(uri.exists()) 

1114 self.assertFalse(butler.exists(refs[-1])) 

1115 

1116 with self.assertRaises(ValueError): 

1117 butler.retrieve_artifacts_zip([], destination=".") 

1118 

1119 def testIngest(self) -> None: 

1120 butler = self.create_empty_butler(run=self.default_run) 

1121 

1122 # Create and register a DatasetType 

1123 dimensions = butler.dimensions.conform(["instrument", "visit", "detector"]) 

1124 

1125 storageClass = self.storageClassFactory.getStorageClass("StructuredDataDictYaml") 

1126 datasetTypeName = "metric" 

1127 

1128 datasetType = self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry) 

1129 

1130 # Add needed Dimensions 

1131 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"}) 

1132 butler.registry.insertDimensionData( 

1133 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"} 

1134 ) 

1135 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20250101}) 

1136 for detector in (1, 2): 

1137 butler.registry.insertDimensionData( 

1138 "detector", {"instrument": "DummyCamComp", "id": detector, "full_name": f"detector{detector}"} 

1139 ) 

1140 

1141 butler.registry.insertDimensionData( 

1142 "visit", 

1143 { 

1144 "instrument": "DummyCamComp", 

1145 "id": 423, 

1146 "name": "fourtwentythree", 

1147 "physical_filter": "d-r", 

1148 "day_obs": 20250101, 

1149 }, 

1150 { 

1151 "instrument": "DummyCamComp", 

1152 "id": 424, 

1153 "name": "fourtwentyfour", 

1154 "physical_filter": "d-r", 

1155 "day_obs": 20250101, 

1156 }, 

1157 ) 

1158 

1159 formatter = doImportType("lsst.daf.butler.formatters.yaml.YamlFormatter") 

1160 dataRoot = os.path.join(TESTDIR, "data", "basic") 

1161 datasets = [] 

1162 # Test one DatasetRef with a run that exists, and the other with a run 

1163 # that doesn't exist, to verify that run collections are created when 

1164 # required. 

1165 runs = {1: self.default_run, 2: "a/new/run"} 

1166 for detector in (1, 2): 

1167 detector_name = f"detector_{detector}" 

1168 metricFile = os.path.join(dataRoot, f"{detector_name}.yaml") 

1169 dataId = butler.registry.expandDataId( 

1170 {"instrument": "DummyCamComp", "visit": 423, "detector": detector} 

1171 ) 

1172 # Create a DatasetRef for ingest 

1173 refIn = DatasetRef(datasetType, dataId, run=runs[detector]) 

1174 

1175 datasets.append(FileDataset(path=metricFile, refs=[refIn], formatter=formatter)) 

1176 

1177 butler.ingest(*datasets, transfer="copy") 

1178 

1179 dataId1 = {"instrument": "DummyCamComp", "detector": 1, "visit": 423} 

1180 dataId2 = {"instrument": "DummyCamComp", "detector": 2, "visit": 423} 

1181 

1182 metrics1 = butler.get(datasetTypeName, dataId1) 

1183 metrics2 = butler.get(datasetTypeName, dataId2, collections="a/new/run") 

1184 self.assertNotEqual(metrics1, metrics2) 

1185 

1186 # Compare URIs 

1187 uri1 = butler.getURI(datasetTypeName, dataId1) 

1188 uri2 = butler.getURI(datasetTypeName, dataId2, collections="a/new/run") 

1189 self.assertFalse(self.are_uris_equivalent(uri1, uri2), f"Cf. {uri1} with {uri2}") 

1190 

1191 # Re-ingesting the same datasets raises an error with 

1192 # skip_existing=False. 

1193 with self.assertRaises(ConflictingDefinitionError): 

1194 butler.ingest(*datasets, transfer="copy") 

1195 # skip_existing=True makes it a no-op to re-ingest the same datasets. 

1196 butler.ingest(*datasets, transfer="copy", skip_existing=True) 

1197 

1198 # Now do a multi-dataset but single file ingest 

1199 metricFile = os.path.join(dataRoot, "detectors.yaml") 

1200 refs = [] 

1201 for detector in (1, 2): 

1202 detector_name = f"detector_{detector}" 

1203 dataId = butler.registry.expandDataId( 

1204 {"instrument": "DummyCamComp", "visit": 424, "detector": detector} 

1205 ) 

1206 # Create a DatasetRef for ingest 

1207 refs.append(DatasetRef(datasetType, dataId, run=self.default_run)) 

1208 

1209 # Test "move" transfer to ensure that the files themselves 

1210 # have disappeared following ingest. 

1211 with ResourcePath.temporary_uri(suffix=".yaml") as tempFile: 

1212 tempFile.transfer_from(ResourcePath(metricFile), transfer="copy") 

1213 

1214 datasets = [] 

1215 datasets.append(FileDataset(path=tempFile, refs=refs, formatter=MultiDetectorFormatter)) 

1216 

1217 # For first ingest use copy. 

1218 butler.ingest(*datasets, transfer="copy", record_validation_info=False) 

1219 

1220 # Now try to ingest again in "execution butler" mode where 

1221 # the registry entries exist but the datastore does not have 

1222 # the files. We also need to strip the dimension records to ensure 

1223 # that they will be re-added by the ingest. 

1224 ref = datasets[0].refs[0] 

1225 datasets[0].refs = [ 

1226 cast( 

1227 DatasetRef, 

1228 butler.find_dataset(ref.datasetType, data_id=ref.dataId, collections=ref.run), 

1229 ) 

1230 for ref in datasets[0].refs 

1231 ] 

1232 all_refs = [] 

1233 for dataset in datasets: 

1234 refs = [] 

1235 for ref in dataset.refs: 

1236 # Create a dict from the dataId to drop the records. 

1237 new_data_id = dict(ref.dataId.required) 

1238 new_ref = butler.find_dataset(ref.datasetType, new_data_id, collections=ref.run) 

1239 assert new_ref is not None 

1240 self.assertFalse(new_ref.dataId.hasRecords()) 

1241 refs.append(new_ref) 

1242 dataset.refs = refs 

1243 all_refs.extend(dataset.refs) 

1244 butler.pruneDatasets(all_refs, disassociate=False, unstore=True, purge=False) 

1245 

1246 # Use move mode to test that the file is deleted. Also 

1247 # disable recording of file size. 

1248 butler.ingest(*datasets, transfer="move", record_validation_info=False) 

1249 

1250 # Check that every ref now has records. 

1251 for dataset in datasets: 

1252 for ref in dataset.refs: 

1253 self.assertTrue(ref.dataId.hasRecords()) 

1254 

1255 # Ensure that the file has disappeared. 

1256 self.assertFalse(tempFile.exists()) 

1257 

1258 # Check that the datastore recorded no file size. 

1259 # Not all datastores can support this. 

1260 try: 

1261 infos = butler._datastore.getStoredItemsInfo(datasets[0].refs[0]) # type: ignore[attr-defined] 

1262 self.assertEqual(infos[0].file_size, -1) 

1263 except AttributeError: 

1264 pass 

1265 

1266 dataId1 = {"instrument": "DummyCamComp", "detector": 1, "visit": 424} 

1267 dataId2 = {"instrument": "DummyCamComp", "detector": 2, "visit": 424} 

1268 

1269 multi1 = butler.get(datasetTypeName, dataId1) 

1270 multi2 = butler.get(datasetTypeName, dataId2) 

1271 

1272 self.assertEqual(multi1, metrics1) 

1273 self.assertEqual(multi2, metrics2) 

1274 

1275 # Compare URIs 

1276 uri1 = butler.getURI(datasetTypeName, dataId1) 

1277 uri2 = butler.getURI(datasetTypeName, dataId2) 

1278 self.assertTrue(self.are_uris_equivalent(uri1, uri2), f"Cf. {uri1} with {uri2}") 

1279 

1280 # Test that removing one does not break the second 

1281 # This line will issue a warning log message for a ChainedDatastore 

1282 # that uses an InMemoryDatastore since in-memory can not ingest 

1283 # files. 

1284 butler.pruneDatasets([datasets[0].refs[0]], unstore=True, disassociate=False) 

1285 self.assertFalse(butler.exists(datasetTypeName, dataId1)) 

1286 self.assertTrue(butler.exists(datasetTypeName, dataId2)) 

1287 multi2b = butler.get(datasetTypeName, dataId2) 

1288 self.assertEqual(multi2, multi2b) 

1289 

1290 # Ensure we can ingest 0 datasets 

1291 datasets = [] 

1292 butler.ingest(*datasets) 

1293 

1294 def testPickle(self) -> None: 

1295 """Test pickle support.""" 

1296 butler = self.create_empty_butler(run=self.default_run) 

1297 assert isinstance(butler, DirectButler), "Expect DirectButler in configuration" 

1298 butlerOut = pickle.loads(pickle.dumps(butler)) 

1299 self.enterContext(butlerOut) 

1300 self.assertIsInstance(butlerOut, Butler) 

1301 self.assertEqual(butlerOut._config, butler._config) 

1302 self.assertEqual(list(butlerOut.collections.defaults), list(butler.collections.defaults)) 

1303 self.assertEqual(butlerOut.run, butler.run) 

1304 

1305 def testGetDatasetTypes(self) -> None: 

1306 butler = self.create_empty_butler(run=self.default_run) 

1307 dimensions = butler.dimensions.conform(["instrument", "visit", "physical_filter"]) 

1308 dimensionEntries: list[tuple[str, list[Mapping[str, Any]]]] = [ 

1309 ( 

1310 "instrument", 

1311 [ 

1312 {"instrument": "DummyCam"}, 

1313 {"instrument": "DummyHSC"}, 

1314 {"instrument": "DummyCamComp"}, 

1315 ], 

1316 ), 

1317 ("physical_filter", [{"instrument": "DummyCam", "name": "d-r", "band": "R"}]), 

1318 ("day_obs", [{"instrument": "DummyCam", "id": 20250101}]), 

1319 ( 

1320 "visit", 

1321 [ 

1322 { 

1323 "instrument": "DummyCam", 

1324 "id": 42, 

1325 "name": "fortytwo", 

1326 "physical_filter": "d-r", 

1327 "day_obs": 20250101, 

1328 } 

1329 ], 

1330 ), 

1331 ] 

1332 storageClass = self.storageClassFactory.getStorageClass("StructuredData") 

1333 # Add needed Dimensions 

1334 for element, data in dimensionEntries: 

1335 butler.registry.insertDimensionData(element, *data) 

1336 

1337 # When a DatasetType is added to the registry entries are not created 

1338 # for components but querying them can return the components. 

1339 datasetTypeNames = {"metric", "metric2", "metric4", "metric33", "pvi", "paramtest"} 

1340 components = set() 

1341 for datasetTypeName in datasetTypeNames: 

1342 # Create and register a DatasetType 

1343 self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry) 

1344 

1345 for componentName in storageClass.components: 

1346 components.add(DatasetType.nameWithComponent(datasetTypeName, componentName)) 

1347 

1348 fromRegistry: set[DatasetType] = set() 

1349 for parent_dataset_type in butler.registry.queryDatasetTypes(): 

1350 fromRegistry.add(parent_dataset_type) 

1351 fromRegistry.update(parent_dataset_type.makeAllComponentDatasetTypes()) 

1352 self.assertEqual({d.name for d in fromRegistry}, datasetTypeNames | components) 

1353 

1354 # Query with wildcard. 

1355 dataset_types = butler.registry.queryDatasetTypes("metric*") 

1356 self.assertEqual(len(dataset_types), 4, f"Got: {dataset_types}") 

1357 # but not regex. 

1358 with self.assertRaises(DatasetTypeExpressionError): 

1359 butler.registry.queryDatasetTypes(["pvi", re.compile("metric.*")]) 

1360 

1361 # Now that we have some dataset types registered, validate them 

1362 butler.validateConfiguration( 

1363 ignore=[ 

1364 "test_metric_comp", 

1365 "metric3", 

1366 "metric5", 

1367 "calexp", 

1368 "DummySC", 

1369 "datasetType.component", 

1370 "random_data", 

1371 "random_data_2", 

1372 ] 

1373 ) 

1374 

1375 # Add a new datasetType that will fail template validation 

1376 self.addDatasetType("test_metric_comp", dimensions, storageClass, butler.registry) 

1377 if self.validationCanFail: 

1378 with self.assertRaises(ValidationError): 

1379 butler.validateConfiguration() 

1380 

1381 # Rerun validation but with a subset of dataset type names 

1382 butler.validateConfiguration(datasetTypeNames=["metric4"]) 

1383 

1384 # Rerun validation but ignore the bad datasetType 

1385 butler.validateConfiguration( 

1386 ignore=[ 

1387 "test_metric_comp", 

1388 "metric3", 

1389 "metric5", 

1390 "calexp", 

1391 "DummySC", 

1392 "datasetType.component", 

1393 "random_data", 

1394 "random_data_2", 

1395 ] 

1396 ) 

1397 

1398 def testTransaction(self) -> None: 

1399 butler = self.create_empty_butler(run=self.default_run) 

1400 datasetTypeName = "test_metric" 

1401 dimensions = butler.dimensions.conform(["instrument", "visit"]) 

1402 dimensionEntries: tuple[tuple[str, Mapping[str, Any]], ...] = ( 

1403 ("instrument", {"instrument": "DummyCam"}), 

1404 ("physical_filter", {"instrument": "DummyCam", "name": "d-r", "band": "R"}), 

1405 ("day_obs", {"instrument": "DummyCam", "id": 20250101}), 

1406 ( 

1407 "visit", 

1408 { 

1409 "instrument": "DummyCam", 

1410 "id": 42, 

1411 "name": "fortytwo", 

1412 "physical_filter": "d-r", 

1413 "day_obs": 20250101, 

1414 }, 

1415 ), 

1416 ) 

1417 storageClass = self.storageClassFactory.getStorageClass("StructuredData") 

1418 metric = makeExampleMetrics() 

1419 dataId = {"instrument": "DummyCam", "visit": 42} 

1420 # Create and register a DatasetType 

1421 datasetType = self.addDatasetType(datasetTypeName, dimensions, storageClass, butler.registry) 

1422 with self.assertRaises(TransactionTestError): 

1423 with butler.transaction(): 

1424 # Add needed Dimensions 

1425 for args in dimensionEntries: 

1426 butler.registry.insertDimensionData(*args) 

1427 # Store a dataset 

1428 ref = butler.put(metric, datasetTypeName, dataId) 

1429 self.assertIsInstance(ref, DatasetRef) 

1430 # Test get of a ref. 

1431 metricOut = butler.get(ref) 

1432 self.assertEqual(metric, metricOut) 

1433 # Test get 

1434 metricOut = butler.get(datasetTypeName, dataId) 

1435 self.assertEqual(metric, metricOut) 

1436 # Check we can get components 

1437 self.assertGetComponents(butler, ref, ("summary", "data", "output"), metric) 

1438 raise TransactionTestError("This should roll back the entire transaction") 

1439 with self.assertRaises(DataIdValueError, msg=f"Check can't expand DataId {dataId}"): 

1440 butler.registry.expandDataId(dataId) 

1441 # Should raise LookupError for missing data ID value 

1442 with self.assertRaises(LookupError, msg=f"Check can't get by {datasetTypeName} and {dataId}"): 

1443 butler.get(datasetTypeName, dataId) 

1444 # Also check explicitly if Dataset entry is missing 

1445 self.assertIsNone(butler.find_dataset(datasetType, dataId, collections=butler.collections.defaults)) 

1446 # Direct retrieval should not find the file in the Datastore 

1447 with self.assertRaises(FileNotFoundError, msg=f"Check {ref} can't be retrieved directly"): 

1448 butler.get(ref) 

1449 

1450 def testMakeRepo(self) -> None: 

1451 """Test that we can write butler configuration to a new repository via 

1452 the Butler.makeRepo interface and then instantiate a butler from the 

1453 repo root. 

1454 """ 

1455 # Do not run the test if we know this datastore configuration does 

1456 # not support a file system root 

1457 if self.fullConfigKey is None: 

1458 return 

1459 

1460 # create two separate directories 

1461 root1 = tempfile.mkdtemp(dir=self.root) 

1462 root2 = tempfile.mkdtemp(dir=self.root) 

1463 

1464 self.assertFalse(Butler.has_repo_config(root1)) 

1465 butlerConfig = Butler.makeRepo(root1, config=Config(self.configFile)) 

1466 self.assertTrue(Butler.has_repo_config(root1)) 

1467 limited = Config(self.configFile) 

1468 butler1 = Butler.from_config(butlerConfig) 

1469 self.enterContext(butler1) 

1470 assert isinstance(butler1, DirectButler), "Expect DirectButler in configuration" 

1471 butlerConfig = Butler.makeRepo(root2, standalone=True, config=Config(self.configFile)) 

1472 full = Config(self.tmpConfigFile) 

1473 butler2 = Butler.from_config(butlerConfig) 

1474 self.enterContext(butler2) 

1475 assert isinstance(butler2, DirectButler), "Expect DirectButler in configuration" 

1476 # Butlers should have the same configuration regardless of whether 

1477 # defaults were expanded. 

1478 self.assertEqual(butler1._config, butler2._config) 

1479 # Config files loaded directly should not be the same. 

1480 self.assertNotEqual(limited, full) 

1481 # Make sure "limited" doesn't have a few keys we know it should be 

1482 # inheriting from defaults. 

1483 self.assertIn(self.fullConfigKey, full) 

1484 self.assertNotIn(self.fullConfigKey, limited) 

1485 

1486 # Collections don't appear until something is put in them 

1487 collections1 = set(butler1.registry.queryCollections()) 

1488 self.assertEqual(collections1, set()) 

1489 self.assertEqual(set(butler2.registry.queryCollections()), collections1) 

1490 

1491 # Check that a config with no associated file name will not 

1492 # work properly with relocatable Butler repo 

1493 butlerConfig.configFile = None 

1494 with self.assertRaises(ValueError): 

1495 Butler.from_config(butlerConfig) 

1496 

1497 with self.assertRaises(FileExistsError): 

1498 Butler.makeRepo(self.root, standalone=True, config=Config(self.configFile), overwrite=False) 

1499 

1500 def testStringification(self) -> None: 

1501 butler = Butler.from_config(self.tmpConfigFile, run=self.default_run) 

1502 self.enterContext(butler) 

1503 butlerStr = str(butler) 

1504 

1505 if self.datastoreStr is not None: 1505 ↛ 1508line 1505 didn't jump to line 1508 because the condition on line 1505 was always true

1506 for testStr in self.datastoreStr: 

1507 self.assertIn(testStr, butlerStr) 

1508 if self.registryStr is not None: 1508 ↛ 1511line 1508 didn't jump to line 1511 because the condition on line 1508 was always true

1509 self.assertIn(self.registryStr, butlerStr) 

1510 

1511 datastoreName = butler._datastore.name 

1512 if self.datastoreName is not None: 1512 ↛ exitline 1512 didn't return from function 'testStringification' because the condition on line 1512 was always true

1513 for testStr in self.datastoreName: 

1514 self.assertIn(testStr, datastoreName) 

1515 

1516 def testButlerRewriteDataId(self) -> None: 

1517 """Test that dataIds can be rewritten based on dimension records.""" 

1518 butler = self.create_empty_butler(run=self.default_run) 

1519 

1520 storageClass = self.storageClassFactory.getStorageClass("StructuredDataDict") 

1521 datasetTypeName = "random_data" 

1522 

1523 # Create dimension records. 

1524 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"}) 

1525 butler.registry.insertDimensionData( 

1526 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"} 

1527 ) 

1528 butler.registry.insertDimensionData( 

1529 "detector", {"instrument": "DummyCamComp", "id": 1, "full_name": "det1"} 

1530 ) 

1531 

1532 dimensions = butler.dimensions.conform(["instrument", "exposure"]) 

1533 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

1534 butler.registry.registerDatasetType(datasetType) 

1535 

1536 n_exposures = 5 

1537 dayobs = 20210530 

1538 

1539 # Create records for multiple day_obs but same seq_num to test that 

1540 # we are constraining gets properly when day_obs/seq_num is used 

1541 # for an exposure. Second day is year in future but is not used. 

1542 for day_obs in (dayobs, dayobs + 1_00_00): 

1543 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": day_obs}) 

1544 

1545 for i in range(n_exposures): 

1546 group_name = f"group_{day_obs}_{i}" 

1547 butler.registry.insertDimensionData( 

1548 "group", {"instrument": "DummyCamComp", "name": group_name} 

1549 ) 

1550 butler.registry.insertDimensionData( 

1551 "exposure", 

1552 { 

1553 "instrument": "DummyCamComp", 

1554 "id": day_obs + i, 

1555 "obs_id": f"exp_{day_obs}_{i}", 

1556 "seq_num": i, 

1557 "day_obs": day_obs, 

1558 "physical_filter": "d-r", 

1559 "group": group_name, 

1560 }, 

1561 ) 

1562 

1563 # Write some data. 

1564 for i in range(n_exposures): 

1565 metric = {"something": i, "other": "metric", "list": [2 * x for x in range(i)]} 

1566 

1567 # Use the seq_num for the put to test rewriting. 

1568 dataId = {"seq_num": i, "day_obs": dayobs, "instrument": "DummyCamComp", "physical_filter": "d-r"} 

1569 ref = butler.put(metric, datasetTypeName, dataId=dataId) 

1570 

1571 # Check that the exposure is correct in the dataId 

1572 self.assertEqual(ref.dataId["exposure"], dayobs + i) 

1573 

1574 # and check that we can get the dataset back with the same dataId 

1575 new_metric = butler.get(datasetTypeName, dataId=dataId) 

1576 self.assertEqual(new_metric, metric) 

1577 

1578 # Check that we can find the datasets using the day_obs or the 

1579 # exposure.day_obs. 

1580 datasets_1 = list( 

1581 butler.registry.queryDatasets( 

1582 datasetType, 

1583 collections=self.default_run, 

1584 where="day_obs = :dayObs AND instrument = :instr", 

1585 bind={"dayObs": dayobs, "instr": "DummyCamComp"}, 

1586 ) 

1587 ) 

1588 datasets_2 = list( 

1589 butler.registry.queryDatasets( 

1590 datasetType, 

1591 collections=self.default_run, 

1592 where="exposure.day_obs = :dayObs AND instrument = :instr", 

1593 bind={"dayObs": dayobs, "instr": "DummyCamComp"}, 

1594 ) 

1595 ) 

1596 self.assertEqual(datasets_1, datasets_2) 

1597 

1598 def testGetDatasetCollectionCaching(self): 

1599 # Prior to DM-41117, there was a bug where get_dataset would throw 

1600 # MissingCollectionError if you tried to fetch a dataset that was added 

1601 # after the collection cache was last updated. 

1602 reader_butler, datasetType = self.create_butler(self.default_run, "int", "datasettypename") 

1603 writer_butler = self.create_empty_butler(writeable=True, run="new_run") 

1604 dataId = {"instrument": "DummyCamComp", "visit": 423} 

1605 put_ref = writer_butler.put(123, datasetType, dataId) 

1606 get_ref = reader_butler.get_dataset(put_ref.id) 

1607 self.assertEqual(get_ref.id, put_ref.id) 

1608 # Also works when looking up via a hexadecimal string instead of a UUID 

1609 # instance. 

1610 hex_ref = reader_butler.get_dataset(put_ref.id.hex) 

1611 self.assertEqual(hex_ref.id, put_ref.id) 

1612 

1613 def testCollectionChainRedefine(self): 

1614 butler = self._setup_to_test_collection_chain() 

1615 

1616 butler.collections.redefine_chain("chain", "a") 

1617 self._check_chain(butler, ["a"]) 

1618 

1619 # Duplicates are removed from the list of children 

1620 butler.collections.redefine_chain("chain", ["c", "b", "c"]) 

1621 self._check_chain(butler, ["c", "b"]) 

1622 

1623 # Empty list clears the chain 

1624 butler.collections.redefine_chain("chain", []) 

1625 self._check_chain(butler, []) 

1626 

1627 self._test_common_chain_functionality(butler, butler.collections.redefine_chain) 

1628 

1629 def testCollectionChainPrepend(self): 

1630 butler = self._setup_to_test_collection_chain() 

1631 

1632 # Duplicates are removed from the list of children 

1633 butler.collections.prepend_chain("chain", ["c", "b", "c"]) 

1634 self._check_chain(butler, ["c", "b"]) 

1635 

1636 # Prepend goes on the front of existing chain 

1637 butler.collections.prepend_chain("chain", ["a"]) 

1638 self._check_chain(butler, ["a", "c", "b"]) 

1639 

1640 # Empty prepend does nothing 

1641 butler.collections.prepend_chain("chain", []) 

1642 self._check_chain(butler, ["a", "c", "b"]) 

1643 

1644 # Prepending children that already exist in the chain removes them from 

1645 # their current position. 

1646 butler.collections.prepend_chain("chain", ["d", "b", "c"]) 

1647 self._check_chain(butler, ["d", "b", "c", "a"]) 

1648 

1649 self._test_common_chain_functionality(butler, butler.collections.prepend_chain) 

1650 

1651 def testCollectionChainExtend(self): 

1652 butler = self._setup_to_test_collection_chain() 

1653 

1654 # Duplicates are removed from the list of children 

1655 butler.collections.extend_chain("chain", ["c", "b", "c"]) 

1656 self._check_chain(butler, ["c", "b"]) 

1657 

1658 # Extend goes on the end of existing chain 

1659 butler.collections.extend_chain("chain", ["a"]) 

1660 self._check_chain(butler, ["c", "b", "a"]) 

1661 

1662 # Empty extend does nothing 

1663 butler.collections.extend_chain("chain", []) 

1664 self._check_chain(butler, ["c", "b", "a"]) 

1665 

1666 # Extending children that already exist in the chain removes them from 

1667 # their current position. 

1668 butler.collections.extend_chain("chain", ["d", "b", "c"]) 

1669 self._check_chain(butler, ["a", "d", "b", "c"]) 

1670 

1671 self._test_common_chain_functionality(butler, butler.collections.extend_chain) 

1672 

1673 def testCollectionChainRemove(self) -> None: 

1674 butler = self._setup_to_test_collection_chain() 

1675 

1676 butler.collections.redefine_chain("chain", ["a", "b", "c", "d"]) 

1677 

1678 butler.collections.remove_from_chain("chain", "c") 

1679 self._check_chain(butler, ["a", "b", "d"]) 

1680 

1681 # Duplicates are allowed in the list of children 

1682 butler.collections.remove_from_chain("chain", ["b", "b", "a"]) 

1683 self._check_chain(butler, ["d"]) 

1684 

1685 # Empty remove does nothing 

1686 butler.collections.remove_from_chain("chain", []) 

1687 self._check_chain(butler, ["d"]) 

1688 

1689 # Removing children that aren't in the chain does nothing 

1690 butler.collections.remove_from_chain("chain", ["a", "chain"]) 

1691 self._check_chain(butler, ["d"]) 

1692 

1693 self._test_common_chain_functionality( 

1694 butler, butler.collections.remove_from_chain, skip_cycle_check=True 

1695 ) 

1696 

1697 def _setup_to_test_collection_chain(self) -> Butler: 

1698 butler = self.create_empty_butler(writeable=True) 

1699 

1700 butler.collections.register("chain", CollectionType.CHAINED) 

1701 

1702 runs = ["a", "b", "c", "d"] 

1703 for run in runs: 

1704 butler.collections.register(run) 

1705 

1706 butler.collections.register("staticchain", CollectionType.CHAINED) 

1707 butler.collections.redefine_chain("staticchain", ["a", "b"]) 

1708 

1709 return butler 

1710 

1711 def _check_chain(self, butler: Butler, expected: list[str]) -> None: 

1712 children = butler.collections.get_info("chain").children 

1713 self.assertEqual(expected, list(children)) 

1714 

1715 def _test_common_chain_functionality( 

1716 self, butler, func: Callable[[str, str | list[str]], Any], *, skip_cycle_check=False 

1717 ) -> None: 

1718 # Missing parent collection 

1719 with self.assertRaises(MissingCollectionError): 

1720 func("doesnotexist", []) 

1721 # Missing child collection 

1722 with self.assertRaises(MissingCollectionError): 

1723 func("chain", ["doesnotexist"]) 

1724 # Forbid operations on non-chained collections 

1725 with self.assertRaises(CollectionTypeError): 

1726 func("d", ["a"]) 

1727 

1728 # Prevent collection cycles 

1729 if not skip_cycle_check: 

1730 butler.collections.register("chain2", CollectionType.CHAINED) 

1731 func("chain2", "chain") 

1732 with self.assertRaises(CollectionCycleError): 

1733 func("chain", "chain2") 

1734 

1735 # Make sure none of the earlier operations interfered with unrelated 

1736 # chains. 

1737 self.assertEqual(["a", "b"], list(butler.collections.get_info("staticchain").children)) 

1738 

1739 with butler._caching_context(): 

1740 with self.assertRaisesRegex(RuntimeError, "Chained collection modification not permitted"): 

1741 func("chain", "a") 

1742 

1743 def test_transfer_dimension_records_from(self) -> None: 

1744 source_butler = self.create_empty_butler(writeable=True) 

1745 source_butler.import_(filename=_get_test_data_path("lsstcam-subset.yaml")) 

1746 

1747 visit_id = 2025120200439 

1748 exposure_id = visit_id 

1749 target_butler = self.enterContext(create_populated_sqlite_registry()) 

1750 target_butler.transfer_dimension_records_from( 

1751 source_butler, 

1752 [ 

1753 # Should trigger the lookup of visit and all its associated 

1754 # "populated_by" records (visit_detector_region, 

1755 # visit_definition, etc.) 

1756 DataCoordinate.standardize( 

1757 {"instrument": "LSSTCam", "visit": visit_id, "detector": 10}, 

1758 universe=source_butler.dimensions, 

1759 ), 

1760 # Shouldn't add any records to the lookup. 

1761 DataCoordinate.make_empty(source_butler.dimensions), 

1762 ], 

1763 ) 

1764 

1765 def _fetch_record(dimension: str) -> DimensionRecord: 

1766 records = target_butler.query_dimension_records(dimension) 

1767 self.assertEqual(len(records), 1) 

1768 return records[0] 

1769 

1770 visit = _fetch_record("visit") 

1771 self.assertEqual(visit.id, visit_id) 

1772 self.assertEqual(visit.day_obs, 20251202) 

1773 self.assertEqual(visit.target_name, "lowdust") 

1774 self.assertEqual(visit.seq_num, 439) 

1775 original_visit = source_butler.query_dimension_records("visit", instrument="LSSTCam", visit=visit_id)[ 

1776 0 

1777 ] 

1778 self.assertEqual(visit.region, original_visit.region) 

1779 self.assertEqual(visit.timespan, original_visit.timespan) 

1780 

1781 visit_detector_region = _fetch_record("visit_detector_region") 

1782 self.assertEqual(visit_detector_region.instrument, "LSSTCam") 

1783 self.assertEqual(visit_detector_region.detector, 10) 

1784 self.assertEqual(visit_detector_region.visit, visit_id) 

1785 original_visit_detector_region = source_butler.query_dimension_records( 

1786 "visit_detector_region", instrument="LSSTCam", visit=visit_id, detector=10 

1787 )[0] 

1788 self.assertEqual(visit_detector_region.region, original_visit_detector_region.region) 

1789 

1790 visit_definition = _fetch_record("visit_definition") 

1791 self.assertEqual(visit_definition.instrument, "LSSTCam") 

1792 self.assertEqual(visit_definition.exposure, 2025120200439) 

1793 self.assertEqual(visit_definition.visit, visit_id) 

1794 

1795 # The matching exposure record should have been pulled in via 

1796 # visit -> visit_definition. 

1797 exposure = _fetch_record("exposure") 

1798 self.assertEqual(exposure.instrument, "LSSTCam") 

1799 self.assertEqual(exposure.id, 2025120200439) 

1800 self.assertEqual(exposure.obs_id, "MC_O_20251202_000439") 

1801 original_exposure = source_butler.query_dimension_records( 

1802 "exposure", instrument="LSSTCam", exposure=exposure_id 

1803 )[0] 

1804 self.assertEqual(exposure.timespan, original_exposure.timespan) 

1805 

1806 group = _fetch_record("group") 

1807 self.assertEqual(group.instrument, "LSSTCam") 

1808 self.assertEqual(group.name, "2025-12-03T07:58:10.858") 

1809 

1810 visit_system_memberships = target_butler.query_dimension_records("visit_system_membership") 

1811 visit_system_memberships.sort(key=lambda record: record.visit_system) 

1812 self.assertEqual(len(visit_system_memberships), 2) 

1813 self.assertEqual(visit_system_memberships[0].visit_system, 0) 

1814 self.assertEqual(visit_system_memberships[1].visit_system, 2) 

1815 self.assertEqual(visit_system_memberships[0].visit, visit_id) 

1816 self.assertEqual(visit_system_memberships[1].visit, visit_id) 

1817 

1818 visit_systems = target_butler.query_dimension_records("visit_system") 

1819 visit_systems.sort(key=lambda record: record.id) 

1820 visit_system_memberships.sort(key=lambda record: record.visit_system) 

1821 self.assertEqual(visit_systems[0].id, 0) 

1822 self.assertEqual(visit_systems[1].id, 2) 

1823 self.assertEqual(visit_systems[0].name, "one-to-one") 

1824 self.assertEqual(visit_systems[1].name, "by-seq-start-end") 

1825 

1826 

1827class FileDatastoreButlerTests(ButlerTests): 

1828 """Common tests and specialization of ButlerTests for butlers backed 

1829 by datastores that inherit from FileDatastore. 

1830 """ 

1831 

1832 trustModeSupported = True 

1833 

1834 def testComponentFromOverriddenStorageClassWarns(self) -> None: 

1835 """Test that getting a component that only the read storage class 

1836 defines warns, since the whole dataset has to be retrieved and 

1837 converted before the component can be extracted. 

1838 """ 

1839 write_sc = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

1840 read_sc = self.storageClassFactory.getStorageClass("MetricsConversion") 

1841 butler, datasetType = self.create_butler(self.default_run, write_sc, "unstructured") 

1842 metric = makeExampleMetrics() 

1843 dataId = {"instrument": "DummyCamComp", "visit": 423} 

1844 ref = butler.put(metric, datasetType, dataId) 

1845 component_ref = ref.overrideStorageClass(read_sc).makeComponentRef("summary") 

1846 

1847 logger = "lsst.daf.butler.datastores.file_datastore.get" 

1848 with self.assertLogs(logger, level="WARNING") as cm: 

1849 self.assertEqual(butler.get(component_ref), metric.summary) 

1850 message = "\n".join(cm.output) 

1851 # The message must name the component, the storage class that lacks it 

1852 # along with the components it does have, and the storage class the 

1853 # dataset has to be converted to. 

1854 self.assertIn("summary", message) 

1855 self.assertIn(write_sc.name, message) 

1856 self.assertIn("components it does define: none", message) 

1857 self.assertIn(read_sc.name, message) 

1858 self.assertIn("less efficient", message) 

1859 

1860 # Reading a component that the write storage class does define must not 

1861 # warn. 

1862 composite_type = self.addDatasetType( 

1863 "composite", datasetType.dimensions, "StructuredData", butler.registry 

1864 ) 

1865 composite_ref = butler.put(metric, composite_type, dataId) 

1866 with self.assertNoLogs(logger, level="WARNING"): 

1867 self.assertEqual(butler.get(composite_ref.makeComponentRef("summary")), metric.summary) 

1868 

1869 def checkFileExists(self, root: str | ResourcePath, relpath: str | ResourcePath) -> bool: 

1870 """Check if file exists at a given path (relative to root). 

1871 

1872 Test testPutTemplates verifies actual physical existance of the files 

1873 in the requested location. 

1874 """ 

1875 uri = ResourcePath(root, forceDirectory=True) 

1876 return uri.join(relpath).exists() 

1877 

1878 def testPutTemplates(self) -> None: 

1879 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

1880 butler = self.create_empty_butler(run=self.default_run) 

1881 

1882 # Add needed Dimensions 

1883 butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"}) 

1884 butler.registry.insertDimensionData( 

1885 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"} 

1886 ) 

1887 butler.registry.insertDimensionData("day_obs", {"instrument": "DummyCamComp", "id": 20250101}) 

1888 butler.registry.insertDimensionData( 

1889 "visit", 

1890 { 

1891 "instrument": "DummyCamComp", 

1892 "id": 423, 

1893 "name": "v423", 

1894 "physical_filter": "d-r", 

1895 "day_obs": 20250101, 

1896 }, 

1897 ) 

1898 butler.registry.insertDimensionData( 

1899 "visit", 

1900 { 

1901 "instrument": "DummyCamComp", 

1902 "id": 425, 

1903 "name": "v425", 

1904 "physical_filter": "d-r", 

1905 "day_obs": 20250101, 

1906 }, 

1907 ) 

1908 

1909 # Create and store a dataset 

1910 metric = makeExampleMetrics() 

1911 

1912 # Create two almost-identical DatasetTypes (both will use default 

1913 # template) 

1914 dimensions = butler.dimensions.conform(["instrument", "visit"]) 

1915 butler.registry.registerDatasetType(DatasetType("metric1", dimensions, storageClass)) 

1916 butler.registry.registerDatasetType(DatasetType("metric2", dimensions, storageClass)) 

1917 butler.registry.registerDatasetType(DatasetType("metric3", dimensions, storageClass)) 

1918 

1919 dataId1 = {"instrument": "DummyCamComp", "visit": 423} 

1920 dataId2 = {"instrument": "DummyCamComp", "visit": 423, "physical_filter": "d-r"} 

1921 

1922 # Put with exactly the data ID keys needed 

1923 ref = butler.put(metric, "metric1", dataId1) 

1924 uri = butler.getURI(ref) 

1925 self.assertTrue(uri.exists()) 

1926 self.assertTrue( 

1927 uri.unquoted_path.endswith(f"{self.default_run}/metric1/??#?/d-r/DummyCamComp_423.pickle") 

1928 ) 

1929 

1930 # Check the template based on dimensions 

1931 if hasattr(butler._datastore, "templates"): 

1932 butler._datastore.templates.validateTemplates([ref]) 

1933 

1934 # Put with extra data ID keys (physical_filter is an optional 

1935 # dependency); should not change template (at least the way we're 

1936 # defining them to behave now; the important thing is that they 

1937 # must be consistent). 

1938 ref = butler.put(metric, "metric2", dataId2) 

1939 uri = butler.getURI(ref) 

1940 self.assertTrue(uri.exists()) 

1941 self.assertTrue( 

1942 uri.unquoted_path.endswith(f"{self.default_run}/metric2/d-r/DummyCamComp_v423.pickle") 

1943 ) 

1944 

1945 # Check the template based on dimensions 

1946 if hasattr(butler._datastore, "templates"): 

1947 butler._datastore.templates.validateTemplates([ref]) 

1948 

1949 # Use a template that has a typo in dimension record metadata. 

1950 # Easier to test with a butler that has a ref with records attached. 

1951 template = FileTemplate("a/{visit.name}/{id}_{visit.namex:?}.fits") 

1952 with self.assertLogs("lsst.daf.butler.datastore.file_templates", "INFO"): 

1953 path = template.format(ref) 

1954 self.assertEqual(path, f"a/v423/{ref.id}_fits") 

1955 

1956 template = FileTemplate("a/{visit.name}/{id}_{visit.namex}.fits") 

1957 with self.assertRaises(KeyError): 

1958 with self.assertLogs("lsst.daf.butler.datastore.file_templates", "INFO"): 

1959 template.format(ref) 

1960 

1961 # Now use a file template that will not result in unique filenames 

1962 with self.assertRaises(FileTemplateValidationError): 

1963 butler.put(metric, "metric3", dataId1) 

1964 

1965 def testImportExport(self) -> None: 

1966 # Run put/get tests just to create and populate a repo. 

1967 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

1968 self.runImportExportTest(storageClass) 

1969 

1970 @unittest.expectedFailure 

1971 def testImportExportVirtualComposite(self) -> None: 

1972 # Run put/get tests just to create and populate a repo. 

1973 storageClass = self.storageClassFactory.getStorageClass("StructuredComposite") 

1974 self.runImportExportTest(storageClass) 

1975 

1976 def runImportExportTest(self, storageClass: StorageClass) -> None: 

1977 """Test exporting and importing. 

1978 

1979 This test does an export to a temp directory and an import back 

1980 into a new temp directory repo. It does not assume a posix datastore. 

1981 """ 

1982 exportButler = self.runPutGetTest(storageClass, "test_metric") 

1983 

1984 # Test that we must have a file extension. 

1985 with self.assertRaises(ValueError): 

1986 with exportButler.export(filename="dump", directory=".") as export: 

1987 pass 

1988 

1989 # Test that unknown format is not allowed. 

1990 with self.assertRaises(ValueError): 

1991 with exportButler.export(filename="dump.fits", directory=".") as export: 

1992 pass 

1993 

1994 # Test that the repo actually has at least one dataset. 

1995 datasets = list(exportButler.registry.queryDatasets(..., collections=...)) 

1996 self.assertGreater(len(datasets), 0) 

1997 # Add a DimensionRecord that's unused by those datasets. 

1998 skymapRecord = {"name": "example_skymap", "hash": (50).to_bytes(8, byteorder="little")} 

1999 exportButler.registry.insertDimensionData("skymap", skymapRecord) 

2000 # Export and then import datasets. 

2001 with safeTestTempDir(TESTDIR) as exportDir: 

2002 exportFile = os.path.join(exportDir, "exports.yaml") 

2003 with exportButler.export(filename=exportFile, directory=exportDir, transfer="auto") as export: 

2004 export.saveDatasets(datasets) 

2005 # Export the same datasets again. This should quietly do 

2006 # nothing because of internal deduplication, and it shouldn't 

2007 # complain about being asked to export the "htm7" elements even 

2008 # though there aren't any in these datasets or in the database. 

2009 export.saveDatasets(datasets, elements=["htm7"]) 

2010 # Save one of the data IDs again; this should be harmless 

2011 # because of internal deduplication. 

2012 export.saveDataIds([datasets[0].dataId]) 

2013 # Save some dimension records directly. 

2014 export.saveDimensionData("skymap", [skymapRecord]) 

2015 self.assertTrue(os.path.exists(exportFile)) 

2016 with safeTestTempDir(TESTDIR) as importDir: 

2017 # We always want this to be a local posix butler 

2018 Butler.makeRepo(importDir, config=Config(os.path.join(TESTDIR, "config/basic/butler.yaml"))) 

2019 # Calling script.butlerImport tests the implementation of the 

2020 # butler command line interface "import" subcommand. Functions 

2021 # in the script folder are generally considered protected and 

2022 # should not be used as public api. 

2023 with open(exportFile) as f: 

2024 script.butlerImport( 

2025 importDir, 

2026 export_file=f, 

2027 directory=exportDir, 

2028 transfer="auto", 

2029 skip_dimensions=None, 

2030 ) 

2031 importButler = Butler.from_config(importDir, run=self.default_run) 

2032 self.enterContext(importButler) 

2033 for ref in datasets: 

2034 with self.subTest(ref=repr(ref)): 

2035 # Test for existence by passing in the DatasetType and 

2036 # data ID separately, to avoid lookup by dataset_id. 

2037 self.assertTrue(importButler.exists(ref.datasetType, ref.dataId)) 

2038 self.assertEqual( 

2039 list(importButler.registry.queryDimensionRecords("skymap")), 

2040 [importButler.dimensions["skymap"].RecordClass(**skymapRecord)], 

2041 ) 

2042 

2043 def testRemoveRuns(self) -> None: 

2044 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2045 butler = self.create_empty_butler(writeable=True) 

2046 # Load registry data with dimensions to hang datasets off of. 

2047 butler.import_(filename=ResourcePath("resource://lsst.daf.butler/tests/registry_data/base.yaml")) 

2048 # Add some RUN-type collection. 

2049 run1 = "run1" 

2050 butler.collections.register(run1) 

2051 run2 = "run2" 

2052 butler.collections.register(run2) 

2053 # put a dataset in each 

2054 metric = makeExampleMetrics() 

2055 dimensions = butler.dimensions.conform(["instrument", "physical_filter"]) 

2056 datasetType = self.addDatasetType( 

2057 "prune_collections_test_dataset", dimensions, storageClass, butler.registry 

2058 ) 

2059 ref1 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run1) 

2060 ref2 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run2) 

2061 uri1 = butler.getURI(ref1) 

2062 uri2 = butler.getURI(ref2) 

2063 

2064 # Put one of the runs in a chain. 

2065 butler.collections.register("Chain", CollectionType.CHAINED) 

2066 butler.collections.extend_chain("Chain", run1) 

2067 

2068 with self.assertRaises(OrphanedRecordError): 

2069 butler.registry.removeDatasetType(datasetType.name) 

2070 

2071 # Remove a non-run. 

2072 with self.assertRaises(TypeError): 

2073 butler.removeRuns(["Chain"]) 

2074 

2075 # Remove without unlinking from chain should fail. 

2076 with self.assertRaises(IntegrityError): 

2077 butler.removeRuns([run1]) 

2078 

2079 # Remove from both runs. No longer use unstore parameter since it 

2080 # always purges. 

2081 butler.removeRuns([run1, run2], unlink_from_chains=True) 

2082 

2083 # Should be nothing in registry for either one, and datastore should 

2084 # not think either exists. 

2085 with self.assertRaises(MissingCollectionError): 

2086 butler.collections.get_info(run1) 

2087 with self.assertRaises(MissingCollectionError): 

2088 butler.collections.get_info(run1) 

2089 self.assertFalse(butler.stored(ref1)) 

2090 self.assertFalse(butler.stored(ref2)) 

2091 # We always unstore so both URIs should be gone. 

2092 self.assertFalse(uri1.exists()) 

2093 self.assertFalse(uri2.exists()) 

2094 

2095 # Now that the collections have been pruned we can remove the 

2096 # dataset type 

2097 butler.registry.removeDatasetType(datasetType.name) 

2098 

2099 with self.assertLogs("lsst.daf.butler.registry", "INFO") as cm: 

2100 butler.registry.removeDatasetType(("test*", "test*")) 

2101 self.assertIn("not defined", "\n".join(cm.output)) 

2102 

2103 def remove_dataset_out_of_band(self, butler: Butler, ref: DatasetRef) -> None: 

2104 """Simulate an external actor removing a file outside of Butler's 

2105 knowledge. 

2106 

2107 Subclasses may override to handle more complicated datastore 

2108 configurations. 

2109 """ 

2110 uri = butler.getURI(ref) 

2111 uri.remove() 

2112 datastore = cast(FileDatastore, butler._datastore) 

2113 datastore.cacheManager.remove_from_cache(ref) 

2114 

2115 def testPruneDatasets(self) -> None: 

2116 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2117 butler = self.create_empty_butler(writeable=True) 

2118 # Load registry data with dimensions to hang datasets off of. 

2119 butler.import_(filename=_get_test_data_path("base.yaml")) 

2120 # Add some RUN-type collections. 

2121 run1 = "run1" 

2122 butler.collections.register(run1) 

2123 run2 = "run2" 

2124 butler.collections.register(run2) 

2125 # put some datasets. ref1 and ref2 have the same data ID, and are in 

2126 # different runs. ref3 has a different data ID. 

2127 metric = makeExampleMetrics() 

2128 dimensions = butler.dimensions.conform(["instrument", "physical_filter"]) 

2129 datasetType = self.addDatasetType( 

2130 "prune_collections_test_dataset", dimensions, storageClass, butler.registry 

2131 ) 

2132 ref1 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run1) 

2133 ref2 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-G"}, run=run2) 

2134 ref3 = butler.put(metric, datasetType, {"instrument": "Cam1", "physical_filter": "Cam1-R1"}, run=run1) 

2135 

2136 many_stored = butler.stored_many([ref1, ref2, ref3]) 

2137 for ref, stored in many_stored.items(): 

2138 self.assertTrue(stored, f"Ref {ref} should be stored") 

2139 

2140 many_exists = butler._exists_many([ref1, ref2, ref3]) 

2141 for ref, exists in many_exists.items(): 

2142 self.assertTrue(exists, f"Checking ref {ref} exists.") 

2143 self.assertEqual(exists, DatasetExistence.VERIFIED, f"Ref {ref} should be stored") 

2144 

2145 # Simple prune. 

2146 butler.pruneDatasets([ref1, ref2, ref3], purge=True, unstore=True) 

2147 self.assertFalse(butler.exists(ref1.datasetType, ref1.dataId, collections=run1)) 

2148 

2149 many_stored = butler.stored_many([ref1, ref2, ref3]) 

2150 for ref, stored in many_stored.items(): 

2151 self.assertFalse(stored, f"Ref {ref} should not be stored") 

2152 

2153 many_exists = butler._exists_many([ref1, ref2, ref3]) 

2154 for ref, exists in many_exists.items(): 

2155 self.assertEqual(exists, DatasetExistence.UNRECOGNIZED, f"Ref {ref} should not be stored") 

2156 

2157 # Put data back. 

2158 ref1_new = butler.put(metric, ref1) 

2159 self.assertEqual(ref1_new, ref1) # Reuses original ID. 

2160 ref2 = butler.put(metric, ref2) 

2161 

2162 many_stored = butler.stored_many([ref1, ref2, ref3]) 

2163 self.assertTrue(many_stored[ref1]) 

2164 self.assertTrue(many_stored[ref2]) 

2165 self.assertFalse(many_stored[ref3]) 

2166 

2167 ref3 = butler.put(metric, ref3) 

2168 

2169 many_exists = butler._exists_many([ref1, ref2, ref3]) 

2170 for ref, exists in many_exists.items(): 

2171 self.assertTrue(exists, f"Ref {ref} should not be stored") 

2172 

2173 # Clear out the datasets from registry and start again. 

2174 refs = [ref1, ref2, ref3] 

2175 butler.pruneDatasets(refs, purge=True, unstore=True) 

2176 for ref in refs: 

2177 butler.put(metric, ref) 

2178 

2179 # Confirm we can retrieve deferred. 

2180 dref1 = butler.getDeferred(ref1) # known and exists 

2181 metric1 = dref1.get() 

2182 self.assertEqual(metric1, metric) 

2183 

2184 # Test different forms of file availability. 

2185 # Need to be in a state where: 

2186 # - one ref just has registry record. 

2187 # - one ref has a missing file but a datastore record. 

2188 # - one ref has a missing datastore record but file is there. 

2189 # - one ref does not exist anywhere. 

2190 # Do not need to test a ref that has everything since that is tested 

2191 # above. 

2192 ref0 = DatasetRef( 

2193 datasetType, 

2194 DataCoordinate.standardize( 

2195 {"instrument": "Cam1", "physical_filter": "Cam1-G"}, universe=butler.dimensions 

2196 ), 

2197 run=run1, 

2198 ) 

2199 

2200 # Delete from datastore and retain in Registry. 

2201 butler.pruneDatasets([ref1], purge=False, unstore=True, disassociate=False) 

2202 

2203 # File has been removed. 

2204 self.remove_dataset_out_of_band(butler, ref2) 

2205 

2206 # Datastore has lost track. 

2207 butler._datastore.forget([ref3]) 

2208 

2209 # First test with a standard butler. 

2210 exists_many = butler._exists_many([ref0, ref1, ref2, ref3], full_check=True) 

2211 self.assertEqual(exists_many[ref0], DatasetExistence.UNRECOGNIZED) 

2212 self.assertEqual(exists_many[ref1], DatasetExistence.RECORDED) 

2213 self.assertEqual(exists_many[ref2], DatasetExistence.RECORDED | DatasetExistence.DATASTORE) 

2214 self.assertEqual(exists_many[ref3], DatasetExistence.RECORDED) 

2215 

2216 exists_many = butler._exists_many([ref0, ref1, ref2, ref3], full_check=False) 

2217 self.assertEqual(exists_many[ref0], DatasetExistence.UNRECOGNIZED) 

2218 self.assertEqual(exists_many[ref1], DatasetExistence.RECORDED | DatasetExistence._ASSUMED) 

2219 self.assertEqual(exists_many[ref2], DatasetExistence.KNOWN) 

2220 self.assertEqual(exists_many[ref3], DatasetExistence.RECORDED | DatasetExistence._ASSUMED) 

2221 self.assertTrue(exists_many[ref2]) 

2222 

2223 # Check that per-ref query gives the same answer as many query. 

2224 for ref, exists in exists_many.items(): 

2225 self.assertEqual(butler.exists(ref, full_check=False), exists) 

2226 

2227 # Get deferred checks for existence before it allows it to be 

2228 # retrieved. 

2229 with self.assertRaises(LookupError): 

2230 butler.getDeferred(ref3) # not known, file exists 

2231 dref2 = butler.getDeferred(ref2) # known but file missing 

2232 with self.assertRaises(FileNotFoundError): 

2233 dref2.get() 

2234 

2235 # Test again with a trusting butler. 

2236 if self.trustModeSupported: 2236 ↛ exitline 2236 didn't return from function 'testPruneDatasets' because the condition on line 2236 was always true

2237 butler._datastore.trustGetRequest = True 

2238 exists_many = butler._exists_many([ref0, ref1, ref2, ref3], full_check=True) 

2239 self.assertEqual(exists_many[ref0], DatasetExistence.UNRECOGNIZED) 

2240 self.assertEqual(exists_many[ref1], DatasetExistence.RECORDED) 

2241 self.assertEqual(exists_many[ref2], DatasetExistence.RECORDED | DatasetExistence.DATASTORE) 

2242 self.assertEqual(exists_many[ref3], DatasetExistence.RECORDED | DatasetExistence._ARTIFACT) 

2243 

2244 # When trusting we can get a deferred dataset handle that is not 

2245 # known but does exist. 

2246 dref3 = butler.getDeferred(ref3) 

2247 metric3 = dref3.get() 

2248 self.assertEqual(metric3, metric) 

2249 

2250 # Check that per-ref query gives the same answer as many query. 

2251 for ref, exists in exists_many.items(): 

2252 self.assertEqual(butler.exists(ref, full_check=True), exists) 

2253 

2254 # Create a ref that surprisingly has the UUID of an existing ref 

2255 # but is not the same. 

2256 ref_bad = DatasetRef(datasetType, dataId=ref3.dataId, run=ref3.run, id=ref2.id) 

2257 with self.assertRaises(ValueError): 

2258 butler.exists(ref_bad) 

2259 

2260 # Create a ref that has a compatible storage class. 

2261 ref_compat = ref2.overrideStorageClass("StructuredDataDict") 

2262 exists = butler.exists(ref_compat) 

2263 self.assertEqual(exists, exists_many[ref2]) 

2264 

2265 # Remove everything and start from scratch. 

2266 butler._datastore.trustGetRequest = False 

2267 butler.pruneDatasets(refs, purge=True, unstore=True) 

2268 for ref in refs: 

2269 butler.put(metric, ref) 

2270 

2271 # These tests mess directly with the trash table and can leave the 

2272 # datastore in an odd state. Do them at the end. 

2273 # Check that in normal mode, deleting the record will lead to 

2274 # trash not touching the file. 

2275 uri1 = butler.getURI(ref1) 

2276 butler._datastore.bridge.moveToTrash( 

2277 [ref1], transaction=None 

2278 ) # Update the dataset_location table 

2279 butler._datastore.forget([ref1]) 

2280 butler._datastore.trash(ref1) 

2281 butler._datastore.emptyTrash() 

2282 self.assertTrue(uri1.exists()) 

2283 uri1.remove() # Clean it up. 

2284 

2285 # Simulate execution butler setup by deleting the datastore 

2286 # record but keeping the file around and trusting. 

2287 butler._datastore.trustGetRequest = True 

2288 uris = butler.get_many_uris([ref2, ref3]) 

2289 uri2 = uris[ref2].primaryURI 

2290 uri3 = uris[ref3].primaryURI 

2291 self.assertTrue(uri2.exists()) 

2292 self.assertTrue(uri3.exists()) 

2293 

2294 # Remove the datastore record. 

2295 butler._datastore.bridge.moveToTrash( 

2296 [ref2], transaction=None 

2297 ) # Update the dataset_location table 

2298 butler._datastore.forget([ref2]) 

2299 self.assertTrue(uri2.exists()) 

2300 butler._datastore.trash([ref2, ref3]) 

2301 # Immediate removal for ref2 file 

2302 self.assertFalse(uri2.exists()) 

2303 # But ref3 has to wait for the empty. 

2304 self.assertTrue(uri3.exists()) 

2305 butler._datastore.emptyTrash() 

2306 self.assertFalse(uri3.exists()) 

2307 

2308 # Clear out the datasets from registry. 

2309 butler.pruneDatasets([ref1, ref2, ref3], purge=True, unstore=True) 

2310 

2311 def test_butler_metrics(self): 

2312 """Test that metrics are collected.""" 

2313 run = "test_run" 

2314 metrics = ButlerMetrics() 

2315 butler, datasetType = self.create_butler( 

2316 run, "MetricsExampleModelProvenance", "prov_metric", metrics=metrics 

2317 ) 

2318 data = MetricsExampleModel( 

2319 summary={"AM1": 5.2, "AM2": 30.6}, 

2320 output={"a": [1, 2, 3], "b": {"blue": 5, "red": "green"}}, 

2321 data=[563, 234, 456.7, 752, 8, 9, 27], 

2322 ) 

2323 

2324 data_ref = butler.put(data, datasetType, visit=424, instrument="DummyCamComp") 

2325 butler.get(data_ref) 

2326 butler.get(data_ref) 

2327 self.assertEqual(metrics.n_get, 2) 

2328 self.assertGreater(metrics.time_in_get, 0.0) 

2329 self.assertEqual(metrics.n_put, 1) 

2330 self.assertGreater(metrics.time_in_put, 0.0) 

2331 

2332 deferred = butler.getDeferred(data_ref) 

2333 deferred.get() 

2334 self.assertEqual(metrics.n_get, 3) 

2335 

2336 with butler.record_metrics() as new: 

2337 data_ref_2 = butler.put(data, datasetType, visit=425, instrument="DummyCamComp") 

2338 butler.get(data_ref) 

2339 

2340 butler.pruneDatasets([data_ref, data_ref_2], purge=True, unstore=True) 

2341 with ResourcePath.temporary_uri(suffix=".json") as tmpFile: 

2342 tmpFile.write(data.model_dump_json().encode()) 

2343 refs = [ 

2344 DatasetRef(datasetType, data_ref_2.dataId, run), 

2345 DatasetRef(datasetType, data_ref.dataId, run), 

2346 ] 

2347 datasets = [FileDataset(path=tmpFile, refs=refs)] 

2348 butler.ingest(*datasets, transfer="copy") 

2349 

2350 self.assertEqual(new.n_get, 1) 

2351 self.assertEqual(new.n_put, 1) 

2352 self.assertEqual(new.n_ingest, 2) 

2353 

2354 

2355class PosixDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2356 """PosixDatastore specialization of a butler""" 

2357 

2358 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2359 fullConfigKey: str | None = ".datastore.formatters" 

2360 validationCanFail = True 

2361 datastoreStr = ["/tmp"] 

2362 datastoreName = [f"FileDatastore@{BUTLER_ROOT_TAG}"] 

2363 registryStr = "/gen3.sqlite3" 

2364 

2365 def testPathConstructor(self) -> None: 

2366 """Independent test of constructor using PathLike.""" 

2367 butler = Butler.from_config(self.tmpConfigFile, run=self.default_run) 

2368 self.enterContext(butler) 

2369 self.assertIsInstance(butler, Butler) 

2370 

2371 # And again with a Path object with the butler yaml 

2372 path = pathlib.Path(self.tmpConfigFile) 

2373 butler = Butler.from_config(path, writeable=False) 

2374 self.enterContext(butler) 

2375 self.assertIsInstance(butler, Butler) 

2376 

2377 # And again with a Path object without the butler yaml 

2378 # (making sure we skip it if the tmp config doesn't end 

2379 # in butler.yaml -- which is the case for a subclass) 

2380 if self.tmpConfigFile.endswith("butler.yaml"): 

2381 path = pathlib.Path(os.path.dirname(self.tmpConfigFile)) 

2382 butler = Butler.from_config(path, writeable=False) 

2383 self.enterContext(butler) 

2384 self.assertIsInstance(butler, Butler) 

2385 

2386 def testExportTransferCopy(self) -> None: 

2387 """Test local export using all transfer modes""" 

2388 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2389 exportButler = self.runPutGetTest(storageClass, "test_metric") 

2390 # Test that the repo actually has at least one dataset. 

2391 datasets = list(exportButler.registry.queryDatasets(..., collections=...)) 

2392 self.assertGreater(len(datasets), 0) 

2393 uris = [exportButler.getURI(d) for d in datasets] 

2394 assert isinstance(exportButler._datastore, FileDatastore) 

2395 datastoreRoot = exportButler.get_datastore_roots()[exportButler.get_datastore_names()[0]] 

2396 

2397 pathsInStore = [uri.relative_to(datastoreRoot) for uri in uris] 

2398 

2399 for path in pathsInStore: 

2400 # Assume local file system 

2401 assert path is not None 

2402 self.assertTrue(self.checkFileExists(datastoreRoot, path), f"Checking path {path}") 

2403 

2404 for transfer in ("copy", "link", "symlink", "relsymlink"): 

2405 with safeTestTempDir(TESTDIR) as exportDir: 

2406 with exportButler.export(directory=exportDir, format="yaml", transfer=transfer) as export: 

2407 export.saveDatasets(datasets) 

2408 for path in pathsInStore: 

2409 assert path is not None 

2410 self.assertTrue( 

2411 self.checkFileExists(exportDir, path), 

2412 f"Check that mode {transfer} exported files", 

2413 ) 

2414 

2415 def testPytypeCoercion(self) -> None: 

2416 """Test python type coercion on Butler.get and put.""" 

2417 # Store some data with the normal example storage class. 

2418 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2419 datasetTypeName = "test_metric" 

2420 butler = self.runPutGetTest(storageClass, datasetTypeName) 

2421 

2422 dataId = {"instrument": "DummyCamComp", "visit": 423} 

2423 metric = butler.get(datasetTypeName, dataId=dataId) 

2424 self.assertEqual(get_full_type_name(metric), "lsst.daf.butler.tests.MetricsExample") 

2425 

2426 datasetType_ori = butler.get_dataset_type(datasetTypeName) 

2427 self.assertEqual(datasetType_ori.storageClass.name, "StructuredDataNoComponents") 

2428 

2429 # Now need to hack the registry dataset type definition. 

2430 # There is no API for this. 

2431 assert isinstance(butler._registry, SqlRegistry) 

2432 manager = butler._registry._managers.datasets 

2433 assert hasattr(manager, "_db") and hasattr(manager, "_static") 

2434 manager._db.update( 

2435 manager._static.dataset_type, 

2436 {"name": datasetTypeName}, 

2437 {datasetTypeName: datasetTypeName, "storage_class": "StructuredDataNoComponentsModel"}, 

2438 ) 

2439 

2440 # Force reset of dataset type cache 

2441 butler.registry.refresh() 

2442 

2443 datasetType_new = butler.get_dataset_type(datasetTypeName) 

2444 self.assertEqual(datasetType_new.name, datasetType_ori.name) 

2445 self.assertEqual(datasetType_new.storageClass.name, "StructuredDataNoComponentsModel") 

2446 

2447 metric_model = butler.get(datasetTypeName, dataId=dataId) 

2448 self.assertNotEqual(type(metric_model), type(metric)) 

2449 self.assertEqual(get_full_type_name(metric_model), "lsst.daf.butler.tests.MetricsExampleModel") 

2450 

2451 # Put the model and read it back to show that everything now 

2452 # works as normal. 

2453 metric_ref = butler.put(metric_model, datasetTypeName, dataId=dataId, visit=424) 

2454 metric_model_new = butler.get(metric_ref) 

2455 self.assertEqual(metric_model_new, metric_model) 

2456 

2457 # Hack the storage class again to something that will fail on the 

2458 # get with no conversion class. 

2459 manager._db.update( 

2460 manager._static.dataset_type, 

2461 {"name": datasetTypeName}, 

2462 {datasetTypeName: datasetTypeName, "storage_class": "StructuredDataListYaml"}, 

2463 ) 

2464 butler.registry.refresh() 

2465 

2466 with self.assertRaises(ValueError): 

2467 butler.get(datasetTypeName, dataId=dataId) 

2468 

2469 def test_provenance(self): 

2470 """Test that provenance is attached on put.""" 

2471 run = "test_run" 

2472 butler, datasetType = self.create_butler(run, "MetricsExampleModelProvenance", "prov_metric") 

2473 metric = MetricsExampleModel( 

2474 summary={"AM1": 5.2, "AM2": 30.6}, 

2475 output={"a": [1, 2, 3], "b": {"blue": 5, "red": "green"}}, 

2476 data=[563, 234, 456.7, 752, 8, 9, 27], 

2477 ) 

2478 # Provenance can be attached to the object being put. Whether 

2479 # it is or not is dependent on the formatter. For this test we 

2480 # copy on adding provenance to ensure they differ. 

2481 self.assertIsNone(metric.dataset_id) 

2482 metric_ref = butler.put(metric, datasetType, visit=424, instrument="DummyCamComp") 

2483 self.assertIsNone(metric.dataset_id) 

2484 metric_2 = butler.get(metric_ref) 

2485 self.assertEqual(metric_2.data, metric.data) 

2486 self.assertEqual(metric_2.dataset_id, metric_ref.id) 

2487 self.assertIsNone(metric_2.provenance) 

2488 

2489 # Put with provenance. 

2490 prov = DatasetProvenance(quantum_id=uuid.uuid4()) 

2491 prov.add_input(metric_ref) 

2492 prov.add_extra_provenance(metric_ref.id, {"answer": 42}) 

2493 metric_ref2 = butler.put(metric, datasetType, visit=423, instrument="DummyCamComp", provenance=prov) 

2494 metric_3 = butler.get(metric_ref2) 

2495 self.assertEqual(metric_3.provenance, prov) 

2496 

2497 # Check that we can extract provenance from dict form. 

2498 prov_dict = prov.to_flat_dict(metric_ref2) 

2499 prov_from_prov, ref_from_prov = DatasetProvenance.from_flat_dict(prov_dict, butler) 

2500 self.assertEqual(ref_from_prov, metric_ref2) 

2501 # Direct __eq__ of the provenance does not work because one side 

2502 # includes dimension records. 

2503 self.assertEqual({ref.id for ref in prov_from_prov.inputs}, {ref.id for ref in prov.inputs}) 

2504 self.assertEqual(prov_from_prov.quantum_id, prov.quantum_id) 

2505 self.assertEqual(prov_from_prov.extras, prov.extras) 

2506 

2507 # Force a bad ID into the dict. 

2508 prov_dict["id"] = uuid.uuid4() 

2509 with self.assertRaises(ValueError): 

2510 DatasetProvenance.from_flat_dict(prov_dict, butler) 

2511 del prov_dict["id"] 

2512 prov_dict["input 0 id"] = uuid.uuid4() 

2513 with self.assertRaises(ValueError): 

2514 DatasetProvenance.from_flat_dict(prov_dict, butler) 

2515 

2516 # Check that simple types can be reconstructed with non-standard 

2517 # separators. 

2518 prov_dict = prov.to_flat_dict(metric_ref2, prefix="XYZ", sep="😎", simple_types=True) 

2519 prov_from_prov, ref_from_prov = DatasetProvenance.from_flat_dict(prov_dict, butler) 

2520 self.assertEqual(ref_from_prov, metric_ref2) 

2521 self.assertEqual({ref.id for ref in prov_from_prov.inputs}, {ref.id for ref in prov.inputs}) 

2522 

2523 with self.assertRaises(ValueError): 

2524 DatasetProvenance.from_flat_dict({"unknown": 42}, butler) 

2525 

2526 def test_specialized_file_datasets_functions(self): 

2527 """Test a workflow used in Prompt Processing where we export datasets 

2528 from one repository and write them in-place to the datastore of 

2529 another, without immediately inserting registry entries for the 

2530 datasets. 

2531 """ 

2532 repo = MetricTestRepo.create_from_butler( 

2533 self.create_empty_butler(writeable=True), 

2534 self.tmpConfigFile, 

2535 "StructuredCompositeReadCompNoDisassembly", 

2536 ) 

2537 source_butler = repo.butler 

2538 

2539 # Test writing outputs to a FileDatastore. 

2540 with tempfile.TemporaryDirectory() as tempdir: 

2541 target_repo_config = Butler.makeRepo(tempdir) 

2542 refs = [repo.ref1, repo.ref2] 

2543 datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), refs) 

2544 self.assertEqual(len(datasets), 2) 

2545 self.assertEqual({ref.id for ref in refs}, {dataset.refs[0].id for dataset in datasets}) 

2546 for dataset in datasets: 

2547 path = ResourcePath(dataset.path, forceAbsolute=False) 

2548 # Paths should be relative paths to the target datastore. 

2549 self.assertFalse(path.isabs()) 

2550 # Files should have been copied into the target datastore 

2551 self.assertTrue(ResourcePath(tempdir).join(path).exists()) 

2552 

2553 # Make sure the target Butler can ingest the datasets. 

2554 target_butler = Butler(target_repo_config, writeable=True) 

2555 self.enterContext(target_butler) 

2556 target_butler.transfer_dimension_records_from(source_butler, refs) 

2557 target_butler.ingest(*datasets, transfer=None) 

2558 self.assertIsNotNone(target_butler.get(repo.ref1)) 

2559 self.assertIsNotNone(target_butler.get(repo.ref2)) 

2560 

2561 # Giving an empty list of files is a no-op. 

2562 no_datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), []) 

2563 self.assertEqual(len(no_datasets), 0) 

2564 

2565 # Test writing outputs to a ChainedDatastore. 

2566 with tempfile.TemporaryDirectory() as tempdir: 

2567 # Set up a second dataset type, so we can split the files across 

2568 # multiple datastore roots. 

2569 dt1 = repo.datasetType 

2570 dt2 = DatasetType("other", dt1.dimensions, dt1.storageClass) 

2571 source_butler.registry.registerDatasetType(dt2) 

2572 other_ref = repo.addDataset(repo.ref1.dataId, datasetType=dt2) 

2573 config = Config.fromString( 

2574 f""" 

2575 datastore: 

2576 cls: lsst.daf.butler.datastores.chainedDatastore.ChainedDatastore 

2577 datastore_constraints: 

2578 - constraints: 

2579 accept: 

2580 - {dt1.name} 

2581 - constraints: 

2582 accept: 

2583 - {dt2.name} 

2584 datastores: 

2585 - datastore: 

2586 cls: lsst.daf.butler.datastores.fileDatastore.FileDatastore 

2587 root: <butlerRoot>/FileDatastore_0 

2588 - datastore: 

2589 cls: lsst.daf.butler.datastores.fileDatastore.FileDatastore 

2590 root: <butlerRoot>/FileDatastore_1 

2591 """ 

2592 ) 

2593 target_repo_config = Butler.makeRepo(tempdir, config) 

2594 refs = [repo.ref1, repo.ref2, other_ref] 

2595 datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), refs) 

2596 self.assertEqual(len(datasets), 3) 

2597 self.assertEqual({ref.id for ref in refs}, {dataset.refs[0].id for dataset in datasets}) 

2598 for dataset in datasets: 

2599 path = ResourcePath(dataset.path, forceAbsolute=False) 

2600 # Paths should be relative paths to the target datastore. 

2601 self.assertFalse(path.isabs()) 

2602 # Files should have been split up between the two datastores 

2603 # in the chain. 

2604 datastore_root = ResourcePath(tempdir) 

2605 if dataset.refs[0].datasetType.name == dt1.name: 

2606 datastore_root = datastore_root.join("FileDatastore_0") 

2607 else: 

2608 datastore_root = datastore_root.join("FileDatastore_1") 

2609 self.assertTrue(datastore_root.join(path).exists()) 

2610 

2611 # Make sure the target Butler can ingest the datasets. 

2612 target_butler = Butler(target_repo_config, writeable=True) 

2613 self.enterContext(target_butler) 

2614 target_butler.transfer_dimension_records_from(source_butler, refs) 

2615 target_butler.ingest(*datasets, transfer=None) 

2616 self.assertIsNotNone(target_butler.get(repo.ref1)) 

2617 self.assertIsNotNone(target_butler.get(repo.ref2)) 

2618 self.assertIsNotNone(target_butler.get(other_ref)) 

2619 

2620 def test_temporary_for_ingest(self) -> None: 

2621 """Test the `lsst.daf.butler._rubin.ingest_from_temporary` module.""" 

2622 with self.create_empty_butler("example_run") as butler: 

2623 dataset_type = DatasetType("example", butler.dimensions.empty, "StructuredDataDict") 

2624 butler.registry.registerDatasetType(dataset_type) 

2625 ref = DatasetRef(dataset_type, DataCoordinate.make_empty(butler.dimensions), "example_run") 

2626 with TemporaryForIngest(butler, ref) as temporary: 

2627 temporary.path.write(b"three: 3") 

2628 found = TemporaryForIngest.find_orphaned_temporaries_by_ref(ref, butler) 

2629 self.assertEqual(found, [temporary.path]) 

2630 self.assertIn(".tmp", temporary.ospath) 

2631 temporary.ingest() 

2632 loaded = butler.get(ref) 

2633 self.assertEqual(loaded, {"three": 3}) 

2634 

2635 

2636class PostgresPosixDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2637 """PosixDatastore specialization of a butler using Postgres""" 

2638 

2639 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2640 fullConfigKey = ".datastore.formatters" 

2641 validationCanFail = True 

2642 datastoreStr = ["/tmp"] 

2643 datastoreName = [f"FileDatastore@{BUTLER_ROOT_TAG}"] 

2644 registryStr = "PostgreSQL@test" 

2645 

2646 @classmethod 

2647 def setUpClass(cls) -> None: 

2648 cls.postgresql = cls.enterClassContext(setup_postgres_test_db()) 

2649 super().setUpClass() 

2650 

2651 def setUp(self) -> None: 

2652 # Need to add a registry section to the config. 

2653 self._temp_config = False 

2654 config = Config(self.configFile) 

2655 self.postgresql.patch_butler_config(config) 

2656 with tempfile.NamedTemporaryFile("w", suffix=".yaml", delete=False) as fh: 

2657 config.dump(fh) 

2658 self.configFile = fh.name 

2659 self._temp_config = True 

2660 super().setUp() 

2661 

2662 def tearDown(self) -> None: 

2663 if self._temp_config and os.path.exists(self.configFile): 

2664 os.remove(self.configFile) 

2665 super().tearDown() 

2666 

2667 def testMakeRepo(self) -> None: 

2668 # The base class test assumes that it's using sqlite and assumes 

2669 # the config file is acceptable to sqlite. 

2670 raise unittest.SkipTest("Postgres config is not compatible with this test.") 

2671 

2672 

2673class ClonedPostgresPosixDatastoreButlerTestCase(PostgresPosixDatastoreButlerTestCase, unittest.TestCase): 

2674 """Test that Butler with a Postgres registry still works after cloning.""" 

2675 

2676 def create_butler( 

2677 self, 

2678 run: str, 

2679 storageClass: StorageClass | str, 

2680 datasetTypeName: str, 

2681 metrics: ButlerMetrics | None = None, 

2682 ) -> tuple[DirectButler, DatasetType]: 

2683 butler, datasetType = super().create_butler(run, storageClass, datasetTypeName, metrics=metrics) 

2684 return butler.clone(run=run, metrics=metrics), datasetType 

2685 

2686 

2687class InMemoryDatastoreButlerTestCase(ButlerTests, unittest.TestCase): 

2688 """InMemoryDatastore specialization of a butler""" 

2689 

2690 configFile = os.path.join(TESTDIR, "config/basic/butler-inmemory.yaml") 

2691 fullConfigKey = None 

2692 useTempRoot = False 

2693 validationCanFail = False 

2694 datastoreStr = ["datastore='InMemory"] 

2695 datastoreName = ["InMemoryDatastore@"] 

2696 registryStr = "/gen3.sqlite3" 

2697 

2698 def testIngest(self) -> None: 

2699 pass 

2700 

2701 def test_ingest_zip(self) -> None: 

2702 pass 

2703 

2704 

2705class ClonedSqliteButlerTestCase(InMemoryDatastoreButlerTestCase, unittest.TestCase): 

2706 """Test that a Butler with a Sqlite registry still works after cloning.""" 

2707 

2708 def create_butler( 

2709 self, 

2710 run: str, 

2711 storageClass: StorageClass | str, 

2712 datasetTypeName: str, 

2713 metrics: ButlerMetrics | None = None, 

2714 ) -> tuple[DirectButler, DatasetType]: 

2715 butler, datasetType = super().create_butler(run, storageClass, datasetTypeName, metrics=metrics) 

2716 return butler.clone(run=run), datasetType 

2717 

2718 

2719class ChainedDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2720 """PosixDatastore specialization""" 

2721 

2722 configFile = os.path.join(TESTDIR, "config/basic/butler-chained.yaml") 

2723 fullConfigKey = ".datastore.datastores.1.formatters" 

2724 validationCanFail = True 

2725 datastoreStr = ["datastore='InMemory", "/FileDatastore_1/,", "/FileDatastore_2/'"] 

2726 datastoreName = [ 

2727 "InMemoryDatastore@", 

2728 f"FileDatastore@{BUTLER_ROOT_TAG}/FileDatastore_1", 

2729 "SecondDatastore", 

2730 ] 

2731 registryStr = "/gen3.sqlite3" 

2732 

2733 def testPruneDatasets(self) -> None: 

2734 # This test relies on manipulating files out-of-band, which is 

2735 # impossible for this configuration because of the InMemoryDatastore in 

2736 # the ChainedDatastore. 

2737 pass 

2738 

2739 def testComponentFromOverriddenStorageClassWarns(self) -> None: 

2740 # The InMemoryDatastore in the ChainedDatastore satisfies the get, so 

2741 # the FileDatastore warning about having to read the whole dataset to 

2742 # extract the component is never issued. 

2743 pass 

2744 

2745 

2746class ButlerExplicitRootTestCase(PosixDatastoreButlerTestCase): 

2747 """Test that a yaml file in one location can refer to a root in another.""" 

2748 

2749 datastoreStr = ["dir1"] 

2750 # Disable the makeRepo test since we are deliberately not using 

2751 # butler.yaml as the config name. 

2752 fullConfigKey = None 

2753 

2754 def setUp(self) -> None: 

2755 self.root = makeTestTempDir(TESTDIR) 

2756 

2757 # Make a new repository in one place 

2758 self.dir1 = os.path.join(self.root, "dir1") 

2759 Butler.makeRepo(self.dir1, config=Config(self.configFile)) 

2760 

2761 # Move the yaml file to a different place and add a "root" 

2762 self.dir2 = os.path.join(self.root, "dir2") 

2763 os.makedirs(self.dir2, exist_ok=True) 

2764 configFile1 = os.path.join(self.dir1, "butler.yaml") 

2765 config = Config(configFile1) 

2766 config["root"] = self.dir1 

2767 configFile2 = os.path.join(self.dir2, "butler2.yaml") 

2768 config.dumpToUri(configFile2) 

2769 os.remove(configFile1) 

2770 self.tmpConfigFile = configFile2 

2771 

2772 def testFileLocations(self) -> None: 

2773 self.assertNotEqual(self.dir1, self.dir2) 

2774 self.assertTrue(os.path.exists(os.path.join(self.dir2, "butler2.yaml"))) 

2775 self.assertFalse(os.path.exists(os.path.join(self.dir1, "butler.yaml"))) 

2776 self.assertTrue(os.path.exists(os.path.join(self.dir1, "gen3.sqlite3"))) 

2777 

2778 

2779class ButlerMakeRepoOutfileTestCase(ButlerPutGetTests, unittest.TestCase): 

2780 """Test that a config file created by makeRepo outside of repo works.""" 

2781 

2782 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2783 

2784 def setUp(self) -> None: 

2785 self.root = makeTestTempDir(TESTDIR) 

2786 self.root2 = makeTestTempDir(TESTDIR) 

2787 

2788 self.tmpConfigFile = os.path.join(self.root2, "different.yaml") 

2789 Butler.makeRepo(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile) 

2790 

2791 def tearDown(self) -> None: 

2792 if os.path.exists(self.root2): 2792 ↛ 2794line 2792 didn't jump to line 2794 because the condition on line 2792 was always true

2793 shutil.rmtree(self.root2, ignore_errors=True) 

2794 super().tearDown() 

2795 

2796 def testConfigExistence(self) -> None: 

2797 c = Config(self.tmpConfigFile) 

2798 uri_config = ResourcePath(c["root"]) 

2799 uri_expected = ResourcePath(self.root, forceDirectory=True) 

2800 self.assertEqual(uri_config.geturl(), uri_expected.geturl()) 

2801 self.assertNotIn(":", uri_config.path, "Check for URI concatenated with normal path") 

2802 

2803 def testPutGet(self) -> None: 

2804 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2805 self.runPutGetTest(storageClass, "test_metric") 

2806 

2807 

2808class ButlerMakeRepoOutfileDirTestCase(ButlerMakeRepoOutfileTestCase): 

2809 """Test that a config file created by makeRepo outside of repo works.""" 

2810 

2811 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2812 

2813 def setUp(self) -> None: 

2814 self.root = makeTestTempDir(TESTDIR) 

2815 self.root2 = makeTestTempDir(TESTDIR) 

2816 

2817 self.tmpConfigFile = self.root2 

2818 Butler.makeRepo(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile) 

2819 

2820 def testConfigExistence(self) -> None: 

2821 # Append the yaml file else Config constructor does not know the file 

2822 # type. 

2823 self.tmpConfigFile = os.path.join(self.tmpConfigFile, "butler.yaml") 

2824 super().testConfigExistence() 

2825 

2826 

2827class ButlerMakeRepoOutfileUriTestCase(ButlerMakeRepoOutfileTestCase): 

2828 """Test that a config file created by makeRepo outside of repo works.""" 

2829 

2830 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2831 

2832 def setUp(self) -> None: 

2833 self.root = makeTestTempDir(TESTDIR) 

2834 self.root2 = makeTestTempDir(TESTDIR) 

2835 

2836 self.tmpConfigFile = ResourcePath(os.path.join(self.root2, "something.yaml")).geturl() 

2837 Butler.makeRepo(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile) 

2838 

2839 

2840class RemoteTestDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2841 """Specialization of a butler using a datastore root that reports itself 

2842 as not local; a remote file datastore + a local SqlRegistry. 

2843 """ 

2844 

2845 configFile = os.path.join(TESTDIR, "config/basic/butler-remotetest-store.yaml") 

2846 fullConfigKey = None 

2847 validationCanFail = True 

2848 

2849 registryStr = "/gen3.sqlite3" 

2850 """Expected format of the Registry string.""" 

2851 

2852 def setUp(self) -> None: 

2853 config = Config(self.configFile) 

2854 

2855 self.root = makeTestTempDir(TESTDIR) 

2856 # The space in the directory name is deliberate. It ensures the URI 

2857 # has to be percent-encoded correctly on the way in and decoded on 

2858 # the way out. 

2859 root_path = os.path.join(self.root, "butler root") 

2860 os.makedirs(root_path) 

2861 rooturi = make_remote_test_uri(root_path) 

2862 config.update({"datastore": {"datastore": {"root": str(rooturi)}}}) 

2863 

2864 # The registry database has to live on a real local file system. 

2865 self.reg_dir = makeTestTempDir(TESTDIR) 

2866 config["registry", "db"] = f"sqlite:///{self.reg_dir}/gen3.sqlite3" 

2867 

2868 self.datastoreStr = [f"datastore='{rooturi}'"] 

2869 self.datastoreName = [f"FileDatastore@{rooturi}"] 

2870 Butler.makeRepo(rooturi, config=config, forceConfigRoot=False) 

2871 self.tmpConfigFile = str(rooturi.join("butler.yaml", forceDirectory=False)) 

2872 

2873 def tearDown(self) -> None: 

2874 removeTestTempDir(self.reg_dir) 

2875 # The base class removes self.root, which contains the datastore. 

2876 super().tearDown() 

2877 

2878 

2879class DatastoreTransfers(TestCaseMixin): 

2880 """Base test setup for data transfers between butlers. The concrete tests 

2881 for specific configurations are in other classes, below. 

2882 """ 

2883 

2884 storageClassFactory: StorageClassFactory 

2885 

2886 @classmethod 

2887 def setUpClass(cls) -> None: 

2888 cls.storageClassFactory = StorageClassFactory() 

2889 

2890 def setUp(self) -> None: 

2891 self.root = makeTestTempDir(TESTDIR) 

2892 self.config = Config(self.configFile) 

2893 

2894 # Some tests cause convertors to be replaced so ensure 

2895 # the storage class factory is reset each time. 

2896 self.storageClassFactory.reset() 

2897 self.storageClassFactory.addFromConfig(self.configFile) 

2898 

2899 def tearDown(self) -> None: 

2900 removeTestTempDir(self.root) 

2901 

2902 def create_butler(self, manager: str | None, label: str, config_file: str | None = None) -> Butler: 

2903 if manager is None: 2903 ↛ 2907line 2903 didn't jump to line 2907 because the condition on line 2903 was always true

2904 manager = ( 

2905 "lsst.daf.butler.registry.datasets.byDimensions.ByDimensionsDatasetRecordStorageManagerUUID" 

2906 ) 

2907 config = Config(config_file if config_file is not None else self.configFile) 

2908 config["registry", "managers", "datasets"] = manager 

2909 butler = Butler.from_config( 

2910 Butler.makeRepo(f"{self.root}/butler{label}", config=config), writeable=True 

2911 ) 

2912 self.enterContext(butler) 

2913 return butler 

2914 

2915 def assertButlerTransfers( 

2916 self, 

2917 purge: bool = False, 

2918 storageClassName: str = "StructuredData", 

2919 storageClassNameTarget: str | None = None, 

2920 ) -> None: 

2921 """Test that a run can be transferred to another butler.""" 

2922 storageClass = self.storageClassFactory.getStorageClass(storageClassName) 

2923 if storageClassNameTarget is not None: 

2924 storageClassTarget = self.storageClassFactory.getStorageClass(storageClassNameTarget) 

2925 else: 

2926 storageClassTarget = storageClass 

2927 

2928 datasetTypeName = "random_data" 

2929 

2930 # Test will create 3 collections and we will want to transfer 

2931 # two of those three. 

2932 runs = ["run1", "run2", "other"] 

2933 

2934 # Also want to use two different dataset types to ensure that 

2935 # grouping works. 

2936 datasetTypeNames = ["random_data", "random_data_2"] 

2937 

2938 # Create the run collections in the source butler. 

2939 for run in runs: 

2940 self.source_butler.collections.register(run) 

2941 

2942 # Create dimensions in source butler. 

2943 n_exposures = 30 

2944 self.source_butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"}) 

2945 self.source_butler.registry.insertDimensionData( 

2946 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"} 

2947 ) 

2948 self.source_butler.registry.insertDimensionData( 

2949 "detector", {"instrument": "DummyCamComp", "id": 1, "full_name": "det1"} 

2950 ) 

2951 self.source_butler.registry.insertDimensionData( 

2952 "day_obs", 

2953 { 

2954 "instrument": "DummyCamComp", 

2955 "id": 20250101, 

2956 }, 

2957 ) 

2958 

2959 for i in range(n_exposures): 

2960 self.source_butler.registry.insertDimensionData( 

2961 "group", {"instrument": "DummyCamComp", "name": f"group{i}"} 

2962 ) 

2963 self.source_butler.registry.insertDimensionData( 

2964 "exposure", 

2965 { 

2966 "instrument": "DummyCamComp", 

2967 "id": i, 

2968 "obs_id": f"exp{i}", 

2969 "physical_filter": "d-r", 

2970 "group": f"group{i}", 

2971 "day_obs": 20250101, 

2972 }, 

2973 ) 

2974 

2975 # Create dataset types in the source butler. 

2976 dimensions = self.source_butler.dimensions.conform(["instrument", "exposure"]) 

2977 for datasetTypeName in datasetTypeNames: 

2978 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

2979 self.source_butler.registry.registerDatasetType(datasetType) 

2980 

2981 # Write a dataset to an unrelated run -- this will ensure that 

2982 # we are rewriting integer dataset ids in the target if necessary. 

2983 # Will not be relevant for UUID. 

2984 run = "distraction" 

2985 butler = Butler.from_config(butler=self.source_butler, run=run) 

2986 self.enterContext(butler) 

2987 butler.put( 

2988 makeExampleMetrics(), 

2989 datasetTypeName, 

2990 exposure=1, 

2991 instrument="DummyCamComp", 

2992 physical_filter="d-r", 

2993 ) 

2994 

2995 # Write some example metrics to the source 

2996 butler = Butler.from_config(butler=self.source_butler) 

2997 self.enterContext(butler) 

2998 

2999 # Set of DatasetRefs that should be in the list of refs to transfer 

3000 # but which will not be transferred. 

3001 deleted: set[DatasetRef] = set() 

3002 

3003 n_expected = 20 # Number of datasets expected to be transferred 

3004 source_refs = [] 

3005 for i in range(n_exposures): 

3006 # Put a third of datasets into each collection, only retain 

3007 # two thirds. 

3008 index = i % 3 

3009 run = runs[index] 

3010 datasetTypeName = datasetTypeNames[i % 2] 

3011 

3012 metric = MetricsExample( 

3013 summary={"counter": i}, output={"text": "metric"}, data=[2 * x for x in range(i)] 

3014 ) 

3015 dataId = {"exposure": i, "instrument": "DummyCamComp", "physical_filter": "d-r"} 

3016 ref = butler.put(metric, datasetTypeName, dataId=dataId, run=run) 

3017 

3018 # Remove the datastore record using low-level API, but only 

3019 # for a specific index. 

3020 if purge and index == 1: 

3021 # For one of these delete the file as well. 

3022 # This allows the "missing" code to filter the 

3023 # file out. 

3024 # Access the individual datastores. 

3025 datastores = [] 

3026 if hasattr(butler._datastore, "datastores"): 

3027 datastores.extend(butler._datastore.datastores) 

3028 else: 

3029 datastores.append(butler._datastore) 

3030 

3031 if not deleted: 

3032 # For a chained datastore we need to remove 

3033 # files in each chain. 

3034 for datastore in datastores: 

3035 # The file might not be known to the datastore 

3036 # if constraints are used. 

3037 try: 

3038 primary, uris = datastore.getURIs(ref) 

3039 except FileNotFoundError: 

3040 continue 

3041 if primary and primary.scheme != "mem": 

3042 primary.remove() 

3043 for uri in uris.values(): 

3044 if uri.scheme != "mem": 3044 ↛ 3043line 3044 didn't jump to line 3043 because the condition on line 3044 was always true

3045 uri.remove() 

3046 n_expected -= 1 

3047 deleted.add(ref) 

3048 

3049 # Remove the datastore record. 

3050 for datastore in datastores: 

3051 if hasattr(datastore, "removeStoredItemInfo"): 3051 ↛ 3050line 3051 didn't jump to line 3050 because the condition on line 3051 was always true

3052 datastore.removeStoredItemInfo(ref) 

3053 

3054 if index < 2: 

3055 source_refs.append(ref) 

3056 if ref not in deleted: 

3057 new_metric = butler.get(ref) 

3058 self.assertEqual(new_metric, metric) 

3059 

3060 # Create some bad dataset types to ensure we check for inconsistent 

3061 # definitions. 

3062 badStorageClass = self.storageClassFactory.getStorageClass("StructuredDataList") 

3063 for datasetTypeName in datasetTypeNames: 

3064 datasetType = DatasetType(datasetTypeName, dimensions, badStorageClass) 

3065 self.target_butler.registry.registerDatasetType(datasetType) 

3066 with self.assertRaises(ConflictingDefinitionError) as cm: 

3067 self.target_butler.transfer_from(self.source_butler, source_refs) 

3068 self.assertIn("dataset type differs", str(cm.exception)) 

3069 

3070 # And remove the bad definitions. 

3071 for datasetTypeName in datasetTypeNames: 

3072 self.target_butler.registry.removeDatasetType(datasetTypeName) 

3073 

3074 # Transfer without creating dataset types should fail. 

3075 with self.assertRaises(KeyError): 

3076 self.target_butler.transfer_from(self.source_butler, source_refs) 

3077 

3078 # Transfer without creating dimensions should fail. 

3079 with self.assertRaises(ConflictingDefinitionError) as cm: 

3080 self.target_butler.transfer_from(self.source_butler, source_refs, register_dataset_types=True) 

3081 self.assertIn("dimension", str(cm.exception)) 

3082 

3083 # The dry run test requires dataset types to exist. If we have 

3084 # been given distinct storage classes for the target we have 

3085 # to redefine at least one of the dataset types in the target butler. 

3086 if storageClass != storageClassTarget: 

3087 self.target_butler.registry.removeDatasetType(datasetTypeNames[0]) 

3088 datasetType = DatasetType(datasetTypeNames[0], dimensions, storageClassTarget) 

3089 self.target_butler.registry.registerDatasetType(datasetType) 

3090 

3091 # The failed transfer above leaves registry in an inconsistent 

3092 # state because the run is created but then rolled back without 

3093 # the collection cache being cleared. For now force a refresh. 

3094 # Can remove with DM-35498. 

3095 self.target_butler.registry.refresh() 

3096 

3097 # Do a dry run -- this should not have any effect on the target butler. 

3098 self.target_butler.transfer_from(self.source_butler, source_refs, dry_run=True) 

3099 

3100 # Transfer the records for one ref to test the alternative API. 

3101 with self.assertLogs(logger="lsst", level=logging.DEBUG) as log_cm: 

3102 self.target_butler.transfer_dimension_records_from(self.source_butler, [source_refs[0]]) 

3103 self.assertIn("number of records transferred: 1", ";".join(log_cm.output)) 

3104 

3105 # Now transfer them to the second butler, including dimensions. 

3106 with self.assertLogs(logger="lsst", level=logging.DEBUG) as log_cm: 

3107 transferred = self.target_butler.transfer_from( 

3108 self.source_butler, 

3109 source_refs, 

3110 register_dataset_types=True, 

3111 transfer_dimensions=True, 

3112 ) 

3113 self.assertEqual(len(transferred), n_expected) 

3114 log_output = ";".join(log_cm.output) 

3115 

3116 # A ChainedDatastore will use the in-memory datastore for mexists 

3117 # so we can not rely on the mexists log message. 

3118 self.assertIn("Number of datastore records found in source", log_output) 

3119 self.assertIn("Creating output run", log_output) 

3120 

3121 # Do the transfer twice to ensure that it will do nothing extra. 

3122 # Only do this if purge=True because it does not work for int 

3123 # dataset_id. 

3124 if purge: 

3125 # This should not need to register dataset types. 

3126 transferred = self.target_butler.transfer_from(self.source_butler, source_refs) 

3127 self.assertEqual(len(transferred), n_expected) 

3128 

3129 with self.assertRaises((TypeError, AttributeError)): 

3130 self.target_butler._datastore.transfer_from(self.source_butler, source_refs) # type: ignore 

3131 

3132 with self.assertRaises(ValueError): 

3133 self.target_butler._datastore.transfer_from( 

3134 self.source_butler._datastore, source_refs, transfer="split" 

3135 ) 

3136 

3137 # Now try to get the same refs from the new butler. 

3138 for ref in source_refs: 

3139 if ref not in deleted: 

3140 new_metric = self.target_butler.get(ref) 

3141 old_metric = self.source_butler.get(ref) 

3142 self.assertEqual(new_metric, old_metric) 

3143 

3144 # Try again without implicit storage class conversion 

3145 # triggered by using the source ref. This will do conversion 

3146 # since the formatter will be returning the source python type. 

3147 target_ref = self.target_butler.get_dataset(ref.id) 

3148 if target_ref.datasetType.storageClass != ref.datasetType.storageClass: 

3149 new_metric = self.target_butler.get(target_ref) 

3150 self.assertNotEqual(type(new_metric), type(old_metric)) 

3151 

3152 # Remove the dataset from the target and put it again 

3153 # as if it was the right type all along for this butler. 

3154 self.target_butler.pruneDatasets( 

3155 [target_ref], unstore=True, purge=True, disassociate=True 

3156 ) 

3157 self.target_butler.put(new_metric, target_ref) 

3158 new_new_metric = self.target_butler.get(target_ref) 

3159 new_old_metric = self.target_butler.get( 

3160 target_ref, storageClass=ref.datasetType.storageClass 

3161 ) 

3162 self.assertEqual(new_new_metric, new_metric) 

3163 self.assertEqual(new_old_metric, old_metric) 

3164 

3165 # Now prune run2 collection and create instead a CHAINED collection. 

3166 # This should block the transfer. 

3167 self.target_butler.removeRuns(["run2"]) 

3168 self.target_butler.collections.register("run2", CollectionType.CHAINED) 

3169 with self.assertRaises(CollectionTypeError): 

3170 # Re-importing the run1 datasets can be problematic if they 

3171 # use integer IDs so filter those out. 

3172 to_transfer = [ref for ref in source_refs if ref.run == "run2"] 

3173 self.target_butler.transfer_from(self.source_butler, to_transfer) 

3174 

3175 

3176class PosixDatastoreTransfers(DatastoreTransfers, unittest.TestCase): 

3177 """Test data transfers between butlers. 

3178 

3179 Test for different managers. UUID to UUID and integer to integer are 

3180 tested. UUID to integer is not supported since we do not currently 

3181 want to allow that. Integer to UUID is supported with the caveat 

3182 that UUID4 will be generated and this will be incorrect for raw 

3183 dataset types. The test ignores that. 

3184 """ 

3185 

3186 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3187 

3188 def create_butlers( 

3189 self, manager1: str | None = None, manager2: str | None = None, source_config: str | None = None 

3190 ) -> None: 

3191 self.source_butler = self.create_butler(manager1, "1", config_file=source_config) 

3192 self.target_butler = self.create_butler(manager2, "2") 

3193 

3194 def testTransferUuidToUuid(self) -> None: 

3195 self.create_butlers() 

3196 self.assertButlerTransfers() 

3197 

3198 def testTransferFromChainedUuidToUuid(self) -> None: 

3199 """Force the source butler to be a ChainedDatastore.""" 

3200 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-chained.yaml")) 

3201 self.assertButlerTransfers() 

3202 

3203 def testTransferFromIncompatibleUuidToUuid(self) -> None: 

3204 """Force the source butler to be a incompatible datastore.""" 

3205 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-inmemory.yaml")) 

3206 with self.assertRaises(NotImplementedError): 

3207 self.assertButlerTransfers() 

3208 

3209 def testTransferFromIncompatibleChainUuidToUuid(self) -> None: 

3210 """Force the source butler to be a incompatible datastore.""" 

3211 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-inmemory-chain.yaml")) 

3212 with self.assertRaises(TypeError): 

3213 self.assertButlerTransfers() 

3214 

3215 def testTransferFromFileUuidToUuid(self) -> None: 

3216 """Force the source butler to be a FileDatastore.""" 

3217 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler.yaml")) 

3218 self.assertButlerTransfers() 

3219 

3220 def testTransferMissing(self) -> None: 

3221 """Test transfers where datastore records are missing. 

3222 

3223 This is how execution butler works. 

3224 """ 

3225 self.create_butlers() 

3226 

3227 # Configure the source butler to allow trust. 

3228 self.source_butler._datastore._set_trust_mode(True) 

3229 

3230 self.assertButlerTransfers(purge=True) 

3231 

3232 def testTransferMissingDisassembly(self) -> None: 

3233 """Test transfers where datastore records are missing. 

3234 

3235 This is how execution butler works. 

3236 """ 

3237 self.create_butlers() 

3238 

3239 # Configure the source butler to allow trust. 

3240 self.source_butler._datastore._set_trust_mode(True) 

3241 

3242 # Test disassembly. 

3243 self.assertButlerTransfers(purge=True, storageClassName="StructuredComposite") 

3244 

3245 def testTransferDifferingStorageClasses(self) -> None: 

3246 """Test transfers when the source butler dataset type has a different 

3247 but compatible storage class. 

3248 """ 

3249 self.create_butlers() 

3250 

3251 self.assertButlerTransfers(storageClassNameTarget="MetricsConversion") 

3252 

3253 def testTransferDifferingStorageClassesDisassembly(self) -> None: 

3254 """Test transfers when the source butler dataset type has a different 

3255 but compatible storage class and where the source butler has 

3256 disassembled. 

3257 """ 

3258 self.create_butlers() 

3259 

3260 self.assertButlerTransfers( 

3261 storageClassName="StructuredComposite", storageClassNameTarget="MetricsConversion" 

3262 ) 

3263 

3264 def testUnsafeDirectTransfer(self) -> None: 

3265 """Test that transfer='unsafe_direct' records the absolute URI of 

3266 source files in the target datastore. 

3267 """ 

3268 self.create_butlers() 

3269 dataset_type = DatasetType("dt", [], "int", universe=self.source_butler.dimensions) 

3270 self.source_butler.registry.registerDatasetType(dataset_type) 

3271 self.source_butler.collections.register("run") 

3272 ref = self.source_butler.put(123, "dt", [], run="run") 

3273 self.target_butler.transfer_from( 

3274 self.source_butler, [ref], transfer="unsafe_direct", register_dataset_types=True 

3275 ) 

3276 self.assertEqual(self.target_butler.get(ref), 123) 

3277 self.assertEqual(self.source_butler.getURI(ref), self.target_butler.getURI(ref)) 

3278 

3279 def testAbsoluteURITransferDirect(self) -> None: 

3280 """Test transfer using an absolute URI.""" 

3281 self._absolute_transfer("auto") 

3282 

3283 def testAbsoluteURITransferUnsafeDirect(self) -> None: 

3284 """Test transfer using an absolute URI.""" 

3285 self._absolute_transfer("unsafe_direct") 

3286 

3287 def testAbsoluteURITransferCopy(self) -> None: 

3288 """Test transfer using an absolute URI.""" 

3289 self._absolute_transfer("copy") 

3290 

3291 def _absolute_transfer(self, transfer: str) -> None: 

3292 self.create_butlers() 

3293 

3294 storageClassName = "StructuredData" 

3295 storageClass = self.storageClassFactory.getStorageClass(storageClassName) 

3296 datasetTypeName = "random_data" 

3297 run = "run1" 

3298 self.source_butler.collections.register(run) 

3299 

3300 dimensions = self.source_butler.dimensions.conform(()) 

3301 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

3302 self.source_butler.registry.registerDatasetType(datasetType) 

3303 

3304 metrics = makeExampleMetrics() 

3305 # Ingest from a URI that reports itself as not local, so that the test 

3306 # distinguishes "the absolute URI was preserved" from "a local path 

3307 # happened to work". 

3308 source_dir = os.path.join(self.root, "source data") 

3309 os.makedirs(source_dir) 

3310 with ResourcePath.temporary_uri(prefix=make_remote_test_uri(source_dir), suffix=".json") as temp: 

3311 self.assertFalse(temp.isLocal) 

3312 dataId = DataCoordinate.make_empty(self.source_butler.dimensions) 

3313 source_refs = [DatasetRef(datasetType, dataId, run=run)] 

3314 temp.write(json.dumps(metrics.exportAsDict()).encode()) 

3315 dataset = FileDataset(path=temp, refs=source_refs) 

3316 self.source_butler.ingest(dataset, transfer="direct") 

3317 

3318 self.target_butler.transfer_from( 

3319 self.source_butler, dataset.refs, register_dataset_types=True, transfer=transfer 

3320 ) 

3321 

3322 uri = self.target_butler.getURI(dataset.refs[0]) 

3323 if transfer == "auto" or transfer == "unsafe_direct": 

3324 self.assertEqual(uri, temp) 

3325 else: 

3326 self.assertNotEqual(uri, temp) 

3327 

3328 def test_shared_dimension_group(self): 

3329 """Test internal logic that divides dataset types by dimension group 

3330 when doing registry updates. 

3331 """ 

3332 self.create_butlers() 

3333 self.source_butler.import_(filename=_get_test_data_path("base.yaml"), without_datastore=True) 

3334 self.source_butler.import_(filename=_get_test_data_path("datasets.yaml"), without_datastore=True) 

3335 

3336 source_butler = self.source_butler 

3337 target_butler = self.target_butler 

3338 

3339 # Create a dataset type with the same dimensions as the 'bias' dataset 

3340 # type from base.yaml 

3341 dataset_type = DatasetType( 

3342 "test_type", ["instrument", "detector"], "int", universe=source_butler.dimensions 

3343 ) 

3344 source_butler.registry.registerDatasetType(dataset_type) 

3345 # This has the same data ID as one of the bias datasets in 

3346 # datasets.yaml. 

3347 test_ref = source_butler.registry.insertDatasets( 

3348 "test_type", [{"instrument": "Cam1", "detector": 2}], run="imported_g" 

3349 )[0] 

3350 

3351 biases = source_butler.query_datasets("bias", ["imported_g", "imported_r"]) 

3352 flats = source_butler.query_datasets("flat", ["imported_g", "imported_r"]) 

3353 refs = [test_ref, *biases, *flats] 

3354 

3355 # Test setup will be even more convoluted if we want the datastore to 

3356 # actually transfer files. For testing the dimension group behavior, 

3357 # we really only care about the registry. 

3358 with unittest.mock.patch.object(target_butler._datastore, "transfer_from") as mock: 

3359 mock.return_value = (set(refs), set()) 

3360 target_butler.transfer_from( 

3361 source_butler, 

3362 refs, 

3363 transfer=None, 

3364 register_dataset_types=True, 

3365 skip_missing=False, 

3366 transfer_dimensions=True, 

3367 ) 

3368 

3369 transferred_test_ref = target_butler.find_dataset( 

3370 "test_type", {"instrument": "Cam1", "detector": 2}, collections="imported_g" 

3371 ) 

3372 self.assertEqual(transferred_test_ref.id, test_ref.id) 

3373 

3374 transferred_bias = target_butler.find_dataset( 

3375 "bias", {"instrument": "Cam1", "detector": 2}, collections="imported_g" 

3376 ) 

3377 self.assertEqual(transferred_bias.id, uuid.UUID("51352db4-a47a-447c-b12d-a50b206b17cd")) 

3378 

3379 transferred_flat = target_butler.find_dataset( 

3380 "flat", 

3381 {"instrument": "Cam1", "detector": 2, "physical_filter": "Cam1-R1", "band": "r"}, 

3382 collections="imported_r", 

3383 ) 

3384 self.assertEqual(transferred_flat.id, uuid.UUID("c1296796-56c5-4acf-9b49-40d920c6f840")) 

3385 

3386 

3387class ChainedDatastoreTransfers(PosixDatastoreTransfers): 

3388 """Test transfers using a chained datastore.""" 

3389 

3390 configFile = os.path.join(TESTDIR, "config/basic/butler-chained.yaml") 

3391 

3392 

3393@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3394class ButlerServerDatastoreTransfers(DatastoreTransfers, unittest.TestCase): 

3395 """Test ``transfer_from`` involving Butler server.""" 

3396 

3397 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3398 

3399 def test_transfers_from_remote_to_direct(self) -> None: 

3400 from lsst.daf.butler.remote_butler._remote_file_transfer_source import ( 

3401 mock_file_transfer_uris_for_unit_test, 

3402 ) 

3403 

3404 self.target_butler = self.create_butler(None, "2") 

3405 with create_test_server(TESTDIR) as server: 

3406 self.source_butler = server.hybrid_butler 

3407 

3408 def _remap_transfer_url(path: HttpResourcePath) -> HttpResourcePath: 

3409 # The Butler server returns HTTP URIs with a domain name that 

3410 # is not resolvable because there is no actual HTTP server 

3411 # involved in these tests. Strip this first layer of 

3412 # indirection, and return the target of the redirect instead. 

3413 response = server.client.get(str(path), follow_redirects=False, headers=path._extra_headers) 

3414 return ResourcePath(str(response.next_request.url)) 

3415 

3416 with mock_file_transfer_uris_for_unit_test(_remap_transfer_url): 

3417 self.assertButlerTransfers() 

3418 

3419 

3420class TransferDatasetsInPlace(unittest.TestCase): 

3421 """Test behavior of transfer_datasets_in_place() specialty function used by 

3422 Prompt Publication service. 

3423 """ 

3424 

3425 def test_file_datastore(self) -> None: 

3426 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3427 with ( 

3428 tempfile.TemporaryDirectory() as datastore_root, 

3429 tempfile.TemporaryDirectory() as other_repo_root, 

3430 ): 

3431 config = Config(configFile) 

3432 config["datastore", "datastore", "name"] = "file_datastore" 

3433 Butler.makeRepo(datastore_root, config=config) 

3434 config["datastore", "datastore", "root"] = datastore_root 

3435 Butler.makeRepo(other_repo_root, config, forceConfigRoot=False) 

3436 with ( 

3437 Butler(datastore_root, writeable=True) as source_butler, 

3438 Butler(other_repo_root, writeable=True) as target_butler, 

3439 ): 

3440 self._test_transfer_datasets_in_place(source_butler, target_butler) 

3441 

3442 def test_chained_datastore(self) -> None: 

3443 configFile = os.path.join(TESTDIR, "config/basic/butler-chained-posix.yaml") 

3444 with ( 

3445 tempfile.TemporaryDirectory() as datastore_root, 

3446 tempfile.TemporaryDirectory() as other_repo_root, 

3447 ): 

3448 config = Config(configFile) 

3449 config["datastore", "datastore", "datastores", 0, "datastore", "root"] = ( 

3450 f"{datastore_root}/butler_test_repository" 

3451 ) 

3452 config["datastore", "datastore", "datastores", 1, "datastore", "root"] = ( 

3453 f"{datastore_root}/butler_test_repository2" 

3454 ) 

3455 Butler.makeRepo(datastore_root, config=config, forceConfigRoot=False) 

3456 Butler.makeRepo(other_repo_root, config=config, forceConfigRoot=False) 

3457 with ( 

3458 Butler(datastore_root, writeable=True) as source_butler, 

3459 Butler(other_repo_root, writeable=True) as target_butler, 

3460 ): 

3461 self._test_transfer_datasets_in_place(source_butler, target_butler) 

3462 

3463 def _test_transfer_datasets_in_place( 

3464 self, source_butler: DirectButler, target_butler: DirectButler 

3465 ) -> None: 

3466 metric_repo = MetricTestRepo.create_from_butler( 

3467 source_butler, 

3468 source_butler._config, 

3469 ) 

3470 target_butler.transfer_dimension_records_from(source_butler, [metric_repo.ref1, metric_repo.ref2]) 

3471 # Verify that the setup was correct and the two repos have 

3472 # independent registries. 

3473 self.assertIsNone(target_butler.get_dataset(metric_repo.ref1.id)) 

3474 # Copy one dataset, and make sure we can load it from the 

3475 # target repo. 

3476 self.assertEqual( 

3477 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1]), 

3478 [metric_repo.ref1], 

3479 ) 

3480 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1)) 

3481 self.assertIsNone(target_butler.get_dataset(metric_repo.ref2.id)) 

3482 self.assertEqual(source_butler.getURIs(metric_repo.ref1), target_butler.getURIs(metric_repo.ref1)) 

3483 # Trying to copy the same dataset again is a no-op. 

3484 self.assertEqual( 

3485 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1]), 

3486 [], 

3487 ) 

3488 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1)) 

3489 # A mix of existing and non-existing datasets. 

3490 self.assertEqual( 

3491 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1, metric_repo.ref2]), 

3492 [metric_repo.ref2], 

3493 ) 

3494 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1)) 

3495 self.assertEqual(target_butler.get(metric_repo.ref2), source_butler.get(metric_repo.ref2)) 

3496 

3497 # For testing datastore chaining, set up a dataset that is only 

3498 # accepted by one of the datastores. 

3499 source_butler.registry.registerDatasetType( 

3500 DatasetType("rejected_by_first", source_butler.dimensions.conform([]), "int") 

3501 ) 

3502 source_butler.registry.registerRun("run") 

3503 ref = source_butler.put(1, "rejected_by_first", dataId={}, run="run") 

3504 self.assertEqual( 

3505 transfer_datasets_in_place(source_butler, target_butler, [ref]), 

3506 [ref], 

3507 ) 

3508 self.assertEqual(1, target_butler.get(ref)) 

3509 

3510 

3511class NullDatastoreTestCase(unittest.TestCase): 

3512 """Test that we can fall back to a null datastore.""" 

3513 

3514 # Need a good config to create the repo. 

3515 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3516 storageClassFactory: StorageClassFactory 

3517 

3518 @classmethod 

3519 def setUpClass(cls) -> None: 

3520 cls.storageClassFactory = StorageClassFactory() 

3521 cls.storageClassFactory.addFromConfig(cls.configFile) 

3522 

3523 def setUp(self) -> None: 

3524 """Create a new butler root for each test.""" 

3525 self.root = makeTestTempDir(TESTDIR) 

3526 Butler.makeRepo(self.root, config=Config(self.configFile)) 

3527 

3528 def tearDown(self) -> None: 

3529 removeTestTempDir(self.root) 

3530 

3531 def test_fallback(self) -> None: 

3532 # Read the butler config and mess with the datastore section. 

3533 config_path = os.path.join(self.root, "butler.yaml") 

3534 bad_config = Config(config_path) 

3535 bad_config["datastore", "cls"] = "lsst.not.a.datastore.Datastore" 

3536 bad_config.dumpToUri(config_path) 

3537 

3538 with self.assertRaises(RuntimeError): 

3539 Butler(self.root, without_datastore=False) 

3540 

3541 with self.assertRaises(RuntimeError): 

3542 Butler.from_config(self.root, without_datastore=False) 

3543 

3544 butler = Butler.from_config(self.root, writeable=True, without_datastore=True) 

3545 self.enterContext(butler) 

3546 self.assertIsInstance(butler._datastore, NullDatastore) 

3547 

3548 # Check that registry is working. 

3549 butler.collections.register("MYRUN") 

3550 collections = butler.collections.query("*") 

3551 self.assertIn("MYRUN", set(collections)) 

3552 

3553 # Create a ref. 

3554 dimensions = butler.dimensions.conform([]) 

3555 storageClass = self.storageClassFactory.getStorageClass("StructuredDataDict") 

3556 datasetTypeName = "metric" 

3557 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

3558 butler.registry.registerDatasetType(datasetType) 

3559 ref = DatasetRef(datasetType, {}, run="MYRUN") 

3560 

3561 # Check that datastore will complain. 

3562 with self.assertRaises(FileNotFoundError): 

3563 butler.get(ref) 

3564 with self.assertRaises(FileNotFoundError): 

3565 butler.getURI(ref) 

3566 

3567 

3568@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3569class ButlerServerTests(FileDatastoreButlerTests): 

3570 """Test RemoteButler and Butler server.""" 

3571 

3572 configFile = None 

3573 predictionSupported = False 

3574 trustModeSupported = False 

3575 

3576 postgres: TemporaryPostgresInstance | None 

3577 

3578 def setUp(self): 

3579 self.server_instance = self.enterContext(create_test_server(TESTDIR)) 

3580 

3581 def tearDown(self): 

3582 pass 

3583 

3584 def are_uris_equivalent(self, uri1: ResourcePath, uri2: ResourcePath) -> bool: 

3585 # S3 pre-signed URLs may end up with differing expiration times in the 

3586 # query parameters, so ignore query parameters when comparing. 

3587 return uri1.scheme == uri2.scheme and uri1.netloc == uri2.netloc and uri1.path == uri2.path 

3588 

3589 def create_empty_butler( 

3590 self, 

3591 run: str | None = None, 

3592 writeable: bool | None = None, 

3593 metrics: ButlerMetrics | None = None, 

3594 cleanup: bool = True, 

3595 ) -> Butler: 

3596 return self.server_instance.hybrid_butler.clone(run=run, metrics=metrics) 

3597 

3598 def remove_dataset_out_of_band(self, butler: Butler, ref: DatasetRef) -> None: 

3599 # Can't delete a file via S3 signed URLs, so we need to reach in 

3600 # through DirectButler to delete the dataset. 

3601 uri = self.server_instance.direct_butler.getURI(ref) 

3602 uri.remove() 

3603 

3604 def testConstructor(self): 

3605 # RemoteButler constructor is tested in test_server.py and 

3606 # test_remote_butler.py. 

3607 pass 

3608 

3609 def testDafButlerRepositories(self): 

3610 # Loading of RemoteButler via repository index is tested in 

3611 # test_server.py. 

3612 pass 

3613 

3614 def testGetDatasetTypes(self) -> None: 

3615 # This is mostly a test of validateConfiguration, which is for 

3616 # validating Datastore configuration and thus isn't relevant to 

3617 # RemoteButler. 

3618 pass 

3619 

3620 def testMakeRepo(self) -> None: 

3621 # Only applies to DirectButler. 

3622 pass 

3623 

3624 # Pickling not yet implemented for RemoteButler/HybridButler. 

3625 @unittest.expectedFailure 

3626 def testPickle(self) -> None: 

3627 return super().testPickle() 

3628 

3629 def testStringification(self) -> None: 

3630 self.assertEqual( 

3631 str(self.server_instance.remote_butler), 

3632 "RemoteButler(https://test.example/api/butler/repo/testrepo/)", 

3633 ) 

3634 

3635 def testTransaction(self) -> None: 

3636 # Transactions will never be supported for RemoteButler. 

3637 pass 

3638 

3639 def testPutTemplates(self) -> None: 

3640 # The Butler server instance is configured with different file naming 

3641 # templates than this test is expecting. 

3642 pass 

3643 

3644 

3645@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3646class ButlerServerSqliteTests(ButlerServerTests, unittest.TestCase): 

3647 """Tests for RemoteButler's registry shim, with a SQLite DB backing the 

3648 server. 

3649 """ 

3650 

3651 postgres = None 

3652 

3653 

3654@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3655class ButlerServerPostgresTests(ButlerServerTests, unittest.TestCase): 

3656 """Tests for RemoteButler's registry shim, with a Postgres DB backing the 

3657 server. 

3658 """ 

3659 

3660 @classmethod 

3661 def setUpClass(cls): 

3662 cls.postgres = cls.enterClassContext(setup_postgres_test_db()) 

3663 super().setUpClass() 

3664 

3665 

3666def setup_module(module: types.ModuleType) -> None: 

3667 """Set up the module for pytest.""" 

3668 clean_environment() 

3669 

3670 

3671def _get_test_data_path(filename: str) -> ResourcePath: 

3672 return ResourcePath(f"resource://lsst.daf.butler/tests/registry_data/{filename}") 

3673 

3674 

3675if __name__ == "__main__": 

3676 clean_environment() 

3677 unittest.main()