Coverage for python/lsst/resources/s3.py: 88%
364 statements
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-22 09:18 +0000
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-22 09:18 +0000
1# This file is part of lsst-resources.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# Use of this source code is governed by a 3-clause BSD-style
10# license that can be found in the LICENSE file.
12from __future__ import annotations
14__all__ = ("S3ResourcePath",)
16import concurrent.futures
17import contextlib
18import datetime
19import io
20import logging
21import os
22import re
23import sys
24import threading
25from collections import defaultdict
26from collections.abc import Generator, Iterable, Iterator
27from functools import cache, cached_property
28from typing import IO, TYPE_CHECKING, cast
30from botocore.client import BaseClient
31from botocore.exceptions import ClientError
33from lsst.utils.iteration import chunk_iterable
34from lsst.utils.timer import time_this
36from ._resourceHandles._baseResourceHandle import ResourceHandleProtocol
37from ._resourceHandles._s3ResourceHandle import S3ResourceHandle
38from ._resourcePath import MBulkResult, ResourceInfo, ResourcePath
39from .s3utils import (
40 _get_s3_connection_parameters,
41 _s3_disable_bucket_validation,
42 _s3_should_validate_bucket,
43 all_retryable_errors,
44 backoff,
45 bucketExists,
46 getS3Client,
47 max_retry_time,
48 retryable_io_errors,
49 s3CheckFileExists,
50 translate_client_error,
51)
52from .utils import _get_num_workers
54try:
55 from boto3.s3.transfer import TransferConfig
56except ImportError:
57 # Hidden from type checkers so that the names above keep the types they
58 # have when the optional dependency is installed.
59 if not TYPE_CHECKING:
60 TransferConfig = None
62try:
63 import s3fs
64 from fsspec.spec import AbstractFileSystem
65except ImportError:
66 if not TYPE_CHECKING:
67 s3fs = None
68 AbstractFileSystem = type
70if TYPE_CHECKING:
71 from .utils import TransactionProtocol
74log = logging.getLogger(__name__)
77class ProgressPercentage:
78 """Progress bar for S3 file uploads.
80 Parameters
81 ----------
82 file : `ResourcePath`
83 Resource that is relevant to the progress percentage. The size of this
84 resource will be used to determine progress. The name will be used
85 in the log messages unless overridden by ``file_for_msg``.
86 file_for_msg : `ResourcePath` or `None`, optional
87 Resource name to include in log messages in preference to ``file``.
88 msg : `str`, optional
89 Message text to be included in every progress log message.
90 """
92 log_level = logging.DEBUG
93 """Default log level to use when issuing a message."""
95 def __init__(self, file: ResourcePath, file_for_msg: ResourcePath | None = None, msg: str = ""):
96 self._filename = file
97 self._file_for_msg = str(file_for_msg) if file_for_msg is not None else str(file)
98 self._size = file.size()
99 self._seen_so_far = 0
100 self._lock = threading.Lock()
101 self._msg = msg
103 def __call__(self, bytes_amount: int) -> None:
104 # To simplify, assume this is hooked up to a single filename
105 with self._lock:
106 self._seen_so_far += bytes_amount
107 percentage = (100 * self._seen_so_far) // self._size
108 log.log(
109 self.log_level,
110 "%s %s %s / %s (%s%%)",
111 self._msg,
112 self._file_for_msg,
113 self._seen_so_far,
114 self._size,
115 percentage,
116 )
119@cache
120def _parse_string_to_maybe_bool(maybe_bool_str: str) -> bool | None:
121 """Map a string to either a boolean value or None.
123 Parameters
124 ----------
125 maybe_bool_str : `str`
126 The value to parse
128 Results
129 -------
130 maybe_bool : `bool` or `None`
131 The parsed value.
132 """
133 if maybe_bool_str.lower() in ["t", "true", "yes", "y", "1"]:
134 maybe_bool = True
135 elif maybe_bool_str.lower() in ["f", "false", "no", "n", "0"]:
136 maybe_bool = False
137 elif maybe_bool_str.lower() in ["none", ""]: 137 ↛ 140line 137 didn't jump to line 140 because the condition on line 137 was always true
138 maybe_bool = None
139 else:
140 raise ValueError(f'Value of "{maybe_bool_str}" is not True, False, or None.')
142 return maybe_bool
145class S3ResourcePath(ResourcePath):
146 """S3 URI resource path implementation class.
148 Notes
149 -----
150 In order to configure the behavior of instances of this class, the
151 environment variable is inspected:
153 - LSST_S3_USE_THREADS: May be True, False, or None. Sets whether threading
154 is used for downloads, with a value of None defaulting to boto's default
155 value. Users may wish to set it to False when the downloads will be started
156 within threads other than python's main thread.
157 """
159 use_threads: bool | None = None
160 """Explicitly turn on or off threading in use of boto's download_fileobj.
161 Setting this to None results in boto's default behavior."""
163 @cached_property
164 def _environ_use_threads(self) -> bool | None:
165 try:
166 use_threads_str = os.environ["LSST_S3_USE_THREADS"]
167 except KeyError:
168 use_threads_str = "None"
170 use_threads = _parse_string_to_maybe_bool(use_threads_str)
172 return use_threads
174 @contextlib.contextmanager
175 def _use_threads_temp_override(self, multithreaded: bool) -> Generator[None]:
176 """Temporarily override the value of use_threads."""
177 original = self.use_threads
178 self.use_threads = multithreaded
179 yield
180 self.use_threads = original
182 @property
183 def _transfer_config(self) -> TransferConfig:
184 if self.use_threads is None:
185 self.use_threads = self._environ_use_threads
187 if self.use_threads is None:
188 transfer_config = TransferConfig()
189 else:
190 transfer_config = TransferConfig(use_threads=self.use_threads)
192 return transfer_config
194 @property
195 def client(self) -> BaseClient:
196 """Client object to address remote resource."""
197 return getS3Client(self._profile)
199 @property
200 def _profile(self) -> str | None:
201 """Profile name to use for looking up S3 credentials and endpoint."""
202 return self._uri.username
204 @property
205 def _bucket(self) -> str:
206 """S3 bucket where the files are stored."""
207 # Notionally the bucket is stored in the 'hostname' part of the URI.
208 # However, Ceph S3 uses a "multi-tenant" syntax for bucket names in the
209 # form 'tenant:bucket'. The part after the colon is parsed as the port
210 # portion of the URI, and urllib throws an exception if you try to read
211 # a non-integer port value. So manually split off this portion of the
212 # URI.
213 split = self._uri.netloc.split("@")
214 num_components = len(split)
215 if num_components == 2:
216 # There is a profile@ portion of the URL, so take the second half.
217 bucket = split[1]
218 elif num_components == 1: 218 ↛ 222line 218 didn't jump to line 222 because the condition on line 218 was always true
219 # There is no profile@, so take the whole netloc.
220 bucket = split[0]
221 else:
222 raise ValueError(f"Unexpected extra '@' in S3 URI: '{str(self)}'")
224 if not bucket: 224 ↛ 225line 224 didn't jump to line 225 because the condition on line 224 was never true
225 raise ValueError(f"S3 URI does not include bucket name: '{str(self)}'")
227 return bucket
229 @classmethod
230 def _mexists(
231 cls, uris: Iterable[ResourcePath], *, num_workers: int | None = None
232 ) -> dict[ResourcePath, bool]:
233 # Force client to be created for each profile before creating threads.
234 profiles = set[str | None]()
235 for path in uris:
236 if path.scheme == "s3": 236 ↛ 235line 236 didn't jump to line 235 because the condition on line 236 was always true
237 path = cast(S3ResourcePath, path)
238 profiles.add(path._profile)
239 for profile in profiles:
240 getS3Client(profile)
242 return super()._mexists(uris, num_workers=num_workers)
244 @classmethod
245 def _mremove(cls, uris: Iterable[ResourcePath]) -> dict[ResourcePath, MBulkResult]:
246 # Delete multiple objects in one API call.
247 # Must group by profile and bucket.
248 grouped_uris: dict[tuple[str | None, str], list[S3ResourcePath]] = defaultdict(list)
249 for uri in uris:
250 uri = cast(S3ResourcePath, uri)
251 grouped_uris[uri._profile, uri._bucket].append(uri)
253 results: dict[ResourcePath, MBulkResult] = {}
254 for related_uris in grouped_uris.values():
255 # API requires no more than 1000 per call.
256 chunk_num = 0
257 chunks: list[tuple[ResourcePath, ...]] = []
258 key_to_uri: dict[str, ResourcePath] = {}
259 for chunk in chunk_iterable(related_uris, chunk_size=1_000):
260 for uri in chunk:
261 key = uri.relativeToPathRoot
262 key_to_uri[key] = uri
263 # Default to assuming everything worked.
264 results[uri] = MBulkResult(True, None)
265 chunk_num += 1
266 chunks.append(chunk)
268 # Bulk remove.
269 with time_this(
270 log,
271 msg="Bulk delete; %d chunk%s; totalling %d dataset%s",
272 args=(
273 len(chunks),
274 "s" if len(chunks) != 1 else "",
275 len(related_uris),
276 "s" if len(related_uris) != 1 else "",
277 ),
278 ):
279 errored = cls._mremove_select(chunks)
281 # Update with error information.
282 results.update(errored)
284 return results
286 @classmethod
287 def _mremove_select(cls, chunks: list[tuple[ResourcePath, ...]]) -> dict[ResourcePath, MBulkResult]:
288 if len(chunks) == 1:
289 # Do the removal directly without futures.
290 return cls._delete_objects_wrapper(chunks[0])
291 return cls._mremove_with_pool(chunks)
293 @classmethod
294 def _mremove_with_pool(cls, chunks: list[tuple[ResourcePath, ...]]) -> dict[ResourcePath, MBulkResult]:
295 # Different name because different API to base class.
296 max_workers = _get_num_workers(cls._max_workers)
297 results: dict[ResourcePath, MBulkResult] = {}
298 with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as remove_executor:
299 future_remove = {
300 remove_executor.submit(cls._delete_objects_wrapper, chunk): i
301 for i, chunk in enumerate(chunks)
302 }
303 for future in concurrent.futures.as_completed(future_remove):
304 try:
305 results.update(future.result())
306 except Exception as e:
307 # The chunk utterly failed.
308 chunk = chunks[future_remove[future]]
309 for uri in chunk:
310 results[uri] = MBulkResult(False, e)
311 return results
313 @classmethod
314 def _delete_objects_wrapper(cls, uris: tuple[ResourcePath, ...]) -> dict[ResourcePath, MBulkResult]:
315 """Convert URIs to keys and call low-level API."""
316 if not uris: 316 ↛ 317line 316 didn't jump to line 317 because the condition on line 316 was never true
317 return {}
318 keys: list[dict[str, str]] = []
319 key_to_uri: dict[str, ResourcePath] = {}
320 for uri in uris:
321 key = uri.relativeToPathRoot
322 key_to_uri[key] = uri
323 keys.append({"Key": key})
325 first_uri = cast(S3ResourcePath, uris[0])
326 results = cls._delete_related_objects(first_uri.client, first_uri._bucket, keys)
328 # Remap error object keys to uris.
329 return {key_to_uri[key]: result for key, result in results.items()}
331 @classmethod
332 @backoff.on_exception(backoff.expo, retryable_io_errors, max_time=max_retry_time)
333 def _delete_related_objects(
334 cls, client: BaseClient, bucket: str, keys: list[dict[str, str]]
335 ) -> dict[str, MBulkResult]:
336 # Delete multiple objects from the same bucket, allowing for backoff
337 # retry.
338 response = client.delete_objects(Bucket=bucket, Delete={"Objects": keys, "Quiet": True})
339 # Use Quiet mode so we assume everything worked unless told otherwise.
340 # Only returning errors -- indexed by Key name.
341 errors: dict[str, MBulkResult] = {}
342 for errored_key in response.get("Errors", []): 342 ↛ 343line 342 didn't jump to line 343 because the loop on line 342 never started
343 errors[errored_key["Key"]] = MBulkResult(
344 False, ClientError({"Error": errored_key}, f"delete_objects: {errored_key['Key']}")
345 )
346 return errors
348 @backoff.on_exception(backoff.expo, retryable_io_errors, max_time=max_retry_time)
349 def exists(self) -> bool:
350 """Check that the S3 resource exists."""
351 if self.is_root:
352 # Only check for the bucket since the path is irrelevant
353 return bucketExists(self._bucket, self.client)
354 exists, _ = s3CheckFileExists(self, bucket=self._bucket, client=self.client)
355 return exists
357 @backoff.on_exception(backoff.expo, retryable_io_errors, max_time=max_retry_time)
358 def size(self) -> int:
359 """Return the size of the resource in bytes."""
360 if self.dirLike:
361 return 0
362 exists, sz = s3CheckFileExists(self, bucket=self._bucket, client=self.client)
363 if not exists:
364 raise FileNotFoundError(f"Resource {self} does not exist")
365 return sz
367 @backoff.on_exception(backoff.expo, retryable_io_errors, max_time=max_retry_time)
368 def get_info(self) -> ResourceInfo:
369 """Return lightweight metadata about this S3 resource."""
370 if self.is_root:
371 if not bucketExists(self._bucket, self.client): 371 ↛ 372line 371 didn't jump to line 372 because the condition on line 371 was never true
372 raise FileNotFoundError(f"Resource {self} does not exist")
373 return ResourceInfo(
374 uri=str(self),
375 is_file=False,
376 size=0,
377 last_modified=None,
378 checksums={},
379 )
381 try:
382 response = self.client.head_object(
383 Bucket=self._bucket,
384 Key=self.relativeToPathRoot,
385 ChecksumMode="ENABLED",
386 )
387 except (self.client.exceptions.NoSuchKey, self.client.exceptions.NoSuchBucket) as err:
388 raise FileNotFoundError(f"No such resource: {self}") from err
389 except ClientError as err:
390 translate_client_error(err, self)
391 raise
393 checksums = {}
394 for response_key, checksum_name in (
395 ("ChecksumCRC32", "crc32"),
396 ("ChecksumCRC32C", "crc32c"),
397 ("ChecksumCRC64NVME", "crc64nvme"),
398 ("ChecksumSHA1", "sha1"),
399 ("ChecksumSHA256", "sha256"),
400 ):
401 if value := response.get(response_key):
402 checksums[checksum_name] = value
404 last_modified = response.get("LastModified")
405 if last_modified is not None: 405 ↛ 415line 405 didn't jump to line 415 because the condition on line 405 was always true
406 if getattr(last_modified, "tzinfo", None) is None: 406 ↛ 407line 406 didn't jump to line 407 because the condition on line 406 was never true
407 last_modified = last_modified.replace(tzinfo=datetime.UTC)
408 else:
409 last_modified = last_modified.astimezone(datetime.UTC)
411 # For ResourcePath usage a dirLike object with zero size is a directory
412 # but in the general case anyone can create an object with a trailing
413 # `/` and treat it as a file. For self-consistency with ResourcePath
414 # call it a file if it has size > 0 even if dirLike.
415 size = response["ContentLength"]
416 is_file = (self.dirLike is not True) or (size > 0)
418 return ResourceInfo(
419 uri=str(self),
420 is_file=is_file,
421 size=size,
422 last_modified=last_modified,
423 checksums=checksums,
424 )
426 @backoff.on_exception(backoff.expo, retryable_io_errors, max_time=max_retry_time)
427 def remove(self) -> None:
428 """Remove the resource."""
429 # https://github.com/boto/boto3/issues/507 - there is no
430 # way of knowing if the file was actually deleted except
431 # for checking all the keys again, reponse is HTTP 204 OK
432 # response all the time
433 try:
434 self.client.delete_object(Bucket=self._bucket, Key=self.relativeToPathRoot)
435 except (self.client.exceptions.NoSuchKey, self.client.exceptions.NoSuchBucket) as err:
436 raise FileNotFoundError("No such resource: {self}") from err
438 @backoff.on_exception(backoff.expo, all_retryable_errors, max_time=max_retry_time)
439 def read(self, size: int = -1) -> bytes:
440 args = {}
441 if size > 0:
442 args["Range"] = f"bytes=0-{size - 1}"
443 try:
444 response = self.client.get_object(Bucket=self._bucket, Key=self.relativeToPathRoot, **args)
445 except (self.client.exceptions.NoSuchKey, self.client.exceptions.NoSuchBucket) as err:
446 raise FileNotFoundError(f"No such resource: {self}") from err
447 except ClientError as err:
448 translate_client_error(err, self)
449 raise
450 with time_this(log, msg="Read from %s", args=(self,)):
451 body = response["Body"].read()
452 response["Body"].close()
453 return body
455 @backoff.on_exception(backoff.expo, all_retryable_errors, max_time=max_retry_time)
456 def write(self, data: bytes, overwrite: bool = True) -> None:
457 if not overwrite and self.exists():
458 raise FileExistsError(f"Remote resource {self} exists and overwrite has been disabled")
459 with time_this(log, msg="Write to %s", args=(self,)):
460 self.client.put_object(Bucket=self._bucket, Key=self.relativeToPathRoot, Body=data)
462 @backoff.on_exception(backoff.expo, all_retryable_errors, max_time=max_retry_time)
463 def mkdir(self) -> None:
464 """Write a directory key to S3."""
465 if not bucketExists(self._bucket, self.client):
466 raise ValueError(f"Bucket {self._bucket} does not exist for {self}!")
468 if not self.dirLike:
469 raise NotADirectoryError(f"Can not create a 'directory' for file-like URI {self}")
471 # don't create S3 key when root is at the top-level of an Bucket
472 if self.path != "/":
473 self.client.put_object(Bucket=self._bucket, Key=self.relativeToPathRoot)
475 @backoff.on_exception(backoff.expo, all_retryable_errors, max_time=max_retry_time)
476 def _download_file(
477 self, local_file: IO | ResourceHandleProtocol, progress: ProgressPercentage | None
478 ) -> None:
479 """Download the remote resource to a local file.
481 Helper routine for _as_local to allow backoff without regenerating
482 the temporary file.
483 """
484 try:
485 self.client.download_fileobj(
486 self._bucket,
487 self.relativeToPathRoot,
488 local_file,
489 Callback=progress,
490 Config=self._transfer_config,
491 )
492 except (
493 self.client.exceptions.NoSuchKey,
494 self.client.exceptions.NoSuchBucket,
495 ) as err:
496 raise FileNotFoundError(f"No such resource: {self}") from err
497 except ClientError as err:
498 translate_client_error(err, self)
499 raise
501 def to_fsspec(self) -> tuple[AbstractFileSystem, str]:
502 """Return an abstract file system and path that can be used by fsspec.
504 Returns
505 -------
506 fs : `fsspec.spec.AbstractFileSystem`
507 A file system object suitable for use with the returned path.
508 path : `str`
509 A path that can be opened by the file system object.
510 """
511 if s3fs is None: 511 ↛ 512line 511 didn't jump to line 512 because the condition on line 511 was never true
512 raise ImportError("s3fs is not available")
513 # Must remove the profile from the URL and form it again.
514 endpoint_config = _get_s3_connection_parameters(self._profile)
515 s3 = s3fs.S3FileSystem(
516 profile=endpoint_config.profile,
517 endpoint_url=endpoint_config.endpoint_url,
518 key=endpoint_config.access_key_id,
519 secret=endpoint_config.secret_access_key,
520 )
521 if not _s3_should_validate_bucket(): 521 ↛ 524line 521 didn't jump to line 524 because the condition on line 521 was never true
522 # Accessing the s3 property forces the boto client to be
523 # constructed and cached and allows the validation to be removed.
524 _s3_disable_bucket_validation(s3.s3)
526 return s3, f"{self._bucket}/{self.relativeToPathRoot}"
528 @contextlib.contextmanager
529 def _as_local(
530 self, multithreaded: bool = True, tmpdir: ResourcePath | None = None
531 ) -> Generator[ResourcePath]:
532 """Download object from S3 and place in temporary directory.
534 Parameters
535 ----------
536 multithreaded : `bool`, optional
537 If `True` the transfer will be allowed to attempt to improve
538 throughput by using parallel download streams. This may of no
539 effect if the URI scheme does not support parallel streams or
540 if a global override has been applied. If `False` parallel
541 streams will be disabled.
542 tmpdir : `ResourcePath` or `None`, optional
543 Explicit override of the temporary directory to use for remote
544 downloads.
546 Returns
547 -------
548 local_uri : `ResourcePath`
549 A URI to a local POSIX file corresponding to a local temporary
550 downloaded copy of the resource.
551 """
552 with (
553 ResourcePath.temporary_uri(prefix=tmpdir, suffix=self.getExtension(), delete=True) as tmp_uri,
554 self._use_threads_temp_override(multithreaded),
555 time_this(log, msg="Downloading %s to local file", args=(self,)),
556 ):
557 progress = (
558 ProgressPercentage(self, msg="Downloading:")
559 if log.isEnabledFor(ProgressPercentage.log_level)
560 else None
561 )
562 with tmp_uri.open("wb") as tmpFile:
563 self._download_file(tmpFile, progress)
564 yield tmp_uri
566 @backoff.on_exception(backoff.expo, all_retryable_errors, max_time=max_retry_time)
567 def _upload_file(self, local_file: ResourcePath, progress: ProgressPercentage | None) -> None:
568 """Upload a local file with backoff.
570 Helper method to wrap file uploading in backoff for transfer_from.
571 """
572 try:
573 self.client.upload_file(
574 local_file.ospath,
575 self._bucket,
576 self.relativeToPathRoot,
577 Callback=progress,
578 Config=self._transfer_config,
579 )
580 except self.client.exceptions.NoSuchBucket as err:
581 raise NotADirectoryError(f"Target does not exist: {err}") from err
582 except ClientError as err:
583 translate_client_error(err, self)
584 raise
586 @backoff.on_exception(backoff.expo, all_retryable_errors, max_time=max_retry_time)
587 def _copy_from(self, src: S3ResourcePath) -> None:
588 copy_source = {
589 "Bucket": src._bucket,
590 "Key": src.relativeToPathRoot,
591 }
592 try:
593 self.client.copy_object(CopySource=copy_source, Bucket=self._bucket, Key=self.relativeToPathRoot)
594 except (self.client.exceptions.NoSuchKey, self.client.exceptions.NoSuchBucket) as err:
595 raise FileNotFoundError(f"No such resource to transfer: {src} -> {self}") from err
596 except ClientError as err:
597 translate_client_error(err, self)
598 raise
600 def transfer_from(
601 self,
602 src: ResourcePath,
603 transfer: str = "copy",
604 overwrite: bool = False,
605 transaction: TransactionProtocol | None = None,
606 multithreaded: bool = True,
607 ) -> None:
608 """Transfer the current resource to an S3 bucket.
610 Parameters
611 ----------
612 src : `ResourcePath`
613 Source URI.
614 transfer : `str`
615 Mode to use for transferring the resource. Supports the following
616 options: copy.
617 overwrite : `bool`, optional
618 Allow an existing file to be overwritten. Defaults to `False`.
619 transaction : `~lsst.resources.utils.TransactionProtocol`, optional
620 Currently unused.
621 multithreaded : `bool`, optional
622 If `True` the transfer will be allowed to attempt to improve
623 throughput by using parallel download streams. This may of no
624 effect if the URI scheme does not support parallel streams or
625 if a global override has been applied. If `False` parallel
626 streams will be disabled.
627 """
628 # Fail early to prevent delays if remote resources are requested
629 if transfer not in self.transferModes:
630 raise ValueError(f"Transfer mode '{transfer}' not supported by URI scheme {self.scheme}")
632 # Existence checks cost time so do not call this unless we know
633 # that debugging is enabled.
634 if log.isEnabledFor(logging.DEBUG): 634 ↛ 645line 634 didn't jump to line 645 because the condition on line 634 was always true
635 log.debug(
636 "Transferring %s [exists: %s] -> %s [exists: %s] (transfer=%s)",
637 src,
638 src.exists(),
639 self,
640 self.exists(),
641 transfer,
642 )
644 # Short circuit if the URIs are identical immediately.
645 if self == src:
646 log.debug(
647 "Target and destination URIs are identical: %s, returning immediately."
648 " No further action required.",
649 self,
650 )
651 return
653 if not overwrite and self.exists():
654 raise FileExistsError(f"Destination path '{self}' already exists.")
656 if transfer == "auto":
657 transfer = self.transferDefault
659 timer_msg = "Transfer from %s to %s"
660 timer_args = (src, self)
662 if isinstance(src, type(self)) and self.client == src.client:
663 # Looks like an S3 remote uri so we can use direct copy.
664 # This only works if the source and destination are using the same
665 # S3 endpoint and profile.
666 with time_this(log, msg=timer_msg, args=timer_args):
667 self._copy_from(src)
669 else:
670 # Use local file and upload it
671 with src.as_local(multithreaded=multithreaded) as local_uri:
672 progress = (
673 ProgressPercentage(local_uri, file_for_msg=src, msg="Uploading:")
674 if log.isEnabledFor(ProgressPercentage.log_level)
675 else None
676 )
677 with (
678 time_this(log, msg=timer_msg, args=timer_args),
679 self._use_threads_temp_override(multithreaded),
680 ):
681 self._upload_file(local_uri, progress)
683 # This was an explicit move requested from a remote resource
684 # try to remove that resource
685 if transfer == "move":
686 # Transactions do not work here
687 src.remove()
689 @backoff.on_exception(backoff.expo, all_retryable_errors, max_time=max_retry_time)
690 def walk(
691 self, file_filter: str | re.Pattern | None = None
692 ) -> Iterator[list | tuple[ResourcePath, list[str], list[str]]]:
693 """Walk the directory tree returning matching files and directories.
695 Parameters
696 ----------
697 file_filter : `str` or `re.Pattern`, optional
698 Regex to filter out files from the list before it is returned.
700 Yields
701 ------
702 dirpath : `ResourcePath`
703 Current directory being examined.
704 dirnames : `list` of `str`
705 Names of subdirectories within dirpath.
706 filenames : `list` of `str`
707 Names of all the files within dirpath.
708 """
709 # We pretend that S3 uses directories and files and not simply keys
710 if not (self.isdir() or self.is_root):
711 raise ValueError(f"Can not walk a non-directory URI: {self}")
713 if isinstance(file_filter, str): 713 ↛ 714line 713 didn't jump to line 714 because the condition on line 713 was never true
714 file_filter = re.compile(file_filter)
716 s3_paginator = self.client.get_paginator("list_objects_v2")
718 # Limit each query to a single "directory" to match os.walk
719 # We could download all keys at once with no delimiter and work
720 # it out locally but this could potentially lead to large memory
721 # usage for millions of keys. It will also make the initial call
722 # to this method potentially very slow. If making this method look
723 # like os.walk was not required, we could query all keys with
724 # pagination and return them in groups of 1000, but that would
725 # be a different interface since we can't guarantee we would get
726 # them all grouped properly across the 1000 limit boundary.
727 prefix = self.relativeToPathRoot if not self.is_root else ""
728 prefix_len = len(prefix)
729 dirnames = []
730 filenames = []
731 files_there = False
733 for page in s3_paginator.paginate(Bucket=self._bucket, Prefix=prefix, Delimiter="/"):
734 # All results are returned as full key names and we must
735 # convert them back to the root form. The prefix is fixed
736 # and delimited so that is a simple trim
738 # Directories are reported in the CommonPrefixes result
739 # which reports the entire key and must be stripped.
740 found_dirs = [dir["Prefix"][prefix_len:] for dir in page.get("CommonPrefixes", ())]
741 dirnames.extend(found_dirs)
743 found_files = [file["Key"][prefix_len:] for file in page.get("Contents", ())]
744 if found_files:
745 files_there = True
746 if file_filter is not None:
747 found_files = [f for f in found_files if file_filter.search(f)]
749 filenames.extend(found_files)
751 # Directories do not exist so we can't test for them. If no files
752 # or directories were found though, this means that it effectively
753 # does not exist and we should match os.walk() behavior and return
754 # immediately.
755 if not dirnames and not files_there:
756 return
757 else:
758 yield self, dirnames, filenames
760 for dir in dirnames:
761 new_uri = self.join(dir)
762 yield from new_uri.walk(file_filter)
764 @contextlib.contextmanager
765 def _openImpl(
766 self,
767 mode: str = "r",
768 *,
769 encoding: str | None = None,
770 ) -> Generator[ResourceHandleProtocol]:
771 with S3ResourceHandle(mode, log, self) as handle:
772 if "b" in mode:
773 yield handle
774 else:
775 if encoding is None:
776 encoding = sys.getdefaultencoding()
777 # cast because the protocol is compatible, but does not have
778 # BytesIO in the inheritance tree
779 with io.TextIOWrapper(cast(io.BytesIO, handle), encoding=encoding, write_through=True) as sub:
780 yield sub
782 def generate_presigned_get_url(self, *, expiration_time_seconds: int) -> str:
783 # Docstring inherited
784 return self._generate_presigned_url("get_object", expiration_time_seconds)
786 def generate_presigned_put_url(self, *, expiration_time_seconds: int) -> str:
787 # Docstring inherited
788 return self._generate_presigned_url("put_object", expiration_time_seconds)
790 def _generate_presigned_url(self, method: str, expiration_time_seconds: int) -> str:
791 url = self.client.generate_presigned_url(
792 method,
793 Params={"Bucket": self._bucket, "Key": self.relativeToPathRoot},
794 ExpiresIn=expiration_time_seconds,
795 )
796 if self.fragment:
797 resource = ResourcePath(url)
798 url = str(resource.replace(fragment=self.fragment))
799 return url