Coverage for python/lsst/daf/butler/registry/sql_registry.py: 89%

409 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-09-09 02:00 -0700

1# This file is part of daf_butler. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

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

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

7# for details of code ownership. 

8# 

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

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

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

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

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

14# 

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

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

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

18# (at your option) any later version. 

19# 

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

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

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

23# GNU General Public License for more details. 

24# 

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

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

27 

28from __future__ import annotations 

29 

30from .. import ddl 

31 

32__all__ = ("SqlRegistry",) 

33 

34import contextlib 

35import logging 

36import warnings 

37from collections.abc import Iterable, Iterator, Mapping, Sequence 

38from typing import TYPE_CHECKING, Any 

39 

40import sqlalchemy 

41 

42from lsst.resources import ResourcePathExpression 

43from lsst.utils.iteration import ensure_iterable 

44 

45from .._collection_type import CollectionType 

46from .._config import Config 

47from .._dataset_ref import DatasetId, DatasetIdGenEnum, DatasetRef 

48from .._dataset_type import DatasetType, validate_dataset_type_name 

49from .._exceptions import DataIdValueError, DimensionNameError, InconsistentDataIdError 

50from .._storage_class import StorageClassFactory 

51from .._timespan import Timespan 

52from ..dimensions import ( 

53 DataCoordinate, 

54 DataId, 

55 DimensionConfig, 

56 DimensionElement, 

57 DimensionGroup, 

58 DimensionRecord, 

59 DimensionUniverse, 

60) 

61from ..dimensions.record_cache import DimensionRecordCache 

62from ..direct_query_driver import DirectQueryDriver 

63from ..progress import Progress 

64from ..queries import Query 

65from ..registry import ( 

66 CollectionExpressionError, 

67 CollectionSummary, 

68 CollectionTypeError, 

69 ConflictingDefinitionError, 

70 NoDefaultCollectionError, 

71 OrphanedRecordError, 

72 RegistryConfig, 

73 RegistryDefaults, 

74) 

75from ..registry.interfaces import ChainedCollectionRecord, ReadOnlyDatabaseError, RunRecord 

76from ..registry.managers import RegistryManagerInstances, RegistryManagerTypes 

77from ..registry.wildcards import CollectionWildcard, DatasetTypeWildcard 

78from ..utils import transactional 

79from .expand_data_ids import expand_data_ids 

80 

81if TYPE_CHECKING: 

82 from .._butler_config import ButlerConfig 

83 from ..datastore._datastore import DatastoreOpaqueTable 

84 from ..datastore.stored_file_info import StoredDatastoreItemInfo 

85 from ..registry.interfaces import ( 

86 CollectionRecord, 

87 Database, 

88 DatastoreRegistryBridgeManager, 

89 ObsCoreTableManager, 

90 ) 

91 

92 

93_LOG = logging.getLogger(__name__) 

94 

95 

96class SqlRegistry: 

97 """Butler Registry implementation that uses SQL database as backend. 

98 

99 Parameters 

100 ---------- 

101 database : `Database` 

102 Database instance to store Registry. 

103 defaults : `RegistryDefaults` 

104 Default collection search path and/or output `~CollectionType.RUN` 

105 collection. 

106 managers : `RegistryManagerInstances` 

107 All the managers required for this registry. 

108 """ 

109 

110 defaultConfigFile: str | None = None 

111 """Path to configuration defaults. Accessed within the ``configs`` resource 

112 or relative to a search path. Can be None if no defaults specified. 

113 """ 

114 

115 @classmethod 

116 def forceRegistryConfig( 

117 cls, config: ButlerConfig | RegistryConfig | Config | str | None 

118 ) -> RegistryConfig: 

119 """Force the supplied config to a `RegistryConfig`. 

120 

121 Parameters 

122 ---------- 

123 config : `RegistryConfig`, `Config` or `str` or `None` 

124 Registry configuration, if missing then default configuration will 

125 be loaded from registry.yaml. 

126 

127 Returns 

128 ------- 

129 registry_config : `RegistryConfig` 

130 A registry config. 

131 """ 

132 if not isinstance(config, RegistryConfig): 

133 if isinstance(config, str | Config) or config is None: 133 ↛ 136line 133 didn't jump to line 136 because the condition on line 133 was always true

134 config = RegistryConfig(config) 

135 else: 

136 raise ValueError(f"Incompatible Registry configuration: {config}") 

137 return config 

138 

139 @classmethod 

140 def createFromConfig( 

141 cls, 

142 config: RegistryConfig | str | None = None, 

143 dimensionConfig: DimensionConfig | str | None = None, 

144 butlerRoot: ResourcePathExpression | None = None, 

145 ) -> SqlRegistry: 

146 """Create registry database and return `SqlRegistry` instance. 

147 

148 This method initializes database contents, database must be empty 

149 prior to calling this method. 

150 

151 Parameters 

152 ---------- 

153 config : `RegistryConfig` or `str`, optional 

154 Registry configuration, if missing then default configuration will 

155 be loaded from registry.yaml. 

156 dimensionConfig : `DimensionConfig` or `str`, optional 

157 Dimensions configuration, if missing then default configuration 

158 will be loaded from dimensions.yaml. 

159 butlerRoot : convertible to `lsst.resources.ResourcePath`, optional 

160 Path to the repository root this `SqlRegistry` will manage. 

161 

162 Returns 

163 ------- 

164 registry : `SqlRegistry` 

165 A new `SqlRegistry` instance. 

166 """ 

167 config = cls.forceRegistryConfig(config) 

168 config.replaceRoot(butlerRoot) 

169 

170 if isinstance(dimensionConfig, str): 170 ↛ 171line 170 didn't jump to line 171 because the condition on line 170 was never true

171 dimensionConfig = DimensionConfig(dimensionConfig) 

172 elif dimensionConfig is None: 

173 dimensionConfig = DimensionConfig() 

174 elif not isinstance(dimensionConfig, DimensionConfig): 174 ↛ 175line 174 didn't jump to line 175 because the condition on line 174 was never true

175 raise TypeError(f"Incompatible Dimension configuration type: {type(dimensionConfig)}") 

176 

177 managerTypes = RegistryManagerTypes.fromConfig(config) 

178 DatabaseClass = config.getDatabaseClass() 

179 database = DatabaseClass.fromUri( 

180 config.connectionString, 

181 origin=config.get("origin", 0), 

182 namespace=config.get("namespace"), 

183 allow_temporary_tables=config.areTemporaryTablesAllowed, 

184 ) 

185 

186 try: 

187 managers = managerTypes.makeRepo(database, dimensionConfig) 

188 return cls(database, RegistryDefaults(), managers) 

189 except Exception: 

190 database.dispose() 

191 raise 

192 

193 @classmethod 

194 def fromConfig( 

195 cls, 

196 config: ButlerConfig | RegistryConfig | Config | str, 

197 butlerRoot: ResourcePathExpression | None = None, 

198 writeable: bool = True, 

199 defaults: RegistryDefaults | None = None, 

200 ) -> SqlRegistry: 

201 """Create `Registry` subclass instance from `config`. 

202 

203 Registry database must be initialized prior to calling this method. 

204 

205 Parameters 

206 ---------- 

207 config : `ButlerConfig`, `RegistryConfig`, `Config` or `str` 

208 Registry configuration. 

209 butlerRoot : `lsst.resources.ResourcePathExpression`, optional 

210 Path to the repository root this `Registry` will manage. 

211 writeable : `bool`, optional 

212 If `True` (default) create a read-write connection to the database. 

213 defaults : `RegistryDefaults`, optional 

214 Default collection search path and/or output `~CollectionType.RUN` 

215 collection. 

216 

217 Returns 

218 ------- 

219 registry : `SqlRegistry` 

220 A new `SqlRegistry` subclass instance. 

221 """ 

222 config = cls.forceRegistryConfig(config) 

223 config.replaceRoot(butlerRoot) 

224 if defaults is None: 

225 defaults = RegistryDefaults() 

226 DatabaseClass = config.getDatabaseClass() 

227 database = DatabaseClass.fromUri( 

228 config.connectionString, 

229 origin=config.get("origin", 0), 

230 namespace=config.get("namespace"), 

231 writeable=writeable, 

232 allow_temporary_tables=config.areTemporaryTablesAllowed, 

233 ) 

234 try: 

235 managerTypes = RegistryManagerTypes.fromConfig(config) 

236 with database.session(): 

237 managers = managerTypes.loadRepo(database) 

238 

239 return cls(database, defaults, managers) 

240 except Exception: 

241 database.dispose() 

242 raise 

243 

244 def __init__( 

245 self, 

246 database: Database, 

247 defaults: RegistryDefaults, 

248 managers: RegistryManagerInstances, 

249 ): 

250 self._db = database 

251 self._managers = managers 

252 if managers.obscore is not None: 

253 managers.obscore.set_query_function(self._query) 

254 self.storageClasses = StorageClassFactory() 

255 # This is public to SqlRegistry's internal-to-daf_butler callers, but 

256 # it is intentionally not part of RegistryShim. 

257 self.dimension_record_cache = DimensionRecordCache( 

258 self._managers.dimensions.universe, 

259 fetch=self._managers.dimensions.fetch_cache_dict, 

260 ) 

261 # Intentionally invoke property setter to initialize defaults. This 

262 # can only be done after most of the rest of Registry has already been 

263 # initialized, and must be done before the property getter is used. 

264 self.defaults = defaults 

265 # TODO: This is currently initialized by `make_datastore_tables`, 

266 # eventually we'll need to do it during construction. 

267 # The mapping is indexed by the opaque table name. 

268 self._datastore_record_classes: Mapping[str, type[StoredDatastoreItemInfo]] = {} 

269 self._is_clone = False 

270 

271 def close(self) -> None: 

272 # Connection pool is shared between cloned instances, so only the root 

273 # instance should close it. 

274 # Note: The underlying SQLAlchemy call will create a fresh connection 

275 # pool, so nothing breaks if the root instance is accidentally closed 

276 # before the clones are finished -- we just have a small performance 

277 # hit from re-creating the connections. 

278 if not self._is_clone: 

279 self._db.dispose() 

280 

281 def __str__(self) -> str: 

282 return str(self._db) 

283 

284 def __repr__(self) -> str: 

285 return f"SqlRegistry({self._db!r}, {self.dimensions!r})" 

286 

287 def isWriteable(self) -> bool: 

288 """Return `True` if this registry allows write operations, and `False` 

289 otherwise. 

290 """ 

291 return self._db.isWriteable() 

292 

293 def copy(self, defaults: RegistryDefaults | None = None) -> SqlRegistry: 

294 """Create a new `SqlRegistry` backed by the same data repository 

295 as this one and sharing a database connection pool with it, but with 

296 independent defaults and database sessions. 

297 

298 Parameters 

299 ---------- 

300 defaults : `~lsst.daf.butler.registry.RegistryDefaults`, optional 

301 Default collections and data ID values for the new registry. If 

302 not provided, ``self.defaults`` will be used (but future changes 

303 to either registry's defaults will not affect the other). 

304 

305 Returns 

306 ------- 

307 copy : `SqlRegistry` 

308 A new `SqlRegistry` instance with its own defaults. 

309 """ 

310 if defaults is None: 

311 # No need to copy, because `RegistryDefaults` is immutable; we 

312 # effectively copy on write. 

313 defaults = self.defaults 

314 db = self._db.clone() 

315 result = SqlRegistry(db, defaults, self._managers.clone(db)) 

316 result._datastore_record_classes = dict(self._datastore_record_classes) 

317 result.dimension_record_cache.load_from(self.dimension_record_cache) 

318 result._is_clone = True 

319 return result 

320 

321 @property 

322 def dimensions(self) -> DimensionUniverse: 

323 """Definitions of all dimensions recognized by this `Registry` 

324 (`DimensionUniverse`). 

325 """ 

326 return self._managers.dimensions.universe 

327 

328 @property 

329 def defaults(self) -> RegistryDefaults: 

330 """Default collection search path and/or output `~CollectionType.RUN` 

331 collection (`~lsst.daf.butler.registry.RegistryDefaults`). 

332 

333 This is an immutable struct whose components may not be set 

334 individually, but the entire struct can be set by assigning to this 

335 property. 

336 """ 

337 return self._defaults 

338 

339 @defaults.setter 

340 def defaults(self, value: RegistryDefaults) -> None: 

341 if value.run is not None: 

342 self.registerRun(value.run) 

343 value.finish(self) 

344 self._defaults = value 

345 

346 def refresh(self) -> None: 

347 """Refresh all in-memory state by querying the database. 

348 

349 This may be necessary to enable querying for entities added by other 

350 registry instances after this one was constructed. 

351 """ 

352 self.dimension_record_cache.reset() 

353 with self._db.transaction(): 

354 self._managers.refresh() 

355 

356 def refresh_collection_summaries(self) -> None: 

357 """Refresh content of the collection summary tables in the database. 

358 

359 This only cleans dataset type summaries, we may want to add cleanup of 

360 governor summaries later. 

361 """ 

362 for dataset_type in self.queryDatasetTypes(): 

363 self._managers.datasets.refresh_collection_summaries(dataset_type) 

364 

365 def caching_context(self) -> contextlib.AbstractContextManager[None]: 

366 """Return context manager that enables caching. 

367 

368 Returns 

369 ------- 

370 manager 

371 A context manager that enables client-side caching. Entering 

372 the context returns `None`. 

373 """ 

374 return self._managers.caching_context_manager() 

375 

376 @contextlib.contextmanager 

377 def transaction(self, *, savepoint: bool = False) -> Iterator[None]: 

378 """Return a context manager that represents a transaction. 

379 

380 Parameters 

381 ---------- 

382 savepoint : `bool` 

383 Whether to issue a SAVEPOINT in the database. 

384 

385 Yields 

386 ------ 

387 `None` 

388 """ 

389 with self._db.transaction(savepoint=savepoint): 

390 yield 

391 

392 def resetConnectionPool(self) -> None: 

393 """Reset SQLAlchemy connection pool for `SqlRegistry` database. 

394 

395 This operation is useful when using registry with fork-based 

396 multiprocessing. To use registry across fork boundary one has to make 

397 sure that there are no currently active connections (no session or 

398 transaction is in progress) and connection pool is reset using this 

399 method. This method should be called by the child process immediately 

400 after the fork. 

401 """ 

402 self._db._engine.dispose() 

403 

404 def registerOpaqueTable(self, tableName: str, spec: ddl.TableSpec) -> None: 

405 """Add an opaque (to the `Registry`) table for use by a `Datastore` or 

406 other data repository client. 

407 

408 Opaque table records can be added via `insertOpaqueData`, retrieved via 

409 `fetchOpaqueData`, and removed via `deleteOpaqueData`. 

410 

411 Parameters 

412 ---------- 

413 tableName : `str` 

414 Logical name of the opaque table. This may differ from the 

415 actual name used in the database by a prefix and/or suffix. 

416 spec : `ddl.TableSpec` 

417 Specification for the table to be added. 

418 """ 

419 self._managers.opaque.register(tableName, spec) 

420 

421 @transactional 

422 def insertOpaqueData(self, tableName: str, *data: dict) -> None: 

423 """Insert records into an opaque table. 

424 

425 Parameters 

426 ---------- 

427 tableName : `str` 

428 Logical name of the opaque table. Must match the name used in a 

429 previous call to `registerOpaqueTable`. 

430 *data 

431 Each additional positional argument is a dictionary that represents 

432 a single row to be added. 

433 """ 

434 self._managers.opaque[tableName].insert(*data) 

435 

436 def fetchOpaqueData(self, tableName: str, **where: Any) -> Iterator[Mapping[str, Any]]: 

437 """Retrieve records from an opaque table. 

438 

439 Parameters 

440 ---------- 

441 tableName : `str` 

442 Logical name of the opaque table. Must match the name used in a 

443 previous call to `registerOpaqueTable`. 

444 **where 

445 Additional keyword arguments are interpreted as equality 

446 constraints that restrict the returned rows (combined with AND); 

447 keyword arguments are column names and values are the values they 

448 must have. 

449 

450 Yields 

451 ------ 

452 row : `dict` 

453 A dictionary representing a single result row. 

454 """ 

455 yield from self._managers.opaque[tableName].fetch(**where) 

456 

457 @transactional 

458 def deleteOpaqueData(self, tableName: str, **where: Any) -> None: 

459 """Remove records from an opaque table. 

460 

461 Parameters 

462 ---------- 

463 tableName : `str` 

464 Logical name of the opaque table. Must match the name used in a 

465 previous call to `registerOpaqueTable`. 

466 **where 

467 Additional keyword arguments are interpreted as equality 

468 constraints that restrict the deleted rows (combined with AND); 

469 keyword arguments are column names and values are the values they 

470 must have. 

471 """ 

472 self._managers.opaque[tableName].delete(where.keys(), where) 

473 

474 def registerCollection( 

475 self, name: str, type: CollectionType = CollectionType.TAGGED, doc: str | None = None 

476 ) -> bool: 

477 """Add a new collection if one with the given name does not exist. 

478 

479 Parameters 

480 ---------- 

481 name : `str` 

482 The name of the collection to create. 

483 type : `CollectionType` 

484 Enum value indicating the type of collection to create. 

485 doc : `str`, optional 

486 Documentation string for the collection. 

487 

488 Returns 

489 ------- 

490 registered : `bool` 

491 Boolean indicating whether the collection was already registered 

492 or was created by this call. 

493 

494 Notes 

495 ----- 

496 This method cannot be called within transactions, as it needs to be 

497 able to perform its own transaction to be concurrent. 

498 """ 

499 _, registered = self._managers.collections.register(name, type, doc=doc) 

500 return registered 

501 

502 def getCollectionType(self, name: str) -> CollectionType: 

503 """Return an enumeration value indicating the type of the given 

504 collection. 

505 

506 Parameters 

507 ---------- 

508 name : `str` 

509 The name of the collection. 

510 

511 Returns 

512 ------- 

513 type : `CollectionType` 

514 Enum value indicating the type of this collection. 

515 

516 Raises 

517 ------ 

518 lsst.daf.butler.registry.MissingCollectionError 

519 Raised if no collection with the given name exists. 

520 """ 

521 return self._managers.collections.find(name).type 

522 

523 def get_collection_record(self, name: str) -> CollectionRecord: 

524 """Return the record for this collection. 

525 

526 Parameters 

527 ---------- 

528 name : `str` 

529 Name of the collection for which the record is to be retrieved. 

530 

531 Returns 

532 ------- 

533 record : `CollectionRecord` 

534 The record for this collection. 

535 """ 

536 return self._managers.collections.find(name) 

537 

538 def registerRun(self, name: str, doc: str | None = None) -> bool: 

539 """Add a new run if one with the given name does not exist. 

540 

541 Parameters 

542 ---------- 

543 name : `str` 

544 The name of the run to create. 

545 doc : `str`, optional 

546 Documentation string for the collection. 

547 

548 Returns 

549 ------- 

550 registered : `bool` 

551 Boolean indicating whether a new run was registered. `False` 

552 if it already existed. 

553 

554 Notes 

555 ----- 

556 This method cannot be called within transactions, as it needs to be 

557 able to perform its own transaction to be concurrent. 

558 """ 

559 _, registered = self._managers.collections.register(name, CollectionType.RUN, doc=doc) 

560 return registered 

561 

562 @transactional 

563 def removeCollection(self, name: str) -> None: 

564 """Remove the given collection from the registry. 

565 

566 Parameters 

567 ---------- 

568 name : `str` 

569 The name of the collection to remove. 

570 

571 Raises 

572 ------ 

573 lsst.daf.butler.registry.MissingCollectionError 

574 Raised if no collection with the given name exists. 

575 sqlalchemy.exc.IntegrityError 

576 Raised if the database rows associated with the collection are 

577 still referenced by some other table, such as a dataset in a 

578 datastore (for `~CollectionType.RUN` collections only) or a 

579 `~CollectionType.CHAINED` collection of which this collection is 

580 a child. 

581 

582 Notes 

583 ----- 

584 If this is a `~CollectionType.RUN` collection, all datasets and quanta 

585 in it will removed from the `Registry` database. This requires that 

586 those datasets be removed (or at least trashed) from any datastores 

587 that hold them first. 

588 

589 A collection may not be deleted as long as it is referenced by a 

590 `~CollectionType.CHAINED` collection; the ``CHAINED`` collection must 

591 be deleted or redefined first. 

592 """ 

593 self._managers.collections.remove(name) 

594 

595 def getCollectionChain(self, parent: str) -> tuple[str, ...]: 

596 """Return the child collections in a `~CollectionType.CHAINED` 

597 collection. 

598 

599 Parameters 

600 ---------- 

601 parent : `str` 

602 Name of the chained collection. Must have already been added via 

603 a call to `Registry.registerCollection`. 

604 

605 Returns 

606 ------- 

607 children : `~collections.abc.Sequence` [ `str` ] 

608 An ordered sequence of collection names that are searched when the 

609 given chained collection is searched. 

610 

611 Raises 

612 ------ 

613 lsst.daf.butler.registry.MissingCollectionError 

614 Raised if ``parent`` does not exist in the `Registry`. 

615 lsst.daf.butler.registry.CollectionTypeError 

616 Raised if ``parent`` does not correspond to a 

617 `~CollectionType.CHAINED` collection. 

618 """ 

619 record = self._managers.collections.find(parent) 

620 if record.type is not CollectionType.CHAINED: 620 ↛ 621line 620 didn't jump to line 621 because the condition on line 620 was never true

621 raise CollectionTypeError(f"Collection '{parent}' has type {record.type.name}, not CHAINED.") 

622 assert isinstance(record, ChainedCollectionRecord) 

623 return record.children 

624 

625 @transactional 

626 def setCollectionChain(self, parent: str, children: Any, *, flatten: bool = False) -> None: 

627 """Define or redefine a `~CollectionType.CHAINED` collection. 

628 

629 Parameters 

630 ---------- 

631 parent : `str` 

632 Name of the chained collection. Must have already been added via 

633 a call to `Registry.registerCollection`. 

634 children : collection expression 

635 An expression defining an ordered search of child collections, 

636 generally an iterable of `str`; see 

637 :ref:`daf_butler_collection_expressions` for more information. 

638 flatten : `bool`, optional 

639 If `True` (`False` is default), recursively flatten out any nested 

640 `~CollectionType.CHAINED` collections in ``children`` first. 

641 

642 Raises 

643 ------ 

644 lsst.daf.butler.registry.MissingCollectionError 

645 Raised when any of the given collections do not exist in the 

646 `Registry`. 

647 lsst.daf.butler.registry.CollectionTypeError 

648 Raised if ``parent`` does not correspond to a 

649 `~CollectionType.CHAINED` collection. 

650 CollectionCycleError 

651 Raised if the given collections contains a cycle. 

652 

653 Notes 

654 ----- 

655 If this function is called within a call to ``Butler.transaction``, it 

656 will hold a lock that prevents other processes from modifying the 

657 parent collection until the end of the transaction. Keep these 

658 transactions short. 

659 """ 

660 children = CollectionWildcard.from_expression(children).require_ordered() 

661 if flatten: 

662 children = self.queryCollections(children, flattenChains=True) 

663 

664 self._managers.collections.update_chain(parent, list(children), allow_use_in_caching_context=True) 

665 

666 def getCollectionParentChains(self, collection: str) -> set[str]: 

667 """Return the CHAINED collections that directly contain the given one. 

668 

669 Parameters 

670 ---------- 

671 collection : `str` 

672 Name of the collection. 

673 

674 Returns 

675 ------- 

676 chains : `set` of `str` 

677 Set of `~CollectionType.CHAINED` collection names. 

678 """ 

679 return self._managers.collections.getParentChains(self._managers.collections.find(collection).key) 

680 

681 def getCollectionDocumentation(self, collection: str) -> str | None: 

682 """Retrieve the documentation string for a collection. 

683 

684 Parameters 

685 ---------- 

686 collection : `str` 

687 Name of the collection. 

688 

689 Returns 

690 ------- 

691 docs : `str` or `None` 

692 Docstring for the collection with the given name. 

693 """ 

694 return self._managers.collections.getDocumentation(self._managers.collections.find(collection).key) 

695 

696 def setCollectionDocumentation(self, collection: str, doc: str | None) -> None: 

697 """Set the documentation string for a collection. 

698 

699 Parameters 

700 ---------- 

701 collection : `str` 

702 Name of the collection. 

703 doc : `str` or `None` 

704 Docstring for the collection with the given name; will replace any 

705 existing docstring. Passing `None` will remove any existing 

706 docstring. 

707 """ 

708 self._managers.collections.setDocumentation(self._managers.collections.find(collection).key, doc) 

709 

710 def getCollectionSummary(self, collection: str) -> CollectionSummary: 

711 """Return a summary for the given collection. 

712 

713 Parameters 

714 ---------- 

715 collection : `str` 

716 Name of the collection for which a summary is to be retrieved. 

717 

718 Returns 

719 ------- 

720 summary : `~lsst.daf.butler.registry.CollectionSummary` 

721 Summary of the dataset types and governor dimension values in 

722 this collection. 

723 """ 

724 record = self._managers.collections.find(collection) 

725 return self._managers.datasets.getCollectionSummary(record) 

726 

727 def registerDatasetType(self, datasetType: DatasetType) -> bool: 

728 """Add a new `DatasetType` to the Registry. 

729 

730 It is not an error to register the same `DatasetType` twice. 

731 

732 Parameters 

733 ---------- 

734 datasetType : `DatasetType` 

735 The `DatasetType` to be added. 

736 

737 Returns 

738 ------- 

739 inserted : `bool` 

740 `True` if ``datasetType`` was inserted, `False` if an identical 

741 existing `DatasetType` was found. Note that in either case the 

742 DatasetType is guaranteed to be defined in the Registry 

743 consistently with the given definition. 

744 

745 Raises 

746 ------ 

747 ValueError 

748 Raised if the dimensions or storage class are invalid. 

749 lsst.daf.butler.registry.ConflictingDefinitionError 

750 Raised if this `DatasetType` is already registered with a different 

751 definition. 

752 

753 Notes 

754 ----- 

755 This method cannot be called within transactions, as it needs to be 

756 able to perform its own transaction to be concurrent. 

757 """ 

758 return self._managers.datasets.register_dataset_type(datasetType) 

759 

760 def removeDatasetType(self, name: str | tuple[str, ...]) -> None: 

761 """Remove the named `DatasetType` from the registry. 

762 

763 .. warning:: 

764 

765 Registry implementations can cache the dataset type definitions. 

766 This means that deleting the dataset type definition may result in 

767 unexpected behavior from other butler processes that are active 

768 that have not seen the deletion. 

769 

770 Parameters 

771 ---------- 

772 name : `str` or `tuple` [`str`] 

773 Name of the type to be removed or tuple containing a list of type 

774 names to be removed. Wildcards are allowed. 

775 

776 Raises 

777 ------ 

778 lsst.daf.butler.registry.OrphanedRecordError 

779 Raised if an attempt is made to remove the dataset type definition 

780 when there are already datasets associated with it. 

781 

782 Notes 

783 ----- 

784 If the dataset type is not registered the method will return without 

785 action. 

786 """ 

787 for datasetTypeExpression in ensure_iterable(name): 

788 # Catch any warnings from the caller specifying a component 

789 # dataset type. This will result in an error later but the 

790 # warning could be confusing when the caller is not querying 

791 # anything. 

792 with warnings.catch_warnings(): 

793 warnings.simplefilter("ignore", category=FutureWarning) 

794 datasetTypes = list(self.queryDatasetTypes(datasetTypeExpression)) 

795 if not datasetTypes: 

796 _LOG.info("Dataset type %r not defined", datasetTypeExpression) 

797 else: 

798 for datasetType in datasetTypes: 

799 self._managers.datasets.remove_dataset_type(datasetType.name) 

800 _LOG.info("Removed dataset type %r", datasetType.name) 

801 

802 def getDatasetType(self, name: str) -> DatasetType: 

803 """Get the `DatasetType`. 

804 

805 Parameters 

806 ---------- 

807 name : `str` 

808 Name of the type. 

809 

810 Returns 

811 ------- 

812 type : `DatasetType` 

813 The `DatasetType` associated with the given name. 

814 

815 Raises 

816 ------ 

817 lsst.daf.butler.DatasetTypeExpressionError 

818 Raised if ``name`` is not a valid dataset type name. 

819 lsst.daf.butler.MissingDatasetTypeError 

820 Raised if the requested dataset type has not been registered. 

821 

822 Notes 

823 ----- 

824 This method handles component dataset types automatically, though most 

825 other registry operations do not. 

826 """ 

827 validate_dataset_type_name(name) 

828 parent_name, component = DatasetType.splitDatasetTypeName(name) 

829 parent_dataset_type = self._managers.datasets.get_dataset_type(parent_name) 

830 if component is None: 

831 return parent_dataset_type 

832 else: 

833 return parent_dataset_type.makeComponentDatasetType(component) 

834 

835 def supportsIdGenerationMode(self, mode: DatasetIdGenEnum) -> bool: 

836 """Test whether the given dataset ID generation mode is supported by 

837 `insertDatasets`. 

838 

839 Parameters 

840 ---------- 

841 mode : `DatasetIdGenEnum` 

842 Enum value for the mode to test. 

843 

844 Returns 

845 ------- 

846 supported : `bool` 

847 Whether the given mode is supported. 

848 """ 

849 return True 

850 

851 @transactional 

852 def insertDatasets( 

853 self, 

854 datasetType: DatasetType | str, 

855 dataIds: Iterable[DataId], 

856 run: str | None = None, 

857 expand: bool = True, 

858 idGenerationMode: DatasetIdGenEnum = DatasetIdGenEnum.UNIQUE, 

859 ) -> list[DatasetRef]: 

860 """Insert one or more datasets into the `Registry`. 

861 

862 This always adds new datasets; to associate existing datasets with 

863 a new collection, use ``associate``. 

864 

865 Parameters 

866 ---------- 

867 datasetType : `DatasetType` or `str` 

868 A `DatasetType` or the name of one. 

869 dataIds : `~collections.abc.Iterable` of `dict` or `DataCoordinate` 

870 Dimension-based identifiers for the new datasets. 

871 run : `str`, optional 

872 The name of the run that produced the datasets. Defaults to 

873 ``self.defaults.run``. 

874 expand : `bool`, optional 

875 If `True` (default), expand data IDs as they are inserted. This is 

876 necessary in general to allow datastore to generate file templates, 

877 but it may be disabled if the caller can guarantee this is 

878 unnecessary. 

879 idGenerationMode : `DatasetIdGenEnum`, optional 

880 Specifies option for generating dataset IDs. By default unique IDs 

881 are generated for each inserted dataset. 

882 

883 Returns 

884 ------- 

885 refs : `list` of `DatasetRef` 

886 Resolved `DatasetRef` instances for all given data IDs (in the same 

887 order). 

888 

889 Raises 

890 ------ 

891 lsst.daf.butler.registry.DatasetTypeError 

892 Raised if ``datasetType`` is not known to registry. 

893 lsst.daf.butler.registry.CollectionTypeError 

894 Raised if ``run`` collection type is not `~CollectionType.RUN`. 

895 lsst.daf.butler.registry.NoDefaultCollectionError 

896 Raised if ``run`` is `None` and ``self.defaults.run`` is `None`. 

897 lsst.daf.butler.registry.ConflictingDefinitionError 

898 If a dataset with the same dataset type and data ID as one of those 

899 given already exists in ``run``. 

900 lsst.daf.butler.registry.MissingCollectionError 

901 Raised if ``run`` does not exist in the registry. 

902 """ 

903 datasetType = self._managers.datasets.conform_exact_dataset_type(datasetType) 

904 if run is None: 

905 if self.defaults.run is None: 

906 raise NoDefaultCollectionError( 

907 "No run provided to insertDatasets, and no default from registry construction." 

908 ) 

909 run = self.defaults.run 

910 runRecord = self._managers.collections.find(run) 

911 if runRecord.type is not CollectionType.RUN: 911 ↛ 912line 911 didn't jump to line 912 because the condition on line 911 was never true

912 raise CollectionTypeError( 

913 f"Given collection is of type {runRecord.type.name}; RUN collection required." 

914 ) 

915 assert isinstance(runRecord, RunRecord) 

916 

917 expandedDataIds = [ 

918 DataCoordinate.standardize(dataId, dimensions=datasetType.dimensions) for dataId in dataIds 

919 ] 

920 if expand: 920 ↛ 925line 920 didn't jump to line 925 because the condition on line 920 was always true

921 _LOG.debug("Expanding %d data IDs", len(expandedDataIds)) 

922 expandedDataIds = self.expand_data_ids(expandedDataIds) 

923 _LOG.debug("Finished expanding data IDs") 

924 

925 try: 

926 refs = list( 

927 self._managers.datasets.insert(datasetType.name, runRecord, expandedDataIds, idGenerationMode) 

928 ) 

929 if self._managers.obscore: 

930 self._managers.obscore.add_datasets(refs) 

931 except sqlalchemy.exc.IntegrityError as err: 

932 raise ConflictingDefinitionError( 

933 "A database constraint failure was triggered by inserting " 

934 f"one or more datasets of type {datasetType} into " 

935 f"collection '{run}'. " 

936 "This probably means a dataset with the same data ID " 

937 "and dataset type already exists, but it may also mean a " 

938 "dimension row is missing." 

939 ) from err 

940 return refs 

941 

942 @transactional 

943 def _importDatasets( 

944 self, 

945 datasets: Iterable[DatasetRef], 

946 expand: bool = True, 

947 assume_new: bool = False, 

948 ) -> list[DatasetRef]: 

949 """Import one or more datasets into the `Registry`. 

950 

951 This differs from `insertDatasets` method in that this method accepts 

952 `DatasetRef` instances, which already have a dataset ID. 

953 

954 Parameters 

955 ---------- 

956 datasets : `~collections.abc.Iterable` of `DatasetRef` 

957 Datasets to be inserted. All `DatasetRef` instances must have 

958 identical ``run`` attributes. ``run`` 

959 attribute can be `None` and defaults to ``self.defaults.run``. 

960 Datasets can specify ``id`` attribute which will be used for 

961 inserted datasets. 

962 Datasets can be of multiple dataset types, but all the dataset 

963 types must have the same set of dimensions. 

964 expand : `bool`, optional 

965 If `True` (default), expand data IDs as they are inserted. This is 

966 necessary in general, but it may be disabled if the caller can 

967 guarantee this is unnecessary. 

968 assume_new : `bool`, optional 

969 If `True`, assume datasets are new. If `False`, datasets that are 

970 identical to an existing one are ignored. 

971 

972 Returns 

973 ------- 

974 refs : `list` of `DatasetRef` 

975 `DatasetRef` instances for all given data IDs (in the same order). 

976 If any of ``datasets`` has an ID which already exists in the 

977 database then it will not be inserted or updated, but a 

978 `DatasetRef` will be returned for it in any case. 

979 

980 Raises 

981 ------ 

982 lsst.daf.butler.registry.NoDefaultCollectionError 

983 Raised if ``run`` is `None` and ``self.defaults.run`` is `None`. 

984 lsst.daf.butler.registry.DatasetTypeError 

985 Raised if a dataset type is not known to registry. 

986 lsst.daf.butler.registry.ConflictingDefinitionError 

987 If a dataset with the same dataset type and data ID as one of those 

988 given already exists in ``run``, or if ``assume_new=True`` and at 

989 least one dataset is not new. 

990 lsst.daf.butler.registry.MissingCollectionError 

991 Raised if ``run`` does not exist in the registry. 

992 

993 Notes 

994 ----- 

995 This method is considered middleware-internal. 

996 """ 

997 datasets = list(datasets) 

998 if not datasets: 998 ↛ 1000line 998 didn't jump to line 1000 because the condition on line 998 was never true

999 # nothing to do 

1000 return [] 

1001 

1002 # find run name 

1003 runs = {dataset.run for dataset in datasets} 

1004 if len(runs) != 1: 1004 ↛ 1005line 1004 didn't jump to line 1005 because the condition on line 1004 was never true

1005 raise ValueError(f"Multiple run names in input datasets: {runs}") 

1006 run = runs.pop() 

1007 

1008 runRecord = self._managers.collections.find(run) 

1009 if runRecord.type is not CollectionType.RUN: 

1010 raise CollectionTypeError( 

1011 f"Given collection '{runRecord.name}' is of type {runRecord.type.name};" 

1012 " RUN collection required." 

1013 ) 

1014 assert isinstance(runRecord, RunRecord) 

1015 

1016 if expand: 

1017 _LOG.debug("Expanding %d data IDs", len(datasets)) 

1018 datasets = self.expand_refs(datasets) 

1019 _LOG.debug("Finished expanding data IDs") 

1020 

1021 try: 

1022 self._managers.datasets.import_(runRecord, datasets, assume_new=assume_new) 

1023 if self._managers.obscore: 

1024 self._managers.obscore.add_datasets(datasets) 

1025 except sqlalchemy.exc.IntegrityError as err: 

1026 raise ConflictingDefinitionError( 

1027 "A database constraint failure was triggered by inserting " 

1028 f"one or more datasets into collection '{run}'. " 

1029 "This probably means a dataset with the same data ID " 

1030 "and dataset type already exists, but it may also mean a " 

1031 "dimension row is missing, or the dataset was assumed to be " 

1032 "new when it was not." 

1033 ) from err 

1034 return datasets 

1035 

1036 def getDataset(self, id: DatasetId) -> DatasetRef | None: 

1037 """Retrieve a Dataset entry. 

1038 

1039 Parameters 

1040 ---------- 

1041 id : `DatasetId` 

1042 The unique identifier for the dataset. 

1043 

1044 Returns 

1045 ------- 

1046 ref : `DatasetRef` or `None` 

1047 A ref to the Dataset, or `None` if no matching Dataset 

1048 was found. 

1049 """ 

1050 refs = self._managers.datasets.get_dataset_refs([id]) 

1051 if len(refs) == 0: 

1052 return None 

1053 else: 

1054 return refs[0] 

1055 

1056 def _fetch_run_dataset_ids(self, run: str) -> list[DatasetId]: 

1057 """Return the IDs of all datasets in the given ``RUN`` 

1058 collection. 

1059 

1060 Parameters 

1061 ---------- 

1062 run : `str` 

1063 Name of the collection. 

1064 

1065 Returns 

1066 ------- 

1067 dataset_ids : `list` [`uuid.UUID`] 

1068 List of dataset IDs. 

1069 

1070 Notes 

1071 ----- 

1072 This is a middleware-internal interface. 

1073 """ 

1074 run_record = self._managers.collections.find(run) 

1075 if not isinstance(run_record, RunRecord): 1075 ↛ 1076line 1075 didn't jump to line 1076 because the condition on line 1075 was never true

1076 raise CollectionTypeError(f"{run!r} is not a RUN collection.") 

1077 return self._managers.datasets.fetch_run_dataset_ids(run_record) 

1078 

1079 @transactional 

1080 def removeDatasets(self, refs: Iterable[DatasetRef]) -> None: 

1081 """Remove datasets from the Registry. 

1082 

1083 The datasets will be removed unconditionally from all collections. 

1084 `Datastore` records will *not* be deleted; the caller is responsible 

1085 for ensuring that the dataset has already been removed from all 

1086 Datastores. 

1087 

1088 Parameters 

1089 ---------- 

1090 refs : `~collections.abc.Iterable` [`DatasetRef`] 

1091 References to the datasets to be removed. Should be considered 

1092 invalidated upon return. 

1093 

1094 Raises 

1095 ------ 

1096 lsst.daf.butler.registry.OrphanedRecordError 

1097 Raised if any dataset is still present in any `Datastore`. 

1098 """ 

1099 try: 

1100 self._managers.datasets.delete(refs) 

1101 except sqlalchemy.exc.IntegrityError as err: 

1102 raise OrphanedRecordError( 

1103 "One or more datasets is still present in one or more Datastores." 

1104 ) from err 

1105 

1106 @transactional 

1107 def associate(self, collection: str, refs: Iterable[DatasetRef]) -> None: 

1108 """Add existing datasets to a `~CollectionType.TAGGED` collection. 

1109 

1110 If a DatasetRef with the same exact ID is already in a collection 

1111 nothing is changed. If a `DatasetRef` with the same `DatasetType` and 

1112 data ID but with different ID exists in the collection, 

1113 `~lsst.daf.butler.registry.ConflictingDefinitionError` is raised. 

1114 

1115 Parameters 

1116 ---------- 

1117 collection : `str` 

1118 Indicates the collection the datasets should be associated with. 

1119 refs : `~collections.abc.Iterable` [ `DatasetRef` ] 

1120 An iterable of resolved `DatasetRef` instances that already exist 

1121 in this `Registry`. 

1122 

1123 Raises 

1124 ------ 

1125 lsst.daf.butler.registry.ConflictingDefinitionError 

1126 If a Dataset with the given `DatasetRef` already exists in the 

1127 given collection. 

1128 lsst.daf.butler.registry.MissingCollectionError 

1129 Raised if ``collection`` does not exist in the registry. 

1130 lsst.daf.butler.registry.CollectionTypeError 

1131 Raise adding new datasets to the given ``collection`` is not 

1132 allowed. 

1133 """ 

1134 progress = Progress("lsst.daf.butler.Registry.associate", level=logging.DEBUG) 

1135 collectionRecord = self._managers.collections.find(collection) 

1136 for datasetType, refsForType in progress.iter_item_chunks( 

1137 DatasetRef.iter_by_type(refs), desc="Associating datasets by type" 

1138 ): 

1139 try: 

1140 self._managers.datasets.associate(datasetType, collectionRecord, refsForType) 

1141 if self._managers.obscore: 

1142 # If a TAGGED collection is being monitored by ObsCore 

1143 # manager then we may need to save the dataset. 

1144 self._managers.obscore.associate(refsForType, collectionRecord) 

1145 except sqlalchemy.exc.IntegrityError as err: 

1146 raise ConflictingDefinitionError( 

1147 f"Constraint violation while associating dataset of type {datasetType.name} with " 

1148 f"collection {collection}. This probably means that one or more datasets with the same " 

1149 "dataset type and data ID already exist in the collection, but it may also indicate " 

1150 "that the datasets do not exist." 

1151 ) from err 

1152 

1153 @transactional 

1154 def disassociate(self, collection: str, refs: Iterable[DatasetRef]) -> None: 

1155 """Remove existing datasets from a `~CollectionType.TAGGED` collection. 

1156 

1157 ``collection`` and ``ref`` combinations that are not currently 

1158 associated are silently ignored. 

1159 

1160 Parameters 

1161 ---------- 

1162 collection : `str` 

1163 The collection the datasets should no longer be associated with. 

1164 refs : `~collections.abc.Iterable` [ `DatasetRef` ] 

1165 An iterable of resolved `DatasetRef` instances that already exist 

1166 in this `Registry`. 

1167 

1168 Raises 

1169 ------ 

1170 lsst.daf.butler.AmbiguousDatasetError 

1171 Raised if any of the given dataset references is unresolved. 

1172 lsst.daf.butler.registry.MissingCollectionError 

1173 Raised if ``collection`` does not exist in the registry. 

1174 lsst.daf.butler.registry.CollectionTypeError 

1175 Raise adding new datasets to the given ``collection`` is not 

1176 allowed. 

1177 """ 

1178 progress = Progress("lsst.daf.butler.Registry.disassociate", level=logging.DEBUG) 

1179 collectionRecord = self._managers.collections.find(collection) 

1180 for datasetType, refsForType in progress.iter_item_chunks( 

1181 DatasetRef.iter_by_type(refs), desc="Disassociating datasets by type" 

1182 ): 

1183 self._managers.datasets.disassociate(datasetType, collectionRecord, refsForType) 

1184 if self._managers.obscore: 

1185 self._managers.obscore.disassociate(refsForType, collectionRecord) 

1186 

1187 @transactional 

1188 def certify(self, collection: str, refs: Iterable[DatasetRef], timespan: Timespan) -> None: 

1189 """Associate one or more datasets with a calibration collection and a 

1190 validity range within it. 

1191 

1192 Parameters 

1193 ---------- 

1194 collection : `str` 

1195 The name of an already-registered `~CollectionType.CALIBRATION` 

1196 collection. 

1197 refs : `~collections.abc.Iterable` [ `DatasetRef` ] 

1198 Datasets to be associated. 

1199 timespan : `Timespan` 

1200 The validity range for these datasets within the collection. 

1201 

1202 Raises 

1203 ------ 

1204 lsst.daf.butler.AmbiguousDatasetError 

1205 Raised if any of the given `DatasetRef` instances is unresolved. 

1206 lsst.daf.butler.registry.ConflictingDefinitionError 

1207 Raised if the collection already contains a different dataset with 

1208 the same `DatasetType` and data ID and an overlapping validity 

1209 range. 

1210 DatasetTypeError 

1211 Raised if ``ref.datasetType.isCalibration() is False`` for any ref. 

1212 CollectionTypeError 

1213 Raised if 

1214 ``collection.type is not CollectionType.CALIBRATION``. 

1215 """ 

1216 progress = Progress("lsst.daf.butler.Registry.certify", level=logging.DEBUG) 

1217 with self._managers.caching_context.enable_collection_record_cache(): 

1218 collectionRecord = self._managers.collections.find(collection) 

1219 for datasetType, refsForType in progress.iter_item_chunks( 

1220 DatasetRef.iter_by_type(refs), desc="Certifying datasets by type" 

1221 ): 

1222 self._managers.datasets.certify( 

1223 datasetType, collectionRecord, refsForType, timespan, self._query 

1224 ) 

1225 

1226 @transactional 

1227 def decertify( 

1228 self, 

1229 collection: str, 

1230 datasetType: str | DatasetType, 

1231 timespan: Timespan, 

1232 *, 

1233 dataIds: Iterable[DataId] | None = None, 

1234 ) -> None: 

1235 """Remove or adjust datasets to clear a validity range within a 

1236 calibration collection. 

1237 

1238 Parameters 

1239 ---------- 

1240 collection : `str` 

1241 The name of an already-registered `~CollectionType.CALIBRATION` 

1242 collection. 

1243 datasetType : `str` or `DatasetType` 

1244 Name or `DatasetType` instance for the datasets to be decertified. 

1245 timespan : `Timespan`, optional 

1246 The validity range to remove datasets from within the collection. 

1247 Datasets that overlap this range but are not contained by it will 

1248 have their validity ranges adjusted to not overlap it, which may 

1249 split a single dataset validity range into two. 

1250 dataIds : `~collections.abc.Iterable` [`dict` or `DataCoordinate`], \ 

1251 optional 

1252 Data IDs that should be decertified within the given validity range 

1253 If `None`, all data IDs for ``self.datasetType`` will be 

1254 decertified. 

1255 

1256 Raises 

1257 ------ 

1258 DatasetTypeError 

1259 Raised if ``datasetType.isCalibration() is False``. 

1260 CollectionTypeError 

1261 Raised if 

1262 ``collection.type is not CollectionType.CALIBRATION``. 

1263 """ 

1264 collectionRecord = self._managers.collections.find(collection) 

1265 if isinstance(datasetType, str): 1265 ↛ 1267line 1265 didn't jump to line 1267 because the condition on line 1265 was always true

1266 datasetType = self.getDatasetType(datasetType) 

1267 standardizedDataIds = None 

1268 if dataIds is not None: 

1269 standardizedDataIds = [ 

1270 DataCoordinate.standardize(d, dimensions=datasetType.dimensions) for d in dataIds 

1271 ] 

1272 self._managers.datasets.decertify( 

1273 datasetType, collectionRecord, timespan, data_ids=standardizedDataIds, query_func=self._query 

1274 ) 

1275 

1276 def getDatastoreBridgeManager(self) -> DatastoreRegistryBridgeManager: 

1277 """Return an object that allows a new `Datastore` instance to 

1278 communicate with this `Registry`. 

1279 

1280 Returns 

1281 ------- 

1282 manager : `~.interfaces.DatastoreRegistryBridgeManager` 

1283 Object that mediates communication between this `Registry` and its 

1284 associated datastores. 

1285 """ 

1286 return self._managers.datastores 

1287 

1288 def getDatasetLocations(self, ref: DatasetRef) -> Iterable[str]: 

1289 """Retrieve datastore locations for a given dataset. 

1290 

1291 Parameters 

1292 ---------- 

1293 ref : `DatasetRef` 

1294 A reference to the dataset for which to retrieve storage 

1295 information. 

1296 

1297 Returns 

1298 ------- 

1299 datastores : `~collections.abc.Iterable` [ `str` ] 

1300 All the matching datastores holding this dataset. 

1301 

1302 Raises 

1303 ------ 

1304 lsst.daf.butler.AmbiguousDatasetError 

1305 Raised if ``ref.id`` is `None`. 

1306 """ 

1307 return self._managers.datastores.findDatastores(ref) 

1308 

1309 def expandDataId( 

1310 self, 

1311 dataId: DataId | None = None, 

1312 *, 

1313 dimensions: Iterable[str] | DimensionGroup | None = None, 

1314 records: Mapping[str, DimensionRecord | None] | None = None, 

1315 withDefaults: bool = True, 

1316 **kwargs: Any, 

1317 ) -> DataCoordinate: 

1318 """Expand a dimension-based data ID to include additional information. 

1319 

1320 Parameters 

1321 ---------- 

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

1323 Data ID to be expanded; augmented and overridden by ``kwargs``. 

1324 dimensions : `~collections.abc.Iterable` [ `str` ], \ 

1325 `DimensionGroup`, optional 

1326 The dimensions to be identified by the new `DataCoordinate`. 

1327 If not provided, will be inferred from the keys of ``dataId`` and 

1328 ``**kwargs``, and ``universe`` must be provided unless ``dataId`` 

1329 is already a `DataCoordinate`. 

1330 records : `~collections.abc.Mapping` [`str`, `DimensionRecord`], \ 

1331 optional 

1332 Dimension record data to use before querying the database for that 

1333 data, keyed by element name. 

1334 withDefaults : `bool`, optional 

1335 Utilize ``self.defaults.dataId`` to fill in missing governor 

1336 dimension key-value pairs. Defaults to `True` (i.e. defaults are 

1337 used). 

1338 **kwargs 

1339 Additional keywords are treated like additional key-value pairs for 

1340 ``dataId``, extending and overriding. 

1341 

1342 Returns 

1343 ------- 

1344 expanded : `DataCoordinate` 

1345 A data ID that includes full metadata for all of the dimensions it 

1346 identifies, i.e. guarantees that ``expanded.hasRecords()`` and 

1347 ``expanded.hasFull()`` both return `True`. 

1348 

1349 Raises 

1350 ------ 

1351 lsst.daf.butler.registry.DataIdError 

1352 Raised when ``dataId`` or keyword arguments specify unknown 

1353 dimensions or values, or when a resulting data ID contains 

1354 contradictory key-value pairs, according to dimension 

1355 relationships. 

1356 

1357 Notes 

1358 ----- 

1359 This method cannot be relied upon to reject invalid data ID values 

1360 for dimensions that do actually not have any record columns. For 

1361 efficiency reasons the records for these dimensions (which have only 

1362 dimension key values that are given by the caller) may be constructed 

1363 directly rather than obtained from the registry database. 

1364 """ 

1365 if not withDefaults: 

1366 defaults = None 

1367 else: 

1368 defaults = self.defaults.dataId 

1369 standardized = DataCoordinate.standardize( 

1370 dataId, 

1371 dimensions=dimensions, 

1372 universe=self.dimensions, 

1373 defaults=defaults, 

1374 **kwargs, 

1375 ) 

1376 if standardized.hasRecords(): 

1377 return standardized 

1378 if records is None: 1378 ↛ 1381line 1378 didn't jump to line 1381 because the condition on line 1378 was always true

1379 records = {} 

1380 else: 

1381 records = dict(records) 

1382 if isinstance(dataId, DataCoordinate) and dataId.hasRecords() and not kwargs: 1382 ↛ 1383line 1382 didn't jump to line 1383 because the condition on line 1382 was never true

1383 for element_name in dataId.dimensions.elements: 

1384 records[element_name] = dataId.records[element_name] 

1385 keys: dict[str, str | int] = dict(standardized.mapping) 

1386 for element_name in standardized.dimensions.lookup_order: 

1387 element = self.dimensions[element_name] 

1388 record = records.get(element_name, ...) # Use ... to mean not found; None might mean NULL 

1389 if record is ...: 1389 ↛ 1399line 1389 didn't jump to line 1399 because the condition on line 1389 was always true

1390 if element_name in self.dimensions.dimensions.names and keys.get(element_name) is None: 1390 ↛ 1391line 1390 didn't jump to line 1391 because the condition on line 1390 was never true

1391 raise DimensionNameError(f"No value or null value for dimension {element_name}.") 

1392 else: 

1393 record = self._managers.dimensions.fetch_one( 

1394 element_name, 

1395 DataCoordinate.standardize(keys, dimensions=element.minimal_group), 

1396 self.dimension_record_cache, 

1397 ) 

1398 records[element_name] = record 

1399 if record is not None: 

1400 for d in element.implied: 

1401 value = getattr(record, d.name) 

1402 if keys.setdefault(d.name, value) != value: 1402 ↛ 1403line 1402 didn't jump to line 1403 because the condition on line 1402 was never true

1403 raise InconsistentDataIdError( 

1404 f"Data ID {standardized} has {d.name}={keys[d.name]!r}, " 

1405 f"but {element_name} implies {d.name}={value!r}." 

1406 ) 

1407 else: 

1408 if element_name in standardized.dimensions.names: 

1409 raise DataIdValueError( 

1410 f"Could not fetch record for dimension {element.name} via keys {keys}." 

1411 ) 

1412 if element.defines_relationships: 

1413 raise InconsistentDataIdError( 

1414 f"Could not fetch record for element {element_name} via keys {keys}, " 

1415 "but it is marked as defining relationships; this means one or more dimensions are " 

1416 "have inconsistent values.", 

1417 ) 

1418 return DataCoordinate.standardize(keys, dimensions=standardized.dimensions).expanded(records=records) 

1419 

1420 def expand_data_ids(self, data_ids: Iterable[DataCoordinate]) -> list[DataCoordinate]: 

1421 return expand_data_ids(data_ids, self.dimensions, self._query, self.dimension_record_cache) 

1422 

1423 def expand_refs(self, dataset_refs: list[DatasetRef]) -> list[DatasetRef]: 

1424 expanded_ids = self.expand_data_ids([ref.dataId for ref in dataset_refs]) 

1425 return [ref.expanded(data_id) for ref, data_id in zip(dataset_refs, expanded_ids)] 

1426 

1427 def insertDimensionData( 

1428 self, 

1429 element: DimensionElement | str, 

1430 *data: Mapping[str, Any] | DimensionRecord, 

1431 conform: bool = True, 

1432 replace: bool = False, 

1433 skip_existing: bool = False, 

1434 ) -> None: 

1435 """Insert one or more dimension records into the database. 

1436 

1437 Parameters 

1438 ---------- 

1439 element : `DimensionElement` or `str` 

1440 The `DimensionElement` or name thereof that identifies the table 

1441 records will be inserted into. 

1442 *data : `dict` or `DimensionRecord` 

1443 One or more records to insert. 

1444 conform : `bool`, optional 

1445 If `False` (`True` is default) perform no checking or conversions, 

1446 and assume that ``element`` is a `DimensionElement` instance and 

1447 ``data`` is a one or more `DimensionRecord` instances of the 

1448 appropriate subclass. 

1449 replace : `bool`, optional 

1450 If `True` (`False` is default), replace existing records in the 

1451 database if there is a conflict. 

1452 skip_existing : `bool`, optional 

1453 If `True` (`False` is default), skip insertion if a record with 

1454 the same primary key values already exists. Unlike 

1455 `syncDimensionData`, this will not detect when the given record 

1456 differs from what is in the database, and should not be used when 

1457 this is a concern. 

1458 """ 

1459 if isinstance(element, str): 

1460 element = self.dimensions[element] 

1461 if conform: 1461 ↛ 1467line 1461 didn't jump to line 1467 because the condition on line 1461 was always true

1462 records = [ 

1463 row if isinstance(row, DimensionRecord) else element.RecordClass(**row) for row in data 

1464 ] 

1465 else: 

1466 # Ignore typing since caller said to trust them with conform=False. 

1467 records = data # type: ignore 

1468 if element.name in self.dimension_record_cache: 

1469 self.dimension_record_cache.reset() 

1470 self._managers.dimensions.insert( 

1471 element, 

1472 *records, 

1473 replace=replace, 

1474 skip_existing=skip_existing, 

1475 ) 

1476 

1477 def syncDimensionData( 

1478 self, 

1479 element: DimensionElement | str, 

1480 row: Mapping[str, Any] | DimensionRecord, 

1481 conform: bool = True, 

1482 update: bool = False, 

1483 ) -> bool | dict[str, Any]: 

1484 """Synchronize the given dimension record with the database, inserting 

1485 if it does not already exist and comparing values if it does. 

1486 

1487 Parameters 

1488 ---------- 

1489 element : `DimensionElement` or `str` 

1490 The `DimensionElement` or name thereof that identifies the table 

1491 records will be inserted into. 

1492 row : `dict` or `DimensionRecord` 

1493 The record to insert. 

1494 conform : `bool`, optional 

1495 If `False` (`True` is default) perform no checking or conversions, 

1496 and assume that ``element`` is a `DimensionElement` instance and 

1497 ``data`` is a `DimensionRecord` instances of the appropriate 

1498 subclass. 

1499 update : `bool`, optional 

1500 If `True` (`False` is default), update the existing record in the 

1501 database if there is a conflict. 

1502 

1503 Returns 

1504 ------- 

1505 inserted_or_updated : `bool` or `dict` 

1506 `True` if a new row was inserted, `False` if no changes were 

1507 needed, or a `dict` mapping updated column names to their old 

1508 values if an update was performed (only possible if 

1509 ``update=True``). 

1510 

1511 Raises 

1512 ------ 

1513 lsst.daf.butler.registry.ConflictingDefinitionError 

1514 Raised if the record exists in the database (according to primary 

1515 key lookup) but is inconsistent with the given one. 

1516 """ 

1517 if conform: 1517 ↛ 1523line 1517 didn't jump to line 1523 because the condition on line 1517 was always true

1518 if isinstance(element, str): 

1519 element = self.dimensions[element] 

1520 record = row if isinstance(row, DimensionRecord) else element.RecordClass(**row) 

1521 else: 

1522 # Ignore typing since caller said to trust them with conform=False. 

1523 record = row # type: ignore 

1524 if record.definition.name in self.dimension_record_cache: 

1525 self.dimension_record_cache.reset() 

1526 return self._managers.dimensions.sync(record, update=update) 

1527 

1528 def queryDatasetTypes( 

1529 self, 

1530 expression: Any = ..., 

1531 *, 

1532 missing: list[str] | None = None, 

1533 ) -> Iterable[DatasetType]: 

1534 """Iterate over the dataset types whose names match an expression. 

1535 

1536 Parameters 

1537 ---------- 

1538 expression : dataset type expression, optional 

1539 An expression that fully or partially identifies the dataset types 

1540 to return, such as a `str`, `re.Pattern`, or iterable thereof. 

1541 ``...`` can be used to return all dataset types, and is the 

1542 default. See :ref:`daf_butler_dataset_type_expressions` for more 

1543 information. 

1544 missing : `list` of `str`, optional 

1545 String dataset type names that were explicitly given (i.e. not 

1546 regular expression patterns) but not found will be appended to this 

1547 list, if it is provided. 

1548 

1549 Returns 

1550 ------- 

1551 dataset_types : `~collections.abc.Iterable` [ `DatasetType`] 

1552 An `~collections.abc.Iterable` of `DatasetType` instances whose 

1553 names match ``expression``. 

1554 

1555 Raises 

1556 ------ 

1557 lsst.daf.butler.DatasetTypeExpressionError 

1558 Raised when ``expression`` is invalid. 

1559 """ 

1560 wildcard = DatasetTypeWildcard.from_expression(expression) 

1561 return self._managers.datasets.resolve_wildcard(wildcard, missing=missing) 

1562 

1563 def queryCollections( 

1564 self, 

1565 expression: Any = ..., 

1566 datasetType: DatasetType | None = None, 

1567 collectionTypes: Iterable[CollectionType] | CollectionType = CollectionType.all(), 

1568 flattenChains: bool = False, 

1569 includeChains: bool | None = None, 

1570 ) -> Sequence[str]: 

1571 """Iterate over the collections whose names match an expression. 

1572 

1573 Parameters 

1574 ---------- 

1575 expression : collection expression, optional 

1576 An expression that identifies the collections to return, such as 

1577 a `str` (for full matches or partial matches via globs), 

1578 `re.Pattern` (for partial matches), or iterable thereof. ``...`` 

1579 can be used to return all collections, and is the default. 

1580 See :ref:`daf_butler_collection_expressions` for more information. 

1581 datasetType : `DatasetType`, optional 

1582 If provided, only yield collections that may contain datasets of 

1583 this type. This is a conservative approximation in general; it may 

1584 yield collections that do not have any such datasets. 

1585 collectionTypes : `~collections.abc.Set` [`CollectionType`] or \ 

1586 `CollectionType`, optional 

1587 If provided, only yield collections of these types. 

1588 flattenChains : `bool`, optional 

1589 If `True` (`False` is default), recursively yield the child 

1590 collections of matching `~CollectionType.CHAINED` collections. 

1591 includeChains : `bool`, optional 

1592 If `True`, yield records for matching `~CollectionType.CHAINED` 

1593 collections. Default is the opposite of ``flattenChains``: include 

1594 either CHAINED collections or their children, but not both. 

1595 

1596 Returns 

1597 ------- 

1598 collections : `~collections.abc.Sequence` [ `str` ] 

1599 The names of collections that match ``expression``. 

1600 

1601 Raises 

1602 ------ 

1603 lsst.daf.butler.registry.CollectionExpressionError 

1604 Raised when ``expression`` is invalid. 

1605 

1606 Notes 

1607 ----- 

1608 The order in which collections are returned is unspecified, except that 

1609 the children of a `~CollectionType.CHAINED` collection are guaranteed 

1610 to be in the order in which they are searched. When multiple parent 

1611 `~CollectionType.CHAINED` collections match the same criteria, the 

1612 order in which the two lists appear is unspecified, and the lists of 

1613 children may be incomplete if a child has multiple parents. 

1614 """ 

1615 # Right now the datasetTypes argument is completely ignored, but that 

1616 # is consistent with its [lack of] guarantees. DM-24939 or a follow-up 

1617 # ticket will take care of that. 

1618 if datasetType is not None: 1618 ↛ 1619line 1618 didn't jump to line 1619 because the condition on line 1618 was never true

1619 warnings.warn( 

1620 "The datasetType parameter should no longer be used. It has" 

1621 " never had any effect. Will be removed after v28", 

1622 FutureWarning, 

1623 ) 

1624 try: 

1625 wildcard = CollectionWildcard.from_expression(expression) 

1626 except TypeError as exc: 

1627 raise CollectionExpressionError(f"Invalid collection expression '{expression}'") from exc 

1628 collectionTypes = ensure_iterable(collectionTypes) 

1629 return [ 

1630 record.name 

1631 for record in self._managers.collections.resolve_wildcard( 

1632 wildcard, 

1633 collection_types=frozenset(collectionTypes), 

1634 flatten_chains=flattenChains, 

1635 include_chains=includeChains, 

1636 ) 

1637 ] 

1638 

1639 @contextlib.contextmanager 

1640 def _query(self) -> Iterator[Query]: 

1641 """Context manager returning a `Query` object used for construction 

1642 and execution of complex queries. 

1643 """ 

1644 with self._query_driver(self.defaults.collections, self.defaults.dataId) as driver: 

1645 yield Query(driver) 

1646 

1647 @contextlib.contextmanager 

1648 def _query_driver( 

1649 self, 

1650 default_collections: Iterable[str], 

1651 default_data_id: DataCoordinate, 

1652 ) -> Iterator[DirectQueryDriver]: 

1653 """Set up a `QueryDriver` instance for query execution.""" 

1654 # Query internals do repeated lookups of the same collections, so it 

1655 # benefits from the collection record cache. 

1656 with self._managers.caching_context.enable_collection_record_cache(): 

1657 driver = DirectQueryDriver( 

1658 self._db, 

1659 self.dimensions, 

1660 self._managers, 

1661 self.dimension_record_cache, 

1662 default_collections=default_collections, 

1663 default_data_id=default_data_id, 

1664 ) 

1665 with driver: 

1666 yield driver 

1667 

1668 def get_datastore_records(self, ref: DatasetRef) -> DatasetRef: 

1669 """Retrieve datastore records for given ref. 

1670 

1671 Parameters 

1672 ---------- 

1673 ref : `DatasetRef` 

1674 Dataset reference for which to retrieve its corresponding datastore 

1675 records. 

1676 

1677 Returns 

1678 ------- 

1679 updated_ref : `DatasetRef` 

1680 Dataset reference with filled datastore records. 

1681 

1682 Notes 

1683 ----- 

1684 If this method is called with the dataset ref that is not known to the 

1685 registry then the reference with an empty set of records is returned. 

1686 """ 

1687 datastore_records: dict[str, list[StoredDatastoreItemInfo]] = {} 

1688 for opaque, record_class in self._datastore_record_classes.items(): 

1689 records = self.fetchOpaqueData(opaque, dataset_id=ref.id) 

1690 datastore_records[opaque] = [record_class.from_record(record) for record in records] 

1691 return ref.replace(datastore_records=datastore_records) 

1692 

1693 def store_datastore_records(self, refs: Mapping[str, DatasetRef]) -> None: 

1694 """Store datastore records for given refs. 

1695 

1696 Parameters 

1697 ---------- 

1698 refs : `~collections.abc.Mapping` [`str`, `DatasetRef`] 

1699 Mapping of a datastore name to dataset reference stored in that 

1700 datastore, reference must include datastore records. 

1701 """ 

1702 for datastore_name, ref in refs.items(): 

1703 # Store ref IDs in the bridge table. 

1704 bridge = self._managers.datastores.register(datastore_name) 

1705 bridge.insert([ref]) 

1706 

1707 # store records in opaque tables 

1708 assert ref._datastore_records is not None, "Dataset ref must have datastore records" 

1709 for table_name, records in ref._datastore_records.items(): 

1710 opaque_table = self._managers.opaque.get(table_name) 

1711 assert opaque_table is not None, f"Unexpected opaque table name {table_name}" 

1712 opaque_table.insert(*(record.to_record(dataset_id=ref.id) for record in records)) 

1713 

1714 def make_datastore_tables(self, tables: Mapping[str, DatastoreOpaqueTable]) -> None: 

1715 """Create opaque tables used by datastores. 

1716 

1717 Parameters 

1718 ---------- 

1719 tables : `~collections.abc.Mapping` 

1720 Maps opaque table name to its definition. 

1721 

1722 Notes 

1723 ----- 

1724 This method should disappear in the future when opaque table 

1725 definitions will be provided during `Registry` construction. 

1726 """ 

1727 datastore_record_classes = {} 

1728 for table_name, table_def in tables.items(): 

1729 datastore_record_classes[table_name] = table_def.record_class 

1730 try: 

1731 self._managers.opaque.register(table_name, table_def.table_spec) 

1732 except ReadOnlyDatabaseError: 

1733 # If the database is read only and we just tried and failed to 

1734 # create a table, it means someone is trying to create a 

1735 # read-only butler client for an empty repo. That should be 

1736 # okay, as long as they then try to get any datasets before 

1737 # some other client creates the table. Chances are they're 

1738 # just validating configuration. 

1739 pass 

1740 self._datastore_record_classes = datastore_record_classes 

1741 

1742 def preload_cache(self, *, load_dimension_record_cache: bool) -> None: 

1743 """Immediately load caches that are used for common operations. 

1744 

1745 Parameters 

1746 ---------- 

1747 load_dimension_record_cache : `bool` 

1748 If True, preload the dimension record cache. When this cache is 

1749 preloaded, subsequent external changes to governor dimension 

1750 records will not be visible to this Butler. 

1751 """ 

1752 self._managers.datasets.preload_cache() 

1753 

1754 if load_dimension_record_cache: 1754 ↛ exitline 1754 didn't return from function 'preload_cache' because the condition on line 1754 was always true

1755 self.dimension_record_cache.preload_cache() 

1756 

1757 @property 

1758 def obsCoreTableManager(self) -> ObsCoreTableManager | None: 

1759 """The ObsCore manager instance for this registry 

1760 (`~.interfaces.ObsCoreTableManager` 

1761 or `None`). 

1762 

1763 ObsCore manager may not be implemented for all registry backend, or 

1764 may not be enabled for many repositories. 

1765 """ 

1766 return self._managers.obscore 

1767 

1768 storageClasses: StorageClassFactory 

1769 """All storage classes known to the registry (`StorageClassFactory`). 

1770 """ 

1771 

1772 _defaults: RegistryDefaults 

1773 """Default collections used for registry queries (`RegistryDefaults`)."""