Coverage for python/lsst/daf/butler/datastores/file_datastore/get.py: 91%

141 statements  

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

1# This file is part of daf_butler. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

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

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

7# for details of code ownership. 

8# 

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

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

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

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

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

14# 

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

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

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

18# (at your option) any later version. 

19# 

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

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

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

23# GNU General Public License for more details. 

24# 

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

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

27 

28__all__ = ( 

29 "DatasetLocationInformation", 

30 "DatastoreFileGetInformation", 

31 "generate_datastore_get_information", 

32 "get_dataset_as_python_object_from_get_info", 

33) 

34 

35from collections.abc import Mapping 

36from dataclasses import dataclass 

37from typing import Any, TypeAlias 

38 

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 

54 

55log = getLogger(__name__) 

56 

57DatasetLocationInformation: TypeAlias = tuple[Location, StoredFileInfo] 

58 

59 

60@dataclass(frozen=True) 

61class DatastoreFileGetInformation: 

62 """Collection of useful parameters needed to retrieve a file from 

63 a Datastore. 

64 """ 

65 

66 location: Location 

67 """The location from which to read the dataset.""" 

68 

69 formatter: Formatter | FormatterV2 

70 """The `Formatter` to use to deserialize the dataset.""" 

71 

72 info: StoredFileInfo 

73 """Stored information about this file and its formatter.""" 

74 

75 assemblerParams: Mapping[str, Any] 

76 """Parameters to use for post-processing the retrieved dataset.""" 

77 

78 formatterParams: Mapping[str, Any] 

79 """Parameters that were understood by the associated formatter.""" 

80 

81 component: str | None 

82 """The component to be retrieved (can be `None`).""" 

83 

84 readStorageClass: StorageClass 

85 """The `StorageClass` that the `Formatter` will return.""" 

86 

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. 

91 

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 """ 

98 

99 

100def _describe_components(storageClass: StorageClass) -> str: 

101 """Return a description of the components a `StorageClass` recognizes, 

102 suitable for inclusion in a log message. 

103 

104 Parameters 

105 ---------- 

106 storageClass : `StorageClass` 

107 Storage class to describe. 

108 

109 Returns 

110 ------- 

111 description : `str` 

112 Comma-separated list of component names, or ``"none"``. 

113 """ 

114 return ", ".join(sorted(storageClass.allComponents())) or "none" 

115 

116 

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. 

124 

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) 

165 

166 

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. 

176 

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. 

193 

194 Returns 

195 ------- 

196 getInfo : `list` [`DatastoreFileGetInformation`] 

197 The parameters needed to retrieve each file. 

198 

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 

207 

208 # Is this a component request? 

209 refComponent = read_ref.datasetType.component() 

210 

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 

214 

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 

222 

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 ) 

242 

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 ) 

255 

256 formatterParams, notFormatterParams = formatter.segregate_parameters() 

257 

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) 

261 

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 

266 

267 fileGetInfo.append( 

268 DatastoreFileGetInformation( 

269 location, 

270 formatter, 

271 storedFileInfo, 

272 assemblerParams, 

273 formatterParams, 

274 component, 

275 thisReadStorageClass, 

276 componentStorageClass, 

277 ) 

278 ) 

279 

280 return fileGetInfo 

281 

282 

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. 

290 

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 

301 

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) 

310 

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 

315 

316 formatter = getInfo.formatter 

317 

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 ) 

326 

327 assert isinstance(formatter, FormatterV2) 

328 

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 

351 

352 return post_process_get( 

353 result, ref.datasetType.storageClass, getInfo.assemblerParams, isComponent=isComponent 

354 ) 

355 

356 

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. 

365 

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. 

377 

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} 

388 

389 # By definition the dataset is disassembled if we have more 

390 # than one record for it. 

391 isDisassembled = len(allGetInfo) > 1 

392 

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 

403 

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 

411 

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 

416 

417 # Check that the supplied parameters are suitable for the type read 

418 refStorageClass.validateParameters(parameters) 

419 

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() 

425 

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)) 

431 

432 component = getInfo.component 

433 

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}") 

436 

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 ) 

444 

445 inMemoryDataset = ref.datasetType.storageClass.delegate().assemble(components) 

446 

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 = {} 

452 

453 # Process parameters 

454 return ref.datasetType.storageClass.delegate().handleParameters( 

455 inMemoryDataset, parameters=unusedParams 

456 ) 

457 

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 ) 

465 

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") 

469 

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)) 

478 

479 # Select the relevant component 

480 rwInfo = allComponents[forwardedComponent] 

481 

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) 

486 

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 

491 

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 ) 

505 

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] = {} 

512 

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 ) 

524 

525 return _read_artifact_into_memory(readInfo, ref, cache_manager, isComponent=True) 

526 

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") 

535 

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 

541 

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) 

552 

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)) 

565 

566 return _read_artifact_into_memory(getInfo, ref, cache_manager, isComponent=isComponent)