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