Coverage for python/lsst/daf/butler/datastores/file_datastore/get.py: 91%
141 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-27 09:25 +0000
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-27 09:25 +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__all__ = (
29 "DatasetLocationInformation",
30 "DatastoreFileGetInformation",
31 "generate_datastore_get_information",
32 "get_dataset_as_python_object_from_get_info",
33)
35from collections.abc import Mapping
36from dataclasses import dataclass
37from typing import Any, TypeAlias
39from lsst.daf.butler import (
40 DatasetRef,
41 FileDescriptor,
42 FileIntegrityError,
43 Formatter,
44 FormatterV1inV2,
45 FormatterV2,
46 Location,
47 StorageClass,
48)
49from lsst.daf.butler.datastore.cache_manager import AbstractDatastoreCacheManager
50from lsst.daf.butler.datastore.generic_base import post_process_get
51from lsst.daf.butler.datastore.stored_file_info import StoredFileInfo
52from lsst.utils.introspection import get_instance_of
53from lsst.utils.logging import getLogger
55log = getLogger(__name__)
57DatasetLocationInformation: TypeAlias = tuple[Location, StoredFileInfo]
60@dataclass(frozen=True)
61class DatastoreFileGetInformation:
62 """Collection of useful parameters needed to retrieve a file from
63 a Datastore.
64 """
66 location: Location
67 """The location from which to read the dataset."""
69 formatter: Formatter | FormatterV2
70 """The `Formatter` to use to deserialize the dataset."""
72 info: StoredFileInfo
73 """Stored information about this file and its formatter."""
75 assemblerParams: Mapping[str, Any]
76 """Parameters to use for post-processing the retrieved dataset."""
78 formatterParams: Mapping[str, Any]
79 """Parameters that were understood by the associated formatter."""
81 component: str | None
82 """The component to be retrieved (can be `None`)."""
84 readStorageClass: StorageClass
85 """The `StorageClass` that the `Formatter` will return."""
87 componentStorageClass: StorageClass | None = None
88 """The `StorageClass` of the component to extract from the object returned
89 by the `Formatter`, or `None` if the `Formatter` returns the requested
90 object directly.
92 This is only set when the requested component is defined by the read
93 `StorageClass` but not by the `StorageClass` the file was written with.
94 The `Formatter` then knows nothing of the component, so the composite is
95 read and converted first (to ``readStorageClass``) and the component is
96 extracted from the converted form.
97 """
100def _describe_components(storageClass: StorageClass) -> str:
101 """Return a description of the components a `StorageClass` recognizes,
102 suitable for inclusion in a log message.
104 Parameters
105 ----------
106 storageClass : `StorageClass`
107 Storage class to describe.
109 Returns
110 -------
111 description : `str`
112 Comma-separated list of component names, or ``"none"``.
113 """
114 return ", ".join(sorted(storageClass.allComponents())) or "none"
117def _warn_about_converted_component_read(
118 registry_ref: DatasetRef,
119 component: str,
120 writeStorageClass: StorageClass,
121 parentReadStorageClass: StorageClass,
122) -> None:
123 """Warn that a component request needs the whole dataset to be read.
125 Parameters
126 ----------
127 registry_ref : `DatasetRef`
128 The dataset as defined in the repository, used to identify the dataset
129 in the message.
130 component : `str`
131 Name of the component that was requested.
132 writeStorageClass : `StorageClass`
133 The `StorageClass` the dataset was written with, which does not define
134 ``component``.
135 parentReadStorageClass : `StorageClass`
136 The `StorageClass` the composite has to be converted to in order to
137 obtain ``component``.
138 """
139 message = (
140 "Component %r was requested from dataset %s but storage class %s, which the dataset was "
141 "written with, does not define it (components it does define: %s). The entire dataset must "
142 "therefore be retrieved and converted to storage class %s before %r can be extracted from "
143 "the result. This is slower than reading the component on its own: the storage class "
144 "override has made this a less efficient request."
145 )
146 args: list[Any] = [
147 component,
148 registry_ref,
149 writeStorageClass.name,
150 _describe_components(writeStorageClass),
151 parentReadStorageClass.name,
152 component,
153 ]
154 registryStorageClass = registry_ref.datasetType.storageClass
155 if registryStorageClass != writeStorageClass: 155 ↛ 159line 155 didn't jump to line 159 because the condition on line 155 was never true
156 # The dataset type definition has changed since the dataset was
157 # written, so the caller may be surprised that the component they were
158 # offered is not available from the file itself.
159 message += (
160 " The dataset type is now defined with storage class %s (components: %s), which is why "
161 "the component appeared to be available."
162 )
163 args += [registryStorageClass.name, _describe_components(registryStorageClass)]
164 log.warning(message, *args)
167def generate_datastore_get_information(
168 fileLocations: list[DatasetLocationInformation],
169 *,
170 registry_ref: DatasetRef,
171 read_ref: DatasetRef,
172 parameters: Mapping[str, Any] | None,
173) -> list[DatastoreFileGetInformation]:
174 """Process parameters and instantiate formatters for in preparation for
175 retrieving an artifact and converting it to a Python object.
177 Parameters
178 ----------
179 fileLocations : `list` [`DatasetLocationInformation`]
180 List of file locations for this artifact and their associated datastore
181 records.
182 registry_ref : `DatasetRef`
183 The registry information associated with this artifact, using the
184 dataset type definition from the repository and never naming a
185 component. This is the ref given to the `Formatter`.
186 read_ref : `DatasetRef`
187 The dataset the caller asked for. Its `StorageClass` is the one to
188 return and its component (if any) is the component to extract. This
189 differs from ``registry_ref`` when a read-time `StorageClass` override
190 has been requested or a component has been requested.
191 parameters : `~collections.abc.Mapping` [`str`, `typing.Any`]
192 `StorageClass` and `Formatter` parameters.
194 Returns
195 -------
196 getInfo : `list` [`DatastoreFileGetInformation`]
197 The parameters needed to retrieve each file.
199 Notes
200 -----
201 The `StorageClass` that each file was written with is taken from the
202 datastore records in ``fileLocations`` and can differ from the one in
203 ``registry_ref`` if the dataset type definition has been changed since the
204 dataset was written.
205 """
206 readStorageClass = read_ref.datasetType.storageClass
208 # Is this a component request?
209 refComponent = read_ref.datasetType.component()
211 # The storage class of the composite that the caller wants to read. This
212 # differs from the write storage class when there is an override.
213 parentReadStorageClass = read_ref.datasetType.parentStorageClass
215 disassembled = len(fileLocations) > 1
216 fileGetInfo = []
217 for location, storedFileInfo in fileLocations:
218 # The storage class used to write the file
219 writeStorageClass = storedFileInfo.storageClass
220 thisReadStorageClass = readStorageClass
221 componentStorageClass = None
223 # If this has been disassembled we need read to match the write
224 # except for if a component has specified an override.
225 if disassembled and storedFileInfo.component != refComponent:
226 thisReadStorageClass = writeStorageClass
227 elif (
228 refComponent is not None
229 and storedFileInfo.component is None
230 and parentReadStorageClass is not None
231 and refComponent not in writeStorageClass.allComponents()
232 ):
233 # The component only exists in the storage class the caller is
234 # converting to, so the formatter can not extract it. Read the
235 # whole composite as the converted type instead and extract the
236 # component from that.
237 componentStorageClass = readStorageClass
238 thisReadStorageClass = parentReadStorageClass
239 _warn_about_converted_component_read(
240 registry_ref, refComponent, writeStorageClass, parentReadStorageClass
241 )
243 formatter = get_instance_of(
244 storedFileInfo.formatter,
245 FileDescriptor(
246 location,
247 readStorageClass=thisReadStorageClass,
248 storageClass=writeStorageClass,
249 parameters=parameters,
250 component=storedFileInfo.component,
251 ),
252 dataId=registry_ref.dataId,
253 ref=registry_ref,
254 )
256 formatterParams, notFormatterParams = formatter.segregate_parameters()
258 # Of the remaining parameters, extract the ones supported by
259 # this StorageClass (for components not all will be handled)
260 assemblerParams = thisReadStorageClass.filterParameters(notFormatterParams)
262 # The ref itself could be a component if the dataset was
263 # disassembled by butler, or we disassembled in datastore and
264 # components came from the datastore records
265 component = storedFileInfo.component if storedFileInfo.component else refComponent
267 fileGetInfo.append(
268 DatastoreFileGetInformation(
269 location,
270 formatter,
271 storedFileInfo,
272 assemblerParams,
273 formatterParams,
274 component,
275 thisReadStorageClass,
276 componentStorageClass,
277 )
278 )
280 return fileGetInfo
283def _read_artifact_into_memory(
284 getInfo: DatastoreFileGetInformation,
285 ref: DatasetRef,
286 cache_manager: AbstractDatastoreCacheManager,
287 isComponent: bool = False,
288) -> Any:
289 """Read the artifact from datastore into in memory object.
291 Parameters
292 ----------
293 getInfo : `DatastoreFileGetInformation`
294 Information about the artifact within the datastore.
295 ref : `DatasetRef`
296 The registry information associated with this artifact.
297 isComponent : `bool`
298 Flag to indicate if a component is being read from this artifact.
299 cache_manager : `AbstractDatastoreCacheManager`
300 The cache manager to use for caching retrieved files
302 Returns
303 -------
304 inMemoryDataset : `object`
305 The artifact as a python object.
306 """
307 location = getInfo.location
308 uri = location.uri
309 log.debug("Accessing data from %s", uri)
311 # Cannot recalculate checksum but can compare size as a quick check
312 # Do not do this if the size is negative since that indicates
313 # we do not know.
314 recorded_size = getInfo.info.file_size
316 formatter = getInfo.formatter
318 if isinstance(formatter, Formatter):
319 formatter = FormatterV1inV2(
320 formatter.file_descriptor,
321 ref=ref,
322 formatter=formatter,
323 write_parameters=formatter.write_parameters,
324 write_recipes=formatter.write_recipes,
325 )
327 assert isinstance(formatter, FormatterV2)
329 try:
330 result = formatter.read(
331 component=getInfo.component if isComponent else None,
332 expected_size=recorded_size,
333 cache_manager=cache_manager,
334 )
335 except (FileNotFoundError, FileIntegrityError):
336 # This is expected for the case where the resource is missing
337 # or the information we passed to the formatter about the file size
338 # is incorrect.
339 # Allow them to propagate up.
340 raise
341 except Exception as e:
342 # For clarity, include any notes that may have been added by the
343 # formatter to this new exception.
344 notes = "\n".join(getattr(e, "__notes__", []))
345 if notes: 345 ↛ 346line 345 didn't jump to line 346 because the condition on line 345 was never true
346 notes = "\n" + notes
347 raise ValueError(
348 f"Failure from formatter '{formatter.name()}' for dataset {ref.id}"
349 f" ({ref.datasetType.name} from {uri}): {e}{notes}"
350 ) from e
352 return post_process_get(
353 result, ref.datasetType.storageClass, getInfo.assemblerParams, isComponent=isComponent
354 )
357def get_dataset_as_python_object_from_get_info(
358 allGetInfo: list[DatastoreFileGetInformation],
359 *,
360 ref: DatasetRef,
361 parameters: Mapping[str, Any] | None,
362 cache_manager: AbstractDatastoreCacheManager,
363) -> Any:
364 """Retrieve an artifact from storage and return it as a Python object.
366 Parameters
367 ----------
368 allGetInfo : `list` [`DatastoreFileGetInformation`]
369 Pre-processed information about each file associated with this
370 artifact.
371 ref : `DatasetRef`
372 The registry information associated with this artifact.
373 parameters : `~collections.abc.Mapping` [`str`, `typing.Any`]
374 `StorageClass` and `Formatter` parameters.
375 cache_manager : `AbstractDatastoreCacheManager`
376 The cache manager to use for caching retrieved files.
378 Returns
379 -------
380 python_object : `typing.Any`
381 The retrieved artifact, converted to a Python object according to the
382 `StorageClass` specified in ``ref``.
383 """
384 refStorageClass = ref.datasetType.storageClass
385 refComponent = ref.datasetType.component()
386 # Create mapping from component name to related info
387 allComponents = {i.component: i for i in allGetInfo}
389 # By definition the dataset is disassembled if we have more
390 # than one record for it.
391 isDisassembled = len(allGetInfo) > 1
393 # Look for the special case where we are disassembled but the
394 # component is a derived component that was not written during
395 # disassembly. For this scenario we need to check that the
396 # component requested is listed as a derived component for the
397 # composite storage class
398 isDisassembledReadOnlyComponent = False
399 if isDisassembled and refComponent:
400 # The composite storage class should be accessible through
401 # the component dataset type
402 compositeStorageClass = ref.datasetType.parentStorageClass
404 # In the unlikely scenario where the composite storage
405 # class is not known, we can only assume that this is a
406 # normal component. If that assumption is wrong then the
407 # branch below that reads a persisted component will fail
408 # so there is no need to complain here.
409 if compositeStorageClass is not None: 409 ↛ 412line 409 didn't jump to line 412 because the condition on line 409 was always true
410 isDisassembledReadOnlyComponent = refComponent in compositeStorageClass.derivedComponents
412 if isDisassembled and not refComponent:
413 # This was a disassembled dataset spread over multiple files
414 # and we need to put them all back together again.
415 # Read into memory and then assemble
417 # Check that the supplied parameters are suitable for the type read
418 refStorageClass.validateParameters(parameters)
420 # We want to keep track of all the parameters that were not used
421 # by formatters. We assume that if any of the component formatters
422 # use a parameter that we do not need to apply it again in the
423 # assembler.
424 usedParams = set()
426 components: dict[str, Any] = {}
427 for getInfo in allGetInfo:
428 # assemblerParams are parameters not understood by the
429 # associated formatter.
430 usedParams.update(set(getInfo.formatterParams))
432 component = getInfo.component
434 if component is None: 434 ↛ 435line 434 didn't jump to line 435 because the condition on line 434 was never true
435 raise RuntimeError(f"Internal error in datastore assembly of {ref}")
437 # We do not want the formatter to think it's reading
438 # a component though because it is really reading a
439 # standalone dataset -- always tell reader it is not a
440 # component.
441 components[component] = _read_artifact_into_memory(
442 getInfo, ref.makeComponentRef(component), cache_manager, isComponent=False
443 )
445 inMemoryDataset = ref.datasetType.storageClass.delegate().assemble(components)
447 # Any unused parameters will have to be passed to the assembler
448 if parameters:
449 unusedParams = {k: v for k, v in parameters.items() if k not in usedParams}
450 else:
451 unusedParams = {}
453 # Process parameters
454 return ref.datasetType.storageClass.delegate().handleParameters(
455 inMemoryDataset, parameters=unusedParams
456 )
458 elif isDisassembledReadOnlyComponent:
459 compositeStorageClass = ref.datasetType.parentStorageClass
460 if compositeStorageClass is None: 460 ↛ 461line 460 didn't jump to line 461 because the condition on line 460 was never true
461 raise RuntimeError(
462 f"Unable to retrieve derived component '{refComponent}' since"
463 "no composite storage class is available."
464 )
466 if refComponent is None: 466 ↛ 468line 466 didn't jump to line 468 because the condition on line 466 was never true
467 # Mainly for mypy
468 raise RuntimeError("Internal error in datastore: component can not be None here")
470 # Assume that every derived component can be calculated by
471 # forwarding the request to a single read/write component.
472 # Rather than guessing which rw component is the right one by
473 # scanning each for a derived component of the same name,
474 # we ask the storage class delegate directly which one is best to
475 # use.
476 compositeDelegate = compositeStorageClass.delegate()
477 forwardedComponent = compositeDelegate.selectResponsibleComponent(refComponent, set(allComponents))
479 # Select the relevant component
480 rwInfo = allComponents[forwardedComponent]
482 # For now assume that read parameters are validated against
483 # the real component and not the requested component
484 forwardedStorageClass = rwInfo.formatter.file_descriptor.readStorageClass
485 forwardedStorageClass.validateParameters(parameters)
487 # Unfortunately the FileDescriptor inside the formatter will have
488 # the wrong write storage class so we need to create a new one
489 # given the immutability constraint.
490 writeStorageClass = rwInfo.info.storageClass
492 # We may need to put some thought into parameters for read
493 # components but for now forward them on as is
494 readFormatter = type(rwInfo.formatter)(
495 FileDescriptor(
496 rwInfo.location,
497 readStorageClass=refStorageClass,
498 storageClass=writeStorageClass,
499 parameters=parameters,
500 component=forwardedComponent,
501 ),
502 dataId=ref.dataId,
503 ref=ref,
504 )
506 # The assembler can not receive any parameter requests for a
507 # derived component at this time since the assembler will
508 # see the storage class of the derived component and those
509 # parameters will have to be handled by the formatter on the
510 # forwarded storage class.
511 assemblerParams: dict[str, Any] = {}
513 # Need to created a new info that specifies the derived
514 # component and associated storage class
515 readInfo = DatastoreFileGetInformation(
516 rwInfo.location,
517 readFormatter,
518 rwInfo.info,
519 assemblerParams,
520 {},
521 refComponent,
522 refStorageClass,
523 )
525 return _read_artifact_into_memory(readInfo, ref, cache_manager, isComponent=True)
527 else:
528 # Single file request or component from that composite file
529 for lookup in (refComponent, None): 529 ↛ 534line 529 didn't jump to line 534 because the loop on line 529 didn't complete
530 if lookup in allComponents: 530 ↛ 529line 530 didn't jump to line 529 because the condition on line 530 was always true
531 getInfo = allComponents[lookup]
532 break
533 else:
534 raise FileNotFoundError(f"Component {refComponent} not found for ref {ref} in datastore")
536 # Do not need the component itself if already disassembled
537 if isDisassembled:
538 isComponent = False
539 else:
540 isComponent = getInfo.component is not None
542 # For a disassembled component we can validate parameters against
543 # the component storage class directly
544 if isDisassembled:
545 refStorageClass.validateParameters(parameters)
546 else:
547 # For an assembled composite this could be a derived
548 # component derived from a real component. The validity
549 # of the parameters is not clear. For now validate against
550 # the composite storage class
551 getInfo.formatter.file_descriptor.storageClass.validateParameters(parameters)
553 if getInfo.componentStorageClass is not None:
554 # The formatter knows nothing of this component because it is
555 # defined by the read storage class and not by the storage class
556 # the file was written with. Read the whole composite, letting the
557 # formatter convert it to the requested composite type, and then
558 # extract the component from the result.
559 if refComponent is None: 559 ↛ 560line 559 didn't jump to line 560 because the condition on line 559 was never true
560 raise RuntimeError("Internal error in datastore: component can not be None here")
561 compositeRef = ref.makeCompositeRef()
562 composite = _read_artifact_into_memory(getInfo, compositeRef, cache_manager, isComponent=False)
563 delegate = compositeRef.datasetType.storageClass.delegate()
564 return getInfo.componentStorageClass.coerce_type(delegate.getComponent(composite, refComponent))
566 return _read_artifact_into_memory(getInfo, ref, cache_manager, isComponent=isComponent)