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-14 02:08 -0700
« prev ^ index » next coverage.py v7.15.4, created at 2026-09-14 02:08 -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/>.
28from __future__ import annotations
30from .. import ddl
32__all__ = ("SqlRegistry",)
34import contextlib
35import logging
36import warnings
37from collections.abc import Iterable, Iterator, Mapping, Sequence
38from typing import TYPE_CHECKING, Any
40import sqlalchemy
42from lsst.resources import ResourcePathExpression
43from lsst.utils.iteration import ensure_iterable
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
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 )
93_LOG = logging.getLogger(__name__)
96class SqlRegistry:
97 """Butler Registry implementation that uses SQL database as backend.
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 """
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 """
115 @classmethod
116 def forceRegistryConfig(
117 cls, config: ButlerConfig | RegistryConfig | Config | str | None
118 ) -> RegistryConfig:
119 """Force the supplied config to a `RegistryConfig`.
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.
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
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.
148 This method initializes database contents, database must be empty
149 prior to calling this method.
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.
162 Returns
163 -------
164 registry : `SqlRegistry`
165 A new `SqlRegistry` instance.
166 """
167 config = cls.forceRegistryConfig(config)
168 config.replaceRoot(butlerRoot)
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)}")
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 )
186 try:
187 managers = managerTypes.makeRepo(database, dimensionConfig)
188 return cls(database, RegistryDefaults(), managers)
189 except Exception:
190 database.dispose()
191 raise
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`.
203 Registry database must be initialized prior to calling this method.
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.
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)
239 return cls(database, defaults, managers)
240 except Exception:
241 database.dispose()
242 raise
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
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()
281 def __str__(self) -> str:
282 return str(self._db)
284 def __repr__(self) -> str:
285 return f"SqlRegistry({self._db!r}, {self.dimensions!r})"
287 def isWriteable(self) -> bool:
288 """Return `True` if this registry allows write operations, and `False`
289 otherwise.
290 """
291 return self._db.isWriteable()
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.
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).
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
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
328 @property
329 def defaults(self) -> RegistryDefaults:
330 """Default collection search path and/or output `~CollectionType.RUN`
331 collection (`~lsst.daf.butler.registry.RegistryDefaults`).
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
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
346 def refresh(self) -> None:
347 """Refresh all in-memory state by querying the database.
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()
356 def refresh_collection_summaries(self) -> None:
357 """Refresh content of the collection summary tables in the database.
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)
365 def caching_context(self) -> contextlib.AbstractContextManager[None]:
366 """Return context manager that enables caching.
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()
376 @contextlib.contextmanager
377 def transaction(self, *, savepoint: bool = False) -> Iterator[None]:
378 """Return a context manager that represents a transaction.
380 Parameters
381 ----------
382 savepoint : `bool`
383 Whether to issue a SAVEPOINT in the database.
385 Yields
386 ------
387 `None`
388 """
389 with self._db.transaction(savepoint=savepoint):
390 yield
392 def resetConnectionPool(self) -> None:
393 """Reset SQLAlchemy connection pool for `SqlRegistry` database.
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()
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.
408 Opaque table records can be added via `insertOpaqueData`, retrieved via
409 `fetchOpaqueData`, and removed via `deleteOpaqueData`.
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)
421 @transactional
422 def insertOpaqueData(self, tableName: str, *data: dict) -> None:
423 """Insert records into an opaque table.
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)
436 def fetchOpaqueData(self, tableName: str, **where: Any) -> Iterator[Mapping[str, Any]]:
437 """Retrieve records from an opaque table.
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.
450 Yields
451 ------
452 row : `dict`
453 A dictionary representing a single result row.
454 """
455 yield from self._managers.opaque[tableName].fetch(**where)
457 @transactional
458 def deleteOpaqueData(self, tableName: str, **where: Any) -> None:
459 """Remove records from an opaque table.
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)
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.
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.
488 Returns
489 -------
490 registered : `bool`
491 Boolean indicating whether the collection was already registered
492 or was created by this call.
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
502 def getCollectionType(self, name: str) -> CollectionType:
503 """Return an enumeration value indicating the type of the given
504 collection.
506 Parameters
507 ----------
508 name : `str`
509 The name of the collection.
511 Returns
512 -------
513 type : `CollectionType`
514 Enum value indicating the type of this collection.
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
523 def get_collection_record(self, name: str) -> CollectionRecord:
524 """Return the record for this collection.
526 Parameters
527 ----------
528 name : `str`
529 Name of the collection for which the record is to be retrieved.
531 Returns
532 -------
533 record : `CollectionRecord`
534 The record for this collection.
535 """
536 return self._managers.collections.find(name)
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.
541 Parameters
542 ----------
543 name : `str`
544 The name of the run to create.
545 doc : `str`, optional
546 Documentation string for the collection.
548 Returns
549 -------
550 registered : `bool`
551 Boolean indicating whether a new run was registered. `False`
552 if it already existed.
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
562 @transactional
563 def removeCollection(self, name: str) -> None:
564 """Remove the given collection from the registry.
566 Parameters
567 ----------
568 name : `str`
569 The name of the collection to remove.
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.
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.
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)
595 def getCollectionChain(self, parent: str) -> tuple[str, ...]:
596 """Return the child collections in a `~CollectionType.CHAINED`
597 collection.
599 Parameters
600 ----------
601 parent : `str`
602 Name of the chained collection. Must have already been added via
603 a call to `Registry.registerCollection`.
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.
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
625 @transactional
626 def setCollectionChain(self, parent: str, children: Any, *, flatten: bool = False) -> None:
627 """Define or redefine a `~CollectionType.CHAINED` collection.
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.
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.
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)
664 self._managers.collections.update_chain(parent, list(children), allow_use_in_caching_context=True)
666 def getCollectionParentChains(self, collection: str) -> set[str]:
667 """Return the CHAINED collections that directly contain the given one.
669 Parameters
670 ----------
671 collection : `str`
672 Name of the collection.
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)
681 def getCollectionDocumentation(self, collection: str) -> str | None:
682 """Retrieve the documentation string for a collection.
684 Parameters
685 ----------
686 collection : `str`
687 Name of the collection.
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)
696 def setCollectionDocumentation(self, collection: str, doc: str | None) -> None:
697 """Set the documentation string for a collection.
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)
710 def getCollectionSummary(self, collection: str) -> CollectionSummary:
711 """Return a summary for the given collection.
713 Parameters
714 ----------
715 collection : `str`
716 Name of the collection for which a summary is to be retrieved.
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)
727 def registerDatasetType(self, datasetType: DatasetType) -> bool:
728 """Add a new `DatasetType` to the Registry.
730 It is not an error to register the same `DatasetType` twice.
732 Parameters
733 ----------
734 datasetType : `DatasetType`
735 The `DatasetType` to be added.
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.
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.
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)
760 def removeDatasetType(self, name: str | tuple[str, ...]) -> None:
761 """Remove the named `DatasetType` from the registry.
763 .. warning::
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.
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.
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.
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)
802 def getDatasetType(self, name: str) -> DatasetType:
803 """Get the `DatasetType`.
805 Parameters
806 ----------
807 name : `str`
808 Name of the type.
810 Returns
811 -------
812 type : `DatasetType`
813 The `DatasetType` associated with the given name.
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.
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)
835 def supportsIdGenerationMode(self, mode: DatasetIdGenEnum) -> bool:
836 """Test whether the given dataset ID generation mode is supported by
837 `insertDatasets`.
839 Parameters
840 ----------
841 mode : `DatasetIdGenEnum`
842 Enum value for the mode to test.
844 Returns
845 -------
846 supported : `bool`
847 Whether the given mode is supported.
848 """
849 return True
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`.
862 This always adds new datasets; to associate existing datasets with
863 a new collection, use ``associate``.
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.
883 Returns
884 -------
885 refs : `list` of `DatasetRef`
886 Resolved `DatasetRef` instances for all given data IDs (in the same
887 order).
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)
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")
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
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`.
951 This differs from `insertDatasets` method in that this method accepts
952 `DatasetRef` instances, which already have a dataset ID.
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.
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.
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.
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 []
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()
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)
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")
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
1036 def getDataset(self, id: DatasetId) -> DatasetRef | None:
1037 """Retrieve a Dataset entry.
1039 Parameters
1040 ----------
1041 id : `DatasetId`
1042 The unique identifier for the dataset.
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]
1056 def _fetch_run_dataset_ids(self, run: str) -> list[DatasetId]:
1057 """Return the IDs of all datasets in the given ``RUN``
1058 collection.
1060 Parameters
1061 ----------
1062 run : `str`
1063 Name of the collection.
1065 Returns
1066 -------
1067 dataset_ids : `list` [`uuid.UUID`]
1068 List of dataset IDs.
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)
1079 @transactional
1080 def removeDatasets(self, refs: Iterable[DatasetRef]) -> None:
1081 """Remove datasets from the Registry.
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.
1088 Parameters
1089 ----------
1090 refs : `~collections.abc.Iterable` [`DatasetRef`]
1091 References to the datasets to be removed. Should be considered
1092 invalidated upon return.
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
1106 @transactional
1107 def associate(self, collection: str, refs: Iterable[DatasetRef]) -> None:
1108 """Add existing datasets to a `~CollectionType.TAGGED` collection.
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.
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`.
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
1153 @transactional
1154 def disassociate(self, collection: str, refs: Iterable[DatasetRef]) -> None:
1155 """Remove existing datasets from a `~CollectionType.TAGGED` collection.
1157 ``collection`` and ``ref`` combinations that are not currently
1158 associated are silently ignored.
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`.
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)
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.
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.
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 )
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.
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.
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 )
1276 def getDatastoreBridgeManager(self) -> DatastoreRegistryBridgeManager:
1277 """Return an object that allows a new `Datastore` instance to
1278 communicate with this `Registry`.
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
1288 def getDatasetLocations(self, ref: DatasetRef) -> Iterable[str]:
1289 """Retrieve datastore locations for a given dataset.
1291 Parameters
1292 ----------
1293 ref : `DatasetRef`
1294 A reference to the dataset for which to retrieve storage
1295 information.
1297 Returns
1298 -------
1299 datastores : `~collections.abc.Iterable` [ `str` ]
1300 All the matching datastores holding this dataset.
1302 Raises
1303 ------
1304 lsst.daf.butler.AmbiguousDatasetError
1305 Raised if ``ref.id`` is `None`.
1306 """
1307 return self._managers.datastores.findDatastores(ref)
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.
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.
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`.
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.
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)
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)
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)]
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.
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 )
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.
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.
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``).
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)
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.
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.
1549 Returns
1550 -------
1551 dataset_types : `~collections.abc.Iterable` [ `DatasetType`]
1552 An `~collections.abc.Iterable` of `DatasetType` instances whose
1553 names match ``expression``.
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)
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.
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.
1596 Returns
1597 -------
1598 collections : `~collections.abc.Sequence` [ `str` ]
1599 The names of collections that match ``expression``.
1601 Raises
1602 ------
1603 lsst.daf.butler.registry.CollectionExpressionError
1604 Raised when ``expression`` is invalid.
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 ]
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)
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
1668 def get_datastore_records(self, ref: DatasetRef) -> DatasetRef:
1669 """Retrieve datastore records for given ref.
1671 Parameters
1672 ----------
1673 ref : `DatasetRef`
1674 Dataset reference for which to retrieve its corresponding datastore
1675 records.
1677 Returns
1678 -------
1679 updated_ref : `DatasetRef`
1680 Dataset reference with filled datastore records.
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)
1693 def store_datastore_records(self, refs: Mapping[str, DatasetRef]) -> None:
1694 """Store datastore records for given refs.
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])
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))
1714 def make_datastore_tables(self, tables: Mapping[str, DatastoreOpaqueTable]) -> None:
1715 """Create opaque tables used by datastores.
1717 Parameters
1718 ----------
1719 tables : `~collections.abc.Mapping`
1720 Maps opaque table name to its definition.
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
1742 def preload_cache(self, *, load_dimension_record_cache: bool) -> None:
1743 """Immediately load caches that are used for common operations.
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()
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()
1757 @property
1758 def obsCoreTableManager(self) -> ObsCoreTableManager | None:
1759 """The ObsCore manager instance for this registry
1760 (`~.interfaces.ObsCoreTableManager`
1761 or `None`).
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
1768 storageClasses: StorageClassFactory
1769 """All storage classes known to the registry (`StorageClassFactory`).
1770 """
1772 _defaults: RegistryDefaults
1773 """Default collections used for registry queries (`RegistryDefaults`)."""