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

298 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-23 09:36 +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 make_read_ref, 

79 normalize_dataset_type_name, 

80 simplify_dataId, 

81 split_dataset_type_name, 

82) 

83from ._registry import RemoteButlerRegistry 

84from ._remote_butler_collections import RemoteButlerCollections 

85from ._remote_file_transfer_source import RemoteFileTransferSource 

86from .server_models import ( 

87 CollectionList, 

88 FileInfoPayload, 

89 FindDatasetRequestModel, 

90 FindDatasetResponseModel, 

91 GetDatasetTypeResponseModel, 

92 GetFileByDataIdRequestModel, 

93 GetFileResponseModel, 

94 GetManyDatasetsRequestModel, 

95 GetManyDatasetsResponseModel, 

96 GetUniverseResponseModel, 

97 QueryAllDatasetsRequestModel, 

98) 

99 

100if TYPE_CHECKING: 

101 from .._dataset_provenance import DatasetProvenance 

102 from .._file_dataset import FileDataset 

103 from .._limited_butler import LimitedButler 

104 from .._timespan import Timespan 

105 from ..dimensions import DataId 

106 from ..transfers import RepoExportContext 

107 

108 

109_LOG = logging.getLogger(__name__) 

110 

111 

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

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

114 

115 Parameters 

116 ---------- 

117 options : `ButlerInstanceOptions` 

118 Default values and other settings for the Butler instance. 

119 connection : `RemoteButlerHttpConnection` 

120 Connection to Butler server. 

121 cache : `RemoteButlerCache` 

122 Cache of data shared between multiple RemoteButler instances connected 

123 to the same server. 

124 use_disabled_datastore_cache : `bool`, optional 

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

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

127 a cache manager will be constructed from the default local 

128 configuration, likely caching by default but only specific storage 

129 classes. 

130 

131 Notes 

132 ----- 

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

134 `Butler.from_config` or `RemoteButlerFactory`. 

135 """ 

136 

137 _registry_defaults: DefaultsHolder 

138 _connection: RemoteButlerHttpConnection 

139 _cache: RemoteButlerCache 

140 _registry: RemoteButlerRegistry 

141 _datastore_cache_manager: AbstractDatastoreCacheManager | None 

142 _use_disabled_datastore_cache: bool 

143 

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

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

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

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

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

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

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

151 # is returned from Butler.__new__() 

152 def __new__( 

153 cls, 

154 *, 

155 connection: RemoteButlerHttpConnection, 

156 defaults: RegistryDefaults, 

157 cache: RemoteButlerCache, 

158 use_disabled_datastore_cache: bool = True, 

159 metrics: ButlerMetrics | None = None, 

160 ) -> RemoteButler: 

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

162 self.storageClasses = StorageClassFactory() 

163 

164 self._connection = connection 

165 self._cache = cache 

166 self._datastore_cache_manager = None 

167 self._use_disabled_datastore_cache = use_disabled_datastore_cache 

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

169 

170 self._registry_defaults = DefaultsHolder(defaults) 

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

172 defaults.finish(self._registry) 

173 

174 return self 

175 

176 def isWriteable(self) -> bool: 

177 # Docstring inherited. 

178 return False 

179 

180 @property 

181 @deprecated( 

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

183 version="v28", 

184 category=FutureWarning, 

185 ) 

186 def collection_chains(self) -> ButlerCollections: 

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

188 return self.collections 

189 

190 @property 

191 def collections(self) -> ButlerCollections: 

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

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

194 

195 @property 

196 def dimensions(self) -> DimensionUniverse: 

197 # Docstring inherited. 

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

199 if cache.dimensions is not None: 

200 return cache.dimensions 

201 

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

203 model = parse_model(response, GetUniverseResponseModel) 

204 

205 config = DimensionConfig.from_simple(model.universe) 

206 universe = DimensionUniverse(config) 

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

208 if cache.dimensions is None: 

209 cache.dimensions = universe 

210 return cache.dimensions 

211 

212 @property 

213 def _cache_manager(self) -> AbstractDatastoreCacheManager: 

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

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

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

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

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

219 # DatastoreConfig no longer include the cache. 

220 if self._datastore_cache_manager is None: 

221 datastore_config = DatastoreConfig() 

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

223 self._datastore_cache_manager = DatastoreCacheManager( 

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

225 ) 

226 else: 

227 self._datastore_cache_manager = DatastoreCacheManager.create_disabled( 

228 universe=self.dimensions 

229 ) 

230 return self._datastore_cache_manager 

231 

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

233 # Docstring inherited. 

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

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

236 raise NotImplementedError() 

237 

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

239 """Will always raise NotImplementedError. 

240 Transactions are not supported by RemoteButler. 

241 """ 

242 raise NotImplementedError() 

243 

244 def put( 

245 self, 

246 obj: Any, 

247 datasetRefOrType: DatasetRef | DatasetType | str, 

248 /, 

249 dataId: DataId | None = None, 

250 *, 

251 run: str | None = None, 

252 provenance: DatasetProvenance | None = None, 

253 **kwargs: Any, 

254 ) -> DatasetRef: 

255 # Docstring inherited. 

256 raise NotImplementedError() 

257 

258 def getDeferred( 

259 self, 

260 datasetRefOrType: DatasetRef | DatasetType | str, 

261 /, 

262 dataId: DataId | None = None, 

263 *, 

264 parameters: dict | None = None, 

265 collections: Any = None, 

266 storageClass: str | StorageClass | None = None, 

267 timespan: Timespan | None = None, 

268 **kwargs: Any, 

269 ) -> DeferredDatasetHandle: 

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

271 # Check that artifact information is available. 

272 _to_file_payload(response) 

273 if isinstance(datasetRefOrType, DatasetRef): 

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

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

276 ref = datasetRefOrType 

277 else: 

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

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

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

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

282 # any component here. 

283 component = get_component_override(datasetRefOrType) 

284 if component is not None: 

285 ref = ref.makeComponentRef(component) 

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

287 

288 def get( 

289 self, 

290 datasetRefOrType: DatasetRef | DatasetType | str, 

291 /, 

292 dataId: DataId | None = None, 

293 *, 

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

295 collections: Any = None, 

296 storageClass: StorageClass | str | None = None, 

297 timespan: Timespan | None = None, 

298 **kwargs: Any, 

299 ) -> Any: 

300 # Docstring inherited. 

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

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

303 

304 # The server is the source of truth for the dataset type definition 

305 # in the repository, so this ref describes what a `get` with no 

306 # overrides would return. The server returns the parent dataset 

307 # type -- component dataset types are never sent to the server, 

308 # because it may not have the storage class definitions needed to 

309 # construct them. 

310 registry_ref = DatasetRef.from_simple(model.dataset_ref, universe=self.dimensions) 

311 

312 # Any component requested by the caller, along with any storage 

313 # class override, is carried by the read ref. Keep the two 

314 # separate so that the Formatter always sees the repository 

315 # definition (as it does with DirectButler). The storage class the 

316 # file was actually written with comes from the datastore records 

317 # in the payload, not from either ref. 

318 read_ref = make_read_ref(registry_ref, datasetRefOrType, storageClass) 

319 

320 return self._get_dataset_as_python_object(registry_ref, read_ref, model, parameters) 

321 

322 def _get_dataset_as_python_object( 

323 self, 

324 registry_ref: DatasetRef, 

325 read_ref: DatasetRef, 

326 model: GetFileResponseModel, 

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

328 ) -> Any: 

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

330 # mimicking DatastoreMock functionality for use in unit tests. 

331 return get_dataset_as_python_object( 

332 registry_ref, 

333 read_ref, 

334 _to_file_payload(model), 

335 auth=self._connection.auth, 

336 parameters=parameters, 

337 cache_manager=self._cache_manager, 

338 ) 

339 

340 def _get_file_info( 

341 self, 

342 datasetRefOrType: DatasetRef | DatasetType | str, 

343 dataId: DataId | None, 

344 collections: CollectionArgType, 

345 timespan: Timespan | None, 

346 kwargs: dict[str, DataIdValue], 

347 ) -> GetFileResponseModel: 

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

349 associated with a dataset. 

350 """ 

351 if isinstance(datasetRefOrType, DatasetRef): 

352 if dataId is not None: 

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

354 return self._get_file_info_for_ref(datasetRefOrType) 

355 else: 

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

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

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

359 # any component to the returned ref. 

360 dataset_type_name, _ = split_dataset_type_name(datasetRefOrType) 

361 request = GetFileByDataIdRequestModel( 

362 dataset_type=dataset_type_name, 

363 collections=self._normalize_collections(collections), 

364 data_id=simplify_dataId(dataId, kwargs), 

365 default_data_id=self._serialize_default_data_id(), 

366 timespan=timespan, 

367 ) 

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

369 return parse_model(response, GetFileResponseModel) 

370 

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

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

373 return parse_model(response, GetFileResponseModel) 

374 

375 def getURIs( 

376 self, 

377 datasetRefOrType: DatasetRef | DatasetType | str, 

378 /, 

379 dataId: DataId | None = None, 

380 *, 

381 predict: bool = False, 

382 collections: Any = None, 

383 run: str | None = None, 

384 **kwargs: Any, 

385 ) -> DatasetRefURIs: 

386 # Docstring inherited. 

387 if predict or run: 

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

389 

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

391 file_info = _to_file_payload(response).file_info 

392 if len(file_info) == 1: 

393 return DatasetRefURIs( 

394 primaryURI=convert_http_url_to_resource_path( 

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

396 ) 

397 ) 

398 else: 

399 components = {} 

400 for f in file_info: 

401 component = f.datastoreRecords.component 

402 if component is None: 

403 raise ValueError( 

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

405 " with no component name defined" 

406 ) 

407 components[component] = convert_http_url_to_resource_path( 

408 f.url, self._connection.auth, f.auth 

409 ) 

410 return DatasetRefURIs(componentURIs=components) 

411 

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

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

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

415 return cached_value 

416 

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

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

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

420 # here from the parent definition. 

421 parent_name, component = split_dataset_type_name(name) 

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

423 model = parse_model(response, GetDatasetTypeResponseModel) 

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

425 if component is not None: 

426 value = value.makeComponentDatasetType(component) 

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

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

429 

430 def get_dataset( 

431 self, 

432 id: DatasetId | str, 

433 *, 

434 storage_class: str | StorageClass | None = None, 

435 dimension_records: bool = False, 

436 datastore_records: bool = False, 

437 ) -> DatasetRef | None: 

438 # datastore_records is intentionally ignored. It is an optimization 

439 # flag that only applies to DirectButler. 

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

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

442 model = parse_model(response, FindDatasetResponseModel) 

443 if model.dataset_ref is None: 

444 return None 

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

446 if storage_class is not None: 

447 ref = ref.overrideStorageClass(storage_class) 

448 return ref 

449 

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

451 result = [] 

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

453 request = GetManyDatasetsRequestModel(dataset_ids=batch) 

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

455 model = parse_model(response, GetManyDatasetsResponseModel) 

456 refs = convert_dataset_ref_results(model, self.dimensions) 

457 result.extend(refs) 

458 return result 

459 

460 def find_dataset( 

461 self, 

462 dataset_type: DatasetType | str, 

463 data_id: DataId | None = None, 

464 *, 

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

466 timespan: Timespan | None = None, 

467 storage_class: str | StorageClass | None = None, 

468 dimension_records: bool = False, 

469 datastore_records: bool = False, 

470 **kwargs: Any, 

471 ) -> DatasetRef | None: 

472 # datastore_records is intentionally ignored. It is an optimization 

473 # flag that only applies to DirectButler. 

474 

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

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

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

478 dataset_type_name, component = split_dataset_type_name(dataset_type) 

479 query = FindDatasetRequestModel( 

480 dataset_type=dataset_type_name, 

481 data_id=simplify_dataId(data_id, kwargs), 

482 default_data_id=self._serialize_default_data_id(), 

483 collections=self._normalize_collections(collections), 

484 timespan=timespan, 

485 dimension_records=dimension_records, 

486 ) 

487 

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

489 

490 model = parse_model(response, FindDatasetResponseModel) 

491 if model.dataset_ref is None: 

492 return None 

493 

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

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

496 ref = ref.expanded(data_id) 

497 if component is not None: 

498 ref = ref.makeComponentRef(component) 

499 return apply_storage_class_override(ref, dataset_type, storage_class) 

500 

501 def _retrieve_artifacts( 

502 self, 

503 refs: Iterable[DatasetRef], 

504 destination: ResourcePathExpression, 

505 transfer: str = "auto", 

506 preserve_path: bool = True, 

507 overwrite: bool = False, 

508 write_index: bool = True, 

509 add_prefix: bool = False, 

510 ) -> dict[ResourcePath, ArtifactIndexInfo]: 

511 destination = ResourcePath(destination).abspath() 

512 if not destination.isdir(): 

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

514 

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

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

517 

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

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

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

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

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

523 for ref in sorted(refs): 

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

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

526 for file in file_info: 

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

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

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

530 # zipped up again if needed. 

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

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

533 if is_zip: 

534 if cleaned_source_uri not in have_copied: 

535 zipped_artifacts = unpack_zips( 

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

537 ) 

538 artifact_map.update(zipped_artifacts) 

539 have_copied[cleaned_source_uri] = cleaned_source_uri 

540 elif cleaned_source_uri not in have_copied: 

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

542 target_uri = determine_destination_for_retrieved_artifact( 

543 destination, relative_path, preserve_path, prefix 

544 ) 

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

546 # soon after retrieving the URL. 

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

548 have_copied[cleaned_source_uri] = target_uri 

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

550 else: 

551 target_uri = have_copied[cleaned_source_uri] 

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

553 

554 if write_index: 

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

556 index.write_index(destination) 

557 

558 return artifact_map 

559 

560 def retrieve_artifacts_zip( 

561 self, 

562 refs: Iterable[DatasetRef], 

563 destination: ResourcePathExpression, 

564 overwrite: bool = True, 

565 ) -> ResourcePath: 

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

567 

568 def retrieveArtifacts( 

569 self, 

570 refs: Iterable[DatasetRef], 

571 destination: ResourcePathExpression, 

572 transfer: str = "auto", 

573 preserve_path: bool = True, 

574 overwrite: bool = False, 

575 ) -> list[ResourcePath]: 

576 artifact_map = self._retrieve_artifacts( 

577 refs, 

578 destination, 

579 transfer, 

580 preserve_path, 

581 overwrite, 

582 ) 

583 return list(artifact_map) 

584 

585 def exists( 

586 self, 

587 dataset_ref_or_type: DatasetRef | DatasetType | str, 

588 /, 

589 data_id: DataId | None = None, 

590 *, 

591 full_check: bool = True, 

592 collections: Any = None, 

593 **kwargs: Any, 

594 ) -> DatasetExistence: 

595 try: 

596 response = self._get_file_info( 

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

598 ) 

599 except DatasetNotFoundError: 

600 return DatasetExistence.UNRECOGNIZED 

601 

602 if response.artifact is None: 

603 if full_check: 

604 return DatasetExistence.RECORDED 

605 else: 

606 return DatasetExistence.RECORDED | DatasetExistence._ASSUMED 

607 

608 if full_check: 

609 for file in response.artifact.file_info: 

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

611 return DatasetExistence.RECORDED | DatasetExistence.DATASTORE 

612 return DatasetExistence.VERIFIED 

613 else: 

614 return DatasetExistence.KNOWN 

615 

616 def _exists_many( 

617 self, 

618 refs: Iterable[DatasetRef], 

619 /, 

620 *, 

621 full_check: bool = True, 

622 ) -> dict[DatasetRef, DatasetExistence]: 

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

624 

625 def removeRuns( 

626 self, 

627 names: Iterable[str], 

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

629 *, 

630 unlink_from_chains: bool = False, 

631 ) -> None: 

632 # Docstring inherited. 

633 raise NotImplementedError() 

634 

635 def ingest( 

636 self, 

637 *datasets: FileDataset, 

638 transfer: str | None = "auto", 

639 record_validation_info: bool = True, 

640 skip_existing: bool = False, 

641 ) -> None: 

642 # Docstring inherited. 

643 raise NotImplementedError() 

644 

645 def ingest_zip( 

646 self, 

647 zip_file: ResourcePathExpression, 

648 transfer: str = "auto", 

649 *, 

650 transfer_dimensions: bool = False, 

651 dry_run: bool = False, 

652 skip_existing: bool = False, 

653 ) -> None: 

654 # Docstring inherited. 

655 raise NotImplementedError() 

656 

657 def export( 

658 self, 

659 *, 

660 directory: str | None = None, 

661 filename: str | None = None, 

662 format: str | None = None, 

663 transfer: str | None = None, 

664 ) -> AbstractContextManager[RepoExportContext]: 

665 # Docstring inherited. 

666 raise NotImplementedError() 

667 

668 def import_( 

669 self, 

670 *, 

671 directory: ResourcePathExpression | None = None, 

672 filename: ResourcePathExpression | TextIO | None = None, 

673 format: str | None = None, 

674 transfer: str | None = None, 

675 skip_dimensions: set | None = None, 

676 record_validation_info: bool = True, 

677 without_datastore: bool = False, 

678 ) -> None: 

679 # Docstring inherited. 

680 raise NotImplementedError() 

681 

682 def transfer_dimension_records_from( 

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

684 ) -> None: 

685 # Docstring inherited. 

686 raise NotImplementedError() 

687 

688 def transfer_from( 

689 self, 

690 source_butler: LimitedButler, 

691 source_refs: Iterable[DatasetRef], 

692 transfer: str = "auto", 

693 skip_missing: bool = True, 

694 register_dataset_types: bool = False, 

695 transfer_dimensions: bool = False, 

696 dry_run: bool = False, 

697 ) -> Collection[DatasetRef]: 

698 # Docstring inherited. 

699 raise NotImplementedError() 

700 

701 def validateConfiguration( 

702 self, 

703 logFailures: bool = False, 

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

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

706 ) -> None: 

707 # Docstring inherited. 

708 raise NotImplementedError() 

709 

710 @property 

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

712 # Docstring inherited. 

713 return self._registry_defaults.get().run 

714 

715 @property 

716 def registry(self) -> Registry: 

717 return self._registry 

718 

719 @contextmanager 

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

721 driver = RemoteQueryDriver(self, self._connection) 

722 with driver: 

723 query = Query(driver) 

724 yield query 

725 

726 @contextmanager 

727 def _query_all_datasets_by_page( 

728 self, args: QueryAllDatasetsParameters 

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

730 universe = self.dimensions 

731 

732 request = QueryAllDatasetsRequestModel( 

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

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

735 find_first=args.find_first, 

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

737 default_data_id=self._serialize_default_data_id(), 

738 where=args.where, 

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

740 limit=args.limit, 

741 with_dimension_records=args.with_dimension_records, 

742 ) 

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

744 pages = read_query_results(response) 

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

746 

747 def pruneDatasets( 

748 self, 

749 refs: Iterable[DatasetRef], 

750 *, 

751 disassociate: bool = True, 

752 unstore: bool = False, 

753 tags: Iterable[str] = (), 

754 purge: bool = False, 

755 ) -> None: 

756 # Docstring inherited. 

757 raise NotImplementedError() 

758 

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

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

761 methods to a standardized format for the REST API. 

762 """ 

763 if collections is None: 

764 if not self.collections.defaults: 

765 raise NoDefaultCollectionError( 

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

767 ) 

768 collections = self.collections.defaults 

769 return convert_collection_arg_to_glob_string_list(collections) 

770 

771 def clone( 

772 self, 

773 *, 

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

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

776 inferDefaults: bool | EllipsisType = ..., 

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

778 metrics: ButlerMetrics | None = None, 

779 ) -> RemoteButler: 

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

781 return RemoteButler( 

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

783 ) 

784 

785 def close(self) -> None: 

786 pass 

787 

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

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

790 

791 @property 

792 def _file_transfer_source(self) -> RemoteFileTransferSource: 

793 return RemoteFileTransferSource(self._connection) 

794 

795 def __str__(self) -> str: 

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

797 

798 def _serialize_default_data_id(self) -> SerializedDataId: 

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

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

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

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

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

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

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

806 # required dimensions are. 

807 

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

809 

810 

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

812 if get_file_response.artifact is None: 

813 ref = get_file_response.dataset_ref 

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

815 

816 return get_file_response.artifact 

817 

818 

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

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

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

822 """ 

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

824 

825 

826class _RemoteButlerCacheData: 

827 def __init__(self) -> None: 

828 self.dimensions: DimensionUniverse | None = None 

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

830 

831 

832class RemoteButlerCache(LockedObject[_RemoteButlerCacheData]): 

833 def __init__(self) -> None: 

834 super().__init__(_RemoteButlerCacheData())