Coverage for python/lsst/resources/http.py: 54%
809 statements
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-22 09:28 +0000
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-22 09:28 +0000
1# This file is part of lsst-resources.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# Use of this source code is governed by a 3-clause BSD-style
10# license that can be found in the LICENSE file.
12from __future__ import annotations
14__all__ = ("HttpResourcePath",)
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
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
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
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
63from urllib.parse import parse_qs
65import requests
66from astropy import units as u
67from requests.adapters import HTTPAdapter
68from requests.auth import AuthBase
69from urllib3.util.retry import Retry
71from lsst.utils.timer import time_this
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
78if TYPE_CHECKING:
79 from .utils import TransactionProtocol
81log = logging.getLogger(__name__)
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.
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.
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
109 if math.isnan(timeout):
110 raise ValueError(f"Unexpected timeout value NaN found in environment variable {env_var}")
112 return timeout
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.
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)
129class HttpResourcePathConfig:
130 """Configuration class to encapsulate the configurable items used by class
131 HttpResourcePath.
132 """
134 # Default timeouts for all HTTP requests (seconds).
135 DEFAULT_TIMEOUT_CONNECT: float = 60.0
136 DEFAULT_TIMEOUT_READ: float = 1_500.0
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
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
149 # Accepted digest algorithms
150 ACCEPTED_DIGESTS: list[str] = ["adler32", "md5", "sha-256", "sha-512"]
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
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
175 default_pool_size = max(_get_num_workers(), self.DEFAULT_FRONTEND_PERSISTENT_CONNECTIONS)
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
184 return self._front_end_connections
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
192 default_pool_size = max(_get_num_workers(), self.DEFAULT_FRONTEND_PERSISTENT_CONNECTIONS)
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
201 return self._back_end_connections
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.
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
217 digest = os.environ.get("LSST_HTTP_DIGEST", "").lower()
218 if digest not in self.ACCEPTED_DIGESTS:
219 digest = ""
221 self._digest_algorithm = digest
222 return self._digest_algorithm
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.
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
236 self._send_expect_on_put = "LSST_HTTP_PUT_SEND_EXPECT_HEADER" in os.environ
237 return self._send_expect_on_put
239 @property
240 def fsspec_is_enabled(self) -> bool:
241 """Return True if `fsspec` is enabled for objects of class
242 HttpResourcePath.
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
250 self._fsspec_is_enabled = "LSST_HTTP_ENABLE_FSSPEC" in os.environ
251 return self._fsspec_is_enabled
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
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
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
277 self._collect_memory_usage = "LSST_HTTP_COLLECT_MEMORY_USAGE" in os.environ
278 return self._collect_memory_usage
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
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
296 return self._backoff_min
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
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
314 return self._backoff_max
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.
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
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
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.
338 Return None if no bearer token is configured in the environment.
339 """
340 if self._client_token != "":
341 return self._client_token
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
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.
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)
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)
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 )
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 )
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)
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
394 tmpdir = get_tempdir()
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)
404 return self._tmpdir_buffersize
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)
419 return self._ssl_context
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.
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.
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)
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)
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
456 return (dav_header, server_header)
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
468class BearerTokenAuth(AuthBase):
469 """Attach a bearer token 'Authorization' header to each request.
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 """
478 def __init__(self, token: str):
479 self._token = self._path = None
480 self._mtime: float = -1.0
481 if not token:
482 return
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()
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
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")
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
516class SessionStore:
517 """Cache a reusable HTTP client session per endpoint.
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 """
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] = {}
549 # Configuration for all instances of HttpResourcePath objects.
550 self._config = config
552 # See documentation of urllib3 PoolManager class:
553 # https://urllib3.readthedocs.io
554 self._num_pools: int = num_pools
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
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
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()
576 self._sessions.clear()
578 def get(self, rpath: ResourcePath) -> requests.Session:
579 """Retrieve a session for accessing the remote resource at rpath.
581 Parameters
582 ----------
583 rpath : `ResourcePath`
584 URL to a resource at the remote server for which a session is to
585 be retrieved.
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".
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)
603 return self._sessions[root_uri]
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 )
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 )
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 )
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
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
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
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
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
708class ActivityCaveat(enum.Enum):
709 """Helper class for enumerating accepted activity caveats for requesting
710 macaroons.
711 """
713 DOWNLOAD = 1
714 UPLOAD = 2
717class HttpResourcePath(ResourcePath):
718 """General HTTP(S) resource.
720 Notes
721 -----
722 In order to configure the behavior of instances of this class, the
723 environment variables below are inspected:
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.
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.
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.
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.
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.
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.
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.
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 """
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.
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).
794 Returns
795 -------
796 instance : `ResourcePath`
797 Newly-created `HttpResourcePath` instance.
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
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")
818 # Configuration items for this class instances.
819 _config: HttpResourcePathConfig = HttpResourcePathConfig()
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 )
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 )
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
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
859 # Additional headers added to every request.
860 _extra_headers: dict[str, str] | None = None
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()
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)
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()
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)
902 def _clear_sessions(self) -> None:
903 """Close the socket connections that are still open.
905 Used only in test suites to avoid warnings.
906 """
907 self._metadata_session_store.clear()
908 self._data_session_store.clear()
910 if hasattr(self, "_metadata_session"):
911 delattr(self, "_metadata_session")
913 if hasattr(self, "_data_session"):
914 delattr(self, "_data_session")
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())
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(",")
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()
946 @property
947 def is_webdav_endpoint(self) -> bool:
948 """Check if the current endpoint implements WebDAV features.
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
956 self._init_server_properties()
957 return self._is_webdav
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.
964 If the remote server does not include that header in its response
965 to an 'OPTIONS' request, server() returns None.
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
972 self._init_server_properties()
973 return self._server
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
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.
987 This is an internal method mainly intended for tests.
988 """
989 HttpResourcePath._config = HttpResourcePathConfig()
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)
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
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
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)
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 )
1037 prop = _parse_propfind_response_body(resp.text)[0]
1038 if not prop.exists:
1039 raise FileNotFoundError(f"Resource {self} does not exist")
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 )
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 )
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
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
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)
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 )
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.
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 )
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 )
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
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 )
1161 if not self.dirLike:
1162 raise NotADirectoryError(f"Can not create a 'directory' for file-like URI {self}")
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 )
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()
1184 log.debug("Creating new directory: %s", self.geturl())
1185 self._mkcol()
1187 def remove(self) -> None:
1188 """Remove the resource."""
1189 self._delete()
1191 def read(self, size: int = -1) -> bytes:
1192 """Open the resource and return the contents in bytes.
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)
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))
1218 def write(self, data: bytes, overwrite: bool = True) -> None:
1219 """Write the supplied bytes to the new resource.
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")
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()
1241 # Upload the data.
1242 log.debug("Writing data to remote resource: %s", self.geturl())
1243 self._put(data=data)
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.
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}")
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 )
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
1298 if not overwrite and self.exists():
1299 raise FileExistsError(f"Destination path {self} already exists.")
1301 if transfer == "auto":
1302 transfer = self.transferDefault
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)
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)
1315 # This was an explicit move, try to remove the source.
1316 if transfer == "move":
1317 src.remove()
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.
1324 Parameters
1325 ----------
1326 file_filter : `str` or `re.Pattern`, optional
1327 Regex to filter out files from the list before it is returned.
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")
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")
1345 if isinstance(file_filter, str):
1346 file_filter = re.compile(file_filter)
1348 resp = self._propfind(depth="1")
1349 if resp.status_code == requests.codes.multi_status: # 207
1350 files: list[str] = []
1351 dirs: list[str] = []
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)
1361 if file_filter is not None:
1362 files = [f for f in files if file_filter.search(f)]
1364 if not dirs and not files:
1365 return
1366 else:
1367 yield type(self)(self, forceAbsolute=False, forceDirectory=True), dirs, files
1369 for dir in dirs:
1370 new_uri = self.join(dir, forceDirectory=True)
1371 yield from new_uri.walk(file_filter)
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.
1377 Parameters
1378 ----------
1379 expiration_time_seconds : `int`
1380 Number of seconds until the generated URL is no longer valid.
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)
1392 return self._sign_with_macaroon(ActivityCaveat.DOWNLOAD, expiration_time_seconds)
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.
1398 Parameters
1399 ----------
1400 expiration_time_seconds : `int`
1401 Number of seconds until the generated URL is no longer valid.
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)
1411 return self._sign_with_macaroon(ActivityCaveat.UPLOAD, expiration_time_seconds)
1413 def to_fsspec(self) -> tuple[AbstractFileSystem, str]:
1414 """Return an abstract file system and path that can be used by fsspec.
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()
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})
1429 if self.isdir():
1430 raise NotImplementedError(
1431 f"method HttpResourcePath.to_fsspec() not implemented for directory {self}"
1432 )
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")
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.
1458 Parameters
1459 ----------
1460 **kwargs : `Any`
1461 Keyword arguments passed unmodified to the contructor of
1462 `aiohttp.ClientSession`.
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 )
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 )
1503 # Retrieve a signed URL for download valid for 2 hours.
1504 url = self.generate_presigned_get_url(expiration_time_seconds=2 * 3_600)
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
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}")
1521 match activity:
1522 case ActivityCaveat.DOWNLOAD:
1523 activity_caveat = "DOWNLOAD,LIST"
1524 case ActivityCaveat.UPLOAD:
1525 activity_caveat = "UPLOAD,LIST"
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 )
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}")
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.
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.
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 )
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)
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)
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 )
1640 yield tmp_uri
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.
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.
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).
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.
1675 See `requests.sessions.SessionRedirectMixin.rebuild_method()` in
1676 https://github.com/psf/requests/blob/main/requests/sessions.py
1678 This behavior of the 'requests' package is meant to be compatible with
1679 what is specified in RFC 9110:
1681 https://www.rfc-editor.org/rfc/rfc9110#name-302-found
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()
1689 if headers is None:
1690 headers = {}
1692 if session is None:
1693 session = self.metadata_session
1695 if timeout is None:
1696 timeout = self._config.timeout
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
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 )
1729 def _propfind(self, body: str | None = None, depth: str = "0") -> requests.Response:
1730 """Send a PROPFIND webDAV request and return the response.
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.
1740 Returns
1741 -------
1742 response : `requests.Response`
1743 Response to the PROPFIND request.
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 )
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
1780 raise ValueError(
1781 f"Unexpected response to OPTIONS request for {self}, status: {resp.status_code} {resp.reason}"
1782 )
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")
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
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}")
1805 def _delete(self) -> None:
1806 """Send a DELETE webDAV request for this resource."""
1807 log.debug("Deleting %s ...", self.geturl())
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 )
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}")
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.
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)
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.
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
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 )
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.
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)
1907 return self._copy_or_move("COPY", src)
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.
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)
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.
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
1940 raise ValueError(f"POST request for {self} failed, status: {resp.status_code} {resp.reason}")
1942 def _put(self, data: BinaryIO | bytes) -> None:
1943 """Perform an HTTP PUT request and handle redirection.
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"}
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"
1976 url = self.geturl()
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"]
2003 # Upload the data to the final destination.
2004 log.debug("Uploading data to %s", url)
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}
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}")
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
2075 def _copy_extra_attributes(self, original_uri: ResourcePath) -> None:
2076 assert isinstance(original_uri, HttpResourcePath)
2077 self._extra_headers = original_uri._extra_headers
2080def _total_size_from_partial_content(resp: requests.Response) -> int | None:
2081 """Return the total size of the resource reported by a 206 response.
2083 Parameters
2084 ----------
2085 resp : `requests.Response`
2086 Response to inspect.
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
2103def _dump_response(resp: requests.Response) -> None:
2104 """Log the contents of a HTTP or webDAV request and its response.
2106 Parameters
2107 ----------
2108 resp : `requests.Response`
2109 The response to log.
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])
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])
2138def _is_protected(filepath: str) -> bool:
2139 """Return true if the permissions of file at filepath only allow for access
2140 by its owner.
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
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.
2160 Parameters
2161 ----------
2162 body : `str`
2163 XML-encoded response body to a PROPFIND request
2165 Returns
2166 -------
2167 responses : `List[DavProperty]`
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>
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))
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}")
2219class DavProperty:
2220 """Helper class to encapsulate select live DAV properties of a single
2221 resource, as retrieved via a PROPFIND request.
2223 Parameters
2224 ----------
2225 response : `~xml.etree.ElementTree.Element` or `None`
2226 The XML response defining the DAV property.
2227 """
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)
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
2240 if response is not None:
2241 self._parse(response)
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 )
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
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
2266 # Parse "getlastmodified".
2267 if (element := prop.find("./{DAV:}getlastmodified")) is not None:
2268 self._getlastmodified = str(element.text)
2270 # Parse "getcontentlength".
2271 if (element := prop.find("./{DAV:}getcontentlength")) is not None:
2272 self._getcontentlength = int(str(element.text))
2274 # Parse "displayname".
2275 if (element := prop.find("./{DAV:}displayname")) is not None:
2276 self._displayname = str(element.text)
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("/"))
2284 # Force a size of 0 for collections.
2285 if self._collection:
2286 self._getcontentlength = 0
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
2293 @property
2294 def is_directory(self) -> bool:
2295 return self._collection
2297 @property
2298 def is_file(self) -> bool:
2299 return not self._collection
2301 @property
2302 def size(self) -> int:
2303 return self._getcontentlength
2305 @property
2306 def last_modified(self) -> datetime.datetime | None:
2307 if not self._getlastmodified:
2308 return None
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
2317 @property
2318 def checksums(self) -> dict[str, str]:
2319 return {}
2321 @property
2322 def name(self) -> str:
2323 return self._displayname
2325 @property
2326 def href(self) -> str:
2327 return self._href
2330class _SessionWrapper(contextlib.AbstractContextManager):
2331 """Wraps a `requests.Session` to allow header values to be injected with
2332 all requests.
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 """
2341 def __init__(self, session: requests.Session, *, extra_headers: dict[str, str] | None) -> None:
2342 self._session = session
2343 self._extra_headers = extra_headers
2345 def __enter__(self) -> _SessionWrapper:
2346 self._session.__enter__()
2347 return self
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)
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 )
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 )
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 )
2412 def _augment_headers(self, headers: dict[str, str] | None) -> dict[str, str]:
2413 if headers is None:
2414 headers = {}
2416 if self._extra_headers is not None:
2417 headers = headers | self._extra_headers
2419 return headers