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

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. 

11 

12from __future__ import annotations 

13 

14__all__ = ("S3ResourcePath",) 

15 

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 

29 

30from botocore.client import BaseClient 

31from botocore.exceptions import ClientError 

32 

33from lsst.utils.iteration import chunk_iterable 

34from lsst.utils.timer import time_this 

35 

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 

60 

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 

68 

69try: 

70 import s3fs 

71 from fsspec.spec import AbstractFileSystem 

72except ImportError: 

73 if not TYPE_CHECKING: 

74 s3fs = None 

75 AbstractFileSystem = type 

76 

77if TYPE_CHECKING: 

78 from .utils import TransactionProtocol 

79 

80 

81log = logging.getLogger(__name__) 

82 

83 

84class ProgressPercentage: 

85 """Progress bar for S3 file uploads. 

86 

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 """ 

98 

99 log_level = logging.DEBUG 

100 """Default log level to use when issuing a message.""" 

101 

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 

109 

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 ) 

124 

125 

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. 

129 

130 Parameters 

131 ---------- 

132 maybe_bool_str : `str` 

133 The value to parse 

134 

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.') 

148 

149 return maybe_bool 

150 

151 

152class S3ResourcePath(ResourcePath): 

153 """S3 URI resource path implementation class. 

154 

155 Notes 

156 ----- 

157 In order to configure the behavior of instances of this class, the 

158 environment variable is inspected: 

159 

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 """ 

165 

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.""" 

169 

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" 

176 

177 use_threads = _parse_string_to_maybe_bool(use_threads_str) 

178 

179 return use_threads 

180 

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 

188 

189 @property 

190 def _transfer_config(self) -> TransferConfig: 

191 if self.use_threads is None: 

192 self.use_threads = self._environ_use_threads 

193 

194 if self.use_threads is None: 

195 transfer_config = TransferConfig() 

196 else: 

197 transfer_config = TransferConfig(use_threads=self.use_threads) 

198 

199 return transfer_config 

200 

201 @property 

202 def client(self) -> BaseClient: 

203 """Client object to address remote resource.""" 

204 return getS3Client(self._profile) 

205 

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 

210 

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)}'") 

230 

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)}'") 

233 

234 return bucket 

235 

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) 

248 

249 return super()._mexists(uris, num_workers=num_workers) 

250 

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) 

259 

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) 

274 

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) 

287 

288 # Update with error information. 

289 results.update(errored) 

290 

291 return results 

292 

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) 

306 

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 

333 

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}) 

345 

346 first_uri = cast(S3ResourcePath, uris[0]) 

347 results = cls._delete_related_objects(first_uri.client, first_uri._bucket, keys) 

348 

349 # Remap error object keys to uris. 

350 return {key_to_uri[key]: result for key, result in results.items()} 

351 

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 

368 

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 

377 

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 

387 

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 ) 

401 

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 

413 

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 

424 

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) 

431 

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) 

438 

439 return ResourceInfo( 

440 uri=str(self), 

441 is_file=is_file, 

442 size=size, 

443 last_modified=last_modified, 

444 checksums=checksums, 

445 ) 

446 

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 

458 

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 

475 

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) 

482 

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}!") 

488 

489 if not self.dirLike: 

490 raise NotADirectoryError(f"Can not create a 'directory' for file-like URI {self}") 

491 

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) 

495 

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. 

501 

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 

521 

522 def to_fsspec(self) -> tuple[AbstractFileSystem, str]: 

523 """Return an abstract file system and path that can be used by fsspec. 

524 

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) 

546 

547 return s3, f"{self._bucket}/{self.relativeToPathRoot}" 

548 

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. 

554 

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. 

566 

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 

586 

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. 

590 

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 

606 

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 

620 

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. 

630 

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}") 

652 

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 ) 

664 

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 

673 

674 if not overwrite and self.exists(): 

675 raise FileExistsError(f"Destination path '{self}' already exists.") 

676 

677 if transfer == "auto": 

678 transfer = self.transferDefault 

679 

680 timer_msg = "Transfer from %s to %s" 

681 timer_args = (src, self) 

682 

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) 

689 

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) 

703 

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() 

709 

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. 

715 

716 Parameters 

717 ---------- 

718 file_filter : `str` or `re.Pattern`, optional 

719 Regex to filter out files from the list before it is returned. 

720 

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}") 

733 

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) 

736 

737 s3_paginator = self.client.get_paginator("list_objects_v2") 

738 

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 

753 

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 

758 

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) 

763 

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)] 

769 

770 filenames.extend(found_files) 

771 

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 

780 

781 for dir in dirnames: 

782 new_uri = self.join(dir) 

783 yield from new_uri.walk(file_filter) 

784 

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 

802 

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) 

806 

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) 

810 

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