Coverage for python/lsst/daf/butler/datastores/fileDatastore.py: 85%
1076 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-16 09:20 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-16 09:20 +0000
1# This file is part of daf_butler.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (http://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This software is dual licensed under the GNU General Public License and also
10# under a 3-clause BSD license. Recipients may choose which of these licenses
11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt,
12# respectively. If you choose the GPL option then the following text applies
13# (but note that there is still no warranty even if you opt for BSD instead):
14#
15# This program is free software: you can redistribute it and/or modify
16# it under the terms of the GNU General Public License as published by
17# the Free Software Foundation, either version 3 of the License, or
18# (at your option) any later version.
19#
20# This program is distributed in the hope that it will be useful,
21# but WITHOUT ANY WARRANTY; without even the implied warranty of
22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
23# GNU General Public License for more details.
24#
25# You should have received a copy of the GNU General Public License
26# along with this program. If not, see <http://www.gnu.org/licenses/>.
28"""Generic file-based datastore code."""
30from __future__ import annotations
32__all__ = ("FileDatastore",)
34import contextlib
35import hashlib
36import logging
37import math
38from collections import defaultdict
39from collections.abc import Callable, Collection, Iterable, Iterator, Mapping, Sequence
40from typing import TYPE_CHECKING, Any, ClassVar, cast
42from sqlalchemy import BigInteger, String
44from lsst.daf.butler import (
45 Config,
46 DatasetDatastoreRecords,
47 DatasetId,
48 DatasetRef,
49 DatasetType,
50 DatasetTypeNotSupportedError,
51 FileDataset,
52 FileDescriptor,
53 Formatter,
54 FormatterFactory,
55 FormatterV1inV2,
56 FormatterV2,
57 Location,
58 LocationFactory,
59 Progress,
60 StorageClass,
61 ddl,
62)
63from lsst.daf.butler.datastore import (
64 DatasetRefURIs,
65 Datastore,
66 DatastoreConfig,
67 DatastoreOpaqueTable,
68 DatastoreValidationError,
69)
70from lsst.daf.butler.datastore.cache_manager import (
71 AbstractDatastoreCacheManager,
72 DatastoreCacheManager,
73 DatastoreDisabledCacheManager,
74)
75from lsst.daf.butler.datastore.composites import CompositesMap
76from lsst.daf.butler.datastore.file_templates import FileTemplates, FileTemplateValidationError
77from lsst.daf.butler.datastore.generic_base import GenericBaseDatastore
78from lsst.daf.butler.datastore.record_data import DatastoreRecordData, DatastoreRecordTable
79from lsst.daf.butler.datastore.stored_file_info import (
80 StoredDatastoreItemInfo,
81 StoredFileInfo,
82 StoredFileInfoTable,
83)
84from lsst.daf.butler.datastores.file_datastore.get import (
85 DatasetLocationInformation,
86 DatastoreFileGetInformation,
87 generate_datastore_get_information,
88 get_dataset_as_python_object_from_get_info,
89)
90from lsst.daf.butler.datastores.file_datastore.retrieve_artifacts import (
91 ArtifactIndexInfo,
92 ZipIndex,
93 determine_destination_for_retrieved_artifact,
94 unpack_zips,
95)
96from lsst.daf.butler.registry.interfaces import (
97 DatabaseInsertMode,
98 DatastoreRegistryBridge,
99 FakeDatasetRef,
100 ReadOnlyDatabaseError,
101)
102from lsst.daf.butler.repo_relocation import replaceRoot
103from lsst.daf.butler.utils import transactional
104from lsst.resources import ResourcePath, ResourcePathExpression
105from lsst.utils.introspection import get_class_of, get_full_type_name
106from lsst.utils.iteration import chunk_iterable
108# For VERBOSE logging usage.
109from lsst.utils.logging import VERBOSE, getLogger
110from lsst.utils.timer import time_this
112from ..datastore import FileTransferMap, FileTransferRecord
114if TYPE_CHECKING:
115 from lsst.daf.butler import DatasetProvenance, LookupKey
116 from lsst.daf.butler.registry.interfaces import DatasetIdRef, DatastoreRegistryBridgeManager
118log = getLogger(__name__)
121class _IngestPrepData(Datastore.IngestPrepData):
122 """Helper class for FileDatastore ingest implementation.
124 Parameters
125 ----------
126 datasets : `~collections.abc.Iterable` of `FileDataset`
127 Files to be ingested by this datastore.
128 """
130 def __init__(self, datasets: Iterable[FileDataset]):
131 super().__init__(ref for dataset in datasets for ref in dataset.refs)
132 self.datasets = datasets
135class FileDatastore(GenericBaseDatastore[StoredFileInfo]):
136 """Generic Datastore for file-based implementations.
138 Should always be sub-classed since key abstract methods are missing.
140 Parameters
141 ----------
142 config : `DatastoreConfig` or `str`
143 Configuration as either a `Config` object or URI to file.
144 bridgeManager : `DatastoreRegistryBridgeManager`
145 Object that manages the interface between `Registry` and datastores.
146 root : `lsst.resources.ResourcePath`
147 Root directory URI of this `Datastore`.
148 formatterFactory : `FormatterFactory`
149 Factory for creating instances of formatters.
150 templates : `FileTemplates`
151 File templates that can be used by this `Datastore`.
152 composites : `CompositesMap`
153 Determines whether a dataset should be disassembled on put.
154 trustGetRequest : `bool`
155 Determine whether we can fall back to configuration if a requested
156 dataset is not known to registry.
158 Raises
159 ------
160 ValueError
161 If root location does not exist and ``create`` is `False` in the
162 configuration.
163 """
165 defaultConfigFile: ClassVar[str | None] = None
166 """Path to configuration defaults. Accessed within the ``config`` resource
167 or relative to a search path. Can be None if no defaults specified.
168 """
170 root: ResourcePath
171 """Root directory URI of this `Datastore`."""
173 locationFactory: LocationFactory
174 """Factory for creating locations relative to the datastore root."""
176 formatterFactory: FormatterFactory
177 """Factory for creating instances of formatters."""
179 templates: FileTemplates
180 """File templates that can be used by this `Datastore`."""
182 composites: CompositesMap
183 """Determines whether a dataset should be disassembled on put."""
185 defaultConfigFile = "datastores/fileDatastore.yaml"
186 """Path to configuration defaults. Accessed within the ``config`` resource
187 or relative to a search path. Can be None if no defaults specified.
188 """
190 _retrieve_dataset_method: Callable[[str], DatasetType | None] | None = None
191 """Callable that is used in trusted mode to retrieve registry definition
192 of a named dataset type.
193 """
195 @classmethod
196 def setConfigRoot(cls, root: str, config: Config, full: Config, overwrite: bool = True) -> None:
197 """Set any filesystem-dependent config options for this Datastore to
198 be appropriate for a new empty repository with the given root.
200 Parameters
201 ----------
202 root : `str`
203 URI to the root of the data repository.
204 config : `Config`
205 A `Config` to update. Only the subset understood by
206 this component will be updated. Will not expand
207 defaults.
208 full : `Config`
209 A complete config with all defaults expanded that can be
210 converted to a `DatastoreConfig`. Read-only and will not be
211 modified by this method.
212 Repository-specific options that should not be obtained
213 from defaults when Butler instances are constructed
214 should be copied from ``full`` to ``config``.
215 overwrite : `bool`, optional
216 If `False`, do not modify a value in ``config`` if the value
217 already exists. Default is always to overwrite with the provided
218 ``root``.
220 Notes
221 -----
222 If a keyword is explicitly defined in the supplied ``config`` it
223 will not be overridden by this method if ``overwrite`` is `False`.
224 This allows explicit values set in external configs to be retained.
225 """
226 Config.updateParameters(
227 DatastoreConfig,
228 config,
229 full,
230 toUpdate={"root": root},
231 toCopy=("cls", ("records", "table")),
232 overwrite=overwrite,
233 )
235 @classmethod
236 def makeTableSpec(cls) -> ddl.TableSpec:
237 return ddl.TableSpec(
238 fields=[
239 ddl.FieldSpec(name="dataset_id", dtype=ddl.GUID, primaryKey=True),
240 ddl.FieldSpec(name="path", dtype=String, length=256, nullable=False),
241 ddl.FieldSpec(name="formatter", dtype=String, length=128, nullable=False),
242 ddl.FieldSpec(name="storage_class", dtype=String, length=64, nullable=False),
243 # Use empty string to indicate no component
244 ddl.FieldSpec(name="component", dtype=String, length=32, primaryKey=True),
245 # TODO: should checksum be Base64Bytes instead?
246 ddl.FieldSpec(name="checksum", dtype=String, length=128, nullable=True),
247 ddl.FieldSpec(name="file_size", dtype=BigInteger, nullable=True),
248 ],
249 unique=frozenset(),
250 indexes=[ddl.IndexSpec("path")],
251 )
253 def __init__(
254 self,
255 config: DatastoreConfig,
256 bridgeManager: DatastoreRegistryBridgeManager,
257 root: ResourcePath,
258 formatterFactory: FormatterFactory,
259 templates: FileTemplates,
260 composites: CompositesMap,
261 trustGetRequest: bool,
262 ):
263 super().__init__(config, bridgeManager)
264 self.root = ResourcePath(root)
265 self.formatterFactory = formatterFactory
266 self.templates = templates
267 self.composites = composites
268 self.trustGetRequest = trustGetRequest
270 # Name ourselves either using an explicit name or a name
271 # derived from the (unexpanded) root
272 if "name" in self.config:
273 self.name = self.config["name"]
274 else:
275 # We use the unexpanded root in the name to indicate that this
276 # datastore can be moved without having to update registry.
277 self.name = "{}@{}".format(type(self).__name__, self.config["root"])
279 self.locationFactory = LocationFactory(self.root)
281 self._opaque_table_name = self.config["records", "table"]
282 try:
283 # Storage of paths and formatters, keyed by dataset_id
284 self._table = bridgeManager.opaque.register(self._opaque_table_name, self.makeTableSpec())
285 # Interface to Registry.
286 self._bridge = bridgeManager.register(self.name)
287 except ReadOnlyDatabaseError:
288 # If the database is read only and we just tried and failed to
289 # create a table, it means someone is trying to create a read-only
290 # butler client for an empty repo. That should be okay, as long
291 # as they then try to get any datasets before some other client
292 # creates the table. Chances are they're just validating
293 # configuration.
294 pass
296 # Determine whether checksums should be used - default to False
297 self.useChecksum = self.config.get("checksum", False)
299 # Create a cache manager
300 self.cacheManager: AbstractDatastoreCacheManager
301 if "cached" in self.config: 301 ↛ 304line 301 didn't jump to line 304 because the condition on line 301 was always true
302 self.cacheManager = DatastoreCacheManager(self.config["cached"], universe=bridgeManager.universe)
303 else:
304 self.cacheManager = DatastoreDisabledCacheManager("", universe=bridgeManager.universe)
306 self.universe = bridgeManager.universe
308 @classmethod
309 def _create_from_config(
310 cls,
311 config: DatastoreConfig,
312 bridgeManager: DatastoreRegistryBridgeManager,
313 butlerRoot: ResourcePathExpression | None,
314 ) -> FileDatastore:
315 if "root" not in config: 315 ↛ 316line 315 didn't jump to line 316 because the condition on line 315 was never true
316 raise ValueError("No root directory specified in configuration")
318 # Support repository relocation in config
319 # Existence of self.root is checked in subclass
320 root = ResourcePath(replaceRoot(config["root"], butlerRoot), forceDirectory=True, forceAbsolute=True)
322 # Now associate formatters with storage classes
323 formatterFactory = FormatterFactory()
324 formatterFactory.registerFormatters(config["formatters"], universe=bridgeManager.universe)
326 # Read the file naming templates
327 templates = FileTemplates(config["templates"], universe=bridgeManager.universe)
329 # See if composites should be disassembled
330 composites = CompositesMap(config["composites"], universe=bridgeManager.universe)
332 # Determine whether we can fall back to configuration if a
333 # requested dataset is not known to registry
334 trustGetRequest = config.get("trust_get_request", False)
336 self = FileDatastore(
337 config, bridgeManager, root, formatterFactory, templates, composites, trustGetRequest
338 )
340 # Check existence and create directory structure if necessary.
341 #
342 # The concept of a 'root directory' is problematic for some resource
343 # path types that don't necessarily support the concept of a directory
344 # (http, s3, gs... basically anything that isn't a local filesystem or
345 # WebDAV.)
346 # On these resource paths an object representing the
347 # "root" directory may not exist even though files under the root do,
348 # and in a read-only repository we will be unable to create it.
349 # So we only immediately verify the root for local filesystems,
350 # the only case where this check will definitely not give a false
351 # negative.
352 if self.root.isLocal and not self.root.exists():
353 if "create" not in self.config or not self.config["create"]: 353 ↛ 354line 353 didn't jump to line 354 because the condition on line 353 was never true
354 raise ValueError(f"No valid root and not allowed to create one at: {self.root}")
355 try:
356 self.root.mkdir()
357 except Exception as e:
358 raise ValueError(
359 f"Can not create datastore root '{self.root}', check permissions. Got error: {e}"
360 ) from e
362 return self
364 def clone(self, bridgeManager: DatastoreRegistryBridgeManager) -> Datastore:
365 return FileDatastore(
366 self.config,
367 bridgeManager,
368 self.root,
369 self.formatterFactory,
370 self.templates,
371 self.composites,
372 self.trustGetRequest,
373 )
375 def __str__(self) -> str:
376 return str(self.root)
378 @property
379 def bridge(self) -> DatastoreRegistryBridge:
380 return self._bridge
382 @property
383 def roots(self) -> dict[str, ResourcePath | None]:
384 # Docstring inherited.
385 return {self.name: self.root}
387 def _set_trust_mode(self, mode: bool) -> None:
388 self.trustGetRequest = mode
390 def _artifact_exists(self, location: Location) -> bool:
391 """Check that an artifact exists in this datastore at the specified
392 location.
394 Parameters
395 ----------
396 location : `Location`
397 Expected location of the artifact associated with this datastore.
399 Returns
400 -------
401 exists : `bool`
402 True if the location can be found, false otherwise.
403 """
404 log.debug("Checking if resource exists: %s", location.uri)
405 return location.uri.exists()
407 def addStoredItemInfo(
408 self,
409 refs: Iterable[DatasetRef],
410 infos: Iterable[StoredFileInfo],
411 insert_mode: DatabaseInsertMode = DatabaseInsertMode.INSERT,
412 ) -> None:
413 """Record internal storage information associated with one or more
414 datasets.
416 Parameters
417 ----------
418 refs : sequence of `DatasetRef`
419 The datasets that have been stored.
420 infos : sequence of `StoredDatastoreItemInfo`
421 Metadata associated with the stored datasets.
422 insert_mode : `~lsst.daf.butler.registry.interfaces.DatabaseInsertMode`
423 Mode to use to insert the new records into the table. The
424 options are ``INSERT`` (error if pre-existing), ``REPLACE``
425 (replace content with new values), and ``ENSURE`` (skip if the row
426 already exists).
427 """
428 records = [
429 info.rebase(ref).to_record(dataset_id=ref.id) for ref, info in zip(refs, infos, strict=True)
430 ]
431 match insert_mode:
432 case DatabaseInsertMode.INSERT:
433 self._table.insert(*records, transaction=self._transaction)
434 case DatabaseInsertMode.ENSURE: 434 ↛ 435line 434 didn't jump to line 435 because the pattern on line 434 never matched
435 self._table.ensure(*records, transaction=self._transaction)
436 case DatabaseInsertMode.REPLACE: 436 ↛ 438line 436 didn't jump to line 438 because the pattern on line 436 always matched
437 self._table.replace(*records, transaction=self._transaction)
438 case _:
439 raise ValueError(f"Unknown insert mode of '{insert_mode}'")
441 def getStoredItemsInfo(
442 self, ref: DatasetIdRef, ignore_datastore_records: bool = False
443 ) -> list[StoredFileInfo]:
444 """Retrieve information associated with files stored in this
445 `Datastore` associated with this dataset ref.
447 Parameters
448 ----------
449 ref : `DatasetRef`
450 The dataset that is to be queried.
451 ignore_datastore_records : `bool`
452 If `True` then do not use datastore records stored in refs.
454 Returns
455 -------
456 items : `~collections.abc.Iterable` [`StoredDatastoreItemInfo`]
457 Stored information about the files and associated formatters
458 associated with this dataset. Only one file will be returned
459 if the dataset has not been disassembled. Can return an empty
460 list if no matching datasets can be found.
461 """
462 # Try to get them from the ref first.
463 if ref._datastore_records is not None and not ignore_datastore_records:
464 ref_records = ref._datastore_records.get(self._table.name, [])
465 # Need to make sure they have correct type.
466 for record in ref_records:
467 if not isinstance(record, StoredFileInfo): 467 ↛ 468line 467 didn't jump to line 468 because the condition on line 467 was never true
468 raise TypeError(f"Datastore record has unexpected type {record.__class__.__name__}")
469 return cast(list[StoredFileInfo], ref_records)
471 # Look for the dataset_id -- there might be multiple matches
472 # if we have disassembled the dataset.
473 records = self._table.fetch(dataset_id=ref.id)
474 return [StoredFileInfo.from_record(record) for record in records]
476 def _register_datasets(
477 self,
478 refsAndInfos: Iterable[tuple[DatasetRef, StoredFileInfo]],
479 insert_mode: DatabaseInsertMode = DatabaseInsertMode.INSERT,
480 ) -> None:
481 """Update registry to indicate that one or more datasets have been
482 stored.
484 Parameters
485 ----------
486 refsAndInfos : sequence `tuple` [`DatasetRef`,
487 `StoredDatastoreItemInfo`]
488 Datasets to register and the internal datastore metadata associated
489 with them.
490 insert_mode : `str`, optional
491 Indicate whether the new records should be new ("insert", default),
492 or allowed to exists ("ensure") or be replaced if already present
493 ("replace").
494 """
495 expandedRefs: list[DatasetRef] = []
496 expandedItemInfos: list[StoredFileInfo] = []
498 for ref, itemInfo in refsAndInfos:
499 expandedRefs.append(ref)
500 expandedItemInfos.append(itemInfo)
502 # Dataset location only cares about registry ID so if we have
503 # disassembled in datastore we have to deduplicate. Since they
504 # will have different datasetTypes we can't use a set
505 registryRefs = {r.id: r for r in expandedRefs}
506 if insert_mode == DatabaseInsertMode.INSERT:
507 self.bridge.insert(registryRefs.values())
508 else:
509 # There are only two columns and all that matters is the
510 # dataset ID.
511 self.bridge.ensure(registryRefs.values())
512 self.addStoredItemInfo(expandedRefs, expandedItemInfos, insert_mode=insert_mode)
514 def _get_stored_records_associated_with_refs(
515 self, refs: Iterable[DatasetIdRef], ignore_datastore_records: bool = False
516 ) -> dict[DatasetId, list[StoredFileInfo]]:
517 """Retrieve all records associated with the provided refs.
519 Parameters
520 ----------
521 refs : `~collections.abc.Iterable` of `DatasetIdRef`
522 The refs for which records are to be retrieved.
523 ignore_datastore_records : `bool`
524 If `True` then do not use datastore records stored in refs.
526 Returns
527 -------
528 records : `dict` of [`DatasetId`, `list` of `StoredFileInfo`]
529 The matching records indexed by the ref ID. The number of entries
530 in the dict can be smaller than the number of requested refs.
531 """
532 # Check datastore records in refs first.
533 records_by_ref: defaultdict[DatasetId, list[StoredFileInfo]] = defaultdict(list)
534 refs_with_no_records = []
535 for ref in refs:
536 if ignore_datastore_records or ref._datastore_records is None: 536 ↛ 539line 536 didn't jump to line 539 because the condition on line 536 was always true
537 refs_with_no_records.append(ref)
538 else:
539 if (ref_records := ref._datastore_records.get(self._table.name)) is not None:
540 # Need to make sure they have correct type.
541 for ref_record in ref_records:
542 if not isinstance(ref_record, StoredFileInfo):
543 raise TypeError(
544 f"Datastore record has unexpected type {ref_record.__class__.__name__}"
545 )
546 records_by_ref[ref.id].append(ref_record)
548 # If there were any refs without datastore records, check opaque table.
549 records = self._table.fetch(dataset_id=[ref.id for ref in refs_with_no_records])
551 # Uniqueness is dataset_id + component so can have multiple records
552 # per ref.
553 for record in records:
554 records_by_ref[record["dataset_id"]].append(StoredFileInfo.from_record(record))
555 return records_by_ref
557 def _refs_associated_with_artifacts(
558 self, paths: Iterable[str | ResourcePath]
559 ) -> dict[str, set[DatasetId]]:
560 """Return paths and associated dataset refs.
562 Parameters
563 ----------
564 paths : `list` of `str` or `lsst.resources.ResourcePath`
565 All the paths to include in search. These are exact matches
566 to the entries in the records table and can include fragments.
568 Returns
569 -------
570 mapping : `dict` of [`str`, `set` [`DatasetId`]]
571 Mapping of each path to a set of associated database IDs.
572 These are artifacts and so any fragments are stripped from the
573 keys.
574 """
575 # Group paths by those that have fragments and those that do not.
576 with_fragment = set()
577 without_fragment = set()
578 for rpath in paths:
579 spath = str(rpath) # Typing says can be ResourcePath so must force to string.
580 if "#" in spath:
581 spath, fragment = spath.rsplit("#", 1)
582 with_fragment.add(spath)
583 else:
584 without_fragment.add(spath)
586 result: dict[str, set[DatasetId]] = defaultdict(set)
587 if without_fragment:
588 records = self._table.fetch(path=without_fragment)
589 for row in records:
590 path = row["path"]
591 result[path].add(row["dataset_id"])
592 if with_fragment:
593 # Do a query per prefix.
594 for path in with_fragment:
595 records = self._table.fetch(path=f"{path}#%")
596 for row in records:
597 # Need to strip fragments before adding to dict.
598 row_path = row["path"]
599 artifact_path = row_path[: row_path.rfind("#")]
600 result[artifact_path].add(row["dataset_id"])
601 return result
603 def _registered_refs_per_artifact(self, pathInStore: ResourcePath) -> set[DatasetId]:
604 """Return all dataset refs associated with the supplied path.
606 Parameters
607 ----------
608 pathInStore : `lsst.resources.ResourcePath`
609 Path of interest in the data store.
611 Returns
612 -------
613 ids : `set` of `int`
614 All `DatasetRef` IDs associated with this path.
615 """
616 records = list(self._table.fetch(path=str(pathInStore)))
617 ids = {r["dataset_id"] for r in records}
618 return ids
620 def removeStoredItemInfo(self, ref: DatasetIdRef) -> None:
621 """Remove information about the file associated with this dataset.
623 Parameters
624 ----------
625 ref : `DatasetRef`
626 The dataset that has been removed.
627 """
628 # Note that this method is actually not used by this implementation,
629 # we depend on bridge to delete opaque records. But there are some
630 # tests that check that this method works, so we keep it for now.
631 self._table.delete(["dataset_id"], {"dataset_id": ref.id})
633 def _get_dataset_locations_info(
634 self, ref: DatasetIdRef, ignore_datastore_records: bool = False
635 ) -> list[DatasetLocationInformation]:
636 r"""Find all the `Location`\ s of the requested dataset in the
637 `Datastore` and the associated stored file information.
639 Parameters
640 ----------
641 ref : `DatasetRef`
642 Reference to the required `Dataset`.
643 ignore_datastore_records : `bool`
644 If `True` then do not use datastore records stored in refs.
646 Returns
647 -------
648 results : `list` [`tuple` [`Location`, `StoredFileInfo` ]]
649 Location of the dataset within the datastore and
650 stored information about each file and its formatter.
651 """
652 # Get the file information (this will fail if no file)
653 records = self.getStoredItemsInfo(ref, ignore_datastore_records)
655 # Use the path to determine the location -- we need to take
656 # into account absolute URIs in the datastore record
657 return [(r.file_location(self.locationFactory), r) for r in records]
659 def _can_remove_dataset_artifact(self, ref: DatasetIdRef, location: Location) -> bool:
660 """Check that there is only one dataset associated with the
661 specified artifact.
663 Parameters
664 ----------
665 ref : `DatasetRef` or `FakeDatasetRef`
666 Dataset to be removed.
667 location : `Location`
668 The location of the artifact to be removed.
670 Returns
671 -------
672 can_remove : `Bool`
673 True if the artifact can be safely removed.
674 """
675 # Can't ever delete absolute URIs.
676 if location.pathInStore.isabs():
677 return False
679 # Get all entries associated with this path
680 allRefs = self._registered_refs_per_artifact(location.pathInStore)
681 if not allRefs:
682 raise RuntimeError(f"Datastore inconsistency error. {location.pathInStore} not in registry")
684 # Remove these refs from all the refs and if there is nothing left
685 # then we can delete
686 remainingRefs = allRefs - {ref.id}
688 if remainingRefs:
689 return False
690 return True
692 def _get_expected_dataset_locations_info(self, ref: DatasetRef) -> list[tuple[Location, StoredFileInfo]]:
693 """Predict the location and related file information of the requested
694 dataset in this datastore.
696 Parameters
697 ----------
698 ref : `DatasetRef`
699 Reference to the required `Dataset`.
701 Returns
702 -------
703 results : `list` [`tuple` [`Location`, `StoredFileInfo` ]]
704 Expected Location of the dataset within the datastore and
705 placeholder information about each file and its formatter.
707 Notes
708 -----
709 Uses the current configuration to determine how we would expect the
710 datastore files to have been written if we couldn't ask registry.
711 This is safe so long as there has been no change to datastore
712 configuration between writing the dataset and wanting to read it.
713 Will not work for files that have been ingested without using the
714 standard file template or default formatter.
715 """
716 # If we have a component ref we always need to ask the questions
717 # of the composite. If the composite is disassembled this routine
718 # should return all components. If the composite was not
719 # disassembled the composite is what is stored regardless of
720 # component request. Note that if the caller has disassembled
721 # a composite there is no way for this guess to know that
722 # without trying both the composite and component ref and seeing
723 # if there is something at the component Location even without
724 # disassembly being enabled.
725 if ref.datasetType.isComponent(): 725 ↛ 726line 725 didn't jump to line 726 because the condition on line 725 was never true
726 ref = ref.makeCompositeRef()
728 # See if the ref is a composite that should be disassembled
729 doDisassembly = self.composites.shouldBeDisassembled(ref)
731 all_info: list[tuple[Location, Formatter | FormatterV2, StorageClass, str | None]] = []
733 if doDisassembly:
734 for component, componentStorage in ref.datasetType.storageClass.components.items():
735 compRef = ref.makeComponentRef(component)
736 location, formatter = self._determine_put_formatter_location(compRef)
737 all_info.append((location, formatter, componentStorage, component))
739 else:
740 # Always use the composite ref if no disassembly
741 location, formatter = self._determine_put_formatter_location(ref)
742 all_info.append((location, formatter, ref.datasetType.storageClass, None))
744 # Convert the list of tuples to have StoredFileInfo as second element
745 return [
746 (
747 location,
748 StoredFileInfo(
749 formatter=formatter,
750 path=location.pathInStore.path,
751 storageClass=storageClass,
752 component=component,
753 checksum=None,
754 file_size=-1,
755 ),
756 )
757 for location, formatter, storageClass, component in all_info
758 ]
760 def _prepare_for_direct_get(
761 self, ref: DatasetRef, parameters: Mapping[str, Any] | None = None
762 ) -> list[DatastoreFileGetInformation]:
763 """Check parameters for ``get`` and obtain formatter and
764 location.
766 Parameters
767 ----------
768 ref : `DatasetRef`
769 Reference to the required Dataset.
770 parameters : `dict`
771 `StorageClass`-specific parameters that specify, for example,
772 a slice of the dataset to be loaded.
774 Returns
775 -------
776 getInfo : `list` [`DatastoreFileGetInformation`]
777 Parameters needed to retrieve each file.
778 """
779 log.debug("Retrieve %s from %s with parameters %s", ref, self.name, parameters)
781 # The ref as supplied describes what the caller wants back, including
782 # any component and read storage class override. Internally the
783 # composite as defined in the repository is used: that is what a get
784 # with no overrides would return and it is the ref given to the
785 # Formatter. Using the registry storage class also resets the storage
786 # class for trusted mode.
787 registry_ref = self._cast_storage_class(ref.makeCompositeRef() if ref.isComponent() else ref)
789 # Get file metadata and internal metadata
790 fileLocations = self._get_dataset_locations_info(registry_ref)
791 if not fileLocations:
792 if not self.trustGetRequest:
793 raise FileNotFoundError(f"Could not retrieve dataset {ref}.")
794 # Assume the dataset is where we think it should be
795 fileLocations = self._get_expected_dataset_locations_info(registry_ref)
797 if len(fileLocations) > 1:
798 # If trust is involved it is possible that there will be
799 # components listed here that do not exist in the datastore.
800 # Explicitly check for file artifact existence and filter out any
801 # that are missing.
802 if self.trustGetRequest:
803 fileLocations = [loc for loc in fileLocations if loc[0].uri.exists()]
805 # For now complain only if we have no components at all. One
806 # component is probably a problem but we can punt that to the
807 # assembler.
808 if not fileLocations:
809 raise FileNotFoundError(f"None of the component files for dataset {ref} exist.")
811 return generate_datastore_get_information(
812 fileLocations,
813 registry_ref=registry_ref,
814 read_ref=ref,
815 parameters=parameters,
816 )
818 def _determine_put_formatter_location(
819 self, ref: DatasetRef, provenance: DatasetProvenance | None = None
820 ) -> tuple[Location, Formatter | FormatterV2]:
821 """Calculate the formatter and output location to use for put.
823 Parameters
824 ----------
825 ref : `DatasetRef`
826 Reference to the associated Dataset.
827 provenance : `DatasetProvenance`
828 Any provenance that should be attached to the serialized dataset.
830 Returns
831 -------
832 location : `Location`
833 The location to write the dataset.
834 formatter : `Formatter`
835 The `Formatter` to use to write the dataset.
836 """
837 # Work out output file name
838 try:
839 template = self.templates.getTemplate(ref)
840 except KeyError as e:
841 raise DatasetTypeNotSupportedError(f"Unable to find template for {ref}") from e
843 # Validate the template to protect against filenames from different
844 # dataIds returning the same and causing overwrite confusion.
845 template.validateTemplate(ref)
847 location = self.locationFactory.fromPath(template.format(ref), trusted_path=True)
849 # Get the formatter based on the storage class
850 storageClass = ref.datasetType.storageClass
851 try:
852 formatter = self.formatterFactory.getFormatter(
853 ref,
854 FileDescriptor(location, storageClass=storageClass, component=ref.datasetType.component()),
855 dataId=ref.dataId,
856 ref=ref,
857 provenance=provenance,
858 )
859 except KeyError as e:
860 raise DatasetTypeNotSupportedError(
861 f"Unable to find formatter for {ref} in datastore {self.name}"
862 ) from e
864 # Now that we know the formatter, update the location
865 location = formatter.make_updated_location(location)
867 return location, formatter
869 def _overrideTransferMode(self, *datasets: FileDataset, transfer: str | None = None) -> str | None:
870 # Docstring inherited from base class
871 if transfer != "auto":
872 return transfer
874 # See if the paths are within the datastore or not
875 inside = [self._pathInStore(d.path) is not None for d in datasets]
877 if all(inside):
878 transfer = None
879 elif not any(inside): 879 ↛ 888line 879 didn't jump to line 888 because the condition on line 879 was always true
880 # Allow ResourcePath to use its own knowledge
881 transfer = "auto"
882 else:
883 # This can happen when importing from a datastore that
884 # has had some datasets ingested using "direct" mode.
885 # Also allow ResourcePath to sort it out but warn about it.
886 # This can happen if you are importing from a datastore
887 # that had some direct transfer datasets.
888 log.warning(
889 "Some datasets are inside the datastore and some are outside. Using 'split' "
890 "transfer mode. This assumes that the files outside the datastore are "
891 "still accessible to the new butler since they will not be copied into "
892 "the target datastore."
893 )
894 transfer = "split"
896 return transfer
898 def _pathInStore(self, path: ResourcePathExpression) -> str | None:
899 """Return path relative to datastore root.
901 Parameters
902 ----------
903 path : `lsst.resources.ResourcePathExpression`
904 Path to dataset. Can be absolute URI. If relative assumed to
905 be relative to the datastore. Returns path in datastore
906 or raises an exception if the path it outside.
908 Returns
909 -------
910 inStore : `str`
911 Path relative to datastore root. Returns `None` if the file is
912 outside the root.
913 """
914 # Relative path will always be relative to datastore
915 pathUri = ResourcePath(path, forceAbsolute=False, forceDirectory=False)
916 return pathUri.relative_to(self.root)
918 def _standardizeIngestPath(
919 self,
920 path: str | ResourcePath,
921 *,
922 transfer: str | None = None,
923 check_existence: bool = False,
924 ) -> str | ResourcePath:
925 """Standardize the path of a to-be-ingested file.
927 Parameters
928 ----------
929 path : `str` or `lsst.resources.ResourcePath`
930 Path of a file to be ingested. This parameter is not expected
931 to be all the types that can be used to construct a
932 `~lsst.resources.ResourcePath`.
933 transfer : `str`, optional
934 How (and whether) the dataset should be added to the datastore.
935 See `ingest` for details of transfer modes.
936 This implementation is provided only so
937 `NotImplementedError` can be raised if the mode is not supported;
938 actual transfers are deferred to `_extractIngestInfo`.
939 check_existence : `bool`, optional
940 If `True` the existence of the file will be checked, otherwise
941 no check will be made.
943 Returns
944 -------
945 path : `str` or `lsst.resources.ResourcePath`
946 New path in what the datastore considers standard form. If an
947 absolute URI was given that will be returned unchanged.
949 Notes
950 -----
951 Subclasses of `FileDatastore` can implement this method instead
952 of `_prepIngest`. It should not modify the data repository or given
953 file in any way.
955 Raises
956 ------
957 NotImplementedError
958 Raised if the datastore does not support the given transfer mode
959 (including the case where ingest is not supported at all).
960 """
961 if transfer not in (None, "direct", "split") + self.root.transferModes:
962 raise NotImplementedError(f"Transfer mode {transfer} not supported.")
964 # A relative URI indicates relative to datastore root
965 srcUri = ResourcePath(path, forceAbsolute=False, forceDirectory=False)
966 if not srcUri.isabs():
967 srcUri = self.root.join(path)
969 if check_existence and not srcUri.exists():
970 raise FileNotFoundError(
971 f"Resource at {srcUri} does not exist; note that paths to ingest "
972 f"are assumed to be relative to {self.root} unless they are absolute."
973 )
975 if transfer is None:
976 relpath = srcUri.relative_to(self.root)
977 if not relpath:
978 raise RuntimeError(
979 f"Transfer is none but source file ({srcUri}) is not within datastore ({self.root})"
980 )
982 # Return the relative path within the datastore for internal
983 # transfer
984 path = relpath
986 return path
988 def _extractIngestInfo(
989 self,
990 path: ResourcePathExpression,
991 ref: DatasetRef,
992 *,
993 formatter: Formatter | FormatterV2 | type[Formatter | FormatterV2],
994 transfer: str | None = None,
995 record_validation_info: bool = True,
996 ) -> StoredFileInfo:
997 """Relocate (if necessary) and extract `StoredFileInfo` from a
998 to-be-ingested file.
1000 Parameters
1001 ----------
1002 path : `lsst.resources.ResourcePathExpression`
1003 URI or path of a file to be ingested.
1004 ref : `DatasetRef`
1005 Reference for the dataset being ingested. Guaranteed to have
1006 ``dataset_id not None`.
1007 formatter : `type` or `Formatter`
1008 `Formatter` subclass to use for this dataset or an instance.
1009 transfer : `str`, optional
1010 How (and whether) the dataset should be added to the datastore.
1011 See `ingest` for details of transfer modes.
1012 record_validation_info : `bool`, optional
1013 If `True`, the default, the datastore can record validation
1014 information associated with the file. If `False` the datastore
1015 will not attempt to track any information such as checksums
1016 or file sizes. This can be useful if such information is tracked
1017 in an external system or if the file is to be compressed in place.
1018 It is up to the datastore whether this parameter is relevant.
1020 Returns
1021 -------
1022 info : `StoredFileInfo`
1023 Internal datastore record for this file. This will be inserted by
1024 the caller; the `_extractIngestInfo` is only responsible for
1025 creating and populating the struct.
1027 Raises
1028 ------
1029 FileNotFoundError
1030 Raised if one of the given files does not exist.
1031 FileExistsError
1032 Raised if transfer is not `None` but the (internal) location the
1033 file would be moved to is already occupied.
1034 """
1035 if self._transaction is None: 1035 ↛ 1036line 1035 didn't jump to line 1036 because the condition on line 1035 was never true
1036 raise RuntimeError("Ingest called without transaction enabled")
1038 # Create URI of the source path, do not need to force a relative
1039 # path to absolute.
1040 srcUri = ResourcePath(path, forceAbsolute=False, forceDirectory=False)
1042 # Track whether we have read the size of the source yet
1043 have_sized = False
1045 tgtLocation: Location | None
1046 if transfer is None or transfer == "split":
1047 # A relative path is assumed to be relative to the datastore
1048 # in this context
1049 if not srcUri.isabs():
1050 tgtLocation = self.locationFactory.fromPath(srcUri.ospath, trusted_path=False)
1051 else:
1052 # Work out the path in the datastore from an absolute URI
1053 # This is required to be within the datastore.
1054 pathInStore = srcUri.relative_to(self.root)
1055 if pathInStore is None and transfer is None: 1055 ↛ 1056line 1055 didn't jump to line 1056 because the condition on line 1055 was never true
1056 raise RuntimeError(
1057 f"Unexpectedly learned that {srcUri} is not within datastore {self.root}"
1058 )
1059 if pathInStore: 1059 ↛ 1061line 1059 didn't jump to line 1061 because the condition on line 1059 was always true
1060 tgtLocation = self.locationFactory.fromPath(pathInStore, trusted_path=True)
1061 elif transfer == "split":
1062 # Outside the datastore but treat that as a direct ingest
1063 # instead.
1064 tgtLocation = None
1065 else:
1066 raise RuntimeError(f"Unexpected transfer mode encountered: {transfer} for URI {srcUri}")
1067 elif transfer == "direct":
1068 # Want to store the full URI to the resource directly in
1069 # datastore. This is useful for referring to permanent archive
1070 # storage for raw data.
1071 # Trust that people know what they are doing.
1072 tgtLocation = None
1073 else:
1074 # Work out the name we want this ingested file to have
1075 # inside the datastore
1076 tgtLocation = self._calculate_ingested_datastore_name(srcUri, ref, formatter)
1078 # if we are transferring from a local file to a remote location
1079 # it may be more efficient to get the size and checksum of the
1080 # local file rather than the transferred one
1081 if record_validation_info and srcUri.isLocal:
1082 size = srcUri.size()
1083 checksum = self.computeChecksum(srcUri) if self.useChecksum else None
1084 have_sized = True
1086 # Transfer the resource to the destination.
1087 # Allow overwrite of an existing file. This matches the behavior
1088 # of datastore.put() in that it trusts that registry would not
1089 # be asking to overwrite unless registry thought that the
1090 # overwrite was allowed.
1091 tgtLocation.uri.transfer_from(
1092 srcUri, transfer=transfer, transaction=self._transaction, overwrite=True
1093 )
1095 if tgtLocation is None:
1096 # This means we are using direct mode
1097 targetUri = srcUri
1098 targetPath = str(srcUri)
1099 else:
1100 targetUri = tgtLocation.uri
1101 targetPath = tgtLocation.pathInStore.path
1103 # the file should exist in the datastore now
1104 if record_validation_info:
1105 if not have_sized:
1106 size = targetUri.size()
1107 checksum = self.computeChecksum(targetUri) if self.useChecksum else None
1108 else:
1109 # Not recording any file information.
1110 size = -1
1111 checksum = None
1113 return StoredFileInfo(
1114 formatter=formatter,
1115 path=targetPath,
1116 storageClass=ref.datasetType.storageClass,
1117 component=ref.datasetType.component(),
1118 file_size=size,
1119 checksum=checksum,
1120 )
1122 def _prepIngest(self, *datasets: FileDataset, transfer: str | None = None) -> _IngestPrepData:
1123 # Docstring inherited from Datastore._prepIngest.
1124 filtered = []
1126 # Ingest could be given tens of thousands of files. It is not efficient
1127 # to check for the existence of every single file (especially if they
1128 # are remote URIs) but in some transfer modes the files will be checked
1129 # anyhow when they are relocated. For direct or None transfer modes
1130 # it is possible to not know if the file is accessible at all.
1131 # Therefore limit number of files that will be checked (but always
1132 # include the first one).
1133 max_checks = 200
1134 n_datasets = len(datasets)
1135 if n_datasets <= max_checks: 1135 ↛ 1137line 1135 didn't jump to line 1137 because the condition on line 1135 was always true
1136 check_every_n = 1
1137 elif transfer in ("direct", None):
1138 check_every_n = int(n_datasets / max_checks + 1) # +1 so that if n < max_checks the answer is 1.
1139 else:
1140 check_every_n = 0
1142 for count, dataset in enumerate(datasets):
1143 acceptable = [ref for ref in dataset.refs if self.constraints.isAcceptable(ref)]
1144 if not acceptable:
1145 continue
1146 else:
1147 dataset.refs = acceptable
1148 if dataset.formatter is None:
1149 dataset.formatter = self.formatterFactory.getFormatterClass(dataset.refs[0])
1150 else:
1151 assert isinstance(dataset.formatter, type | str)
1152 formatter_class = get_class_of(dataset.formatter)
1153 if not issubclass(formatter_class, Formatter | FormatterV2): 1153 ↛ 1154line 1153 didn't jump to line 1154 because the condition on line 1153 was never true
1154 raise TypeError(f"Requested formatter {dataset.formatter} is not a Formatter class.")
1155 dataset.formatter = formatter_class
1157 # Decide whether the file should be checked.
1158 check_existence = False
1159 if check_every_n != 0: 1159 ↛ 1164line 1159 didn't jump to line 1164 because the condition on line 1159 was always true
1160 # First time through count is 0 so we guarantee to check
1161 # the first file but not necessarily the final one.
1162 check_existence = count % check_every_n == 0
1164 if check_existence: 1164 ↛ 1173line 1164 didn't jump to line 1173 because the condition on line 1164 was always true
1165 log.debug(
1166 "Checking file existence: %s (%d/%d) [%s]",
1167 check_existence,
1168 count + 1,
1169 n_datasets,
1170 transfer,
1171 )
1173 dataset.path = self._standardizeIngestPath(
1174 dataset.path, transfer=transfer, check_existence=check_existence
1175 )
1176 filtered.append(dataset)
1177 return _IngestPrepData(filtered)
1179 @transactional
1180 def _finishIngest(
1181 self,
1182 prepData: Datastore.IngestPrepData,
1183 *,
1184 transfer: str | None = None,
1185 record_validation_info: bool = True,
1186 ) -> None:
1187 # Docstring inherited from Datastore._finishIngest.
1188 refsAndInfos = []
1189 progress = Progress("lsst.daf.butler.datastores.FileDatastore.ingest", level=logging.DEBUG)
1190 for dataset in progress.wrap(prepData.datasets, desc="Ingesting dataset files"):
1191 # Do ingest as if the first dataset ref is associated with the file
1192 info = self._extractIngestInfo(
1193 dataset.path,
1194 dataset.refs[0],
1195 formatter=dataset.formatter,
1196 transfer=transfer,
1197 record_validation_info=record_validation_info,
1198 )
1199 refsAndInfos.extend([(ref, info) for ref in dataset.refs])
1201 # In direct mode we can allow repeated ingests of the same thing
1202 # if we are sure that the external dataset is immutable. We use
1203 # UUIDv5 to indicate this. If there is a mix of v4 and v5 they are
1204 # separated.
1205 refs_and_infos_replace = []
1206 refs_and_infos_insert = []
1207 if transfer == "direct":
1208 for entry in refsAndInfos:
1209 if entry[0].id.version == 5:
1210 refs_and_infos_replace.append(entry)
1211 else:
1212 refs_and_infos_insert.append(entry)
1213 else:
1214 refs_and_infos_insert = refsAndInfos
1216 if refs_and_infos_insert:
1217 self._register_datasets(refs_and_infos_insert, insert_mode=DatabaseInsertMode.INSERT)
1218 if refs_and_infos_replace:
1219 self._register_datasets(refs_and_infos_replace, insert_mode=DatabaseInsertMode.REPLACE)
1221 def _calculate_ingested_datastore_name(
1222 self,
1223 srcUri: ResourcePath,
1224 ref: DatasetRef,
1225 formatter: Formatter | FormatterV2 | type[Formatter | FormatterV2] | None = None,
1226 ) -> Location:
1227 """Given a source URI and a DatasetRef, determine the name the
1228 dataset will have inside datastore.
1230 Parameters
1231 ----------
1232 srcUri : `lsst.resources.ResourcePath`
1233 URI to the source dataset file.
1234 ref : `DatasetRef`
1235 Ref associated with the newly-ingested dataset artifact. This
1236 is used to determine the name within the datastore.
1237 formatter : `Formatter` or Formatter class.
1238 Formatter to use for validation. Can be a class or an instance.
1239 No validation of the file extension is performed if the
1240 ``formatter`` is `None`. This can be used if the caller knows
1241 that the source URI and target URI will use the same formatter.
1243 Returns
1244 -------
1245 location : `Location`
1246 Target location for the newly-ingested dataset.
1247 """
1248 # Ingesting a file from outside the datastore.
1249 # This involves a new name.
1250 template = self.templates.getTemplate(ref)
1251 location = self.locationFactory.fromPath(template.format(ref), trusted_path=True)
1253 # Get the extension
1254 ext = srcUri.getExtension()
1256 # Update the destination to include that extension
1257 location.updateExtension(ext)
1259 # Ask the formatter to validate this extension
1260 if formatter is not None:
1261 formatter.validate_extension(location)
1263 return location
1265 def _write_in_memory_to_artifact(
1266 self, inMemoryDataset: Any, ref: DatasetRef, provenance: DatasetProvenance | None = None
1267 ) -> StoredFileInfo:
1268 """Write out in memory dataset to datastore.
1270 Parameters
1271 ----------
1272 inMemoryDataset : `object`
1273 Dataset to write to datastore.
1274 ref : `DatasetRef`
1275 Registry information associated with this dataset.
1276 provenance : `DatasetProvenance` or `None`, optional
1277 Any provenance that should be attached to the serialized dataset.
1278 Not supported by all formatters.
1280 Returns
1281 -------
1282 info : `StoredFileInfo`
1283 Information describing the artifact written to the datastore.
1284 """
1285 # May need to coerce the in memory dataset to the correct
1286 # python type, but first we need to make sure the storage class
1287 # reflects the one defined in the data repository.
1288 ref = self._cast_storage_class(ref)
1290 # Confirm that we can accept this dataset
1291 if not self.constraints.isAcceptable(ref):
1292 # Raise rather than use boolean return value.
1293 raise DatasetTypeNotSupportedError(
1294 f"Dataset {ref} has been rejected by this datastore via configuration."
1295 )
1297 location, formatter = self._determine_put_formatter_location(ref)
1299 # The external storage class can differ from the registry storage
1300 # class AND the given in-memory dataset might not match any of the
1301 # storage class definitions.
1302 if formatter.can_accept(inMemoryDataset):
1303 # Do not need to coerce. Must assume that the formatter can handle
1304 # it without further checking of types.
1305 pass
1306 else:
1307 # Coerce to a type that it can accept.
1308 inMemoryDataset = ref.datasetType.storageClass.coerce_type(inMemoryDataset)
1309 required_pytype = ref.datasetType.storageClass.pytype
1311 if not isinstance(inMemoryDataset, required_pytype): 1311 ↛ 1312line 1311 didn't jump to line 1312 because the condition on line 1311 was never true
1312 raise TypeError(
1313 f"Inconsistency between supplied object ({type(inMemoryDataset)}) "
1314 f"and storage class type ({required_pytype})"
1315 )
1317 if self._transaction is None: 1317 ↛ 1318line 1317 didn't jump to line 1318 because the condition on line 1317 was never true
1318 raise RuntimeError("Attempting to write artifact without transaction enabled")
1320 def _removeFileExists(uri: ResourcePath) -> None:
1321 """Remove a file and do not complain if it is not there.
1323 This is important since a formatter might fail before the file
1324 is written and we should not confuse people by writing spurious
1325 error messages to the log.
1326 """
1327 with contextlib.suppress(FileNotFoundError):
1328 uri.remove()
1330 # Register a callback to try to delete the uploaded data if
1331 # something fails below
1332 uri = location.uri
1333 self._transaction.registerUndo("artifactWrite", _removeFileExists, uri)
1335 # Need to record the specified formatter but if this is a V1 formatter
1336 # we need to convert it to a V2 compatible shim to do the write.
1337 if not isinstance(formatter, Formatter):
1338 formatter_compat = formatter
1339 else:
1340 formatter_compat = FormatterV1inV2(
1341 formatter.file_descriptor,
1342 ref=ref,
1343 formatter=formatter,
1344 write_parameters=formatter.write_parameters,
1345 write_recipes=formatter.write_recipes,
1346 )
1348 assert isinstance(formatter_compat, FormatterV2)
1350 with time_this(log, msg="Writing dataset %s with formatter %s", args=(ref, formatter.name())):
1351 try:
1352 formatter_compat.write(
1353 inMemoryDataset, cache_manager=self.cacheManager, provenance=provenance
1354 )
1355 except Exception as e:
1356 raise RuntimeError(
1357 f"Failed to serialize dataset {ref} of type {get_full_type_name(inMemoryDataset)} "
1358 f"using formatter {formatter.name()}."
1359 ) from e
1361 # URI is needed to resolve what ingest case are we dealing with
1362 return self._extractIngestInfo(uri, ref, formatter=formatter)
1364 def knows(self, ref: DatasetRef) -> bool:
1365 """Check if the dataset is known to the datastore.
1367 Does not check for existence of any artifact.
1369 Parameters
1370 ----------
1371 ref : `DatasetRef`
1372 Reference to the required dataset.
1374 Returns
1375 -------
1376 exists : `bool`
1377 `True` if the dataset is known to the datastore.
1378 """
1379 fileLocations = self._get_dataset_locations_info(ref)
1380 if fileLocations:
1381 return True
1382 return False
1384 def knows_these(self, refs: Iterable[DatasetRef]) -> dict[DatasetRef, bool]:
1385 # Docstring inherited from the base class.
1386 refs = list(refs)
1388 # The records themselves. Could be missing some entries.
1389 records = self._get_stored_records_associated_with_refs(refs, ignore_datastore_records=True)
1391 return {ref: ref.id in records for ref in refs}
1393 def _process_mexists_records(
1394 self,
1395 id_to_ref: dict[DatasetId, DatasetRef],
1396 records: dict[DatasetId, list[StoredFileInfo]],
1397 all_required: bool,
1398 artifact_existence: dict[ResourcePath, bool] | None = None,
1399 ) -> dict[DatasetRef, bool]:
1400 """Check given records for existence.
1402 Helper function for `mexists()`.
1404 Parameters
1405 ----------
1406 id_to_ref : `dict` of [`DatasetId`, `DatasetRef`]
1407 Mapping of the dataset ID to the dataset ref itself.
1408 records : `dict` of [`DatasetId`, `list` of `StoredFileInfo`]
1409 Records as generally returned by
1410 ``_get_stored_records_associated_with_refs``.
1411 all_required : `bool`
1412 Flag to indicate whether existence requires all artifacts
1413 associated with a dataset ID to exist or not for existence.
1414 artifact_existence : `dict` [`lsst.resources.ResourcePath`, `bool`]
1415 Optional mapping of datastore artifact to existence. Updated by
1416 this method with details of all artifacts tested. Can be `None`
1417 if the caller is not interested.
1419 Returns
1420 -------
1421 existence : `dict` of [`DatasetRef`, `bool`]
1422 Mapping from dataset to boolean indicating existence.
1423 """
1424 # The URIs to be checked and a mapping of those URIs to
1425 # the dataset ID.
1426 uris_to_check: list[ResourcePath] = []
1427 location_map: dict[ResourcePath, DatasetId] = {}
1429 location_factory = self.locationFactory
1431 uri_existence: dict[ResourcePath, bool] = {}
1432 for ref_id, infos in records.items():
1433 # Key is the dataset Id, value is list of StoredItemInfo
1434 uris = [info.file_location(location_factory).uri for info in infos]
1435 location_map.update({uri: ref_id for uri in uris})
1437 # Check the local cache directly for a dataset corresponding
1438 # to the remote URI.
1439 if self.cacheManager.file_count > 0:
1440 ref = id_to_ref[ref_id]
1441 for uri, storedFileInfo in zip(uris, infos, strict=True):
1442 check_ref = ref
1443 if not ref.datasetType.isComponent() and (component := storedFileInfo.component): 1443 ↛ 1444line 1443 didn't jump to line 1444 because the condition on line 1443 was never true
1444 check_ref = ref.makeComponentRef(component)
1445 if self.cacheManager.known_to_cache(check_ref, uri.getExtension()):
1446 # Proxy for URI existence.
1447 uri_existence[uri] = True
1448 else:
1449 uris_to_check.append(uri)
1450 else:
1451 # Check all of them.
1452 uris_to_check.extend(uris)
1454 if artifact_existence is not None:
1455 # If a URI has already been checked remove it from the list
1456 # and immediately add the status to the output dict.
1457 filtered_uris_to_check = []
1458 for uri in uris_to_check:
1459 if uri in artifact_existence: 1459 ↛ 1460line 1459 didn't jump to line 1460 because the condition on line 1459 was never true
1460 uri_existence[uri] = artifact_existence[uri]
1461 else:
1462 filtered_uris_to_check.append(uri)
1463 uris_to_check = filtered_uris_to_check
1465 # Results.
1466 dataset_existence: dict[DatasetRef, bool] = {}
1468 uri_existence.update(ResourcePath.mexists(uris_to_check))
1469 for uri, exists in uri_existence.items():
1470 dataset_id = location_map[uri]
1471 ref = id_to_ref[dataset_id]
1473 # Disassembled composite needs to check all locations.
1474 # all_required indicates whether all need to exist or not.
1475 if ref in dataset_existence:
1476 if all_required:
1477 exists = dataset_existence[ref] and exists
1478 else:
1479 exists = dataset_existence[ref] or exists
1480 dataset_existence[ref] = exists
1482 if artifact_existence is not None:
1483 artifact_existence.update(uri_existence)
1485 return dataset_existence
1487 def mexists(
1488 self, refs: Iterable[DatasetRef], artifact_existence: dict[ResourcePath, bool] | None = None
1489 ) -> dict[DatasetRef, bool]:
1490 """Check the existence of multiple datasets at once.
1492 Parameters
1493 ----------
1494 refs : `~collections.abc.Iterable` of `DatasetRef`
1495 The datasets to be checked.
1496 artifact_existence : `dict` [`lsst.resources.ResourcePath`, `bool`]
1497 Optional mapping of datastore artifact to existence. Updated by
1498 this method with details of all artifacts tested. Can be `None`
1499 if the caller is not interested.
1501 Returns
1502 -------
1503 existence : `dict` of [`DatasetRef`, `bool`]
1504 Mapping from dataset to boolean indicating existence.
1506 Notes
1507 -----
1508 To minimize potentially costly remote existence checks, the local
1509 cache is checked as a proxy for existence. If a file for this
1510 `DatasetRef` does exist no check is done for the actual URI. This
1511 could result in possibly unexpected behavior if the dataset itself
1512 has been removed from the datastore by another process whilst it is
1513 still in the cache.
1514 """
1515 chunk_size = 50_000
1516 dataset_existence: dict[DatasetRef, bool] = {}
1517 log.debug("Checking for the existence of multiple artifacts in datastore in chunks of %d", chunk_size)
1518 n_found_total = 0
1519 n_checked = 0
1520 n_chunks = 0
1521 for chunk in chunk_iterable(refs, chunk_size=chunk_size):
1522 chunk_result = self._mexists(chunk, artifact_existence)
1524 # The log message level and content depend on how many
1525 # datasets we are processing.
1526 n_results = len(chunk_result)
1528 # Use verbose logging to ensure that messages can be seen
1529 # easily if many refs are being checked.
1530 log_threshold = VERBOSE
1531 n_checked += n_results
1533 # This sum can take some time so only do it if we know the
1534 # result is going to be used.
1535 n_found = 0
1536 if log.isEnabledFor(log_threshold):
1537 # Can treat the booleans as 0, 1 integers and sum them.
1538 n_found = sum(chunk_result.values())
1539 n_found_total += n_found
1541 # We are deliberately not trying to count the number of refs
1542 # provided in case it's in the millions. This means there is a
1543 # situation where the number of refs exactly matches the chunk
1544 # size and we will switch to the multi-chunk path even though
1545 # we only have a single chunk.
1546 if n_results < chunk_size and n_chunks == 0: 1546 ↛ 1568line 1546 didn't jump to line 1568 because the condition on line 1546 was always true
1547 # Single chunk will be processed so we can provide more detail.
1548 if n_results == 1:
1549 ref = list(chunk_result)[0]
1550 # Use debug logging to be consistent with `exists()`.
1551 log.debug(
1552 "Calling mexists() with single ref that does%s exist (%s).",
1553 "" if chunk_result[ref] else " not",
1554 ref,
1555 )
1556 else:
1557 # Single chunk but multiple files. Summarize.
1558 log.log(
1559 log_threshold,
1560 "Number of datasets found in datastore %s: %d out of %d datasets checked.",
1561 self.name,
1562 n_found,
1563 n_checked,
1564 )
1566 else:
1567 # Use incremental verbose logging when we have multiple chunks.
1568 log.log(
1569 log_threshold,
1570 "Number of datasets found in datastore for chunk %d: %d out of %d checked "
1571 "(running total from all chunks so far: %d found out of %d checked)",
1572 n_chunks,
1573 n_found,
1574 n_results,
1575 n_found_total,
1576 n_checked,
1577 )
1578 dataset_existence.update(chunk_result)
1579 n_chunks += 1
1581 return dataset_existence
1583 def _mexists(
1584 self, refs: Sequence[DatasetRef], artifact_existence: dict[ResourcePath, bool] | None = None
1585 ) -> dict[DatasetRef, bool]:
1586 """Check the existence of multiple datasets at once.
1588 Parameters
1589 ----------
1590 refs : `~collections.abc.Iterable` of `DatasetRef`
1591 The datasets to be checked.
1592 artifact_existence : `dict` [`lsst.resources.ResourcePath`, `bool`]
1593 Optional mapping of datastore artifact to existence. Updated by
1594 this method with details of all artifacts tested. Can be `None`
1595 if the caller is not interested.
1597 Returns
1598 -------
1599 existence : `dict` of [`DatasetRef`, `bool`]
1600 Mapping from dataset to boolean indicating existence.
1601 """
1602 # Make a mapping from refs with the internal storage class to the given
1603 # refs that may have a different one. We'll use the internal refs
1604 # throughout this method and convert back at the very end.
1605 internal_ref_to_input_ref = {self._cast_storage_class(ref): ref for ref in refs}
1607 # Need a mapping of dataset_id to (internal) dataset ref since some
1608 # internal APIs work with dataset_id.
1609 id_to_ref = {ref.id: ref for ref in internal_ref_to_input_ref}
1611 # Set of all IDs we are checking for.
1612 requested_ids = set(id_to_ref.keys())
1614 # The records themselves. Could be missing some entries.
1615 records = self._get_stored_records_associated_with_refs(
1616 id_to_ref.values(), ignore_datastore_records=True
1617 )
1619 dataset_existence = self._process_mexists_records(
1620 id_to_ref, records, True, artifact_existence=artifact_existence
1621 )
1623 # Set of IDs that have been handled.
1624 handled_ids = {ref.id for ref in dataset_existence}
1626 missing_ids = requested_ids - handled_ids
1627 if missing_ids:
1628 dataset_existence.update(
1629 self._mexists_check_expected(
1630 [id_to_ref[missing] for missing in missing_ids], artifact_existence
1631 )
1632 )
1634 return {
1635 internal_ref_to_input_ref[internal_ref]: existence
1636 for internal_ref, existence in dataset_existence.items()
1637 }
1639 def _mexists_check_expected(
1640 self, refs: Sequence[DatasetRef], artifact_existence: dict[ResourcePath, bool] | None = None
1641 ) -> dict[DatasetRef, bool]:
1642 """Check existence of refs that are not known to datastore.
1644 Parameters
1645 ----------
1646 refs : `~collections.abc.Iterable` of `DatasetRef`
1647 The datasets to be checked. These are assumed not to be known
1648 to datastore.
1649 artifact_existence : `dict` [`lsst.resources.ResourcePath`, `bool`]
1650 Optional mapping of datastore artifact to existence. Updated by
1651 this method with details of all artifacts tested. Can be `None`
1652 if the caller is not interested.
1654 Returns
1655 -------
1656 existence : `dict` of [`DatasetRef`, `bool`]
1657 Mapping from dataset to boolean indicating existence.
1658 """
1659 dataset_existence: dict[DatasetRef, bool] = {}
1660 if not self.trustGetRequest:
1661 # Must assume these do not exist
1662 for ref in refs:
1663 dataset_existence[ref] = False
1664 else:
1665 log.debug(
1666 "%d datasets were not known to datastore during initial existence check.",
1667 len(refs),
1668 )
1670 # Construct data structure identical to that returned
1671 # by _get_stored_records_associated_with_refs() but using
1672 # guessed names.
1673 records = {}
1674 id_to_ref = {}
1675 for missing_ref in refs:
1676 expected = self._get_expected_dataset_locations_info(missing_ref)
1677 dataset_id = missing_ref.id
1678 records[dataset_id] = [info for _, info in expected]
1679 id_to_ref[dataset_id] = missing_ref
1681 dataset_existence.update(
1682 self._process_mexists_records(
1683 id_to_ref,
1684 records,
1685 False,
1686 artifact_existence=artifact_existence,
1687 )
1688 )
1690 return dataset_existence
1692 def exists(self, ref: DatasetRef) -> bool:
1693 """Check if the dataset exists in the datastore.
1695 Parameters
1696 ----------
1697 ref : `DatasetRef`
1698 Reference to the required dataset.
1700 Returns
1701 -------
1702 exists : `bool`
1703 `True` if the entity exists in the `Datastore`.
1705 Notes
1706 -----
1707 The local cache is checked as a proxy for existence in the remote
1708 object store. It is possible that another process on a different
1709 compute node could remove the file from the object store even
1710 though it is present in the local cache.
1711 """
1712 ref = self._cast_storage_class(ref)
1713 # We cannot trust datastore records from ref, as many unit tests delete
1714 # datasets and check their existence.
1715 fileLocations = self._get_dataset_locations_info(ref, ignore_datastore_records=True)
1717 # if we are being asked to trust that registry might not be correct
1718 # we ask for the expected locations and check them explicitly
1719 if not fileLocations:
1720 if not self.trustGetRequest:
1721 return False
1723 # First check the cache. If it is not found we must check
1724 # the datastore itself. Assume that any component in the cache
1725 # means that the dataset does exist somewhere.
1726 if self.cacheManager.known_to_cache(ref):
1727 return True
1729 # When we are guessing a dataset location we can not check
1730 # for the existence of every component since we can not
1731 # know if every component was written. Instead we check
1732 # for the existence of any of the expected locations.
1733 for location, _ in self._get_expected_dataset_locations_info(ref):
1734 if self._artifact_exists(location):
1735 return True
1736 return False
1738 # All listed artifacts must exist.
1739 for location, storedFileInfo in fileLocations:
1740 # Checking in cache needs the component ref.
1741 check_ref = ref
1742 if not ref.datasetType.isComponent() and (component := storedFileInfo.component):
1743 check_ref = ref.makeComponentRef(component)
1744 if self.cacheManager.known_to_cache(check_ref, location.getExtension()):
1745 continue
1747 if not self._artifact_exists(location):
1748 return False
1750 return True
1752 def getURIs(self, ref: DatasetRef, predict: bool = False) -> DatasetRefURIs:
1753 """Return URIs associated with dataset.
1755 Parameters
1756 ----------
1757 ref : `DatasetRef`
1758 Reference to the required dataset.
1759 predict : `bool`, optional
1760 If the datastore does not know about the dataset, controls whether
1761 it should return a predicted URI or not.
1763 Returns
1764 -------
1765 uris : `DatasetRefURIs`
1766 The URI to the primary artifact associated with this dataset (if
1767 the dataset was disassembled within the datastore this may be
1768 `None`), and the URIs to any components associated with the dataset
1769 artifact. (can be empty if there are no components).
1770 """
1771 many = self.getManyURIs([ref], predict=predict, allow_missing=False)
1772 return many[ref]
1774 def getURI(self, ref: DatasetRef, predict: bool = False) -> ResourcePath:
1775 """URI to the Dataset.
1777 Parameters
1778 ----------
1779 ref : `DatasetRef`
1780 Reference to the required Dataset.
1781 predict : `bool`
1782 If `True`, allow URIs to be returned of datasets that have not
1783 been written.
1785 Returns
1786 -------
1787 uri : `str`
1788 URI pointing to the dataset within the datastore. If the
1789 dataset does not exist in the datastore, and if ``predict`` is
1790 `True`, the URI will be a prediction and will include a URI
1791 fragment "#predicted".
1792 If the datastore does not have entities that relate well
1793 to the concept of a URI the returned URI will be
1794 descriptive. The returned URI is not guaranteed to be obtainable.
1796 Raises
1797 ------
1798 FileNotFoundError
1799 Raised if a URI has been requested for a dataset that does not
1800 exist and guessing is not allowed.
1801 RuntimeError
1802 Raised if a request is made for a single URI but multiple URIs
1803 are associated with this dataset.
1805 Notes
1806 -----
1807 When a predicted URI is requested an attempt will be made to form
1808 a reasonable URI based on file templates and the expected formatter.
1809 """
1810 primary, components = self.getURIs(ref, predict)
1811 if primary is None or components: 1811 ↛ 1812line 1811 didn't jump to line 1812 because the condition on line 1811 was never true
1812 raise RuntimeError(
1813 f"Dataset ({ref}) includes distinct URIs for components. Use Datastore.getURIs() instead."
1814 )
1815 return primary
1817 def _predict_URIs(
1818 self,
1819 ref: DatasetRef,
1820 ) -> DatasetRefURIs:
1821 """Predict the URIs of a dataset ref.
1823 Parameters
1824 ----------
1825 ref : `DatasetRef`
1826 Reference to the required Dataset.
1828 Returns
1829 -------
1830 URI : DatasetRefUris
1831 Primary and component URIs. URIs will contain a URI fragment
1832 "#predicted".
1833 """
1834 uris = DatasetRefURIs()
1836 if self.composites.shouldBeDisassembled(ref):
1837 for component, _ in ref.datasetType.storageClass.components.items():
1838 comp_ref = ref.makeComponentRef(component)
1839 comp_location, _ = self._determine_put_formatter_location(comp_ref)
1841 # Add the "#predicted" URI fragment to indicate this is a
1842 # guess
1843 uris.componentURIs[component] = ResourcePath(
1844 comp_location.uri.geturl() + "#predicted", forceDirectory=comp_location.uri.dirLike
1845 )
1847 else:
1848 location, _ = self._determine_put_formatter_location(ref)
1850 # Add the "#predicted" URI fragment to indicate this is a guess
1851 uris.primaryURI = ResourcePath(
1852 location.uri.geturl() + "#predicted", forceDirectory=location.uri.dirLike
1853 )
1855 return uris
1857 def getManyURIs(
1858 self,
1859 refs: Iterable[DatasetRef],
1860 predict: bool = False,
1861 allow_missing: bool = False,
1862 ) -> dict[DatasetRef, DatasetRefURIs]:
1863 # Docstring inherited
1865 uris: dict[DatasetRef, DatasetRefURIs] = {}
1867 records = self._get_stored_records_associated_with_refs(refs)
1868 records_keys = records.keys()
1870 existing_refs = tuple(ref for ref in refs if ref.id in records_keys)
1871 missing_refs = tuple(ref for ref in refs if ref.id not in records_keys)
1873 # Have to handle trustGetRequest mode by checking for the existence
1874 # of the missing refs on disk.
1875 if missing_refs and not predict:
1876 dataset_existence = self._mexists_check_expected(missing_refs, None)
1877 really_missing = set()
1878 not_missing = set()
1879 for ref, exists in dataset_existence.items():
1880 if exists:
1881 not_missing.add(ref)
1882 else:
1883 really_missing.add(ref)
1885 if not_missing:
1886 # Need to recalculate the missing/existing split.
1887 existing_refs = existing_refs + tuple(not_missing)
1888 missing_refs = tuple(really_missing)
1890 for ref in missing_refs:
1891 # if this has never been written then we have to guess
1892 if not predict:
1893 if not allow_missing:
1894 raise FileNotFoundError(f"Dataset {ref} not in this datastore.")
1895 else:
1896 uris[ref] = self._predict_URIs(ref)
1898 for ref in existing_refs:
1899 file_infos = records[ref.id]
1900 file_locations = [(i.file_location(self.locationFactory), i) for i in file_infos]
1901 uris[ref] = self._locations_to_URI(ref, file_locations)
1903 return uris
1905 def _locations_to_URI(
1906 self,
1907 ref: DatasetRef,
1908 file_locations: Sequence[tuple[Location, StoredFileInfo]],
1909 ) -> DatasetRefURIs:
1910 """Convert one or more file locations associated with a DatasetRef
1911 to a DatasetRefURIs.
1913 Parameters
1914 ----------
1915 ref : `DatasetRef`
1916 Reference to the dataset.
1917 file_locations : Sequence[Tuple[Location, StoredFileInfo]]
1918 Each item in the sequence is the location of the dataset within the
1919 datastore and stored information about the file and its formatter.
1920 If there is only one item in the sequence then it is treated as the
1921 primary URI. If there is more than one item then they are treated
1922 as component URIs. If there are no items then an error is raised
1923 unless ``self.trustGetRequest`` is `True`.
1925 Returns
1926 -------
1927 uris: DatasetRefURIs
1928 Represents the primary URI or component URIs described by the
1929 inputs.
1931 Raises
1932 ------
1933 RuntimeError
1934 If no file locations are passed in and ``self.trustGetRequest`` is
1935 `False`.
1936 FileNotFoundError
1937 If the a passed-in URI does not exist, and ``self.trustGetRequest``
1938 is `False`.
1939 RuntimeError
1940 If a passed in `StoredFileInfo`'s ``component`` is `None` (this is
1941 unexpected).
1942 """
1943 guessing = False
1944 uris = DatasetRefURIs()
1946 if not file_locations:
1947 if not self.trustGetRequest: 1947 ↛ 1948line 1947 didn't jump to line 1948 because the condition on line 1947 was never true
1948 raise RuntimeError(f"Unexpectedly got no artifacts for dataset {ref}")
1949 file_locations = self._get_expected_dataset_locations_info(ref)
1950 guessing = True
1952 if len(file_locations) == 1:
1953 # No disassembly so this is the primary URI
1954 uris.primaryURI = file_locations[0][0].uri
1955 if guessing and not uris.primaryURI.exists(): 1955 ↛ 1956line 1955 didn't jump to line 1956 because the condition on line 1955 was never true
1956 raise FileNotFoundError(f"Expected URI ({uris.primaryURI}) does not exist")
1957 else:
1958 for location, file_info in file_locations:
1959 if file_info.component is None: 1959 ↛ 1960line 1959 didn't jump to line 1960 because the condition on line 1959 was never true
1960 raise RuntimeError(f"Unexpectedly got no component name for a component at {location}")
1961 if guessing and not location.uri.exists(): 1961 ↛ 1965line 1961 didn't jump to line 1965 because the condition on line 1961 was never true
1962 # If we are trusting then it is entirely possible for
1963 # some components to be missing. In that case we skip
1964 # to the next component.
1965 if self.trustGetRequest:
1966 continue
1967 raise FileNotFoundError(f"Expected URI ({location.uri}) does not exist")
1968 uris.componentURIs[file_info.component] = location.uri
1970 return uris
1972 def _find_missing_records(
1973 self,
1974 refs: Iterable[DatasetRef],
1975 missing_ids: set[DatasetId],
1976 artifact_existence: dict[ResourcePath, bool] | None = None,
1977 warn_for_missing: bool = True,
1978 ) -> dict[DatasetId, list[StoredFileInfo]]:
1979 if not missing_ids:
1980 return {}
1982 if artifact_existence is None: 1982 ↛ 1983line 1982 didn't jump to line 1983 because the condition on line 1982 was never true
1983 artifact_existence = {}
1985 found_records: dict[DatasetId, list[StoredFileInfo]] = defaultdict(list)
1986 id_to_ref = {ref.id: ref for ref in refs if ref.id in missing_ids}
1988 # This should be chunked in case we end up having to check
1989 # the file store since we need some log output to show
1990 # progress.
1991 chunk_size = 50_000
1992 for missing_ids_chunk in chunk_iterable(missing_ids, chunk_size=chunk_size):
1993 records = {}
1994 for missing in missing_ids_chunk:
1995 # Ask the source datastore where the missing artifacts
1996 # should be. An execution butler might not know about the
1997 # artifacts even if they are there.
1998 expected = self._get_expected_dataset_locations_info(id_to_ref[missing])
1999 records[missing] = [info for _, info in expected]
2001 # Call the mexist helper method in case we have not already
2002 # checked these artifacts such that artifact_existence is
2003 # empty. This allows us to benefit from parallelism.
2004 # datastore.mexists() itself does not give us access to the
2005 # derived datastore record.
2006 log.verbose("Checking existence of %d datasets unknown to datastore", len(records))
2007 ref_exists = self._process_mexists_records(
2008 id_to_ref, records, False, artifact_existence=artifact_existence
2009 )
2011 # Now go through the records and propagate the ones that exist.
2012 location_factory = self.locationFactory
2013 for missing, record_list in records.items():
2014 # Skip completely if the ref does not exist.
2015 ref = id_to_ref[missing]
2016 if not ref_exists[ref]:
2017 if warn_for_missing: 2017 ↛ 2018line 2017 didn't jump to line 2018 because the condition on line 2017 was never true
2018 log.warning("Asked to transfer dataset %s but no file artifacts exist for it.", ref)
2019 continue
2020 # Check for file artifact to decide which parts of a
2021 # disassembled composite do exist. If there is only a
2022 # single record we don't even need to look because it can't
2023 # be a composite and must exist.
2024 if len(record_list) == 1:
2025 dataset_records = record_list
2026 else:
2027 dataset_records = [
2028 record
2029 for record in record_list
2030 if artifact_existence[record.file_location(location_factory).uri]
2031 ]
2032 assert len(dataset_records) > 0, "Disassembled composite should have had some files."
2034 # Rely on source_records being a defaultdict.
2035 found_records[missing].extend(dataset_records)
2036 log.verbose("Completed scan for missing data files")
2037 return found_records
2039 def retrieveArtifacts(
2040 self,
2041 refs: Iterable[DatasetRef],
2042 destination: ResourcePath,
2043 transfer: str = "auto",
2044 preserve_path: bool = True,
2045 overwrite: bool = False,
2046 write_index: bool = True,
2047 add_prefix: bool = False,
2048 ) -> dict[ResourcePath, ArtifactIndexInfo]:
2049 """Retrieve the file artifacts associated with the supplied refs.
2051 Parameters
2052 ----------
2053 refs : `~collections.abc.Iterable` of `DatasetRef`
2054 The datasets for which file artifacts are to be retrieved.
2055 A single ref can result in multiple files. The refs must
2056 be resolved.
2057 destination : `lsst.resources.ResourcePath`
2058 Location to write the file artifacts.
2059 transfer : `str`, optional
2060 Method to use to transfer the artifacts. Must be one of the options
2061 supported by `lsst.resources.ResourcePath.transfer_from`.
2062 "move" is not allowed.
2063 preserve_path : `bool`, optional
2064 If `True` the full path of the file artifact within the datastore
2065 is preserved. If `False` the final file component of the path
2066 is used.
2067 overwrite : `bool`, optional
2068 If `True` allow transfers to overwrite existing files at the
2069 destination.
2070 write_index : `bool`, optional
2071 If `True` write a file at the top level containing a serialization
2072 of a `ZipIndex` for the downloaded datasets.
2073 add_prefix : `bool`, optional
2074 If `True` and if ``preserve_path`` is `False`, apply a prefix to
2075 the filenames corresponding to some part of the dataset ref ID.
2076 This can be used to guarantee uniqueness.
2078 Returns
2079 -------
2080 artifact_map : `dict` [ `lsst.resources.ResourcePath`, \
2081 `ArtifactIndexInfo` ]
2082 Mapping of retrieved file to associated index information.
2083 """
2084 if not destination.isdir():
2085 raise ValueError(f"Destination location must refer to a directory. Given {destination}")
2087 if transfer == "move":
2088 raise ValueError("Can not move artifacts out of datastore. Use copy instead.")
2090 # Source -> Destination
2091 # This also helps filter out duplicate DatasetRef in the request
2092 # that will map to the same underlying file transfer.
2093 to_transfer: dict[ResourcePath, ResourcePath] = {}
2094 zips_to_transfer: set[ResourcePath] = set()
2096 # Retrieve all the records in bulk indexed by ref.id.
2097 records = self._get_stored_records_associated_with_refs(refs, ignore_datastore_records=True)
2099 # Check for missing records.
2100 known_ids = set(records)
2101 log.debug("Number of datastore records found in database: %d", len(known_ids))
2102 requested_ids = {ref.id for ref in refs}
2103 missing_ids = requested_ids - known_ids
2105 if missing_ids and not self.trustGetRequest: 2105 ↛ 2106line 2105 didn't jump to line 2106 because the condition on line 2105 was never true
2106 raise ValueError(f"Number of datasets missing from this datastore: {len(missing_ids)}")
2108 missing_records = self._find_missing_records(refs, missing_ids)
2109 records.update(missing_records)
2111 # One artifact can be used by multiple DatasetRef.
2112 # e.g. DECam.
2113 artifact_map: dict[ResourcePath, ArtifactIndexInfo] = {}
2114 # Sort to ensure that in many refs to one file situation the same
2115 # ref is used for any prefix that might be added.
2116 for ref in sorted(refs):
2117 prefix = str(ref.id)[:8] + "-" if add_prefix else ""
2118 for info in records[ref.id]:
2119 location = info.file_location(self.locationFactory)
2120 source_uri = location.uri
2121 # For DECam/zip we only want to copy once.
2122 # For zip files we need to unpack so that they can be
2123 # zipped up again if needed.
2124 is_zip = source_uri.getExtension() == ".zip" and "zip-path" in source_uri.fragment
2125 # We need to remove fragments for consistency.
2126 cleaned_source_uri = source_uri.replace(fragment="", query="", params="")
2127 if is_zip: 2127 ↛ 2130line 2127 didn't jump to line 2130 because the condition on line 2127 was never true
2128 # Assume the DatasetRef definitions are within the Zip
2129 # file itself and so can be dropped from loop.
2130 zips_to_transfer.add(cleaned_source_uri)
2131 elif cleaned_source_uri not in to_transfer: 2131 ↛ 2138line 2131 didn't jump to line 2138 because the condition on line 2131 was always true
2132 target_uri = determine_destination_for_retrieved_artifact(
2133 destination, location.pathInStore, preserve_path, prefix
2134 )
2135 to_transfer[cleaned_source_uri] = target_uri
2136 artifact_map[target_uri] = ArtifactIndexInfo.from_single(info.to_simple(), ref.id)
2137 else:
2138 target_uri = to_transfer[cleaned_source_uri]
2139 artifact_map[target_uri].append(ref.id)
2141 # Parallelize the transfer. Re-raise as a single exception if
2142 # a FileExistsError is encountered anywhere.
2143 log.debug("Number of artifacts to transfer to %s: %d", str(destination), len(to_transfer))
2144 try:
2145 ResourcePath.mtransfer(transfer, tuple(to_transfer.items()), overwrite=overwrite)
2146 except* FileExistsError as egroup:
2147 raise FileExistsError(
2148 "Some files already exist in destination directory and overwrite is False"
2149 ) from egroup
2151 # Transfer the Zip files and unpack them.
2152 zipped_artifacts = unpack_zips(zips_to_transfer, requested_ids, destination, preserve_path, overwrite)
2153 artifact_map.update(zipped_artifacts)
2155 if write_index:
2156 index = ZipIndex.from_artifact_map(refs, artifact_map, destination)
2157 index.write_index(destination)
2159 return artifact_map
2161 def ingest_zip(
2162 self,
2163 zip_path: ResourcePath,
2164 transfer: str | None,
2165 *,
2166 dry_run: bool = False,
2167 ) -> None:
2168 """Ingest an indexed Zip file and contents.
2170 The Zip file must have an index file as created by `retrieveArtifacts`.
2172 Parameters
2173 ----------
2174 zip_path : `lsst.resources.ResourcePath`
2175 Path to the Zip file.
2176 transfer : `str`
2177 Method to use for transferring the Zip file into the datastore.
2178 dry_run : `bool`, optional
2179 If `True` the ingest will be processed without any modifications
2180 made to the target datastore and as if the target datastore did not
2181 have any of the datasets.
2183 Notes
2184 -----
2185 Datastore constraints are bypassed with Zip ingest. A zip file can
2186 contain multiple dataset types. Should the entire Zip be rejected
2187 if one dataset type is in the constraints list?
2189 If any dataset is already present in the datastore the entire ingest
2190 will fail.
2191 """
2192 index = ZipIndex.from_zip_file(zip_path)
2194 # Refs indexed by UUID.
2195 refs = index.refs.to_refs(universe=self.universe)
2196 id_to_ref = {ref.id: ref for ref in refs}
2198 # Any failing constraints trigger entire failure.
2199 if any(not self.constraints.isAcceptable(ref) for ref in refs): 2199 ↛ 2200line 2199 didn't jump to line 2200 because the condition on line 2199 was never true
2200 raise DatasetTypeNotSupportedError(
2201 "Some refs in the Zip file are not supported by this datastore"
2202 )
2204 # Transfer the Zip file into the datastore file system.
2205 # There is no RUN as such to use for naming.
2206 # Potentially could use the RUN from the first ref in the index
2207 # There is no requirement that the contents of the Zip files share
2208 # the same RUN.
2209 # Could use the Zip UUID from the index + special "zips/" prefix.
2210 if transfer is None: 2210 ↛ 2212line 2210 didn't jump to line 2212 because the condition on line 2210 was never true
2211 # Indicated that the zip file is already in the right place.
2212 if not zip_path.isabs():
2213 tgtLocation = self.locationFactory.fromPath(zip_path.ospath, trusted_path=False)
2214 else:
2215 pathInStore = zip_path.relative_to(self.root)
2216 if pathInStore is None:
2217 raise RuntimeError(
2218 f"Unexpectedly learned that {zip_path} is not within datastore {self.root}"
2219 )
2220 tgtLocation = self.locationFactory.fromPath(pathInStore, trusted_path=True)
2221 elif transfer == "direct": 2221 ↛ 2223line 2221 didn't jump to line 2223 because the condition on line 2221 was never true
2222 # Reference in original location.
2223 tgtLocation = None
2224 else:
2225 # Name the zip file based on index contents.
2226 tgtLocation = self.locationFactory.fromPath(index.calculate_zip_file_path_in_store())
2228 # Transfer the Zip file into the datastore.
2229 if not dry_run:
2230 tgtLocation.uri.transfer_from(
2231 zip_path, transfer=transfer, transaction=self._transaction, overwrite=True
2232 )
2233 else:
2234 log.info("Would be copying Zip from %s to %s", zip_path, tgtLocation)
2236 if tgtLocation is None: 2236 ↛ 2237line 2236 didn't jump to line 2237 because the condition on line 2236 was never true
2237 path_in_store = str(zip_path)
2238 else:
2239 path_in_store = tgtLocation.pathInStore.path
2241 # Associate each file with a (DatasetRef, StoredFileInfo) tuple.
2242 artifacts: list[tuple[DatasetRef, StoredFileInfo]] = []
2243 for path_in_zip, index_info in index.artifact_map.items():
2244 # Need to modify the info to include the path to the Zip file
2245 # that was previously written to the datastore.
2246 index_info.info.path = f"{path_in_store}#zip-path={path_in_zip}"
2248 info = StoredFileInfo.from_simple(index_info.info)
2249 for id_ in index_info.ids:
2250 artifacts.append((id_to_ref[id_], info))
2252 if not dry_run:
2253 self._register_datasets(artifacts, insert_mode=DatabaseInsertMode.INSERT)
2254 else:
2255 log.info("Would be registering %d artifacts from Zip into datastore", len(artifacts))
2257 def get(
2258 self,
2259 ref: DatasetRef,
2260 parameters: Mapping[str, Any] | None = None,
2261 storageClass: StorageClass | str | None = None,
2262 ) -> Any:
2263 """Load an InMemoryDataset from the store.
2265 Parameters
2266 ----------
2267 ref : `DatasetRef`
2268 Reference to the required Dataset.
2269 parameters : `dict`
2270 `StorageClass`-specific parameters that specify, for example,
2271 a slice of the dataset to be loaded.
2272 storageClass : `StorageClass` or `str`, optional
2273 The storage class to be used to override the Python type
2274 returned by this method. By default the returned type matches
2275 the dataset type definition for this dataset. Specifying a
2276 read `StorageClass` can force a different type to be returned.
2277 This type must be compatible with the original type.
2279 Returns
2280 -------
2281 inMemoryDataset : `object`
2282 Requested dataset or slice thereof as an InMemoryDataset.
2284 Raises
2285 ------
2286 FileNotFoundError
2287 Requested dataset can not be retrieved.
2288 TypeError
2289 Return value from formatter has unexpected type.
2290 ValueError
2291 Formatter failed to process the dataset.
2292 """
2293 # Supplied storage class for the component being read is either
2294 # from the ref itself or some an override if we want to force
2295 # type conversion.
2296 if storageClass is not None:
2297 ref = ref.overrideStorageClass(storageClass)
2299 allGetInfo = self._prepare_for_direct_get(ref, parameters)
2300 return get_dataset_as_python_object_from_get_info(
2301 allGetInfo, ref=ref, parameters=parameters, cache_manager=self.cacheManager
2302 )
2304 def prepare_get_for_external_client(self, ref: DatasetRef) -> list[DatasetLocationInformation] | None:
2305 # Docstring inherited
2307 locations = self._get_dataset_locations_info(ref)
2308 if len(locations) == 0: 2308 ↛ 2311line 2308 didn't jump to line 2311 because the condition on line 2308 was always true
2309 return None
2311 return locations
2313 @transactional
2314 def put(self, inMemoryDataset: Any, ref: DatasetRef, provenance: DatasetProvenance | None = None) -> None:
2315 """Write a InMemoryDataset with a given `DatasetRef` to the store.
2317 Parameters
2318 ----------
2319 inMemoryDataset : `object`
2320 The dataset to store.
2321 ref : `DatasetRef`
2322 Reference to the associated Dataset.
2323 provenance : `DatasetProvenance` or `None`, optional
2324 Any provenance that should be attached to the serialized dataset.
2325 Can be ignored by a formatter or delegate.
2327 Raises
2328 ------
2329 TypeError
2330 Supplied object and storage class are inconsistent.
2331 DatasetTypeNotSupportedError
2332 The associated `DatasetType` is not handled by this datastore.
2334 Notes
2335 -----
2336 If the datastore is configured to reject certain dataset types it
2337 is possible that the put will fail and raise a
2338 `DatasetTypeNotSupportedError`. The main use case for this is to
2339 allow `ChainedDatastore` to put to multiple datastores without
2340 requiring that every datastore accepts the dataset.
2341 """
2342 doDisassembly = self.composites.shouldBeDisassembled(ref)
2343 # doDisassembly = True
2345 artifacts = []
2346 if doDisassembly:
2347 inMemoryDataset = ref.datasetType.storageClass.delegate().add_provenance(
2348 inMemoryDataset, ref, provenance=provenance
2349 )
2350 components = ref.datasetType.storageClass.delegate().disassemble(inMemoryDataset)
2351 if components is None: 2351 ↛ 2352line 2351 didn't jump to line 2352 because the condition on line 2351 was never true
2352 raise RuntimeError(
2353 f"Inconsistent configuration: dataset type {ref.datasetType.name} "
2354 f"with storage class {ref.datasetType.storageClass.name} "
2355 "is configured to be disassembled, but cannot be."
2356 )
2357 for component, componentInfo in components.items():
2358 # Don't recurse because we want to take advantage of
2359 # bulk insert -- need a new DatasetRef that refers to the
2360 # same dataset_id but has the component DatasetType
2361 # DatasetType does not refer to the types of components
2362 # So we construct one ourselves.
2363 compRef = ref.makeComponentRef(component)
2364 # Provenance has already been attached above.
2365 storedInfo = self._write_in_memory_to_artifact(componentInfo.component, compRef)
2366 artifacts.append((compRef, storedInfo))
2367 else:
2368 # Write the entire thing out
2369 storedInfo = self._write_in_memory_to_artifact(inMemoryDataset, ref, provenance=provenance)
2370 artifacts.append((ref, storedInfo))
2372 self._register_datasets(artifacts, insert_mode=DatabaseInsertMode.INSERT)
2374 @transactional
2375 def put_new(self, in_memory_dataset: Any, ref: DatasetRef) -> Mapping[str, DatasetRef]:
2376 doDisassembly = self.composites.shouldBeDisassembled(ref)
2377 # doDisassembly = True
2379 artifacts = []
2380 if doDisassembly:
2381 components = ref.datasetType.storageClass.delegate().disassemble(in_memory_dataset)
2382 if components is None:
2383 raise RuntimeError(
2384 f"Inconsistent configuration: dataset type {ref.datasetType.name} "
2385 f"with storage class {ref.datasetType.storageClass.name} "
2386 "is configured to be disassembled, but cannot be."
2387 )
2388 for component, componentInfo in components.items():
2389 # Don't recurse because we want to take advantage of
2390 # bulk insert -- need a new DatasetRef that refers to the
2391 # same dataset_id but has the component DatasetType
2392 # DatasetType does not refer to the types of components
2393 # So we construct one ourselves.
2394 compRef = ref.makeComponentRef(component)
2395 storedInfo = self._write_in_memory_to_artifact(componentInfo.component, compRef)
2396 artifacts.append((compRef, storedInfo))
2397 else:
2398 # Write the entire thing out
2399 storedInfo = self._write_in_memory_to_artifact(in_memory_dataset, ref)
2400 artifacts.append((ref, storedInfo))
2402 ref_records: DatasetDatastoreRecords = {self._opaque_table_name: [info for _, info in artifacts]}
2403 ref = ref.replace(datastore_records=ref_records)
2404 return {self.name: ref}
2406 @transactional
2407 def trash(self, ref: DatasetRef | Iterable[DatasetRef], ignore_errors: bool = True) -> None:
2408 # At this point can safely remove these datasets from the cache
2409 # to avoid confusion later on. If they are not trashed later
2410 # the cache will simply be refilled.
2411 self.cacheManager.remove_from_cache(ref)
2413 # If we are in trust mode there will be nothing to move to
2414 # the trash table and we will have to try to delete the file
2415 # immediately.
2416 if self.trustGetRequest:
2417 # Try to keep the logic below for a single file trash.
2418 if isinstance(ref, DatasetRef):
2419 refs = {ref}
2420 else:
2421 # Will recreate ref at the end of this branch.
2422 refs = set(ref)
2424 # Determine which datasets are known to datastore directly.
2425 id_to_ref = {ref.id: ref for ref in refs}
2426 existing_ids = self._get_stored_records_associated_with_refs(refs, ignore_datastore_records=True)
2427 existing_refs = {id_to_ref[ref_id] for ref_id in existing_ids}
2429 missing = refs - existing_refs
2430 if missing:
2431 # Do an explicit existence check on these refs.
2432 # We only care about the artifacts at this point and not
2433 # the dataset existence.
2434 artifact_existence: dict[ResourcePath, bool] = {}
2435 _ = self.mexists(missing, artifact_existence)
2436 uris = [uri for uri, exists in artifact_existence.items() if exists]
2438 # FUTURE UPGRADE: Implement a parallelized bulk remove.
2439 log.debug("Removing %d artifacts from datastore that are unknown to datastore", len(uris))
2440 for uri in uris:
2441 try:
2442 uri.remove()
2443 except Exception as e:
2444 if ignore_errors:
2445 log.debug("Artifact %s could not be removed: %s", uri, e)
2446 continue
2447 raise
2449 # There is no point asking the code below to remove refs we
2450 # know are missing so update it with the list of existing
2451 # records. Try to retain one vs many logic.
2452 if not existing_refs:
2453 # Nothing more to do since none of the datasets were
2454 # known to the datastore record table.
2455 return
2456 ref = list(existing_refs)
2457 if len(ref) == 1:
2458 ref = ref[0]
2460 # Get file metadata and internal metadata
2461 if not isinstance(ref, DatasetRef):
2462 log.debug("Doing multi-dataset trash in datastore %s", self.name)
2463 # Assumed to be an iterable of refs so bulk mode enabled.
2464 try:
2465 self.bridge.moveToTrash(ref, transaction=self._transaction)
2466 except Exception as e:
2467 if ignore_errors:
2468 log.warning("Unexpected issue moving multiple datasets to trash: %s", e)
2469 else:
2470 raise
2471 return
2473 log.debug("Trashing dataset %s in datastore %s", ref, self.name)
2475 fileLocations = self._get_dataset_locations_info(ref)
2477 if not fileLocations:
2478 err_msg = f"Requested dataset to trash ({ref}) is not known to datastore {self.name}"
2479 if ignore_errors:
2480 log.warning(err_msg)
2481 return
2482 else:
2483 raise FileNotFoundError(err_msg)
2485 for location, _ in fileLocations:
2486 if not self._artifact_exists(location): 2486 ↛ 2487line 2486 didn't jump to line 2487 because the condition on line 2486 was never true
2487 err_msg = (
2488 f"Dataset is known to datastore {self.name} but "
2489 f"associated artifact ({location.uri}) is missing"
2490 )
2491 if ignore_errors:
2492 log.warning(err_msg)
2493 return
2494 else:
2495 raise FileNotFoundError(err_msg)
2497 # Mark dataset as trashed
2498 try:
2499 self.bridge.moveToTrash([ref], transaction=self._transaction)
2500 except Exception as e:
2501 if ignore_errors:
2502 log.warning(
2503 "Attempted to mark dataset (%s) to be trashed in datastore %s "
2504 "but encountered an error: %s",
2505 ref,
2506 self.name,
2507 e,
2508 )
2509 pass
2510 else:
2511 raise
2513 def emptyTrash(
2514 self, ignore_errors: bool = True, refs: Collection[DatasetRef] | None = None, dry_run: bool = False
2515 ) -> set[ResourcePath]:
2516 """Remove all datasets from the trash.
2518 Parameters
2519 ----------
2520 ignore_errors : `bool`
2521 If `True` return without error even if something went wrong.
2522 Problems could occur if another process is simultaneously trying
2523 to delete.
2524 refs : `collections.abc.Collection` [ `DatasetRef` ] or `None`
2525 Explicit list of datasets that can be removed from trash. If listed
2526 datasets are not already stored in the trash table they will be
2527 ignored. If `None` every entry in the trash table will be
2528 processed.
2529 dry_run : `bool`, optional
2530 If `True`, the trash table will be queried and results reported
2531 but no artifacts will be removed.
2533 Returns
2534 -------
2535 removed : `set` [ `lsst.resources.ResourcePath` ]
2536 List of artifacts that were removed.
2538 Notes
2539 -----
2540 Will empty the records from the trash tables only if this call finishes
2541 without raising.
2542 """
2543 removed = set()
2544 if refs:
2545 selected_ids = {ref.id for ref in refs}
2546 chunk_size = 50_000
2547 n_chunks = math.ceil(len(selected_ids) / chunk_size)
2548 chunk_num = 0
2549 for chunk in chunk_iterable(selected_ids, chunk_size=chunk_size):
2550 chunk_num += 1
2551 if n_chunks == 1: 2551 ↛ 2558line 2551 didn't jump to line 2558 because the condition on line 2551 was always true
2552 log.verbose(
2553 "Emptying datastore trash for %d dataset%s",
2554 len(chunk),
2555 "s" if len(chunk) != 1 else "",
2556 )
2557 else:
2558 log.verbose(
2559 "Emptying datastore trash for chunk %d out of %d of size %d",
2560 chunk_num,
2561 n_chunks,
2562 len(chunk),
2563 )
2564 removed.update(
2565 self._empty_trash_subset(ignore_errors=ignore_errors, selected_ids=chunk, dry_run=dry_run)
2566 )
2567 else:
2568 log.verbose("Emptying all trash in datastore %s", self.name)
2569 removed = self._empty_trash_subset(ignore_errors=ignore_errors, dry_run=dry_run)
2570 log.info(
2571 "%sRemoved %d file artifact%s from datastore %s",
2572 "Would have " if dry_run else "",
2573 len(removed),
2574 "s" if len(removed) != 1 else "",
2575 self.name,
2576 )
2577 return removed
2579 @transactional
2580 def _empty_trash_subset(
2581 self,
2582 *,
2583 ignore_errors: bool = True,
2584 selected_ids: Collection[DatasetId] | None = None,
2585 dry_run: bool = False,
2586 ) -> set[ResourcePath]:
2587 """Empty trash table in transaction.
2589 Parameters
2590 ----------
2591 ignore_errors : `bool`
2592 If `True` return without error even if something went wrong.
2593 Problems could occur if another process is simultaneously trying
2594 to delete.
2595 selected_ids : `collections.abc.collection` [`DatasetId`] or `None`
2596 Explicit list of dataset IDs that can be removed from the trash.
2597 If listed datasets are not already included in the trash table
2598 they will be ignored. If `None` every entry in the trash table
2599 will be processed.
2600 dry_run : `bool`, optional
2601 If `True`, the trash table will be queried and results reported
2602 but no artifacts will be removed.
2604 Returns
2605 -------
2606 removed : `set` [ `lsst.resources.ResourcePath` ]
2607 Artifacts successfully removed.
2609 Notes
2610 -----
2611 Will empty the records from the trash tables only if this call finishes
2612 without raising.
2613 """
2614 # Context manager will empty trash iff we finish it without raising.
2615 # It will also automatically delete the relevant rows from the
2616 # trash table and the records table.
2617 with self.bridge.emptyTrash(
2618 self._table,
2619 record_class=StoredFileInfo,
2620 record_column="path",
2621 selected_ids=selected_ids,
2622 dry_run=dry_run,
2623 ) as trash_data:
2624 # Removing the artifacts themselves requires that the files are
2625 # not also associated with refs that are not to be trashed.
2626 # Therefore need to do a query with the file paths themselves
2627 # and return all the refs associated with them. Can only delete
2628 # a file if the refs to be trashed are the only refs associated
2629 # with the file.
2630 # This requires multiple copies of the trashed items
2631 trashed, artifacts_to_keep = trash_data
2633 # Assume that # in path means there are fragments involved. The
2634 # fragments can not be handled by the emptyTrash bridge call
2635 # so need to be processed independently.
2636 # The generator has to be converted to a list for multiple
2637 # iterations. Clean up the typing so that multiple isinstance
2638 # tests aren't needed later.
2639 trashed_list = [(ref, ninfo) for ref, ninfo in trashed if isinstance(ninfo, StoredFileInfo)]
2641 if artifacts_to_keep is None or any("#" in info[1].path for info in trashed_list):
2642 # The bridge is not helping us so have to work it out
2643 # ourselves. This is not going to be as efficient.
2644 # This mapping does not include the fragments.
2645 if artifacts_to_keep is not None:
2646 # This means we have already checked for non-fragment
2647 # examples so can filter.
2648 paths_to_check = {info.path for _, info in trashed_list if "#" in info.path}
2649 else:
2650 paths_to_check = {info.path for _, info in trashed_list}
2652 path_map = self._refs_associated_with_artifacts(paths_to_check)
2654 for ref, info in trashed_list:
2655 path = info.artifact_path
2656 # For disassembled composites in a Zip it is possible
2657 # for the same path to correspond to the same dataset ref
2658 # multiple times so trap for that.
2659 if ref.id in path_map[path]:
2660 path_map[path].remove(ref.id)
2661 if not path_map[path]:
2662 del path_map[path]
2664 slow_artifacts_to_keep = set(path_map)
2665 if artifacts_to_keep is not None:
2666 artifacts_to_keep.update(slow_artifacts_to_keep)
2667 else:
2668 artifacts_to_keep = slow_artifacts_to_keep
2670 n_direct = 0
2671 artifacts_to_delete: set[ResourcePath] = set()
2672 for ref, info in trashed_list:
2673 # Should not happen for this implementation but need
2674 # to keep mypy happy.
2675 assert info is not None, f"Internal logic error in emptyTrash with ref {ref}."
2677 if info.artifact_path in artifacts_to_keep:
2678 # This is a multi-dataset artifact and we are not
2679 # removing all associated refs.
2680 continue
2682 # Only trashed refs still known to datastore will be returned.
2683 location = info.file_location(self.locationFactory)
2685 if location.pathInStore.isabs(): 2685 ↛ 2686line 2685 didn't jump to line 2686 because the condition on line 2685 was never true
2686 n_direct += 1
2687 continue
2689 # Strip fragment before storing since it is the artifact
2690 # we are deleting and we do not want repeats for every member
2691 # in a zip.
2692 artifacts_to_delete.add(location.uri.replace(fragment=""))
2694 if n_direct > 0: 2694 ↛ 2695line 2694 didn't jump to line 2695 because the condition on line 2694 was never true
2695 s = "s" if n_direct != 1 else ""
2696 log.verbose("Not deleting %d artifact%s using absolute URI%s", n_direct, s, s)
2698 if artifacts_to_keep:
2699 log.verbose(
2700 "%d artifact%s %s not deleted because of association with other datasets",
2701 len(artifacts_to_keep),
2702 "s" if len(artifacts_to_keep) != 1 else "",
2703 "were" if len(artifacts_to_keep) != 1 else "was",
2704 )
2706 if not artifacts_to_delete:
2707 return set()
2709 # Now do the deleting. Special case the log message for a single
2710 # artifact.
2711 if len(artifacts_to_delete) == 1:
2712 log.verbose(
2713 "%s removing file artifact %s from datastore %s",
2714 "Would be" if dry_run else "Now",
2715 list(artifacts_to_delete)[0],
2716 self.name,
2717 )
2718 else:
2719 log.verbose(
2720 "%s removing %d file artifacts from datastore %s",
2721 "Would be" if dry_run else "Now",
2722 len(artifacts_to_delete),
2723 self.name,
2724 )
2726 # For dry-run mode do not attempt to search the file store for
2727 # the artifacts to determine whether they exist or not. Simply
2728 # report that an attempt would be made to delete them. Never
2729 # report direct imports.
2730 if dry_run:
2731 return artifacts_to_delete
2733 # Now remove the actual file artifacts.
2734 remove_result = ResourcePath.mremove(artifacts_to_delete, do_raise=False)
2736 removed: set[ResourcePath] = set()
2737 exceptions: list[Exception] = []
2738 for uri, result in remove_result.items():
2739 if result.exception is None or isinstance(result.exception, FileNotFoundError): 2739 ↛ 2746line 2739 didn't jump to line 2746 because the condition on line 2739 was always true
2740 # File not existing is not an error since some other
2741 # process might have been trying to clean it and we do not
2742 # want to raise an error for a situation where the file
2743 # is not there and we do not want it to be there.
2744 removed.add(uri)
2745 else:
2746 exceptions.append(result.exception)
2748 if exceptions: 2748 ↛ 2749line 2748 didn't jump to line 2749 because the condition on line 2748 was never true
2749 s_err = "s" if len(exceptions) != 1 else ""
2750 e = ExceptionGroup(f"Error{s_err} removing {len(exceptions)} artifact{s_err}", exceptions)
2751 if ignore_errors:
2752 # Use a debug message here even though it's not
2753 # a good situation. In some cases this can be
2754 # caused by a race between user A and user B
2755 # and neither of them has permissions for the
2756 # other's files. Butler does not know about users
2757 # and trash has no idea what collections these
2758 # files were in (without guessing from a path).
2759 log.debug(
2760 "Encountered %d error%s removing %d artifact%s from datastore %s: %s",
2761 len(exceptions),
2762 s_err,
2763 len(artifacts_to_delete),
2764 "s" if len(artifacts_to_delete) != 1 else "",
2765 self.name,
2766 e,
2767 )
2768 else:
2769 raise e
2770 return removed
2772 @transactional
2773 def transfer_from(
2774 self,
2775 source_records: FileTransferMap,
2776 refs: Collection[DatasetRef],
2777 transfer: str = "auto",
2778 artifact_existence: dict[ResourcePath, bool] | None = None,
2779 dry_run: bool = False,
2780 ) -> tuple[set[DatasetRef], set[DatasetRef]]:
2781 log.verbose("Transferring %d datasets to %s", len(refs), self.name)
2783 # Stop early if "direct" transfer mode is requested. That would
2784 # require that the URI inside the source datastore should be stored
2785 # directly in the target datastore, which seems unlikely to be useful
2786 # since at any moment the source datastore could delete the file.
2787 if transfer in ("direct", "split"):
2788 raise ValueError(
2789 f"Can not transfer from a source datastore using {transfer} mode since"
2790 " those files are controlled by the other datastore."
2791 )
2793 if not refs: 2793 ↛ 2794line 2793 didn't jump to line 2794 because the condition on line 2793 was never true
2794 return set(), set()
2796 # Empty existence lookup if none given.
2797 if artifact_existence is None:
2798 artifact_existence = {}
2800 # In order to handle disassembled composites the code works
2801 # at the records level since it can assume that internal APIs
2802 # can be used.
2803 # - If the record already exists in the destination this is assumed
2804 # to be okay.
2805 # - If there is no record but the source and destination URIs are
2806 # identical no transfer is done but the record is added.
2807 # - If the source record refers to an absolute URI currently assume
2808 # that that URI should remain absolute and will be visible to the
2809 # destination butler. May need to have a flag to indicate whether
2810 # the dataset should be transferred. This will only happen if
2811 # the detached Butler has had a local ingest.
2813 # See if we already have these records
2814 log.verbose("Looking up existing datastore records in target %s for %d refs", self.name, len(refs))
2815 target_records = self._get_stored_records_associated_with_refs(refs, ignore_datastore_records=True)
2817 # The artifacts to register
2818 artifacts = []
2820 # Refs that already exist
2821 already_present = []
2823 # Refs that were rejected by this datastore.
2824 rejected = set()
2826 # Refs that were transferred successfully.
2827 accepted = set()
2829 # Record each time we have done a "direct" transfer.
2830 direct_transfers = []
2832 # Keep track of all the file transfers that are required.
2833 from_to: list[tuple[ResourcePath, ResourcePath]] = []
2835 # Now can transfer the artifacts
2836 log.verbose("Transferring artifacts")
2837 for ref in refs:
2838 if not self.constraints.isAcceptable(ref): 2838 ↛ 2840line 2838 didn't jump to line 2840 because the condition on line 2838 was never true
2839 # This datastore should not be accepting this dataset.
2840 rejected.add(ref)
2841 continue
2843 accepted.add(ref)
2845 if ref.id in target_records:
2846 # Already have an artifact for this.
2847 already_present.append(ref)
2848 continue
2850 # mypy needs to know these are always resolved refs
2851 for transfer_info in source_records.get(ref.id, []):
2852 info = transfer_info.file_info
2853 source_location = transfer_info.location
2854 target_location = info.file_location(self.locationFactory)
2855 if transfer == "unsafe_direct":
2856 # Use the existing file from the source location in place,
2857 # by recording the absolute URI in the target DB. This is
2858 # "unsafe" because the file could be deleted from the
2859 # source Butler at any time, leaving a dangling reference.
2860 source_location = source_location.toAbsolute()
2861 direct_transfers.append(source_location)
2862 info = info.update(path=str(source_location.uri))
2863 elif source_location == target_location and not source_location.pathInStore.isabs(): 2863 ↛ 2866line 2863 didn't jump to line 2866 because the condition on line 2863 was never true
2864 # Artifact is already in the target location.
2865 # (which is how execution butler currently runs)
2866 pass
2867 else:
2868 if target_location.pathInStore.isabs():
2869 # Just because we can see the artifact when running
2870 # the transfer doesn't mean it will be generally
2871 # accessible to a user of this butler. Need to decide
2872 # what to do about an absolute path.
2873 if transfer == "auto":
2874 # For "auto" transfers we allow the absolute URI
2875 # to be recorded in the target datastore.
2876 direct_transfers.append(source_location)
2877 else:
2878 # The user is explicitly requesting a transfer
2879 # even for an absolute URI. This requires us to
2880 # calculate the target path.
2881 template_ref = ref
2882 if info.component: 2882 ↛ 2883line 2882 didn't jump to line 2883 because the condition on line 2882 was never true
2883 template_ref = ref.makeComponentRef(info.component)
2884 target_location = self._calculate_ingested_datastore_name(
2885 source_location.uri,
2886 template_ref,
2887 )
2889 info = info.update(path=target_location.pathInStore.path)
2891 # Need to transfer it to the new location.
2892 from_to.append((source_location.uri, target_location.uri))
2894 artifacts.append((ref, info))
2896 # Do the file transfers in bulk.
2897 # Assume we should always overwrite. If the artifact
2898 # is there this might indicate that a previous transfer
2899 # was interrupted but was not able to be rolled back
2900 # completely (eg pre-emption) so follow Datastore default
2901 # and overwrite. Do not copy if we are in dry-run mode.
2902 if dry_run:
2903 log.info("Would be copying %d file artifacts", len(from_to))
2904 else:
2905 log.verbose("Copying %d file artifacts", len(from_to))
2906 with time_this(log, msg="Transferring datasets into datastore", level=VERBOSE):
2907 ResourcePath.mtransfer(
2908 transfer,
2909 from_to,
2910 overwrite=True,
2911 transaction=self._transaction,
2912 )
2914 if direct_transfers:
2915 log.info(
2916 "Transfer request for an outside-datastore artifact with absolute URI done %d time%s",
2917 len(direct_transfers),
2918 "" if len(direct_transfers) == 1 else "s",
2919 )
2921 # We are overwriting previous datasets that may have already
2922 # existed. We therefore should ensure that we force the
2923 # datastore records to agree. Note that this can potentially lead
2924 # to difficulties if the dataset has previously been ingested
2925 # disassembled and is somehow now assembled, or vice versa.
2926 if not dry_run:
2927 log.verbose("Registering datastore records in database")
2928 self._register_datasets(artifacts, insert_mode=DatabaseInsertMode.REPLACE)
2930 if already_present:
2931 n_skipped = len(already_present)
2932 log.info(
2933 "Skipped transfer of %d dataset%s already present in datastore",
2934 n_skipped,
2935 "" if n_skipped == 1 else "s",
2936 )
2938 log.verbose(
2939 "Finished transfer_from to %s with %d accepted, %d rejected",
2940 self.name,
2941 len(accepted),
2942 len(rejected),
2943 )
2944 return accepted, rejected
2946 def get_file_info_for_transfer(self, dataset_ids: Iterable[DatasetId]) -> FileTransferMap:
2947 source_records = self._get_stored_records_associated_with_refs(
2948 [FakeDatasetRef(id) for id in dataset_ids], ignore_datastore_records=True
2949 )
2950 return self._convert_stored_file_info_to_file_transfer_record(source_records)
2952 def locate_missing_files_for_transfer(
2953 self, refs: Iterable[DatasetRef], artifact_existence: dict[ResourcePath, bool]
2954 ) -> FileTransferMap:
2955 missing_ids = {ref.id for ref in refs}
2956 # Missing IDs can be okay if that datastore has allowed
2957 # gets based on file existence. Should we transfer what we can
2958 # or complain about it and warn?
2959 if not self.trustGetRequest:
2960 return {}
2962 found_records = self._find_missing_records(
2963 refs, missing_ids, artifact_existence, warn_for_missing=False
2964 )
2965 return self._convert_stored_file_info_to_file_transfer_record(found_records)
2967 def _convert_stored_file_info_to_file_transfer_record(
2968 self, info_map: dict[DatasetId, list[StoredFileInfo]]
2969 ) -> FileTransferMap:
2970 output: dict[DatasetId, list[FileTransferRecord]] = {}
2971 for k, file_info_list in info_map.items():
2972 output[k] = [
2973 FileTransferRecord(file_info=info, location=info.file_location(self.locationFactory))
2974 for info in file_info_list
2975 ]
2976 return output
2978 @transactional
2979 def forget(self, refs: Iterable[DatasetRef]) -> None:
2980 # Docstring inherited.
2981 refs = list(refs)
2982 self.bridge.forget(refs)
2983 self._table.delete(["dataset_id"], *[{"dataset_id": ref.id} for ref in refs])
2985 def validateConfiguration(
2986 self, entities: Iterable[DatasetRef | DatasetType | StorageClass], logFailures: bool = False
2987 ) -> None:
2988 """Validate some of the configuration for this datastore.
2990 Parameters
2991 ----------
2992 entities : `~collections.abc.Iterable` [`DatasetRef` | `DatasetType` \
2993 | `StorageClass`]
2994 Entities to test against this configuration. Can be differing
2995 types.
2996 logFailures : `bool`, optional
2997 If `True`, output a log message for every validation error
2998 detected.
3000 Returns
3001 -------
3002 None
3004 Raises
3005 ------
3006 DatastoreValidationError
3007 Raised if there is a validation problem with a configuration.
3008 All the problems are reported in a single exception.
3010 Notes
3011 -----
3012 This method checks that all the supplied entities have valid file
3013 templates and also have formatters defined.
3014 """
3015 templateFailed = None
3016 try:
3017 self.templates.validateTemplates(entities, logFailures=logFailures)
3018 except FileTemplateValidationError as e:
3019 templateFailed = str(e)
3021 formatterFailed = []
3022 for entity in entities:
3023 try:
3024 self.formatterFactory.getFormatterClass(entity)
3025 except KeyError as e:
3026 formatterFailed.append(str(e))
3027 if logFailures: 3027 ↛ 3022line 3027 didn't jump to line 3022 because the condition on line 3027 was always true
3028 log.critical("Formatter failure: %s", e)
3030 if templateFailed or formatterFailed:
3031 messages = []
3032 if templateFailed: 3032 ↛ 3033line 3032 didn't jump to line 3033 because the condition on line 3032 was never true
3033 messages.append(templateFailed)
3034 if formatterFailed: 3034 ↛ 3036line 3034 didn't jump to line 3036 because the condition on line 3034 was always true
3035 messages.append(",".join(formatterFailed))
3036 msg = ";\n".join(messages)
3037 raise DatastoreValidationError(msg)
3039 def getLookupKeys(self) -> set[LookupKey]:
3040 # Docstring is inherited from base class
3041 return (
3042 self.templates.getLookupKeys()
3043 | self.formatterFactory.getLookupKeys()
3044 | self.constraints.getLookupKeys()
3045 )
3047 def validateKey(self, lookupKey: LookupKey, entity: DatasetRef | DatasetType | StorageClass) -> None:
3048 # Docstring is inherited from base class
3049 # The key can be valid in either formatters or templates so we can
3050 # only check the template if it exists
3051 if lookupKey in self.templates:
3052 try:
3053 self.templates[lookupKey].validateTemplate(entity)
3054 except FileTemplateValidationError as e:
3055 raise DatastoreValidationError(e) from e
3057 def export(
3058 self,
3059 refs: Iterable[DatasetRef],
3060 *,
3061 directory: ResourcePathExpression | None = None,
3062 transfer: str | None = "auto",
3063 ) -> Iterable[FileDataset]:
3064 # Docstring inherited from Datastore.export.
3065 if transfer == "auto" and directory is None:
3066 transfer = None
3068 if transfer is not None and transfer != "direct" and directory is None:
3069 raise TypeError(f"Cannot export using transfer mode {transfer} with no export directory given")
3071 if transfer == "move":
3072 raise TypeError("Can not export by moving files out of datastore.")
3074 # Force the directory to be a URI object
3075 directoryUri: ResourcePath | None = None
3076 if directory is not None:
3077 directoryUri = ResourcePath(directory, forceDirectory=True)
3079 if transfer is not None and directoryUri is not None and not directoryUri.exists(): 3079 ↛ 3081line 3079 didn't jump to line 3081 because the condition on line 3079 was never true
3080 # mypy needs the second test
3081 raise FileNotFoundError(f"Export location {directory} does not exist")
3083 progress = Progress("lsst.daf.butler.datastores.FileDatastore.export", level=logging.DEBUG)
3084 for ref in progress.wrap(refs, "Exporting dataset files"):
3085 fileLocations = self._get_dataset_locations_info(ref)
3086 if not fileLocations:
3087 raise FileNotFoundError(f"Could not retrieve dataset {ref}.")
3088 # For now we can not export disassembled datasets
3089 if len(fileLocations) > 1:
3090 raise NotImplementedError(f"Can not export disassembled datasets such as {ref}")
3091 location, storedFileInfo = fileLocations[0]
3093 pathInStore = location.pathInStore.path
3094 if transfer is None:
3095 # TODO: do we also need to return the readStorageClass somehow?
3096 # We will use the path in store directly. If this is an
3097 # absolute URI, preserve it.
3098 if location.pathInStore.isabs(): 3098 ↛ 3099line 3098 didn't jump to line 3099 because the condition on line 3098 was never true
3099 pathInStore = str(location.uri)
3100 elif transfer == "direct":
3101 # Use full URIs to the remote store in the export
3102 pathInStore = str(location.uri)
3103 else:
3104 # mypy needs help
3105 assert directoryUri is not None, "directoryUri must be defined to get here"
3106 storeUri = ResourcePath(location.uri, forceDirectory=False)
3108 # if the datastore has an absolute URI to a resource, we
3109 # have two options:
3110 # 1. Keep the absolute URI in the exported YAML
3111 # 2. Allocate a new name in the local datastore and transfer
3112 # it.
3113 # For now go with option 2
3114 if location.pathInStore.isabs(): 3114 ↛ 3115line 3114 didn't jump to line 3115 because the condition on line 3114 was never true
3115 template = self.templates.getTemplate(ref)
3116 newURI = ResourcePath(template.format(ref), forceAbsolute=False, forceDirectory=False)
3117 pathInStore = str(newURI.updatedExtension(location.pathInStore.getExtension()))
3119 exportUri = directoryUri.join(pathInStore)
3120 exportUri.transfer_from(storeUri, transfer=transfer)
3122 yield FileDataset(refs=[ref], path=pathInStore, formatter=storedFileInfo.formatter)
3124 @staticmethod
3125 def computeChecksum(uri: ResourcePath, algorithm: str = "blake2b", block_size: int = 8192) -> str | None:
3126 """Compute the checksum of the supplied file.
3128 Parameters
3129 ----------
3130 uri : `lsst.resources.ResourcePath`
3131 Name of resource to calculate checksum from.
3132 algorithm : `str`, optional
3133 Name of algorithm to use. Must be one of the algorithms supported
3134 by :py:class`hashlib`.
3135 block_size : `int`
3136 Number of bytes to read from file at one time.
3138 Returns
3139 -------
3140 hexdigest : `str`
3141 Hex digest of the file.
3143 Notes
3144 -----
3145 Currently returns None if the URI is for a remote resource.
3146 """
3147 if algorithm not in hashlib.algorithms_guaranteed: 3147 ↛ 3148line 3147 didn't jump to line 3148 because the condition on line 3147 was never true
3148 raise NameError(f"The specified algorithm '{algorithm}' is not supported by hashlib")
3150 if not uri.isLocal: 3150 ↛ 3151line 3150 didn't jump to line 3151 because the condition on line 3150 was never true
3151 return None
3153 hasher = hashlib.new(algorithm)
3155 with uri.as_local() as local_uri, open(local_uri.ospath, "rb") as f:
3156 for chunk in iter(lambda: f.read(block_size), b""):
3157 hasher.update(chunk)
3159 return hasher.hexdigest()
3161 def needs_expanded_data_ids(
3162 self,
3163 transfer: str | None,
3164 entity: DatasetRef | DatasetType | StorageClass | None = None,
3165 ) -> bool:
3166 # Docstring inherited.
3167 # This _could_ also use entity to inspect whether the filename template
3168 # involves placeholders other than the required dimensions for its
3169 # dataset type, but that's not necessary for correctness; it just
3170 # enables more optimizations (perhaps only in theory).
3171 return transfer not in ("direct", None)
3173 def import_records(self, data: Mapping[str, DatastoreRecordData]) -> None:
3174 # Docstring inherited from the base class.
3175 record_data = data.get(self.name)
3176 if not record_data: 3176 ↛ 3177line 3176 didn't jump to line 3177 because the condition on line 3176 was never true
3177 return
3179 self._bridge.insert(FakeDatasetRef(dataset_id) for dataset_id in record_data.records)
3181 # TODO: Verify that there are no unexpected table names in the dict?
3182 unpacked_records = []
3183 for dataset_id, dataset_data in record_data.records.items():
3184 records = dataset_data.get(self._table.name)
3185 if records: 3185 ↛ 3183line 3185 didn't jump to line 3183 because the condition on line 3185 was always true
3186 for info in records:
3187 assert isinstance(info, StoredFileInfo), "Expecting StoredFileInfo records"
3188 unpacked_records.append(info.to_record(dataset_id=dataset_id))
3189 if unpacked_records:
3190 self._table.insert(*unpacked_records, transaction=self._transaction)
3192 def export_records(self, refs: Iterable[DatasetIdRef]) -> Mapping[str, DatastoreRecordData]:
3193 # Docstring inherited from the base class.
3195 records: dict[DatasetId, dict[str, list[StoredDatastoreItemInfo]]] = {}
3196 for batch in self._export_rows([ref.id for ref in refs]):
3197 for row in batch:
3198 info: StoredDatastoreItemInfo = StoredFileInfo.from_record(row)
3199 dataset_records = records.setdefault(row["dataset_id"], {})
3200 dataset_records.setdefault(self._table.name, []).append(info)
3202 record_data = DatastoreRecordData(records=records)
3203 return {self.name: record_data}
3205 def _export_rows(self, datasets: Collection[DatasetId]) -> Iterator[Sequence[Mapping[str, Any]]]:
3206 # This call to 'bridge.check' filters out "partially deleted" datasets.
3207 # Specifically, ones in the unusual edge state that:
3208 # 1. They have an entry in the registry dataset tables
3209 # 2. They were "trashed" from the datastore, so they are not
3210 # present in the "dataset_location" table.)
3211 # 3. But the trash has not been "emptied", so there are still entries
3212 # in the "opaque" datastore records table.
3213 #
3214 # As far as I can tell, this can only occur in the case of a concurrent
3215 # or aborted call to `Butler.pruneDatasets(unstore=True, purge=False)`.
3216 # Datasets (with or without files existing on disk) can persist in
3217 # this zombie state indefinitely, until someone manually empties
3218 # the trash.
3219 found_ids = self._bridge.check(datasets)
3220 return self._table.fetch_batches(dataset_id=found_ids)
3222 def export_table(self, datasets: Collection[DatasetId]) -> DatastoreRecordTable:
3223 # Docstring inherited from the base class.
3225 tables: list[DatastoreRecordTable] = []
3226 for batch in self._export_rows(datasets):
3227 file_info = StoredFileInfoTable.from_records(batch)
3228 tables.append(DatastoreRecordTable.from_stored_file_info_table(self.name, file_info))
3229 return DatastoreRecordTable.combine(tables)
3231 def import_table(self, table: DatastoreRecordTable) -> None:
3232 # Docstring inherited from the base class.
3234 records = table.to_stored_file_info_table().to_records()
3235 dataset_ids = [FakeDatasetRef(record["dataset_id"]) for record in records]
3236 if len(records) > 0: 3236 ↛ exitline 3236 didn't return from function 'import_table' because the condition on line 3236 was always true
3237 self._bridge.insert(dataset_ids)
3238 self._table.insert(*records, transaction=self._transaction)
3240 def export_predicted_records(self, refs: Iterable[DatasetRef]) -> dict[str, DatastoreRecordData]:
3241 # Docstring inherited from the base class.
3242 refs = [self._cast_storage_class(ref) for ref in refs]
3243 records: dict[DatasetId, dict[str, list[StoredDatastoreItemInfo]]] = {}
3244 for ref in refs:
3245 if not self.constraints.isAcceptable(ref): 3245 ↛ 3246line 3245 didn't jump to line 3246 because the condition on line 3245 was never true
3246 continue
3247 fileLocations = self._get_expected_dataset_locations_info(ref)
3248 if not fileLocations: 3248 ↛ 3249line 3248 didn't jump to line 3249 because the condition on line 3248 was never true
3249 continue
3250 dataset_records = records.setdefault(ref.id, {})
3251 dataset_records.setdefault(self._table.name, [])
3252 for _, storedFileInfo in fileLocations:
3253 dataset_records[self._table.name].append(storedFileInfo)
3255 record_data = DatastoreRecordData(records=records)
3256 return {self.name: record_data}
3258 def set_retrieve_dataset_type_method(self, method: Callable[[str], DatasetType | None] | None) -> None:
3259 # Docstring inherited from the base class.
3260 self._retrieve_dataset_method = method
3262 def _cast_storage_class(self, ref: DatasetRef) -> DatasetRef:
3263 """Update dataset reference to use the storage class from registry."""
3264 if self._retrieve_dataset_method is None:
3265 # We could raise an exception here but unit tests do not define
3266 # this method.
3267 return ref
3268 dataset_type = self._retrieve_dataset_method(ref.datasetType.name)
3269 if dataset_type is not None: 3269 ↛ 3271line 3269 didn't jump to line 3271 because the condition on line 3269 was always true
3270 ref = ref.overrideStorageClass(dataset_type.storageClass_name)
3271 return ref
3273 def get_opaque_table_definitions(self) -> Mapping[str, DatastoreOpaqueTable]:
3274 # Docstring inherited from the base class.
3275 return {self._opaque_table_name: DatastoreOpaqueTable(self.makeTableSpec(), StoredFileInfo)}