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