Coverage for python/lsst/daf/butler/remote_butler/_remote_butler.py: 0%

301 statements  

« prev     ^ index     » next       coverage.py v7.15.3, created at 2026-08-07 09:51 +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 

28from __future__ import annotations 

29 

30__all__ = ("RemoteButler",) 

31 

32import logging 

33import uuid 

34from collections.abc import Collection, Iterable, Iterator, Sequence 

35from contextlib import AbstractContextManager, contextmanager 

36from types import EllipsisType 

37from typing import TYPE_CHECKING, Any, TextIO, cast 

38 

39from deprecated.sphinx import deprecated 

40 

41from lsst.daf.butler.datastores.file_datastore.retrieve_artifacts import ( 

42 ArtifactIndexInfo, 

43 ZipIndex, 

44 determine_destination_for_retrieved_artifact, 

45 retrieve_and_zip, 

46 unpack_zips, 

47) 

48from lsst.resources import ResourcePath, ResourcePathExpression 

49from lsst.utils.iteration import chunk_iterable 

50 

51from .._butler import Butler, _DeprecatedDefault 

52from .._butler_collections import ButlerCollections 

53from .._butler_metrics import ButlerMetrics 

54from .._dataset_existence import DatasetExistence 

55from .._dataset_ref import DatasetId, DatasetRef 

56from .._dataset_type import DatasetType 

57from .._deferredDatasetHandle import DeferredDatasetHandle 

58from .._exceptions import DatasetNotFoundError 

59from .._query_all_datasets import QueryAllDatasetsParameters 

60from .._storage_class import StorageClass, StorageClassFactory 

61from .._utilities.locked_object import LockedObject 

62from ..datastore import DatasetRefURIs, DatastoreConfig 

63from ..datastore.cache_manager import AbstractDatastoreCacheManager, DatastoreCacheManager 

64from ..dimensions import DataCoordinate, DataIdValue, DimensionConfig, DimensionUniverse, SerializedDataId 

65from ..queries import Query 

66from ..queries.tree import make_column_literal 

67from ..registry import CollectionArgType, NoDefaultCollectionError, Registry, RegistryDefaults 

68from ..registry.expand_data_ids import expand_data_ids 

69from ._collection_args import convert_collection_arg_to_glob_string_list 

70from ._defaults import DefaultsHolder 

71from ._get import convert_http_url_to_resource_path, get_dataset_as_python_object 

72from ._http_connection import RemoteButlerHttpConnection, parse_model, quote_path_variable 

73from ._query_driver import RemoteQueryDriver 

74from ._query_results import convert_dataset_ref_results, read_query_results 

75from ._ref_utils import ( 

76 apply_storage_class_override, 

77 get_component_override, 

78 normalize_dataset_type_name, 

79 simplify_dataId, 

80 split_dataset_type_name, 

81) 

82from ._registry import RemoteButlerRegistry 

83from ._remote_butler_collections import RemoteButlerCollections 

84from ._remote_file_transfer_source import RemoteFileTransferSource 

85from .server_models import ( 

86 CollectionList, 

87 FileInfoPayload, 

88 FindDatasetRequestModel, 

89 FindDatasetResponseModel, 

90 GetDatasetTypeResponseModel, 

91 GetFileByDataIdRequestModel, 

92 GetFileResponseModel, 

93 GetManyDatasetsRequestModel, 

94 GetManyDatasetsResponseModel, 

95 GetUniverseResponseModel, 

96 QueryAllDatasetsRequestModel, 

97) 

98 

99if TYPE_CHECKING: 

100 from .._dataset_provenance import DatasetProvenance 

101 from .._file_dataset import FileDataset 

102 from .._limited_butler import LimitedButler 

103 from .._timespan import Timespan 

104 from ..dimensions import DataId 

105 from ..transfers import RepoExportContext 

106 

107 

108_LOG = logging.getLogger(__name__) 

109 

110 

111class RemoteButler(Butler): # numpydoc ignore=PR02 

112 """A `Butler` that can be used to connect through a remote server. 

113 

114 Parameters 

115 ---------- 

116 options : `ButlerInstanceOptions` 

117 Default values and other settings for the Butler instance. 

118 connection : `RemoteButlerHttpConnection` 

119 Connection to Butler server. 

120 cache : `RemoteButlerCache` 

121 Cache of data shared between multiple RemoteButler instances connected 

122 to the same server. 

123 use_disabled_datastore_cache : `bool`, optional 

124 If `True`, a datastore cache manager will be created with a default 

125 disabled state which can be enabled by the environment. If `False` 

126 a cache manager will be constructed from the default local 

127 configuration, likely caching by default but only specific storage 

128 classes. 

129 

130 Notes 

131 ----- 

132 Instead of using this constructor, most users should use either 

133 `Butler.from_config` or `RemoteButlerFactory`. 

134 """ 

135 

136 _registry_defaults: DefaultsHolder 

137 _connection: RemoteButlerHttpConnection 

138 _cache: RemoteButlerCache 

139 _registry: RemoteButlerRegistry 

140 _datastore_cache_manager: AbstractDatastoreCacheManager | None 

141 _use_disabled_datastore_cache: bool 

142 

143 # This is __new__ instead of __init__ because we have to support 

144 # instantiation via the legacy constructor Butler.__new__(), which 

145 # reads the configuration and selects which subclass to instantiate. The 

146 # interaction between __new__ and __init__ is kind of wacky in Python. If 

147 # we were using __init__ here, __init__ would be called twice (once when 

148 # the RemoteButler instance is constructed inside Butler.from_config(), and 

149 # a second time with the original arguments to Butler() when the instance 

150 # is returned from Butler.__new__() 

151 def __new__( 

152 cls, 

153 *, 

154 connection: RemoteButlerHttpConnection, 

155 defaults: RegistryDefaults, 

156 cache: RemoteButlerCache, 

157 use_disabled_datastore_cache: bool = True, 

158 metrics: ButlerMetrics | None = None, 

159 ) -> RemoteButler: 

160 self = cast(RemoteButler, super().__new__(cls)) 

161 self.storageClasses = StorageClassFactory() 

162 

163 self._connection = connection 

164 self._cache = cache 

165 self._datastore_cache_manager = None 

166 self._use_disabled_datastore_cache = use_disabled_datastore_cache 

167 self._metrics = metrics if metrics is not None else ButlerMetrics() 

168 

169 self._registry_defaults = DefaultsHolder(defaults) 

170 self._registry = RemoteButlerRegistry(self, self._registry_defaults, self._connection) 

171 defaults.finish(self._registry) 

172 

173 return self 

174 

175 def isWriteable(self) -> bool: 

176 # Docstring inherited. 

177 return False 

178 

179 @property 

180 @deprecated( 

181 "Please use 'collections' instead. collection_chains will be removed after v28.", 

182 version="v28", 

183 category=FutureWarning, 

184 ) 

185 def collection_chains(self) -> ButlerCollections: 

186 """Object with methods for modifying collection chains.""" 

187 return self.collections 

188 

189 @property 

190 def collections(self) -> ButlerCollections: 

191 """Object with methods for modifying and querying collections.""" 

192 return RemoteButlerCollections(self._registry_defaults, self._connection) 

193 

194 @property 

195 def dimensions(self) -> DimensionUniverse: 

196 # Docstring inherited. 

197 with self._cache.access() as cache: 

198 if cache.dimensions is not None: 

199 return cache.dimensions 

200 

201 response = self._connection.get("universe") 

202 model = parse_model(response, GetUniverseResponseModel) 

203 

204 config = DimensionConfig.from_simple(model.universe) 

205 universe = DimensionUniverse(config) 

206 with self._cache.access() as cache: 

207 if cache.dimensions is None: 

208 cache.dimensions = universe 

209 return cache.dimensions 

210 

211 @property 

212 def _cache_manager(self) -> AbstractDatastoreCacheManager: 

213 """Cache manager to use when reading files from the butler.""" 

214 # RemoteButler does not get any cache configuration from the server. 

215 # Either create a disabled cache manager which can be enabled via the 

216 # environment, or create a cache manager from the default FileDatastore 

217 # config. This will not work properly if the defaults for 

218 # DatastoreConfig no longer include the cache. 

219 if self._datastore_cache_manager is None: 

220 datastore_config = DatastoreConfig() 

221 if not self._use_disabled_datastore_cache and "cached" in datastore_config: 

222 self._datastore_cache_manager = DatastoreCacheManager( 

223 datastore_config["cached"], universe=self.dimensions 

224 ) 

225 else: 

226 self._datastore_cache_manager = DatastoreCacheManager.create_disabled( 

227 universe=self.dimensions 

228 ) 

229 return self._datastore_cache_manager 

230 

231 def _caching_context(self) -> AbstractContextManager[None]: 

232 # Docstring inherited. 

233 # Not implemented for now, will have to think whether this needs to 

234 # do something on client side and/or remote side. 

235 raise NotImplementedError() 

236 

237 def transaction(self) -> AbstractContextManager[None]: 

238 """Will always raise NotImplementedError. 

239 Transactions are not supported by RemoteButler. 

240 """ 

241 raise NotImplementedError() 

242 

243 def put( 

244 self, 

245 obj: Any, 

246 datasetRefOrType: DatasetRef | DatasetType | str, 

247 /, 

248 dataId: DataId | None = None, 

249 *, 

250 run: str | None = None, 

251 provenance: DatasetProvenance | None = None, 

252 **kwargs: Any, 

253 ) -> DatasetRef: 

254 # Docstring inherited. 

255 raise NotImplementedError() 

256 

257 def getDeferred( 

258 self, 

259 datasetRefOrType: DatasetRef | DatasetType | str, 

260 /, 

261 dataId: DataId | None = None, 

262 *, 

263 parameters: dict | None = None, 

264 collections: Any = None, 

265 storageClass: str | StorageClass | None = None, 

266 timespan: Timespan | None = None, 

267 **kwargs: Any, 

268 ) -> DeferredDatasetHandle: 

269 response = self._get_file_info(datasetRefOrType, dataId, collections, timespan, kwargs) 

270 # Check that artifact information is available. 

271 _to_file_payload(response) 

272 if isinstance(datasetRefOrType, DatasetRef): 

273 # Use the ref provided by the caller, which may include component 

274 # or storage class overrides that are not known to the server. 

275 ref = datasetRefOrType 

276 else: 

277 ref = DatasetRef.from_simple(response.dataset_ref, universe=self.dimensions) 

278 # The server returns the parent dataset type -- component dataset 

279 # types are never sent to the server, because it may not have the 

280 # storage class definitions needed to construct them. Re-apply 

281 # any component here. 

282 component = get_component_override(datasetRefOrType) 

283 if component is not None: 

284 ref = ref.makeComponentRef(component) 

285 return DeferredDatasetHandle(butler=self, ref=ref, parameters=parameters, storageClass=storageClass) 

286 

287 def get( 

288 self, 

289 datasetRefOrType: DatasetRef | DatasetType | str, 

290 /, 

291 dataId: DataId | None = None, 

292 *, 

293 parameters: dict[str, Any] | None = None, 

294 collections: Any = None, 

295 storageClass: StorageClass | str | None = None, 

296 timespan: Timespan | None = None, 

297 **kwargs: Any, 

298 ) -> Any: 

299 # Docstring inherited. 

300 with self._metrics.instrument_get(log=_LOG, msg="Retrieved remote dataset"): 

301 model = self._get_file_info(datasetRefOrType, dataId, collections, timespan, kwargs) 

302 

303 written_ref = DatasetRef.from_simple(model.dataset_ref, universe=self.dimensions) 

304 # The server returns the parent dataset type -- component dataset 

305 # types are never sent to the server, because it may not have the 

306 # storage class definitions needed to construct them. Re-apply 

307 # any component here. 

308 componentOverride = get_component_override(datasetRefOrType) 

309 if componentOverride is not None: 

310 written_ref = written_ref.makeComponentRef(componentOverride) 

311 # The written ref carries the storage class the file was written 

312 # with; the read ref carries any storage class override requested 

313 # by the caller. Keep them separate so the Formatter always sees 

314 # the written ref (as it does with DirectButler). 

315 read_ref = apply_storage_class_override(written_ref, datasetRefOrType, storageClass) 

316 

317 return self._get_dataset_as_python_object(written_ref, read_ref, model, parameters) 

318 

319 def _get_dataset_as_python_object( 

320 self, 

321 written_ref: DatasetRef, 

322 read_ref: DatasetRef, 

323 model: GetFileResponseModel, 

324 parameters: dict[str, Any] | None, 

325 ) -> Any: 

326 # This thin wrapper method is here to provide a place to hook in a mock 

327 # mimicking DatastoreMock functionality for use in unit tests. 

328 return get_dataset_as_python_object( 

329 written_ref, 

330 read_ref, 

331 _to_file_payload(model), 

332 auth=self._connection.auth, 

333 parameters=parameters, 

334 cache_manager=self._cache_manager, 

335 ) 

336 

337 def _get_file_info( 

338 self, 

339 datasetRefOrType: DatasetRef | DatasetType | str, 

340 dataId: DataId | None, 

341 collections: CollectionArgType, 

342 timespan: Timespan | None, 

343 kwargs: dict[str, DataIdValue], 

344 ) -> GetFileResponseModel: 

345 """Send a request to the server for the file URLs and metadata 

346 associated with a dataset. 

347 """ 

348 if isinstance(datasetRefOrType, DatasetRef): 

349 if dataId is not None: 

350 raise ValueError("DatasetRef given, cannot use dataId as well") 

351 return self._get_file_info_for_ref(datasetRefOrType) 

352 else: 

353 # Only the parent dataset type is sent to the server -- it may 

354 # not have the storage class definitions needed to construct a 

355 # component DatasetType. Callers are responsible for re-applying 

356 # any component to the returned ref. 

357 dataset_type_name, _ = split_dataset_type_name(datasetRefOrType) 

358 request = GetFileByDataIdRequestModel( 

359 dataset_type=dataset_type_name, 

360 collections=self._normalize_collections(collections), 

361 data_id=simplify_dataId(dataId, kwargs), 

362 default_data_id=self._serialize_default_data_id(), 

363 timespan=timespan, 

364 ) 

365 response = self._connection.post("get_file_by_data_id", request) 

366 return parse_model(response, GetFileResponseModel) 

367 

368 def _get_file_info_for_ref(self, ref: DatasetRef) -> GetFileResponseModel: 

369 response = self._connection.get(f"get_file/{_to_uuid_string(ref.id)}") 

370 return parse_model(response, GetFileResponseModel) 

371 

372 def getURIs( 

373 self, 

374 datasetRefOrType: DatasetRef | DatasetType | str, 

375 /, 

376 dataId: DataId | None = None, 

377 *, 

378 predict: bool = False, 

379 collections: Any = None, 

380 run: str | None = None, 

381 **kwargs: Any, 

382 ) -> DatasetRefURIs: 

383 # Docstring inherited. 

384 if predict or run: 

385 raise NotImplementedError("Predict mode is not supported by RemoteButler") 

386 

387 response = self._get_file_info(datasetRefOrType, dataId, collections, None, kwargs) 

388 file_info = _to_file_payload(response).file_info 

389 if len(file_info) == 1: 

390 return DatasetRefURIs( 

391 primaryURI=convert_http_url_to_resource_path( 

392 file_info[0].url, self._connection.auth, file_info[0].auth 

393 ) 

394 ) 

395 else: 

396 components = {} 

397 for f in file_info: 

398 component = f.datastoreRecords.component 

399 if component is None: 

400 raise ValueError( 

401 f"DatasetId {response.dataset_ref.id} has a component file" 

402 " with no component name defined" 

403 ) 

404 components[component] = convert_http_url_to_resource_path( 

405 f.url, self._connection.auth, f.auth 

406 ) 

407 return DatasetRefURIs(componentURIs=components) 

408 

409 def get_dataset_type(self, name: str) -> DatasetType: 

410 with self._cache.access() as cache: 

411 if (cached_value := cache.dataset_types.get(name)) is not None: 

412 return cached_value 

413 

414 # Only the parent dataset type name is sent to the server -- it may 

415 # not have the storage class definitions needed to construct a 

416 # component DatasetType, so the component dataset type is constructed 

417 # here from the parent definition. 

418 parent_name, component = split_dataset_type_name(name) 

419 response = self._connection.get(f"dataset_type/{quote_path_variable(parent_name)}") 

420 model = parse_model(response, GetDatasetTypeResponseModel) 

421 value = DatasetType.from_simple(model.dataset_type, universe=self.dimensions) 

422 if component is not None: 

423 value = value.makeComponentDatasetType(component) 

424 with self._cache.access() as cache: 

425 return cache.dataset_types.setdefault(name, value) 

426 

427 def get_dataset( 

428 self, 

429 id: DatasetId | str, 

430 *, 

431 storage_class: str | StorageClass | None = None, 

432 dimension_records: bool = False, 

433 datastore_records: bool = False, 

434 ) -> DatasetRef | None: 

435 # datastore_records is intentionally ignored. It is an optimization 

436 # flag that only applies to DirectButler. 

437 path = f"dataset/{_to_uuid_string(id)}" 

438 response = self._connection.get(path, params={"dimension_records": bool(dimension_records)}) 

439 model = parse_model(response, FindDatasetResponseModel) 

440 if model.dataset_ref is None: 

441 return None 

442 ref = DatasetRef.from_simple(model.dataset_ref, universe=self.dimensions) 

443 if storage_class is not None: 

444 ref = ref.overrideStorageClass(storage_class) 

445 return ref 

446 

447 def get_many_datasets(self, ids: Iterable[DatasetId | str]) -> list[DatasetRef]: 

448 result = [] 

449 for batch in chunk_iterable(ids, GetManyDatasetsRequestModel.MAX_ITEMS_PER_REQUEST): 

450 request = GetManyDatasetsRequestModel(dataset_ids=batch) 

451 response = self._connection.post("datasets", request) 

452 model = parse_model(response, GetManyDatasetsResponseModel) 

453 refs = convert_dataset_ref_results(model, self.dimensions) 

454 result.extend(refs) 

455 return result 

456 

457 def find_dataset( 

458 self, 

459 dataset_type: DatasetType | str, 

460 data_id: DataId | None = None, 

461 *, 

462 collections: str | Sequence[str] | None = None, 

463 timespan: Timespan | None = None, 

464 storage_class: str | StorageClass | None = None, 

465 dimension_records: bool = False, 

466 datastore_records: bool = False, 

467 **kwargs: Any, 

468 ) -> DatasetRef | None: 

469 # datastore_records is intentionally ignored. It is an optimization 

470 # flag that only applies to DirectButler. 

471 

472 # Only the parent dataset type is sent to the server -- it may not 

473 # have the storage class definitions needed to construct a component 

474 # DatasetType. The component is re-applied to the returned ref below. 

475 dataset_type_name, component = split_dataset_type_name(dataset_type) 

476 query = FindDatasetRequestModel( 

477 dataset_type=dataset_type_name, 

478 data_id=simplify_dataId(data_id, kwargs), 

479 default_data_id=self._serialize_default_data_id(), 

480 collections=self._normalize_collections(collections), 

481 timespan=timespan, 

482 dimension_records=dimension_records, 

483 ) 

484 

485 response = self._connection.post("find_dataset", query) 

486 

487 model = parse_model(response, FindDatasetResponseModel) 

488 if model.dataset_ref is None: 

489 return None 

490 

491 ref = DatasetRef.from_simple(model.dataset_ref, universe=self.dimensions) 

492 if isinstance(data_id, DataCoordinate) and data_id.hasRecords(): 

493 ref = ref.expanded(data_id) 

494 if component is not None: 

495 ref = ref.makeComponentRef(component) 

496 return apply_storage_class_override(ref, dataset_type, storage_class) 

497 

498 def _retrieve_artifacts( 

499 self, 

500 refs: Iterable[DatasetRef], 

501 destination: ResourcePathExpression, 

502 transfer: str = "auto", 

503 preserve_path: bool = True, 

504 overwrite: bool = False, 

505 write_index: bool = True, 

506 add_prefix: bool = False, 

507 ) -> dict[ResourcePath, ArtifactIndexInfo]: 

508 destination = ResourcePath(destination).abspath() 

509 if not destination.isdir(): 

510 raise ValueError(f"Destination location must refer to a directory. Given {destination}.") 

511 

512 if transfer not in ("auto", "copy"): 

513 raise ValueError("Only 'copy' and 'auto' transfer modes are supported.") 

514 

515 requested_ids = {ref.id for ref in refs} 

516 have_copied: dict[ResourcePath, ResourcePath] = {} 

517 artifact_map: dict[ResourcePath, ArtifactIndexInfo] = {} 

518 # Sort to ensure that in many refs to one file situation the same 

519 # ref is used for any prefix that might be added. 

520 for ref in sorted(refs): 

521 prefix = str(ref.id)[:8] + "-" if add_prefix else "" 

522 file_info = _to_file_payload(self._get_file_info_for_ref(ref)).file_info 

523 for file in file_info: 

524 source_uri = ResourcePath(str(file.url)) 

525 # For DECam/zip we only want to copy once. 

526 # For zip files we need to unpack so that they can be 

527 # zipped up again if needed. 

528 is_zip = source_uri.getExtension() == ".zip" and "zip-path" in source_uri.fragment 

529 cleaned_source_uri = source_uri.replace(fragment="", query="", params="") 

530 if is_zip: 

531 if cleaned_source_uri not in have_copied: 

532 zipped_artifacts = unpack_zips( 

533 [cleaned_source_uri], requested_ids, destination, preserve_path, overwrite 

534 ) 

535 artifact_map.update(zipped_artifacts) 

536 have_copied[cleaned_source_uri] = cleaned_source_uri 

537 elif cleaned_source_uri not in have_copied: 

538 relative_path = ResourcePath(file.datastoreRecords.path, forceAbsolute=False) 

539 target_uri = determine_destination_for_retrieved_artifact( 

540 destination, relative_path, preserve_path, prefix 

541 ) 

542 # Because signed URLs expire, we want to do the transfer 

543 # soon after retrieving the URL. 

544 target_uri.transfer_from(source_uri, transfer="copy", overwrite=overwrite) 

545 have_copied[cleaned_source_uri] = target_uri 

546 artifact_map[target_uri] = ArtifactIndexInfo.from_single(file.datastoreRecords, ref.id) 

547 else: 

548 target_uri = have_copied[cleaned_source_uri] 

549 artifact_map[target_uri].append(ref.id) 

550 

551 if write_index: 

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

553 index.write_index(destination) 

554 

555 return artifact_map 

556 

557 def retrieve_artifacts_zip( 

558 self, 

559 refs: Iterable[DatasetRef], 

560 destination: ResourcePathExpression, 

561 overwrite: bool = True, 

562 ) -> ResourcePath: 

563 return retrieve_and_zip(refs, destination, self._retrieve_artifacts, overwrite) 

564 

565 def retrieveArtifacts( 

566 self, 

567 refs: Iterable[DatasetRef], 

568 destination: ResourcePathExpression, 

569 transfer: str = "auto", 

570 preserve_path: bool = True, 

571 overwrite: bool = False, 

572 ) -> list[ResourcePath]: 

573 artifact_map = self._retrieve_artifacts( 

574 refs, 

575 destination, 

576 transfer, 

577 preserve_path, 

578 overwrite, 

579 ) 

580 return list(artifact_map) 

581 

582 def exists( 

583 self, 

584 dataset_ref_or_type: DatasetRef | DatasetType | str, 

585 /, 

586 data_id: DataId | None = None, 

587 *, 

588 full_check: bool = True, 

589 collections: Any = None, 

590 **kwargs: Any, 

591 ) -> DatasetExistence: 

592 try: 

593 response = self._get_file_info( 

594 dataset_ref_or_type, dataId=data_id, collections=collections, timespan=None, kwargs=kwargs 

595 ) 

596 except DatasetNotFoundError: 

597 return DatasetExistence.UNRECOGNIZED 

598 

599 if response.artifact is None: 

600 if full_check: 

601 return DatasetExistence.RECORDED 

602 else: 

603 return DatasetExistence.RECORDED | DatasetExistence._ASSUMED 

604 

605 if full_check: 

606 for file in response.artifact.file_info: 

607 if not ResourcePath(str(file.url)).exists(): 

608 return DatasetExistence.RECORDED | DatasetExistence.DATASTORE 

609 return DatasetExistence.VERIFIED 

610 else: 

611 return DatasetExistence.KNOWN 

612 

613 def _exists_many( 

614 self, 

615 refs: Iterable[DatasetRef], 

616 /, 

617 *, 

618 full_check: bool = True, 

619 ) -> dict[DatasetRef, DatasetExistence]: 

620 return {ref: self.exists(ref, full_check=full_check) for ref in refs} 

621 

622 def removeRuns( 

623 self, 

624 names: Iterable[str], 

625 unstore: bool | type[_DeprecatedDefault] = _DeprecatedDefault, 

626 *, 

627 unlink_from_chains: bool = False, 

628 ) -> None: 

629 # Docstring inherited. 

630 raise NotImplementedError() 

631 

632 def ingest( 

633 self, 

634 *datasets: FileDataset, 

635 transfer: str | None = "auto", 

636 record_validation_info: bool = True, 

637 skip_existing: bool = False, 

638 ) -> None: 

639 # Docstring inherited. 

640 raise NotImplementedError() 

641 

642 def ingest_zip( 

643 self, 

644 zip_file: ResourcePathExpression, 

645 transfer: str = "auto", 

646 *, 

647 transfer_dimensions: bool = False, 

648 dry_run: bool = False, 

649 skip_existing: bool = False, 

650 ) -> None: 

651 # Docstring inherited. 

652 raise NotImplementedError() 

653 

654 def export( 

655 self, 

656 *, 

657 directory: str | None = None, 

658 filename: str | None = None, 

659 format: str | None = None, 

660 transfer: str | None = None, 

661 ) -> AbstractContextManager[RepoExportContext]: 

662 # Docstring inherited. 

663 raise NotImplementedError() 

664 

665 def import_( 

666 self, 

667 *, 

668 directory: ResourcePathExpression | None = None, 

669 filename: ResourcePathExpression | TextIO | None = None, 

670 format: str | None = None, 

671 transfer: str | None = None, 

672 skip_dimensions: set | None = None, 

673 record_validation_info: bool = True, 

674 without_datastore: bool = False, 

675 ) -> None: 

676 # Docstring inherited. 

677 raise NotImplementedError() 

678 

679 def transfer_dimension_records_from( 

680 self, source_butler: LimitedButler | Butler, source_refs: Iterable[DatasetRef | DataCoordinate] 

681 ) -> None: 

682 # Docstring inherited. 

683 raise NotImplementedError() 

684 

685 def transfer_from( 

686 self, 

687 source_butler: LimitedButler, 

688 source_refs: Iterable[DatasetRef], 

689 transfer: str = "auto", 

690 skip_missing: bool = True, 

691 register_dataset_types: bool = False, 

692 transfer_dimensions: bool = False, 

693 dry_run: bool = False, 

694 ) -> Collection[DatasetRef]: 

695 # Docstring inherited. 

696 raise NotImplementedError() 

697 

698 def validateConfiguration( 

699 self, 

700 logFailures: bool = False, 

701 datasetTypeNames: Iterable[str] | None = None, 

702 ignore: Iterable[str] | None = None, 

703 ) -> None: 

704 # Docstring inherited. 

705 raise NotImplementedError() 

706 

707 @property 

708 def run(self) -> str | None: 

709 # Docstring inherited. 

710 return self._registry_defaults.get().run 

711 

712 @property 

713 def registry(self) -> Registry: 

714 return self._registry 

715 

716 @contextmanager 

717 def query(self) -> Iterator[Query]: 

718 driver = RemoteQueryDriver(self, self._connection) 

719 with driver: 

720 query = Query(driver) 

721 yield query 

722 

723 @contextmanager 

724 def _query_all_datasets_by_page( 

725 self, args: QueryAllDatasetsParameters 

726 ) -> Iterator[Iterator[list[DatasetRef]]]: 

727 universe = self.dimensions 

728 

729 request = QueryAllDatasetsRequestModel( 

730 collections=self._normalize_collections(args.collections), 

731 name=[normalize_dataset_type_name(name) for name in args.name], 

732 find_first=args.find_first, 

733 data_id=simplify_dataId(args.data_id, args.kwargs), 

734 default_data_id=self._serialize_default_data_id(), 

735 where=args.where, 

736 bind={k: make_column_literal(v) for k, v in args.bind.items()}, 

737 limit=args.limit, 

738 with_dimension_records=args.with_dimension_records, 

739 ) 

740 with self._connection.post_with_stream_response("query/all_datasets", request) as response: 

741 pages = read_query_results(response) 

742 yield (convert_dataset_ref_results(page, universe) for page in pages) 

743 

744 def pruneDatasets( 

745 self, 

746 refs: Iterable[DatasetRef], 

747 *, 

748 disassociate: bool = True, 

749 unstore: bool = False, 

750 tags: Iterable[str] = (), 

751 purge: bool = False, 

752 ) -> None: 

753 # Docstring inherited. 

754 raise NotImplementedError() 

755 

756 def _normalize_collections(self, collections: CollectionArgType | None) -> CollectionList: 

757 """Convert the ``collections`` parameter in the format used by Butler 

758 methods to a standardized format for the REST API. 

759 """ 

760 if collections is None: 

761 if not self.collections.defaults: 

762 raise NoDefaultCollectionError( 

763 "No collections provided, and no defaults from butler construction." 

764 ) 

765 collections = self.collections.defaults 

766 return convert_collection_arg_to_glob_string_list(collections) 

767 

768 def clone( 

769 self, 

770 *, 

771 collections: CollectionArgType | None | EllipsisType = ..., 

772 run: str | None | EllipsisType = ..., 

773 inferDefaults: bool | EllipsisType = ..., 

774 dataId: dict[str, str] | EllipsisType = ..., 

775 metrics: ButlerMetrics | None = None, 

776 ) -> RemoteButler: 

777 defaults = self._registry_defaults.get().clone(collections, run, inferDefaults, dataId) 

778 return RemoteButler( 

779 connection=self._connection, cache=self._cache, defaults=defaults, metrics=metrics 

780 ) 

781 

782 def close(self) -> None: 

783 pass 

784 

785 def _expand_data_ids(self, data_ids: Iterable[DataCoordinate]) -> list[DataCoordinate]: 

786 return expand_data_ids(data_ids, self.dimensions, self.query, None) 

787 

788 @property 

789 def _file_transfer_source(self) -> RemoteFileTransferSource: 

790 return RemoteFileTransferSource(self._connection) 

791 

792 def __str__(self) -> str: 

793 return f"RemoteButler({self._connection.server_url})" 

794 

795 def _serialize_default_data_id(self) -> SerializedDataId: 

796 """Convert the default data ID to a serializable format.""" 

797 # In an ideal world, the default data ID would just get combined with 

798 # the rest of the data ID on the client side instead of being sent 

799 # separately to the server. Unfortunately, that requires knowledge of 

800 # the DatasetType's dimensions which we don't always have available on 

801 # the client. Data ID values can be specified indirectly by "implied" 

802 # dimensions, but knowing what things are implied depends on what the 

803 # required dimensions are. 

804 

805 return self._registry_defaults.get().dataId.to_simple(minimal=True).dataId 

806 

807 

808def _to_file_payload(get_file_response: GetFileResponseModel) -> FileInfoPayload: 

809 if get_file_response.artifact is None: 

810 ref = get_file_response.dataset_ref 

811 raise DatasetNotFoundError(f"Dataset is known, but artifact is not available. (datasetId='{ref.id}')") 

812 

813 return get_file_response.artifact 

814 

815 

816def _to_uuid_string(id: uuid.UUID | str) -> str: 

817 """Convert a UUID, or string parseable as a UUID, into a string formatted 

818 like '1481269e-4c8d-4696-bcca-d1b4c9005d06' 

819 """ 

820 return str(uuid.UUID(str(id))) 

821 

822 

823class _RemoteButlerCacheData: 

824 def __init__(self) -> None: 

825 self.dimensions: DimensionUniverse | None = None 

826 self.dataset_types: dict[str, DatasetType] = {} 

827 

828 

829class RemoteButlerCache(LockedObject[_RemoteButlerCacheData]): 

830 def __init__(self) -> None: 

831 super().__init__(_RemoteButlerCacheData())