Coverage for python/lsst/resources/s3utils.py: 79%

150 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__ = ( 

15 "_TooManyRequestsError", 

16 "all_retryable_errors", 

17 "backoff", 

18 "bucketExists", 

19 "clean_test_environment_for_s3", 

20 "getS3Client", 

21 "max_retry_time", 

22 "retryable_client_errors", 

23 "retryable_io_errors", 

24 "s3CheckFileExists", 

25) 

26 

27import functools 

28import os 

29import re 

30import urllib.parse 

31from collections.abc import Callable, Generator 

32from contextlib import contextmanager 

33from http.client import HTTPException, ImproperConnectionState 

34from types import ModuleType 

35from typing import TYPE_CHECKING, Any, NamedTuple, cast 

36from unittest.mock import patch 

37 

38from botocore.client import BaseClient 

39from botocore.exceptions import ClientError 

40from botocore.handlers import validate_bucket_name 

41from urllib3.exceptions import HTTPError, RequestError 

42from urllib3.util import Url, parse_url 

43 

44try: 

45 import boto3 

46except ImportError: 

47 # Hidden from type checkers so that ``boto3`` keeps the type it has when 

48 # the optional dependency is installed. 

49 if not TYPE_CHECKING: 

50 boto3 = None 

51 

52try: 

53 import botocore.config 

54except ImportError: 

55 if not TYPE_CHECKING: 

56 botocore = None 

57 

58 

59from ._resourcePath import ResourcePath 

60from .location import Location 

61from .utils import _get_num_workers 

62 

63# https://pypi.org/project/backoff/ 

64try: 

65 import backoff 

66except ImportError: 

67 

68 class Backoff: 

69 """Mock implementation of the backoff class.""" 

70 

71 @staticmethod 

72 def expo(func: Callable, *args: Any, **kwargs: Any) -> Callable: 

73 return func 

74 

75 @staticmethod 

76 def on_exception(func: Callable, *args: Any, **kwargs: Any) -> Callable: 

77 return func 

78 

79 backoff = cast(ModuleType, Backoff) 

80 

81 

82class _TooManyRequestsError(Exception): 

83 """Private exception that can be used for 429 retry. 

84 

85 botocore refuses to deal with 429 error itself so issues a generic 

86 ClientError. 

87 """ 

88 

89 pass 

90 

91 

92# settings for "backoff" retry decorators. these retries are belt-and- 

93# suspenders along with the retries built into Boto3, to account for 

94# semantic differences in errors between S3-like providers. 

95retryable_io_errors = ( 

96 # http.client 

97 ImproperConnectionState, 

98 HTTPException, 

99 # urllib3.exceptions 

100 RequestError, 

101 HTTPError, 

102 # built-ins 

103 TimeoutError, 

104 ConnectionError, 

105 # private 

106 _TooManyRequestsError, 

107) 

108 

109# Client error can include NoSuchKey so retry may not be the right 

110# thing. This may require more consideration if it is to be used. 

111retryable_client_errors = ( 

112 # botocore.exceptions 

113 ClientError, 

114 # built-ins 

115 PermissionError, 

116) 

117 

118 

119# Combine all errors into an easy package. For now client errors 

120# are not included. 

121all_retryable_errors = retryable_io_errors 

122max_retry_time = 60 

123 

124 

125@contextmanager 

126def clean_test_environment_for_s3() -> Generator[None]: 

127 """Reset S3 environment to ensure that unit tests with a mock S3 can't 

128 accidentally reference real infrastructure. 

129 """ 

130 with patch.dict( 

131 os.environ, 

132 { 

133 "AWS_ACCESS_KEY_ID": "test-access-key", 

134 "AWS_SECRET_ACCESS_KEY": "test-secret-access-key", 

135 "AWS_DEFAULT_REGION": "us-east-1", 

136 }, 

137 ) as patched_environ: 

138 for var in ( 

139 "S3_ENDPOINT_URL", 

140 "AWS_SECURITY_TOKEN", 

141 "AWS_SESSION_TOKEN", 

142 "AWS_PROFILE", 

143 "AWS_SHARED_CREDENTIALS_FILE", 

144 "AWS_CONFIG_FILE", 

145 ): 

146 patched_environ.pop(var, None) 

147 # Clear the cached boto3 S3 client instances. 

148 # This helps us avoid a potential situation where the client could be 

149 # instantiated before moto mocks are installed, which would prevent the 

150 # mocks from taking effect. 

151 _get_s3_client.cache_clear() 

152 yield 

153 

154 

155def getS3Client(profile: str | None = None) -> BaseClient: 

156 """Create a S3 client with AWS (default) or the specified endpoint. 

157 

158 Parameters 

159 ---------- 

160 profile : `str`, optional 

161 The name of an S3 profile describing which S3 service to use. 

162 

163 Returns 

164 ------- 

165 s3client : `botocore.client.S3` 

166 A client of the S3 service. 

167 

168 Notes 

169 ----- 

170 If an explicit profile name is specified, its configuration will be read 

171 from an environment variable named ``LSST_RESOURCES_S3_PROFILE_<profile>`` 

172 if it exists. Note that the name of the profile is case sensitive. This 

173 configuration is specified in the format: ``https://<access key ID>:<secret 

174 key>@<s3 endpoint hostname>``. If the access key ID or secret key values 

175 contain slashes, the slashes must be URI-encoded (replace "/" with "%2F"). 

176 

177 If profile is `None` or the profile environment variable was not set, the 

178 configuration is read from the environment variable ``S3_ENDPOINT_URL``. 

179 If it is not specified, the default AWS endpoint is used. 

180 

181 The access key ID and secret key are optional -- if not specified, they 

182 will be looked up via the `AWS credentials file 

183 <https://boto3.amazonaws.com/v1/documentation/api/latest/guide/credentials.html>`_. 

184 

185 If the environment variable LSST_DISABLE_BUCKET_VALIDATION exists 

186 and has a value that is not empty, "0", "f", "n", or "false" 

187 (case-insensitive), then bucket name validation is disabled. This 

188 disabling allows Ceph multi-tenancy colon separators to appear in 

189 bucket names. 

190 """ 

191 if boto3 is None: 191 ↛ 192line 191 didn't jump to line 192 because the condition on line 191 was never true

192 raise ModuleNotFoundError("Could not find boto3. Are you sure it is installed?") 

193 if botocore is None: 193 ↛ 194line 193 didn't jump to line 194 because the condition on line 193 was never true

194 raise ModuleNotFoundError("Could not find botocore. Are you sure it is installed?") 

195 

196 endpoint_config = _get_s3_connection_parameters(profile) 

197 

198 return _get_s3_client(endpoint_config, not _s3_should_validate_bucket()) 

199 

200 

201def _s3_should_validate_bucket() -> bool: 

202 """Indicate whether bucket validation should be enabled. 

203 

204 Returns 

205 ------- 

206 validate : `bool` 

207 If `True` bucket names should be validated. 

208 """ 

209 disable_value = os.environ.get("LSST_DISABLE_BUCKET_VALIDATION", "0") 

210 return bool(re.search(r"^(0|f|n|false)?$", disable_value, re.I)) 

211 

212 

213def _get_s3_connection_parameters(profile: str | None = None) -> _EndpointConfig: 

214 """Calculate the connection details. 

215 

216 Parameters 

217 ---------- 

218 profile : `str`, optional 

219 The name of an S3 profile describing which S3 service to use. 

220 

221 Returns 

222 ------- 

223 config : _EndPointConfig 

224 All the information necessary to connect to the bucket. 

225 """ 

226 endpoint = None 

227 if profile is not None: 

228 var_name = f"LSST_RESOURCES_S3_PROFILE_{profile}" 

229 endpoint = os.environ.get(var_name, None) 

230 if not endpoint: 

231 endpoint = os.environ.get("S3_ENDPOINT_URL", None) 

232 if not endpoint: 

233 endpoint = None # Handle "" 

234 

235 return _parse_endpoint_config(endpoint, profile) 

236 

237 

238def _s3_disable_bucket_validation(client: BaseClient) -> None: 

239 """Disable the bucket name validation in the client. 

240 

241 This removes the ``validate_bucket_name`` handler from the handlers 

242 registered for this client. 

243 

244 Parameters 

245 ---------- 

246 client : `boto3.client` 

247 The client to modify. 

248 """ 

249 client.meta.events.unregister("before-parameter-build.s3", validate_bucket_name) 

250 

251 

252@functools.lru_cache 

253def _get_s3_client(endpoint_config: _EndpointConfig, skip_validation: bool) -> BaseClient: 

254 # Helper function to cache the client for this endpoint 

255 # boto seems to assume it will always have at least 10 available. 

256 max_pool_size = max(_get_num_workers(), 10) 

257 config = botocore.config.Config( 

258 read_timeout=180, 

259 max_pool_connections=max_pool_size, 

260 retries={"mode": "adaptive", "max_attempts": 10}, 

261 ) 

262 

263 session = boto3.Session(profile_name=endpoint_config.profile) 

264 

265 client = session.client( 

266 "s3", 

267 endpoint_url=endpoint_config.endpoint_url, 

268 aws_access_key_id=endpoint_config.access_key_id, 

269 aws_secret_access_key=endpoint_config.secret_access_key, 

270 config=config, 

271 ) 

272 if skip_validation: 

273 _s3_disable_bucket_validation(client) 

274 return client 

275 

276 

277class _EndpointConfig(NamedTuple): 

278 endpoint_url: str | None = None 

279 access_key_id: str | None = None 

280 secret_access_key: str | None = None 

281 profile: str | None = None 

282 

283 

284def _parse_endpoint_config(endpoint: str | None, profile: str | None = None) -> _EndpointConfig: 

285 if not endpoint: 

286 return _EndpointConfig(profile=profile) 

287 

288 parsed = parse_url(endpoint) 

289 

290 # Strip the username/password portion of the URL from the result. 

291 endpoint_url = Url(host=parsed.host, path=parsed.path, port=parsed.port, scheme=parsed.scheme).url 

292 

293 access_key_id = None 

294 secret_access_key = None 

295 if parsed.auth: 

296 split = parsed.auth.split(":") 

297 if len(split) != 2: 

298 raise ValueError("S3 access key and secret not in expected format.") 

299 access_key_id, secret_access_key = split 

300 access_key_id = urllib.parse.unquote(access_key_id) 

301 secret_access_key = urllib.parse.unquote(secret_access_key) 

302 

303 if access_key_id is not None and secret_access_key is not None: 

304 # We already have the necessary configuration for the profile, so do 

305 # not pass the profile to boto3. boto3 will raise an exception if the 

306 # profile is not defined in its configuration file, whether or not it 

307 # needs to read the configuration from it. 

308 profile = None 

309 

310 return _EndpointConfig( 

311 endpoint_url=endpoint_url, 

312 access_key_id=access_key_id, 

313 secret_access_key=secret_access_key, 

314 profile=profile, 

315 ) 

316 

317 

318def s3CheckFileExists( 

319 path: Location | ResourcePath | str, 

320 bucket: str | None = None, 

321 client: BaseClient | None = None, 

322) -> tuple[bool, int]: 

323 """Return if the file exists in the bucket or not. 

324 

325 Parameters 

326 ---------- 

327 path : `Location`, `ResourcePath` or `str` 

328 Location or ResourcePath containing the bucket name and filepath. 

329 bucket : `str`, optional 

330 Name of the bucket in which to look. If provided, path will be assumed 

331 to correspond to be relative to the given bucket. 

332 client : `boto3.client`, optional 

333 S3 Client object to query, if not supplied boto3 will try to resolve 

334 the credentials as in order described in its manual_. 

335 

336 Returns 

337 ------- 

338 exists : `bool` 

339 True if key exists, False otherwise. 

340 size : `int` 

341 Size of the key, if key exists, in bytes, otherwise -1. 

342 

343 Notes 

344 ----- 

345 S3 Paths are sensitive to leading and trailing path separators. 

346 

347 .. _manual: https://boto3.amazonaws.com/v1/documentation/api/latest/guide/\ 

348 configuration.html#configuring-credentials 

349 """ 

350 if boto3 is None: 350 ↛ 351line 350 didn't jump to line 351 because the condition on line 350 was never true

351 raise ModuleNotFoundError("Could not find boto3. Are you sure it is installed?") 

352 

353 if client is None: 

354 client = getS3Client() 

355 

356 if isinstance(path, str): 

357 if bucket is not None: 

358 filepath = path 

359 else: 

360 uri = ResourcePath(path) 

361 bucket = uri.netloc 

362 filepath = uri.relativeToPathRoot 

363 elif isinstance(path, ResourcePath | Location): 363 ↛ 368line 363 didn't jump to line 368 because the condition on line 363 was always true

364 if bucket is None: 

365 bucket = path.netloc 

366 filepath = path.relativeToPathRoot 

367 else: 

368 raise TypeError(f"Unsupported path type: {path!r}.") 

369 

370 try: 

371 obj = client.head_object(Bucket=bucket, Key=filepath) 

372 return (True, obj["ContentLength"]) 

373 except client.exceptions.ClientError as err: 

374 # resource unreachable error means key does not exist 

375 errcode = err.response["ResponseMetadata"]["HTTPStatusCode"] 

376 if errcode == 404: 376 ↛ 384line 376 didn't jump to line 384 because the condition on line 376 was always true

377 return (False, -1) 

378 # head_object returns 404 when object does not exist only when user has 

379 # s3:ListBucket permission. If list permission does not exist a 403 is 

380 # returned. In practical terms this generally means that the file does 

381 # not exist, but it could also mean user lacks s3:GetObject permission: 

382 # https://docs.aws.amazon.com/AmazonS3/latest/API/RESTObjectHEAD.html 

383 # I don't think its possible to discern which case is it with certainty 

384 if errcode == 403: 

385 raise PermissionError( 

386 "Forbidden HEAD operation error occurred. " 

387 "Verify s3:ListBucket and s3:GetObject " 

388 "permissions are granted for your IAM user. " 

389 ) from err 

390 if errcode == 429: 

391 # boto3, incorrectly, does not automatically retry with 429 

392 # so instead we raise an explicit retry exception for backoff. 

393 raise _TooManyRequestsError(str(err)) from err 

394 raise 

395 

396 

397def bucketExists(bucketName: str, client: BaseClient | None = None) -> bool: 

398 """Check if the S3 bucket with the given name actually exists. 

399 

400 Parameters 

401 ---------- 

402 bucketName : `str` 

403 Name of the S3 Bucket. 

404 client : `boto3.client`, optional 

405 S3 Client object to query, if not supplied boto3 will try to resolve 

406 the credentials by calling `getS3Client`. 

407 

408 Returns 

409 ------- 

410 exists : `bool` 

411 True if it exists, False if no Bucket with specified parameters is 

412 found. 

413 """ 

414 if boto3 is None: 414 ↛ 415line 414 didn't jump to line 415 because the condition on line 414 was never true

415 raise ModuleNotFoundError("Could not find boto3. Are you sure it is installed?") 

416 

417 if client is None: 

418 client = getS3Client() 

419 try: 

420 client.get_bucket_location(Bucket=bucketName) 

421 return True 

422 except client.exceptions.NoSuchBucket: 

423 return False 

424 

425 

426def translate_client_error(err: ClientError, uri: ResourcePath) -> None: 

427 """Translate a ClientError into a specialist error if relevant. 

428 

429 Parameters 

430 ---------- 

431 err : `ClientError` 

432 Exception to translate. 

433 uri : `ResourcePath` 

434 The URI of the resource that is resulting in the error. 

435 

436 Raises 

437 ------ 

438 _TooManyRequestsError 

439 Raised if the `ClientError` looks like a 429 retry request. 

440 """ 

441 if "(429)" in str(err): 441 ↛ 445line 441 didn't jump to line 445 because the condition on line 441 was never true

442 # ClientError includes the error code in the message 

443 # but no direct way to access it without looking inside the 

444 # response. 

445 raise _TooManyRequestsError(f"{err} when accessing {uri}") from err 

446 elif "(404)" in str(err): 446 ↛ exitline 446 didn't return from function 'translate_client_error' because the condition on line 446 was always true

447 # Some systems can generate this rather than NoSuchKey. 

448 raise FileNotFoundError(f"Resource not found (permission denied): {uri}")