Coverage for python/lsst/resources/s3.py: 88%

364 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-23 09:29 +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. 

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

53 

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 

61 

62try: 

63 import s3fs 

64 from fsspec.spec import AbstractFileSystem 

65except ImportError: 

66 if not TYPE_CHECKING: 

67 s3fs = None 

68 AbstractFileSystem = type 

69 

70if TYPE_CHECKING: 

71 from .utils import TransactionProtocol 

72 

73 

74log = logging.getLogger(__name__) 

75 

76 

77class ProgressPercentage: 

78 """Progress bar for S3 file uploads. 

79 

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

91 

92 log_level = logging.DEBUG 

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

94 

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 

102 

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 ) 

117 

118 

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. 

122 

123 Parameters 

124 ---------- 

125 maybe_bool_str : `str` 

126 The value to parse 

127 

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

141 

142 return maybe_bool 

143 

144 

145class S3ResourcePath(ResourcePath): 

146 """S3 URI resource path implementation class. 

147 

148 Notes 

149 ----- 

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

151 environment variable is inspected: 

152 

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

158 

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

162 

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" 

169 

170 use_threads = _parse_string_to_maybe_bool(use_threads_str) 

171 

172 return use_threads 

173 

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 

181 

182 @property 

183 def _transfer_config(self) -> TransferConfig: 

184 if self.use_threads is None: 

185 self.use_threads = self._environ_use_threads 

186 

187 if self.use_threads is None: 

188 transfer_config = TransferConfig() 

189 else: 

190 transfer_config = TransferConfig(use_threads=self.use_threads) 

191 

192 return transfer_config 

193 

194 @property 

195 def client(self) -> BaseClient: 

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

197 return getS3Client(self._profile) 

198 

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 

203 

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

223 

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

226 

227 return bucket 

228 

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) 

241 

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

243 

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) 

252 

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) 

267 

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) 

280 

281 # Update with error information. 

282 results.update(errored) 

283 

284 return results 

285 

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) 

292 

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 

312 

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

324 

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

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

327 

328 # Remap error object keys to uris. 

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

330 

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 

347 

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 

356 

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 

366 

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 ) 

380 

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 

392 

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 

403 

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) 

410 

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) 

417 

418 return ResourceInfo( 

419 uri=str(self), 

420 is_file=is_file, 

421 size=size, 

422 last_modified=last_modified, 

423 checksums=checksums, 

424 ) 

425 

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 

437 

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 

454 

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) 

461 

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

467 

468 if not self.dirLike: 

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

470 

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) 

474 

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. 

480 

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 

500 

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

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

503 

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) 

525 

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

527 

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. 

533 

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. 

545 

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 

565 

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. 

569 

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 

585 

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 

599 

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. 

609 

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

631 

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 ) 

643 

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 

652 

653 if not overwrite and self.exists(): 

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

655 

656 if transfer == "auto": 

657 transfer = self.transferDefault 

658 

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

660 timer_args = (src, self) 

661 

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) 

668 

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) 

682 

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

688 

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. 

694 

695 Parameters 

696 ---------- 

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

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

699 

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

712 

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) 

715 

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

717 

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 

732 

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 

737 

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) 

742 

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

748 

749 filenames.extend(found_files) 

750 

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 

759 

760 for dir in dirnames: 

761 new_uri = self.join(dir) 

762 yield from new_uri.walk(file_filter) 

763 

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 

781 

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) 

785 

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) 

789 

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