Coverage for python/lsst/resources/davutils.py: 31%
1134 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-08-05 01:24 -0700
« prev ^ index » next coverage.py v7.15.2, created at 2026-08-05 01:24 -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.
12from __future__ import annotations
14import base64
15import enum
16import io
17import json
18import logging
19import os
20import posixpath
21import random
22import re
23import stat
24import threading
25import time
26import uuid
27import xml.etree.ElementTree as eTree
28from datetime import UTC, datetime
29from http import HTTPStatus
30from typing import Any, BinaryIO
32try:
33 from typing import override # Python 3.12+
34except ImportError:
35 from typing_extensions import override # Python 3.11
37from urllib.parse import parse_qsl, urlencode, urlparse, urlunparse
39try:
40 import fsspec
41 from fsspec.spec import AbstractFileSystem
42except ImportError:
43 fsspec = None
44 AbstractFileSystem = type
46import yaml
47from astropy import units as u
48from urllib3 import PoolManager, make_headers
49from urllib3.response import HTTPResponse
50from urllib3.util import Retry, Timeout, Url, parse_url
52from lsst.utils.logging import getLogger
53from lsst.utils.timer import time_this
55# Use the same logger than `dav.py`.
56log = getLogger(f"""{__name__.replace(".davutils", ".dav")}""")
59def normalize_path(path: str | None) -> str:
60 """Normalize a path intended to be part of a URL.
62 A path of the form "///a/b/c///../d/e/" would be normalized as "/a/b/d/e".
63 The returned path is always absolute, i.e. starts by "/" and never
64 ends by "/" except when the path is exactly "/" and does not contain
65 "." nor "..". It does not contain consecutive "/" either.
67 Parameters
68 ----------
69 path : `str`, optional
70 Path to normalize (e.g., '/path/to/..///normalize/').
72 Returns
73 -------
74 url : `str`
75 Normalized URL (e.g., '/path/normalize').
76 """
77 return "/" if not path else "/" + posixpath.normpath(path).lstrip("/")
80def normalize_url(url: str, preserve_scheme: bool = False, preserve_path: bool = True) -> str:
81 """Normalize a URL so that scheme be 'http' or 'https' and the URL path
82 is normalized.
84 Parameters
85 ----------
86 url : `str`
87 URL to normalize (e.g., 'davs://example.org:1234///path/to//../dir/').
88 preserve_scheme : `bool`
89 If True the scheme of `url` will be preserved. Otherwise the scheme
90 of the returned normalized URL will be 'http' or 'https'.
91 preserve_path : `bool`
92 If True, the path of `url` will be preserved in the returned
93 normalized URL, otherwise, the returned URL will have '/' as path.
95 Returns
96 -------
97 url : `str`
98 Normalized URL (e.g. 'https://example.org:1234/path/to/dir').
99 """
100 parsed = parse_url(url)
101 if parsed.scheme is None: 101 ↛ 102line 101 didn't jump to line 102 because the condition on line 101 was never true
102 scheme = "http"
103 else:
104 scheme = parsed.scheme if preserve_scheme else parsed.scheme.replace("dav", "http")
105 path = normalize_path(parsed.path) if preserve_path else "/"
106 return Url(scheme=scheme, host=parsed.host, port=parsed.port, path=path).url
109def redact_url(url: str) -> str:
110 """Return a modified `url` with authorization query redacted.
112 The goal is that this method should be used for logging URLs to avoid
113 leaking authorization tokens.
115 Parameters
116 ----------
117 url : `str`
118 URL to redact.
120 Returns
121 -------
122 redacted_url : `str`
123 For instance, when called with an URL like:
125 https://host.example.org:1234/a/b/c/file.data?key1=value1&key2=value2&authz=token#fragment
127 the returned value would be:
129 https://host.example.org:1234/a/b/c/file.data?key1=value1&key2=value2&authz=....#fragment
130 """
131 parsed_url = urlparse(url)
132 redacted_query: list[tuple[str, str]] = []
133 for pair in parse_qsl(parsed_url.query):
134 redacted_query.append((pair[0], "...." if pair[0] == "authz" else pair[1]))
136 redacted_url = parsed_url._replace(query=urlencode(redacted_query))
137 return urlunparse(redacted_url)
140class DavConfig:
141 """Configurable settings a webDAV client must use when interacting with a
142 particular storage endpoint.
144 Parameters
145 ----------
146 config : `dict[str, str]`
147 Dictionary of configurable settings for the webdav endpoint which
148 base URL is `config["base_url"]`.
150 For instance, if `config["base_url"]` is
152 "davs://webdav.example.org:1234/"
154 any object of class `DavResourcePath` like
156 "davs://webdav.example.org:1234/path/to/any/file"
158 will use the settings in this configuration to configure its client.
159 """
161 # Timeout in seconds to establish a network connection with the remote
162 # server.
163 DEFAULT_TIMEOUT_CONNECT: float = 10.0
165 # Timeout in seconds to read the response to a request sent to a server.
166 # This is total time for reading both the headers and the response body.
167 # It must be large enough to allow for upload and download of files
168 # of typical size the webdav client supports.
169 DEFAULT_TIMEOUT_READ: float = 300.0
171 # Maximum number of network connections to persist against a single
172 # "host:port" pair. If this endpoint client needs to issue more
173 # simultaneous requests than this number, additional network connections
174 # will be created but won't be persisted after use.
175 DEFAULT_PERSISTENT_CONNECTIONS_PER_HOST: int = 20
177 # Size of the buffer (in mebibytes, i.e. 1024*1024 bytes) the webdav
178 # client of this endpoint will use when sending requests and receiving
179 # responses.
180 DEFAULT_BUFFER_SIZE: int = 5
182 # Size of the block (in mebibytes, i.e. 1024*1024 bytes) the webdav
183 # client of this endpoint will use for making partial reads. Each partial
184 # read will request at least this number of bytes, unless the total size
185 # of the file is lower than this value.
186 DEFAULT_BLOCK_SIZE: int = 1
188 # Number of times to retry requests before failing. Retry happens only
189 # under certain conditions.
190 DEFAULT_RETRIES: int = 3
192 # Minimal and maximal retry backoff (in seconds) for the client to compute
193 # the wait time before retrying a request.
194 # A value in this interval is randomly selected as the backoff factor
195 # every time a request is retried.
196 DEFAULT_RETRY_BACKOFF_MIN: float = 1.0
197 DEFAULT_RETRY_BACKOFF_MAX: float = 3.0
199 # Path to a directory or certificate bundle file where the certificates
200 # of the trusted certificate authorities can be found.
201 # Those certificates will be used by the client of the webdav endpoint
202 # to verify the server's host certificate.
203 # If None, the certificates trusted by the system are used.
204 DEFAULT_TRUSTED_AUTHORITIES: str | None = None
206 # User name and password for the client to authenticate to the server.
207 # If specified, HTTP basic authentication is used on all requests.
208 DEFAULT_USER_NAME: str | None = None
209 DEFAULT_USER_PASSWORD: str | None = None
211 # Path to the client certificate and associated private key the webdav
212 # client must present to the server for authentication purposes.
213 # If None, no client certificate is presented.
214 DEFAULT_USER_CERT: str | None = None
215 DEFAULT_USER_KEY: str | None = None
217 # Token the webdav client must sent to the server for authentication
218 # purposes. The token may be the value of the token itself or the path
219 # to a file where the token can be found.
220 DEFAULT_TOKEN: str | None = None
222 # If this option is set to True, the webdav client attempts to reuse
223 # the network connection to the server as long as possible. Note that
224 # the server can unitaleraly decide to close the connection.
225 # If disabled, the connection is closed after each request.
226 DEFAULT_REUSE_CONNECTION: bool = True
228 # Default checksum algorithm to request the server to compute on every
229 # file upload. Not al servers support this.
230 # See RFC 3230 for details.
231 DEFAULT_REQUEST_CHECKSUM: str | None = None
233 # If this option is set to True, the webdav client can return objects
234 # compliant to the fsspec specification.
235 # See: https://filesystem-spec.readthedocs.io
236 DEFAULT_ENABLE_FSSPEC: bool = True
238 # If this option is set to True, memory usage is computed and reported
239 # when executing in debug mode. Computing memory usage is costly, so only
240 # set this when debugging.
241 DEFAULT_COLLECT_MEMORY_USAGE: bool = False
243 # Accepted checksum algorithms. Must be lowercase.
244 ACCEPTED_CHECKSUMS: list[str] = ["adler32", "md5", "sha-256", "sha-512"]
246 def __init__(self, config: dict | None = None) -> None:
247 if config is None:
248 config = {}
250 if (base_url := expand_vars(config.get("base_url"))) is None:
251 self._base_url = "_default_"
252 else:
253 self._base_url = normalize_url(base_url, preserve_path=False)
255 self._timeout_connect: float = float(config.get("timeout_connect", DavConfig.DEFAULT_TIMEOUT_CONNECT))
256 self._timeout_read: float = float(config.get("timeout_read", DavConfig.DEFAULT_TIMEOUT_READ))
257 self._persistent_connections_per_host: int = int(
258 config.get(
259 "persistent_connections_per_host",
260 DavConfig.DEFAULT_PERSISTENT_CONNECTIONS_PER_HOST,
261 )
262 )
263 self._buffer_size: int = 1_048_576 * int(config.get("buffer_size", DavConfig.DEFAULT_BUFFER_SIZE))
264 self._block_size: int = 1_048_576 * int(config.get("block_size", DavConfig.DEFAULT_BLOCK_SIZE))
265 self._retries: int = int(config.get("retries", DavConfig.DEFAULT_RETRIES))
266 self._retry_backoff_min: float = float(
267 config.get("retry_backoff_min", DavConfig.DEFAULT_RETRY_BACKOFF_MIN)
268 )
269 self._retry_backoff_max: float = float(
270 config.get("retry_backoff_max", DavConfig.DEFAULT_RETRY_BACKOFF_MAX)
271 )
272 self._trusted_authorities: str | None = expand_vars(
273 config.get("trusted_authorities", DavConfig.DEFAULT_TRUSTED_AUTHORITIES)
274 )
275 self._user_name: str | None = expand_vars(config.get("user_name", DavConfig.DEFAULT_USER_NAME))
276 self._user_password: str | None = expand_vars(
277 config.get("user_password", DavConfig.DEFAULT_USER_PASSWORD)
278 )
279 self._user_cert: str | None = expand_vars(config.get("user_cert", DavConfig.DEFAULT_USER_CERT))
280 self._user_key: str | None = expand_vars(config.get("user_key", DavConfig.DEFAULT_USER_KEY))
281 self._token: str | None = expand_vars(config.get("token", DavConfig.DEFAULT_TOKEN))
282 self._reuse_connection: bool = config.get("reuse_connection", DavConfig.DEFAULT_REUSE_CONNECTION)
283 self._enable_fsspec: bool = config.get("enable_fsspec", DavConfig.DEFAULT_ENABLE_FSSPEC)
284 self._frontend_urls: list[str] = self._init_frontend_urls(config=config)
285 self._collect_memory_usage: bool = config.get(
286 "collect_memory_usage", DavConfig.DEFAULT_COLLECT_MEMORY_USAGE
287 )
288 self._request_checksum: str | None = config.get(
289 "request_checksum", DavConfig.DEFAULT_REQUEST_CHECKSUM
290 )
291 if self._request_checksum is not None:
292 self._request_checksum = self._request_checksum.lower()
293 if self._request_checksum not in DavConfig.ACCEPTED_CHECKSUMS: 293 ↛ 294line 293 didn't jump to line 294 because the condition on line 293 was never true
294 raise ValueError(
295 f"""Value for checksum algorithm {self._request_checksum} for storage endpoint """
296 f"""{self._base_url} is not among the accepted values: {DavConfig.ACCEPTED_CHECKSUMS}"""
297 )
299 def _init_frontend_urls(self, config: dict | None = None) -> list[str]:
300 if config is None: 300 ↛ 301line 300 didn't jump to line 301 because the condition on line 300 was never true
301 return []
303 # Initialize the URLs of the frontend servers, if present in
304 # the configuration.
305 frontend_urls: list[str] = []
306 for url in config.get("frontend_base_urls", []):
307 # Expand environment variables in this URL
308 if (expanded_url := expand_vars(url)) is not None: 308 ↛ 306line 308 didn't jump to line 306 because the condition on line 308 was always true
309 frontend_urls.append(normalize_url(expanded_url, preserve_path=False))
311 # Eliminate duplicate URLs.
312 frontend_urls = list(set(frontend_urls))
314 # Check that the scheme of this client's base URL is identical to
315 # the scheme of the frontend server URLs.
316 base_url_scheme = parse_url(self._base_url).scheme
317 for url in frontend_urls:
318 if base_url_scheme != parse_url(url).scheme:
319 raise ValueError(
320 f"""inconsistent scheme in frontend URL {url} for endpoint """
321 f"""with base URL {self._base_url}"""
322 )
324 return frontend_urls
326 @property
327 def base_url(self) -> str:
328 return self._base_url
330 @property
331 def timeout_connect(self) -> float:
332 return self._timeout_connect
334 @property
335 def timeout_read(self) -> float:
336 return self._timeout_read
338 @property
339 def persistent_connections_per_host(self) -> int:
340 return self._persistent_connections_per_host
342 @property
343 def buffer_size(self) -> int:
344 return self._buffer_size
346 @property
347 def block_size(self) -> int:
348 return self._block_size
350 @property
351 def retries(self) -> int:
352 return self._retries
354 @property
355 def retry_backoff_min(self) -> float:
356 return self._retry_backoff_min
358 @property
359 def retry_backoff_max(self) -> float:
360 return self._retry_backoff_max
362 @property
363 def trusted_authorities(self) -> str | None:
364 return self._trusted_authorities
366 @property
367 def token(self) -> str | None:
368 return self._token
370 @property
371 def reuse_connection(self) -> bool:
372 return self._reuse_connection
374 @property
375 def request_checksum(self) -> str | None:
376 return self._request_checksum
378 @property
379 def user_cert(self) -> str | None:
380 return self._user_cert
382 @property
383 def user_key(self) -> str | None:
384 # If no user certificate was specified in the configuration,
385 # ignore the private key, even if it was provided.
386 if self._user_cert is None:
387 return None
389 # If we have a user certificate but not a private key, assume the
390 # private key is included in the same file as the user certificate.
391 # That is typically the case when using a X.509 grid proxy as
392 # client certificate.
393 return self._user_cert if self._user_key is None else self._user_key
395 @property
396 def user_name(self) -> str | None:
397 return self._user_name
399 @property
400 def user_password(self) -> str | None:
401 return self._user_password
403 @property
404 def enable_fsspec(self) -> bool:
405 return self._enable_fsspec
407 @property
408 def collect_memory_usage(self) -> bool:
409 return self._collect_memory_usage
411 @property
412 def frontend_urls(self) -> list[str]:
413 return self._frontend_urls
416class DavConfigPool:
417 """Registry of configurable settings for all known webDAV endpoints.
419 Parameters
420 ----------
421 filename : `list` [ `str` ]
422 List of environment variables or file names to load the configuration
423 from. The first file found in the list will be read and the
424 configuration settings for all webDAV endpoints will be extracted
425 from it. Other files will be ignored.
427 Each component of `filenames` can be an environment variable or
428 the path of a file which itself can include an environment variable,
429 e.g. '$HOME/path/to/config.yaml'.
431 The configuration file is a YAML file with the structure below:
433 - base_url: "davs://webdav1.example.org:1234/"
434 persistent_connections_per_host: 10
435 timeout_connect: 20.0
436 timeout_read: 120.0
437 retries: 3
438 retry_backoff_min: 1.0
439 retry_backoff_max: 3.0
440 user_cert: "${X509_USER_PROXY}"
441 user_key: "${X509_USER_PROXY}"
442 token: "/path/to/bearer/token/file"
443 trusted_authorities: "/etc/grid-security/certificates"
444 buffer_size: 5
445 enable_fsspec: false
446 request_checksum: "md5"
447 collect_memory_usage: false
449 - base_url: "davs://webdav2.example.org:1234/"
450 user_name: "user"
451 user_password: "password"
452 persistent_connections_per_host: 5
453 reuse_connection: false
454 ...
456 All settings are optional. If no settings are found in the
457 configuration file for a particular webDAV endpoint, sensible
458 defaults will be used.
460 There is only a single instance of this class. This thead-safe
461 singleton is intended to be initialized when the module is imported
462 the first time.
463 """
465 _instance = None
466 _lock = threading.Lock()
468 def __new__(cls, filename: str | None = None) -> DavConfigPool:
469 if cls._instance is None:
470 with cls._lock:
471 if cls._instance is None: 471 ↛ 474line 471 didn't jump to line 474
472 cls._instance = super().__new__(cls)
474 return cls._instance
476 def __init__(self, filename: str | None = None) -> None:
477 # Create a default configuration. This configuration is
478 # used when a URL doest not match any of the endpoints in the
479 # configuration.
480 self._default_config: DavConfig = DavConfig()
482 # The key of this dictionary is the URL of the webDAV endpoint,
483 # e.g. "davs://host.example.org:1234/"
484 self._configs: dict[str, DavConfig] = {}
486 # Load the configuration from the file we have been provided with,
487 # if any.
488 if filename is None:
489 return
491 # filename can be the name of an environment variable or a path.
492 # A path can include environment variables
493 # (e.g. "$HOME/path/to/config.yaml") or "~"
494 # (e.g. "~/path/to/config.yaml")
495 if (filename := os.getenv(filename)) is not None:
496 # Expand environment variables and '~' in the file name, if any.
497 filename = os.path.expandvars(filename)
498 filename = os.path.expanduser(filename)
499 with open(filename) as file:
500 for config_item in yaml.safe_load(file):
501 config = DavConfig(config_item)
502 if config.base_url not in self._configs:
503 self._configs[config.base_url] = config
504 else:
505 # We already have a configuration for the same
506 # endpoint. That is likely a human error in
507 # the configuration file.
508 raise ValueError(
509 f"""configuration file {filename} contains two configurations for """
510 f"""endpoint {config.base_url}"""
511 )
513 def get_config_for_url(self, url: str) -> DavConfig:
514 """Return the configuration to use a webDAV client when interacting
515 with the server which hosts the resource at `url`.
517 Parameters
518 ----------
519 url : `str`
520 URL for which to obtain a configuration.
521 """
522 # Select the configuration for the endpoint of the provided URL.
523 normalized_url: str = normalize_url(url, preserve_path=False)
524 if (config := self._configs.get(normalized_url)) is not None:
525 return config
527 # No config was found for the specified URL. Use the default.
528 return self._default_config
530 def _destroy(self) -> None:
531 """Destroy this class singleton instance.
533 Helper method to be used in tests to reset global configuration.
534 """
535 with DavConfigPool._lock:
536 DavConfigPool._instance = None
539def make_retry(config: DavConfig) -> Retry:
540 """Create a ``urllib3.util.Retry`` object from settings in `config`.
542 Parameters
543 ----------
544 config : `DavConfig`
545 Configurable settings for a webDAV storage endpoint.
547 Returns
548 -------
549 retry : `urllib3.util.Retry`
550 Retry object to he used when creating a ``urllib3.PoolManager``.
551 """
552 backoff_min: float = config.retry_backoff_min
553 backoff_max: float = config.retry_backoff_max
554 retry = Retry(
555 # Total number of retries to allow. Takes precedence over other
556 # counts.
557 total=3 * config.retries,
558 # How many connection-related errors to retry on.
559 connect=config.retries,
560 # How many times to retry on read errors.
561 read=config.retries,
562 # How many times to retry on bad status codes.
563 status=config.retries,
564 # How many times to retry on other errors.
565 other=config.retries,
566 # Backoff factor to apply between attempts after the second try
567 # (seconds). Compute a random jitter to prevent all the clients which
568 # started at the same time (even on different hosts) to overwhelm the
569 # server by sending requests at the same time.
570 backoff_factor=backoff_min + (backoff_max - backoff_min) * random.random(),
571 # How many redirects to perform. Set to a finite value to avoid
572 # infinite redirect loops.
573 redirect=3,
574 # Set of uppercased HTTP method verbs that we should retry on.
575 # By default, we automatically retry idempotent requests. Specific
576 # retry configuration may be set for some requests such as
577 # non-idempotent `PUT` requests.
578 allowed_methods=frozenset(
579 [
580 "COPY",
581 "DELETE",
582 "GET",
583 "HEAD",
584 "MKCOL",
585 "OPTIONS",
586 "PROPFIND",
587 "PUT",
588 ]
589 ),
590 # HTTP status codes that we should force a retry on.
591 status_forcelist=frozenset(
592 [
593 HTTPStatus.TOO_MANY_REQUESTS, # 429
594 HTTPStatus.INTERNAL_SERVER_ERROR, # 500
595 HTTPStatus.BAD_GATEWAY, # 502
596 HTTPStatus.SERVICE_UNAVAILABLE, # 503
597 HTTPStatus.GATEWAY_TIMEOUT, # 504
598 ]
599 ),
600 # Whether to respect "Retry-After" header on status codes defined
601 # above.
602 respect_retry_after_header=True,
603 )
604 return retry
607class DavClientPool:
608 """Container of reusable webDAV clients, each one specifically configured
609 to talk to a single storage endpoint.
611 Parameters
612 ----------
613 config_pool : `DavConfigPool`
614 Pool of all known webDAV client configurations.
616 Notes
617 -----
618 There is a single instance of this class. This thead-safe singleton is
619 intended to be initialized when the module is imported the first time.
620 """
622 _instance = None
623 _lock = threading.Lock()
625 def __new__(cls, config_pool: DavConfigPool) -> DavClientPool:
626 if cls._instance is None: 626 ↛ 631line 626 didn't jump to line 631 because the condition on line 626 was always true
627 with cls._lock:
628 if cls._instance is None: 628 ↛ 631line 628 didn't jump to line 631
629 cls._instance = super().__new__(cls)
631 return cls._instance
633 def __init__(self, config_pool: DavConfigPool) -> None:
634 self._config_pool: DavConfigPool = config_pool
636 # The key of this dictionnary is a path-stripped URL of the form
637 # "davs://host.example.org:1234/". The value is a reusable
638 # DavClient to interact with that endpoint.
639 self._clients: dict[str, DavClient] = {}
641 def get_client_for_url(self, url: str) -> DavClient:
642 """Return a client for interacting with the endpoint where `url`
643 is hosted.
645 Parameters
646 ----------
647 url : `str`
648 URL for which to obtain a client.
650 Notes
651 -----
652 The returned client is thread-safe. If a client for that endpoint
653 already exists it is reused, otherwise a new client is created
654 with the appropriate configuration for interacting with the storage
655 endpoint.
656 """
657 # If we already have a client for this endpoint reuse it.
658 url = normalize_url(url, preserve_path=False)
659 if (client := self._clients.get(url)) is not None:
660 return client
662 # No client for this endpoint was found. Create a new one and save it
663 # for serving subsequent requests.
664 with DavClientPool._lock:
665 # If another client was created in the meantime by another thread
666 # reuse it.
667 if (client := self._clients.get(url)) is not None:
668 return client
670 config: DavConfig = self._config_pool.get_config_for_url(url)
671 self._clients[url] = self._make_client(url, config)
673 return self._clients[url]
675 def _make_client(self, url: str, config: DavConfig) -> DavClient:
676 """Make a webDAV client for interacting with the server at `url`."""
677 # Check the server implements webDAV protocol and retrieve its
678 # identity so that we can build a client for that specific
679 # server implementation.
680 client = DavClient(url, config)
681 server_details = client.get_server_details(url)
682 server_id = server_details.get("Server", None)
683 accepts_ranges: bool | str | None = server_details.get("Accept-Ranges", None)
684 if accepts_ranges is not None:
685 accepts_ranges = accepts_ranges == "bytes"
687 if server_id is None:
688 # Create a generic webDAV client
689 return DavClient(url, config, accepts_ranges)
690 server_id = server_id.lower()
691 if server_id.startswith("dcache"):
692 # Create a client for a dCache webDAV server
693 return DavClientDCache(url, config, accepts_ranges)
694 elif server_id.startswith("xrootd"):
695 # Create a client for a XrootD webDAV server
696 return DavClientXrootD(url, config, accepts_ranges)
697 else:
698 # Return a generic webDAV client
699 return DavClient(url, config, accepts_ranges)
701 def _destroy(self) -> None:
702 """Destroy this class singleton instance.
704 Helper method to be used in tests to reset global configuration.
705 """
706 with DavClientPool._lock:
707 DavClientPool._instance = None
710class DavFileSizeCache:
711 """Helper class to cache file sizes of recently uploaded files.
713 Parameters
714 ----------
715 default_timeout : `float`, optional
716 Default validity period, in seconds, of the entries in this cache.
717 The validity period for a specific entry can be specified when the
718 entry is added to the cache (see `update_size` method).
720 Notes
721 -----
722 There is a single instance of this class shared by several `DavClient`
723 objects. This singleton is thread safe.
725 Caching file sizes helps preventing sending requests to the server for
726 retrieving the size of recently uploaded files. This is in particular
727 intended to efficiently serve `Butler` requests for the size of a file it
728 just wrote to the datastore.
729 """
731 _instance = None
732 _lock = threading.Lock()
734 def __new__(cls) -> DavFileSizeCache:
735 if cls._instance is None:
736 with cls._lock:
737 if cls._instance is None:
738 cls._instance = super().__new__(cls)
740 return cls._instance
742 def __init__(self, default_timeout: float = 60.0) -> None:
743 # The key of the cache dictionnary is a URL of the form
744 #
745 # "https://host.example.org:1234/path/to/file".
746 #
747 # The value is a triplet (file_size, last_updated, timeout) where:
748 # - 'file_size' is the size of the file in bytes,
749 # - 'last_updated' is the time when this entry was added to the cache
750 # or last updated, in seconds since epoch,
751 # - 'timeout' is the validity period of this cache entry, in seconds,
752 # understood from the moment the cache entry was created.
753 with DavFileSizeCache._lock:
754 if not hasattr(self, "_cache"):
755 self._default_timeout: float = default_timeout
756 self._cache: dict[str, tuple[int, float, float]] = {}
758 def invalidate(self, url: str) -> None:
759 """Invalidate the cache entry for `url`, if any.
761 Parameters
762 ----------
763 url : `str`
764 URL of the file to invalidate which cache entry must be
765 invalidated.
766 """
767 with DavFileSizeCache._lock:
768 self._cache.pop(url, None)
770 def update_size(self, url: str, size: int | None, timeout: float | None = None) -> None:
771 """Update the cache with an entry for `url` which has a size of `size`
772 bytes. This entry is considered valid for a period of `timeout`
773 seconds from now.
775 Parameters
776 ----------
777 url : `str`
778 URL of the file the size to be cached.
779 size : `size` or `None`, optional
780 Size in bytes of the file at `url`. If this value is `None`, the
781 cache is not modified.
782 timeout : `float` or `None`, optional
783 The validity period, in seconds, this size is to be considered
784 valid. If not specified, the default value specified when this
785 object was created will be used for this cache entry.
786 """
787 if size is None:
788 return
790 timeout = self._default_timeout if timeout is None else timeout
791 with DavFileSizeCache._lock:
792 self._cache[url] = (size, time.time(), timeout)
794 def get_size(self, url: str) -> int | None:
795 """Retrieve the cached valued of the size of file at `url`.
797 Parameters
798 ----------
799 url : `str`
800 URL of the file to retrieve the size for.
802 Returns
803 -------
804 `size`: `int` or `None`
805 The cached value of the size of file at `url` if any value was
806 found in the cache, `None` otherwise.
807 `None` is also returned if there is a cached value but its
808 validity period has expired. In this case, the entry associated to
809 `url` is removed from the cache.
810 """
811 with DavFileSizeCache._lock:
812 if (entry := self._cache.get(url, None)) is None:
813 # There is no entry in the cache for this URL
814 return None
816 # There is an entry in the cache for this URL. Check that
817 # its validity period has not yet expired.
818 size, last_updated, timeout = entry
819 if time.time() <= last_updated + timeout:
820 # This entry is stil valid
821 return size
822 else:
823 # This entry is no longer valid. Remove it from the cache.
824 self._cache.pop(url)
825 return None
828def unexpected_status_error(method: str, url: str, resp: HTTPResponse) -> Exception:
829 """Raise an exception from `resp`.
831 Parameters
832 ----------
833 method : `str`
834 The method name triggering the error.
835 url : `str`
836 The URL that cause the error.
837 resp : `resp`
838 The error response.
839 """
840 message = f"Unexpected response to HTTP request {method} {redact_url(url)}: {resp.status} {resp.reason}"
841 body = resp.data.decode()
842 if len(body) > 0:
843 message += f" [response body: {body}]"
845 return ValueError(message)
848class DavClient:
849 """WebDAV client, configured to talk to a single storage endpoint.
851 Instances of this class are thread-safe.
853 Parameters
854 ----------
855 url : `str`
856 Root URL of the storage endpoint (e.g.
857 "https://host.example.org:1234/").
858 config : `DavConfig`
859 Configuration to initialize this client.
860 accepts_ranges : `bool` | `None`
861 Indicate whether the remote server accepts the ``Range`` header in GET
862 requests.
863 """
865 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None:
866 # Lock to protect this client fields from concurrent modification.
867 self._lock = threading.Lock()
869 # Base URL of the server this client will interact with.
870 # It is of the form: "davs://host.example.org:1234/"
871 self._base_url: str = url
873 # Configuration settings for the storage endpoint this client
874 # will interact with.
875 self._config: DavConfig = config
877 # Make the authorizer for this client's requests.
878 self._authorizer: Authorizer | None = self._make_authorizer(config=self._config)
880 # Make the pool manager for this client to use for sending
881 # requests to the server.
882 self._pool_manager: PoolManager = self._make_pool_manager(config=self._config)
884 # Parser of PROPFIND responses.
885 self._propfind_parser: DavPropfindParser = DavPropfindParser()
887 # Does the remote server accept a "Range" header in GET requests?
888 # This field is lazy initialized.
889 self._accepts_ranges: bool | None = accepts_ranges
891 # Can this client use a COPY request to duplicate files within a
892 # single webDAV server?
893 # Subclasses can overwrite this setting according to the server
894 # capabilities and compliance to webDAV RFC.
895 self._can_duplicate: bool = True
897 # Cache to store sizes of files this client has recently uploaded
898 # to the server.
899 self._file_size_cache = DavFileSizeCache()
901 def _make_authorizer(self, config: DavConfig) -> Authorizer | None:
902 # If a token was specified in the configuration settings for this
903 # endpoint, prefer it as the authentication method, even if other
904 # authentication settings were also specified.
905 if config.token is not None:
906 return TokenAuthorizer(token=config.token)
907 elif config.user_name is not None and config.user_password is not None:
908 return BasicAuthorizer(user_name=config.user_name, user_password=config.user_password)
910 return None
912 def _make_pool_manager(self, config: DavConfig) -> PoolManager:
913 # Prepare the trusted authorities certificates
914 ca_certs, ca_cert_dir = None, None
915 if config.trusted_authorities is not None:
916 if os.path.isdir(config.trusted_authorities):
917 ca_cert_dir = config.trusted_authorities
918 elif os.path.isfile(config.trusted_authorities):
919 ca_certs = config.trusted_authorities
920 else:
921 raise FileNotFoundError(
922 f"Trusted authorities file or directory {config.trusted_authorities} does not exist"
923 )
925 # If a token was specified for this endpoint don't use the
926 # <user certificate, private key> pair, even if they were also
927 # specified.
928 user_cert, user_key = None, None
929 if config.token is None:
930 user_cert = config.user_cert
931 user_key = config.user_key
933 # Pool manager for sending requests. Connections in this pool manager
934 # are generally left open by the client but the front-end server may
935 # choose to close them in some specific situations. For instance,
936 # whe serving a PUT request, the front server may redirect to a
937 # backend server and close the network connection making it
938 # unsuable for subsequent requests.
939 #
940 # In addition, the client may also choose to explicitly close the
941 # network connection after receiving a response.
942 return PoolManager(
943 # Number of connection pools to cache before discarding the least
944 # recently used pool. Each connection pool manages network
945 # connections to a single host, so this is basically the number
946 # of "host:port" we persist network connections to.
947 num_pools=200,
948 # Number of connections to the same "host:port" to persist for
949 # later reuse. More than 1 is useful in multithreaded situations.
950 # If more than this number of network connections are needed at
951 # a particular moment, they will be created and discarded after
952 # use.
953 maxsize=config.persistent_connections_per_host,
954 # Retry configuration to use by default with requests sent to
955 # host in the front end.
956 retries=make_retry(config),
957 # Socket timeout in seconds for each individual connection.
958 timeout=Timeout(
959 connect=config.timeout_connect,
960 read=config.timeout_read,
961 ),
962 # Size in bytes of the buffer for reading/writing data from/to
963 # the underlying socket.
964 blocksize=config.buffer_size,
965 # Client certificate and private key for esablishing TLS
966 # connections. If None, no client certificate is sent to the
967 # server. Only relevant for endpoints using secure HTTP protocol.
968 cert_file=user_cert,
969 key_file=user_key,
970 # We require verification of the server certificate.
971 cert_reqs="CERT_REQUIRED",
972 # Directory where the certificates of the trusted certificate
973 # authorities can be found. The contents of that directory
974 # must be as expected by OpenSSL.
975 ca_cert_dir=ca_cert_dir,
976 # Path to a file of concatenated CA certificates in PEM format.
977 ca_certs=ca_certs,
978 )
980 def get_server_details(self, url: str) -> dict[str, str]:
981 """Retrieve the details of the server and check it advertises
982 compliance to class 1 of webDAV protocol.
984 Parameters
985 ----------
986 url : `str`
987 URL to check.
989 Returns
990 -------
991 details: `dic[str, str]`
992 The keys of the returned dictionary can be "Server" and
993 "Accept-Ranges". Any of those keys may not exist in the returned
994 dictionary if the server did not include it in its response.
996 The values are the values of the corresponding
997 headers found in the response to the OPTIONS request.
998 Examples of values for the "Server" header are 'dCache/9.2.4' or
999 'XrootD/v5.7.1'.
1000 """
1001 # Check that the value "1" is part of the value of the "DAV" header in
1002 # the response to an 'OPTIONS' request.
1003 #
1004 # We don't rely on webDAV locks, so a server complying to class 1 is
1005 # enough for our purposes. All webDAV servers must advertise at least
1006 # compliance class "1".
1007 #
1008 # Compliance classes are documented in
1009 # http://www.webdav.org/specs/rfc4918.html#dav.compliance.classes
1010 #
1011 # Examples of values for header DAV are:
1012 # DAV: 1, 2
1013 # DAV: 1, <http://apache.org/dav/propset/fs/1>
1014 resp = self.options(url)
1015 if "DAV" not in resp.headers:
1016 raise ValueError(f"Server of {resp.geturl()} does not implement webDAV protocol")
1018 if "1" not in resp.headers.get("DAV").replace(" ", "").split(","):
1019 raise ValueError(
1020 f"Server of {resp.geturl()} does not advertise required compliance to webDAV protocol class 1"
1021 )
1023 # The value of 'Server' header is expected to be of the form
1024 # 'dCache/9.2.4' or 'XrootD/v5.7.1'. Not all servers include such a
1025 # header in their response to an OPTIONS request.
1026 details: dict[str, str] = {}
1027 for header in ("Server", "Accept-Ranges"):
1028 value = resp.headers.get(header, None)
1029 if value is not None:
1030 details[header] = value
1032 return details
1034 def _get_response_url(self, resp: HTTPResponse, default_url: str) -> str:
1035 """Return the URL that response `resp` was obtained from.
1037 If `resp` contains no redirection history, return `default_url`.
1038 """
1039 if resp.retries is None:
1040 return default_url
1042 if len(resp.retries.history) == 0:
1043 return default_url
1045 return str(resp.retries.history[-1].redirect_location)
1047 def _rewrite_url_for_frontend(self, url: str) -> str:
1048 """Return a URL to reach one of the frontend servers that serves
1049 requests sent against `url`.
1051 Parameters
1052 ----------
1053 url : `str`
1054 Target URL.
1056 Returns
1057 -------
1058 url: `str`
1059 URL to reach one of this client's frontend servers. If `url` does
1060 not target this client's frontend servers, the returned value
1061 is `url` unmodified.
1062 """
1063 # Do nothing if this URL does not match this client's base URL. This
1064 # happens, for instance, when `url` is a redirection to a backend
1065 # server, so we don't want to rewrite it.
1066 #
1067 # Also, don't rewrite the URL if we don't have frontends configured
1068 # for this client.
1069 if not self._config.frontend_urls or not url.startswith(self._base_url):
1070 return url
1072 # Randomly select one of the configured frontends and return a modified
1073 # URL which uses the selected frontend instead of the original one.
1074 return random.choice(self._config.frontend_urls) + url.removeprefix(self._base_url)
1076 def _request(
1077 self,
1078 method: str,
1079 url: str,
1080 headers: dict[str, str] | None = None,
1081 body: BinaryIO | bytes | str | None = None,
1082 preload_content: bool = True,
1083 redirect: bool = True,
1084 pool_manager: PoolManager | None = None,
1085 **kwargs: dict[Any, Any],
1086 ) -> HTTPResponse:
1087 """Send a generic HTTP request and return the response.
1089 Parameters
1090 ----------
1091 method : `str`
1092 Request method, e.g. 'GET', 'PUT', 'PROPFIND'.
1093 url : `str`
1094 Target URL.
1095 headers : `dict[str, str]`, optional
1096 Headers to sent with the request.
1097 body : `bytes` or `str` or `None`, optional
1098 Request body.
1099 preload_content : `bool`, optional
1100 If True, the response body is downloaded and can be retrieved
1101 via the returned response `.data` property. If False, the
1102 caller needs to call `.read()` on the returned response object to
1103 download the body, either entirely in one call or by chunks.
1104 redirect : `bool`, optional
1105 If True, automatically handle redirects. If False, the returned
1106 response may contain a redirection to another location.
1107 pool_manager : `PoolManager`, optional
1108 Pool manager to use for sending this request. If not provided,
1109 this client's pool manager is used.
1110 kwargs : `dict[Any, Any]`, optional
1111 Keyword arguments to pass unmodified to
1112 `urllib3.PoolManager.request()`.
1114 Returns
1115 -------
1116 resp: `HTTPResponse`
1117 Response to the request as received from the server.
1118 """
1119 # Retrieve the URL we must use to send this request to one of this
1120 # client's configured frontend servers.
1121 url = self._rewrite_url_for_frontend(url)
1123 # If this client is configured not to reuse the network connection
1124 # with the server, add a "Connection: close" header to this request.
1125 #
1126 # However, if the caller has explicitly specified a "Connection"
1127 # header, whatever its value, don't modify it.
1128 headers = {} if headers is None else dict(headers)
1129 if "Connection" not in headers and not self._config.reuse_connection:
1130 headers.update({"Connection": "close"})
1132 # If an authorizer (basic or token) is configured for this client,
1133 # allow it to set the "Authorization" header to this outgoing request.
1134 if self._authorizer is not None:
1135 self._authorizer.set_authorization(headers)
1137 if log.isEnabledFor(logging.DEBUG):
1138 annotation = ""
1139 if method == "GET" and "Range" in headers:
1140 byte_range = headers.get("Range", "").removeprefix("bytes=")
1141 annotation = f" (byte range: {byte_range})"
1143 log.debug("sending request %s %s%s", method, redact_url(url), annotation)
1145 if pool_manager is None:
1146 pool_manager = self._pool_manager
1148 with time_this(
1149 log,
1150 msg="%s %s",
1151 args=(method, url),
1152 mem_usage=self._config.collect_memory_usage,
1153 mem_unit=u.mebibyte,
1154 ):
1155 return pool_manager.request(
1156 method,
1157 url,
1158 body=body,
1159 headers=headers,
1160 preload_content=preload_content,
1161 redirect=redirect,
1162 **kwargs,
1163 )
1165 def _options(
1166 self,
1167 url: str,
1168 headers: dict[str, str] | None = None,
1169 pool_manager: PoolManager | None = None,
1170 ) -> HTTPResponse:
1171 """Send a HTTP OPTIONS request and return the response unmodified.
1173 Parameters
1174 ----------
1175 url : `str`
1176 Target URL.
1177 headers : `dict[str, str]`, optional
1178 Headers to sent with the request.
1179 pool_manager : `PoolManager`, optional
1180 Pool manager to use to send this request.
1182 Returns
1183 -------
1184 resp: `HTTPResponse`
1185 Response to the request as received from the server.
1187 Notes
1188 -----
1189 This method is intended for subclasses to override when needed.
1190 """
1191 return self._request("OPTIONS", url=url, headers=headers, pool_manager=pool_manager)
1193 def _copy(
1194 self,
1195 url: str,
1196 headers: dict[str, str] | None = None,
1197 preload_content: bool = True,
1198 pool_manager: PoolManager | None = None,
1199 ) -> HTTPResponse:
1200 """Send a webDAV COPY request and return the response unmodified.
1202 Parameters
1203 ----------
1204 url : `str`
1205 Target URL.
1206 headers : `dict[str, str]`, optional
1207 Headers to sent with the request.
1208 pool_manager : `PoolManager`, optional
1209 Pool manager to use to send this request.
1211 Notes
1212 -----
1213 This method is intended for subclasses to override when needed.
1214 """
1215 return self._request(
1216 "COPY", url=url, headers=headers, preload_content=preload_content, pool_manager=pool_manager
1217 )
1219 def _delete(
1220 self,
1221 url: str,
1222 headers: dict[str, str] | None = None,
1223 pool_manager: PoolManager | None = None,
1224 ) -> HTTPResponse:
1225 """Send a HTTP DELETE request and return the response unmodified.
1227 Parameters
1228 ----------
1229 url : `str`
1230 Target URL.
1231 headers : `dict[str, str]`, optional
1232 Headers to sent with the request.
1233 pool_manager : `PoolManager`, optional
1234 Pool manager to use to send this request.
1236 Notes
1237 -----
1238 This method is intended for subclasses to override when needed.
1239 """
1240 return self._request("DELETE", url=url, headers=headers, pool_manager=pool_manager)
1242 def _get(
1243 self,
1244 url: str,
1245 headers: dict[str, str] | None = None,
1246 preload_content: bool = True,
1247 redirect: bool = True,
1248 pool_manager: PoolManager | None = None,
1249 ) -> HTTPResponse:
1250 """Send a HTTP GET request and return the response unmodified.
1252 Parameters
1253 ----------
1254 url : `str`
1255 Target URL.
1256 headers : `dict[str, str]`, optional
1257 Headers to sent with the request.
1258 preload_content : `bool`, optional
1259 If True, the response body is downloaded and can be retrieved
1260 via the returned response `.data` property. If False, the
1261 caller needs to call the `.read()` on the returned response
1262 object to download the body.
1263 redirect : `bool`, optional
1264 If True, follow redirections.
1265 pool_manager : `PoolManager`, optional
1266 Pool manager to send the request through.
1268 Returns
1269 -------
1270 resp: `HTTPResponse`
1271 Response to the GET request as received from the server.
1273 Notes
1274 -----
1275 This method is intended for subclasses to override when needed.
1276 """
1277 return self._request(
1278 "GET",
1279 url=url,
1280 headers=headers,
1281 preload_content=preload_content,
1282 redirect=redirect,
1283 pool_manager=pool_manager,
1284 )
1286 def _head(
1287 self,
1288 url: str,
1289 headers: dict[str, str] | None = None,
1290 pool_manager: PoolManager | None = None,
1291 ) -> HTTPResponse:
1292 """Send a HTTP HEAD request and return the response.
1294 Parameters
1295 ----------
1296 url : `str`
1297 Target URL.
1298 headers : `bool`
1299 If the target URL is not found, raise an exception. Otherwise
1300 just return the response.
1301 pool_manager : `PoolManager`, optional
1302 Pool manager to use to send this request.
1304 Notes
1305 -----
1306 This method is intended for subclasses to override when needed.
1307 """
1308 return self._request("HEAD", url=url, headers=headers, pool_manager=pool_manager)
1310 def _mkcol(
1311 self,
1312 url: str,
1313 headers: dict[str, str] | None = None,
1314 pool_manager: PoolManager | None = None,
1315 ) -> HTTPResponse:
1316 """Send a webDAV MKCOL request and return the response unmodified.
1318 Parameters
1319 ----------
1320 url : `str`
1321 Target URL.
1322 headers : `dict[str, str]`, optional
1323 Headers to sent with the request.
1324 pool_manager : `PoolManager`, optional
1325 Pool manager to use to send this request.
1327 Notes
1328 -----
1329 This method is intended for subclasses to override when needed.
1330 """
1331 return self._request("MKCOL", url=url, headers=headers, pool_manager=pool_manager)
1333 def _move(
1334 self,
1335 url: str,
1336 headers: dict[str, str] | None = None,
1337 pool_manager: PoolManager | None = None,
1338 ) -> HTTPResponse:
1339 """Send a webDAV MOVE request and return the response unmodified.
1341 Parameters
1342 ----------
1343 url : `str`
1344 Target URL.
1345 headers : `dict[str, str]`, optional
1346 Headers to sent with the request.
1347 pool_manager : `PoolManager`, optional
1348 Pool manager to use to send this request.
1350 Notes
1351 -----
1352 This method is intended for subclasses to override when needed.
1353 """
1354 return self._request("MOVE", url=url, headers=headers, pool_manager=pool_manager)
1356 def _propfind(
1357 self,
1358 url: str,
1359 headers: dict[str, str] | None = None,
1360 body: str = "",
1361 pool_manager: PoolManager | None = None,
1362 ) -> HTTPResponse:
1363 """Send a webDAV PROPFIND request and return the response unmodified.
1365 Parameters
1366 ----------
1367 url : `str`
1368 Target URL.
1369 headers : `dict[str, str]`, optional
1370 Headers to sent with the request.
1371 body : `str`, optional
1372 Request body.
1373 pool_manager : `PoolManager`, optional
1374 Pool manager to use to send this request.
1376 Notes
1377 -----
1378 This method is intended for subclasses to override when needed.
1379 """
1380 return self._request("PROPFIND", url=url, headers=headers, body=body, pool_manager=pool_manager)
1382 def _put(
1383 self,
1384 url: str,
1385 headers: dict[str, str] | None = None,
1386 body: BinaryIO | bytes = b"",
1387 preload_content: bool = True,
1388 redirect: bool = True,
1389 pool_manager: PoolManager | None = None,
1390 ) -> HTTPResponse:
1391 """Send a HTTP PUT request and return the response unmodified.
1393 Parameters
1394 ----------
1395 url : `str`
1396 Target URL.
1397 headers : `dict[str, str]`, optional
1398 Headers to sent with the request.
1399 body : `BinaryIO` or `bytes`, optional
1400 Request body.
1401 preload_content : `bool`, optional
1402 If True, the response body is downloaded and can be retrieved
1403 via the returned response `.data` property. If False, the
1404 caller needs to call the `.read()` on the returned response
1405 object to download the body.
1406 redirect : `bool`, optional
1407 If True, follow redirections.
1408 pool_manager : `PoolManager`, optional
1409 Pool manager to send the request through.
1411 Returns
1412 -------
1413 resp: `HTTPResponse`
1414 Response to the PUT request as received from the server.
1416 Notes
1417 -----
1418 This method is intended for subclasses to override when needed.
1419 """
1420 # Disable retries when we know the request is not idempotent. In
1421 # particular, when the body of the request is an `io.BufferedReader`,
1422 # any attempt to use that body may totally or partially consume it.
1423 # That means that in case of a retry, the last successful attempt may
1424 # end up uploading an incomplete body and, as a consequence, the
1425 # resulting uploaded file may be either incomplete or have a length
1426 # of zero.
1427 #
1428 # So we only retry a PUT request when the body is an instance of
1429 # `bytes`.
1430 #
1431 # Note that we cannot set `retries` to False since that setting would
1432 # also disable redirection. To disable retries only, we must explicitly
1433 # set the `retries` keyword argument to integer value zero (0).
1434 #
1435 # See documentation:
1436 # https://urllib3.readthedocs.io/en/stable/user-guide.html#retrying-requests
1437 kwargs: dict[Any, Any] = {} if isinstance(body, bytes) else {"retries": 0}
1439 return self._request(
1440 "PUT",
1441 url=url,
1442 headers=headers,
1443 body=body,
1444 preload_content=preload_content,
1445 redirect=redirect,
1446 pool_manager=pool_manager,
1447 **kwargs,
1448 )
1450 def head(
1451 self,
1452 url: str,
1453 headers: dict[str, str] | None = None,
1454 ) -> HTTPResponse:
1455 """Send a HTTP HEAD request, process and return the response
1456 only if successful.
1458 Parameters
1459 ----------
1460 url : `str`
1461 Target URL.
1462 headers : `bool`
1463 If the target URL is not found, raise an exception. Otherwise
1464 just return the response.
1465 """
1466 headers = {} if headers is None else dict(headers)
1467 resp = self._head(url=url, headers=headers)
1468 match resp.status:
1469 case HTTPStatus.OK:
1470 return resp
1471 case HTTPStatus.NOT_FOUND:
1472 raise FileNotFoundError(f"No file found at {resp.geturl()}")
1473 case _:
1474 raise unexpected_status_error("HEAD", url, resp)
1476 def get(
1477 self,
1478 url: str,
1479 headers: dict[str, str] | None = None,
1480 preload_content: bool = True,
1481 redirect: bool = True,
1482 ) -> tuple[str, HTTPResponse]:
1483 """Send a HTTP GET request.
1485 Parameters
1486 ----------
1487 url : `str`
1488 Target URL.
1489 headers : `dict[str, str]`, optional
1490 Headers to sent with the request.
1491 preload_content : `bool`, optional
1492 If True, the response body is downloaded and can be retrieved
1493 via the returned response `.data` property. If False, the
1494 caller needs to call the `.read()` on the returned response
1495 object to download the body.
1496 redirect : `bool`, optional
1497 If True, follow redirections.
1499 Returns
1500 -------
1501 url: `str`
1502 The URL we used to obtain this response. It may be different from
1503 the URL passed as argument in case of redirection.
1504 resp: `HTTPResponse`
1505 Response to the GET request as received from the server.
1506 """
1507 # Send the GET request to the frontend servers.
1508 headers = {} if headers is None else dict(headers)
1509 resp = self._get(
1510 url,
1511 headers=headers,
1512 preload_content=preload_content,
1513 redirect=redirect,
1514 )
1515 match resp.status:
1516 case HTTPStatus.OK | HTTPStatus.PARTIAL_CONTENT:
1517 return self._get_response_url(resp, default_url=url), resp
1518 case HTTPStatus.NOT_FOUND:
1519 raise FileNotFoundError(f"No file found at {resp.geturl()}")
1520 case status if status in resp.REDIRECT_STATUSES and not redirect:
1521 # This response is a redirection but we are asked not to
1522 # follow redirections, so return this response as is.
1523 return self._get_response_url(resp, default_url=url), resp
1524 case _:
1525 raise unexpected_status_error("GET", url, resp)
1527 def options(
1528 self,
1529 url: str,
1530 headers: dict[str, str] | None = None,
1531 ) -> HTTPResponse:
1532 """Send a HTTP OPTIONS request and return the response on success.
1534 Parameters
1535 ----------
1536 url : `str`
1537 Target URL.
1538 headers : `dict` [`str`, `str`], optional
1539 Headers to sent with the request.
1541 Returns
1542 -------
1543 resp: `HTTPResponse`
1544 Response to the request as received from the server.
1545 """
1546 resp = self._options(url=url, headers=headers)
1547 match resp.status:
1548 case HTTPStatus.OK | HTTPStatus.CREATED:
1549 return resp
1550 case _:
1551 raise unexpected_status_error("OPTIONS", url, resp)
1553 def propfind(
1554 self,
1555 url: str,
1556 headers: dict[str, str] | None = None,
1557 body: str = "",
1558 depth: str = "0",
1559 ) -> HTTPResponse:
1560 """Send a HTTP PROPFIND request and return the unmodified response on
1561 success.
1563 Parameters
1564 ----------
1565 url : `str`
1566 Target URL.
1567 headers : `dict[str, str]`, optional
1568 Headers to sent with the request.
1569 body : `str`, optional
1570 Request body.
1571 depth : `str`, optional
1572 ???.
1573 """
1574 headers = {} if headers is None else dict(headers)
1575 headers.update(
1576 {
1577 "Depth": depth,
1578 "Content-Type": 'application/xml; charset="utf-8"',
1579 "Content-Length": str(len(body)),
1580 }
1581 )
1582 resp = self._propfind(url=url, headers=headers, body=body)
1583 match resp.status:
1584 case HTTPStatus.MULTI_STATUS | HTTPStatus.NOT_FOUND:
1585 return resp
1586 case _:
1587 raise unexpected_status_error("PROPFIND", url, resp)
1589 def put(
1590 self,
1591 url: str,
1592 headers: dict[str, str] | None = None,
1593 data: BinaryIO | bytes = b"",
1594 ) -> int | None:
1595 """Send a HTTP PUT request.
1597 Parameters
1598 ----------
1599 url : `str`
1600 Target URL.
1601 headers : `dict[str, str]`, optional
1602 Headers to sent with the request.
1603 data : `BinaryIO` or `bytes`
1604 Request body.
1606 Returns
1607 -------
1608 size : `int | None`
1609 The size in bytes of the file uploaded. Can be `None` if the size
1610 could not be retrieved.
1611 """
1612 # Send a PUT request with empty body and handle redirection. This
1613 # is useful if the server redirects us; since we cannot rewind the
1614 # data we are uploading, we don't start uploading data until we
1615 # connect to the server that will actually serve our request.
1616 frontend_headers = {} if headers is None else dict(headers)
1617 frontend_headers.update({"Content-Length": "0"})
1618 resp = self._put(url, headers=frontend_headers, body=b"", redirect=False)
1619 match resp.status:
1620 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT:
1621 redirect_url = url
1622 case status if status in resp.REDIRECT_STATUSES:
1623 redirect_url = resp.headers.get("Location")
1624 case _:
1625 raise unexpected_status_error("PUT", url, resp)
1627 # We may have been redirectred. Upload the file contents to
1628 # its final destination.
1630 # Ask the server to compute and record a checksum of the uploaded
1631 # file contents, for later integrity checks. Since we don't compute
1632 # the digest ourselves while uploading the data, we cannot control
1633 # after the request is complete that the data we uploaded is
1634 # identical to the data recorded by the server, but at least the
1635 # server has recorded a digest of the data it stored.
1636 #
1637 # See RFC-3230 for details and
1638 # https://www.iana.org/assignments/http-dig-alg/http-dig-alg.xhtml
1639 # for the list of supported digest algorithhms.
1640 #
1641 # In addition, note that not all servers implement this RFC so
1642 # the checksum reqquest may be ignored by the server.
1643 backend_headers = {} if headers is None else dict(headers)
1644 if (checksum := self._config.request_checksum) is not None:
1645 backend_headers.update({"Want-Digest": checksum})
1647 resp = self._put(redirect_url, body=data, headers=backend_headers)
1648 match resp.status:
1649 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT:
1650 # Send a HEAD request to retrieve the size of the file we
1651 # just uploaded
1652 resp = self.head(redirect_url)
1653 size = int(resp.headers.get("Content-Length", -1))
1654 return None if size == -1 else size
1655 case _:
1656 raise unexpected_status_error("PUT", redirect_url, resp)
1658 def _get_temporary_basename(self, basename: str, prefix: str) -> str:
1659 """Return a basename for a temporary file."""
1660 unique_id = str(uuid.uuid4())
1661 return f"{prefix}.{unique_id}.{basename}"
1663 def _split_parent_and_basename(self, url: str) -> tuple[str, str]:
1664 """Return the URL of the parent directory and the basename from
1665 `url`.
1666 """
1667 parsed: Url = parse_url(url)
1668 normalized_path = normalize_path(parsed.path)
1669 parent_path = posixpath.dirname(normalized_path)
1670 basename = posixpath.basename(normalized_path)
1671 parent_url = Url(
1672 scheme=parsed.scheme,
1673 auth=parsed.auth,
1674 host=parsed.host,
1675 port=parsed.port,
1676 path=parent_path,
1677 query=parsed.query,
1678 fragment=parsed.fragment,
1679 ).url
1680 return parent_url, basename
1682 def _parent(self, url: str) -> str:
1683 """Return the URL of the parent directory to `url`."""
1684 parent_url, _ = self._split_parent_and_basename(url)
1685 return parent_url
1687 def _make_temporary_url(self, url: str, prefix: str = ".tmp") -> str:
1688 """Return the URL of a temporary file based on `url`."""
1689 parent_url, basename = self._split_parent_and_basename(url)
1690 temporary_basename = self._get_temporary_basename(basename=basename, prefix=prefix)
1691 return f"{parent_url}/{temporary_basename}"
1693 def exists(self, url: str) -> bool:
1694 """Return True if a file or directory exists at `url`.
1696 Parameters
1697 ----------
1698 url : `str`
1699 Target URL.
1701 Returns
1702 -------
1703 result: `bool`
1704 True if there is an object at `url`.
1705 """
1706 return self.stat(url).exists
1708 def size(self, url: str) -> int:
1709 """Return the size in bytes of resource at `url`.
1711 If `url` designates a directory, the size is zero.
1713 Parameters
1714 ----------
1715 url : `str`
1716 Target URL.
1718 Returns
1719 -------
1720 size: `int`
1721 The number of bytes of the resource located at `url`.
1722 """
1723 # Check if we have the size of this URL in our cache
1724 if (size := self._file_size_cache.get_size(url)) is not None:
1725 return size
1727 stat = self.stat(url)
1728 if not stat.exists:
1729 raise FileNotFoundError(f"No file or directory found at {url}")
1730 else:
1731 return stat.size
1733 def is_dir(self, url: str) -> bool:
1734 """Return True if a directory exists at `url`.
1736 Parameters
1737 ----------
1738 url : `str`
1739 Target URL.
1741 Returns
1742 -------
1743 result: `bool`
1744 True if there is a directory at `url`.
1745 """
1746 return self.stat(url).is_dir
1748 def mkcol(self, url: str) -> None:
1749 """Create a directory at `url`.
1751 If a directory already exists at `url` no error is returned nor
1752 exception is raised. An exception is raised if a file exists at `url`.
1754 Parameters
1755 ----------
1756 url : `str`
1757 Target URL.
1758 """
1759 resp = self._mkcol(url=url)
1760 match resp.status:
1761 case HTTPStatus.CREATED | HTTPStatus.METHOD_NOT_ALLOWED:
1762 return
1763 case HTTPStatus.CONFLICT:
1764 # The parent directory does not exist. Create it first except
1765 # if the parent's path is "/".
1766 parent = self._parent(url)
1767 if not parent.endswith("/"):
1768 self.mkcol(parent)
1769 resp = self._mkcol(url=url)
1770 case _:
1771 raise ValueError(
1772 f"Can not create directory {resp.geturl()}: status {resp.status} {resp.reason}"
1773 )
1775 def stat(self, url: str) -> DavFileMetadata:
1776 """Return some properties of file or directory located at `url`.
1778 Parameters
1779 ----------
1780 url : `str`
1781 Target URL.
1783 Returns
1784 -------
1785 result: `DavResourceMetadata`
1786 Details of the resources at `url`. If no resource was found at
1787 that URL no exception is raised. Instead the returned details allow
1788 for detecting that the resource does not exist.
1790 The returned value should include fields to determine
1791 if there is a file or a directory at that `url` and if so, its
1792 size and kind (file or directory). Other fields may also be
1793 included depending on the implementation of the webDAV protocol
1794 by the server.
1795 """
1796 # Request the minimum set of DAV properties.
1797 body = (
1798 """<?xml version="1.0" encoding="utf-8"?>"""
1799 """<D:propfind xmlns:D="DAV:">"""
1800 """<D:prop>"""
1801 """<D:resourcetype/>"""
1802 """<D:getcontentlength/>"""
1803 """<D:getlastmodified/>"""
1804 """</D:prop>"""
1805 """</D:propfind>"""
1806 )
1807 resp = self.propfind(url, body=body, depth="0")
1808 match resp.status:
1809 case HTTPStatus.NOT_FOUND:
1810 href = url.replace(self._base_url, "", 1)
1811 return DavFileMetadata(base_url=self._base_url, href=href)
1812 case HTTPStatus.MULTI_STATUS:
1813 property = self._propfind_parser.parse(resp.data)[0]
1814 return DavFileMetadata.from_property(base_url=self._base_url, property=property)
1815 case _:
1816 raise unexpected_status_error("PROPFIND", url, resp)
1818 def info(self, url: str, name: str | None = None) -> dict[str, Any]:
1819 """Return the details about the file or directory at `url`.
1821 Parameters
1822 ----------
1823 url : `str`
1824 Target URL.
1825 name : `str`
1826 Name of the object to be included in the returned value. If None,
1827 the `url` is used as name.
1829 Returns
1830 -------
1831 result: `dict`
1832 For an existing file, the returned value has the form:
1834 .. code-block:: json
1836 {
1837 "name": name,
1838 "size": 1234,
1839 "type": "file",
1840 "last_modified":
1841 datetime.datetime(2025, 4, 10, 15, 12, 51, 227854),
1842 "checksums": {
1843 "adler32": "0fc5f83f",
1844 "md5": "1f57339acdec099c6c0a41f8e3d5fcd0",
1845 }
1846 }
1848 For an existing directory, the returned value has the form:
1850 .. code-block:: json
1852 {
1853 "name": name,
1854 "size": 0,
1855 "type": "directory",
1856 "last_modified":
1857 datetime.datetime(2025, 4, 10, 15, 12, 51, 227854),
1858 "checksums": {},
1859 }
1861 For a non-existing file or directory, the returned value has the
1862 form:
1864 .. code-block:: json
1866 {
1867 "name": name,
1868 "size": None,
1869 "type": None,
1870 "last_modified": datetime.datetime(1, 1, 1, 0, 0),
1871 "checksums": {},
1872 }
1874 Notes
1875 -----
1876 The format of the returned directory is inspired and compatible with
1877 `fsspec`.
1879 The size of existing directories is always zero. The `checksums`
1880 dictionary is empty for directories and may be empty for files if the
1881 server does not compute and store the checksum of the files it stores.
1882 """
1883 result: dict[str, Any] = {
1884 "name": name if name is not None else url,
1885 "type": None,
1886 "size": None,
1887 "last_modified": datetime.min,
1888 "checksums": {},
1889 }
1890 metadata = self.stat(url)
1891 if not metadata.exists:
1892 return result
1894 result.update(
1895 {
1896 "type": "directory" if metadata.is_dir else "file",
1897 "size": metadata.size,
1898 "last_modified": metadata.last_modified,
1899 "checksums": metadata.checksums,
1900 }
1901 )
1902 return result
1904 def move(self, source_url: str, destination_url: str, overwrite: bool = False) -> HTTPResponse:
1905 """Send a webDAV MOVE request and return the response unmodified.
1907 Parameters
1908 ----------
1909 source_url : `str`
1910 Source URL.
1911 destination_url : `str`
1912 Destination URL.
1913 overwrite : `bool`, optional
1914 Overwrite the destination if it exists.
1916 Returns
1917 -------
1918 resp : `HTTPResponse`
1919 The unmodified response received from the server.
1920 """
1921 headers = {
1922 "Destination": destination_url,
1923 "Overwrite": "T" if overwrite else "F",
1924 }
1925 return self._move(source_url, headers=headers)
1927 def read_dir(self, url: str) -> list[DavFileMetadata]:
1928 """Return the properties of the files or directories contained in
1929 directory located at `url`.
1931 If `url` designates a file, only the details of itself are returned.
1933 Parameters
1934 ----------
1935 url : `str`
1936 Target URL.
1938 Returns
1939 -------
1940 result: `list[DavResourceMetadata]`
1941 List of details of each file or directory within `url`.
1942 """
1943 body = (
1944 """<?xml version="1.0" encoding="utf-8"?>"""
1945 """<D:propfind xmlns:D="DAV:"><D:prop>"""
1946 """<D:resourcetype/><D:getcontentlength/><D:getlastmodified/><D:displayname/>"""
1947 """</D:prop></D:propfind>"""
1948 )
1949 resp = self.propfind(url, body=body, depth="1")
1950 match resp.status:
1951 case HTTPStatus.MULTI_STATUS:
1952 pass
1953 case HTTPStatus.NOT_FOUND:
1954 raise FileNotFoundError(f"No directory found at {resp.geturl()}")
1955 case _:
1956 raise unexpected_status_error("PROPFIND", url, resp)
1958 if (path := parse_url(url).path) is not None:
1959 this_dir_href = path.rstrip("/") + "/"
1960 else:
1961 this_dir_href = "/"
1963 result = []
1964 for property in self._propfind_parser.parse(resp.data):
1965 # Don't include in the results the metadata of the directory we
1966 # traversing.
1967 # Some webDAV servers do not append a "/" to the href of a
1968 # directory in their response to PROPFIND, so we must take into
1969 # account that.
1970 if property.is_file:
1971 result.append(DavFileMetadata.from_property(base_url=self._base_url, property=property))
1972 elif property.is_dir and property.href != this_dir_href:
1973 result.append(DavFileMetadata.from_property(base_url=self._base_url, property=property))
1975 return result
1977 def read(self, url: str) -> tuple[str, bytes]:
1978 """Download the contents of file located at `url`.
1980 Parameters
1981 ----------
1982 url : `str`
1983 Target URL.
1985 Returns
1986 -------
1987 url: `str`
1988 Backend URL from which the data was obtained.
1989 data: `bytes`
1990 Contents of the file.
1992 Notes
1993 -----
1994 The caller must ensure that the resource at `url` is a file, not
1995 a directory.
1996 """
1997 backend_url, resp = self.get(url)
1998 return backend_url, resp.data
2000 def read_range(
2001 self,
2002 url: str,
2003 start: int,
2004 end: int | None,
2005 headers: dict[str, str] | None = None,
2006 ) -> tuple[str, bytes]:
2007 """Download partial content of file located at `url`.
2009 Parameters
2010 ----------
2011 url : `str`
2012 Target URL.
2013 start : `int`
2014 Starting byte offset of the range to download.
2015 end : `int`, optional
2016 Ending byte offset of the range to download.
2017 headers : `dict[str,str]`, optional
2018 Specific headers to sent with the GET request.
2020 Returns
2021 -------
2022 backend_url: `str`
2023 URL used to retrieve this data. If the server redirected us
2024 this is the URL we were redirected to.
2025 data: `bytes`
2026 Partial contents of the file.
2028 Notes
2029 -----
2030 The caller must ensure that the resource at `url` is a file, not
2031 a directory. This is important because some webDAV servers respond
2032 with an HTML document when asked for reading a directory.
2033 """
2034 get_headers = {} if headers is None else dict(headers)
2035 get_headers.update({"Accept-Encoding": "identity"})
2036 if end is None:
2037 get_headers.update({"Range": f"bytes={start}-"})
2038 else:
2039 get_headers.update({"Range": f"bytes={start}-{end}"})
2041 final_url, resp = self.get(url, headers=get_headers, redirect=True)
2042 match resp.status:
2043 case HTTPStatus.PARTIAL_CONTENT:
2044 return final_url, resp.data
2045 case _:
2046 raise unexpected_status_error("GET (with 'Range' header)", url, resp)
2048 def _close(self, url: str) -> None:
2049 """Signal the server hosting `url` that no more `read_range` requests
2050 will be issued, so that the resources allocated for serving those
2051 requests can be released.
2053 This helper method is intended to be overwritten by subclasses that
2054 need such functionality according to the behavior of the specific
2055 remote storage endpoint.
2057 Parameters
2058 ----------
2059 url : `str`
2060 Target URL.
2061 """
2062 # This is a NOP for a generic webDAV server.
2063 pass
2065 def _write_response_body_to_file(self, resp: HTTPResponse, filename: str, chunk_size: int) -> int:
2066 """Write the response body to a local file.
2068 Parameters
2069 ----------
2070 resp : `HTTPResponse`
2071 The HTTP Response to read the body from.
2072 filename : `str`
2073 Local file to write the content to. If the file already exists,
2074 it will be rewritten.
2075 chunk_size : `int`
2076 Size of the chunks to write to `filename`.
2078 Returns
2079 -------
2080 count: `int`
2081 Number of bytes written to `filename`.
2082 """
2083 try:
2084 # Read the response body into a pre-allocated memory buffer and
2085 # write the buffer content to the destination file avoiding
2086 # copies if possible.
2087 content_length = 0
2088 with open(filename, "wb", buffering=0) as file:
2089 view = memoryview(bytearray(chunk_size))
2090 while True:
2091 if (count := resp.readinto(view)) > 0: # type: ignore
2092 content_length += count
2093 file.write(view[:count])
2094 else:
2095 break
2097 # Check that the expected and actual content lengths match.
2098 # Perform this check only when the body of the response was not
2099 # encoded by the server.
2100 expected_length: int = int(resp.headers.get("Content-Length", -1))
2101 if (
2102 "Content-Encoding" not in resp.headers
2103 and expected_length != -1
2104 and expected_length != content_length
2105 ):
2106 raise ValueError(
2107 f"Size of downloaded file does not match value in Content-Length header for "
2108 f"{resp.geturl()}: expecting {expected_length} and got {content_length} bytes"
2109 )
2111 return content_length
2112 finally:
2113 # Release the connection
2114 resp.drain_conn()
2115 resp.release_conn()
2117 def download(self, url: str, filename: str, chunk_size: int) -> int:
2118 """Download the content of a file and write it to local file.
2120 Parameters
2121 ----------
2122 url : `str`
2123 Target URL.
2124 filename : `str`
2125 Local file to write the content to. If the file already exists,
2126 it will be rewritten.
2127 chunk_size : `int`
2128 Size of the chunks to write to `filename`.
2130 Returns
2131 -------
2132 count: `int`
2133 Number of bytes written to `filename`.
2135 Notes
2136 -----
2137 The caller must ensure that the resource at `url` is a file, not
2138 a directory.
2139 """
2140 _, resp = self.get(url, preload_content=False)
2141 return self._write_response_body_to_file(resp, filename, chunk_size)
2143 def write(self, url: str, data: BinaryIO | bytes) -> int | None:
2144 """Create or rewrite a remote file at `url` with `data` as its
2145 contents.
2147 Parameters
2148 ----------
2149 url : `str`
2150 Target URL.
2151 data : `bytes`
2152 Sequence of bytes to upload.
2154 Returns
2155 -------
2156 size : `int | None`
2157 The size in bytes of the file uploaded. Can be `None` if the size
2158 could not be retrieved.
2160 Notes
2161 -----
2162 If a file already exists at `url` it will be rewritten.
2163 """
2164 # According to RFC 4918, the parent directory of the file must
2165 # exist before we can write to it. So create it first and then
2166 # upload.
2167 self.mkcol(self._parent(url))
2169 try:
2170 # Upload to a temporary file and rename to the final name.
2171 temporary_url = self._make_temporary_url(url)
2172 size = self.put(temporary_url, data=data)
2173 self.rename(temporary_url, url, overwrite=True, create_parent=False)
2175 # Update the file size cache with this size
2176 self._file_size_cache.update_size(url, size)
2177 return size
2178 except Exception:
2179 # Upload failed. Attempt to remove the temporary file.
2180 self.delete(temporary_url)
2181 raise
2183 def checksums(self, url: str) -> dict[str, str]:
2184 """Return the checksums of the contents of file located at `url`.
2186 The checksums are retrieved from the storage endpoint. There may be
2187 none if the storage endpoint does not automatically expose the
2188 checksums it computes.
2190 Parameters
2191 ----------
2192 url : `str`
2193 Target URL.
2195 Returns
2196 -------
2197 checksums: `dict[str, str]`
2198 A file exists at `url`.
2199 The key of the dictionary is the lowercased name of the checksum
2200 algorithm (e.g. "md5", "adler32"). The value is the lowercased
2201 checksum itself (e.g. "78441cec2479ec8b545c4d6699f542da").
2202 """
2203 stat = self.stat(url)
2204 if not stat.exists:
2205 raise FileNotFoundError(f"No file found at {url}")
2207 return stat.checksums if stat.is_file else {}
2209 def delete(self, url: str) -> None:
2210 """Delete the file or directory at `url`.
2212 If there is no file or directory at `url` is not considered an error.
2214 Parameters
2215 ----------
2216 url : `str`
2217 Target URL.
2219 Notes
2220 -----
2221 If `url` designates a directory, some webDAV servers recursively
2222 remove the directory and its contents. Others, only remove the
2223 directory if it is empty.
2225 For a consisten behavior, the caller must check what kind of object
2226 the target URL is and walk the hierarchy removing all objects.
2227 """
2228 resp = self._delete(url)
2229 match resp.status:
2230 case HTTPStatus.OK | HTTPStatus.ACCEPTED | HTTPStatus.NO_CONTENT | HTTPStatus.NOT_FOUND:
2231 # Invalidate the entry for this file in our cache, if any
2232 self._file_size_cache.invalidate(url)
2233 case _:
2234 raise ValueError(
2235 f"Unable to delete resource {resp.geturl()}: status {resp.status} {resp.reason}"
2236 )
2238 def accepts_ranges(self, url: str) -> bool:
2239 """Return `True` if the server supports a 'Range' header in
2240 GET requests against `url`.
2242 Parameters
2243 ----------
2244 url : `str`
2245 Target URL.
2246 """
2247 # If we have already determined that the server accepts "Range" for
2248 # another URL, we assume that it implements that feature for any
2249 # file it serves, so reuse that information.
2250 if self._accepts_ranges is not None:
2251 return self._accepts_ranges
2253 with self._lock:
2254 if self._accepts_ranges is None:
2255 self._accepts_ranges = self.head(url).headers.get("Accept-Ranges", "") == "bytes"
2257 return self._accepts_ranges
2259 @property
2260 def supports_duplicate(self) -> bool:
2261 """Return True if the server this client interacts with implements
2262 webDAV COPY method.
2263 """
2264 return self._can_duplicate
2266 def copy(self, source_url: str, destination_url: str, overwrite: bool = False) -> None:
2267 """Copy the file at `source_url` to `destination_url` in the same
2268 storage endpoint.
2270 Parameters
2271 ----------
2272 source_url : `str`
2273 URL of the source file.
2274 destination_url : `str`
2275 URL of the destination file. Its parent directory must exist.
2276 overwrite : `bool`
2277 If True and a file exists at `destination_url` it will be
2278 overwritten. Otherwise an exception is raised.
2279 """
2280 headers = {
2281 "Destination": destination_url,
2282 "Overwrite": "T" if overwrite else "F",
2283 }
2284 resp = self._copy(source_url, headers=headers)
2285 match resp.status:
2286 case HTTPStatus.CREATED | HTTPStatus.NO_CONTENT:
2287 self._file_size_cache.invalidate(destination_url)
2288 case _:
2289 raise ValueError(
2290 f"Could not copy {resp.geturl()} to {destination_url}: status {resp.status} {resp.reason}"
2291 )
2293 def duplicate(self, source_url: str, destination_url: str, overwrite: bool = False) -> None:
2294 """Copy the file at `source_url` to `destination_url` in the same
2295 storage endpoint.
2297 Parameters
2298 ----------
2299 source_url : `str`
2300 URL of the source file.
2301 destination_url : `str`
2302 URL of the destination file. Its parent directory is created if
2303 necessary.
2304 overwrite : `bool`
2305 If True and a file exists at `destination_url` it will be
2306 overwritten. Otherwise an exception is raised.
2307 """
2308 # Check the source is a file
2309 if self.is_dir(source_url):
2310 raise NotImplementedError(f"copy is not implemented for directory {source_url}")
2312 # Create the destination's parent directory first because COPY may
2313 # fail if it does not exist, depending on the server implementation
2314 # of RFC 4918.
2315 destination_parent = self._parent(destination_url)
2316 self.mkcol(destination_parent)
2317 self.copy(source_url=source_url, destination_url=destination_url, overwrite=overwrite)
2319 def rename(
2320 self,
2321 source_url: str,
2322 destination_url: str,
2323 overwrite: bool = False,
2324 create_parent: bool = True,
2325 ) -> None:
2326 """Rename (move) the file at `source_url` to `destination_url` in the
2327 same storage endpoint.
2329 Parameters
2330 ----------
2331 source_url : `str`
2332 URL of the source file.
2333 destination_url : `str`
2334 URL of the destination file. Its parent directory must exist.
2335 overwrite : `bool`, optional
2336 If True and a file exists at `destination_url` it will be
2337 overwritten. Otherwise an exception is raised.
2338 create_parent : `bool`, optional
2339 Whether to create the parent.
2340 """
2341 # Create the destination's parent directory first because MOVE may
2342 # fail if it does not exist, depending on the server implementation
2343 # of RFC 4918.
2344 if create_parent:
2345 destination_parent = self._parent(destination_url)
2346 self.mkcol(destination_parent)
2348 resp = self.move(source_url=source_url, destination_url=destination_url, overwrite=overwrite)
2349 match resp.status:
2350 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT:
2351 self._file_size_cache.invalidate(destination_url)
2352 case _:
2353 raise ValueError(
2354 f"""Could not move file {resp.geturl()} to {destination_url}: status {resp.status} """
2355 f"""{resp.reason}"""
2356 )
2358 def generate_presigned_get_url(self, url: str, expiration_time_seconds: int) -> str:
2359 """Return a pre-signed URL that can be used to retrieve this resource
2360 using an HTTP GET without supplying any access credentials.
2362 Parameters
2363 ----------
2364 url : `str`
2365 Target URL.
2366 expiration_time_seconds : `int`
2367 Number of seconds until the generated URL is no longer valid.
2369 Returns
2370 -------
2371 url : `str`
2372 HTTP URL signed for GET.
2373 """
2374 raise NotImplementedError(f"URL signing is not supported by server for {self}")
2376 def generate_presigned_put_url(self, url: str, expiration_time_seconds: int) -> str:
2377 """Return a pre-signed URL that can be used to upload a file to this
2378 path using an HTTP PUT without supplying any access credentials.
2380 Parameters
2381 ----------
2382 url : `str`
2383 Target URL.
2384 expiration_time_seconds : `int`
2385 Number of seconds until the generated URL is no longer valid.
2387 Returns
2388 -------
2389 url : `str`
2390 HTTP URL signed for PUT.
2391 """
2392 raise NotImplementedError(f"URL signing is not supported by server for {self}")
2395class ActivityCaveat(enum.Enum):
2396 """Helper class for enumerating accepted activity caveats for requesting
2397 macaroons for dCache or XRootD webDAV servers.
2398 """
2400 DOWNLOAD = 1
2401 UPLOAD = 2
2404class DavClientURLSigner(DavClient):
2405 """WebDAV client which supports signing of URL for upload and download.
2407 Instances of this class are thread-safe.
2409 Parameters
2410 ----------
2411 url : `str`
2412 Root URL of the storage endpoint
2413 (e.g. "https://host.example.org:1234/").
2414 config : `DavConfig`
2415 Configuration to initialize this client.
2416 accepts_ranges : `bool` | `None`
2417 Indicate whether the remote server accepts the ``Range`` header in GET
2418 requests.
2419 """
2421 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None:
2422 super().__init__(url=url, config=config, accepts_ranges=accepts_ranges)
2424 def generate_presigned_get_url(self, url: str, expiration_time_seconds: int) -> str:
2425 """Return a pre-signed URL that can be used to retrieve the resource
2426 at `url` using an HTTP GET without supplying any access credentials.
2428 Parameters
2429 ----------
2430 url : `str`
2431 URL of an existing file.
2432 expiration_time_seconds : `int`
2433 Number of seconds until the generated URL is no longer valid.
2435 Returns
2436 -------
2437 url : `str`
2438 HTTP URL signed for GET.
2440 Notes
2441 -----
2442 Although the returned URL allows for downloading the file at `url`
2443 without supplying credentials, the HTTP client must be configured
2444 to accept the certificate the server will present if the client wants
2445 validate it. The server's certificate may be issued by a certificate
2446 authority unknown to the client.
2447 """
2448 macaroon: str = self._get_macaroon(url, ActivityCaveat.DOWNLOAD, expiration_time_seconds)
2449 return f"{url}?authz={macaroon}"
2451 def generate_presigned_put_url(self, url: str, expiration_time_seconds: int) -> str:
2452 """Return a pre-signed URL that can be used to upload a file to `url`
2453 using an HTTP PUT without supplying any access credentials.
2455 Parameters
2456 ----------
2457 url : `str`
2458 URL of an existing file.
2459 expiration_time_seconds : `int`
2460 Number of seconds until the generated URL is no longer valid.
2462 Returns
2463 -------
2464 url : `str`
2465 HTTP URL signed for PUT.
2467 Notes
2468 -----
2469 Although the returned URL allows for uploading a file to `url`
2470 without supplying credentials, the HTTP client must be configured
2471 to accept the certificate the server will present if the client wants
2472 validate it. The server's certificate may be issued by a certificate
2473 authority unknown to the client.
2474 """
2475 macaroon: str = self._get_macaroon(url, ActivityCaveat.UPLOAD, expiration_time_seconds)
2476 return f"{url}?authz={macaroon}"
2478 def _get_macaroon(self, url: str, activity: ActivityCaveat, expiration_time_seconds: int) -> str:
2479 """Return a macaroon for uploading or downloading the file at `url`.
2481 Parameters
2482 ----------
2483 url : `str`
2484 URL of an existing file.
2485 activity : `ActivityCaveat`
2486 the activity the macaroon is requested for.
2487 expiration_time_seconds : `int`
2488 Requested duration of the macaroon, in seconds.
2490 Returns
2491 -------
2492 macaroon : `str`
2493 Macaroon to be used with `url` in a GET or PUT request.
2494 """
2495 # dCache and XRootD webDAV servers support delivery of macaroons.
2496 #
2497 # For details about dCache macaroons see:
2498 # https://www.dcache.org/manuals/UserGuide-9.2/macaroons.shtml
2499 match activity:
2500 case ActivityCaveat.DOWNLOAD:
2501 activity_caveat = "DOWNLOAD,LIST"
2502 case ActivityCaveat.UPLOAD:
2503 activity_caveat = "UPLOAD,LIST,DELETE,MANAGE"
2505 # Retrieve a macaroon for the requested activities and duration
2506 headers = {"Content-Type": "application/macaroon-request"}
2507 body = {
2508 "caveats": [
2509 f"activity:{activity_caveat}",
2510 ],
2511 "validity": f"PT{expiration_time_seconds}S",
2512 }
2513 resp = self._request("POST", url, headers=headers, body=json.dumps(body))
2514 if resp.status != HTTPStatus.OK:
2515 raise ValueError(
2516 f"Could not retrieve a macaroon for URL {resp.geturl()}, status: {resp.status} {resp.reason}"
2517 )
2519 # We are expecting the body of the response to be formatted in JSON.
2520 # dCache sets the 'Content-Type' of the response to 'application/json'
2521 # but XRootD does not set any 'Content-Type' header 8-[
2522 #
2523 # An example of a response body returned by dCache is shown below:
2524 # {
2525 # "macaroon": "MDA[...]Qo",
2526 # "uri": {
2527 # "targetWithMacaroon": "https://dcache.example.org/?authz=MD...",
2528 # "baseWithMacaroon": "https://dcache.example.org/?authz=MD...",
2529 # "target": "https://dcache.example.org/",
2530 # "base": "https://dcache.example.org/"
2531 # }
2532 # }
2533 #
2534 # An example of a response body returned by XRootD is shown below:
2535 # {
2536 # "macaroon": "MDA[...]Qo",
2537 # "expires_in": 86400
2538 # }
2539 try:
2540 response_body = json.loads(resp.data.decode())
2541 except json.JSONDecodeError:
2542 raise ValueError(f"Could not deserialize response to POST request for URL {resp.geturl()}")
2544 if "macaroon" in response_body:
2545 return response_body["macaroon"]
2547 raise ValueError(f"Could not retrieve macaroon for URL {resp.geturl()}")
2549 @override
2550 def duplicate(self, source_url: str, destination_url: str, overwrite: bool = False) -> None:
2551 """Copy the file at `source_url` to `destination_url` in the same
2552 storage endpoint.
2554 Parameters
2555 ----------
2556 source_url : `str`
2557 URL of the source file.
2558 destination_url : `str`
2559 URL of the destination file. Its parent directory must exist.
2560 overwrite : `bool`
2561 If True and a file exists at `destination_url` it will be
2562 overwritten. Otherwise an exception is raised.
2563 """
2564 # Check the source is a file
2565 if self.is_dir(source_url):
2566 raise NotImplementedError(f"copy is not implemented for directory {source_url}")
2568 # Neither dCache nor XrootD currently implement the COPY
2569 # webDAV method as documented in
2570 #
2571 # http://www.webdav.org/specs/rfc4918.html#METHOD_COPY
2572 #
2573 # (See issues DM-37603 and DM-37651 for details)
2574 # With those servers use third-party copy instead.
2575 return self._copy_via_third_party(source_url, destination_url, overwrite)
2577 def _copy_via_third_party(self, source_url: str, destination_url: str, overwrite: bool = False) -> None:
2578 """Copy the file at `source_url` to `destination_url` in the same
2579 storage endpoint using the third-party copy functionality
2580 implemented by dCache and XRootD servers.
2582 Parameters
2583 ----------
2584 source_url : `str`
2585 URL of the source file.
2586 destination_url : `str`
2587 URL of the destination file. Its parent directory must exist.
2588 overwrite : `bool`
2589 If True and a file exists at `destination_url` it will be
2590 overwritten. Otherwise an exception is raised.
2591 """
2592 # To implement COPY we use dCache's third-party copy mechanism
2593 # documented at:
2594 #
2595 # https://www.dcache.org/manuals/UserGuide-10.2/webdav.shtml#third-party-transfers
2596 #
2597 # The reason is that dCache does not correctly implement webDAV's COPY
2598 # method. See https://github.com/dCache/dcache/issues/6950
2600 # Create the destination's parent directory first because COPY may
2601 # fail if it does not exist, depending on the server implementation
2602 # of RFC 4918.
2603 destination_parent = self._parent(destination_url)
2604 self.mkcol(destination_parent)
2606 # Retrieve a macaroon for downloading the source
2607 download_macaroon = self._get_macaroon(source_url, ActivityCaveat.DOWNLOAD, 300)
2609 # Prepare and send the COPY request
2610 try:
2611 headers = {
2612 "Source": source_url,
2613 "TransferHeaderAuthorization": f"Bearer {download_macaroon}",
2614 "Credential": "none",
2615 "Depth": "0",
2616 "Overwrite": "T" if overwrite else "F",
2617 "RequireChecksumVerification": "false",
2618 }
2619 resp = self._copy(destination_url, headers=headers, preload_content=False)
2620 match resp.status:
2621 case HTTPStatus.CREATED:
2622 return
2623 case HTTPStatus.ACCEPTED:
2624 pass
2625 case _:
2626 raise ValueError(
2627 f"Unable to copy resource {resp.geturl()}; status: {resp.status} {resp.reason}"
2628 )
2630 # Analyse the response to the COPY request that the server has
2631 # not completed yet.
2632 content_type = resp.headers.get("Content-Type")
2633 if content_type != "text/perf-marker-stream":
2634 raise ValueError(
2635 f"""Unexpected Content-Type {content_type} in response to COPY request from """
2636 f"""{source_url} to {destination_url}"""
2637 )
2639 # Read the performance markers in the response body until we get
2640 # a "success" or "failure" notification.
2641 #
2642 # Documentation:
2643 # https://dcache.org/manuals/UserGuide-10.2/webdav.shtml#third-party-transfers
2644 for marker in io.TextIOWrapper(resp): # type: ignore
2645 marker = marker.rstrip("\n")
2646 if marker == "": # EOF
2647 raise ValueError(
2648 f"""Copying file from {source_url} to {destination_url} failed: """
2649 """could not get response from server"""
2650 )
2651 elif marker.startswith("failure:"):
2652 raise ValueError(
2653 f"""Copying file from {source_url} to {destination_url} failed with error: """
2654 f"""{marker}"""
2655 )
2656 elif marker.startswith("success:"):
2657 return
2658 finally:
2659 resp.drain_conn()
2662class DavClientDCache(DavClientURLSigner):
2663 """Client for interacting with a dCache webDAV server.
2665 Instances of this class are thread-safe.
2667 Parameters
2668 ----------
2669 url : `str`
2670 Root URL of the storage endpoint
2671 (e.g. "https://host.example.org:1234/").
2672 config : `DavConfig`
2673 Configuration to initialize this client.
2674 accepts_ranges : `bool` | `None`
2675 Indicate whether the remote server accepts the ``Range`` header in GET
2676 requests.
2677 """
2679 # Regular expression to parse dCache's response body of a successful
2680 # PUT request. Such a response body is of the form:
2681 #
2682 # "104857600 bytes uploaded\r\n\r\n"
2683 #
2684 rex: re.Pattern = re.compile(r"^(\d*) bytes uploaded", re.IGNORECASE | re.ASCII)
2686 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None:
2687 super().__init__(url=url, config=config, accepts_ranges=accepts_ranges)
2689 # Create a specialized pool manager for sending requests to dCache
2690 # webdav door, in particular for retrieving metadata.
2691 #
2692 # As of dCache v10.2.14, the webDAV door leaves the network connection
2693 # unusable for us for sending subsequent requests after serving
2694 # GET, PUT, DELETE, etc., but leaves the connection intact after
2695 # serving MKCOL, MOVE and PROPFIND requests.
2696 # We take advantage of that by using a dedicated pool manager for
2697 # those requests, so that the network connections managed by that pool
2698 # be reused. This avoids establishing the TCP+TLS connection for each
2699 # request.
2700 pool_manager = self._make_pool_manager(self._config)
2701 self._propfind_pool_manager = pool_manager
2702 self._move_pool_manager = pool_manager
2703 self._mkcol_pool_manager = pool_manager
2705 # dCache does not deliver macaroons when we are not using a secure
2706 # channel to interact with the door. In that case, we can not use
2707 # third party copy and dCache does not correctly support the COPY
2708 # method as stated in RFC-4918.
2709 self._can_duplicate = self._base_url.startswith("https://")
2711 @override
2712 def _mkcol(
2713 self,
2714 url: str,
2715 headers: dict[str, str] | None = None,
2716 pool_manager: PoolManager | None = None,
2717 ) -> HTTPResponse:
2718 # Docstring inherited.
2719 return self._request("MKCOL", url=url, headers=headers, pool_manager=self._mkcol_pool_manager)
2721 @override
2722 def _move(
2723 self,
2724 url: str,
2725 headers: dict[str, str] | None = None,
2726 pool_manager: PoolManager | None = None,
2727 ) -> HTTPResponse:
2728 # Docstring inherited.
2729 return self._request("MOVE", url=url, headers=headers, pool_manager=self._move_pool_manager)
2731 @override
2732 def _propfind(
2733 self,
2734 url: str,
2735 headers: dict[str, str] | None = None,
2736 body: str = "",
2737 pool_manager: PoolManager | None = None,
2738 ) -> HTTPResponse:
2739 # Docstring inherited.
2740 return self._request(
2741 "PROPFIND", url=url, headers=headers, body=body, pool_manager=self._propfind_pool_manager
2742 )
2744 @override
2745 def put(
2746 self,
2747 url: str,
2748 headers: dict[str, str] | None = None,
2749 data: BinaryIO | bytes = b"",
2750 ) -> int | None:
2751 # Docstring inherited.
2753 # Send a PUT request with empty body to the dCache frontend server to
2754 # get redirected to the backend.
2755 #
2756 # Details:
2757 # https://www.dcache.org/manuals/UserGuide-10.2/webdav.shtml#redirection
2758 frontend_headers = {} if headers is None else dict(headers)
2759 frontend_headers.update({"Content-Length": "0", "Expect": "100-continue"})
2760 if is_zero_length := isinstance(data, bytes) and len(data) == 0:
2761 # We are uploading an empty file. Don't send the "Expect" header
2762 # so that the dCache door handles this PUT request itself without
2763 # redirecting us to a pool.
2764 frontend_headers.pop("Expect")
2766 resp = self._put(url, headers=frontend_headers, body=b"", redirect=False)
2767 match resp.status:
2768 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT:
2769 redirect_url = url
2770 case status if status in resp.REDIRECT_STATUSES:
2771 redirect_url = resp.headers.get("Location")
2772 case _:
2773 raise unexpected_status_error("PUT", url, resp)
2775 # If we are uploading an empty file, there is nothing more to do.
2776 if is_zero_length:
2777 return 0
2779 # We may have beend redirected to a backend server. Upload the file
2780 # contents to its final destination. Explicitly ask the server to close
2781 # this network connection after serving this PUT request to release
2782 # the associated dCache mover.
2783 backend_headers = {} if headers is None else dict(headers)
2784 backend_headers.update({"Connection": "close"})
2786 # Ask dCache to compute and record a checksum of the uploaded
2787 # file contents, for later integrity checks. Since we don't compute
2788 # the digest ourselves while uploading the data, we cannot control
2789 # after the request is complete that the data we uploaded is
2790 # identical to the data recorded by the server, but at least the
2791 # server has recorded a digest of the data it stored.
2792 #
2793 # See RFC-3230 for details and
2794 # https://www.iana.org/assignments/http-dig-alg/http-dig-alg.xhtml
2795 # for the list of supported digest algorithhms.
2796 if (checksum := self._config.request_checksum) is not None:
2797 backend_headers.update({"Want-Digest": checksum})
2799 resp = self._put(redirect_url, body=data, headers=backend_headers)
2800 match resp.status:
2801 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT:
2802 # Parse the response body and extract the number of bytes
2803 # uploaded. This allows us to avoid sending a HEAD request
2804 # to retrieve the file size.
2805 response_body = resp.data.decode()
2806 if match := DavClientDCache.rex.match(response_body):
2807 return int(match.group(1))
2808 else:
2809 return None
2810 case _:
2811 raise unexpected_status_error("PUT", redirect_url, resp)
2813 @override
2814 def download(self, url: str, filename: str, chunk_size: int) -> int:
2815 """Download the content of a file and write it to local file.
2817 Parameters
2818 ----------
2819 url : `str`
2820 Target URL.
2821 filename : `str`
2822 Local file to write the content to. If the file already exists,
2823 it will be rewritten.
2824 chunk_size : `int`
2825 Size of the chunks to write to `filename`.
2827 Returns
2828 -------
2829 count: `int`
2830 Number of bytes written to `filename`.
2832 Notes
2833 -----
2834 The caller must ensure that the resource at `url` is a file, not
2835 a directory.
2836 """
2837 # Send a GET request without following redirection to get redirected
2838 # to the backend server.
2839 _, resp = self.get(url, preload_content=False, redirect=False)
2840 match resp.status:
2841 case HTTPStatus.OK:
2842 # We were not redirected. Consume this response.
2843 return self._write_response_body_to_file(resp, filename, chunk_size)
2844 case status if status not in resp.REDIRECT_STATUSES:
2845 raise unexpected_status_error("GET", url, resp)
2846 case _:
2847 # We were redirected. Follow this redirection.
2848 pass
2850 # Drain and release the response we received from the frontend server
2851 # so that the connection can be reused.
2852 resp.drain_conn()
2853 resp.release_conn()
2855 # We were redirected to a backend server. Send a GET request to the
2856 # backend server and ask it to close the HTTP connection to force
2857 # closing the network connection.
2858 redirect_url = resp.headers.get("Location")
2859 _, resp = self.get(redirect_url, headers={"Connection": "close"}, preload_content=False)
2860 match resp.status:
2861 case HTTPStatus.OK:
2862 return self._write_response_body_to_file(resp, filename, chunk_size)
2863 case _:
2864 raise unexpected_status_error("GET", redirect_url, resp)
2866 @override
2867 def read(self, url: str) -> tuple[str, bytes]:
2868 """Download the contents of file located at `url`.
2870 Parameters
2871 ----------
2872 url : `str`
2873 Target URL.
2875 Returns
2876 -------
2877 url: `str`
2878 Backend URL from which the data was obtained.
2879 data: `bytes`
2880 Contents of the file.
2882 Notes
2883 -----
2884 The caller must ensure that the resource at `url` is a file, not
2885 a directory.
2886 """
2887 # Send a GET request without following redirection to get redirected
2888 # to the backend server.
2889 backend_url, resp = self.get(url, redirect=False)
2890 match resp.status:
2891 case HTTPStatus.OK:
2892 return backend_url, resp.data
2893 case status if status in resp.REDIRECT_STATUSES:
2894 redirect_url = resp.headers.get("Location")
2895 case _:
2896 raise unexpected_status_error("GET", url, resp)
2898 # We were redirected. Send a GET request to the backend server
2899 # and ask it to close the HTTP connection to force closing the
2900 # network connection.
2901 final_url, resp = self.get(redirect_url, headers={"Connection": "close"})
2902 match resp.status:
2903 case HTTPStatus.OK:
2904 return final_url, resp.data
2905 case _:
2906 raise unexpected_status_error("GET", redirect_url, resp)
2908 @override
2909 def write(self, url: str, data: BinaryIO | bytes) -> int | None:
2910 """Create or rewrite a remote file at `url` with `data` as its
2911 contents.
2913 Parameters
2914 ----------
2915 url : `str`
2916 Target URL.
2917 data : `bytes`
2918 Sequence of bytes to upload.
2920 Returns
2921 -------
2922 size : `int | None`
2923 The size in bytes of the file uploaded. Can be `None` if the size
2924 could not be retrieved.
2926 Notes
2927 -----
2928 If a file already exists at `url` it will be rewritten.
2929 """
2930 # dCache will automatically create all the parent directories so we
2931 # don't need to explicitly create them. Although this is not compliant
2932 # to RFC 4918, this is advantageous because it avoids several
2933 # round-trips to the server for creating all the directories
2934 # before actually uploading the data.
2935 try:
2936 # Upload to a temporary file and rename to the final name.
2937 temporary_url = self._make_temporary_url(url)
2938 size = self.put(temporary_url, data=data)
2939 self.rename(temporary_url, url, overwrite=True, create_parent=False)
2941 # Update the file size cache with this size
2942 self._file_size_cache.update_size(url, size)
2943 return size
2944 except Exception:
2945 # Upload failed. Attempt to remove the temporary file.
2946 self.delete(temporary_url)
2947 raise
2949 @override
2950 def mkcol(self, url: str) -> None:
2951 """Create a directory at `url`.
2953 If a directory already exists at `url` no error is returned nor
2954 exception is raised. An exception is raised if a file exists at `url`.
2956 Parameters
2957 ----------
2958 url : `str`
2959 Target URL.
2960 """
2961 # A "MKCOL" request to dCache does not automatically create all
2962 # the intermediate directories if they do not exist. However, a
2963 # "PUT" request of a file does create the directory hierarchy.
2964 #
2965 # We exploit that to create directory hierarchies: we first create an
2966 # empty file with a random name and then we remove it. As a side
2967 # effect, the target directory will be created.
2968 #
2969 # Creating a directory this way implies two requests to the server
2970 # ("PUT" and "DELETE"), while using "MKCOL" would on average imply
2971 # one request per inexisting directory in the hierarchy. When
2972 # directory hierarchies are relatively deep, requiring two
2973 # requests per hierarchy is better than sending a "MKCOL" request
2974 # per directory in the hierarchy.
2975 try:
2976 temporary_url = self._make_temporary_url(url=f"{url}/mkcol")
2977 self.put(temporary_url, data=b"")
2978 finally:
2979 self.delete(temporary_url)
2981 @override
2982 def info(self, url: str, name: str | None = None) -> dict[str, Any]:
2983 # Docstring inherited.
2984 result: dict[str, Any] = {
2985 "name": name if name is not None else url,
2986 "type": None,
2987 "size": None,
2988 "last_modified": datetime.min,
2989 "checksums": {},
2990 }
2992 # Request live DAV properties as well as the checksums that dCache
2993 # recorded about this file.
2994 body = (
2995 """<?xml version="1.0" encoding="utf-8"?>"""
2996 """<D:propfind xmlns:D="DAV:" xmlns:dcache="http://www.dcache.org/2013/webdav">"""
2997 """<D:prop>"""
2998 """<D:resourcetype/>"""
2999 """<D:getcontentlength/>"""
3000 """<D:getlastmodified/>"""
3001 """<D:displayname/>"""
3002 """<dcache:Checksums/>"""
3003 """</D:prop>"""
3004 """</D:propfind>"""
3005 )
3006 resp = self.propfind(url, body=body, depth="0")
3007 match resp.status:
3008 case HTTPStatus.NOT_FOUND:
3009 return result
3010 case HTTPStatus.MULTI_STATUS:
3011 property = self._propfind_parser.parse(resp.data)[0]
3012 metadata = DavFileMetadata.from_property(base_url=self._base_url, property=property)
3013 result.update(
3014 {
3015 "type": "directory" if metadata.is_dir else "file",
3016 "size": metadata.size,
3017 "last_modified": metadata.last_modified,
3018 "checksums": metadata.checksums,
3019 }
3020 )
3021 return result
3022 case _:
3023 raise unexpected_status_error("PROPFIND", url, resp)
3025 @override
3026 def read_range(
3027 self,
3028 url: str,
3029 start: int,
3030 end: int | None,
3031 headers: dict[str, str] | None = None,
3032 ) -> tuple[str, bytes]:
3033 # Docstring inherited.
3034 range_headers = {"Accept-Encoding": "identity"}
3035 if end is None:
3036 range_headers.update({"Range": f"bytes={start}-"})
3037 else:
3038 range_headers.update({"Range": f"bytes={start}-{end}"})
3040 frontend_headers = {} if headers is None else dict(headers)
3041 frontend_headers.update(range_headers)
3043 # Send the GET request to the dCache door but don't follow redirections
3044 # automatically. We need to be able to add a `Connection: close`
3045 # request header when sending the request to the dCache pool we
3046 # will be redirected to.
3047 #
3048 # We don't send that header to the door since we want to keep the
3049 # network connection with the door open for later reuse.
3050 final_url, resp = self.get(url, headers=frontend_headers, redirect=False)
3051 match resp.status:
3052 case HTTPStatus.PARTIAL_CONTENT:
3053 return final_url, resp.data
3054 case status if status not in resp.REDIRECT_STATUSES:
3055 raise unexpected_status_error("GET (with 'Range' header)", url, resp)
3056 case _:
3057 pass
3059 # We were redirected to the dCache pool. Follow the redirection and
3060 # add a `Connection: close` header to notify the mover to stop its
3061 # execution after serving this request.
3062 backend_headers = {} if headers is None else dict(headers)
3063 backend_headers.update(range_headers)
3064 backend_headers.update({"Connection": "close"})
3066 redirect_url = resp.headers.get("Location")
3067 _, resp = self.get(redirect_url, headers=backend_headers, redirect=True)
3068 match resp.status:
3069 case HTTPStatus.PARTIAL_CONTENT:
3070 # Return the door URL so that subsequent requests (if any)
3071 # go through the dCache door instead first of going directly to
3072 # the pool since we asked the mover to be stopped.
3073 return url, resp.data
3074 case _:
3075 raise unexpected_status_error("GET (with 'Range' header)", redirect_url, resp)
3077 @override
3078 def _close(self, url: str) -> None:
3079 # Docstring inherited.
3081 # For dCache this is a NOP since `read_range` does not keep the
3082 # connection with the dCache pool open.
3083 pass
3086class DavClientXrootD(DavClientURLSigner):
3087 """Client for interacting with a XrootD webDAV server.
3089 Instances of this class are thread-safe.
3091 Parameters
3092 ----------
3093 url : `str`
3094 Root URL of the storage endpoint
3095 (e.g. "https://host.example.org:1234/").
3096 config : `DavConfig`
3097 Configuration to initialize this client.
3098 accepts_ranges : `bool` | `None`
3099 Indicate whether the remote server accepts the ``Range`` header in GET
3100 requests.
3101 """
3103 def __init__(self, url: str, config: DavConfig, accepts_ranges: bool | None = None) -> None:
3104 super().__init__(url=url, config=config, accepts_ranges=accepts_ranges)
3106 @override
3107 def put(
3108 self,
3109 url: str,
3110 headers: dict[str, str] | None = None,
3111 data: BinaryIO | bytes = b"",
3112 ) -> int | None:
3113 # Docstring inherited.
3115 # Send a PUT request with empty body to the XRootD frontend server to
3116 # get redirected to the backend.
3117 frontend_headers = {} if headers is None else dict(headers)
3118 frontend_headers.update({"Content-Length": "0", "Expect": "100-continue"})
3119 for attempt in range(max_attempts := 3):
3120 resp = self._put(url, headers=frontend_headers, body=b"", redirect=False)
3121 if resp.status in (
3122 HTTPStatus.OK,
3123 HTTPStatus.CREATED,
3124 HTTPStatus.NO_CONTENT,
3125 ):
3126 redirect_url = url
3127 break
3128 elif resp.status in resp.REDIRECT_STATUSES:
3129 redirect_url = resp.headers.get("Location")
3130 break
3131 elif resp.status == HTTPStatus.LOCKED:
3132 # Sometimes XRootD servers respond with status code LOCKED and
3133 # response body of the form:
3134 #
3135 # "Output file /path/to/file is already opened by 1 writer;
3136 # open denied."
3137 #
3138 # If we get such a response, try again, unless we reached
3139 # the maximum number of attempts.
3140 if attempt == max_attempts - 1:
3141 raise ValueError(
3142 f"""Unexpected response to HTTP request PUT {resp.geturl()}: status {resp.status} """
3143 f"""{resp.reason} [{resp.data.decode()}] after {max_attempts} attempts"""
3144 )
3146 # Wait a bit and try again
3147 log.warning(
3148 f"""got unexpected response status {HTTPStatus.LOCKED} Locked for PUT {resp.geturl()} """
3149 f"""(attempt {attempt}/{max_attempts}), retrying..."""
3150 )
3151 time.sleep((attempt + 1) * 0.100)
3152 continue
3153 else:
3154 raise unexpected_status_error("PUT", url, resp)
3156 # We were redirected to a backend server. Upload the file contents to
3157 # its final destination.
3159 # XRootD backend servers typically use a single port number for
3160 # accepting connections from clients. It is therefore beneficial
3161 # to keep those connections open, if the server allows.
3163 # Ask the server to compute and record a checksum of the uploaded
3164 # file contents, for later integrity checks. Since we don't compute
3165 # the digest ourselves while uploading the data, we cannot control
3166 # after the request is complete that the data we uploaded is
3167 # identical to the data recorded by the server, but at least the
3168 # server has recorded a digest of the data it stored.
3169 #
3170 # See RFC-3230 for details and
3171 # https://www.iana.org/assignments/http-dig-alg/http-dig-alg.xhtml
3172 # for the list of supported digest algorithhms.
3173 #
3174 # In addition, note that not all servers implement this RFC so
3175 # the checksum reqquest may be ignored by the server.
3176 backend_headers = {} if headers is None else dict(headers)
3177 if (checksum := self._config.request_checksum) is not None:
3178 backend_headers.update({"Want-Digest": checksum})
3180 resp = self._put(redirect_url, body=data, headers=backend_headers)
3181 match resp.status:
3182 case HTTPStatus.OK | HTTPStatus.CREATED | HTTPStatus.NO_CONTENT:
3183 # Send a HEAD request to retrieve the size of the file we
3184 # just uploaded.
3185 resp = self.head(redirect_url)
3186 size = int(resp.headers.get("Content-Length", -1))
3187 return None if size == -1 else size
3188 case _:
3189 raise unexpected_status_error("PUT", redirect_url, resp)
3191 @override
3192 def info(self, url: str, name: str | None = None) -> dict[str, Any]:
3193 # XRootD does not include checksums in the response to PROPFIND
3194 # request. We need to send a specific HEAD request to retrieve
3195 # the ADLER32 checksum.
3196 #
3197 # If found, the checksum is included in the response header "Digest",
3198 # which is of the form:
3199 #
3200 # Digest: adler32=0e4709f2
3201 result = super().info(url, name)
3202 if result["type"] == "file":
3203 headers: dict[str, str] = {"Want-Digest": "adler32"}
3204 resp = self.head(url=url, headers=headers)
3205 if (digest := resp.headers.get("Digest")) is not None:
3206 value = digest.split("=")[1]
3207 result["checksums"].update({"adler32": value})
3209 return result
3211 @override
3212 def write(self, url: str, data: BinaryIO | bytes) -> int | None:
3213 """Create or rewrite a remote file at `url` with `data` as its
3214 contents.
3216 Parameters
3217 ----------
3218 url : `str`
3219 Target URL.
3220 data : `bytes`
3221 Sequence of bytes to upload.
3223 Returns
3224 -------
3225 size : `int | None`
3226 The size in bytes of the file uploaded. Can be `None` if the size
3227 could not be retrieved.
3229 Notes
3230 -----
3231 If a file already exists at `url` it will be rewritten.
3232 """
3233 # XRootD will automatically create all the parent directories so we
3234 # don't need to explicitly create them. Although this is not compliant
3235 # to RFC 4918, this is advantageous because it avoids several
3236 # round-trips to the server for creating all the directories
3237 # before actually uploading the data.
3238 try:
3239 # Upload to a temporary file and rename to the final name.
3240 temporary_url = self._make_temporary_url(url)
3241 size = self.put(temporary_url, data=data)
3242 self.rename(temporary_url, url, overwrite=True, create_parent=False)
3244 # Update the file size cache with this size
3245 self._file_size_cache.update_size(url, size)
3246 return size
3247 except Exception:
3248 # Upload failed. Attempt to remove the temporary file.
3249 self.delete(temporary_url)
3250 raise
3252 @override
3253 def mkcol(self, url: str) -> None:
3254 """Create a directory at `url`.
3256 If a directory already exists at `url` no error is returned nor
3257 exception is raised. An exception is raised if a file exists at `url`.
3259 Parameters
3260 ----------
3261 url : `str`
3262 Target URL.
3263 """
3264 # XRootD automatically creates all the intermediate directories.
3265 resp = self._mkcol(url)
3266 match resp.status:
3267 case HTTPStatus.CREATED:
3268 return
3269 case HTTPStatus.METHOD_NOT_ALLOWED:
3270 # XRootD returns "405 Method Not Allowed" when either a file
3271 # or a directory already exists at `url`
3272 stat = self.stat(url)
3273 if stat.is_dir:
3274 # A directory exists at `url`. Nothing more to do.
3275 return
3276 elif stat.is_file:
3277 raise NotADirectoryError(
3278 f"Can not create a directory because a file already exists at {resp.geturl()}"
3279 )
3280 case _:
3281 raise ValueError(
3282 f"Can not create directory {resp.geturl()}: status {resp.status} {resp.reason}"
3283 )
3285 @override
3286 def stat(self, url: str) -> DavFileMetadata:
3287 # Docstring inherited.
3289 # XRootD v5.9.1 responds "200 OK" to a HEAD request against an
3290 # existing file. When the target URL is a directory, it also responds
3291 # "200 OK". In both cases the response header "Content-Length"
3292 # is present but has different meaning. If the target URL is a file,
3293 # the header value is the size in bytes of the file. If the target
3294 # URL is a directory, the header value is the number of items in
3295 # the directory.
3296 #
3297 # So there is not an easy way to determine if the target URL is a
3298 # file or a directory from the response to a HEAD request.
3299 #
3300 # When the target URL is a directory and we ask for a digest, the
3301 # server responds "409 Conflict". We use this behavior to
3302 # discriminate between a file and a directory.
3303 #
3304 # Note that XRootD does not include the "Last-Modified" header in the
3305 # response to a HEAD request so we cannot include the last modified
3306 # time in the value returned by this method.
3307 resp = self._head(url, headers={"Want-Digest": "adler32"})
3308 match resp.status:
3309 case HTTPStatus.OK:
3310 # There is a file at target URL
3311 if "Content-Length" in resp.headers:
3312 href = url.replace(self._base_url, "", 1)
3313 size = int(resp.headers.get("Content-Length"))
3314 return DavFileMetadata(self._base_url, href=href, exists=True, is_dir=False, size=size)
3315 else:
3316 raise ValueError(
3317 f"""Expecting Content-Length header to be present in """
3318 f"""response to HTTP HEAD {resp.geturl()}: status {resp.status} """
3319 f"""{resp.reason} [{resp.data.decode()}] but could not find it"""
3320 )
3321 case HTTPStatus.CONFLICT:
3322 # There is a directory at target URL
3323 href = url.replace(self._base_url, "", 1)
3324 return DavFileMetadata(self._base_url, href=href, exists=True, is_dir=True)
3325 case HTTPStatus.NOT_FOUND:
3326 # There is neither a file nor a directory at target URL
3327 return DavFileMetadata(base_url=url, exists=False)
3328 case _:
3329 raise unexpected_status_error("HEAD", url, resp)
3331 @override
3332 def read_range(
3333 self,
3334 url: str,
3335 start: int,
3336 end: int | None,
3337 headers: dict[str, str] | None = None,
3338 ) -> tuple[str, bytes]:
3339 # Docstring inherited.
3341 # Send the request to the XRootD redirector and follow
3342 # redirections automatically.
3343 #
3344 # The network connection with the redirector and with the backend
3345 # file server are left open for later reuse.
3346 range_headers = {"Accept-Encoding": "identity"}
3347 if end is None:
3348 range_headers.update({"Range": f"bytes={start}-"})
3349 else:
3350 range_headers.update({"Range": f"bytes={start}-{end}"})
3352 get_headers = {} if headers is None else dict(headers)
3353 get_headers.update(range_headers)
3355 final_url, resp = self.get(url, headers=get_headers, redirect=True)
3356 match resp.status:
3357 case HTTPStatus.PARTIAL_CONTENT:
3358 return final_url, resp.data
3359 case _:
3360 raise unexpected_status_error("GET (with 'Range' header)", url, resp)
3362 @override
3363 def _close(self, url: str) -> None:
3364 # Docstring inherited.
3366 # Send a `HEAD` request with a `Connection: close` header to notify
3367 # the remote server that we are not sending other GET requests with
3368 # `Range` header for this URL.
3369 self._request("HEAD", url=url, headers={"Connection": "close"}, redirect=False)
3372class DavFileMetadata:
3373 """Container for attributes of interest of a webDAV file or directory.
3375 Parameters
3376 ----------
3377 base_url : `str`
3378 Base URL.
3379 href : `str`, optional
3380 Path component that can be added to the base URL.
3381 name : `str`, optional
3382 Name.
3383 exists : `bool`, optional
3384 Whether file or directory exist.
3385 size : `int`, optional
3386 Size of file.
3387 is_dir : `bool`, optional
3388 Whether the URL points to a directory or file.
3389 last_modified : `bool`, optional
3390 Last modified date.
3391 checksums : `dict` [ `str`, `str` ] | `None`, optional
3392 Checksums.
3393 """
3395 def __init__(
3396 self,
3397 base_url: str,
3398 href: str = "",
3399 name: str = "",
3400 exists: bool = False,
3401 size: int = -1,
3402 is_dir: bool = False,
3403 last_modified: datetime = datetime.min,
3404 checksums: dict[str, str] | None = None,
3405 ):
3406 self._url: str = base_url if not href else base_url.rstrip("/") + href
3407 self._href: str = href
3408 self._name: str = name
3409 self._exists: bool = exists
3410 self._size: int = size
3411 self._is_dir: bool = is_dir
3412 self._last_modified: datetime = last_modified
3413 self._checksums: dict[str, str] = {} if checksums is None else dict(checksums)
3415 @staticmethod
3416 def from_property(base_url: str, property: DavProperty) -> DavFileMetadata:
3417 """Create an instance from the values in `property`.
3419 Parameters
3420 ----------
3421 base_url : `str`
3422 Base URL.
3423 property : `DavProperty`
3424 Properties to associate with URL.
3425 """
3426 return DavFileMetadata(
3427 base_url=base_url,
3428 href=property.href,
3429 name=property.name,
3430 exists=property.exists,
3431 size=property.size,
3432 is_dir=property.is_dir,
3433 last_modified=property.last_modified,
3434 checksums=dict(property.checksums),
3435 )
3437 def __str__(self) -> str:
3438 return (
3439 f"""{self._url} {self._href} {self._name} {self._exists} {self._size} {self._is_dir} """
3440 f"""{self._checksums}"""
3441 )
3443 @property
3444 def url(self) -> str:
3445 return self._url
3447 @property
3448 def href(self) -> str:
3449 return self._href
3451 @property
3452 def name(self) -> str:
3453 return self._name
3455 @property
3456 def exists(self) -> bool:
3457 return self._exists
3459 @property
3460 def size(self) -> int:
3461 if not self._exists:
3462 return -1
3464 return 0 if self._is_dir else self._size
3466 @property
3467 def is_dir(self) -> bool:
3468 return self._exists and self._is_dir
3470 @property
3471 def is_file(self) -> bool:
3472 return self._exists and not self._is_dir
3474 @property
3475 def last_modified(self) -> datetime:
3476 return self._last_modified
3478 @property
3479 def checksums(self) -> dict[str, str]:
3480 return self._checksums
3483class DavProperty:
3484 """Helper class to encapsulate select live DAV properties of a single
3485 resource, as retrieved via a PROPFIND request.
3487 Parameters
3488 ----------
3489 response : `eTree.Element` or `None`
3490 The XML response defining the DAV property.
3491 """
3493 # Regular expression to compare against the 'status' element of a
3494 # PROPFIND response's 'propstat' element.
3495 _status_ok_rex = re.compile(r"^HTTP/.* 200 .*$", re.IGNORECASE)
3497 def __init__(self, response: eTree.Element | None):
3498 self._href: str = ""
3499 self._displayname: str = ""
3500 self._collection: bool = False
3501 self._getlastmodified: str = ""
3502 self._getcontentlength: int = -1
3503 self._checksums: dict[str, str] = {}
3505 if response is not None:
3506 self._parse(response)
3508 def _parse(self, response: eTree.Element) -> None:
3509 # Extract 'href'.
3510 if (element := response.find("./{DAV:}href")) is not None:
3511 # We need to use "str(element.text)"" instead of "element.text" to
3512 # keep mypy happy.
3513 self._href = str(element.text).strip()
3514 else:
3515 raise ValueError(
3516 "Property 'href' expected but not found in PROPFIND response: "
3517 f"{eTree.tostring(response, encoding='unicode')}"
3518 )
3520 for propstat in response.findall("./{DAV:}propstat"):
3521 # Only extract properties of interest with status OK.
3522 status = propstat.find("./{DAV:}status")
3523 if status is None or not self._status_ok_rex.match(str(status.text)):
3524 continue
3526 for prop in propstat.findall("./{DAV:}prop"):
3527 # Parse "collection".
3528 if (element := prop.find("./{DAV:}resourcetype/{DAV:}collection")) is not None:
3529 self._collection = True
3531 # Parse "getlastmodified".
3532 if (element := prop.find("./{DAV:}getlastmodified")) is not None:
3533 self._getlastmodified = str(element.text)
3535 # Parse "getcontentlength".
3536 if (element := prop.find("./{DAV:}getcontentlength")) is not None:
3537 self._getcontentlength = int(str(element.text))
3539 # Parse "displayname".
3540 if (element := prop.find("./{DAV:}displayname")) is not None:
3541 self._displayname = str(element.text)
3543 # Parse "Checksums"
3544 if (element := prop.find("./{http://www.dcache.org/2013/webdav}Checksums")) is not None:
3545 self._checksums = self._parse_checksums(element.text)
3547 # Some webDAV servers don't include the 'displayname' property in the
3548 # response so try to infer it from the value of the 'href' property.
3549 # Depending on the server the href value may end with '/'.
3550 if not self._displayname:
3551 self._displayname = os.path.basename(self._href.rstrip("/"))
3553 # Some webDAV servers do not append a "/" to the href of directories.
3554 # Ensure we include a single final "/" in our response.
3555 if self._collection:
3556 self._href = self._href.rstrip("/") + "/"
3558 # Force a size of 0 for collections.
3559 if self._collection:
3560 self._getcontentlength = 0
3562 def _parse_checksums(self, checksums: str | None) -> dict[str, str]:
3563 # checksums argument is of the form
3564 # md5=MyS/wljSzI9WYiyrsuyoxw==,adler32=23b104f2
3565 result: dict[str, str] = {}
3566 if checksums is not None:
3567 for checksum in checksums.split(","):
3568 if (pos := checksum.find("=")) != -1:
3569 algorithm, value = (checksum[:pos].lower(), checksum[pos + 1 :])
3570 if algorithm == "md5":
3571 # dCache documentation about how it encodes the
3572 # MD5 checksum:
3573 #
3574 # https://www.dcache.org/manuals/UserGuide-10.2/webdav.shtml#checksums
3575 result[algorithm] = bytes.hex(base64.standard_b64decode(value))
3576 else:
3577 result[algorithm] = value
3579 return result
3581 @property
3582 def exists(self) -> bool:
3583 # It is either a directory or a file with length of at least zero
3584 return self._collection or self._getcontentlength >= 0
3586 @property
3587 def is_dir(self) -> bool:
3588 return self._collection
3590 @property
3591 def is_file(self) -> bool:
3592 return not self._collection
3594 @property
3595 def last_modified(self) -> datetime:
3596 if not self._getlastmodified:
3597 return datetime.min
3599 # Last modified timestamp is of the form:
3600 # 'Wed, 12 Mar 2025 10:11:13 GMT'
3601 return datetime.strptime(self._getlastmodified, "%a, %d %b %Y %H:%M:%S %Z").replace(tzinfo=UTC)
3603 @property
3604 def size(self) -> int:
3605 return self._getcontentlength
3607 @property
3608 def name(self) -> str:
3609 return self._displayname
3611 @property
3612 def href(self) -> str:
3613 return self._href
3615 @property
3616 def checksums(self) -> dict[str, str]:
3617 return self._checksums
3620class DavPropfindParser:
3621 """Helper class to parse the response body of a PROPFIND request."""
3623 def __init__(self) -> None:
3624 return
3626 def parse(self, body: bytes) -> list[DavProperty]:
3627 """Parse the XML-encoded contents of the response body to a webDAV
3628 PROPFIND request.
3630 Parameters
3631 ----------
3632 body : `bytes`
3633 XML-encoded response body to a PROPFIND request.
3635 Returns
3636 -------
3637 responses : `list` [ `DavProperty` ]
3638 Parsed content of the response.
3640 Notes
3641 -----
3642 Is is expected that there is at least one reponse in `body`, otherwise
3643 this function raises.
3644 """
3645 # A response body to a PROPFIND request is of the form (indented for
3646 # readability):
3647 #
3648 # <?xml version="1.0" encoding="UTF-8"?>
3649 # <D:multistatus xmlns:D="DAV:">
3650 # <D:response>
3651 # <D:href>path/to/resource</D:href>
3652 # <D:propstat>
3653 # <D:prop>
3654 # <D:resourcetype>
3655 # <D:collection xmlns:D="DAV:"/>
3656 # </D:resourcetype>
3657 # <D:getlastmodified>
3658 # Fri, 27 Jan 2 023 13:59:01 GMT
3659 # </D:getlastmodified>
3660 # <D:getcontentlength>
3661 # 12345
3662 # </D:getcontentlength>
3663 # </D:prop>
3664 # <D:status>
3665 # HTTP/1.1 200 OK
3666 # </D:status>
3667 # </D:propstat>
3668 # </D:response>
3669 # <D:response>
3670 # ...
3671 # </D:response>
3672 # <D:response>
3673 # ...
3674 # </D:response>
3675 # </D:multistatus>
3677 # Scan all the 'response' elements and extract the relevant properties
3678 decoded_body: str = body.decode().strip()
3679 responses = []
3680 multistatus = eTree.fromstring(decoded_body)
3681 for response in multistatus.findall("./{DAV:}response"):
3682 responses.append(DavProperty(response))
3684 if responses:
3685 return responses
3686 else:
3687 # Could not parse the body
3688 raise ValueError(f"Unable to parse response for PROPFIND request: {decoded_body}")
3691class Authorizer:
3692 """Base class for attaching an 'Authorization' header to a HTTP request."""
3694 def set_authorization(self, headers: dict[str, str]) -> None:
3695 """Add the 'Authorization' header to `headers`.
3697 Parameters
3698 ----------
3699 headers : `dict` [ `str`, `str` ]
3700 Dict to augment with authorization information.
3702 Notes
3703 -----
3704 This method must be implemented by concrete subclasses.
3705 """
3706 raise NotImplementedError
3708 def _is_file_protected(self, filepath: str) -> bool:
3709 """Return true if the permissions of file at `filepath` only allow for
3710 access by its owner.
3712 Parameters
3713 ----------
3714 filepath : `str`
3715 Path of a local file.
3716 """
3717 if not os.path.isfile(filepath): 3717 ↛ 3718line 3717 didn't jump to line 3718 because the condition on line 3717 was never true
3718 return False
3720 mode = stat.S_IMODE(os.stat(filepath).st_mode)
3721 owner_accessible = bool(mode & stat.S_IRWXU)
3722 group_accessible = bool(mode & stat.S_IRWXG)
3723 other_accessible = bool(mode & stat.S_IRWXO)
3724 return owner_accessible and not group_accessible and not other_accessible
3726 def _read_if_modified_since(
3727 self, filename: str | None, timestamp: float
3728 ) -> tuple[str, float] | tuple[None, None]:
3729 """Read local file `filename` if its modification time is more
3730 recent than `timestamp`.
3732 Parameters
3733 ----------
3734 filename : `str`, optional
3735 Path of a local file.
3737 timestamp: `float`, optional
3738 Timestamp to compare against the last modification time of
3739 `filename`. The contents of file at `filename` is only read if its
3740 modification time is more recent than `timestamp`.
3742 Returns
3743 -------
3744 result: `tuple[str, float]`
3745 tuple of (contents of file `filename`, timestamp of the read
3746 operation).
3748 If `filename` is `None`, the returned value is `tuple[None, None]`.
3749 """
3750 if filename is None: 3750 ↛ 3751line 3750 didn't jump to line 3751 because the condition on line 3750 was never true
3751 return (None, None)
3753 if os.stat(filename).st_mtime < timestamp:
3754 return (None, None)
3756 with open(filename) as file:
3757 time_of_last_read = time.time()
3758 return (file.read().rstrip("\n"), time_of_last_read)
3761class TokenAuthorizer(Authorizer):
3762 """Attach a bearer token 'Authorization' header to each request.
3764 Parameters
3765 ----------
3766 token : `str`
3767 Can be either the path to a local file which contains the
3768 value of the token or the token itself. If `token` is a file
3769 it must be protected so that only the owner can read and write it.
3770 """
3772 def __init__(self, token: str | None = None) -> None:
3773 self._token = self._token_path = None
3774 self._time_of_last_read: float = -1.0
3775 if token is None:
3776 return
3778 self._token = token
3779 if os.path.isfile(token):
3780 self._token_path = os.path.abspath(token)
3781 if not self._is_file_protected(self._token_path):
3782 raise PermissionError(
3783 f"""Authorization token file at {self._token_path} must be protected for access only """
3784 """by its owner"""
3785 )
3786 self._update_token()
3788 def _update_token(self) -> None:
3789 """Read the token file (if any) if its modification time is more recent
3790 than the last time we read it.
3791 """
3792 if self._token_path is None:
3793 return None
3795 token, time_of_last_read = self._read_if_modified_since(self._token_path, self._time_of_last_read)
3796 if token is None or time_of_last_read is None:
3797 return
3799 # Update the token value and the last time we read it.
3800 self._token = token
3801 self._time_of_last_read = time_of_last_read
3803 @override
3804 def set_authorization(self, headers: dict[str, str]) -> None:
3805 """Add the 'Authorization' header to `headers`.
3807 Parameters
3808 ----------
3809 headers : `dict` [ `str`, `str` ]
3810 Dict to augment with authorization information.
3811 """
3812 if self._token is None:
3813 return
3815 self._update_token()
3816 headers["Authorization"] = f"Bearer {self._token}"
3819class BasicAuthorizer(Authorizer):
3820 """Attach a 'Authorization' header to each request using Basic
3821 authentication.
3823 Parameters
3824 ----------
3825 user_name : `str`
3826 Can be either the path to a local file which contains the
3827 user name or the user name itself. If `user_name` is a file
3828 it must be protected so that only the owner can read and write it.
3829 user_password : `str`
3830 Can be either the path to a local file which contains the
3831 value of the password or the password itself. If `user_password` is a
3832 file it must be protected so that only the owner can read and write it.
3833 """
3835 def __init__(self, user_name: str | None = None, user_password: str | None = None) -> None:
3836 if user_name is None or user_password is None:
3837 return
3839 self._user_name: str | None = user_name
3840 self._user_password: str | None = user_password
3841 self._user_password_path: str | None = None
3842 self._time_of_last_read: float = -1.0
3843 self._header_value: str = ""
3845 if os.path.isfile(self._user_password):
3846 # The value in `user_password` is the path to a file. Check
3847 # the file is protected and read its contents.
3848 self._user_password_path = os.path.abspath(self._user_password)
3849 if not self._is_file_protected(self._user_password_path):
3850 raise PermissionError(
3851 f"""Password file at {self._user_password_path} must be protected for access only """
3852 """by its owner"""
3853 )
3854 self._update_password()
3855 else:
3856 self._update_header_value()
3858 def _update_header_value(self) -> None:
3859 """Compute the value of the 'Authorization' header using HTTP basic
3860 authorization.
3861 """
3862 basic_auth_header = make_headers(basic_auth=f"{self._user_name}:{self._user_password}")
3863 self._header_value = basic_auth_header["authorization"]
3865 def _update_password(self) -> None:
3866 """Update the password of this authorizer if the file it is stored in
3867 has been modified since the last time we read it.
3868 """
3869 if self._user_password_path is None:
3870 return None
3872 password, time_of_last_read = self._read_if_modified_since(
3873 self._user_password_path, self._time_of_last_read
3874 )
3875 if password is None or time_of_last_read is None:
3876 return
3878 # Update the password, the last time we read it and re-compute the
3879 # value of the "Authorization" header.
3880 self._user_password = password
3881 self._time_of_last_read = time_of_last_read
3882 self._update_header_value()
3884 @override
3885 def set_authorization(self, headers: dict[str, str]) -> None:
3886 """Add the 'Authorization' header to `headers`.
3888 Parameters
3889 ----------
3890 headers : `dict` [ `str`, `str` ]
3891 Dict to augment with authorization information.
3892 """
3893 if self._user_name is None or self._user_password is None:
3894 return
3896 self._update_password()
3897 headers["Authorization"] = self._header_value
3900def expand_vars(path: str | None) -> str | None:
3901 """Expand the environment variables in `path` and return the path with
3902 the value of the variable expanded.
3904 Parameters
3905 ----------
3906 path : `str` or `None`
3907 Abolute or relative path which may include an environment variable
3908 (e.g. '$HOME/path/to/my/file').
3910 Returns
3911 -------
3912 path: `str`
3913 The path with the values of the environment variables expanded.
3914 """
3915 return None if path is None else os.path.expandvars(path)
3918def dump_response(method: str, resp: HTTPResponse, dump_body: bool = False) -> None:
3919 """Dump response for debugging purposes.
3921 Parameters
3922 ----------
3923 method : `str`
3924 Method name to include in log output.
3925 resp : `HTTPResponse`
3926 Response to dump.
3927 dump_body : `bool`, optional
3928 Whether or not to issue a debug log message.
3929 """
3930 log.debug("%s %s", method, resp.geturl())
3931 log.debug(" %s %s", resp.status, resp.reason)
3933 for header, value in resp.headers.items():
3934 log.debug(" %s: %s", header, value)
3936 if dump_body:
3937 log.debug(" response body length: %d", len(resp.data.decode()))