Coverage for python/lsst/daf/butler/remote_butler/_remote_butler.py: 0%
298 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-29 02:09 -0700
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-29 02:09 -0700
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/>.
28from __future__ import annotations
30__all__ = ("RemoteButler",)
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
39from deprecated.sphinx import deprecated
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
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)
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
109_LOG = logging.getLogger(__name__)
112class RemoteButler(Butler): # numpydoc ignore=PR02
113 """A `Butler` that can be used to connect through a remote server.
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.
131 Notes
132 -----
133 Instead of using this constructor, most users should use either
134 `Butler.from_config` or `RemoteButlerFactory`.
135 """
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
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()
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()
170 self._registry_defaults = DefaultsHolder(defaults)
171 self._registry = RemoteButlerRegistry(self, self._registry_defaults, self._connection)
172 defaults.finish(self._registry)
174 return self
176 def isWriteable(self) -> bool:
177 # Docstring inherited.
178 return False
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
190 @property
191 def collections(self) -> ButlerCollections:
192 """Object with methods for modifying and querying collections."""
193 return RemoteButlerCollections(self._registry_defaults, self._connection)
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
202 response = self._connection.get("universe")
203 model = parse_model(response, GetUniverseResponseModel)
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
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
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()
238 def transaction(self) -> AbstractContextManager[None]:
239 """Will always raise NotImplementedError.
240 Transactions are not supported by RemoteButler.
241 """
242 raise NotImplementedError()
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()
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)
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)
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)
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)
320 return self._get_dataset_as_python_object(registry_ref, read_ref, model, parameters)
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 )
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)
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)
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")
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)
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
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)
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
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
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.
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 )
488 response = self._connection.post("find_dataset", query)
490 model = parse_model(response, FindDatasetResponseModel)
491 if model.dataset_ref is None:
492 return None
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)
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}.")
515 if transfer not in ("auto", "copy"):
516 raise ValueError("Only 'copy' and 'auto' transfer modes are supported.")
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)
554 if write_index:
555 index = ZipIndex.from_artifact_map(refs, artifact_map, destination)
556 index.write_index(destination)
558 return artifact_map
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)
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)
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
602 if response.artifact is None:
603 if full_check:
604 return DatasetExistence.RECORDED
605 else:
606 return DatasetExistence.RECORDED | DatasetExistence._ASSUMED
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
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}
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()
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()
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()
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()
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()
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()
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()
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()
710 @property
711 def run(self) -> str | None:
712 # Docstring inherited.
713 return self._registry_defaults.get().run
715 @property
716 def registry(self) -> Registry:
717 return self._registry
719 @contextmanager
720 def query(self) -> Iterator[Query]:
721 driver = RemoteQueryDriver(self, self._connection)
722 with driver:
723 query = Query(driver)
724 yield query
726 @contextmanager
727 def _query_all_datasets_by_page(
728 self, args: QueryAllDatasetsParameters
729 ) -> Iterator[Iterator[list[DatasetRef]]]:
730 universe = self.dimensions
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)
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()
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)
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 )
785 def close(self) -> None:
786 pass
788 def _expand_data_ids(self, data_ids: Iterable[DataCoordinate]) -> list[DataCoordinate]:
789 return expand_data_ids(data_ids, self.dimensions, self.query, None)
791 @property
792 def _file_transfer_source(self) -> RemoteFileTransferSource:
793 return RemoteFileTransferSource(self._connection)
795 def __str__(self) -> str:
796 return f"RemoteButler({self._connection.server_url})"
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.
808 return self._registry_defaults.get().dataId.to_simple(minimal=True).dataId
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}')")
816 return get_file_response.artifact
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)))
826class _RemoteButlerCacheData:
827 def __init__(self) -> None:
828 self.dimensions: DimensionUniverse | None = None
829 self.dataset_types: dict[str, DatasetType] = {}
832class RemoteButlerCache(LockedObject[_RemoteButlerCacheData]):
833 def __init__(self) -> None:
834 super().__init__(_RemoteButlerCacheData())