Coverage for python/lsst/daf/butler/datastores/fileDatastore.py: 85%

1076 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-27 09:06 +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"""Generic file-based datastore code.""" 

29 

30from __future__ import annotations 

31 

32__all__ = ("FileDatastore",) 

33 

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 

41 

42from sqlalchemy import BigInteger, String 

43 

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 

107 

108# For VERBOSE logging usage. 

109from lsst.utils.logging import VERBOSE, getLogger 

110from lsst.utils.timer import time_this 

111 

112from ..datastore import FileTransferMap, FileTransferRecord 

113 

114if TYPE_CHECKING: 

115 from lsst.daf.butler import DatasetProvenance, LookupKey 

116 from lsst.daf.butler.registry.interfaces import DatasetIdRef, DatastoreRegistryBridgeManager 

117 

118log = getLogger(__name__) 

119 

120 

121class _IngestPrepData(Datastore.IngestPrepData): 

122 """Helper class for FileDatastore ingest implementation. 

123 

124 Parameters 

125 ---------- 

126 datasets : `~collections.abc.Iterable` of `FileDataset` 

127 Files to be ingested by this datastore. 

128 """ 

129 

130 def __init__(self, datasets: Iterable[FileDataset]): 

131 super().__init__(ref for dataset in datasets for ref in dataset.refs) 

132 self.datasets = datasets 

133 

134 

135class FileDatastore(GenericBaseDatastore[StoredFileInfo]): 

136 """Generic Datastore for file-based implementations. 

137 

138 Should always be sub-classed since key abstract methods are missing. 

139 

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. 

157 

158 Raises 

159 ------ 

160 ValueError 

161 If root location does not exist and ``create`` is `False` in the 

162 configuration. 

163 """ 

164 

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

169 

170 root: ResourcePath 

171 """Root directory URI of this `Datastore`.""" 

172 

173 locationFactory: LocationFactory 

174 """Factory for creating locations relative to the datastore root.""" 

175 

176 formatterFactory: FormatterFactory 

177 """Factory for creating instances of formatters.""" 

178 

179 templates: FileTemplates 

180 """File templates that can be used by this `Datastore`.""" 

181 

182 composites: CompositesMap 

183 """Determines whether a dataset should be disassembled on put.""" 

184 

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

189 

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

194 

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. 

199 

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``. 

219 

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 ) 

234 

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 ) 

252 

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 

269 

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

278 

279 self.locationFactory = LocationFactory(self.root) 

280 

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 

295 

296 # Determine whether checksums should be used - default to False 

297 self.useChecksum = self.config.get("checksum", False) 

298 

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) 

305 

306 self.universe = bridgeManager.universe 

307 

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

317 

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) 

321 

322 # Now associate formatters with storage classes 

323 formatterFactory = FormatterFactory() 

324 formatterFactory.registerFormatters(config["formatters"], universe=bridgeManager.universe) 

325 

326 # Read the file naming templates 

327 templates = FileTemplates(config["templates"], universe=bridgeManager.universe) 

328 

329 # See if composites should be disassembled 

330 composites = CompositesMap(config["composites"], universe=bridgeManager.universe) 

331 

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) 

335 

336 self = FileDatastore( 

337 config, bridgeManager, root, formatterFactory, templates, composites, trustGetRequest 

338 ) 

339 

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 

361 

362 return self 

363 

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 ) 

374 

375 def __str__(self) -> str: 

376 return str(self.root) 

377 

378 @property 

379 def bridge(self) -> DatastoreRegistryBridge: 

380 return self._bridge 

381 

382 @property 

383 def roots(self) -> dict[str, ResourcePath | None]: 

384 # Docstring inherited. 

385 return {self.name: self.root} 

386 

387 def _set_trust_mode(self, mode: bool) -> None: 

388 self.trustGetRequest = mode 

389 

390 def _artifact_exists(self, location: Location) -> bool: 

391 """Check that an artifact exists in this datastore at the specified 

392 location. 

393 

394 Parameters 

395 ---------- 

396 location : `Location` 

397 Expected location of the artifact associated with this datastore. 

398 

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

406 

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. 

415 

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

440 

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. 

446 

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. 

453 

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) 

470 

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] 

475 

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. 

483 

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] = [] 

497 

498 for ref, itemInfo in refsAndInfos: 

499 expandedRefs.append(ref) 

500 expandedItemInfos.append(itemInfo) 

501 

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) 

513 

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. 

518 

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. 

525 

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) 

547 

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

550 

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 

556 

557 def _refs_associated_with_artifacts( 

558 self, paths: Iterable[str | ResourcePath] 

559 ) -> dict[str, set[DatasetId]]: 

560 """Return paths and associated dataset refs. 

561 

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. 

567 

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) 

585 

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 

602 

603 def _registered_refs_per_artifact(self, pathInStore: ResourcePath) -> set[DatasetId]: 

604 """Return all dataset refs associated with the supplied path. 

605 

606 Parameters 

607 ---------- 

608 pathInStore : `lsst.resources.ResourcePath` 

609 Path of interest in the data store. 

610 

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 

619 

620 def removeStoredItemInfo(self, ref: DatasetIdRef) -> None: 

621 """Remove information about the file associated with this dataset. 

622 

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

632 

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. 

638 

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. 

645 

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) 

654 

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] 

658 

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. 

662 

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. 

669 

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 

678 

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

683 

684 # Remove these refs from all the refs and if there is nothing left 

685 # then we can delete 

686 remainingRefs = allRefs - {ref.id} 

687 

688 if remainingRefs: 

689 return False 

690 return True 

691 

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. 

695 

696 Parameters 

697 ---------- 

698 ref : `DatasetRef` 

699 Reference to the required `Dataset`. 

700 

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. 

706 

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

727 

728 # See if the ref is a composite that should be disassembled 

729 doDisassembly = self.composites.shouldBeDisassembled(ref) 

730 

731 all_info: list[tuple[Location, Formatter | FormatterV2, StorageClass, str | None]] = [] 

732 

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

738 

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

743 

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 ] 

759 

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. 

765 

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. 

773 

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) 

780 

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) 

788 

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) 

796 

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

804 

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

810 

811 return generate_datastore_get_information( 

812 fileLocations, 

813 registry_ref=registry_ref, 

814 read_ref=ref, 

815 parameters=parameters, 

816 ) 

817 

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. 

822 

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. 

829 

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 

842 

843 # Validate the template to protect against filenames from different 

844 # dataIds returning the same and causing overwrite confusion. 

845 template.validateTemplate(ref) 

846 

847 location = self.locationFactory.fromPath(template.format(ref), trusted_path=True) 

848 

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 

863 

864 # Now that we know the formatter, update the location 

865 location = formatter.make_updated_location(location) 

866 

867 return location, formatter 

868 

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 

873 

874 # See if the paths are within the datastore or not 

875 inside = [self._pathInStore(d.path) is not None for d in datasets] 

876 

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" 

895 

896 return transfer 

897 

898 def _pathInStore(self, path: ResourcePathExpression) -> str | None: 

899 """Return path relative to datastore root. 

900 

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. 

907 

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) 

917 

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. 

926 

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. 

942 

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. 

948 

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. 

954 

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

963 

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) 

968 

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 ) 

974 

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 ) 

981 

982 # Return the relative path within the datastore for internal 

983 # transfer 

984 path = relpath 

985 

986 return path 

987 

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. 

999 

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. 

1019 

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. 

1026 

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

1037 

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) 

1041 

1042 # Track whether we have read the size of the source yet 

1043 have_sized = False 

1044 

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) 

1077 

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 

1085 

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 ) 

1094 

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 

1102 

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 

1112 

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 ) 

1121 

1122 def _prepIngest(self, *datasets: FileDataset, transfer: str | None = None) -> _IngestPrepData: 

1123 # Docstring inherited from Datastore._prepIngest. 

1124 filtered = [] 

1125 

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 

1141 

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 

1156 

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 

1163 

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 ) 

1172 

1173 dataset.path = self._standardizeIngestPath( 

1174 dataset.path, transfer=transfer, check_existence=check_existence 

1175 ) 

1176 filtered.append(dataset) 

1177 return _IngestPrepData(filtered) 

1178 

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

1200 

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 

1215 

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) 

1220 

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. 

1229 

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. 

1242 

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) 

1252 

1253 # Get the extension 

1254 ext = srcUri.getExtension() 

1255 

1256 # Update the destination to include that extension 

1257 location.updateExtension(ext) 

1258 

1259 # Ask the formatter to validate this extension 

1260 if formatter is not None: 

1261 formatter.validate_extension(location) 

1262 

1263 return location 

1264 

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. 

1269 

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. 

1279 

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) 

1289 

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 ) 

1296 

1297 location, formatter = self._determine_put_formatter_location(ref) 

1298 

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 

1310 

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 ) 

1316 

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

1319 

1320 def _removeFileExists(uri: ResourcePath) -> None: 

1321 """Remove a file and do not complain if it is not there. 

1322 

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

1329 

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) 

1334 

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 ) 

1347 

1348 assert isinstance(formatter_compat, FormatterV2) 

1349 

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 

1360 

1361 # URI is needed to resolve what ingest case are we dealing with 

1362 return self._extractIngestInfo(uri, ref, formatter=formatter) 

1363 

1364 def knows(self, ref: DatasetRef) -> bool: 

1365 """Check if the dataset is known to the datastore. 

1366 

1367 Does not check for existence of any artifact. 

1368 

1369 Parameters 

1370 ---------- 

1371 ref : `DatasetRef` 

1372 Reference to the required dataset. 

1373 

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 

1383 

1384 def knows_these(self, refs: Iterable[DatasetRef]) -> dict[DatasetRef, bool]: 

1385 # Docstring inherited from the base class. 

1386 refs = list(refs) 

1387 

1388 # The records themselves. Could be missing some entries. 

1389 records = self._get_stored_records_associated_with_refs(refs, ignore_datastore_records=True) 

1390 

1391 return {ref: ref.id in records for ref in refs} 

1392 

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. 

1401 

1402 Helper function for `mexists()`. 

1403 

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. 

1418 

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

1428 

1429 location_factory = self.locationFactory 

1430 

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

1436 

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) 

1453 

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 

1464 

1465 # Results. 

1466 dataset_existence: dict[DatasetRef, bool] = {} 

1467 

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] 

1472 

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 

1481 

1482 if artifact_existence is not None: 

1483 artifact_existence.update(uri_existence) 

1484 

1485 return dataset_existence 

1486 

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. 

1491 

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. 

1500 

1501 Returns 

1502 ------- 

1503 existence : `dict` of [`DatasetRef`, `bool`] 

1504 Mapping from dataset to boolean indicating existence. 

1505 

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) 

1523 

1524 # The log message level and content depend on how many 

1525 # datasets we are processing. 

1526 n_results = len(chunk_result) 

1527 

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 

1532 

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 

1540 

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 ) 

1565 

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 

1580 

1581 return dataset_existence 

1582 

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. 

1587 

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. 

1596 

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} 

1606 

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} 

1610 

1611 # Set of all IDs we are checking for. 

1612 requested_ids = set(id_to_ref.keys()) 

1613 

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 ) 

1618 

1619 dataset_existence = self._process_mexists_records( 

1620 id_to_ref, records, True, artifact_existence=artifact_existence 

1621 ) 

1622 

1623 # Set of IDs that have been handled. 

1624 handled_ids = {ref.id for ref in dataset_existence} 

1625 

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 ) 

1633 

1634 return { 

1635 internal_ref_to_input_ref[internal_ref]: existence 

1636 for internal_ref, existence in dataset_existence.items() 

1637 } 

1638 

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. 

1643 

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. 

1653 

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 ) 

1669 

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 

1680 

1681 dataset_existence.update( 

1682 self._process_mexists_records( 

1683 id_to_ref, 

1684 records, 

1685 False, 

1686 artifact_existence=artifact_existence, 

1687 ) 

1688 ) 

1689 

1690 return dataset_existence 

1691 

1692 def exists(self, ref: DatasetRef) -> bool: 

1693 """Check if the dataset exists in the datastore. 

1694 

1695 Parameters 

1696 ---------- 

1697 ref : `DatasetRef` 

1698 Reference to the required dataset. 

1699 

1700 Returns 

1701 ------- 

1702 exists : `bool` 

1703 `True` if the entity exists in the `Datastore`. 

1704 

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) 

1716 

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 

1722 

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 

1728 

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 

1737 

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 

1746 

1747 if not self._artifact_exists(location): 

1748 return False 

1749 

1750 return True 

1751 

1752 def getURIs(self, ref: DatasetRef, predict: bool = False) -> DatasetRefURIs: 

1753 """Return URIs associated with dataset. 

1754 

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. 

1762 

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] 

1773 

1774 def getURI(self, ref: DatasetRef, predict: bool = False) -> ResourcePath: 

1775 """URI to the Dataset. 

1776 

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. 

1784 

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. 

1795 

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. 

1804 

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 

1816 

1817 def _predict_URIs( 

1818 self, 

1819 ref: DatasetRef, 

1820 ) -> DatasetRefURIs: 

1821 """Predict the URIs of a dataset ref. 

1822 

1823 Parameters 

1824 ---------- 

1825 ref : `DatasetRef` 

1826 Reference to the required Dataset. 

1827 

1828 Returns 

1829 ------- 

1830 URI : DatasetRefUris 

1831 Primary and component URIs. URIs will contain a URI fragment 

1832 "#predicted". 

1833 """ 

1834 uris = DatasetRefURIs() 

1835 

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) 

1840 

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 ) 

1846 

1847 else: 

1848 location, _ = self._determine_put_formatter_location(ref) 

1849 

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 ) 

1854 

1855 return uris 

1856 

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 

1864 

1865 uris: dict[DatasetRef, DatasetRefURIs] = {} 

1866 

1867 records = self._get_stored_records_associated_with_refs(refs) 

1868 records_keys = records.keys() 

1869 

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) 

1872 

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) 

1884 

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) 

1889 

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) 

1897 

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) 

1902 

1903 return uris 

1904 

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. 

1912 

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`. 

1924 

1925 Returns 

1926 ------- 

1927 uris: DatasetRefURIs 

1928 Represents the primary URI or component URIs described by the 

1929 inputs. 

1930 

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

1945 

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 

1951 

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 

1969 

1970 return uris 

1971 

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

1981 

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

1984 

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} 

1987 

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] 

2000 

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 ) 

2010 

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

2033 

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 

2038 

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. 

2050 

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. 

2077 

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

2086 

2087 if transfer == "move": 

2088 raise ValueError("Can not move artifacts out of datastore. Use copy instead.") 

2089 

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

2095 

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) 

2098 

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 

2104 

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

2107 

2108 missing_records = self._find_missing_records(refs, missing_ids) 

2109 records.update(missing_records) 

2110 

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) 

2140 

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 

2150 

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) 

2154 

2155 if write_index: 

2156 index = ZipIndex.from_artifact_map(refs, artifact_map, destination) 

2157 index.write_index(destination) 

2158 

2159 return artifact_map 

2160 

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. 

2169 

2170 The Zip file must have an index file as created by `retrieveArtifacts`. 

2171 

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. 

2182 

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? 

2188 

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) 

2193 

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} 

2197 

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 ) 

2203 

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

2227 

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) 

2235 

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 

2240 

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

2247 

2248 info = StoredFileInfo.from_simple(index_info.info) 

2249 for id_ in index_info.ids: 

2250 artifacts.append((id_to_ref[id_], info)) 

2251 

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

2256 

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. 

2264 

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. 

2278 

2279 Returns 

2280 ------- 

2281 inMemoryDataset : `object` 

2282 Requested dataset or slice thereof as an InMemoryDataset. 

2283 

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) 

2298 

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 ) 

2303 

2304 def prepare_get_for_external_client(self, ref: DatasetRef) -> list[DatasetLocationInformation] | None: 

2305 # Docstring inherited 

2306 

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 

2310 

2311 return locations 

2312 

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. 

2316 

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. 

2326 

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. 

2333 

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 

2344 

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

2371 

2372 self._register_datasets(artifacts, insert_mode=DatabaseInsertMode.INSERT) 

2373 

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 

2378 

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

2401 

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} 

2405 

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) 

2412 

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) 

2423 

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} 

2428 

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] 

2437 

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 

2448 

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] 

2459 

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 

2472 

2473 log.debug("Trashing dataset %s in datastore %s", ref, self.name) 

2474 

2475 fileLocations = self._get_dataset_locations_info(ref) 

2476 

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) 

2484 

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) 

2496 

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 

2512 

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. 

2517 

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. 

2532 

2533 Returns 

2534 ------- 

2535 removed : `set` [ `lsst.resources.ResourcePath` ] 

2536 List of artifacts that were removed. 

2537 

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 

2578 

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. 

2588 

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. 

2603 

2604 Returns 

2605 ------- 

2606 removed : `set` [ `lsst.resources.ResourcePath` ] 

2607 Artifacts successfully removed. 

2608 

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 

2632 

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

2640 

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} 

2651 

2652 path_map = self._refs_associated_with_artifacts(paths_to_check) 

2653 

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] 

2663 

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 

2669 

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

2676 

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 

2681 

2682 # Only trashed refs still known to datastore will be returned. 

2683 location = info.file_location(self.locationFactory) 

2684 

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 

2688 

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

2693 

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) 

2697 

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 ) 

2705 

2706 if not artifacts_to_delete: 

2707 return set() 

2708 

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 ) 

2725 

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 

2732 

2733 # Now remove the actual file artifacts. 

2734 remove_result = ResourcePath.mremove(artifacts_to_delete, do_raise=False) 

2735 

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) 

2747 

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 

2771 

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) 

2782 

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 ) 

2792 

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

2795 

2796 # Empty existence lookup if none given. 

2797 if artifact_existence is None: 

2798 artifact_existence = {} 

2799 

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. 

2812 

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) 

2816 

2817 # The artifacts to register 

2818 artifacts = [] 

2819 

2820 # Refs that already exist 

2821 already_present = [] 

2822 

2823 # Refs that were rejected by this datastore. 

2824 rejected = set() 

2825 

2826 # Refs that were transferred successfully. 

2827 accepted = set() 

2828 

2829 # Record each time we have done a "direct" transfer. 

2830 direct_transfers = [] 

2831 

2832 # Keep track of all the file transfers that are required. 

2833 from_to: list[tuple[ResourcePath, ResourcePath]] = [] 

2834 

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 

2842 

2843 accepted.add(ref) 

2844 

2845 if ref.id in target_records: 

2846 # Already have an artifact for this. 

2847 already_present.append(ref) 

2848 continue 

2849 

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 ) 

2888 

2889 info = info.update(path=target_location.pathInStore.path) 

2890 

2891 # Need to transfer it to the new location. 

2892 from_to.append((source_location.uri, target_location.uri)) 

2893 

2894 artifacts.append((ref, info)) 

2895 

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 ) 

2913 

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 ) 

2920 

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) 

2929 

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 ) 

2937 

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 

2945 

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) 

2951 

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

2961 

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) 

2966 

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 

2977 

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

2984 

2985 def validateConfiguration( 

2986 self, entities: Iterable[DatasetRef | DatasetType | StorageClass], logFailures: bool = False 

2987 ) -> None: 

2988 """Validate some of the configuration for this datastore. 

2989 

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. 

2999 

3000 Returns 

3001 ------- 

3002 None 

3003 

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. 

3009 

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) 

3020 

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) 

3029 

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) 

3038 

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 ) 

3046 

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 

3056 

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 

3067 

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

3070 

3071 if transfer == "move": 

3072 raise TypeError("Can not export by moving files out of datastore.") 

3073 

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) 

3078 

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

3082 

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] 

3092 

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) 

3107 

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

3118 

3119 exportUri = directoryUri.join(pathInStore) 

3120 exportUri.transfer_from(storeUri, transfer=transfer) 

3121 

3122 yield FileDataset(refs=[ref], path=pathInStore, formatter=storedFileInfo.formatter) 

3123 

3124 @staticmethod 

3125 def computeChecksum(uri: ResourcePath, algorithm: str = "blake2b", block_size: int = 8192) -> str | None: 

3126 """Compute the checksum of the supplied file. 

3127 

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. 

3137 

3138 Returns 

3139 ------- 

3140 hexdigest : `str` 

3141 Hex digest of the file. 

3142 

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

3149 

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 

3152 

3153 hasher = hashlib.new(algorithm) 

3154 

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) 

3158 

3159 return hasher.hexdigest() 

3160 

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) 

3172 

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 

3178 

3179 self._bridge.insert(FakeDatasetRef(dataset_id) for dataset_id in record_data.records) 

3180 

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) 

3191 

3192 def export_records(self, refs: Iterable[DatasetIdRef]) -> Mapping[str, DatastoreRecordData]: 

3193 # Docstring inherited from the base class. 

3194 

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) 

3201 

3202 record_data = DatastoreRecordData(records=records) 

3203 return {self.name: record_data} 

3204 

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) 

3221 

3222 def export_table(self, datasets: Collection[DatasetId]) -> DatastoreRecordTable: 

3223 # Docstring inherited from the base class. 

3224 

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) 

3230 

3231 def import_table(self, table: DatastoreRecordTable) -> None: 

3232 # Docstring inherited from the base class. 

3233 

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) 

3239 

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) 

3254 

3255 record_data = DatastoreRecordData(records=records) 

3256 return {self.name: record_data} 

3257 

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 

3261 

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 

3272 

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