Coverage for python/lsst/resources/http.py: 54%

809 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__ = ("HttpResourcePath",) 

15 

16import contextlib 

17import datetime 

18import enum 

19import functools 

20import io 

21import json 

22import logging 

23import math 

24import os 

25import os.path 

26import random 

27import re 

28import ssl 

29import stat 

30from collections.abc import Generator, Iterator 

31from email.utils import parsedate_to_datetime 

32from typing import TYPE_CHECKING, Any, BinaryIO, cast 

33 

34if TYPE_CHECKING: 

35 # defusedxml ships no type information, so let type checkers see the 

36 # standard library module that it hardens and mirrors. 

37 import xml.etree.ElementTree as eTree 

38else: 

39 try: 

40 # Prefer 'defusedxml' (not part of standard library) if available, 

41 # since 'xml' is vulnerable to XML bombs. 

42 import defusedxml.ElementTree as eTree 

43 except ImportError: 

44 import xml.etree.ElementTree as eTree 

45 

46# defusedxml hardens the parser but still builds trees out of the standard 

47# library element type, which is the only one it re-exports. 

48from xml.etree.ElementTree import Element 

49 

50try: 

51 import fsspec 

52 from aiohttp import ClientSession, ClientTimeout, TCPConnector 

53 from fsspec.implementations.http import HTTPFileSystem 

54 from fsspec.spec import AbstractFileSystem 

55except ImportError: 

56 # Hidden from type checkers so that the names above keep the types they 

57 # have when fsspec is installed. 

58 if not TYPE_CHECKING: 

59 fsspec = None 

60 AbstractFileSystem = type 

61 HTTPFileSystem = type 

62 

63from urllib.parse import parse_qs 

64 

65import requests 

66from astropy import units as u 

67from requests.adapters import HTTPAdapter 

68from requests.auth import AuthBase 

69from urllib3.util.retry import Retry 

70 

71from lsst.utils.timer import time_this 

72 

73from ._resourceHandles import ResourceHandleProtocol 

74from ._resourceHandles._httpResourceHandle import HttpReadResourceHandle, parse_content_range_header 

75from ._resourcePath import ResourceInfo, ResourcePath 

76from .utils import _get_num_workers, get_tempdir 

77 

78if TYPE_CHECKING: 

79 from .utils import TransactionProtocol 

80 

81log = logging.getLogger(__name__) 

82 

83 

84def _timeout_from_environment(env_var: str, default_value: float) -> float: 

85 """Convert and return a timeout from the value of an environment variable 

86 or a default value if the environment variable is not initialized. The 

87 value of `env_var` must be a valid `float` otherwise this function raises. 

88 

89 Parameters 

90 ---------- 

91 env_var : `str` 

92 Environment variable to look for. 

93 default_value : `float`` 

94 Value to return if `env_var` is not defined in the environment. 

95 

96 Returns 

97 ------- 

98 _timeout_from_environment : `float` 

99 Converted value. 

100 """ 

101 try: 

102 timeout = float(os.environ.get(env_var, default_value)) 

103 except ValueError: 

104 raise ValueError( 

105 f"Expecting valid timeout value in environment variable {env_var} but found " 

106 f"{os.environ.get(env_var)}" 

107 ) from None 

108 

109 if math.isnan(timeout): 

110 raise ValueError(f"Unexpected timeout value NaN found in environment variable {env_var}") 

111 

112 return timeout 

113 

114 

115@functools.lru_cache 

116def _calc_tmpdir_buffer_size(tmpdir: str) -> int: 

117 """Compute the block size as 256 blocks of typical size 

118 (i.e. 4096 bytes) or 10 times the file system block size, 

119 whichever is higher. 

120 

121 This is a reasonable compromise between 

122 using memory for buffering and the number of system calls 

123 issued to read from or write to temporary files. 

124 """ 

125 fsstats = os.statvfs(tmpdir) 

126 return max(10 * fsstats.f_bsize, 256 * 4096) 

127 

128 

129class HttpResourcePathConfig: 

130 """Configuration class to encapsulate the configurable items used by class 

131 HttpResourcePath. 

132 """ 

133 

134 # Default timeouts for all HTTP requests (seconds). 

135 DEFAULT_TIMEOUT_CONNECT: float = 60.0 

136 DEFAULT_TIMEOUT_READ: float = 1_500.0 

137 

138 # Default lower and upper bounds for the backoff interval (seconds). 

139 # A value in this interval is randomly selected as the backoff factor when 

140 # requests need to be retried. 

141 DEFAULT_BACKOFF_MIN: float = 1.0 

142 DEFAULT_BACKOFF_MAX: float = 3.0 

143 

144 # Default number of connections to persist with both the front end and 

145 # back end servers. 

146 DEFAULT_FRONTEND_PERSISTENT_CONNECTIONS: int = 2 

147 DEFAULT_BACKEND_PERSISTENT_CONNECTIONS: int = 1 

148 

149 # Accepted digest algorithms 

150 ACCEPTED_DIGESTS: list[str] = ["adler32", "md5", "sha-256", "sha-512"] 

151 

152 def __init__(self) -> None: 

153 self._front_end_connections: int | None = None 

154 self._back_end_connections: int | None = None 

155 self._digest_algorithm: str | None = None 

156 self._send_expect_on_put: bool | None = None 

157 self._fsspec_is_enabled: bool | None = None 

158 self._timeout: tuple[float, float] | None = None 

159 self._collect_memory_usage: bool | None = None 

160 self._backoff_min: float | None = None 

161 self._backoff_max: float | None = None 

162 self._ca_bundle: str | None = "" 

163 self._client_token: str | None = "" 

164 self._client_cert: str | None = "" 

165 self._client_key: str | None = "" 

166 self._tmpdir_buffersize: tuple[str, int] | None = None 

167 self._ssl_context: ssl.SSLContext | None = None 

168 

169 @property 

170 def front_end_connections(self) -> int: 

171 """Number of persistent connections to the front end server.""" 

172 if self._front_end_connections is not None: 172 ↛ 173line 172 didn't jump to line 173 because the condition on line 172 was never true

173 return self._front_end_connections 

174 

175 default_pool_size = max(_get_num_workers(), self.DEFAULT_FRONTEND_PERSISTENT_CONNECTIONS) 

176 

177 try: 

178 self._front_end_connections = int( 

179 os.environ.get("LSST_HTTP_FRONTEND_PERSISTENT_CONNECTIONS", default_pool_size) 

180 ) 

181 except ValueError: 

182 self._front_end_connections = default_pool_size 

183 

184 return self._front_end_connections 

185 

186 @property 

187 def back_end_connections(self) -> int: 

188 """Number of persistent connections to the back end servers.""" 

189 if self._back_end_connections is not None: 189 ↛ 190line 189 didn't jump to line 190 because the condition on line 189 was never true

190 return self._back_end_connections 

191 

192 default_pool_size = max(_get_num_workers(), self.DEFAULT_FRONTEND_PERSISTENT_CONNECTIONS) 

193 

194 try: 

195 self._back_end_connections = int( 

196 os.environ.get("LSST_HTTP_BACKEND_PERSISTENT_CONNECTIONS", default_pool_size) 

197 ) 

198 except ValueError: 

199 self._back_end_connections = default_pool_size 

200 

201 return self._back_end_connections 

202 

203 @property 

204 def digest_algorithm(self) -> str: 

205 """Algorithm to ask the server to use for computing and recording 

206 digests of each file contents in PUT requests. 

207 

208 Returns 

209 ------- 

210 digest_algorithm: `str` 

211 The name of a digest algorithm or the empty string if no algotihm 

212 is configured. 

213 """ 

214 if self._digest_algorithm is not None: 

215 return self._digest_algorithm 

216 

217 digest = os.environ.get("LSST_HTTP_DIGEST", "").lower() 

218 if digest not in self.ACCEPTED_DIGESTS: 

219 digest = "" 

220 

221 self._digest_algorithm = digest 

222 return self._digest_algorithm 

223 

224 @property 

225 def send_expect_on_put(self) -> bool: 

226 """Return True if a "Expect: 100-continue" header is to be sent to 

227 the server on each PUT request. 

228 

229 Some servers (e.g. dCache) uses this information as an indication that 

230 the client knows how to handle redirects to the specific server that 

231 will actually receive the data for PUT requests. 

232 """ 

233 if self._send_expect_on_put is not None: 

234 return self._send_expect_on_put 

235 

236 self._send_expect_on_put = "LSST_HTTP_PUT_SEND_EXPECT_HEADER" in os.environ 

237 return self._send_expect_on_put 

238 

239 @property 

240 def fsspec_is_enabled(self) -> bool: 

241 """Return True if `fsspec` is enabled for objects of class 

242 HttpResourcePath. 

243 

244 To determine if `fsspec` is enabled, this method inspects the presence 

245 of the environment variable `LSST_HTTP_ENABLE_FSSPEC` (with any value). 

246 """ 

247 if self._fsspec_is_enabled is not None: 247 ↛ 248line 247 didn't jump to line 248 because the condition on line 247 was never true

248 return self._fsspec_is_enabled 

249 

250 self._fsspec_is_enabled = "LSST_HTTP_ENABLE_FSSPEC" in os.environ 

251 return self._fsspec_is_enabled 

252 

253 @property 

254 def timeout(self) -> tuple[float, float]: 

255 """Return a tuple with the values of timeouts for connecting to the 

256 server and reading its response, respectively. Both values are in 

257 seconds. 

258 """ 

259 if self._timeout is not None: 

260 return self._timeout 

261 

262 self._timeout = ( 

263 _timeout_from_environment("LSST_HTTP_TIMEOUT_CONNECT", self.DEFAULT_TIMEOUT_CONNECT), 

264 _timeout_from_environment("LSST_HTTP_TIMEOUT_READ", self.DEFAULT_TIMEOUT_READ), 

265 ) 

266 return self._timeout 

267 

268 @property 

269 def collect_memory_usage(self) -> bool: 

270 """Return true if we want to collect memory usage when timing 

271 operations against the remote server via the `lsst.utils.time_this` 

272 context manager. 

273 """ 

274 if self._collect_memory_usage is not None: 

275 return self._collect_memory_usage 

276 

277 self._collect_memory_usage = "LSST_HTTP_COLLECT_MEMORY_USAGE" in os.environ 

278 return self._collect_memory_usage 

279 

280 @property 

281 def backoff_min(self) -> float: 

282 """Lower bound of the interval from which a backoff factor is randomly 

283 selected when retrying requests (seconds). 

284 """ 

285 if self._backoff_min is not None: 

286 return self._backoff_min 

287 

288 self._backoff_min = self.DEFAULT_BACKOFF_MIN 

289 try: 

290 backoff_min = float(os.environ.get("LSST_HTTP_BACKOFF_MIN", self.DEFAULT_BACKOFF_MIN)) 

291 if not math.isnan(backoff_min): 

292 self._backoff_min = backoff_min 

293 except ValueError: 

294 pass 

295 

296 return self._backoff_min 

297 

298 @property 

299 def backoff_max(self) -> float: 

300 """Upper bound of the interval from which a backoff factor is randomly 

301 selected when retrying requests (seconds). 

302 """ 

303 if self._backoff_max is not None: 

304 return self._backoff_max 

305 

306 self._backoff_max = self.DEFAULT_BACKOFF_MAX 

307 try: 

308 backoff_max = float(os.environ.get("LSST_HTTP_BACKOFF_MAX", self.DEFAULT_BACKOFF_MAX)) 

309 if not math.isnan(backoff_max): 

310 self._backoff_max = backoff_max 

311 except ValueError: 

312 pass 

313 

314 return self._backoff_max 

315 

316 @property 

317 def ca_bundle(self) -> str | None: 

318 """Local path to the certificate bundle file or directory where the 

319 certifcates of the trusted authorities are located. 

320 

321 Return None if this host's system certificate bundle should be 

322 used for authenticating remote servers' certificates. 

323 """ 

324 if self._ca_bundle != "": 

325 return self._ca_bundle 

326 

327 # If a bundle was specified via the environment variable 

328 # 'LSST_HTTP_CACERT_BUNDLE' use it. 

329 self._ca_bundle = os.getenv("LSST_HTTP_CACERT_BUNDLE") 

330 return self._ca_bundle 

331 

332 @property 

333 def client_token(self) -> str | None: 

334 """Value of a bearer token or path to a local file which contains 

335 the bearer token to use for authenticating the client when sending 

336 requests to the webDAV or HTTP server. 

337 

338 Return None if no bearer token is configured in the environment. 

339 """ 

340 if self._client_token != "": 

341 return self._client_token 

342 

343 # If environment variable LSST_HTTP_AUTH_BEARER_TOKEN is 

344 # initialized use its value as the bearer token. 

345 self._client_token = os.getenv("LSST_HTTP_AUTH_BEARER_TOKEN") 

346 return self._client_token 

347 

348 @property 

349 def client_cert_key(self) -> tuple[str | None, str | None]: 

350 """Paths to a local file where the client certificate and associated 

351 private key are located. 

352 

353 Return a tuple (client certificate, private key) or (None, None) if no 

354 client certificate is configured via environment variables. 

355 """ 

356 if self._client_cert != "" and self._client_key != "": 

357 return (self._client_cert, self._client_key) 

358 

359 # If the environment variables LSST_HTTP_AUTH_CLIENT_CERT 

360 # and LSST_HTTP_AUTH_CLIENT_KEY are initialized use their values. 

361 self._client_cert = os.getenv("LSST_HTTP_AUTH_CLIENT_CERT") 

362 self._client_key = os.getenv("LSST_HTTP_AUTH_CLIENT_KEY") 

363 if self._client_cert and self._client_key: 

364 if not _is_protected(self._client_key): 

365 raise PermissionError( 

366 f"Private key file at {self._client_key} must be protected for access only by its owner" 

367 ) 

368 return (self._client_cert, self._client_key) 

369 

370 # If only the certificate was provided raise. 

371 if self._client_cert: 

372 raise ValueError( 

373 "Environment variable LSST_HTTP_AUTH_CLIENT_KEY must be set to client private key file path" 

374 ) 

375 

376 # If only the private key was provided raise. 

377 if self._client_key: 

378 raise ValueError( 

379 "Environment variable LSST_HTTP_AUTH_CLIENT_CERT must be set to client certificate file path" 

380 ) 

381 

382 # If a X.509 user proxy is available, use it as client credentials. 

383 self._client_cert = self._client_key = os.getenv("X509_USER_PROXY") 

384 return (self._client_cert, self._client_key) 

385 

386 @property 

387 def tmpdir_buffersize(self) -> tuple[str, int]: 

388 """Return the path to a temporary directory and the preferred buffer 

389 size to use when reading or writing files in that directory. 

390 """ 

391 if self._tmpdir_buffersize is not None: 

392 return self._tmpdir_buffersize 

393 

394 tmpdir = get_tempdir() 

395 

396 # Compute the block size as 256 blocks of typical size 

397 # (i.e. 4096 bytes) or 10 times the file system block size, 

398 # whichever is higher. This is a reasonable compromise between 

399 # using memory for buffering and the number of system calls 

400 # issued to read from or write to temporary files. 

401 bufsize = _calc_tmpdir_buffer_size(tmpdir) 

402 self._tmpdir_buffersize = (tmpdir, bufsize) 

403 

404 return self._tmpdir_buffersize 

405 

406 @property 

407 def ssl_context(self) -> ssl.SSLContext: 

408 """Return an SSL context equiped with the certificates of the trusted 

409 authorities. 

410 """ 

411 if self._ssl_context is None: 

412 self._ssl_context = ssl.create_default_context() 

413 if self.ca_bundle is not None: 

414 if os.path.isdir(self.ca_bundle): 

415 self._ssl_context.load_verify_locations(capath=self.ca_bundle) 

416 elif os.path.isfile(self.ca_bundle): 

417 self._ssl_context.load_verify_locations(cafile=self.ca_bundle) 

418 

419 return self._ssl_context 

420 

421 

422@functools.lru_cache 

423def _get_dav_and_server_headers(path: ResourcePath | str) -> tuple[str | None, str | None]: 

424 """Retrieve the "DAV" and "Server" headers sent by the remote server as 

425 part of the response to a single "OPTIONS" HTTP request. 

426 

427 Parameters 

428 ---------- 

429 path : `ResourcePath` or `str` 

430 URL to the resource to be checked. 

431 Should preferably refer to the root since the status is shared 

432 by all paths in that server. 

433 

434 Returns 

435 ------- 

436 _get_dav_and_server_headers : `tuple[str|None, str|None]` 

437 Values of the "DAV" and "Server" headers found in the response or 

438 None if any of those headers was not part of the response. 

439 """ 

440 try: 

441 if not isinstance(path, HttpResourcePath): 441 ↛ 442line 441 didn't jump to line 442 because the condition on line 441 was never true

442 path = HttpResourcePath(path) 

443 

444 config = HttpResourcePathConfig() 

445 with SessionStore(config=config).get(path) as session: 

446 # ResourcePath.__new__ is a scheme-dispatching factory declared as 

447 # returning the base class, which ty honors and mypy does not. 

448 headers = path._extra_headers # ty: ignore[unresolved-attribute] 

449 resp = session.options(str(path), stream=False, timeout=config.timeout, headers=headers) 

450 

451 dav_header = server_header = None 

452 if resp.status_code == requests.codes.ok: 452 ↛ 456line 452 didn't jump to line 456 because the condition on line 452 was always true

453 dav_header = resp.headers.get("DAV") if "DAV" in resp.headers else None 

454 server_header = resp.headers.get("Server") if "Server" in resp.headers else None 

455 

456 return (dav_header, server_header) 

457 

458 except requests.exceptions.SSLError as e: 

459 log.warning( 

460 "Environment variable LSST_HTTP_CACERT_BUNDLE can be used to " 

461 "specify the path to a bundle of certificate authorities you trust " 

462 "which are not included in the default set of trusted authorities " 

463 "of this system." 

464 ) 

465 raise e 

466 

467 

468class BearerTokenAuth(AuthBase): 

469 """Attach a bearer token 'Authorization' header to each request. 

470 

471 Parameters 

472 ---------- 

473 token : `str` 

474 Can be either the path to a local protected file which contains the 

475 value of the token or the token itself. 

476 """ 

477 

478 def __init__(self, token: str): 

479 self._token = self._path = None 

480 self._mtime: float = -1.0 

481 if not token: 

482 return 

483 

484 self._token = token 

485 if os.path.isfile(token): 

486 self._path = os.path.abspath(token) 

487 if not _is_protected(self._path): 

488 raise PermissionError( 

489 f"Bearer token file at {self._path} must be protected for access only by its owner" 

490 ) 

491 self._refresh() 

492 

493 def _refresh(self) -> None: 

494 """Read the token file (if any) if its modification time is more recent 

495 than the last time we read it. 

496 """ 

497 if not self._path: 

498 return 

499 

500 if (mtime := os.stat(self._path).st_mtime) > self._mtime: 

501 log.debug("Reading bearer token file at %s", self._path) 

502 self._mtime = mtime 

503 with open(self._path) as f: 

504 self._token = f.read().rstrip("\n") 

505 

506 def __call__(self, r: requests.PreparedRequest) -> requests.PreparedRequest: 

507 # Parameter is named to match requests.auth.AuthBase.__call__, which 

508 # callers may invoke by keyword. 

509 # Only add a bearer token to a request when using secure HTTP. 

510 if r.url and r.url.lower().startswith("https://") and self._token: 

511 self._refresh() 

512 r.headers["Authorization"] = f"Bearer {self._token}" 

513 return r 

514 

515 

516class SessionStore: 

517 """Cache a reusable HTTP client session per endpoint. 

518 

519 Parameters 

520 ---------- 

521 config : `HttpResourcePathConfig` 

522 Configuration items shared by all instances of HttpResourcePath. 

523 num_pools : `int`, optional 

524 Number of connection pools to keep: there is one pool per remote 

525 host. 

526 max_persistent_connections : `int`, optional 

527 Maximum number of connections per remote host to persist in each 

528 connection pool. 

529 backoff_min : `float`, optional 

530 Minimum value of the interval to compute the exponential 

531 backoff factor when retrying requests (seconds). 

532 backoff_max : `float`, optional 

533 Maximum value of the interval to compute the exponential 

534 backoff factor when retrying requests (seconds). 

535 """ 

536 

537 def __init__( 

538 self, 

539 config: HttpResourcePathConfig, 

540 num_pools: int = 10, 

541 max_persistent_connections: int = 1, 

542 backoff_min: float = 1.0, 

543 backoff_max: float = 3.0, 

544 ) -> None: 

545 # Dictionary to store the session associated to a given URI. The key 

546 # of the dictionary is a root URI and the value is the session. 

547 self._sessions: dict[str, requests.Session] = {} 

548 

549 # Configuration for all instances of HttpResourcePath objects. 

550 self._config = config 

551 

552 # See documentation of urllib3 PoolManager class: 

553 # https://urllib3.readthedocs.io 

554 self._num_pools: int = num_pools 

555 

556 # See urllib3 Advanced Usage documentation: 

557 # https://urllib3.readthedocs.io/en/stable/advanced-usage.html 

558 self._max_persistent_connections: int = max_persistent_connections 

559 

560 # Minimum and maximum values of the interval to compute the exponential 

561 # backoff factor when retrying requests (seconds). 

562 self._backoff_min: float = backoff_min 

563 self._backoff_max: float = backoff_max if backoff_max > backoff_min else backoff_min + 1.0 

564 

565 def clear(self) -> None: 

566 """Destroy all previously created sessions and attempt to close 

567 underlying idle network connections. 

568 """ 

569 # Close all sessions and empty the store. Idle network connections 

570 # should be closed as a consequence. We don't have means through 

571 # the API exposed by Requests to actually force closing the 

572 # underlying open sockets. 

573 for session in self._sessions.values(): 

574 session.close() 

575 

576 self._sessions.clear() 

577 

578 def get(self, rpath: ResourcePath) -> requests.Session: 

579 """Retrieve a session for accessing the remote resource at rpath. 

580 

581 Parameters 

582 ---------- 

583 rpath : `ResourcePath` 

584 URL to a resource at the remote server for which a session is to 

585 be retrieved. 

586 

587 Notes 

588 ----- 

589 Once a session is created for a given endpoint it is cached and 

590 returned every time a session is requested for any path under that same 

591 endpoint. For instance, a single session will be cached and shared 

592 for paths "https://www.example.org/path/to/file" and 

593 "https://www.example.org/any/other/path". 

594 

595 Note that "https://www.example.org" and "https://www.example.org:12345" 

596 will have different sessions since the port number is not identical. 

597 """ 

598 root_uri = str(rpath.root_uri()) 

599 if root_uri not in self._sessions: 

600 # We don't have yet a session for this endpoint: create a new one. 

601 self._sessions[root_uri] = self._make_session(rpath) 

602 

603 return self._sessions[root_uri] 

604 

605 def _make_session(self, rpath: ResourcePath) -> requests.Session: 

606 """Make a new session configured from values from the environment.""" 

607 session = requests.Session() 

608 root_uri = str(rpath.root_uri()) 

609 log.debug("Creating new HTTP session for endpoint %s ...", root_uri) 

610 retries = Retry( 

611 # Total number of retries to allow. Takes precedence over other 

612 # counts. 

613 total=6, 

614 # How many connection-related errors to retry on. 

615 connect=3, 

616 # How many times to retry on read errors. 

617 read=3, 

618 # Backoff factor to apply between attempts after the second try 

619 # (seconds). Compute a random jitter to prevent all the clients 

620 # to overwhelm the server by sending requests at the same time. 

621 backoff_factor=self._backoff_min + (self._backoff_max - self._backoff_min) * random.random(), 

622 # How many times to retry on bad status codes. 

623 status=5, 

624 # Set of uppercased HTTP method verbs that we should retry on. 

625 # We only automatically retry idempotent requests. 

626 allowed_methods=frozenset( 

627 [ 

628 "COPY", 

629 "DELETE", 

630 "GET", 

631 "HEAD", 

632 "MKCOL", 

633 "OPTIONS", 

634 "PROPFIND", 

635 "PUT", 

636 ] 

637 ), 

638 # HTTP status codes that we should force a retry on. 

639 status_forcelist=frozenset( 

640 [ 

641 requests.codes.too_many_requests, # 429 

642 requests.codes.internal_server_error, # 500 

643 requests.codes.bad_gateway, # 502 

644 requests.codes.service_unavailable, # 503 

645 requests.codes.gateway_timeout, # 504 

646 ] 

647 ), 

648 # Whether to respect Retry-After header on status codes defined 

649 # above. 

650 respect_retry_after_header=True, 

651 ) 

652 

653 # Persist the specified number of connections to the front end server. 

654 session.mount( 

655 root_uri, 

656 HTTPAdapter( 

657 pool_connections=self._num_pools, 

658 pool_maxsize=self._max_persistent_connections, 

659 pool_block=False, 

660 max_retries=retries, 

661 ), 

662 ) 

663 

664 # Do not persist the connections to back end servers which may vary 

665 # from request to request. Systematically persisting connections to 

666 # those servers may exhaust their capabilities when there are thousands 

667 # of simultaneous clients. 

668 session.mount( 

669 f"{rpath.scheme}://", 

670 HTTPAdapter( 

671 pool_connections=self._num_pools, 

672 pool_maxsize=0, 

673 pool_block=False, 

674 max_retries=retries, 

675 ), 

676 ) 

677 

678 # If the remote endpoint doesn't use secure HTTP we don't include 

679 # bearer tokens in the requests nor need to authenticate the remote 

680 # server. 

681 if rpath.scheme != "https": 

682 return session 

683 

684 # Set the trusted CA certificates bundle for authenticating remote 

685 # servers. 

686 session.verify = True if self._config.ca_bundle is None else self._config.ca_bundle 

687 

688 # Should we use a bearer token for client authentication? 

689 if (token := self._config.client_token) is not None: 

690 log.debug("... using bearer token authentication") 

691 session.auth = BearerTokenAuth(token) 

692 return session 

693 

694 # Should we instead use client certificate and private key? 

695 client_cert, client_key = self._config.client_cert_key 

696 if client_cert and client_key: 

697 log.debug("... using client certificate authentication.") 

698 session.cert = (client_cert, client_key) 

699 return session 

700 

701 log.debug( 

702 "Neither LSST_HTTP_AUTH_BEARER_TOKEN nor (LSST_HTTP_AUTH_CLIENT_CERT and " 

703 "LSST_HTTP_AUTH_CLIENT_KEY) are initialized. Client authentication is disabled." 

704 ) 

705 return session 

706 

707 

708class ActivityCaveat(enum.Enum): 

709 """Helper class for enumerating accepted activity caveats for requesting 

710 macaroons. 

711 """ 

712 

713 DOWNLOAD = 1 

714 UPLOAD = 2 

715 

716 

717class HttpResourcePath(ResourcePath): 

718 """General HTTP(S) resource. 

719 

720 Notes 

721 ----- 

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

723 environment variables below are inspected: 

724 

725 - LSST_HTTP_CACERT_BUNDLE: path to a .pem file or to a directory which 

726 contains the .pem files of the trusted certificate authorities's 

727 certificates. If the remote server presents a server certificate 

728 issued by one of those trusted authorities, we trust it. 

729 If this environment variable is not initialized, the default 

730 authorities of the the execution host are trusted. 

731 

732 - LSST_HTTP_AUTH_BEARER_TOKEN: value of a bearer token or path to a 

733 local file containing a bearer token to be used as the client 

734 authentication mechanism with all requests. 

735 The permissions of the token file must be set so that only its 

736 owner can access it. 

737 If initialized, takes precedence over LSST_HTTP_AUTH_CLIENT_CERT 

738 and LSST_HTTP_AUTH_CLIENT_KEY. 

739 

740 - LSST_HTTP_AUTH_CLIENT_CERT: path to a .pem file which contains the 

741 client certificate for authenticating to the server. 

742 If initialized, the variable LSST_HTTP_AUTH_CLIENT_KEY must also be 

743 initialized with the path of the client private key file. 

744 The permissions of the client private key must be set so that only 

745 its owner can access it, at least for reading. 

746 

747 - LSST_HTTP_PUT_SEND_EXPECT_HEADER: if set (with any value), a 

748 "Expect: 100-Continue" header will be added to all HTTP PUT requests. 

749 This header is required by some servers to detect if the client 

750 knows how to handle redirections. In case of redirection, the body 

751 of the PUT request is sent to the redirected location and not to 

752 the front end server. 

753 

754 - LSST_HTTP_TIMEOUT_CONNECT and LSST_HTTP_TIMEOUT_READ: if set to a 

755 numeric value, they are interpreted as the number of seconds to wait 

756 for establishing a connection with the server and for reading its 

757 response, respectively. 

758 

759 - LSST_HTTP_FRONTEND_PERSISTENT_CONNECTIONS and 

760 LSST_HTTP_BACKEND_PERSISTENT_CONNECTIONS: contain the maximum number 

761 of connections to attempt to persist with both the front end servers 

762 and the back end servers. 

763 Default values: DEFAULT_FRONTEND_PERSISTENT_CONNECTIONS and 

764 DEFAULT_BACKEND_PERSISTENT_CONNECTIONS. 

765 

766 - LSST_HTTP_DIGEST: case-insensitive name of the digest algorithm to 

767 ask the server to compute for every file's content sent to the server 

768 via a PUT request. No digest is requested if this variable is not set 

769 or is set to an invalid value. 

770 Valid values are those in ACCEPTED_DIGESTS. 

771 

772 - LSST_HTTP_ENABLE_FSSPEC: the presence of this environment variable 

773 activates the usage of `fsspec` compatible file system to read 

774 a HTTP URL. The value of the variable is not inspected. 

775 """ 

776 

777 @staticmethod 

778 def create_http_resource_path( 

779 path: str, *, extra_headers: dict[str, str] | None = None 

780 ) -> HttpResourcePath: 

781 """Create an instance of `HttpResourcePath` with additional 

782 HTTP-specific configuration. 

783 

784 Parameters 

785 ---------- 

786 path : `str` 

787 HTTP URL to be wrapped in a `ResourcePath` instance. 

788 extra_headers : `dict` [ `str`, `str` ], optional 

789 Additional headers that will be sent with every HTTP request made 

790 by this `ResourcePath`. These override any headers that may be 

791 generated internally by `HttpResourcePath` (e.g. authentication 

792 headers). 

793 

794 Returns 

795 ------- 

796 instance : `ResourcePath` 

797 Newly-created `HttpResourcePath` instance. 

798 

799 Notes 

800 ----- 

801 Most users should use the `ResourcePath` constructor, instead. 

802 """ 

803 # Make sure we instantiate ResourcePath using a string to guarantee we 

804 # get a new ResourcePath. If we accidentally provided a ResourcePath 

805 # instance instead, the ResourcePath constructor sometimes returns 

806 # the original object and we would be modifying an object that is 

807 # supposed to be immutable. 

808 instance = ResourcePath(str(path)) 

809 assert isinstance(instance, HttpResourcePath) 

810 instance._extra_headers = extra_headers 

811 return instance 

812 

813 # WebDAV servers known to be able to sign URLs. The values are lowercased 

814 # server identifiers retrieved from the 'Server' header included in 

815 # the response to a HTTP OPTIONS request. 

816 SUPPORTED_URL_SIGNERS = ("dcache", "xrootd") 

817 

818 # Configuration items for this class instances. 

819 _config: HttpResourcePathConfig = HttpResourcePathConfig() 

820 

821 # The session for metadata requests is used for interacting with 

822 # the front end servers for requests such as PROPFIND, HEAD, etc. Those 

823 # interactions are typically served by the front end servers. We want to 

824 # keep the connection to the front end servers open, to reduce the cost 

825 # associated to TCP and TLS handshaking for each new request. 

826 _metadata_session_store = SessionStore( 

827 config=_config, 

828 num_pools=5, 

829 max_persistent_connections=_config.front_end_connections, 

830 backoff_min=_config.backoff_min, 

831 backoff_max=_config.backoff_max, 

832 ) 

833 

834 # The data session is used for interaction with the front end servers which 

835 # typically redirect to the back end servers for serving our PUT and GET 

836 # requests. We attempt to keep a single connection open with the front end 

837 # server, if possible. This depends on how the server behaves and the 

838 # kind of request. Some servers close the connection when redirecting 

839 # the client to a back end server, for instance when serving a PUT 

840 # request. 

841 _data_session_store = SessionStore( 

842 config=_config, 

843 num_pools=25, 

844 max_persistent_connections=_config.back_end_connections, 

845 backoff_min=_config.backoff_min, 

846 backoff_max=_config.backoff_max, 

847 ) 

848 

849 # Process ID which created the session stores above. We need to store this 

850 # to replace sessions created by a parent process and inherited by a 

851 # child process after a fork, to avoid confusing the SSL layer. 

852 _pid: int = -1 

853 

854 # Connector used by a session pool to establish network connections to 

855 # remote servers. This connector is exclusively used by fsspec file system 

856 # and is shared by all instances of this class. 

857 _tcp_connector: TCPConnector | None = None 

858 

859 # Additional headers added to every request. 

860 _extra_headers: dict[str, str] | None = None 

861 

862 @property 

863 def metadata_session(self) -> _SessionWrapper: 

864 """Client session to send requests which do not require upload or 

865 download of data, i.e. mostly metadata requests. 

866 """ 

867 session = None 

868 if hasattr(self, "_metadata_session"): 

869 if HttpResourcePath._pid == os.getpid(): 869 ↛ 874line 869 didn't jump to line 874 because the condition on line 869 was always true

870 session = self._metadata_session 

871 else: 

872 # The metadata session we have in cache was likely created by 

873 # a parent process. Discard all the sessions in that store. 

874 self._metadata_session_store.clear() 

875 

876 # Retrieve a new metadata session. 

877 if session is None: 

878 HttpResourcePath._pid = os.getpid() 

879 session = self._metadata_session_store.get(self) 

880 self._metadata_session: requests.Session = session 

881 return _SessionWrapper(session, extra_headers=self._extra_headers) 

882 

883 @property 

884 def data_session(self) -> _SessionWrapper: 

885 """Client session for uploading and downloading data.""" 

886 session = None 

887 if hasattr(self, "_data_session"): 

888 if HttpResourcePath._pid == os.getpid(): 888 ↛ 893line 888 didn't jump to line 893 because the condition on line 888 was always true

889 session = self._data_session 

890 else: 

891 # The data session we have in cache was likely created by 

892 # a parent process. Discard all the sessions in that store. 

893 self._data_session_store.clear() 

894 

895 # Retrieve a new data session. 

896 if session is None: 

897 HttpResourcePath._pid = os.getpid() 

898 session = self._data_session_store.get(self) 

899 self._data_session: requests.Session = session 

900 return _SessionWrapper(session, extra_headers=self._extra_headers) 

901 

902 def _clear_sessions(self) -> None: 

903 """Close the socket connections that are still open. 

904 

905 Used only in test suites to avoid warnings. 

906 """ 

907 self._metadata_session_store.clear() 

908 self._data_session_store.clear() 

909 

910 if hasattr(self, "_metadata_session"): 

911 delattr(self, "_metadata_session") 

912 

913 if hasattr(self, "_data_session"): 

914 delattr(self, "_data_session") 

915 

916 def _init_server_properties(self) -> None: 

917 """Initialize instance variables '_is_webdav' and '_server' by 

918 sending a single OPTIONS request to the remote server and 

919 saving the results. 

920 """ 

921 # Retrieve the "DAV" and the "Server" headers for the root URL of this 

922 # path 

923 dav_header, server_header = _get_dav_and_server_headers(self.root_uri()) 

924 

925 # Check that "1" is part of the value of the "DAV" header. We don't 

926 # use locks, so a server complying to class 1 is enough for our 

927 # purposes. All webDAV servers must advertise at least compliance 

928 # class "1". 

929 # 

930 # Compliance classes are documented in 

931 # http://www.webdav.org/specs/rfc4918.html#dav.compliance.classes 

932 # 

933 # Examples of values for header DAV are: 

934 # DAV: 1, 2 

935 # DAV: 1, <http://apache.org/dav/propset/fs/1> 

936 self._is_webdav: bool = False 

937 if dav_header is not None: 937 ↛ 938line 937 didn't jump to line 938 because the condition on line 937 was never true

938 self._is_webdav = "1" in dav_header.replace(" ", "").split(",") 

939 

940 self._server: str | None = None 

941 if server_header is not None: 

942 # Server header is expected to be of the form 'dCache/9.2.4' 

943 # or 'XrootD/v5.7.1'. Strip version and put in lowercase. 

944 self._server = server_header.split("/")[0].lower() 

945 

946 @property 

947 def is_webdav_endpoint(self) -> bool: 

948 """Check if the current endpoint implements WebDAV features. 

949 

950 This is stored per URI but cached by root so there is only one check 

951 per hostname. 

952 """ 

953 if hasattr(self, "_is_webdav"): 

954 return self._is_webdav 

955 

956 self._init_server_properties() 

957 return self._is_webdav 

958 

959 @property 

960 def server(self) -> str | None: 

961 """Return the lowercased identifier of the remote server, retrieved 

962 from the response header 'Server' from an 'OPTIONS' HTTP request. 

963 

964 If the remote server does not include that header in its response 

965 to an 'OPTIONS' request, server() returns None. 

966 

967 Examples of return values are "dcache", "xrootd". 

968 """ 

969 if hasattr(self, "_server"): 969 ↛ 972line 969 didn't jump to line 972 because the condition on line 969 was always true

970 return self._server 

971 

972 self._init_server_properties() 

973 return self._server 

974 

975 @property 

976 def server_signs_urls(self) -> bool: 

977 """Return true if the remote server support signing or URLs for 

978 download and upload. 

979 """ 

980 return self.server in HttpResourcePath.SUPPORTED_URL_SIGNERS 

981 

982 @classmethod 

983 def _reload_config(cls) -> None: 

984 """Reload the configuration for all instances of this class. That 

985 configuration is instantiated from the environment. 

986 

987 This is an internal method mainly intended for tests. 

988 """ 

989 HttpResourcePath._config = HttpResourcePathConfig() 

990 

991 def exists(self) -> bool: 

992 """Check that a remote HTTP resource exists.""" 

993 log.debug("Checking if resource exists: %s", self.geturl()) 

994 if not self.is_webdav_endpoint: 994 ↛ 1004line 994 didn't jump to line 1004 because the condition on line 994 was always true

995 # The remote is a plain HTTP server. Let's attempt a HEAD 

996 # request, even if the behavior for such a request against a 

997 # directory is not specified, so it depends on the server 

998 # implementation. 

999 resp = self._head_non_webdav_url() 

1000 return self._is_successful_non_webdav_head_request(resp) 

1001 

1002 # The remote endpoint is a webDAV server: send a PROPFIND request 

1003 # to determine if it exists. 

1004 resp = self._propfind() 

1005 if resp.status_code == requests.codes.multi_status: # 207 

1006 prop = _parse_propfind_response_body(resp.text)[0] 

1007 return prop.exists 

1008 else: # 404 Not Found 

1009 return False 

1010 

1011 def size(self) -> int: 

1012 """Return the size of the remote resource in bytes.""" 

1013 if self.dirLike: 1013 ↛ 1014line 1013 didn't jump to line 1014 because the condition on line 1013 was never true

1014 return 0 

1015 info = self.get_info() 

1016 # dirLike can be None if we are unsure. Only flag if we are certain 

1017 # we have been told this is a directory but webDAV reports it as a 

1018 # file. 

1019 if not info.is_file and self.dirLike is False: 1019 ↛ 1020line 1019 didn't jump to line 1020 because the condition on line 1019 was never true

1020 raise IsADirectoryError( 

1021 f"Resource {self} is reported by server as a directory but has a file path" 

1022 ) 

1023 return info.size 

1024 

1025 def get_info(self) -> ResourceInfo: 

1026 """Return lightweight metadata about this HTTP resource.""" 

1027 if not self.is_webdav_endpoint: 1027 ↛ 1031line 1027 didn't jump to line 1031 because the condition on line 1027 was always true

1028 resp = self._head_non_webdav_url() 

1029 return self._get_info_from_non_webdav_head(resp) 

1030 

1031 resp = self._propfind() 

1032 if resp.status_code != requests.codes.multi_status: 

1033 raise FileNotFoundError( 

1034 f"Resource {self} does not exist, status: {resp.status_code} {resp.reason}" 

1035 ) 

1036 

1037 prop = _parse_propfind_response_body(resp.text)[0] 

1038 if not prop.exists: 

1039 raise FileNotFoundError(f"Resource {self} does not exist") 

1040 

1041 return ResourceInfo( 

1042 uri=str(self), 

1043 is_file=prop.is_file, 

1044 size=prop.size, 

1045 last_modified=prop.last_modified, 

1046 checksums=dict(prop.checksums), 

1047 ) 

1048 

1049 def _get_info_from_non_webdav_head(self, resp: requests.Response) -> ResourceInfo: 

1050 """Build `ResourceInfo` from a non-WebDAV HEAD-like response.""" 

1051 if not self._is_successful_non_webdav_head_request(resp): 

1052 if resp.status_code == requests.codes.not_found: 1052 ↛ 1056line 1052 didn't jump to line 1056 because the condition on line 1052 was always true

1053 raise FileNotFoundError( 

1054 f"Resource {self} does not exist, status: {resp.status_code} {resp.reason}" 

1055 ) 

1056 raise ValueError( 

1057 f"Unexpected response for HEAD request for {self}, status: {resp.status_code} {resp.reason}" 

1058 ) 

1059 

1060 if self.dirLike: 1060 ↛ 1061line 1060 didn't jump to line 1061 because the condition on line 1060 was never true

1061 size = 0 

1062 elif resp.status_code == requests.codes.ok: # 200 

1063 if "Content-Length" not in resp.headers: 1063 ↛ 1064line 1063 didn't jump to line 1064 because the condition on line 1063 was never true

1064 raise ValueError( 

1065 f"Response to HEAD request to {self} does not contain 'Content-Length' header" 

1066 ) 

1067 size = int(resp.headers["Content-Length"]) 

1068 elif resp.status_code == requests.codes.partial_content: 

1069 # 206 Partial Content, returned from a GET request with a Range 

1070 # header (used to emulate HEAD for presigned S3 URLs). 

1071 content_range_header = resp.headers.get("Content-Range") 

1072 if content_range_header is None: 1072 ↛ 1073line 1072 didn't jump to line 1073 because the condition on line 1072 was never true

1073 raise ValueError(f"Response to GET request to {self} did not contain 'Content-Range' header") 

1074 content_range = parse_content_range_header(content_range_header) 

1075 size_total = content_range.total 

1076 if size_total is None: 1076 ↛ 1077line 1076 didn't jump to line 1077 because the condition on line 1076 was never true

1077 raise ValueError(f"Content-Range header for {self} did not include a total file size") 

1078 size = size_total 

1079 else: 

1080 # 416 Range Not Satisfiable can occur on a GET for a 0-byte file. 

1081 size = 0 

1082 

1083 checksums = {} 

1084 digest_header = resp.headers.get("Digest") 

1085 if digest_header is not None: 

1086 for digest in digest_header.split(","): 

1087 algorithm, separator, value = digest.strip().partition("=") 

1088 if separator: 1088 ↛ 1086line 1088 didn't jump to line 1086 because the condition on line 1088 was always true

1089 checksums[algorithm.lower()] = value 

1090 

1091 last_modified = None 

1092 if last_modified_header := resp.headers.get("Last-Modified"): 

1093 last_modified = parsedate_to_datetime(last_modified_header) 

1094 if last_modified.tzinfo is None: 1094 ↛ 1095line 1094 didn't jump to line 1095 because the condition on line 1094 was never true

1095 last_modified = last_modified.replace(tzinfo=datetime.UTC) 

1096 else: 

1097 last_modified = last_modified.astimezone(datetime.UTC) 

1098 

1099 return ResourceInfo( 

1100 uri=str(self), 

1101 is_file=not self.dirLike, 

1102 size=size, 

1103 last_modified=last_modified, 

1104 checksums=checksums, 

1105 ) 

1106 

1107 def _head_non_webdav_url(self) -> requests.Response: 

1108 """Return a response from a HTTP HEAD request for a non-WebDAV HTTP 

1109 URL. 

1110 

1111 Emulates HEAD using a 1-byte GET for presigned S3 URLs. 

1112 """ 

1113 if self._looks_like_presigned_s3_url(): 

1114 # Presigned S3 URLs are signed for a single method only, so you 

1115 # can't call HEAD on a URL signed for GET. However, S3 does 

1116 # support Range requests, so you can ask for a 1-byte range with 

1117 # GET for a similar effect to HEAD. 

1118 # 

1119 # Note that some headers differ between a true HEAD request and the 

1120 # response returned by this GET, e.g. Content-Length will always be 

1121 # 1, and the status code is 206 instead of 200. 

1122 return self.metadata_session.get( 

1123 self.geturl(), 

1124 timeout=self._config.timeout, 

1125 allow_redirects=True, 

1126 stream=False, 

1127 headers={"Range": "bytes=0-0"}, 

1128 ) 

1129 else: 

1130 return self.metadata_session.head( 

1131 self.geturl(), timeout=self._config.timeout, allow_redirects=True, stream=False 

1132 ) 

1133 

1134 def _is_successful_non_webdav_head_request(self, resp: requests.Response) -> bool: 

1135 """Return `True` if the status code in the response indicates a 

1136 successful response to ``_head_non_webdav_url``. 

1137 """ 

1138 return resp.status_code in ( 

1139 requests.codes.ok, # 200, from a normal HEAD or GET request 

1140 requests.codes.partial_content, # 206, returned from a GET request with a Range header. 

1141 # 416, returned from a GET request with a 1-byte Range header that 

1142 # is longer than the 0-byte file. 

1143 requests.codes.range_not_satisfiable, 

1144 ) 

1145 

1146 def _looks_like_presigned_s3_url(self) -> bool: 

1147 """Return `True` if this ResourcePath's URL is likely to be a presigned 

1148 S3 URL. 

1149 """ 

1150 query_params = parse_qs(self._uri.query) 

1151 return "Signature" in query_params and "Expires" in query_params 

1152 

1153 def mkdir(self) -> None: 

1154 """Create the directory resource if it does not already exist.""" 

1155 # Creating directories is only available on WebDAV back ends. 

1156 if not self.is_webdav_endpoint: 

1157 raise NotImplementedError( 

1158 f"Creation of directory {self} is not implemented by plain HTTP servers" 

1159 ) 

1160 

1161 if not self.dirLike: 

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

1163 

1164 # Check if the target directory already exists. 

1165 resp = self._propfind() 

1166 if resp.status_code == requests.codes.multi_status: # 207 

1167 prop = _parse_propfind_response_body(resp.text)[0] 

1168 if prop.exists: 

1169 if prop.is_directory: 

1170 return 

1171 else: 

1172 # A file exists at this path 

1173 raise NotADirectoryError( 

1174 f"Can not create a directory for {self} because a file already exists at that path" 

1175 ) 

1176 

1177 # Target directory does not exist. Create it and its ancestors as 

1178 # needed. We need to test if parent URL is different from self URL, 

1179 # otherwise we could be stuck in a recursive loop 

1180 # where self == parent. 

1181 if self.geturl() != self.parent().geturl(): 

1182 self.parent().mkdir() 

1183 

1184 log.debug("Creating new directory: %s", self.geturl()) 

1185 self._mkcol() 

1186 

1187 def remove(self) -> None: 

1188 """Remove the resource.""" 

1189 self._delete() 

1190 

1191 def read(self, size: int = -1) -> bytes: 

1192 """Open the resource and return the contents in bytes. 

1193 

1194 Parameters 

1195 ---------- 

1196 size : `int`, optional 

1197 The number of bytes to read. Negative or omitted indicates 

1198 that all data should be read. 

1199 """ 

1200 # Use the data session as a context manager to ensure that the 

1201 # network connections to both the front end and back end servers are 

1202 # closed after downloading the data. 

1203 log.debug("Reading from remote resource: %s", self.geturl()) 

1204 stream = size > 0 

1205 with self.data_session as session: 

1206 with time_this(log, msg="GET %s", args=(self,)): 

1207 resp = session.get(self.geturl(), stream=stream, timeout=self._config.timeout) 

1208 

1209 if resp.status_code != requests.codes.ok: # 200 1209 ↛ 1210line 1209 didn't jump to line 1210 because the condition on line 1209 was never true

1210 raise FileNotFoundError( 

1211 f"Unable to read resource {self}; status: {resp.status_code} {resp.reason}" 

1212 ) 

1213 if not stream: 1213 ↛ 1216line 1213 didn't jump to line 1216 because the condition on line 1213 was always true

1214 return resp.content 

1215 else: 

1216 return next(resp.iter_content(chunk_size=size)) 

1217 

1218 def write(self, data: bytes, overwrite: bool = True) -> None: 

1219 """Write the supplied bytes to the new resource. 

1220 

1221 Parameters 

1222 ---------- 

1223 data : `bytes` 

1224 The bytes to write to the resource. The entire contents of the 

1225 resource will be replaced. 

1226 overwrite : `bool`, optional 

1227 If `True` the resource will be overwritten if it exists. Otherwise 

1228 the write will fail. 

1229 """ 

1230 log.debug("Writing to remote resource: %s", self.geturl()) 

1231 if not overwrite and self.exists(): 1231 ↛ 1232line 1231 didn't jump to line 1232 because the condition on line 1231 was never true

1232 raise FileExistsError(f"Remote resource {self} exists and overwrite has been disabled") 

1233 

1234 # Ensure the parent directory exists. 

1235 # This is only meaningful and appropriate for WebDAV, not the general 

1236 # HTTP case. e.g. for S3 HTTP URLs, the underlying service has no 

1237 # concept of 'directories' at all. 

1238 if self.is_webdav_endpoint: 1238 ↛ 1239line 1238 didn't jump to line 1239 because the condition on line 1238 was never true

1239 self.parent().mkdir() 

1240 

1241 # Upload the data. 

1242 log.debug("Writing data to remote resource: %s", self.geturl()) 

1243 self._put(data=data) 

1244 

1245 def transfer_from( 

1246 self, 

1247 src: ResourcePath, 

1248 transfer: str = "copy", 

1249 overwrite: bool = False, 

1250 transaction: TransactionProtocol | None = None, 

1251 multithreaded: bool = True, 

1252 ) -> None: 

1253 """Transfer the current resource to a Webdav repository. 

1254 

1255 Parameters 

1256 ---------- 

1257 src : `ResourcePath` 

1258 Source URI. 

1259 transfer : `str` 

1260 Mode to use for transferring the resource. Supports the following 

1261 options: copy. 

1262 overwrite : `bool`, optional 

1263 Whether overwriting the remote resource is allowed or not. 

1264 transaction : `~lsst.resources.utils.TransactionProtocol`, optional 

1265 Currently unused. 

1266 multithreaded : `bool`, optional 

1267 If `True` the transfer will be allowed to attempt to improve 

1268 throughput by using parallel download streams. This may of no 

1269 effect if the URI scheme does not support parallel streams or 

1270 if a global override has been applied. If `False` parallel 

1271 streams will be disabled. 

1272 """ 

1273 # Fail early to prevent delays if remote resources are requested. 

1274 if transfer not in self.transferModes: 

1275 raise ValueError(f"Transfer mode {transfer} not supported by URI scheme {self.scheme}") 

1276 

1277 # Existence checks cost time so do not call this unless we know 

1278 # that debugging is enabled. 

1279 if log.isEnabledFor(logging.DEBUG): 

1280 log.debug( 

1281 "Transferring %s [exists: %s] -> %s [exists: %s] (transfer=%s)", 

1282 src, 

1283 src.exists(), 

1284 self, 

1285 self.exists(), 

1286 transfer, 

1287 ) 

1288 

1289 # Short circuit immediately if the URIs are identical. 

1290 if self == src: 

1291 log.debug( 

1292 "Target and destination URIs are identical: %s, returning immediately." 

1293 " No further action required.", 

1294 self, 

1295 ) 

1296 return 

1297 

1298 if not overwrite and self.exists(): 

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

1300 

1301 if transfer == "auto": 

1302 transfer = self.transferDefault 

1303 

1304 # We can use webDAV 'COPY' or 'MOVE' if both the current and source 

1305 # resources are located in the same server. 

1306 if isinstance(src, type(self)) and self.root_uri() == src.root_uri() and self.is_webdav_endpoint: 

1307 log.debug("Transfer from %s to %s directly", src, self) 

1308 return self._move(src) if transfer == "move" else self._copy(src) 

1309 

1310 # For resources of different classes or for plain HTTP resources we can 

1311 # perform the copy or move operation by downloading to a local file 

1312 # and uploading to the destination. 

1313 self._copy_via_local(src) 

1314 

1315 # This was an explicit move, try to remove the source. 

1316 if transfer == "move": 

1317 src.remove() 

1318 

1319 def walk( 

1320 self, file_filter: str | re.Pattern | None = None 

1321 ) -> Iterator[list | tuple[ResourcePath, list[str], list[str]]]: 

1322 """Walk the directory tree returning matching files and directories. 

1323 

1324 Parameters 

1325 ---------- 

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

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

1328 

1329 Yields 

1330 ------ 

1331 dirpath : `ResourcePath` 

1332 Current directory being examined. 

1333 dirnames : `list` of `str` 

1334 Names of subdirectories within dirpath. 

1335 filenames : `list` of `str` 

1336 Names of all the files within dirpath. 

1337 """ 

1338 if not self.dirLike: 

1339 raise ValueError("Can not walk a non-directory URI") 

1340 

1341 # Walking directories is only available on WebDAV back ends. 

1342 if not self.is_webdav_endpoint: 

1343 raise NotImplementedError(f"Walking directory {self} is not implemented by plain HTTP servers") 

1344 

1345 if isinstance(file_filter, str): 

1346 file_filter = re.compile(file_filter) 

1347 

1348 resp = self._propfind(depth="1") 

1349 if resp.status_code == requests.codes.multi_status: # 207 

1350 files: list[str] = [] 

1351 dirs: list[str] = [] 

1352 

1353 for prop in _parse_propfind_response_body(resp.text): 

1354 if prop.is_file: 

1355 files.append(prop.name) 

1356 elif not prop.href.rstrip("/").endswith(self.path.rstrip("/")): 

1357 # Only include the names of sub-directories not the name of 

1358 # the directory being walked. 

1359 dirs.append(prop.name) 

1360 

1361 if file_filter is not None: 

1362 files = [f for f in files if file_filter.search(f)] 

1363 

1364 if not dirs and not files: 

1365 return 

1366 else: 

1367 yield type(self)(self, forceAbsolute=False, forceDirectory=True), dirs, files 

1368 

1369 for dir in dirs: 

1370 new_uri = self.join(dir, forceDirectory=True) 

1371 yield from new_uri.walk(file_filter) 

1372 

1373 def generate_presigned_get_url(self, *, expiration_time_seconds: int) -> str: 

1374 """Return a pre-signed URL that can be used to retrieve this resource 

1375 using an HTTP GET without supplying any access credentials. 

1376 

1377 Parameters 

1378 ---------- 

1379 expiration_time_seconds : `int` 

1380 Number of seconds until the generated URL is no longer valid. 

1381 

1382 Returns 

1383 ------- 

1384 url : `str` 

1385 HTTP URL signed for GET. 

1386 """ 

1387 if not self.is_webdav_endpoint: 

1388 # This is already an HTTP URL readable without any authentication 

1389 # credentials, so return it as-is. 

1390 return str(self) 

1391 

1392 return self._sign_with_macaroon(ActivityCaveat.DOWNLOAD, expiration_time_seconds) 

1393 

1394 def generate_presigned_put_url(self, *, expiration_time_seconds: int) -> str: 

1395 """Return a pre-signed URL that can be used to upload a file to this 

1396 path using an HTTP PUT without supplying any access credentials. 

1397 

1398 Parameters 

1399 ---------- 

1400 expiration_time_seconds : `int` 

1401 Number of seconds until the generated URL is no longer valid. 

1402 

1403 Returns 

1404 ------- 

1405 url : `str` 

1406 HTTP URL signed for PUT. 

1407 """ 

1408 if not self.is_webdav_endpoint: 

1409 return super().generate_presigned_put_url(expiration_time_seconds=expiration_time_seconds) 

1410 

1411 return self._sign_with_macaroon(ActivityCaveat.UPLOAD, expiration_time_seconds) 

1412 

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

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

1415 

1416 Returns 

1417 ------- 

1418 fs : `fsspec.spec.AbstractFileSystem` 

1419 A file system object suitable for use with the returned path. 

1420 path : `str` 

1421 A path that can be opened by the file system object. 

1422 """ 

1423 if fsspec is None: 1423 ↛ 1424line 1423 didn't jump to line 1424 because the condition on line 1423 was never true

1424 return super().to_fsspec() 

1425 

1426 if not self.is_webdav_endpoint or self.server not in HttpResourcePath.SUPPORTED_URL_SIGNERS: 1426 ↛ 1429line 1426 didn't jump to line 1429 because the condition on line 1426 was always true

1427 return fsspec.url_to_fs(self.geturl(), client_kwargs={"headers": self._extra_headers}) 

1428 

1429 if self.isdir(): 

1430 raise NotImplementedError( 

1431 f"method HttpResourcePath.to_fsspec() not implemented for directory {self}" 

1432 ) 

1433 

1434 # If usage of fsspec-compatible file system is disabled in the 

1435 # configuration we raise an exception which signals the caller 

1436 # that it cannot use fsspec. An example of such a caller is 

1437 # `lsst.daf.butler.formatters.ParquetFormatter`. 

1438 # 

1439 # Note that we don't call super().to_fsspec() since that method 

1440 # assumes that fsspec can be used provided fsspec package is 

1441 # importable. 

1442 # 

1443 # The motivation for making this configurable is that for HTTP 

1444 # URLs fsspec.HTTPFileSystem uses async I/O and we have found 

1445 # unexpected behavior by clients when used against dCache for reading 

1446 # parquet files via a ParquetFormatter instance. That behavior cannot 

1447 # be reproduced when using other callers. 

1448 # 

1449 # This needs more investigation to discard the possibility that async 

1450 # I/O, used by fsspec.HTTPFileSystem, is related to this behavior. 

1451 if not self._config.fsspec_is_enabled: 

1452 raise ImportError("fsspec is disabled for HttpResourcePath objects with webDAV back end") 

1453 

1454 async def get_client_session(**kwargs: Any) -> ClientSession: 

1455 """Return a aiohttp.ClientSession configured to use an 

1456 `aiohttp.TCPConnector` shared by all instances of this class. 

1457 

1458 Parameters 

1459 ---------- 

1460 **kwargs : `Any` 

1461 Keyword arguments passed unmodified to the contructor of 

1462 `aiohttp.ClientSession`. 

1463 

1464 Returns 

1465 ------- 

1466 session : `aiohttp.ClientSession` 

1467 Client session that `aiohttp.HTTPFileSystem` will use to pool 

1468 TCP connections to the server. 

1469 """ 

1470 if HttpResourcePath._tcp_connector is None: 

1471 HttpResourcePath._tcp_connector = TCPConnector( 

1472 # SSL context equipped with client credentials and 

1473 # configured to validate server certificates. 

1474 ssl=self._config.ssl_context, 

1475 # Total number of simultaneous connections this connector 

1476 # keeps open with any host. 

1477 # 

1478 # The default is 100 but we deliberately reduced it to 

1479 # avoid keeping a large number of open connexions to file 

1480 # servers when thousands of quanta execute simultaneously. 

1481 # 

1482 # In any case, new connexions are automatically established 

1483 # when needed. 

1484 limit=10, 

1485 # Number of simultaneous connections to a single host:port. 

1486 limit_per_host=1, 

1487 # Close network connection after usage 

1488 force_close=True, 

1489 ) 

1490 

1491 connect_timeout, read_timeout = self._config.timeout 

1492 return ClientSession( 

1493 connector=HttpResourcePath._tcp_connector, 

1494 timeout=ClientTimeout( 

1495 connect=connect_timeout, 

1496 sock_connect=connect_timeout, 

1497 sock_read=read_timeout, 

1498 total=2 * read_timeout, 

1499 ), 

1500 **kwargs, 

1501 ) 

1502 

1503 # Retrieve a signed URL for download valid for 2 hours. 

1504 url = self.generate_presigned_get_url(expiration_time_seconds=2 * 3_600) 

1505 

1506 # HTTPFileSystem constructor accepts the argument 'block_size'. The 

1507 # default value is 'fsspec.utils.DEFAULT_BLOCK_SIZE' which is 5 MB. 

1508 # That seems to be a reasonable block size for downloading files. 

1509 return HTTPFileSystem(get_client=get_client_session), url 

1510 

1511 def _sign_with_macaroon(self, activity: ActivityCaveat, expiration_time_seconds: int) -> str: 

1512 # dCache and XRootD webDAV servers support delivery of macaroons. 

1513 # 

1514 # For details about dCache macaroons see: 

1515 # https://www.dcache.org/manuals/UserGuide-9.2/macaroons.shtml 

1516 if self.server is None: 

1517 raise NotImplementedError(f"server for '{self}' does not support signing URLs") 

1518 elif self.server not in HttpResourcePath.SUPPORTED_URL_SIGNERS: 

1519 raise NotImplementedError(f"server '{self.server}' does not support signing for {self}") 

1520 

1521 match activity: 

1522 case ActivityCaveat.DOWNLOAD: 

1523 activity_caveat = "DOWNLOAD,LIST" 

1524 case ActivityCaveat.UPLOAD: 

1525 activity_caveat = "UPLOAD,LIST" 

1526 

1527 # Retrieve a macaroon for the requested activities and duration 

1528 headers = {"Content-Type": "application/macaroon-request"} 

1529 body = { 

1530 "caveats": [ 

1531 f"activity:{activity_caveat}", 

1532 ], 

1533 "validity": f"PT{expiration_time_seconds}S", 

1534 } 

1535 resp = self._post(data=json.dumps(body), headers=headers) 

1536 if resp.status_code != requests.codes.ok: 

1537 raise ValueError( 

1538 f"could not retrieve a macaroon for URL {self}, status: {resp.status_code} {resp.reason}" 

1539 ) 

1540 

1541 # We are expecting the body of the response to be formatted in JSON. 

1542 # dCache sets the 'Content-Type' of the response to 'application/json' 

1543 # but XRootD does not set any 'Content-Type' header 8-[ 

1544 # 

1545 # An example of a response body returned by dCache is shown below: 

1546 # { 

1547 # "macaroon": "MDA[...]Qo", 

1548 # "uri": { 

1549 # "targetWithMacaroon": "https://dcache.example.org/?authz=MD...", 

1550 # "baseWithMacaroon": "https://dcache.example.org/?authz=MD...", 

1551 # "target": "https://dcache.example.org/", 

1552 # "base": "https://dcache.example.org/" 

1553 # } 

1554 # } 

1555 # 

1556 # An example of a response body returned by XRootD is shown below: 

1557 # { 

1558 # "macaroon": "MDA[...]Qo", 

1559 # "expires_in": 86400 

1560 # } 

1561 try: 

1562 response_body = json.loads(resp.text) 

1563 if "macaroon" in response_body: 

1564 return str(self.replace(query=f"authz={response_body['macaroon']}")) 

1565 else: 

1566 raise ValueError(f"could not retrieve macaroon for URL {self}") 

1567 except json.JSONDecodeError: 

1568 raise ValueError(f"could not deserialize response to POST request for URL {self}") 

1569 

1570 @contextlib.contextmanager 

1571 def _as_local( 

1572 self, multithreaded: bool = True, tmpdir: ResourcePath | None = None 

1573 ) -> Generator[ResourcePath]: 

1574 """Download object over HTTP and place in temporary directory. 

1575 

1576 Parameters 

1577 ---------- 

1578 multithreaded : `bool`, optional 

1579 If `True` the transfer will be allowed to attempt to improve 

1580 throughput by using parallel download streams. This may of no 

1581 effect if the URI scheme does not support parallel streams or 

1582 if a global override has been applied. If `False` parallel 

1583 streams will be disabled. 

1584 tmpdir : `ResourcePath` or `None`, optional 

1585 Explicit override of the temporary directory to use for remote 

1586 downloads. 

1587 

1588 Returns 

1589 ------- 

1590 local_uri : `ResourcePath` 

1591 A URI to a local POSIX file corresponding to a local temporary 

1592 downloaded copy of the resource. 

1593 """ 

1594 # Use the session as a context manager to ensure that connections 

1595 # to both the front end and back end servers are closed after the 

1596 # download operation is finished. 

1597 with self.data_session as session: 

1598 resp = session.get(self.geturl(), stream=True, timeout=self._config.timeout) 

1599 if resp.status_code != requests.codes.ok: 

1600 raise FileNotFoundError( 

1601 f"Unable to download resource {self}; status: {resp.status_code} {resp.reason}" 

1602 ) 

1603 

1604 if tmpdir is None: 

1605 temp_dir, buffer_size = self._config.tmpdir_buffersize 

1606 tmpdir = ResourcePath(temp_dir, forceDirectory=True) 

1607 else: 

1608 buffer_size = _calc_tmpdir_buffer_size(tmpdir.ospath) 

1609 

1610 with ResourcePath.temporary_uri( 

1611 suffix=self.getExtension(), prefix=tmpdir, delete=True 

1612 ) as tmp_uri: 

1613 expected_length = int(resp.headers.get("Content-Length", "-1")) 

1614 with time_this( 

1615 log, 

1616 msg="GET %s [length=%d] to local file %s [chunk_size=%d]", 

1617 args=(self, expected_length, tmp_uri, buffer_size), 

1618 mem_usage=self._config.collect_memory_usage, 

1619 mem_unit=u.mebibyte, 

1620 ): 

1621 content_length = 0 

1622 with open(tmp_uri.ospath, "wb", buffering=buffer_size) as tmpFile: 

1623 for chunk in resp.iter_content(chunk_size=buffer_size): 

1624 tmpFile.write(chunk) 

1625 content_length += len(chunk) 

1626 

1627 # Check that the expected and actual content lengths match. 

1628 # Perform this check only when the contents of the file was not 

1629 # encoded by the server. 

1630 if ( 

1631 "Content-Encoding" not in resp.headers 

1632 and expected_length >= 0 

1633 and expected_length != content_length 

1634 ): 

1635 raise ValueError( 

1636 f"Size of downloaded file does not match value in Content-Length header for {self}: " 

1637 f"expecting {expected_length} and got {content_length} bytes" 

1638 ) 

1639 

1640 yield tmp_uri 

1641 

1642 def _send_webdav_request( 

1643 self, 

1644 method: str, 

1645 url: str | None = None, 

1646 headers: dict[str, str] | None = None, 

1647 body: str | None = None, 

1648 session: _SessionWrapper | None = None, 

1649 timeout: tuple[float, float] | None = None, 

1650 ) -> requests.Response: 

1651 """Send a webDAV request and correctly handle redirects. 

1652 

1653 Parameters 

1654 ---------- 

1655 method : `str` 

1656 The mthod of the HTTP request to be sent, e.g. PROPFIND, MKCOL. 

1657 headers : `dict`, optional 

1658 A dictionary of key-value pairs (both strings) to include as 

1659 headers in the request. 

1660 body : `str`, optional 

1661 The body of the request. 

1662 

1663 Notes 

1664 ----- 

1665 This way of sending webDAV requests is necessary for handling 

1666 redirection ourselves, since the 'requests' package changes the method 

1667 of the redirected request when the server responds with status 302 and 

1668 the method of the original request is not HEAD (which is the case for 

1669 webDAV requests). 

1670 

1671 That means that when the webDAV server we interact with responds with 

1672 a redirection to a PROPFIND or MKCOL request, the request gets 

1673 converted to a GET request when sent to the redirected location. 

1674 

1675 See `requests.sessions.SessionRedirectMixin.rebuild_method()` in 

1676 https://github.com/psf/requests/blob/main/requests/sessions.py 

1677 

1678 This behavior of the 'requests' package is meant to be compatible with 

1679 what is specified in RFC 9110: 

1680 

1681 https://www.rfc-editor.org/rfc/rfc9110#name-302-found 

1682 

1683 For our purposes, we do need to follow the redirection and send a new 

1684 request using the same HTTP verb. 

1685 """ 

1686 if url is None: 

1687 url = self.geturl() 

1688 

1689 if headers is None: 

1690 headers = {} 

1691 

1692 if session is None: 

1693 session = self.metadata_session 

1694 

1695 if timeout is None: 

1696 timeout = self._config.timeout 

1697 

1698 with time_this( 

1699 log, 

1700 msg="%s %s", 

1701 args=( 

1702 method, 

1703 url, 

1704 ), 

1705 mem_usage=self._config.collect_memory_usage, 

1706 mem_unit=u.mebibyte, 

1707 ): 

1708 for _ in range(max_redirects := 5): 

1709 resp = session.request( 

1710 method, 

1711 url, 

1712 data=body, 

1713 headers=headers, 

1714 stream=False, 

1715 timeout=timeout, 

1716 allow_redirects=False, 

1717 ) 

1718 if resp.is_redirect: 

1719 url = resp.headers["Location"] 

1720 else: 

1721 return resp 

1722 

1723 # We reached the maximum allowed number of redirects. 

1724 # Stop trying. 

1725 raise ValueError( 

1726 f"Could not get a response to {method} request for {self} after {max_redirects} redirections" 

1727 ) 

1728 

1729 def _propfind(self, body: str | None = None, depth: str = "0") -> requests.Response: 

1730 """Send a PROPFIND webDAV request and return the response. 

1731 

1732 Parameters 

1733 ---------- 

1734 body : `str`, optional 

1735 The body of the PROPFIND request to send to the server. If 

1736 provided, it is expected to be a XML document. 

1737 depth : `str`, optional 

1738 The value of the 'Depth' header to include in the request. 

1739 

1740 Returns 

1741 ------- 

1742 response : `requests.Response` 

1743 Response to the PROPFIND request. 

1744 

1745 Notes 

1746 ----- 

1747 It raises `ValueError` if the status code of the PROPFIND request 

1748 is different from "207 Multistatus" or "404 Not Found". 

1749 """ 

1750 if body is None: 

1751 # Request only the DAV live properties we are explicitly interested 

1752 # in namely 'resourcetype', 'getcontentlength', 'getlastmodified' 

1753 # and 'displayname'. 

1754 body = ( 

1755 """<?xml version="1.0" encoding="utf-8" ?>""" 

1756 """<D:propfind xmlns:D="DAV:"><D:prop>""" 

1757 """<D:resourcetype/><D:getcontentlength/><D:getlastmodified/><D:displayname/>""" 

1758 """</D:prop></D:propfind>""" 

1759 ) 

1760 headers = { 

1761 "Depth": depth, 

1762 "Content-Type": 'application/xml; charset="utf-8"', 

1763 "Content-Length": str(len(body)), 

1764 } 

1765 resp = self._send_webdav_request("PROPFIND", headers=headers, body=body) 

1766 if resp.status_code in (requests.codes.multi_status, requests.codes.not_found): 

1767 return resp 

1768 else: 

1769 raise ValueError( 

1770 f"Unexpected response for PROPFIND request for {self}, status: {resp.status_code} " 

1771 f"{resp.reason}" 

1772 ) 

1773 

1774 def _options(self) -> requests.Response: 

1775 """Send a OPTIONS webDAV request for this resource.""" 

1776 resp = self._send_webdav_request("OPTIONS") 

1777 if resp.status_code in (requests.codes.ok, requests.codes.created): 

1778 return resp 

1779 

1780 raise ValueError( 

1781 f"Unexpected response to OPTIONS request for {self}, status: {resp.status_code} {resp.reason}" 

1782 ) 

1783 

1784 def _head(self) -> requests.Response: 

1785 """Send a HEAD request for this resource.""" 

1786 if not self.is_webdav_endpoint: 1786 ↛ 1789line 1786 didn't jump to line 1789 because the condition on line 1786 was always true

1787 # The remote is a plain HTTP server. 

1788 return self._head_non_webdav_url() 

1789 return self._send_webdav_request("HEAD") 

1790 

1791 def _mkcol(self) -> None: 

1792 """Send a MKCOL webDAV request to create a collection. The collection 

1793 may already exist. 

1794 """ 

1795 resp = self._send_webdav_request("MKCOL") 

1796 if resp.status_code == requests.codes.created: # 201 

1797 return 

1798 

1799 if resp.status_code == requests.codes.method_not_allowed: # 405 

1800 # The remote directory already exists 

1801 log.debug("Can not create directory: %s may already exist: skipping.", self.geturl()) 

1802 else: 

1803 raise ValueError(f"Can not create directory {self}, status: {resp.status_code} {resp.reason}") 

1804 

1805 def _delete(self) -> None: 

1806 """Send a DELETE webDAV request for this resource.""" 

1807 log.debug("Deleting %s ...", self.geturl()) 

1808 

1809 # If this is a directory, ensure the remote is a webDAV server because 

1810 # plain HTTP servers don't support DELETE requests on non-file 

1811 # paths. 

1812 if self.dirLike and not self.is_webdav_endpoint: 

1813 raise NotImplementedError( 

1814 f"Deletion of directory {self} is not implemented by plain HTTP servers" 

1815 ) 

1816 

1817 # Deleting non-empty directories may take some time, so increase 

1818 # the timeout for getting a response from the server. 

1819 timeout = self._config.timeout 

1820 if self.dirLike: 

1821 timeout = (timeout[0], timeout[1] * 100) 

1822 resp = self._send_webdav_request("DELETE", timeout=timeout) 

1823 if resp.status_code in ( 

1824 requests.codes.ok, 

1825 requests.codes.accepted, 

1826 requests.codes.no_content, 

1827 requests.codes.not_found, 

1828 ): 

1829 # We can get a "404 Not Found" error when the file or directory 

1830 # does not exist or when the DELETE request was retried several 

1831 # times and a previous attempt actually deleted the resource. 

1832 # Therefore we consider that a "Not Found" response is not an 

1833 # error since we reached the state desired by the user. 

1834 return 

1835 else: 

1836 # TODO: the response to a DELETE request against a webDAV server 

1837 # may be multistatus. If so, we need to parse the reponse body to 

1838 # determine more precisely the reason of the failure (e.g. a lock) 

1839 # and provide a more helpful error message. 

1840 raise ValueError(f"Unable to delete resource {self}; status: {resp.status_code} {resp.reason}") 

1841 

1842 def _copy_via_local(self, src: ResourcePath) -> None: 

1843 """Replace the contents of this resource with the contents of a remote 

1844 resource by using a local temporary file. 

1845 

1846 Parameters 

1847 ---------- 

1848 src : `HttpResourcePath` 

1849 The source of the contents to copy to `self`. 

1850 """ 

1851 with src.as_local() as local_uri: 

1852 log.debug("Transfer from %s to %s via local file %s", src, self, local_uri) 

1853 with open(local_uri.ospath, "rb") as f: 

1854 self._put(data=f) 

1855 

1856 def _copy_or_move(self, method: str, src: HttpResourcePath) -> None: 

1857 """Send a COPY or MOVE webDAV request to copy or replace the contents 

1858 of this resource with the contents of another resource located in the 

1859 same server. 

1860 

1861 Parameters 

1862 ---------- 

1863 method : `str` 

1864 The method to perform. Valid values are "COPY" or "MOVE" (in 

1865 uppercase). 

1866 src : `HttpResourcePath` 

1867 The source of the contents to move to `self`. 

1868 """ 

1869 headers = {"Destination": self.geturl()} 

1870 resp = self._send_webdav_request(method, url=src.geturl(), headers=headers, session=self.data_session) 

1871 if resp.status_code in (requests.codes.created, requests.codes.no_content): 

1872 return 

1873 

1874 if resp.status_code == requests.codes.multi_status: 

1875 tree = eTree.fromstring(resp.content) 

1876 status_element = tree.find("./{DAV:}response/{DAV:}status") 

1877 status = status_element.text if status_element is not None else "unknown" 

1878 error = tree.find("./{DAV:}response/{DAV:}error") 

1879 raise ValueError(f"{method} returned multistatus reponse with status {status} and error {error}") 

1880 else: 

1881 raise ValueError( 

1882 f"{method} operation from {src} to {self} failed, status: {resp.status_code} {resp.reason}" 

1883 ) 

1884 

1885 def _copy(self, src: HttpResourcePath) -> None: 

1886 """Send a COPY webDAV request to replace the contents of this resource 

1887 (if any) with the contents of another resource located in the same 

1888 server. 

1889 

1890 Parameters 

1891 ---------- 

1892 src : `HttpResourcePath` 

1893 The source of the contents to copy to `self`. 

1894 """ 

1895 # Neither dCache nor XrootD currently implement the COPY 

1896 # webDAV method as documented in 

1897 # http://www.webdav.org/specs/rfc4918.html#METHOD_COPY 

1898 # (See issues DM-37603 and DM-37651 for details) 

1899 # 

1900 # For the time being, we use a temporary local file to 

1901 # perform the copy client side. 

1902 # TODO: when those 2 issues above are solved remove the 3 lines below. 

1903 must_use_local = True 

1904 if must_use_local: 

1905 return self._copy_via_local(src) 

1906 

1907 return self._copy_or_move("COPY", src) 

1908 

1909 def _move(self, src: HttpResourcePath) -> None: 

1910 """Send a MOVE webDAV request to replace the contents of this resource 

1911 with the contents of another resource located in the same server. 

1912 

1913 Parameters 

1914 ---------- 

1915 src : `HttpResourcePath` 

1916 The source of the contents to move to `self`. 

1917 """ 

1918 return self._copy_or_move("MOVE", src) 

1919 

1920 def _post(self, data: str | None = None, headers: dict[str, str] | None = None) -> requests.Response: 

1921 """Perform an HTTP POST request and returns the received response. 

1922 

1923 Parameters 

1924 ---------- 

1925 body : `bytes` 

1926 The contents of the request body. 

1927 """ 

1928 resp = self.metadata_session.request( 

1929 "POST", 

1930 self.geturl(), 

1931 data=data, 

1932 headers=headers, 

1933 stream=False, 

1934 timeout=self._config.timeout, 

1935 allow_redirects=True, 

1936 ) 

1937 if resp.status_code == requests.codes.ok: 

1938 return resp 

1939 

1940 raise ValueError(f"POST request for {self} failed, status: {resp.status_code} {resp.reason}") 

1941 

1942 def _put(self, data: BinaryIO | bytes) -> None: 

1943 """Perform an HTTP PUT request and handle redirection. 

1944 

1945 Parameters 

1946 ---------- 

1947 data : `Union[BinaryIO, bytes]` 

1948 The data to be included in the body of the PUT request. 

1949 """ 

1950 # Retrieve the final URL for this upload by sending a PUT request with 

1951 # no content. Follow a single server redirection to retrieve the 

1952 # final URL. 

1953 headers = {"Content-Length": "0"} 

1954 

1955 # If we are explicitly configured for or if we know the remote server 

1956 # is dCache, send an "Expect" header to signal the server that this 

1957 # client knows how to handle redirection in PUT requests. 

1958 # 

1959 # The goal is that the contents of the file we want to upload is sent 

1960 # directly to the dCache pool without transiting through the dCache 

1961 # webDAV door. Otherwise, the uploaded data would transit twice over 

1962 # the network: first from this client to the dCache webDAV door and 

1963 # second from webDAV door to the target pool (i.e. the dCache file 

1964 # server which will ultimately store the data we will upload). 

1965 # 

1966 # Systematically uploading data via dCache webDAV door could add 

1967 # unnecessary load to the door which we can avoid by instead uploading 

1968 # directly to the dCache pool. 

1969 # 

1970 # For further details see section "Redirection on upload": 

1971 # 

1972 # https://www.dcache.org/manuals/UserGuide-9.2/webdav.shtml#redirection 

1973 if self._config.send_expect_on_put or self.server == "dcache": 1973 ↛ 1974line 1973 didn't jump to line 1974 because the condition on line 1973 was never true

1974 headers["Expect"] = "100-continue" 

1975 

1976 url = self.geturl() 

1977 

1978 # Use the session as a context manager to ensure the underlying 

1979 # connections are closed after finishing uploading the data. 

1980 with self.data_session as session: 

1981 # Send an empty PUT request to get redirected to the final 

1982 # destination. 

1983 log.debug("Sending empty PUT request to %s", url) 

1984 with time_this( 

1985 log, 

1986 msg="PUT (no data) %s", 

1987 args=(url,), 

1988 mem_usage=self._config.collect_memory_usage, 

1989 mem_unit=u.mebibyte, 

1990 ): 

1991 resp = session.request( 

1992 "PUT", 

1993 url, 

1994 data=None, 

1995 headers=headers, 

1996 stream=False, 

1997 timeout=self._config.timeout, 

1998 allow_redirects=False, 

1999 ) 

2000 if resp.is_redirect: 2000 ↛ 2001line 2000 didn't jump to line 2001 because the condition on line 2000 was never true

2001 url = resp.headers["Location"] 

2002 

2003 # Upload the data to the final destination. 

2004 log.debug("Uploading data to %s", url) 

2005 

2006 # Ask the server to compute and record a checksum of the uploaded 

2007 # file contents, for later integrity checks. Since we don't compute 

2008 # the digest ourselves while uploading the data, we cannot control 

2009 # after the request is complete that the data we uploaded is 

2010 # identical to the data recorded by the server, but at least the 

2011 # server has recorded a digest of the data it stored. 

2012 # 

2013 # See RFC-3230 for details and 

2014 # https://www.iana.org/assignments/http-dig-alg/http-dig-alg.xhtml 

2015 # for the list of supported digest algorithhms. 

2016 # In addition, note that not all servers implement this RFC so 

2017 # the checksum may not be computed by the server. 

2018 put_headers: dict[str, str] | None = None 

2019 if digest := self._config.digest_algorithm: 2019 ↛ 2020line 2019 didn't jump to line 2020 because the condition on line 2019 was never true

2020 put_headers = {"Want-Digest": digest} 

2021 

2022 with time_this( 

2023 log, 

2024 msg="PUT %s", 

2025 args=(url,), 

2026 mem_usage=self._config.collect_memory_usage, 

2027 mem_unit=u.mebibyte, 

2028 ): 

2029 resp = session.request( 

2030 "PUT", 

2031 url, 

2032 data=data, 

2033 headers=put_headers, 

2034 stream=False, 

2035 timeout=self._config.timeout, 

2036 allow_redirects=False, 

2037 ) 

2038 if resp.status_code in ( 2038 ↛ 2045line 2038 didn't jump to line 2045 because the condition on line 2038 was always true

2039 requests.codes.ok, 

2040 requests.codes.created, 

2041 requests.codes.no_content, 

2042 ): 

2043 return 

2044 else: 

2045 raise ValueError(f"Can not write file {self}, status: {resp.status_code} {resp.reason}") 

2046 

2047 @contextlib.contextmanager 

2048 def _openImpl( 

2049 self, 

2050 mode: str = "r", 

2051 *, 

2052 encoding: str | None = None, 

2053 ) -> Generator[ResourceHandleProtocol]: 

2054 resp = self._head() 

2055 # A presigned S3 URL is signed for a single method, so _head() emulates 

2056 # HEAD with a one-byte ranged GET, which is answered with 206 rather 

2057 # than 200. 

2058 range_capable = (requests.codes.ok, requests.codes.partial_content) 

2059 accepts_range = resp.status_code in range_capable and resp.headers.get("Accept-Ranges") == "bytes" 

2060 handle: ResourceHandleProtocol 

2061 if mode in ("rb", "r") and accepts_range: 

2062 handle = HttpReadResourceHandle( 

2063 mode, log, self, timeout=self._config.timeout, size=_total_size_from_partial_content(resp) 

2064 ) 

2065 if mode == "r": 2065 ↛ 2068line 2065 didn't jump to line 2068 because the condition on line 2065 was never true

2066 # cast because the protocol is compatible, but does not have 

2067 # BytesIO in the inheritance tree 

2068 yield io.TextIOWrapper(cast(Any, handle), encoding=encoding) 

2069 else: 

2070 yield handle 

2071 else: 

2072 with super()._openImpl(mode, encoding=encoding) as http_handle: 

2073 yield http_handle 

2074 

2075 def _copy_extra_attributes(self, original_uri: ResourcePath) -> None: 

2076 assert isinstance(original_uri, HttpResourcePath) 

2077 self._extra_headers = original_uri._extra_headers 

2078 

2079 

2080def _total_size_from_partial_content(resp: requests.Response) -> int | None: 

2081 """Return the total size of the resource reported by a 206 response. 

2082 

2083 Parameters 

2084 ---------- 

2085 resp : `requests.Response` 

2086 Response to inspect. 

2087 

2088 Returns 

2089 ------- 

2090 size : `int` or `None` 

2091 Total size of the resource in bytes, or `None` if the response does 

2092 not report one. A 200 response is deliberately ignored because its 

2093 'Content-Length' describes the transferred body, which may be 

2094 content-encoded, rather than the resource itself. 

2095 """ 

2096 if resp.status_code != requests.codes.partial_content: 2096 ↛ 2097line 2096 didn't jump to line 2097 because the condition on line 2096 was never true

2097 return None 

2098 if (content_range_header := resp.headers.get("Content-Range")) is None: 2098 ↛ 2099line 2098 didn't jump to line 2099 because the condition on line 2098 was never true

2099 return None 

2100 return parse_content_range_header(content_range_header).total 

2101 

2102 

2103def _dump_response(resp: requests.Response) -> None: 

2104 """Log the contents of a HTTP or webDAV request and its response. 

2105 

2106 Parameters 

2107 ---------- 

2108 resp : `requests.Response` 

2109 The response to log. 

2110 

2111 Notes 

2112 ----- 

2113 Intended for development purposes only. 

2114 """ 

2115 log.debug("-----------------------------------------------") 

2116 log.debug("Request") 

2117 log.debug(" method=%s", resp.request.method) 

2118 log.debug(" URL=%s", resp.request.url) 

2119 log.debug(" headers=%s", resp.request.headers) 

2120 if resp.request.method == "PUT": 

2121 log.debug(" body=<data>") 

2122 elif resp.request.body is None: 

2123 log.debug(" body=<empty>") 

2124 else: 

2125 log.debug(" body=%r", resp.request.body[:120]) 

2126 

2127 log.debug("Response:") 

2128 log.debug(" status_code=%d", resp.status_code) 

2129 log.debug(" headers=%s", resp.headers) 

2130 if not resp.content: 

2131 log.debug(" body=<empty>") 

2132 elif "Content-Type" in resp.headers and resp.headers["Content-Type"] == "text/plain": 

2133 log.debug(" body=%r", resp.content) 

2134 else: 

2135 log.debug(" body=%r", resp.content[:80]) 

2136 

2137 

2138def _is_protected(filepath: str) -> bool: 

2139 """Return true if the permissions of file at filepath only allow for access 

2140 by its owner. 

2141 

2142 Parameters 

2143 ---------- 

2144 filepath : `str` 

2145 Path of a local file. 

2146 """ 

2147 if not os.path.isfile(filepath): 

2148 return False 

2149 mode = stat.S_IMODE(os.stat(filepath).st_mode) 

2150 owner_accessible = bool(mode & stat.S_IRWXU) 

2151 group_accessible = bool(mode & stat.S_IRWXG) 

2152 other_accessible = bool(mode & stat.S_IRWXO) 

2153 return owner_accessible and not group_accessible and not other_accessible 

2154 

2155 

2156def _parse_propfind_response_body(body: str) -> list[DavProperty]: 

2157 """Parse the XML-encoded contents of the response body to a webDAV PROPFIND 

2158 request. 

2159 

2160 Parameters 

2161 ---------- 

2162 body : `str` 

2163 XML-encoded response body to a PROPFIND request 

2164 

2165 Returns 

2166 ------- 

2167 responses : `List[DavProperty]` 

2168 

2169 Notes 

2170 ----- 

2171 Is is expected that there is at least one reponse in `body`, otherwise 

2172 this function raises. 

2173 """ 

2174 # A response body to a PROPFIND request is of the form (indented for 

2175 # readability): 

2176 # 

2177 # <?xml version="1.0" encoding="UTF-8"?> 

2178 # <D:multistatus xmlns:D="DAV:"> 

2179 # <D:response> 

2180 # <D:href>path/to/resource</D:href> 

2181 # <D:propstat> 

2182 # <D:prop> 

2183 # <D:resourcetype> 

2184 # <D:collection xmlns:D="DAV:"/> 

2185 # </D:resourcetype> 

2186 # <D:getlastmodified> 

2187 # Fri, 27 Jan 2 023 13:59:01 GMT 

2188 # </D:getlastmodified> 

2189 # <D:getcontentlength> 

2190 # 12345 

2191 # </D:getcontentlength> 

2192 # </D:prop> 

2193 # <D:status> 

2194 # HTTP/1.1 200 OK 

2195 # </D:status> 

2196 # </D:propstat> 

2197 # </D:response> 

2198 # <D:response> 

2199 # ... 

2200 # </D:response> 

2201 # <D:response> 

2202 # ... 

2203 # </D:response> 

2204 # </D:multistatus> 

2205 

2206 # Scan all the 'response' elements and extract the relevant properties 

2207 responses = [] 

2208 multistatus = eTree.fromstring(body.strip()) 

2209 for response in multistatus.findall("./{DAV:}response"): 

2210 responses.append(DavProperty(response)) 

2211 

2212 if responses: 

2213 return responses 

2214 else: 

2215 # Could not parse the body 

2216 raise ValueError(f"Unable to parse response for PROPFIND request: {body}") 

2217 

2218 

2219class DavProperty: 

2220 """Helper class to encapsulate select live DAV properties of a single 

2221 resource, as retrieved via a PROPFIND request. 

2222 

2223 Parameters 

2224 ---------- 

2225 response : `~xml.etree.ElementTree.Element` or `None` 

2226 The XML response defining the DAV property. 

2227 """ 

2228 

2229 # Regular expression to compare against the 'status' element of a 

2230 # PROPFIND response's 'propstat' element. 

2231 _status_ok_rex = re.compile(r"^HTTP/.* 200 .*$", re.IGNORECASE) 

2232 

2233 def __init__(self, response: Element | None): 

2234 self._href: str = "" 

2235 self._displayname: str = "" 

2236 self._collection: bool = False 

2237 self._getlastmodified: str = "" 

2238 self._getcontentlength: int = -1 

2239 

2240 if response is not None: 

2241 self._parse(response) 

2242 

2243 def _parse(self, response: Element) -> None: 

2244 # Extract 'href'. 

2245 if (element := response.find("./{DAV:}href")) is not None: 

2246 # We need to use "str(element.text)"" instead of "element.text" to 

2247 # keep mypy happy. 

2248 self._href = str(element.text).strip() 

2249 else: 

2250 raise ValueError( 

2251 "Property 'href' expected but not found in PROPFIND response: " 

2252 f"{eTree.tostring(response, encoding='unicode')}" 

2253 ) 

2254 

2255 for propstat in response.findall("./{DAV:}propstat"): 

2256 # Only extract properties of interest with status OK. 

2257 status = propstat.find("./{DAV:}status") 

2258 if status is None or not self._status_ok_rex.match(str(status.text)): 

2259 continue 

2260 

2261 for prop in propstat.findall("./{DAV:}prop"): 

2262 # Parse "collection". 

2263 if (element := prop.find("./{DAV:}resourcetype/{DAV:}collection")) is not None: 

2264 self._collection = True 

2265 

2266 # Parse "getlastmodified". 

2267 if (element := prop.find("./{DAV:}getlastmodified")) is not None: 

2268 self._getlastmodified = str(element.text) 

2269 

2270 # Parse "getcontentlength". 

2271 if (element := prop.find("./{DAV:}getcontentlength")) is not None: 

2272 self._getcontentlength = int(str(element.text)) 

2273 

2274 # Parse "displayname". 

2275 if (element := prop.find("./{DAV:}displayname")) is not None: 

2276 self._displayname = str(element.text) 

2277 

2278 # Some webDAV servers don't include the 'displayname' property in the 

2279 # response so try to infer it from the value of the 'href' property. 

2280 # Depending on the server the href value may end with '/'. 

2281 if not self._displayname: 

2282 self._displayname = os.path.basename(self._href.rstrip("/")) 

2283 

2284 # Force a size of 0 for collections. 

2285 if self._collection: 

2286 self._getcontentlength = 0 

2287 

2288 @property 

2289 def exists(self) -> bool: 

2290 # It is either a directory or a file with length of at least zero 

2291 return self._collection or self._getcontentlength >= 0 

2292 

2293 @property 

2294 def is_directory(self) -> bool: 

2295 return self._collection 

2296 

2297 @property 

2298 def is_file(self) -> bool: 

2299 return not self._collection 

2300 

2301 @property 

2302 def size(self) -> int: 

2303 return self._getcontentlength 

2304 

2305 @property 

2306 def last_modified(self) -> datetime.datetime | None: 

2307 if not self._getlastmodified: 

2308 return None 

2309 

2310 last_modified = parsedate_to_datetime(self._getlastmodified) 

2311 if last_modified.tzinfo is None: 

2312 last_modified = last_modified.replace(tzinfo=datetime.UTC) 

2313 else: 

2314 last_modified = last_modified.astimezone(datetime.UTC) 

2315 return last_modified 

2316 

2317 @property 

2318 def checksums(self) -> dict[str, str]: 

2319 return {} 

2320 

2321 @property 

2322 def name(self) -> str: 

2323 return self._displayname 

2324 

2325 @property 

2326 def href(self) -> str: 

2327 return self._href 

2328 

2329 

2330class _SessionWrapper(contextlib.AbstractContextManager): 

2331 """Wraps a `requests.Session` to allow header values to be injected with 

2332 all requests. 

2333 

2334 Notes 

2335 ----- 

2336 `requests.Session` already has a feature for setting headers globally, but 

2337 our session objects are global and authorization headers can vary for each 

2338 HttpResourcePath instance. 

2339 """ 

2340 

2341 def __init__(self, session: requests.Session, *, extra_headers: dict[str, str] | None) -> None: 

2342 self._session = session 

2343 self._extra_headers = extra_headers 

2344 

2345 def __enter__(self) -> _SessionWrapper: 

2346 self._session.__enter__() 

2347 return self 

2348 

2349 def __exit__( 

2350 self, 

2351 exc_type: Any, 

2352 exc_value: Any, 

2353 traceback: Any, 

2354 ) -> None: 

2355 return self._session.__exit__(exc_type, exc_value, traceback) 

2356 

2357 def get( 

2358 self, 

2359 url: str, 

2360 *, 

2361 timeout: tuple[float, float], 

2362 allow_redirects: bool = True, 

2363 stream: bool, 

2364 headers: dict[str, str] | None = None, 

2365 ) -> requests.Response: 

2366 return self._session.get( 

2367 url, 

2368 timeout=timeout, 

2369 allow_redirects=allow_redirects, 

2370 stream=stream, 

2371 headers=self._augment_headers(headers), 

2372 ) 

2373 

2374 def head( 

2375 self, 

2376 url: str, 

2377 *, 

2378 timeout: tuple[float, float], 

2379 allow_redirects: bool, 

2380 stream: bool, 

2381 headers: dict[str, str] | None = None, 

2382 ) -> requests.Response: 

2383 return self._session.head( 

2384 url, 

2385 timeout=timeout, 

2386 allow_redirects=allow_redirects, 

2387 stream=stream, 

2388 headers=self._augment_headers(headers), 

2389 ) 

2390 

2391 def request( 

2392 self, 

2393 method: str, 

2394 url: str, 

2395 *, 

2396 data: str | bytes | BinaryIO | None, 

2397 timeout: tuple[float, float], 

2398 allow_redirects: bool, 

2399 stream: bool, 

2400 headers: dict[str, str] | None = None, 

2401 ) -> requests.Response: 

2402 return self._session.request( 

2403 method, 

2404 url, 

2405 data=data, 

2406 timeout=timeout, 

2407 allow_redirects=allow_redirects, 

2408 stream=stream, 

2409 headers=self._augment_headers(headers), 

2410 ) 

2411 

2412 def _augment_headers(self, headers: dict[str, str] | None) -> dict[str, str]: 

2413 if headers is None: 

2414 headers = {} 

2415 

2416 if self._extra_headers is not None: 

2417 headers = headers | self._extra_headers 

2418 

2419 return headers