Coverage for python/lsst/daf/butler/direct_butler/_direct_butler.py: 88%

974 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-08-29 09:13 +0000

1# This file is part of daf_butler. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

5# (http://www.lsst.org). 

6# See the COPYRIGHT file at the top-level directory of this distribution 

7# for details of code ownership. 

8# 

9# This software is dual licensed under the GNU General Public License and also 

10# under a 3-clause BSD license. Recipients may choose which of these licenses 

11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt, 

12# respectively. If you choose the GPL option then the following text applies 

13# (but note that there is still no warranty even if you opt for BSD instead): 

14# 

15# This program is free software: you can redistribute it and/or modify 

16# it under the terms of the GNU General Public License as published by 

17# the Free Software Foundation, either version 3 of the License, or 

18# (at your option) any later version. 

19# 

20# This program is distributed in the hope that it will be useful, 

21# but WITHOUT ANY WARRANTY; without even the implied warranty of 

22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the 

23# GNU General Public License for more details. 

24# 

25# You should have received a copy of the GNU General Public License 

26# along with this program. If not, see <http://www.gnu.org/licenses/>. 

27 

28"""Butler top level classes.""" 

29 

30from __future__ import annotations 

31 

32__all__ = ( 

33 "ButlerValidationError", 

34 "DirectButler", 

35) 

36 

37import collections.abc 

38import contextlib 

39import io 

40import itertools 

41import math 

42import numbers 

43import os 

44import uuid 

45import warnings 

46from collections import Counter, defaultdict 

47from collections.abc import Collection, Iterable, Iterator, Mapping, MutableMapping, Sequence 

48from functools import partial 

49from types import EllipsisType 

50from typing import TYPE_CHECKING, Any, ClassVar, NamedTuple, TextIO, cast 

51 

52from deprecated.sphinx import deprecated 

53from sqlalchemy.exc import IntegrityError 

54 

55from lsst.resources import ResourcePath, ResourcePathExpression 

56from lsst.utils.introspection import find_outside_stacklevel, get_class_of 

57from lsst.utils.iteration import chunk_iterable 

58from lsst.utils.logging import VERBOSE, getLogger 

59from lsst.utils.timer import time_this 

60 

61from .._butler import Butler, _DeprecatedDefault 

62from .._butler_config import ButlerConfig 

63from .._butler_instance_options import ButlerInstanceOptions 

64from .._butler_metrics import ButlerMetrics 

65from .._collection_type import CollectionType 

66from .._dataset_existence import DatasetExistence 

67from .._dataset_ref import DatasetRef 

68from .._dataset_type import DatasetType 

69from .._deferredDatasetHandle import DeferredDatasetHandle 

70from .._exceptions import ( 

71 DatasetNotFoundError, 

72 DimensionValueError, 

73 EmptyQueryResultError, 

74 ValidationError, 

75) 

76from .._file_dataset import FileDataset 

77from .._limited_butler import LimitedButler 

78from .._query_all_datasets import QueryAllDatasetsParameters, query_all_datasets 

79from .._registry_shim import RegistryShim 

80from .._storage_class import StorageClass, StorageClassFactory 

81from .._timespan import Timespan 

82from ..datastore import Datastore, NullDatastore 

83from ..datastores.file_datastore.retrieve_artifacts import ZipIndex, retrieve_and_zip 

84from ..datastores.file_datastore.transfer import retrieve_file_transfer_records 

85from ..dimensions import DataCoordinate, Dimension, DimensionGroup 

86from ..direct_query_driver import DirectQueryDriver 

87from ..progress import Progress 

88from ..queries import Query 

89from ..registry import ( 

90 ConflictingDefinitionError, 

91 DataIdError, 

92 MissingDatasetTypeError, 

93 RegistryDefaults, 

94 _RegistryFactory, 

95) 

96from ..registry.sql_registry import SqlRegistry 

97from ..transfers import RepoExportContext 

98from ..utils import transactional 

99from ._direct_butler_collections import DirectButlerCollections 

100 

101if TYPE_CHECKING: 

102 from lsst.resources import ResourceHandleProtocol 

103 

104 from .._dataset_provenance import DatasetProvenance 

105 from .._dataset_ref import DatasetId 

106 from ..datastore import DatasetRefURIs 

107 from ..dimensions import DataId, DataIdValue, DimensionElement, DimensionRecord, DimensionUniverse 

108 from ..registry import CollectionArgType, Registry 

109 from ..transfers import RepoImportBackend 

110 

111_LOG = getLogger(__name__) 

112 

113 

114class ButlerValidationError(ValidationError): 

115 """There is a problem with the Butler configuration.""" 

116 

117 pass 

118 

119 

120class DirectButler(Butler): # numpydoc ignore=PR02 

121 """Main entry point for the data access system. 

122 

123 Parameters 

124 ---------- 

125 config : `ButlerConfig` 

126 The configuration for this Butler instance. 

127 registry : `SqlRegistry` 

128 The object that manages dataset metadata and relationships. 

129 datastore : Datastore 

130 The object that manages actual dataset storage. 

131 storageClasses : StorageClassFactory 

132 An object that maps known storage class names to objects that fully 

133 describe them. 

134 

135 Notes 

136 ----- 

137 Most users should call the top-level `Butler`.``from_config`` instead of 

138 using this constructor directly. 

139 """ 

140 

141 # This is __new__ instead of __init__ because we have to support 

142 # instantiation via the legacy constructor Butler.__new__(), which 

143 # reads the configuration and selects which subclass to instantiate. The 

144 # interaction between __new__ and __init__ is kind of wacky in Python. If 

145 # we were using __init__ here, __init__ would be called twice (once when 

146 # the DirectButler instance is constructed inside Butler.from_config(), and 

147 # a second time with the original arguments to Butler() when the instance 

148 # is returned from Butler.__new__() 

149 def __new__( 

150 cls, 

151 *, 

152 config: ButlerConfig, 

153 registry: SqlRegistry, 

154 datastore: Datastore, 

155 storageClasses: StorageClassFactory, 

156 metrics: ButlerMetrics | None = None, 

157 ) -> DirectButler: 

158 self = cast(DirectButler, super().__new__(cls)) 

159 self._config = config 

160 self._registry = registry 

161 self._datastore = datastore 

162 self.storageClasses = storageClasses 

163 self._metrics = metrics if metrics is not None else ButlerMetrics() 

164 

165 # For execution butler the datastore needs a special 

166 # dependency-inversion trick. This is not used by regular butler, 

167 # but we do not have a way to distinguish regular butler from execution 

168 # butler. 

169 self._datastore.set_retrieve_dataset_type_method(partial(_retrieve_dataset_type, registry)) 

170 

171 self._closed = False 

172 

173 return self 

174 

175 @classmethod 

176 def create_from_config( 

177 cls, 

178 config: ButlerConfig, 

179 *, 

180 options: ButlerInstanceOptions, 

181 without_datastore: bool = False, 

182 ) -> DirectButler: 

183 """Construct a Butler instance from a configuration file. 

184 

185 Parameters 

186 ---------- 

187 config : `ButlerConfig` 

188 The configuration for this Butler instance. 

189 options : `ButlerInstanceOptions` 

190 Default values and other settings for the Butler instance. 

191 without_datastore : `bool`, optional 

192 If `True` do not attach a datastore to this butler. Any attempts 

193 to use a datastore will fail. 

194 

195 Notes 

196 ----- 

197 Most users should call the top-level `Butler`.``from_config`` 

198 instead of using this function directly. 

199 """ 

200 if "run" in config or "collection" in config: 200 ↛ 201line 200 didn't jump to line 201 because the condition on line 200 was never true

201 raise ValueError("Passing a run or collection via configuration is no longer supported.") 

202 

203 defaults = RegistryDefaults.from_butler_instance_options(options) 

204 try: 

205 butlerRoot = config.get("root", config.configDir) 

206 writeable = options.writeable 

207 if writeable is None: 

208 writeable = options.run is not None 

209 registry = _RegistryFactory(config).from_config( 

210 butlerRoot=butlerRoot, writeable=writeable, defaults=defaults 

211 ) 

212 if without_datastore: 

213 datastore: Datastore = NullDatastore(None, None) 

214 else: 

215 datastore = Datastore.fromConfig( 

216 config, registry.getDatastoreBridgeManager(), butlerRoot=butlerRoot 

217 ) 

218 # TODO: Once datastore drops dependency on registry we can 

219 # construct datastore first and pass opaque tables to registry 

220 # constructor. 

221 registry.make_datastore_tables(datastore.get_opaque_table_definitions()) 

222 storageClasses = StorageClassFactory() 

223 storageClasses.addFromConfig(config) 

224 

225 return DirectButler( 

226 config=config, 

227 registry=registry, 

228 datastore=datastore, 

229 storageClasses=storageClasses, 

230 metrics=options.metrics, 

231 ) 

232 except Exception: 

233 # Failures here usually mean that configuration is incomplete, 

234 # just issue an error message which includes config file URI. 

235 _LOG.error(f"Failed to instantiate Butler from config {config.configFile}.") 

236 raise 

237 

238 def clone( 

239 self, 

240 *, 

241 collections: CollectionArgType | None | EllipsisType = ..., 

242 run: str | None | EllipsisType = ..., 

243 inferDefaults: bool | EllipsisType = ..., 

244 dataId: dict[str, str] | EllipsisType = ..., 

245 metrics: ButlerMetrics | None = None, 

246 ) -> DirectButler: 

247 # Docstring inherited 

248 defaults = self._registry.defaults.clone(collections, run, inferDefaults, dataId) 

249 registry = self._registry.copy(defaults) 

250 

251 return DirectButler( 

252 registry=registry, 

253 config=self._config, 

254 datastore=self._datastore.clone(registry.getDatastoreBridgeManager()), 

255 storageClasses=self.storageClasses, 

256 metrics=metrics, 

257 ) 

258 

259 def close(self) -> None: 

260 if not self._closed: 

261 self._closed = True 

262 self._registry.close() 

263 # Cause exceptions to be raised if a user attempts to use the 

264 # instance after closing it. Without this, Butler would still 

265 # work after being closed because of implementation details 

266 # of SqlAlchemy, but this may not continue to be the case in the 

267 # future and we don't want users to get in the habit of doing this. 

268 self._registry = _BUTLER_CLOSED_INSTANCE 

269 self._datastore = _BUTLER_CLOSED_INSTANCE 

270 

271 GENERATION: ClassVar[int] = 3 

272 """This is a Generation 3 Butler. 

273 

274 This attribute may be removed in the future, once the Generation 2 Butler 

275 interface has been fully retired; it should only be used in transitional 

276 code. 

277 """ 

278 

279 @classmethod 

280 def _unpickle( 

281 cls, 

282 config: ButlerConfig, 

283 collections: tuple[str, ...] | None, 

284 run: str | None, 

285 defaultDataId: dict[str, str], 

286 writeable: bool, 

287 ) -> DirectButler: 

288 """Callable used to unpickle a Butler. 

289 

290 We prefer not to use ``Butler.__init__`` directly so we can force some 

291 of its many arguments to be keyword-only (note that ``__reduce__`` 

292 can only invoke callables with positional arguments). 

293 

294 Parameters 

295 ---------- 

296 config : `ButlerConfig` 

297 Butler configuration, already coerced into a true `ButlerConfig` 

298 instance (and hence after any search paths for overrides have been 

299 utilized). 

300 collections : `tuple` [ `str` ] 

301 Names of the default collections to read from. 

302 run : `str`, optional 

303 Name of the default `~CollectionType.RUN` collection to write to. 

304 defaultDataId : `dict` [ `str`, `str` ] 

305 Default data ID values. 

306 writeable : `bool` 

307 Whether the Butler should support write operations. 

308 

309 Returns 

310 ------- 

311 butler : `Butler` 

312 A new `Butler` instance. 

313 """ 

314 return cls.create_from_config( 

315 config=config, 

316 options=ButlerInstanceOptions( 

317 collections=collections, run=run, writeable=writeable, kwargs=defaultDataId 

318 ), 

319 ) 

320 

321 def __reduce__(self) -> tuple: 

322 """Support pickling.""" 

323 return ( 

324 DirectButler._unpickle, 

325 ( 

326 self._config, 

327 self.collections.defaults, 

328 self.run, 

329 dict(self._registry.defaults.dataId.required), 

330 self._registry.isWriteable(), 

331 ), 

332 ) 

333 

334 def __str__(self) -> str: 

335 return ( 

336 f"Butler(collections={self.collections}, run={self.run}, " 

337 f"datastore='{self._datastore}', registry='{self._registry}')" 

338 ) 

339 

340 def isWriteable(self) -> bool: 

341 # Docstring inherited. 

342 return self._registry.isWriteable() 

343 

344 def _caching_context(self) -> contextlib.AbstractContextManager[None]: 

345 """Context manager that enables caching.""" 

346 return self._registry.caching_context() 

347 

348 @contextlib.contextmanager 

349 def transaction(self) -> Iterator[None]: 

350 """Context manager supporting `Butler` transactions. 

351 

352 Transactions can be nested. 

353 """ 

354 with self._registry.transaction(), self._datastore.transaction(): 

355 yield 

356 

357 def _standardizeArgs( 

358 self, 

359 datasetRefOrType: DatasetRef | DatasetType | str, 

360 dataId: DataId | None = None, 

361 for_put: bool = True, 

362 **kwargs: Any, 

363 ) -> tuple[DatasetType, DataId | None]: 

364 """Standardize the arguments passed to several Butler APIs. 

365 

366 Parameters 

367 ---------- 

368 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str` 

369 When `DatasetRef` the `dataId` should be `None`. 

370 Otherwise the `DatasetType` or name thereof. 

371 dataId : `dict` or `DataCoordinate` 

372 A `dict` of `Dimension` link name, value pairs that label the 

373 `DatasetRef` within a Collection. When `None`, a `DatasetRef` 

374 should be provided as the second argument. 

375 for_put : `bool`, optional 

376 If `True` this call is invoked as part of a `Butler.put`. 

377 Otherwise it is assumed to be part of a `Butler.get()`. This 

378 parameter is only relevant if there is dataset type 

379 inconsistency. 

380 **kwargs 

381 Additional keyword arguments used to augment or construct a 

382 `DataCoordinate`. See `DataCoordinate.standardize` 

383 parameters. 

384 

385 Returns 

386 ------- 

387 datasetType : `DatasetType` 

388 A `DatasetType` instance extracted from ``datasetRefOrType``. 

389 dataId : `dict` or `DataId`, optional 

390 Argument that can be used (along with ``kwargs``) to construct a 

391 `DataId`. 

392 

393 Notes 

394 ----- 

395 Butler APIs that conceptually need a DatasetRef also allow passing a 

396 `DatasetType` (or the name of one) and a `DataId` (or a dict and 

397 keyword arguments that can be used to construct one) separately. This 

398 method accepts those arguments and always returns a true `DatasetType` 

399 and a `DataId` or `dict`. 

400 

401 Standardization of `dict` vs `DataId` is best handled by passing the 

402 returned ``dataId`` (and ``kwargs``) to `Registry` APIs, which are 

403 generally similarly flexible. 

404 """ 

405 externalDatasetType: DatasetType | None = None 

406 internalDatasetType: DatasetType | None = None 

407 if isinstance(datasetRefOrType, DatasetRef): 

408 if dataId is not None or kwargs: 

409 raise ValueError("DatasetRef given, cannot use dataId as well") 

410 externalDatasetType = datasetRefOrType.datasetType 

411 dataId = datasetRefOrType.dataId 

412 else: 

413 # Don't check whether DataId is provided, because Registry APIs 

414 # can usually construct a better error message when it wasn't. 

415 if isinstance(datasetRefOrType, DatasetType): 

416 externalDatasetType = datasetRefOrType 

417 else: 

418 internalDatasetType = self.get_dataset_type(datasetRefOrType) 

419 

420 # Check that they are self-consistent 

421 if externalDatasetType is not None: 

422 registryDatasetType = self._get_registry_dataset_type(externalDatasetType) 

423 if registryDatasetType is None: 

424 # The caller has asked for a component that only their own 

425 # storage class defines, so there is no registry definition to 

426 # check against and their definition has to be used as given. 

427 internalDatasetType = externalDatasetType 

428 else: 

429 internalDatasetType = registryDatasetType 

430 if externalDatasetType != internalDatasetType: 

431 # We can allow differences if they are compatible, 

432 # depending on whether this is a get or a put. A get 

433 # requires that the python type associated with the 

434 # datastore can be converted to the user type. A put 

435 # requires that the user supplied python type can be 

436 # converted to the internal type expected by registry. 

437 relevantDatasetType = internalDatasetType 

438 if for_put: 

439 is_compatible = internalDatasetType.is_compatible_with(externalDatasetType) 

440 else: 

441 is_compatible = externalDatasetType.is_compatible_with(internalDatasetType) 

442 relevantDatasetType = externalDatasetType 

443 if not is_compatible: 

444 raise ValueError( 

445 f"Supplied dataset type ({externalDatasetType}) inconsistent with " 

446 f"registry definition ({internalDatasetType})" 

447 ) 

448 # Override the internal definition. 

449 internalDatasetType = relevantDatasetType 

450 

451 assert internalDatasetType is not None 

452 return internalDatasetType, dataId 

453 

454 def _get_registry_dataset_type(self, datasetType: DatasetType) -> DatasetType | None: 

455 """Return the registry definition corresponding to the given dataset 

456 type. 

457 

458 Parameters 

459 ---------- 

460 datasetType : `DatasetType` 

461 Dataset type, possibly a component, supplied by the caller. 

462 

463 Returns 

464 ------- 

465 registry_type : `DatasetType` or `None` 

466 The registry definition of ``datasetType``, or `None` if the 

467 registry has no definition for it. 

468 

469 Raises 

470 ------ 

471 MissingDatasetTypeError 

472 Raised if the dataset type is not registered. For a component 

473 dataset type this refers to the composite. 

474 

475 Notes 

476 ----- 

477 Only composites are registered: a component dataset type is derived 

478 from the storage class of its composite. A component that is defined 

479 only by a read-time storage class override therefore has no registry 

480 definition at all, even though its composite is registered, and `None` 

481 is returned to say so rather than something that only resembles a 

482 registry definition. 

483 """ 

484 parent_name, component = DatasetType.splitDatasetTypeName(datasetType.name) 

485 parent = self.get_dataset_type(parent_name) 

486 if component is None: 

487 return parent 

488 if component in parent.storageClass.allComponents(): 

489 return parent.makeComponentDatasetType(component) 

490 return None 

491 

492 def _rewrite_data_id( 

493 self, dataId: DataId | None, datasetType: DatasetType, **kwargs: Any 

494 ) -> tuple[DataId | None, dict[str, Any]]: 

495 """Rewrite a data ID taking into account dimension records. 

496 

497 Take a Data ID and keyword args and rewrite it if necessary to 

498 allow the user to specify dimension records rather than dimension 

499 primary values. 

500 

501 This allows a user to include a dataId dict with keys of 

502 ``exposure.day_obs`` and ``exposure.seq_num`` instead of giving 

503 the integer exposure ID. It also allows a string to be given 

504 for a dimension value rather than the integer ID if that is more 

505 convenient. For example, rather than having to specifying the 

506 detector with ``detector.full_name``, a string given for ``detector`` 

507 will be interpreted as the full name and converted to the integer 

508 value. 

509 

510 Keyword arguments can also use strings for dimensions like detector 

511 and exposure but python does not allow them to include ``.`` and 

512 so the ``exposure.day_obs`` syntax can not be used in a keyword 

513 argument. 

514 

515 Parameters 

516 ---------- 

517 dataId : `dict` or `DataCoordinate` 

518 A `dict` of `Dimension` link name, value pairs that will label the 

519 `DatasetRef` within a Collection. 

520 datasetType : `DatasetType` 

521 The dataset type associated with this dataId. Required to 

522 determine the relevant dimensions. 

523 **kwargs 

524 Additional keyword arguments used to augment or construct a 

525 `DataId`. See `DataId` parameters. 

526 

527 Returns 

528 ------- 

529 dataId : `dict` or `DataCoordinate` 

530 The, possibly rewritten, dataId. If given a `DataCoordinate` and 

531 no keyword arguments, the original dataId will be returned 

532 unchanged. 

533 **kwargs : `dict` 

534 Any unused keyword arguments (would normally be empty dict). 

535 """ 

536 # Process dimension records that are using record information 

537 # rather than ids 

538 newDataId: dict[str, DataIdValue] = {} 

539 byRecord: dict[str, dict[str, Any]] = defaultdict(dict) 

540 

541 if isinstance(dataId, DataCoordinate): 

542 # Do nothing if we have a DataCoordinate and no kwargs. 

543 if not kwargs: 543 ↛ 547line 543 didn't jump to line 547 because the condition on line 543 was always true

544 return dataId, kwargs 

545 # If we have a DataCoordinate with kwargs, we know the 

546 # DataCoordinate only has values for real dimensions. 

547 newDataId.update(dataId.mapping) 

548 elif dataId: 

549 # The data is mapping, which means it might have keys like 

550 # "exposure.obs_id" (unlike kwargs, because a "." is not allowed in 

551 # a keyword parameter). 

552 for k, v in dataId.items(): 

553 if isinstance(k, str) and "." in k: 

554 # Someone is using a more human-readable dataId 

555 dimensionName, record = k.split(".", 1) 

556 byRecord[dimensionName][record] = v 

557 else: 

558 newDataId[k] = v 

559 

560 # Go through the updated dataId and check the type in case someone is 

561 # using an alternate key. We have already filtered out the compound 

562 # keys dimensions.record format. 

563 not_dimensions = {} 

564 

565 # Will need to look in the dataId and the keyword arguments 

566 # and will remove them if they need to be fixed or are unrecognized. 

567 for dataIdDict in (newDataId, kwargs): 

568 # Use a list so we can adjust the dict safely in the loop 

569 for dimensionName in list(dataIdDict): 

570 value = dataIdDict[dimensionName] 

571 try: 

572 dimension = self.dimensions.dimensions[dimensionName] 

573 except KeyError: 

574 # This is not a real dimension 

575 not_dimensions[dimensionName] = value 

576 del dataIdDict[dimensionName] 

577 continue 

578 

579 # Convert an integral type to an explicit int to simplify 

580 # comparisons here 

581 if isinstance(value, numbers.Integral): 

582 value = int(value) 

583 

584 if not isinstance(value, dimension.primaryKey.getPythonType()): 

585 for alternate in dimension.alternateKeys: 585 ↛ 598line 585 didn't jump to line 598 because the loop on line 585 didn't complete

586 if isinstance(value, alternate.getPythonType()): 586 ↛ 585line 586 didn't jump to line 585 because the condition on line 586 was always true

587 byRecord[dimensionName][alternate.name] = value 

588 del dataIdDict[dimensionName] 

589 _LOG.debug( 

590 "Converting dimension %s to %s.%s=%s", 

591 dimensionName, 

592 dimensionName, 

593 alternate.name, 

594 value, 

595 ) 

596 break 

597 else: 

598 _LOG.warning( 

599 "Type mismatch found for value '%r' provided for dimension %s. " 

600 "Could not find matching alternative (primary key has type %s) " 

601 "so attempting to use as-is.", 

602 value, 

603 dimensionName, 

604 dimension.primaryKey.getPythonType(), 

605 ) 

606 

607 # By this point kwargs and newDataId should only include valid 

608 # dimensions. Merge kwargs in to the new dataId and log if there 

609 # are dimensions in both (rather than calling update). 

610 for k, v in kwargs.items(): 

611 if k in newDataId and newDataId[k] != v: 

612 _LOG.debug( 

613 "Keyword arg %s overriding explicit value in dataId of %s with %s", k, newDataId[k], v 

614 ) 

615 newDataId[k] = v 

616 # No need to retain any values in kwargs now. 

617 kwargs = {} 

618 

619 # If we have some unrecognized dimensions we have to try to connect 

620 # them to records in other dimensions. This is made more complicated 

621 # by some dimensions having records with clashing names. A mitigation 

622 # is that we can tell by this point which dimensions are missing 

623 # for the DatasetType but this does not work for calibrations 

624 # where additional dimensions can be used to constrain the temporal 

625 # axis. 

626 if not_dimensions: 

627 # Search for all dimensions even if we have been given a value 

628 # explicitly. In some cases records are given as well as the 

629 # actually dimension and this should not be an error if they 

630 # match. 

631 mandatoryDimensions = datasetType.dimensions.names # - provided 

632 

633 candidateDimensions: set[str] = set() 

634 candidateDimensions.update(mandatoryDimensions) 

635 

636 # For calibrations we may well be needing temporal dimensions 

637 # so rather than always including all dimensions in the scan 

638 # restrict things a little. It is still possible for there 

639 # to be confusion over day_obs in visit vs exposure for example. 

640 # If we are not searching calibration collections things may 

641 # fail but they are going to fail anyway because of the 

642 # ambiguousness of the dataId... 

643 if datasetType.isCalibration(): 

644 for dim in self.dimensions.dimensions: 

645 if dim.temporal: 

646 candidateDimensions.add(str(dim)) 

647 

648 # Look up table for the first association with a dimension 

649 guessedAssociation: dict[str, dict[str, Any]] = defaultdict(dict) 

650 

651 # Keep track of whether an item is associated with multiple 

652 # dimensions. 

653 counter: Counter[str] = Counter() 

654 assigned: dict[str, set[str]] = defaultdict(set) 

655 

656 # Go through the missing dimensions and associate the 

657 # given names with records within those dimensions 

658 matched_dims = set() 

659 for dimensionName in candidateDimensions: 

660 dimension = self.dimensions.dimensions[dimensionName] 

661 fields = dimension.metadata.names | dimension.uniqueKeys.names 

662 for field in not_dimensions: 

663 if field in fields: 

664 guessedAssociation[dimensionName][field] = not_dimensions[field] 

665 counter[dimensionName] += 1 

666 assigned[field].add(dimensionName) 

667 matched_dims.add(field) 

668 

669 # Calculate the fields that matched nothing. 

670 never_found = set(not_dimensions) - matched_dims 

671 

672 if never_found: 

673 raise DimensionValueError(f"Unrecognized keyword args given: {never_found}") 

674 

675 # There is a chance we have allocated a single dataId item 

676 # to multiple dimensions. Need to decide which should be retained. 

677 # For now assume that the most popular alternative wins. 

678 # This means that day_obs with seq_num will result in 

679 # exposure.day_obs and not visit.day_obs 

680 # Also prefer an explicitly missing dimension over an inferred 

681 # temporal dimension. 

682 for fieldName, assignedDimensions in assigned.items(): 

683 if len(assignedDimensions) > 1: 

684 # Pick the most popular (preferring mandatory dimensions) 

685 requiredButMissing = assignedDimensions.intersection(mandatoryDimensions) 

686 if requiredButMissing: 686 ↛ 687line 686 didn't jump to line 687 because the condition on line 686 was never true

687 candidateDimensions = requiredButMissing 

688 else: 

689 candidateDimensions = assignedDimensions 

690 

691 # If this is a choice between visit and exposure and 

692 # neither was a required part of the dataset type, 

693 # (hence in this branch) always prefer exposure over 

694 # visit since exposures are always defined and visits 

695 # are defined from exposures. 

696 if candidateDimensions == {"exposure", "visit"}: 696 ↛ 701line 696 didn't jump to line 701 because the condition on line 696 was always true

697 candidateDimensions = {"exposure"} 

698 

699 # Select the relevant items and get a new restricted 

700 # counter. 

701 theseCounts = {k: v for k, v in counter.items() if k in candidateDimensions} 

702 duplicatesCounter: Counter[str] = Counter() 

703 duplicatesCounter.update(theseCounts) 

704 

705 # Choose the most common. If they are equally common 

706 # we will pick the one that was found first. 

707 # Returns a list of tuples 

708 selected = duplicatesCounter.most_common(1)[0][0] 

709 

710 _LOG.debug( 

711 "Ambiguous dataId entry '%s' associated with multiple dimensions: %s." 

712 " Removed ambiguity by choosing dimension %s.", 

713 fieldName, 

714 ", ".join(assignedDimensions), 

715 selected, 

716 ) 

717 

718 for candidateDimension in assignedDimensions: 

719 if candidateDimension != selected: 

720 del guessedAssociation[candidateDimension][fieldName] 

721 

722 # Update the record look up dict with the new associations 

723 for dimensionName, values in guessedAssociation.items(): 

724 if values: # A dict might now be empty 

725 _LOG.debug( 

726 "Assigned non-dimension dataId keys to dimension %s: %s", dimensionName, values 

727 ) 

728 byRecord[dimensionName].update(values) 

729 

730 if byRecord: 

731 # Some record specifiers were found so we need to convert 

732 # them to the Id form 

733 for dimensionName, values in byRecord.items(): 

734 if dimensionName in newDataId: 

735 _LOG.debug( 

736 "DataId specified explicit %s dimension value of %s in addition to" 

737 " general record specifiers for it of %s. Checking for self-consistency.", 

738 dimensionName, 

739 newDataId[dimensionName], 

740 str(values), 

741 ) 

742 # Get the actual record and compare with these values. 

743 # Only query with relevant data ID values. 

744 filtered_data_id = { 

745 k: v for k, v in newDataId.items() if k in self.dimensions[dimensionName].required 

746 } 

747 try: 

748 recs = self.query_dimension_records( 

749 dimensionName, 

750 data_id=filtered_data_id, 

751 ) 

752 except (DataIdError, EmptyQueryResultError): 

753 raise DimensionValueError( 

754 f"Could not find dimension '{dimensionName}'" 

755 f" with dataId {filtered_data_id} as part of comparing with" 

756 f" record values {byRecord[dimensionName]}" 

757 ) from None 

758 if len(recs) == 1: 758 ↛ 772line 758 didn't jump to line 772 because the condition on line 758 was always true

759 errmsg: list[str] = [] 

760 for k, v in values.items(): 

761 if (recval := getattr(recs[0], k)) != v: 

762 errmsg.append(f"{k} ({recval} != {v})") 

763 if errmsg: 

764 raise DimensionValueError( 

765 f"Dimension {dimensionName} in dataId has explicit value" 

766 f" {newDataId[dimensionName]} inconsistent with" 

767 f" {dimensionName} dimension record: " + ", ".join(errmsg) 

768 ) 

769 else: 

770 # Multiple matches for an explicit dimension 

771 # should never happen but let downstream complain. 

772 pass 

773 continue 

774 

775 # Do not use data ID keys in query that aren't relevant. 

776 # Otherwise we can have detector queries being constrained 

777 # by an exposure ID that doesn't exist and return no matches 

778 # for a detector even though it's a good detector name. 

779 filtered_data_id = { 

780 k: v 

781 for k, v in newDataId.items() 

782 if k in self.dimensions[dimensionName].minimal_group.names 

783 } 

784 

785 def _get_attr(obj: Any, attr: str) -> Any: 

786 # Used to implement x.exposure.seq_num when given 

787 # x and "exposure.seq_num". 

788 for component in attr.split("."): 

789 obj = getattr(obj, component) 

790 return obj 

791 

792 with self.query() as q: 

793 x = q.expression_factory 

794 # Build up a WHERE expression. 

795 predicates = tuple(_get_attr(x, f"{dimensionName}.{k}") == v for k, v in values.items()) 

796 extra_args: dict[str, Any] = {} # For mypy. 

797 extra_args.update(filtered_data_id) 

798 extra_args.update(kwargs) 

799 q = q.where(x.all(*predicates), **extra_args) 

800 records = set(q.dimension_records(dimensionName)) 

801 

802 if len(records) != 1: 

803 if len(records) > 1: 

804 # visit can have an ambiguous answer without involving 

805 # visit_system. The default visit_system is defined 

806 # by the instrument. 

807 if ( 807 ↛ 812line 807 didn't jump to line 812 because the condition on line 807 was never true

808 dimensionName == "visit" 

809 and "visit_system_membership" in self.dimensions 

810 and "visit_system" in self.dimensions["instrument"].metadata 

811 ): 

812 instrument_records = self.query_dimension_records( 

813 "instrument", 

814 data_id=newDataId, 

815 explain=False, 

816 **kwargs, 

817 ) 

818 if len(instrument_records) == 1: 

819 visit_system = instrument_records[0].visit_system 

820 if visit_system is None: 

821 # Set to a value that will never match. 

822 visit_system = -1 

823 

824 # Look up each visit in the 

825 # visit_system_membership records. 

826 for rec in records: 

827 membership = self.query_dimension_records( 

828 # Use bind to allow zero results. 

829 # This is a fully-specified query. 

830 "visit_system_membership", 

831 instrument=instrument_records[0].name, 

832 visit_system=visit_system, 

833 visit=rec.id, 

834 explain=False, 

835 ) 

836 if membership: 

837 # This record is the right answer. 

838 records = {rec} 

839 break 

840 

841 # The ambiguity may have been resolved so check again. 

842 if len(records) > 1: 842 ↛ 860line 842 didn't jump to line 860 because the condition on line 842 was always true

843 _LOG.debug( 

844 "Received %d records from constraints of %s", len(records), str(values) 

845 ) 

846 for r in records: 

847 _LOG.debug("- %s", str(r)) 

848 raise DimensionValueError( 

849 f"DataId specification for dimension {dimensionName} is not" 

850 f" uniquely constrained to a single dataset by {values}." 

851 f" Got {len(records)} results." 

852 ) 

853 else: 

854 raise DimensionValueError( 

855 f"DataId specification for dimension {dimensionName} matched no" 

856 f" records when constrained by {values}" 

857 ) 

858 

859 # Get the primary key from the real dimension object 

860 dimension = self.dimensions.dimensions[dimensionName] 

861 if not isinstance(dimension, Dimension): 861 ↛ 862line 861 didn't jump to line 862 because the condition on line 861 was never true

862 raise RuntimeError( 

863 f"{dimension.name} is not a true dimension, and cannot be used in data IDs." 

864 ) 

865 newDataId[dimensionName] = getattr(records.pop(), dimension.primaryKey.name) 

866 

867 return newDataId, kwargs 

868 

869 def _findDatasetRef( 

870 self, 

871 datasetRefOrType: DatasetRef | DatasetType | str, 

872 dataId: DataId | None = None, 

873 *, 

874 collections: Any = None, 

875 predict: bool = False, 

876 run: str | None = None, 

877 datastore_records: bool = False, 

878 timespan: Timespan | None = None, 

879 **kwargs: Any, 

880 ) -> DatasetRef: 

881 """Shared logic for methods that start with a search for a dataset in 

882 the registry. 

883 

884 Parameters 

885 ---------- 

886 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str` 

887 When `DatasetRef` the `dataId` should be `None`. 

888 Otherwise the `DatasetType` or name thereof. 

889 dataId : `dict` or `DataCoordinate`, optional 

890 A `dict` of `Dimension` link name, value pairs that label the 

891 `DatasetRef` within a Collection. When `None`, a `DatasetRef` 

892 should be provided as the first argument. 

893 collections : Any, optional 

894 Collections to be searched, overriding ``self.collections``. 

895 Can be any of the types supported by the ``collections`` argument 

896 to butler construction. 

897 predict : `bool`, optional 

898 If `True`, return a newly created `DatasetRef` with a unique 

899 dataset ID if finding a reference in the `Registry` fails. 

900 Defaults to `False`. 

901 run : `str`, optional 

902 Run collection name to use for creating `DatasetRef` for predicted 

903 datasets. Only used if ``predict`` is `True`. 

904 datastore_records : `bool`, optional 

905 If `True` add datastore records to returned `DatasetRef`. 

906 timespan : `Timespan` or `None`, optional 

907 A timespan that the validity range of the dataset must overlap. 

908 If not provided and this is a calibration dataset type, an attempt 

909 will be made to find the timespan from any temporal coordinate 

910 in the data ID. 

911 **kwargs 

912 Additional keyword arguments used to augment or construct a 

913 `DataId`. See `DataId` parameters. 

914 

915 Returns 

916 ------- 

917 ref : `DatasetRef` 

918 A reference to the dataset identified by the given arguments. 

919 This can be the same dataset reference as given if it was 

920 resolved. 

921 

922 Raises 

923 ------ 

924 LookupError 

925 Raised if no matching dataset exists in the `Registry` (and 

926 ``predict`` is `False`). 

927 ValueError 

928 Raised if a resolved `DatasetRef` was passed as an input, but it 

929 differs from the one found in the registry. 

930 TypeError 

931 Raised if no collections were provided. 

932 """ 

933 datasetType, dataId = self._standardizeArgs(datasetRefOrType, dataId, for_put=False, **kwargs) 

934 if isinstance(datasetRefOrType, DatasetRef): 

935 if collections is not None: 935 ↛ 936line 935 didn't jump to line 936 because the condition on line 935 was never true

936 warnings.warn("Collections should not be specified with DatasetRef", stacklevel=3) 

937 if predict and not datasetRefOrType.dataId.hasRecords(): 

938 return datasetRefOrType.expanded(self.registry.expandDataId(datasetRefOrType.dataId)) 

939 # May need to retrieve datastore records if requested. 

940 if datastore_records and datasetRefOrType._datastore_records is None: 

941 datasetRefOrType = self._registry.get_datastore_records(datasetRefOrType) 

942 return datasetRefOrType 

943 

944 dataId, kwargs = self._rewrite_data_id(dataId, datasetType, **kwargs) 

945 

946 if datasetType.isCalibration(): 

947 # Because this is a calibration dataset, first try to make a 

948 # standardize the data ID without restricting the dimensions to 

949 # those of the dataset type requested, because there may be extra 

950 # dimensions that provide temporal information for a validity-range 

951 # lookup. 

952 dataId = DataCoordinate.standardize( 

953 dataId, universe=self.dimensions, defaults=self._registry.defaults.dataId, **kwargs 

954 ) 

955 if timespan is None: 

956 if dataId.dimensions.temporal: 

957 dataId = self._registry.expandDataId(dataId) 

958 # Use the timespan from the data ID to constrain the 

959 # calibration lookup, but only if the caller has not 

960 # specified an explicit timespan. 

961 timespan = dataId.timespan 

962 else: 

963 # Try an arbitrary timespan. Downstream will fail if this 

964 # results in more than one matching dataset. 

965 timespan = Timespan(None, None) 

966 else: 

967 # Standardize the data ID to just the dimensions of the dataset 

968 # type instead of letting registry.findDataset do it, so we get the 

969 # result even if no dataset is found. 

970 dataId = DataCoordinate.standardize( 

971 dataId, 

972 dimensions=datasetType.dimensions, 

973 defaults=self._registry.defaults.dataId, 

974 **kwargs, 

975 ) 

976 # Always lookup the DatasetRef, even if one is given, to ensure it is 

977 # present in the current collection. 

978 ref = self.find_dataset( 

979 datasetType, 

980 dataId, 

981 collections=collections, 

982 timespan=timespan, 

983 datastore_records=datastore_records, 

984 ) 

985 if ref is None: 

986 if predict: 

987 if run is None: 987 ↛ 991line 987 didn't jump to line 991 because the condition on line 987 was always true

988 run = self.run 

989 if run is None: 989 ↛ 990line 989 didn't jump to line 990 because the condition on line 989 was never true

990 raise TypeError("Cannot predict dataset ID/location with run=None.") 

991 dataId = self.registry.expandDataId(dataId) 

992 return DatasetRef(datasetType, dataId, run=run) 

993 else: 

994 if collections is None: 

995 collections = self._registry.defaults.collections 

996 raise DatasetNotFoundError( 

997 f"Dataset {datasetType.name} with data ID {dataId} " 

998 f"could not be found in collections {collections}." 

999 ) 

1000 if datasetType != ref.datasetType: 

1001 # If they differ it is because the user explicitly specified 

1002 # a compatible dataset type to this call rather than using the 

1003 # registry definition. The DatasetRef must therefore be recreated 

1004 # using the user definition such that the expected type is 

1005 # returned. 

1006 ref = DatasetRef( 

1007 datasetType, ref.dataId, run=ref.run, id=ref.id, datastore_records=ref._datastore_records 

1008 ) 

1009 

1010 return ref 

1011 

1012 @transactional 

1013 def put( 

1014 self, 

1015 obj: Any, 

1016 datasetRefOrType: DatasetRef | DatasetType | str, 

1017 /, 

1018 dataId: DataId | None = None, 

1019 *, 

1020 run: str | None = None, 

1021 provenance: DatasetProvenance | None = None, 

1022 **kwargs: Any, 

1023 ) -> DatasetRef: 

1024 """Store and register a dataset. 

1025 

1026 Parameters 

1027 ---------- 

1028 obj : `object` 

1029 The dataset. 

1030 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str` 

1031 When `DatasetRef` is provided, ``dataId`` should be `None`. 

1032 Otherwise the `DatasetType` or name thereof. If a fully resolved 

1033 `DatasetRef` is given the run and ID are used directly. 

1034 dataId : `dict` or `DataCoordinate` 

1035 A `dict` of `Dimension` link name, value pairs that label the 

1036 `DatasetRef` within a Collection. When `None`, a `DatasetRef` 

1037 should be provided as the second argument. 

1038 run : `str`, optional 

1039 The name of the run the dataset should be added to, overriding 

1040 ``self.run``. Not used if a resolved `DatasetRef` is provided. 

1041 provenance : `DatasetProvenance` or `None`, optional 

1042 Any provenance that should be attached to the serialized dataset. 

1043 Not supported by all serialization mechanisms. 

1044 **kwargs 

1045 Additional keyword arguments used to augment or construct a 

1046 `DataCoordinate`. See `DataCoordinate.standardize` 

1047 parameters. Not used if a resolve `DatasetRef` is provided. 

1048 

1049 Returns 

1050 ------- 

1051 ref : `DatasetRef` 

1052 A reference to the stored dataset, updated with the correct id if 

1053 given. 

1054 

1055 Raises 

1056 ------ 

1057 TypeError 

1058 Raised if the butler is read-only or if no run has been provided. 

1059 """ 

1060 if isinstance(datasetRefOrType, DatasetRef): 

1061 # This is a direct put of predefined DatasetRef. 

1062 _LOG.debug("Butler put direct: %s", datasetRefOrType) 

1063 if run is not None: 1063 ↛ 1064line 1063 didn't jump to line 1064 because the condition on line 1063 was never true

1064 warnings.warn("Run collection is not used for DatasetRef", stacklevel=3) 

1065 

1066 with self._metrics.instrument_put(_LOG, msg="Dataset put direct"): 

1067 # If registry already has a dataset with the same dataset ID, 

1068 # dataset type and DataId, then _importDatasets will do 

1069 # nothing and just return an original ref. We have to raise in 

1070 # this case, there is a datastore check below for that. 

1071 self._registry._importDatasets([datasetRefOrType], expand=True) 

1072 # Before trying to write to the datastore check that it does 

1073 # not know this dataset. This is prone to races, of course. 

1074 if self._datastore.knows(datasetRefOrType): 

1075 raise ConflictingDefinitionError( 

1076 f"Datastore already contains dataset: {datasetRefOrType}" 

1077 ) 

1078 # Try to write dataset to the datastore, if it fails due to a 

1079 # race with another write, the content of stored data may be 

1080 # unpredictable. 

1081 try: 

1082 self._datastore.put(obj, datasetRefOrType, provenance=provenance) 

1083 except IntegrityError as e: 

1084 raise ConflictingDefinitionError(f"Datastore already contains dataset: {e}") from e 

1085 

1086 return datasetRefOrType 

1087 

1088 _LOG.debug("Butler put: %s, dataId=%s, run=%s", datasetRefOrType, dataId, run) 

1089 if not self.isWriteable(): 1089 ↛ 1090line 1089 didn't jump to line 1090 because the condition on line 1089 was never true

1090 raise TypeError("Butler is read-only.") 

1091 

1092 with self._metrics.instrument_put(_LOG, msg="Dataset put with dataID"): 

1093 datasetType, dataId = self._standardizeArgs(datasetRefOrType, dataId, **kwargs) 

1094 

1095 # Handle dimension records in dataId 

1096 dataId, kwargs = self._rewrite_data_id(dataId, datasetType, **kwargs) 

1097 

1098 # Add Registry Dataset entry. 

1099 dataId = self._registry.expandDataId(dataId, dimensions=datasetType.dimensions, **kwargs) 

1100 (ref,) = self._registry.insertDatasets(datasetType, run=run, dataIds=[dataId]) 

1101 self._datastore.put(obj, ref, provenance=provenance) 

1102 

1103 return ref 

1104 

1105 def getDeferred( 

1106 self, 

1107 datasetRefOrType: DatasetRef | DatasetType | str, 

1108 /, 

1109 dataId: DataId | None = None, 

1110 *, 

1111 parameters: dict | None = None, 

1112 collections: Any = None, 

1113 storageClass: str | StorageClass | None = None, 

1114 timespan: Timespan | None = None, 

1115 **kwargs: Any, 

1116 ) -> DeferredDatasetHandle: 

1117 """Create a `DeferredDatasetHandle` which can later retrieve a dataset, 

1118 after an immediate registry lookup. 

1119 

1120 Parameters 

1121 ---------- 

1122 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str` 

1123 When `DatasetRef` the `dataId` should be `None`. 

1124 Otherwise the `DatasetType` or name thereof. 

1125 dataId : `dict` or `DataCoordinate`, optional 

1126 A `dict` of `Dimension` link name, value pairs that label the 

1127 `DatasetRef` within a Collection. When `None`, a `DatasetRef` 

1128 should be provided as the first argument. 

1129 parameters : `dict` 

1130 Additional StorageClass-defined options to control reading, 

1131 typically used to efficiently read only a subset of the dataset. 

1132 collections : Any, optional 

1133 Collections to be searched, overriding ``self.collections``. 

1134 Can be any of the types supported by the ``collections`` argument 

1135 to butler construction. 

1136 storageClass : `StorageClass` or `str`, optional 

1137 The storage class to be used to override the Python type 

1138 returned by this method. By default the returned type matches 

1139 the dataset type definition for this dataset. Specifying a 

1140 read `StorageClass` can force a different type to be returned. 

1141 This type must be compatible with the original type. 

1142 timespan : `Timespan` or `None`, optional 

1143 A timespan that the validity range of the dataset must overlap. 

1144 If not provided and this is a calibration dataset type, an attempt 

1145 will be made to find the timespan from any temporal coordinate 

1146 in the data ID. 

1147 **kwargs 

1148 Additional keyword arguments used to augment or construct a 

1149 `DataId`. See `DataId` parameters. 

1150 

1151 Returns 

1152 ------- 

1153 obj : `DeferredDatasetHandle` 

1154 A handle which can be used to retrieve a dataset at a later time. 

1155 

1156 Raises 

1157 ------ 

1158 LookupError 

1159 Raised if no matching dataset exists in the `Registry` or 

1160 datastore. 

1161 ValueError 

1162 Raised if a resolved `DatasetRef` was passed as an input, but it 

1163 differs from the one found in the registry. 

1164 TypeError 

1165 Raised if no collections were provided. 

1166 """ 

1167 if isinstance(datasetRefOrType, DatasetRef): 

1168 # Do the quick check first and if that fails, check for artifact 

1169 # existence. This is necessary for datastores that are configured 

1170 # in trust mode where there won't be a record but there will be 

1171 # a file. 

1172 if self._datastore.knows(datasetRefOrType) or self._datastore.exists(datasetRefOrType): 

1173 ref = datasetRefOrType 

1174 else: 

1175 raise LookupError(f"Dataset reference {datasetRefOrType} does not exist.") 

1176 else: 

1177 ref = self._findDatasetRef( 

1178 datasetRefOrType, dataId, collections=collections, timespan=timespan, **kwargs 

1179 ) 

1180 return DeferredDatasetHandle(butler=self, ref=ref, parameters=parameters, storageClass=storageClass) 

1181 

1182 def get( 

1183 self, 

1184 datasetRefOrType: DatasetRef | DatasetType | str, 

1185 /, 

1186 dataId: DataId | None = None, 

1187 *, 

1188 parameters: dict[str, Any] | None = None, 

1189 collections: Any = None, 

1190 storageClass: StorageClass | str | None = None, 

1191 timespan: Timespan | None = None, 

1192 **kwargs: Any, 

1193 ) -> Any: 

1194 """Retrieve a stored dataset. 

1195 

1196 Parameters 

1197 ---------- 

1198 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str` 

1199 When `DatasetRef` the `dataId` should be `None`. 

1200 Otherwise the `DatasetType` or name thereof. 

1201 If a resolved `DatasetRef`, the associated dataset 

1202 is returned directly without additional querying. 

1203 dataId : `dict` or `DataCoordinate` 

1204 A `dict` of `Dimension` link name, value pairs that label the 

1205 `DatasetRef` within a Collection. When `None`, a `DatasetRef` 

1206 should be provided as the first argument. 

1207 parameters : `dict` 

1208 Additional StorageClass-defined options to control reading, 

1209 typically used to efficiently read only a subset of the dataset. 

1210 collections : Any, optional 

1211 Collections to be searched, overriding ``self.collections``. 

1212 Can be any of the types supported by the ``collections`` argument 

1213 to butler construction. 

1214 storageClass : `StorageClass` or `str`, optional 

1215 The storage class to be used to override the Python type 

1216 returned by this method. By default the returned type matches 

1217 the dataset type definition for this dataset. Specifying a 

1218 read `StorageClass` can force a different type to be returned. 

1219 This type must be compatible with the original type. 

1220 timespan : `Timespan` or `None`, optional 

1221 A timespan that the validity range of the dataset must overlap. 

1222 If not provided and this is a calibration dataset type, an attempt 

1223 will be made to find the timespan from any temporal coordinate 

1224 in the data ID. 

1225 **kwargs 

1226 Additional keyword arguments used to augment or construct a 

1227 `DataCoordinate`. See `DataCoordinate.standardize` 

1228 parameters. 

1229 

1230 Returns 

1231 ------- 

1232 obj : `object` 

1233 The dataset. 

1234 

1235 Raises 

1236 ------ 

1237 LookupError 

1238 Raised if no matching dataset exists in the `Registry`. 

1239 TypeError 

1240 Raised if no collections were provided. 

1241 

1242 Notes 

1243 ----- 

1244 When looking up datasets in a `~CollectionType.CALIBRATION` collection, 

1245 this method requires that the given data ID include temporal dimensions 

1246 beyond the dimensions of the dataset type itself, in order to find the 

1247 dataset with the appropriate validity range. For example, a "bias" 

1248 dataset with native dimensions ``{instrument, detector}`` could be 

1249 fetched with a ``{instrument, detector, exposure}`` data ID, because 

1250 ``exposure`` is a temporal dimension. 

1251 """ 

1252 _LOG.debug("Butler get: %s, dataId=%s, parameters=%s", datasetRefOrType, dataId, parameters) 

1253 with self._metrics.instrument_get(_LOG, msg="Retrieved dataset"): 

1254 ref = self._findDatasetRef( 

1255 datasetRefOrType, 

1256 dataId, 

1257 collections=collections, 

1258 datastore_records=True, 

1259 timespan=timespan, 

1260 **kwargs, 

1261 ) 

1262 return self._datastore.get(ref, parameters=parameters, storageClass=storageClass) 

1263 

1264 def getURIs( 

1265 self, 

1266 datasetRefOrType: DatasetRef | DatasetType | str, 

1267 /, 

1268 dataId: DataId | None = None, 

1269 *, 

1270 predict: bool = False, 

1271 collections: Any = None, 

1272 run: str | None = None, 

1273 **kwargs: Any, 

1274 ) -> DatasetRefURIs: 

1275 """Return the URIs associated with the dataset. 

1276 

1277 Parameters 

1278 ---------- 

1279 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str` 

1280 When `DatasetRef` the `dataId` should be `None`. 

1281 Otherwise the `DatasetType` or name thereof. 

1282 dataId : `dict` or `DataCoordinate` 

1283 A `dict` of `Dimension` link name, value pairs that label the 

1284 `DatasetRef` within a Collection. When `None`, a `DatasetRef` 

1285 should be provided as the first argument. 

1286 predict : `bool` 

1287 If `True`, allow URIs to be returned of datasets that have not 

1288 been written. 

1289 collections : Any, optional 

1290 Collections to be searched, overriding ``self.collections``. 

1291 Can be any of the types supported by the ``collections`` argument 

1292 to butler construction. 

1293 run : `str`, optional 

1294 Run to use for predictions, overriding ``self.run``. 

1295 **kwargs 

1296 Additional keyword arguments used to augment or construct a 

1297 `DataCoordinate`. See `DataCoordinate.standardize` 

1298 parameters. 

1299 

1300 Returns 

1301 ------- 

1302 uris : `DatasetRefURIs` 

1303 The URI to the primary artifact associated with this dataset (if 

1304 the dataset was disassembled within the datastore this may be 

1305 `None`), and the URIs to any components associated with the dataset 

1306 artifact. (can be empty if there are no components). 

1307 """ 

1308 ref = self._findDatasetRef( 

1309 datasetRefOrType, dataId, predict=predict, run=run, collections=collections, **kwargs 

1310 ) 

1311 return self._datastore.getURIs(ref, predict) 

1312 

1313 def get_dataset_type(self, name: str) -> DatasetType: 

1314 return self._registry.getDatasetType(name) 

1315 

1316 def get_dataset( 

1317 self, 

1318 id: DatasetId | str, 

1319 *, 

1320 storage_class: str | StorageClass | None = None, 

1321 dimension_records: bool = False, 

1322 datastore_records: bool = False, 

1323 ) -> DatasetRef | None: 

1324 id = _to_uuid(id) 

1325 ref = self._registry.getDataset(id) 

1326 if ref is not None: 

1327 if dimension_records: 1327 ↛ 1328line 1327 didn't jump to line 1328 because the condition on line 1327 was never true

1328 ref = ref.expanded( 

1329 self._registry.expandDataId(ref.dataId, dimensions=ref.datasetType.dimensions) 

1330 ) 

1331 if storage_class: 1331 ↛ 1332line 1331 didn't jump to line 1332 because the condition on line 1331 was never true

1332 ref = ref.overrideStorageClass(storage_class) 

1333 if datastore_records: 1333 ↛ 1334line 1333 didn't jump to line 1334 because the condition on line 1333 was never true

1334 ref = self._registry.get_datastore_records(ref) 

1335 return ref 

1336 

1337 def get_many_datasets(self, ids: Iterable[DatasetId | str]) -> list[DatasetRef]: 

1338 uuids = [_to_uuid(id) for id in ids] 

1339 return self._registry._managers.datasets.get_dataset_refs(uuids) 

1340 

1341 def find_dataset( 

1342 self, 

1343 dataset_type: DatasetType | str, 

1344 data_id: DataId | None = None, 

1345 *, 

1346 collections: str | Sequence[str] | None = None, 

1347 timespan: Timespan | None = None, 

1348 storage_class: str | StorageClass | None = None, 

1349 dimension_records: bool = False, 

1350 datastore_records: bool = False, 

1351 **kwargs: Any, 

1352 ) -> DatasetRef | None: 

1353 # Handle any parts of the dataID that are not using primary dimension 

1354 # keys. 

1355 if isinstance(dataset_type, str): 

1356 actual_type = self.get_dataset_type(dataset_type) 

1357 else: 

1358 actual_type = dataset_type 

1359 

1360 # Store the component for later. 

1361 component_name = actual_type.component() 

1362 if actual_type.isComponent(): 

1363 parent_type = actual_type.makeCompositeDatasetType() 

1364 else: 

1365 parent_type = actual_type 

1366 

1367 data_id, kwargs = self._rewrite_data_id(data_id, parent_type, **kwargs) 

1368 

1369 ref = self.registry.findDataset( 

1370 parent_type, 

1371 data_id, 

1372 collections=collections, 

1373 timespan=timespan, 

1374 datastore_records=datastore_records, 

1375 **kwargs, 

1376 ) 

1377 if ref is not None and dimension_records: 1377 ↛ 1378line 1377 didn't jump to line 1378 because the condition on line 1377 was never true

1378 ref = ref.expanded(self._registry.expandDataId(ref.dataId, dimensions=ref.datasetType.dimensions)) 

1379 if ref is not None and component_name: 

1380 ref = ref.makeComponentRef(component_name) 

1381 if ref is not None and storage_class is not None: 1381 ↛ 1382line 1381 didn't jump to line 1382 because the condition on line 1381 was never true

1382 ref = ref.overrideStorageClass(storage_class) 

1383 

1384 return ref 

1385 

1386 def retrieve_artifacts_zip( 

1387 self, 

1388 refs: Iterable[DatasetRef], 

1389 destination: ResourcePathExpression, 

1390 overwrite: bool = True, 

1391 ) -> ResourcePath: 

1392 return retrieve_and_zip(refs, destination, self._datastore.retrieveArtifacts, overwrite) 

1393 

1394 def retrieveArtifacts( 

1395 self, 

1396 refs: Iterable[DatasetRef], 

1397 destination: ResourcePathExpression, 

1398 transfer: str = "auto", 

1399 preserve_path: bool = True, 

1400 overwrite: bool = False, 

1401 ) -> list[ResourcePath]: 

1402 # Docstring inherited. 

1403 outdir = ResourcePath(destination) 

1404 artifact_map = self._datastore.retrieveArtifacts( 

1405 refs, 

1406 outdir, 

1407 transfer=transfer, 

1408 preserve_path=preserve_path, 

1409 overwrite=overwrite, 

1410 write_index=True, 

1411 ) 

1412 return list(artifact_map) 

1413 

1414 def exists( 

1415 self, 

1416 dataset_ref_or_type: DatasetRef | DatasetType | str, 

1417 /, 

1418 data_id: DataId | None = None, 

1419 *, 

1420 full_check: bool = True, 

1421 collections: Any = None, 

1422 **kwargs: Any, 

1423 ) -> DatasetExistence: 

1424 # Docstring inherited. 

1425 existence = DatasetExistence.UNRECOGNIZED 

1426 

1427 if isinstance(dataset_ref_or_type, DatasetRef): 

1428 if collections is not None: 1428 ↛ 1429line 1428 didn't jump to line 1429 because the condition on line 1428 was never true

1429 warnings.warn("Collections should not be specified with DatasetRef", stacklevel=2) 

1430 if data_id is not None: 1430 ↛ 1431line 1430 didn't jump to line 1431 because the condition on line 1430 was never true

1431 warnings.warn("A DataID should not be specified with DatasetRef", stacklevel=2) 

1432 ref = dataset_ref_or_type 

1433 registry_ref = self._registry.getDataset(dataset_ref_or_type.id) 

1434 if registry_ref is not None: 

1435 existence |= DatasetExistence.RECORDED 

1436 

1437 if dataset_ref_or_type != registry_ref: 

1438 # This could mean that storage classes differ, so we should 

1439 # check for that but use the registry ref for the rest of 

1440 # the method. 

1441 if registry_ref.is_compatible_with(dataset_ref_or_type): 

1442 # Use the registry version from now on. 

1443 ref = registry_ref 

1444 else: 

1445 raise ValueError( 

1446 f"The ref given to exists() ({ref}) has the same dataset ID as one " 

1447 f"in registry but has different incompatible values ({registry_ref})." 

1448 ) 

1449 else: 

1450 try: 

1451 ref = self._findDatasetRef(dataset_ref_or_type, data_id, collections=collections, **kwargs) 

1452 except (LookupError, TypeError): 

1453 return existence 

1454 existence |= DatasetExistence.RECORDED 

1455 

1456 if self._datastore.knows(ref): 

1457 existence |= DatasetExistence.DATASTORE 

1458 

1459 if full_check: 

1460 if self._datastore.exists(ref): 

1461 existence |= DatasetExistence._ARTIFACT 

1462 elif existence.value != DatasetExistence.UNRECOGNIZED.value: 

1463 # Do not add this flag if we have no other idea about a dataset. 

1464 existence |= DatasetExistence(DatasetExistence._ASSUMED) 

1465 

1466 return existence 

1467 

1468 def _exists_many( 

1469 self, 

1470 refs: Iterable[DatasetRef], 

1471 /, 

1472 *, 

1473 full_check: bool = True, 

1474 ) -> dict[DatasetRef, DatasetExistence]: 

1475 # Docstring inherited. 

1476 existence = {ref: DatasetExistence.UNRECOGNIZED for ref in refs} 

1477 

1478 # Check which refs exist in the registry. 

1479 id_map = {ref.id: ref for ref in existence.keys()} 

1480 for registry_ref in self.get_many_datasets(id_map.keys()): 

1481 # Consistency between the given DatasetRef and the information 

1482 # recorded in the registry is not verified. 

1483 existence[id_map[registry_ref.id]] |= DatasetExistence.RECORDED 

1484 

1485 # Ask datastore if it knows about these refs. 

1486 knows = self._datastore.knows_these(refs) 

1487 for ref, known in knows.items(): 

1488 if known: 

1489 existence[ref] |= DatasetExistence.DATASTORE 

1490 

1491 if full_check: 

1492 mexists = self._datastore.mexists(refs) 

1493 for ref, exists in mexists.items(): 

1494 if exists: 

1495 existence[ref] |= DatasetExistence._ARTIFACT 

1496 else: 

1497 # Do not set this flag if nothing is known about the dataset. 

1498 for ref in existence: 

1499 if existence[ref] != DatasetExistence.UNRECOGNIZED: 

1500 existence[ref] |= DatasetExistence._ASSUMED 

1501 

1502 return existence 

1503 

1504 def removeRuns( 

1505 self, 

1506 names: Iterable[str], 

1507 unstore: bool | type[_DeprecatedDefault] = _DeprecatedDefault, 

1508 *, 

1509 unlink_from_chains: bool = False, 

1510 ) -> None: 

1511 # Docstring inherited. 

1512 if not self.isWriteable(): 1512 ↛ 1513line 1512 didn't jump to line 1513 because the condition on line 1512 was never true

1513 raise TypeError("Butler is read-only.") 

1514 

1515 if unstore is not _DeprecatedDefault: 1515 ↛ 1518line 1515 didn't jump to line 1518 because the condition on line 1515 was never true

1516 # The value was passed in by a user. Must report it is now 

1517 # ignored. 

1518 if unstore is True: 

1519 msg = "The unstore parameter is deprecated and is now always treated as True. " 

1520 else: 

1521 msg = "The unstore parameter for removeRuns can no longer be False and is now ignored. " 

1522 warnings.warn( 

1523 msg + " The parameter will be removed after v30.", 

1524 category=FutureWarning, 

1525 stacklevel=find_outside_stacklevel("lsst.daf.butler"), 

1526 ) 

1527 

1528 names = list(names) 

1529 refs: list[DatasetRef] = [] 

1530 # Map of the chained collections to the RUN children. 

1531 parents_to_children: dict[str, set[str]] = defaultdict(set) 

1532 

1533 with self._caching_context(): 

1534 # Get information about these RUNs. 

1535 collections_info = self.collections.query_info(names, include_parents=unlink_from_chains) 

1536 for info in collections_info: 

1537 if info.type is not CollectionType.RUN: 

1538 raise TypeError(f"The collection type of '{info.name}' is {info.type.name}, not RUN.") 

1539 if unlink_from_chains: 

1540 if info.parents is None: # For mypy. 

1541 raise AssertionError("Internal error: Collection parents required but not received") 

1542 for parent in info.parents: 

1543 parents_to_children[parent].add(info.name) 

1544 

1545 # Update the names in case the query unexpectedly had a wildcard. 

1546 names = [info.name for info in collections_info] 

1547 

1548 # Get all the datasets from these runs. 

1549 refs = self.query_all_datasets(names, find_first=False, limit=None) 

1550 

1551 # Call pruneDatasets since we are deliberately removing 

1552 # datasets in chunks from the RUN collections rather than 

1553 # attempting to remove everything at once. 

1554 with time_this( 

1555 _LOG, 

1556 msg="Removing %d dataset%s from %s", 

1557 args=(len(refs), "s" if len(refs) != 1 else "", ", ".join(names)), 

1558 ): 

1559 self.pruneDatasets(refs, unstore=True, purge=True, disassociate=True) 

1560 

1561 # Now can remove the actual RUN collection and unlink from chains. 

1562 with self._registry.transaction(): 

1563 # This will fail if caller is not unlinking from chains but the 

1564 # RUN is in a chain -- but we have already deleted all the datasets 

1565 # by this point. 

1566 if unlink_from_chains: 

1567 # Use deterministic order for deletions to attempt to minimize 

1568 # risk of deadlocks for parallel deletes. 

1569 for parent in sorted(parents_to_children): 

1570 self.collections.remove_from_chain(parent, sorted(parents_to_children[parent])) 

1571 # Sort to avoid potential deadlocks. 

1572 for name in sorted(names): 

1573 # This should be fast since the collection should be empty. 

1574 with time_this(_LOG, msg="Removing RUN collection %s", args=(name,)): 

1575 self._registry.removeCollection(name) 

1576 _LOG.info("Completely removed the following RUN collections: %s", ", ".join(names)) 

1577 

1578 def pruneDatasets( 

1579 self, 

1580 refs: Iterable[DatasetRef], 

1581 *, 

1582 disassociate: bool = True, 

1583 unstore: bool = False, 

1584 tags: Iterable[str] = (), 

1585 purge: bool = False, 

1586 ) -> None: 

1587 # docstring inherited from LimitedButler 

1588 

1589 if not self.isWriteable(): 1589 ↛ 1590line 1589 didn't jump to line 1590 because the condition on line 1589 was never true

1590 raise TypeError("Butler is read-only.") 

1591 if purge: 

1592 if not disassociate: 1592 ↛ 1593line 1592 didn't jump to line 1593 because the condition on line 1592 was never true

1593 raise TypeError("Cannot pass purge=True without disassociate=True.") 

1594 if not unstore: 1594 ↛ 1595line 1594 didn't jump to line 1595 because the condition on line 1594 was never true

1595 raise TypeError("Cannot pass purge=True without unstore=True.") 

1596 elif disassociate: 

1597 tags = tuple(tags) 

1598 if not tags: 1598 ↛ 1599line 1598 didn't jump to line 1599 because the condition on line 1598 was never true

1599 raise TypeError("No tags provided but disassociate=True.") 

1600 for tag in tags: 

1601 collectionType = self._registry.getCollectionType(tag) 

1602 if collectionType is not CollectionType.TAGGED: 1602 ↛ 1603line 1602 didn't jump to line 1603 because the condition on line 1602 was never true

1603 raise TypeError( 

1604 f"Cannot disassociate from collection '{tag}' " 

1605 f"of non-TAGGED type {collectionType.name}." 

1606 ) 

1607 # Transform possibly-single-pass iterable into something we can iterate 

1608 # over multiple times. 

1609 refs = list(refs) 

1610 # Pruning a component of a DatasetRef makes no sense since registry 

1611 # doesn't know about components and datastore might not store 

1612 # components in a separate file 

1613 for ref in refs: 

1614 if ref.datasetType.component(): 1614 ↛ 1615line 1614 didn't jump to line 1615 because the condition on line 1614 was never true

1615 raise ValueError(f"Can not prune a component of a dataset (ref={ref})") 

1616 

1617 # Chunk the deletions using a size that is reasonably efficient whilst 

1618 # also giving reasonable feedback to the user. Chunking also minimizes 

1619 # what needs to rollback if there is a failure and should allow 

1620 # incremental re-running of the pruning (assuming the query is 

1621 # repeated). The only issue will be if the Ctrl-C comes during 

1622 # emptyTrash since an admin command would need to run to finish the 

1623 # emptying of that chunk. 

1624 progress = Progress("lsst.daf.butler.Butler.pruneDatasets", level=_LOG.INFO) 

1625 chunk_size = 50_000 

1626 n_chunks = math.ceil(len(refs) / chunk_size) 

1627 if n_chunks > 1: 1627 ↛ 1628line 1627 didn't jump to line 1628 because the condition on line 1627 was never true

1628 _LOG.verbose("Pruning a total of %d datasets", len(refs)) 

1629 chunk_num = 0 

1630 for chunked_refs in progress.wrap( 

1631 chunk_iterable(refs, chunk_size=chunk_size), desc="Deleting datasets", total=n_chunks 

1632 ): 

1633 chunk_num += 1 

1634 _LOG.verbose( 

1635 "Pruning %d dataset%s in chunk %d/%d", 

1636 len(chunked_refs), 

1637 "s" if len(chunked_refs) != 1 else "", 

1638 chunk_num, 

1639 n_chunks, 

1640 ) 

1641 with time_this( 

1642 _LOG, 

1643 msg="Removing %d datasets for chunk %d/%d", 

1644 args=(len(chunked_refs), chunk_num, n_chunks), 

1645 ): 

1646 self._prune_datasets( 

1647 chunked_refs, tags=tags, unstore=unstore, purge=purge, disassociate=disassociate 

1648 ) 

1649 

1650 def _prune_datasets( 

1651 self, 

1652 refs: Collection[DatasetRef], 

1653 *, 

1654 disassociate: bool = True, 

1655 unstore: bool = False, 

1656 tags: Iterable[str] = (), 

1657 purge: bool = False, 

1658 ) -> None: 

1659 # We don't need an unreliable Datastore transaction for this, because 

1660 # we've been extra careful to ensure that Datastore.trash only involves 

1661 # mutating the Registry (it can _look_ at Datastore-specific things, 

1662 # but shouldn't change them), and hence all operations here are 

1663 # Registry operations. 

1664 with self.transaction(): 

1665 plural = "s" if len(refs) != 1 else "" 

1666 if unstore: 

1667 with time_this( 

1668 _LOG, 

1669 msg="Marking %d dataset%s for removal during pruneDatasets", 

1670 args=(len(refs), plural), 

1671 ): 

1672 self._datastore.trash(refs) 

1673 if purge: 

1674 with time_this( 

1675 _LOG, msg="Removing %d pruned dataset%s from registry", args=(len(refs), plural) 

1676 ): 

1677 self._registry.removeDatasets(refs) 

1678 elif disassociate: 

1679 assert tags, "Guaranteed by earlier logic in this function." 

1680 with time_this( 

1681 _LOG, msg="Disassociating %d dataset%ss from tagged collections", args=(len(refs), plural) 

1682 ): 

1683 for tag in tags: 

1684 self._registry.disassociate(tag, refs) 

1685 # We've exited the Registry transaction, and apparently committed. 

1686 # (if there was an exception, everything rolled back, and it's as if 

1687 # nothing happened - and we never get here). 

1688 # Datastore artifacts are not yet gone, but they're clearly marked 

1689 # as trash, so if we fail to delete now because of (e.g.) filesystem 

1690 # problems we can try again later, and if manual administrative 

1691 # intervention is required, it's pretty clear what that should entail: 

1692 # deleting everything on disk and in private Datastore tables that is 

1693 # in the dataset_location_trash table. 

1694 if unstore: 

1695 # Point of no return for removing artifacts. Restrict the trash 

1696 # emptying to the refs that this call trashed. 

1697 with time_this( 

1698 _LOG, 

1699 msg="Attempting to remove artifacts for %d dataset%s associated with pruning", 

1700 args=(len(refs), plural), 

1701 ): 

1702 self._datastore.emptyTrash(refs=refs) 

1703 

1704 def ingest_zip( 

1705 self, 

1706 zip_file: ResourcePathExpression, 

1707 transfer: str = "auto", 

1708 *, 

1709 transfer_dimensions: bool = False, 

1710 dry_run: bool = False, 

1711 skip_existing: bool = False, 

1712 ) -> None: 

1713 # Docstring inherited. 

1714 if not self.isWriteable(): 1714 ↛ 1715line 1714 didn't jump to line 1715 because the condition on line 1714 was never true

1715 raise TypeError("Butler is read-only.") 

1716 

1717 zip_path = ResourcePath(zip_file) 

1718 index = ZipIndex.from_zip_file(zip_path) 

1719 _LOG.verbose( 

1720 "Ingesting %s containing %d datasets and %d files.", zip_path, len(index.refs), len(index) 

1721 ) 

1722 

1723 # Need to ingest the refs into registry. Re-use the standard ingest 

1724 # code by reconstructing FileDataset from the index. 

1725 refs = index.refs.to_refs(universe=self.dimensions) 

1726 id_to_ref = {ref.id: ref for ref in refs} 

1727 datasets: list[FileDataset] = [] 

1728 processed_ids: set[uuid.UUID] = set() 

1729 for path_in_zip, index_info in index.artifact_map.items(): 

1730 # Disassembled composites need to check this ref isn't already 

1731 # included. 

1732 unprocessed = {id_ for id_ in index_info.ids if id_ not in processed_ids} 

1733 if not unprocessed: 1733 ↛ 1734line 1733 didn't jump to line 1734 because the condition on line 1733 was never true

1734 continue 

1735 dataset = FileDataset(refs=[id_to_ref[id_] for id_ in unprocessed], path=path_in_zip) 

1736 datasets.append(dataset) 

1737 processed_ids.update(unprocessed) 

1738 

1739 new_datasets, existing_datasets = self._partition_datasets_by_known(datasets) 

1740 if existing_datasets: 

1741 if skip_existing: 

1742 _LOG.info( 

1743 "Skipping %d datasets from zip file %s which already exist in the repository.", 

1744 len(existing_datasets), 

1745 zip_file, 

1746 ) 

1747 else: 

1748 raise ConflictingDefinitionError( 

1749 f"Datastore already contains {len(existing_datasets)} of the given datasets." 

1750 f" Example: {existing_datasets[0]}" 

1751 ) 

1752 if new_datasets: 1752 ↛ 1755line 1752 didn't jump to line 1755 because the condition on line 1752 was never true

1753 # Can not yet support partial zip ingests where a zip contains 

1754 # some datasets that are already in another zip. 

1755 raise ValueError( 

1756 f"The given zip file from {zip_file} contains {len(new_datasets)} datasets not known " 

1757 f"to this butler but also contains {len(existing_datasets)} datasets already known to " 

1758 "this butler. Currently butler can not ingest zip files with overlapping content." 

1759 ) 

1760 return 

1761 

1762 # Ingest doesn't create the RUN collections so we have to do that 

1763 # here. 

1764 # 

1765 # Sort by run collection name to ensure Postgres takes locks in the 

1766 # same order between different processes, to mitigate an issue 

1767 # where Postgres can deadlock due to the unique index on collection 

1768 # name. (See DM-47543). 

1769 runs = {ref.run for ref in refs} 

1770 for run in sorted(runs): 

1771 registered = self.collections.register(run) 

1772 if registered: 

1773 _LOG.verbose("Created RUN collection %s as part of zip ingest", run) 

1774 

1775 progress = Progress("lsst.daf.butler.Butler.ingest", level=VERBOSE) 

1776 import_info = self._prepare_ingest_file_datasets( 

1777 datasets, progress, dry_run=dry_run, transfer_dimensions=transfer_dimensions 

1778 ) 

1779 

1780 # Calculate some statistics based on the given list of datasets. 

1781 n_datasets = 0 

1782 for d in datasets: 

1783 n_datasets += len(d.refs) 

1784 srefs = "s" if n_datasets != 1 else "" 

1785 

1786 with ( 

1787 self._metrics.instrument_ingest( 

1788 n_datasets, 

1789 _LOG, 

1790 msg="Ingesting zip file %s with %s dataset%s", 

1791 args=(zip_file, n_datasets, srefs), 

1792 ), 

1793 self.transaction(), 

1794 ): 

1795 # Do not need expanded dataset refs so can ignore the return value. 

1796 self._ingest_file_datasets(datasets, import_info, progress, dry_run=dry_run) 

1797 

1798 try: 

1799 self._datastore.ingest_zip(zip_path, transfer=transfer, dry_run=dry_run) 

1800 except IntegrityError as e: 

1801 raise ConflictingDefinitionError( 

1802 f"Datastore already contains one or more datasets: {e}" 

1803 ) from e 

1804 

1805 def _prepare_ingest_file_datasets( 

1806 self, 

1807 datasets: Sequence[FileDataset], 

1808 progress: Progress, 

1809 *, 

1810 transfer_dimensions: bool = False, 

1811 dry_run: bool = False, 

1812 ) -> _ImportDatasetsInfo: 

1813 # Track DataIDs that are being ingested so we can spot issues early 

1814 # with duplication. Retain previous FileDataset so we can report it. 

1815 groupedDataIds: MutableMapping[tuple[DatasetType, str], dict[DataCoordinate, FileDataset]] = ( 

1816 defaultdict(dict) 

1817 ) 

1818 

1819 # All the refs we need to import. 

1820 refs: list[DatasetRef] = [] 

1821 

1822 for dataset in progress.wrap(datasets, desc="Validating dataIDs"): 

1823 for ref in dataset.refs: 

1824 group_key = (ref.datasetType, ref.run) 

1825 

1826 if ref.dataId in groupedDataIds[group_key]: 1826 ↛ 1827line 1826 didn't jump to line 1827 because the condition on line 1826 was never true

1827 raise ConflictingDefinitionError( 

1828 f"Ingest conflict. Dataset {dataset.path} has same" 

1829 " DataId as other ingest dataset" 

1830 f" {groupedDataIds[group_key][ref.dataId].path} " 

1831 f" ({ref.dataId})" 

1832 ) 

1833 

1834 groupedDataIds[group_key][ref.dataId] = dataset 

1835 refs.extend(dataset.refs) 

1836 

1837 # Ensure that dataset types are created and all ref information 

1838 # extracted. 

1839 import_info = self._prepare_for_import_refs( 

1840 self, 

1841 refs, 

1842 register_dataset_types=True, 

1843 dry_run=dry_run, 

1844 transfer_dimensions=transfer_dimensions, 

1845 ) 

1846 return import_info 

1847 

1848 def _ingest_file_datasets( 

1849 self, 

1850 datasets: Sequence[FileDataset], 

1851 import_info: _ImportDatasetsInfo, 

1852 progress: Progress, 

1853 *, 

1854 dry_run: bool = False, 

1855 ) -> None: 

1856 self._import_dimension_records(import_info.dimension_records, dry_run=dry_run) 

1857 imported_refs = self._import_grouped_refs( 

1858 import_info.grouped_refs, None, progress, dry_run=dry_run, expand_refs=True 

1859 ) 

1860 

1861 # The expanded refs need to be attached back to the original 

1862 # FileDatasets for datastore to use. 

1863 id_to_ref = {ref.id: ref for ref in imported_refs} 

1864 

1865 for dataset in progress.wrap(datasets, desc="Re-attaching expanded refs"): 

1866 dataset.refs = [id_to_ref[ref.id] for ref in dataset.refs] 

1867 

1868 def ingest( 

1869 self, 

1870 *datasets: FileDataset, 

1871 transfer: str | None = "auto", 

1872 record_validation_info: bool = True, 

1873 skip_existing: bool = False, 

1874 ) -> None: 

1875 # Docstring inherited. 

1876 if not datasets: 

1877 return 

1878 if not self.isWriteable(): 1878 ↛ 1879line 1878 didn't jump to line 1879 because the condition on line 1878 was never true

1879 raise TypeError("Butler is read-only.") 

1880 _LOG.verbose("Ingesting %d file dataset%s.", len(datasets), "" if len(datasets) == 1 else "s") 

1881 progress = Progress("lsst.daf.butler.Butler.ingest", level=VERBOSE) 

1882 

1883 new_datasets, existing_datasets = self._partition_datasets_by_known(datasets) 

1884 if existing_datasets: 

1885 if skip_existing: 

1886 _LOG.info( 

1887 "Skipping %d datasets which already exist in the repository.", len(existing_datasets) 

1888 ) 

1889 else: 

1890 raise ConflictingDefinitionError( 

1891 f"Datastore already contains {len(existing_datasets)} of the given datasets." 

1892 f" Example: {existing_datasets[0]}" 

1893 ) 

1894 

1895 # Calculate some statistics based on the given list of datasets. 

1896 n_files = len(datasets) 

1897 n_datasets = 0 

1898 for d in datasets: 

1899 n_datasets += len(d.refs) 

1900 sfiles = "s" if n_files != 1 else "" 

1901 srefs = "s" if n_datasets != 1 else "" 

1902 

1903 # We use `datasets` rather `new_datasets` for the Registry 

1904 # portion of this, to let it confirm that everything matches the 

1905 # existing datasets. 

1906 import_info = self._prepare_ingest_file_datasets(datasets, progress) 

1907 

1908 with ( 

1909 self._metrics.instrument_ingest( 

1910 n_datasets, 

1911 _LOG, 

1912 msg="Ingesting %s file%s with %s dataset%s", 

1913 args=(n_files, sfiles, n_datasets, srefs), 

1914 ), 

1915 self.transaction(), 

1916 ): 

1917 self._ingest_file_datasets(datasets, import_info, progress) 

1918 

1919 # Bulk-insert everything into Datastore. 

1920 # We do not know if any of the registry entries already existed 

1921 # (_importDatasets only complains if they exist but differ). 

1922 # The _partition_datasets_by_known logic above should catch most 

1923 # instances where we attempt to re-ingest files that were already 

1924 # ingested, but a concurrent writer could cause a unique constraint 

1925 # violation here. 

1926 try: 

1927 self._datastore.ingest( 

1928 *new_datasets, transfer=transfer, record_validation_info=record_validation_info 

1929 ) 

1930 except IntegrityError as e: 

1931 raise ConflictingDefinitionError( 

1932 f"Datastore already contains one or more datasets: {e}" 

1933 ) from e 

1934 

1935 def _partition_datasets_by_known( 

1936 self, datasets: Iterable[FileDataset] 

1937 ) -> tuple[list[FileDataset], list[FileDataset]]: 

1938 """Divides the given `FileDataset` objects into two groups: those for 

1939 which the Datastore already has an entry, and those for which it does 

1940 not. 

1941 """ 

1942 new_datasets = [] 

1943 existing_datasets = [] 

1944 

1945 refs = itertools.chain.from_iterable(dataset.refs for dataset in datasets) 

1946 known_refs = self._datastore.knows_these(refs) 

1947 

1948 for dataset in datasets: 

1949 if any(known_refs[ref] for ref in dataset.refs): 

1950 existing_datasets.append(dataset) 

1951 else: 

1952 new_datasets.append(dataset) 

1953 

1954 return new_datasets, existing_datasets 

1955 

1956 @contextlib.contextmanager 

1957 def export( 

1958 self, 

1959 *, 

1960 directory: str | None = None, 

1961 filename: str | None = None, 

1962 format: str | None = None, 

1963 transfer: str | None = None, 

1964 ) -> Iterator[RepoExportContext]: 

1965 # Docstring inherited. 

1966 if directory is None and transfer is not None: 1966 ↛ 1967line 1966 didn't jump to line 1967 because the condition on line 1966 was never true

1967 raise TypeError("Cannot transfer without providing a directory.") 

1968 if transfer == "move": 1968 ↛ 1969line 1968 didn't jump to line 1969 because the condition on line 1968 was never true

1969 raise TypeError("Transfer may not be 'move': export is read-only") 

1970 if format is None: 

1971 if filename is None: 1971 ↛ 1972line 1971 didn't jump to line 1972 because the condition on line 1971 was never true

1972 raise TypeError("At least one of 'filename' or 'format' must be provided.") 

1973 else: 

1974 _, format = os.path.splitext(filename) 

1975 if not format: 

1976 raise ValueError("Please specify a file extension to determine export format.") 

1977 format = format[1:] # Strip leading "."" 

1978 elif filename is None: 1978 ↛ 1980line 1978 didn't jump to line 1980 because the condition on line 1978 was always true

1979 filename = f"export.{format}" 

1980 if directory is not None: 

1981 filename = os.path.join(directory, filename) 

1982 formats = self._config["repo_transfer_formats"] 

1983 if format not in formats: 

1984 raise ValueError(f"Unknown export format {format!r}, allowed: {','.join(formats.keys())}") 

1985 BackendClass = get_class_of(formats[format, "export"]) 

1986 with open(filename, "w") as stream: 

1987 backend = BackendClass(stream, universe=self.dimensions) 

1988 try: 

1989 helper = RepoExportContext(self, backend=backend, directory=directory, transfer=transfer) 

1990 with self._caching_context(): 

1991 yield helper 

1992 except BaseException: 

1993 raise 

1994 else: 

1995 helper._finish() 

1996 

1997 def import_( 

1998 self, 

1999 *, 

2000 directory: ResourcePathExpression | None = None, 

2001 filename: ResourcePathExpression | TextIO | None = None, 

2002 format: str | None = None, 

2003 transfer: str | None = None, 

2004 skip_dimensions: set | None = None, 

2005 record_validation_info: bool = True, 

2006 without_datastore: bool = False, 

2007 ) -> None: 

2008 # Docstring inherited. 

2009 if not self.isWriteable(): 2009 ↛ 2010line 2009 didn't jump to line 2010 because the condition on line 2009 was never true

2010 raise TypeError("Butler is read-only.") 

2011 if filename is None and format is not None: 2011 ↛ 2012line 2011 didn't jump to line 2012 because the condition on line 2011 was never true

2012 filename = ResourcePath(f"export.{format}", forceAbsolute=False) 

2013 if directory is not None: 

2014 directory = ResourcePath(directory, forceDirectory=True) 

2015 # mypy doesn't think this will work but it does in python >= 3.10. 

2016 if isinstance(filename, ResourcePathExpression): # type: ignore 

2017 filename = ResourcePath(filename, forceAbsolute=False) # type: ignore 

2018 if format is None: 2018 ↛ 2020line 2018 didn't jump to line 2020 because the condition on line 2018 was always true

2019 format = filename.getExtension() 

2020 if not filename.isabs() and directory is not None: 2020 ↛ 2021line 2020 didn't jump to line 2021 because the condition on line 2020 was never true

2021 potential = directory.join(filename) 

2022 exists_in_cwd = filename.exists() 

2023 exists_in_dir = potential.exists() 

2024 if exists_in_cwd and exists_in_dir: 

2025 _LOG.warning( 

2026 "A relative path for filename was specified (%s) which exists relative to cwd. " 

2027 "Additionally, the file exists relative to the given search directory (%s). " 

2028 "Using the export file in the given directory.", 

2029 filename, 

2030 potential, 

2031 ) 

2032 # Given they specified an explicit directory and that 

2033 # directory has the export file in it, assume that that 

2034 # is what was meant despite the file in cwd. 

2035 filename = potential 

2036 elif exists_in_dir: 

2037 filename = potential 

2038 elif not exists_in_cwd and not exists_in_dir: 

2039 # Raise early. 

2040 raise FileNotFoundError( 

2041 f"Export file could not be found in {filename.abspath()} or {potential.abspath()}." 

2042 ) 

2043 elif format is None: 2043 ↛ 2044line 2043 didn't jump to line 2044 because the condition on line 2043 was never true

2044 format = ".yaml" 

2045 BackendClass: type[RepoImportBackend] = get_class_of( 

2046 self._config["repo_transfer_formats"][format]["import"] 

2047 ) 

2048 

2049 def doImport(importStream: TextIO | ResourceHandleProtocol) -> None: 

2050 with self._caching_context(): 

2051 backend = BackendClass(importStream, self) # type: ignore[call-arg] 

2052 backend.register() 

2053 with self.transaction(): 

2054 backend.load( 

2055 datastore=self._datastore if not without_datastore else None, 

2056 directory=directory, 

2057 transfer=transfer, 

2058 skip_dimensions=skip_dimensions, 

2059 record_validation_info=record_validation_info, 

2060 ) 

2061 

2062 if isinstance(filename, ResourcePath): 

2063 # We can not use open() here at the moment because of 

2064 # DM-38589 since yaml does stream.read(8192) in a loop. 

2065 stream = io.StringIO(filename.read().decode()) 

2066 doImport(stream) 

2067 else: 

2068 doImport(filename) # type: ignore 

2069 

2070 def transfer_dimension_records_from( 

2071 self, source_butler: LimitedButler | Butler, source_refs: Iterable[DatasetRef | DataCoordinate] 

2072 ) -> None: 

2073 # Allowed dimensions in the target butler. 

2074 elements = frozenset(element for element in self.dimensions.elements if element.has_own_table) 

2075 

2076 data_ids = {ref.dataId for ref in source_refs} 

2077 

2078 dimension_records = self._extract_all_dimension_records_from_data_ids( 

2079 source_butler, data_ids, elements 

2080 ) 

2081 

2082 # Insert order is important. 

2083 for element in self.dimensions.sorted(dimension_records.keys()): 

2084 records = [r for r in dimension_records[element].values()] 

2085 # Assume that if the record is already present that we can 

2086 # use it without having to check that the record metadata 

2087 # is consistent. 

2088 self._registry.insertDimensionData(element, *records, skip_existing=True) 

2089 _LOG.debug("Dimension '%s' -- number of records transferred: %d", element.name, len(records)) 

2090 

2091 def _extract_all_dimension_records_from_data_ids( 

2092 self, 

2093 source_butler: LimitedButler | Butler, 

2094 data_ids: set[DataCoordinate], 

2095 allowed_elements: frozenset[DimensionElement], 

2096 ) -> dict[DimensionElement, dict[DataCoordinate, DimensionRecord]]: 

2097 primary_records = self._extract_dimension_records_from_data_ids( 

2098 source_butler, data_ids, allowed_elements 

2099 ) 

2100 

2101 additional_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]] = defaultdict(dict) 

2102 for original_element, record_mapping in primary_records.items(): 

2103 # Get dimensions that depend on this dimension. 

2104 populated_by = self.dimensions.get_elements_populated_by( 

2105 self.dimensions[original_element.name] # type: ignore 

2106 ) 

2107 if populated_by: 

2108 for element in populated_by: 

2109 if element not in allowed_elements: 2109 ↛ 2110line 2109 didn't jump to line 2110 because the condition on line 2109 was never true

2110 continue 

2111 if element.name == original_element.name: 

2112 continue 

2113 

2114 if element.name in primary_records: 

2115 # If this element has already been stored avoid 

2116 # re-finding records since that may lead to additional 

2117 # spurious records. e.g. visit is populated_by 

2118 # visit_detector_region but querying 

2119 # visit_detector_region by visit will return all the 

2120 # detectors for this visit -- the visit dataId does not 

2121 # constrain this. 

2122 # To constrain the query the original dataIds would 

2123 # have to be scanned. 

2124 continue 

2125 

2126 if record_mapping: 2126 ↛ 2108line 2126 didn't jump to line 2108 because the condition on line 2126 was always true

2127 if not isinstance(source_butler, Butler): 2127 ↛ 2128line 2127 didn't jump to line 2128 because the condition on line 2127 was never true

2128 raise RuntimeError( 

2129 f"Transferring populated_by records like {element.name}" 

2130 " requires a full Butler." 

2131 ) 

2132 

2133 with source_butler.query() as query: 

2134 records = query.join_data_coordinates(record_mapping.keys()).dimension_records( 

2135 element.name 

2136 ) 

2137 for record in records: 

2138 additional_records[record.definition].setdefault(record.dataId, record) 

2139 

2140 # The next step is to walk back through the additional records to 

2141 # pick up any missing content (such as visit_definition needing to 

2142 # know the exposure). Want to ensure we do not request records we 

2143 # already have. 

2144 missing_data_ids = set() 

2145 for record_mapping in additional_records.values(): 

2146 for data_id in record_mapping.keys(): 

2147 for dimension in data_id.dimensions.required: 

2148 element = source_butler.dimensions[dimension] 

2149 dimension_key = data_id.subset(dimension) 

2150 if dimension_key not in primary_records[element]: 

2151 missing_data_ids.add(dimension_key) 

2152 

2153 # Fill out the new records. Assume that these new records do not 

2154 # also need to carry over additional populated_by records. 

2155 secondary_records = self._extract_dimension_records_from_data_ids( 

2156 source_butler, missing_data_ids, allowed_elements 

2157 ) 

2158 

2159 # Merge the extra sets of records in with the original. 

2160 for name, record_mapping in itertools.chain(additional_records.items(), secondary_records.items()): 

2161 primary_records[name].update(record_mapping) 

2162 

2163 return primary_records 

2164 

2165 def _extract_dimension_records_from_data_ids( 

2166 self, 

2167 source_butler: LimitedButler | Butler, 

2168 data_ids: Iterable[DataCoordinate], 

2169 allowed_elements: frozenset[DimensionElement], 

2170 ) -> dict[DimensionElement, dict[DataCoordinate, DimensionRecord]]: 

2171 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]] = defaultdict(dict) 

2172 

2173 data_ids = set(data_ids) 

2174 if not all(data_id.hasRecords() for data_id in data_ids): 

2175 if isinstance(source_butler, Butler): 2175 ↛ 2178line 2175 didn't jump to line 2178 because the condition on line 2175 was always true

2176 data_ids = source_butler._expand_data_ids(data_ids) 

2177 else: 

2178 raise TypeError("Input butler needs to be a full butler to expand DataId.") 

2179 

2180 for data_id in data_ids: 

2181 # If this butler doesn't know about a dimension in the source 

2182 # butler things will break later. 

2183 for element_name in data_id.dimensions.elements: 

2184 record = data_id.records[element_name] 

2185 if record is not None and record.definition in allowed_elements: 

2186 dimension_records[record.definition].setdefault(record.dataId, record) 

2187 

2188 return dimension_records 

2189 

2190 def _cast_universe_for_import_refs( 

2191 self, source_refs: Iterable[DatasetRef] 

2192 ) -> Mapping[DatasetType, list[DatasetRef]]: 

2193 """Try to cast imported refs to the target universe if possible. 

2194 

2195 Parameters 

2196 ---------- 

2197 source_refs 

2198 The refs to be imported. 

2199 

2200 Returns 

2201 ------- 

2202 refs 

2203 The refs to be imported, grouped by dataset type, with the dataset 

2204 types cast to the target universe. 

2205 

2206 Raises 

2207 ------ 

2208 InconsistentUniverseError 

2209 Raised if any reference cannot be converted to target universe. 

2210 

2211 Notes 

2212 ----- 

2213 Potentially this method can perform a non-trivial migrations of the 

2214 datasets by modifying dimensions and dataIds. Presently though it can 

2215 only perform a trivial validation of the dataset types compatibility. 

2216 Returned mapping will contain dataset types in the new universe, but 

2217 returned references will still have the original dataset types as 

2218 there is presently no easy way to replace dataset type in a reference. 

2219 """ 

2220 # In theory input refs could come from multiple universes, but in 

2221 # practice this will not happen, so just group everything by dataset 

2222 # type. 

2223 refs_by_source_type: defaultdict[DatasetType, list[DatasetRef]] = defaultdict(list) 

2224 for ref in source_refs: 

2225 refs_by_source_type[ref.datasetType].append(ref) 

2226 

2227 refs_by_type: defaultdict[DatasetType, list[DatasetRef]] = defaultdict(list) 

2228 for source_type, refs in refs_by_source_type.items(): 

2229 refs_by_type[source_type.conform_to(self.dimensions)] = refs 

2230 

2231 return refs_by_type 

2232 

2233 def _prepare_for_import_refs( 

2234 self, 

2235 source_butler: LimitedButler, 

2236 source_refs: Iterable[DatasetRef], 

2237 *, 

2238 register_dataset_types: bool = False, 

2239 transfer_dimensions: bool = False, 

2240 dry_run: bool = False, 

2241 ) -> _ImportDatasetsInfo: 

2242 # Docstring inherited. 

2243 if not self.isWriteable() and not dry_run: 2243 ↛ 2244line 2243 didn't jump to line 2244 because the condition on line 2243 was never true

2244 raise TypeError("Butler is read-only.") 

2245 

2246 # Will iterate through the refs multiple times so need to convert 

2247 # to a list if this isn't a collection. 

2248 if not isinstance(source_refs, collections.abc.Collection): 2248 ↛ 2249line 2248 didn't jump to line 2249 because the condition on line 2248 was never true

2249 source_refs = list(source_refs) 

2250 

2251 original_count = len(source_refs) 

2252 log_level = _LOG.INFO if original_count > 1 else _LOG.VERBOSE 

2253 _LOG.log( 

2254 log_level, 

2255 "Importing %d dataset%s into %s", 

2256 original_count, 

2257 "s" if original_count != 1 else "", 

2258 str(self), 

2259 ) 

2260 

2261 refs_by_type = self._cast_universe_for_import_refs(source_refs) 

2262 

2263 # Importing requires that we group the refs by dimension group and run 

2264 # before doing the import. 

2265 grouped_refs: defaultdict[_RefGroup, list[DatasetRef]] = defaultdict(list) 

2266 for ref in source_refs: 

2267 grouped_refs[_RefGroup(ref.datasetType.dimensions, ref.run)].append(ref) 

2268 

2269 # Check to see if the dataset type in the source butler has 

2270 # the same definition in the target butler and register missing 

2271 # ones if requested. Registration must happen outside a transaction. 

2272 newly_registered_dataset_types = set() 

2273 for datasetType in refs_by_type: 

2274 if register_dataset_types: 

2275 # Let this raise immediately if inconsistent. Continuing 

2276 # on to find additional inconsistent dataset types 

2277 # might result in additional unwanted dataset types being 

2278 # registered. 

2279 try: 

2280 if not dry_run and self._registry.registerDatasetType(datasetType): 

2281 newly_registered_dataset_types.add(datasetType) 

2282 except ConflictingDefinitionError as e: 

2283 # Be safe and require that conversions be bidirectional 

2284 # when there are storage class mismatches. This is because 

2285 # get() will have to support conversion from source to 

2286 # target python type (the source formatter will be 

2287 # returning source python type) but there also is an 

2288 # expectation that people will want to be able to get() in 

2289 # the target using the source python type, which will not 

2290 # require conversion for transferred datasets but might 

2291 # for target-native types. Additionally, butler.get does 

2292 # not know that the formatter will return the wrong 

2293 # python type and so will always check that the conversion 

2294 # works even though it won't need it. 

2295 target_dataset_type = self.get_dataset_type(datasetType.name) 

2296 target_compatible_with_source = target_dataset_type.is_compatible_with(datasetType) 

2297 source_compatible_with_target = datasetType.is_compatible_with(target_dataset_type) 

2298 if not (target_compatible_with_source and source_compatible_with_target): 2298 ↛ 2299line 2298 didn't jump to line 2299 because the condition on line 2298 was never true

2299 if target_compatible_with_source: 

2300 e.add_note( 

2301 "Target dataset type storage class is compatible with source " 

2302 "but the reverse is not true." 

2303 ) 

2304 elif source_compatible_with_target: 

2305 e.add_note( 

2306 "Source dataset type storage class is compatible with target " 

2307 "but the reverse is not true." 

2308 ) 

2309 else: 

2310 e.add_note("If storage classes differ, please register converters.") 

2311 raise 

2312 else: 

2313 # If the dataset type is missing, let it fail immediately. 

2314 target_dataset_type = self.get_dataset_type(datasetType.name) 

2315 if target_dataset_type != datasetType: 

2316 target_compatible_with_source = target_dataset_type.is_compatible_with(datasetType) 

2317 source_compatible_with_target = datasetType.is_compatible_with(target_dataset_type) 

2318 # Both conversion directions are currently required. 

2319 if not (target_compatible_with_source and source_compatible_with_target): 

2320 msg = "" 

2321 if target_compatible_with_source: 2321 ↛ 2322line 2321 didn't jump to line 2322 because the condition on line 2321 was never true

2322 msg = ( 

2323 "Target storage class is compatible with the source storage class " 

2324 "but the reverse is not true." 

2325 ) 

2326 elif source_compatible_with_target: 2326 ↛ 2327line 2326 didn't jump to line 2327 because the condition on line 2326 was never true

2327 msg = ( 

2328 "Source storage class is compatible with the target storage class" 

2329 " but the reverse is not true." 

2330 ) 

2331 else: 

2332 msg = "If storage classes differ register converters." 

2333 if msg: 2333 ↛ 2335line 2333 didn't jump to line 2335 because the condition on line 2333 was always true

2334 msg = f"({msg})" 

2335 raise ConflictingDefinitionError( 

2336 "Source butler dataset type differs from definition" 

2337 f" in target butler: {datasetType} !=" 

2338 f" {target_dataset_type} {msg}" 

2339 ) 

2340 if newly_registered_dataset_types: 

2341 # We may have registered some even if there were inconsistencies 

2342 # but should let people know (or else remove them again). 

2343 _LOG.verbose( 

2344 "Registered the following dataset types in the target Butler: %s", 

2345 ", ".join(d.name for d in newly_registered_dataset_types), 

2346 ) 

2347 else: 

2348 _LOG.verbose("All required dataset types are known to the target Butler") 

2349 

2350 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]] = defaultdict(dict) 

2351 if transfer_dimensions: 

2352 # Collect all the dimension records for these refs. 

2353 # All dimensions are to be copied but the list of valid dimensions 

2354 # come from this butler's universe. 

2355 elements = frozenset(element for element in self.dimensions.elements if element.has_own_table) 

2356 dataIds = {ref.dataId for ref in source_refs} 

2357 dimension_records = self._extract_all_dimension_records_from_data_ids( 

2358 source_butler, dataIds, elements 

2359 ) 

2360 return _ImportDatasetsInfo(grouped_refs, dimension_records) 

2361 

2362 def _import_dimension_records( 

2363 self, 

2364 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]], 

2365 *, 

2366 dry_run: bool, 

2367 ) -> None: 

2368 """Import dimension records collected during import pre-process.""" 

2369 if dimension_records and not dry_run: 

2370 _LOG.verbose("Ensuring that dimension records exist for transferred datasets.") 

2371 # Order matters. 

2372 for element in self.dimensions.sorted(dimension_records.keys()): 

2373 records = list(dimension_records[element].values()) 

2374 # Assume that if the record is already present that we can 

2375 # use it without having to check that the record metadata 

2376 # is consistent. 

2377 self._registry.insertDimensionData(element, *records, skip_existing=True) 

2378 

2379 def _import_grouped_refs( 

2380 self, 

2381 grouped_refs: defaultdict[_RefGroup, list[DatasetRef]], 

2382 source_butler: LimitedButler | None, 

2383 progress: Progress, 

2384 *, 

2385 dry_run: bool = False, 

2386 expand_refs: bool = False, 

2387 ) -> list[DatasetRef]: 

2388 handled_collections: set[str] = set() 

2389 n_to_import = 0 

2390 all_imported_refs: list[DatasetRef] = [] 

2391 # Sort by run collection name to ensure Postgres takes locks in the 

2392 # same order between different processes, to mitigate an issue 

2393 # where Postgres can deadlock due to the unique index on collection 

2394 # name. (See DM-47543). 

2395 groups = sorted(grouped_refs.items(), key=lambda item: item[0].run) 

2396 for (dimension_group, run), refs_to_import in progress.iter_item_chunks( 

2397 groups, desc="Importing to registry by run and dataset type" 

2398 ): 

2399 if run not in handled_collections: 

2400 # May need to create output collection. If source butler 

2401 # has a registry, ask for documentation string. 

2402 run_doc = None 

2403 if source_butler is not None and (registry := getattr(source_butler, "registry", None)): 

2404 run_doc = registry.getCollectionDocumentation(run) 

2405 if not dry_run: 

2406 registered = self.collections.register(run, doc=run_doc) 

2407 else: 

2408 registered = True 

2409 handled_collections.add(run) 

2410 if registered: 

2411 _LOG.verbose("Creating output run %s", run) 

2412 

2413 n_refs = len(refs_to_import) 

2414 n_to_import += n_refs 

2415 _LOG.verbose( 

2416 "Importing %d ref%s with dimensions %s into run %s", 

2417 n_refs, 

2418 "" if n_refs == 1 else "s", 

2419 dimension_group.names, 

2420 run, 

2421 ) 

2422 

2423 # Assume we are using UUIDs and the source refs will match 

2424 # those imported. 

2425 if not dry_run: 

2426 imported_refs = self._registry._importDatasets(refs_to_import, expand=expand_refs) 

2427 else: 

2428 imported_refs = refs_to_import 

2429 

2430 all_imported_refs.extend(imported_refs) 

2431 

2432 assert n_to_import == len(all_imported_refs) 

2433 _LOG.verbose("Imported %d datasets into destination butler", n_to_import) 

2434 return all_imported_refs 

2435 

2436 def transfer_from( 

2437 self, 

2438 source_butler: LimitedButler, 

2439 source_refs: Iterable[DatasetRef], 

2440 transfer: str = "auto", 

2441 skip_missing: bool = True, 

2442 register_dataset_types: bool = False, 

2443 transfer_dimensions: bool = False, 

2444 dry_run: bool = False, 

2445 ) -> collections.abc.Collection[DatasetRef]: 

2446 # Docstring inherited. 

2447 source_refs = list(source_refs) 

2448 if not self.isWriteable() and not dry_run: 2448 ↛ 2449line 2448 didn't jump to line 2449 because the condition on line 2448 was never true

2449 raise TypeError("Butler is read-only.") 

2450 

2451 progress = Progress("lsst.daf.butler.Butler.transfer_from", level=VERBOSE) 

2452 

2453 artifact_existence: dict[ResourcePath, bool] = {} 

2454 file_transfer_source = source_butler._file_transfer_source 

2455 transfer_records = retrieve_file_transfer_records( 

2456 file_transfer_source, source_refs, artifact_existence 

2457 ) 

2458 # In some situations the datastore artifact may be missing and we do 

2459 # not want that registry entry to be imported. For example, this can 

2460 # happen if a file was removed but the dataset was left in the registry 

2461 # for provenance, or if a pipeline task didn't create all of the 

2462 # possible files in a QuantumBackedButler. 

2463 if skip_missing: 

2464 original_ids = {ref.id for ref in source_refs} 

2465 missing_ids = original_ids - transfer_records.keys() 

2466 if missing_ids: 

2467 original_count = len(source_refs) 

2468 source_refs = [ref for ref in source_refs if ref.id not in missing_ids] 

2469 filtered_count = len(source_refs) 

2470 n_missing = original_count - filtered_count 

2471 _LOG.verbose( 

2472 "%d dataset%s removed because the artifact does not exist. Now have %d.", 

2473 n_missing, 

2474 "" if n_missing == 1 else "s", 

2475 filtered_count, 

2476 ) 

2477 

2478 import_info = self._prepare_for_import_refs( 

2479 source_butler, 

2480 source_refs, 

2481 register_dataset_types=register_dataset_types, 

2482 dry_run=dry_run, 

2483 transfer_dimensions=transfer_dimensions, 

2484 ) 

2485 

2486 # Do all the importing in a single transaction. 

2487 with self.transaction(): 

2488 self._import_dimension_records(import_info.dimension_records, dry_run=dry_run) 

2489 imported_refs = self._import_grouped_refs( 

2490 import_info.grouped_refs, source_butler, progress, dry_run=dry_run 

2491 ) 

2492 

2493 # Ask the datastore to transfer. The datastore has to check that 

2494 # the source datastore is compatible with the target datastore. 

2495 _LOG.verbose("Transferring %d datasets from %s", len(transfer_records), file_transfer_source.name) 

2496 accepted, rejected = self._datastore.transfer_from( 

2497 transfer_records, 

2498 imported_refs, 

2499 transfer=transfer, 

2500 artifact_existence=artifact_existence, 

2501 dry_run=dry_run, 

2502 ) 

2503 if rejected: 2503 ↛ 2505line 2503 didn't jump to line 2505 because the condition on line 2503 was never true

2504 # For now, accept the registry entries but not the files. 

2505 _LOG.warning( 

2506 "%d datasets were rejected and %d accepted for transfer.", 

2507 len(rejected), 

2508 len(accepted), 

2509 ) 

2510 

2511 return imported_refs 

2512 

2513 def validateConfiguration( 

2514 self, 

2515 logFailures: bool = False, 

2516 datasetTypeNames: Iterable[str] | None = None, 

2517 ignore: Iterable[str] | None = None, 

2518 ) -> None: 

2519 # Docstring inherited. 

2520 if datasetTypeNames: 

2521 datasetTypes = [self.get_dataset_type(name) for name in datasetTypeNames] 

2522 else: 

2523 datasetTypes = list(self._registry.queryDatasetTypes()) 

2524 

2525 # filter out anything from the ignore list 

2526 if ignore: 

2527 ignore = set(ignore) 

2528 datasetTypes = [ 

2529 e for e in datasetTypes if e.name not in ignore and e.nameAndComponent()[0] not in ignore 

2530 ] 

2531 else: 

2532 ignore = set() 

2533 

2534 # For each datasetType that has an instrument dimension, create 

2535 # a DatasetRef for each defined instrument 

2536 datasetRefs = [] 

2537 

2538 # Find all the registered instruments (if "instrument" is in the 

2539 # universe). 

2540 instruments: set[str] = set() 

2541 if "instrument" in self.dimensions: 2541 ↛ 2561line 2541 didn't jump to line 2561 because the condition on line 2541 was always true

2542 instruments = {rec.name for rec in self.query_dimension_records("instrument", explain=False)} 

2543 

2544 for datasetType in datasetTypes: 

2545 if "instrument" in datasetType.dimensions: 2545 ↛ 2544line 2545 didn't jump to line 2544 because the condition on line 2545 was always true

2546 # In order to create a conforming dataset ref, create 

2547 # fake DataCoordinate values for the non-instrument 

2548 # dimensions. The type of the value does not matter here. 

2549 dataId = {dim: 1 for dim in datasetType.dimensions.names if dim != "instrument"} 

2550 

2551 for instrument in instruments: 

2552 datasetRef = DatasetRef( 

2553 datasetType, 

2554 DataCoordinate.standardize( 

2555 dataId, instrument=instrument, dimensions=datasetType.dimensions 

2556 ), 

2557 run="validate", 

2558 ) 

2559 datasetRefs.append(datasetRef) 

2560 

2561 entities: list[DatasetType | DatasetRef] = [] 

2562 entities.extend(datasetTypes) 

2563 entities.extend(datasetRefs) 

2564 

2565 datastoreErrorStr = None 

2566 try: 

2567 self._datastore.validateConfiguration(entities, logFailures=logFailures) 

2568 except ValidationError as e: 

2569 datastoreErrorStr = str(e) 

2570 

2571 # Also check that the LookupKeys used by the datastores match 

2572 # registry and storage class definitions 

2573 keys = self._datastore.getLookupKeys() 

2574 

2575 failedNames = set() 

2576 failedDataId = set() 

2577 for key in keys: 

2578 if key.name is not None: 

2579 if key.name in ignore: 

2580 continue 

2581 

2582 # skip if specific datasetType names were requested and this 

2583 # name does not match 

2584 if datasetTypeNames and key.name not in datasetTypeNames: 

2585 continue 

2586 

2587 # See if it is a StorageClass or a DatasetType 

2588 if key.name in self.storageClasses: 

2589 pass 

2590 else: 

2591 try: 

2592 self.get_dataset_type(key.name) 

2593 except KeyError: 

2594 if logFailures: 2594 ↛ 2595line 2594 didn't jump to line 2595 because the condition on line 2594 was never true

2595 _LOG.critical( 

2596 "Key '%s' does not correspond to a DatasetType or StorageClass", key 

2597 ) 

2598 failedNames.add(key) 

2599 else: 

2600 # Dimensions are checked for consistency when the Butler 

2601 # is created and rendezvoused with a universe. 

2602 pass 

2603 

2604 # Check that the instrument is a valid instrument 

2605 # Currently only support instrument so check for that 

2606 if key.dataId: 

2607 dataIdKeys = set(key.dataId) 

2608 if {"instrument"} != dataIdKeys: 2608 ↛ 2609line 2608 didn't jump to line 2609 because the condition on line 2608 was never true

2609 if logFailures: 

2610 _LOG.critical("Key '%s' has unsupported DataId override", key) 

2611 failedDataId.add(key) 

2612 elif key.dataId["instrument"] not in instruments: 2612 ↛ 2613line 2612 didn't jump to line 2613 because the condition on line 2612 was never true

2613 if logFailures: 

2614 _LOG.critical("Key '%s' has unknown instrument", key) 

2615 failedDataId.add(key) 

2616 

2617 messages = [] 

2618 

2619 if datastoreErrorStr: 2619 ↛ 2620line 2619 didn't jump to line 2620 because the condition on line 2619 was never true

2620 messages.append(datastoreErrorStr) 

2621 

2622 for failed, msg in ( 

2623 (failedNames, "Keys without corresponding DatasetType or StorageClass entry: "), 

2624 (failedDataId, "Keys with bad DataId entries: "), 

2625 ): 

2626 if failed: 

2627 msg += ", ".join(str(k) for k in failed) 

2628 messages.append(msg) 

2629 

2630 if messages: 

2631 raise ValidationError(";\n".join(messages)) 

2632 

2633 @property 

2634 @deprecated( 

2635 "Please use 'collections' instead. collection_chains will be removed after v28.", 

2636 version="v28", 

2637 category=FutureWarning, 

2638 ) 

2639 def collection_chains(self) -> DirectButlerCollections: 

2640 """Object with methods for modifying collection chains.""" 

2641 return DirectButlerCollections(self._registry) 

2642 

2643 @property 

2644 def collections(self) -> DirectButlerCollections: 

2645 """Object with methods for modifying and inspecting collections.""" 

2646 return DirectButlerCollections(self._registry) 

2647 

2648 @property 

2649 def run(self) -> str | None: 

2650 """Name of the run this butler writes outputs to by default (`str` or 

2651 `None`). 

2652 

2653 This is an alias for ``self.registry.defaults.run``. It cannot be set 

2654 directly in isolation, but all defaults may be changed together by 

2655 assigning a new `RegistryDefaults` instance to 

2656 ``self.registry.defaults``. 

2657 """ 

2658 return self._registry.defaults.run 

2659 

2660 @property 

2661 def registry(self) -> Registry: 

2662 """The object that manages dataset metadata and relationships 

2663 (`Registry`). 

2664 

2665 Many operations that don't involve reading or writing butler datasets 

2666 are accessible only via `Registry` methods. Eventually these methods 

2667 will be replaced by equivalent `Butler` methods. 

2668 """ 

2669 return RegistryShim(self) 

2670 

2671 @property 

2672 def dimensions(self) -> DimensionUniverse: 

2673 # Docstring inherited. 

2674 return self._registry.dimensions 

2675 

2676 def query(self) -> contextlib.AbstractContextManager[Query]: 

2677 # Docstring inherited. 

2678 return self._registry._query() 

2679 

2680 def _query_driver( 

2681 self, 

2682 default_collections: Iterable[str], 

2683 default_data_id: DataCoordinate, 

2684 ) -> contextlib.AbstractContextManager[DirectQueryDriver]: 

2685 """Set up a QueryDriver instance for use with this Butler. Although 

2686 this is marked as a private method, it is also used by Butler server. 

2687 """ 

2688 return self._registry._query_driver(default_collections, default_data_id) 

2689 

2690 @contextlib.contextmanager 

2691 def _query_all_datasets_by_page( 

2692 self, args: QueryAllDatasetsParameters 

2693 ) -> Iterator[Iterator[list[DatasetRef]]]: 

2694 with self.query() as query: 

2695 pages = query_all_datasets(self, query, args) 

2696 yield iter(page.data for page in pages) 

2697 

2698 def _preload_cache(self, *, load_dimension_record_cache: bool = True) -> None: 

2699 """Immediately load caches that are used for common operations.""" 

2700 self._registry.preload_cache(load_dimension_record_cache=load_dimension_record_cache) 

2701 

2702 def _expand_data_ids(self, data_ids: Iterable[DataCoordinate]) -> list[DataCoordinate]: 

2703 return self._registry.expand_data_ids(data_ids) 

2704 

2705 _config: ButlerConfig 

2706 """Configuration for this Butler instance.""" 

2707 

2708 _registry: SqlRegistry 

2709 """The object that manages dataset metadata and relationships 

2710 (`SqlRegistry`). 

2711 

2712 Most operations that don't involve reading or writing butler datasets are 

2713 accessible only via `SqlRegistry` methods. 

2714 """ 

2715 

2716 storageClasses: StorageClassFactory 

2717 """An object that maps known storage class names to objects that fully 

2718 describe them (`StorageClassFactory`). 

2719 """ 

2720 

2721 _closed: bool 

2722 """`True` if close() has already been called on this instance; `False` 

2723 otherwise. 

2724 """ 

2725 

2726 

2727class _RefGroup(NamedTuple): 

2728 """Key identifying a batch of DatasetRefs to be inserted in 

2729 `Butler.transfer_from`. 

2730 """ 

2731 

2732 dimensions: DimensionGroup 

2733 run: str 

2734 

2735 

2736class _ImportDatasetsInfo(NamedTuple): 

2737 """Information extracted from datasets to be imported.""" 

2738 

2739 grouped_refs: defaultdict[_RefGroup, list[DatasetRef]] 

2740 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]] 

2741 

2742 

2743def _to_uuid(id: DatasetId | str) -> uuid.UUID: 

2744 if isinstance(id, uuid.UUID): 

2745 return id 

2746 else: 

2747 return uuid.UUID(id) 

2748 

2749 

2750class _ButlerClosed: 

2751 def __getattr__(self, name: str) -> Any: 

2752 raise RuntimeError("Attempted to use a Butler instance which has been closed.") 

2753 

2754 

2755_BUTLER_CLOSED_INSTANCE: Any = _ButlerClosed() 

2756 

2757 

2758def _retrieve_dataset_type(registry: SqlRegistry, name: str) -> DatasetType | None: 

2759 """Return DatasetType defined in registry given dataset type name.""" 

2760 try: 

2761 return registry.getDatasetType(name) 

2762 except MissingDatasetTypeError: 

2763 return None