Coverage for python/lsst/resources/davutils.py: 31%

1135 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-19 09:03 +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. 

11 

12from __future__ import annotations 

13 

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 

32 

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 

37 

38from urllib.parse import parse_qsl, urlencode, urlparse, urlunparse 

39 

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 

49 

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 

55 

56from lsst.utils.logging import getLogger 

57from lsst.utils.timer import time_this 

58 

59# Use the same logger than `dav.py`. 

60log = getLogger(f"""{__name__.replace(".davutils", ".dav")}""") 

61 

62 

63def normalize_path(path: str | None) -> str: 

64 """Normalize a path intended to be part of a URL. 

65 

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. 

70 

71 Parameters 

72 ---------- 

73 path : `str`, optional 

74 Path to normalize (e.g., '/path/to/..///normalize/'). 

75 

76 Returns 

77 ------- 

78 url : `str` 

79 Normalized URL (e.g., '/path/normalize'). 

80 """ 

81 return "/" if not path else "/" + posixpath.normpath(path).lstrip("/") 

82 

83 

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. 

87 

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. 

98 

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 

111 

112 

113def redact_url(url: str) -> str: 

114 """Return a modified `url` with authorization query redacted. 

115 

116 The goal is that this method should be used for logging URLs to avoid 

117 leaking authorization tokens. 

118 

119 Parameters 

120 ---------- 

121 url : `str` 

122 URL to redact. 

123 

124 Returns 

125 ------- 

126 redacted_url : `str` 

127 For instance, when called with an URL like: 

128 

129 https://host.example.org:1234/a/b/c/file.data?key1=value1&key2=value2&authz=token#fragment 

130 

131 the returned value would be: 

132 

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])) 

139 

140 redacted_url = parsed_url._replace(query=urlencode(redacted_query)) 

141 return urlunparse(redacted_url) 

142 

143 

144class DavConfig: 

145 """Configurable settings a webDAV client must use when interacting with a 

146 particular storage endpoint. 

147 

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"]`. 

153 

154 For instance, if `config["base_url"]` is 

155 

156 "davs://webdav.example.org:1234/" 

157 

158 any object of class `DavResourcePath` like 

159 

160 "davs://webdav.example.org:1234/path/to/any/file" 

161 

162 will use the settings in this configuration to configure its client. 

163 """ 

164 

165 # Timeout in seconds to establish a network connection with the remote 

166 # server. 

167 DEFAULT_TIMEOUT_CONNECT: float = 10.0 

168 

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 

174 

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 

180 

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 

185 

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 

191 

192 # Number of times to retry requests before failing. Retry happens only 

193 # under certain conditions. 

194 DEFAULT_RETRIES: int = 3 

195 

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 

202 

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 

209 

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 

214 

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 

220 

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 

225 

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 

231 

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 

236 

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 

241 

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 

246 

247 # Accepted checksum algorithms. Must be lowercase. 

248 ACCEPTED_CHECKSUMS: list[str] = ["adler32", "md5", "sha-256", "sha-512"] 

249 

250 def __init__(self, config: dict | None = None) -> None: 

251 if config is None: 

252 config = {} 

253 

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) 

258 

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 ) 

302 

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 [] 

306 

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)) 

314 

315 # Eliminate duplicate URLs. 

316 frontend_urls = list(set(frontend_urls)) 

317 

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 ) 

327 

328 return frontend_urls 

329 

330 @property 

331 def base_url(self) -> str: 

332 return self._base_url 

333 

334 @property 

335 def timeout_connect(self) -> float: 

336 return self._timeout_connect 

337 

338 @property 

339 def timeout_read(self) -> float: 

340 return self._timeout_read 

341 

342 @property 

343 def persistent_connections_per_host(self) -> int: 

344 return self._persistent_connections_per_host 

345 

346 @property 

347 def buffer_size(self) -> int: 

348 return self._buffer_size 

349 

350 @property 

351 def block_size(self) -> int: 

352 return self._block_size 

353 

354 @property 

355 def retries(self) -> int: 

356 return self._retries 

357 

358 @property 

359 def retry_backoff_min(self) -> float: 

360 return self._retry_backoff_min 

361 

362 @property 

363 def retry_backoff_max(self) -> float: 

364 return self._retry_backoff_max 

365 

366 @property 

367 def trusted_authorities(self) -> str | None: 

368 return self._trusted_authorities 

369 

370 @property 

371 def token(self) -> str | None: 

372 return self._token 

373 

374 @property 

375 def reuse_connection(self) -> bool: 

376 return self._reuse_connection 

377 

378 @property 

379 def request_checksum(self) -> str | None: 

380 return self._request_checksum 

381 

382 @property 

383 def user_cert(self) -> str | None: 

384 return self._user_cert 

385 

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 

392 

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 

398 

399 @property 

400 def user_name(self) -> str | None: 

401 return self._user_name 

402 

403 @property 

404 def user_password(self) -> str | None: 

405 return self._user_password 

406 

407 @property 

408 def enable_fsspec(self) -> bool: 

409 return self._enable_fsspec 

410 

411 @property 

412 def collect_memory_usage(self) -> bool: 

413 return self._collect_memory_usage 

414 

415 @property 

416 def frontend_urls(self) -> list[str]: 

417 return self._frontend_urls 

418 

419 

420class DavConfigPool: 

421 """Registry of configurable settings for all known webDAV endpoints. 

422 

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. 

430 

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'. 

434 

435 The configuration file is a YAML file with the structure below: 

436 

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 

452 

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 ... 

459 

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. 

463 

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 """ 

468 

469 _instance = None 

470 _lock = threading.Lock() 

471 

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) 

477 

478 return cls._instance 

479 

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() 

485 

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] = {} 

489 

490 # Load the configuration from the file we have been provided with, 

491 # if any. 

492 if filename is None: 

493 return 

494 

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 ) 

516 

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`. 

520 

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 

530 

531 # No config was found for the specified URL. Use the default. 

532 return self._default_config 

533 

534 def _destroy(self) -> None: 

535 """Destroy this class singleton instance. 

536 

537 Helper method to be used in tests to reset global configuration. 

538 """ 

539 with DavConfigPool._lock: 

540 DavConfigPool._instance = None 

541 

542 

543def make_retry(config: DavConfig) -> Retry: 

544 """Create a ``urllib3.util.Retry`` object from settings in `config`. 

545 

546 Parameters 

547 ---------- 

548 config : `DavConfig` 

549 Configurable settings for a webDAV storage endpoint. 

550 

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 

609 

610 

611class DavClientPool: 

612 """Container of reusable webDAV clients, each one specifically configured 

613 to talk to a single storage endpoint. 

614 

615 Parameters 

616 ---------- 

617 config_pool : `DavConfigPool` 

618 Pool of all known webDAV client configurations. 

619 

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 """ 

625 

626 _instance = None 

627 _lock = threading.Lock() 

628 

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) 

634 

635 return cls._instance 

636 

637 def __init__(self, config_pool: DavConfigPool) -> None: 

638 self._config_pool: DavConfigPool = config_pool 

639 

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] = {} 

644 

645 def get_client_for_url(self, url: str) -> DavClient: 

646 """Return a client for interacting with the endpoint where `url` 

647 is hosted. 

648 

649 Parameters 

650 ---------- 

651 url : `str` 

652 URL for which to obtain a client. 

653 

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 

665 

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 

673 

674 config: DavConfig = self._config_pool.get_config_for_url(url) 

675 self._clients[url] = self._make_client(url, config) 

676 

677 return self._clients[url] 

678 

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" 

690 

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) 

704 

705 def _destroy(self) -> None: 

706 """Destroy this class singleton instance. 

707 

708 Helper method to be used in tests to reset global configuration. 

709 """ 

710 with DavClientPool._lock: 

711 DavClientPool._instance = None 

712 

713 

714class DavFileSizeCache: 

715 """Helper class to cache file sizes of recently uploaded files. 

716 

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). 

723 

724 Notes 

725 ----- 

726 There is a single instance of this class shared by several `DavClient` 

727 objects. This singleton is thread safe. 

728 

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 """ 

734 

735 _instance = None 

736 _lock = threading.Lock() 

737 

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) 

743 

744 return cls._instance 

745 

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]] = {} 

761 

762 def invalidate(self, url: str) -> None: 

763 """Invalidate the cache entry for `url`, if any. 

764 

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) 

773 

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. 

778 

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 

793 

794 timeout = self._default_timeout if timeout is None else timeout 

795 with DavFileSizeCache._lock: 

796 self._cache[url] = (size, time.time(), timeout) 

797 

798 def get_size(self, url: str) -> int | None: 

799 """Retrieve the cached valued of the size of file at `url`. 

800 

801 Parameters 

802 ---------- 

803 url : `str` 

804 URL of the file to retrieve the size for. 

805 

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 

819 

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 

830 

831 

832def unexpected_status_error(method: str, url: str, resp: HTTPResponse) -> Exception: 

833 """Raise an exception from `resp`. 

834 

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}]" 

848 

849 return ValueError(message) 

850 

851 

852class DavClient: 

853 """WebDAV client, configured to talk to a single storage endpoint. 

854 

855 Instances of this class are thread-safe. 

856 

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 """ 

868 

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() 

872 

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 

876 

877 # Configuration settings for the storage endpoint this client 

878 # will interact with. 

879 self._config: DavConfig = config 

880 

881 # Make the authorizer for this client's requests. 

882 self._authorizer: Authorizer | None = self._make_authorizer(config=self._config) 

883 

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) 

887 

888 # Parser of PROPFIND responses. 

889 self._propfind_parser: DavPropfindParser = DavPropfindParser() 

890 

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 

894 

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 

900 

901 # Cache to store sizes of files this client has recently uploaded 

902 # to the server. 

903 self._file_size_cache = DavFileSizeCache() 

904 

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) 

913 

914 return None 

915 

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 ) 

928 

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 

936 

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 ) 

983 

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. 

987 

988 Parameters 

989 ---------- 

990 url : `str` 

991 URL to check. 

992 

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. 

999 

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") 

1021 

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 ) 

1026 

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 

1035 

1036 return details 

1037 

1038 def _get_response_url(self, resp: HTTPResponse, default_url: str) -> str: 

1039 """Return the URL that response `resp` was obtained from. 

1040 

1041 If `resp` contains no redirection history, return `default_url`. 

1042 """ 

1043 if resp.retries is None: 

1044 return default_url 

1045 

1046 if len(resp.retries.history) == 0: 

1047 return default_url 

1048 

1049 return str(resp.retries.history[-1].redirect_location) 

1050 

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`. 

1054 

1055 Parameters 

1056 ---------- 

1057 url : `str` 

1058 Target URL. 

1059 

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 

1075 

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) 

1079 

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. 

1092 

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()`. 

1117 

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) 

1126 

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"}) 

1135 

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) 

1140 

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})" 

1146 

1147 log.debug("sending request %s %s%s", method, redact_url(url), annotation) 

1148 

1149 if pool_manager is None: 

1150 pool_manager = self._pool_manager 

1151 

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 ) 

1168 

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. 

1176 

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. 

1185 

1186 Returns 

1187 ------- 

1188 resp: `HTTPResponse` 

1189 Response to the request as received from the server. 

1190 

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) 

1196 

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. 

1205 

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. 

1214 

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 ) 

1222 

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. 

1230 

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. 

1239 

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) 

1245 

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. 

1255 

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. 

1271 

1272 Returns 

1273 ------- 

1274 resp: `HTTPResponse` 

1275 Response to the GET request as received from the server. 

1276 

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 ) 

1289 

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. 

1297 

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. 

1307 

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) 

1313 

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. 

1321 

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. 

1330 

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) 

1336 

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. 

1344 

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. 

1353 

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) 

1359 

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. 

1368 

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. 

1379 

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) 

1385 

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. 

1396 

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. 

1414 

1415 Returns 

1416 ------- 

1417 resp: `HTTPResponse` 

1418 Response to the PUT request as received from the server. 

1419 

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} 

1442 

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 ) 

1453 

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. 

1461 

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) 

1479 

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. 

1488 

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. 

1502 

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) 

1530 

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. 

1537 

1538 Parameters 

1539 ---------- 

1540 url : `str` 

1541 Target URL. 

1542 headers : `dict` [`str`, `str`], optional 

1543 Headers to sent with the request. 

1544 

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) 

1556 

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. 

1566 

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) 

1592 

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. 

1600 

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. 

1609 

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) 

1630 

1631 # We may have been redirectred. Upload the file contents to 

1632 # its final destination. 

1633 

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}) 

1650 

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) 

1661 

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}" 

1666 

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 

1685 

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 

1690 

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}" 

1696 

1697 def exists(self, url: str) -> bool: 

1698 """Return True if a file or directory exists at `url`. 

1699 

1700 Parameters 

1701 ---------- 

1702 url : `str` 

1703 Target URL. 

1704 

1705 Returns 

1706 ------- 

1707 result: `bool` 

1708 True if there is an object at `url`. 

1709 """ 

1710 return self.stat(url).exists 

1711 

1712 def size(self, url: str) -> int: 

1713 """Return the size in bytes of resource at `url`. 

1714 

1715 If `url` designates a directory, the size is zero. 

1716 

1717 Parameters 

1718 ---------- 

1719 url : `str` 

1720 Target URL. 

1721 

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 

1730 

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 

1736 

1737 def is_dir(self, url: str) -> bool: 

1738 """Return True if a directory exists at `url`. 

1739 

1740 Parameters 

1741 ---------- 

1742 url : `str` 

1743 Target URL. 

1744 

1745 Returns 

1746 ------- 

1747 result: `bool` 

1748 True if there is a directory at `url`. 

1749 """ 

1750 return self.stat(url).is_dir 

1751 

1752 def mkcol(self, url: str) -> None: 

1753 """Create a directory at `url`. 

1754 

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`. 

1757 

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 ) 

1778 

1779 def stat(self, url: str) -> DavFileMetadata: 

1780 """Return some properties of file or directory located at `url`. 

1781 

1782 Parameters 

1783 ---------- 

1784 url : `str` 

1785 Target URL. 

1786 

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. 

1793 

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) 

1821 

1822 def info(self, url: str, name: str | None = None) -> dict[str, Any]: 

1823 """Return the details about the file or directory at `url`. 

1824 

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. 

1832 

1833 Returns 

1834 ------- 

1835 result: `dict` 

1836 For an existing file, the returned value has the form: 

1837 

1838 .. code-block:: json 

1839 

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 } 

1851 

1852 For an existing directory, the returned value has the form: 

1853 

1854 .. code-block:: json 

1855 

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 } 

1864 

1865 For a non-existing file or directory, the returned value has the 

1866 form: 

1867 

1868 .. code-block:: json 

1869 

1870 { 

1871 "name": name, 

1872 "size": None, 

1873 "type": None, 

1874 "last_modified": datetime.datetime(1, 1, 1, 0, 0), 

1875 "checksums": {}, 

1876 } 

1877 

1878 Notes 

1879 ----- 

1880 The format of the returned directory is inspired and compatible with 

1881 `fsspec`. 

1882 

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 

1897 

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 

1907 

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. 

1910 

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. 

1919 

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) 

1930 

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`. 

1934 

1935 If `url` designates a file, only the details of itself are returned. 

1936 

1937 Parameters 

1938 ---------- 

1939 url : `str` 

1940 Target URL. 

1941 

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) 

1961 

1962 if (path := parse_url(url).path) is not None: 

1963 this_dir_href = path.rstrip("/") + "/" 

1964 else: 

1965 this_dir_href = "/" 

1966 

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)) 

1978 

1979 return result 

1980 

1981 def read(self, url: str) -> tuple[str, bytes]: 

1982 """Download the contents of file located at `url`. 

1983 

1984 Parameters 

1985 ---------- 

1986 url : `str` 

1987 Target URL. 

1988 

1989 Returns 

1990 ------- 

1991 url: `str` 

1992 Backend URL from which the data was obtained. 

1993 data: `bytes` 

1994 Contents of the file. 

1995 

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 

2003 

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`. 

2012 

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. 

2023 

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. 

2031 

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}"}) 

2044 

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) 

2051 

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. 

2056 

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. 

2060 

2061 Parameters 

2062 ---------- 

2063 url : `str` 

2064 Target URL. 

2065 """ 

2066 # This is a NOP for a generic webDAV server. 

2067 pass 

2068 

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. 

2071 

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`. 

2081 

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 

2100 

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 ) 

2114 

2115 return content_length 

2116 finally: 

2117 # Release the connection 

2118 resp.drain_conn() 

2119 resp.release_conn() 

2120 

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. 

2123 

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`. 

2133 

2134 Returns 

2135 ------- 

2136 count: `int` 

2137 Number of bytes written to `filename`. 

2138 

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) 

2146 

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. 

2150 

2151 Parameters 

2152 ---------- 

2153 url : `str` 

2154 Target URL. 

2155 data : `bytes` 

2156 Sequence of bytes to upload. 

2157 

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. 

2163 

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)) 

2172 

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) 

2178 

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 

2186 

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

2188 """Return the checksums of the contents of file located at `url`. 

2189 

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. 

2193 

2194 Parameters 

2195 ---------- 

2196 url : `str` 

2197 Target URL. 

2198 

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}") 

2210 

2211 return stat.checksums if stat.is_file else {} 

2212 

2213 def delete(self, url: str) -> None: 

2214 """Delete the file or directory at `url`. 

2215 

2216 If there is no file or directory at `url` is not considered an error. 

2217 

2218 Parameters 

2219 ---------- 

2220 url : `str` 

2221 Target URL. 

2222 

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. 

2228 

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 ) 

2241 

2242 def accepts_ranges(self, url: str) -> bool: 

2243 """Return `True` if the server supports a 'Range' header in 

2244 GET requests against `url`. 

2245 

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 

2256 

2257 with self._lock: 

2258 if self._accepts_ranges is None: 

2259 self._accepts_ranges = self.head(url).headers.get("Accept-Ranges", "") == "bytes" 

2260 

2261 return self._accepts_ranges 

2262 

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 

2269 

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. 

2273 

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 ) 

2296 

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. 

2300 

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}") 

2315 

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) 

2322 

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. 

2332 

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) 

2351 

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 ) 

2361 

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. 

2365 

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. 

2372 

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}") 

2379 

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. 

2383 

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. 

2390 

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}") 

2397 

2398 

2399class ActivityCaveat(enum.Enum): 

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

2401 macaroons for dCache or XRootD webDAV servers. 

2402 """ 

2403 

2404 DOWNLOAD = 1 

2405 UPLOAD = 2 

2406 

2407 

2408class DavClientURLSigner(DavClient): 

2409 """WebDAV client which supports signing of URL for upload and download. 

2410 

2411 Instances of this class are thread-safe. 

2412 

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 """ 

2424 

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) 

2427 

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. 

2431 

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. 

2438 

2439 Returns 

2440 ------- 

2441 url : `str` 

2442 HTTP URL signed for GET. 

2443 

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}" 

2454 

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. 

2458 

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. 

2465 

2466 Returns 

2467 ------- 

2468 url : `str` 

2469 HTTP URL signed for PUT. 

2470 

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}" 

2481 

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`. 

2484 

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. 

2493 

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" 

2508 

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 ) 

2522 

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()}") 

2547 

2548 if "macaroon" in response_body: 

2549 return response_body["macaroon"] 

2550 

2551 raise ValueError(f"Could not retrieve macaroon for URL {resp.geturl()}") 

2552 

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. 

2557 

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}") 

2571 

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) 

2580 

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. 

2585 

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 

2603 

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) 

2609 

2610 # Retrieve a macaroon for downloading the source 

2611 download_macaroon = self._get_macaroon(source_url, ActivityCaveat.DOWNLOAD, 300) 

2612 

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 ) 

2633 

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 ) 

2642 

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() 

2664 

2665 

2666class DavClientDCache(DavClientURLSigner): 

2667 """Client for interacting with a dCache webDAV server. 

2668 

2669 Instances of this class are thread-safe. 

2670 

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 """ 

2682 

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) 

2689 

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) 

2692 

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 

2708 

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://") 

2714 

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) 

2724 

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) 

2734 

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 ) 

2747 

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. 

2756 

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") 

2769 

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) 

2778 

2779 # If we are uploading an empty file, there is nothing more to do. 

2780 if is_zero_length: 

2781 return 0 

2782 

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"}) 

2789 

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}) 

2802 

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) 

2816 

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. 

2820 

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`. 

2830 

2831 Returns 

2832 ------- 

2833 count: `int` 

2834 Number of bytes written to `filename`. 

2835 

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 

2853 

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() 

2858 

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) 

2869 

2870 @override 

2871 def read(self, url: str) -> tuple[str, bytes]: 

2872 """Download the contents of file located at `url`. 

2873 

2874 Parameters 

2875 ---------- 

2876 url : `str` 

2877 Target URL. 

2878 

2879 Returns 

2880 ------- 

2881 url: `str` 

2882 Backend URL from which the data was obtained. 

2883 data: `bytes` 

2884 Contents of the file. 

2885 

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) 

2901 

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) 

2911 

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. 

2916 

2917 Parameters 

2918 ---------- 

2919 url : `str` 

2920 Target URL. 

2921 data : `bytes` 

2922 Sequence of bytes to upload. 

2923 

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. 

2929 

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) 

2944 

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 

2952 

2953 @override 

2954 def mkcol(self, url: str) -> None: 

2955 """Create a directory at `url`. 

2956 

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`. 

2959 

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) 

2984 

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 } 

2995 

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) 

3028 

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}"}) 

3043 

3044 frontend_headers = {} if headers is None else dict(headers) 

3045 frontend_headers.update(range_headers) 

3046 

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 

3062 

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"}) 

3069 

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) 

3080 

3081 @override 

3082 def _close(self, url: str) -> None: 

3083 # Docstring inherited. 

3084 

3085 # For dCache this is a NOP since `read_range` does not keep the 

3086 # connection with the dCache pool open. 

3087 pass 

3088 

3089 

3090class DavClientXrootD(DavClientURLSigner): 

3091 """Client for interacting with a XrootD webDAV server. 

3092 

3093 Instances of this class are thread-safe. 

3094 

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 """ 

3106 

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) 

3109 

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. 

3118 

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 ) 

3149 

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) 

3159 

3160 # We were redirected to a backend server. Upload the file contents to 

3161 # its final destination. 

3162 

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. 

3166 

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}) 

3183 

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) 

3194 

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}) 

3212 

3213 return result 

3214 

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. 

3219 

3220 Parameters 

3221 ---------- 

3222 url : `str` 

3223 Target URL. 

3224 data : `bytes` 

3225 Sequence of bytes to upload. 

3226 

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. 

3232 

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) 

3247 

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 

3255 

3256 @override 

3257 def mkcol(self, url: str) -> None: 

3258 """Create a directory at `url`. 

3259 

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`. 

3262 

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 ) 

3288 

3289 @override 

3290 def stat(self, url: str) -> DavFileMetadata: 

3291 # Docstring inherited. 

3292 

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) 

3334 

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. 

3344 

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}"}) 

3355 

3356 get_headers = {} if headers is None else dict(headers) 

3357 get_headers.update(range_headers) 

3358 

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) 

3365 

3366 @override 

3367 def _close(self, url: str) -> None: 

3368 # Docstring inherited. 

3369 

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) 

3374 

3375 

3376class DavFileMetadata: 

3377 """Container for attributes of interest of a webDAV file or directory. 

3378 

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 """ 

3398 

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) 

3418 

3419 @staticmethod 

3420 def from_property(base_url: str, property: DavProperty) -> DavFileMetadata: 

3421 """Create an instance from the values in `property`. 

3422 

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 ) 

3440 

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 ) 

3446 

3447 @property 

3448 def url(self) -> str: 

3449 return self._url 

3450 

3451 @property 

3452 def href(self) -> str: 

3453 return self._href 

3454 

3455 @property 

3456 def name(self) -> str: 

3457 return self._name 

3458 

3459 @property 

3460 def exists(self) -> bool: 

3461 return self._exists 

3462 

3463 @property 

3464 def size(self) -> int: 

3465 if not self._exists: 

3466 return -1 

3467 

3468 return 0 if self._is_dir else self._size 

3469 

3470 @property 

3471 def is_dir(self) -> bool: 

3472 return self._exists and self._is_dir 

3473 

3474 @property 

3475 def is_file(self) -> bool: 

3476 return self._exists and not self._is_dir 

3477 

3478 @property 

3479 def last_modified(self) -> datetime: 

3480 return self._last_modified 

3481 

3482 @property 

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

3484 return self._checksums 

3485 

3486 

3487class DavProperty: 

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

3489 resource, as retrieved via a PROPFIND request. 

3490 

3491 Parameters 

3492 ---------- 

3493 response : `eTree.Element` or `None` 

3494 The XML response defining the DAV property. 

3495 """ 

3496 

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) 

3500 

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] = {} 

3508 

3509 if response is not None: 

3510 self._parse(response) 

3511 

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 ) 

3523 

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 

3529 

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 

3534 

3535 # Parse "getlastmodified". 

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

3537 self._getlastmodified = str(element.text) 

3538 

3539 # Parse "getcontentlength". 

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

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

3542 

3543 # Parse "displayname". 

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

3545 self._displayname = str(element.text) 

3546 

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) 

3550 

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("/")) 

3556 

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("/") + "/" 

3561 

3562 # Force a size of 0 for collections. 

3563 if self._collection: 

3564 self._getcontentlength = 0 

3565 

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 

3582 

3583 return result 

3584 

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 

3589 

3590 @property 

3591 def is_dir(self) -> bool: 

3592 return self._collection 

3593 

3594 @property 

3595 def is_file(self) -> bool: 

3596 return not self._collection 

3597 

3598 @property 

3599 def last_modified(self) -> datetime: 

3600 if not self._getlastmodified: 

3601 return datetime.min 

3602 

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) 

3606 

3607 @property 

3608 def size(self) -> int: 

3609 return self._getcontentlength 

3610 

3611 @property 

3612 def name(self) -> str: 

3613 return self._displayname 

3614 

3615 @property 

3616 def href(self) -> str: 

3617 return self._href 

3618 

3619 @property 

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

3621 return self._checksums 

3622 

3623 

3624class DavPropfindParser: 

3625 """Helper class to parse the response body of a PROPFIND request.""" 

3626 

3627 def __init__(self) -> None: 

3628 return 

3629 

3630 def parse(self, body: bytes) -> list[DavProperty]: 

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

3632 PROPFIND request. 

3633 

3634 Parameters 

3635 ---------- 

3636 body : `bytes` 

3637 XML-encoded response body to a PROPFIND request. 

3638 

3639 Returns 

3640 ------- 

3641 responses : `list` [ `DavProperty` ] 

3642 Parsed content of the response. 

3643 

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> 

3680 

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)) 

3687 

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}") 

3693 

3694 

3695class Authorizer: 

3696 """Base class for attaching an 'Authorization' header to a HTTP request.""" 

3697 

3698 def set_authorization(self, headers: dict[str, str]) -> None: 

3699 """Add the 'Authorization' header to `headers`. 

3700 

3701 Parameters 

3702 ---------- 

3703 headers : `dict` [ `str`, `str` ] 

3704 Dict to augment with authorization information. 

3705 

3706 Notes 

3707 ----- 

3708 This method must be implemented by concrete subclasses. 

3709 """ 

3710 raise NotImplementedError 

3711 

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. 

3715 

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 

3723 

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 

3729 

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`. 

3735 

3736 Parameters 

3737 ---------- 

3738 filename : `str`, optional 

3739 Path of a local file. 

3740 

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`. 

3745 

3746 Returns 

3747 ------- 

3748 result: `tuple[str, float]` 

3749 tuple of (contents of file `filename`, timestamp of the read 

3750 operation). 

3751 

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) 

3756 

3757 if os.stat(filename).st_mtime < timestamp: 

3758 return (None, None) 

3759 

3760 with open(filename) as file: 

3761 time_of_last_read = time.time() 

3762 return (file.read().rstrip("\n"), time_of_last_read) 

3763 

3764 

3765class TokenAuthorizer(Authorizer): 

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

3767 

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 """ 

3775 

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 

3781 

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() 

3791 

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 

3798 

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 

3802 

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 

3806 

3807 @override 

3808 def set_authorization(self, headers: dict[str, str]) -> None: 

3809 """Add the 'Authorization' header to `headers`. 

3810 

3811 Parameters 

3812 ---------- 

3813 headers : `dict` [ `str`, `str` ] 

3814 Dict to augment with authorization information. 

3815 """ 

3816 if self._token is None: 

3817 return 

3818 

3819 self._update_token() 

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

3821 

3822 

3823class BasicAuthorizer(Authorizer): 

3824 """Attach a 'Authorization' header to each request using Basic 

3825 authentication. 

3826 

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 """ 

3838 

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 

3842 

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 = "" 

3848 

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() 

3861 

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"] 

3868 

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 

3875 

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 

3881 

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() 

3887 

3888 @override 

3889 def set_authorization(self, headers: dict[str, str]) -> None: 

3890 """Add the 'Authorization' header to `headers`. 

3891 

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 

3899 

3900 self._update_password() 

3901 headers["Authorization"] = self._header_value 

3902 

3903 

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. 

3907 

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'). 

3913 

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) 

3920 

3921 

3922def dump_response(method: str, resp: HTTPResponse, dump_body: bool = False) -> None: 

3923 """Dump response for debugging purposes. 

3924 

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) 

3936 

3937 for header, value in resp.headers.items(): 

3938 log.debug(" %s: %s", header, value) 

3939 

3940 if dump_body: 

3941 log.debug(" response body length: %d", len(resp.data.decode()))