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:50 +0000
« prev ^ index » next coverage.py v7.15.3, created at 2026-08-07 09:50 +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/>.
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 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)
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
108_LOG = logging.getLogger(__name__)
111class RemoteButler(Butler): # numpydoc ignore=PR02
112 """A `Butler` that can be used to connect through a remote server.
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.
130 Notes
131 -----
132 Instead of using this constructor, most users should use either
133 `Butler.from_config` or `RemoteButlerFactory`.
134 """
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
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()
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()
169 self._registry_defaults = DefaultsHolder(defaults)
170 self._registry = RemoteButlerRegistry(self, self._registry_defaults, self._connection)
171 defaults.finish(self._registry)
173 return self
175 def isWriteable(self) -> bool:
176 # Docstring inherited.
177 return False
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
189 @property
190 def collections(self) -> ButlerCollections:
191 """Object with methods for modifying and querying collections."""
192 return RemoteButlerCollections(self._registry_defaults, self._connection)
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
201 response = self._connection.get("universe")
202 model = parse_model(response, GetUniverseResponseModel)
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
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
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()
237 def transaction(self) -> AbstractContextManager[None]:
238 """Will always raise NotImplementedError.
239 Transactions are not supported by RemoteButler.
240 """
241 raise NotImplementedError()
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()
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)
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)
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)
317 return self._get_dataset_as_python_object(written_ref, read_ref, model, parameters)
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 )
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)
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)
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")
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)
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
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)
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
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
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.
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 )
485 response = self._connection.post("find_dataset", query)
487 model = parse_model(response, FindDatasetResponseModel)
488 if model.dataset_ref is None:
489 return None
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)
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}.")
512 if transfer not in ("auto", "copy"):
513 raise ValueError("Only 'copy' and 'auto' transfer modes are supported.")
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)
551 if write_index:
552 index = ZipIndex.from_artifact_map(refs, artifact_map, destination)
553 index.write_index(destination)
555 return artifact_map
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)
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)
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
599 if response.artifact is None:
600 if full_check:
601 return DatasetExistence.RECORDED
602 else:
603 return DatasetExistence.RECORDED | DatasetExistence._ASSUMED
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
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}
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()
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()
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()
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()
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()
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()
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()
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()
707 @property
708 def run(self) -> str | None:
709 # Docstring inherited.
710 return self._registry_defaults.get().run
712 @property
713 def registry(self) -> Registry:
714 return self._registry
716 @contextmanager
717 def query(self) -> Iterator[Query]:
718 driver = RemoteQueryDriver(self, self._connection)
719 with driver:
720 query = Query(driver)
721 yield query
723 @contextmanager
724 def _query_all_datasets_by_page(
725 self, args: QueryAllDatasetsParameters
726 ) -> Iterator[Iterator[list[DatasetRef]]]:
727 universe = self.dimensions
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)
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()
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)
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 )
782 def close(self) -> None:
783 pass
785 def _expand_data_ids(self, data_ids: Iterable[DataCoordinate]) -> list[DataCoordinate]:
786 return expand_data_ids(data_ids, self.dimensions, self.query, None)
788 @property
789 def _file_transfer_source(self) -> RemoteFileTransferSource:
790 return RemoteFileTransferSource(self._connection)
792 def __str__(self) -> str:
793 return f"RemoteButler({self._connection.server_url})"
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.
805 return self._registry_defaults.get().dataId.to_simple(minimal=True).dataId
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}')")
813 return get_file_response.artifact
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)))
823class _RemoteButlerCacheData:
824 def __init__(self) -> None:
825 self.dimensions: DimensionUniverse | None = None
826 self.dataset_types: dict[str, DatasetType] = {}
829class RemoteButlerCache(LockedObject[_RemoteButlerCacheData]):
830 def __init__(self) -> None:
831 super().__init__(_RemoteButlerCacheData())