Coverage for tests/test_butler.py: 96%

1944 statements  

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

38import random 

39import re 

40import shutil 

41import string 

42import tempfile 

43import unittest 

44import unittest.mock 

45import uuid 

46import warnings 

47import weakref 

48from collections.abc import Callable, Mapping 

49from typing import TYPE_CHECKING, Any, cast 

50 

51try: 

52 import boto3 

53 import botocore 

54 

55 from lsst.resources.s3utils import clean_test_environment_for_s3 

56 

57 try: 

58 from moto import mock_aws # v5 

59 except ImportError: 

60 from moto import mock_s3 as mock_aws 

61except ImportError: 

62 boto3 = None 

63 

64 def mock_aws(*args: Any, **kwargs: Any) -> Any: # type: ignore[no-untyped-def] 

65 """No-op decorator in case moto mock_aws can not be imported.""" 

66 return None 

67 

68 

69import astropy.time 

70from sqlalchemy.exc import IntegrityError 

71 

72from lsst.daf.butler import ( 

73 Butler, 

74 ButlerConfig, 

75 ButlerMetrics, 

76 ButlerRepoIndex, 

77 CollectionCycleError, 

78 CollectionType, 

79 Config, 

80 DataCoordinate, 

81 DatasetExistence, 

82 DatasetNotFoundError, 

83 DatasetProvenance, 

84 DatasetRef, 

85 DatasetType, 

86 DimensionRecord, 

87 FileDataset, 

88 NoDefaultCollectionError, 

89 StorageClassFactory, 

90 ValidationError, 

91 script, 

92) 

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

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

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

96from lsst.daf.butler.datastore import NullDatastore 

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

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

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

100from lsst.daf.butler.direct_butler import DirectButler 

101from lsst.daf.butler.registry import ( 

102 CollectionError, 

103 CollectionTypeError, 

104 ConflictingDefinitionError, 

105 DataIdValueError, 

106 DatasetTypeExpressionError, 

107 MissingCollectionError, 

108 OrphanedRecordError, 

109) 

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

111from lsst.daf.butler.repo_relocation import BUTLER_ROOT_TAG 

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

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

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

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

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

117 MetricTestRepo, 

118 TestCaseMixin, 

119 create_populated_sqlite_registry, 

120 makeTestTempDir, 

121 removeTestTempDir, 

122 safeTestTempDir, 

123) 

124from lsst.resources import ResourcePath 

125from lsst.resources.http import HttpResourcePath 

126from lsst.utils import doImportType 

127from lsst.utils.introspection import get_full_type_name 

128 

129if butler_server_is_available: 

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

131 

132 

133if TYPE_CHECKING: 

134 import types 

135 

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

137 

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

139 

140 

141def clean_environment() -> None: 

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

143 for k in ("DAF_BUTLER_REPOSITORY_INDEX",): 

144 os.environ.pop(k, None) 

145 

146 

147def makeExampleMetrics() -> MetricsExample: 

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

149 return MetricsExample( 

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

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

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

153 ) 

154 

155 

156class TransactionTestError(Exception): 

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

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

159 """ 

160 

161 pass 

162 

163 

164class ButlerConfigTests(unittest.TestCase): 

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

166 cases. 

167 """ 

168 

169 def testSearchPath(self) -> None: 

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

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

172 config1 = ButlerConfig(configFile) 

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

174 

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

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

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

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

179 

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

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

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

183 

184 

185class ButlerPutGetTests(TestCaseMixin): 

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

187 butler configurations. 

188 """ 

189 

190 root: str 

191 default_run = "ingésτ😺" 

192 storageClassFactory: StorageClassFactory 

193 configFile: str | None 

194 tmpConfigFile: str 

195 

196 @staticmethod 

197 def addDatasetType( 

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

199 ) -> DatasetType: 

200 """Create a DatasetType and register it""" 

201 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

202 registry.registerDatasetType(datasetType) 

203 return datasetType 

204 

205 @classmethod 

206 def setUpClass(cls) -> None: 

207 cls.storageClassFactory = StorageClassFactory() 

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

209 cls.storageClassFactory.addFromConfig(cls.configFile) 

210 

211 def assertGetComponents( 

212 self, 

213 butler: Butler, 

214 datasetRef: DatasetRef, 

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

216 reference: Any, 

217 collections: Any = None, 

218 ) -> None: 

219 datasetType = datasetRef.datasetType 

220 dataId = datasetRef.dataId 

221 deferred = butler.getDeferred(datasetRef) 

222 

223 for component in components: 

224 compTypeName = datasetType.componentTypeName(component) 

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

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

227 result_deferred = deferred.get(component=component) 

228 self.assertEqual(result_deferred, result) 

229 

230 def tearDown(self) -> None: 

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

232 removeTestTempDir(self.root) 

233 

234 def create_empty_butler( 

235 self, 

236 run: str | None = None, 

237 writeable: bool | None = None, 

238 metrics: ButlerMetrics | None = None, 

239 cleanup: bool = True, 

240 ): 

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

242 data. 

243 """ 

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

245 if cleanup: 

246 self.enterContext(butler) 

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

248 return butler 

249 

250 def create_butler( 

251 self, 

252 run: str, 

253 storageClass: StorageClass | str, 

254 datasetTypeName: str, 

255 metrics: ButlerMetrics | None = None, 

256 ) -> tuple[Butler, DatasetType]: 

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

258 into it. 

259 """ 

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

261 

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

263 self.assertEqual(collections, {run}) 

264 # Create and register a DatasetType 

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

266 

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

268 

269 # Add needed Dimensions 

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

271 butler.registry.insertDimensionData( 

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

273 ) 

274 butler.registry.insertDimensionData( 

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

276 ) 

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

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

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

280 butler.registry.insertDimensionData( 

281 "visit", 

282 { 

283 "instrument": "DummyCamComp", 

284 "id": 423, 

285 "name": "fourtwentythree", 

286 "physical_filter": "d-r", 

287 "datetime_begin": visit_start, 

288 "datetime_end": visit_end, 

289 "day_obs": 20200101, 

290 }, 

291 ) 

292 

293 # Add more visits for some later tests 

294 for visit_id in (424, 425): 

295 butler.registry.insertDimensionData( 

296 "visit", 

297 { 

298 "instrument": "DummyCamComp", 

299 "id": visit_id, 

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

301 "physical_filter": "d-r", 

302 "day_obs": 20200101, 

303 }, 

304 ) 

305 return butler, datasetType 

306 

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

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

309 # tag when looking up datasets. 

310 run = self.default_run 

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

312 assert butler.run is not None 

313 

314 # Create and store a dataset 

315 metric = makeExampleMetrics() 

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

317 

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

319 with self.assertRaises(DatasetNotFoundError): 

320 butler.get(datasetTypeName, dataId) 

321 

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

323 # and once with a DatasetType 

324 

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

326 expected_collections = {run} 

327 

328 counter = 0 

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

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

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

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

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

334 # immediately because the dataset already exists. Work around 

335 # this by using a distinct run collection each time 

336 counter += 1 

337 this_run = f"put_run_{counter}" 

338 butler.collections.register(this_run) 

339 expected_collections.update({this_run}) 

340 

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

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

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

344 kwargs["run"] = this_run 

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

346 self.assertIsInstance(ref, DatasetRef) 

347 

348 # Test get of a ref. 

349 metricOut = butler.get(ref) 

350 self.assertEqual(metric, metricOut) 

351 # Test get 

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

353 self.assertEqual(metric, metricOut) 

354 # Test get with a datasetRef 

355 metricOut = butler.get(ref) 

356 self.assertEqual(metric, metricOut) 

357 # Test getDeferred with dataId 

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

359 self.assertEqual(metric, metricOut) 

360 # Test getDeferred with a ref 

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

362 self.assertEqual(metric, metricOut) 

363 

364 # Check we can get components 

365 if storageClass.isComposite(): 

366 self.assertGetComponents( 

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

368 ) 

369 

370 primary_uri, secondary_uris = butler.getURIs(ref) 

371 n_uris = len(secondary_uris) 

372 if primary_uri: 

373 n_uris += 1 

374 

375 # Can the artifacts themselves be retrieved? 

376 if not butler._datastore.isEphemeral: 

377 # Create a temporary directory to hold the retrieved 

378 # artifacts. 

379 with tempfile.TemporaryDirectory( 

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

381 ) as artifact_root: 

382 root_uri = ResourcePath(artifact_root, forceDirectory=True) 

383 

384 for preserve_path in (True, False): 

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

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

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

388 # Use copy so that we can test that overwrite 

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

390 # would use hard links and subsequent transfer 

391 # would work because it knows they are the same 

392 # file). 

393 transferred = butler.retrieveArtifacts( 

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

395 ) 

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

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

398 # Filter out the index file. 

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

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

401 

402 for artifact in transferred: 

403 path_in_destination = artifact.relative_to(destination) 

404 self.assertIsNotNone(path_in_destination) 

405 assert path_in_destination is not None 

406 

407 # When path is not preserved there should not 

408 # be any path separators. 

409 num_seps = path_in_destination.count("/") 

410 if preserve_path: 

411 self.assertGreater(num_seps, 0) 

412 else: 

413 self.assertEqual(num_seps, 0) 

414 

415 self.assertEqual( 

416 len(artifacts), 

417 n_uris, 

418 "Comparing expected artifacts vs actual:" 

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

420 ) 

421 

422 if preserve_path: 

423 # No need to run these twice 

424 with self.assertRaises(ValueError): 

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

426 

427 with self.assertRaisesRegex( 

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

429 ): 

430 butler.retrieveArtifacts( 

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

432 ) 

433 

434 with self.assertRaises(FileExistsError): 

435 butler.retrieveArtifacts([ref], destination) 

436 

437 transferred_again = butler.retrieveArtifacts( 

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

439 ) 

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

441 

442 # Now remove the dataset completely. 

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

444 # Lookup with original args should still fail. 

445 kwargs = {"collections": this_run} 

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

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

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

449 # get() should still fail. 

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

451 butler.get(ref) 

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

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

454 

455 # Do explicit registry removal since we know they are 

456 # empty 

457 butler.collections.x_remove(this_run) 

458 expected_collections.remove(this_run) 

459 

460 # Create DatasetRef for put using default run. 

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

462 

463 # Check that getDeferred fails with standalone ref. 

464 with self.assertRaises(LookupError): 

465 butler.getDeferred(refIn) 

466 

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

468 # and we want to use the default collection. 

469 ref = butler.put(metric, refIn) 

470 

471 # Get with parameters 

472 stop = 4 

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

474 self.assertNotEqual(metric, sliced) 

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

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

477 assert metric.data is not None # for mypy 

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

479 # getDeferred with parameters 

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

481 self.assertNotEqual(metric, sliced) 

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

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

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

485 # getDeferred with deferred parameters 

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

487 self.assertNotEqual(metric, sliced) 

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

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

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

491 

492 if storageClass.isComposite(): 

493 # Check that components can be retrieved 

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

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

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

497 summary = butler.get(compNameS, dataId) 

498 self.assertEqual(summary, metric.summary) 

499 data = butler.get(compNameD, dataId) 

500 self.assertEqual(data, metric.data) 

501 

502 if "counter" in storageClass.derivedComponents: 

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

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

505 

506 count = butler.get( 

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

508 ) 

509 self.assertEqual(count, stop) 

510 

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

512 assert compRef is not None 

513 summary = butler.get(compRef) 

514 self.assertEqual(summary, metric.summary) 

515 

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

517 inconsistentDatasetType = DatasetType( 

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

519 ) 

520 

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

522 with self.assertRaisesRegex( 

523 ValueError, 

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

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

526 ): 

527 butler.get(inconsistentDatasetType, dataId) 

528 

529 # Combining a DatasetRef with a dataId should fail 

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

531 butler.get(ref, dataId) 

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

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

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

535 

536 # Getting a dataset with unknown parameters should fail 

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

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

539 

540 # Check we have a collection 

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

542 self.assertEqual(collections, expected_collections) 

543 

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

545 # already had a component removed 

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

547 

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

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

550 

551 # Repeat put will fail. 

552 with self.assertRaisesRegex( 

553 ConflictingDefinitionError, "A database constraint failure was triggered" 

554 ): 

555 butler.put(metric, datasetType, dataId) 

556 

557 # Remove the datastore entry. 

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

559 

560 # Put will still fail 

561 with self.assertRaisesRegex( 

562 ConflictingDefinitionError, "A database constraint failure was triggered" 

563 ): 

564 butler.put(metric, datasetType, dataId) 

565 

566 # Repeat the same sequence with resolved ref. 

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

568 ref = butler.put(metric, refIn) 

569 

570 # Repeat put will fail. 

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

572 butler.put(metric, refIn) 

573 

574 # Remove the datastore entry. 

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

576 

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

578 ref = butler.put(metric, refIn) 

579 

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

581 # something to be present 

582 

583 return butler 

584 

585 def testDeferredCollectionPassing(self) -> None: 

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

587 butler = self.create_empty_butler(writeable=True) 

588 # Create and register a DatasetType 

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

590 datasetType = self.addDatasetType( 

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

592 ) 

593 # Add needed Dimensions 

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

595 butler.registry.insertDimensionData( 

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

597 ) 

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

599 butler.registry.insertDimensionData( 

600 "visit", 

601 { 

602 "instrument": "DummyCamComp", 

603 "id": 423, 

604 "name": "fourtwentythree", 

605 "physical_filter": "d-r", 

606 "day_obs": 20250101, 

607 }, 

608 ) 

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

610 # Create dataset. 

611 metric = makeExampleMetrics() 

612 # Register a new run and put dataset. 

613 run = "deferred" 

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

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

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

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

618 # Putting with no run should fail with TypeError. 

619 with self.assertRaises(CollectionError): 

620 butler.put(metric, datasetType, dataId) 

621 # Dataset should exist. 

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

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

624 # a deferred dataset handle. 

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

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

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

628 with self.assertRaises(NoDefaultCollectionError): 

629 butler.exists(datasetType, dataId) 

630 with self.assertRaises(CollectionError): 

631 butler.get(datasetType, dataId) 

632 # Associate the dataset with a different collection. 

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

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

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

636 # in the original collection. 

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

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

639 

640 

641class ButlerTests(ButlerPutGetTests): 

642 """Tests for Butler.""" 

643 

644 useTempRoot = True 

645 validationCanFail: bool 

646 fullConfigKey: str | None 

647 registryStr: str | None 

648 datastoreName: list[str] | None 

649 datastoreStr: list[str] 

650 predictionSupported = True 

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

652 

653 def setUp(self) -> None: 

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

655 self.root = makeTestTempDir(TESTDIR) 

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

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

658 

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

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

661 

662 Subclasses may override to handle unique requirements. 

663 """ 

664 return uri1 == uri2 

665 

666 def testConstructor(self) -> None: 

667 """Independent test of constructor.""" 

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

669 self.enterContext(butler) 

670 self.assertIsInstance(butler, Butler) 

671 

672 # Check that butler.yaml is added automatically. 

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

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

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

676 self.enterContext(butler) 

677 self.assertIsInstance(butler, Butler) 

678 

679 # Even with a ResourcePath. 

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

681 self.enterContext(butler) 

682 self.assertIsInstance(butler, Butler) 

683 

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

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

686 

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

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

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

690 self.enterContext(butler_special) 

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

692 self.assertEqual(collections, {special_run}) 

693 

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

695 self.enterContext(butler2) 

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

697 self.assertIsNone(butler2.run) 

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

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

700 

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

702 # repository. 

703 butler_index = Config() 

704 butler_index["label"] = self.tmpConfigFile 

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

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

707 # we aren't reusing the cache. 

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

709 butler_index["bad_label"] = bad_label 

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

711 butler_index.dumpToUri(temp_file) 

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

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

714 uri = Butler.get_repo_uri("bad_label") 

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

716 uri = Butler.get_repo_uri("label") 

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

718 self.assertIsInstance(butler, Butler) 

719 butler.close() 

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

721 self.assertIsInstance(butler, Butler) 

722 butler.close() 

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

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

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

726 Butler.from_config("bad_label") 

727 with self.assertRaises(FileNotFoundError): 

728 # Should ignore aliases. 

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

730 with self.assertRaises(KeyError) as cm: 

731 Butler.get_repo_uri("missing") 

732 self.assertEqual( 

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

734 ) 

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

736 # Should report no failure. 

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

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

739 # Now with empty configuration. 

740 butler_index = Config() 

741 butler_index.dumpToUri(temp_file) 

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

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

744 Butler.from_config("label") 

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

746 # Now with bad contents. 

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

748 print("'", file=fh) 

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

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

751 Butler.from_config("label") 

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

753 with self.assertRaises(FileNotFoundError): 

754 Butler.get_repo_uri("label") 

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

756 

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

758 Butler.from_config("label") 

759 

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

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

762 self.enterContext(butler) 

763 self.assertIsInstance(butler, Butler) 

764 with self.assertRaises(RuntimeError) as cm: 

765 # No environment variable set. 

766 Butler.get_repo_uri("label") 

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

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

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

770 # No aliases registered. 

771 Butler.from_config("not_there") 

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

773 

774 def testClose(self): 

775 butler = self.create_empty_butler(cleanup=False) 

776 is_direct_butler = isinstance(butler, DirectButler) 

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

778 self.assertFalse(butler._closed) 

779 

780 with butler as butler_from_context_manager: 

781 self.assertIs(butler, butler_from_context_manager) 

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

783 self.assertTrue(butler._closed) 

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

785 butler.get_dataset_type("raw") 

786 

787 # Close may be called multiple times. 

788 butler.close() 

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

790 self.assertTrue(butler._closed) 

791 

792 def testGarbageCollection(self): 

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

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

795 """ 

796 butler = self.create_empty_butler(cleanup=False) 

797 is_direct_butler = isinstance(butler, DirectButler) 

798 butler_ref = weakref.ref(butler) 

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

800 registry_ref = weakref.ref(butler._registry) 

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

802 datastore_ref = weakref.ref(butler._datastore) 

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

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

805 

806 with warnings.catch_warnings(): 

807 # Hide warnings from unclosed database handles. 

808 warnings.simplefilter("ignore", ResourceWarning) 

809 del butler 

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

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

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

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

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

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

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

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

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

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

820 engine = engine_ref() 

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

822 engine.dispose() 

823 

824 def testDafButlerRepositories(self): 

825 with unittest.mock.patch.dict( 

826 os.environ, 

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

828 ): 

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

830 

831 with unittest.mock.patch.dict( 

832 os.environ, 

833 { 

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

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

836 }, 

837 ): 

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

839 Butler.get_repo_uri("label") 

840 

841 with unittest.mock.patch.dict( 

842 os.environ, 

843 {"DAF_BUTLER_REPOSITORIES": "invalid"}, 

844 ): 

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

846 Butler.get_repo_uri("label") 

847 

848 def testBasicPutGet(self) -> None: 

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

850 self.runPutGetTest(storageClass, "test_metric") 

851 

852 def testCompositePutGetConcrete(self) -> None: 

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

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

855 

856 # Should *not* be disassembled 

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

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

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

860 self.assertIsInstance(uri, ResourcePath) 

861 self.assertFalse(components) 

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

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

864 

865 # Predicted dataset 

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

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

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

869 self.assertFalse(components) 

870 self.assertIsInstance(uri, ResourcePath) 

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

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

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

874 ref = DatasetRef( 

875 datasets[0].datasetType, 

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

877 run=self.default_run, 

878 ) 

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

880 self.assertFalse(components2) 

881 self.assertEqual(uri, uri2) 

882 

883 def testCompositePutGetVirtual(self) -> None: 

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

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

886 

887 # Should be disassembled 

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

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

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

891 

892 if butler._datastore.isEphemeral: 

893 # Never disassemble in-memory datastore 

894 self.assertIsInstance(uri, ResourcePath) 

895 self.assertFalse(components) 

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

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

898 else: 

899 self.assertIsNone(uri) 

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

901 for compuri in components.values(): 

902 self.assertIsInstance(compuri, ResourcePath) 

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

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

905 

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

907 # Predicted dataset 

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

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

910 

911 if butler._datastore.isEphemeral: 

912 # Never disassembled 

913 self.assertIsInstance(uri, ResourcePath) 

914 self.assertFalse(components) 

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

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

917 else: 

918 self.assertIsNone(uri) 

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

920 for compuri in components.values(): 

921 self.assertIsInstance(compuri, ResourcePath) 

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

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

924 

925 def testStorageClassOverrideGet(self) -> None: 

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

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

928 datasetTypeName = "anything" 

929 run = self.default_run 

930 

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

932 

933 # Create and store a dataset. 

934 metric = makeExampleMetrics() 

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

936 

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

938 

939 # Return native type. 

940 retrieved = butler.get(ref) 

941 self.assertEqual(retrieved, metric) 

942 

943 # Specify an override. 

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

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

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

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

948 self.assertEqual(retrieved, model) 

949 

950 # Defer but override later. 

951 deferred = butler.getDeferred(ref) 

952 model = deferred.get(storageClass=new_sc) 

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

954 self.assertEqual(retrieved, model) 

955 

956 # Defer but override up front. 

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

958 model = deferred.get() 

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

960 self.assertEqual(retrieved, model) 

961 

962 # Retrieve a component. Should be a tuple. 

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

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

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

966 

967 # Parameter on the write storage class should work regardless 

968 # of read storage class. 

969 data = butler.get( 

970 "anything.data", 

971 dataId, 

972 storageClass="StructuredDataDataTestTuple", 

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

974 ) 

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

976 

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

978 # the write storage class. 

979 with self.assertRaises(KeyError): 

980 butler.get( 

981 "anything.data", 

982 dataId, 

983 storageClass="StructuredDataDataTestTuple", 

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

985 ) 

986 

987 def testComponentFromOverriddenStorageClass(self) -> None: 

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

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

990 """ 

991 # StructuredDataNoComponents defines no components at all, whereas 

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

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

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

995 self.assertFalse(write_sc.allComponents()) 

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

997 

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

999 

1000 metric = makeExampleMetrics() 

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

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

1003 

1004 # The composite conversion on its own must work. 

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

1006 

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

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

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

1010 

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

1012 # the storage class override up front. 

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

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

1015 

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

1017 # storage class the read composite declares for it. 

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

1019 self.assertIsInstance(converted, DictConvertibleModel) 

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

1021 

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

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

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

1025 self.assertIsInstance(converted, DictConvertibleModel) 

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

1027 

1028 def testPytypePutCoercion(self) -> None: 

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

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

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

1032 datasetTypeName = "test_metric" 

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

1034 

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

1036 

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

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

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

1040 test_metric = butler.get(metric_ref) 

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

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

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

1044 

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

1046 # a definition matching this python type. 

1047 registry_type = butler.get_dataset_type(datasetTypeName) 

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

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

1050 self.assertEqual(metric2_ref.datasetType, registry_type) 

1051 

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

1053 test_metric2 = butler.get(metric2_ref) 

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

1055 

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

1057 # This should now return a dict. 

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

1059 test_dict2 = butler.get(new_ref) 

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

1061 

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

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

1064 # behavior and return the type of the DatasetType. 

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

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

1067 

1068 def test_ingest_zip(self) -> None: 

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

1070 butler, dataset_type = self.create_butler( 

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

1072 ) 

1073 

1074 metric = makeExampleMetrics() 

1075 refs = [] 

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

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

1078 refs.append(ref) 

1079 

1080 # Retrieve a Zip file. 

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

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

1083 

1084 # Ingest will fail. 

1085 with self.assertRaises(ConflictingDefinitionError): 

1086 butler.ingest_zip(zip) 

1087 

1088 # Clear out the collection. 

1089 butler.removeRuns([self.default_run]) 

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

1091 

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

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

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

1095 

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

1097 with self.assertRaises(ConflictingDefinitionError): 

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

1099 

1100 # This will be a no-op. 

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

1102 

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

1104 new_butler_cfg = Butler.makeRepo(tmpdir) 

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

1106 self.enterContext(new_butler) 

1107 

1108 # This will fail since dimensions records are missing. 

1109 with self.assertRaises(ConflictingDefinitionError): 

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

1111 

1112 # Dry run should work. 

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

1114 

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

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

1117 

1118 # Check that the refs can be read again. 

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

1120 

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

1122 self.assertTrue(uri.exists()) 

1123 

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

1125 # remaining refs to be read. 

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

1127 self.assertTrue(uri.exists()) 

1128 

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

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

1131 

1132 butler.removeRuns([self.default_run]) 

1133 self.assertFalse(uri.exists()) 

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

1135 

1136 with self.assertRaises(ValueError): 

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

1138 

1139 def testIngest(self) -> None: 

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

1141 

1142 # Create and register a DatasetType 

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

1144 

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

1146 datasetTypeName = "metric" 

1147 

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

1149 

1150 # Add needed Dimensions 

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

1152 butler.registry.insertDimensionData( 

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

1154 ) 

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

1156 for detector in (1, 2): 

1157 butler.registry.insertDimensionData( 

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

1159 ) 

1160 

1161 butler.registry.insertDimensionData( 

1162 "visit", 

1163 { 

1164 "instrument": "DummyCamComp", 

1165 "id": 423, 

1166 "name": "fourtwentythree", 

1167 "physical_filter": "d-r", 

1168 "day_obs": 20250101, 

1169 }, 

1170 { 

1171 "instrument": "DummyCamComp", 

1172 "id": 424, 

1173 "name": "fourtwentyfour", 

1174 "physical_filter": "d-r", 

1175 "day_obs": 20250101, 

1176 }, 

1177 ) 

1178 

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

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

1181 datasets = [] 

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

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

1184 # required. 

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

1186 for detector in (1, 2): 

1187 detector_name = f"detector_{detector}" 

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

1189 dataId = butler.registry.expandDataId( 

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

1191 ) 

1192 # Create a DatasetRef for ingest 

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

1194 

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

1196 

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

1198 

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

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

1201 

1202 metrics1 = butler.get(datasetTypeName, dataId1) 

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

1204 self.assertNotEqual(metrics1, metrics2) 

1205 

1206 # Compare URIs 

1207 uri1 = butler.getURI(datasetTypeName, dataId1) 

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

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

1210 

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

1212 # skip_existing=False. 

1213 with self.assertRaises(ConflictingDefinitionError): 

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

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

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

1217 

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

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

1220 refs = [] 

1221 for detector in (1, 2): 

1222 detector_name = f"detector_{detector}" 

1223 dataId = butler.registry.expandDataId( 

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

1225 ) 

1226 # Create a DatasetRef for ingest 

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

1228 

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

1230 # have disappeared following ingest. 

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

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

1233 

1234 datasets = [] 

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

1236 

1237 # For first ingest use copy. 

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

1239 

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

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

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

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

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

1245 datasets[0].refs = [ 

1246 cast( 

1247 DatasetRef, 

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

1249 ) 

1250 for ref in datasets[0].refs 

1251 ] 

1252 all_refs = [] 

1253 for dataset in datasets: 

1254 refs = [] 

1255 for ref in dataset.refs: 

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

1257 new_data_id = dict(ref.dataId.required) 

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

1259 assert new_ref is not None 

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

1261 refs.append(new_ref) 

1262 dataset.refs = refs 

1263 all_refs.extend(dataset.refs) 

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

1265 

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

1267 # disable recording of file size. 

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

1269 

1270 # Check that every ref now has records. 

1271 for dataset in datasets: 

1272 for ref in dataset.refs: 

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

1274 

1275 # Ensure that the file has disappeared. 

1276 self.assertFalse(tempFile.exists()) 

1277 

1278 # Check that the datastore recorded no file size. 

1279 # Not all datastores can support this. 

1280 try: 

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

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

1283 except AttributeError: 

1284 pass 

1285 

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

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

1288 

1289 multi1 = butler.get(datasetTypeName, dataId1) 

1290 multi2 = butler.get(datasetTypeName, dataId2) 

1291 

1292 self.assertEqual(multi1, metrics1) 

1293 self.assertEqual(multi2, metrics2) 

1294 

1295 # Compare URIs 

1296 uri1 = butler.getURI(datasetTypeName, dataId1) 

1297 uri2 = butler.getURI(datasetTypeName, dataId2) 

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

1299 

1300 # Test that removing one does not break the second 

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

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

1303 # files. 

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

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

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

1307 multi2b = butler.get(datasetTypeName, dataId2) 

1308 self.assertEqual(multi2, multi2b) 

1309 

1310 # Ensure we can ingest 0 datasets 

1311 datasets = [] 

1312 butler.ingest(*datasets) 

1313 

1314 def testPickle(self) -> None: 

1315 """Test pickle support.""" 

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

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

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

1319 self.enterContext(butlerOut) 

1320 self.assertIsInstance(butlerOut, Butler) 

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

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

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

1324 

1325 def testGetDatasetTypes(self) -> None: 

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

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

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

1329 ( 

1330 "instrument", 

1331 [ 

1332 {"instrument": "DummyCam"}, 

1333 {"instrument": "DummyHSC"}, 

1334 {"instrument": "DummyCamComp"}, 

1335 ], 

1336 ), 

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

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

1339 ( 

1340 "visit", 

1341 [ 

1342 { 

1343 "instrument": "DummyCam", 

1344 "id": 42, 

1345 "name": "fortytwo", 

1346 "physical_filter": "d-r", 

1347 "day_obs": 20250101, 

1348 } 

1349 ], 

1350 ), 

1351 ] 

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

1353 # Add needed Dimensions 

1354 for element, data in dimensionEntries: 

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

1356 

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

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

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

1360 components = set() 

1361 for datasetTypeName in datasetTypeNames: 

1362 # Create and register a DatasetType 

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

1364 

1365 for componentName in storageClass.components: 

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

1367 

1368 fromRegistry: set[DatasetType] = set() 

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

1370 fromRegistry.add(parent_dataset_type) 

1371 fromRegistry.update(parent_dataset_type.makeAllComponentDatasetTypes()) 

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

1373 

1374 # Query with wildcard. 

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

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

1377 # but not regex. 

1378 with self.assertRaises(DatasetTypeExpressionError): 

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

1380 

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

1382 butler.validateConfiguration( 

1383 ignore=[ 

1384 "test_metric_comp", 

1385 "metric3", 

1386 "metric5", 

1387 "calexp", 

1388 "DummySC", 

1389 "datasetType.component", 

1390 "random_data", 

1391 "random_data_2", 

1392 ] 

1393 ) 

1394 

1395 # Add a new datasetType that will fail template validation 

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

1397 if self.validationCanFail: 

1398 with self.assertRaises(ValidationError): 

1399 butler.validateConfiguration() 

1400 

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

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

1403 

1404 # Rerun validation but ignore the bad datasetType 

1405 butler.validateConfiguration( 

1406 ignore=[ 

1407 "test_metric_comp", 

1408 "metric3", 

1409 "metric5", 

1410 "calexp", 

1411 "DummySC", 

1412 "datasetType.component", 

1413 "random_data", 

1414 "random_data_2", 

1415 ] 

1416 ) 

1417 

1418 def testTransaction(self) -> None: 

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

1420 datasetTypeName = "test_metric" 

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

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

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

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

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

1426 ( 

1427 "visit", 

1428 { 

1429 "instrument": "DummyCam", 

1430 "id": 42, 

1431 "name": "fortytwo", 

1432 "physical_filter": "d-r", 

1433 "day_obs": 20250101, 

1434 }, 

1435 ), 

1436 ) 

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

1438 metric = makeExampleMetrics() 

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

1440 # Create and register a DatasetType 

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

1442 with self.assertRaises(TransactionTestError): 

1443 with butler.transaction(): 

1444 # Add needed Dimensions 

1445 for args in dimensionEntries: 

1446 butler.registry.insertDimensionData(*args) 

1447 # Store a dataset 

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

1449 self.assertIsInstance(ref, DatasetRef) 

1450 # Test get of a ref. 

1451 metricOut = butler.get(ref) 

1452 self.assertEqual(metric, metricOut) 

1453 # Test get 

1454 metricOut = butler.get(datasetTypeName, dataId) 

1455 self.assertEqual(metric, metricOut) 

1456 # Check we can get components 

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

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

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

1460 butler.registry.expandDataId(dataId) 

1461 # Should raise LookupError for missing data ID value 

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

1463 butler.get(datasetTypeName, dataId) 

1464 # Also check explicitly if Dataset entry is missing 

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

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

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

1468 butler.get(ref) 

1469 

1470 def testMakeRepo(self) -> None: 

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

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

1473 repo root. 

1474 """ 

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

1476 # not support a file system root 

1477 if self.fullConfigKey is None: 

1478 return 

1479 

1480 # create two separate directories 

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

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

1483 

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

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

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

1487 limited = Config(self.configFile) 

1488 butler1 = Butler.from_config(butlerConfig) 

1489 self.enterContext(butler1) 

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

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

1492 full = Config(self.tmpConfigFile) 

1493 butler2 = Butler.from_config(butlerConfig) 

1494 self.enterContext(butler2) 

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

1496 # Butlers should have the same configuration regardless of whether 

1497 # defaults were expanded. 

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

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

1500 self.assertNotEqual(limited, full) 

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

1502 # inheriting from defaults. 

1503 self.assertIn(self.fullConfigKey, full) 

1504 self.assertNotIn(self.fullConfigKey, limited) 

1505 

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

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

1508 self.assertEqual(collections1, set()) 

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

1510 

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

1512 # work properly with relocatable Butler repo 

1513 butlerConfig.configFile = None 

1514 with self.assertRaises(ValueError): 

1515 Butler.from_config(butlerConfig) 

1516 

1517 with self.assertRaises(FileExistsError): 

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

1519 

1520 def testStringification(self) -> None: 

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

1522 self.enterContext(butler) 

1523 butlerStr = str(butler) 

1524 

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

1526 for testStr in self.datastoreStr: 

1527 self.assertIn(testStr, butlerStr) 

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

1529 self.assertIn(self.registryStr, butlerStr) 

1530 

1531 datastoreName = butler._datastore.name 

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

1533 for testStr in self.datastoreName: 

1534 self.assertIn(testStr, datastoreName) 

1535 

1536 def testButlerRewriteDataId(self) -> None: 

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

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

1539 

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

1541 datasetTypeName = "random_data" 

1542 

1543 # Create dimension records. 

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

1545 butler.registry.insertDimensionData( 

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

1547 ) 

1548 butler.registry.insertDimensionData( 

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

1550 ) 

1551 

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

1553 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

1554 butler.registry.registerDatasetType(datasetType) 

1555 

1556 n_exposures = 5 

1557 dayobs = 20210530 

1558 

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

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

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

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

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

1564 

1565 for i in range(n_exposures): 

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

1567 butler.registry.insertDimensionData( 

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

1569 ) 

1570 butler.registry.insertDimensionData( 

1571 "exposure", 

1572 { 

1573 "instrument": "DummyCamComp", 

1574 "id": day_obs + i, 

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

1576 "seq_num": i, 

1577 "day_obs": day_obs, 

1578 "physical_filter": "d-r", 

1579 "group": group_name, 

1580 }, 

1581 ) 

1582 

1583 # Write some data. 

1584 for i in range(n_exposures): 

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

1586 

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

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

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

1590 

1591 # Check that the exposure is correct in the dataId 

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

1593 

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

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

1596 self.assertEqual(new_metric, metric) 

1597 

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

1599 # exposure.day_obs. 

1600 datasets_1 = list( 

1601 butler.registry.queryDatasets( 

1602 datasetType, 

1603 collections=self.default_run, 

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

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

1606 ) 

1607 ) 

1608 datasets_2 = list( 

1609 butler.registry.queryDatasets( 

1610 datasetType, 

1611 collections=self.default_run, 

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

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

1614 ) 

1615 ) 

1616 self.assertEqual(datasets_1, datasets_2) 

1617 

1618 def testGetDatasetCollectionCaching(self): 

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

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

1621 # after the collection cache was last updated. 

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

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

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

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

1626 get_ref = reader_butler.get_dataset(put_ref.id) 

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

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

1629 # instance. 

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

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

1632 

1633 def testCollectionChainRedefine(self): 

1634 butler = self._setup_to_test_collection_chain() 

1635 

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

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

1638 

1639 # Duplicates are removed from the list of children 

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

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

1642 

1643 # Empty list clears the chain 

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

1645 self._check_chain(butler, []) 

1646 

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

1648 

1649 def testCollectionChainPrepend(self): 

1650 butler = self._setup_to_test_collection_chain() 

1651 

1652 # Duplicates are removed from the list of children 

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

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

1655 

1656 # Prepend goes on the front of existing chain 

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

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

1659 

1660 # Empty prepend does nothing 

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

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

1663 

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

1665 # their current position. 

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

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

1668 

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

1670 

1671 def testCollectionChainExtend(self): 

1672 butler = self._setup_to_test_collection_chain() 

1673 

1674 # Duplicates are removed from the list of children 

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

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

1677 

1678 # Extend goes on the end of existing chain 

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

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

1681 

1682 # Empty extend does nothing 

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

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

1685 

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

1687 # their current position. 

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

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

1690 

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

1692 

1693 def testCollectionChainRemove(self) -> None: 

1694 butler = self._setup_to_test_collection_chain() 

1695 

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

1697 

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

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

1700 

1701 # Duplicates are allowed in the list of children 

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

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

1704 

1705 # Empty remove does nothing 

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

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

1708 

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

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

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

1712 

1713 self._test_common_chain_functionality( 

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

1715 ) 

1716 

1717 def _setup_to_test_collection_chain(self) -> Butler: 

1718 butler = self.create_empty_butler(writeable=True) 

1719 

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

1721 

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

1723 for run in runs: 

1724 butler.collections.register(run) 

1725 

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

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

1728 

1729 return butler 

1730 

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

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

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

1734 

1735 def _test_common_chain_functionality( 

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

1737 ) -> None: 

1738 # Missing parent collection 

1739 with self.assertRaises(MissingCollectionError): 

1740 func("doesnotexist", []) 

1741 # Missing child collection 

1742 with self.assertRaises(MissingCollectionError): 

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

1744 # Forbid operations on non-chained collections 

1745 with self.assertRaises(CollectionTypeError): 

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

1747 

1748 # Prevent collection cycles 

1749 if not skip_cycle_check: 

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

1751 func("chain2", "chain") 

1752 with self.assertRaises(CollectionCycleError): 

1753 func("chain", "chain2") 

1754 

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

1756 # chains. 

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

1758 

1759 with butler._caching_context(): 

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

1761 func("chain", "a") 

1762 

1763 def test_transfer_dimension_records_from(self) -> None: 

1764 source_butler = self.create_empty_butler(writeable=True) 

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

1766 

1767 visit_id = 2025120200439 

1768 exposure_id = visit_id 

1769 target_butler = self.enterContext(create_populated_sqlite_registry()) 

1770 target_butler.transfer_dimension_records_from( 

1771 source_butler, 

1772 [ 

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

1774 # "populated_by" records (visit_detector_region, 

1775 # visit_definition, etc.) 

1776 DataCoordinate.standardize( 

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

1778 universe=source_butler.dimensions, 

1779 ), 

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

1781 DataCoordinate.make_empty(source_butler.dimensions), 

1782 ], 

1783 ) 

1784 

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

1786 records = target_butler.query_dimension_records(dimension) 

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

1788 return records[0] 

1789 

1790 visit = _fetch_record("visit") 

1791 self.assertEqual(visit.id, visit_id) 

1792 self.assertEqual(visit.day_obs, 20251202) 

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

1794 self.assertEqual(visit.seq_num, 439) 

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

1796 0 

1797 ] 

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

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

1800 

1801 visit_detector_region = _fetch_record("visit_detector_region") 

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

1803 self.assertEqual(visit_detector_region.detector, 10) 

1804 self.assertEqual(visit_detector_region.visit, visit_id) 

1805 original_visit_detector_region = source_butler.query_dimension_records( 

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

1807 )[0] 

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

1809 

1810 visit_definition = _fetch_record("visit_definition") 

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

1812 self.assertEqual(visit_definition.exposure, 2025120200439) 

1813 self.assertEqual(visit_definition.visit, visit_id) 

1814 

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

1816 # visit -> visit_definition. 

1817 exposure = _fetch_record("exposure") 

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

1819 self.assertEqual(exposure.id, 2025120200439) 

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

1821 original_exposure = source_butler.query_dimension_records( 

1822 "exposure", instrument="LSSTCam", exposure=exposure_id 

1823 )[0] 

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

1825 

1826 group = _fetch_record("group") 

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

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

1829 

1830 visit_system_memberships = target_butler.query_dimension_records("visit_system_membership") 

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

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

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

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

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

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

1837 

1838 visit_systems = target_butler.query_dimension_records("visit_system") 

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

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

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

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

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

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

1845 

1846 

1847class FileDatastoreButlerTests(ButlerTests): 

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

1849 by datastores that inherit from FileDatastore. 

1850 """ 

1851 

1852 trustModeSupported = True 

1853 

1854 def testComponentFromOverriddenStorageClassWarns(self) -> None: 

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

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

1857 converted before the component can be extracted. 

1858 """ 

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

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

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

1862 metric = makeExampleMetrics() 

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

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

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

1866 

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

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

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

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

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

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

1873 # dataset has to be converted to. 

1874 self.assertIn("summary", message) 

1875 self.assertIn(write_sc.name, message) 

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

1877 self.assertIn(read_sc.name, message) 

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

1879 

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

1881 # warn. 

1882 composite_type = self.addDatasetType( 

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

1884 ) 

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

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

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

1888 

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

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

1891 

1892 Test testPutTemplates verifies actual physical existance of the files 

1893 in the requested location. 

1894 """ 

1895 uri = ResourcePath(root, forceDirectory=True) 

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

1897 

1898 def testPutTemplates(self) -> None: 

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

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

1901 

1902 # Add needed Dimensions 

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

1904 butler.registry.insertDimensionData( 

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

1906 ) 

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

1908 butler.registry.insertDimensionData( 

1909 "visit", 

1910 { 

1911 "instrument": "DummyCamComp", 

1912 "id": 423, 

1913 "name": "v423", 

1914 "physical_filter": "d-r", 

1915 "day_obs": 20250101, 

1916 }, 

1917 ) 

1918 butler.registry.insertDimensionData( 

1919 "visit", 

1920 { 

1921 "instrument": "DummyCamComp", 

1922 "id": 425, 

1923 "name": "v425", 

1924 "physical_filter": "d-r", 

1925 "day_obs": 20250101, 

1926 }, 

1927 ) 

1928 

1929 # Create and store a dataset 

1930 metric = makeExampleMetrics() 

1931 

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

1933 # template) 

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

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

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

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

1938 

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

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

1941 

1942 # Put with exactly the data ID keys needed 

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

1944 uri = butler.getURI(ref) 

1945 self.assertTrue(uri.exists()) 

1946 self.assertTrue( 

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

1948 ) 

1949 

1950 # Check the template based on dimensions 

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

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

1953 

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

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

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

1957 # must be consistent). 

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

1959 uri = butler.getURI(ref) 

1960 self.assertTrue(uri.exists()) 

1961 self.assertTrue( 

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

1963 ) 

1964 

1965 # Check the template based on dimensions 

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

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

1968 

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

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

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

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

1973 path = template.format(ref) 

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

1975 

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

1977 with self.assertRaises(KeyError): 

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

1979 template.format(ref) 

1980 

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

1982 with self.assertRaises(FileTemplateValidationError): 

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

1984 

1985 def testImportExport(self) -> None: 

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

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

1988 self.runImportExportTest(storageClass) 

1989 

1990 @unittest.expectedFailure 

1991 def testImportExportVirtualComposite(self) -> None: 

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

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

1994 self.runImportExportTest(storageClass) 

1995 

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

1997 """Test exporting and importing. 

1998 

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

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

2001 """ 

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

2003 

2004 # Test that we must have a file extension. 

2005 with self.assertRaises(ValueError): 

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

2007 pass 

2008 

2009 # Test that unknown format is not allowed. 

2010 with self.assertRaises(ValueError): 

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

2012 pass 

2013 

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

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

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

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

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

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

2020 # Export and then import datasets. 

2021 with safeTestTempDir(TESTDIR) as exportDir: 

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

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

2024 export.saveDatasets(datasets) 

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

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

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

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

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

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

2031 # because of internal deduplication. 

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

2033 # Save some dimension records directly. 

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

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

2036 with safeTestTempDir(TESTDIR) as importDir: 

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

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

2039 # Calling script.butlerImport tests the implementation of the 

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

2041 # in the script folder are generally considered protected and 

2042 # should not be used as public api. 

2043 with open(exportFile) as f: 

2044 script.butlerImport( 

2045 importDir, 

2046 export_file=f, 

2047 directory=exportDir, 

2048 transfer="auto", 

2049 skip_dimensions=None, 

2050 ) 

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

2052 self.enterContext(importButler) 

2053 for ref in datasets: 

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

2055 # Test for existence by passing in the DatasetType and 

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

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

2058 self.assertEqual( 

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

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

2061 ) 

2062 

2063 def testRemoveRuns(self) -> None: 

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

2065 butler = self.create_empty_butler(writeable=True) 

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

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

2068 # Add some RUN-type collection. 

2069 run1 = "run1" 

2070 butler.collections.register(run1) 

2071 run2 = "run2" 

2072 butler.collections.register(run2) 

2073 # put a dataset in each 

2074 metric = makeExampleMetrics() 

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

2076 datasetType = self.addDatasetType( 

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

2078 ) 

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

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

2081 uri1 = butler.getURI(ref1) 

2082 uri2 = butler.getURI(ref2) 

2083 

2084 # Put one of the runs in a chain. 

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

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

2087 

2088 with self.assertRaises(OrphanedRecordError): 

2089 butler.registry.removeDatasetType(datasetType.name) 

2090 

2091 # Remove a non-run. 

2092 with self.assertRaises(TypeError): 

2093 butler.removeRuns(["Chain"]) 

2094 

2095 # Remove without unlinking from chain should fail. 

2096 with self.assertRaises(IntegrityError): 

2097 butler.removeRuns([run1]) 

2098 

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

2100 # always purges. 

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

2102 

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

2104 # not think either exists. 

2105 with self.assertRaises(MissingCollectionError): 

2106 butler.collections.get_info(run1) 

2107 with self.assertRaises(MissingCollectionError): 

2108 butler.collections.get_info(run1) 

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

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

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

2112 self.assertFalse(uri1.exists()) 

2113 self.assertFalse(uri2.exists()) 

2114 

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

2116 # dataset type 

2117 butler.registry.removeDatasetType(datasetType.name) 

2118 

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

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

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

2122 

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

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

2125 knowledge. 

2126 

2127 Subclasses may override to handle more complicated datastore 

2128 configurations. 

2129 """ 

2130 uri = butler.getURI(ref) 

2131 uri.remove() 

2132 datastore = cast(FileDatastore, butler._datastore) 

2133 datastore.cacheManager.remove_from_cache(ref) 

2134 

2135 def testPruneDatasets(self) -> None: 

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

2137 butler = self.create_empty_butler(writeable=True) 

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

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

2140 # Add some RUN-type collections. 

2141 run1 = "run1" 

2142 butler.collections.register(run1) 

2143 run2 = "run2" 

2144 butler.collections.register(run2) 

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

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

2147 metric = makeExampleMetrics() 

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

2149 datasetType = self.addDatasetType( 

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

2151 ) 

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

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

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

2155 

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

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

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

2159 

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

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

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

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

2164 

2165 # Simple prune. 

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

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

2168 

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

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

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

2172 

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

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

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

2176 

2177 # Put data back. 

2178 ref1_new = butler.put(metric, ref1) 

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

2180 ref2 = butler.put(metric, ref2) 

2181 

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

2183 self.assertTrue(many_stored[ref1]) 

2184 self.assertTrue(many_stored[ref2]) 

2185 self.assertFalse(many_stored[ref3]) 

2186 

2187 ref3 = butler.put(metric, ref3) 

2188 

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

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

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

2192 

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

2194 refs = [ref1, ref2, ref3] 

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

2196 for ref in refs: 

2197 butler.put(metric, ref) 

2198 

2199 # Confirm we can retrieve deferred. 

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

2201 metric1 = dref1.get() 

2202 self.assertEqual(metric1, metric) 

2203 

2204 # Test different forms of file availability. 

2205 # Need to be in a state where: 

2206 # - one ref just has registry record. 

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

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

2209 # - one ref does not exist anywhere. 

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

2211 # above. 

2212 ref0 = DatasetRef( 

2213 datasetType, 

2214 DataCoordinate.standardize( 

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

2216 ), 

2217 run=run1, 

2218 ) 

2219 

2220 # Delete from datastore and retain in Registry. 

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

2222 

2223 # File has been removed. 

2224 self.remove_dataset_out_of_band(butler, ref2) 

2225 

2226 # Datastore has lost track. 

2227 butler._datastore.forget([ref3]) 

2228 

2229 # First test with a standard butler. 

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

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

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

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

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

2235 

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

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

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

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

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

2241 self.assertTrue(exists_many[ref2]) 

2242 

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

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

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

2246 

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

2248 # retrieved. 

2249 with self.assertRaises(LookupError): 

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

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

2252 with self.assertRaises(FileNotFoundError): 

2253 dref2.get() 

2254 

2255 # Test again with a trusting butler. 

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

2257 butler._datastore.trustGetRequest = True 

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

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

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

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

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

2263 

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

2265 # known but does exist. 

2266 dref3 = butler.getDeferred(ref3) 

2267 metric3 = dref3.get() 

2268 self.assertEqual(metric3, metric) 

2269 

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

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

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

2273 

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

2275 # but is not the same. 

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

2277 with self.assertRaises(ValueError): 

2278 butler.exists(ref_bad) 

2279 

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

2281 ref_compat = ref2.overrideStorageClass("StructuredDataDict") 

2282 exists = butler.exists(ref_compat) 

2283 self.assertEqual(exists, exists_many[ref2]) 

2284 

2285 # Remove everything and start from scratch. 

2286 butler._datastore.trustGetRequest = False 

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

2288 for ref in refs: 

2289 butler.put(metric, ref) 

2290 

2291 # These tests mess directly with the trash table and can leave the 

2292 # datastore in an odd state. Do them at the end. 

2293 # Check that in normal mode, deleting the record will lead to 

2294 # trash not touching the file. 

2295 uri1 = butler.getURI(ref1) 

2296 butler._datastore.bridge.moveToTrash( 

2297 [ref1], transaction=None 

2298 ) # Update the dataset_location table 

2299 butler._datastore.forget([ref1]) 

2300 butler._datastore.trash(ref1) 

2301 butler._datastore.emptyTrash() 

2302 self.assertTrue(uri1.exists()) 

2303 uri1.remove() # Clean it up. 

2304 

2305 # Simulate execution butler setup by deleting the datastore 

2306 # record but keeping the file around and trusting. 

2307 butler._datastore.trustGetRequest = True 

2308 uris = butler.get_many_uris([ref2, ref3]) 

2309 uri2 = uris[ref2].primaryURI 

2310 uri3 = uris[ref3].primaryURI 

2311 self.assertTrue(uri2.exists()) 

2312 self.assertTrue(uri3.exists()) 

2313 

2314 # Remove the datastore record. 

2315 butler._datastore.bridge.moveToTrash( 

2316 [ref2], transaction=None 

2317 ) # Update the dataset_location table 

2318 butler._datastore.forget([ref2]) 

2319 self.assertTrue(uri2.exists()) 

2320 butler._datastore.trash([ref2, ref3]) 

2321 # Immediate removal for ref2 file 

2322 self.assertFalse(uri2.exists()) 

2323 # But ref3 has to wait for the empty. 

2324 self.assertTrue(uri3.exists()) 

2325 butler._datastore.emptyTrash() 

2326 self.assertFalse(uri3.exists()) 

2327 

2328 # Clear out the datasets from registry. 

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

2330 

2331 def test_butler_metrics(self): 

2332 """Test that metrics are collected.""" 

2333 run = "test_run" 

2334 metrics = ButlerMetrics() 

2335 butler, datasetType = self.create_butler( 

2336 run, "MetricsExampleModelProvenance", "prov_metric", metrics=metrics 

2337 ) 

2338 data = MetricsExampleModel( 

2339 summary={"AM1": 5.2, "AM2": 30.6}, 

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

2341 data=[563, 234, 456.7, 752, 8, 9, 27], 

2342 ) 

2343 

2344 data_ref = butler.put(data, datasetType, visit=424, instrument="DummyCamComp") 

2345 butler.get(data_ref) 

2346 butler.get(data_ref) 

2347 self.assertEqual(metrics.n_get, 2) 

2348 self.assertGreater(metrics.time_in_get, 0.0) 

2349 self.assertEqual(metrics.n_put, 1) 

2350 self.assertGreater(metrics.time_in_put, 0.0) 

2351 

2352 deferred = butler.getDeferred(data_ref) 

2353 deferred.get() 

2354 self.assertEqual(metrics.n_get, 3) 

2355 

2356 with butler.record_metrics() as new: 

2357 data_ref_2 = butler.put(data, datasetType, visit=425, instrument="DummyCamComp") 

2358 butler.get(data_ref) 

2359 

2360 butler.pruneDatasets([data_ref, data_ref_2], purge=True, unstore=True) 

2361 with ResourcePath.temporary_uri(suffix=".json") as tmpFile: 

2362 tmpFile.write(data.model_dump_json().encode()) 

2363 refs = [ 

2364 DatasetRef(datasetType, data_ref_2.dataId, run), 

2365 DatasetRef(datasetType, data_ref.dataId, run), 

2366 ] 

2367 datasets = [FileDataset(path=tmpFile, refs=refs)] 

2368 butler.ingest(*datasets, transfer="copy") 

2369 

2370 self.assertEqual(new.n_get, 1) 

2371 self.assertEqual(new.n_put, 1) 

2372 self.assertEqual(new.n_ingest, 2) 

2373 

2374 

2375class PosixDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2376 """PosixDatastore specialization of a butler""" 

2377 

2378 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2379 fullConfigKey: str | None = ".datastore.formatters" 

2380 validationCanFail = True 

2381 datastoreStr = ["/tmp"] 

2382 datastoreName = [f"FileDatastore@{BUTLER_ROOT_TAG}"] 

2383 registryStr = "/gen3.sqlite3" 

2384 

2385 def testPathConstructor(self) -> None: 

2386 """Independent test of constructor using PathLike.""" 

2387 butler = Butler.from_config(self.tmpConfigFile, run=self.default_run) 

2388 self.enterContext(butler) 

2389 self.assertIsInstance(butler, Butler) 

2390 

2391 # And again with a Path object with the butler yaml 

2392 path = pathlib.Path(self.tmpConfigFile) 

2393 butler = Butler.from_config(path, writeable=False) 

2394 self.enterContext(butler) 

2395 self.assertIsInstance(butler, Butler) 

2396 

2397 # And again with a Path object without the butler yaml 

2398 # (making sure we skip it if the tmp config doesn't end 

2399 # in butler.yaml -- which is the case for a subclass) 

2400 if self.tmpConfigFile.endswith("butler.yaml"): 

2401 path = pathlib.Path(os.path.dirname(self.tmpConfigFile)) 

2402 butler = Butler.from_config(path, writeable=False) 

2403 self.enterContext(butler) 

2404 self.assertIsInstance(butler, Butler) 

2405 

2406 def testExportTransferCopy(self) -> None: 

2407 """Test local export using all transfer modes""" 

2408 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2409 exportButler = self.runPutGetTest(storageClass, "test_metric") 

2410 # Test that the repo actually has at least one dataset. 

2411 datasets = list(exportButler.registry.queryDatasets(..., collections=...)) 

2412 self.assertGreater(len(datasets), 0) 

2413 uris = [exportButler.getURI(d) for d in datasets] 

2414 assert isinstance(exportButler._datastore, FileDatastore) 

2415 datastoreRoot = exportButler.get_datastore_roots()[exportButler.get_datastore_names()[0]] 

2416 

2417 pathsInStore = [uri.relative_to(datastoreRoot) for uri in uris] 

2418 

2419 for path in pathsInStore: 

2420 # Assume local file system 

2421 assert path is not None 

2422 self.assertTrue(self.checkFileExists(datastoreRoot, path), f"Checking path {path}") 

2423 

2424 for transfer in ("copy", "link", "symlink", "relsymlink"): 

2425 with safeTestTempDir(TESTDIR) as exportDir: 

2426 with exportButler.export(directory=exportDir, format="yaml", transfer=transfer) as export: 

2427 export.saveDatasets(datasets) 

2428 for path in pathsInStore: 

2429 assert path is not None 

2430 self.assertTrue( 

2431 self.checkFileExists(exportDir, path), 

2432 f"Check that mode {transfer} exported files", 

2433 ) 

2434 

2435 def testPytypeCoercion(self) -> None: 

2436 """Test python type coercion on Butler.get and put.""" 

2437 # Store some data with the normal example storage class. 

2438 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2439 datasetTypeName = "test_metric" 

2440 butler = self.runPutGetTest(storageClass, datasetTypeName) 

2441 

2442 dataId = {"instrument": "DummyCamComp", "visit": 423} 

2443 metric = butler.get(datasetTypeName, dataId=dataId) 

2444 self.assertEqual(get_full_type_name(metric), "lsst.daf.butler.tests.MetricsExample") 

2445 

2446 datasetType_ori = butler.get_dataset_type(datasetTypeName) 

2447 self.assertEqual(datasetType_ori.storageClass.name, "StructuredDataNoComponents") 

2448 

2449 # Now need to hack the registry dataset type definition. 

2450 # There is no API for this. 

2451 assert isinstance(butler._registry, SqlRegistry) 

2452 manager = butler._registry._managers.datasets 

2453 assert hasattr(manager, "_db") and hasattr(manager, "_static") 

2454 manager._db.update( 

2455 manager._static.dataset_type, 

2456 {"name": datasetTypeName}, 

2457 {datasetTypeName: datasetTypeName, "storage_class": "StructuredDataNoComponentsModel"}, 

2458 ) 

2459 

2460 # Force reset of dataset type cache 

2461 butler.registry.refresh() 

2462 

2463 datasetType_new = butler.get_dataset_type(datasetTypeName) 

2464 self.assertEqual(datasetType_new.name, datasetType_ori.name) 

2465 self.assertEqual(datasetType_new.storageClass.name, "StructuredDataNoComponentsModel") 

2466 

2467 metric_model = butler.get(datasetTypeName, dataId=dataId) 

2468 self.assertNotEqual(type(metric_model), type(metric)) 

2469 self.assertEqual(get_full_type_name(metric_model), "lsst.daf.butler.tests.MetricsExampleModel") 

2470 

2471 # Put the model and read it back to show that everything now 

2472 # works as normal. 

2473 metric_ref = butler.put(metric_model, datasetTypeName, dataId=dataId, visit=424) 

2474 metric_model_new = butler.get(metric_ref) 

2475 self.assertEqual(metric_model_new, metric_model) 

2476 

2477 # Hack the storage class again to something that will fail on the 

2478 # get with no conversion class. 

2479 manager._db.update( 

2480 manager._static.dataset_type, 

2481 {"name": datasetTypeName}, 

2482 {datasetTypeName: datasetTypeName, "storage_class": "StructuredDataListYaml"}, 

2483 ) 

2484 butler.registry.refresh() 

2485 

2486 with self.assertRaises(ValueError): 

2487 butler.get(datasetTypeName, dataId=dataId) 

2488 

2489 def test_provenance(self): 

2490 """Test that provenance is attached on put.""" 

2491 run = "test_run" 

2492 butler, datasetType = self.create_butler(run, "MetricsExampleModelProvenance", "prov_metric") 

2493 metric = MetricsExampleModel( 

2494 summary={"AM1": 5.2, "AM2": 30.6}, 

2495 output={"a": [1, 2, 3], "b": {"blue": 5, "red": "green"}}, 

2496 data=[563, 234, 456.7, 752, 8, 9, 27], 

2497 ) 

2498 # Provenance can be attached to the object being put. Whether 

2499 # it is or not is dependent on the formatter. For this test we 

2500 # copy on adding provenance to ensure they differ. 

2501 self.assertIsNone(metric.dataset_id) 

2502 metric_ref = butler.put(metric, datasetType, visit=424, instrument="DummyCamComp") 

2503 self.assertIsNone(metric.dataset_id) 

2504 metric_2 = butler.get(metric_ref) 

2505 self.assertEqual(metric_2.data, metric.data) 

2506 self.assertEqual(metric_2.dataset_id, metric_ref.id) 

2507 self.assertIsNone(metric_2.provenance) 

2508 

2509 # Put with provenance. 

2510 prov = DatasetProvenance(quantum_id=uuid.uuid4()) 

2511 prov.add_input(metric_ref) 

2512 prov.add_extra_provenance(metric_ref.id, {"answer": 42}) 

2513 metric_ref2 = butler.put(metric, datasetType, visit=423, instrument="DummyCamComp", provenance=prov) 

2514 metric_3 = butler.get(metric_ref2) 

2515 self.assertEqual(metric_3.provenance, prov) 

2516 

2517 # Check that we can extract provenance from dict form. 

2518 prov_dict = prov.to_flat_dict(metric_ref2) 

2519 prov_from_prov, ref_from_prov = DatasetProvenance.from_flat_dict(prov_dict, butler) 

2520 self.assertEqual(ref_from_prov, metric_ref2) 

2521 # Direct __eq__ of the provenance does not work because one side 

2522 # includes dimension records. 

2523 self.assertEqual({ref.id for ref in prov_from_prov.inputs}, {ref.id for ref in prov.inputs}) 

2524 self.assertEqual(prov_from_prov.quantum_id, prov.quantum_id) 

2525 self.assertEqual(prov_from_prov.extras, prov.extras) 

2526 

2527 # Force a bad ID into the dict. 

2528 prov_dict["id"] = uuid.uuid4() 

2529 with self.assertRaises(ValueError): 

2530 DatasetProvenance.from_flat_dict(prov_dict, butler) 

2531 del prov_dict["id"] 

2532 prov_dict["input 0 id"] = uuid.uuid4() 

2533 with self.assertRaises(ValueError): 

2534 DatasetProvenance.from_flat_dict(prov_dict, butler) 

2535 

2536 # Check that simple types can be reconstructed with non-standard 

2537 # separators. 

2538 prov_dict = prov.to_flat_dict(metric_ref2, prefix="XYZ", sep="😎", simple_types=True) 

2539 prov_from_prov, ref_from_prov = DatasetProvenance.from_flat_dict(prov_dict, butler) 

2540 self.assertEqual(ref_from_prov, metric_ref2) 

2541 self.assertEqual({ref.id for ref in prov_from_prov.inputs}, {ref.id for ref in prov.inputs}) 

2542 

2543 with self.assertRaises(ValueError): 

2544 DatasetProvenance.from_flat_dict({"unknown": 42}, butler) 

2545 

2546 def test_specialized_file_datasets_functions(self): 

2547 """Test a workflow used in Prompt Processing where we export datasets 

2548 from one repository and write them in-place to the datastore of 

2549 another, without immediately inserting registry entries for the 

2550 datasets. 

2551 """ 

2552 repo = MetricTestRepo.create_from_butler( 

2553 self.create_empty_butler(writeable=True), 

2554 self.tmpConfigFile, 

2555 "StructuredCompositeReadCompNoDisassembly", 

2556 ) 

2557 source_butler = repo.butler 

2558 

2559 # Test writing outputs to a FileDatastore. 

2560 with tempfile.TemporaryDirectory() as tempdir: 

2561 target_repo_config = Butler.makeRepo(tempdir) 

2562 refs = [repo.ref1, repo.ref2] 

2563 datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), refs) 

2564 self.assertEqual(len(datasets), 2) 

2565 self.assertEqual({ref.id for ref in refs}, {dataset.refs[0].id for dataset in datasets}) 

2566 for dataset in datasets: 

2567 path = ResourcePath(dataset.path, forceAbsolute=False) 

2568 # Paths should be relative paths to the target datastore. 

2569 self.assertFalse(path.isabs()) 

2570 # Files should have been copied into the target datastore 

2571 self.assertTrue(ResourcePath(tempdir).join(path).exists()) 

2572 

2573 # Make sure the target Butler can ingest the datasets. 

2574 target_butler = Butler(target_repo_config, writeable=True) 

2575 self.enterContext(target_butler) 

2576 target_butler.transfer_dimension_records_from(source_butler, refs) 

2577 target_butler.ingest(*datasets, transfer=None) 

2578 self.assertIsNotNone(target_butler.get(repo.ref1)) 

2579 self.assertIsNotNone(target_butler.get(repo.ref2)) 

2580 

2581 # Giving an empty list of files is a no-op. 

2582 no_datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), []) 

2583 self.assertEqual(len(no_datasets), 0) 

2584 

2585 # Test writing outputs to a ChainedDatastore. 

2586 with tempfile.TemporaryDirectory() as tempdir: 

2587 # Set up a second dataset type, so we can split the files across 

2588 # multiple datastore roots. 

2589 dt1 = repo.datasetType 

2590 dt2 = DatasetType("other", dt1.dimensions, dt1.storageClass) 

2591 source_butler.registry.registerDatasetType(dt2) 

2592 other_ref = repo.addDataset(repo.ref1.dataId, datasetType=dt2) 

2593 config = Config.fromString( 

2594 f""" 

2595 datastore: 

2596 cls: lsst.daf.butler.datastores.chainedDatastore.ChainedDatastore 

2597 datastore_constraints: 

2598 - constraints: 

2599 accept: 

2600 - {dt1.name} 

2601 - constraints: 

2602 accept: 

2603 - {dt2.name} 

2604 datastores: 

2605 - datastore: 

2606 cls: lsst.daf.butler.datastores.fileDatastore.FileDatastore 

2607 root: <butlerRoot>/FileDatastore_0 

2608 - datastore: 

2609 cls: lsst.daf.butler.datastores.fileDatastore.FileDatastore 

2610 root: <butlerRoot>/FileDatastore_1 

2611 """ 

2612 ) 

2613 target_repo_config = Butler.makeRepo(tempdir, config) 

2614 refs = [repo.ref1, repo.ref2, other_ref] 

2615 datasets = transfer_datasets_to_datastore(source_butler, ButlerConfig(target_repo_config), refs) 

2616 self.assertEqual(len(datasets), 3) 

2617 self.assertEqual({ref.id for ref in refs}, {dataset.refs[0].id for dataset in datasets}) 

2618 for dataset in datasets: 

2619 path = ResourcePath(dataset.path, forceAbsolute=False) 

2620 # Paths should be relative paths to the target datastore. 

2621 self.assertFalse(path.isabs()) 

2622 # Files should have been split up between the two datastores 

2623 # in the chain. 

2624 datastore_root = ResourcePath(tempdir) 

2625 if dataset.refs[0].datasetType.name == dt1.name: 

2626 datastore_root = datastore_root.join("FileDatastore_0") 

2627 else: 

2628 datastore_root = datastore_root.join("FileDatastore_1") 

2629 self.assertTrue(datastore_root.join(path).exists()) 

2630 

2631 # Make sure the target Butler can ingest the datasets. 

2632 target_butler = Butler(target_repo_config, writeable=True) 

2633 self.enterContext(target_butler) 

2634 target_butler.transfer_dimension_records_from(source_butler, refs) 

2635 target_butler.ingest(*datasets, transfer=None) 

2636 self.assertIsNotNone(target_butler.get(repo.ref1)) 

2637 self.assertIsNotNone(target_butler.get(repo.ref2)) 

2638 self.assertIsNotNone(target_butler.get(other_ref)) 

2639 

2640 def test_temporary_for_ingest(self) -> None: 

2641 """Test the `lsst.daf.butler._rubin.ingest_from_temporary` module.""" 

2642 with self.create_empty_butler("example_run") as butler: 

2643 dataset_type = DatasetType("example", butler.dimensions.empty, "StructuredDataDict") 

2644 butler.registry.registerDatasetType(dataset_type) 

2645 ref = DatasetRef(dataset_type, DataCoordinate.make_empty(butler.dimensions), "example_run") 

2646 with TemporaryForIngest(butler, ref) as temporary: 

2647 temporary.path.write(b"three: 3") 

2648 found = TemporaryForIngest.find_orphaned_temporaries_by_ref(ref, butler) 

2649 self.assertEqual(found, [temporary.path]) 

2650 self.assertIn(".tmp", temporary.ospath) 

2651 temporary.ingest() 

2652 loaded = butler.get(ref) 

2653 self.assertEqual(loaded, {"three": 3}) 

2654 

2655 

2656class PostgresPosixDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2657 """PosixDatastore specialization of a butler using Postgres""" 

2658 

2659 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2660 fullConfigKey = ".datastore.formatters" 

2661 validationCanFail = True 

2662 datastoreStr = ["/tmp"] 

2663 datastoreName = [f"FileDatastore@{BUTLER_ROOT_TAG}"] 

2664 registryStr = "PostgreSQL@test" 

2665 

2666 @classmethod 

2667 def setUpClass(cls) -> None: 

2668 cls.postgresql = cls.enterClassContext(setup_postgres_test_db()) 

2669 super().setUpClass() 

2670 

2671 def setUp(self) -> None: 

2672 # Need to add a registry section to the config. 

2673 self._temp_config = False 

2674 config = Config(self.configFile) 

2675 self.postgresql.patch_butler_config(config) 

2676 with tempfile.NamedTemporaryFile("w", suffix=".yaml", delete=False) as fh: 

2677 config.dump(fh) 

2678 self.configFile = fh.name 

2679 self._temp_config = True 

2680 super().setUp() 

2681 

2682 def tearDown(self) -> None: 

2683 if self._temp_config and os.path.exists(self.configFile): 

2684 os.remove(self.configFile) 

2685 super().tearDown() 

2686 

2687 def testMakeRepo(self) -> None: 

2688 # The base class test assumes that it's using sqlite and assumes 

2689 # the config file is acceptable to sqlite. 

2690 raise unittest.SkipTest("Postgres config is not compatible with this test.") 

2691 

2692 

2693class ClonedPostgresPosixDatastoreButlerTestCase(PostgresPosixDatastoreButlerTestCase, unittest.TestCase): 

2694 """Test that Butler with a Postgres registry still works after cloning.""" 

2695 

2696 def create_butler( 

2697 self, 

2698 run: str, 

2699 storageClass: StorageClass | str, 

2700 datasetTypeName: str, 

2701 metrics: ButlerMetrics | None = None, 

2702 ) -> tuple[DirectButler, DatasetType]: 

2703 butler, datasetType = super().create_butler(run, storageClass, datasetTypeName, metrics=metrics) 

2704 return butler.clone(run=run, metrics=metrics), datasetType 

2705 

2706 

2707class InMemoryDatastoreButlerTestCase(ButlerTests, unittest.TestCase): 

2708 """InMemoryDatastore specialization of a butler""" 

2709 

2710 configFile = os.path.join(TESTDIR, "config/basic/butler-inmemory.yaml") 

2711 fullConfigKey = None 

2712 useTempRoot = False 

2713 validationCanFail = False 

2714 datastoreStr = ["datastore='InMemory"] 

2715 datastoreName = ["InMemoryDatastore@"] 

2716 registryStr = "/gen3.sqlite3" 

2717 

2718 def testIngest(self) -> None: 

2719 pass 

2720 

2721 def test_ingest_zip(self) -> None: 

2722 pass 

2723 

2724 

2725class ClonedSqliteButlerTestCase(InMemoryDatastoreButlerTestCase, unittest.TestCase): 

2726 """Test that a Butler with a Sqlite registry still works after cloning.""" 

2727 

2728 def create_butler( 

2729 self, 

2730 run: str, 

2731 storageClass: StorageClass | str, 

2732 datasetTypeName: str, 

2733 metrics: ButlerMetrics | None = None, 

2734 ) -> tuple[DirectButler, DatasetType]: 

2735 butler, datasetType = super().create_butler(run, storageClass, datasetTypeName, metrics=metrics) 

2736 return butler.clone(run=run), datasetType 

2737 

2738 

2739class ChainedDatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2740 """PosixDatastore specialization""" 

2741 

2742 configFile = os.path.join(TESTDIR, "config/basic/butler-chained.yaml") 

2743 fullConfigKey = ".datastore.datastores.1.formatters" 

2744 validationCanFail = True 

2745 datastoreStr = ["datastore='InMemory", "/FileDatastore_1/,", "/FileDatastore_2/'"] 

2746 datastoreName = [ 

2747 "InMemoryDatastore@", 

2748 f"FileDatastore@{BUTLER_ROOT_TAG}/FileDatastore_1", 

2749 "SecondDatastore", 

2750 ] 

2751 registryStr = "/gen3.sqlite3" 

2752 

2753 def testPruneDatasets(self) -> None: 

2754 # This test relies on manipulating files out-of-band, which is 

2755 # impossible for this configuration because of the InMemoryDatastore in 

2756 # the ChainedDatastore. 

2757 pass 

2758 

2759 def testComponentFromOverriddenStorageClassWarns(self) -> None: 

2760 # The InMemoryDatastore in the ChainedDatastore satisfies the get, so 

2761 # the FileDatastore warning about having to read the whole dataset to 

2762 # extract the component is never issued. 

2763 pass 

2764 

2765 

2766class ButlerExplicitRootTestCase(PosixDatastoreButlerTestCase): 

2767 """Test that a yaml file in one location can refer to a root in another.""" 

2768 

2769 datastoreStr = ["dir1"] 

2770 # Disable the makeRepo test since we are deliberately not using 

2771 # butler.yaml as the config name. 

2772 fullConfigKey = None 

2773 

2774 def setUp(self) -> None: 

2775 self.root = makeTestTempDir(TESTDIR) 

2776 

2777 # Make a new repository in one place 

2778 self.dir1 = os.path.join(self.root, "dir1") 

2779 Butler.makeRepo(self.dir1, config=Config(self.configFile)) 

2780 

2781 # Move the yaml file to a different place and add a "root" 

2782 self.dir2 = os.path.join(self.root, "dir2") 

2783 os.makedirs(self.dir2, exist_ok=True) 

2784 configFile1 = os.path.join(self.dir1, "butler.yaml") 

2785 config = Config(configFile1) 

2786 config["root"] = self.dir1 

2787 configFile2 = os.path.join(self.dir2, "butler2.yaml") 

2788 config.dumpToUri(configFile2) 

2789 os.remove(configFile1) 

2790 self.tmpConfigFile = configFile2 

2791 

2792 def testFileLocations(self) -> None: 

2793 self.assertNotEqual(self.dir1, self.dir2) 

2794 self.assertTrue(os.path.exists(os.path.join(self.dir2, "butler2.yaml"))) 

2795 self.assertFalse(os.path.exists(os.path.join(self.dir1, "butler.yaml"))) 

2796 self.assertTrue(os.path.exists(os.path.join(self.dir1, "gen3.sqlite3"))) 

2797 

2798 

2799class ButlerMakeRepoOutfileTestCase(ButlerPutGetTests, unittest.TestCase): 

2800 """Test that a config file created by makeRepo outside of repo works.""" 

2801 

2802 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2803 

2804 def setUp(self) -> None: 

2805 self.root = makeTestTempDir(TESTDIR) 

2806 self.root2 = makeTestTempDir(TESTDIR) 

2807 

2808 self.tmpConfigFile = os.path.join(self.root2, "different.yaml") 

2809 Butler.makeRepo(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile) 

2810 

2811 def tearDown(self) -> None: 

2812 if os.path.exists(self.root2): 2812 ↛ 2814line 2812 didn't jump to line 2814 because the condition on line 2812 was always true

2813 shutil.rmtree(self.root2, ignore_errors=True) 

2814 super().tearDown() 

2815 

2816 def testConfigExistence(self) -> None: 

2817 c = Config(self.tmpConfigFile) 

2818 uri_config = ResourcePath(c["root"]) 

2819 uri_expected = ResourcePath(self.root, forceDirectory=True) 

2820 self.assertEqual(uri_config.geturl(), uri_expected.geturl()) 

2821 self.assertNotIn(":", uri_config.path, "Check for URI concatenated with normal path") 

2822 

2823 def testPutGet(self) -> None: 

2824 storageClass = self.storageClassFactory.getStorageClass("StructuredDataNoComponents") 

2825 self.runPutGetTest(storageClass, "test_metric") 

2826 

2827 

2828class ButlerMakeRepoOutfileDirTestCase(ButlerMakeRepoOutfileTestCase): 

2829 """Test that a config file created by makeRepo outside of repo works.""" 

2830 

2831 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2832 

2833 def setUp(self) -> None: 

2834 self.root = makeTestTempDir(TESTDIR) 

2835 self.root2 = makeTestTempDir(TESTDIR) 

2836 

2837 self.tmpConfigFile = self.root2 

2838 Butler.makeRepo(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile) 

2839 

2840 def testConfigExistence(self) -> None: 

2841 # Append the yaml file else Config constructor does not know the file 

2842 # type. 

2843 self.tmpConfigFile = os.path.join(self.tmpConfigFile, "butler.yaml") 

2844 super().testConfigExistence() 

2845 

2846 

2847class ButlerMakeRepoOutfileUriTestCase(ButlerMakeRepoOutfileTestCase): 

2848 """Test that a config file created by makeRepo outside of repo works.""" 

2849 

2850 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

2851 

2852 def setUp(self) -> None: 

2853 self.root = makeTestTempDir(TESTDIR) 

2854 self.root2 = makeTestTempDir(TESTDIR) 

2855 

2856 self.tmpConfigFile = ResourcePath(os.path.join(self.root2, "something.yaml")).geturl() 

2857 Butler.makeRepo(self.root, config=Config(self.configFile), outfile=self.tmpConfigFile) 

2858 

2859 

2860@unittest.skipIf(not boto3, "Warning: boto3 AWS SDK not found!") 

2861class S3DatastoreButlerTestCase(FileDatastoreButlerTests, unittest.TestCase): 

2862 """S3Datastore specialization of a butler; an S3 storage Datastore + 

2863 a local in-memory SqlRegistry. 

2864 """ 

2865 

2866 configFile = os.path.join(TESTDIR, "config/basic/butler-s3store.yaml") 

2867 fullConfigKey = None 

2868 validationCanFail = True 

2869 

2870 bucketName = "anybucketname" 

2871 """Name of the Bucket that will be used in the tests. The name is read from 

2872 the config file used with the tests during set-up. 

2873 """ 

2874 

2875 root = "butlerRoot/" 

2876 """Root repository directory expected to be used in case useTempRoot=False. 

2877 Otherwise the root is set to a 20 characters long randomly generated string 

2878 during set-up. 

2879 """ 

2880 

2881 datastoreStr = [f"datastore={root}"] 

2882 """Contains all expected root locations in a format expected to be 

2883 returned by Butler stringification. 

2884 """ 

2885 

2886 datastoreName = ["FileDatastore@s3://{bucketName}/{root}"] 

2887 """The expected format of the S3 Datastore string.""" 

2888 

2889 registryStr = "/gen3.sqlite3" 

2890 """Expected format of the Registry string.""" 

2891 

2892 mock_aws = mock_aws() 

2893 """The mocked s3 interface from moto.""" 

2894 

2895 def genRoot(self) -> str: 

2896 """Return a random string of len 20 to serve as a root 

2897 name for the temporary bucket repo. 

2898 

2899 This is equivalent to tempfile.mkdtemp as this is what self.root 

2900 becomes when useTempRoot is True. 

2901 """ 

2902 rndstr = "".join(random.choice(string.ascii_uppercase + string.digits) for _ in range(20)) 

2903 return rndstr + "/" 

2904 

2905 def setUp(self) -> None: 

2906 config = Config(self.configFile) 

2907 uri = ResourcePath(config[".datastore.datastore.root"]) 

2908 self.bucketName = uri.netloc 

2909 

2910 # Enable S3 mocking of tests. 

2911 self.enterContext(clean_test_environment_for_s3()) 

2912 self.mock_aws.start() 

2913 

2914 if self.useTempRoot: 2914 ↛ 2916line 2914 didn't jump to line 2916 because the condition on line 2914 was always true

2915 self.root = self.genRoot() 

2916 rooturi = f"s3://{self.bucketName}/{self.root}" 

2917 config.update({"datastore": {"datastore": {"root": rooturi}}}) 

2918 

2919 # need local folder to store registry database 

2920 self.reg_dir = makeTestTempDir(TESTDIR) 

2921 config["registry", "db"] = f"sqlite:///{self.reg_dir}/gen3.sqlite3" 

2922 

2923 # MOTO needs to know that we expect Bucket bucketname to exist 

2924 # (this used to be the class attribute bucketName) 

2925 s3 = boto3.resource("s3") 

2926 s3.create_bucket(Bucket=self.bucketName) 

2927 

2928 self.datastoreStr = [f"datastore='{rooturi}'"] 

2929 self.datastoreName = [f"FileDatastore@{rooturi}"] 

2930 Butler.makeRepo(rooturi, config=config, forceConfigRoot=False) 

2931 self.tmpConfigFile = posixpath.join(rooturi, "butler.yaml") 

2932 

2933 def tearDown(self) -> None: 

2934 s3 = boto3.resource("s3") 

2935 bucket = s3.Bucket(self.bucketName) 

2936 try: 

2937 bucket.objects.all().delete() 

2938 except botocore.exceptions.ClientError as e: 

2939 if e.response["Error"]["Code"] == "404": 

2940 # the key was not reachable - pass 

2941 pass 

2942 else: 

2943 raise 

2944 

2945 bucket = s3.Bucket(self.bucketName) 

2946 bucket.delete() 

2947 

2948 # Stop the S3 mock. 

2949 self.mock_aws.stop() 

2950 

2951 if self.reg_dir is not None and os.path.exists(self.reg_dir): 2951 ↛ 2954line 2951 didn't jump to line 2954 because the condition on line 2951 was always true

2952 shutil.rmtree(self.reg_dir, ignore_errors=True) 

2953 

2954 if self.useTempRoot and os.path.exists(self.root): 2954 ↛ 2955line 2954 didn't jump to line 2955 because the condition on line 2954 was never true

2955 shutil.rmtree(self.root, ignore_errors=True) 

2956 

2957 super().tearDown() 

2958 

2959 

2960class DatastoreTransfers(TestCaseMixin): 

2961 """Base test setup for data transfers between butlers. The concrete tests 

2962 for specific configurations are in other classes, below. 

2963 """ 

2964 

2965 storageClassFactory: StorageClassFactory 

2966 

2967 @classmethod 

2968 def setUpClass(cls) -> None: 

2969 cls.storageClassFactory = StorageClassFactory() 

2970 

2971 def setUp(self) -> None: 

2972 self.root = makeTestTempDir(TESTDIR) 

2973 self.config = Config(self.configFile) 

2974 

2975 # Some tests cause convertors to be replaced so ensure 

2976 # the storage class factory is reset each time. 

2977 self.storageClassFactory.reset() 

2978 self.storageClassFactory.addFromConfig(self.configFile) 

2979 

2980 def tearDown(self) -> None: 

2981 removeTestTempDir(self.root) 

2982 

2983 def create_butler(self, manager: str | None, label: str, config_file: str | None = None) -> Butler: 

2984 if manager is None: 2984 ↛ 2988line 2984 didn't jump to line 2988 because the condition on line 2984 was always true

2985 manager = ( 

2986 "lsst.daf.butler.registry.datasets.byDimensions.ByDimensionsDatasetRecordStorageManagerUUID" 

2987 ) 

2988 config = Config(config_file if config_file is not None else self.configFile) 

2989 config["registry", "managers", "datasets"] = manager 

2990 butler = Butler.from_config( 

2991 Butler.makeRepo(f"{self.root}/butler{label}", config=config), writeable=True 

2992 ) 

2993 self.enterContext(butler) 

2994 return butler 

2995 

2996 def assertButlerTransfers( 

2997 self, 

2998 purge: bool = False, 

2999 storageClassName: str = "StructuredData", 

3000 storageClassNameTarget: str | None = None, 

3001 ) -> None: 

3002 """Test that a run can be transferred to another butler.""" 

3003 storageClass = self.storageClassFactory.getStorageClass(storageClassName) 

3004 if storageClassNameTarget is not None: 

3005 storageClassTarget = self.storageClassFactory.getStorageClass(storageClassNameTarget) 

3006 else: 

3007 storageClassTarget = storageClass 

3008 

3009 datasetTypeName = "random_data" 

3010 

3011 # Test will create 3 collections and we will want to transfer 

3012 # two of those three. 

3013 runs = ["run1", "run2", "other"] 

3014 

3015 # Also want to use two different dataset types to ensure that 

3016 # grouping works. 

3017 datasetTypeNames = ["random_data", "random_data_2"] 

3018 

3019 # Create the run collections in the source butler. 

3020 for run in runs: 

3021 self.source_butler.collections.register(run) 

3022 

3023 # Create dimensions in source butler. 

3024 n_exposures = 30 

3025 self.source_butler.registry.insertDimensionData("instrument", {"name": "DummyCamComp"}) 

3026 self.source_butler.registry.insertDimensionData( 

3027 "physical_filter", {"instrument": "DummyCamComp", "name": "d-r", "band": "R"} 

3028 ) 

3029 self.source_butler.registry.insertDimensionData( 

3030 "detector", {"instrument": "DummyCamComp", "id": 1, "full_name": "det1"} 

3031 ) 

3032 self.source_butler.registry.insertDimensionData( 

3033 "day_obs", 

3034 { 

3035 "instrument": "DummyCamComp", 

3036 "id": 20250101, 

3037 }, 

3038 ) 

3039 

3040 for i in range(n_exposures): 

3041 self.source_butler.registry.insertDimensionData( 

3042 "group", {"instrument": "DummyCamComp", "name": f"group{i}"} 

3043 ) 

3044 self.source_butler.registry.insertDimensionData( 

3045 "exposure", 

3046 { 

3047 "instrument": "DummyCamComp", 

3048 "id": i, 

3049 "obs_id": f"exp{i}", 

3050 "physical_filter": "d-r", 

3051 "group": f"group{i}", 

3052 "day_obs": 20250101, 

3053 }, 

3054 ) 

3055 

3056 # Create dataset types in the source butler. 

3057 dimensions = self.source_butler.dimensions.conform(["instrument", "exposure"]) 

3058 for datasetTypeName in datasetTypeNames: 

3059 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

3060 self.source_butler.registry.registerDatasetType(datasetType) 

3061 

3062 # Write a dataset to an unrelated run -- this will ensure that 

3063 # we are rewriting integer dataset ids in the target if necessary. 

3064 # Will not be relevant for UUID. 

3065 run = "distraction" 

3066 butler = Butler.from_config(butler=self.source_butler, run=run) 

3067 self.enterContext(butler) 

3068 butler.put( 

3069 makeExampleMetrics(), 

3070 datasetTypeName, 

3071 exposure=1, 

3072 instrument="DummyCamComp", 

3073 physical_filter="d-r", 

3074 ) 

3075 

3076 # Write some example metrics to the source 

3077 butler = Butler.from_config(butler=self.source_butler) 

3078 self.enterContext(butler) 

3079 

3080 # Set of DatasetRefs that should be in the list of refs to transfer 

3081 # but which will not be transferred. 

3082 deleted: set[DatasetRef] = set() 

3083 

3084 n_expected = 20 # Number of datasets expected to be transferred 

3085 source_refs = [] 

3086 for i in range(n_exposures): 

3087 # Put a third of datasets into each collection, only retain 

3088 # two thirds. 

3089 index = i % 3 

3090 run = runs[index] 

3091 datasetTypeName = datasetTypeNames[i % 2] 

3092 

3093 metric = MetricsExample( 

3094 summary={"counter": i}, output={"text": "metric"}, data=[2 * x for x in range(i)] 

3095 ) 

3096 dataId = {"exposure": i, "instrument": "DummyCamComp", "physical_filter": "d-r"} 

3097 ref = butler.put(metric, datasetTypeName, dataId=dataId, run=run) 

3098 

3099 # Remove the datastore record using low-level API, but only 

3100 # for a specific index. 

3101 if purge and index == 1: 

3102 # For one of these delete the file as well. 

3103 # This allows the "missing" code to filter the 

3104 # file out. 

3105 # Access the individual datastores. 

3106 datastores = [] 

3107 if hasattr(butler._datastore, "datastores"): 

3108 datastores.extend(butler._datastore.datastores) 

3109 else: 

3110 datastores.append(butler._datastore) 

3111 

3112 if not deleted: 

3113 # For a chained datastore we need to remove 

3114 # files in each chain. 

3115 for datastore in datastores: 

3116 # The file might not be known to the datastore 

3117 # if constraints are used. 

3118 try: 

3119 primary, uris = datastore.getURIs(ref) 

3120 except FileNotFoundError: 

3121 continue 

3122 if primary and primary.scheme != "mem": 

3123 primary.remove() 

3124 for uri in uris.values(): 

3125 if uri.scheme != "mem": 3125 ↛ 3124line 3125 didn't jump to line 3124 because the condition on line 3125 was always true

3126 uri.remove() 

3127 n_expected -= 1 

3128 deleted.add(ref) 

3129 

3130 # Remove the datastore record. 

3131 for datastore in datastores: 

3132 if hasattr(datastore, "removeStoredItemInfo"): 3132 ↛ 3131line 3132 didn't jump to line 3131 because the condition on line 3132 was always true

3133 datastore.removeStoredItemInfo(ref) 

3134 

3135 if index < 2: 

3136 source_refs.append(ref) 

3137 if ref not in deleted: 

3138 new_metric = butler.get(ref) 

3139 self.assertEqual(new_metric, metric) 

3140 

3141 # Create some bad dataset types to ensure we check for inconsistent 

3142 # definitions. 

3143 badStorageClass = self.storageClassFactory.getStorageClass("StructuredDataList") 

3144 for datasetTypeName in datasetTypeNames: 

3145 datasetType = DatasetType(datasetTypeName, dimensions, badStorageClass) 

3146 self.target_butler.registry.registerDatasetType(datasetType) 

3147 with self.assertRaises(ConflictingDefinitionError) as cm: 

3148 self.target_butler.transfer_from(self.source_butler, source_refs) 

3149 self.assertIn("dataset type differs", str(cm.exception)) 

3150 

3151 # And remove the bad definitions. 

3152 for datasetTypeName in datasetTypeNames: 

3153 self.target_butler.registry.removeDatasetType(datasetTypeName) 

3154 

3155 # Transfer without creating dataset types should fail. 

3156 with self.assertRaises(KeyError): 

3157 self.target_butler.transfer_from(self.source_butler, source_refs) 

3158 

3159 # Transfer without creating dimensions should fail. 

3160 with self.assertRaises(ConflictingDefinitionError) as cm: 

3161 self.target_butler.transfer_from(self.source_butler, source_refs, register_dataset_types=True) 

3162 self.assertIn("dimension", str(cm.exception)) 

3163 

3164 # The dry run test requires dataset types to exist. If we have 

3165 # been given distinct storage classes for the target we have 

3166 # to redefine at least one of the dataset types in the target butler. 

3167 if storageClass != storageClassTarget: 

3168 self.target_butler.registry.removeDatasetType(datasetTypeNames[0]) 

3169 datasetType = DatasetType(datasetTypeNames[0], dimensions, storageClassTarget) 

3170 self.target_butler.registry.registerDatasetType(datasetType) 

3171 

3172 # The failed transfer above leaves registry in an inconsistent 

3173 # state because the run is created but then rolled back without 

3174 # the collection cache being cleared. For now force a refresh. 

3175 # Can remove with DM-35498. 

3176 self.target_butler.registry.refresh() 

3177 

3178 # Do a dry run -- this should not have any effect on the target butler. 

3179 self.target_butler.transfer_from(self.source_butler, source_refs, dry_run=True) 

3180 

3181 # Transfer the records for one ref to test the alternative API. 

3182 with self.assertLogs(logger="lsst", level=logging.DEBUG) as log_cm: 

3183 self.target_butler.transfer_dimension_records_from(self.source_butler, [source_refs[0]]) 

3184 self.assertIn("number of records transferred: 1", ";".join(log_cm.output)) 

3185 

3186 # Now transfer them to the second butler, including dimensions. 

3187 with self.assertLogs(logger="lsst", level=logging.DEBUG) as log_cm: 

3188 transferred = self.target_butler.transfer_from( 

3189 self.source_butler, 

3190 source_refs, 

3191 register_dataset_types=True, 

3192 transfer_dimensions=True, 

3193 ) 

3194 self.assertEqual(len(transferred), n_expected) 

3195 log_output = ";".join(log_cm.output) 

3196 

3197 # A ChainedDatastore will use the in-memory datastore for mexists 

3198 # so we can not rely on the mexists log message. 

3199 self.assertIn("Number of datastore records found in source", log_output) 

3200 self.assertIn("Creating output run", log_output) 

3201 

3202 # Do the transfer twice to ensure that it will do nothing extra. 

3203 # Only do this if purge=True because it does not work for int 

3204 # dataset_id. 

3205 if purge: 

3206 # This should not need to register dataset types. 

3207 transferred = self.target_butler.transfer_from(self.source_butler, source_refs) 

3208 self.assertEqual(len(transferred), n_expected) 

3209 

3210 with self.assertRaises((TypeError, AttributeError)): 

3211 self.target_butler._datastore.transfer_from(self.source_butler, source_refs) # type: ignore 

3212 

3213 with self.assertRaises(ValueError): 

3214 self.target_butler._datastore.transfer_from( 

3215 self.source_butler._datastore, source_refs, transfer="split" 

3216 ) 

3217 

3218 # Now try to get the same refs from the new butler. 

3219 for ref in source_refs: 

3220 if ref not in deleted: 

3221 new_metric = self.target_butler.get(ref) 

3222 old_metric = self.source_butler.get(ref) 

3223 self.assertEqual(new_metric, old_metric) 

3224 

3225 # Try again without implicit storage class conversion 

3226 # triggered by using the source ref. This will do conversion 

3227 # since the formatter will be returning the source python type. 

3228 target_ref = self.target_butler.get_dataset(ref.id) 

3229 if target_ref.datasetType.storageClass != ref.datasetType.storageClass: 

3230 new_metric = self.target_butler.get(target_ref) 

3231 self.assertNotEqual(type(new_metric), type(old_metric)) 

3232 

3233 # Remove the dataset from the target and put it again 

3234 # as if it was the right type all along for this butler. 

3235 self.target_butler.pruneDatasets( 

3236 [target_ref], unstore=True, purge=True, disassociate=True 

3237 ) 

3238 self.target_butler.put(new_metric, target_ref) 

3239 new_new_metric = self.target_butler.get(target_ref) 

3240 new_old_metric = self.target_butler.get( 

3241 target_ref, storageClass=ref.datasetType.storageClass 

3242 ) 

3243 self.assertEqual(new_new_metric, new_metric) 

3244 self.assertEqual(new_old_metric, old_metric) 

3245 

3246 # Now prune run2 collection and create instead a CHAINED collection. 

3247 # This should block the transfer. 

3248 self.target_butler.removeRuns(["run2"]) 

3249 self.target_butler.collections.register("run2", CollectionType.CHAINED) 

3250 with self.assertRaises(CollectionTypeError): 

3251 # Re-importing the run1 datasets can be problematic if they 

3252 # use integer IDs so filter those out. 

3253 to_transfer = [ref for ref in source_refs if ref.run == "run2"] 

3254 self.target_butler.transfer_from(self.source_butler, to_transfer) 

3255 

3256 

3257class PosixDatastoreTransfers(DatastoreTransfers, unittest.TestCase): 

3258 """Test data transfers between butlers. 

3259 

3260 Test for different managers. UUID to UUID and integer to integer are 

3261 tested. UUID to integer is not supported since we do not currently 

3262 want to allow that. Integer to UUID is supported with the caveat 

3263 that UUID4 will be generated and this will be incorrect for raw 

3264 dataset types. The test ignores that. 

3265 """ 

3266 

3267 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3268 

3269 def create_butlers( 

3270 self, manager1: str | None = None, manager2: str | None = None, source_config: str | None = None 

3271 ) -> None: 

3272 self.source_butler = self.create_butler(manager1, "1", config_file=source_config) 

3273 self.target_butler = self.create_butler(manager2, "2") 

3274 

3275 def testTransferUuidToUuid(self) -> None: 

3276 self.create_butlers() 

3277 self.assertButlerTransfers() 

3278 

3279 def testTransferFromChainedUuidToUuid(self) -> None: 

3280 """Force the source butler to be a ChainedDatastore.""" 

3281 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-chained.yaml")) 

3282 self.assertButlerTransfers() 

3283 

3284 def testTransferFromIncompatibleUuidToUuid(self) -> None: 

3285 """Force the source butler to be a incompatible datastore.""" 

3286 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-inmemory.yaml")) 

3287 with self.assertRaises(NotImplementedError): 

3288 self.assertButlerTransfers() 

3289 

3290 def testTransferFromIncompatibleChainUuidToUuid(self) -> None: 

3291 """Force the source butler to be a incompatible datastore.""" 

3292 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler-inmemory-chain.yaml")) 

3293 with self.assertRaises(TypeError): 

3294 self.assertButlerTransfers() 

3295 

3296 def testTransferFromFileUuidToUuid(self) -> None: 

3297 """Force the source butler to be a FileDatastore.""" 

3298 self.create_butlers(source_config=os.path.join(TESTDIR, "config/basic/butler.yaml")) 

3299 self.assertButlerTransfers() 

3300 

3301 def testTransferMissing(self) -> None: 

3302 """Test transfers where datastore records are missing. 

3303 

3304 This is how execution butler works. 

3305 """ 

3306 self.create_butlers() 

3307 

3308 # Configure the source butler to allow trust. 

3309 self.source_butler._datastore._set_trust_mode(True) 

3310 

3311 self.assertButlerTransfers(purge=True) 

3312 

3313 def testTransferMissingDisassembly(self) -> None: 

3314 """Test transfers where datastore records are missing. 

3315 

3316 This is how execution butler works. 

3317 """ 

3318 self.create_butlers() 

3319 

3320 # Configure the source butler to allow trust. 

3321 self.source_butler._datastore._set_trust_mode(True) 

3322 

3323 # Test disassembly. 

3324 self.assertButlerTransfers(purge=True, storageClassName="StructuredComposite") 

3325 

3326 def testTransferDifferingStorageClasses(self) -> None: 

3327 """Test transfers when the source butler dataset type has a different 

3328 but compatible storage class. 

3329 """ 

3330 self.create_butlers() 

3331 

3332 self.assertButlerTransfers(storageClassNameTarget="MetricsConversion") 

3333 

3334 def testTransferDifferingStorageClassesDisassembly(self) -> None: 

3335 """Test transfers when the source butler dataset type has a different 

3336 but compatible storage class and where the source butler has 

3337 disassembled. 

3338 """ 

3339 self.create_butlers() 

3340 

3341 self.assertButlerTransfers( 

3342 storageClassName="StructuredComposite", storageClassNameTarget="MetricsConversion" 

3343 ) 

3344 

3345 def testUnsafeDirectTransfer(self) -> None: 

3346 """Test that transfer='unsafe_direct' records the absolute URI of 

3347 source files in the target datastore. 

3348 """ 

3349 self.create_butlers() 

3350 dataset_type = DatasetType("dt", [], "int", universe=self.source_butler.dimensions) 

3351 self.source_butler.registry.registerDatasetType(dataset_type) 

3352 self.source_butler.collections.register("run") 

3353 ref = self.source_butler.put(123, "dt", [], run="run") 

3354 self.target_butler.transfer_from( 

3355 self.source_butler, [ref], transfer="unsafe_direct", register_dataset_types=True 

3356 ) 

3357 self.assertEqual(self.target_butler.get(ref), 123) 

3358 self.assertEqual(self.source_butler.getURI(ref), self.target_butler.getURI(ref)) 

3359 

3360 def testAbsoluteURITransferDirect(self) -> None: 

3361 """Test transfer using an absolute URI.""" 

3362 self._absolute_transfer("auto") 

3363 

3364 def testAbsoluteURITransferUnsafeDirect(self) -> None: 

3365 """Test transfer using an absolute URI.""" 

3366 self._absolute_transfer("unsafe_direct") 

3367 

3368 def testAbsoluteURITransferCopy(self) -> None: 

3369 """Test transfer using an absolute URI.""" 

3370 self._absolute_transfer("copy") 

3371 

3372 def _absolute_transfer(self, transfer: str) -> None: 

3373 self.create_butlers() 

3374 

3375 storageClassName = "StructuredData" 

3376 storageClass = self.storageClassFactory.getStorageClass(storageClassName) 

3377 datasetTypeName = "random_data" 

3378 run = "run1" 

3379 self.source_butler.collections.register(run) 

3380 

3381 dimensions = self.source_butler.dimensions.conform(()) 

3382 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

3383 self.source_butler.registry.registerDatasetType(datasetType) 

3384 

3385 metrics = makeExampleMetrics() 

3386 with ResourcePath.temporary_uri(suffix=".json") as temp: 

3387 dataId = DataCoordinate.make_empty(self.source_butler.dimensions) 

3388 source_refs = [DatasetRef(datasetType, dataId, run=run)] 

3389 temp.write(json.dumps(metrics.exportAsDict()).encode()) 

3390 dataset = FileDataset(path=temp, refs=source_refs) 

3391 self.source_butler.ingest(dataset, transfer="direct") 

3392 

3393 self.target_butler.transfer_from( 

3394 self.source_butler, dataset.refs, register_dataset_types=True, transfer=transfer 

3395 ) 

3396 

3397 uri = self.target_butler.getURI(dataset.refs[0]) 

3398 if transfer == "auto" or transfer == "unsafe_direct": 

3399 self.assertEqual(uri, temp) 

3400 else: 

3401 self.assertNotEqual(uri, temp) 

3402 

3403 def test_shared_dimension_group(self): 

3404 """Test internal logic that divides dataset types by dimension group 

3405 when doing registry updates. 

3406 """ 

3407 self.create_butlers() 

3408 self.source_butler.import_(filename=_get_test_data_path("base.yaml"), without_datastore=True) 

3409 self.source_butler.import_(filename=_get_test_data_path("datasets.yaml"), without_datastore=True) 

3410 

3411 source_butler = self.source_butler 

3412 target_butler = self.target_butler 

3413 

3414 # Create a dataset type with the same dimensions as the 'bias' dataset 

3415 # type from base.yaml 

3416 dataset_type = DatasetType( 

3417 "test_type", ["instrument", "detector"], "int", universe=source_butler.dimensions 

3418 ) 

3419 source_butler.registry.registerDatasetType(dataset_type) 

3420 # This has the same data ID as one of the bias datasets in 

3421 # datasets.yaml. 

3422 test_ref = source_butler.registry.insertDatasets( 

3423 "test_type", [{"instrument": "Cam1", "detector": 2}], run="imported_g" 

3424 )[0] 

3425 

3426 biases = source_butler.query_datasets("bias", ["imported_g", "imported_r"]) 

3427 flats = source_butler.query_datasets("flat", ["imported_g", "imported_r"]) 

3428 refs = [test_ref, *biases, *flats] 

3429 

3430 # Test setup will be even more convoluted if we want the datastore to 

3431 # actually transfer files. For testing the dimension group behavior, 

3432 # we really only care about the registry. 

3433 with unittest.mock.patch.object(target_butler._datastore, "transfer_from") as mock: 

3434 mock.return_value = (set(refs), set()) 

3435 target_butler.transfer_from( 

3436 source_butler, 

3437 refs, 

3438 transfer=None, 

3439 register_dataset_types=True, 

3440 skip_missing=False, 

3441 transfer_dimensions=True, 

3442 ) 

3443 

3444 transferred_test_ref = target_butler.find_dataset( 

3445 "test_type", {"instrument": "Cam1", "detector": 2}, collections="imported_g" 

3446 ) 

3447 self.assertEqual(transferred_test_ref.id, test_ref.id) 

3448 

3449 transferred_bias = target_butler.find_dataset( 

3450 "bias", {"instrument": "Cam1", "detector": 2}, collections="imported_g" 

3451 ) 

3452 self.assertEqual(transferred_bias.id, uuid.UUID("51352db4-a47a-447c-b12d-a50b206b17cd")) 

3453 

3454 transferred_flat = target_butler.find_dataset( 

3455 "flat", 

3456 {"instrument": "Cam1", "detector": 2, "physical_filter": "Cam1-R1", "band": "r"}, 

3457 collections="imported_r", 

3458 ) 

3459 self.assertEqual(transferred_flat.id, uuid.UUID("c1296796-56c5-4acf-9b49-40d920c6f840")) 

3460 

3461 

3462class ChainedDatastoreTransfers(PosixDatastoreTransfers): 

3463 """Test transfers using a chained datastore.""" 

3464 

3465 configFile = os.path.join(TESTDIR, "config/basic/butler-chained.yaml") 

3466 

3467 

3468@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3469class ButlerServerDatastoreTransfers(DatastoreTransfers, unittest.TestCase): 

3470 """Test ``transfer_from`` involving Butler server.""" 

3471 

3472 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3473 

3474 def test_transfers_from_remote_to_direct(self) -> None: 

3475 from lsst.daf.butler.remote_butler._remote_file_transfer_source import ( 

3476 mock_file_transfer_uris_for_unit_test, 

3477 ) 

3478 

3479 self.target_butler = self.create_butler(None, "2") 

3480 with create_test_server(TESTDIR) as server: 

3481 self.source_butler = server.hybrid_butler 

3482 

3483 def _remap_transfer_url(path: HttpResourcePath) -> HttpResourcePath: 

3484 # The Butler server returns HTTP URIs with a domain name that 

3485 # is not resolvable because there is no actual HTTP server 

3486 # involved in these tests. Strip this first layer of 

3487 # indirection, and return the target of the redirect instead. 

3488 response = server.client.get(str(path), follow_redirects=False, headers=path._extra_headers) 

3489 return ResourcePath(str(response.next_request.url)) 

3490 

3491 with mock_file_transfer_uris_for_unit_test(_remap_transfer_url): 

3492 self.assertButlerTransfers() 

3493 

3494 

3495class TransferDatasetsInPlace(unittest.TestCase): 

3496 """Test behavior of transfer_datasets_in_place() specialty function used by 

3497 Prompt Publication service. 

3498 """ 

3499 

3500 def test_file_datastore(self) -> None: 

3501 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3502 with ( 

3503 tempfile.TemporaryDirectory() as datastore_root, 

3504 tempfile.TemporaryDirectory() as other_repo_root, 

3505 ): 

3506 config = Config(configFile) 

3507 config["datastore", "datastore", "name"] = "file_datastore" 

3508 Butler.makeRepo(datastore_root, config=config) 

3509 config["datastore", "datastore", "root"] = datastore_root 

3510 Butler.makeRepo(other_repo_root, config, forceConfigRoot=False) 

3511 with ( 

3512 Butler(datastore_root, writeable=True) as source_butler, 

3513 Butler(other_repo_root, writeable=True) as target_butler, 

3514 ): 

3515 self._test_transfer_datasets_in_place(source_butler, target_butler) 

3516 

3517 def test_chained_datastore(self) -> None: 

3518 configFile = os.path.join(TESTDIR, "config/basic/butler-chained-posix.yaml") 

3519 with ( 

3520 tempfile.TemporaryDirectory() as datastore_root, 

3521 tempfile.TemporaryDirectory() as other_repo_root, 

3522 ): 

3523 config = Config(configFile) 

3524 config["datastore", "datastore", "datastores", 0, "datastore", "root"] = ( 

3525 f"{datastore_root}/butler_test_repository" 

3526 ) 

3527 config["datastore", "datastore", "datastores", 1, "datastore", "root"] = ( 

3528 f"{datastore_root}/butler_test_repository2" 

3529 ) 

3530 Butler.makeRepo(datastore_root, config=config, forceConfigRoot=False) 

3531 Butler.makeRepo(other_repo_root, config=config, forceConfigRoot=False) 

3532 with ( 

3533 Butler(datastore_root, writeable=True) as source_butler, 

3534 Butler(other_repo_root, writeable=True) as target_butler, 

3535 ): 

3536 self._test_transfer_datasets_in_place(source_butler, target_butler) 

3537 

3538 def _test_transfer_datasets_in_place( 

3539 self, source_butler: DirectButler, target_butler: DirectButler 

3540 ) -> None: 

3541 metric_repo = MetricTestRepo.create_from_butler( 

3542 source_butler, 

3543 source_butler._config, 

3544 ) 

3545 target_butler.transfer_dimension_records_from(source_butler, [metric_repo.ref1, metric_repo.ref2]) 

3546 # Verify that the setup was correct and the two repos have 

3547 # independent registries. 

3548 self.assertIsNone(target_butler.get_dataset(metric_repo.ref1.id)) 

3549 # Copy one dataset, and make sure we can load it from the 

3550 # target repo. 

3551 self.assertEqual( 

3552 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1]), 

3553 [metric_repo.ref1], 

3554 ) 

3555 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1)) 

3556 self.assertIsNone(target_butler.get_dataset(metric_repo.ref2.id)) 

3557 self.assertEqual(source_butler.getURIs(metric_repo.ref1), target_butler.getURIs(metric_repo.ref1)) 

3558 # Trying to copy the same dataset again is a no-op. 

3559 self.assertEqual( 

3560 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1]), 

3561 [], 

3562 ) 

3563 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1)) 

3564 # A mix of existing and non-existing datasets. 

3565 self.assertEqual( 

3566 transfer_datasets_in_place(source_butler, target_butler, [metric_repo.ref1, metric_repo.ref2]), 

3567 [metric_repo.ref2], 

3568 ) 

3569 self.assertEqual(target_butler.get(metric_repo.ref1), source_butler.get(metric_repo.ref1)) 

3570 self.assertEqual(target_butler.get(metric_repo.ref2), source_butler.get(metric_repo.ref2)) 

3571 

3572 # For testing datastore chaining, set up a dataset that is only 

3573 # accepted by one of the datastores. 

3574 source_butler.registry.registerDatasetType( 

3575 DatasetType("rejected_by_first", source_butler.dimensions.conform([]), "int") 

3576 ) 

3577 source_butler.registry.registerRun("run") 

3578 ref = source_butler.put(1, "rejected_by_first", dataId={}, run="run") 

3579 self.assertEqual( 

3580 transfer_datasets_in_place(source_butler, target_butler, [ref]), 

3581 [ref], 

3582 ) 

3583 self.assertEqual(1, target_butler.get(ref)) 

3584 

3585 

3586class NullDatastoreTestCase(unittest.TestCase): 

3587 """Test that we can fall back to a null datastore.""" 

3588 

3589 # Need a good config to create the repo. 

3590 configFile = os.path.join(TESTDIR, "config/basic/butler.yaml") 

3591 storageClassFactory: StorageClassFactory 

3592 

3593 @classmethod 

3594 def setUpClass(cls) -> None: 

3595 cls.storageClassFactory = StorageClassFactory() 

3596 cls.storageClassFactory.addFromConfig(cls.configFile) 

3597 

3598 def setUp(self) -> None: 

3599 """Create a new butler root for each test.""" 

3600 self.root = makeTestTempDir(TESTDIR) 

3601 Butler.makeRepo(self.root, config=Config(self.configFile)) 

3602 

3603 def tearDown(self) -> None: 

3604 removeTestTempDir(self.root) 

3605 

3606 def test_fallback(self) -> None: 

3607 # Read the butler config and mess with the datastore section. 

3608 config_path = os.path.join(self.root, "butler.yaml") 

3609 bad_config = Config(config_path) 

3610 bad_config["datastore", "cls"] = "lsst.not.a.datastore.Datastore" 

3611 bad_config.dumpToUri(config_path) 

3612 

3613 with self.assertRaises(RuntimeError): 

3614 Butler(self.root, without_datastore=False) 

3615 

3616 with self.assertRaises(RuntimeError): 

3617 Butler.from_config(self.root, without_datastore=False) 

3618 

3619 butler = Butler.from_config(self.root, writeable=True, without_datastore=True) 

3620 self.enterContext(butler) 

3621 self.assertIsInstance(butler._datastore, NullDatastore) 

3622 

3623 # Check that registry is working. 

3624 butler.collections.register("MYRUN") 

3625 collections = butler.collections.query("*") 

3626 self.assertIn("MYRUN", set(collections)) 

3627 

3628 # Create a ref. 

3629 dimensions = butler.dimensions.conform([]) 

3630 storageClass = self.storageClassFactory.getStorageClass("StructuredDataDict") 

3631 datasetTypeName = "metric" 

3632 datasetType = DatasetType(datasetTypeName, dimensions, storageClass) 

3633 butler.registry.registerDatasetType(datasetType) 

3634 ref = DatasetRef(datasetType, {}, run="MYRUN") 

3635 

3636 # Check that datastore will complain. 

3637 with self.assertRaises(FileNotFoundError): 

3638 butler.get(ref) 

3639 with self.assertRaises(FileNotFoundError): 

3640 butler.getURI(ref) 

3641 

3642 

3643@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3644class ButlerServerTests(FileDatastoreButlerTests): 

3645 """Test RemoteButler and Butler server.""" 

3646 

3647 configFile = None 

3648 predictionSupported = False 

3649 trustModeSupported = False 

3650 

3651 postgres: TemporaryPostgresInstance | None 

3652 

3653 def setUp(self): 

3654 self.server_instance = self.enterContext(create_test_server(TESTDIR)) 

3655 

3656 def tearDown(self): 

3657 pass 

3658 

3659 def are_uris_equivalent(self, uri1: ResourcePath, uri2: ResourcePath) -> bool: 

3660 # S3 pre-signed URLs may end up with differing expiration times in the 

3661 # query parameters, so ignore query parameters when comparing. 

3662 return uri1.scheme == uri2.scheme and uri1.netloc == uri2.netloc and uri1.path == uri2.path 

3663 

3664 def create_empty_butler( 

3665 self, 

3666 run: str | None = None, 

3667 writeable: bool | None = None, 

3668 metrics: ButlerMetrics | None = None, 

3669 cleanup: bool = True, 

3670 ) -> Butler: 

3671 return self.server_instance.hybrid_butler.clone(run=run, metrics=metrics) 

3672 

3673 def remove_dataset_out_of_band(self, butler: Butler, ref: DatasetRef) -> None: 

3674 # Can't delete a file via S3 signed URLs, so we need to reach in 

3675 # through DirectButler to delete the dataset. 

3676 uri = self.server_instance.direct_butler.getURI(ref) 

3677 uri.remove() 

3678 

3679 def testConstructor(self): 

3680 # RemoteButler constructor is tested in test_server.py and 

3681 # test_remote_butler.py. 

3682 pass 

3683 

3684 def testDafButlerRepositories(self): 

3685 # Loading of RemoteButler via repository index is tested in 

3686 # test_server.py. 

3687 pass 

3688 

3689 def testGetDatasetTypes(self) -> None: 

3690 # This is mostly a test of validateConfiguration, which is for 

3691 # validating Datastore configuration and thus isn't relevant to 

3692 # RemoteButler. 

3693 pass 

3694 

3695 def testMakeRepo(self) -> None: 

3696 # Only applies to DirectButler. 

3697 pass 

3698 

3699 # Pickling not yet implemented for RemoteButler/HybridButler. 

3700 @unittest.expectedFailure 

3701 def testPickle(self) -> None: 

3702 return super().testPickle() 

3703 

3704 def testStringification(self) -> None: 

3705 self.assertEqual( 

3706 str(self.server_instance.remote_butler), 

3707 "RemoteButler(https://test.example/api/butler/repo/testrepo/)", 

3708 ) 

3709 

3710 def testTransaction(self) -> None: 

3711 # Transactions will never be supported for RemoteButler. 

3712 pass 

3713 

3714 def testPutTemplates(self) -> None: 

3715 # The Butler server instance is configured with different file naming 

3716 # templates than this test is expecting. 

3717 pass 

3718 

3719 

3720@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3721class ButlerServerSqliteTests(ButlerServerTests, unittest.TestCase): 

3722 """Tests for RemoteButler's registry shim, with a SQLite DB backing the 

3723 server. 

3724 """ 

3725 

3726 postgres = None 

3727 

3728 

3729@unittest.skipIf(not butler_server_is_available, butler_server_import_error) 

3730class ButlerServerPostgresTests(ButlerServerTests, unittest.TestCase): 

3731 """Tests for RemoteButler's registry shim, with a Postgres DB backing the 

3732 server. 

3733 """ 

3734 

3735 @classmethod 

3736 def setUpClass(cls): 

3737 cls.postgres = cls.enterClassContext(setup_postgres_test_db()) 

3738 super().setUpClass() 

3739 

3740 

3741def setup_module(module: types.ModuleType) -> None: 

3742 """Set up the module for pytest.""" 

3743 clean_environment() 

3744 

3745 

3746def _get_test_data_path(filename: str) -> ResourcePath: 

3747 return ResourcePath(f"resource://lsst.daf.butler/tests/registry_data/{filename}") 

3748 

3749 

3750if __name__ == "__main__": 

3751 clean_environment() 

3752 unittest.main()